From bc7ea37f312e02f7788ec1c328498ef53ad168ec Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 18 Mar 2024 16:19:00 -0400 Subject: [PATCH] GH-9001: Revise Observation propagation over the channel Fixes: #9001 The `ObservationPropagationChannelInterceptor` does not propagate an observation properly. And it fully cannot when the message channel is persistent. * Deprecate `ObservationPropagationChannelInterceptor` in favor of enabled observation on the channel and target `MessageHandler` which is a consumer of this channel. * Remove tests with an `ObservationPropagationChannelInterceptor` * Mention a correct behavior in the `metrics.adoc` and `ObservationPropagationChannelInterceptor` Javadocs --- ...ervationPropagationChannelInterceptor.java | 11 + ...ionPropagationChannelInterceptorTests.java | 356 ------------------ .../IntegrationObservabilityZipkinTests.java | 9 - .../WebFluxObservationPropagationTests.java | 18 +- src/reference/asciidoc/metrics.adoc | 10 +- 5 files changed, 18 insertions(+), 386 deletions(-) delete mode 100644 spring-integration-core/src/test/java/org/springframework/integration/channel/interceptor/ObservationPropagationChannelInterceptorTests.java 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