using System.Diagnostics; namespace GB5Shared.Telemetry.Dapr; /// /// Helpers for instrumenting Dapr pub/sub and state-store operations with OpenTelemetry spans. /// Use these factory methods to wrap DaprClient.PublishEventAsync, /// GetStateAsync, and event-handler entry points so every event can be traced /// end-to-end across service boundaries. /// /// /// Publishing: /// /// using var publishActivity = DaprEventTracing.StartPublishActivity("order-placed", "pubsub", metadata); /// try { /// await DaprClient.PublishEventAsync("pubsub", "order-placed", payload, metadata, ct); /// publishActivity?.SetStatus(ActivityStatusCode.Ok); /// } catch (Exception ex) { /// publishActivity?.RecordException(ex); /// publishActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); /// throw; /// } /// /// Subscribing (in a controller action): /// /// using var subActivity = DaprEventTracing.StartSubscribeActivity("order-placed", "pubsub", eventMetadata); /// try { /// await ProcessOrderAsync(cloudEvent.Data, ct); /// subActivity?.SetStatus(ActivityStatusCode.Ok); /// } catch (Exception ex) { /// subActivity?.RecordException(ex); /// subActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); /// throw; /// } /// /// public static class DaprEventTracing { private const string TraceParentKey = "traceparent"; private const string TraceStateKey = "tracestate"; // ── Producer (publish) ──────────────────────────────────────────────────── /// /// Opens a producer span for a Dapr publish operation and injects the W3C traceparent /// into so downstream consumers can link the trace. /// /// The Dapr pub/sub topic name. /// The Dapr pub/sub component name (e.g. "pubsub"). /// /// Optional metadata dictionary passed to DaprClient.PublishEventAsync. /// When supplied, trace context is injected automatically. /// public static Activity? StartPublishActivity( string topic, string pubSubName, IDictionary? metadata = null) { var activity = GB5ActivitySources.DaprPublish.StartActivity( $"dapr publish {pubSubName}/{topic}", ActivityKind.Producer); if (activity is null) return null; activity.SetTag("messaging.system", "dapr"); activity.SetTag("messaging.operation", "publish"); activity.SetTag("messaging.destination", topic); activity.SetTag("messaging.destination_kind", "topic"); activity.SetTag("dapr.pubsub.name", pubSubName); if (metadata is not null) InjectTraceContext(activity, metadata); return activity; } // ── Consumer (subscribe) ────────────────────────────────────────────────── /// /// Opens a consumer span for a Dapr subscribe event handler and extracts the remote /// W3C traceparent from to link back to the producer's trace. /// /// The Dapr pub/sub topic name. /// The Dapr pub/sub component name. /// /// Metadata from the received CloudEvent. When it contains a traceparent key /// the new span is created as a child of the producer's span automatically. /// public static Activity? StartSubscribeActivity( string topic, string pubSubName, IDictionary? metadata = null) { var parentContext = metadata is not null ? ExtractTraceContext(metadata) : default; return StartSubscribeActivityCore(topic, pubSubName, parentContext); } /// /// Opens a consumer span for a Dapr subscribe event handler, parenting it directly from a /// W3C traceparent/tracestate pair carried in the message payload itself /// (e.g. OutboxEventMessage.TraceParent/TraceState). Use this overload instead /// of the metadata-dictionary one when the caller already has a strongly-typed DTO — /// it sidesteps any uncertainty about whether the Dapr broker/CloudEvent envelope forwards /// custom metadata keys through to the subscriber. /// public static Activity? StartSubscribeActivity( string topic, string pubSubName, string? traceParent, string? traceState = null) { var parentContext = traceParent is not null && ActivityContext.TryParse(traceParent, traceState, isRemote: true, out var ctx) ? ctx : default; return StartSubscribeActivityCore(topic, pubSubName, parentContext); } private static Activity? StartSubscribeActivityCore(string topic, string pubSubName, ActivityContext parentContext) { var activity = parentContext != default ? GB5ActivitySources.DaprSubscribe.StartActivity( $"dapr receive {pubSubName}/{topic}", ActivityKind.Consumer, parentContext) : GB5ActivitySources.DaprSubscribe.StartActivity( $"dapr receive {pubSubName}/{topic}", ActivityKind.Consumer); if (activity is null) return null; activity.SetTag("messaging.system", "dapr"); activity.SetTag("messaging.operation", "receive"); activity.SetTag("messaging.destination", topic); activity.SetTag("messaging.destination_kind", "topic"); activity.SetTag("dapr.pubsub.name", pubSubName); return activity; } // ── State store (cache) ─────────────────────────────────────────────────── /// /// Opens a client-side span for a Dapr state-store read, write, or delete operation. /// /// One of: GET, SET, DELETE. /// The Dapr state store component name (e.g. "statestore"). /// The state key being accessed. public static Activity? StartCacheActivity( string operation, string stateStoreName, string key) { var activity = GB5ActivitySources.Cache.StartActivity( $"dapr statestore {operation.ToUpperInvariant()} {stateStoreName}", ActivityKind.Client); if (activity is null) return null; activity.SetTag("db.system", "dapr_statestore"); activity.SetTag("db.operation", operation.ToUpperInvariant()); activity.SetTag("dapr.statestore.name", stateStoreName); activity.SetTag("dapr.statestore.key", key); return activity; } // ── W3C trace context propagation ──────────────────────────────────────── /// /// Writes the current activity's W3C traceparent (and optional tracestate) /// into a Dapr metadata dictionary so they survive the pub/sub message boundary. /// public static void InjectTraceContext(Activity activity, IDictionary metadata) { if (activity.Id is not null) metadata[TraceParentKey] = activity.Id; if (!string.IsNullOrEmpty(activity.TraceStateString)) metadata[TraceStateKey] = activity.TraceStateString; } /// /// Reads traceparent / tracestate from a Dapr metadata dictionary and /// returns the parsed , or default if not present. /// public static ActivityContext ExtractTraceContext(IDictionary metadata) { if (!metadata.TryGetValue(TraceParentKey, out var traceParent)) return default; metadata.TryGetValue(TraceStateKey, out var traceState); return ActivityContext.TryParse(traceParent, traceState, isRemote: true, out var ctx) ? ctx : default; } }