Merge PR 664 observability into big-main for validation

# Conflicts:
#	server/skillhub-app/src/main/java/com/iflytek/skillhub/exception/GlobalExceptionHandler.java
#	server/skillhub-app/src/main/java/com/iflytek/skillhub/filter/RequestLoggingFilter.java
#	server/skillhub-app/src/test/java/com/iflytek/skillhub/exception/GlobalExceptionHandlerTest.java
#	server/skillhub-app/src/test/java/com/iflytek/skillhub/filter/RequestLoggingFilterTest.java
This commit is contained in:
XiaoSeS 2026-08-03 19:35:49 +08:00
commit b459b1a82c
28 changed files with 941 additions and 115 deletions

View file

@ -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 业务结果。

View file

@ -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 镜像的远端原型验证。

View file

@ -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 不丢关联。
运行后端验证使用:

View file

@ -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
);
}
}

View file

@ -3,14 +3,13 @@ package com.iflytek.skillhub.exception;
import com.iflytek.skillhub.auth.exception.AuthFlowException;
import com.iflytek.skillhub.auth.config.AccountMergeRouteRequestMatcher;
import com.iflytek.skillhub.auth.config.IdentityLinkRouteRequestMatcher;
import com.iflytek.skillhub.auth.merge.AccountMergeException;
import com.iflytek.skillhub.auth.merge.AccountMergeFailureCode;
import com.iflytek.skillhub.auth.identity.IdentityLinkException;
import com.iflytek.skillhub.auth.identity.IdentityLinkFailureCode;
import com.iflytek.skillhub.auth.rbac.PlatformPrincipal;
import com.iflytek.skillhub.auth.merge.AccountMergeException;
import com.iflytek.skillhub.auth.merge.AccountMergeFailureCode;
import com.iflytek.skillhub.dto.AccountMergeErrorResponse;
import com.iflytek.skillhub.dto.ApiResponse;
import com.iflytek.skillhub.dto.ApiResponseFactory;
import com.iflytek.skillhub.dto.AccountMergeErrorResponse;
import com.iflytek.skillhub.dto.IdentityLinkErrorResponse;
import com.iflytek.skillhub.domain.shared.exception.LocalizedDomainException;
import com.iflytek.skillhub.domain.shared.exception.LocalizedMessage;
@ -244,11 +243,11 @@ public class GlobalExceptionHandler {
public ResponseEntity<ApiResponse<Void>> 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
@ -273,11 +272,11 @@ public class GlobalExceptionHandler {
@ExceptionHandler(Exception.class)
public ResponseEntity<ApiResponse<Void>> 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(
@ -286,12 +285,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
);
}
@ -304,13 +303,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";
}
}

View file

@ -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();
}

View file

@ -17,7 +17,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;
/**
@ -28,18 +27,12 @@ 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<String> SKIP_PREFIXES = Set.of(
"/actuator", "/favicon.ico", "/assets/"
);
private static final Set<String> SKIP_SUFFIXES = Set.of(
"/sse"
);
private static final Set<String> SENSITIVE_BODY_PREFIXES = Set.of(
"/api/v1/auth/",
"/api/v1/account/merge"
);
private final SensitiveLogSanitizer sensitiveLogSanitizer;
public RequestLoggingFilter(
@ -95,13 +88,6 @@ public class RequestLoggingFilter extends OncePerRequestFilter {
sb.append(" | UA: ").append(truncate(userAgent, 80));
}
String requestBody = shouldLogBody(request.getRequestURI())
? getRequestBody(request)
: null;
if (requestBody != null && !requestBody.isBlank()) {
sb.append(" | Body: ").append(requestBody);
}
log.info(sb.toString());
}
@ -119,11 +105,6 @@ public class RequestLoggingFilter extends OncePerRequestFilter {
return false;
}
private boolean shouldLogBody(String uri) {
return SENSITIVE_BODY_PREFIXES.stream()
.noneMatch(uri::startsWith);
}
private boolean isNotificationSse(String uri) {
return uri != null && uri.endsWith("/notifications/sse");
}
@ -134,18 +115,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;

View file

@ -0,0 +1,19 @@
package com.iflytek.skillhub.observability;
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 before trusted context is injected.
*/
void remove(C carrier, String key);
}

View file

@ -0,0 +1,178 @@
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.
*
* <p>The message carrier owns transport metadata. Business payloads remain independent from
* 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",
"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 <C, T> T observePublish(
String messagingSystem,
String destination,
C carrier,
MessageCarrierAdapter<C> carrierAdapter,
Supplier<T> action
) {
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);
}
SenderContext<C> senderContext = new SenderContext<>(carrierAdapter, Kind.PRODUCER);
senderContext.setCarrier(carrier);
senderContext.setRemoteServiceName(messagingSystem);
return observe(
"skillhub.message.publish",
"publish " + destination,
messagingSystem,
destination,
"publish",
"send",
senderContext,
action
);
}
/**
* Extracts a message transport context and observes processing inside its scope.
*/
public <C, T> T observeProcess(
String messagingSystem,
String destination,
C carrier,
MessageCarrierAdapter<C> carrierAdapter,
Supplier<T> action
) {
validateArguments(messagingSystem, destination, carrier, action);
Objects.requireNonNull(carrierAdapter, "carrierAdapter must not be null");
ReceiverContext<C> receiverContext = new ReceiverContext<>(carrierAdapter, Kind.CONSUMER);
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
);
try (requestIdScope) {
return observe(
"skillhub.message.process",
"process " + destination,
messagingSystem,
destination,
"process",
"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> T observe(
String observationName,
String contextualName,
String messagingSystem,
String destination,
String operationName,
String operationType,
Observation.Context transportContext,
Supplier<T> action
) {
Observation observation = Observation
.createNotStarted(observationName, () -> transportContext, observationRegistry)
.contextualName(contextualName)
.lowCardinalityKeyValue("messaging.system", messagingSystem)
.lowCardinalityKeyValue("messaging.operation.name", operationName)
.lowCardinalityKeyValue("messaging.operation.type", operationType)
.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");
}
}

View file

@ -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<String> 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);

View file

@ -19,7 +19,10 @@ public final class TracingModeAutoConfigurationImportFilter
private static final Set<String> 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<String> 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();
}
}

View file

@ -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<T> {
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<String, String> stream;
@ -47,8 +49,18 @@ public abstract class AbstractStreamConsumer<T> {
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<T> {
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<T> {
this.reclaimMinIdle = reclaimMinIdle;
this.reclaimBatchSize = reclaimBatchSize;
this.reclaimInterval = reclaimInterval;
this.messageObservationSupport = messageObservationSupport;
this.consumerName = consumerPrefix() + "-" + UUID.randomUUID().toString().substring(0, 8);
}
@ -191,10 +205,24 @@ 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);
}
void handleMessage(StreamMessageId messageId, Map<String, String> data) {
messageObservationSupport.observeProcess(
RedisStreamMessageCarrier.MESSAGING_SYSTEM,
streamKey,
data,
RedisStreamMessageCarrier.ADAPTER,
() -> {
handleMessageInScope(messageId, data);
return null;
}
);
}
private void handleMessageInScope(StreamMessageId messageId, Map<String, String> data) {
T payload = parsePayload(messageId.toString(), data);
if (payload == null) {
acknowledge(messageId);
@ -208,6 +236,7 @@ public abstract class AbstractStreamConsumer<T> {
markCompleted(payload);
acknowledge(messageId);
} catch (Exception e) {
messageObservationSupport.recordCurrentError(e);
handleFailure(payload, retryCount, e);
acknowledge(messageId);
}
@ -215,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

@ -0,0 +1,38 @@
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.
*
* <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 {
static final String MESSAGING_SYSTEM = "redis";
static final MessageCarrierAdapter<Map<String, String>> ADAPTER =
new MessageCarrierAdapter<>() {
@Override
public String get(Map<String, String> carrier, String key) {
return carrier.get(key);
}
@Override
public void set(Map<String, String> carrier, String key, String value) {
carrier.put(key, value);
}
@Override
public void remove(Map<String, String> carrier, String key) {
carrier.remove(key);
}
};
private RedisStreamMessageCarrier() {
}
}

View file

@ -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,14 @@ public class RedissonScanTaskProducer implements ScanTaskProducer {
}
RStream<String, String> stream = redissonClient.getStream(streamKey, StringCodec.INSTANCE);
StreamMessageId messageId = stream.add(StreamAddArgs.entries(fields));
// Starting the producer Observation injects trusted propagation fields before XADD.
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);
}

View file

@ -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<ScanTaskConsumer.Sc
SecurityScanService securityScanService,
SkillVersionRepository skillVersionRepository,
ScanTaskProducer scanTaskProducer,
ObjectStorageService objectStorageService) {
super(redissonClient, streamKey, groupName);
ObjectStorageService objectStorageService,
MessageObservationSupport messageObservationSupport) {
super(redissonClient, streamKey, groupName, messageObservationSupport);
this.securityScanner = securityScanner;
this.securityScanService = securityScanService;
this.skillVersionRepository = skillVersionRepository;
@ -58,8 +60,18 @@ public class ScanTaskConsumer extends AbstractStreamConsumer<ScanTaskConsumer.Sc
boolean reclaimEnabled,
Duration reclaimMinIdle,
int reclaimBatchSize,
Duration reclaimInterval) {
super(redissonClient, streamKey, groupName, reclaimEnabled, reclaimMinIdle, reclaimBatchSize, reclaimInterval);
Duration reclaimInterval,
MessageObservationSupport messageObservationSupport) {
super(
redissonClient,
streamKey,
groupName,
reclaimEnabled,
reclaimMinIdle,
reclaimBatchSize,
reclaimInterval,
messageObservationSupport
);
this.securityScanner = securityScanner;
this.securityScanService = securityScanService;
this.skillVersionRepository = skillVersionRepository;

View file

@ -0,0 +1,4 @@
/**
* Redis Stream transport adapters and background consumers.
*/
package com.iflytek.skillhub.stream;

View file

@ -4,6 +4,11 @@ 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.dto.IdentityLinkErrorResponse;
@ -11,23 +16,34 @@ import com.iflytek.skillhub.auth.exception.AuthFlowException;
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;
@ -38,6 +54,8 @@ class GlobalExceptionHandlerTest {
private HttpServletRequest request;
private GlobalExceptionHandler handler;
private ListAppender<ILoggingEvent> appender;
private RequestIdAccessor requestIdAccessor;
@BeforeEach
void setUp() {
@ -47,7 +65,7 @@ class GlobalExceptionHandlerTest {
"error.auth.local.invalidCredentials",
java.util.Locale.getDefault(),
"Invalid username or password");
RequestIdAccessor requestIdAccessor = new RequestIdAccessor();
requestIdAccessor = new RequestIdAccessor();
ApiResponseFactory responseFactory = new ApiResponseFactory(
messageSource,
Clock.fixed(Instant.parse("2026-03-20T00:00:00Z"), ZoneOffset.UTC),
@ -61,6 +79,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");
@ -73,6 +99,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");
@ -84,6 +111,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<ApiResponse<Void>> 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<ApiResponse<Void>> 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
@ -135,4 +209,30 @@ class GlobalExceptionHandlerTest {
assertThat(body.reasonCode())
.isEqualTo("REAUTHENTICATION_REQUIRED");
}
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<String> loggedMessages() {
return appender.list.stream()
.map(ILoggingEvent::getFormattedMessage)
.toList();
}
}

View file

@ -39,16 +39,17 @@ class RequestLoggingFilterTest {
}
@Test
void doFilterInternal_truncatesLongRequestBodyAndOmitsResponseBody()
void doFilterInternal_omitsRequestAndResponseBodies()
throws ServletException, IOException {
RequestLoggingFilter filter = filter();
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());
@ -56,17 +57,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<String> 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

View file

@ -0,0 +1,287 @@
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;
import org.junit.jupiter.api.Test;
import org.slf4j.MDC;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
class MessageObservationSupportTest {
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withUserConfiguration(TestApplication.class)
.withPropertyValues(
"spring.flyway.enabled=false",
"spring.jpa.hibernate.ddl-auto=none"
);
@Test
void shouldPropagateTraceAndRequestIdAcrossMessageBoundaryWithoutLeakingWorkerContext() {
contextRunner
.withPropertyValues(
"skillhub.observability.tracing-mode=otel-sdk",
"management.tracing.sampling.probability=1.0"
)
.run(context -> {
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");
}
@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(
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<TestCarrier> 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<Header> 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 {
}
}

View file

@ -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

View file

@ -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[]{

View file

@ -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<String, String> 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<String> {
private final RStream<String, String> stream;
private final RequestIdAccessor requestIdAccessor;
private boolean fail;
private String processedRequestId;
private TestConsumer(RStream<String, String> stream) {
super(mock(RedissonClient.class), "scan-stream", "scan-group", true, Duration.ofMinutes(2), 20, Duration.ofSeconds(30));
this(stream, new RequestIdAccessor());
}
private TestConsumer(RStream<String, String> 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");
}

View file

@ -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(

View file

@ -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<StreamAddArgs<String, String>> argsCaptor = ArgumentCaptor.forClass(StreamAddArgs.class);
verify(stream).add(argsCaptor.capture());
assertThat(argsCaptor.getValue()).isNotNull();
assertThat(argsCaptor.getValue()).isInstanceOf(StreamAddParams.class);
@SuppressWarnings("unchecked")
StreamAddParams<String, String> params = (StreamAddParams<String, String>) argsCaptor.getValue();
assertThat(params.getEntries())
.containsEntry(MessageObservationSupport.REQUEST_ID_FIELD, "request-stream-1");
}
}

View file

@ -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())
);
}

View file

@ -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");

View file

@ -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);
}

View file

@ -7,12 +7,14 @@ import com.iflytek.skillhub.auth.merge.AccountMergeBrowserFlowReference;
import com.iflytek.skillhub.auth.merge.AccountMergeBrowserPhase;
import com.iflytek.skillhub.auth.merge.AccountMergeFailureCode;
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;
@ -30,6 +32,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);