Fix race condition in DebeziumMessProducerTests
The `then(debeziumEngineMock).should().run()` cannot be just checked after `debeziumMessageProducer.start()`: the `DebeziumEngine` is really started on a separate thread. * Check for a `run()` interaction with the mock already after calling `debeziumMessageProducer.stop()`. The `stop()` waits for an internal `latch` which is fulfilled when `DebeziumEngine` exists from its `run()` cycle * Rename `DebeziumMessageProducer.latch` to `lifecycleLatch` to give it more sense.
This commit is contained in:
@@ -77,7 +77,7 @@ public class DebeziumMessageProducer extends MessageProducerSupport {
|
||||
|
||||
private ThreadFactory threadFactory;
|
||||
|
||||
private volatile CountDownLatch latch = new CountDownLatch(0);
|
||||
private volatile CountDownLatch lifecycleLatch = new CountDownLatch(0);
|
||||
|
||||
/**
|
||||
* Create new Debezium message producer inbound channel adapter.
|
||||
@@ -174,10 +174,10 @@ public class DebeziumMessageProducer extends MessageProducerSupport {
|
||||
|
||||
@Override
|
||||
protected void doStart() {
|
||||
if (this.latch.getCount() > 0) {
|
||||
if (this.lifecycleLatch.getCount() > 0) {
|
||||
return;
|
||||
}
|
||||
this.latch = new CountDownLatch(1);
|
||||
this.lifecycleLatch = new CountDownLatch(1);
|
||||
this.executorService.execute(() -> {
|
||||
try {
|
||||
// Runs the debezium connector and deliver database changes to the registered consumer. This method
|
||||
@@ -191,7 +191,7 @@ public class DebeziumMessageProducer extends MessageProducerSupport {
|
||||
this.debeziumEngine.run();
|
||||
}
|
||||
finally {
|
||||
this.latch.countDown();
|
||||
this.lifecycleLatch.countDown();
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -205,7 +205,7 @@ public class DebeziumMessageProducer extends MessageProducerSupport {
|
||||
logger.warn(e, "Debezium failed to close!");
|
||||
}
|
||||
try {
|
||||
if (!this.latch.await(5, TimeUnit.SECONDS)) {
|
||||
if (!this.lifecycleLatch.await(5, TimeUnit.SECONDS)) {
|
||||
throw new IllegalStateException("Failed to stop " + this);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -38,6 +38,7 @@ import static org.mockito.Mockito.reset;
|
||||
|
||||
/**
|
||||
* @author Christian Tzolov
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 6.2
|
||||
*/
|
||||
@@ -76,12 +77,15 @@ public class DebeziumMessageProducerTests {
|
||||
|
||||
await().atMost(5, TimeUnit.SECONDS).until(() -> debeziumMessageProducer.isRunning());
|
||||
assertThat(debeziumMessageProducer.isActive()).isEqualTo(true);
|
||||
then(debeziumEngineMock).should().run();
|
||||
|
||||
debeziumMessageProducer.stop(); // STOP
|
||||
|
||||
assertThat(debeziumMessageProducer.isActive()).isEqualTo(false);
|
||||
assertThat(debeziumMessageProducer.isRunning()).isEqualTo(false);
|
||||
|
||||
// The DebeziumEngine is started on a different thread.
|
||||
// Only the way to catch the run() mock is to stop DebeziumMessageProducer and wait for its internal latch
|
||||
then(debeziumEngineMock).should().run();
|
||||
then(debeziumEngineMock).should().close();
|
||||
|
||||
reset(debeziumEngineMock);
|
||||
@@ -90,10 +94,10 @@ public class DebeziumMessageProducerTests {
|
||||
|
||||
await().atMost(5, TimeUnit.SECONDS).until(() -> debeziumMessageProducer.isRunning());
|
||||
assertThat(debeziumMessageProducer.isActive()).isEqualTo(true);
|
||||
then(debeziumEngineMock).should().run();
|
||||
|
||||
debeziumMessageProducer.destroy(); // DESTROY
|
||||
|
||||
then(debeziumEngineMock).should().run();
|
||||
then(debeziumEngineMock).should().close();
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user