Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -121,6 +122,21 @@ public Flux<DeviceMessage> decode(@Nullable ClientConnection connection,
return decoder;
}

@Override
public Flux<DeviceMessage> handleUpstream(@Nullable ClientConnection connection,
DeviceSession session,
EncodedMessage origin,
Flux<DeviceMessage> decoder,
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
// 平台处理器作为链尾只组合一次,同时保留 beforeSendToPlatform 的注册顺序。
UnaryOperator<Flux<DeviceMessage>> handler = platformHandler;
for (DeviceGatewayMonitor monitor : monitors) {
UnaryOperator<Flux<DeviceMessage>> next = handler;
handler = source -> monitor.handleUpstream(connection, session, origin, source, next);
}
return handler.apply(decoder);
}

@Override
public Flux<DeviceMessage> beforeSendToPlatform(@Nullable ClientConnection connection,
DeviceSession session,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import javax.annotation.Nullable;
import java.util.function.Supplier;
import java.util.function.UnaryOperator;

class LazyDeviceGatewayMonitor implements DeviceGatewayMonitor {

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

@Override
public Flux<DeviceMessage> handleUpstream(@Nullable ClientConnection connection,
DeviceSession session,
EncodedMessage origin,
Flux<DeviceMessage> decoder,
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
return getTarget().handleUpstream(connection, session, origin, decoder, platformHandler);
}

@Override
public Flux<DeviceMessage> beforeSendToPlatform(@Nullable ClientConnection connection,
DeviceSession session,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,107 @@ public void rejected() {
assertEquals(1, rejected.get());
}

@Test
void shouldDelegateCompositeHandleUpstreamWithoutDuplicatingPlatformHandler() {
List<String> 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<DeviceMessage> 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<DeviceMessage> 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<String> 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<DeviceMessage> 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<String> signals = new ArrayList<>();
Expand Down Expand Up @@ -420,4 +521,35 @@ public Mono<Void> downstream(ClientConnection connection,
return sender.doOnSuccess(ignore -> signals.add(id + ":downstream"));
}
}

private static class HandleUpstreamRecordingMonitor implements DeviceGatewayMonitor {

private final String id;
private final List<String> signals;
private final AtomicInteger handleUpstreamInvocations = new AtomicInteger();

private HandleUpstreamRecordingMonitor(String id, List<String> signals) {
this.id = id;
this.signals = signals;
}

@Override
public Flux<DeviceMessage> handleUpstream(ClientConnection connection,
DeviceSession session,
EncodedMessage origin,
Flux<DeviceMessage> decoder,
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
handleUpstreamInvocations.incrementAndGet();
return DeviceGatewayMonitor.super
.handleUpstream(connection, session, origin, decoder, platformHandler);
}

@Override
public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
DeviceSession session,
EncodedMessage origin,
Flux<DeviceMessage> handler) {
return handler.doOnNext(ignore -> signals.add(id));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -132,6 +133,7 @@ public Mono<DeviceMessageCodec> 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));
Expand All @@ -146,6 +148,7 @@ public Mono<DeviceMessageCodec> 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);
Expand Down Expand Up @@ -173,6 +176,17 @@ public Flux<DeviceMessage> decode(ClientConnection connection,
return decoder;
}

@Override
public Flux<DeviceMessage> handleUpstream(ClientConnection connection,
DeviceSession session,
EncodedMessage origin,
Flux<DeviceMessage> decoder,
UnaryOperator<Flux<DeviceMessage>> platformHandler) {
handleUpstream.incrementAndGet();
return DeviceGatewayMonitor.super
.handleUpstream(connection, session, origin, decoder, platformHandler);
}

@Override
public Flux<DeviceMessage> beforeSendToPlatform(ClientConnection connection,
DeviceSession session,
Expand Down
Loading