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 678c0e07e47 MINOR: Kafka Streams should log broker provided KIP-1071 
config values (#22730)
678c0e07e47 is described below

commit 678c0e07e4733c5a592e52046dc2c4e1625587f1
Author: Matthias J. Sax <[email protected]>
AuthorDate: Fri Jul 17 14:03:20 2026 -0700

    MINOR: Kafka Streams should log broker provided KIP-1071 config values 
(#22730)
    
    In the "streams" protocol, certain configs moved from the client to the
    broker, and the broker provides some of these configs to the client via
    the heartbeat response. The client should log these values for better
    client-side observability.
    
    Reviewers: Lianet Magrans <[email protected]>, Bill Bejeck
     <[email protected]>
---
 .../StreamsGroupHeartbeatRequestManager.java       |  28 +-
 .../StreamsGroupHeartbeatRequestManagerTest.java   | 309 ++++++++++++++++++++-
 2 files changed, 333 insertions(+), 4 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index bd554051f30..5ba7d1deb20 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -650,9 +650,24 @@ public class StreamsGroupHeartbeatRequestManager 
implements RequestManager {
         
heartbeatRequestState.updateHeartbeatIntervalMs(data.heartbeatIntervalMs());
         heartbeatRequestState.onSuccessfulAttempt(currentTimeMs);
         
heartbeatState.setEndpointInformationEpoch(data.endpointInformationEpoch());
-        
streamsRebalanceData.setHeartbeatIntervalMs(data.heartbeatIntervalMs());
-        
streamsRebalanceData.setTaskOffsetIntervalMs(data.taskOffsetIntervalMs());
-        
streamsRebalanceData.setAcceptableRecoveryLag(data.acceptableRecoveryLag());
+        // A leaving member's response carries no group configuration (the 
fields hold protocol defaults), so do not
+        // log or store it while shutting down. Normal responses have 
memberEpoch >= 0; leave responses use the
+        // negative leave sentinels (fenced members take the error path 
instead). Log only when a value changes, to
+        // avoid repeating it on every heartbeat: this fires on first receipt 
(values start unset) and on any later change.
+        if (data.memberEpoch() >= 0) {
+            if (data.heartbeatIntervalMs() != 
streamsRebalanceData.heartbeatIntervalMs()
+                    || data.taskOffsetIntervalMs() != 
streamsRebalanceData.taskOffsetIntervalMs()
+                    || data.acceptableRecoveryLag() != 
streamsRebalanceData.acceptableRecoveryLag()) {
+                logger.info("Received Streams group configuration from the 
group coordinator: "
+                        + "heartbeatIntervalMs={}, taskOffsetIntervalMs={}, 
acceptableRecoveryLag={}",
+                    describeConfig(data.heartbeatIntervalMs(), 1),
+                    describeConfig(data.taskOffsetIntervalMs(), 1),
+                    describeConfig(data.acceptableRecoveryLag(), 0));
+            }
+            
streamsRebalanceData.setHeartbeatIntervalMs(data.heartbeatIntervalMs());
+            
streamsRebalanceData.setTaskOffsetIntervalMs(data.taskOffsetIntervalMs());
+            
streamsRebalanceData.setAcceptableRecoveryLag(data.acceptableRecoveryLag());
+        }
 
         if (data.topologyDescriptionRequired() && 
streamsRebalanceData.wireTopologyDescription() != null) {
             logger.info("Broker requested topology description push");
@@ -677,6 +692,13 @@ public class StreamsGroupHeartbeatRequestManager 
implements RequestManager {
         membershipManager.onHeartbeatSuccess(response);
     }
 
+    // Renders a coordinator-provided config value for logging, or a note when 
the broker did not provide it. An older
+    // broker leaves these at their protocol defaults (intervals 0, 
acceptableRecoveryLag -1); a value below minValid
+    // means "not provided".
+    private static String describeConfig(final long value, final long 
minValid) {
+        return value < minValid ? "not provided (older broker)" : 
Long.toString(value);
+    }
+
     private void onErrorResponse(final StreamsGroupHeartbeatResponse response, 
final long currentTimeMs) {
         final Errors error = Errors.forCode(response.data().errorCode());
         final String errorMessage = response.data().errorMessage();
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
index 7a9c5ff86c6..860d3439f8f 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
@@ -43,6 +43,7 @@ import org.apache.kafka.common.utils.Time;
 import org.apache.kafka.common.utils.Timer;
 import org.apache.kafka.common.utils.internals.LogContext;
 
+import org.apache.logging.log4j.Level;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.junit.jupiter.params.ParameterizedTest;
@@ -92,6 +93,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
 
     private static final LogContext LOG_CONTEXT = new LogContext("test");
     private static final long RECEIVED_HEARTBEAT_INTERVAL_MS = 1200;
+    private static final int RECEIVED_TASK_OFFSET_INTERVAL_MS = 7000;
+    private static final long RECEIVED_ACCEPTABLE_RECOVERY_LAG = 4242;
     private static final int DEFAULT_MAX_POLL_INTERVAL_MS = 10000;
     private static final String GROUP_ID = "group-id";
     private static final String MEMBER_ID = "member-id";
@@ -645,6 +648,135 @@ class StreamsGroupHeartbeatRequestManagerTest {
         }
     }
 
+    @Test
+    public void testLogsReceivedGroupConfigWhenAnyValueChanges() {
+        try (
+            final MockedConstruction<HeartbeatRequestState> 
heartbeatRequestStateMockedConstruction = mockConstruction(
+                HeartbeatRequestState.class,
+                (mock, context) -> 
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+            final LogCaptureAppender logAppender = 
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+        ) {
+            
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class, 
Level.INFO);
+            final StreamsGroupHeartbeatRequestManager heartbeatRequestManager 
= createStreamsGroupHeartbeatRequestManager();
+            
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+            when(membershipManager.groupId()).thenReturn(GROUP_ID);
+            when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+            when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+            
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+            final int heartbeatIntervalMs = (int) 
RECEIVED_HEARTBEAT_INTERVAL_MS;
+            final int taskOffsetIntervalMs = RECEIVED_TASK_OFFSET_INTERVAL_MS;
+            final long acceptableRecoveryLag = 
RECEIVED_ACCEPTABLE_RECOVERY_LAG;
+
+            // First response: all three values change from their initial 
unset (-1) state, so they are logged.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs, 
taskOffsetIntervalMs, acceptableRecoveryLag));
+            assertEquals(1, countConfigLogs(logAppender), "Config should be 
logged on first receipt.");
+            assertTrue(logAppender.getMessages().stream().anyMatch(m -> 
m.contains(
+                    "heartbeatIntervalMs=" + heartbeatIntervalMs
+                        + ", taskOffsetIntervalMs=" + taskOffsetIntervalMs
+                        + ", acceptableRecoveryLag=" + acceptableRecoveryLag)),
+                "The logged message should contain the received config 
values.");
+
+            // Identical values: nothing changed, so it is NOT logged again.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs, 
taskOffsetIntervalMs, acceptableRecoveryLag));
+            assertEquals(1, countConfigLogs(logAppender), "Unchanged config 
must not be logged again.");
+
+            // Only heartbeatIntervalMs changes -> logged.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs + 1, 
taskOffsetIntervalMs, acceptableRecoveryLag));
+            assertEquals(2, countConfigLogs(logAppender), "A change to 
heartbeatIntervalMs must be logged.");
+            // Same values again -> not logged.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs + 1, 
taskOffsetIntervalMs, acceptableRecoveryLag));
+            assertEquals(2, countConfigLogs(logAppender));
+
+            // Only taskOffsetIntervalMs changes -> logged.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs + 1, 
taskOffsetIntervalMs + 1, acceptableRecoveryLag));
+            assertEquals(3, countConfigLogs(logAppender), "A change to 
taskOffsetIntervalMs must be logged.");
+
+            // Only acceptableRecoveryLag changes -> logged.
+            completeSuccessfulHeartbeat(heartbeatRequestManager,
+                buildClientResponseWithConfig(heartbeatIntervalMs + 1, 
taskOffsetIntervalMs + 1, acceptableRecoveryLag + 1));
+            assertEquals(4, countConfigLogs(logAppender), "A change to 
acceptableRecoveryLag must be logged.");
+        }
+    }
+
+    @Test
+    public void testLogsNotProvidedWhenOlderBrokerOmitsConfig() {
+        try (
+            final MockedConstruction<HeartbeatRequestState> 
heartbeatRequestStateMockedConstruction = mockConstruction(
+                HeartbeatRequestState.class,
+                (mock, context) -> 
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+            final LogCaptureAppender logAppender = 
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+        ) {
+            
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class, 
Level.INFO);
+            final StreamsGroupHeartbeatRequestManager heartbeatRequestManager 
= createStreamsGroupHeartbeatRequestManager();
+            
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+            when(membershipManager.groupId()).thenReturn(GROUP_ID);
+            when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+            when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+            
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+            // An older broker sets only heartbeatIntervalMs; 
taskOffsetIntervalMs arrives as 0 and acceptableRecoveryLag
+            // as -1. Those must be rendered as "not provided (older broker)" 
rather than the raw defaults.
+            completeSuccessfulHeartbeat(heartbeatRequestManager, 
buildClientResponseWithConfig(5000, 0, -1L));
+            assertTrue(logAppender.getMessages().stream().anyMatch(m -> 
m.contains(
+                    "heartbeatIntervalMs=5000, "
+                        + "taskOffsetIntervalMs=not provided (older broker), "
+                        + "acceptableRecoveryLag=not provided (older 
broker)")),
+                "An older broker's unset config should be logged as 'not 
provided (older broker)'.");
+        }
+    }
+
+    @Test
+    public void testDoesNotLogOrUpdateGroupConfigOnLeaveResponse() {
+        try (
+            final MockedConstruction<HeartbeatRequestState> 
heartbeatRequestStateMockedConstruction = mockConstruction(
+                HeartbeatRequestState.class,
+                (mock, context) -> 
when(mock.canSendRequest(time.milliseconds())).thenReturn(true));
+            final LogCaptureAppender logAppender = 
LogCaptureAppender.createAndRegister(StreamsGroupHeartbeatRequestManager.class)
+        ) {
+            
logAppender.setClassLogger(StreamsGroupHeartbeatRequestManager.class, 
Level.INFO);
+            final StreamsGroupHeartbeatRequestManager heartbeatRequestManager 
= createStreamsGroupHeartbeatRequestManager();
+            
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(coordinatorNode));
+            when(membershipManager.groupId()).thenReturn(GROUP_ID);
+            when(membershipManager.memberId()).thenReturn(MEMBER_ID);
+            when(membershipManager.memberEpoch()).thenReturn(MEMBER_EPOCH);
+            
when(membershipManager.groupInstanceId()).thenReturn(Optional.of(INSTANCE_ID));
+
+            // A normal response applies and logs the config once.
+            completeSuccessfulHeartbeat(heartbeatRequestManager, 
buildClientResponse());
+            assertEquals(1, countConfigLogs(logAppender));
+            assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS, 
streamsRebalanceData.heartbeatIntervalMs());
+            assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS, 
streamsRebalanceData.taskOffsetIntervalMs());
+            assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG, 
streamsRebalanceData.acceptableRecoveryLag());
+
+            // A leave response (memberEpoch < 0) carries no real config; it 
must neither be logged again nor
+            // overwrite the stored config with the response's protocol 
defaults.
+            completeSuccessfulHeartbeat(heartbeatRequestManager, 
buildClientLeaveResponse());
+            assertEquals(1, countConfigLogs(logAppender), "A leave response 
must not log group config.");
+            assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS, 
streamsRebalanceData.heartbeatIntervalMs());
+            assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS, 
streamsRebalanceData.taskOffsetIntervalMs());
+            assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG, 
streamsRebalanceData.acceptableRecoveryLag());
+        }
+    }
+
+    private static long countConfigLogs(final LogCaptureAppender logAppender) {
+        return logAppender.getMessages().stream()
+            .filter(m -> m.contains("Received Streams group configuration from 
the group coordinator:"))
+            .count();
+    }
+
+    private void completeSuccessfulHeartbeat(final 
StreamsGroupHeartbeatRequestManager heartbeatRequestManager,
+                                             final ClientResponse response) {
+        final NetworkClientDelegate.PollResult result = 
heartbeatRequestManager.poll(time.milliseconds());
+        assertEquals(1, result.unsentRequests.size());
+        result.unsentRequests.get(0).handler().onComplete(response);
+    }
+
     @ParameterizedTest
     @ValueSource(booleans = {false, true})
     public void testBuildingHeartbeatRequestFieldsThatAreAlwaysSent(final 
boolean instanceIdPresent) {
@@ -1006,6 +1138,69 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertEquals(500L, reportedOffsets.get(coldWarmup));
     }
 
+    @Test
+    public void testTaskOffsetsReportedForStandbyAndWarmupTasks() {
+        // Per KIP-1071, a member reports cumulative changelog offsets for all 
of its tasks with local state.
+        // Running active tasks are intentionally excluded (trivially caught 
up), but standby and warm-up tasks
+        // both carry local state, so a single heartbeat must report offsets 
for BOTH, not just the warm-up.
+        final StreamsRebalanceData.TaskId standbyTask = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final StreamsRebalanceData.TaskId warmupTask = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_2, 0);
+
+        final Map<StreamsRebalanceData.TaskId, Long> offsets = Map.of(
+            standbyTask, 300L,
+            warmupTask, 950L
+        );
+        final Map<StreamsRebalanceData.TaskId, Long> endOffsets = Map.of(
+            standbyTask, 400L,
+            warmupTask, 1000L
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            () -> offsets,
+            () -> endOffsets
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(standbyTask),
+            Set.of(warmupTask),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(100L);
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // First STABLE heartbeat: the assignment-changed trigger reports 
offsets for all reported tasks.
+        final StreamsGroupHeartbeatRequestData request = 
heartbeatState.buildRequestData();
+
+        final Map<StreamsRebalanceData.TaskId, Long> reportedOffsets = 
request.taskOffsets().stream()
+            .collect(Collectors.toMap(
+                t -> new StreamsRebalanceData.TaskId(t.subtopologyId(), 
t.partition()),
+                StreamsGroupHeartbeatRequestData.TaskOffset::offset
+            ));
+        assertEquals(300L, reportedOffsets.get(standbyTask));
+        assertEquals(950L, reportedOffsets.get(warmupTask));
+
+        final Map<StreamsRebalanceData.TaskId, Long> reportedEndOffsets = 
request.taskEndOffsets().stream()
+            .collect(Collectors.toMap(
+                t -> new StreamsRebalanceData.TaskId(t.subtopologyId(), 
t.partition()),
+                StreamsGroupHeartbeatRequestData.TaskOffset::offset
+            ));
+        assertEquals(400L, reportedEndOffsets.get(standbyTask));
+        assertEquals(1000L, reportedEndOffsets.get(warmupTask));
+    }
+
     @Test
     public void testHotWarmupTaskDisabledWhenEndOffsetMissing() {
         // Without an end-offset entry for the warmup task, lag cannot be 
computed.
@@ -1307,6 +1502,68 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertNull(followUp.taskEndOffsets());
     }
 
+    @Test
+    public void testTaskOffsetReportingCadenceConditionsInterplay() {
+        // Stages the three task-offset reporting conditions in one sequence 
to validate their interplay
+        // (not just in isolation):
+        //   (1) report only if the offsets changed since the last heartbeat;
+        //   (2) when NOT "hot" (warm-up lag above acceptable.recovery.lag), 
report only when
+        //       task.offset.interval.ms has elapsed (not on every heartbeat);
+        //   (3) when "hot" (warm-up lag at/below acceptable.recovery.lag), 
report on every heartbeat,
+        //       without waiting for the interval -- but condition (1) still 
applies (unchanged => not sent).
+        final StreamsRebalanceData.TaskId warmup = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> offsets =
+            new AtomicReference<>(Map.of(warmup, 500L)); // lag = 1000 - 500 = 
500 > 100 -> not hot
+        final Map<StreamsRebalanceData.TaskId, Long> endOffsets = 
Map.of(warmup, 1000L);
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            offsets::get,
+            () -> endOffsets
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(),
+            Set.of(warmup),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(100L);
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // Baseline: first STABLE heartbeat sends offsets via the 
assignment-changed trigger.
+        assertEquals(500L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+        // (1) not hot, unchanged, interval not elapsed -> not sent.
+        assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+        // (2) negative: not hot, CHANGED, but interval not elapsed -> still 
not sent.
+        offsets.set(Map.of(warmup, 550L)); // lag = 450 > 100 -> not hot
+        time.sleep(500); // < task.offset.interval.ms (1000)
+        assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+        // (2) positive: not hot, changed, interval now elapsed -> sent.
+        time.sleep(500); // total 1000 since last send -> interval passed
+        assertEquals(550L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+        // (3) hot, changed, interval NOT elapsed -> sent immediately (hot 
trigger).
+        offsets.set(Map.of(warmup, 950L)); // lag = 50 <= 100 -> hot
+        assertEquals(950L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+        // (3) + (1): hot but unchanged, interval not elapsed -> not sent 
(changed-gate still applies).
+        assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+        // (3) again: hot and changed, interval not elapsed -> sent again 
(every heartbeat while hot and progressing).
+        offsets.set(Map.of(warmup, 980L)); // lag = 20 <= 100 -> still hot
+        assertEquals(980L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+    }
+
     private StreamsRebalanceData newRebalanceDataWithStandbyOffsets(
             final StreamsRebalanceData.TaskId standbyTaskId,
             final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> 
taskOffsetSum,
@@ -2258,7 +2515,7 @@ class StreamsGroupHeartbeatRequestManagerTest {
     }
 
     @Test
-    public void testStreamsRebalanceDataHeartbeatIntervalMsUpdatedOnSuccess() {
+    public void testStreamsRebalanceDataUpdatedOnSuccess() {
         try (
                 final MockedConstruction<HeartbeatRequestState> ignored = 
mockConstruction(
                         HeartbeatRequestState.class,
@@ -2273,6 +2530,11 @@ class StreamsGroupHeartbeatRequestManagerTest {
 
             // Initially, heartbeatIntervalMs should be -1
             assertEquals(-1, streamsRebalanceData.heartbeatIntervalMs());
+            // The broker returns task.offset.interval.ms and 
acceptable.recovery.lag (KAFKA-18652)
+            // so the client knows how often to report task changelog offsets;
+            // the client must store it in StreamsRebalanceData.
+            assertEquals(-1, streamsRebalanceData.taskOffsetIntervalMs());
+            assertEquals(-1, streamsRebalanceData.acceptableRecoveryLag());
 
             final NetworkClientDelegate.PollResult result = 
heartbeatRequestManager.poll(time.milliseconds());
             assertEquals(1, result.unsentRequests.size());
@@ -2283,6 +2545,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
 
             // After successful response, heartbeatIntervalMs should be updated
             assertEquals(RECEIVED_HEARTBEAT_INTERVAL_MS, 
streamsRebalanceData.heartbeatIntervalMs());
+            assertEquals(RECEIVED_TASK_OFFSET_INTERVAL_MS, 
streamsRebalanceData.taskOffsetIntervalMs());
+            assertEquals(RECEIVED_ACCEPTABLE_RECOVERY_LAG, 
streamsRebalanceData.acceptableRecoveryLag());
         }
     }
 
@@ -2378,6 +2642,49 @@ class StreamsGroupHeartbeatRequestManagerTest {
                 new StreamsGroupHeartbeatResponseData()
                     .setPartitionsByUserEndpoint(ENDPOINT_TO_PARTITIONS)
                     .setHeartbeatIntervalMs((int) 
RECEIVED_HEARTBEAT_INTERVAL_MS)
+                    .setTaskOffsetIntervalMs(RECEIVED_TASK_OFFSET_INTERVAL_MS)
+                    .setAcceptableRecoveryLag(RECEIVED_ACCEPTABLE_RECOVERY_LAG)
+            )
+        );
+    }
+
+    // A successful (non-leave) response carrying specific group-config values.
+    private ClientResponse buildClientResponseWithConfig(final int 
heartbeatIntervalMs,
+                                                         final int 
taskOffsetIntervalMs,
+                                                         final long 
acceptableRecoveryLag) {
+        return new ClientResponse(
+            new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1, "", 
1),
+            null,
+            "-1",
+            time.milliseconds(),
+            time.milliseconds(),
+            false,
+            null,
+            null,
+            new StreamsGroupHeartbeatResponse(
+                new StreamsGroupHeartbeatResponseData()
+                    .setHeartbeatIntervalMs(heartbeatIntervalMs)
+                    .setTaskOffsetIntervalMs(taskOffsetIntervalMs)
+                    .setAcceptableRecoveryLag(acceptableRecoveryLag)
+            )
+        );
+    }
+
+    // A successful response to a leaving member: memberEpoch is negative and 
no group configuration is set, so the
+    // config fields carry their protocol defaults (as the coordinator sends 
for a leave).
+    private ClientResponse buildClientLeaveResponse() {
+        return new ClientResponse(
+            new RequestHeader(ApiKeys.STREAMS_GROUP_HEARTBEAT, (short) 1, "", 
1),
+            null,
+            "-1",
+            time.milliseconds(),
+            time.milliseconds(),
+            false,
+            null,
+            null,
+            new StreamsGroupHeartbeatResponse(
+                new StreamsGroupHeartbeatResponseData()
+                    .setMemberEpoch(-1)
             )
         );
     }

Reply via email to