Skip to content

Commit e485401

Browse files
committed
fix(设备网关): 补齐上行监控包装器委托
1 parent 523aee7 commit e485401

4 files changed

Lines changed: 172 additions & 0 deletions

File tree

jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/CompositeDeviceGatewayMonitor.java

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import java.util.Collection;
3030
import java.util.List;
3131
import java.util.function.Consumer;
32+
import java.util.function.UnaryOperator;
3233

3334
class CompositeDeviceGatewayMonitor implements DeviceGatewayMonitor {
3435

@@ -121,6 +122,21 @@ public Flux<DeviceMessage> decode(@Nullable ClientConnection connection,
121122
return decoder;
122123
}
123124

125+
@Override
126+
public Flux<DeviceMessage> handleUpstream(@Nullable ClientConnection connection,
127+
DeviceSession session,
128+
EncodedMessage origin,
129+
Flux<DeviceMessage> decoder,
130+
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
131+
// 平台处理器作为链尾只组合一次,同时保留 beforeSendToPlatform 的注册顺序。
132+
UnaryOperator<Flux<DeviceMessage>> handler = platformHandler;
133+
for (DeviceGatewayMonitor monitor : monitors) {
134+
UnaryOperator<Flux<DeviceMessage>> next = handler;
135+
handler = source -> monitor.handleUpstream(connection, session, origin, source, next);
136+
}
137+
return handler.apply(decoder);
138+
}
139+
124140
@Override
125141
public Flux<DeviceMessage> beforeSendToPlatform(@Nullable ClientConnection connection,
126142
DeviceSession session,

jetlinks-components/gateway-component/src/main/java/org/jetlinks/community/gateway/monitor/LazyDeviceGatewayMonitor.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424

2525
import javax.annotation.Nullable;
2626
import java.util.function.Supplier;
27+
import java.util.function.UnaryOperator;
2728

2829
class LazyDeviceGatewayMonitor implements DeviceGatewayMonitor {
2930

@@ -101,6 +102,15 @@ public Flux<DeviceMessage> decode(@Nullable ClientConnection connection,
101102
return getTarget().decode(connection, session, origin, decoder);
102103
}
103104

105+
@Override
106+
public Flux<DeviceMessage> handleUpstream(@Nullable ClientConnection connection,
107+
DeviceSession session,
108+
EncodedMessage origin,
109+
Flux<DeviceMessage> decoder,
110+
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
111+
return getTarget().handleUpstream(connection, session, origin, decoder, platformHandler);
112+
}
113+
104114
@Override
105115
public Flux<DeviceMessage> beforeSendToPlatform(@Nullable ClientConnection connection,
106116
DeviceSession session,

jetlinks-components/gateway-component/src/test/java/org/jetlinks/community/gateway/monitor/DeviceGatewayMonitorTest.java

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -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
}

jetlinks-components/network-component/tcp-component/src/test/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGatewayMonitorTest.java

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@
5353
import java.util.concurrent.CountDownLatch;
5454
import java.util.concurrent.TimeUnit;
5555
import java.util.concurrent.atomic.AtomicInteger;
56+
import java.util.function.UnaryOperator;
5657

5758
import static org.junit.jupiter.api.Assertions.assertEquals;
5859
import static org.junit.jupiter.api.Assertions.assertSame;
@@ -132,6 +133,7 @@ public Mono<DeviceMessageCodec> getMessageCodec(Transport transport) {
132133

133134
assertEquals(1, monitor.beforeDecode.get());
134135
assertEquals(1, monitor.decode.get());
136+
assertEquals(2, monitor.handleUpstream.get());
135137
assertEquals(2, monitor.beforeSend.get());
136138
assertEquals(2, monitor.received.get());
137139
assertEquals(1, Collections.frequency(monitor.monitored, manualMessage));
@@ -146,6 +148,7 @@ public Mono<DeviceMessageCodec> getMessageCodec(Transport transport) {
146148
private static class RecordingMonitor implements DeviceGatewayMonitor {
147149
private final AtomicInteger beforeDecode = new AtomicInteger();
148150
private final AtomicInteger decode = new AtomicInteger();
151+
private final AtomicInteger handleUpstream = new AtomicInteger();
149152
private final AtomicInteger beforeSend = new AtomicInteger();
150153
private final AtomicInteger received = new AtomicInteger();
151154
private final CountDownLatch completed = new CountDownLatch(2);
@@ -173,6 +176,17 @@ public Flux<DeviceMessage> decode(ClientConnection connection,
173176
return decoder;
174177
}
175178

179+
@Override
180+
public Flux<DeviceMessage> handleUpstream(ClientConnection connection,
181+
DeviceSession session,
182+
EncodedMessage origin,
183+
Flux<DeviceMessage> decoder,
184+
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
185+
handleUpstream.incrementAndGet();
186+
return DeviceGatewayMonitor.super
187+
.handleUpstream(connection, session, origin, decoder, platformHandler);
188+
}
189+
176190
@Override
177191
public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
178192
DeviceSession session,

0 commit comments

Comments
 (0)