This is an automated email from the ASF dual-hosted git repository.

mjsax pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new d9b215da915 MINOR: handle invalid StreamsHeartbeatRequest (#22728)
d9b215da915 is described below

commit d9b215da9156a5ce080cea30c0eca73851c9d5e8
Author: Matthias J. Sax <[email protected]>
AuthorDate: Thu Jul 2 23:32:33 2026 -0700

    MINOR: handle invalid StreamsHeartbeatRequest (#22728)
    
    A StreamsHeartbeatRequest must not contain duplicate task-offset
    entries, and should be rejected as invalid request. We also need to
    avoid to mutate in-memory state for an invalid request.
    
    Reviewers: David Jacot <[email protected]>
---
 .../coordinator/group/GroupMetadataManager.java    |  20 ++--
 .../group/streams/MemberTaskOffsets.java           |  19 +++-
 .../group/GroupMetadataManagerTest.java            | 106 +++++++++++++++++++++
 .../group/streams/MemberTaskOffsetsTest.java       |  19 ++++
 4 files changed, 151 insertions(+), 13 deletions(-)

diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index 9643099e088..38fee045c3d 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -2116,14 +2116,6 @@ public class GroupMetadataManager {
             );
         }
 
-        // Store the latest task changelog offsets/end-offsets reported by the 
member. These are transient telemetry
-        // (used by the assignor to estimate task lag) and are not persisted. 
Task offsets and end-offsets are reported
-        // independently: a null list means "unchanged since the last 
heartbeat", so we retain the previously reported
-        // value for whichever of the two is null and only update when at 
least one is reported.
-        if (taskOffsets != null || taskEndOffsets != null) {
-            group.updateTaskOffsets(memberId, 
group.taskOffsets(memberId).update(taskOffsets, taskEndOffsets));
-        }
-
         // 1. Create or update the member.
         StreamsGroupMember.Builder updatedMemberBuilder = new 
StreamsGroupMember.Builder(member)
             .maybeUpdateInstanceId(Optional.ofNullable(instanceId))
@@ -2194,6 +2186,18 @@ public class GroupMetadataManager {
             throwIfRequestContainsInvalidTasks(subtopologySortedMap, 
ownedStandbyTasks);
             throwIfRequestContainsInvalidTasks(subtopologySortedMap, 
ownedWarmupTasks);
         }
+
+        // Store the latest task changelog offsets/end-offsets reported by the 
member. These are transient telemetry
+        // (used by the assignor to estimate task lag) and are not persisted. 
Task offsets and end-offsets are reported
+        // independently: a null list means "unchanged since the last 
heartbeat", so we retain the previously reported
+        // value for whichever of the two is null and only update when at 
least one is reported.
+        // This must run after the task validation above: it mutates a 
non-timeline map that the coordinator runtime
+        // does not roll back, so updating it before validation could leave an 
orphaned entry for a member whose
+        // (joining) heartbeat is then rejected.
+        if (taskOffsets != null || taskEndOffsets != null) {
+            group.updateTaskOffsets(memberId, 
group.taskOffsets(memberId).update(taskOffsets, taskEndOffsets));
+        }
+
         // We validated a topology that was not validated before, so bump the 
group epoch as we may have to reassign tasks.
         if (validatedTopologyEpoch != group.validatedTopologyEpoch()) {
             bumpGroupEpoch = true;
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
index 40c863b4b20..ed2434a79bf 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
@@ -16,12 +16,13 @@
  */
 package org.apache.kafka.coordinator.group.streams;
 
+import org.apache.kafka.common.errors.InvalidRequestException;
 import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData;
 import org.apache.kafka.coordinator.group.streams.assignor.TaskId;
 
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
-import java.util.stream.Collectors;
 
 /**
  * The latest per-task cumulative changelog offsets and end-offsets that a 
member reported in its heartbeat.
@@ -57,9 +58,17 @@ public record MemberTaskOffsets(Map<TaskId, Long> 
taskOffsets, Map<TaskId, Long>
     }
 
     private static Map<TaskId, Long> toTaskIdMap(final 
List<StreamsGroupHeartbeatRequestData.TaskOffset> taskOffsets) {
-        return taskOffsets.stream().collect(Collectors.toMap(
-            taskOffset -> new TaskId(taskOffset.subtopologyId(), 
taskOffset.partition()),
-            StreamsGroupHeartbeatRequestData.TaskOffset::offset
-        ));
+        final Map<TaskId, Long> result = new HashMap<>(taskOffsets.size());
+        for (final StreamsGroupHeartbeatRequestData.TaskOffset taskOffset : 
taskOffsets) {
+            final TaskId taskId = new TaskId(taskOffset.subtopologyId(), 
taskOffset.partition());
+            // The reported values come straight from the client heartbeat and 
the protocol does not enforce uniqueness
+            // of (subtopologyId, partition). Reject a duplicate with a clear 
client error rather than silently picking
+            // one of the values.
+            if (result.putIfAbsent(taskId, taskOffset.offset()) != null) {
+                throw new InvalidRequestException("Task offsets contain a 
duplicate entry for subtopology "
+                    + taskOffset.subtopologyId() + " and partition " + 
taskOffset.partition() + ".");
+            }
+        }
+        return result;
     }
 }
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index 9f348600f4c..f3568dc089b 100644
--- 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -134,6 +134,7 @@ import 
org.apache.kafka.coordinator.group.modern.share.ShareGroup.InitMapValue;
 import org.apache.kafka.coordinator.group.modern.share.ShareGroupBuilder;
 import org.apache.kafka.coordinator.group.modern.share.ShareGroupConfig;
 import org.apache.kafka.coordinator.group.modern.share.ShareGroupMember;
+import org.apache.kafka.coordinator.group.streams.MemberTaskOffsets;
 import org.apache.kafka.coordinator.group.streams.MockTaskAssignor;
 import 
org.apache.kafka.coordinator.group.streams.StreamsCoordinatorRecordHelpers;
 import org.apache.kafka.coordinator.group.streams.StreamsGroup;
@@ -18576,6 +18577,111 @@ public class GroupMetadataManagerTest {
         assertEquals("Task 3 for subtopology subtopology1 is invalid. Number 
of tasks for this subtopology: 3", e2.getMessage());
     }
 
+    @Test
+    public void 
testStreamsGroupHeartbeatDoesNotStoreTaskOffsetsWhenOwnedTasksAreInvalid() {
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String subtopology1 = "subtopology1";
+        String fooTopicName = "foo";
+        Uuid fooTopicId = Uuid.randomUuid();
+        Topology topology = new Topology().setSubtopologies(List.of(
+            new 
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+        ));
+
+        MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(assignor))
+            .withMetadataImage(new MetadataImageBuilder()
+                .addTopic(fooTopicId, fooTopicName, 3)
+                .buildCoordinatorMetadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, 10)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+                    .setMemberEpoch(10)
+                    .setPreviousMemberEpoch(10)
+                    
.setAssignedTasks(mkTasksTupleWithCommonEpoch(TaskRole.ACTIVE, 10,
+                        TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+                    .build())
+                .withTopology(StreamsTopology.fromHeartbeatRequest(topology))
+                .withTargetAssignment(memberId, 
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+                    TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+                .withTargetAssignmentEpoch(10)
+            )
+            .build();
+
+        // The heartbeat reports task offsets but also claims an invalid owned 
task, so it is rejected during
+        // validation. The reported offsets must not be stored: updating the 
(non-timeline) task-offsets map before
+        // validation would leave an orphaned entry that the coordinator 
runtime does not roll back.
+        assertThrows(InvalidRequestException.class, () -> 
context.streamsGroupHeartbeat(
+            new StreamsGroupHeartbeatRequestData()
+                .setGroupId(groupId)
+                .setMemberId(memberId)
+                .setMemberEpoch(10)
+                .setActiveTasks(List.of(new 
StreamsGroupHeartbeatRequestData.TaskIds()
+                    .setSubtopologyId(subtopology1)
+                    .setPartitions(List.of(3)))) // partition 3 is out of 
range (only 0, 1, 2 exist)
+                .setStandbyTasks(List.of())
+                .setWarmupTasks(List.of())
+                .setTaskOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()
+                    
.setSubtopologyId(subtopology1).setPartition(0).setOffset(10L)))
+                .setTaskEndOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()
+                    
.setSubtopologyId(subtopology1).setPartition(0).setOffset(20L)))));
+
+        StreamsGroup group = 
context.groupMetadataManager.streamsGroup(groupId);
+        assertEquals(MemberTaskOffsets.EMPTY, group.taskOffsets(memberId));
+    }
+
+    @Test
+    public void testStreamsGroupHeartbeatRejectsDuplicateTaskOffsets() {
+        String groupId = "fooup";
+        String memberId = Uuid.randomUuid().toString();
+        String subtopology1 = "subtopology1";
+        String fooTopicName = "foo";
+        Uuid fooTopicId = Uuid.randomUuid();
+        Topology topology = new Topology().setSubtopologies(List.of(
+            new 
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+        ));
+
+        MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+        GroupMetadataManagerTestContext context = new 
GroupMetadataManagerTestContext.Builder()
+            .withStreamsGroupTaskAssignors(List.of(assignor))
+            .withMetadataImage(new MetadataImageBuilder()
+                .addTopic(fooTopicId, fooTopicName, 3)
+                .buildCoordinatorMetadataImage())
+            .withStreamsGroup(new StreamsGroupBuilder(groupId, 10)
+                .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+                    .setMemberEpoch(10)
+                    .setPreviousMemberEpoch(10)
+                    
.setAssignedTasks(mkTasksTupleWithCommonEpoch(TaskRole.ACTIVE, 10,
+                        TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+                    .build())
+                .withTopology(StreamsTopology.fromHeartbeatRequest(topology))
+                .withTargetAssignment(memberId, 
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+                    TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+                .withTargetAssignmentEpoch(10)
+            )
+            .build();
+
+        // The owned tasks are valid, but the reported task offsets contain a 
duplicate (subtopologyId, partition)
+        // entry. This is a malformed request and must be rejected with a 
clear client error; nothing is stored.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+            new StreamsGroupHeartbeatRequestData()
+                .setGroupId(groupId)
+                .setMemberId(memberId)
+                .setMemberEpoch(10)
+                .setActiveTasks(List.of(new 
StreamsGroupHeartbeatRequestData.TaskIds()
+                    .setSubtopologyId(subtopology1)
+                    .setPartitions(List.of(0, 1, 2))))
+                .setStandbyTasks(List.of())
+                .setWarmupTasks(List.of())
+                .setTaskOffsets(List.of(
+                    new 
StreamsGroupHeartbeatRequestData.TaskOffset().setSubtopologyId(subtopology1).setPartition(0).setOffset(10L),
+                    new 
StreamsGroupHeartbeatRequestData.TaskOffset().setSubtopologyId(subtopology1).setPartition(0).setOffset(11L)))));
+        assertEquals("Task offsets contain a duplicate entry for subtopology " 
+ subtopology1 + " and partition 0.", e.getMessage());
+
+        StreamsGroup group = 
context.groupMetadataManager.streamsGroup(groupId);
+        assertEquals(MemberTaskOffsets.EMPTY, group.taskOffsets(memberId));
+    }
+
     @Test
     public void testStreamsGroupHeartbeatStoresTaskOffsetsWithoutPersisting() {
         String groupId = "fooup";
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
index 68ff6b0e370..941c5c3cc97 100644
--- 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
@@ -16,6 +16,7 @@
  */
 package org.apache.kafka.coordinator.group.streams;
 
+import org.apache.kafka.common.errors.InvalidRequestException;
 import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData;
 import org.apache.kafka.coordinator.group.streams.assignor.TaskId;
 
@@ -25,6 +26,7 @@ import java.util.List;
 import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 
 public class MemberTaskOffsetsTest {
 
@@ -54,6 +56,23 @@ public class MemberTaskOffsetsTest {
         );
     }
 
+    @Test
+    public void shouldRejectDuplicateTaskEntries() {
+        // The protocol does not enforce uniqueness of (subtopologyId, 
partition) within a reported list. A duplicate
+        // entry is a malformed request and must be rejected with a clear 
client error rather than silently resolved.
+        InvalidRequestException e = 
assertThrows(InvalidRequestException.class, () -> 
MemberTaskOffsets.EMPTY.update(
+            List.of(taskOffset("sub-1", 0, 10L), taskOffset("sub-1", 0, 42L)),
+            null
+        ));
+        assertEquals("Task offsets contain a duplicate entry for subtopology 
sub-1 and partition 0.", e.getMessage());
+
+        // The check also applies to the end-offsets list.
+        assertThrows(InvalidRequestException.class, () -> 
MemberTaskOffsets.EMPTY.update(
+            null,
+            List.of(taskOffset("sub-1", 0, 15L), taskOffset("sub-1", 0, 55L))
+        ));
+    }
+
     @Test
     public void shouldRetainBothMapsWhenBothListsAreNull() {
         MemberTaskOffsets previous = new MemberTaskOffsets(

Reply via email to