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


##########
storm-client/src/jvm/org/apache/storm/executor/bolt/BoltOutputCollectorImpl.java:
##########
@@ -117,6 +138,48 @@ private List<Integer> boltEmit(String streamId, 
Collection<Tuple> anchors, List<
         return outTasks;
     }
 
+    /**
+     * Runs on the emitting thread. If the traced anchors share one span, 
returns their context. If
+     * they carry several, returns a new root linked to each (the SDK keeps up 
to 128 links by
+     * default), or null when no SDK is registered on this worker.
+     */
+    private Context traceContextFor(Collection<Tuple> anchors) {
+        Context first = null;
+        Set<SpanContext> linkedSpans = null;
+        for (Tuple anchor : anchors) {
+            Context context = anchor instanceof TupleImpl impl ? 
impl.getTraceContext() : null;
+            if (context == null) {
+                continue;
+            }
+            if (first == null) {
+                first = context;
+            } else {
+                if (linkedSpans == null) {
+                    linkedSpans = new LinkedHashSet<>();
+                    linkedSpans.add(Span.fromContext(first).getSpanContext());
+                }
+                linkedSpans.add(Span.fromContext(context).getSpanContext());
+            }
+        }
+        if (linkedSpans == null || linkedSpans.size() == 1) {
+            return first;
+        }
+        return executor.newRootContext(emitSpanName, linkedSpans);

Review Comment:
   Documented, so fine as is. Still: if all anchors share one trace id (two 
executions of the same upstream bolt within one tree), we could parent on one 
and link the rest instead of starting a new root. Otherwise the tail of a 
sampled tree is re-sampled by the root sampler, which with 
`parentbased_traceidratio` drops most of it. Windowed bolts will hit this on 
nearly every emit.



##########
storm-client/src/jvm/org/apache/storm/executor/Executor.java:
##########
@@ -781,6 +793,58 @@ public String getComponentId() {
         return componentId;
     }
 
+    /**
+     * Returns the tracer, or null until an OpenTelemetry SDK is registered as 
the global instance.
+     * Checking isSet() instead of calling get() leaves the global unset, so 
an SDK registered later
+     * is still used. Safe to call from any thread.
+     */
+    protected Tracer tracer() {
+        Tracer current = tracer;
+        if (current == null && GlobalOpenTelemetry.isSet()) {

Review Comment:
   The agent only bridges `isSet()` since 2.23.0 
(open-telemetry/opentelemetry-java-instrumentation#15620). Before that it 
intercepts `get()` only, the app-side global stays unset, and this returns 
false forever. With an older agent tracing is silently off and nothing is 
logged. `Tracing.md` should state the minimum agent version, and a one-time 
INFO/WARN when tracing is enabled but no SDK is found would make this 
diagnosable.
   
   Minor: until an SDK is found, every spout emit goes through the synchronized 
`isSet()`.



##########
storm-client/src/jvm/org/apache/storm/utils/TupleUtils.java:
##########
@@ -34,6 +36,18 @@ public static boolean isTick(Tuple tuple) {
                && 
Constants.SYSTEM_TICK_STREAM_ID.equals(tuple.getSourceStreamId());
     }
 
+    /**
+     * Returns the OpenTelemetry context to run work for this tuple under, so 
that spans created
+     * there join the tuple's trace, or {@link Context#root()} when the tuple 
carries none (see
+     * {@link Config#TOPOLOGY_TRACING_ENABLED}). Never null. The context holds 
span ids only: it
+     * parents new spans but gives no access to the execute span itself. 
Example:
+     * {@code pool.submit(TupleUtils.traceContext(input).wrap(task))}.
+     */
+    public static Context traceContext(Tuple tuple) {

Review Comment:
   This puts `io.opentelemetry.context.Context` into our public API (see the 
summary for a way around it).
   
   Nit: the description says this is the only new public method. 
`TupleImpl#get/setTraceContext`, `TupleInfo#get/setTraceContext` and 
`Executor#isTracingEnabled/newRootContext/recordOutcome` are public too.



##########
storm-client/src/jvm/org/apache/storm/executor/spout/SpoutOutputCollectorImpl.java:
##########
@@ -123,6 +128,12 @@ private List<Integer> sendSpoutMsg(String stream, 
List<Object> values, Object me
 
         final long rootId = needAck ? MessageId.generateId(random) : 0;
 
+        // checkpoint tuples of stateful bolts are system tuples: no trace
+        boolean traced = executor.isTracingEnabled()
+            && !CheckpointSpout.CHECKPOINT_STREAM_ID.equals(stream);
+        final Context traceContext =
+            traced ? executor.newRootContext(emitSpanName, 
Collections.emptyList()) : null;

Review Comment:
   Always a new root, so a Kafka/JMS spout cannot continue a trace coming from 
message headers. Fine for a first step, but should be in the limits section of 
`Tracing.md` (baggage is dropped as well).
   
   The Trident coordinator streams (`$batch`, `$commit`, `$success`) are spout 
emits too, so each one starts its own trace. Only the checkpoint stream is 
excluded above.



##########
storm-client/src/jvm/org/apache/storm/serialization/KryoTupleSerializer.java:
##########
@@ -64,4 +77,30 @@ public byte[] serialize(Tuple tuple) {
             throw new RuntimeException(e);
         }
     }
+
+    /**
+     * Appends the trace context after the values: header byte (version, 
tracestate bit), 16-byte
+     * trace id, 8-byte span id, trace flags byte, then the tracestate if not 
empty. Readers that
+     * stop after the values ignore these bytes.
+     */
+    private static void writeTraceContext(Output out, Context traceContext) {
+        if (traceContext == null) {
+            return;
+        }
+        SpanContext span = Span.fromContext(traceContext).getSpanContext();
+        // isValid() ignores the sampled flag: unsampled contexts propagate too
+        if (!span.isValid()) {
+            return;
+        }
+        TraceState traceState = span.getTraceState();
+        out.writeByte(TRACE_CONTEXT_VERSION | (traceState.isEmpty() ? 0 : 
HAS_TRACE_STATE));

Review Comment:
   `DefaultStateSerializer.TupleSerializer` uses this serializer, so persistent 
windowed bolts now store the context in state (e.g. Redis). After a restore, 
emits get parented to spans that may be days old. Either clear it there or 
document it.
   
   Small one: the first byte is type and version at once, without a length, so 
a reader cannot skip an unknown extension. Cheap to reserve a scheme now while 
the format is new.



##########
storm-client/src/jvm/org/apache/storm/serialization/KryoTupleDeserializer.java:
##########
@@ -84,12 +94,55 @@ private TupleImpl deserializeTuple(byte[] data) {
             String streamName = ids.getStreamName(componentName, streamId);
             MessageId id = MessageId.deserialize(kryoInput);
             List<Object> values = kryo.deserializeFrom(kryoInput);
-            return new TupleImpl(context, values, componentName, taskId, 
streamName, id);
+            TupleImpl tuple = new TupleImpl(context, values, componentName, 
taskId, streamName, id);
+            tuple.setTraceContext(readTraceContext(kryoInput));
+            return tuple;
         } catch (IOException e) {
             throw new RuntimeException(FAILED_TO_DESERIALIZE_TUPLE, e);
         }
     }
 
+    /**
+     * Reads the trace context written after the values. Null when absent, of 
an unknown version
+     * or unreadable, so a bad extension never fails the tuple. Invalid 
tracestate entries are
+     * dropped.
+     */
+    private static Context readTraceContext(Input in) {
+        if (in.position() == in.limit()) {
+            return null;
+        }
+        try {
+            int header = in.readByte() & 0xFF;
+            int version = header & ~KryoTupleSerializer.HAS_TRACE_STATE;
+            if (version != KryoTupleSerializer.TRACE_CONTEXT_VERSION) {
+                return null;
+            }
+            String traceId = TraceId.fromBytes(in.readBytes(TRACE_ID_BYTES));
+            String spanId = SpanId.fromBytes(in.readBytes(SPAN_ID_BYTES));
+            TraceFlags flags = TraceFlags.fromByte(in.readByte());
+            TraceState traceState = TraceState.getDefault();
+            if ((header & KryoTupleSerializer.HAS_TRACE_STATE) != 0) {
+                String[] entries = in.readString().split(",");

Review Comment:
   Length comes from the wire and Kryo allocates `char[length]` before checking 
remaining bytes. An OOM is not caught by the `catch (RuntimeException)` below. 
Inter-worker traffic is trusted, so no security issue, but capping at the W3C 
limit (512) is cheap.



##########
storm-client/src/jvm/org/apache/storm/executor/Executor.java:
##########
@@ -781,6 +793,58 @@ public String getComponentId() {
         return componentId;
     }
 
+    /**
+     * Returns the tracer, or null until an OpenTelemetry SDK is registered as 
the global instance.
+     * Checking isSet() instead of calling get() leaves the global unset, so 
an SDK registered later
+     * is still used. Safe to call from any thread.
+     */
+    protected Tracer tracer() {
+        Tracer current = tracer;
+        if (current == null && GlobalOpenTelemetry.isSet()) {
+            current = GlobalOpenTelemetry.get().getTracer("org.apache.storm");
+            tracer = current;
+        }
+        return current;
+    }
+
+    public boolean isTracingEnabled() {
+        return tracingEnabled;
+    }
+
+    /**
+     * Starts and immediately ends a root span linked to {@code links} and 
returns its context, or
+     * null when no SDK is registered or the span is not valid.
+     */
+    public Context newRootContext(String spanName, Collection<SpanContext> 
links) {
+        Tracer current = tracer();
+        if (current == null) {
+            return null;
+        }
+        SpanBuilder builder = current.spanBuilder(spanName).setNoParent();
+        links.forEach(builder::addLink);
+        Span span = builder.startSpan();
+        span.end();
+        // keep only the ids: pending tuples hold this context until their 
tree completes
+        SpanContext ids = span.getSpanContext();
+        return ids.isValid() ? Context.root().with(Span.wrap(ids)) : null;
+    }
+
+    /**
+     * Records a span under {@code parent}, started and ended at once, with 
status ERROR when
+     * {@code error}. Nothing is recorded until an SDK is registered.
+     */
+    public void recordOutcome(Context parent, String spanName, boolean error) {
+        Tracer current = tracer();
+        if (current == null) {
+            return;
+        }
+        Span span = 
current.spanBuilder(spanName).setParent(parent).startSpan();
+        if (error) {
+            span.setStatus(StatusCode.ERROR);

Review Comment:
   Zero duration, and the complete latency is known in the spout at this point. 
Recording it as an attribute would make these spans more useful.



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

Review Comment:
   Importing the BOM here also moves hbase-client's transitive 
`opentelemetry-api` from 1.49 to 1.66 in storm-autocreds, hdfs etc. Should be 
fine given their compat policy, but please mention it in the description. Goes 
away with the split proposed in the summary.



##########
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:
   Does not count the spout ack spans. `executeInSpan` ends the span in 
`finally`, after the bolt already acked, so the wait can be satisfied before 
the last sink execute span is exported. Narrow window, but on a loaded CI node 
the exact count asserts fail. Same for the expected counts passed to 
`runThroughMiddle`.



-- 
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