From 7e2dc525e0670279ecb82f4e1deec5031da11a39 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 28 Jul 2026 20:23:38 -0700 Subject: [PATCH 1/3] [ISSUE #10672] Sanitize proxy telemetry failure logs --- .../grpc/v2/channel/GrpcClientChannel.java | 40 +++++++++++++++++-- .../v2/channel/GrpcClientChannelTest.java | 28 ++++++++++++- 2 files changed, 63 insertions(+), 5 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java index 0135818fb3b..a718c161619 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java @@ -25,7 +25,6 @@ import com.google.common.base.MoreObjects; import com.google.common.collect.ComparisonChain; import com.google.protobuf.InvalidProtocolBufferException; -import com.google.protobuf.TextFormat; import com.google.protobuf.util.JsonFormat; import io.grpc.StatusRuntimeException; import io.grpc.stub.StreamObserver; @@ -263,24 +262,57 @@ public String getClientId() { public void writeTelemetryCommand(TelemetryCommand command) { StreamObserver observer = this.telemetryCommandRef.get(); if (observer == null) { - log.warn("telemetry command observer is null when try to write data. command:{}, channel:{}", TextFormat.shortDebugString(command), this); + log.warn("telemetry command observer is null when try to write data. command:{}, channel:{}", + summarizeTelemetryCommand(command), this); return; } synchronized (this.telemetryWriteLock) { observer = this.telemetryCommandRef.get(); if (observer == null) { - log.warn("telemetry command observer is null when try to write data. command:{}, channel:{}", TextFormat.shortDebugString(command), this); + log.warn("telemetry command observer is null when try to write data. command:{}, channel:{}", + summarizeTelemetryCommand(command), this); return; } try { observer.onNext(command); } catch (StatusRuntimeException | IllegalStateException exception) { - log.warn("write telemetry failed. command:{}", command, exception); + log.warn("write telemetry failed. command:{}, channel:{}", summarizeTelemetryCommand(command), this, exception); this.clearClientObserver(observer); } } } + static String summarizeTelemetryCommand(TelemetryCommand command) { + if (command == null) { + return "null"; + } + + StringBuilder builder = new StringBuilder(command.getCommandCase().name()); + switch (command.getCommandCase()) { + case PRINT_THREAD_STACK_TRACE_COMMAND: + builder.append(", nonce=").append(command.getPrintThreadStackTraceCommand().getNonce()); + break; + case VERIFY_MESSAGE_COMMAND: + builder.append(", nonce=").append(command.getVerifyMessageCommand().getNonce()); + break; + case RECOVER_ORPHANED_TRANSACTION_COMMAND: + builder.append(", transactionId=") + .append(command.getRecoverOrphanedTransactionCommand().getTransactionId()); + break; + case NOTIFY_UNSUBSCRIBE_LITE_COMMAND: + builder.append(", liteTopic=") + .append(command.getNotifyUnsubscribeLiteCommand().getLiteTopic()); + break; + case SETTINGS: + builder.append(", clientType=").append(command.getSettings().getClientType()) + .append(", pubSubCase=").append(command.getSettings().getPubSubCase()); + break; + default: + break; + } + return builder.toString(); + } + @Override public String toString() { return MoreObjects.toStringHelper(this) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java index 1bdbdd9befe..d8bd9f6b6b3 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java @@ -17,9 +17,13 @@ package org.apache.rocketmq.proxy.grpc.v2.channel; +import apache.rocketmq.v2.Message; import apache.rocketmq.v2.Publishing; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.Settings; +import apache.rocketmq.v2.TelemetryCommand; +import apache.rocketmq.v2.VerifyMessageCommand; +import com.google.protobuf.ByteString; import org.apache.commons.lang3.RandomStringUtils; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.InitConfigTest; @@ -35,7 +39,9 @@ import org.mockito.junit.MockitoJUnitRunner; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -79,4 +85,24 @@ public void testChannelExtendAttributeParse() { assertEquals(clientSettings, GrpcClientChannel.parseChannelExtendAttribute(this.grpcClientChannel)); assertNull(GrpcClientChannel.parseChannelExtendAttribute(mock(RemotingChannel.class))); } -} \ No newline at end of file + + @Test + public void testSummarizeTelemetryCommandDoesNotIncludeMessagePayload() { + TelemetryCommand command = TelemetryCommand.newBuilder() + .setVerifyMessageCommand(VerifyMessageCommand.newBuilder() + .setNonce("nonce-1") + .setMessage(Message.newBuilder() + .setBody(ByteString.copyFromUtf8("secret-body")) + .build()) + .build()) + .build(); + + String summary = GrpcClientChannel.summarizeTelemetryCommand(command); + + assertTrue(summary.contains("VERIFY_MESSAGE_COMMAND")); + assertTrue(summary.contains("nonce-1")); + assertFalse(summary.contains("secret-body")); + assertFalse(summary.contains("message")); + assertFalse(summary.contains("body")); + } +} From a94c20bb74b2753f2d16990e74c352482913519f Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 28 Jul 2026 20:41:20 -0700 Subject: [PATCH 2/3] [ISSUE #10672] Address telemetry log review feedback --- .../grpc/v2/channel/GrpcClientChannel.java | 4 +++- .../v2/channel/GrpcClientChannelTest.java | 22 +++++++++++++++++++ 2 files changed, 25 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java index a718c161619..b6c08b9f02c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java @@ -297,7 +297,8 @@ static String summarizeTelemetryCommand(TelemetryCommand command) { break; case RECOVER_ORPHANED_TRANSACTION_COMMAND: builder.append(", transactionId=") - .append(command.getRecoverOrphanedTransactionCommand().getTransactionId()); + .append(command.getRecoverOrphanedTransactionCommand().getTransactionId()) + .append(", details omitted"); break; case NOTIFY_UNSUBSCRIBE_LITE_COMMAND: builder.append(", liteTopic=") @@ -308,6 +309,7 @@ static String summarizeTelemetryCommand(TelemetryCommand command) { .append(", pubSubCase=").append(command.getSettings().getPubSubCase()); break; default: + // Add explicit cases for future commands that carry sensitive payloads. break; } return builder.toString(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java index d8bd9f6b6b3..7f811c18f0f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java @@ -19,6 +19,7 @@ import apache.rocketmq.v2.Message; import apache.rocketmq.v2.Publishing; +import apache.rocketmq.v2.RecoverOrphanedTransactionCommand; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.Settings; import apache.rocketmq.v2.TelemetryCommand; @@ -105,4 +106,25 @@ public void testSummarizeTelemetryCommandDoesNotIncludeMessagePayload() { assertFalse(summary.contains("message")); assertFalse(summary.contains("body")); } + + @Test + public void testSummarizeRecoverTransactionCommandDoesNotIncludeMessagePayload() { + TelemetryCommand command = TelemetryCommand.newBuilder() + .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() + .setTransactionId("transaction-id") + .setMessage(Message.newBuilder() + .setBody(ByteString.copyFromUtf8("secret-body")) + .build()) + .build()) + .build(); + + String summary = GrpcClientChannel.summarizeTelemetryCommand(command); + + assertTrue(summary.contains("RECOVER_ORPHANED_TRANSACTION_COMMAND")); + assertTrue(summary.contains("transaction-id")); + assertTrue(summary.contains("details omitted")); + assertFalse(summary.contains("secret-body")); + assertFalse(summary.contains("message")); + assertFalse(summary.contains("body")); + } } From b5c898b4fe266215be1c5afb96b544f938a2362b Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 28 Jul 2026 23:02:39 -0700 Subject: [PATCH 3/3] [ISSUE #10672] Improve telemetry summary test coverage --- .../v2/channel/GrpcClientChannelTest.java | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java index 7f811c18f0f..369ebdbaf19 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannelTest.java @@ -18,7 +18,9 @@ package org.apache.rocketmq.proxy.grpc.v2.channel; import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.NotifyUnsubscribeLiteCommand; import apache.rocketmq.v2.Publishing; +import apache.rocketmq.v2.PrintThreadStackTraceCommand; import apache.rocketmq.v2.RecoverOrphanedTransactionCommand; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.Settings; @@ -127,4 +129,38 @@ public void testSummarizeRecoverTransactionCommandDoesNotIncludeMessagePayload() assertFalse(summary.contains("message")); assertFalse(summary.contains("body")); } + + @Test + public void testSummarizeTelemetryCommandDiagnosticFields() { + assertEquals("null", GrpcClientChannel.summarizeTelemetryCommand(null)); + assertEquals("COMMAND_NOT_SET", GrpcClientChannel.summarizeTelemetryCommand(TelemetryCommand.getDefaultInstance())); + + TelemetryCommand settingsCommand = TelemetryCommand.newBuilder() + .setSettings(Settings.newBuilder() + .setPublishing(Publishing.getDefaultInstance()) + .build()) + .build(); + String settingsSummary = GrpcClientChannel.summarizeTelemetryCommand(settingsCommand); + assertTrue(settingsSummary.contains("SETTINGS")); + assertTrue(settingsSummary.contains("clientType=")); + assertTrue(settingsSummary.contains("pubSubCase=PUBLISHING")); + + TelemetryCommand threadStackCommand = TelemetryCommand.newBuilder() + .setPrintThreadStackTraceCommand(PrintThreadStackTraceCommand.newBuilder() + .setNonce("stack-nonce") + .build()) + .build(); + String threadStackSummary = GrpcClientChannel.summarizeTelemetryCommand(threadStackCommand); + assertTrue(threadStackSummary.contains("PRINT_THREAD_STACK_TRACE_COMMAND")); + assertTrue(threadStackSummary.contains("stack-nonce")); + + TelemetryCommand liteCommand = TelemetryCommand.newBuilder() + .setNotifyUnsubscribeLiteCommand(NotifyUnsubscribeLiteCommand.newBuilder() + .setLiteTopic("lite-topic") + .build()) + .build(); + String liteSummary = GrpcClientChannel.summarizeTelemetryCommand(liteCommand); + assertTrue(liteSummary.contains("NOTIFY_UNSUBSCRIBE_LITE_COMMAND")); + assertTrue(liteSummary.contains("lite-topic")); + } }