From 3f14d6e3ead11374075768db0b7210d66f96c60e Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Mon, 3 Aug 2026 11:57:18 +0800 Subject: [PATCH 1/5] fix(observability): harden operational log privacy Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- .../exception/GlobalExceptionHandler.java | 25 ++--- .../skillhub/filter/RequestLoggingFilter.java | 20 ---- .../exception/GlobalExceptionHandlerTest.java | 102 +++++++++++++++++- .../filter/RequestLoggingFilterTest.java | 17 +-- 4 files changed, 121 insertions(+), 43 deletions(-) diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/exception/GlobalExceptionHandler.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/exception/GlobalExceptionHandler.java index 71241186..c2c807b6 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/exception/GlobalExceptionHandler.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/exception/GlobalExceptionHandler.java @@ -1,7 +1,6 @@ package com.iflytek.skillhub.exception; import com.iflytek.skillhub.auth.exception.AuthFlowException; -import com.iflytek.skillhub.auth.rbac.PlatformPrincipal; import com.iflytek.skillhub.dto.ApiResponse; import com.iflytek.skillhub.dto.ApiResponseFactory; import com.iflytek.skillhub.domain.shared.exception.LocalizedDomainException; @@ -113,11 +112,11 @@ public class GlobalExceptionHandler { public ResponseEntity> handleStorageAccess(StorageAccessException ex, HttpServletRequest request) { metrics.incrementStorageAccessFailure(ex.getOperation()); logger.warn( - "Object storage unavailable [requestId={}, method={}, path={}, userId={}, operation={}, key={}]", + "Object storage unavailable [requestId={}, method={}, path={}, authentication={}, operation={}, key={}]", requestIdAccessor.current(), request.getMethod(), sensitiveLogSanitizer.sanitizeRequestTarget(request), - resolveUserId(request), + resolveAuthenticationState(request), ex.getOperation(), ex.getKey(), ex @@ -142,11 +141,11 @@ public class GlobalExceptionHandler { @ExceptionHandler(Exception.class) public ResponseEntity> handleGlobalException(Exception ex, HttpServletRequest request) { logger.error( - "Unhandled API exception [requestId={}, method={}, path={}, userId={}]", + "Unhandled API exception [requestId={}, method={}, path={}, authentication={}]", requestIdAccessor.current(), request.getMethod(), sensitiveLogSanitizer.sanitizeRequestTarget(request), - resolveUserId(request), + resolveAuthenticationState(request), ex ); return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body( @@ -155,12 +154,12 @@ public class GlobalExceptionHandler { private void logHandledException(HttpStatus status, String messageCode, HttpServletRequest request) { logger.info( - "API request failed [requestId={}, status={}, method={}, path={}, userId={}, code={}]", + "API request failed [requestId={}, status={}, method={}, path={}, authentication={}, code={}]", requestIdAccessor.current(), status.value(), request.getMethod(), sensitiveLogSanitizer.sanitizeRequestTarget(request), - resolveUserId(request), + resolveAuthenticationState(request), messageCode ); } @@ -173,13 +172,11 @@ public class GlobalExceptionHandler { apiResponseFactory.error(status.value(), error.messageCode(), error.messageArgs())); } - private String resolveUserId(HttpServletRequest request) { - if (!(request.getUserPrincipal() instanceof Authentication authentication)) { - return "anonymous"; + private String resolveAuthenticationState(HttpServletRequest request) { + if (request.getUserPrincipal() instanceof Authentication authentication + && authentication.isAuthenticated()) { + return "authenticated"; } - if (authentication.getPrincipal() instanceof PlatformPrincipal principal) { - return principal.userId(); - } - return authentication.getName(); + return "anonymous"; } } diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestLoggingFilter.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestLoggingFilter.java index d46a1e49..1de3809b 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestLoggingFilter.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestLoggingFilter.java @@ -16,7 +16,6 @@ import org.springframework.web.util.ContentCachingRequestWrapper; import org.springframework.web.util.ContentCachingResponseWrapper; import java.io.IOException; -import java.io.UnsupportedEncodingException; import java.util.Set; /** @@ -27,8 +26,6 @@ import java.util.Set; public class RequestLoggingFilter extends OncePerRequestFilter { private static final Logger log = LoggerFactory.getLogger(RequestLoggingFilter.class); - private static final int MAX_LOG_BODY_LENGTH = 200; - private static final Set SKIP_PREFIXES = Set.of( "/actuator", "/favicon.ico", "/assets/" ); @@ -85,11 +82,6 @@ public class RequestLoggingFilter extends OncePerRequestFilter { sb.append(" | UA: ").append(truncate(userAgent, 80)); } - String requestBody = getRequestBody(request); - if (requestBody != null && !requestBody.isBlank()) { - sb.append(" | Body: ").append(requestBody); - } - log.info(sb.toString()); } @@ -117,18 +109,6 @@ public class RequestLoggingFilter extends OncePerRequestFilter { response.setHeader("X-Accel-Buffering", "no"); } - private String getRequestBody(ContentCachingRequestWrapper request) { - byte[] buf = request.getContentAsByteArray(); - if (buf.length > 0) { - try { - return truncate(new String(buf, request.getCharacterEncoding()), MAX_LOG_BODY_LENGTH); - } catch (UnsupportedEncodingException e) { - return "[unknown encoding]"; - } - } - return null; - } - private String truncate(String value, int maxLength) { if (value == null || value.length() <= maxLength) { return value; diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/exception/GlobalExceptionHandlerTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/exception/GlobalExceptionHandlerTest.java index 6aae4cb0..b29c5428 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/exception/GlobalExceptionHandlerTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/exception/GlobalExceptionHandlerTest.java @@ -4,28 +4,44 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.when; +import ch.qos.logback.classic.Level; +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.spi.ILoggingEvent; +import ch.qos.logback.core.read.ListAppender; +import com.iflytek.skillhub.auth.rbac.PlatformPrincipal; import com.iflytek.skillhub.dto.ApiResponse; import com.iflytek.skillhub.dto.ApiResponseFactory; import com.iflytek.skillhub.metrics.SkillHubMetrics; import com.iflytek.skillhub.observability.RequestIdAccessor; import com.iflytek.skillhub.security.SensitiveLogSanitizer; +import com.iflytek.skillhub.storage.StorageAccessException; import jakarta.servlet.http.HttpServletRequest; import java.time.Clock; import java.time.Instant; import java.time.ZoneOffset; +import java.util.List; +import java.util.Set; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.slf4j.LoggerFactory; import org.springframework.context.support.StaticMessageSource; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; +import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; import org.springframework.web.context.request.async.AsyncRequestTimeoutException; @ExtendWith(MockitoExtension.class) class GlobalExceptionHandlerTest { + private static final String STABLE_USER_ID = "stable-user-123"; + + private final Logger logger = + (Logger) LoggerFactory.getLogger(GlobalExceptionHandler.class); + @Mock private SensitiveLogSanitizer sensitiveLogSanitizer; @@ -36,12 +52,14 @@ class GlobalExceptionHandlerTest { private HttpServletRequest request; private GlobalExceptionHandler handler; + private ListAppender appender; + private RequestIdAccessor requestIdAccessor; @BeforeEach void setUp() { StaticMessageSource messageSource = new StaticMessageSource(); messageSource.addMessage("error.request.timeout", java.util.Locale.getDefault(), "Request timed out"); - RequestIdAccessor requestIdAccessor = new RequestIdAccessor(); + requestIdAccessor = new RequestIdAccessor(); ApiResponseFactory responseFactory = new ApiResponseFactory( messageSource, Clock.fixed(Instant.parse("2026-03-20T00:00:00Z"), ZoneOffset.UTC), @@ -55,6 +73,14 @@ class GlobalExceptionHandlerTest { ); } + @AfterEach + void tearDown() { + if (appender != null) { + logger.detachAppender(appender); + appender.stop(); + } + } + @Test void handleAsyncRequestTimeout_shouldReturnNoContentForSseRequests() { when(request.getRequestURI()).thenReturn("/api/v1/notifications/sse"); @@ -67,6 +93,7 @@ class GlobalExceptionHandlerTest { @Test void handleAsyncRequestTimeout_shouldReturnApiEnvelopeForNonSseRequests() { + attachAppender(); when(request.getRequestURI()).thenReturn("/api/v1/publish"); when(request.getMethod()).thenReturn("POST"); when(sensitiveLogSanitizer.sanitizeRequestTarget(request)).thenReturn("/api/v1/publish"); @@ -78,6 +105,53 @@ class GlobalExceptionHandlerTest { ApiResponse body = (ApiResponse) response.getBody(); assertThat(body.code()).isEqualTo(408); assertThat(body.msg()).isEqualTo("Request timed out"); + assertThat(loggedMessages()).anySatisfy(message -> assertThat(message) + .contains("authentication=anonymous") + .doesNotContain("userId=")); + } + + @Test + void handleGlobalException_shouldLogAuthenticationWithoutStableUserId() { + authenticateRequest(); + attachAppender(); + when(request.getMethod()).thenReturn("GET"); + when(sensitiveLogSanitizer.sanitizeRequestTarget(request)) + .thenReturn("/api/v1/skills/sensitive"); + + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("request-123")) { + ResponseEntity> response = handler.handleGlobalException( + new RuntimeException("boom"), request); + assertThat(response.getStatusCode()).isEqualTo(HttpStatus.INTERNAL_SERVER_ERROR); + } + + assertThat(loggedMessages()).anySatisfy(message -> assertThat(message) + .contains("requestId=request-123") + .contains("authentication=authenticated") + .doesNotContain(STABLE_USER_ID) + .doesNotContain("userId=")); + } + + @Test + void handleStorageAccess_shouldLogAuthenticationWithoutStableUserId() { + authenticateRequest(); + attachAppender(); + when(request.getMethod()).thenReturn("GET"); + when(sensitiveLogSanitizer.sanitizeRequestTarget(request)) + .thenReturn("/api/v1/skills/test/download"); + + StorageAccessException exception = new StorageAccessException( + "download", "skills/test.zip", new RuntimeException("unavailable")); + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("request-123")) { + ResponseEntity> response = + handler.handleStorageAccess(exception, request); + assertThat(response.getStatusCode()).isEqualTo(HttpStatus.SERVICE_UNAVAILABLE); + } + + assertThat(loggedMessages()).anySatisfy(message -> assertThat(message) + .contains("requestId=request-123") + .contains("authentication=authenticated") + .doesNotContain(STABLE_USER_ID) + .doesNotContain("userId=")); } @Test @@ -100,4 +174,30 @@ class GlobalExceptionHandlerTest { assertThatThrownBy(() -> handler.handleSessionInvalidated(ex, request)) .isSameAs(ex); } + + private void authenticateRequest() { + PlatformPrincipal principal = new PlatformPrincipal( + STABLE_USER_ID, + "User", + "user@example.com", + null, + "local", + Set.of("USER") + ); + when(request.getUserPrincipal()).thenReturn( + new UsernamePasswordAuthenticationToken(principal, null, List.of())); + } + + private void attachAppender() { + logger.setLevel(Level.INFO); + appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + } + + private List loggedMessages() { + return appender.list.stream() + .map(ILoggingEvent::getFormattedMessage) + .toList(); + } } diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/filter/RequestLoggingFilterTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/filter/RequestLoggingFilterTest.java index 11ec0aec..18e6f88a 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/filter/RequestLoggingFilterTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/filter/RequestLoggingFilterTest.java @@ -36,16 +36,17 @@ class RequestLoggingFilterTest { } @Test - void doFilterInternal_truncatesLongRequestBodyAndOmitsResponseBody() + void doFilterInternal_omitsRequestAndResponseBodies() throws ServletException, IOException { RequestLoggingFilter filter = new RequestLoggingFilter(); - String longBody = "x".repeat(5_000); + String requestBody = "{\"username\":\"alice\",\"password\":\"super-secret\"}"; + String responseBody = "x".repeat(5_000); attachAppender(); MockHttpServletRequest request = new MockHttpServletRequest("POST", "/api/test"); request.setCharacterEncoding(StandardCharsets.UTF_8.name()); request.setContentType("application/json"); - request.setContent(longBody.getBytes(StandardCharsets.UTF_8)); + request.setContent(requestBody.getBytes(StandardCharsets.UTF_8)); MockHttpServletResponse response = new MockHttpServletResponse(); response.setCharacterEncoding(StandardCharsets.UTF_8.name()); @@ -53,17 +54,17 @@ class RequestLoggingFilterTest { FilterChain filterChain = (req, res) -> { req.getReader().lines().count(); res.setContentType("application/json"); - res.getWriter().write(longBody); + res.getWriter().write(responseBody); }; filter.doFilter(request, response, filterChain); List loggedMessages = loggedMessages(); - assertThat(loggedMessages).anySatisfy(message -> - assertThat(message).contains("Body: " + "x".repeat(200) + "...[truncated]")); - assertThat(loggedMessages).noneMatch(message -> message.contains("Body: " + longBody)); + assertThat(loggedMessages).anyMatch(message -> message.contains("POST /api/test")); + assertThat(loggedMessages).noneMatch(message -> message.contains("Body:")); + assertThat(loggedMessages).noneMatch(message -> message.contains("super-secret")); assertThat(loggedMessages).noneMatch(message -> message.contains("Response Body:")); - assertThat(response.getContentAsString()).isEqualTo(longBody); + assertThat(response.getContentAsString()).isEqualTo(responseBody); } @Test From 5058cc3387554d85e87e54862461ceb73c997c18 Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:46:03 +0800 Subject: [PATCH 2/5] feat(observability): propagate message trace context Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- ...6-07-31-observability-construction-plan.md | 16 +- docs/observability-decision-map.md | 25 +- docs/observability-developer-guide.md | 28 ++- .../skillhub/config/RedisStreamConfig.java | 14 +- .../skillhub/filter/RequestIdFilter.java | 5 +- .../observability/MessageCarrierAdapter.java | 15 ++ .../MessageObservationSupport.java | 161 +++++++++++++ .../observability/RequestIdAccessor.java | 15 ++ .../stream/AbstractStreamConsumer.java | 34 ++- .../stream/RedisStreamMessageCarrier.java | 34 +++ .../stream/RedissonScanTaskProducer.java | 17 +- .../skillhub/stream/ScanTaskConsumer.java | 20 +- .../iflytek/skillhub/stream/package-info.java | 4 + .../MessageObservationSupportTest.java | 217 ++++++++++++++++++ .../stream/AbstractStreamConsumerTest.java | 41 +++- .../RedissonScanTaskProducerLoggingTest.java | 9 +- .../stream/RedissonScanTaskProducerTest.java | 37 ++- .../stream/ScanTaskConsumerLoggingTest.java | 6 +- .../ScanTaskConsumerPathSafetyTest.java | 6 +- .../skillhub/stream/ScanTaskConsumerTest.java | 6 +- 20 files changed, 655 insertions(+), 55 deletions(-) create mode 100644 server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java create mode 100644 server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java create mode 100644 server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java create mode 100644 server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/package-info.java create mode 100644 server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/MessageObservationSupportTest.java diff --git a/docs/2026-07-31-observability-construction-plan.md b/docs/2026-07-31-observability-construction-plan.md index cbbe85fe..d03651a0 100644 --- a/docs/2026-07-31-observability-construction-plan.md +++ b/docs/2026-07-31-observability-construction-plan.md @@ -282,12 +282,16 @@ MANAGEMENT_OTLP_TRACING_ENDPOINT=http://otel-collector:4318/v1/traces 测试必须重复复用同一工作线程,证明不同请求之间不会串号。 -### 10.2 长生命周期后台线程 +### 10.2 消息队列与长生命周期后台线程 -Redis Stream 消费循环、Reclaimer 和其他长生命周期线程不继承应用启动线程或任意请求的 -MDC。需要追踪具体任务时,由通用任务执行边界建立新的上下文。 +Redis Stream 消费循环和 Reclaimer 不继承应用启动线程或任意请求的 MDC。Producer 通过 +通用消息 Observation 把 W3C Trace Context 与受控 Request ID 注入 transport metadata; +Consumer/Reclaimer 逐条提取、建立 Scope,并在处理结束后清理。Scanner HTTP 调用自然成为 +Consumer Span 的子调用。 -一期不把 HTTP Trace Context 写入搜索业务 payload,也不改造可靠任务状态机。 +上下文不写入 `ScanTask` 或搜索业务 payload,也不改变可靠任务状态机。普通定时任务没有 +上游 carrier,仍建立独立执行上下文;长期延迟任务使用稳定任务 ID 或 Span Link,不维持 +超长父 Span。 ### 10.3 HTTP 出站 @@ -434,6 +438,8 @@ Kibana 使用 `trace.id` 查询日志,SkyWalking 使用同一个 Trace ID 查 - `none`、`otel-sdk`、`external-agent` 的 Spring Context 互斥。 - 未配置 OTLP endpoint 时不创建网络导出。 - 内部 Scanner 请求携带 `traceparent`。 +- Redis Stream Producer/Consumer 保持同一 Trace 和 Request ID,处理结束后线程不串号。 +- 重试发布和 Reclaimer 重新消费仍能恢复消息关联上下文。 - 外部 HTTP 请求不携带 `traceparent`。 ### 13.2 远端原型 @@ -463,6 +469,7 @@ Kibana 使用 `trace.id` 查询日志,SkyWalking 使用同一个 Trace ID 查 - HTTP 成功、4xx、5xx 和未认证请求。 - Scanner 成功、超时和失败。 - 异步事件正常执行和抛出异常。 +- Redis Stream 正常消费、失败重试、Pending Reclaim 和重复投递。 - 并发请求重复使用线程池。 - Collector 启动、停止和恢复。 - 日志采集器停止或消费变慢。 @@ -479,6 +486,7 @@ Kibana 使用 `trace.id` 查询日志,SkyWalking 使用同一个 Trace ID 查 - [ ] Request ID 校验、响应和审计关联测试通过。 - [ ] 日志字段符合约定,且不输出完整 MDC。 - [ ] Spring 异步执行器上下文传播和隔离测试通过。 +- [ ] Redis Stream 消息上下文传播、重试、Reclaimer 和隔离测试通过。 - [ ] 内外部 HTTP 传播边界测试通过。 - [ ] 无 OTLP endpoint 时不存在外部连接尝试。 - [ ] Collector 中断不影响 SkillHub 业务结果。 diff --git a/docs/observability-decision-map.md b/docs/observability-decision-map.md index 3f9ad918..fc9c039a 100644 --- a/docs/observability-decision-map.md +++ b/docs/observability-decision-map.md @@ -1,14 +1,15 @@ # 通用可观测性决策图 目标:为 SkillHub 建立独立、通用、可插拔的日志关联、指标和链路追踪基础设施。 -当前实现覆盖 HTTP、SkillHub 管理的线程池和明确接入的内部 HTTP Client;定时任务、 -Redis Stream 消费循环和 Reclaimer 是独立后台边界,不继承 HTTP 上下文。 -这些边界不进入业务模型和业务载荷。 +当前实现覆盖 HTTP、SkillHub 管理的线程池、Redis Stream 和明确接入的内部 HTTP Client; +普通定时任务没有上游 carrier,仍是独立后台边界。Redis Stream 上下文只进入 transport +metadata,不进入业务模型和业务载荷。 边界: -- Servlet Filter、执行器装饰器和明确接入的 Client Builder 负责建立/恢复上下文。 -- 定时任务和 Redis Stream 目前只输出自身执行日志,不自动继承请求或 Trace 上下文。 +- Servlet Filter、执行器装饰器、消息 Observation 和明确接入的 Client Builder 负责 + 建立/恢复上下文。 +- Redis Stream Producer 注入、Consumer/Reclaimer 逐条提取;定时任务不继承任意请求。 - 业务代码不读写 MDC,不负责创建通用 Span,也不负责统计任务生命周期指标。 - 使用 W3C Trace Context;日志后端、Metrics 后端和 Trace Exporter 均可替换。 - 上下文是有长度限制的基础设施元数据,不进入业务 payload。 @@ -54,7 +55,8 @@ Type: Research - `requestId` 是 SkillHub 的请求/审计关联标识,不冒充分布式 Trace。 - `traceId/spanId` 由 Tracer 生成,跨进程只使用 W3C `traceparent/tracestate`。 - 定时任务或可靠任务的执行资源 ID 只作为当前执行作用域属性,不进入业务 payload。 -- HTTP 和线程池 carrier 的注入/提取位于基础设施拦截器;持久化任务不携带 HTTP Trace。 +- HTTP、线程池和消息 carrier 的注入/提取位于基础设施层;消息上下文是 transport + metadata,不是任务业务字段。 - 不传播任意 MDC Map;baggage 默认关闭,任何允许项都必须低敏、限长、显式配置。 - 无效或不可信的公网 Trace Context 按 W3C 规则丢弃,服务端控制采样。 @@ -88,9 +90,9 @@ Type: Prototype ### Question -验证线程复用隔离、嵌套作用域、异步任务边界、采样、Exporter 超时、Collector 中断、 -队列打满和关闭观测能力等场景。调度任务和 Redis Stream 的独立后台边界只验证不继承 -请求上下文。 +验证线程复用隔离、嵌套作用域、异步任务边界、消息传播、采样、Exporter 超时、 +Collector 中断、队列打满和关闭观测能力等场景。调度任务验证不继承请求上下文; +Redis Stream 验证逐条注入、提取和清理。 ### Answer @@ -99,10 +101,13 @@ Type: Prototype - Request ID Scope 在线程复用、嵌套 Scope、异常退出和 `CallerRunsPolicy` 下均能恢复并 清理。 - Micrometer 手工 Span 和 Observation 均能随 `skillhubEventExecutor` 传播。 +- Redis Stream Producer/Consumer 通过通用消息 Observation 传播 W3C Trace Context 和 + 受控 Request ID,Reclaimer 从原消息重新提取。 - Scanner 使用 Spring 管理的 `WebClient.Builder` 传播 W3C `traceparent`。 - 面向用户配置的 GitLab 外部 Client 不传播 Trace Context。 - `none / otel-sdk / external-agent` 的应用上下文和 Exporter 条件符合设计。 -- `@Scheduled` 和 Redis Stream/Reclaimer 不继承请求上下文,保持独立后台执行边界。 +- `@Scheduled` 保持独立后台执行边界;Redis Stream/Reclaimer 不继承线程上下文,而是 + 从每条消息的 transport metadata 恢复。 Collector 中断、日志背压、采样率和关闭行为仍由 `big-main` 精确 SHA 镜像的远端原型验证。 diff --git a/docs/observability-developer-guide.md b/docs/observability-developer-guide.md index 205a921c..405105ed 100644 --- a/docs/observability-developer-guide.md +++ b/docs/observability-developer-guide.md @@ -91,17 +91,27 @@ Scanner 是当前已接入的内部客户端。新增内部客户端时,应补 删除 Header。使用明确不接入 SkillHub Observation 的客户端,并补测试断言请求不包含 `traceparent`。 -### 3.5 定时任务和 Redis Stream +### 3.5 Redis Stream 和定时任务 -当前实现把定时任务、Redis Stream 消费循环和 Reclaimer 视为独立后台执行边界: +Redis Stream 已通过 `MessageObservationSupport` 接入通用消息传播: -- 不继承任意 HTTP 请求的 `request.id` 或 `trace.id`; -- 不把 HTTP Trace Context 写入 Redis 业务载荷; -- 日志仍可使用 ECS 格式和固定服务字段; -- 若未来需要任务级关联,应增加独立的任务执行 ID/Observation carrier,并单独设计 - 持久化与重试语义。 +- Producer 把 `traceparent`、`tracestate` 和受控的 `skillhub.request_id` 写入 Stream + transport metadata,不修改 `ScanTask` 等业务对象; +- `AbstractStreamConsumer` 逐条提取上下文并建立 `CONSUMER` Observation,在 `finally` + 中恢复线程原状态; +- Consumer 内部调用 Scanner 时,Spring 管理的 `WebClient` 自动创建同一 Trace 的子 Span; +- 重试发布发生在当前 Consumer Scope 内,新消息继续携带关联上下文;Reclaimer 处理原消息 + 时重新从消息提取,不继承 Reclaimer 线程的上下文; +- `none` 和 `external-agent` 模式仍传播 Request ID;应用保证完整 W3C Trace 的模式是 + `otel-sdk`,外部 Agent 的跨 Stream Trace 能力取决于对应 Agent 插件。 -因此,不要假设在 `@Scheduled` 或 Stream consumer 中能自动查到发起 HTTP 请求的 Trace。 +新增 Redis Stream Consumer 应继承 `AbstractStreamConsumer`,新增 Producer 应调用 +`MessageObservationSupport.observePublish`。其他消息中间件只实现自身 carrier 的 +`MessageCarrierAdapter`;传播核心不依赖 Redis、Redisson 或 `Map`。不要在业务 DTO、MDC +或日志代码中复制上下文。 + +普通 `@Scheduled` 任务没有上游消息 carrier,仍是独立后台边界;需要长期任务关联时应使用 +稳定任务 ID,而不是把任意历史 HTTP Span 保持为超长父 Span。 ## 4. 可扩展点 @@ -113,6 +123,7 @@ Scanner 是当前已接入的内部客户端。新增内部客户端时,应补 | 新增日志字段 | `SkillHubEcsEncoder` 白名单 | “输出全部 MDC” | | 新增内部 HTTP 客户端 | Spring `WebClient.Builder` + propagation test | URL 正则删 Header | | 新增外部 HTTP 客户端 | 独立客户端构建入口 + no-propagation test | 依赖全局默认行为 | +| 新增消息队列边界 | `MessageObservationSupport` + `MessageCarrierAdapter` | 业务 DTO、手工 MDC/OTel API | ## 5. 接入验收清单 @@ -124,6 +135,7 @@ Scanner 是当前已接入的内部客户端。新增内部客户端时,应补 4. 线程复用后上下文被清理,不发生串号; 5. 日志只出现 `request.id`、`trace.id`、`span.id` 等白名单字段; 6. Collector 不可用时不影响业务结果。 +7. 消息 Producer/Consumer 使用同一 Trace,Request ID 不串号,重试和 Reclaimer 不丢关联。 运行后端验证使用: diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/config/RedisStreamConfig.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/config/RedisStreamConfig.java index 8c479cfe..87ca0762 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/config/RedisStreamConfig.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/config/RedisStreamConfig.java @@ -4,6 +4,7 @@ import com.iflytek.skillhub.domain.security.ScanTaskProducer; import com.iflytek.skillhub.domain.security.SecurityScanService; import com.iflytek.skillhub.domain.security.SecurityScanner; import com.iflytek.skillhub.domain.skill.SkillVersionRepository; +import com.iflytek.skillhub.observability.MessageObservationSupport; import com.iflytek.skillhub.storage.ObjectStorageService; import com.iflytek.skillhub.stream.RedissonScanTaskProducer; import com.iflytek.skillhub.stream.ScanTaskConsumer; @@ -37,8 +38,11 @@ public class RedisStreamConfig { private Duration reclaimInterval; @Bean - public RedissonScanTaskProducer redisScanTaskProducer(RedissonClient redissonClient) { - return new RedissonScanTaskProducer(redissonClient, streamKey); + public RedissonScanTaskProducer redisScanTaskProducer( + RedissonClient redissonClient, + MessageObservationSupport messageObservationSupport + ) { + return new RedissonScanTaskProducer(redissonClient, streamKey, messageObservationSupport); } @Bean @@ -47,7 +51,8 @@ public class RedisStreamConfig { SecurityScanService securityScanService, SkillVersionRepository skillVersionRepository, ScanTaskProducer scanTaskProducer, - ObjectStorageService objectStorageService) { + ObjectStorageService objectStorageService, + MessageObservationSupport messageObservationSupport) { return new ScanTaskConsumer( redissonClient, streamKey, @@ -60,7 +65,8 @@ public class RedisStreamConfig { reclaimEnabled, reclaimMinIdle, reclaimBatchSize, - reclaimInterval + reclaimInterval, + messageObservationSupport ); } } diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestIdFilter.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestIdFilter.java index f5b88bc3..d3122e44 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestIdFilter.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestIdFilter.java @@ -12,7 +12,6 @@ import org.springframework.web.filter.OncePerRequestFilter; import java.io.IOException; import java.util.UUID; -import java.util.regex.Pattern; /** * Ensures every request has a request identifier for logs, responses, and downstream audit @@ -23,8 +22,6 @@ import java.util.regex.Pattern; public class RequestIdFilter extends OncePerRequestFilter { private static final String REQUEST_ID_HEADER = "X-Request-Id"; - private static final Pattern VALID_REQUEST_ID = - Pattern.compile("^[A-Za-z0-9][A-Za-z0-9._:-]{0,63}$"); private final RequestIdAccessor requestIdAccessor; @@ -36,7 +33,7 @@ public class RequestIdFilter extends OncePerRequestFilter { protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain filterChain) throws ServletException, IOException { String requestId = request.getHeader(REQUEST_ID_HEADER); - if (requestId == null || !VALID_REQUEST_ID.matcher(requestId).matches()) { + if (!RequestIdAccessor.isValid(requestId)) { requestId = UUID.randomUUID().toString(); } 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 new file mode 100644 index 00000000..2fb26318 --- /dev/null +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageCarrierAdapter.java @@ -0,0 +1,15 @@ +package com.iflytek.skillhub.observability; + +import io.micrometer.observation.transport.Propagator; + +/** + * Adapts transport-specific message headers to the common observation boundary. + */ +public interface MessageCarrierAdapter + extends Propagator.Getter, Propagator.Setter { + + /** + * Removes every value associated with a transport header. + */ + 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 new file mode 100644 index 00000000..4e84e1bd --- /dev/null +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/MessageObservationSupport.java @@ -0,0 +1,161 @@ +package com.iflytek.skillhub.observability; + +import io.micrometer.observation.Observation; +import io.micrometer.observation.ObservationRegistry; +import io.micrometer.observation.transport.Kind; +import io.micrometer.observation.transport.ReceiverContext; +import io.micrometer.observation.transport.SenderContext; +import org.springframework.stereotype.Component; + +import java.util.Objects; +import java.util.Set; +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.

+ */ +@Component +public class MessageObservationSupport { + + public static final String REQUEST_ID_FIELD = "skillhub.request_id"; + + private static final Set OWNED_TRANSPORT_FIELDS = Set.of( + REQUEST_ID_FIELD, + "traceparent", + "tracestate", + "baggage" + ); + + private final ObservationRegistry observationRegistry; + private final RequestIdAccessor requestIdAccessor; + + public MessageObservationSupport( + ObservationRegistry observationRegistry, + RequestIdAccessor requestIdAccessor + ) { + this.observationRegistry = observationRegistry; + this.requestIdAccessor = requestIdAccessor; + } + + /** + * Observes a message publish operation and injects the current transport context. + */ + public T observePublish( + String messagingSystem, + String destination, + C carrier, + MessageCarrierAdapter carrierAdapter, + Supplier action + ) { + validateArguments(messagingSystem, destination, carrier, action); + Objects.requireNonNull(carrierAdapter, "carrierAdapter must not be null"); + OWNED_TRANSPORT_FIELDS.forEach(field -> carrierAdapter.remove(carrier, field)); + String requestId = requestIdAccessor.current(); + if (RequestIdAccessor.isValid(requestId)) { + carrierAdapter.set(carrier, REQUEST_ID_FIELD, requestId); + } + + SenderContext senderContext = new SenderContext<>(carrierAdapter, Kind.PRODUCER); + senderContext.setCarrier(carrier); + senderContext.setRemoteServiceName(messagingSystem); + return observe( + "skillhub.message.publish", + destination + " publish", + messagingSystem, + destination, + "publish", + senderContext, + action + ); + } + + /** + * Extracts a message transport context and observes processing inside its scope. + */ + public T observeProcess( + String messagingSystem, + String destination, + C carrier, + MessageCarrierAdapter carrierAdapter, + Supplier action + ) { + validateArguments(messagingSystem, destination, carrier, action); + Objects.requireNonNull(carrierAdapter, "carrierAdapter must not be null"); + ReceiverContext receiverContext = new ReceiverContext<>(carrierAdapter, Kind.CONSUMER); + receiverContext.setCarrier(carrier); + receiverContext.setRemoteServiceName(messagingSystem); + + String propagatedRequestId = carrierAdapter.get(carrier, REQUEST_ID_FIELD); + RequestIdAccessor.Scope requestIdScope = requestIdAccessor.openNullable( + RequestIdAccessor.isValid(propagatedRequestId) ? propagatedRequestId : null + ); + try (requestIdScope) { + return observe( + "skillhub.message.process", + destination + " process", + messagingSystem, + destination, + "process", + receiverContext, + action + ); + } + } + + /** + * Marks the currently active message Observation as failed when processing handles the + * exception without rethrowing it. + */ + public void recordCurrentError(Throwable error) { + Objects.requireNonNull(error, "error must not be null"); + Observation currentObservation = observationRegistry.getCurrentObservation(); + if (currentObservation != null) { + currentObservation.error(error); + } + } + + private T observe( + String observationName, + String contextualName, + String messagingSystem, + String destination, + String operation, + Observation.Context transportContext, + Supplier action + ) { + Observation observation = Observation + .createNotStarted(observationName, () -> transportContext, observationRegistry) + .contextualName(contextualName) + .lowCardinalityKeyValue("messaging.system", messagingSystem) + .lowCardinalityKeyValue("messaging.operation.type", operation) + .highCardinalityKeyValue("messaging.destination.name", destination) + .start(); + try (Observation.Scope ignored = observation.openScope()) { + return action.get(); + } catch (RuntimeException | Error error) { + observation.error(error); + throw error; + } finally { + observation.stop(); + } + } + + private void validateArguments( + String messagingSystem, + String destination, + Object carrier, + Supplier action + ) { + if (messagingSystem == null || messagingSystem.isBlank()) { + throw new IllegalArgumentException("messagingSystem must not be blank"); + } + if (destination == null || destination.isBlank()) { + throw new IllegalArgumentException("destination must not be blank"); + } + Objects.requireNonNull(carrier, "carrier must not be null"); + Objects.requireNonNull(action, "action must not be null"); + } +} diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/RequestIdAccessor.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/RequestIdAccessor.java index ecdfb803..d6589b88 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/RequestIdAccessor.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/RequestIdAccessor.java @@ -4,6 +4,7 @@ import org.slf4j.MDC; import org.springframework.stereotype.Component; import java.util.Objects; +import java.util.regex.Pattern; /** * Holds the current SkillHub request identifier independently from the logging implementation. @@ -16,6 +17,9 @@ public class RequestIdAccessor { public static final String MDC_KEY = "requestId"; + private static final Pattern VALID_REQUEST_ID = + Pattern.compile("^[A-Za-z0-9][A-Za-z0-9._:-]{0,63}$"); + private final ThreadLocal currentRequestId = new ThreadLocal<>(); /** @@ -25,6 +29,13 @@ public class RequestIdAccessor { return currentRequestId.get(); } + /** + * Returns whether a value is safe to use as a request identifier across transport boundaries. + */ + public static boolean isValid(String requestId) { + return requestId != null && VALID_REQUEST_ID.matcher(requestId).matches(); + } + /** * Opens a nested request identifier scope on the current thread. */ @@ -34,6 +45,10 @@ public class RequestIdAccessor { throw new IllegalArgumentException("requestId must not be blank"); } + return openNullable(requestId); + } + + Scope openNullable(String requestId) { String previousRequestId = currentRequestId.get(); replace(requestId); return new Scope(previousRequestId, requestId); 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 229af851..01a73e77 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 @@ -1,5 +1,6 @@ package com.iflytek.skillhub.stream; +import com.iflytek.skillhub.observability.MessageObservationSupport; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import org.redisson.api.AutoClaimResult; @@ -39,6 +40,7 @@ public abstract class AbstractStreamConsumer { private final Duration reclaimMinIdle; private final int reclaimBatchSize; private final Duration reclaimInterval; + private final MessageObservationSupport messageObservationSupport; private final AtomicBoolean running = new AtomicBoolean(false); private RStream stream; @@ -47,8 +49,18 @@ public abstract class AbstractStreamConsumer { protected AbstractStreamConsumer(RedissonClient redissonClient, String streamKey, - String groupName) { - this(redissonClient, streamKey, groupName, true, Duration.ofMinutes(2), 20, Duration.ofSeconds(30)); + String groupName, + MessageObservationSupport messageObservationSupport) { + this( + redissonClient, + streamKey, + groupName, + true, + Duration.ofMinutes(2), + 20, + Duration.ofSeconds(30), + messageObservationSupport + ); } protected AbstractStreamConsumer(RedissonClient redissonClient, @@ -57,7 +69,8 @@ public abstract class AbstractStreamConsumer { boolean reclaimEnabled, Duration reclaimMinIdle, int reclaimBatchSize, - Duration reclaimInterval) { + Duration reclaimInterval, + MessageObservationSupport messageObservationSupport) { this.redissonClient = redissonClient; this.streamKey = streamKey; this.groupName = groupName; @@ -65,6 +78,7 @@ public abstract class AbstractStreamConsumer { this.reclaimMinIdle = reclaimMinIdle; this.reclaimBatchSize = reclaimBatchSize; this.reclaimInterval = reclaimInterval; + this.messageObservationSupport = messageObservationSupport; this.consumerName = consumerPrefix() + "-" + UUID.randomUUID().toString().substring(0, 8); } @@ -195,6 +209,19 @@ public abstract class AbstractStreamConsumer { } void handleMessage(StreamMessageId messageId, Map data) { + messageObservationSupport.observeProcess( + RedisStreamMessageCarrier.MESSAGING_SYSTEM, + streamKey, + data, + RedisStreamMessageCarrier.ADAPTER, + () -> { + handleMessageInScope(messageId, data); + return null; + } + ); + } + + private void handleMessageInScope(StreamMessageId messageId, Map data) { T payload = parsePayload(messageId.toString(), data); if (payload == null) { acknowledge(messageId); @@ -208,6 +235,7 @@ public abstract class AbstractStreamConsumer { markCompleted(payload); acknowledge(messageId); } catch (Exception e) { + messageObservationSupport.recordCurrentError(e); handleFailure(payload, retryCount, e); acknowledge(messageId); } 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 new file mode 100644 index 00000000..acbb349b --- /dev/null +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/RedisStreamMessageCarrier.java @@ -0,0 +1,34 @@ +package com.iflytek.skillhub.stream; + +import com.iflytek.skillhub.observability.MessageCarrierAdapter; + +import java.util.Map; + +/** + * Adapts Redis Stream field maps to the transport-neutral message observation boundary. + */ +final class RedisStreamMessageCarrier { + + static final String MESSAGING_SYSTEM = "redis"; + + static final MessageCarrierAdapter> ADAPTER = + new MessageCarrierAdapter<>() { + @Override + public String get(Map carrier, String key) { + return carrier.get(key); + } + + @Override + public void set(Map carrier, String key, String value) { + carrier.put(key, value); + } + + @Override + public void remove(Map carrier, String key) { + carrier.remove(key); + } + }; + + private 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 047215d3..144693f1 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 @@ -2,6 +2,7 @@ package com.iflytek.skillhub.stream; import com.iflytek.skillhub.domain.security.ScanTask; import com.iflytek.skillhub.domain.security.ScanTaskProducer; +import com.iflytek.skillhub.observability.MessageObservationSupport; import org.redisson.api.RStream; import org.redisson.api.RedissonClient; import org.redisson.api.StreamMessageId; @@ -19,10 +20,16 @@ public class RedissonScanTaskProducer implements ScanTaskProducer { private final RedissonClient redissonClient; private final String streamKey; + private final MessageObservationSupport messageObservationSupport; - public RedissonScanTaskProducer(RedissonClient redissonClient, String streamKey) { + public RedissonScanTaskProducer( + RedissonClient redissonClient, + String streamKey, + MessageObservationSupport messageObservationSupport + ) { this.redissonClient = redissonClient; this.streamKey = streamKey; + this.messageObservationSupport = messageObservationSupport; } @Override @@ -43,7 +50,13 @@ public class RedissonScanTaskProducer implements ScanTaskProducer { } RStream stream = redissonClient.getStream(streamKey, StringCodec.INSTANCE); - StreamMessageId messageId = stream.add(StreamAddArgs.entries(fields)); + StreamMessageId messageId = messageObservationSupport.observePublish( + RedisStreamMessageCarrier.MESSAGING_SYSTEM, + streamKey, + fields, + RedisStreamMessageCarrier.ADAPTER, + () -> stream.add(StreamAddArgs.entries(fields)) + ); log.info("Published scan task: taskId={}, versionId={}, bundleKey={}, hasSkillPath={}, recordId={}", task.taskId(), task.versionId(), task.bundleKey(), task.skillPath() != null && !task.skillPath().isBlank(), messageId); } diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java index 19722d87..fce9bcc2 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/stream/ScanTaskConsumer.java @@ -9,6 +9,7 @@ import com.iflytek.skillhub.domain.security.SecurityScanService; import com.iflytek.skillhub.domain.security.SecurityScanner; import com.iflytek.skillhub.domain.skill.SkillVersionRepository; import com.iflytek.skillhub.domain.skill.SkillVersionStatus; +import com.iflytek.skillhub.observability.MessageObservationSupport; import com.iflytek.skillhub.storage.ObjectStorageService; import org.redisson.api.RedissonClient; @@ -38,8 +39,9 @@ public class ScanTaskConsumer extends AbstractStreamConsumer { + ObservationRegistry observationRegistry = context.getBean(ObservationRegistry.class); + RequestIdAccessor requestIdAccessor = context.getBean(RequestIdAccessor.class); + Tracer tracer = context.getBean(Tracer.class); + MessageObservationSupport support = new MessageObservationSupport( + observationRegistry, + requestIdAccessor + ); + TestCarrier carrier = new TestCarrier(); + Observation parent = Observation.start("publish-request", observationRegistry); + String parentTraceId; + String producerSpanId; + + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("request-async-1"); + Observation.Scope observationScope = parent.openScope()) { + parentTraceId = tracer.currentSpan().context().traceId(); + producerSpanId = support.observePublish( + "redis", + "skillhub:scan:requests", + carrier, + TEST_CARRIER_ADAPTER, + () -> tracer.currentSpan().context().spanId() + ); + } finally { + parent.stop(); + } + + assertThat(carrier.get(MessageObservationSupport.REQUEST_ID_FIELD)) + .isEqualTo("request-async-1"); + assertThat(carrier.get("traceparent")) + .matches("^00-" + parentTraceId + "-[0-9a-f]{16}-0[01]$"); + + ExecutorService worker = Executors.newSingleThreadExecutor(); + try { + ContextValues consumed = worker.submit(() -> support.observeProcess( + "redis", + "skillhub:scan:requests", + carrier, + TEST_CARRIER_ADAPTER, + () -> currentValues(requestIdAccessor, tracer) + )).get(5, TimeUnit.SECONDS); + + assertThat(consumed.requestId()).isEqualTo("request-async-1"); + assertThat(consumed.mdcRequestId()).isEqualTo("request-async-1"); + assertThat(consumed.traceId()).isEqualTo(parentTraceId); + assertThat(consumed.spanId()).isNotEqualTo(producerSpanId); + + ContextValues clean = worker.submit( + () -> currentValues(requestIdAccessor, tracer) + ).get(5, TimeUnit.SECONDS); + assertThat(clean.requestId()).isNull(); + assertThat(clean.mdcRequestId()).isNull(); + assertThat(clean.traceId()).isNull(); + assertThat(clean.spanId()).isNull(); + } finally { + worker.shutdownNow(); + MDC.clear(); + } + }); + } + + @Test + void shouldClearMissingOrInvalidMessageContextAndRestoreOuterScope() { + RequestIdAccessor requestIdAccessor = new RequestIdAccessor(); + MessageObservationSupport support = new MessageObservationSupport( + ObservationRegistry.NOOP, + requestIdAccessor + ); + TestCarrier carrier = new TestCarrier(); + carrier.set(MessageObservationSupport.REQUEST_ID_FIELD, "invalid request id"); + + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("outer-request")) { + String valueInsideMessage = support.observeProcess( + "test-broker", + "jobs", + carrier, + TEST_CARRIER_ADAPTER, + requestIdAccessor::current + ); + + assertThat(valueInsideMessage).isNull(); + assertThat(requestIdAccessor.current()).isEqualTo("outer-request"); + } finally { + MDC.clear(); + } + } + + @Test + void shouldRemoveCallerSuppliedTransportContextBeforePublishingWithoutTracer() { + RequestIdAccessor requestIdAccessor = new RequestIdAccessor(); + MessageObservationSupport support = new MessageObservationSupport( + ObservationRegistry.NOOP, + requestIdAccessor + ); + TestCarrier carrier = new TestCarrier(); + carrier.set("traceparent", "caller-controlled"); + carrier.set(MessageObservationSupport.REQUEST_ID_FIELD, "caller-controlled"); + + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("trusted-request")) { + support.observePublish( + "test-broker", + "jobs", + carrier, + TEST_CARRIER_ADAPTER, + () -> null + ); + } finally { + MDC.clear(); + } + + assertThat(carrier.get("traceparent")).isNull(); + assertThat(carrier.get(MessageObservationSupport.REQUEST_ID_FIELD)) + .isEqualTo("trusted-request"); + } + + private ContextValues currentValues(RequestIdAccessor requestIdAccessor, Tracer tracer) { + Span currentSpan = tracer.currentSpan(); + return new ContextValues( + requestIdAccessor.current(), + MDC.get(RequestIdAccessor.MDC_KEY), + currentSpan == null ? null : currentSpan.context().traceId(), + currentSpan == null ? null : currentSpan.context().spanId() + ); + } + + private record ContextValues( + String requestId, + String mdcRequestId, + String traceId, + String spanId + ) { + } + + private static final MessageCarrierAdapter TEST_CARRIER_ADAPTER = + new MessageCarrierAdapter<>() { + @Override + public String get(TestCarrier carrier, String key) { + return carrier.get(key); + } + + @Override + public void set(TestCarrier carrier, String key, String value) { + carrier.set(key, value); + } + + @Override + public void remove(TestCarrier carrier, String key) { + carrier.set(key, null); + } + }; + + private static final class TestCarrier { + private final List
headers = new ArrayList<>(); + + private void set(String key, String value) { + headers.removeIf(header -> header.key().equals(key)); + if (value != null) { + headers.add(new Header(key, value)); + } + } + + private String get(String key) { + return headers.stream() + .filter(header -> header.key().equals(key)) + .map(Header::value) + .findFirst() + .orElse(null); + } + } + + private record Header(String key, String value) { + } + + @Configuration(proxyBeanMethods = false) + @EnableAutoConfiguration + @Import({SkillHubTracingConfiguration.class, RequestIdAccessor.class}) + static class TestApplication { + } +} diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/AbstractStreamConsumerTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/AbstractStreamConsumerTest.java index 6a5436f0..bbb58a70 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/AbstractStreamConsumerTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/AbstractStreamConsumerTest.java @@ -1,5 +1,8 @@ package com.iflytek.skillhub.stream; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; +import io.micrometer.observation.ObservationRegistry; import io.lettuce.core.RedisBusyException; import org.junit.jupiter.api.Test; import org.redisson.api.AutoClaimResult; @@ -104,6 +107,25 @@ class AbstractStreamConsumerTest { assertThat(consumer.streamCreationCount.get()).isEqualTo(1); } + @Test + void handleMessage_restoresPropagatedRequestIdOnlyWhileProcessing() { + @SuppressWarnings("unchecked") + RStream stream = mock(RStream.class); + RequestIdAccessor requestIdAccessor = new RequestIdAccessor(); + TestConsumer consumer = new TestConsumer(stream, requestIdAccessor); + + consumer.handleMessage( + new StreamMessageId(8, 0), + Map.of( + "payload", "correlated", + MessageObservationSupport.REQUEST_ID_FIELD, "request-stream-1" + ) + ); + + assertThat(consumer.processedRequestId).isEqualTo("request-stream-1"); + assertThat(requestIdAccessor.current()).isNull(); + } + @Test void detectsBusyGroupWhenWrappedInRedisSystemException() { RedisSystemException wrapped = new RedisSystemException( @@ -116,11 +138,27 @@ class AbstractStreamConsumerTest { private static class TestConsumer extends AbstractStreamConsumer { private final RStream stream; + private final RequestIdAccessor requestIdAccessor; private boolean fail; + private String processedRequestId; private TestConsumer(RStream stream) { - super(mock(RedissonClient.class), "scan-stream", "scan-group", true, Duration.ofMinutes(2), 20, Duration.ofSeconds(30)); + this(stream, new RequestIdAccessor()); + } + + private TestConsumer(RStream stream, RequestIdAccessor requestIdAccessor) { + super( + mock(RedissonClient.class), + "scan-stream", + "scan-group", + true, + Duration.ofMinutes(2), + 20, + Duration.ofSeconds(30), + new MessageObservationSupport(ObservationRegistry.NOOP, requestIdAccessor) + ); this.stream = stream; + this.requestIdAccessor = requestIdAccessor; } @Override @@ -154,6 +192,7 @@ class AbstractStreamConsumerTest { @Override protected void processBusiness(String payload) { + processedRequestId = requestIdAccessor.current(); if (fail) { throw new IllegalStateException("boom"); } diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerLoggingTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerLoggingTest.java index 542908df..a184aa30 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerLoggingTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerLoggingTest.java @@ -5,6 +5,9 @@ import ch.qos.logback.classic.Logger; import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.core.read.ListAppender; import com.iflytek.skillhub.domain.security.ScanTask; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; +import io.micrometer.observation.ObservationRegistry; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; import org.redisson.api.RStream; @@ -43,7 +46,11 @@ class RedissonScanTaskProducerLoggingTest { RedissonClient redissonClient = mock(RedissonClient.class); doReturn(typedStream).when(redissonClient).getStream("skillhub:scan:requests", StringCodec.INSTANCE); when(stream.add(any())).thenReturn(new StreamMessageId(1, 0)); - RedissonScanTaskProducer producer = new RedissonScanTaskProducer(redissonClient, "skillhub:scan:requests"); + RedissonScanTaskProducer producer = new RedissonScanTaskProducer( + redissonClient, + "skillhub:scan:requests", + new MessageObservationSupport(ObservationRegistry.NOOP, new RequestIdAccessor()) + ); attachAppender(); producer.publishScanTask(new ScanTask( diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerTest.java index 003c1158..ccf22ea3 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/RedissonScanTaskProducerTest.java @@ -1,12 +1,16 @@ package com.iflytek.skillhub.stream; import com.iflytek.skillhub.domain.security.ScanTask; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; +import io.micrometer.observation.ObservationRegistry; import org.junit.jupiter.api.Test; import org.redisson.api.RStream; import org.redisson.api.RedissonClient; import org.redisson.api.StreamMessageId; import org.redisson.client.codec.StringCodec; import org.redisson.api.stream.StreamAddArgs; +import org.redisson.api.stream.StreamAddParams; import org.mockito.ArgumentCaptor; import java.util.Map; @@ -29,21 +33,32 @@ class RedissonScanTaskProducerTest { RedissonClient redissonClient = mock(RedissonClient.class); doReturn(typedStream).when(redissonClient).getStream("skillhub:scan:requests", StringCodec.INSTANCE); when(stream.add(any())).thenReturn(new StreamMessageId(1, 0)); - RedissonScanTaskProducer producer = new RedissonScanTaskProducer(redissonClient, "skillhub:scan:requests"); + RequestIdAccessor requestIdAccessor = new RequestIdAccessor(); + RedissonScanTaskProducer producer = new RedissonScanTaskProducer( + redissonClient, + "skillhub:scan:requests", + new MessageObservationSupport(ObservationRegistry.NOOP, requestIdAccessor) + ); - producer.publishScanTask(new ScanTask( - "task-1", - 42L, - "/tmp/skill", - null, - "publisher-1", - 1711260000000L, - Map.of("scannerType", "skill-scanner") - )); + try (RequestIdAccessor.Scope ignored = requestIdAccessor.open("request-stream-1")) { + producer.publishScanTask(new ScanTask( + "task-1", + 42L, + "/tmp/skill", + null, + "publisher-1", + 1711260000000L, + Map.of("scannerType", "skill-scanner") + )); + } verify(redissonClient).getStream("skillhub:scan:requests", StringCodec.INSTANCE); ArgumentCaptor> argsCaptor = ArgumentCaptor.forClass(StreamAddArgs.class); verify(stream).add(argsCaptor.capture()); - assertThat(argsCaptor.getValue()).isNotNull(); + assertThat(argsCaptor.getValue()).isInstanceOf(StreamAddParams.class); + @SuppressWarnings("unchecked") + StreamAddParams params = (StreamAddParams) argsCaptor.getValue(); + assertThat(params.getEntries()) + .containsEntry(MessageObservationSupport.REQUEST_ID_FIELD, "request-stream-1"); } } diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerLoggingTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerLoggingTest.java index 293a3dd7..8f473b65 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerLoggingTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerLoggingTest.java @@ -13,8 +13,11 @@ import com.iflytek.skillhub.domain.security.SecurityScanner; import com.iflytek.skillhub.domain.skill.SkillVersion; import com.iflytek.skillhub.domain.skill.SkillVersionRepository; import com.iflytek.skillhub.domain.skill.SkillVersionStatus; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; import com.iflytek.skillhub.storage.ObjectMetadata; import com.iflytek.skillhub.storage.ObjectStorageService; +import io.micrometer.observation.ObservationRegistry; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; import org.redisson.api.RStream; @@ -163,7 +166,8 @@ class ScanTaskConsumerLoggingTest { securityScanService, skillVersionRepository, scanTaskProducer, - objectStorageService + objectStorageService, + new MessageObservationSupport(ObservationRegistry.NOOP, new RequestIdAccessor()) ); } diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerPathSafetyTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerPathSafetyTest.java index 984acbd2..914ffd2f 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerPathSafetyTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerPathSafetyTest.java @@ -6,7 +6,10 @@ import com.iflytek.skillhub.domain.security.ScanTaskProducer; import com.iflytek.skillhub.domain.security.SecurityScanService; import com.iflytek.skillhub.domain.security.SecurityScanner; import com.iflytek.skillhub.domain.skill.SkillVersionRepository; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; import com.iflytek.skillhub.storage.ObjectStorageService; +import io.micrometer.observation.ObservationRegistry; import java.lang.reflect.Method; import java.nio.file.Files; import java.nio.file.Path; @@ -25,7 +28,8 @@ class ScanTaskConsumerPathSafetyTest { org.mockito.Mockito.mock(SecurityScanService.class), org.mockito.Mockito.mock(SkillVersionRepository.class), org.mockito.Mockito.mock(ScanTaskProducer.class), - org.mockito.Mockito.mock(ObjectStorageService.class) + org.mockito.Mockito.mock(ObjectStorageService.class), + new MessageObservationSupport(ObservationRegistry.NOOP, new RequestIdAccessor()) ); Path outsideFile = Files.createTempFile("scan-cleanup-", ".txt"); Files.writeString(outsideFile, "keep"); diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerTest.java index c8c55f0f..7f8cabaa 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/stream/ScanTaskConsumerTest.java @@ -14,8 +14,11 @@ import com.iflytek.skillhub.domain.security.SecurityVerdict; import com.iflytek.skillhub.domain.skill.SkillVersion; import com.iflytek.skillhub.domain.skill.SkillVersionRepository; import com.iflytek.skillhub.domain.skill.SkillVersionStatus; +import com.iflytek.skillhub.observability.MessageObservationSupport; +import com.iflytek.skillhub.observability.RequestIdAccessor; import com.iflytek.skillhub.storage.ObjectStorageService; import com.iflytek.skillhub.storage.ObjectMetadata; +import io.micrometer.observation.ObservationRegistry; import org.junit.jupiter.api.Test; import org.redisson.api.RStream; import org.redisson.api.RedissonClient; @@ -303,7 +306,8 @@ class ScanTaskConsumerTest { securityScanService, skillVersionRepository, scanTaskProducer, - objectStorageService + objectStorageService, + new MessageObservationSupport(ObservationRegistry.NOOP, new RequestIdAccessor()) ); this.stream = mock(RStream.class); } 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 3/5] 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( From 41e0776b4c2815b832d475440d31a9c99c1b6e69 Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Mon, 3 Aug 2026 16:36:52 +0800 Subject: [PATCH 4/5] test(auth): isolate security context between tests Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- .../skillhub/auth/oauth/OAuth2LoginHandlersTest.java | 7 +++++++ .../auth/token/ApiTokenAuthenticationFilterTest.java | 6 ++++++ 2 files changed, 13 insertions(+) diff --git a/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/oauth/OAuth2LoginHandlersTest.java b/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/oauth/OAuth2LoginHandlersTest.java index 52c0077b..750d29d3 100644 --- a/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/oauth/OAuth2LoginHandlersTest.java +++ b/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/oauth/OAuth2LoginHandlersTest.java @@ -1,12 +1,14 @@ package com.iflytek.skillhub.auth.oauth; import jakarta.servlet.http.HttpSession; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; import org.springframework.mock.web.MockHttpServletRequest; import org.springframework.mock.web.MockHttpServletResponse; import org.springframework.security.authentication.UsernamePasswordAuthenticationToken; import org.springframework.security.core.Authentication; import org.springframework.security.core.context.SecurityContext; +import org.springframework.security.core.context.SecurityContextHolder; import org.springframework.security.oauth2.core.OAuth2AuthenticationException; import org.springframework.security.oauth2.core.OAuth2Error; import org.springframework.security.oauth2.core.user.DefaultOAuth2User; @@ -22,6 +24,11 @@ import static org.mockito.Mockito.mock; class OAuth2LoginHandlersTest { + @AfterEach + void clearSecurityContext() { + SecurityContextHolder.clearContext(); + } + @Test void successHandler_redirectsToStoredReturnTo() throws Exception { OAuthLoginFlowService oauthLoginFlowService = mock(OAuthLoginFlowService.class); diff --git a/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/token/ApiTokenAuthenticationFilterTest.java b/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/token/ApiTokenAuthenticationFilterTest.java index d9f030a9..fa9fb2cf 100644 --- a/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/token/ApiTokenAuthenticationFilterTest.java +++ b/server/skillhub-auth/src/test/java/com/iflytek/skillhub/auth/token/ApiTokenAuthenticationFilterTest.java @@ -10,6 +10,7 @@ import com.iflytek.skillhub.domain.user.UserAccount; import com.iflytek.skillhub.domain.user.UserAccountRepository; import com.iflytek.skillhub.domain.user.UserStatus; import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.mock.web.MockFilterChain; import org.springframework.mock.web.MockHttpServletRequest; @@ -43,6 +44,11 @@ class ApiTokenAuthenticationFilterTest { scopeService ); + @BeforeEach + void initializeSecurityContext() { + SecurityContextHolder.clearContext(); + } + @AfterEach void clearSecurityContext() { SecurityContextHolder.clearContext(); From 9290fe3ca222964335f94949506a40849fffc5aa Mon Sep 17 00:00:00 2001 From: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> Date: Mon, 3 Aug 2026 19:25:19 +0800 Subject: [PATCH 5/5] fix(observability): skip otlp exporter without endpoint Signed-off-by: XiaoSeS <87064762+XiaoSeS@users.noreply.github.com> --- ...cingModeAutoConfigurationImportFilter.java | 17 ++++++++++--- .../SkillHubTracingConfigurationTest.java | 18 ++++++++++++++ ...ModeAutoConfigurationImportFilterTest.java | 24 +++++++++++++++++++ 3 files changed, 56 insertions(+), 3 deletions(-) diff --git a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilter.java b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilter.java index 55a33fda..d2aa538a 100644 --- a/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilter.java +++ b/server/skillhub-app/src/main/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilter.java @@ -19,7 +19,10 @@ public final class TracingModeAutoConfigurationImportFilter private static final Set OTEL_AUTO_CONFIGURATIONS = Set.of( "org.springframework.boot.actuate.autoconfigure.opentelemetry.OpenTelemetryAutoConfiguration", - "org.springframework.boot.actuate.autoconfigure.tracing.OpenTelemetryAutoConfiguration", + "org.springframework.boot.actuate.autoconfigure.tracing.OpenTelemetryAutoConfiguration" + ); + + private static final Set OTLP_EXPORT_AUTO_CONFIGURATIONS = Set.of( "org.springframework.boot.actuate.autoconfigure.tracing.otlp.OtlpAutoConfiguration" ); @@ -34,12 +37,16 @@ public final class TracingModeAutoConfigurationImportFilter && "otel-sdk".equalsIgnoreCase( environment.getProperty(TRACING_MODE_PROPERTY, "none") ); + boolean otlpExportEnabled = otelSdkEnabled + && hasText(environment.getProperty("management.otlp.tracing.endpoint")); boolean[] matches = new boolean[autoConfigurationClasses.length]; for (int index = 0; index < autoConfigurationClasses.length; index++) { String autoConfigurationClass = autoConfigurationClasses[index]; matches[index] = autoConfigurationClass != null - && (otelSdkEnabled - || !OTEL_AUTO_CONFIGURATIONS.contains(autoConfigurationClass)); + && ((otelSdkEnabled + || !OTEL_AUTO_CONFIGURATIONS.contains(autoConfigurationClass)) + && (otlpExportEnabled + || !OTLP_EXPORT_AUTO_CONFIGURATIONS.contains(autoConfigurationClass))); } return matches; } @@ -48,4 +55,8 @@ public final class TracingModeAutoConfigurationImportFilter public void setEnvironment(Environment environment) { this.environment = environment; } + + private static boolean hasText(String value) { + return value != null && !value.isBlank(); + } } diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/SkillHubTracingConfigurationTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/SkillHubTracingConfigurationTest.java index 1c5913ce..99e72a0e 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/SkillHubTracingConfigurationTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/SkillHubTracingConfigurationTest.java @@ -63,6 +63,24 @@ class SkillHubTracingConfigurationTest { }); } + @Test + void otelSdkModeWithEmptyEndpointShouldCreateInProcessTracerOnly() { + contextRunner + .withPropertyValues( + "skillhub.observability.tracing-mode=otel-sdk", + "management.otlp.tracing.endpoint=", + "management.tracing.sampling.probability=1.0", + "management.tracing.baggage.enabled=false", + "management.tracing.propagation.type=W3C" + ) + .run(context -> { + assertThat(context).hasNotFailed(); + assertThat(context.getBean(Tracer.class)).isInstanceOf(OtelTracer.class); + assertThat(context).hasSingleBean(OpenTelemetry.class); + assertThat(context).doesNotHaveBean(OtlpHttpSpanExporter.class); + }); + } + @Test void otelSdkModeShouldCreateExporterOnlyWhenEndpointIsConfigured() { contextRunner diff --git a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilterTest.java b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilterTest.java index 6cfbb7d3..e465010a 100644 --- a/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilterTest.java +++ b/server/skillhub-app/src/test/java/com/iflytek/skillhub/observability/tracing/TracingModeAutoConfigurationImportFilterTest.java @@ -47,9 +47,33 @@ class TracingModeAutoConfigurationImportFilterTest { "otel-sdk" )); + assertThat(matches()).containsExactly(true, true, false, true, false); + } + + @Test + void shouldEnableOtlpExporterOnlyWhenEndpointHasText() { + filter.setEnvironment(new MockEnvironment() + .withProperty( + TracingModeAutoConfigurationImportFilter.TRACING_MODE_PROPERTY, + "otel-sdk" + ) + .withProperty("management.otlp.tracing.endpoint", "http://127.0.0.1:4318/v1/traces")); + assertThat(matches()).containsExactly(true, true, true, true, false); } + @Test + void shouldExcludeOtlpExporterWhenEndpointIsEmpty() { + filter.setEnvironment(new MockEnvironment() + .withProperty( + TracingModeAutoConfigurationImportFilter.TRACING_MODE_PROPERTY, + "otel-sdk" + ) + .withProperty("management.otlp.tracing.endpoint", "")); + + assertThat(matches()).containsExactly(true, true, false, true, false); + } + private boolean[] matches() { return filter.match( new String[]{