From 801818bba2a7376327bd0c2bc79e1088deed5ee0 Mon Sep 17 00:00:00 2001
From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
Date: Mon, 3 Aug 2026 15:09:55 +0800
Subject: [PATCH] fix(observability): document message propagation semantics
Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
---
.../observability/MessageCarrierAdapter.java | 6 +-
.../MessageObservationSupport.java | 27 +++++--
.../stream/AbstractStreamConsumer.java | 3 +
.../stream/RedisStreamMessageCarrier.java | 4 ++
.../stream/RedissonScanTaskProducer.java | 1 +
.../MessageObservationSupportTest.java | 70 +++++++++++++++++++
6 files changed, 105 insertions(+), 6 deletions(-)
diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java
index 2fb26318..721d4167 100644
--- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java
+++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java
@@ -4,12 +4,16 @@ import io.micrometer.observation.transport.Propagator;
/**
* Adapts transport-specific message headers to the common observation boundary.
+ *
+ *
Implementations should operate on a transport envelope or header collection, never on a
+ * business DTO. Micrometer uses the inherited getter and setter to extract and inject propagation
+ * fields without exposing a concrete tracing implementation to the transport.
*/
public interface MessageCarrierAdapter
extends Propagator.Getter, Propagator.Setter {
/**
- * Removes every value associated with a transport header.
+ * Removes every value associated with a transport header before trusted context is injected.
*/
void remove(C carrier, String key);
}
diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java
index 4e84e1bd..48f25995 100644
--- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java
+++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java
@@ -15,13 +15,20 @@ import java.util.function.Supplier;
* Propagates tracing and request correlation across asynchronous message transports.
*
* The message carrier owns transport metadata. Business payloads remain independent from
- * Micrometer, OpenTelemetry, MDC, and a concrete tracing backend.
+ * Micrometer, OpenTelemetry, MDC, and a concrete tracing backend. A transport integrates by
+ * providing a {@link MessageCarrierAdapter}, then wrapping the actual send and per-message
+ * processing operations with this component.
*/
@Component
public class MessageObservationSupport {
+ /**
+ * Reserved transport field for log correlation when distributed tracing is disabled.
+ */
public static final String REQUEST_ID_FIELD = "skillhub.request_id";
+ // The propagation boundary owns these fields. Removing existing values prevents business
+ // metadata from forging a parent trace or leaking stale context into a newly published message.
private static final Set OWNED_TRANSPORT_FIELDS = Set.of(
REQUEST_ID_FIELD,
"traceparent",
@@ -53,6 +60,9 @@ public class MessageObservationSupport {
validateArguments(messagingSystem, destination, carrier, action);
Objects.requireNonNull(carrierAdapter, "carrierAdapter must not be null");
OWNED_TRANSPORT_FIELDS.forEach(field -> carrierAdapter.remove(carrier, field));
+
+ // Request ID is propagated independently because it must remain useful in modes where no
+ // tracing handler is registered. Micrometer injects W3C trace fields when tracing is active.
String requestId = requestIdAccessor.current();
if (RequestIdAccessor.isValid(requestId)) {
carrierAdapter.set(carrier, REQUEST_ID_FIELD, requestId);
@@ -63,10 +73,11 @@ public class MessageObservationSupport {
senderContext.setRemoteServiceName(messagingSystem);
return observe(
"skillhub.message.publish",
- destination + " publish",
+ "publish " + destination,
messagingSystem,
destination,
"publish",
+ "send",
senderContext,
action
);
@@ -88,6 +99,9 @@ public class MessageObservationSupport {
receiverContext.setCarrier(carrier);
receiverContext.setRemoteServiceName(messagingSystem);
+ // Micrometer scopes the extracted trace context when the Observation starts. Scope the
+ // independently propagated Request ID over the same processing boundary and restore both
+ // before the worker thread is reused.
String propagatedRequestId = carrierAdapter.get(carrier, REQUEST_ID_FIELD);
RequestIdAccessor.Scope requestIdScope = requestIdAccessor.openNullable(
RequestIdAccessor.isValid(propagatedRequestId) ? propagatedRequestId : null
@@ -95,10 +109,11 @@ public class MessageObservationSupport {
try (requestIdScope) {
return observe(
"skillhub.message.process",
- destination + " process",
+ "process " + destination,
messagingSystem,
destination,
"process",
+ "process",
receiverContext,
action
);
@@ -122,7 +137,8 @@ public class MessageObservationSupport {
String contextualName,
String messagingSystem,
String destination,
- String operation,
+ String operationName,
+ String operationType,
Observation.Context transportContext,
Supplier action
) {
@@ -130,7 +146,8 @@ public class MessageObservationSupport {
.createNotStarted(observationName, () -> transportContext, observationRegistry)
.contextualName(contextualName)
.lowCardinalityKeyValue("messaging.system", messagingSystem)
- .lowCardinalityKeyValue("messaging.operation.type", operation)
+ .lowCardinalityKeyValue("messaging.operation.name", operationName)
+ .lowCardinalityKeyValue("messaging.operation.type", operationType)
.highCardinalityKeyValue("messaging.destination.name", destination)
.start();
try (Observation.Scope ignored = observation.openScope()) {
diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/AbstractStreamConsumer.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/AbstractStreamConsumer.java
index 01a73e77..67ab8fef 100644
--- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/AbstractStreamConsumer.java
+++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/AbstractStreamConsumer.java
@@ -205,6 +205,7 @@ public abstract class AbstractStreamConsumer {
if (messages == null || messages.isEmpty()) {
return;
}
+ // Scope each entry independently because one XREADGROUP batch may contain unrelated traces.
messages.forEach(this::handleMessage);
}
@@ -243,6 +244,8 @@ public abstract class AbstractStreamConsumer {
private void handleFailure(T payload, int retryCount, Exception e) {
if (retryCount < MAX_RETRY_COUNT) {
+ // Retry publication remains inside the current consumer scope, so the new producer
+ // span and message carrier continue the original trace.
retryMessage(payload, retryCount + 1);
return;
}
diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java
index acbb349b..0752d1fc 100644
--- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java
+++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java
@@ -6,6 +6,10 @@ import java.util.Map;
/**
* Adapts Redis Stream field maps to the transport-neutral message observation boundary.
+ *
+ * Redis Stream entries do not have a separate header collection, so reserved propagation fields
+ * share the entry map with business fields. {@code MessageObservationSupport} owns and sanitizes
+ * those reserved fields; business payload types remain unaware of them.
*/
final class RedisStreamMessageCarrier {
diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedissonScanTaskProducer.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedissonScanTaskProducer.java
index 144693f1..691669c5 100644
--- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedissonScanTaskProducer.java
+++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedissonScanTaskProducer.java
@@ -50,6 +50,7 @@ public class RedissonScanTaskProducer implements ScanTaskProducer {
}
RStream stream = redissonClient.getStream(streamKey, StringCodec.INSTANCE);
+ // Starting the producer Observation injects trusted propagation fields before XADD.
StreamMessageId messageId = messageObservationSupport.observePublish(
RedisStreamMessageCarrier.MESSAGING_SYSTEM,
streamKey,
diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/MessageObservationSupportTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/MessageObservationSupportTest.java
index 02f114df..cc1b585d 100644
--- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/MessageObservationSupportTest.java
+++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/MessageObservationSupportTest.java
@@ -1,7 +1,9 @@
package com.iflytek.skillhub.observability;
import com.iflytek.skillhub.observability.tracing.SkillHubTracingConfiguration;
+import io.micrometer.common.KeyValue;
import io.micrometer.observation.Observation;
+import io.micrometer.observation.ObservationHandler;
import io.micrometer.observation.ObservationRegistry;
import io.micrometer.tracing.Span;
import io.micrometer.tracing.Tracer;
@@ -151,6 +153,74 @@ class MessageObservationSupportTest {
.isEqualTo("trusted-request");
}
+ @Test
+ void shouldUseOpenTelemetryMessagingOperationSemantics() {
+ ObservationRegistry observationRegistry = ObservationRegistry.create();
+ List stoppedContexts = new ArrayList<>();
+ observationRegistry.observationConfig().observationHandler(
+ new ObservationHandler() {
+ @Override
+ public void onStop(Observation.Context context) {
+ stoppedContexts.add(context);
+ }
+
+ @Override
+ public boolean supportsContext(Observation.Context context) {
+ return true;
+ }
+ }
+ );
+ MessageObservationSupport support = new MessageObservationSupport(
+ observationRegistry,
+ new RequestIdAccessor()
+ );
+ TestCarrier carrier = new TestCarrier();
+
+ support.observePublish(
+ "redis",
+ "skillhub:scan:requests",
+ carrier,
+ TEST_CARRIER_ADAPTER,
+ () -> null
+ );
+ support.observeProcess(
+ "redis",
+ "skillhub:scan:requests",
+ carrier,
+ TEST_CARRIER_ADAPTER,
+ () -> null
+ );
+
+ Observation.Context publish = findContext(stoppedContexts, "skillhub.message.publish");
+ assertThat(publish.getContextualName()).isEqualTo("publish skillhub:scan:requests");
+ assertThat(lowCardinalityValue(publish, "messaging.operation.name")).isEqualTo("publish");
+ assertThat(lowCardinalityValue(publish, "messaging.operation.type")).isEqualTo("send");
+
+ Observation.Context process = findContext(stoppedContexts, "skillhub.message.process");
+ assertThat(process.getContextualName()).isEqualTo("process skillhub:scan:requests");
+ assertThat(lowCardinalityValue(process, "messaging.operation.name")).isEqualTo("process");
+ assertThat(lowCardinalityValue(process, "messaging.operation.type")).isEqualTo("process");
+ }
+
+ private Observation.Context findContext(
+ List contexts,
+ String observationName
+ ) {
+ return contexts.stream()
+ .filter(context -> observationName.equals(context.getName()))
+ .findFirst()
+ .orElseThrow();
+ }
+
+ private String lowCardinalityValue(Observation.Context context, String key) {
+ for (KeyValue keyValue : context.getLowCardinalityKeyValues()) {
+ if (key.equals(keyValue.getKey())) {
+ return keyValue.getValue();
+ }
+ }
+ return null;
+ }
+
private ContextValues currentValues(RequestIdAccessor requestIdAccessor, Tracer tracer) {
Span currentSpan = tracer.currentSpan();
return new ContextValues(