diff --git a/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/CompositeDeviceGatewayMonitor.java b/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/CompositeDeviceGatewayMonitor.java index 72c19d9c5..b393b8704 100644 --- a/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/CompositeDeviceGatewayMonitor.java +++ b/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/CompositeDeviceGatewayMonitor.java @@ -29,6 +29,7 @@ import java.util.Collection; import java.util.List; import java.util.function.Consumer; +import java.util.function.UnaryOperator; class CompositeDeviceGatewayMonitor implements DeviceGatewayMonitor { @@ -121,6 +122,21 @@ public Flux decode(@Nullable ClientConnection connection, return decoder; } + @Override + public Flux handleUpstream(@Nullable ClientConnection connection, + DeviceSession session, + EncodedMessage origin, + Flux decoder, + UnaryOperator> platformHandler) { + // 平台处理器作为链尾只组合一次,同时保留 beforeSendToPlatform 的注册顺序。 + UnaryOperator> handler = platformHandler; + for (DeviceGatewayMonitor monitor : monitors) { + UnaryOperator> next = handler; + handler = source -> monitor.handleUpstream(connection, session, origin, source, next); + } + return handler.apply(decoder); + } + @Override public Flux beforeSendToPlatform(@Nullable ClientConnection connection, DeviceSession session, diff --git a/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/LazyDeviceGatewayMonitor.java b/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/LazyDeviceGatewayMonitor.java index ee023ba82..c9785d03e 100644 --- a/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/LazyDeviceGatewayMonitor.java +++ b/jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/LazyDeviceGatewayMonitor.java @@ -24,6 +24,7 @@ import javax.annotation.Nullable; import java.util.function.Supplier; +import java.util.function.UnaryOperator; class LazyDeviceGatewayMonitor implements DeviceGatewayMonitor { @@ -101,6 +102,15 @@ public Flux decode(@Nullable ClientConnection connection, return getTarget().decode(connection, session, origin, decoder); } + @Override + public Flux handleUpstream(@Nullable ClientConnection connection, + DeviceSession session, + EncodedMessage origin, + Flux decoder, + UnaryOperator> platformHandler) { + return getTarget().handleUpstream(connection, session, origin, decoder, platformHandler); + } + @Override public Flux beforeSendToPlatform(@Nullable ClientConnection connection, DeviceSession session, diff --git a/jetlinks-components/gateway-component/src/test/java/org/jetlinks/community/gateway/monitor/DeviceGatewayMonitorTest.java b/jetlinks-components/gateway-component/src/test/java/org/jetlinks/community/gateway/monitor/DeviceGatewayMonitorTest.java index b3a572f0d..a05da369c 100644 --- a/jetlinks-components/gateway-component/src/test/java/org/jetlinks/community/gateway/monitor/DeviceGatewayMonitorTest.java +++ b/jetlinks-components/gateway-component/src/test/java/org/jetlinks/community/gateway/monitor/DeviceGatewayMonitorTest.java @@ -240,6 +240,107 @@ public void rejected() { assertEquals(1, rejected.get()); } + @Test + void shouldDelegateCompositeHandleUpstreamWithoutDuplicatingPlatformHandler() { + List signals = new ArrayList<>(); + AtomicInteger platformApplied = new AtomicInteger(); + AtomicInteger platformHandled = new AtomicInteger(); + HandleUpstreamRecordingMonitor first = new HandleUpstreamRecordingMonitor("first", signals); + HandleUpstreamRecordingMonitor second = new HandleUpstreamRecordingMonitor("second", signals); + CompositeDeviceGatewayMonitor monitor = new CompositeDeviceGatewayMonitor() + .add(first, second); + + ClientConnection connection = mock(ClientConnection.class); + DeviceSession session = mock(DeviceSession.class); + EncodedMessage origin = mock(EncodedMessage.class); + DeviceMessage message = mock(DeviceMessage.class); + + Flux upstream = monitor.handleUpstream( + connection, + session, + origin, + Flux.just(message), + decoder -> { + platformApplied.incrementAndGet(); + return decoder.doOnNext(ignore -> { + platformHandled.incrementAndGet(); + signals.add("platform"); + }); + } + ); + + assertEquals(1, platformApplied.get()); + StepVerifier + .create(upstream) + .expectNext(message) + .verifyComplete(); + + assertEquals(1, first.handleUpstreamInvocations.get()); + assertEquals(1, second.handleUpstreamInvocations.get()); + assertEquals(1, platformHandled.get()); + assertEquals(Arrays.asList("platform", "first", "second"), signals); + } + + @Test + void shouldApplyPlatformHandlerOnceForEmptyComposite() { + AtomicInteger platformApplied = new AtomicInteger(); + AtomicInteger platformHandled = new AtomicInteger(); + DeviceMessage message = mock(DeviceMessage.class); + + Flux upstream = new CompositeDeviceGatewayMonitor() + .handleUpstream( + mock(ClientConnection.class), + mock(DeviceSession.class), + mock(EncodedMessage.class), + Flux.just(message), + decoder -> { + platformApplied.incrementAndGet(); + return decoder.doOnNext(ignore -> platformHandled.incrementAndGet()); + } + ); + + assertEquals(1, platformApplied.get()); + StepVerifier + .create(upstream) + .expectNext(message) + .verifyComplete(); + assertEquals(1, platformHandled.get()); + } + + @Test + void shouldDelegateLazyHandleUpstreamToTarget() { + List signals = new ArrayList<>(); + AtomicInteger resolved = new AtomicInteger(); + AtomicInteger platformHandled = new AtomicInteger(); + HandleUpstreamRecordingMonitor target = new HandleUpstreamRecordingMonitor("target", signals); + LazyDeviceGatewayMonitor monitor = new LazyDeviceGatewayMonitor(() -> { + resolved.incrementAndGet(); + return target; + }); + DeviceMessage message = mock(DeviceMessage.class); + + Flux upstream = monitor.handleUpstream( + mock(ClientConnection.class), + mock(DeviceSession.class), + mock(EncodedMessage.class), + Flux.just(message), + decoder -> decoder.doOnNext(ignore -> { + platformHandled.incrementAndGet(); + signals.add("platform"); + }) + ); + + StepVerifier + .create(upstream) + .expectNext(message) + .verifyComplete(); + + assertEquals(1, resolved.get()); + assertEquals(1, target.handleUpstreamInvocations.get()); + assertEquals(1, platformHandled.get()); + assertEquals(Arrays.asList("platform", "target"), signals); + } + @Test void shouldComposeAllMonitorDecisionsAndReactiveWrappersInOrder() { List signals = new ArrayList<>(); @@ -420,4 +521,35 @@ public Mono downstream(ClientConnection connection, return sender.doOnSuccess(ignore -> signals.add(id + ":downstream")); } } + + private static class HandleUpstreamRecordingMonitor implements DeviceGatewayMonitor { + + private final String id; + private final List signals; + private final AtomicInteger handleUpstreamInvocations = new AtomicInteger(); + + private HandleUpstreamRecordingMonitor(String id, List signals) { + this.id = id; + this.signals = signals; + } + + @Override + public Flux handleUpstream(ClientConnection connection, + DeviceSession session, + EncodedMessage origin, + Flux decoder, + UnaryOperator> platformHandler) { + handleUpstreamInvocations.incrementAndGet(); + return DeviceGatewayMonitor.super + .handleUpstream(connection, session, origin, decoder, platformHandler); + } + + @Override + public Flux beforeSendToPlatform(ClientConnection connection, + DeviceSession session, + EncodedMessage origin, + Flux handler) { + return handler.doOnNext(ignore -> signals.add(id)); + } + } } diff --git a/jetlinks-components/network-component/tcp-component/src/test/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGatewayMonitorTest.java b/jetlinks-components/network-component/tcp-component/src/test/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGatewayMonitorTest.java index fae840596..93f5c1dea 100644 --- a/jetlinks-components/network-component/tcp-component/src/test/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGatewayMonitorTest.java +++ b/jetlinks-components/network-component/tcp-component/src/test/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGatewayMonitorTest.java @@ -53,6 +53,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.UnaryOperator; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertSame; @@ -132,6 +133,7 @@ public Mono getMessageCodec(Transport transport) { assertEquals(1, monitor.beforeDecode.get()); assertEquals(1, monitor.decode.get()); + assertEquals(2, monitor.handleUpstream.get()); assertEquals(2, monitor.beforeSend.get()); assertEquals(2, monitor.received.get()); assertEquals(1, Collections.frequency(monitor.monitored, manualMessage)); @@ -146,6 +148,7 @@ public Mono getMessageCodec(Transport transport) { private static class RecordingMonitor implements DeviceGatewayMonitor { private final AtomicInteger beforeDecode = new AtomicInteger(); private final AtomicInteger decode = new AtomicInteger(); + private final AtomicInteger handleUpstream = new AtomicInteger(); private final AtomicInteger beforeSend = new AtomicInteger(); private final AtomicInteger received = new AtomicInteger(); private final CountDownLatch completed = new CountDownLatch(2); @@ -173,6 +176,17 @@ public Flux decode(ClientConnection connection, return decoder; } + @Override + public Flux handleUpstream(ClientConnection connection, + DeviceSession session, + EncodedMessage origin, + Flux decoder, + UnaryOperator> platformHandler) { + handleUpstream.incrementAndGet(); + return DeviceGatewayMonitor.super + .handleUpstream(connection, session, origin, decoder, platformHandler); + } + @Override public Flux beforeSendToPlatform(ClientConnection connection, DeviceSession session,