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 1b44ab1c54d KAFKA-20846:Make the streams assignor GroupSpec deeply 
unmodifiable (#22970)
1b44ab1c54d is described below

commit 1b44ab1c54d9e2cf41f513f5020c916c34192eec
Author: gabriellafu <[email protected]>
AuthorDate: Sat Aug 1 18:55:52 2026 -0400

    KAFKA-20846:Make the streams assignor GroupSpec deeply unmodifiable (#22970)
    
    With KIP-1357, the "streams" task assignor is now pluggable and could be
    a custom assignor. Therefore, we need to protect the input parameters of
    the `assign()` call: we need to avoid that the assignor modifies GC
    internal data structures, and also need to ensure that the GC does not
    modify what it passed into the assignor accidentally.
    
    This PR updates GroupSpecImpl and MemberMetadataAndStateImpl to use
    proper deep-copies and unmodifiable views to protect against any
    undesired modifications.
    
    Reviewers: Matthias J. Sax <[email protected]>
---
 .../group/api/streams/assignor/GroupSpec.java      |  2 +-
 .../streams/assignor/MemberAssignmentMetadata.java |  2 +-
 .../streams/assignor/MemberAssignmentState.java    |  2 +
 .../group/streams/TargetAssignmentBuilder.java     |  3 +-
 .../group/streams/assignor/GroupSpecImpl.java      |  6 +-
 .../assignor/MemberMetadataAndStateImpl.java       | 26 ++++++--
 .../group/streams/assignor/MockAssignor.java       |  6 +-
 .../group/streams/assignor/GroupSpecImplTest.java  |  6 ++
 .../assignor/MemberMetadataAndStateImplTest.java   | 76 ++++++++++++++++++++++
 9 files changed, 115 insertions(+), 14 deletions(-)

diff --git 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java
 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java
index 486b2f8e179..35aeb291bdf 100644
--- 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java
+++ 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/GroupSpec.java
@@ -51,7 +51,7 @@ public interface GroupSpec {
     MemberAssignmentState memberAssignmentState(String memberId);
 
     /**
-     * @return Any configurations passed to the assignor.
+     * @return Any configurations passed to the assignor. The map is 
unmodifiable.
      */
     Map<String, String> configs();
 
diff --git 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java
 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java
index ef725075d6d..a13e0620e60 100644
--- 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java
+++ 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentMetadata.java
@@ -49,7 +49,7 @@ public interface MemberAssignmentMetadata {
     String processId();
 
     /**
-     * @return The client tags for a rack-aware assignment.
+     * @return The client tags for a rack-aware assignment. The map is 
unmodifiable.
      */
     Map<String, String> clientTags();
 
diff --git 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java
 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java
index df0154d994e..c136ad6e5cc 100644
--- 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java
+++ 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/MemberAssignmentState.java
@@ -25,6 +25,8 @@ import java.util.Set;
 /**
  * The active, standby, and warm-up tasks that a streams group member 
currently has,
  * used by the {@link TaskAssignor} to compute a new target assignment.
+ *
+ * <p>All maps returned by this interface, including their nested collections, 
are unmodifiable.
  */
 @InterfaceAudience.Public
 @InterfaceStability.Evolving
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
index 22cd1745808..124c392ba42 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TargetAssignmentBuilder.java
@@ -28,7 +28,6 @@ import 
org.apache.kafka.coordinator.group.streams.assignor.MemberMetadataAndStat
 import org.apache.kafka.coordinator.group.streams.topics.ConfiguredTopology;
 
 import java.util.ArrayList;
-import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -241,7 +240,7 @@ public class TargetAssignmentBuilder {
             }
             newGroupAssignment = assignor.assign(
                 new GroupSpecImpl(
-                    Collections.unmodifiableMap(memberMetadataMap),
+                    memberMetadataMap,
                     assignmentConfigs
                 ),
                 new TopologyMetadata(metadataImage, 
topology.subtopologies().get())
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java
index bdd07eeba61..19b4d08556b 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImpl.java
@@ -21,6 +21,7 @@ import 
org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentM
 import 
org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentState;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.Map;
 import java.util.Objects;
 
@@ -37,8 +38,9 @@ public record GroupSpecImpl(
 ) implements GroupSpec {
 
     public GroupSpecImpl {
-        Objects.requireNonNull(members);
-        Objects.requireNonNull(configs);
+        // Both maps are exposed to a custom assignor through the public 
GroupSpec interface.
+        members = Collections.unmodifiableMap(Objects.requireNonNull(members));
+        configs = Collections.unmodifiableMap(Objects.requireNonNull(configs));
     }
 
     @Override
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImpl.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImpl.java
index ec04f4804f0..d567153aabb 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImpl.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImpl.java
@@ -19,11 +19,11 @@ package org.apache.kafka.coordinator.group.streams.assignor;
 import 
org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentMetadata;
 import 
org.apache.kafka.coordinator.group.api.streams.assignor.MemberAssignmentState;
 
-import java.util.Collections;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
+import java.util.stream.Collectors;
 
 /**
  * Implementation of both the {@link MemberAssignmentMetadata} and the {@link 
MemberAssignmentState}
@@ -55,12 +55,24 @@ public record MemberMetadataAndStateImpl(
         Objects.requireNonNull(instanceId);
         Objects.requireNonNull(rackId);
         Objects.requireNonNull(processId);
-        clientTags = 
Collections.unmodifiableMap(Objects.requireNonNull(clientTags));
-        activeTasks = 
Collections.unmodifiableMap(Objects.requireNonNull(activeTasks));
-        standbyTasks = 
Collections.unmodifiableMap(Objects.requireNonNull(standbyTasks));
-        warmupTasks = 
Collections.unmodifiableMap(Objects.requireNonNull(warmupTasks));
-        taskOffsets = 
Collections.unmodifiableMap(Objects.requireNonNull(taskOffsets));
-        taskEndOffsets = 
Collections.unmodifiableMap(Objects.requireNonNull(taskEndOffsets));
+        clientTags = Map.copyOf(Objects.requireNonNull(clientTags));
+        // These collections belong to the coordinator (the member's current 
assignment and the offsets reported in
+        // the last heartbeat), so deep-copy them: what the assignor sees is a 
snapshot it cannot reach through.
+        activeTasks = deepCopyTasks(activeTasks);
+        standbyTasks = deepCopyTasks(standbyTasks);
+        warmupTasks = deepCopyTasks(warmupTasks);
+        taskOffsets = deepCopyOffsets(taskOffsets);
+        taskEndOffsets = deepCopyOffsets(taskEndOffsets);
+    }
+
+    private static Map<String, Set<Integer>> deepCopyTasks(Map<String, 
Set<Integer>> tasks) {
+        return Objects.requireNonNull(tasks).entrySet().stream()
+            .collect(Collectors.toUnmodifiableMap(Map.Entry::getKey, entry -> 
Set.copyOf(entry.getValue())));
+    }
+
+    private static Map<String, Map<Integer, Long>> deepCopyOffsets(Map<String, 
Map<Integer, Long>> offsets) {
+        return Objects.requireNonNull(offsets).entrySet().stream()
+            .collect(Collectors.toUnmodifiableMap(Map.Entry::getKey, entry -> 
Map.copyOf(entry.getValue())));
     }
 
 }
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java
index 2f81b2bd547..fc621921127 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/MockAssignor.java
@@ -64,7 +64,11 @@ public class MockAssignor implements TaskAssignor {
 
         // Copy existing assignment and fill temporary data structures
         for (final String memberId : groupSpec.memberIds()) {
-            Map<String, Set<Integer>> activeTasks = new 
HashMap<>(groupSpec.memberAssignmentState(memberId).activeTasks());
+            // Deep-copy: the partition sets are grown below when assigning 
unassigned tasks, and the ones owned by
+            // the group spec are unmodifiable.
+            Map<String, Set<Integer>> activeTasks = new HashMap<>();
+            
groupSpec.memberAssignmentState(memberId).activeTasks().forEach((subtopologyId, 
partitions) ->
+                activeTasks.put(subtopologyId, new HashSet<>(partitions)));
 
             newTargetAssignment.put(memberId, new 
MemberAssignment(activeTasks, new HashMap<>()));
             for (Map.Entry<String, Set<Integer>> entry : 
activeTasks.entrySet()) {
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java
index 4d30558192b..976289b5e8e 100644
--- 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/GroupSpecImplTest.java
@@ -83,4 +83,10 @@ public class GroupSpecImplTest {
         assertTrue(groupSpec.configs().isEmpty());
     }
 
+    @Test
+    void testMembersAndConfigsAreUnmodifiable() {
+        assertThrows(UnsupportedOperationException.class, () -> 
groupSpec.members().put("other-member", member));
+        assertThrows(UnsupportedOperationException.class, () -> 
groupSpec.configs().put("key", "value"));
+    }
+
 }
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImplTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImplTest.java
new file mode 100644
index 00000000000..f7bc9124c21
--- /dev/null
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/MemberMetadataAndStateImplTest.java
@@ -0,0 +1,76 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.coordinator.group.streams.assignor;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+public class MemberMetadataAndStateImplTest {
+
+    private static MemberMetadataAndStateImpl memberWith(
+        Map<String, Set<Integer>> tasks,
+        Map<String, Map<Integer, Long>> offsets
+    ) {
+        return new MemberMetadataAndStateImpl(
+            Optional.of("test-instance"),
+            Optional.of("test-rack"),
+            "test-process",
+            new HashMap<>(Map.of("tag1", "value1")),
+            tasks,
+            tasks,
+            tasks,
+            offsets,
+            offsets
+        );
+    }
+
+    @Test
+    void testCollectionsAreDeeplyUnmodifiable() {
+        MemberMetadataAndStateImpl member = memberWith(
+            new HashMap<>(Map.of("subtopology-1", new HashSet<>(Set.of(0, 
1)))),
+            new HashMap<>()
+        );
+
+        assertThrows(UnsupportedOperationException.class, () -> 
member.clientTags().put("tag2", "value2"));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.activeTasks().put("subtopology-2", Set.of(0)));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.activeTasks().get("subtopology-1").add(2));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.standbyTasks().put("subtopology-2", Set.of(0)));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.standbyTasks().get("subtopology-1").add(2));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.warmupTasks().put("subtopology-2", Set.of(0)));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.warmupTasks().get("subtopology-1").clear());
+    }
+
+    @Test
+    void testTaskOffsetsAreDeeplyUnmodifiable() {
+        MemberMetadataAndStateImpl member = memberWith(
+            new HashMap<>(),
+            new HashMap<>(Map.of("subtopology-1", new HashMap<>(Map.of(0, 
100L))))
+        );
+
+        assertThrows(UnsupportedOperationException.class, () -> 
member.taskOffsets().put("subtopology-2", Map.of()));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.taskOffsets().get("subtopology-1").put(0, 200L));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.taskEndOffsets().put("subtopology-2", Map.of()));
+        assertThrows(UnsupportedOperationException.class, () -> 
member.taskEndOffsets().get("subtopology-1").put(1, 200L));
+    }
+}

Reply via email to