aho135 commented on code in PR #20387:
URL: https://github.com/apache/druid/pull/20387#discussion_r4067124642


##########
extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisorTest.java:
##########
@@ -6554,6 +6554,85 @@ public void 
testBoundedStreamConfig_allPartitionsEmptyRange_completesImmediately
     );
   }
 
+  @Test
+  public void testBoundedStreamConfig_completedGroup_doesNotTopUpReplicas() 
throws Exception
+  {
+    // A bounded task group whose offsets have already reached the configured 
end offsets must not have
+    // replacement replicas created for it. Otherwise, as completed replicas 
exit, the supervisor keeps
+    // spawning replicas that start at the end offset and complete instantly, 
churning tasks forever.
+    Map<String, Object> startOffsets = ImmutableMap.of("0", 0, "1", 0, "2", 0);
+    Map<String, Object> endOffsets = ImmutableMap.of("0", 100, "1", 100, "2", 
100);
+    BoundedStreamConfig boundedConfig = new BoundedStreamConfig(startOffsets, 
endOffsets);
+
+    supervisor = getTestableSupervisorWithBoundedConfig(2, 1, "PT1H", 
boundedConfig);
+
+    addSomeEvents(100);
+
+    // One replica is still around, so the group has fewer tasks (1) than the 
configured replica count (2),
+    // which would normally trigger a replica top-up.
+    Task task = createKafkaIndexTask(
+        "id1",
+        DATASOURCE,
+        0,
+        new SeekableStreamStartSequenceNumbers<>(
+            "topic",
+            singlePartitionMap(topic, 0, 0L, 1, 0L, 2, 0L),
+            ImmutableSet.of()
+        ),
+        new SeekableStreamEndSequenceNumbers<>(
+            "topic",
+            singlePartitionMap(topic, 0, 100L, 1, 100L, 2, 100L)
+        ),
+        null,
+        null,
+        supervisor.getTuningConfig()
+    );
+
+    Map<String, Task> existingTasks = Map.of(task.getId(), task);
+
+    // Metadata store reports offsets that have reached the bounded end, with 
a matching bounded config, so
+    // hasTaskGroupReachedBoundedEnd() returns true for the group.
+    KafkaDataSourceMetadata reachedEndMetadata = new KafkaDataSourceMetadata(
+        new SeekableStreamEndSequenceNumbers<>(topic, 
singlePartitionMap(topic, 0, 100L, 1, 100L, 2, 100L)),
+        boundedConfig
+    );
+
+    
EasyMock.expect(taskMaster.getTaskQueue()).andReturn(Optional.of(taskQueue)).anyTimes();
+    
EasyMock.expect(taskMaster.getTaskRunner()).andReturn(Optional.of(taskRunner)).anyTimes();
+    
EasyMock.expect(taskRunner.getRunningTasks()).andReturn(Collections.emptyList()).anyTimes();
+    
EasyMock.expect(taskQueue.getActiveTasksForDatasource(DATASOURCE)).andReturn(existingTasks).anyTimes();
+    
EasyMock.expect(taskStorage.getActiveTasksByDatasource(DATASOURCE)).andReturn(ImmutableList.of(task)).anyTimes();
+    
EasyMock.expect(taskStorage.getStatus("id1")).andReturn(Optional.of(TaskStatus.running("id1"))).anyTimes();
+    
EasyMock.expect(taskStorage.getTask("id1")).andReturn(Optional.of(task)).anyTimes();
+    EasyMock.expect(taskClient.getStatusAsync(EasyMock.anyString()))
+            .andReturn(Futures.immediateFuture(Status.READING))
+            .anyTimes();
+    EasyMock.expect(taskClient.getStartTimeAsync(EasyMock.anyString()))
+            .andReturn(Futures.immediateFuture(DateTimes.nowUtc()))
+            .anyTimes();
+    
EasyMock.expect(indexerMetadataStorageCoordinator.retrieveDataSourceMetadata(DATASOURCE))
+            .andReturn(reachedEndMetadata)
+            .anyTimes();
+
+    // Checkpoints consistent with the metadata end offsets so the existing 
task is not killed as inconsistent.
+    TreeMap<Integer, Map<KafkaTopicPartition, Long>> checkpoints = new 
TreeMap<>();
+    checkpoints.put(0, singlePartitionMap(topic, 0, 100L, 1, 100L, 2, 100L));
+    EasyMock.expect(taskClient.getCheckpointsAsync(EasyMock.contains("id1"), 
EasyMock.anyBoolean()))
+            .andReturn(Futures.immediateFuture(checkpoints))
+            .anyTimes();
+
+    taskRunner.registerListener(EasyMock.anyObject(TaskRunnerListener.class), 
EasyMock.anyObject(Executor.class));
+    EasyMock.expectLastCall().anyTimes();
+
+    // Intentionally NO expectation for taskQueue.add(...): the guard must 
prevent any replica top-up for the
+    // completed bounded group. The strict mock fails the test if add() is 
invoked.
+
+    replayAll();
+    supervisor.start();
+    supervisor.runInternal();
+    verifyAll();

Review Comment:
   You're right — that test never reached the top-up loop, so it passed 
vacuously. I've replaced it with 
`SeekableStreamSupervisorStateTest#testCreateNewTasks_boundedGroupReachedEnd_doesNotTopUpReplicas`
 (aa68cff16d), which flows through `runInternal()`: it sets up an 
actively-reading bounded task group that has reached its end offsets with fewer 
tasks than the replica count, and asserts no replacement replica is submitted. 
I verified it fails without the guard (a top-up task gets submitted). Bonus: 
since it lives in the `indexing-service` module where the change is, it also 
covers the new code path for the diff-coverage check — so the earlier coverage 
note is now moot.



-- 
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