diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptor.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptor.java
index 543cc2aa6e..97f32885fa 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptor.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptor.java
@@ -32,11 +32,22 @@ import org.springframework.util.Assert;
* implementation responsible for an {@link Observation} propagation from one message
* flow's thread to another through the {@link MessageChannel}s involved in the flow.
* Opens a new {@link Observation.Scope} on another thread and cleans up it in the end.
+ *
+ * NOTE: This interceptor is proven to be wrong since an existing observation usually is closed
+ * on the sender side before the message is consumed on the receiver side.
+ * Therefore, it is better to have a {@code sender} observation on this channel,
+ * and then {@code receiver} observation on a subscriber for this channel.
+ * This way a tracing information is stored into message headers passing this channel.
+ * Such an approach also eliminate a problem with persistent message channels where
+ * an {@link Observation} is not serializable to be stored into database as a part of the message.
*
* @author Artem Bilan
*
* @since 6.0
+ *
+ * @deprecated since 6.1.7 for removal in 6.4 in favor of enabling observation on the channel and its consumer.
*/
+@Deprecated(since = "6.1.7", forRemoval = true)
public class ObservationPropagationChannelInterceptor extends ThreadStatePropagationChannelInterceptor {
private final ThreadLocal scopes = new ThreadLocal<>();
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptorTests.java
deleted file mode 100644
index 7295d8a8ab..0000000000
--- a/spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptorTests.java
+++ /dev/null
@@ -1,356 +0,0 @@
-/*
- * Copyright 2022 the original author or authors.
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * https://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.springframework.integration.channel.interceptor;
-
-import java.util.Arrays;
-import java.util.List;
-import java.util.concurrent.CountDownLatch;
-import java.util.concurrent.Executors;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
-
-import io.micrometer.common.KeyValues;
-import io.micrometer.core.instrument.MeterRegistry;
-import io.micrometer.core.instrument.observation.DefaultMeterObservationHandler;
-import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
-import io.micrometer.core.tck.MeterRegistryAssert;
-import io.micrometer.observation.Observation;
-import io.micrometer.observation.ObservationHandler;
-import io.micrometer.observation.ObservationRegistry;
-import io.micrometer.observation.tck.TestObservationRegistry;
-import io.micrometer.observation.tck.TestObservationRegistryAssert;
-import io.micrometer.tracing.Span;
-import io.micrometer.tracing.TraceContext;
-import io.micrometer.tracing.Tracer;
-import io.micrometer.tracing.handler.DefaultTracingObservationHandler;
-import io.micrometer.tracing.handler.PropagatingReceiverTracingObservationHandler;
-import io.micrometer.tracing.handler.PropagatingSenderTracingObservationHandler;
-import io.micrometer.tracing.propagation.Propagator;
-import io.micrometer.tracing.test.simple.SimpleTracer;
-import io.micrometer.tracing.test.simple.SpansAssert;
-import io.micrometer.tracing.test.simple.TracerAssert;
-import org.assertj.core.api.InstanceOfAssertFactories;
-import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Test;
-
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.context.annotation.Bean;
-import org.springframework.context.annotation.Configuration;
-import org.springframework.integration.annotation.BridgeTo;
-import org.springframework.integration.annotation.Poller;
-import org.springframework.integration.channel.DirectChannel;
-import org.springframework.integration.channel.ExecutorChannel;
-import org.springframework.integration.channel.QueueChannel;
-import org.springframework.integration.config.EnableIntegration;
-import org.springframework.integration.config.GlobalChannelInterceptor;
-import org.springframework.integration.handler.BridgeHandler;
-import org.springframework.integration.support.MessageBuilder;
-import org.springframework.integration.support.management.observation.IntegrationObservation;
-import org.springframework.lang.Nullable;
-import org.springframework.messaging.Message;
-import org.springframework.messaging.MessageHeaders;
-import org.springframework.messaging.PollableChannel;
-import org.springframework.messaging.SubscribableChannel;
-import org.springframework.messaging.support.ChannelInterceptor;
-import org.springframework.messaging.support.GenericMessage;
-import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
-
-import static org.assertj.core.api.Assertions.assertThat;
-
-/**
- * @author Artem Bilan
- *
- * @since 6.0
- */
-@SpringJUnitConfig
-public class ObservationPropagationChannelInterceptorTests {
-
- @Autowired
- ObservationRegistry observationRegistry;
-
- @Autowired
- MeterRegistry meterRegistry;
-
- @Autowired
- SimpleTracer simpleTracer;
-
- @Autowired
- SubscribableChannel directChannel;
-
- @Autowired
- SubscribableChannel executorChannel;
-
- @Autowired
- PollableChannel queueChannel;
-
- @Autowired
- DirectChannel testConsumer;
-
- @Autowired
- ExecutorChannel testTracingChannel;
-
- @BeforeEach
- void setup() {
- this.simpleTracer.getSpans().clear();
- }
-
- @Test
- void observationPropagatedOverDirectChannel() throws InterruptedException {
- AtomicReference scopeReference = new AtomicReference<>();
- CountDownLatch handleLatch = new CountDownLatch(1);
- this.directChannel.subscribe(m -> {
- scopeReference.set(this.observationRegistry.getCurrentObservationScope());
- handleLatch.countDown();
- });
-
- AtomicReference originalScope = new AtomicReference<>();
-
- Observation.createNotStarted("test1", this.observationRegistry)
- .observe(() -> {
- originalScope.set(this.observationRegistry.getCurrentObservationScope());
- this.directChannel.send(new GenericMessage<>("test"));
- });
-
- assertThat(handleLatch.await(10, TimeUnit.SECONDS)).isTrue();
- assertThat(scopeReference.get())
- .isNotNull()
- .isSameAs(originalScope.get());
-
- TestObservationRegistryAssert.assertThat(this.observationRegistry)
- .doesNotHaveAnyRemainingCurrentObservation();
-
- TracerAssert.assertThat(this.simpleTracer)
- .onlySpan()
- .hasNameEqualTo("test1");
- }
-
- @Test
- void observationPropagatedOverExecutorChannel() throws InterruptedException {
- AtomicReference scopeReference = new AtomicReference<>();
- CountDownLatch handleLatch = new CountDownLatch(1);
- this.executorChannel.subscribe(m -> {
- scopeReference.set(this.observationRegistry.getCurrentObservationScope());
- handleLatch.countDown();
- });
-
- AtomicReference originalScope = new AtomicReference<>();
-
- Observation.createNotStarted("test2", this.observationRegistry)
- .observe(() -> {
- originalScope.set(this.observationRegistry.getCurrentObservationScope());
- this.executorChannel.send(new GenericMessage<>("test"));
- });
-
- assertThat(handleLatch.await(10, TimeUnit.SECONDS)).isTrue();
- assertThat(scopeReference.get())
- .isNotNull()
- .isNotSameAs(originalScope.get());
-
- assertThat(scopeReference.get().getCurrentObservation())
- .isSameAs(originalScope.get().getCurrentObservation());
-
- TestObservationRegistryAssert.assertThat(this.observationRegistry)
- .doesNotHaveAnyRemainingCurrentObservation();
-
- TracerAssert.assertThat(this.simpleTracer)
- .onlySpan()
- .hasNameEqualTo("test2");
- }
-
- @Test
- void observationPropagatedOverQueueChannel() throws InterruptedException {
- AtomicReference scopeReference = new AtomicReference<>();
- CountDownLatch handleLatch = new CountDownLatch(1);
- this.testConsumer.subscribe(m -> {
- scopeReference.set(this.observationRegistry.getCurrentObservationScope());
- handleLatch.countDown();
- });
-
- AtomicReference originalScope = new AtomicReference<>();
-
- Observation.createNotStarted("test3", this.observationRegistry)
- .observe(() -> {
- originalScope.set(this.observationRegistry.getCurrentObservationScope());
- this.queueChannel.send(new GenericMessage<>("test"));
- });
-
- assertThat(handleLatch.await(10, TimeUnit.SECONDS)).isTrue();
- assertThat(scopeReference.get())
- .isNotNull()
- .isNotSameAs(originalScope.get());
-
- assertThat(scopeReference.get().getCurrentObservation())
- .isSameAs(originalScope.get().getCurrentObservation());
-
- TestObservationRegistryAssert.assertThat(this.observationRegistry)
- .doesNotHaveAnyRemainingCurrentObservation();
-
- TracerAssert.assertThat(this.simpleTracer)
- .onlySpan()
- .hasNameEqualTo("test3");
- }
-
- @Test
- void observationContextPropagatedOverExecutorChannel() {
- BridgeHandler handler = new BridgeHandler();
- handler.registerObservationRegistry(this.observationRegistry);
- handler.setBeanName("testBridge");
- this.testTracingChannel.subscribe(handler);
-
- QueueChannel replyChannel = new QueueChannel();
-
- Message message =
- MessageBuilder.withPayload("test")
- .setHeader(MessageHeaders.REPLY_CHANNEL, replyChannel)
- .build();
-
- this.testTracingChannel.send(message);
-
- Message> receive = replyChannel.receive();
-
- assertThat(receive).isNotNull()
- .extracting(Message::getHeaders)
- .asInstanceOf(InstanceOfAssertFactories.MAP)
- .containsEntry("foo", "some foo value")
- .containsEntry("bar", "some bar value");
-
- TestObservationRegistryAssert.assertThat(this.observationRegistry)
- .doesNotHaveAnyRemainingCurrentObservation();
-
- TracerAssert.assertThat(this.simpleTracer)
- .reportedSpans()
- .hasSize(2)
- .satisfies(simpleSpans -> SpansAssert.assertThat(simpleSpans)
- .assertThatASpanWithNameEqualTo("testTracingChannel send")
- .hasTag("spring.integration.type", "producer")
- .hasTag("spring.integration.name", "testTracingChannel")
- .hasKindEqualTo(Span.Kind.PRODUCER)
- .backToSpans()
- .assertThatASpanWithNameEqualTo("testBridge receive")
- .hasTag("foo", "some foo value")
- .hasTag("bar", "some bar value")
- .hasTag("spring.integration.type", "handler")
- .hasTag("spring.integration.name", "testBridge")
- .hasKindEqualTo(Span.Kind.CONSUMER));
-
-
- MeterRegistryAssert.assertThat(this.meterRegistry)
- .hasTimerWithNameAndTags("spring.integration.handler",
- KeyValues.of(IntegrationObservation.HandlerTags.COMPONENT_NAME.asString(), "testBridge",
- IntegrationObservation.HandlerTags.COMPONENT_TYPE.asString(), "handler",
- "error", "none"));
-
- assertThat(this.meterRegistry.get("spring.integration.handler").timer().count()).isEqualTo(1);
- }
-
- @Configuration
- @EnableIntegration
- public static class ContextConfiguration {
-
- @Bean
- SimpleTracer simpleTracer() {
- return new SimpleTracer();
- }
-
- @Bean
- MeterRegistry meterRegistry() {
- return new SimpleMeterRegistry();
- }
-
- @Bean
- ObservationRegistry observationRegistry(Tracer tracer, Propagator propagator, MeterRegistry meterRegistry) {
- TestObservationRegistry observationRegistry = TestObservationRegistry.create();
- observationRegistry.observationConfig()
- .observationHandler(new DefaultMeterObservationHandler(meterRegistry))
- .observationHandler(
- // Composite will pick the first matching handler
- new ObservationHandler.FirstMatchingCompositeObservationHandler(
- // This is responsible for creating a child span on the sender side
- new PropagatingSenderTracingObservationHandler<>(tracer, propagator),
- // This is responsible for creating a span on the receiver side
- new PropagatingReceiverTracingObservationHandler<>(tracer, propagator),
- // This is responsible for creating a default span
- new DefaultTracingObservationHandler(tracer)));
- return observationRegistry;
- }
-
- @Bean
- @GlobalChannelInterceptor(patterns = "*Channel")
- public ChannelInterceptor observationPropagationInterceptor(ObservationRegistry observationRegistry) {
- return new ObservationPropagationChannelInterceptor(observationRegistry);
- }
-
- @Bean
- @BridgeTo(value = "testConsumer", poller = @Poller(fixedDelay = "100"))
- public PollableChannel queueChannel() {
- return new QueueChannel();
- }
-
- @Bean
- public SubscribableChannel executorChannel() {
- return new ExecutorChannel(Executors.newSingleThreadExecutor());
- }
-
- @Bean
- public SubscribableChannel directChannel() {
- return new DirectChannel();
- }
-
- @Bean
- public DirectChannel testConsumer() {
- return new DirectChannel();
- }
-
- @Bean
- public ExecutorChannel testTracingChannel(ObservationRegistry observationRegistry) {
- ExecutorChannel channel = new ExecutorChannel(Executors.newSingleThreadExecutor());
- channel.registerObservationRegistry(observationRegistry);
- return channel;
- }
-
- @Bean
- public Propagator propagator(Tracer tracer) {
- return new Propagator() {
-
- // List of headers required for tracing propagation
- @Override
- public List fields() {
- return Arrays.asList("foo", "bar");
- }
-
- // This is called on the producer side when the message is being sent
- // Normally we would pass information from tracing context - for tests we don't need to
- @Override
- public void inject(TraceContext context, @Nullable C carrier, Setter setter) {
- setter.set(carrier, "foo", "some foo value");
- setter.set(carrier, "bar", "some bar value");
- }
-
- // This is called on the consumer side when the message is consumed
- // Normally we would use tools like Extractor from tracing but for tests we are just manually creating a span
- @Override
- public Span.Builder extract(C carrier, Getter getter) {
- String foo = getter.get(carrier, "foo");
- String bar = getter.get(carrier, "bar");
- return tracer.spanBuilder().tag("foo", foo).tag("bar", bar);
- }
- };
- }
-
- }
-
-}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/support/management/observation/IntegrationObservabilityZipkinTests.java b/spring-integration-core/src/test/java/org/springframework/integration/support/management/observation/IntegrationObservabilityZipkinTests.java
index e572ebf29d..530eafdf7a 100644
--- a/spring-integration-core/src/test/java/org/springframework/integration/support/management/observation/IntegrationObservabilityZipkinTests.java
+++ b/spring-integration-core/src/test/java/org/springframework/integration/support/management/observation/IntegrationObservabilityZipkinTests.java
@@ -35,17 +35,14 @@ import org.springframework.integration.annotation.Poller;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.NullChannel;
import org.springframework.integration.channel.QueueChannel;
-import org.springframework.integration.channel.interceptor.ObservationPropagationChannelInterceptor;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.config.EnableIntegrationManagement;
-import org.springframework.integration.config.GlobalChannelInterceptor;
import org.springframework.integration.gateway.MessagingGatewaySupport;
import org.springframework.integration.handler.BridgeHandler;
import org.springframework.integration.handler.advice.HandleMessageAdvice;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.PollableChannel;
-import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.messaging.support.GenericMessage;
import static org.assertj.core.api.Assertions.assertThat;
@@ -136,12 +133,6 @@ public class IntegrationObservabilityZipkinTests extends SampleTestRunner {
CountDownLatch observedHandlerLatch = new CountDownLatch(1);
- @Bean
- @GlobalChannelInterceptor
- public ChannelInterceptor observationPropagationInterceptor(ObservationRegistry observationRegistry) {
- return new ObservationPropagationChannelInterceptor(observationRegistry);
- }
-
@Bean
TestMessagingGatewaySupport testInboundGateway(@Qualifier("queueChannel") PollableChannel queueChannel) {
TestMessagingGatewaySupport messagingGatewaySupport = new TestMessagingGatewaySupport();
diff --git a/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/observation/WebFluxObservationPropagationTests.java b/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/observation/WebFluxObservationPropagationTests.java
index b757bb600e..7d811cb7d0 100644
--- a/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/observation/WebFluxObservationPropagationTests.java
+++ b/spring-integration-webflux/src/test/java/org/springframework/integration/webflux/observation/WebFluxObservationPropagationTests.java
@@ -23,7 +23,6 @@ import brave.propagation.ThreadLocalCurrentTraceContext;
import brave.test.TestSpanHandler;
import io.micrometer.observation.ObservationHandler;
import io.micrometer.observation.ObservationRegistry;
-import io.micrometer.observation.tck.TestObservationRegistryAssert;
import io.micrometer.tracing.Tracer;
import io.micrometer.tracing.brave.bridge.BraveBaggageManager;
import io.micrometer.tracing.brave.bridge.BraveCurrentTraceContext;
@@ -44,15 +43,12 @@ import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.HttpMethod;
import org.springframework.integration.channel.FluxMessageChannel;
-import org.springframework.integration.channel.interceptor.ObservationPropagationChannelInterceptor;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.config.EnableIntegrationManagement;
-import org.springframework.integration.config.GlobalChannelInterceptor;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.webflux.dsl.WebFlux;
import org.springframework.messaging.Message;
import org.springframework.messaging.PollableChannel;
-import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.junit.jupiter.web.SpringJUnitWebConfig;
import org.springframework.test.web.reactive.server.WebTestClient;
@@ -100,12 +96,8 @@ public class WebFluxObservationPropagationTests {
.extracting(Message::getPayload)
.isEqualTo("Received data: " + testData);
- TestObservationRegistryAssert.assertThat(this.observationRegistry)
- .hasRemainingCurrentObservation();
-
- this.observationRegistry.getCurrentObservation().stop();
-
- assertThat(SPANS.spans()).hasSize(6);
+ // There is a race condition when we already have a reply, but the span in the last channel is not closed yet.
+ await().untilAsserted(() -> assertThat(SPANS.spans()).hasSize(6));
SpansAssert.assertThat(SPANS.spans().stream().map(BraveFinishedSpan::fromBrave).collect(Collectors.toList()))
.haveSameTraceId();
}
@@ -175,12 +167,6 @@ public class WebFluxObservationPropagationTests {
return new ServerHttpObservationFilter(registry);
}
- @Bean
- @GlobalChannelInterceptor
- public ChannelInterceptor observationPropagationInterceptor(ObservationRegistry observationRegistry) {
- return new ObservationPropagationChannelInterceptor(observationRegistry);
- }
-
@Bean
IntegrationFlow webFluxFlow() {
return IntegrationFlow
diff --git a/src/reference/asciidoc/metrics.adoc b/src/reference/asciidoc/metrics.adoc
index 455584a269..c2df3ed1d4 100644
--- a/src/reference/asciidoc/metrics.adoc
+++ b/src/reference/asciidoc/metrics.adoc
@@ -167,7 +167,7 @@ It uses the `IntegrationObservation.GATEWAY` API;
* An `AbstractMessageChannel.send()` operation is the only Spring Integration API where it produces messages.
So, it is treated as a `PRODUCER` span type and uses the `IntegrationObservation.PRODCUER` API.
This makes more sense when a channel is a distributed implementation (e.g. `PublishSubscribeKafkaChannel` or `ZeroMqChannel`) and trace information has to be added to the message.
-So, the `IntegrationObservation.PRODUCER` observation is based on a `MessageSenderContext` where Spring Integration supplies a `MutableMessage` to allow a subsequent tracing `Propagator` to add headers so they are available to the consumer;
+So, the `IntegrationObservation.PRODUCER` observation is based on a `MessageSenderContext` where Spring Integration supplies a `MutableMessage` to allow a subsequent tracing `Propagator` to add headers, so they are available to the consumer;
* An `AbstractMessageHandler` is a `CONSUMER` span type and uses the `IntegrationObservation.HANDLER` API.
An observation production on the `IntegrationManagement` components can be customized via `ObservationConvention` configuration.
@@ -183,10 +183,10 @@ include::./generated/conventions.adoc[leveloffset=+2]
==== Observation Propagation
-To supply a connected chain of spans in one trace, independently of the nature of the messaging flow, Spring Integration provides an `ObservationPropagationChannelInterceptor` implementation.
-This can be configured on `MessageChannnel` beans individually or as a `@GlobalChannelInterceptor` with respective `MessageChannnel` bean names pattern matching.
-The goal of this interceptor is to propagate an `Observation` from the producer thread to the consumer one independently of the `MessageChannnel` implementation and nature.
-A `DirectChannel`, though, is ignored since its consumer is executed directly on the producer thread.
+To supply a connected chain of spans in one trace, independently of the nature of the messaging flow, even if a `MessageChannel` is persistent and distributed, the observation must be enabled on this channel and on consumers (subscribers) for this channel.
+This way, the tracing information is stored in the message headers before it is propagated to a consumer thread or persisted into the database.
+This is done via mentioned above `MessageSenderContext`.
+The consumer (a `MessageHandler`) side restores tracing information from those headers using a `MessageReceiverContext` and starts a new child `Observation`.
==== Spring Integration JMX Support