@@ -240,6 +240,107 @@ public void rejected() {
240240 assertEquals (1 , rejected .get ());
241241 }
242242
243+ @ Test
244+ void shouldDelegateCompositeHandleUpstreamWithoutDuplicatingPlatformHandler () {
245+ List <String > signals = new ArrayList <>();
246+ AtomicInteger platformApplied = new AtomicInteger ();
247+ AtomicInteger platformHandled = new AtomicInteger ();
248+ HandleUpstreamRecordingMonitor first = new HandleUpstreamRecordingMonitor ("first" , signals );
249+ HandleUpstreamRecordingMonitor second = new HandleUpstreamRecordingMonitor ("second" , signals );
250+ CompositeDeviceGatewayMonitor monitor = new CompositeDeviceGatewayMonitor ()
251+ .add (first , second );
252+
253+ ClientConnection connection = mock (ClientConnection .class );
254+ DeviceSession session = mock (DeviceSession .class );
255+ EncodedMessage origin = mock (EncodedMessage .class );
256+ DeviceMessage message = mock (DeviceMessage .class );
257+
258+ Flux <DeviceMessage > upstream = monitor .handleUpstream (
259+ connection ,
260+ session ,
261+ origin ,
262+ Flux .just (message ),
263+ decoder -> {
264+ platformApplied .incrementAndGet ();
265+ return decoder .doOnNext (ignore -> {
266+ platformHandled .incrementAndGet ();
267+ signals .add ("platform" );
268+ });
269+ }
270+ );
271+
272+ assertEquals (1 , platformApplied .get ());
273+ StepVerifier
274+ .create (upstream )
275+ .expectNext (message )
276+ .verifyComplete ();
277+
278+ assertEquals (1 , first .handleUpstreamInvocations .get ());
279+ assertEquals (1 , second .handleUpstreamInvocations .get ());
280+ assertEquals (1 , platformHandled .get ());
281+ assertEquals (Arrays .asList ("platform" , "first" , "second" ), signals );
282+ }
283+
284+ @ Test
285+ void shouldApplyPlatformHandlerOnceForEmptyComposite () {
286+ AtomicInteger platformApplied = new AtomicInteger ();
287+ AtomicInteger platformHandled = new AtomicInteger ();
288+ DeviceMessage message = mock (DeviceMessage .class );
289+
290+ Flux <DeviceMessage > upstream = new CompositeDeviceGatewayMonitor ()
291+ .handleUpstream (
292+ mock (ClientConnection .class ),
293+ mock (DeviceSession .class ),
294+ mock (EncodedMessage .class ),
295+ Flux .just (message ),
296+ decoder -> {
297+ platformApplied .incrementAndGet ();
298+ return decoder .doOnNext (ignore -> platformHandled .incrementAndGet ());
299+ }
300+ );
301+
302+ assertEquals (1 , platformApplied .get ());
303+ StepVerifier
304+ .create (upstream )
305+ .expectNext (message )
306+ .verifyComplete ();
307+ assertEquals (1 , platformHandled .get ());
308+ }
309+
310+ @ Test
311+ void shouldDelegateLazyHandleUpstreamToTarget () {
312+ List <String > signals = new ArrayList <>();
313+ AtomicInteger resolved = new AtomicInteger ();
314+ AtomicInteger platformHandled = new AtomicInteger ();
315+ HandleUpstreamRecordingMonitor target = new HandleUpstreamRecordingMonitor ("target" , signals );
316+ LazyDeviceGatewayMonitor monitor = new LazyDeviceGatewayMonitor (() -> {
317+ resolved .incrementAndGet ();
318+ return target ;
319+ });
320+ DeviceMessage message = mock (DeviceMessage .class );
321+
322+ Flux <DeviceMessage > upstream = monitor .handleUpstream (
323+ mock (ClientConnection .class ),
324+ mock (DeviceSession .class ),
325+ mock (EncodedMessage .class ),
326+ Flux .just (message ),
327+ decoder -> decoder .doOnNext (ignore -> {
328+ platformHandled .incrementAndGet ();
329+ signals .add ("platform" );
330+ })
331+ );
332+
333+ StepVerifier
334+ .create (upstream )
335+ .expectNext (message )
336+ .verifyComplete ();
337+
338+ assertEquals (1 , resolved .get ());
339+ assertEquals (1 , target .handleUpstreamInvocations .get ());
340+ assertEquals (1 , platformHandled .get ());
341+ assertEquals (Arrays .asList ("platform" , "target" ), signals );
342+ }
343+
243344 @ Test
244345 void shouldComposeAllMonitorDecisionsAndReactiveWrappersInOrder () {
245346 List <String > signals = new ArrayList <>();
@@ -420,4 +521,35 @@ public Mono<Void> downstream(ClientConnection connection,
420521 return sender .doOnSuccess (ignore -> signals .add (id + ":downstream" ));
421522 }
422523 }
524+
525+ private static class HandleUpstreamRecordingMonitor implements DeviceGatewayMonitor {
526+
527+ private final String id ;
528+ private final List <String > signals ;
529+ private final AtomicInteger handleUpstreamInvocations = new AtomicInteger ();
530+
531+ private HandleUpstreamRecordingMonitor (String id , List <String > signals ) {
532+ this .id = id ;
533+ this .signals = signals ;
534+ }
535+
536+ @ Override
537+ public Flux <DeviceMessage > handleUpstream (ClientConnection connection ,
538+ DeviceSession session ,
539+ EncodedMessage origin ,
540+ Flux <DeviceMessage > decoder ,
541+ UnaryOperator <Flux <DeviceMessage >> platformHandler ) {
542+ handleUpstreamInvocations .incrementAndGet ();
543+ return DeviceGatewayMonitor .super
544+ .handleUpstream (connection , session , origin , decoder , platformHandler );
545+ }
546+
547+ @ Override
548+ public Flux <DeviceMessage > beforeSendToPlatform (ClientConnection connection ,
549+ DeviceSession session ,
550+ EncodedMessage origin ,
551+ Flux <DeviceMessage > handler ) {
552+ return handler .doOnNext (ignore -> signals .add (id ));
553+ }
554+ }
423555}
0 commit comments