View Javadoc
1   /*
2    * Copyright 2019-2026 the original author or authors.
3    *
4    * Licensed under the Apache License, Version 2.0 (the "License");
5    * you may not use this file except in compliance with the License.
6    * You may obtain a copy of the License at
7    *
8    *      http://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS,
12   * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13   * See the License for the specific language governing permissions and
14   * limitations under the License.
15   */
16  
17  package org.bremersee.ldaptive.reactive;
18  
19  import static org.ldaptive.handler.ResultPredicate.NOT_SUCCESS;
20  
21  import java.util.Objects;
22  import java.util.concurrent.CompletableFuture;
23  import java.util.function.Function;
24  import org.apache.commons.logging.Log;
25  import org.apache.commons.logging.LogFactory;
26  import org.bremersee.exception.ServiceException;
27  import org.bremersee.ldaptive.DefaultLdaptiveErrorHandler;
28  import org.bremersee.ldaptive.LdaptiveEntryMapper;
29  import org.bremersee.ldaptive.LdaptiveErrorHandler;
30  import org.bremersee.ldaptive.LdaptiveTemplate;
31  import org.ldaptive.AddOperation;
32  import org.ldaptive.AddRequest;
33  import org.ldaptive.AttributeModification;
34  import org.ldaptive.BindRequest;
35  import org.ldaptive.CompareOperation;
36  import org.ldaptive.CompareRequest;
37  import org.ldaptive.ConnectionFactory;
38  import org.ldaptive.DeleteOperation;
39  import org.ldaptive.DeleteRequest;
40  import org.ldaptive.LdapAttribute;
41  import org.ldaptive.LdapEntry;
42  import org.ldaptive.LdapException;
43  import org.ldaptive.ModifyDnOperation;
44  import org.ldaptive.ModifyDnRequest;
45  import org.ldaptive.ModifyOperation;
46  import org.ldaptive.ModifyRequest;
47  import org.ldaptive.Result;
48  import org.ldaptive.ResultCode;
49  import org.ldaptive.SearchOperation;
50  import org.ldaptive.SearchRequest;
51  import org.ldaptive.extended.ExtendedOperation;
52  import org.ldaptive.extended.ExtendedRequest;
53  import org.ldaptive.extended.ExtendedResponse;
54  import org.ldaptive.handler.ResultHandler;
55  import org.ldaptive.handler.ResultPredicate;
56  import reactor.core.publisher.Flux;
57  import reactor.core.publisher.FluxSink;
58  import reactor.core.publisher.Mono;
59  
60  /**
61   * The reactive ldaptive template.
62   *
63   * @author Christian Bremer
64   */
65  public class ReactiveLdaptiveTemplate implements ReactiveLdaptiveOperations {
66  
67    private static final Log log = LogFactory.getLog(ReactiveLdaptiveTemplate.class);
68  
69    private static final ResultPredicate NOT_COMPARE_RESULT = result -> !result.isSuccess()
70        && result.getResultCode() != ResultCode.COMPARE_TRUE
71        && result.getResultCode() != ResultCode.COMPARE_FALSE;
72  
73    private static final ResultPredicate NOT_DELETE_RESULT = result -> !result.isSuccess()
74        && result.getResultCode() != ResultCode.NO_SUCH_OBJECT;
75  
76    private static final ResultPredicate NOT_FIND_RESULT = NOT_DELETE_RESULT;
77  
78    private final ConnectionFactory connectionFactory;
79  
80    private LdaptiveErrorHandler errorHandler = new DefaultLdaptiveErrorHandler();
81  
82    /**
83     * Instantiates a new Reactive ldaptive template.
84     *
85     * @param connectionFactory the connection factory
86     */
87    public ReactiveLdaptiveTemplate(ConnectionFactory connectionFactory) {
88      this.connectionFactory = connectionFactory;
89    }
90  
91    @Override
92    public ConnectionFactory getConnectionFactory() {
93      return connectionFactory;
94    }
95  
96    /**
97     * Sets error handler.
98     *
99     * @param errorHandler the error handler
100    */
101   public void setErrorHandler(LdaptiveErrorHandler errorHandler) {
102     if (errorHandler != null) {
103       this.errorHandler = errorHandler;
104     }
105   }
106 
107   @Override
108   public ReactiveLdaptiveTemplate copy() {
109     return copy(null);
110   }
111 
112   @Override
113   public ReactiveLdaptiveTemplate copy(final LdaptiveErrorHandler errorHandler) {
114     final ReactiveLdaptiveTemplate template = new ReactiveLdaptiveTemplate(connectionFactory);
115     template.setErrorHandler(errorHandler);
116     return template;
117   }
118 
119   @Override
120   public Mono<Result> add(AddRequest addRequest) {
121     CompletableFuture<Result> future = new CompletableFuture<>();
122     try {
123       AddOperation.builder()
124           .factory(connectionFactory)
125           .onResult(new FutureAwareResultHandler<>(future, NOT_SUCCESS, errorHandler, r -> r))
126           .onException(ldapException -> future
127               .completeExceptionally(errorHandler.map(ldapException)))
128           .build()
129           .send(addRequest);
130 
131     } catch (LdapException e) {
132       future.completeExceptionally(errorHandler.map(e));
133     }
134     return Mono.fromFuture(future);
135   }
136 
137   private <T> Mono<T> add(T domainObject, LdaptiveEntryMapper<T> entryMapper) {
138     String[] objectClasses = entryMapper.getObjectClasses();
139     if (objectClasses == null || objectClasses.length == 0) {
140       final ServiceException se = ServiceException.internalServerError(
141           "Object classes must be specified to save a new ldap entry.",
142           "org.bremersee:ldaptive-integration:d7aa5699-fd2e-45df-a863-97960e8095b8");
143       log.error("Saving domain object failed.", se);
144       throw se;
145     }
146     String dn = entryMapper.mapDn(domainObject);
147     LdapEntry entry = new LdapEntry();
148     entryMapper.map(domainObject, entry);
149     entry.setDn(dn);
150     entry.addAttributes(new LdapAttribute("objectclass", objectClasses));
151     return add(new AddRequest(dn, entry.getAttributes()))
152         .then(Mono.just(Objects.requireNonNull(entryMapper.map(entry))));
153   }
154 
155   @Override
156   public Mono<Boolean> bind(BindRequest bindRequest) {
157     // Bind requests are synchronous
158     LdaptiveTemplate template = new LdaptiveTemplate(getConnectionFactory());
159     template.setErrorHandler(errorHandler);
160     return Mono.just(template.bind(bindRequest));
161   }
162 
163   @Override
164   public Mono<Boolean> compare(CompareRequest compareRequest) {
165     CompletableFuture<Boolean> future = new CompletableFuture<>();
166     try {
167       CompareOperation.builder()
168           .factory(connectionFactory)
169           .onCompare(
170               future::complete) // this will be only called, if the result is COMPARE_TRUE or COMPARE_FALSE
171           .onResult(new FutureAwareResultHandler<>(future, NOT_COMPARE_RESULT, errorHandler,
172               Result::isSuccess))
173           .onException(
174               ldapException -> future.completeExceptionally(errorHandler.map(ldapException)))
175           .build()
176           .send(compareRequest);
177 
178     } catch (LdapException e) {
179       future.completeExceptionally(errorHandler.map(e));
180     }
181     return Mono.fromFuture(future);
182   }
183 
184   @Override
185   public Mono<Result> delete(DeleteRequest deleteRequest) {
186     CompletableFuture<Result> future = new CompletableFuture<>();
187     try {
188       DeleteOperation.builder()
189           .factory(connectionFactory)
190           .onResult(new FutureAwareResultHandler<>(future, NOT_DELETE_RESULT, errorHandler, r -> r))
191           .onException(
192               ldapException -> future.completeExceptionally(errorHandler.map(ldapException)))
193           .build()
194           .send(deleteRequest)
195           .await();
196 
197     } catch (LdapException e) {
198       future.completeExceptionally(errorHandler.map(e));
199     }
200     return Mono.fromFuture(future);
201   }
202 
203   @Override
204   public Mono<ExtendedResponse> executeExtension(ExtendedRequest request) {
205 
206     CompletableFuture<ExtendedResponse> future = new CompletableFuture<>();
207     try {
208       ExtendedOperation.builder()
209           .factory(connectionFactory)
210           .onExtended((name, value) -> future.complete(ExtendedResponse.builder()
211               .responseName(name)
212               .responseValue(value)
213               .resultCode(ResultCode.SUCCESS)
214               .build()))
215           .onResult(new FutureAwareResultHandler<>(
216               future,
217               NOT_SUCCESS,
218               errorHandler,
219               r -> ExtendedResponse.builder().resultCode(r.getResultCode()).build()))
220           .onException(
221               ldapException -> future.completeExceptionally(errorHandler.map(ldapException)))
222           .build()
223           .send(request);
224 
225     } catch (LdapException e) {
226       future.completeExceptionally(errorHandler.map(e));
227     }
228     return Mono.fromFuture(future);
229   }
230 
231   @Override
232   public Mono<Result> modify(ModifyRequest modifyRequest) {
233     CompletableFuture<Result> future = new CompletableFuture<>();
234     try {
235       ModifyOperation.builder()
236           .factory(connectionFactory)
237           .onResult(new FutureAwareResultHandler<>(future, NOT_SUCCESS, errorHandler, r -> r))
238           .onException(
239               ldapException -> future.completeExceptionally(errorHandler.map(ldapException)))
240           .build()
241           .send(modifyRequest);
242 
243     } catch (LdapException e) {
244       future.completeExceptionally(errorHandler.map(e));
245     }
246     return Mono.fromFuture(future);
247   }
248 
249   private <T> Mono<T> modify(T domainObject, LdapEntry entry, LdaptiveEntryMapper<T> entryMapper) {
250     String dn = entryMapper.mapDn(domainObject);
251     AttributeModification[] modifications = entryMapper.mapAndComputeModifications(domainObject,
252         entry);
253     return modify(new ModifyRequest(dn, modifications))
254         .then(Mono.just(Objects.requireNonNull(entryMapper.map(entry))));
255   }
256 
257   @Override
258   public Mono<Result> modifyDn(ModifyDnRequest modifyDnRequest) {
259     CompletableFuture<Result> future = new CompletableFuture<>();
260     try {
261       ModifyDnOperation.builder()
262           .factory(connectionFactory)
263           .onResult(new FutureAwareResultHandler<>(future, NOT_SUCCESS, errorHandler, r -> r))
264           .onException(
265               ldapException -> future.completeExceptionally(errorHandler.map(ldapException)))
266           .build()
267           .send(modifyDnRequest);
268 
269     } catch (LdapException e) {
270       future.completeExceptionally(errorHandler.map(e));
271     }
272     return Mono.fromFuture(future);
273   }
274 
275   @Override
276   public Mono<LdapEntry> findOne(SearchRequest searchRequest) {
277     CompletableFuture<LdapEntry> future = new CompletableFuture<>();
278     try {
279       SearchOperation.builder()
280           .factory(connectionFactory)
281           .onEntry(ldapEntry -> {
282             future.complete(ldapEntry);
283             return ldapEntry;
284           })
285           .onResult(new FutureAwareResultHandler<>(future, NOT_FIND_RESULT, errorHandler, null))
286           .onException(ldapException -> future.obtrudeException(errorHandler.map(ldapException)))
287           .build()
288           .send(searchRequest);
289 
290     } catch (LdapException e) {
291       future.completeExceptionally(errorHandler.map(e));
292     }
293     return Mono.fromFuture(future);
294   }
295 
296   @Override
297   public Flux<LdapEntry> findAll(SearchRequest searchRequest) {
298     return Flux.create((FluxSink<LdapEntry> fluxSink) -> {
299       try {
300         SearchOperation.builder()
301             .factory(connectionFactory)
302             .onEntry(ldapEntry -> {
303               fluxSink.next(ldapEntry);
304               return ldapEntry;
305             })
306             .onResult(new FluxSinkAwareResultHandler<>(fluxSink, NOT_FIND_RESULT, errorHandler))
307             .onException(ldapException -> fluxSink.error(errorHandler.map(ldapException)))
308             .build()
309             .send(searchRequest);
310 
311       } catch (LdapException e) {
312         fluxSink.error(errorHandler.map(e));
313       }
314     });
315   }
316 
317   @Override
318   public <T> Mono<T> save(T domainObject, LdaptiveEntryMapper<T> entryMapper) {
319     return findOne(SearchRequest.objectScopeSearchRequest(entryMapper.mapDn(domainObject)))
320         .singleOptional()
321         .flatMap(optExisting -> optExisting
322             .map(existing -> modify(domainObject, existing, entryMapper))
323             .orElseGet(() -> add(domainObject, entryMapper)));
324   }
325 
326   private static class FutureAwareResultHandler<T> implements ResultHandler {
327 
328     private final CompletableFuture<T> future;
329 
330     private final ResultPredicate throwErrorPredicate;
331 
332     private final LdaptiveErrorHandler errorHandler;
333 
334     private final Function<Result, T> resultValueFn;
335 
336     /**
337      * Instantiates a new Future aware result handler.
338      *
339      * @param future the future
340      * @param throwErrorPredicate the throw error predicate
341      * @param errorHandler the error handler
342      * @param resultValueFn the result value fn
343      */
344     public FutureAwareResultHandler(
345         CompletableFuture<T> future,
346         ResultPredicate throwErrorPredicate,
347         LdaptiveErrorHandler errorHandler,
348         Function<Result, T> resultValueFn) {
349       this.throwErrorPredicate = throwErrorPredicate;
350       this.errorHandler = errorHandler;
351       this.future = future;
352       this.resultValueFn = resultValueFn;
353     }
354 
355     @Override
356     public void accept(Result result) {
357       if (!future.isDone()) {
358         if (throwErrorPredicate != null && throwErrorPredicate.test(result)) {
359           future.completeExceptionally(errorHandler.map(new LdapException(result)));
360         } else {
361           future.complete(resultValueFn != null ? resultValueFn.apply(result) : null);
362         }
363       }
364     }
365   }
366 
367   private static class FluxSinkAwareResultHandler<T> implements ResultHandler {
368 
369     private final FluxSink<T> fluxSink;
370 
371     private final ResultPredicate throwErrorPredicate;
372 
373     private final LdaptiveErrorHandler errorHandler;
374 
375     /**
376      * Instantiates a new Flux sink aware result handler.
377      *
378      * @param fluxSink the flux sink
379      * @param throwErrorPredicate the throw error predicate
380      * @param errorHandler the error handler
381      */
382     public FluxSinkAwareResultHandler(
383         FluxSink<T> fluxSink,
384         ResultPredicate throwErrorPredicate,
385         LdaptiveErrorHandler errorHandler) {
386       this.throwErrorPredicate = throwErrorPredicate;
387       this.errorHandler = errorHandler;
388       this.fluxSink = fluxSink;
389     }
390 
391     @Override
392     public void accept(Result result) {
393       if (throwErrorPredicate != null && throwErrorPredicate.test(result)) {
394         fluxSink.error(errorHandler.map(new LdapException(result)));
395       } else {
396         fluxSink.complete();
397       }
398     }
399   }
400 
401 }