dpol1 commented on code in PR #9154:
URL: https://github.com/apache/storm/pull/9154#discussion_r4186224463


##########
storm-server/src/test/java/org/apache/storm/TopologyTracingTest.java:
##########
@@ -0,0 +1,569 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one or more 
contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.  The 
ASF licenses this file to you 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
+ *
+ * http://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.apache.storm;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.api.trace.StatusCode;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.Scope;
+import io.opentelemetry.sdk.OpenTelemetrySdk;
+import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter;
+import io.opentelemetry.sdk.testing.junit5.OpenTelemetryExtension;
+import io.opentelemetry.sdk.trace.SdkTracerProvider;
+import io.opentelemetry.sdk.trace.data.SpanData;
+import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor;
+import io.opentelemetry.sdk.trace.samplers.Sampler;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.BooleanSupplier;
+import java.util.function.Consumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+import org.apache.storm.ILocalCluster.ILocalTopology;
+import org.apache.storm.generated.StormTopology;
+import org.apache.storm.task.OutputCollector;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.testing.AckFailMapTracker;
+import org.apache.storm.testing.FeederSpout;
+import org.apache.storm.topology.OutputFieldsDeclarer;
+import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.topology.base.BaseRichBolt;
+import org.apache.storm.tuple.Fields;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.Values;
+import org.apache.storm.utils.TupleUtils;
+import org.apache.storm.utils.Utils;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Runs topologies on a two-worker local cluster with tracing on or off and 
checks the spans
+ * Storm exports.
+ */
+public class TopologyTracingTest {
+
+    @RegisterExtension
+    static final OpenTelemetryExtension OTEL = OpenTelemetryExtension.create();
+
+    private static final Set<String> RECEIVED_TRACE_IDS = 
ConcurrentHashMap.newKeySet();
+    /** Span ids that were current on the sink's thread while its execute() 
ran. */
+    private static final Set<String> CURRENT_IN_EXECUTE = 
ConcurrentHashMap.newKeySet();
+    private static final AtomicInteger SINK_TUPLES_RECEIVED = new 
AtomicInteger();
+    private static final AtomicInteger UNSAMPLED_CONTEXTS_RECEIVED = new 
AtomicInteger();
+    private static final Set<String> MIDDLE_TRACE_IDS = 
ConcurrentHashMap.newKeySet();
+    /** Tuple value to the id of the span current while middle, then sink, 
handled it. */
+    private static final Map<Object, String> MIDDLE_SPAN_BY_VALUE = new 
ConcurrentHashMap<>();
+    private static final Map<Object, String> SINK_SPAN_BY_VALUE = new 
ConcurrentHashMap<>();
+    private static final Map<String, Integer> WORKER_PORT_BY_COMPONENT = new 
ConcurrentHashMap<>();
+    private static final AtomicReference<Throwable> EMITTER_THREAD_FAILURE =
+        new AtomicReference<>();
+    private static final AtomicInteger TICK_TUPLES_RECEIVED = new 
AtomicInteger();
+    private static final AtomicBoolean SPAN_CURRENT_DURING_TICK = new 
AtomicBoolean();
+
+    private static ILocalCluster cluster;
+    private static int topologyCount;
+    private static String topologyName;
+    private static volatile int sinkTaskId;
+
+    @BeforeAll
+    public static void startCluster() throws Exception {
+        cluster = new LocalCluster();
+    }
+
+    @AfterAll
+    public static void stopCluster() throws Exception {
+        cluster.close();
+    }
+
+    @Test
+    public void testEachSpoutEmitStartsARootSpan() throws Exception {
+        List<SpanData> spans = runSpoutToSink(true, 3, 1, 0); // 3 tuples, 1 
sink task, no ticks
+
+        List<SpanData> emits = named(spans, "spout emit");
+        assertEquals(3, emits.size());
+        for (SpanData emit : emits) {
+            assertFalse(emit.getParentSpanContext().isValid(), "a spout emit 
starts a new trace");
+        }
+        Set<String> emitTraceIds =
+            
emits.stream().map(SpanData::getTraceId).collect(Collectors.toSet());
+        assertEquals(3, emitTraceIds.size());
+        assertEquals(emitTraceIds, RECEIVED_TRACE_IDS, "each tuple carries its 
emit context");
+    }
+
+    @Test
+    public void testNoSpansWhenTracingIsOff() throws Exception {
+        assertTrue(runSpoutToSink(false, 1, 1, 0).isEmpty()); // 1 tuple, 1 
sink task, no ticks
+        assertTrue(RECEIVED_TRACE_IDS.isEmpty(), "tuples carry no context");
+    }
+
+    @Test
+    public void testExecuteSpanIsChildOfTheEmitAndCurrentDuringExecute() 
throws Exception {
+        // two sink tasks with all grouping: on two workers, a copy of each 
tuple crosses workers
+        List<SpanData> spans = runSpoutToSink(true, 2, 2, 0); // 2 tuples, 2 
sink tasks, no ticks
+
+        Map<String, SpanData> emits = byId(named(spans, "spout emit"));
+        List<SpanData> executes = named(spans, "sink execute");
+        assertEquals(2, emits.size());
+        assertEquals(4, executes.size());
+        for (SpanData execute : executes) {
+            SpanData emit = emits.get(execute.getParentSpanId());
+            assertNotNull(emit, "an execute span is a child of the emit that 
produced its tuple");
+            assertEquals(emit.getTraceId(), execute.getTraceId());
+        }
+        assertTrue(executes.stream().anyMatch(s -> 
s.getParentSpanContext().isRemote()),
+            "at least one tuple crossed workers, so its context went through 
the serializer");
+        Set<String> executeIds =
+            
executes.stream().map(SpanData::getSpanId).collect(Collectors.toSet());
+        assertEquals(executeIds, CURRENT_IN_EXECUTE, "the execute span is 
current in the bolt");
+    }
+
+    @Test
+    public void testTickTuplesGetNoSpanAndSeeNoLeftoverContext() throws 
Exception {
+        List<SpanData> spans = runSpoutToSink(true, 1, 1, 1); // 1 tuple, 1 
sink task, 1 s ticks
+
+        assertEquals(1, named(spans, "sink execute").size());
+        assertFalse(SPAN_CURRENT_DURING_TICK.get(), "the execute span's scope 
was closed");
+    }
+
+    @Test
+    public void testExecuteSpanCarriesStormAttributes() throws Exception {
+        List<SpanData> spans = runSpoutToSink(true, 1, 1, 0); // 1 tuple, 1 
sink task, no ticks
+
+        Attributes attributes = named(spans, "sink 
execute").get(0).getAttributes();
+        assertEquals(topologyName, 
attributes.get(AttributeKey.stringKey("storm.topology.name")));
+        String topologyId = 
attributes.get(AttributeKey.stringKey("storm.topology.id"));
+        assertTrue(topologyId.startsWith(topologyName), topologyId);
+        assertEquals("sink", 
attributes.get(AttributeKey.stringKey("storm.component.id")));
+        assertEquals(sinkTaskId, 
attributes.get(AttributeKey.longKey("storm.task.id")));
+        assertEquals("spout", 
attributes.get(AttributeKey.stringKey("storm.source.component.id")));
+        assertEquals("default", 
attributes.get(AttributeKey.stringKey("storm.source.stream.id")));
+        assertEquals(Utils.hostname(), 
attributes.get(AttributeKey.stringKey("storm.worker.host")));
+        assertEquals(WORKER_PORT_BY_COMPONENT.get("sink").longValue(),
+            attributes.get(AttributeKey.longKey("storm.worker.port")));
+    }
+
+    @Test
+    public void testAckRecordsAnOutcomeSpanUnderTheRoot() throws Exception {
+        // spout emit, sink execute, spout ack
+        List<SpanData> spans = runWithSink(SinkOutcome.ACK, conf(true), 3);
+
+        SpanData emit = named(spans, "spout emit").get(0);
+        SpanData ack = named(spans, "spout ack").get(0);
+        assertEquals(emit.getSpanId(), ack.getParentSpanId());
+        assertEquals(StatusCode.UNSET, ack.getStatus().getStatusCode());
+    }
+
+    @Test
+    public void testFailRecordsErrorSpansInTheBoltAndAtTheSpout() throws 
Exception {
+        // spout emit, sink execute, sink fail, spout fail
+        List<SpanData> spans = runWithSink(SinkOutcome.FAIL, conf(true), 4);
+
+        SpanData emit = named(spans, "spout emit").get(0);
+        SpanData execute = named(spans, "sink execute").get(0);
+        assertEquals(1, named(spans, "sink fail").size());
+        SpanData sinkFail = named(spans, "sink fail").get(0);
+        SpanData spoutFail = named(spans, "spout fail").get(0);
+        assertEquals(execute.getSpanId(), sinkFail.getParentSpanId());
+        assertEquals(emit.getSpanId(), spoutFail.getParentSpanId());
+        assertEquals(StatusCode.ERROR, sinkFail.getStatus().getStatusCode());
+        assertEquals(StatusCode.ERROR, spoutFail.getStatus().getStatusCode());
+    }
+
+    @Test
+    public void testBoltFailIsRecordedWithoutAckers() throws Exception {
+        Config conf = conf(true);
+        conf.put(Config.TOPOLOGY_ACKER_EXECUTORS, 0);
+        // spout emit, sink execute, sink fail; without ackers the spout 
records no outcome
+        List<SpanData> spans = runWithSink(SinkOutcome.FAIL, conf, 3);
+
+        SpanData sinkFail = named(spans, "sink fail").get(0);
+        assertEquals(StatusCode.ERROR, sinkFail.getStatus().getStatusCode());
+        assertTrue(named(spans, "spout ack").isEmpty());
+    }
+
+    @Test
+    public void testTimeoutRecordsAnErrorSpanAtTheSpout() throws Exception {
+        Config conf = conf(true);
+        conf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 2);
+        // spout emit, sink execute, spout timeout
+        List<SpanData> spans = runWithSink(SinkOutcome.HOLD, conf, 3);
+
+        SpanData emit = named(spans, "spout emit").get(0);
+        SpanData timeout = named(spans, "spout timeout").get(0);
+        assertEquals(emit.getSpanId(), timeout.getParentSpanId());
+        assertEquals(StatusCode.ERROR, timeout.getStatus().getStatusCode());
+        assertTrue(named(spans, "spout fail").isEmpty());
+    }
+
+    @Test
+    public void testAnchoredEmitContinuesTheTrace() throws Exception {
+        // per tuple: spout emit, middle execute, sink execute
+        List<SpanData> spans = runThroughMiddle(EmitMode.ANCHORED, 2, 6, 2);
+
+        assertSinkExecutesAreChildrenOfMiddleExecutes(spans, 2);
+    }
+
+    @Test
+    public void testAnchorsCarryingTheSameSpanAddNoMergeSpan() throws 
Exception {
+        // per tuple: spout emit, middle execute, sink execute; the emit 
anchors the input twice
+        List<SpanData> spans = runThroughMiddle(EmitMode.ANCHORED_TWICE, 2, 6, 
2);
+
+        assertTrue(named(spans, "middle emit").isEmpty());
+        assertSinkExecutesAreChildrenOfMiddleExecutes(spans, 2);
+    }
+
+    @Test
+    public void testEmitAnchoredToTwoTracedTuplesStartsARootWithTwoLinks() 
throws Exception {
+        // 2 spout emits, 2 middle executes, 1 merge span, 1 sink execute
+        List<SpanData> spans = runThroughMiddle(EmitMode.JOIN, 2, 6, 1);
+
+        List<SpanData> merges = named(spans, "middle emit");
+        assertEquals(1, merges.size());
+        SpanData merge = merges.get(0);
+        assertFalse(merge.getParentSpanContext().isValid(), "a merge starts a 
new trace");
+        Set<String> linked = merge.getLinks().stream()
+            .map(link -> 
link.getSpanContext().getSpanId()).collect(Collectors.toSet());
+        assertEquals(byId(named(spans, "middle execute")).keySet(), linked);
+        List<SpanData> sinks = named(spans, "sink execute");
+        assertEquals(1, sinks.size());
+        assertEquals(merge.getSpanId(), sinks.get(0).getParentSpanId());
+    }
+
+    @Test
+    public void testUnanchoredEmitCarriesNoContext() throws Exception {
+        // per tuple: spout emit, middle execute; the sink gets untraced tuples
+        List<SpanData> spans = runThroughMiddle(EmitMode.UNANCHORED, 2, 4, 2);
+
+        assertEquals(2, named(spans, "middle execute").size());
+        assertTrue(named(spans, "sink execute").isEmpty());
+        assertTrue(RECEIVED_TRACE_IDS.isEmpty(), "the sink's tuples carry no 
context");
+    }
+
+    @Test
+    public void testDelayedEmitsFromAnotherThreadKeepTheirOwnParents() throws 
Exception {
+        // per tuple: spout emit, middle execute, sink execute
+        List<SpanData> spans = runThroughMiddle(EmitMode.ASYNC_REVERSED, 2, 6, 
2);
+
+        assertNull(EMITTER_THREAD_FAILURE.get());
+        Map<String, SpanData> byId = byId(spans);
+        for (Object value : MIDDLE_SPAN_BY_VALUE.keySet()) {
+            SpanData sink = byId.get(SINK_SPAN_BY_VALUE.get(value));
+            assertEquals(MIDDLE_SPAN_BY_VALUE.get(value), 
sink.getParentSpanId(),
+                "the sink span of " + value + " is a child of the middle span 
of " + value);
+        }
+    }
+
+    @Test
+    public void testUnsampledContextsPropagateAndNothingIsExported() throws 
Exception {
+        // parent-based: a sampled flag flipped on the way would export the 
middle or sink span
+        InMemorySpanExporter exporter = InMemorySpanExporter.create();
+        SdkTracerProvider tracerProvider = SdkTracerProvider.builder()
+            .setSampler(Sampler.parentBased(Sampler.alwaysOff()))
+            .addSpanProcessor(SimpleSpanProcessor.create(exporter))
+            .build();
+        try (OpenTelemetrySdk sdk =
+                 
OpenTelemetrySdk.builder().setTracerProvider(tracerProvider).build()) {
+            GlobalOpenTelemetry.resetForTest();
+            GlobalOpenTelemetry.set(sdk);
+            // spans go to this SDK, not OTEL: wait for the sink only
+            runThroughMiddle(EmitMode.ANCHORED, 2, 0, 2);
+        } finally {
+            GlobalOpenTelemetry.resetForTest();
+            GlobalOpenTelemetry.set(OTEL.getOpenTelemetry());
+        }
+
+        assertTrue(exporter.getFinishedSpanItems().isEmpty());
+        assertEquals(2, UNSAMPLED_CONTEXTS_RECEIVED.get(), "the sink got 
unsampled contexts");
+        assertEquals(MIDDLE_TRACE_IDS, RECEIVED_TRACE_IDS, "the traces 
continue to the sink");
+        // different workers: the contexts went through the serializer
+        assertNotEquals(WORKER_PORT_BY_COMPONENT.get("middle"),
+            WORKER_PORT_BY_COMPONENT.get("sink"));
+    }
+
+    private static void 
assertSinkExecutesAreChildrenOfMiddleExecutes(List<SpanData> spans,
+        int count) {
+        Map<String, SpanData> middles = byId(named(spans, "middle execute"));
+        List<SpanData> sinks = named(spans, "sink execute");
+        assertEquals(count, middles.size());
+        assertEquals(count, sinks.size());
+        for (SpanData sink : sinks) {
+            SpanData middle = middles.get(sink.getParentSpanId());
+            assertNotNull(middle, "the emit carries the middle execute span as 
parent");
+            assertEquals(middle.getTraceId(), sink.getTraceId());
+        }
+    }
+
+    /**
+     * Spout to sink (all grouping). With {@code tickSecs} positive, also 
waits for two ticks.
+     */
+    private List<SpanData> runSpoutToSink(boolean tracing, int count, int 
sinkTasks, int tickSecs)
+        throws Exception {
+        Config conf = conf(tracing);
+        if (tickSecs > 0) {
+            conf.put(Config.TOPOLOGY_TICK_TUPLE_FREQ_SECS, tickSecs);
+        }
+        // one emit span per tuple and one execute span per tuple and sink task
+        int expectedSpans = tracing ? count * (1 + sinkTasks) : 0;

Review Comment:
   Right, the counts now include the spout ack spans.



##########
pom.xml:
##########
@@ -759,6 +760,13 @@
                 <type>pom</type>
                 <scope>import</scope>
             </dependency>
+            <dependency>
+                <groupId>io.opentelemetry</groupId>

Review Comment:
   Removed. Only the new module uses OTel, so hbase-client stays on 1.49.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to