github-advanced-security[bot] commented on code in PR #20387:
URL: https://github.com/apache/druid/pull/20387#discussion_r4067261503
##########
indexing-service/src/test/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisorStateTest.java:
##########
@@ -2961,6 +2961,115 @@
verifyAll();
}
+ /**
+ * In bounded mode, a task group that has already reached its end offsets
must not have replacement replicas
+ * created for it. Without the guard in {@code createNewTasks()}'s replica
top-up loop, {@code replicas > tasks}
+ * would submit a replacement replica whose start offset is already at the
bounded end, so it completes instantly
+ * and re-triggers the top-up, churning tasks endlessly. This test flows
through {@code runInternal()} and asserts
+ * that no top-up task is submitted; it fails if the guard is removed.
+ */
+ @Test
+ public void testCreateNewTasks_boundedGroupReachedEnd_doesNotTopUpReplicas()
+ {
+ // replicas = 2, taskCount = 1, bounded with an empty range (start == end)
so the group has reached its end.
+ final BoundedStreamConfig boundedConfig = new BoundedStreamConfig(
+ ImmutableMap.of("0", "100"),
+ ImmutableMap.of("0", "100")
+ );
+ final SeekableStreamSupervisorIOConfig ioConfig =
+ new SupervisorIOConfigBuilder.DefaultSupervisorIOConfigBuilder()
+ .withStream(STREAM)
+ .withInputFormat(new JsonInputFormat(new JSONPathSpec(true,
List.of()), Map.of(), false, false, false))
+ .withReplicas(2)
+ .withTaskCount(1)
+ .withTaskDuration(new Period("PT1H"))
+ .withStartDelay(new Period("P1D"))
+ .withSupervisorRunPeriod(new Period("PT30S"))
+ .withUseEarliestSequenceNumber(false)
+ .withCompletionTimeout(new Period("PT30M"))
+ .withLagAggregator(LagAggregator.DEFAULT)
+ .withBoundedStreamConfig(boundedConfig)
+ .build();
+
+ // A single already-running task for group 0, so the group is discovered
with one task, i.e. fewer than the
+ // configured replica count (2). That would normally trigger a replica
top-up.
+ // minMsgTime/maxMsgTime must match the injected task group's (null) so
isTaskCurrent() keeps the task.
+ final SeekableStreamIndexTaskIOConfig taskIoConfig = createTaskIoConfigExt(
+ 0,
+ Map.of("0", "0"),
+ Map.of("0", "100"),
+ "test",
+ null,
+ null,
+ Set.of(),
+ ioConfig
+ );
+ final TestSeekableStreamIndexTask task1 = createTestTask("task1", "0",
null, taskIoConfig, recordSupplier);
+
+ EasyMock.reset(spec);
+ EasyMock.expect(spec.getId()).andReturn(SUPERVISOR_ID).anyTimes();
+
EasyMock.expect(spec.getSupervisorStateManagerConfig()).andReturn(supervisorConfig).anyTimes();
+
EasyMock.expect(spec.getDataSchema()).andReturn(getDataSchema()).anyTimes();
Review Comment:
## CodeQL / Deprecated method or constructor invocation
Invoking [SeekableStreamSupervisorSpec.getDataSchema](1) should be avoided
because it has been deprecated.
[Show more
details](https://github.com/apache/druid/security/code-scanning/10848)
--
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]