private void connectAndAssert(Supplier<Matcher<List<String>>> requestMatcher)
throws InterruptedException, ExecutionException, IOException {
connectNode();
+ readMessage(requestMatcher);
+ }
+
+ private void readMessage(Supplier<Matcher<List<String>>> requestMatcher) throws IOException {
lines = fcpServer.collectUntil(is("EndMessage"));
identifier = extractIdentifier(lines);
assertThat(lines, requestMatcher.get());
public class UskSubscriptionCommands {
- private static final String URI = "SSK@some,uri/file.txt";
+ private static final String URI = "USK@some,uri/file.txt";
@Test
public void subscriptionWorks() throws InterruptedException, ExecutionException, IOException {
Future<Optional<UskSubscription>> uskSubscription = fcpClient.subscribeUsk().uri(URI).execute();
connectAndAssert(() -> matchesFcpMessage("SubscribeUSK", "URI=" + URI, "EndMessage"));
- fcpServer.writeLine(
- "SubscribedUSK",
- "Identifier=" + identifier,
- "URI=" + URI,
- "DontPoll=false",
- "EndMessage"
- );
+ replyWithSubscribed();
+ assertThat(uskSubscription.get().get().getUri(), is(URI));
+ AtomicInteger edition = new AtomicInteger();
+ CountDownLatch updated = new CountDownLatch(2);
+ uskSubscription.get().get().onUpdate(e -> {
+ edition.set(e);
+ updated.countDown();
+ });
+ sendUpdateNotification(23);
+ sendUpdateNotification(24);
+ assertThat("updated in time", updated.await(5, TimeUnit.SECONDS), is(true));
+ assertThat(edition.get(), is(24));
+ }
+
+ @Test
+ public void subscriptionUpdatesMultipleTimes() throws InterruptedException, ExecutionException, IOException {
+ Future<Optional<UskSubscription>> uskSubscription = fcpClient.subscribeUsk().uri(URI).execute();
+ connectAndAssert(() -> matchesFcpMessage("SubscribeUSK", "URI=" + URI, "EndMessage"));
+ replyWithSubscribed();
assertThat(uskSubscription.get().get().getUri(), is(URI));
AtomicInteger edition = new AtomicInteger();
CountDownLatch updated = new CountDownLatch(2);
edition.set(e);
updated.countDown();
});
+ uskSubscription.get().get().onUpdate(e -> updated.countDown());
+ sendUpdateNotification(23);
+ assertThat("updated in time", updated.await(5, TimeUnit.SECONDS), is(true));
+ assertThat(edition.get(), is(23));
+ }
+
+ private void replyWithSubscribed() throws IOException {
fcpServer.writeLine(
- "SubscribedUSKUpdate",
+ "SubscribedUSK",
"Identifier=" + identifier,
"URI=" + URI,
- "Edition=23",
+ "DontPoll=false",
"EndMessage"
);
+ }
+
+ private void sendUpdateNotification(int edition, String... additionalLines) throws IOException {
fcpServer.writeLine(
"SubscribedUSKUpdate",
"Identifier=" + identifier,
"URI=" + URI,
- "Edition=24",
- "EndMessage"
+ "Edition=" + edition
);
- assertThat("updated in time", updated.await(5, TimeUnit.SECONDS), is(true));
- assertThat(edition.get(), is(24));
+ fcpServer.writeLine(additionalLines);
+ fcpServer.writeLine("EndMessage");
}
}