Skip to content

Commit 5293d48

Browse files
committed
fix(设备网关): 完善协议上行监控
1 parent db894b8 commit 5293d48

6 files changed

Lines changed: 220 additions & 25 deletions

File tree

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

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,14 @@
1717

1818
import org.jetlinks.core.message.DeviceMessage;
1919
import org.jetlinks.core.message.codec.EncodedMessage;
20+
import org.jetlinks.core.message.codec.FromDeviceMessageContext;
2021
import org.jetlinks.core.server.ClientConnection;
2122
import org.jetlinks.core.server.session.DeviceSession;
2223
import reactor.core.publisher.Flux;
2324
import reactor.core.publisher.Mono;
2425

2526
import javax.annotation.Nullable;
27+
import java.util.function.UnaryOperator;
2628

2729
/**
2830
* 设备网关监控扩展点。
@@ -153,6 +155,36 @@ default Flux<DeviceMessage> decode(@Nullable ClientConnection connection,
153155
return decoder;
154156
}
155157

158+
/**
159+
* 组装平台上行处理任务。
160+
*
161+
* <p>本方法固定按 {@code platformHandler}、
162+
* {@link #beforeSendToPlatform(ClientConnection, DeviceSession, EncodedMessage, Flux)} 的顺序组装处理链。
163+
* 调用方应再使用 {@link #decode(ClientConnection, DeviceSession, EncodedMessage, Flux)} 包装返回任务,
164+
* 使解码监控覆盖协议解码、平台处理和发送前处理的完整链路。
165+
* {@code platformHandler} 只能组合传入的解码任务,不得主动订阅。</p>
166+
*
167+
* <p>{@link FromDeviceMessageContext#handleMessage(DeviceMessage)} 在协议解码任务中手动处理消息时,
168+
* 仍会继承发送前监控写入的 Reactor Context;协议随后返回空流表示没有额外的返回值消息。</p>
169+
*
170+
* @param connection 客户端连接,短连接场景可能为 {@code null}
171+
* @param session 当前设备会话
172+
* @param origin 原始报文
173+
* @param decoder 协议解码任务
174+
* @param platformHandler 将解码任务转换为包含平台消息处理逻辑的任务
175+
* @return 待解码监控包装的上行处理任务
176+
* @since 2.12
177+
* @see FromDeviceMessageContext
178+
*/
179+
default Flux<DeviceMessage> handleUpstream(@Nullable ClientConnection connection,
180+
DeviceSession session,
181+
EncodedMessage origin,
182+
Flux<DeviceMessage> decoder,
183+
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
184+
Flux<DeviceMessage> handler = platformHandler.apply(decoder);
185+
return beforeSendToPlatform(connection, session, origin, handler);
186+
}
187+
156188
/**
157189
* 包装解码完成后、发送到平台前的上行处理任务。
158190
*

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

Lines changed: 157 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,10 @@
1515
*/
1616
package org.jetlinks.community.gateway.monitor;
1717

18+
import org.jetlinks.core.device.DeviceRegistry;
1819
import org.jetlinks.core.message.DeviceMessage;
1920
import org.jetlinks.core.message.codec.EncodedMessage;
21+
import org.jetlinks.core.message.codec.FromDeviceMessageContext;
2022
import org.jetlinks.core.server.ClientConnection;
2123
import org.jetlinks.core.server.session.DeviceSession;
2224
import org.junit.jupiter.api.Test;
@@ -37,6 +39,161 @@
3739

3840
class DeviceGatewayMonitorTest {
3941

42+
@Test
43+
void shouldWrapWholeUpstreamWithDecodeAndKeepPlatformHandlerLazy() {
44+
List<String> signals = new ArrayList<>();
45+
AtomicInteger platformHandled = new AtomicInteger();
46+
DeviceGatewayMonitor monitor = new DeviceGatewayMonitor() {
47+
@Override
48+
public Flux<DeviceMessage> decode(ClientConnection connection,
49+
DeviceSession session,
50+
EncodedMessage origin,
51+
Flux<DeviceMessage> decoder) {
52+
signals.add("decode");
53+
return decoder.doOnNext(current -> {
54+
signals.add("decodeOnNext");
55+
assertEquals(1, platformHandled.get());
56+
});
57+
}
58+
59+
@Override
60+
public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
61+
DeviceSession session,
62+
EncodedMessage origin,
63+
Flux<DeviceMessage> handler) {
64+
signals.add("beforeSend");
65+
return handler;
66+
}
67+
};
68+
69+
ClientConnection connection = mock(ClientConnection.class);
70+
DeviceSession session = mock(DeviceSession.class);
71+
EncodedMessage origin = mock(EncodedMessage.class);
72+
DeviceMessage message = mock(DeviceMessage.class);
73+
74+
Flux<DeviceMessage> upstream = monitor.decode(
75+
connection,
76+
session,
77+
origin,
78+
monitor.handleUpstream(
79+
connection,
80+
session,
81+
origin,
82+
Flux.defer(() -> {
83+
signals.add("decoder");
84+
return Flux.just(message);
85+
}),
86+
decoded -> {
87+
signals.add("platformHandler");
88+
return decoded.concatMap(current -> Mono
89+
.fromRunnable(() -> {
90+
signals.add("platform");
91+
platformHandled.incrementAndGet();
92+
})
93+
.thenReturn(current));
94+
}
95+
)
96+
);
97+
98+
assertEquals(Arrays.asList("platformHandler", "beforeSend", "decode"), signals);
99+
assertEquals(0, platformHandled.get());
100+
101+
StepVerifier
102+
.create(upstream)
103+
.expectNext(message)
104+
.verifyComplete();
105+
106+
assertEquals(1, platformHandled.get());
107+
assertEquals(
108+
Arrays.asList(
109+
"platformHandler",
110+
"beforeSend",
111+
"decode",
112+
"decoder",
113+
"platform",
114+
"decodeOnNext"),
115+
signals
116+
);
117+
}
118+
119+
@Test
120+
void shouldWrapManualProtocolOutputWithExistingMonitorMethods() {
121+
String monitorContextKey = DeviceGatewayMonitorTest.class.getName();
122+
AtomicInteger decode = new AtomicInteger();
123+
AtomicInteger beforeSend = new AtomicInteger();
124+
AtomicInteger received = new AtomicInteger();
125+
AtomicInteger manualHandled = new AtomicInteger();
126+
AtomicInteger returnedHandled = new AtomicInteger();
127+
DeviceGatewayMonitor monitor = new DeviceGatewayMonitor() {
128+
@Override
129+
public void receivedMessage() {
130+
received.incrementAndGet();
131+
}
132+
133+
@Override
134+
public Flux<DeviceMessage> decode(ClientConnection connection,
135+
DeviceSession session,
136+
EncodedMessage origin,
137+
Flux<DeviceMessage> decoder) {
138+
decode.incrementAndGet();
139+
return decoder;
140+
}
141+
142+
@Override
143+
public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
144+
DeviceSession session,
145+
EncodedMessage origin,
146+
Flux<DeviceMessage> handler) {
147+
beforeSend.incrementAndGet();
148+
return handler.contextWrite(context -> context.put(monitorContextKey, true));
149+
}
150+
};
151+
152+
ClientConnection connection = mock(ClientConnection.class);
153+
DeviceSession session = mock(DeviceSession.class);
154+
EncodedMessage origin = mock(EncodedMessage.class);
155+
DeviceRegistry registry = mock(DeviceRegistry.class);
156+
DeviceMessage message = mock(DeviceMessage.class);
157+
FromDeviceMessageContext context = FromDeviceMessageContext.of(
158+
session,
159+
origin,
160+
registry,
161+
connection,
162+
current -> Mono.deferContextual(ctx -> {
163+
assertTrue(ctx.getOrDefault(monitorContextKey, false));
164+
assertSame(message, current);
165+
monitor.receivedMessage();
166+
manualHandled.incrementAndGet();
167+
return Mono.empty();
168+
})
169+
);
170+
171+
Flux<DeviceMessage> upstream = monitor.decode(
172+
connection,
173+
session,
174+
origin,
175+
monitor.handleUpstream(
176+
connection,
177+
session,
178+
origin,
179+
Flux.defer(() -> context.handleMessage(message).thenMany(Flux.empty())),
180+
decoded -> decoded.concatMap(current -> Mono
181+
.fromRunnable(returnedHandled::incrementAndGet)
182+
.thenReturn(current))
183+
)
184+
);
185+
186+
StepVerifier
187+
.create(upstream)
188+
.verifyComplete();
189+
190+
assertEquals(1, decode.get());
191+
assertEquals(1, beforeSend.get());
192+
assertEquals(1, received.get());
193+
assertEquals(1, manualHandled.get());
194+
assertEquals(0, returnedHandled.get());
195+
}
196+
40197
@Test
41198
void shouldKeepDefaultMonitorCompatibleAndTransparent() {
42199
AtomicInteger connected = new AtomicInteger();

jetlinks-components/network-component/http-component/src/main/java/org/jetlinks/community/network/http/device/HttpServerDeviceGateway.java

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -179,15 +179,16 @@ private Mono<Void> handleWebsocketRequest(WebSocketExchange exchange, WebSocketM
179179
.cast(DeviceMessage.class);
180180
});
181181

182-
decodeTask = monitor.decode(exchange, session, msg, decodeTask);
183-
decodeTask = monitor.beforeSendToPlatform(
182+
decodeTask = monitor.handleUpstream(
184183
exchange,
185184
session,
186185
msg,
187-
decodeTask.concatMap(deviceMessage -> handleWebsocketMessage(deviceMessage, exchange, session)
186+
decodeTask,
187+
task -> task.concatMap(deviceMessage -> handleWebsocketMessage(deviceMessage, exchange, session)
188188
.doOnNext(session::setOperator)
189189
.thenReturn(deviceMessage))
190190
);
191+
decodeTask = monitor.decode(exchange, session, msg, decodeTask);
191192

192193
return decodeTask
193194
.onErrorResume(err -> {
@@ -261,15 +262,16 @@ private Mono<Void> handleHttpRequest(HttpExchange exchange) {
261262
.flatMapMany(codec -> codec.decode(FromDeviceMessageContext.of(
262263
session, httpMessage, registry, msg -> handleMessage(msg, exchange, httpMessage))))
263264
.cast(DeviceMessage.class);
264-
decodeTask = monitor.decode(null, session, httpMessage, decodeTask);
265-
decodeTask = monitor.beforeSendToPlatform(
265+
decodeTask = monitor.handleUpstream(
266266
null,
267267
session,
268268
httpMessage,
269-
decodeTask.concatMap(deviceMessage ->
269+
decodeTask,
270+
task -> task.concatMap(deviceMessage ->
270271
handleMessage(deviceMessage, exchange, httpMessage)
271272
.thenReturn(deviceMessage))
272273
);
274+
decodeTask = monitor.decode(null, session, httpMessage, decodeTask);
273275
return decodeTask
274276
.then(completeHttpRequest(exchange))
275277
.onErrorResume(err -> {

jetlinks-components/network-component/mqtt-component/src/main/java/org/jetlinks/community/network/mqtt/gateway/device/DeviceMqttConnection.java

Lines changed: 15 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -237,21 +237,22 @@ protected Mono<Void> decodeAndHandleMessage(MqttMessage message) {
237237
.flatMapMany(codec -> codec.decode(context))
238238
.cast(DeviceMessage.class);
239239

240+
decodeTask = monitor.handleUpstream(
241+
connection,
242+
session,
243+
message,
244+
decodeTask,
245+
task -> task
246+
.concatMap(this::handleMessage, 0)
247+
.doOnComplete(() -> {
248+
if (message instanceof MqttPublishing) {
249+
((MqttPublishing) message).acknowledge();
250+
}
251+
})
252+
);
240253
decodeTask = monitor.decode(connection, session, message, decodeTask);
241254

242-
return monitor
243-
.beforeSendToPlatform(
244-
connection,
245-
session,
246-
message,
247-
decodeTask
248-
.concatMap(this::handleMessage, 0)
249-
.doOnComplete(() -> {
250-
if (message instanceof MqttPublishing) {
251-
((MqttPublishing) message).acknowledge();
252-
}
253-
})
254-
)
255+
return decodeTask
255256
.as(FluxTracer
256257
.create(DeviceTracer.SpanName.decode0(operator.getDeviceId()),
257258
(span) -> span
@@ -274,6 +275,7 @@ protected Mono<Void> decodeAndHandleMessage(MqttMessage message) {
274275
}
275276

276277
private Mono<DeviceMessage> handleMessage(DeviceMessage message) {
278+
monitor.receivedMessage();
277279

278280
DeviceOperator mainDevice = session.getOperator();
279281

jetlinks-components/network-component/mqtt-component/src/main/java/org/jetlinks/community/network/mqtt/gateway/device/MqttClientDeviceGateway.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -175,14 +175,15 @@ private Mono<Void> decodeAndHandleMessage(MqttMessage mqttMessage) {
175175
registry,
176176
msg -> handleMessage(mqttMessage, msg).then())))
177177
.cast(DeviceMessage.class);
178-
decodeTask = monitor.decode(null, session, mqttMessage, decodeTask);
179-
decodeTask = monitor.beforeSendToPlatform(
178+
decodeTask = monitor.handleUpstream(
180179
null,
181180
session,
182181
mqttMessage,
183-
decodeTask.concatMap(message ->
182+
decodeTask,
183+
task -> task.concatMap(message ->
184184
handleMessage(mqttMessage, message).thenReturn(message))
185185
);
186+
decodeTask = monitor.decode(null, session, mqttMessage, decodeTask);
186187
return decodeTask.then();
187188
}
188189

jetlinks-components/network-component/tcp-component/src/main/java/org/jetlinks/community/network/tcp/gateway/device/TcpServerDeviceGateway.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -216,13 +216,14 @@ Mono<Void> handleTcpMessage0(EncodedMessage message) {
216216
msg -> handleDeviceMessage(msg).then())))
217217
.cast(DeviceMessage.class);
218218

219-
decodeTask = parent.monitor.decode(client, deviceSession, message, decodeTask);
220-
decodeTask = parent.monitor.beforeSendToPlatform(
219+
decodeTask = parent.monitor.handleUpstream(
221220
client,
222221
deviceSession,
223222
message,
224-
decodeTask.concatMap(this::handleDeviceMessage, 0)
223+
decodeTask,
224+
task -> task.concatMap(this::handleDeviceMessage, 0)
225225
);
226+
decodeTask = parent.monitor.decode(client, deviceSession, message, decodeTask);
226227

227228
return decodeTask
228229
.as(FluxTracer.create(

0 commit comments

Comments
 (0)