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]

Reply via email to