Skip to content

Commit 5456669

Browse files
committed
fix(设备网关): 补齐协议手动上报监控
1 parent 5293d48 commit 5456669

8 files changed

Lines changed: 155 additions & 45 deletions

File tree

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -164,8 +164,9 @@ default Flux<DeviceMessage> decode(@Nullable ClientConnection connection,
164164
* 使解码监控覆盖协议解码、平台处理和发送前处理的完整链路。
165165
* {@code platformHandler} 只能组合传入的解码任务,不得主动订阅。</p>
166166
*
167-
* <p>{@link FromDeviceMessageContext#handleMessage(DeviceMessage)} 在协议解码任务中手动处理消息时,
168-
* 仍会继承发送前监控写入的 Reactor Context;协议随后返回空流表示没有额外的返回值消息。</p>
167+
* <p>{@link FromDeviceMessageContext#handleMessage(DeviceMessage)} 手动输出的消息不会进入协议返回的
168+
* {@code decoder},调用方应使用本方法单独包装该消息的平台处理任务。协议随后返回空流表示没有额外的
169+
* 返回值消息,不能再次处理已经手动输出的消息。</p>
169170
*
170171
* @param connection 客户端连接,短连接场景可能为 {@code null}
171172
* @param session 当前设备会话

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

Lines changed: 22 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
import java.util.Arrays;
3131
import java.util.List;
3232
import java.util.concurrent.atomic.AtomicInteger;
33+
import java.util.function.UnaryOperator;
3334

3435
import static org.junit.jupiter.api.Assertions.assertEquals;
3536
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -121,15 +122,10 @@ void shouldWrapManualProtocolOutputWithExistingMonitorMethods() {
121122
String monitorContextKey = DeviceGatewayMonitorTest.class.getName();
122123
AtomicInteger decode = new AtomicInteger();
123124
AtomicInteger beforeSend = new AtomicInteger();
124-
AtomicInteger received = new AtomicInteger();
125+
AtomicInteger monitored = new AtomicInteger();
125126
AtomicInteger manualHandled = new AtomicInteger();
126127
AtomicInteger returnedHandled = new AtomicInteger();
127128
DeviceGatewayMonitor monitor = new DeviceGatewayMonitor() {
128-
@Override
129-
public void receivedMessage() {
130-
received.incrementAndGet();
131-
}
132-
133129
@Override
134130
public Flux<DeviceMessage> decode(ClientConnection connection,
135131
DeviceSession session,
@@ -145,7 +141,9 @@ public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
145141
EncodedMessage origin,
146142
Flux<DeviceMessage> handler) {
147143
beforeSend.incrementAndGet();
148-
return handler.contextWrite(context -> context.put(monitorContextKey, true));
144+
return handler
145+
.doOnNext(ignore -> monitored.incrementAndGet())
146+
.contextWrite(context -> context.put(monitorContextKey, true));
149147
}
150148
};
151149

@@ -154,18 +152,26 @@ public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
154152
EncodedMessage origin = mock(EncodedMessage.class);
155153
DeviceRegistry registry = mock(DeviceRegistry.class);
156154
DeviceMessage message = mock(DeviceMessage.class);
155+
UnaryOperator<Flux<DeviceMessage>> platformHandler = decoded -> decoded
156+
.concatMap(current -> Mono.deferContextual(ctx -> {
157+
assertTrue(ctx.getOrDefault(monitorContextKey, false));
158+
assertSame(message, current);
159+
manualHandled.incrementAndGet();
160+
return Mono.just(current);
161+
}));
157162
FromDeviceMessageContext context = FromDeviceMessageContext.of(
158163
session,
159164
origin,
160165
registry,
161166
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-
})
167+
current -> monitor
168+
.handleUpstream(
169+
connection,
170+
session,
171+
origin,
172+
Flux.just(current),
173+
platformHandler)
174+
.then()
169175
);
170176

171177
Flux<DeviceMessage> upstream = monitor.decode(
@@ -188,8 +194,8 @@ public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
188194
.verifyComplete();
189195

190196
assertEquals(1, decode.get());
191-
assertEquals(1, beforeSend.get());
192-
assertEquals(1, received.get());
197+
assertEquals(2, beforeSend.get());
198+
assertEquals(1, monitored.get());
193199
assertEquals(1, manualHandled.get());
194200
assertEquals(0, returnedHandled.get());
195201
}

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

Lines changed: 36 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@
5555
import java.util.List;
5656
import java.util.Map;
5757
import java.util.concurrent.ConcurrentHashMap;
58+
import java.util.function.UnaryOperator;
5859

5960
/**
6061
* Http 服务设备网关,使用指定的协议包,将网络组件中Http服务的请求处理为设备消息
@@ -166,6 +167,11 @@ private Mono<Void> handleWebsocketRequest(WebSocketExchange exchange, WebSocketM
166167

167168
WebSocketDeviceSession session = new WebSocketDeviceSession(monitor, device, exchange);
168169

170+
UnaryOperator<Flux<DeviceMessage>> platformHandler = task -> task
171+
.concatMap(deviceMessage -> handleWebsocketMessage(deviceMessage, exchange, session)
172+
.doOnNext(session::setOperator)
173+
.thenReturn(deviceMessage));
174+
169175
Flux<DeviceMessage> decodeTask = protocol
170176
.flatMapMany(protocol -> {
171177
if (log.isDebugEnabled()) {
@@ -175,7 +181,18 @@ private Mono<Void> handleWebsocketRequest(WebSocketExchange exchange, WebSocketM
175181
return protocol
176182
.getMessageCodec(DefaultTransport.WebSocket)
177183
.flatMapMany(codec -> codec.decode(FromDeviceMessageContext.of(
178-
session, msg, registry, deviceMessage -> handleWebsocketMessage(deviceMessage, exchange, session).then())))
184+
session,
185+
msg,
186+
registry,
187+
// 手动输出不进入 codec 返回值,单独复用同一平台处理与发送前监控链。
188+
deviceMessage -> monitor
189+
.handleUpstream(
190+
exchange,
191+
session,
192+
msg,
193+
Flux.just(deviceMessage),
194+
platformHandler)
195+
.then())))
179196
.cast(DeviceMessage.class);
180197
});
181198

@@ -184,9 +201,7 @@ private Mono<Void> handleWebsocketRequest(WebSocketExchange exchange, WebSocketM
184201
session,
185202
msg,
186203
decodeTask,
187-
task -> task.concatMap(deviceMessage -> handleWebsocketMessage(deviceMessage, exchange, session)
188-
.doOnNext(session::setOperator)
189-
.thenReturn(deviceMessage))
204+
platformHandler
190205
);
191206
decodeTask = monitor.decode(exchange, session, msg, decodeTask);
192207

@@ -256,20 +271,33 @@ private Mono<Void> handleHttpRequest(HttpExchange exchange) {
256271
if (!monitor.beforeDecode(null, httpMessage)) {
257272
return completeHttpRequest(exchange);
258273
}
274+
UnaryOperator<Flux<DeviceMessage>> platformHandler = task -> task
275+
.concatMap(deviceMessage ->
276+
handleMessage(deviceMessage, exchange, httpMessage)
277+
.thenReturn(deviceMessage));
259278
//调用协议执行解码
260279
Flux<DeviceMessage> decodeTask = protocol
261280
.getMessageCodec(getTransport())
262281
.flatMapMany(codec -> codec.decode(FromDeviceMessageContext.of(
263-
session, httpMessage, registry, msg -> handleMessage(msg, exchange, httpMessage))))
282+
session,
283+
httpMessage,
284+
registry,
285+
// 手动输出不进入 codec 返回值,单独复用同一平台处理与发送前监控链。
286+
deviceMessage -> monitor
287+
.handleUpstream(
288+
null,
289+
session,
290+
httpMessage,
291+
Flux.just(deviceMessage),
292+
platformHandler)
293+
.then())))
264294
.cast(DeviceMessage.class);
265295
decodeTask = monitor.handleUpstream(
266296
null,
267297
session,
268298
httpMessage,
269299
decodeTask,
270-
task -> task.concatMap(deviceMessage ->
271-
handleMessage(deviceMessage, exchange, httpMessage)
272-
.thenReturn(deviceMessage))
300+
platformHandler
273301
);
274302
decodeTask = monitor.decode(null, session, httpMessage, decodeTask);
275303
return decodeTask

jetlinks-components/network-component/http-component/src/test/java/org/jetlinks/community/network/http/device/HttpServerDeviceGatewayTest.java

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import org.jetlinks.core.message.codec.DefaultTransport;
2828
import org.jetlinks.core.message.codec.DeviceMessageCodec;
2929
import org.jetlinks.core.message.codec.EncodedMessage;
30+
import org.jetlinks.core.message.codec.FromDeviceMessageContext;
3031
import org.jetlinks.core.message.codec.http.HttpExchangeMessage;
3132
import org.jetlinks.core.route.HttpRoute;
3233
import org.jetlinks.core.server.ClientConnection;
@@ -40,7 +41,9 @@
4041
import reactor.test.StepVerifier;
4142

4243
import java.time.Duration;
44+
import java.util.List;
4345
import java.util.concurrent.CountDownLatch;
46+
import java.util.concurrent.CopyOnWriteArrayList;
4447
import java.util.concurrent.atomic.AtomicInteger;
4548

4649
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -54,7 +57,7 @@
5457
class HttpServerDeviceGatewayTest {
5558

5659
@Test
57-
void httpUpstreamShouldUseDecodeMonitors() throws InterruptedException {
60+
void manualHttpUpstreamShouldUseMonitorChain() throws InterruptedException {
5861
String gatewayId = "http-monitor-test";
5962
RecordingMonitor monitor = new RecordingMonitor();
6063
GatewayMonitors.register((id, tags) -> gatewayId.equals(id) ? monitor : null);
@@ -68,6 +71,7 @@ void httpUpstreamShouldUseDecodeMonitors() throws InterruptedException {
6871
HttpRoute route = mock(HttpRoute.class);
6972
HttpExchange exchange = mock(HttpExchange.class);
7073
HttpExchangeMessage message = mock(HttpExchangeMessage.class);
74+
DeviceMessage deviceMessage = mock(DeviceMessage.class);
7175
HttpRequest request = mock(HttpRequest.class);
7276
Sinks.Many<HttpExchange> requests = Sinks.many().unicast().onBackpressureBuffer();
7377

@@ -77,7 +81,12 @@ void httpUpstreamShouldUseDecodeMonitors() throws InterruptedException {
7781
when(protocol.getRoutes(DefaultTransport.WebSocket)).thenReturn(Flux.empty());
7882
when(protocol.getMessageCodec(DefaultTransport.HTTP))
7983
.thenAnswer(ignore -> Mono.just(codec));
80-
when(codec.decode(any())).thenReturn(Flux.empty());
84+
when(codec.decode(any())).thenAnswer(invocation -> {
85+
FromDeviceMessageContext context = invocation.getArgument(0);
86+
return context
87+
.handleMessage(deviceMessage)
88+
.thenMany(Flux.empty());
89+
});
8190
when(server.handleRequest(HttpMethod.POST, "/test")).thenReturn(requests.asFlux());
8291
when(exchange.toExchangeMessage()).thenReturn(Mono.just(message));
8392
when(exchange.request()).thenReturn(request);
@@ -99,7 +108,9 @@ void httpUpstreamShouldUseDecodeMonitors() throws InterruptedException {
99108

100109
assertEquals(1, monitor.beforeDecode.get());
101110
assertEquals(1, monitor.decode.get());
102-
assertEquals(1, monitor.beforeSend.get());
111+
assertEquals(2, monitor.beforeSend.get());
112+
assertEquals(1, monitor.received.get());
113+
assertEquals(List.of(deviceMessage), monitor.monitored);
103114
assertSame(message, monitor.origin);
104115

105116
StepVerifier.create(gateway.shutdown()).verifyComplete();
@@ -109,9 +120,16 @@ private static class RecordingMonitor implements DeviceGatewayMonitor {
109120
private final AtomicInteger beforeDecode = new AtomicInteger();
110121
private final AtomicInteger decode = new AtomicInteger();
111122
private final AtomicInteger beforeSend = new AtomicInteger();
123+
private final AtomicInteger received = new AtomicInteger();
112124
private final CountDownLatch completed = new CountDownLatch(1);
125+
private final List<DeviceMessage> monitored = new CopyOnWriteArrayList<>();
113126
private EncodedMessage origin;
114127

128+
@Override
129+
public void receivedMessage() {
130+
received.incrementAndGet();
131+
}
132+
115133
@Override
116134
public boolean beforeDecode(ClientConnection connection, EncodedMessage message) {
117135
assertNull(connection);
@@ -137,8 +155,10 @@ public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
137155
Flux<DeviceMessage> handler) {
138156
assertNull(connection);
139157
beforeSend.incrementAndGet();
140-
completed.countDown();
141-
return handler;
158+
return handler.doOnNext(message -> {
159+
monitored.add(message);
160+
completed.countDown();
161+
});
142162
}
143163
}
144164
}

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

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@
5656
import java.net.InetSocketAddress;
5757
import java.util.function.Consumer;
5858
import java.util.function.Function;
59+
import java.util.function.UnaryOperator;
5960

6061
import static org.jetlinks.community.network.mqtt.gateway.device.MqttServerDeviceGateway.clientId;
6162

@@ -221,14 +222,25 @@ protected Mono<Void> decodeAndHandleMessage(MqttMessage message) {
221222
return Mono.empty();
222223
}
223224

225+
UnaryOperator<Flux<DeviceMessage>> platformHandler = task -> task
226+
.concatMap(this::handleMessage, 0);
227+
224228
// 上下文
225229
FromDeviceMessageContext context =
226230
FromDeviceMessageContext
227231
.of(session,
228232
message,
229233
helper.getRegistry(),
230234
connection,
231-
this);
235+
// 手动输出不进入 codec 返回值,单独复用同一平台处理与发送前监控链。
236+
deviceMessage -> monitor
237+
.handleUpstream(
238+
connection,
239+
session,
240+
message,
241+
Flux.just(deviceMessage),
242+
platformHandler)
243+
.then());
232244

233245
Flux<DeviceMessage> decodeTask = operator
234246
.getProtocol()
@@ -242,8 +254,8 @@ protected Mono<Void> decodeAndHandleMessage(MqttMessage message) {
242254
session,
243255
message,
244256
decodeTask,
245-
task -> task
246-
.concatMap(this::handleMessage, 0)
257+
task -> platformHandler
258+
.apply(task)
247259
.doOnComplete(() -> {
248260
if (message instanceof MqttPublishing) {
249261
((MqttPublishing) message).acknowledge();

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

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@
4444

4545
import java.util.*;
4646
import java.util.concurrent.ConcurrentHashMap;
47+
import java.util.function.UnaryOperator;
4748

4849
/**
4950
* MQTT Client 设备网关,使用网络组件中的MQTT Client来处理设备数据
@@ -168,20 +169,29 @@ private Mono<Void> decodeAndHandleMessage(MqttMessage mqttMessage) {
168169
}
169170
UnknownDeviceMqttClientSession session =
170171
new UnknownDeviceMqttClientSession(getId(), mqttClient, monitor);
172+
UnaryOperator<Flux<DeviceMessage>> platformHandler = task -> task
173+
.concatMap(message -> handleMessage(mqttMessage, message).thenReturn(message));
171174
Flux<DeviceMessage> decodeTask = codecMono
172175
.flatMapMany(codec -> codec.decode(FromDeviceMessageContext.of(
173176
session,
174177
mqttMessage,
175178
registry,
176-
msg -> handleMessage(mqttMessage, msg).then())))
179+
// 手动输出不进入 codec 返回值,单独复用同一平台处理与发送前监控链。
180+
message -> monitor
181+
.handleUpstream(
182+
null,
183+
session,
184+
mqttMessage,
185+
Flux.just(message),
186+
platformHandler)
187+
.then())))
177188
.cast(DeviceMessage.class);
178189
decodeTask = monitor.handleUpstream(
179190
null,
180191
session,
181192
mqttMessage,
182193
decodeTask,
183-
task -> task.concatMap(message ->
184-
handleMessage(mqttMessage, message).thenReturn(message))
194+
platformHandler
185195
);
186196
decodeTask = monitor.decode(null, session, mqttMessage, decodeTask);
187197
return decodeTask.then();

0 commit comments

Comments
 (0)