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 ossTenantKafkaLis factory.getContainerProperties().setLogContainerConfig(listenerProperties.getLogContainerConfig()); } + factory.getContainerProperties().setObservationEnabled(true); + return factory; } diff --git a/openframe-data-kafka/src/test/java/com/openframe/kafka/config/OssTenantKafkaAutoConfigurationTest.java b/openframe-data-kafka/src/test/java/com/openframe/kafka/config/OssTenantKafkaAutoConfigurationTest.java new file mode 100644 index 000000000..5966a720a --- /dev/null +++ b/openframe-data-kafka/src/test/java/com/openframe/kafka/config/OssTenantKafkaAutoConfigurationTest.java @@ -0,0 +1,40 @@ +package com.openframe.kafka.config; + +import org.junit.jupiter.api.Test; +import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; +import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.ProducerFactory; +import org.springframework.test.util.ReflectionTestUtils; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; + +/** + * Guards the observation contract: template and listener observation must stay + * enabled so the traceparent header is propagated through Kafka records. + */ +class OssTenantKafkaAutoConfigurationTest { + + private final OssTenantKafkaAutoConfiguration config = new OssTenantKafkaAutoConfiguration(); + + @Test + @SuppressWarnings("unchecked") + void templateHasObservationEnabled() { + KafkaTemplate template = config.ossTenantKafkaTemplate( + (ProducerFactory) mock(ProducerFactory.class), + new OssTenantKafkaProperties()); + assertEquals(Boolean.TRUE, ReflectionTestUtils.getField(template, "observationEnabled")); + } + + @Test + @SuppressWarnings("unchecked") + void listenerFactoryHasObservationEnabled() { + ConcurrentKafkaListenerContainerFactory factory = + config.ossTenantKafkaListenerContainerFactory( + (ConsumerFactory) mock(ConsumerFactory.class), + new OssTenantKafkaProperties()); + assertTrue(factory.getContainerProperties().isObservationEnabled()); + } +} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/config/NotificationChannelExecutorConfig.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/config/NotificationChannelExecutorConfig.java index 289ebee92..837d831cd 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/config/NotificationChannelExecutorConfig.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/config/NotificationChannelExecutorConfig.java @@ -1,10 +1,10 @@ package com.openframe.data.nats.config; +import com.openframe.core.async.TracedExecutorFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.concurrent.Executor; -import java.util.concurrent.Executors; @Configuration public class NotificationChannelExecutorConfig { @@ -14,6 +14,6 @@ public class NotificationChannelExecutorConfig { /** Lib-owned so @Async does not fall back to unbounded platform threads in a consumer with no default executor. */ @Bean(CHANNEL_EXECUTOR) public Executor notificationChannelExecutor() { - return Executors.newVirtualThreadPerTaskExecutor(); + return TracedExecutorFactory.newVirtualThreadPerTaskExecutor(); } } diff --git a/openframe-gateway-service-core/pom.xml b/openframe-gateway-service-core/pom.xml index ca7794fa1..6a68205ea 100644 --- a/openframe-gateway-service-core/pom.xml +++ b/openframe-gateway-service-core/pom.xml @@ -80,6 +80,10 @@ io.micrometer micrometer-core + + 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)); + } +}