From 3f0a8a70c02650e4ece2203e7ce3c53e48f3d0cd Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 27 Jul 2026 14:52:59 +0300 Subject: [PATCH 1/4] feat(logging): add traceId/spanId MDC fields to shared logback patterns Co-Authored-By: Claude Fable 5 --- .../logging/shared-logback-includes-logfmt-shortened.xml | 5 +++-- .../resources/logging/shared-logback-includes-logfmt.xml | 5 +++-- .../src/main/resources/logging/shared-logback-includes.xml | 4 ++-- 3 files changed, 8 insertions(+), 6 deletions(-) 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 78e4949cf..36fd6f159 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 @@ -18,7 +18,7 @@ 30 - 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 @@ -38,6 +38,7 @@ + 30 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 221c9a0e6..04c0e5790 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 @@ -17,7 +17,7 @@ 30 - 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 @@ -37,6 +37,7 @@ + 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 b02559b06..195477253 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 @@ -17,7 +17,7 @@ 30 - %d{yyyy-MM-dd HH:mm:ss} %-5level [%thread] %logger{36} - %msg%n + %d{yyyy-MM-dd HH:mm:ss} %-5level [%X{traceId:-}] [%thread] %logger{36} - %msg%n From 70fddaaff9dbf41a8787314a609603a271fac34a Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 27 Jul 2026 14:54:20 +0300 Subject: [PATCH 2/4] feat(gateway): return X-Trace-Id response header Co-Authored-By: Claude Fable 5 --- openframe-gateway-service-core/pom.xml | 4 ++ .../filter/TraceIdResponseHeaderFilter.java | 35 +++++++++++++ .../TraceIdResponseHeaderFilterTest.java | 52 +++++++++++++++++++ 3 files changed, 91 insertions(+) create mode 100644 openframe-gateway-service-core/src/main/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilter.java create mode 100644 openframe-gateway-service-core/src/test/java/com/openframe/gateway/filter/TraceIdResponseHeaderFilterTest.java 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)); + } +} From d7c1b09dcfd274e73c858dbd11ea7292d1e30828 Mon Sep 17 00:00:00 2001 From: andrii Date: Mon, 27 Jul 2026 14:55:27 +0300 Subject: [PATCH 3/4] feat(kafka): enable observation for OSS producer and listener (traceparent propagation) Co-Authored-By: Claude Fable 5 --- .../OssTenantKafkaAutoConfiguration.java | 4 ++ .../OssTenantKafkaAutoConfigurationTest.java | 40 +++++++++++++++++++ 2 files changed, 44 insertions(+) create mode 100644 openframe-data-kafka/src/test/java/com/openframe/kafka/config/OssTenantKafkaAutoConfigurationTest.java 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()); + } +} From 040da3a829581a045cb5f8a7c016aff9e6b598d0 Mon Sep 17 00:00:00 2001 From: andrii Date: Wed, 12 Aug 2026 17:09:30 +0200 Subject: [PATCH 4/4] fix(tracing): propagate trace context across @Async boundaries in libs Adds TracedExecutorFactory in openframe-core and uses it for the tool-install executor (client-core) and the notification channel executor (data-nats). Both created plain virtual-thread executors, so traceId/spanId were lost as soon as work crossed the @Async boundary. Co-Authored-By: Claude Fable 5 --- .../openframe/client/config/AsyncConfig.java | 6 +- openframe-core/pom.xml | 4 ++ .../core/async/TracedExecutorFactory.java | 29 +++++++++ .../core/async/TracedExecutorFactoryTest.java | 60 +++++++++++++++++++ .../NotificationChannelExecutorConfig.java | 4 +- 5 files changed, 99 insertions(+), 4 deletions(-) create mode 100644 openframe-core/src/main/java/com/openframe/core/async/TracedExecutorFactory.java create mode 100644 openframe-core/src/test/java/com/openframe/core/async/TracedExecutorFactoryTest.java 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-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-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(); } }