diff --git a/docs/TELEMETRY.md b/docs/TELEMETRY.md index 001cb20ff..1906b020a 100644 --- a/docs/TELEMETRY.md +++ b/docs/TELEMETRY.md @@ -140,6 +140,121 @@ changing connection behavior. The phase spans intentionally do not trace signing, serialization, response parsing, capability details, or token persistence as separate operations. +## Chat lifecycle + +The tray exports native chat lifecycle diagnostics when an endpoint is configured: + +- traces: `openclaw.chat.turn`, `openclaw.chat.queue.wait`, `openclaw.chat.send`, + `openclaw.chat.response.wait`, `openclaw.chat.response.receive`, + `openclaw.chat.history.load`, and `openclaw.chat.history.backfill` +- counters: `openclaw.chat.turns`, `openclaw.chat.send.attempts`, + `openclaw.chat.history.loads`, `openclaw.chat.history.backfills`, and + `openclaw.chat.remote_turns.dropped`, and + `openclaw.chat.terminal_events.dropped` +- duration histograms: `openclaw.chat.turn.duration`, + `openclaw.chat.queue.wait.duration`, `openclaw.chat.send.duration`, + `openclaw.chat.response.wait.duration`, + `openclaw.chat.response.receive.duration`, + `openclaw.chat.history.load.duration`, and + `openclaw.chat.history.backfill.duration` + +A local turn starts when the tray admits a valid request to direct dispatch or its +local queue. An observed remote turn starts at a gateway lifecycle start carrying +a run ID. Turn correlation uses message, run, and thread identifiers only inside +the process; those identifiers are never attached to exported signals. Each turn +span is explicitly created as a root, while each sampled local send attempt is +explicitly parented to its turn. + +Turn completion is exactly once. Assistant final, lifecycle end, lifecycle error, +send rejection, queue cancellation, explicit abort, reset/supersession, +disconnect, and disposal race through an atomic tracker; the first applicable +terminal transition removes correlation state, and later duplicate signals are +ignored. Assistant-final events do not contain a run ID, so the provider captures +the active run under its existing state lock before completing telemetry. A remote +lifecycle start without a run ID is not traced using an unsafe thread fallback; +it increments `openclaw.chat.remote_turns.dropped` with the finite reason +`missing_run_id`. A missing-run start for an already-dispatched local turn is not +misclassified as a dropped remote turn. + +Terminal lifecycle events are never matched by thread alone. If a terminal event +has no run ID or conflicts with the active run, it cannot complete a potentially +newer turn. The provider drops malformed lifecycle and legacy job terminals +before they can clear active-run, timeline, queue, or telemetry state, logs a +content-free warning, and increments `openclaw.chat.terminal_events.dropped`. +A later terminal carrying the exact active run ID can still complete the turn; +otherwise unresolved turns remain eligible for safe reset, disconnect, or +disposal cleanup. Assistant-final chat events are separate: their protocol shape +does not carry a run ID, so the provider captures the authoritative active run +under its state lock rather than accepting a thread-only agent terminal. + +Each `openclaw.chat.send` span represents one `chat.send` RPC attempt. A valid +retryable deferral is `outcome=success` with admission status `deferred`, because +the RPC completed and returned a recognized decision. Local requeue is not a +separate exported admission status. Accepted responses use `accepted`; terminal +rejection, cancellation, and exceptions use `rejected`, `canceled`, and +`exception`. Unknown values map to `other`. + +Response timing is split into two sibling child spans under the turn: + +- `openclaw.chat.response.wait` starts at accepted local admission or observed + lifecycle start and ends at the first recognized assistant, reasoning, or tool + output. +- `openclaw.chat.response.receive` starts at that first inbound event and ends + with the turn's authoritative terminal transition. + +If a turn terminates before visible output, the wait span closes with +`openclaw.chat.response.first_output=none` and no receive span is emitted. +Repeated chunks do not create additional spans. Phase duration metrics are +recorded at turn completion so they carry the final bounded turn outcome. A wait +span that reaches first output reports its own phase outcome as `success`; its +duration metric still uses the enclosing turn's final outcome for aggregation. +Output received before accepted admission or lifecycle start does not synthesize +a wait or receive phase. +Unknown admission statuses, routine status/error events, and unknown future +event types do not start or transition response phases. The `other` output value +is reserved for future event types only after they are explicitly reviewed and +classified as visible response output. + +Each contiguous local queue or requeue period emits an +`openclaw.chat.queue.wait` sibling span under the turn. A segment that reaches +dispatch completes with `outcome=success`; a segment still queued when the turn +terminates uses the turn's final outcome. Deferred sends therefore show multiple +queue-wait spans rather than one span that incorrectly includes intervening send +attempts. + +The queue-wait duration metric remains cumulative across all queue/retry segments. +The tray adds each segment when dispatch begins and emits the total when the turn +completes so it can carry the final outcome. Direct sends accepted on their first +attempt emit neither queue-wait spans nor queue-wait measurements. Consequently, +the metric timestamp is the turn completion time, not the instant queue congestion +occurred. + +Full transcript loads and targeted remote-message backfills are separate +operations. Full loads use source `initial` or `forced`. Backfills use the finite +reason `remote_turn` or `reset_reconciliation`. + +Chat attributes are restricted to: + +- `openclaw.source`: `local`, `remote`, `initial`, or `forced`, as applicable +- `openclaw.outcome`: `success`, `failure`, or `canceled` +- `openclaw.reason`: `assistant_final`, `lifecycle_end`, `lifecycle_error`, + `send_rejected`, `queued_canceled`, `abort_requested`, `reset`, `superseded`, + `disconnected`, `disposed`, or `other` +- `openclaw.chat.admission.status`: `accepted`, `deferred`, `rejected`, + `canceled`, `exception`, or `other` +- `openclaw.chat.backfill.reason`: `remote_turn` or `reset_reconciliation` +- `openclaw.chat.remote_turn.drop.reason`: `missing_run_id` +- `openclaw.chat.terminal_event.drop.reason`: `missing_run_id` or + `mismatched_run_id` +- `openclaw.chat.response.first_output`: `none`, `assistant`, `reasoning`, + `tool`, or `other` +- `error.type`: exception type only, never the exception message + +Chat telemetry does not export prompts, responses, transcript contents, IDs, +model/provider names, attachment metadata, filenames, tool names, token usage, +URLs, error messages, or local chat log text. No chat log category is added to +the OpenTelemetry log allowlist. + ## Endpoint handling The endpoint setting is a collector endpoint, not a credential or request-parameter store. Accept plain `http` and `https` collector URLs with optional path prefixes. Reject URLs with embedded user info, query strings, or fragments. diff --git a/src/OpenClaw.Shared/Telemetry/OpenClawTelemetry.cs b/src/OpenClaw.Shared/Telemetry/OpenClawTelemetry.cs index d88d376a3..092014659 100644 --- a/src/OpenClaw.Shared/Telemetry/OpenClawTelemetry.cs +++ b/src/OpenClaw.Shared/Telemetry/OpenClawTelemetry.cs @@ -82,6 +82,7 @@ public static class OpenClawTelemetry var previous = Activity.Current; try { + Activity.Current = null; var activity = source.ToActivitySource().StartActivity(spanName, kind, parentContext); ApplyTags(activity, tags); return activity; diff --git a/src/OpenClaw.Tray.WinUI/Chat/ChatTelemetryTracker.cs b/src/OpenClaw.Tray.WinUI/Chat/ChatTelemetryTracker.cs new file mode 100644 index 000000000..e888ab2fb --- /dev/null +++ b/src/OpenClaw.Tray.WinUI/Chat/ChatTelemetryTracker.cs @@ -0,0 +1,1129 @@ +using System.Diagnostics; +using System.Diagnostics.Metrics; +using OpenClaw.Shared.Telemetry; + +namespace OpenClawTray.Chat; + +internal enum ChatTelemetryOutcome +{ + Success, + Failure, + Canceled, +} + +internal enum ChatTurnTelemetryReason +{ + AssistantFinal, + LifecycleEnd, + LifecycleError, + SendRejected, + QueuedCanceled, + AbortRequested, + Reset, + Superseded, + Disconnected, + Disposed, + Other, +} + +internal enum ChatAdmissionTelemetryStatus +{ + Accepted, + Deferred, + Rejected, + Canceled, + Exception, + Other, +} + +internal enum ChatHistoryTelemetrySource +{ + Initial, + Forced, +} + +internal enum ChatBackfillTelemetryReason +{ + RemoteTurn, + ResetReconciliation, +} + +internal enum ChatTerminalEventDropReason +{ + MissingRunId, + MismatchedRunId, +} + +internal enum ChatResponseOutputKind +{ + None, + Assistant, + Reasoning, + Tool, + Other, +} + +internal sealed class ChatTelemetryTracker +{ + internal const string TurnSpanName = "openclaw.chat.turn"; + internal const string QueueWaitSpanName = "openclaw.chat.queue.wait"; + internal const string SendSpanName = "openclaw.chat.send"; + internal const string ResponseWaitSpanName = "openclaw.chat.response.wait"; + internal const string ResponseReceiveSpanName = "openclaw.chat.response.receive"; + internal const string HistoryLoadSpanName = "openclaw.chat.history.load"; + internal const string HistoryBackfillSpanName = "openclaw.chat.history.backfill"; + + internal const string TurnsMetricName = "openclaw.chat.turns"; + internal const string TurnDurationMetricName = "openclaw.chat.turn.duration"; + internal const string QueueWaitDurationMetricName = "openclaw.chat.queue.wait.duration"; + internal const string SendAttemptsMetricName = "openclaw.chat.send.attempts"; + internal const string SendDurationMetricName = "openclaw.chat.send.duration"; + internal const string ResponseWaitDurationMetricName = "openclaw.chat.response.wait.duration"; + internal const string ResponseReceiveDurationMetricName = "openclaw.chat.response.receive.duration"; + internal const string HistoryLoadsMetricName = "openclaw.chat.history.loads"; + internal const string HistoryLoadDurationMetricName = "openclaw.chat.history.load.duration"; + internal const string HistoryBackfillsMetricName = "openclaw.chat.history.backfills"; + internal const string HistoryBackfillDurationMetricName = "openclaw.chat.history.backfill.duration"; + internal const string DroppedRemoteTurnsMetricName = "openclaw.chat.remote_turns.dropped"; + internal const string DroppedTerminalEventsMetricName = "openclaw.chat.terminal_events.dropped"; + + internal const string AdmissionStatusTag = "openclaw.chat.admission.status"; + internal const string BackfillReasonTag = "openclaw.chat.backfill.reason"; + internal const string DroppedRemoteTurnReasonTag = "openclaw.chat.remote_turn.drop.reason"; + internal const string DroppedTerminalEventReasonTag = "openclaw.chat.terminal_event.drop.reason"; + internal const string FirstOutputKindTag = "openclaw.chat.response.first_output"; + + private const string SourceLocal = "local"; + private const string SourceRemote = "remote"; + private const string MissingRunId = "missing_run_id"; + private const string MismatchedRunId = "mismatched_run_id"; + + private static readonly Counter Turns = OpenClawTelemetry.CreateCounter( + TurnsMetricName, + unit: "{turn}", + description: "Number of observed OpenClaw chat turns."); + private static readonly Histogram TurnDuration = OpenClawTelemetry.CreateHistogram( + TurnDurationMetricName, + unit: "ms", + description: "Duration of observed OpenClaw chat turns."); + private static readonly Histogram QueueWaitDuration = OpenClawTelemetry.CreateHistogram( + QueueWaitDurationMetricName, + unit: "ms", + description: "Cumulative local queue dwell for completed OpenClaw chat turns."); + private static readonly Counter SendAttempts = OpenClawTelemetry.CreateCounter( + SendAttemptsMetricName, + unit: "{attempt}", + description: "Number of OpenClaw chat.send RPC attempts."); + private static readonly Histogram SendDuration = OpenClawTelemetry.CreateHistogram( + SendDurationMetricName, + unit: "ms", + description: "Duration of OpenClaw chat.send RPC attempts."); + private static readonly Histogram ResponseWaitDuration = OpenClawTelemetry.CreateHistogram( + ResponseWaitDurationMetricName, + unit: "ms", + description: "Duration from chat admission or lifecycle start to the first accepted inbound output."); + private static readonly Histogram ResponseReceiveDuration = OpenClawTelemetry.CreateHistogram( + ResponseReceiveDurationMetricName, + unit: "ms", + description: "Duration from the first accepted inbound output to terminal chat completion."); + private static readonly Counter HistoryLoads = OpenClawTelemetry.CreateCounter( + HistoryLoadsMetricName, + unit: "{load}", + description: "Number of full OpenClaw chat history loads."); + private static readonly Histogram HistoryLoadDuration = OpenClawTelemetry.CreateHistogram( + HistoryLoadDurationMetricName, + unit: "ms", + description: "Duration of full OpenClaw chat history loads."); + private static readonly Counter HistoryBackfills = OpenClawTelemetry.CreateCounter( + HistoryBackfillsMetricName, + unit: "{backfill}", + description: "Number of targeted OpenClaw chat history backfills."); + private static readonly Histogram HistoryBackfillDuration = OpenClawTelemetry.CreateHistogram( + HistoryBackfillDurationMetricName, + unit: "ms", + description: "Duration of targeted OpenClaw chat history backfills."); + private static readonly Counter DroppedRemoteTurns = OpenClawTelemetry.CreateCounter( + DroppedRemoteTurnsMetricName, + unit: "{turn}", + description: "Number of remote OpenClaw chat turns not traced because safe correlation was unavailable."); + private static readonly Counter DroppedTerminalEvents = OpenClawTelemetry.CreateCounter( + DroppedTerminalEventsMetricName, + unit: "{event}", + description: "Number of terminal OpenClaw chat events not applied because safe correlation was unavailable."); + + private readonly object _gate = new(); + private readonly Dictionary _turnsByMessageId = new(StringComparer.Ordinal); + private readonly Dictionary _turnsByRunId = new(StringComparer.Ordinal); + + public void StartLocalTurn(string messageId, string threadId, bool queued) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + + lock (_gate) + { + if (_turnsByMessageId.ContainsKey(messageId)) + return; + + var state = new TurnState( + messageId, + threadId, + SourceLocal, + OpenClawTelemetry.StartDetachedActivity( + TurnSpanName, + default(ActivityContext), + [OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, SourceLocal)]), + Stopwatch.GetTimestamp()); + if (queued) + StartQueueSegmentLocked(state); + + _turnsByMessageId.Add(messageId, state); + } + } + + public void DispatchLocalTurn(string messageId, string provisionalRunId) + { + CompleteQueueDispatch(PrepareDispatchLocalTurn(messageId, provisionalRunId)); + } + + public QueuePhaseCompletion? PrepareDispatchLocalTurn( + string messageId, + string provisionalRunId) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + ArgumentException.ThrowIfNullOrWhiteSpace(provisionalRunId); + + lock (_gate) + { + if (!_turnsByMessageId.TryGetValue(messageId, out var state)) + return null; + + var queueCompletion = EndQueueSegmentLocked( + state, + Stopwatch.GetTimestamp(), + ChatTelemetryOutcome.Success); + state.PendingQueueCompletion = queueCompletion; + state.IsDispatched = true; + BindRunLocked(state, provisionalRunId); + return queueCompletion; + } + } + + public void CompleteQueueDispatch(QueuePhaseCompletion? completion) => + CompleteQueuePhase(completion); + + public void RequeueLocalTurn(string messageId) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + + lock (_gate) + { + if (!_turnsByMessageId.TryGetValue(messageId, out var state)) + return; + + RemoveRunMappingsLocked(state); + state.IsDispatched = false; + StartQueueSegmentLocked(state); + } + } + + public void BindAcceptedRun(string messageId, string? runId) + { + if (string.IsNullOrWhiteSpace(runId)) + return; + + lock (_gate) + { + if (_turnsByMessageId.TryGetValue(messageId, out var state)) + BindRunLocked(state, runId); + } + } + + public void ObserveAdmissionAccepted(string messageId) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + lock (_gate) + { + if (_turnsByMessageId.TryGetValue(messageId, out var state)) + StartResponseWaitLocked(state); + } + } + + public void ObserveLifecycleStart(string threadId, string? runId, bool allowRemoteTurn = true) + { + ArgumentException.ThrowIfNullOrWhiteSpace(threadId); + if (string.IsNullOrWhiteSpace(runId)) + { + if (!allowRemoteTurn) + return; + + lock (_gate) + { + if (_turnsByMessageId.Values.Any( + state => state.Source == SourceLocal && + state.ThreadId == threadId && + state.IsDispatched)) + { + return; + } + } + + OpenClawTelemetry.Add( + DroppedRemoteTurns, + tags: [OpenClawTelemetryTag.String(DroppedRemoteTurnReasonTag, MissingRunId)]); + return; + } + + lock (_gate) + { + if (_turnsByRunId.TryGetValue(runId, out var existing)) + { + StartResponseWaitLocked(existing); + return; + } + + var pendingLocal = _turnsByMessageId.Values.FirstOrDefault( + state => state.Source == SourceLocal && + state.ThreadId == threadId && + state.IsDispatched); + if (pendingLocal is not null) + { + BindRunLocked(pendingLocal, runId); + StartResponseWaitLocked(pendingLocal); + return; + } + if (!allowRemoteTurn) + return; + + var remote = new TurnState( + messageId: null, + threadId, + SourceRemote, + OpenClawTelemetry.StartDetachedActivity( + TurnSpanName, + default(ActivityContext), + [OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, SourceRemote)]), + Stopwatch.GetTimestamp()); + BindRunLocked(remote, runId); + StartResponseWaitLocked(remote); + } + } + + public bool ObserveInboundOutput( + string threadId, + string? runId, + ChatResponseOutputKind outputKind) + { + ArgumentException.ThrowIfNullOrWhiteSpace(threadId); + ResponsePhaseCompletion? waitCompletion; + lock (_gate) + { + var state = ResolveTurnForOutputLocked(threadId, runId); + if (state is null || state.ReceivePhase is not null) + return false; + + if (state.WaitPhase is not { } waitPhase) + return false; + + var now = Stopwatch.GetTimestamp(); + state.FirstOutputKind = outputKind; + state.ResponseWaitDurationMilliseconds = + Stopwatch.GetElapsedTime(waitPhase.StartTimestamp, now).TotalMilliseconds; + waitCompletion = new ResponsePhaseCompletion( + waitPhase.Activity, + state.Source, + outputKind, + ChatTelemetryOutcome.Success, + Stopwatch.GetElapsedTime(waitPhase.StartTimestamp, now)); + state.PendingWaitCompletion = waitCompletion; + state.WaitPhase = null; + + state.ReceivePhase = StartPhase(state, ResponseReceiveSpanName, now); + } + + CompleteResponsePhase(waitCompletion); + return true; + } + + public ChatTelemetryOperation? StartSendAttempt(string messageId) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + + lock (_gate) + { + if (!_turnsByMessageId.TryGetValue(messageId, out var state)) + return null; + + var activity = state.Activity is not null + ? OpenClawTelemetry.StartDetachedActivity( + SendSpanName, + state.Activity.Context, + [OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, SourceLocal)]) + : OpenClawTelemetry.StartDetachedActivity( + SendSpanName, + default(ActivityContext), + [OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, SourceLocal)]); + return new ChatTelemetryOperation(activity, Stopwatch.GetTimestamp()); + } + } + + public void FinishSendAttempt( + ChatTelemetryOperation? operation, + ChatAdmissionTelemetryStatus status, + ChatTelemetryOutcome outcome, + Exception? exception = null) + { + if (operation is null || !operation.TryFinish()) + return; + + var tags = new[] + { + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(outcome)), + OpenClawTelemetryTag.String(AdmissionStatusTag, ToTelemetryValue(status)), + }; + FinishActivity(operation.Activity, outcome, tags, exception); + OpenClawTelemetry.Add(SendAttempts, tags: tags); + OpenClawTelemetry.Record( + SendDuration, + Stopwatch.GetElapsedTime(operation.StartTimestamp).TotalMilliseconds, + tags); + } + + public bool FinishByMessageId( + string messageId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + var completion = PrepareFinishByMessageId(messageId, outcome, reason); + return CompletePreparedTurn(completion); + } + + public PreparedTurnCompletion? PrepareFinishByMessageId( + string messageId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + ArgumentException.ThrowIfNullOrWhiteSpace(messageId); + lock (_gate) + { + if (!_turnsByMessageId.TryGetValue(messageId, out var state)) + return null; + + var completion = PrepareTurnCompletion(state, outcome, reason); + RemoveTurnLocked(state); + return completion; + } + } + + public bool CompletePreparedTurn(PreparedTurnCompletion? completion) + { + if (completion is null || !completion.TryComplete()) + return false; + + FinishTurn(completion); + return true; + } + + public bool FinishByRunId( + string? runId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + var completion = PrepareFinishByRunId(runId, outcome, reason); + return CompletePreparedTurn(completion); + } + + public PreparedTurnCompletion? PrepareFinishByRunId( + string? runId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + if (string.IsNullOrWhiteSpace(runId)) + return null; + + lock (_gate) + { + if (!_turnsByRunId.TryGetValue(runId, out var state)) + return null; + + var completion = PrepareTurnCompletion(state, outcome, reason); + RemoveTurnLocked(state); + return completion; + } + } + + public bool FinishActiveTurn( + string threadId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + ArgumentException.ThrowIfNullOrWhiteSpace(threadId); + TurnState? state; + lock (_gate) + { + var active = _turnsByRunId.Values + .Concat(_turnsByMessageId.Values) + .Distinct() + .Where(candidate => candidate.ThreadId == threadId && candidate.IsDispatched) + .Take(2) + .ToArray(); + if (active.Length != 1) + return false; + + state = active[0]; + RemoveTurnLocked(state); + } + + FinishTurn(state, outcome, reason); + return true; + } + + public void FinishThread( + string threadId, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + ArgumentException.ThrowIfNullOrWhiteSpace(threadId); + FinishStates( + RemoveWhere(state => state.ThreadId == threadId), + outcome, + reason); + } + + public void FinishAll(ChatTelemetryOutcome outcome, ChatTurnTelemetryReason reason) => + FinishStates(RemoveWhere(static _ => true), outcome, reason); + + public ChatTelemetryOperation StartHistoryLoad(ChatHistoryTelemetrySource source) + { + var tags = new[] + { + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, ToTelemetryValue(source)), + }; + return new ChatTelemetryOperation( + OpenClawTelemetry.StartDetachedActivity( + HistoryLoadSpanName, + default(ActivityContext), + tags), + Stopwatch.GetTimestamp(), + tags); + } + + public void FinishHistoryLoad( + ChatTelemetryOperation operation, + ChatTelemetryOutcome outcome, + Exception? exception = null) + { + if (!operation.TryFinish()) + return; + + var tags = AppendOutcome(operation.Tags, outcome); + FinishActivity(operation.Activity, outcome, tags, exception); + OpenClawTelemetry.Add(HistoryLoads, tags: tags); + OpenClawTelemetry.Record( + HistoryLoadDuration, + Stopwatch.GetElapsedTime(operation.StartTimestamp).TotalMilliseconds, + tags); + } + + public ChatTelemetryOperation StartHistoryBackfill(ChatBackfillTelemetryReason reason) + { + var tags = new[] + { + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, SourceRemote), + OpenClawTelemetryTag.String(BackfillReasonTag, ToTelemetryValue(reason)), + }; + return new ChatTelemetryOperation( + OpenClawTelemetry.StartDetachedActivity( + HistoryBackfillSpanName, + default(ActivityContext), + tags), + Stopwatch.GetTimestamp(), + tags); + } + + public void FinishHistoryBackfill( + ChatTelemetryOperation operation, + ChatTelemetryOutcome outcome, + Exception? exception = null) + { + if (!operation.TryFinish()) + return; + + var tags = AppendOutcome(operation.Tags, outcome); + FinishActivity(operation.Activity, outcome, tags, exception); + OpenClawTelemetry.Add(HistoryBackfills, tags: tags); + OpenClawTelemetry.Record( + HistoryBackfillDuration, + Stopwatch.GetElapsedTime(operation.StartTimestamp).TotalMilliseconds, + tags); + } + + public void RecordDroppedTerminalEvent(ChatTerminalEventDropReason reason) + { + OpenClawTelemetry.Add( + DroppedTerminalEvents, + tags: + [ + OpenClawTelemetryTag.String( + DroppedTerminalEventReasonTag, + ToTelemetryValue(reason)), + ]); + } + + private TurnState[] RemoveWhere(Func predicate) + { + lock (_gate) + { + var states = _turnsByMessageId.Values + .Concat(_turnsByRunId.Values) + .Distinct() + .Where(predicate) + .ToArray(); + foreach (var state in states) + RemoveTurnLocked(state); + return states; + } + } + + private static void FinishStates( + IEnumerable states, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + foreach (var state in states) + FinishTurn(state, outcome, reason); + } + + private static void FinishTurn( + TurnState state, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + FinishTurn(PrepareTurnCompletion(state, outcome, reason)); + } + + private static PreparedTurnCompletion PrepareTurnCompletion( + TurnState state, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + var endTimestamp = Stopwatch.GetTimestamp(); + var queueCompletion = state.PendingQueueCompletion; + if (state.QueuePhase is not null) + { + queueCompletion = EndQueueSegmentLocked(state, endTimestamp, outcome); + state.PendingQueueCompletion = queueCompletion; + } + ResponsePhaseCompletion? pendingWaitCompletion = state.PendingWaitCompletion; + if (state.WaitPhase is { } waitPhase) + { + state.ResponseWaitDurationMilliseconds = + Stopwatch.GetElapsedTime(waitPhase.StartTimestamp, endTimestamp).TotalMilliseconds; + pendingWaitCompletion = new ResponsePhaseCompletion( + waitPhase.Activity, + state.Source, + ChatResponseOutputKind.None, + outcome, + Stopwatch.GetElapsedTime(waitPhase.StartTimestamp, endTimestamp)); + } + ResponsePhaseCompletion? receiveCompletion = null; + if (state.ReceivePhase is { } receivePhase) + { + state.ResponseReceiveDurationMilliseconds = + Stopwatch.GetElapsedTime(receivePhase.StartTimestamp, endTimestamp).TotalMilliseconds; + receiveCompletion = new ResponsePhaseCompletion( + receivePhase.Activity, + state.Source, + state.FirstOutputKind, + outcome, + Stopwatch.GetElapsedTime(receivePhase.StartTimestamp, endTimestamp)); + } + return new PreparedTurnCompletion( + state.Activity, + state.StartTimestamp, + endTimestamp, + state.Source, + state.WasQueued, + state.QueuedDurationMilliseconds, + state.ResponseWaitStarted, + state.ResponseWaitDurationMilliseconds, + state.ReceivePhase is not null, + state.ResponseReceiveDurationMilliseconds, + state.FirstOutputKind, + queueCompletion, + pendingWaitCompletion, + receiveCompletion, + outcome, + reason); + } + + private static void FinishTurn(PreparedTurnCompletion completion) + { + CompleteQueuePhase(completion.QueueCompletion); + CompleteResponsePhase(completion.WaitCompletion); + CompleteResponsePhase(completion.ReceiveCompletion); + var tags = new[] + { + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, completion.Source), + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(completion.Outcome)), + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Reason, ToTelemetryValue(completion.Reason)), + }; + FinishActivity(completion.Activity, completion.Outcome, tags); + OpenClawTelemetry.Add(Turns, tags: tags); + OpenClawTelemetry.Record( + TurnDuration, + Stopwatch.GetElapsedTime(completion.StartTimestamp, completion.EndTimestamp).TotalMilliseconds, + [ + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, completion.Source), + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(completion.Outcome)), + ]); + if (completion.WasQueued) + { + OpenClawTelemetry.Record( + QueueWaitDuration, + completion.QueuedDurationMilliseconds, + [OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(completion.Outcome))]); + } + if (completion.ResponseWaitStarted) + { + OpenClawTelemetry.Record( + ResponseWaitDuration, + completion.ResponseWaitDurationMilliseconds, + ResponseMetricTags(completion)); + } + if (completion.ResponseReceiveStarted) + { + OpenClawTelemetry.Record( + ResponseReceiveDuration, + completion.ResponseReceiveDurationMilliseconds, + ResponseMetricTags(completion)); + } + } + + private static OpenClawTelemetryTag[] ResponseMetricTags(PreparedTurnCompletion completion) => + [ + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, completion.Source), + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(completion.Outcome)), + OpenClawTelemetryTag.String(FirstOutputKindTag, ToTelemetryValue(completion.FirstOutputKind)), + ]; + + private static void StartQueueSegmentLocked(TurnState state) + { + if (state.QueuePhase is not null) + return; + + var startTimestamp = Stopwatch.GetTimestamp(); + state.StartQueueSegment( + StartPhase(state, QueueWaitSpanName, startTimestamp), + startTimestamp); + } + + private static QueuePhaseCompletion? EndQueueSegmentLocked( + TurnState state, + long endTimestamp, + ChatTelemetryOutcome outcome) + { + var phase = state.EndQueueSegment(endTimestamp); + return phase is null + ? null + : new QueuePhaseCompletion( + phase.Activity, + state.Source, + outcome, + Stopwatch.GetElapsedTime(phase.StartTimestamp, endTimestamp)); + } + + private static void CompleteQueuePhase(QueuePhaseCompletion? completion) + { + if (completion is null) + return; + + completion.Complete(() => + { + SetActivityDuration(completion.Activity, completion.Duration); + FinishActivity( + completion.Activity, + completion.Outcome, + [ + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, completion.Source), + OpenClawTelemetryTag.String( + OpenClawTelemetryTagKey.Outcome, + ToTelemetryValue(completion.Outcome)), + ]); + }); + } + + private static void CompleteResponsePhase(ResponsePhaseCompletion? completion) + { + if (completion is null) + return; + + completion.Complete(() => + { + SetActivityDuration(completion.Activity, completion.Duration); + FinishActivity( + completion.Activity, + completion.Outcome, + [ + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, completion.Source), + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(completion.Outcome)), + OpenClawTelemetryTag.String(FirstOutputKindTag, ToTelemetryValue(completion.FirstOutputKind)), + ]); + }); + } + + private static void SetActivityDuration(Activity? activity, TimeSpan duration) + { + if (activity is not null) + activity.SetEndTime(activity.StartTimeUtc + duration); + } + + private static void FinishActivity( + Activity? activity, + ChatTelemetryOutcome outcome, + IEnumerable tags, + Exception? exception = null) + { + if (activity is null) + return; + + foreach (var tag in tags) + activity.SetTag(tag.Key, tag.Value); + + switch (outcome) + { + case ChatTelemetryOutcome.Success: + activity.SetStatus(ActivityStatusCode.Ok); + break; + case ChatTelemetryOutcome.Failure: + activity.SetStatus(ActivityStatusCode.Error, exception?.GetType().Name); + if (exception is not null) + { + activity.SetTag( + OpenClawTelemetryTagKey.ErrorType.ToTelemetryName(), + exception.GetType().FullName); + } + break; + case ChatTelemetryOutcome.Canceled: + break; + default: + throw new ArgumentOutOfRangeException(nameof(outcome), outcome, "Unknown chat telemetry outcome."); + } + + OpenClawTelemetry.StopDetachedActivity(activity); + } + + private static OpenClawTelemetryTag[] AppendOutcome( + IReadOnlyList tags, + ChatTelemetryOutcome outcome) => + [.. tags, OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Outcome, ToTelemetryValue(outcome))]; + + private void BindRunLocked(TurnState state, string runId) + { + state.RunIds.Add(runId); + _turnsByRunId[runId] = state; + } + + private static void StartResponseWaitLocked(TurnState state) + { + if (state.ResponseWaitStarted) + return; + + state.ResponseWaitStarted = true; + state.WaitPhase = StartPhase( + state, + ResponseWaitSpanName, + Stopwatch.GetTimestamp()); + } + + private static TimedPhase StartPhase( + TurnState state, + string spanName, + long startTimestamp) + { + var tags = new[] + { + OpenClawTelemetryTag.String(OpenClawTelemetryTagKey.Source, state.Source), + }; + var activity = state.Activity is not null + ? OpenClawTelemetry.StartDetachedActivity(spanName, state.Activity.Context, tags) + : OpenClawTelemetry.StartDetachedActivity(spanName, default(ActivityContext), tags); + return new TimedPhase(activity, startTimestamp); + } + + private TurnState? ResolveTurnForOutputLocked(string threadId, string? runId) + { + if (!string.IsNullOrWhiteSpace(runId)) + return _turnsByRunId.GetValueOrDefault(runId); + + var active = _turnsByRunId.Values + .Concat(_turnsByMessageId.Values) + .Distinct() + .Where(candidate => candidate.ThreadId == threadId && candidate.IsDispatched) + .Take(2) + .ToArray(); + return active.Length == 1 ? active[0] : null; + } + + private void RemoveTurnLocked(TurnState state) + { + if (state.MessageId is not null) + _turnsByMessageId.Remove(state.MessageId); + RemoveRunMappingsLocked(state); + } + + private void RemoveRunMappingsLocked(TurnState state) + { + foreach (var runId in state.RunIds) + { + if (_turnsByRunId.TryGetValue(runId, out var mapped) && ReferenceEquals(mapped, state)) + _turnsByRunId.Remove(runId); + } + state.RunIds.Clear(); + } + + internal static string ToTelemetryValue(ChatTelemetryOutcome outcome) => + outcome switch + { + ChatTelemetryOutcome.Success => "success", + ChatTelemetryOutcome.Failure => "failure", + ChatTelemetryOutcome.Canceled => "canceled", + _ => throw new ArgumentOutOfRangeException(nameof(outcome), outcome, "Unknown chat telemetry outcome."), + }; + + internal static string ToTelemetryValue(ChatTurnTelemetryReason reason) => + reason switch + { + ChatTurnTelemetryReason.AssistantFinal => "assistant_final", + ChatTurnTelemetryReason.LifecycleEnd => "lifecycle_end", + ChatTurnTelemetryReason.LifecycleError => "lifecycle_error", + ChatTurnTelemetryReason.SendRejected => "send_rejected", + ChatTurnTelemetryReason.QueuedCanceled => "queued_canceled", + ChatTurnTelemetryReason.AbortRequested => "abort_requested", + ChatTurnTelemetryReason.Reset => "reset", + ChatTurnTelemetryReason.Superseded => "superseded", + ChatTurnTelemetryReason.Disconnected => "disconnected", + ChatTurnTelemetryReason.Disposed => "disposed", + ChatTurnTelemetryReason.Other => "other", + _ => "other", + }; + + internal static string ToTelemetryValue(ChatAdmissionTelemetryStatus status) => + status switch + { + ChatAdmissionTelemetryStatus.Accepted => "accepted", + ChatAdmissionTelemetryStatus.Deferred => "deferred", + ChatAdmissionTelemetryStatus.Rejected => "rejected", + ChatAdmissionTelemetryStatus.Canceled => "canceled", + ChatAdmissionTelemetryStatus.Exception => "exception", + ChatAdmissionTelemetryStatus.Other => "other", + _ => "other", + }; + + internal static string ToTelemetryValue(ChatHistoryTelemetrySource source) => + source switch + { + ChatHistoryTelemetrySource.Initial => "initial", + ChatHistoryTelemetrySource.Forced => "forced", + _ => throw new ArgumentOutOfRangeException(nameof(source), source, "Unknown chat history telemetry source."), + }; + + internal static string ToTelemetryValue(ChatBackfillTelemetryReason reason) => + reason switch + { + ChatBackfillTelemetryReason.RemoteTurn => "remote_turn", + ChatBackfillTelemetryReason.ResetReconciliation => "reset_reconciliation", + _ => throw new ArgumentOutOfRangeException(nameof(reason), reason, "Unknown chat backfill telemetry reason."), + }; + + internal static string ToTelemetryValue(ChatTerminalEventDropReason reason) => + reason switch + { + ChatTerminalEventDropReason.MissingRunId => MissingRunId, + ChatTerminalEventDropReason.MismatchedRunId => MismatchedRunId, + _ => throw new ArgumentOutOfRangeException(nameof(reason), reason, "Unknown terminal event drop reason."), + }; + + internal static string ToTelemetryValue(ChatResponseOutputKind kind) => + kind switch + { + ChatResponseOutputKind.None => "none", + ChatResponseOutputKind.Assistant => "assistant", + ChatResponseOutputKind.Reasoning => "reasoning", + ChatResponseOutputKind.Tool => "tool", + ChatResponseOutputKind.Other => "other", + _ => "other", + }; + + private sealed class TurnState( + string? messageId, + string threadId, + string source, + Activity? activity, + long startTimestamp) + { + private long? _queueSegmentStart; + + public string? MessageId { get; } = messageId; + public string ThreadId { get; } = threadId; + public string Source { get; } = source; + public Activity? Activity { get; } = activity; + public long StartTimestamp { get; } = startTimestamp; + public HashSet RunIds { get; } = new(StringComparer.Ordinal); + public bool IsDispatched { get; set; } + public bool WasQueued { get; private set; } + public double QueuedDurationMilliseconds { get; private set; } + public bool ResponseWaitStarted { get; set; } + public double ResponseWaitDurationMilliseconds { get; set; } + public double ResponseReceiveDurationMilliseconds { get; set; } + public ChatResponseOutputKind FirstOutputKind { get; set; } + public TimedPhase? QueuePhase { get; private set; } + public TimedPhase? WaitPhase { get; set; } + public TimedPhase? ReceivePhase { get; set; } + public QueuePhaseCompletion? PendingQueueCompletion { get; set; } + public ResponsePhaseCompletion? PendingWaitCompletion { get; set; } + + public void StartQueueSegment(TimedPhase phase, long startTimestamp) + { + if (_queueSegmentStart.HasValue) + return; + WasQueued = true; + _queueSegmentStart = startTimestamp; + QueuePhase = phase; + } + + public TimedPhase? EndQueueSegment(long endTimestamp) + { + if (_queueSegmentStart is not { } started) + return null; + QueuedDurationMilliseconds += + Stopwatch.GetElapsedTime(started, endTimestamp).TotalMilliseconds; + _queueSegmentStart = null; + var phase = QueuePhase; + QueuePhase = null; + return phase; + } + } + + private sealed record TimedPhase(Activity? Activity, long StartTimestamp); + + internal sealed class QueuePhaseCompletion( + Activity? activity, + string source, + ChatTelemetryOutcome outcome, + TimeSpan duration) : PhaseCompletion + { + public Activity? Activity { get; } = activity; + public string Source { get; } = source; + public ChatTelemetryOutcome Outcome { get; } = outcome; + public TimeSpan Duration { get; } = duration; + } + + internal sealed class ResponsePhaseCompletion( + Activity? activity, + string source, + ChatResponseOutputKind firstOutputKind, + ChatTelemetryOutcome outcome, + TimeSpan duration) : PhaseCompletion + { + public Activity? Activity { get; } = activity; + public string Source { get; } = source; + public ChatResponseOutputKind FirstOutputKind { get; } = firstOutputKind; + public ChatTelemetryOutcome Outcome { get; } = outcome; + public TimeSpan Duration { get; } = duration; + } + + internal abstract class PhaseCompletion + { + private readonly object _completionGate = new(); + private int _completionState; + private int _completingThreadId; + + public void Complete(Action complete) + { + lock (_completionGate) + { + if (_completionState == 2) + return; + if (_completionState == 1) + { + if (_completingThreadId == Environment.CurrentManagedThreadId) + return; + while (_completionState != 2) + Monitor.Wait(_completionGate); + return; + } + + _completionState = 1; + _completingThreadId = Environment.CurrentManagedThreadId; + } + + try + { + complete(); + } + finally + { + lock (_completionGate) + { + _completionState = 2; + _completingThreadId = 0; + Monitor.PulseAll(_completionGate); + } + } + } + } + + internal sealed class PreparedTurnCompletion( + Activity? activity, + long startTimestamp, + long endTimestamp, + string source, + bool wasQueued, + double queuedDurationMilliseconds, + bool responseWaitStarted, + double responseWaitDurationMilliseconds, + bool responseReceiveStarted, + double responseReceiveDurationMilliseconds, + ChatResponseOutputKind firstOutputKind, + QueuePhaseCompletion? queueCompletion, + ResponsePhaseCompletion? waitCompletion, + ResponsePhaseCompletion? receiveCompletion, + ChatTelemetryOutcome outcome, + ChatTurnTelemetryReason reason) + { + private int _completed; + + public Activity? Activity { get; } = activity; + public long StartTimestamp { get; } = startTimestamp; + public long EndTimestamp { get; } = endTimestamp; + public string Source { get; } = source; + public bool WasQueued { get; } = wasQueued; + public double QueuedDurationMilliseconds { get; } = queuedDurationMilliseconds; + public bool ResponseWaitStarted { get; } = responseWaitStarted; + public double ResponseWaitDurationMilliseconds { get; } = responseWaitDurationMilliseconds; + public bool ResponseReceiveStarted { get; } = responseReceiveStarted; + public double ResponseReceiveDurationMilliseconds { get; } = responseReceiveDurationMilliseconds; + public ChatResponseOutputKind FirstOutputKind { get; } = firstOutputKind; + internal QueuePhaseCompletion? QueueCompletion { get; } = queueCompletion; + internal ResponsePhaseCompletion? WaitCompletion { get; } = waitCompletion; + internal ResponsePhaseCompletion? ReceiveCompletion { get; } = receiveCompletion; + public ChatTelemetryOutcome Outcome { get; } = outcome; + public ChatTurnTelemetryReason Reason { get; } = reason; + + public bool TryComplete() => Interlocked.Exchange(ref _completed, 1) == 0; + } +} + +internal sealed class ChatTelemetryOperation( + Activity? activity, + long startTimestamp, + IReadOnlyList? tags = null) +{ + private int _finished; + + public Activity? Activity { get; } = activity; + public long StartTimestamp { get; } = startTimestamp; + public IReadOnlyList Tags { get; } = tags ?? []; + + public bool TryFinish() => Interlocked.Exchange(ref _finished, 1) == 0; +} diff --git a/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs b/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs index 4ddcedf45..b67e43ce0 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs @@ -88,6 +88,7 @@ public sealed class OpenClawChatDataProvider : IChatDataProvider public static readonly ConcurrentDictionary ImagePreviewCache = new(); private readonly IChatGatewayBridge _bridge; + private readonly ChatTelemetryTracker _telemetry = new(); private readonly Action? _post; private readonly object _gate = new(); private readonly object _toolMetaSaveGate = new(); @@ -175,6 +176,7 @@ private sealed record QueuedSendDispatch( long ResetVersion, long StartedLifecycleSequence, long StartedRunStartSequence, + ChatTelemetryTracker.QueuePhaseCompletion? QueueCompletion, bool StartedDirectly); private enum AssistantQueueFrameDisposition { @@ -414,7 +416,9 @@ public async Task SendMessageAsync(string threadId, string message, Cancellation nonce, attachments?.ToArray()); - if (CanSendDirectlyLocked(threadId)) + var sendDirectly = CanSendDirectlyLocked(threadId); + _telemetry.StartLocalTurn(request.Id, threadId, queued: !sendDirectly); + if (sendDirectly) { dispatch = StartDirectSendLocked(request); } @@ -446,15 +450,23 @@ public Task CancelQueuedMessageAsync(string threadId, string queuedMessage throw new ArgumentException("Queued message id is required.", nameof(queuedMessageId)); ChatDataSnapshot? snapshot = null; + ChatTelemetryTracker.PreparedTurnCompletion? telemetryCompletion = null; var canceled = false; lock (_gate) { ObjectDisposedException.ThrowIf(_disposed, this); canceled = CancelQueuedMessageLocked(threadId, queuedMessageId); if (canceled) + { + telemetryCompletion = _telemetry.PrepareFinishByMessageId( + queuedMessageId, + ChatTelemetryOutcome.Canceled, + ChatTurnTelemetryReason.QueuedCanceled); snapshot = BuildSnapshotLocked(); + } } + _telemetry.CompletePreparedTurn(telemetryCompletion); if (snapshot is not null) Publish(snapshot); @@ -466,9 +478,11 @@ private async Task DispatchQueuedSendAsync( bool rethrow, CancellationToken cancellationToken = default) { + _telemetry.CompleteQueueDispatch(dispatch.QueueCompletion); var request = dispatch.Request; var threadId = request.ThreadId; var hasAttachments = request.Attachments is { Count: > 0 }; + ChatTelemetryOperation? sendOperation = null; try { @@ -480,14 +494,36 @@ private async Task DispatchQueuedSendAsync( if (GetResetVersionLocked(threadId) == dispatch.ResetVersion) TrackQueuedMessageRunLocked(threadId, request.SendRunId, request.Id); } + sendOperation = _telemetry.StartSendAttempt(request.Id); var sendResult = await _bridge.SendChatMessageForRunAsync( request.Text, threadId, dispatch.SessionId, request.Attachments, idempotencyKey: request.SendRunId); + var admissionStatus = MapAdmissionTelemetryStatus(sendResult); + var admissionOutcome = admissionStatus == ChatAdmissionTelemetryStatus.Canceled + ? ChatTelemetryOutcome.Canceled + : sendResult.IsTerminalFailure + ? ChatTelemetryOutcome.Failure + : ChatTelemetryOutcome.Success; + _telemetry.FinishSendAttempt( + sendOperation, + admissionStatus, + admissionOutcome); + if (admissionStatus == ChatAdmissionTelemetryStatus.Accepted) + _telemetry.ObserveAdmissionAccepted(request.Id); if (sendResult.IsTerminalFailure) { + ChatTelemetryTracker.PreparedTurnCompletion? rejectedCompletion; + lock (_gate) + { + rejectedCompletion = _telemetry.PrepareFinishByMessageId( + request.Id, + admissionOutcome, + ChatTurnTelemetryReason.SendRejected); + } + _telemetry.CompletePreparedTurn(rejectedCompletion); var failure = !string.IsNullOrWhiteSpace(sendResult.Error) ? sendResult.Error! : string.Format( @@ -499,6 +535,7 @@ private async Task DispatchQueuedSendAsync( bool sendStillCurrent; string? staleRunIdToAbort = null; + ChatTelemetryTracker.PreparedTurnCompletion? staleCompletion = null; ChatDataSnapshot? acceptedSnapshot = null; ChatDataSnapshot? requeuedSnapshot = null; var retryDeferredSend = false; @@ -512,6 +549,10 @@ private async Task DispatchQueuedSendAsync( if (!sendStillCurrent) { staleRunIdToAbort = acceptedRunId ?? request.SendRunId; + staleCompletion = _telemetry.PrepareFinishByMessageId( + request.Id, + ChatTelemetryOutcome.Canceled, + ChatTurnTelemetryReason.Superseded); AddResetIgnoredRunIdLocked(threadId, staleRunIdToAbort); } else if (IsDeferredAdmissionStatus(sendResult.Status)) @@ -523,6 +564,7 @@ private async Task DispatchQueuedSendAsync( && activeStartSequence > dispatch.StartedRunStartSequence; if (runAlreadyStarted) { + _telemetry.BindAcceptedRun(request.Id, acceptedRunId); TrackQueuedMessageRunLocked(threadId, acceptedRunId!, request.Id); AddResetAcceptedRunIdLocked(threadId, acceptedRunId!); if (PromoteQueuedMessageLocked(threadId, request.Id)) @@ -536,6 +578,7 @@ private async Task DispatchQueuedSendAsync( } else if (RequeueDeferredAdmissionLocked(threadId, request.Id, out deferredRetryDelay)) { + _telemetry.RequeueLocalTurn(request.Id); if (!string.IsNullOrEmpty(acceptedRunId)) { TrackQueuedMessageRunLocked(threadId, acceptedRunId, request.Id); @@ -552,6 +595,7 @@ private async Task DispatchQueuedSendAsync( } else if (!string.IsNullOrEmpty(acceptedRunId)) { + _telemetry.BindAcceptedRun(request.Id, acceptedRunId); TrackQueuedMessageRunLocked(threadId, acceptedRunId, request.Id); AddResetAcceptedRunIdLocked(threadId, acceptedRunId); var runAlreadyStarted = _activeRunIds.TryGetValue(threadId, out var activeRunId) @@ -600,6 +644,7 @@ private async Task DispatchQueuedSendAsync( if (staleRunIdToAbort is not null) { + _telemetry.CompletePreparedTurn(staleCompletion); try { Logger.Info($"[Reset] Aborting late pre-reset send runId='{staleRunIdToAbort}' threadId='{threadId}'"); @@ -616,13 +661,27 @@ private async Task DispatchQueuedSendAsync( } catch (Exception ex) { + _telemetry.FinishSendAttempt( + sendOperation, + ChatAdmissionTelemetryStatus.Exception, + ex is OperationCanceledException + ? ChatTelemetryOutcome.Canceled + : ChatTelemetryOutcome.Failure, + ex); bool sendStillCurrent; + ChatTelemetryTracker.PreparedTurnCompletion? rejectedCompletion = null; ChatDataSnapshot? failureSnapshot = null; lock (_gate) { sendStillCurrent = GetResetVersionLocked(threadId) == dispatch.ResetVersion; if (sendStillCurrent) { + rejectedCompletion = _telemetry.PrepareFinishByMessageId( + request.Id, + ex is OperationCanceledException + ? ChatTelemetryOutcome.Canceled + : ChatTelemetryOutcome.Failure, + ChatTurnTelemetryReason.SendRejected); RemovePendingLocalEchoLocked(threadId, request.Id); MarkQueuedMessageFailedLocked(threadId, request.Id, ex.Message); RemoveQueuedSendRequestLocked(threadId, request.Id); @@ -643,6 +702,7 @@ private async Task DispatchQueuedSendAsync( if (!sendStillCurrent) return; + _telemetry.CompletePreparedTurn(rejectedCompletion); Logger.Warn($"[Queue] chat.send failed threadId='{threadId}' queuedMessageId='{request.Id}' sendRunId='{request.SendRunId}': {ex.Message}"); // Surface as an error in the timeline + notification, while the // failed queue card keeps the attempted text visible for retry/edit. @@ -676,6 +736,11 @@ public async Task StopResponseAsync(string threadId, CancellationToken cancellat _pendingAbortCounts.TryGetValue(threadId, out var count); _pendingAbortCounts[threadId] = count + 1; } + + _telemetry.FinishActiveTurn( + threadId, + ChatTelemetryOutcome.Canceled, + ChatTurnTelemetryReason.AbortRequested); } Logger.Info($"[ABORT] StopResponseAsync threadId='{threadId}' runId='{runId ?? "(null)"}' hadActiveTurn={hadActiveTurn} deferred={string.IsNullOrEmpty(runId)}"); @@ -761,6 +826,10 @@ public async Task LoadHistoryAsync(string threadId, bool force = false, Cancella requestResetVersion = GetResetVersionLocked(threadId); } + var historyOperation = _telemetry.StartHistoryLoad( + force ? ChatHistoryTelemetrySource.Forced : ChatHistoryTelemetrySource.Initial); + var historyOutcome = ChatTelemetryOutcome.Success; + Exception? historyException = null; try { var history = await _bridge.RequestChatHistoryAsync(threadId); @@ -1190,6 +1259,10 @@ ChatTimelineState ApplyAndCaptureMeta(ChatTimelineState s, ChatEvent e, ChatEntr } catch (Exception ex) { + historyOutcome = ex is OperationCanceledException + ? ChatTelemetryOutcome.Canceled + : ChatTelemetryOutcome.Failure; + historyException = ex; RaiseNotification(new ChatProviderNotification( ChatProviderNotificationKind.Error, threadId, LocalizationHelper.GetString("Chat_Notification_LoadHistoryFailed"), ex.Message)); @@ -1218,6 +1291,7 @@ ChatTimelineState ApplyAndCaptureMeta(ChatTimelineState s, ChatEvent e, ChatEntr finally { lock (_gate) { _historyInFlight.Remove(threadId); } + _telemetry.FinishHistoryLoad(historyOperation, historyOutcome, historyException); } } @@ -1679,6 +1753,7 @@ public ValueTask DisposeAsync() List pendingLocalApprovals; lock (_gate) { + _telemetry.FinishAll(ChatTelemetryOutcome.Canceled, ChatTurnTelemetryReason.Disposed); timerToDispose = _toolMetaSaveTimer; _toolMetaSaveTimer = null; _toolMetaSaveVersion++; @@ -1778,6 +1853,7 @@ private void OnStatusChanged(object? sender, ConnectionStatus status) // still pending (the request ID is stale after reconnect). if (justReconnected) { + _telemetry.FinishAll(ChatTelemetryOutcome.Canceled, ChatTurnTelemetryReason.Disconnected); var reload = new HashSet(_historyLoaded); foreach (var key in _timelines.Keys) reload.Add(key); @@ -1809,6 +1885,7 @@ private void OnStatusChanged(object? sender, ConnectionStatus status) if (justDisconnected) { + _telemetry.FinishAll(ChatTelemetryOutcome.Canceled, ChatTurnTelemetryReason.Disconnected); var list = new List(); foreach (var (key, tl) in _timelines) { @@ -2204,6 +2281,7 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) if (string.IsNullOrEmpty(message.Text)) return; var trThread = message.SessionKey; ChatEntryMetadata? trMeta; + string? trRunId; lock (_gate) { trMeta = BuildLiveMetaLocked( @@ -2211,10 +2289,15 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) message.Ts, message.OpenClawId, message.OpenClawSeq); + _activeRunIds.TryGetValue(trThread, out trRunId); } var capped = TruncateForChatEntry(message.Text); var kind = ClassifyFlattenedToolOutput(capped); var label = ExtractFlattenedToolSummary(capped); + _telemetry.ObserveInboundOutput( + trThread, + trRunId, + ChatResponseOutputKind.Tool); ApplyEventAndPublish(trThread, new ChatToolStartEvent(label, kind), trMeta); ApplyEventAndPublish(trThread, new ChatToolOutputEvent(capped), trMeta); return; @@ -2246,6 +2329,7 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) PromoteOldestQueuedMessageBeforeAssistantIfNeeded(threadId); ChatEntryMetadata? meta; + string? telemetryRunId; var hasUsage = message.InputTokens is not null || message.OutputTokens is not null || message.ResponseTokens is not null || message.ContextPercent is not null; lock (_gate) @@ -2255,6 +2339,7 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) message.Ts, message.OpenClawId, message.OpenClawSeq); + _activeRunIds.TryGetValue(threadId, out telemetryRunId); // If the gateway included a usage block on this chat event, // attach it so the assistant footer pills (↑/↓/R/ctx%) can // render. Mostly arrives on state="final" frames. @@ -2278,6 +2363,10 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) return; } + _telemetry.ObserveInboundOutput( + threadId, + telemetryRunId, + ChatResponseOutputKind.Assistant); // Both `state: "delta"` and `state: "final"` carry the cumulative // assistant text (the gateway's EmbeddedBlockChunker emits completed // blocks, not token deltas — see spec §"Block Streaming"). Map both @@ -2301,10 +2390,15 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) if (message.IsFinal) { + ChatTelemetryTracker.PreparedTurnCompletion? turnCompletion = null; lock (_gate) { if (_activeRunIds.Remove(threadId, out var completedRunId)) { + turnCompletion = _telemetry.PrepareFinishByRunId( + completedRunId, + ChatTelemetryOutcome.Success, + ChatTurnTelemetryReason.AssistantFinal); RememberTerminalRunIdLocked(threadId, completedRunId); _abortedRunIds.Remove(completedRunId); } @@ -2313,6 +2407,7 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) if (!HasSendingQueuedMessagesLocked(threadId)) _locallyInitiatedThreads.Remove(threadId); } + _telemetry.CompletePreparedTurn(turnCompletion); SnapshotLatestAssistantUsage(threadId); ApplyEventAndPublish(threadId, new ChatTurnEndEvent()); RaiseNotification(new ChatProviderNotification( @@ -2360,13 +2455,14 @@ private void OnAgentEventReceived(object? sender, AgentEventInfo evt) var reloadHistoryAfterResetDrop = false; var shouldProcessEvent = false; + ChatTerminalEventDropReason? droppedTerminalReason = null; lock (_gate) { if (ShouldDropAgentEventAfterResetLocked(evt, threadId, out reloadHistoryAfterResetDrop)) { Logger.Debug($"[Reset] Dropping stale agent event after reset for threadId='{threadId}' stream='{evt.Stream}' runId='{evt.RunId}'"); } - else if (ShouldDropTerminalAgentEventLocked(evt, threadId)) + else if (ShouldDropTerminalAgentEventLocked(evt, threadId, out droppedTerminalReason)) { Logger.Debug($"[Queue] Dropping stale terminal agent event for threadId='{threadId}' stream='{evt.Stream}' runId='{evt.RunId}'"); } @@ -2377,18 +2473,22 @@ private void OnAgentEventReceived(object? sender, AgentEventInfo evt) } if (!shouldProcessEvent) { + if (droppedTerminalReason.HasValue) + RecordDroppedTerminalEvent(droppedTerminalReason.Value); if (reloadHistoryAfterResetDrop) _ = LoadHistoryAsync(threadId, force: true); return; } // Always update run tracking first (state maintenance must not be skipped). - UpdateActiveRunId(evt, threadId); + var deferredAbort = UpdateActiveRunId(evt, threadId); + if (deferredAbort.DroppedTerminalReason.HasValue) + RecordDroppedTerminalEvent(deferredAbort.DroppedTerminalReason.Value); ClearQueuedMessageOnLocalTurnStart(evt, threadId); // Fire deferred chat.abort and persist if pending aborts were queued. - var deferredRunId = _deferredAbortRunId; - var shouldPersist = _deferredAbortCount > 0; + var deferredRunId = deferredAbort.RunId; + var shouldPersist = deferredAbort.Count > 0; if (deferredRunId is not null || shouldPersist) { _ = Task.Run(async () => @@ -2505,6 +2605,15 @@ private void OnAgentEventReceived(object? sender, AgentEventInfo evt) return; } + var outputKind = ClassifyInboundOutput(evt, mapped); + if (outputKind.HasValue) + { + _telemetry.ObserveInboundOutput( + threadId, + evt.RunId, + outputKind.Value); + } + // Cache tool metadata from live SSE events so it survives app restarts. if (mapped is ChatToolStartEvent toolStart && !string.IsNullOrEmpty(toolStart.ToolName)) { @@ -2522,6 +2631,29 @@ private void OnAgentEventReceived(object? sender, AgentEventInfo evt) ScheduleQueuedSendDrain(threadId); } + private static ChatResponseOutputKind? ClassifyInboundOutput( + AgentEventInfo evt, + ChatEvent mapped) + { + if (string.Equals(evt.Stream, "lifecycle", StringComparison.OrdinalIgnoreCase) || + string.Equals(evt.Stream, "job", StringComparison.OrdinalIgnoreCase)) + { + return null; + } + + return mapped switch + { + ChatMessageEvent or ChatMessageDeltaEvent => ChatResponseOutputKind.Assistant, + ChatThinkingEvent or ChatReasoningEvent or ChatReasoningDeltaEvent or + ChatIntentEvent => ChatResponseOutputKind.Reasoning, + ChatToolStartEvent or ChatToolOutputEvent or ChatToolErrorEvent or + ChatPermissionRequestEvent => ChatResponseOutputKind.Tool, + ChatStatusEvent or ChatErrorEvent or ChatReasoningEndEvent or + ChatTurnEndEvent or ChatUserMessageEvent => null, + _ => null, + }; + } + private void RaiseKeylessEventDiagnosticOnce() { if (System.Threading.Interlocked.Exchange(ref _keylessEventDiagnosticRaised, 1) != 0) @@ -2550,13 +2682,14 @@ private void RaiseKeylessEventDiagnosticOnce() } } - private string? _deferredAbortRunId; // set inside lock when pending abort fires; read outside lock to send RPC - private int _deferredAbortCount; // how many user messages to force-persist as aborted - - private void UpdateActiveRunId(AgentEventInfo evt, string threadId) + private (string? RunId, int Count, ChatTerminalEventDropReason? DroppedTerminalReason) UpdateActiveRunId( + AgentEventInfo evt, + string threadId) { - _deferredAbortRunId = null; - _deferredAbortCount = 0; + string? deferredAbortRunId = null; + var deferredAbortCount = 0; + ChatTerminalEventDropReason? droppedTerminalReason = null; + ChatTelemetryTracker.PreparedTurnCompletion? turnCompletion = null; if (string.Equals(evt.Stream, "lifecycle", StringComparison.OrdinalIgnoreCase) && evt.Data.ValueKind == System.Text.Json.JsonValueKind.Object && @@ -2565,33 +2698,56 @@ private void UpdateActiveRunId(AgentEventInfo evt, string threadId) var phase = phaseProp.GetString()?.ToLowerInvariant(); lock (_gate) { - if (phase == "start" && !string.IsNullOrEmpty(evt.RunId)) + if (phase == "start") { - _activeRunIds[threadId] = evt.RunId; - _activeRunStartSequences[threadId] = ++_lifecycleStartSequence; - - // Detect remote turn: if the turn was NOT locally initiated, - // a remote client (e.g. gateway web UI) sent the message. - // Fetch the last user message from history so it appears in - // the timeline before the assistant response. - if (!_locallyInitiatedThreads.Contains(threadId)) + _telemetry.ObserveLifecycleStart( + threadId, + evt.RunId, + allowRemoteTurn: !_locallyInitiatedThreads.Contains(threadId) && + !_abortedThreads.Contains(threadId) && + !_pendingAbortCounts.ContainsKey(threadId)); + if (!string.IsNullOrEmpty(evt.RunId)) { - _ = FetchRemoteUserMessageAsync(threadId); - } + _activeRunIds[threadId] = evt.RunId; + _activeRunStartSequences[threadId] = ++_lifecycleStartSequence; + + // Detect remote turn: if the turn was NOT locally initiated, + // a remote client (e.g. gateway web UI) sent the message. + // Fetch the last user message from history so it appears in + // the timeline before the assistant response. + if (!_locallyInitiatedThreads.Contains(threadId)) + { + _ = FetchRemoteUserMessageAsync(threadId); + } - // Deferred abort: if user clicked stop before lifecycle.start, - // fire chat.abort now that we have the runId. - if (_pendingAbortCounts.TryGetValue(threadId, out var pendingCount) && pendingCount > 0) - { - _pendingAbortCounts.Remove(threadId); - _abortedRunIds.Add(evt.RunId); - _deferredAbortRunId = evt.RunId; - _deferredAbortCount = pendingCount; - Logger.Info($"[ABORT] Deferred abort fired — lifecycle.start arrived with runId='{evt.RunId}' for threadId='{threadId}' (pendingCount={pendingCount})"); + // Deferred abort: if user clicked stop before lifecycle.start, + // fire chat.abort now that we have the runId. + if (_pendingAbortCounts.TryGetValue(threadId, out var pendingCount) && pendingCount > 0) + { + _pendingAbortCounts.Remove(threadId); + _abortedRunIds.Add(evt.RunId); + deferredAbortRunId = evt.RunId; + deferredAbortCount = pendingCount; + Logger.Info($"[ABORT] Deferred abort fired — lifecycle.start arrived with runId='{evt.RunId}' for threadId='{threadId}' (pendingCount={pendingCount})"); + } } } else if (phase == "end" || phase == "error") { + var wasAborted = !string.IsNullOrWhiteSpace(evt.RunId) && + _abortedRunIds.Contains(evt.RunId); + turnCompletion = _telemetry.PrepareFinishByRunId( + evt.RunId, + phase == "error" ? ChatTelemetryOutcome.Failure : ChatTelemetryOutcome.Success, + phase == "error" + ? ChatTurnTelemetryReason.LifecycleError + : ChatTurnTelemetryReason.LifecycleEnd); + if (turnCompletion is null && !wasAborted) + { + droppedTerminalReason = string.IsNullOrWhiteSpace(evt.RunId) + ? ChatTerminalEventDropReason.MissingRunId + : ChatTerminalEventDropReason.MismatchedRunId; + } // Clean up: remove aborted runId tracking on terminal events. if (!string.IsNullOrEmpty(evt.RunId)) _abortedRunIds.Remove(evt.RunId); @@ -2616,8 +2772,8 @@ private void UpdateActiveRunId(AgentEventInfo evt, string threadId) if (_pendingAbortCounts.TryGetValue(threadId, out var lateCount) && lateCount > 0) { _pendingAbortCounts.Remove(threadId); - _deferredAbortRunId = evt.RunId; // may be null, that's ok — persist doesn't need it - _deferredAbortCount = lateCount; + deferredAbortRunId = evt.RunId; // may be null, that's ok — persist doesn't need it + deferredAbortCount = lateCount; Logger.Info($"[ABORT] Late deferred abort — lifecycle.end arrived with pending aborts for threadId='{threadId}' (pendingCount={lateCount})"); } } @@ -2631,21 +2787,50 @@ private void UpdateActiveRunId(AgentEventInfo evt, string threadId) var state = stateProp.GetString()?.ToLowerInvariant(); lock (_gate) { - if ((state == "done" || state == "error") && !string.IsNullOrEmpty(evt.RunId)) + if (state == "done" || state == "error") { - _abortedRunIds.Remove(evt.RunId); + var wasAborted = !string.IsNullOrWhiteSpace(evt.RunId) && + _abortedRunIds.Contains(evt.RunId); + turnCompletion = _telemetry.PrepareFinishByRunId( + evt.RunId, + state == "error" ? ChatTelemetryOutcome.Failure : ChatTelemetryOutcome.Success, + state == "error" + ? ChatTurnTelemetryReason.LifecycleError + : ChatTurnTelemetryReason.LifecycleEnd); + if (turnCompletion is null && !wasAborted) + { + droppedTerminalReason = string.IsNullOrWhiteSpace(evt.RunId) + ? ChatTerminalEventDropReason.MissingRunId + : ChatTerminalEventDropReason.MismatchedRunId; + } + if (!string.IsNullOrWhiteSpace(evt.RunId)) + { + _abortedRunIds.Remove(evt.RunId); + RemoveQueuedRunMappingByRunIdLocked(threadId, evt.RunId); + } _activeRunIds.Remove(threadId); _activeRunStartSequences.Remove(threadId); - RemoveQueuedRunMappingByRunIdLocked(threadId, evt.RunId); } } } + + _telemetry.CompletePreparedTurn(turnCompletion); + return (deferredAbortRunId, deferredAbortCount, droppedTerminalReason); } - private bool ShouldDropTerminalAgentEventLocked(AgentEventInfo evt, string threadId) + private bool ShouldDropTerminalAgentEventLocked( + AgentEventInfo evt, + string threadId, + out ChatTerminalEventDropReason? droppedTerminalReason) { - if (!TryGetTerminalAgentRunId(evt, out var runId) || string.IsNullOrWhiteSpace(runId)) + droppedTerminalReason = null; + if (!TryGetTerminalAgentRunId(evt, out var runId)) return false; + if (string.IsNullOrWhiteSpace(runId)) + { + droppedTerminalReason = ChatTerminalEventDropReason.MissingRunId; + return true; + } if (_terminalRunIdsByThread.TryGetValue(threadId, out var terminalRunIds) && terminalRunIds.Contains(runId, StringComparer.Ordinal)) @@ -2656,6 +2841,7 @@ private bool ShouldDropTerminalAgentEventLocked(AgentEventInfo evt, string threa if (_activeRunIds.TryGetValue(threadId, out var activeRunId) && !string.Equals(activeRunId, runId, StringComparison.Ordinal)) { + droppedTerminalReason = ChatTerminalEventDropReason.MismatchedRunId; return true; } @@ -2666,6 +2852,7 @@ private bool ShouldDropTerminalAgentEventLocked(AgentEventInfo evt, string threa _timelines.TryGetValue(threadId, out var timeline) && timeline.TurnActive) { + droppedTerminalReason = ChatTerminalEventDropReason.MismatchedRunId; return true; } @@ -2673,6 +2860,14 @@ private bool ShouldDropTerminalAgentEventLocked(AgentEventInfo evt, string threa return false; } + private void RecordDroppedTerminalEvent(ChatTerminalEventDropReason reason) + { + _telemetry.RecordDroppedTerminalEvent(reason); + Logger.Warn( + $"[ChatTelemetry] Dropped terminal chat event because safe run correlation was unavailable " + + $"(reason='{ChatTelemetryTracker.ToTelemetryValue(reason)}')."); + } + private void RememberTerminalRunIdLocked(string threadId, string runId) { if (!_terminalRunIdsByThread.TryGetValue(threadId, out var terminalRunIds)) @@ -2894,12 +3089,14 @@ private QueuedSendDispatch StartDirectSendLocked(QueuedSendRequest request) EnqueueLocalEchoLocked(threadId, request.Text, request.Id); _locallyInitiatedThreads.Add(threadId); _assistantFallbackPromotedThreads.Add(threadId); + var queueCompletion = _telemetry.PrepareDispatchLocalTurn(request.Id, request.SendRunId); return new QueuedSendDispatch( request, sessionId, resetVersion, startedLifecycleSequence, startedRunStartSequence, + queueCompletion, StartedDirectly: true); } @@ -2952,12 +3149,14 @@ private QueuedSendDispatch StartDirectSendLocked(QueuedSendRequest request) EnqueueLocalEchoLocked(threadId, request.Text, request.Id); _locallyInitiatedThreads.Add(threadId); + var queueCompletion = _telemetry.PrepareDispatchLocalTurn(request.Id, request.SendRunId); return new QueuedSendDispatch( request, sessionId, resetVersion, startedLifecycleSequence, startedRunStartSequence, + queueCompletion, StartedDirectly: false); } @@ -3041,6 +3240,29 @@ private void ScheduleQueuedSendDrain(string threadId, TimeSpan delay) private static bool IsDeferredAdmissionStatus(string? status) => string.Equals(status, "in_flight", StringComparison.OrdinalIgnoreCase); + private static ChatAdmissionTelemetryStatus MapAdmissionTelemetryStatus(ChatSendResult result) + { + if (IsDeferredAdmissionStatus(result.Status)) + return ChatAdmissionTelemetryStatus.Deferred; + if (result.IsTerminalFailure) + { + return IsCanceledAdmissionStatus(result.Status) + ? ChatAdmissionTelemetryStatus.Canceled + : ChatAdmissionTelemetryStatus.Rejected; + } + if (string.IsNullOrWhiteSpace(result.Status) || + string.Equals(result.Status, "started", StringComparison.OrdinalIgnoreCase)) + { + return ChatAdmissionTelemetryStatus.Accepted; + } + return ChatAdmissionTelemetryStatus.Other; + } + + private static bool IsCanceledAdmissionStatus(string? status) => + string.Equals(status, "aborted", StringComparison.OrdinalIgnoreCase) || + string.Equals(status, "cancelled", StringComparison.OrdinalIgnoreCase) || + string.Equals(status, "canceled", StringComparison.OrdinalIgnoreCase); + private static TimeSpan DeferredAdmissionRetryDelay(int retryCount) { var exponent = Math.Min(Math.Max(retryCount - 1, 0), 5); @@ -3379,6 +3601,12 @@ private void RemovePendingLocalEchoLocked(string threadId, string messageId) /// private async Task FetchRemoteUserMessageAsync(string threadId, bool openResetGateOnSuccess = false) { + var telemetryReason = openResetGateOnSuccess + ? ChatBackfillTelemetryReason.ResetReconciliation + : ChatBackfillTelemetryReason.RemoteTurn; + var historyOperation = _telemetry.StartHistoryBackfill(telemetryReason); + var historyOutcome = ChatTelemetryOutcome.Success; + Exception? historyException = null; long requestResetVersion; long resetCutoffUtcMs; lock (_gate) @@ -3456,6 +3684,10 @@ private async Task FetchRemoteUserMessageAsync(string threadId, bool openResetGa } catch (Exception ex) { + historyOutcome = ex is OperationCanceledException + ? ChatTelemetryOutcome.Canceled + : ChatTelemetryOutcome.Failure; + historyException = ex; Logger.Warn($"[REMOTE] Failed to fetch remote user message for threadId='{threadId}': {ex.Message}"); } finally @@ -3464,6 +3696,7 @@ private async Task FetchRemoteUserMessageAsync(string threadId, bool openResetGa { lock (_gate) { _resetRemoteBackfillInFlight.Remove(threadId); } } + _telemetry.FinishHistoryBackfill(historyOperation, historyOutcome, historyException); } } @@ -4683,6 +4916,7 @@ private long GetResetCutoffUtcMsLocked(string threadId) => private ResetClearPersistence ClearThreadHistoryAfterResetLocked(string threadId) { + _telemetry.FinishThread(threadId, ChatTelemetryOutcome.Canceled, ChatTurnTelemetryReason.Reset); var oldSessionId = _sessionIds.TryGetValue(threadId, out var sid) ? sid : null; var saveToolMeta = false; var saveAttachmentMeta = false; diff --git a/tests/OpenClaw.Shared.Tests/Telemetry/OpenClawTelemetryTests.cs b/tests/OpenClaw.Shared.Tests/Telemetry/OpenClawTelemetryTests.cs index 1113dd347..2b6c3a860 100644 --- a/tests/OpenClaw.Shared.Tests/Telemetry/OpenClawTelemetryTests.cs +++ b/tests/OpenClaw.Shared.Tests/Telemetry/OpenClawTelemetryTests.cs @@ -93,6 +93,22 @@ public void StartDetachedActivity_WithExplicitParent_CreatesChildAndPreservesAmb Assert.Same(ambient, Activity.Current); } + [Fact] + public void StartDetachedActivity_WithEmptyExplicitParent_CreatesRootAndPreservesAmbientActivity() + { + using var collector = ActivityCollector.Listen(OpenClawActivitySourceName.OpenClaw.ToTelemetryName()); + using var ambient = new Activity("ambient").Start(); + + using var root = OpenClawTelemetry.StartDetachedActivity( + "test.root", + default(ActivityContext)); + + Assert.NotNull(root); + Assert.Equal(default, root.ParentSpanId); + Assert.NotEqual(ambient.TraceId, root.TraceId); + Assert.Same(ambient, Activity.Current); + } + [Fact] public void StopDetachedActivity_PreservesNewerAmbientActivity() { diff --git a/tests/OpenClaw.Tray.Tests/ChatTelemetryTrackerTests.cs b/tests/OpenClaw.Tray.Tests/ChatTelemetryTrackerTests.cs new file mode 100644 index 000000000..f1f1f6421 --- /dev/null +++ b/tests/OpenClaw.Tray.Tests/ChatTelemetryTrackerTests.cs @@ -0,0 +1,452 @@ +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Diagnostics.Metrics; +using OpenClaw.Shared.Telemetry; +using OpenClawTray.Chat; + +namespace OpenClaw.Tray.Tests; + +[CollectionDefinition("Chat telemetry", DisableParallelization = true)] +public sealed class ChatTelemetryCollection; + +[Collection("Chat telemetry")] +public sealed class ChatTelemetryTrackerTests +{ + [Fact] + public void ConstantsAndFiniteValues_AreStable() + { + Assert.Equal("openclaw.chat.turn", ChatTelemetryTracker.TurnSpanName); + Assert.Equal("openclaw.chat.queue.wait", ChatTelemetryTracker.QueueWaitSpanName); + Assert.Equal("openclaw.chat.send", ChatTelemetryTracker.SendSpanName); + Assert.Equal("openclaw.chat.response.wait", ChatTelemetryTracker.ResponseWaitSpanName); + Assert.Equal("openclaw.chat.response.receive", ChatTelemetryTracker.ResponseReceiveSpanName); + Assert.Equal("openclaw.chat.history.load", ChatTelemetryTracker.HistoryLoadSpanName); + Assert.Equal("openclaw.chat.history.backfill", ChatTelemetryTracker.HistoryBackfillSpanName); + Assert.Equal("openclaw.chat.turns", ChatTelemetryTracker.TurnsMetricName); + Assert.Equal("openclaw.chat.response.wait.duration", ChatTelemetryTracker.ResponseWaitDurationMetricName); + Assert.Equal("openclaw.chat.response.receive.duration", ChatTelemetryTracker.ResponseReceiveDurationMetricName); + Assert.Equal("openclaw.chat.remote_turns.dropped", ChatTelemetryTracker.DroppedRemoteTurnsMetricName); + Assert.Equal("openclaw.chat.terminal_events.dropped", ChatTelemetryTracker.DroppedTerminalEventsMetricName); + Assert.Equal("success", ChatTelemetryTracker.ToTelemetryValue(ChatTelemetryOutcome.Success)); + Assert.Equal("assistant_final", ChatTelemetryTracker.ToTelemetryValue(ChatTurnTelemetryReason.AssistantFinal)); + Assert.Equal("other", ChatTelemetryTracker.ToTelemetryValue((ChatTurnTelemetryReason)999)); + Assert.Equal("deferred", ChatTelemetryTracker.ToTelemetryValue(ChatAdmissionTelemetryStatus.Deferred)); + Assert.Equal("other", ChatTelemetryTracker.ToTelemetryValue((ChatAdmissionTelemetryStatus)999)); + Assert.Equal("forced", ChatTelemetryTracker.ToTelemetryValue(ChatHistoryTelemetrySource.Forced)); + Assert.Equal("reset_reconciliation", ChatTelemetryTracker.ToTelemetryValue(ChatBackfillTelemetryReason.ResetReconciliation)); + Assert.Equal("missing_run_id", ChatTelemetryTracker.ToTelemetryValue(ChatTerminalEventDropReason.MissingRunId)); + Assert.Equal("mismatched_run_id", ChatTelemetryTracker.ToTelemetryValue(ChatTerminalEventDropReason.MismatchedRunId)); + Assert.Equal("assistant", ChatTelemetryTracker.ToTelemetryValue(ChatResponseOutputKind.Assistant)); + Assert.Equal("other", ChatTelemetryTracker.ToTelemetryValue((ChatResponseOutputKind)999)); + } + + [Fact] + public void LocalTurn_ParentsSendAndCompletesExactlyOnce() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + using var ambient = new Activity("ambient").Start(); + + tracker.StartLocalTurn("private-message", "private-thread", queued: false); + tracker.DispatchLocalTurn("private-message", "private-provisional-run"); + var send = tracker.StartSendAttempt("private-message"); + tracker.FinishSendAttempt( + send, + ChatAdmissionTelemetryStatus.Accepted, + ChatTelemetryOutcome.Success); + tracker.BindAcceptedRun("private-message", "private-accepted-run"); + tracker.ObserveAdmissionAccepted("private-message"); + Assert.True(tracker.ObserveInboundOutput( + "private-thread", + "private-accepted-run", + ChatResponseOutputKind.Assistant)); + Assert.False(tracker.ObserveInboundOutput( + "private-thread", + "private-accepted-run", + ChatResponseOutputKind.Tool)); + + Assert.True(tracker.FinishByRunId( + "private-accepted-run", + ChatTelemetryOutcome.Success, + ChatTurnTelemetryReason.AssistantFinal)); + Assert.False(tracker.FinishByRunId( + "private-provisional-run", + ChatTelemetryOutcome.Failure, + ChatTurnTelemetryReason.LifecycleError)); + + var turn = Assert.Single(activities.Stopped, activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + var sendSpan = Assert.Single(activities.Stopped, activity => activity.OperationName == ChatTelemetryTracker.SendSpanName); + var waitSpan = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + var receiveSpan = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + Assert.Equal(default, turn.ParentSpanId); + Assert.Equal(turn.TraceId, sendSpan.TraceId); + Assert.Equal(turn.SpanId, sendSpan.ParentSpanId); + Assert.Equal(turn.TraceId, waitSpan.TraceId); + Assert.Equal(turn.SpanId, waitSpan.ParentSpanId); + Assert.Equal(turn.TraceId, receiveSpan.TraceId); + Assert.Equal(turn.SpanId, receiveSpan.ParentSpanId); + Assert.Equal("assistant", waitSpan.GetTagItem(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.Equal("assistant", receiveSpan.GetTagItem(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.Equal("success", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("assistant_final", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + Assert.DoesNotContain(turn.Tags, tag => tag.Value?.Contains("private-", StringComparison.Ordinal) == true); + Assert.DoesNotContain(sendSpan.Tags, tag => tag.Value?.Contains("private-", StringComparison.Ordinal) == true); + + Assert.Single(metrics.For(ChatTelemetryTracker.TurnsMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.TurnDurationMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.SendAttemptsMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.ResponseWaitDurationMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.ResponseReceiveDurationMetricName)); + Assert.Empty(metrics.For(ChatTelemetryTracker.QueueWaitDurationMetricName)); + } + + [Fact] + public void TerminalBeforeOutput_RecordsWaitOnly() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: false); + tracker.DispatchLocalTurn("message", "run"); + tracker.ObserveAdmissionAccepted("message"); + + tracker.FinishByRunId( + "run", + ChatTelemetryOutcome.Failure, + ChatTurnTelemetryReason.LifecycleError); + + var wait = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + Assert.Equal("none", wait.GetTagItem(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.Equal("failure", wait.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + var waitMetric = Assert.Single(metrics.For(ChatTelemetryTracker.ResponseWaitDurationMetricName)); + Assert.Equal("none", waitMetric.Tag(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.Empty(metrics.For(ChatTelemetryTracker.ResponseReceiveDurationMetricName)); + } + + [Fact] + public void OutputBeforeAdmission_DoesNotStartResponsePhases() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: false); + tracker.DispatchLocalTurn("message", "run"); + + Assert.False(tracker.ObserveInboundOutput( + "thread", + "run", + ChatResponseOutputKind.Assistant)); + Assert.True(tracker.FinishByRunId( + "run", + ChatTelemetryOutcome.Success, + ChatTurnTelemetryReason.AssistantFinal)); + + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + Assert.Empty(metrics.For(ChatTelemetryTracker.ResponseWaitDurationMetricName)); + Assert.Empty(metrics.For(ChatTelemetryTracker.ResponseReceiveDurationMetricName)); + } + + [Fact] + public void DeferredSend_AccumulatesQueueSegmentsAndRecordsEachAttempt() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + + tracker.StartLocalTurn("message", "thread", queued: true); + tracker.DispatchLocalTurn("message", "attempt-1"); + var first = tracker.StartSendAttempt("message"); + tracker.FinishSendAttempt(first, ChatAdmissionTelemetryStatus.Deferred, ChatTelemetryOutcome.Success); + tracker.RequeueLocalTurn("message"); + tracker.DispatchLocalTurn("message", "attempt-2"); + var second = tracker.StartSendAttempt("message"); + tracker.FinishSendAttempt(second, ChatAdmissionTelemetryStatus.Accepted, ChatTelemetryOutcome.Success); + tracker.BindAcceptedRun("message", "accepted"); + tracker.FinishByRunId("accepted", ChatTelemetryOutcome.Success, ChatTurnTelemetryReason.LifecycleEnd); + + var attempts = metrics.For(ChatTelemetryTracker.SendAttemptsMetricName); + Assert.Equal(2, attempts.Count); + Assert.Contains(attempts, measurement => measurement.Tag(ChatTelemetryTracker.AdmissionStatusTag) == "deferred"); + Assert.Contains(attempts, measurement => measurement.Tag(ChatTelemetryTracker.AdmissionStatusTag) == "accepted"); + Assert.All(attempts, measurement => + Assert.Equal("success", measurement.Tag(OpenClawTelemetryTagKey.Outcome.ToTelemetryName()))); + Assert.Single(metrics.For(ChatTelemetryTracker.QueueWaitDurationMetricName)); + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + var queueWaits = activities.Stopped + .Where(activity => activity.OperationName == ChatTelemetryTracker.QueueWaitSpanName) + .ToArray(); + Assert.Equal(2, queueWaits.Length); + Assert.All(queueWaits, queueWait => + { + Assert.Equal(turn.TraceId, queueWait.TraceId); + Assert.Equal(turn.SpanId, queueWait.ParentSpanId); + Assert.Equal("local", queueWait.GetTagItem(OpenClawTelemetryTagKey.Source.ToTelemetryName())); + Assert.Equal("success", queueWait.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + }); + } + + [Fact] + public void QueuedTurnCanceledBeforeDispatch_ClosesQueueWaitWithTurnOutcome() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: true); + + Assert.True(tracker.FinishByMessageId( + "message", + ChatTelemetryOutcome.Canceled, + ChatTurnTelemetryReason.QueuedCanceled)); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + var queueWait = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.QueueWaitSpanName); + Assert.Equal(turn.TraceId, queueWait.TraceId); + Assert.Equal(turn.SpanId, queueWait.ParentSpanId); + Assert.Equal("canceled", queueWait.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + var queueMetric = Assert.Single(metrics.For(ChatTelemetryTracker.QueueWaitDurationMetricName)); + Assert.Equal("canceled", queueMetric.Tag(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + } + + [Fact] + public async Task ConcurrentTerminalSignals_RecordOneTurn() + { + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: false); + tracker.DispatchLocalTurn("message", "run"); + + await Task.WhenAll( + Task.Run(() => tracker.FinishByRunId( + "run", + ChatTelemetryOutcome.Success, + ChatTurnTelemetryReason.AssistantFinal)), + Task.Run(() => tracker.FinishByRunId( + "run", + ChatTelemetryOutcome.Failure, + ChatTurnTelemetryReason.LifecycleError))); + + Assert.Single(metrics.For(ChatTelemetryTracker.TurnsMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.TurnDurationMetricName)); + } + + [Fact] + public async Task DispatchAndTerminalRace_StopsQueueWaitBeforeTurn() + { + using var queueStopEntered = new ManualResetEventSlim(); + using var releaseQueueStop = new ManualResetEventSlim(); + using var terminalStarted = new ManualResetEventSlim(); + using var activities = new ActivityCollector(activity => + { + if (activity.OperationName != ChatTelemetryTracker.QueueWaitSpanName) + return; + queueStopEntered.Set(); + Assert.True(releaseQueueStop.Wait(TimeSpan.FromSeconds(5))); + }); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: true); + + var dispatch = Task.Run(() => tracker.DispatchLocalTurn("message", "run")); + Assert.True(queueStopEntered.Wait(TimeSpan.FromSeconds(5))); + var terminal = Task.Run(() => + { + terminalStarted.Set(); + return tracker.FinishByRunId( + "run", + ChatTelemetryOutcome.Success, + ChatTurnTelemetryReason.LifecycleEnd); + }); + Assert.True(terminalStarted.Wait(TimeSpan.FromSeconds(5))); + Assert.NotSame(terminal, await Task.WhenAny(terminal, Task.Delay(TimeSpan.FromMilliseconds(100)))); + + releaseQueueStop.Set(); + await Task.WhenAll(dispatch, terminal); + + var stoppedNames = activities.Stopped.Select(activity => activity.OperationName).ToArray(); + Assert.True( + Array.IndexOf(stoppedNames, ChatTelemetryTracker.QueueWaitSpanName) < + Array.IndexOf(stoppedNames, ChatTelemetryTracker.TurnSpanName)); + } + + [Fact] + public void RemoteTurnWithoutRunId_RecordsDropButNoTurn() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + + tracker.ObserveLifecycleStart("private-thread", runId: null); + + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + var dropped = Assert.Single(metrics.For(ChatTelemetryTracker.DroppedRemoteTurnsMetricName)); + Assert.Equal("missing_run_id", dropped.Tag(ChatTelemetryTracker.DroppedRemoteTurnReasonTag)); + } + + [Fact] + public void LocalTurnWithoutLifecycleRunId_DoesNotRecordRemoteDrop() + { + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: false); + tracker.DispatchLocalTurn("message", "provisional-run"); + + tracker.ObserveLifecycleStart("thread", runId: null); + + Assert.Empty(metrics.For(ChatTelemetryTracker.DroppedRemoteTurnsMetricName)); + tracker.FinishAll(ChatTelemetryOutcome.Canceled, ChatTurnTelemetryReason.Disposed); + } + + [Fact] + public void PreparedCompletion_ReservesUnderLockAndEmitsAfterward() + { + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + tracker.StartLocalTurn("message", "thread", queued: false); + + var completion = tracker.PrepareFinishByMessageId( + "message", + ChatTelemetryOutcome.Failure, + ChatTurnTelemetryReason.SendRejected); + + Assert.NotNull(completion); + Assert.Null(tracker.PrepareFinishByMessageId( + "message", + ChatTelemetryOutcome.Canceled, + ChatTurnTelemetryReason.Disconnected)); + Assert.Empty(metrics.For(ChatTelemetryTracker.TurnsMetricName)); + Assert.True(tracker.CompletePreparedTurn(completion)); + Assert.False(tracker.CompletePreparedTurn(completion)); + var turn = Assert.Single(metrics.For(ChatTelemetryTracker.TurnsMetricName)); + Assert.Equal("send_rejected", turn.Tag(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + } + + [Fact] + public void DroppedTerminalEvents_RecordOnlyFiniteReasons() + { + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + + tracker.RecordDroppedTerminalEvent(ChatTerminalEventDropReason.MissingRunId); + tracker.RecordDroppedTerminalEvent(ChatTerminalEventDropReason.MismatchedRunId); + + var dropped = metrics.For(ChatTelemetryTracker.DroppedTerminalEventsMetricName); + Assert.Equal(2, dropped.Count); + Assert.Contains( + dropped, + measurement => measurement.Tag(ChatTelemetryTracker.DroppedTerminalEventReasonTag) == "missing_run_id"); + Assert.Contains( + dropped, + measurement => measurement.Tag(ChatTelemetryTracker.DroppedTerminalEventReasonTag) == "mismatched_run_id"); + } + + [Fact] + public void HistoryOperations_RecordOnlyAllowlistedTags() + { + using var activities = new ActivityCollector(); + using var metrics = new MetricCollector(); + var tracker = new ChatTelemetryTracker(); + using var ambient = new Activity("ambient").Start(); + + var load = tracker.StartHistoryLoad(ChatHistoryTelemetrySource.Forced); + tracker.FinishHistoryLoad(load, ChatTelemetryOutcome.Success); + var backfill = tracker.StartHistoryBackfill(ChatBackfillTelemetryReason.RemoteTurn); + tracker.FinishHistoryBackfill(backfill, ChatTelemetryOutcome.Failure, new InvalidOperationException("private-error")); + + var loadSpan = Assert.Single(activities.Stopped, activity => activity.OperationName == ChatTelemetryTracker.HistoryLoadSpanName); + Assert.Equal(default, loadSpan.ParentSpanId); + Assert.NotEqual(ambient.TraceId, loadSpan.TraceId); + Assert.Equal(["openclaw.outcome", "openclaw.source"], loadSpan.Tags.Select(tag => tag.Key).Order().ToArray()); + var backfillSpan = Assert.Single(activities.Stopped, activity => activity.OperationName == ChatTelemetryTracker.HistoryBackfillSpanName); + Assert.Equal(default, backfillSpan.ParentSpanId); + Assert.NotEqual(ambient.TraceId, backfillSpan.TraceId); + Assert.Equal( + ["error.type", "openclaw.chat.backfill.reason", "openclaw.outcome", "openclaw.source"], + backfillSpan.Tags.Select(tag => tag.Key).Order().ToArray()); + Assert.DoesNotContain(backfillSpan.Tags, tag => tag.Value?.Contains("private-error", StringComparison.Ordinal) == true); + Assert.Single(metrics.For(ChatTelemetryTracker.HistoryLoadsMetricName)); + Assert.Single(metrics.For(ChatTelemetryTracker.HistoryBackfillsMetricName)); + } + + private sealed class ActivityCollector : IDisposable + { + private readonly ActivityListener _listener; + + public ActivityCollector(Action? activityStopped = null) + { + _listener = new ActivityListener + { + ShouldListenTo = source => source.Name == OpenClawActivitySourceName.OpenClaw.ToTelemetryName(), + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + ActivityStopped = activity => + { + activityStopped?.Invoke(activity); + Stopped.Enqueue(activity); + }, + }; + ActivitySource.AddActivityListener(_listener); + } + + public ConcurrentQueue Stopped { get; } = []; + + public void Dispose() => _listener.Dispose(); + } + + private sealed class MetricCollector : IDisposable + { + private readonly MeterListener _listener = new(); + private readonly ConcurrentBag _measurements = []; + + public MetricCollector() + { + _listener.InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == OpenClawMeterName.OpenClaw.ToTelemetryName() && + instrument.Name.StartsWith("openclaw.chat.", StringComparison.Ordinal)) + { + listener.EnableMeasurementEvents(instrument); + } + }; + _listener.SetMeasurementEventCallback((instrument, value, tags, _) => + _measurements.Add(new Measurement(instrument.Name, value, tags.ToArray()))); + _listener.SetMeasurementEventCallback((instrument, value, tags, _) => + _measurements.Add(new Measurement(instrument.Name, value, tags.ToArray()))); + _listener.Start(); + } + + public List For(string name) => + _measurements.Where(measurement => measurement.Name == name).ToList(); + + public void Dispose() => _listener.Dispose(); + } + + private sealed record Measurement( + string Name, + object Value, + KeyValuePair[] Tags) + { + public string? Tag(string key) => + Tags.FirstOrDefault(tag => tag.Key == key).Value?.ToString(); + } +} diff --git a/tests/OpenClaw.Tray.Tests/OpenClaw.Tray.Tests.csproj b/tests/OpenClaw.Tray.Tests/OpenClaw.Tray.Tests.csproj index e5394fad9..a1e16fd58 100644 --- a/tests/OpenClaw.Tray.Tests/OpenClaw.Tray.Tests.csproj +++ b/tests/OpenClaw.Tray.Tests/OpenClaw.Tray.Tests.csproj @@ -39,6 +39,7 @@ + diff --git a/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs b/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs index 4e1aa1f74..76a19bda8 100644 --- a/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs +++ b/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs @@ -1,12 +1,67 @@ using OpenClaw.Chat; using OpenClaw.Shared; +using OpenClaw.Shared.Telemetry; using OpenClawTray.Chat; +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Diagnostics.Metrics; using System.Text.Json; namespace OpenClaw.Tray.Tests; +[Collection("Chat telemetry")] public class OpenClawChatDataProviderTests { + private sealed class ChatActivityCollector : IDisposable + { + private readonly ActivityListener _listener; + + public ChatActivityCollector() + { + _listener = new ActivityListener + { + ShouldListenTo = source => source.Name == OpenClawActivitySourceName.OpenClaw.ToTelemetryName(), + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded, + ActivityStopped = activity => Stopped.Enqueue(activity), + }; + ActivitySource.AddActivityListener(_listener); + } + + public ConcurrentQueue Stopped { get; } = []; + + public void Dispose() => _listener.Dispose(); + } + + private sealed class ChatMetricCollector : IDisposable + { + private readonly MeterListener _listener = new(); + private readonly ConcurrentQueue<(string Name, KeyValuePair[] Tags)> _measurements = []; + + public ChatMetricCollector() + { + _listener.InstrumentPublished = (instrument, listener) => + { + if (instrument.Meter.Name == OpenClawMeterName.OpenClaw.ToTelemetryName() && + instrument.Name.StartsWith("openclaw.chat.", StringComparison.Ordinal)) + { + listener.EnableMeasurementEvents(instrument); + } + }; + _listener.SetMeasurementEventCallback((instrument, _, tags, _) => + _measurements.Enqueue((instrument.Name, tags.ToArray()))); + _listener.Start(); + } + + public string[] TagsFor(string metricName, string tagName) => + _measurements + .Where(measurement => measurement.Name == metricName) + .Select(measurement => + measurement.Tags.First(tag => tag.Key == tagName).Value?.ToString() ?? string.Empty) + .ToArray(); + + public void Dispose() => _listener.Dispose(); + } + private sealed class FakeBridge : IChatGatewayBridge { public bool IsConnected { get; set; } @@ -160,6 +215,352 @@ private static AgentEventInfo MakeAgentEvent(string stream, string json, string }; } + [Fact] + public async Task Telemetry_LocalSendAndLifecycle_EmitCorrelatedAllowlistedSpans() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "private-run", Status = "started" }); + + await provider.SendMessageAsync("main", "private prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "private-run")); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", + Role = "assistant", + Text = "private response", + State = "delta", + }); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "private-run")); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", + Role = "assistant", + Text = "private response", + State = "final", + }); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + var send = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.SendSpanName); + var wait = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + var receive = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + Assert.Equal(turn.TraceId, send.TraceId); + Assert.Equal(turn.SpanId, send.ParentSpanId); + Assert.Equal(turn.TraceId, wait.TraceId); + Assert.Equal(turn.SpanId, wait.ParentSpanId); + Assert.Equal(turn.TraceId, receive.TraceId); + Assert.Equal(turn.SpanId, receive.ParentSpanId); + Assert.Equal("local", turn.GetTagItem(OpenClawTelemetryTagKey.Source.ToTelemetryName())); + Assert.Equal("success", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("lifecycle_end", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + Assert.Equal("accepted", send.GetTagItem(ChatTelemetryTracker.AdmissionStatusTag)); + Assert.Equal("assistant", wait.GetTagItem(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.Equal("assistant", receive.GetTagItem(ChatTelemetryTracker.FirstOutputKindTag)); + Assert.DoesNotContain(turn.Tags, tag => tag.Value?.Contains("private", StringComparison.Ordinal) == true); + Assert.DoesNotContain(send.Tags, tag => tag.Value?.Contains("private", StringComparison.Ordinal) == true); + Assert.DoesNotContain(wait.Tags, tag => tag.Value?.Contains("private", StringComparison.Ordinal) == true); + Assert.DoesNotContain(receive.Tags, tag => tag.Value?.Contains("private", StringComparison.Ordinal) == true); + + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_UnknownAdmissionStatus_MapsToOther() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "future_status" }); + + await provider.SendMessageAsync("main", "prompt"); + + var send = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.SendSpanName); + Assert.Equal("other", send.GetTagItem(ChatTelemetryTracker.AdmissionStatusTag)); + await provider.DisposeAsync(); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + } + + [Fact] + public async Task Telemetry_RemoteLifecycle_EmitsRemoteTurn() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "remote-run")); + bridge.RaiseAgent(MakeAgentEvent("reasoning", """{"delta":"response"}""", runId: "remote-run")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "remote-run")); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("remote", turn.GetTagItem(OpenClawTelemetryTagKey.Source.ToTelemetryName())); + Assert.Equal("lifecycle_end", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseWaitSpanName); + Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.ResponseReceiveSpanName); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_LoadHistory_EmitsBoundedHistorySpan() + { + using var activities = new ChatActivityCollector(); + var (_, provider, _, _) = CreateProvider(new[] { MainSession() }); + + await provider.LoadHistoryAsync("main", force: true); + + var history = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.HistoryLoadSpanName); + Assert.Equal("forced", history.GetTagItem(OpenClawTelemetryTagKey.Source.ToTelemetryName())); + Assert.Equal("success", history.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal( + ["openclaw.outcome", "openclaw.source"], + history.Tags.Select(tag => tag.Key).Order().ToArray()); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_AbortIntent_CompletesTurnAsCanceled() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + await provider.StopResponseAsync("main"); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("canceled", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("abort_requested", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_AbortBeforeLifecycleStart_DoesNotCreateRemoteTurn() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { Status = "started" }); + + await provider.SendMessageAsync("main", "prompt"); + await provider.StopResponseAsync("main"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run")); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("local", turn.GetTagItem(OpenClawTelemetryTagKey.Source.ToTelemetryName())); + Assert.Equal("canceled", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("abort_requested", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_Disconnect_CompletesOutstandingTurn() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + bridge.RaiseStatus(ConnectionStatus.Connected); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + bridge.RaiseStatus(ConnectionStatus.Disconnected); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("canceled", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("disconnected", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_UncorrelatedTerminalEvents_AreDiagnosedWithoutGuessing() + { + using var activities = new ChatActivityCollector(); + using var metrics = new ChatMetricCollector(); + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + bridge.RaiseStatus(ConnectionStatus.Connected); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + var snapshotsBeforeMismatchedTerminal = snapshots.Count; + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "different-run")); + + Assert.Equal(snapshotsBeforeMismatchedTerminal, snapshots.Count); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal( + ["mismatched_run_id"], + metrics.TagsFor( + ChatTelemetryTracker.DroppedTerminalEventsMetricName, + ChatTelemetryTracker.DroppedTerminalEventReasonTag)); + + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run")); + Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_LifecycleTerminalWithoutRunId_PreservesActiveRunUntilExactTerminal() + { + using var activities = new ChatActivityCollector(); + using var metrics = new ChatMetricCollector(); + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "next-run", Status = "started" }); + bridge.RaiseStatus(ConnectionStatus.Connected); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + await provider.SendMessageAsync("main", "queued"); + var snapshotsBeforeMalformedTerminal = snapshots.Count; + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: null)); + + Assert.Equal(snapshotsBeforeMalformedTerminal, snapshots.Count); + Assert.Equal(["prompt"], bridge.SentMessages); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal( + ["missing_run_id"], + metrics.TagsFor( + ChatTelemetryTracker.DroppedTerminalEventsMetricName, + ChatTelemetryTracker.DroppedTerminalEventReasonTag)); + + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run")); + Assert.True(SpinWait.SpinUntil( + () => bridge.SentMessages.Count == 2, + TimeSpan.FromSeconds(5))); + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("success", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("lifecycle_end", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + Assert.Equal(["prompt", "queued"], bridge.SentMessages); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_LegacyJobTerminalWithoutRunId_PreservesActiveRunUntilExactTerminal() + { + using var activities = new ChatActivityCollector(); + using var metrics = new ChatMetricCollector(); + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + bridge.RaiseStatus(ConnectionStatus.Connected); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + var snapshotsBeforeMalformedTerminal = snapshots.Count; + bridge.RaiseAgent(MakeAgentEvent("job", """{"state":"done"}""", runId: null)); + + Assert.Equal(snapshotsBeforeMalformedTerminal, snapshots.Count); + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal( + ["missing_run_id"], + metrics.TagsFor( + ChatTelemetryTracker.DroppedTerminalEventsMetricName, + ChatTelemetryTracker.DroppedTerminalEventReasonTag)); + + bridge.RaiseAgent(MakeAgentEvent("job", """{"state":"done"}""", runId: "run")); + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("success", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("lifecycle_end", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_LifecycleTerminalWithoutRunId_RemainsEligibleForDisconnectCleanup() + { + using var activities = new ChatActivityCollector(); + using var metrics = new ChatMetricCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + bridge.RaiseStatus(ConnectionStatus.Connected); + + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: null)); + + Assert.DoesNotContain( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + bridge.RaiseStatus(ConnectionStatus.Disconnected); + + var turn = Assert.Single( + activities.Stopped, + activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName); + Assert.Equal("canceled", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("disconnected", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + Assert.Equal( + ["missing_run_id"], + metrics.TagsFor( + ChatTelemetryTracker.DroppedTerminalEventsMetricName, + ChatTelemetryTracker.DroppedTerminalEventReasonTag)); + await provider.DisposeAsync(); + } + + [Fact] + public async Task Telemetry_Reset_CompletesQueuedAndActiveTurns() + { + using var activities = new ChatActivityCollector(); + var (bridge, provider, _, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); + + await provider.SendMessageAsync("main", "active"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + await provider.SendMessageAsync("main", "queued"); + bridge.RaiseSessionCommandCompleted(new SessionCommandResult + { + Method = "sessions.reset", + Ok = true, + Key = "main", + }); + + var turns = activities.Stopped + .Where(activity => activity.OperationName == ChatTelemetryTracker.TurnSpanName) + .ToArray(); + Assert.Equal(2, turns.Length); + Assert.All(turns, turn => + { + Assert.Equal("canceled", turn.GetTagItem(OpenClawTelemetryTagKey.Outcome.ToTelemetryName())); + Assert.Equal("reset", turn.GetTagItem(OpenClawTelemetryTagKey.Reason.ToTelemetryName())); + }); + await provider.DisposeAsync(); + } + [Fact] public async Task LoadAsync_ReturnsSeededSessionsAsThreads() { @@ -1164,10 +1565,14 @@ public async Task AgentEvent_ToolStartThenResult_MarksToolSuccess() public async Task AgentEvent_JobError_EmitsErrorEntry() { var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); await provider.LoadAsync(); + bridge.RaiseStatus(ConnectionStatus.Connected); + await provider.SendMessageAsync("main", "prompt"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); snapshots.Clear(); - var evt = MakeAgentEvent("job", """{"state":"error"}"""); + var evt = MakeAgentEvent("job", """{"state":"error"}""", runId: "run"); evt.Summary = "kaboom"; bridge.RaiseAgent(evt); @@ -1211,11 +1616,13 @@ await WaitForConditionAsync(() => public async Task AgentEvent_JobDone_ClearsTurnActive() { var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run", Status = "started" }); await provider.LoadAsync(); - // Kick off a turn - _ = provider.SendMessageAsync("main", "hi"); + bridge.RaiseStatus(ConnectionStatus.Connected); + await provider.SendMessageAsync("main", "hi"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); - bridge.RaiseAgent(MakeAgentEvent("job", """{"state":"done"}""")); + bridge.RaiseAgent(MakeAgentEvent("job", """{"state":"done"}""", runId: "run")); // Snapshot the timeline directly. var snap = await provider.LoadAsync(); @@ -3023,6 +3430,7 @@ public async Task QueuedSend_LifecycleStartBeforeAck_PromotesByIdempotencyKey() [Fact] public async Task QueuedSend_InFlightAckWithoutLifecycle_RequeuesAndRetriesSameIdempotencyKey() { + using var activities = new ChatActivityCollector(); var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-1", Status = "started" }); bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-2", Status = "in_flight" }); @@ -3076,6 +3484,16 @@ public async Task QueuedSend_InFlightAckWithoutLifecycle_RequeuesAndRetriesSameI Assert.Empty(GetQueuedMessages(snapshots[^1], "main")); Assert.Single(snapshots[^1].Timelines["main"].Entries, e => e.Kind == ChatTimelineItemKind.User && e.Text == "Hello"); + + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-2")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run-2")); + + var admissionStatuses = activities.Stopped + .Where(activity => activity.OperationName == ChatTelemetryTracker.SendSpanName) + .Select(activity => activity.GetTagItem(ChatTelemetryTracker.AdmissionStatusTag)) + .ToArray(); + Assert.Equal(new object?[] { "accepted", "deferred", "accepted" }, admissionStatuses); + await provider.DisposeAsync(); } [Fact] @@ -5615,7 +6033,10 @@ public async Task LoadHistoryAsync_AfterLiveActivity_PreservesNonDuplicateLiveEn await provider.LoadAsync(); // A live event the history will NOT carry — must survive the rebuild. - bridge.RaiseAgent(MakeAgentEvent("lifecycle", "{\"phase\":\"error\",\"message\":\"net glitch\"}")); + bridge.RaiseAgent(MakeAgentEvent( + "lifecycle", + "{\"phase\":\"error\",\"message\":\"net glitch\"}", + runId: "run")); await provider.LoadHistoryAsync("main");