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


##########
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:
   Done: `isSet()` at most once a second until an SDK is there, plus one WARN. 
Tracing.md says agent 2.23.0+.



##########
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:
   Done: anchors from one trace now give a child of the first one, linked to 
the others. Different traces still get a new root.



##########
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:
   Gone: storm-client only sees an opaque `Object`; the helper moved to the 
module as `StormTracing.context`.



##########
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:
   Added to the limits in Tracing.md, with baggage and the Trident 
`$batch`/`$commit`/`$success` streams.



##########
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:
   Done: the length is checked against the remaining bytes. A tracestate over 
512 isn't sent and reads as no context.



##########
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:
   Done: tag + varint length + bytes, unknown tags skipped. State serializers 
have no tracer, so state keeps no context (test added).



##########
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:
   Done: outcome spans now run from emit to outcome and carry 
`storm.tuple.latency_ms`.



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