mjsax commented on code in PR #22778:
URL: https://github.com/apache/kafka/pull/22778#discussion_r3686931337


##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +183,68 @@ public void 
shouldGetCorrectHostPartitionInformation(final String groupProtocolC
                     }, TestUtils.DEFAULT_MAX_WAIT_MS,
                             "Kafka Streams clients 1 and 2 never got metadata 
about standby tasks");
 
-                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 2,
+                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 4,

Review Comment:
   Claude says, that this leads to a race condition in the test:
   ```
   The final waitForCondition polls streamsOne's view, but verifyClientMetadata 
asserts on streamsTwo's view:
   
   expected: <4> but was: <2>
       at 
...verifyHostMetadata(IQv2EndpointToPartitionsIntegrationTest.java:233)
       at 
...verifyClientMetadata(IQv2EndpointToPartitionsIntegrationTest.java:208)   // 
port 3030
   
   For the two no-standby parameterizations the preceding standby wait is 0 == 
0, so the only real constraint on streamsTwo's view is metadata.size() == 2 — 
and rebuildMetadataForSingleTopology emits an entry for a host even with an 
empty partition set, so that holds before any task has migrated. The per-thread 
waits are on local task counts, which lead the metadata view. Two threads per 
instance stretch the migration, which is why it surfaces now. Fix: wrap 
verifyClientMetadata in TestUtils.retryOnExceptionWithTimeout, or make the last 
wait check the end state on the view being asserted.
   ```
   
   It suggest to change it too:
   ```
   waitForCondition(
       () -> {
           final List<StreamsMetadata> metadata = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
           return metadata.size() == 2 && metadata.stream().allMatch(m -> 
m.topicPartitions().size() == 4);
       },
       IntegrationTestUtils.DEFAULT_TIMEOUT,
       () -> "Kafka Streams one didn't give up active tasks"
   );
   ```



##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -80,8 +79,8 @@ public void setUp() throws InterruptedException {
         appId = safeUniqueTestName("endpointIntegrationTest");
         inputTopicTwoPartitions = appId + "-input-two";
         outputTopicTwoPartitions = appId + "-output-two";

Review Comment:
   Nit: seems we should update these name to `outputTopicFourPartitions`? Also 
the topic name `...-two` -> `...-four`
   
   Same for input topic -- not sure if there is other similar stuff that needs 
to be updated?



##########
clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java:
##########
@@ -879,8 +879,18 @@ private static Map<StreamsRebalanceData.HostInfo, 
StreamsRebalanceData.EndpointP
             List<TopicPartition> activeTopicPartitions = 
getTopicPartitionList(endpoint.activePartitions());
             List<TopicPartition> standbyTopicPartitions = 
getTopicPartitionList(endpoint.standbyPartitions());
             StreamsGroupHeartbeatResponseData.Endpoint userEndpoint = 
endpoint.userEndpoint();
-            StreamsRebalanceData.EndpointPartitions endpointPartitions = new 
StreamsRebalanceData.EndpointPartitions(activeTopicPartitions, 
standbyTopicPartitions);
-            partitionsByHost.put(new 
StreamsRebalanceData.HostInfo(userEndpoint.host(), userEndpoint.port()), 
endpointPartitions);
+            StreamsRebalanceData.HostInfo hostInfo = new 
StreamsRebalanceData.HostInfo(userEndpoint.host(), userEndpoint.port());
+            partitionsByHost.merge(
+                hostInfo,
+                new 
StreamsRebalanceData.EndpointPartitions(activeTopicPartitions, 
standbyTopicPartitions),
+                (existing, newPartitions) -> {
+                    List<TopicPartition> mergedActive = 
existing.activePartitions();
+                    mergedActive.addAll(newPartitions.activePartitions());

Review Comment:
   Style: if we mutate `mergedActive` here, we technically change what we get 
back from `StreamsRebalanceData.EndpointPartitions#activePartitions()` -- this 
sounds like an anti-pattern as it might mutate an internal field from another 
object. We should rather make a copy of the list first:
   ```
   List<TopicPartition> mergedActive = new 
ArrayList<>(existing.activePartitions());
   ```
   Atm, it's not a problem because `activePartitions()` provides us with an 
copy, but from an API contract POV we should not rely on it -- 
`activePartitions()` could change, breaking the code (especially if it would 
return a immutable list).
   
   Same below for standby case.
   
   With courtesy from Claude :) 



##########
clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java:
##########
@@ -2922,4 +2922,67 @@ private ClientResponse 
buildClientResponseWithTopologyRequired(final boolean top
             )
         );
     }
+
+    @Test
+    public void testPartitionsByUserEndpointMergedForDuplicateUserEndpoints() {
+        try (
+            final MockedConstruction<HeartbeatRequestState> ignored = 
mockConstruction(
+                HeartbeatRequestState.class,
+                (mock, context) -> 
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+            final LogCaptureAppender logAppender = 
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+        ) {
+            
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class, 
Level.WARN);

Review Comment:
   Seems we don't assert anything about WARN logs -- the `logAppender` can be 
removed.



##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +183,68 @@ public void 
shouldGetCorrectHostPartitionInformation(final String groupProtocolC
                     }, TestUtils.DEFAULT_MAX_WAIT_MS,
                             "Kafka Streams clients 1 and 2 never got metadata 
about standby tasks");
 
-                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 2,
+                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 4,
                             IntegrationTestUtils.DEFAULT_TIMEOUT,
                             () -> "Kafka Streams one didn't give up active 
tasks");
 
-                    final List<StreamsMetadata> allClientMetadataUpdated = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
-
-                    final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
-                    final Set<TopicPartition> streamsOneActiveTopicPartitions 
= streamsOneMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsOneStandbyTopicPartitions 
= streamsOneMetadata.standbyTopicPartitions();
-                    final Set<String> streamsOneStoreNames = 
streamsOneMetadata.stateStoreNames();
-                    final Set<String> streamsOneStandbyStoreNames = 
streamsOneMetadata.standbyStateStoreNames();
-
-                    assertEquals(2020, streamsOneMetadata.hostInfo().port());
-                    assertEquals(2, streamsOneActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyTopicPartitions.size());
-                    assertEquals(1, streamsOneStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsOneStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsOneStandbyStoreNames.iterator().next());
-                    }
-
-                    final long streamsOneRepartitionTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsOneSourceTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsOneRepartitionTopicCount);
-                    assertEquals(1, streamsOneSourceTopicCount);
-
-                    final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
-                    final Set<TopicPartition> streamsTwoActiveTopicPartitions 
= streamsTwoMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsTwoStandbyTopicPartitions 
= streamsTwoMetadata.standbyTopicPartitions();
-                    final Set<String> streamsTwoStateStoreNames = 
streamsTwoMetadata.stateStoreNames();
-                    final Set<String> streamsTwoStandbyStateStoreNames = 
streamsTwoMetadata.standbyStateStoreNames();
-
-                    assertEquals(3030, streamsTwoMetadata.hostInfo().port());
-                    assertEquals(2, streamsTwoActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyTopicPartitions.size());
-                    assertEquals(1, streamsTwoStateStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyStateStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStateStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStandbyStateStoreNames.iterator().next());
-                    }
-
-                    final long streamsTwoRepartitionTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsTwoSourceTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsTwoRepartitionTopicCount);
-                    assertEquals(1, streamsTwoSourceTopicCount);
-
-                    if (usingStandbyReplicas) {
-                        final TopicPartition streamsOneStandbyTopicPartition = 
streamsOneStandbyTopicPartitions.iterator().next();
-                        final TopicPartition streamsTwoStandbyTopicPartition = 
streamsTwoStandbyTopicPartitions.iterator().next();
-                        final String streamsOneStandbyTopicName = 
streamsOneStandbyTopicPartition.topic();
-                        final String streamsTwoStandbyTopicName = 
streamsTwoStandbyTopicPartition.topic();
-                        assertEquals(streamsOneStandbyTopicName, 
streamsTwoStandbyTopicName);
-                        
assertNotEquals(streamsOneStandbyTopicPartition.partition(), 
streamsTwoStandbyTopicPartition.partition());
-                    }
+                    verifyClientMetadata(usingStandbyReplicas, new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients()), expectedStandbyCount, 
expectedStandbyStoreCount);
                 }
             }
         } finally {
             closeCluster();
         }
     }
 
+    private static void verifyClientMetadata(
+            final boolean usingStandbyReplicas,
+            final List<StreamsMetadata> allClientMetadataUpdated,
+            final int expectedStandbyCount,
+            final int expectedStandbyStoreCount
+    ) {
+        final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
+        final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
+
+        verifyHostMetadata(streamsOneMetadata, 2020, expectedStandbyCount, 
expectedStandbyStoreCount, usingStandbyReplicas);
+        verifyHostMetadata(streamsTwoMetadata, 3030, expectedStandbyCount, 
expectedStandbyStoreCount, usingStandbyReplicas);
+
+        if (usingStandbyReplicas) {
+            final Set<TopicPartition> streamsOneActiveRepartition = 
streamsOneMetadata.topicPartitions().stream()
+                .filter(tp -> 
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+            final Set<TopicPartition> streamsTwoActiveRepartition = 
streamsTwoMetadata.topicPartitions().stream()
+                .filter(tp -> 
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+            assertEquals(streamsTwoActiveRepartition, 
streamsOneMetadata.standbyTopicPartitions());
+            assertEquals(streamsOneActiveRepartition, 
streamsTwoMetadata.standbyTopicPartitions());
+        }
+    }
+
+    private static void verifyHostMetadata(
+            final StreamsMetadata metadata,
+            final int expectedPort,
+            final int expectedStandbyCount,
+            final int expectedStandbyStoreCount,
+            final boolean usingStandbyReplicas
+    ) {
+        final Set<TopicPartition> activeTopicPartitions = 
metadata.topicPartitions();
+        final Set<TopicPartition> standbyTopicPartitions = 
metadata.standbyTopicPartitions();
+        final Set<String> storeNames = metadata.stateStoreNames();
+        final Set<String> standbyStoreNames = 
metadata.standbyStateStoreNames();
+
+        assertEquals(expectedPort, metadata.hostInfo().port());
+        assertEquals(4, activeTopicPartitions.size());
+        assertEquals(expectedStandbyCount, standbyTopicPartitions.size());
+        assertEquals(1, storeNames.size());
+        assertEquals(expectedStandbyStoreCount, standbyStoreNames.size());
+        assertEquals(EXPECTED_STORE_NAME, storeNames.iterator().next());
+        if (usingStandbyReplicas) {
+            assertEquals(EXPECTED_STORE_NAME, 
standbyStoreNames.iterator().next());
+        }

Review Comment:
   ```suggestion
           assertEquals(
               usingStandbyReplicas ? Set.of(EXPECTED_STORE_NAME) : Set.of(),
               metadata.stateStoreNames()
           );
   ```



##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +183,68 @@ public void 
shouldGetCorrectHostPartitionInformation(final String groupProtocolC
                     }, TestUtils.DEFAULT_MAX_WAIT_MS,
                             "Kafka Streams clients 1 and 2 never got metadata 
about standby tasks");
 
-                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 2,
+                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 4,
                             IntegrationTestUtils.DEFAULT_TIMEOUT,
                             () -> "Kafka Streams one didn't give up active 
tasks");
 
-                    final List<StreamsMetadata> allClientMetadataUpdated = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
-
-                    final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
-                    final Set<TopicPartition> streamsOneActiveTopicPartitions 
= streamsOneMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsOneStandbyTopicPartitions 
= streamsOneMetadata.standbyTopicPartitions();
-                    final Set<String> streamsOneStoreNames = 
streamsOneMetadata.stateStoreNames();
-                    final Set<String> streamsOneStandbyStoreNames = 
streamsOneMetadata.standbyStateStoreNames();
-
-                    assertEquals(2020, streamsOneMetadata.hostInfo().port());
-                    assertEquals(2, streamsOneActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyTopicPartitions.size());
-                    assertEquals(1, streamsOneStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsOneStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsOneStandbyStoreNames.iterator().next());
-                    }
-
-                    final long streamsOneRepartitionTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsOneSourceTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsOneRepartitionTopicCount);
-                    assertEquals(1, streamsOneSourceTopicCount);
-
-                    final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
-                    final Set<TopicPartition> streamsTwoActiveTopicPartitions 
= streamsTwoMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsTwoStandbyTopicPartitions 
= streamsTwoMetadata.standbyTopicPartitions();
-                    final Set<String> streamsTwoStateStoreNames = 
streamsTwoMetadata.stateStoreNames();
-                    final Set<String> streamsTwoStandbyStateStoreNames = 
streamsTwoMetadata.standbyStateStoreNames();
-
-                    assertEquals(3030, streamsTwoMetadata.hostInfo().port());
-                    assertEquals(2, streamsTwoActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyTopicPartitions.size());
-                    assertEquals(1, streamsTwoStateStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyStateStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStateStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStandbyStateStoreNames.iterator().next());
-                    }
-
-                    final long streamsTwoRepartitionTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsTwoSourceTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsTwoRepartitionTopicCount);
-                    assertEquals(1, streamsTwoSourceTopicCount);
-
-                    if (usingStandbyReplicas) {
-                        final TopicPartition streamsOneStandbyTopicPartition = 
streamsOneStandbyTopicPartitions.iterator().next();
-                        final TopicPartition streamsTwoStandbyTopicPartition = 
streamsTwoStandbyTopicPartitions.iterator().next();
-                        final String streamsOneStandbyTopicName = 
streamsOneStandbyTopicPartition.topic();
-                        final String streamsTwoStandbyTopicName = 
streamsTwoStandbyTopicPartition.topic();
-                        assertEquals(streamsOneStandbyTopicName, 
streamsTwoStandbyTopicName);
-                        
assertNotEquals(streamsOneStandbyTopicPartition.partition(), 
streamsTwoStandbyTopicPartition.partition());
-                    }
+                    verifyClientMetadata(usingStandbyReplicas, new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients()), expectedStandbyCount, 
expectedStandbyStoreCount);
                 }
             }
         } finally {
             closeCluster();
         }
     }
 
+    private static void verifyClientMetadata(
+            final boolean usingStandbyReplicas,
+            final List<StreamsMetadata> allClientMetadataUpdated,
+            final int expectedStandbyCount,
+            final int expectedStandbyStoreCount
+    ) {
+        final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
+        final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);

Review Comment:
   Is the order deterministic? Do we know that `get(0)` returns the metadata 
from `streamOne` ?



##########
streams/integration-tests/src/test/java/org/apache/kafka/streams/integration/IQv2EndpointToPartitionsIntegrationTest.java:
##########
@@ -175,69 +183,68 @@ public void 
shouldGetCorrectHostPartitionInformation(final String groupProtocolC
                     }, TestUtils.DEFAULT_MAX_WAIT_MS,
                             "Kafka Streams clients 1 and 2 never got metadata 
about standby tasks");
 
-                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 2,
+                    waitForCondition(() -> 
streamsOne.metadataForAllStreamsClients().iterator().next().topicPartitions().size()
 == 4,
                             IntegrationTestUtils.DEFAULT_TIMEOUT,
                             () -> "Kafka Streams one didn't give up active 
tasks");
 
-                    final List<StreamsMetadata> allClientMetadataUpdated = new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients());
-
-                    final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
-                    final Set<TopicPartition> streamsOneActiveTopicPartitions 
= streamsOneMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsOneStandbyTopicPartitions 
= streamsOneMetadata.standbyTopicPartitions();
-                    final Set<String> streamsOneStoreNames = 
streamsOneMetadata.stateStoreNames();
-                    final Set<String> streamsOneStandbyStoreNames = 
streamsOneMetadata.standbyStateStoreNames();
-
-                    assertEquals(2020, streamsOneMetadata.hostInfo().port());
-                    assertEquals(2, streamsOneActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyTopicPartitions.size());
-                    assertEquals(1, streamsOneStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsOneStandbyStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsOneStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsOneStandbyStoreNames.iterator().next());
-                    }
-
-                    final long streamsOneRepartitionTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsOneSourceTopicCount = 
streamsOneActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsOneRepartitionTopicCount);
-                    assertEquals(1, streamsOneSourceTopicCount);
-
-                    final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
-                    final Set<TopicPartition> streamsTwoActiveTopicPartitions 
= streamsTwoMetadata.topicPartitions();
-                    final Set<TopicPartition> streamsTwoStandbyTopicPartitions 
= streamsTwoMetadata.standbyTopicPartitions();
-                    final Set<String> streamsTwoStateStoreNames = 
streamsTwoMetadata.stateStoreNames();
-                    final Set<String> streamsTwoStandbyStateStoreNames = 
streamsTwoMetadata.standbyStateStoreNames();
-
-                    assertEquals(3030, streamsTwoMetadata.hostInfo().port());
-                    assertEquals(2, streamsTwoActiveTopicPartitions.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyTopicPartitions.size());
-                    assertEquals(1, streamsTwoStateStoreNames.size());
-                    assertEquals(expectedStandbyCount, 
streamsTwoStandbyStateStoreNames.size());
-                    assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStateStoreNames.iterator().next());
-                    if (usingStandbyReplicas) {
-                        assertEquals(EXPECTED_STORE_NAME, 
streamsTwoStandbyStateStoreNames.iterator().next());
-                    }
-
-                    final long streamsTwoRepartitionTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-repartition")).count();
-                    final long streamsTwoSourceTopicCount = 
streamsTwoActiveTopicPartitions.stream().filter(tp -> 
tp.topic().contains("-input-two")).count();
-                    assertEquals(1, streamsTwoRepartitionTopicCount);
-                    assertEquals(1, streamsTwoSourceTopicCount);
-
-                    if (usingStandbyReplicas) {
-                        final TopicPartition streamsOneStandbyTopicPartition = 
streamsOneStandbyTopicPartitions.iterator().next();
-                        final TopicPartition streamsTwoStandbyTopicPartition = 
streamsTwoStandbyTopicPartitions.iterator().next();
-                        final String streamsOneStandbyTopicName = 
streamsOneStandbyTopicPartition.topic();
-                        final String streamsTwoStandbyTopicName = 
streamsTwoStandbyTopicPartition.topic();
-                        assertEquals(streamsOneStandbyTopicName, 
streamsTwoStandbyTopicName);
-                        
assertNotEquals(streamsOneStandbyTopicPartition.partition(), 
streamsTwoStandbyTopicPartition.partition());
-                    }
+                    verifyClientMetadata(usingStandbyReplicas, new 
ArrayList<>(streamsTwo.metadataForAllStreamsClients()), expectedStandbyCount, 
expectedStandbyStoreCount);
                 }
             }
         } finally {
             closeCluster();
         }
     }
 
+    private static void verifyClientMetadata(
+            final boolean usingStandbyReplicas,
+            final List<StreamsMetadata> allClientMetadataUpdated,
+            final int expectedStandbyCount,
+            final int expectedStandbyStoreCount
+    ) {
+        final StreamsMetadata streamsOneMetadata = 
allClientMetadataUpdated.get(0);
+        final StreamsMetadata streamsTwoMetadata = 
allClientMetadataUpdated.get(1);
+
+        verifyHostMetadata(streamsOneMetadata, 2020, expectedStandbyCount, 
expectedStandbyStoreCount, usingStandbyReplicas);
+        verifyHostMetadata(streamsTwoMetadata, 3030, expectedStandbyCount, 
expectedStandbyStoreCount, usingStandbyReplicas);
+
+        if (usingStandbyReplicas) {
+            final Set<TopicPartition> streamsOneActiveRepartition = 
streamsOneMetadata.topicPartitions().stream()
+                .filter(tp -> 
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+            final Set<TopicPartition> streamsTwoActiveRepartition = 
streamsTwoMetadata.topicPartitions().stream()
+                .filter(tp -> 
tp.topic().contains("-repartition")).collect(Collectors.toSet());
+            assertEquals(streamsTwoActiveRepartition, 
streamsOneMetadata.standbyTopicPartitions());
+            assertEquals(streamsOneActiveRepartition, 
streamsTwoMetadata.standbyTopicPartitions());
+        }
+    }
+
+    private static void verifyHostMetadata(
+            final StreamsMetadata metadata,
+            final int expectedPort,
+            final int expectedStandbyCount,
+            final int expectedStandbyStoreCount,
+            final boolean usingStandbyReplicas
+    ) {
+        final Set<TopicPartition> activeTopicPartitions = 
metadata.topicPartitions();
+        final Set<TopicPartition> standbyTopicPartitions = 
metadata.standbyTopicPartitions();
+        final Set<String> storeNames = metadata.stateStoreNames();
+        final Set<String> standbyStoreNames = 
metadata.standbyStateStoreNames();
+
+        assertEquals(expectedPort, metadata.hostInfo().port());
+        assertEquals(4, activeTopicPartitions.size());
+        assertEquals(expectedStandbyCount, standbyTopicPartitions.size());
+        assertEquals(1, storeNames.size());
+        assertEquals(expectedStandbyStoreCount, standbyStoreNames.size());

Review Comment:
   ```suggestion
           assertEquals(Set.of(EXPECTED_STORE_NAME), 
metadata.stateStoreNames());
   ```



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

Reply via email to