diff --git a/openframe-client-core/src/main/java/com/openframe/client/config/AsyncConfig.java b/openframe-client-core/src/main/java/com/openframe/client/config/AsyncConfig.java
index b0f91526b..4224b7fd6 100644
--- a/openframe-client-core/src/main/java/com/openframe/client/config/AsyncConfig.java
+++ b/openframe-client-core/src/main/java/com/openframe/client/config/AsyncConfig.java
@@ -1,12 +1,13 @@
package com.openframe.client.config;
+import com.openframe.core.async.TracedExecutorFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.AsyncTaskExecutor;
import org.springframework.core.task.support.TaskExecutorAdapter;
import org.springframework.scheduling.annotation.EnableAsync;
-import java.util.concurrent.Executors;
+import java.util.concurrent.ExecutorService;
@Configuration
@EnableAsync
@@ -16,6 +17,7 @@ public class AsyncConfig {
@Bean(TOOL_INSTALL_EXECUTOR)
public AsyncTaskExecutor toolInstallExecutor() {
- return new TaskExecutorAdapter(Executors.newVirtualThreadPerTaskExecutor());
+ ExecutorService executor = TracedExecutorFactory.newVirtualThreadPerTaskExecutor();
+ return new TaskExecutorAdapter(executor);
}
}
diff --git a/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt-shortened.xml b/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt-shortened.xml
index 1c3f88612..dde1d818a 100644
--- a/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt-shortened.xml
+++ b/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt-shortened.xml
@@ -4,7 +4,7 @@
- ts=%d{yyyy-MM-dd'T'HH:mm:ss.SSSX} level=%-5level service=${SPRING_APP_NAME} logger=%logger{36} thread=%thread msg="%replace(%msg){'\r?\n',' '}" stack="%replace(%exShort{30,full,2048,rootFirst,inlineHash,FastClassByCGLIB,EnhancerBySpringCGLIB,^sun\.reflect\..*\.invoke,^com\.sun\.,^sun\.net\.,^net\.sf\.cglib\.proxy\.MethodProxy\.invoke,^org\.springframework\.cglib\.,^org\.springframework\.transaction\.,^org\.springframework\.validation\.,^org\.springframework\.app\.,^org\.springframework\.aop\.,^java\.lang\.reflect\.Method\.invoke,^org\.springframework\.ws\..*\.invoke,^org\.springframework\.ws\.transport\.,^org\.springframework\.ws\.soap\.saaj\.SaajSoapMessage\.,^org\.springframework\.ws\.client\.core\.WebServiceTemplate\.,^org\.springframework\.web\.filter\.,^org\.apache\.tomcat\.,^org\.apache\.catalina\.,^org\.apache\.coyote\.,^java\.util\.concurrent\.ThreadPoolExecutor\.runWorker,^java\.lang\.Thread\.run}){'\r?\n',' | '}"%n
+ ts=%d{yyyy-MM-dd'T'HH:mm:ss.SSSX} level=%-5level service=${SPRING_APP_NAME} traceId=%X{traceId:-} spanId=%X{spanId:-} logger=%logger{36} thread=%thread msg="%replace(%msg){'\r?\n',' '}" stack="%replace(%exShort{30,full,2048,rootFirst,inlineHash,FastClassByCGLIB,EnhancerBySpringCGLIB,^sun\.reflect\..*\.invoke,^com\.sun\.,^sun\.net\.,^net\.sf\.cglib\.proxy\.MethodProxy\.invoke,^org\.springframework\.cglib\.,^org\.springframework\.transaction\.,^org\.springframework\.validation\.,^org\.springframework\.app\.,^org\.springframework\.aop\.,^java\.lang\.reflect\.Method\.invoke,^org\.springframework\.ws\..*\.invoke,^org\.springframework\.ws\.transport\.,^org\.springframework\.ws\.soap\.saaj\.SaajSoapMessage\.,^org\.springframework\.ws\.client\.core\.WebServiceTemplate\.,^org\.springframework\.web\.filter\.,^org\.apache\.tomcat\.,^org\.apache\.catalina\.,^org\.apache\.coyote\.,^java\.util\.concurrent\.ThreadPoolExecutor\.runWorker,^java\.lang\.Thread\.run}){'\r?\n',' | '}"%n
diff --git a/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt.xml b/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt.xml
index eb40595f1..b771f6303 100644
--- a/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt.xml
+++ b/openframe-config-core/src/main/resources/logging/shared-logback-includes-logfmt.xml
@@ -3,7 +3,7 @@
- ts=%d{yyyy-MM-dd'T'HH:mm:ss.SSSX} level=%-5level service=${SPRING_APP_NAME} logger=%logger{36} thread=%thread msg="%replace(%msg){'\r?\n',' '}" stack="%replace(%ex){'\r?\n',' | '}"%n
+ ts=%d{yyyy-MM-dd'T'HH:mm:ss.SSSX} level=%-5level service=${SPRING_APP_NAME} traceId=%X{traceId:-} spanId=%X{spanId:-} logger=%logger{36} thread=%thread msg="%replace(%msg){'\r?\n',' '}" stack="%replace(%ex){'\r?\n',' | '}"%n
diff --git a/openframe-config-core/src/main/resources/logging/shared-logback-includes.xml b/openframe-config-core/src/main/resources/logging/shared-logback-includes.xml
index 4e04199e2..383d631c1 100644
--- a/openframe-config-core/src/main/resources/logging/shared-logback-includes.xml
+++ b/openframe-config-core/src/main/resources/logging/shared-logback-includes.xml
@@ -3,7 +3,7 @@
- %d{yyyy-MM-dd HH:mm:ss} %-5level [${SPRING_APP_NAME}] [%thread] %logger{36} - %msg%n
+ %d{yyyy-MM-dd HH:mm:ss} %-5level [${SPRING_APP_NAME}] [%X{traceId:-}] [%thread] %logger{36} - %msg%n
diff --git a/openframe-core/pom.xml b/openframe-core/pom.xml
index cac8047ec..d86535914 100644
--- a/openframe-core/pom.xml
+++ b/openframe-core/pom.xml
@@ -50,6 +50,10 @@
io.micrometer
micrometer-observation
+
+ io.micrometer
+ context-propagation
+
diff --git a/openframe-core/src/main/java/com/openframe/core/async/TracedExecutorFactory.java b/openframe-core/src/main/java/com/openframe/core/async/TracedExecutorFactory.java
new file mode 100644
index 000000000..fa1d582e3
--- /dev/null
+++ b/openframe-core/src/main/java/com/openframe/core/async/TracedExecutorFactory.java
@@ -0,0 +1,29 @@
+package com.openframe.core.async;
+
+import io.micrometer.context.ContextExecutorService;
+import io.micrometer.context.ContextSnapshot;
+import io.micrometer.context.ContextSnapshotFactory;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.function.Supplier;
+
+public final class TracedExecutorFactory {
+
+ private static final ContextSnapshotFactory SNAPSHOT_FACTORY = ContextSnapshotFactory.builder().build();
+
+ private TracedExecutorFactory() {
+ }
+
+ // A fresh thread starts with empty thread locals, so without this wrapper traceId/spanId
+ // and the observation scope are lost the moment work crosses an async boundary.
+ public static ExecutorService newVirtualThreadPerTaskExecutor() {
+ ExecutorService delegate = Executors.newVirtualThreadPerTaskExecutor();
+ return trace(delegate);
+ }
+
+ public static ExecutorService trace(ExecutorService delegate) {
+ Supplier snapshotSupplier = SNAPSHOT_FACTORY::captureAll;
+ return ContextExecutorService.wrap(delegate, snapshotSupplier);
+ }
+}
diff --git a/openframe-core/src/test/java/com/openframe/core/async/TracedExecutorFactoryTest.java b/openframe-core/src/test/java/com/openframe/core/async/TracedExecutorFactoryTest.java
new file mode 100644
index 000000000..e72229399
--- /dev/null
+++ b/openframe-core/src/test/java/com/openframe/core/async/TracedExecutorFactoryTest.java
@@ -0,0 +1,60 @@
+package com.openframe.core.async;
+
+import io.micrometer.context.ContextRegistry;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.time.Duration;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class TracedExecutorFactoryTest {
+
+ private static final String ACCESSOR_KEY = "traced-executor-factory-test";
+ private static final String CONTEXT_VALUE = "trace-context-of-caller";
+ private static final ThreadLocal CONTEXT_HOLDER = new ThreadLocal<>();
+ private static final Duration TASK_TIMEOUT = Duration.ofSeconds(5);
+
+ @BeforeEach
+ void setUp() {
+ ContextRegistry registry = ContextRegistry.getInstance();
+ registry.registerThreadLocalAccessor(ACCESSOR_KEY, CONTEXT_HOLDER);
+ }
+
+ @AfterEach
+ void tearDown() {
+ ContextRegistry registry = ContextRegistry.getInstance();
+ registry.removeThreadLocalAccessor(ACCESSOR_KEY);
+ CONTEXT_HOLDER.remove();
+ }
+
+ @Test
+ void newVirtualThreadPerTaskExecutor_callerHasThreadLocalContext_contextRestoredInsideTask() {
+ // setup
+ CONTEXT_HOLDER.set(CONTEXT_VALUE);
+ ExecutorService executor = TracedExecutorFactory.newVirtualThreadPerTaskExecutor();
+ CompletableFuture contextInsideTask = new CompletableFuture<>();
+
+ // execution
+ executor.execute(() -> contextInsideTask.complete(CONTEXT_HOLDER.get()));
+
+ // verifications
+ assertThat(contextInsideTask).succeedsWithin(TASK_TIMEOUT).isEqualTo(CONTEXT_VALUE);
+ }
+
+ @Test
+ void newVirtualThreadPerTaskExecutor_callerHasNoContext_taskStillRuns() {
+ // setup
+ ExecutorService executor = TracedExecutorFactory.newVirtualThreadPerTaskExecutor();
+ CompletableFuture taskExecuted = new CompletableFuture<>();
+
+ // execution
+ executor.execute(() -> taskExecuted.complete(true));
+
+ // verifications
+ assertThat(taskExecuted).succeedsWithin(TASK_TIMEOUT).isEqualTo(true);
+ }
+}
diff --git a/openframe-data-kafka/src/main/java/com/openframe/kafka/config/OssTenantKafkaAutoConfiguration.java b/openframe-data-kafka/src/main/java/com/openframe/kafka/config/OssTenantKafkaAutoConfiguration.java
index 49cd17ce5..4442f9859 100644
--- a/openframe-data-kafka/src/main/java/com/openframe/kafka/config/OssTenantKafkaAutoConfiguration.java
+++ b/openframe-data-kafka/src/main/java/com/openframe/kafka/config/OssTenantKafkaAutoConfiguration.java
@@ -59,6 +59,8 @@ public KafkaTemplate ossTenantKafkaTemplate(
template.setDefaultTopic(templateProperties.getDefaultTopic());
}
+ template.setObservationEnabled(true);
+
return template;
}
@@ -112,6 +114,8 @@ public ConcurrentKafkaListenerContainerFactory
+
+ io.micrometer
+ micrometer-tracing
+
diff --git a/openframe-gateway-service-core/src/main/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilter.java b/openframe-gateway-service-core/src/main/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilter.java
new file mode 100644
index 000000000..39809178c
--- /dev/null
+++ b/openframe-gateway-service-core/src/main/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilter.java
@@ -0,0 +1,35 @@
+package com.openframe.gateway.filter;
+
+import io.micrometer.tracing.Span;
+import io.micrometer.tracing.Tracer;
+import lombok.RequiredArgsConstructor;
+import org.springframework.stereotype.Component;
+import org.springframework.web.server.ServerWebExchange;
+import org.springframework.web.server.WebFilter;
+import org.springframework.web.server.WebFilterChain;
+import reactor.core.publisher.Mono;
+
+/**
+ * Exposes the current trace id as an X-Trace-Id response header so a failing
+ * request can be looked up in Loki without access to server logs.
+ */
+@Component
+@RequiredArgsConstructor
+public class TraceIdResponseHeaderFilter implements WebFilter {
+
+ public static final String TRACE_ID_HEADER = "X-Trace-Id";
+
+ private final Tracer tracer;
+
+ @Override
+ public Mono filter(ServerWebExchange exchange, WebFilterChain chain) {
+ exchange.getResponse().beforeCommit(() -> Mono.deferContextual(ctx -> {
+ Span span = tracer.currentSpan();
+ if (span != null) {
+ exchange.getResponse().getHeaders().set(TRACE_ID_HEADER, span.context().traceId());
+ }
+ return Mono.empty();
+ }));
+ return chain.filter(exchange);
+ }
+}
diff --git a/openframe-gateway-service-core/src/test/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilterTest.java b/openframe-gateway-service-core/src/test/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilterTest.java
new file mode 100644
index 000000000..7fe959b24
--- /dev/null
+++ b/openframe-gateway-service-core/src/test/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilterTest.java
@@ -0,0 +1,52 @@
+package com.openframe.gateway.filter;
+
+import io.micrometer.tracing.Span;
+import io.micrometer.tracing.TraceContext;
+import io.micrometer.tracing.Tracer;
+import org.junit.jupiter.api.Test;
+import org.springframework.mock.http.server.reactive.MockServerHttpRequest;
+import org.springframework.mock.web.server.MockServerWebExchange;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+class TraceIdResponseHeaderFilterTest {
+
+ private static final String TRACE_ID = "4bf92f3577b34da6a3ce929d0e0e4736";
+
+ @Test
+ void addsTraceIdHeaderWhenSpanPresent() {
+ Tracer tracer = mock(Tracer.class);
+ Span span = mock(Span.class);
+ TraceContext context = mock(TraceContext.class);
+ when(tracer.currentSpan()).thenReturn(span);
+ when(span.context()).thenReturn(context);
+ when(context.traceId()).thenReturn(TRACE_ID);
+
+ MockServerWebExchange exchange =
+ MockServerWebExchange.from(MockServerHttpRequest.get("/api/test"));
+ TraceIdResponseHeaderFilter filter = new TraceIdResponseHeaderFilter(tracer);
+
+ filter.filter(exchange, ex -> ex.getResponse().setComplete()).block();
+
+ assertEquals(TRACE_ID,
+ exchange.getResponse().getHeaders().getFirst(TraceIdResponseHeaderFilter.TRACE_ID_HEADER));
+ }
+
+ @Test
+ void noHeaderWhenNoCurrentSpan() {
+ Tracer tracer = mock(Tracer.class);
+ when(tracer.currentSpan()).thenReturn(null);
+
+ MockServerWebExchange exchange =
+ MockServerWebExchange.from(MockServerHttpRequest.get("/api/test"));
+ TraceIdResponseHeaderFilter filter = new TraceIdResponseHeaderFilter(tracer);
+
+ filter.filter(exchange, ex -> ex.getResponse().setComplete()).block();
+
+ assertNull(exchange.getResponse().getHeaders()
+ .getFirst(TraceIdResponseHeaderFilter.TRACE_ID_HEADER));
+ }
+}