github-advanced-security[bot] commented on code in PR #19858:
URL: https://github.com/apache/druid/pull/19858#discussion_r4083540631
##########
extensions-core/kafka-extraction-namespace/src/test/java/org/apache/druid/query/lookup/KafkaLookupExtractorFactoryTest.java:
##########
@@ -268,13 +290,96 @@
Assertions.assertTrue(factory.start());
Assertions.assertTrue(factory.close());
Assertions.assertTrue(factory.getFuture().isDone());
- EasyMock.verify(cacheManager);
+ verifyCacheManagerAfterExecutorTerminates(factory);
+ }
+
+ @Test
+ public void testStartWaitsForInitialEndOffsets() throws Exception
+ {
+ final MockConsumer<String, String> kafkaConsumer = new
MockConsumer<>(OffsetResetStrategy.EARLIEST);
+ final TopicPartition topicPartition = new TopicPartition(TOPIC, 0);
+ final CountDownLatch firstPollComplete = new CountDownLatch(1);
+ final CountDownLatch allowCatchUp = new CountDownLatch(1);
+
+ kafkaConsumer.schedulePollTask(() -> {
+ kafkaConsumer.updateBeginningOffsets(ImmutableMap.of(topicPartition,
0L));
+ kafkaConsumer.updateEndOffsets(ImmutableMap.of(topicPartition, 2L));
+ kafkaConsumer.rebalance(Collections.singletonList(topicPartition));
+ kafkaConsumer.addRecord(new ConsumerRecord<>(TOPIC, 0, 0L, "key-0",
"value-0"));
+ firstPollComplete.countDown();
+ });
+ kafkaConsumer.schedulePollTask(() -> {
+ try {
+ if (!allowCatchUp.await(10, TimeUnit.SECONDS)) {
+ throw new RuntimeException("Timed out waiting to finish the startup
catch-up");
+ }
+ }
+ catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ kafkaConsumer.addRecord(new ConsumerRecord<>(TOPIC, 0, 1L, "key-1",
"value-1"));
+ });
+
+ EasyMock.replay(cacheManager);
+ final KafkaLookupExtractorFactory factory = new
KafkaLookupExtractorFactory(
+ cacheManager,
+ TOPIC,
+ ImmutableMap.of("bootstrap.servers", "localhost"),
+ 10_000L,
+ false
+ )
+ {
+ @Override
+ Consumer<String, String> getConsumer()
+ {
+ return kafkaConsumer;
+ }
+ };
+ final ExecutorService startExecutor =
Execs.singleThreaded("kafka-lookup-start-test");
+ final Future<Boolean> startFuture = startExecutor.submit(factory::start);
+
+ try {
+ Assertions.assertTrue(firstPollComplete.await(10, TimeUnit.SECONDS));
+ Assertions.assertThrows(
+ TimeoutException.class,
+ () -> startFuture.get(100, TimeUnit.MILLISECONDS),
+ "start returned before the consumer reached its initial end offsets"
+ );
+ allowCatchUp.countDown();
+ Assertions.assertTrue(startFuture.get(10, TimeUnit.SECONDS));
+ Assertions.assertEquals("value-0", factory.get().apply("key-0"));
+ Assertions.assertEquals("value-1", factory.get().apply("key-1"));
+ }
+ finally {
+ allowCatchUp.countDown();
+ factory.close();
+ startExecutor.shutdownNow();
+ }
+ verifyCacheManagerAfterExecutorTerminates(factory);
}
@Test
- public void testStartFailsFromTimeout()
+ public void testStartTimeoutReturnsBeforeConsumerStops() throws Exception
{
+ final CountDownLatch pollStarted = new CountDownLatch(1);
+ final CountDownLatch allowPollToFinish = new CountDownLatch(1);
+ final CountDownLatch consumerClosed = new CountDownLatch(1);
+ final MockConsumer<String, String> kafkaConsumer = new
MockConsumer<>(OffsetResetStrategy.EARLIEST)
Review Comment:
## CodeQL / Deprecated method or constructor invocation
Invoking [MockConsumer.MockConsumer](1) should be avoided because it has
been deprecated.
[Show more
details](https://github.com/apache/druid/security/code-scanning/12003)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]