Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The reported edge case is specifically about the subscriptionId changing, so should we only skip setting the terminating flag for the ID check? I agree it's more consistent to include all validation checks, but I wanted to raise this for discussion

// 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Loading