fix(observability): document message propagation semantics

Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com>
This commit is contained in:
XiaoSeS 2026-08-03 15:09:55 +08:00
parent 5058cc3387
commit 801818bba2
6 changed files with 105 additions and 6 deletions

View file

@ -4,12 +4,16 @@ import io.micrometer.observation.transport.Propagator;
/**
* Adapts transport-specific message headers to the common observation boundary.
*
* <p>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.</p>
*/
public interface MessageCarrierAdapter<C>
extends Propagator.Getter<C>, Propagator.Setter<C> {
/**
* 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);
}

View file

@ -15,13 +15,20 @@ import java.util.function.Supplier;
* Propagates tracing and request correlation across asynchronous message transports.
*
* <p>The message carrier owns transport metadata. Business payloads remain independent from
* Micrometer, OpenTelemetry, MDC, and a concrete tracing backend.</p>
* 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.</p>
*/
@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<String> 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<T> 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()) {

View file

@ -205,6 +205,7 @@ public abstract class AbstractStreamConsumer<T> {
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<T> {
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;
}

View file

@ -6,6 +6,10 @@ import java.util.Map;
/**
* Adapts Redis Stream field maps to the transport-neutral message observation boundary.
*
* <p>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.</p>
*/
final class RedisStreamMessageCarrier {

View file

@ -50,6 +50,7 @@ public class RedissonScanTaskProducer implements ScanTaskProducer {
}
RStream<String, String> stream = redissonClient.getStream(streamKey, StringCodec.INSTANCE);
// Starting the producer Observation injects trusted propagation fields before XADD.
StreamMessageId messageId = messageObservationSupport.observePublish(
RedisStreamMessageCarrier.MESSAGING_SYSTEM,
streamKey,

View file

@ -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<Observation.Context> stoppedContexts = new ArrayList<>();
observationRegistry.observationConfig().observationHandler(
new ObservationHandler<Observation.Context>() {
@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<Observation.Context> 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(