diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java index abc23a53a3e..5b4abb53b17 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivity.java @@ -362,11 +362,39 @@ protected void processTelemetryException(TelemetryCommand request, Throwable t, } } if (exception.getStatus().getCode().equals(io.grpc.Status.Code.INTERNAL)) { - log.warn("process client telemetryCommand failed. request:{}", request, t); + log.warn("process client telemetryCommand failed. request:{}", summarizeTelemetryCommand(request), t); } responseObserver.onError(exception); } + static String summarizeTelemetryCommand(TelemetryCommand request) { + if (request == null) { + return "null"; + } + + StringBuilder builder = new StringBuilder(request.getCommandCase().name()); + if (request.hasStatus()) { + builder.append(", statusCode=").append(request.getStatus().getCode()); + } + switch (request.getCommandCase()) { + case SETTINGS: + builder.append(", clientType=").append(request.getSettings().getClientType()) + .append(", pubSubCase=").append(request.getSettings().getPubSubCase()); + break; + case THREAD_STACK_TRACE: + builder.append(", nonce=").append(request.getThreadStackTrace().getNonce()) + .append(", details omitted"); + break; + case VERIFY_MESSAGE_RESULT: + builder.append(", nonce=").append(request.getVerifyMessageResult().getNonce()); + break; + default: + // Add explicit cases for future telemetry requests that carry sensitive details. + break; + } + return builder.toString(); + } + protected void processAndWriteClientSettings(ProxyContext ctx, TelemetryCommand request, StreamObserver responseObserver) { GrpcClientChannel grpcClientChannel = null; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java index e215c6efaba..126e37b6f26 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java @@ -71,6 +71,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.mockito.ArgumentMatchers.any; @@ -416,6 +417,25 @@ public void onCompleted() { assertThat(result.getResult().getConsumeResult()).isEqualTo(CMResult.CR_SUCCESS); } + @Test + public void testSummarizeTelemetryCommandDoesNotIncludeThreadStackTrace() { + TelemetryCommand command = TelemetryCommand.newBuilder() + .setThreadStackTrace(ThreadStackTrace.newBuilder() + .setNonce("nonce-1") + .setThreadStackTrace("secret-stack-trace") + .build()) + .setStatus(ResponseBuilder.getInstance().buildStatus(Code.OK, Code.OK.name())) + .build(); + + String summary = ClientActivity.summarizeTelemetryCommand(command); + + assertTrue(summary.contains("THREAD_STACK_TRACE")); + assertTrue(summary.contains("statusCode=OK")); + assertTrue(summary.contains("nonce-1")); + assertTrue(summary.contains("details omitted")); + assertFalse(summary.contains("secret-stack-trace")); + } + protected CompletableFuture sendClientTelemetry(ProxyContext ctx, Settings settings) { when(grpcClientSettingsManager.getClientSettings(any())).thenReturn(settings);