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(