Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

CIT: to verify ReceivePump releases ThreadPool upon Close #322

Merged
merged 9 commits into from
May 18, 2018
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,57 @@ public void VerifyTaskQueueEmptyOnMsgFactoryGracefulClose() throws Exception {
}
}

@Test()
public void VerifyTaskQueueEmptyOnMsgFactoryWithPumpGracefulClose() throws Exception {

final LinkedBlockingQueue<Runnable> blockingQueue = new LinkedBlockingQueue<>();
final ThreadPoolExecutor executor = new ThreadPoolExecutor(
1, 1, 1, TimeUnit.MINUTES, blockingQueue);
try {
final EventHubClient ehClient = EventHubClient.createSync(
TestContext.getConnectionString().toString(),
executor);

final PartitionReceiver receiver = ehClient.createReceiverSync(
TestContext.getConsumerGroupName(), PARTITION_ID, EventPosition.fromEnqueuedTime(Instant.now()));

final CompletableFuture<Iterable<EventData>> signalReceive = new CompletableFuture<>();
receiver.setReceiveHandler(new PartitionReceiveHandler() {
@Override
public int getMaxEventCount() {
return 10;
}

@Override
public void onReceive(Iterable<EventData> events) {
signalReceive.complete(events);
}

@Override
public void onError(Throwable error) {
}
}, false);

final PartitionSender sender = ehClient.createPartitionSenderSync(PARTITION_ID);
sender.sendSync(EventData.create("test data - string".getBytes()));

final Iterable<EventData> events = signalReceive.get();
Assert.assertTrue(events.iterator().hasNext());

receiver.setReceiveHandler(null).get();

sender.closeSync();
receiver.closeSync();

ehClient.closeSync();

Assert.assertEquals(blockingQueue.size(), 0);
Assert.assertEquals(executor.getTaskCount(), executor.getCompletedTaskCount());
} finally {
executor.shutdown();
}
}

@Test()
public void VerifyThreadReleaseOnMsgFactoryOpenError() throws Exception {

Expand Down