|
|
@ -100,44 +100,43 @@ public class TbRpcServiceTest { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void sameRpcIdNotificationsRunInSubmissionOrderOnStripedExecutor() throws InterruptedException { |
|
|
public void sameRpcIdUpdateNotificationsRunInSubmissionOrderOnStripedExecutor() throws InterruptedException { |
|
|
// The count is incidental here: a single rpcId always maps to one stripe. What the test really
|
|
|
// A single rpcId always maps to one stripe; this checks that two update notifications for the
|
|
|
// checks is that two notifications for the SAME rpcId are serialized on that one stripe.
|
|
|
// SAME rpcId are serialized on that stripe (DELIVERED before SUCCESSFUL).
|
|
|
tbRpcService = new TbRpcService(rpcService, clusterService, 3); |
|
|
tbRpcService = new TbRpcService(rpcService, clusterService, 3); |
|
|
|
|
|
|
|
|
RpcId rpcId = new RpcId(UUID.randomUUID()); |
|
|
RpcId rpcId = new RpcId(UUID.randomUUID()); |
|
|
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
|
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|
|
Rpc queued = newRpc(rpcId, deviceId, RpcStatus.QUEUED); |
|
|
|
|
|
Rpc delivered = newRpc(rpcId, deviceId, RpcStatus.DELIVERED); |
|
|
Rpc delivered = newRpc(rpcId, deviceId, RpcStatus.DELIVERED); |
|
|
when(rpcService.createAsync(queued)).thenReturn(Futures.immediateFuture(true)); |
|
|
Rpc successful = newRpc(rpcId, deviceId, RpcStatus.SUCCESSFUL); |
|
|
when(rpcService.updateAsync(delivered)).thenReturn(Futures.immediateFuture(true)); |
|
|
when(rpcService.updateAsync(delivered)).thenReturn(Futures.immediateFuture(true)); |
|
|
|
|
|
when(rpcService.updateAsync(successful)).thenReturn(Futures.immediateFuture(true)); |
|
|
|
|
|
|
|
|
// Make the QUEUED notification block while it runs: it signals that it has started, then sleeps -
|
|
|
// Make the DELIVERED notification block while it runs (signal start, then sleep) - holding the
|
|
|
// holding the stripe. If the two callbacks for this rpcId were NOT serialized on one stripe, the
|
|
|
// stripe. If the two callbacks for this rpcId were NOT serialized on one stripe, the fast
|
|
|
// fast DELIVERED callback would overtake the sleeping QUEUED one and be recorded first.
|
|
|
// SUCCESSFUL one would overtake the sleeping DELIVERED one and be recorded first.
|
|
|
CountDownLatch queuedStarted = new CountDownLatch(1); |
|
|
CountDownLatch deliveredStarted = new CountDownLatch(1); |
|
|
doAnswer(invocation -> { |
|
|
doAnswer(invocation -> { |
|
|
TbMsg msg = invocation.getArgument(2); |
|
|
TbMsg msg = invocation.getArgument(2); |
|
|
if (msg.getInternalType() == TbMsgType.RPC_QUEUED) { |
|
|
if (msg.getInternalType() == TbMsgType.RPC_DELIVERED) { |
|
|
queuedStarted.countDown(); |
|
|
deliveredStarted.countDown(); |
|
|
Thread.sleep(300); |
|
|
Thread.sleep(300); |
|
|
} |
|
|
} |
|
|
return null; |
|
|
return null; |
|
|
}).when(clusterService).pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), any(TbMsg.class), isNull()); |
|
|
}).when(clusterService).pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), any(TbMsg.class), isNull()); |
|
|
|
|
|
|
|
|
tbRpcService.create(queued.getTenantId(), queued); |
|
|
|
|
|
// Don't submit DELIVERED until QUEUED is actually in flight (and now sleeping) on the stripe -
|
|
|
|
|
|
// this makes the test about stripe serialization, not about submission timing.
|
|
|
|
|
|
assertTrue(queuedStarted.await(5, TimeUnit.SECONDS)); |
|
|
|
|
|
tbRpcService.update(delivered.getTenantId(), delivered); |
|
|
tbRpcService.update(delivered.getTenantId(), delivered); |
|
|
|
|
|
// Don't submit SUCCESSFUL until DELIVERED is in flight (and now sleeping) on the stripe.
|
|
|
|
|
|
assertTrue(deliveredStarted.await(5, TimeUnit.SECONDS)); |
|
|
|
|
|
tbRpcService.update(successful.getTenantId(), successful); |
|
|
|
|
|
|
|
|
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
verify(clusterService, timeout(5000).times(2)) |
|
|
verify(clusterService, timeout(5000).times(2)) |
|
|
.pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), msgCaptor.capture(), isNull()); |
|
|
.pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), msgCaptor.capture(), isNull()); |
|
|
|
|
|
|
|
|
List<TbMsg> msgs = msgCaptor.getAllValues(); |
|
|
List<TbMsg> msgs = msgCaptor.getAllValues(); |
|
|
assertEquals(TbMsgType.RPC_QUEUED, msgs.get(0).getInternalType()); |
|
|
assertEquals(TbMsgType.RPC_DELIVERED, msgs.get(0).getInternalType()); |
|
|
assertEquals(TbMsgType.RPC_DELIVERED, msgs.get(1).getInternalType()); |
|
|
assertEquals(TbMsgType.RPC_SUCCESSFUL, msgs.get(1).getInternalType()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private Rpc newRpc() { |
|
|
private Rpc newRpc() { |
|
|
|