From ffae9e557ad24244d93f15909a706b83c1301281 Mon Sep 17 00:00:00 2001 From: Ken Huang Date: Thu, 6 Aug 2026 09:00:29 +0800 Subject: [PATCH] fix the bug --- .../kafka/server/ClientMetricsManager.java | 8 ++-- .../server/ClientMetricsManagerTest.java | 47 +++++++++++++++++++ 2 files changed, 52 insertions(+), 3 deletions(-) diff --git a/server/src/main/java/org/apache/kafka/server/ClientMetricsManager.java b/server/src/main/java/org/apache/kafka/server/ClientMetricsManager.java index d938e0dd4cf07..1cbf5eaeebce1 100644 --- a/server/src/main/java/org/apache/kafka/server/ClientMetricsManager.java +++ b/server/src/main/java/org/apache/kafka/server/ClientMetricsManager.java @@ -204,11 +204,13 @@ public PushTelemetryResponse processPushTelemetryRequest(PushTelemetryRequest re log.debug("Error validating push telemetry request from client [{}]", clientInstanceId, exception); clientInstance.lastKnownError(Errors.forException(exception)); return request.getErrorResponse(0, exception); - } finally { - // Update the client instance with the latest push request parameters. - clientInstance.terminating(request.data().terminating()); } + // Update the client instance with the latest push request parameters only + // after successful validation. Setting terminating on validation failure would + // permanently lock out the client instance from future requests. + clientInstance.terminating(request.data().terminating()); + // Push the metrics to the external client receiver plugin. ByteBuffer metrics = request.data().metrics(); if (metrics != null && metrics.limit() > 0) { diff --git a/server/src/test/java/org/apache/kafka/server/ClientMetricsManagerTest.java b/server/src/test/java/org/apache/kafka/server/ClientMetricsManagerTest.java index f6420a0edde96..bfe0e3bd399e5 100644 --- a/server/src/test/java/org/apache/kafka/server/ClientMetricsManagerTest.java +++ b/server/src/test/java/org/apache/kafka/server/ClientMetricsManagerTest.java @@ -1413,6 +1413,53 @@ public void testRemoveConnectionUnknownConnectionId() throws Exception { assertEquals((double) 1, getMetric(ClientMetricsManager.ClientMetricsStats.INSTANCE_COUNT).metricValue()); } + @Test + public void testPushTelemetryTerminatingFlagNotSetOnValidationFailure() throws UnknownHostException { + clientMetricsManager.updateSubscription("sub-1", ClientMetricsTestUtils.defaultTestProperties()); + + GetTelemetrySubscriptionsRequest subscriptionsRequest = new GetTelemetrySubscriptionsRequest.Builder( + new GetTelemetrySubscriptionsRequestData(), true).build(); + + GetTelemetrySubscriptionsResponse subscriptionsResponse = clientMetricsManager.processGetTelemetrySubscriptionRequest( + subscriptionsRequest, ClientMetricsTestUtils.requestContext()); + + ClientMetricsInstance instance = clientMetricsManager.clientInstance(subscriptionsResponse.data().clientInstanceId()); + assertNotNull(instance); + assertFalse(instance.terminating()); + + // Send a push request with terminating=true but an INVALID subscriptionId. + // This simulates the race where the subscription was updated between + // GetTelemetrySubscriptions and PushTelemetry calls. + PushTelemetryRequest request = new PushTelemetryRequest.Builder( + new PushTelemetryRequestData() + .setClientInstanceId(subscriptionsResponse.data().clientInstanceId()) + .setSubscriptionId(1234) // wrong subscription id + .setTerminating(true), true).build(); + + PushTelemetryResponse response = clientMetricsManager.processPushTelemetryRequest( + request, ClientMetricsTestUtils.requestContext()); + + // Validation should fail with UNKNOWN_SUBSCRIPTION_ID + assertEquals(Errors.UNKNOWN_SUBSCRIPTION_ID, response.error()); + + assertFalse(instance.terminating(), "terminating flag should not be set when push validation fails"); + + time.sleep(ClientMetricsTestUtils.INTERVAL_MS_TEST_DEFAULT); + + PushTelemetryRequest validRequest = new PushTelemetryRequest.Builder( + new PushTelemetryRequestData() + .setClientInstanceId(subscriptionsResponse.data().clientInstanceId()) + .setSubscriptionId(subscriptionsResponse.data().subscriptionId()) + .setCompressionType(CompressionType.NONE.id) + .setTerminating(true), true).build(); + + PushTelemetryResponse validResponse = clientMetricsManager.processPushTelemetryRequest( + validRequest, ClientMetricsTestUtils.requestContext()); + + assertEquals(Errors.NONE, validResponse.error()); + assertTrue(instance.terminating()); + } + private KafkaMetric getMetric(String name) throws Exception { return getMetric(kafkaMetrics, name); }