1
2
3
4
5
6
7
8
9
10
11
12
13
14
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
62
63
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
84
85
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
98
99
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
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)
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
338
339
340
341
342
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
377
378
379
380
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 }