From 896d323601866c74014d83ec8c70d9fa1fd1e4f7 Mon Sep 17 00:00:00 2001 From: "bakudies@microsoft.com" Date: Thu, 17 Sep 2026 16:02:47 -0700 Subject: [PATCH] fix(chat): keep cancelled responses from ending newer turns Preserve gateway run identity, isolate stale chat and tool output, and gate notifications on chat admission. Keep trailing finals, bounded terminal retention, and legitimate lifecycle retries consistent with the pinned upstream contract. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- docs/CONNECTION_ARCHITECTURE.md | 31 ++ src/OpenClaw.Chat/ChatTimelineReducer.cs | 22 ++ src/OpenClaw.Shared/Models.cs | 13 + src/OpenClaw.Shared/OpenClawGatewayClient.cs | 34 +- .../Chat/ChatConversationState.cs | 67 +++- .../Chat/ChatEventMapper.cs | 48 +++ .../Chat/ChatLifecycleState.cs | 124 ++++++- .../Chat/ChatResetState.cs | 12 +- .../Chat/OpenClawChatDataProvider.cs | 15 +- .../OpenClawGatewayClientTests.cs | 124 +++++++ .../ChatRunCorrelationTests.cs | 325 ++++++++++++++++++ .../OpenClawChatDataProviderTests.cs | 287 ++++++++++++++++ 12 files changed, 1071 insertions(+), 31 deletions(-) create mode 100644 tests/OpenClaw.Tray.Tests/ChatRunCorrelationTests.cs diff --git a/docs/CONNECTION_ARCHITECTURE.md b/docs/CONNECTION_ARCHITECTURE.md index dee65b1cb..32761bebb 100644 --- a/docs/CONNECTION_ARCHITECTURE.md +++ b/docs/CONNECTION_ARCHITECTURE.md @@ -89,6 +89,37 @@ The operator client is received through the `OperatorClientChanged` event. The a Inbound chat and agent timeline events must include the gateway's canonical `sessionKey`. The tray client must not synthesize a literal `main` key for keyless inbound events, because that can merge unrelated events into the wrong timeline. When a keyless chat or agent event arrives, the tray drops it and raises a one-shot diagnostic so the protocol issue is visible without exposing the dropped message contents. +Live `chat` and `session.message` events also retain the payload's optional +`runId`, separately from message identity. `ChatConversationState` consults +`ChatLifecycleState` before admitting assistant or tool chat output, so an +aborted, completed, or mismatched run cannot replace output or end a newer +turn. Reset admission also checks the supplied run against ignored reset runs. +The most recent non-aborted lifecycle completion may still receive its final +text while idle. Suppressed terminals do not change that eligibility, and the +bounded cache evicts by arrival order rather than suppression priority. The +terminal cache also suppresses late nonterminal agent output after abort +cleanup. Late tool repair requires a retained correlation for the same run; +late approvals must match the pending approval. Known late child/parent repair +does not reactivate an idle turn. An aborted lifecycle start cannot replace +the current run. Frames without a run ID retain legacy thread-level handling; +they cannot be reliably correlated to an older run. + +Upstream lifecycle errors may be followed by a same-run retry. The most recent +non-aborted failed run may restart while idle, before any newer active or +completed run or accepted assistant final. Explicit execution settlement, +exhausted fallback, timeout, cancellation and terminal liveness facts remain +closed, matching upstream +[`isDefinitiveRunLifecycle`](https://github.com/openclaw/openclaw/blob/eb82ef8b80a05058619557dff09f758a6710d4d5/packages/normalization-core/src/agent-run-terminal-outcome.ts#L145-L235). +Successful completion and abort fences remain closed to repeated starts. + +Chat admission also controls notification delivery. The synchronous provider +marks rejected `ChatMessageInfo` frames with `SuppressNotification()` before +the gateway parser emits its separate notification event. That per-frame veto +is monotonic and does not suppress delivery to other chat subscribers. It +prevents rejected finals from reaching notification history, toasts or TTS +without adding another run cache or moving notification ownership into `App`. +Consumers that do not apply a veto retain the existing notification behavior. + ## Startup wiring (App.xaml.cs) ``` diff --git a/src/OpenClaw.Chat/ChatTimelineReducer.cs b/src/OpenClaw.Chat/ChatTimelineReducer.cs index a39bebcac..7b746eee9 100644 --- a/src/OpenClaw.Chat/ChatTimelineReducer.cs +++ b/src/OpenClaw.Chat/ChatTimelineReducer.cs @@ -44,6 +44,28 @@ public static ChatTimelineState RebuildActiveToolTracking(ChatTimelineState stat }; } + /// + /// A completed run may repair a retained tool row or materialize a known + /// terminal child, but must not create a new tool activity from a stale frame. + /// + public static bool CanReconcileToolAfterTurnEnd(ChatTimelineState state, ChatEvent? evt) + { + var (runId, toolCallId) = evt switch + { + ChatToolStartEvent e => (e.RunId, e.ToolCallId), + ChatToolPresentationEvent e => (e.RunId, e.ParentToolCallId), + ChatToolOutputEvent e => (e.RunId, e.ToolCallId), + ChatToolErrorEvent e => (e.RunId, e.ToolCallId), + _ => ((string?)null, (string?)null), + }; + if (string.IsNullOrWhiteSpace(runId) || string.IsNullOrWhiteSpace(toolCallId)) + return false; + + var key = CurrentCorrelationKey(state, runId, toolCallId); + return FindToolEntryIndex(state, key, allowTerminalLegacy: true) >= 0 || + state.TerminalToolCorrelations?.ContainsKey(key) == true; + } + public static ChatTimelineState Apply(ChatTimelineState state, ChatEvent evt) { return evt switch diff --git a/src/OpenClaw.Shared/Models.cs b/src/OpenClaw.Shared/Models.cs index 5abb60dc0..f88bc3852 100644 --- a/src/OpenClaw.Shared/Models.cs +++ b/src/OpenClaw.Shared/Models.cs @@ -1798,6 +1798,19 @@ public static bool IsSilentAssistantDirective(string? role, string? text) => /// Session this message belongs to (e.g. "main"). public string SessionKey { get; set; } = ""; + /// Optional run identity from the live chat event payload, not the message ID. + public string? RunId { get; set; } + + /// Set by synchronous chat consumers when this frame must not produce a notification. + [System.Text.Json.Serialization.JsonIgnore] + public bool IsNotificationSuppressed { get; private set; } + + /// + /// Suppresses the parser's subsequent notification without changing other + /// consumers' delivery. Suppression is monotonic for this received frame. + /// + public void SuppressNotification() => IsNotificationSuppressed = true; + /// "user", "assistant", "system", etc. public string Role { get; set; } = ""; diff --git a/src/OpenClaw.Shared/OpenClawGatewayClient.cs b/src/OpenClaw.Shared/OpenClawGatewayClient.cs index 151563fd3..1fc37fdd7 100644 --- a/src/OpenClaw.Shared/OpenClawGatewayClient.cs +++ b/src/OpenClaw.Shared/OpenClawGatewayClient.cs @@ -3851,6 +3851,11 @@ private void HandleChatEvent(JsonElement root, int rawMessageLength) if (string.IsNullOrEmpty(sessionKey)) _logger.Warn("[GatewayClient] Chat event missing sessionKey; will be dropped downstream."); + var runId = payload.TryGetProperty("runId", out var runIdProperty) && + runIdProperty.ValueKind == JsonValueKind.String + ? runIdProperty.GetString() + : null; + // Best-effort usage extraction — gateway emits this only on terminal // (state="final") events in practice; we still read it defensively // from common locations so any reasonable shape lights up the chat @@ -3888,7 +3893,7 @@ private void HandleChatEvent(JsonElement root, int rawMessageLength) var messageOpenClawMetadata = ExtractOpenClawMetadata(message); var payloadOpenClawMetadata = ExtractOpenClawMetadata(payload); - EmitChatMessageReceived( + var notify = EmitChatMessageReceived( sessionKey, role, text, @@ -3903,9 +3908,10 @@ private void HandleChatEvent(JsonElement root, int rawMessageLength) messageOpenClawMetadata.Kind ?? payloadOpenClawMetadata.Kind, messageOpenClawMetadata.TokensBefore ?? payloadOpenClawMetadata.TokensBefore, messageOpenClawMetadata.TokensAfter ?? payloadOpenClawMetadata.TokensAfter, - contentParts); + contentParts, + runId); - if (role == "assistant" && string.Equals(state, "final", StringComparison.OrdinalIgnoreCase)) + if (notify && role == "assistant" && string.Equals(state, "final", StringComparison.OrdinalIgnoreCase)) { // HIGH 4: log shape only — content previously // surfaced in the operator log. @@ -3929,7 +3935,7 @@ private void HandleChatEvent(JsonElement root, int rawMessageLength) if (ChatMessageInfo.IsSilentAssistantDirective(role, text)) return; var openClawMetadata = ExtractOpenClawMetadata(payload); - EmitChatMessageReceived( + var notify = EmitChatMessageReceived( sessionKey, role, text, @@ -3944,9 +3950,10 @@ private void HandleChatEvent(JsonElement root, int rawMessageLength) openClawMetadata.Kind, openClawMetadata.TokensBefore, openClawMetadata.TokensAfter, - projection.ContentParts); + projection.ContentParts, + runId); - if (role == "assistant" && + if (notify && role == "assistant" && (string.IsNullOrWhiteSpace(state) || string.Equals(state, "final", StringComparison.OrdinalIgnoreCase))) { @@ -4009,7 +4016,7 @@ JsonValueKind.Number when v.TryGetInt32(out var i) => i, return (input, output, response, ctx); } - private void EmitChatMessageReceived( + private bool EmitChatMessageReceived( string sessionKey, string role, string text, @@ -4024,16 +4031,18 @@ private void EmitChatMessageReceived( string? openClawKind = null, long? compactionTokensBefore = null, long? compactionTokensAfter = null, - IReadOnlyList? contentParts = null) + IReadOnlyList? contentParts = null, + string? runId = null) { if (ChatMessageInfo.IsSilentAssistantDirective(role, text)) - return; + return false; try { - ChatMessageReceived?.Invoke(this, new ChatMessageInfo + var message = new ChatMessageInfo { SessionKey = sessionKey, + RunId = runId, Role = role, Text = text, ContentParts = contentParts ?? Array.Empty(), @@ -4048,11 +4057,14 @@ private void EmitChatMessageReceived( OpenClawKind = openClawKind, CompactionTokensBefore = compactionTokensBefore, CompactionTokensAfter = compactionTokensAfter - }); + }; + ChatMessageReceived?.Invoke(this, message); + return !message.IsNotificationSuppressed; } catch (Exception ex) { _logger.Warn($"ChatMessageReceived handler threw: {ex.Message}"); + return false; } } diff --git a/src/OpenClaw.Tray.WinUI/Chat/ChatConversationState.cs b/src/OpenClaw.Tray.WinUI/Chat/ChatConversationState.cs index b88285bf6..56f291d95 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/ChatConversationState.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/ChatConversationState.cs @@ -569,6 +569,7 @@ internal ChatQueuedAdmission AdmitMessage( _lifecycle.ClearThreadSuppression(threadId); _lifecycle.TakePendingAbortCount(threadId); + _lifecycle.ClearRetryEligibility(threadId); var request = new ChatQueuedSendRequest( messageId, Guid.NewGuid().ToString(), @@ -1327,6 +1328,22 @@ internal ChatIncomingMessageGate GateIncomingChatMessage( var hasMediaEnvelope = projection?.HasMediaEnvelope ?? false; lock (_gate) { + if (role is "assistant" or "toolresult" or "tool_result" && + _lifecycle.ShouldSuppressChatMessage( + threadId, + message.RunId, + message.IsFinal, + _queue.RunIdsForThread(threadId), + _timelines.TryGetValue(threadId, out var timeline) && timeline.TurnActive)) + { + return new( + Drop: false, + Suppressed: true, + RequestRemoteBackfill: false, + Snapshot: null, + OpenedLifecycle: null, + CurrentRuntimeGenerationLocked(threadId)); + } _lifecycle.TryGetActiveRun( threadId, out var activeRunId); @@ -1340,7 +1357,8 @@ internal ChatIncomingMessageGate GateIncomingChatMessage( text, attachmentCorrelationSignature, hasMediaEnvelope), - activeRunId); + activeRunId, + message.RunId); var openedLifecycle = ApplyBufferedLifecycleOpenLocked( threadId, @@ -1692,11 +1710,11 @@ message.ResponseTokens is not null || } } - internal string? CompleteAssistantFinal(string threadId) + internal string? CompleteAssistantFinal(string threadId, string? runId = null) { lock (_gate) { - var completedRunId = _lifecycle.CompleteAssistantFinal(threadId); + var completedRunId = _lifecycle.CompleteAssistantFinal(threadId, runId); _reset.CompleteRun(threadId, completedRunId); if (!_queue.HasSendingMessages(threadId)) _queue.ClearLocallyInitiated(threadId); @@ -1745,6 +1763,11 @@ internal ChatAgentEventTransition ProcessAgentEvent( { var mapping = ChatEventMapper.Map(evt); mapped = mapping.Event; + if (mapped is ChatToolPresentationEvent presentation && + _lifecycle.IsCompletedRun(threadId, evt.RunId)) + { + mapped = presentation with { ActivatesTurn = false }; + } if (mapping.Approval is { } approval && !_approval.MarkSeen(approval.RequestId, approval.AlternateId)) { @@ -1914,12 +1937,35 @@ private ChatAgentEventGate GateAgentEventLocked( } if (resetGate.Drop) { + // Reset consumes an ignored terminal to request history reconciliation. + // Retain its identity for chat finals that can arrive after that terminal. + if (resetGate.ReloadHistory && !string.IsNullOrWhiteSpace(evt.RunId)) + _lifecycle.RememberSuppressedTerminal(threadId, evt.RunId); return new( false, resetGate.ReloadHistory, null, openedLifecycle); } + if (ChatEventMapper.IsLifecycleStart(evt) && _lifecycle.IsRunAborted(evt.RunId)) + { + Logger.Debug($"[ChatProvider] Dropping aborted-run lifecycle start for threadId='{threadId}'"); + return new(false, false, null, openedLifecycle); + } + if (ChatEventMapper.IsLifecycleStart(evt) && + (!_timelines.TryGetValue(threadId, out var currentTimeline) || !currentTimeline.TurnActive) && + _lifecycle.TryRestartFailedRun(threadId, evt.RunId)) + { + Logger.Debug($"[ChatProvider] Admitting same-run lifecycle retry for threadId='{threadId}'"); + } + if (!ChatEventMapper.IsTerminalRunEvent(evt) && + _lifecycle.IsCompletedRun(threadId, evt.RunId) && + _lifecycle.ShouldSuppressCompletedAgentEvent( + threadId, evt, CanReconcileCompletedAgentEventLocked(threadId, evt))) + { + Logger.Debug($"[ChatProvider] Dropping completed-run agent event for threadId='{threadId}' stream='{evt.Stream}'"); + return new(false, false, null, openedLifecycle); + } if (ShouldDropTerminalAgentEventLocked( evt, threadId, @@ -2038,6 +2084,18 @@ private ChatRunTransition UpdateRunTrackingLocked( snapshot); } + private bool CanReconcileCompletedAgentEventLocked(string threadId, AgentEventInfo evt) + { + if (!_timelines.TryGetValue(threadId, out var timeline)) + return false; + if (ChatEventMapper.CanReconcileToolAfterRunEnd(evt, timeline)) + return true; + var approval = ChatEventMapper.MapTerminalApproval(evt); + return approval is not null && + timeline.PendingPermission is { } pending && + _approval.Matches(pending.RequestId, approval.ApprovalSlug, approval.ApprovalId); + } + private bool TryResolveTerminalApprovalLocked( AgentEventInfo evt, string threadId) @@ -2156,7 +2214,8 @@ private bool ShouldDropTerminalAgentEventLocked( _queue.RunIdsForThread(threadId), _timelines.TryGetValue(threadId, out var timeline) && timeline.TurnActive, - out droppedReason); + out droppedReason, + retryableError: ChatEventMapper.IsRetryableLifecycleError(evt)); } private static bool TryGetTerminalAgentRunId( diff --git a/src/OpenClaw.Tray.WinUI/Chat/ChatEventMapper.cs b/src/OpenClaw.Tray.WinUI/Chat/ChatEventMapper.cs index 9cd58266d..bffec55af 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/ChatEventMapper.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/ChatEventMapper.cs @@ -45,6 +45,49 @@ internal static bool IsLifecycleStart(AgentEventInfo evt) => evt.Data.TryGetProperty("phase", out var phase) && string.Equals(phase.GetString(), "start", StringComparison.OrdinalIgnoreCase); + internal static bool IsLifecycleError(AgentEventInfo evt) => + string.Equals(evt.Stream, "lifecycle", StringComparison.OrdinalIgnoreCase) && + evt.Data.ValueKind == JsonValueKind.Object && + evt.Data.TryGetProperty("phase", out var phase) && + string.Equals(phase.GetString(), "error", StringComparison.OrdinalIgnoreCase); + + internal static bool IsRetryableLifecycleError(AgentEventInfo evt) + { + if (!IsLifecycleError(evt)) + return false; + var data = evt.Data; + if (IsTrue(data, "executionSettled") || IsTrue(data, "fallbackExhaustedFailure")) + return false; + + // Match upstream isDefinitiveRunLifecycle's failed/non-failed projection. + // ProviderStarted only distinguishes timeout kinds, not retry eligibility. + var status = StringProperty(data, "status").ToLowerInvariant(); + var stopReason = StringProperty(data, "stopReason"); + if (string.IsNullOrWhiteSpace(stopReason)) + stopReason = string.Empty; + var timeoutPhase = StringProperty(data, "timeoutPhase").Trim(); + var timedOut = stopReason == "timeout" || + status is "timeout" or "timed_out" || + timeoutPhase is "queue" or "preflight" or "provider" or "post_turn" or "gateway_draining"; + if (timedOut) + return false; + + var aborted = IsTrue(data, "aborted") || status == "aborted"; + var cancellationStatus = status is "cancelled" or "canceled" or "aborted" or "superseded"; + if ((aborted || cancellationStatus) && + stopReason is not ("aborted" or "restart" or "rpc" or "stop" or "superseded") && + (stopReason.Length == 0 || cancellationStatus)) + { + stopReason = aborted ? "aborted" : "stop"; + } + var liveness = StringProperty(data, "livenessState").Trim().ToLowerInvariant(); + return stopReason is not ("aborted" or "restart" or "rpc" or "stop" or "superseded") && + liveness is not ("blocked" or "abandoned"); + } + + private static bool IsTrue(JsonElement data, string property) => + data.TryGetProperty(property, out var value) && value.ValueKind == JsonValueKind.True; + internal static bool IsTerminalRunEvent(AgentEventInfo evt) { if (evt.Data.ValueKind != JsonValueKind.Object) @@ -90,6 +133,11 @@ ChatStatusEvent or ChatErrorEvent or ChatReasoningEndEvent or }; } + internal static bool CanReconcileToolAfterRunEnd( + AgentEventInfo evt, + ChatTimelineState timeline) => + ChatTimelineReducer.CanReconcileToolAfterTurnEnd(timeline, Map(evt).Event); + internal static bool IsTerminalApprovalPhase(string phase) { if (string.IsNullOrEmpty(phase)) diff --git a/src/OpenClaw.Tray.WinUI/Chat/ChatLifecycleState.cs b/src/OpenClaw.Tray.WinUI/Chat/ChatLifecycleState.cs index 6e0aadbba..8f43ca553 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/ChatLifecycleState.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/ChatLifecycleState.cs @@ -15,9 +15,15 @@ internal sealed class ChatLifecycleState private readonly Dictionary _pendingAbortCounts = new(); private readonly HashSet _abortedRunIds = new(); private readonly HashSet _abortedThreads = new(); - private readonly Dictionary> _terminalRunIdsByThread = + private readonly Dictionary> _terminalRunIdsByThread = new(); + private sealed record TerminalRun( + string RunId, + bool SuppressOutput, + bool AssistantFinalReceived, + bool RetryableError); + private long _lifecycleStartSequence; internal bool IsResponseSuppressed => _abortedThreads.Count > 0; @@ -81,9 +87,72 @@ internal bool ShouldSuppress(string threadId, string? runId) => internal bool IsThreadSuppressed(string threadId) => _abortedThreads.Contains(threadId); + internal bool ShouldSuppressChatMessage( + string threadId, + string? runId, + bool isFinal, + IReadOnlyCollection queuedRunIds, + bool turnActive) + { + if (ShouldSuppress(threadId, runId)) + return true; + // Older gateways omit runId. Keep their existing thread-level behavior. + if (string.IsNullOrWhiteSpace(runId)) + return false; + if (_activeRunIds.TryGetValue(threadId, out var activeRunId) && + !string.Equals(activeRunId, runId, StringComparison.Ordinal)) + { + return true; + } + if (_terminalRunIdsByThread.TryGetValue(threadId, out var terminalRuns) && + terminalRuns.Find(run => run.RunId == runId) is { } terminal) + { + // A lifecycle end can precede the final text. Only the most recent, + // non-aborted run may finish that text, before another turn starts. + return !isFinal || terminal.SuppressOutput || terminal.AssistantFinalReceived || + turnActive || terminal != terminalRuns.FindLast(run => !run.SuppressOutput); + } + return !_activeRunIds.ContainsKey(threadId) && + turnActive && + queuedRunIds.Count > 0 && + !queuedRunIds.Contains(runId, StringComparer.Ordinal); + } + internal bool IsRunAborted(string? runId) => !string.IsNullOrWhiteSpace(runId) && _abortedRunIds.Contains(runId); + internal bool IsCompletedRun(string threadId, string? runId) => + !string.IsNullOrWhiteSpace(runId) && + _terminalRunIdsByThread.TryGetValue(threadId, out var runs) && + runs.Exists(run => run.RunId == runId); + + internal bool TryRestartFailedRun(string threadId, string? runId) + { + if (string.IsNullOrWhiteSpace(runId) || HasActiveRun(threadId) || + !_terminalRunIdsByThread.TryGetValue(threadId, out var runs) || + runs.FindLast(run => !run.SuppressOutput) is not { RetryableError: true, AssistantFinalReceived: false } failed || + !string.Equals(failed.RunId, runId, StringComparison.Ordinal) || + IsRunAborted(runId)) + { + return false; + } + return runs.Remove(failed); + } + + internal bool ShouldSuppressCompletedAgentEvent( + string threadId, + AgentEventInfo evt, + bool canReconcile = false) + { + if (string.IsNullOrWhiteSpace(evt.RunId) || + !_terminalRunIdsByThread.TryGetValue(threadId, out var terminalRuns) || + terminalRuns.Find(run => run.RunId == evt.RunId) is not { } terminal) + { + return false; + } + return terminal.SuppressOutput || !canReconcile; + } + internal bool HasPendingAbort(string threadId) => _pendingAbortCounts.ContainsKey(threadId); @@ -96,6 +165,7 @@ internal int TakePendingAbortCount(string threadId) internal long StartRun(string threadId, string runId) { + ClearRetryEligibility(threadId); _activeRunIds[threadId] = runId; var sequence = ++_lifecycleStartSequence; _activeRunStartSequences[threadId] = sequence; @@ -117,12 +187,14 @@ internal void RemoveAbortedRun(string? runId) internal void ClearThreadSuppression(string threadId) => _abortedThreads.Remove(threadId); - internal string? CompleteAssistantFinal(string threadId) + internal string? CompleteAssistantFinal(string threadId, string? runId = null) { - _activeRunIds.Remove(threadId, out var completedRunId); + ClearRetryEligibility(threadId); + _activeRunIds.Remove(threadId, out var activeRunId); + var completedRunId = string.IsNullOrWhiteSpace(runId) ? activeRunId : runId; if (!string.IsNullOrEmpty(completedRunId)) { - RememberTerminalRun(threadId, completedRunId); + RememberTerminalRun(threadId, completedRunId, assistantFinalReceived: true); _abortedRunIds.Remove(completedRunId); } _activeRunStartSequences.Remove(threadId); @@ -136,12 +208,27 @@ internal void RemoveActiveRun(string threadId) _activeRunStartSequences.Remove(threadId); } + internal void RememberSuppressedTerminal(string threadId, string runId) => + RememberTerminalRun(threadId, runId, suppressOutput: true); + + internal void ClearRetryEligibility(string threadId) + { + if (!_terminalRunIdsByThread.TryGetValue(threadId, out var runs)) + return; + for (var i = 0; i < runs.Count; i++) + { + if (runs[i].RetryableError) + runs[i] = runs[i] with { RetryableError = false }; + } + } + internal bool ShouldDropTerminal( string threadId, string runId, IReadOnlyCollection queuedRunIds, bool turnActive, - out ChatTerminalEventDropReason? droppedReason) + out ChatTerminalEventDropReason? droppedReason, + bool retryableError = false) { droppedReason = null; if (string.IsNullOrWhiteSpace(runId)) @@ -150,7 +237,7 @@ internal bool ShouldDropTerminal( return true; } if (_terminalRunIdsByThread.TryGetValue(threadId, out var terminalRunIds) && - terminalRunIds.Contains(runId, StringComparer.Ordinal)) + terminalRunIds.Exists(run => run.RunId == runId)) { return true; } @@ -168,7 +255,7 @@ internal bool ShouldDropTerminal( droppedReason = ChatTerminalEventDropReason.MismatchedRunId; return true; } - RememberTerminalRun(threadId, runId); + RememberTerminalRun(threadId, runId, retryableError: retryableError); return false; } @@ -201,16 +288,29 @@ internal void ClearActiveRuns(IEnumerable threadIds) RemoveActiveRun(threadId); } - private void RememberTerminalRun(string threadId, string runId) + private void RememberTerminalRun( + string threadId, + string runId, + bool assistantFinalReceived = false, + bool suppressOutput = false, + bool retryableError = false) { if (!_terminalRunIdsByThread.TryGetValue(threadId, out var runIds)) { runIds = []; _terminalRunIdsByThread[threadId] = runIds; } - runIds.Remove(runId); - runIds.Add(runId); - if (runIds.Count > TerminalRunCapacity) - runIds.RemoveRange(0, runIds.Count - TerminalRunCapacity); + var previous = runIds.Find(run => run.RunId == runId); + runIds.RemoveAll(run => run.RunId == runId); + if (runIds.Count >= TerminalRunCapacity) + runIds.RemoveAt(0); + var terminal = new TerminalRun( + runId, + suppressOutput || _abortedRunIds.Contains(runId) || previous?.SuppressOutput == true, + assistantFinalReceived || previous?.AssistantFinalReceived == true, + retryableError); + // Retention follows arrival order. Trailing-final eligibility separately + // ignores suppressed terminals, including ones learned from abort state. + runIds.Add(terminal); } } diff --git a/src/OpenClaw.Tray.WinUI/Chat/ChatResetState.cs b/src/OpenClaw.Tray.WinUI/Chat/ChatResetState.cs index 6055ef0c9..cc893b4fc 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/ChatResetState.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/ChatResetState.cs @@ -227,8 +227,16 @@ internal ChatResetMessageGate EvaluateChatMessage( string rawText, long timestampMs, bool hasPendingLocalEcho, - string? activeRunId = null) + string? activeRunId = null, + string? runId = null) { + if (!string.IsNullOrWhiteSpace(runId) && + _ignoredRunIds.TryGetValue(threadId, out var ignoredRuns) && + ignoredRuns.Contains(runId)) + { + return new(true, null, false, null); + } + var isNormalUserText = role == "user" && !ChatContentFormatting.LooksLikeApprovalSlashCommand(rawText) && !NativeToolProjector.LooksLikeSystemControlNote(rawText); @@ -290,7 +298,7 @@ internal ChatResetMessageGate EvaluateChatMessage( ? !IsPreResetTimestamp(threadId, timestampMs) : IsTimestampAcceptedForRun( threadId, - activeRunId, + string.IsNullOrWhiteSpace(runId) ? activeRunId : runId, timestampMs); return new( !timestampAccepted, diff --git a/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs b/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs index aff1e0b63..8d899a50d 100644 --- a/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs +++ b/src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs @@ -1240,9 +1240,13 @@ private void OnModelsListUpdated(object? sender, ModelsListInfo info) private void OnChatMessageReceived(object? sender, ChatMessageInfo message) { if (message is null || _state.IsDisposed) + { + message?.SuppressNotification(); return; + } if (string.IsNullOrEmpty(message.SessionKey)) { + message.SuppressNotification(); Logger.Warn($"[ChatProvider] Dropping chat message with empty sessionKey (role={message.Role})"); RaiseKeylessEventDiagnosticOnce(); return; @@ -1265,11 +1269,13 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) gate.RuntimeGeneration); if (gate.Suppressed) { - Logger.Debug($"[ABORT] Suppressed ChatMessage for threadId='{threadId}' (role={message.Role})"); + message.SuppressNotification(); + Logger.Debug($"[ChatProvider] Suppressed aborted or stale chat message for threadId='{threadId}' (role={message.Role})"); return; } if (gate.Drop) { + message.SuppressNotification(); if (gate.Snapshot is not null) Publish(gate.Snapshot); if (gate.RequestRemoteBackfill) @@ -1382,6 +1388,7 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) ChatMessageInfo.IsSilentAssistantDirective(role, message.Text) || (string.IsNullOrEmpty(message.Text) && assistantContent is null)) { + message.SuppressNotification(); return; } @@ -1395,9 +1402,13 @@ private void OnChatMessageReceived(object? sender, ChatMessageInfo message) if (preparation.PromotionSnapshot is not null) Publish(preparation.PromotionSnapshot); if (preparation.Disposition != AssistantQueueFrameDisposition.Render) + { + message.SuppressNotification(); return; + } if (!message.IsFinal && _state.IsLateNonFinalAssistantFrame(threadId)) { + message.SuppressNotification(); Logger.Warn($"[ChatProvider] Dropping late non-final assistant frame after completed turn for threadId='{threadId}' len={traceText.Length}"); return; } @@ -1429,7 +1440,7 @@ message.ResponseTokens is not null || if (!message.IsFinal) return; - var completedRunId = _state.CompleteAssistantFinal(threadId); + var completedRunId = _state.CompleteAssistantFinal(threadId, message.RunId); var completion = completedRunId is null ? null : _telemetry.PrepareFinishByRunId( diff --git a/tests/OpenClaw.Shared.Tests/OpenClawGatewayClientTests.cs b/tests/OpenClaw.Shared.Tests/OpenClawGatewayClientTests.cs index 94fb69e94..e4ee31efe 100644 --- a/tests/OpenClaw.Shared.Tests/OpenClawGatewayClientTests.cs +++ b/tests/OpenClaw.Shared.Tests/OpenClawGatewayClientTests.cs @@ -1646,6 +1646,130 @@ public void ProcessRawMessage_SessionMessageWithStringContent_EmitsChatMessage() Assert.Equal(1781631273567, received.Ts); } + [Theory] + [InlineData("chat", false)] + [InlineData("chat", true)] + [InlineData("session.message", false)] + [InlineData("session.message", true)] + public void ProcessRawMessage_ChatRunId_PreservesPayloadIdentity(string eventName, bool legacy) + { + var helper = new GatewayClientTestHelper(); + using var client = helper.Client; + ChatMessageInfo? received = null; + client.ChatMessageReceived += (_, message) => received = message; + var content = legacy + ? """ + "role":"assistant","text":"answer","__openclaw":{"id":"message-id"} + """ + : """ + "message":{"role":"assistant","content":[{"type":"text","text":"answer"}],"runId":"not-the-envelope-run","__openclaw":{"id":"message-id"}} + """; + + helper.ProcessRawMessage($$$""" + {"type":"event","event":"{{{eventName}}}","payload":{ + "sessionKey":"main","runId":"wire-run","state":"final",{{{content}}} + }} + """); + + Assert.NotNull(received); + Assert.Equal("wire-run", received.RunId); + Assert.Equal("message-id", received.OpenClawId); + Assert.True(received.IsFinal); + } + + [Theory] + [InlineData("")] + [InlineData("\"runId\":null,")] + [InlineData("\"runId\":123,")] + public void ProcessRawMessage_ChatRunId_DoesNotGuessMissingIdentity(string runProperty) + { + var helper = new GatewayClientTestHelper(); + using var client = helper.Client; + ChatMessageInfo? received = null; + client.ChatMessageReceived += (_, message) => received = message; + + helper.ProcessRawMessage($$$""" + {"type":"event","event":"chat","payload":{ + {{{runProperty}}}"sessionKey":"main","state":"final", + "message":{"role":"assistant","content":"answer","runId":"not-an-envelope-run","__openclaw":{"id":"message-id"}} + }} + """); + + Assert.NotNull(received); + Assert.Null(received.RunId); + Assert.Equal("message-id", received.OpenClawId); + } + + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public void ProcessRawMessage_ChatNotificationHonorsSynchronousSuppression(bool legacy, bool suppress) + { + var helper = new GatewayClientTestHelper(); + using var client = helper.Client; + var delivered = new List(); + var notifications = new List(); + client.ChatMessageReceived += (_, message) => + { + if (suppress) + message.SuppressNotification(); + }; + client.ChatMessageReceived += (_, message) => delivered.Add(message); + client.NotificationReceived += (_, notification) => notifications.Add(notification); + var content = legacy + ? """ + "role":"assistant","text":"final answer" + """ + : """ + "message":{"role":"assistant","content":"final answer"} + """; + helper.ProcessRawMessage($$$""" + {"type":"event","event":"chat","payload":{"sessionKey":"main","runId":"run","state":"final",{{{content}}}}} + """); + Assert.Equal(suppress, Assert.Single(delivered).IsNotificationSuppressed); + Assert.Equal(suppress ? 0 : 1, notifications.Count); + Assert.DoesNotContain("IsNotificationSuppressed", JsonSerializer.Serialize(delivered[0]), StringComparison.Ordinal); + } + + [Fact] + public void ProcessRawMessage_FailedChatConsumerDoesNotEmitSuccessNotification() + { + var logger = new TestLogger(); + var helper = new GatewayClientTestHelper(logger); + using var client = helper.Client; + var notifications = new List(); + client.ChatMessageReceived += (_, _) => throw new InvalidOperationException("consumer failed"); + client.NotificationReceived += (_, notification) => notifications.Add(notification); + helper.ProcessRawMessage(""" + {"type":"event","event":"chat","payload":{"sessionKey":"main","runId":"run","state":"final","role":"assistant","text":"answer"}} + """); + Assert.Empty(notifications); + Assert.Contains(logger.Logs, log => log.Contains("ChatMessageReceived handler threw", StringComparison.Ordinal)); + } + + [Fact] + public void ProcessRawMessage_AgentLifecyclePreservesDefinitiveFacts() + { + var helper = new GatewayClientTestHelper(); + using var client = helper.Client; + AgentEventInfo? received = null; + client.AgentEventReceived += (_, evt) => received = evt; + helper.ProcessRawMessage(""" + {"type":"event","event":"agent","payload":{"sessionKey":"main","runId":"run","stream":"lifecycle", + "data":{"phase":"error","executionSettled":true,"fallbackExhaustedFailure":true, + "status":"timeout","stopReason":"timeout","timeoutPhase":"provider","livenessState":"blocked"}}} + """); + Assert.NotNull(received); + Assert.True(received.Data.GetProperty("executionSettled").GetBoolean()); + Assert.True(received.Data.GetProperty("fallbackExhaustedFailure").GetBoolean()); + Assert.Equal("timeout", received.Data.GetProperty("status").GetString()); + Assert.Equal("timeout", received.Data.GetProperty("stopReason").GetString()); + Assert.Equal("provider", received.Data.GetProperty("timeoutPhase").GetString()); + Assert.Equal("blocked", received.Data.GetProperty("livenessState").GetString()); + } + [Fact] public void ProcessRawMessage_SessionMessageWithOpenClawMetadata_EmitsMessageIdentity() { diff --git a/tests/OpenClaw.Tray.Tests/ChatRunCorrelationTests.cs b/tests/OpenClaw.Tray.Tests/ChatRunCorrelationTests.cs new file mode 100644 index 000000000..62524d2ec --- /dev/null +++ b/tests/OpenClaw.Tray.Tests/ChatRunCorrelationTests.cs @@ -0,0 +1,325 @@ +using System.Text.Json; +using OpenClaw.Shared; +using OpenClawTray.Chat; + +namespace OpenClaw.Tray.Tests; + +public sealed class ChatRunCorrelationTests +{ + [Theory] + [InlineData(false)] + [InlineData(true)] + public void Lifecycle_AbortedTerminalPreservesNewerTrailingFinal(bool legacyJob) + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context); + state.BeginAbort("main"); + state.CompleteAbort("main", "old"); + state.ProcessAgentEvent(Lifecycle("new", "start"), "main", context); + state.ProcessAgentEvent(Lifecycle("new", "end"), "main", context); + var oldTerminal = legacyJob + ? new AgentEventInfo + { + SessionKey = "main", RunId = "old", Stream = "job", + Data = JsonSerializer.SerializeToElement(new { state = "done" }), + } + : Lifecycle("old", "end"); + state.ProcessAgentEvent(oldTerminal, "main", context); + var current = state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", RunId = "new", Role = "assistant", + Text = "Complete newer answer", State = "final", + }, context); + Assert.False(current.Suppressed || current.Drop, + "A cancelled terminal displaced the newer run's trailing final."); + } + + [Fact] + public void Lifecycle_TerminalCacheEvictsOldestNotMostRecentSuppressedRun() + { + var state = new ChatLifecycleState(); + for (var i = 0; i < 63; i++) + state.CompleteAssistantFinal("main", $"normal-{i}"); + state.RememberSuppressedTerminal("main", "reset-recent"); + state.RememberSuppressedTerminal("main", "reset-next"); + + Assert.True(state.ShouldSuppressChatMessage("main", "reset-recent", true, [], false), + "The next terminal evicted the most recent suppressed run."); + Assert.True(state.ShouldSuppressChatMessage("main", "reset-next", true, [], false)); + Assert.False(state.ShouldSuppressChatMessage("main", "normal-0", true, [], false)); + Assert.True(state.ShouldSuppressChatMessage("main", "normal-1", true, [], false)); + Assert.False(state.ShouldSuppressChatMessage("other", "reset-recent", true, [], false)); + } + + [Fact] + public void Lifecycle_AbortedStartReplayCannotReplaceCurrentRun() + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context); + state.BeginAbort("main"); + state.CompleteAbort("main", "old"); + state.ProcessAgentEvent(Lifecycle("new", "start"), "main", context); + state.ProcessAgentEvent(Lifecycle("old", "end"), "main", context); + Assert.False(state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context).Process); + var current = state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", RunId = "new", Role = "assistant", + Text = "Current answer", State = "final", + }, context); + Assert.False(current.Suppressed || current.Drop); + Assert.True(state.ProcessAgentEvent(Lifecycle("new", "end"), "main", context).Process); + } + + [Theory] + [InlineData("""{}""", true)] + [InlineData("""{"executionSettled":false,"fallbackExhaustedFailure":false}""", true)] + [InlineData("""{"executionSettled":"true","fallbackExhaustedFailure":1}""", true)] + [InlineData("""{"executionSettled":true}""", false)] + [InlineData("""{"fallbackExhaustedFailure":true}""", false)] + [InlineData("""{"stopReason":"timeout"}""", false)] + [InlineData("""{"timeoutPhase":"queue"}""", false)] + [InlineData("""{"timeoutPhase":"provider"}""", false)] + [InlineData("""{"timeoutPhase":" preflight "}""", false)] + [InlineData("""{"timeoutPhase":"post_turn"}""", false)] + [InlineData("""{"timeoutPhase":"gateway_draining"}""", false)] + [InlineData("""{"timeoutPhase":"unknown"}""", true)] + [InlineData("""{"status":" TIMEOUT "}""", true)] + [InlineData("""{"stopReason":" timeout "}""", true)] + [InlineData("""{"stopReason":" stop "}""", true)] + [InlineData("""{"aborted":true,"stopReason":" "}""", false)] + [InlineData("""{"status":" CANCELLED "}""", true)] + [InlineData("""{"status":"timed_out"}""", false)] + [InlineData("""{"status":"cancelled"}""", false)] + [InlineData("""{"status":"canceled","stopReason":"custom"}""", false)] + [InlineData("""{"status":"aborted"}""", false)] + [InlineData("""{"status":"superseded"}""", false)] + [InlineData("""{"status":"failed"}""", true)] + [InlineData("""{"aborted":true}""", false)] + [InlineData("""{"aborted":true,"stopReason":"custom"}""", true)] + [InlineData("""{"aborted":"true"}""", true)] + [InlineData("""{"stopReason":"rpc"}""", false)] + [InlineData("""{"stopReason":"stop"}""", false)] + [InlineData("""{"stopReason":"restart"}""", false)] + [InlineData("""{"stopReason":"aborted"}""", false)] + [InlineData("""{"stopReason":"superseded"}""", false)] + [InlineData("""{"stopReason":"error"}""", true)] + [InlineData("""{"stopReason":"TIMEOUT"}""", true)] + [InlineData("""{"livenessState":" Blocked "}""", false)] + [InlineData("""{"livenessState":"abandoned"}""", false)] + [InlineData("""{"providerStarted":true}""", true)] + public void Lifecycle_ErrorRetryMatchesDefinitiveUpstreamFacts(string json, bool retryable) + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("run", "start"), "main", context); + var error = Lifecycle("run", "error"); + var data = JsonSerializer.Deserialize>(json)!; + data["phase"] = "error"; + error.Data = JsonSerializer.SerializeToElement(data); + state.ProcessAgentEvent(error, "main", context); + Assert.Equal(retryable, state.ProcessAgentEvent(Lifecycle("run", "start"), "main", context).Process); + } + + [Fact] + public void Lifecycle_ErrorCanRetrySameRunWithoutWeakeningCompletedRunFence() + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("run", "start"), "main", context); + state.ProcessAgentEvent(Lifecycle("run", "error"), "main", context); + Assert.True(state.ProcessAgentEvent(Lifecycle("run", "start"), "main", context).Process, + "A retryable lifecycle error prevented the same run from restarting."); + Assert.False(state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", RunId = "run", Role = "assistant", Text = "Retry result", State = "final", + }, context).Suppressed); + state.CompleteAssistantFinal("main", "run"); + Assert.False(state.ProcessAgentEvent(Lifecycle("run", "start"), "main", context).Process); + } + + [Theory] + [InlineData("start")] + [InlineData("end")] + [InlineData("abort")] + public void Lifecycle_ErrorRetryCannotDisplaceNewerRun(string newerPhase) + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context); + state.ProcessAgentEvent(Lifecycle("old", "error"), "main", context); + state.ProcessAgentEvent(Lifecycle("new", "start"), "main", context); + if (newerPhase == "end") + state.ProcessAgentEvent(Lifecycle("new", "end"), "main", context); + if (newerPhase == "abort") + { + state.BeginAbort("main"); + state.CompleteAbort("main", "new"); + state.ProcessAgentEvent(Lifecycle("new", "end"), "main", context); + } + Assert.False(state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context).Process); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public void Reset_TerminalDoesNotForgetSuppressionBeforeLateFinal(bool terminalAfterNewRun) + { + var state = new ChatConversationState(ConnectionStatus.Connected, null, null); + var context = new ChatProjectionContext("main", HasHandshakeSnapshot: true); + state.Load([new SessionInfo { Key = "main", IsMain = true }], context); + state.ProcessAgentEvent(Lifecycle("old", "start"), "main", context); + state.ResetThread("main", context); + if (!terminalAfterNewRun) + Assert.True(state.ProcessAgentEvent(Lifecycle("old", "end"), "main", context).ReloadHistory); + + var timestamp = DateTimeOffset.UtcNow.AddSeconds(2).ToUnixTimeMilliseconds(); + state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", Role = "user", Text = "Fresh question", Ts = timestamp, + }, context); + Assert.True(state.ProcessAgentEvent(Lifecycle("new", "start", timestamp), "main", context).Process); + Assert.True(state.ProcessAgentEvent(Lifecycle("new", "end", timestamp + 1), "main", context).Process); + if (terminalAfterNewRun) + Assert.True(state.ProcessAgentEvent(Lifecycle("old", "end"), "main", context).ReloadHistory); + + var stale = state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", RunId = "old", Role = "assistant", + Text = "Late pre-reset final", State = "final", Ts = timestamp + 2, + }, context); + Assert.True(stale.Suppressed || stale.Drop, + "A reset-invalidated run's terminal forgot suppression before its late final."); + var current = state.GateIncomingChatMessage(new ChatMessageInfo + { + SessionKey = "main", RunId = "new", Role = "assistant", + Text = "Current trailing final", State = "final", Ts = timestamp + 3, + }, context); + Assert.False(current.Suppressed || current.Drop); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public void Lifecycle_SuppressesCancelledChatAfterThreadSuppressionClears(bool final) + { + var state = new ChatLifecycleState(); + state.StartRun("main", "old"); + state.BeginAbort("main", hadActiveTurn: true); + state.CompleteAbort("main", "old"); + Assert.False(state.IsThreadSuppressed("main")); + Assert.True(state.ShouldSuppressChatMessage("main", "old", final, [], turnActive: false)); + state.StartRun("main", "new"); + Assert.True(state.ShouldSuppressChatMessage("main", "old", final, [], turnActive: true)); + Assert.False(state.ShouldSuppressChatMessage("main", "new", final, [], turnActive: true)); + } + + [Fact] + public void Lifecycle_FinalWithoutStartRemembersWireRunAndRejectsReplay() + { + var state = new ChatLifecycleState(); + Assert.Equal("wire-run", state.CompleteAssistantFinal("main", "wire-run")); + Assert.True(state.ShouldSuppressChatMessage("main", "wire-run", true, [], turnActive: false)); + Assert.True(state.ShouldDropTerminal("main", "wire-run", [], false, out _)); + Assert.False(state.ShouldSuppressChatMessage("main", "remote-next", true, [], turnActive: false)); + } + + [Theory] + [InlineData("assistant", """{"delta":"late"}""")] + [InlineData("reasoning", """{"delta":"late"}""")] + [InlineData("lifecycle", """{"phase":"start"}""")] + public void Lifecycle_CompletedRunCannotRestartThroughAgentOutput(string stream, string data) + { + var state = new ChatLifecycleState(); + state.CompleteAssistantFinal("main", "completed"); + var evt = new AgentEventInfo + { + SessionKey = "main", RunId = "completed", Stream = stream, + Data = JsonSerializer.Deserialize(data), + }; + Assert.True(state.ShouldSuppressCompletedAgentEvent("main", evt)); + Assert.False(state.ShouldSuppressCompletedAgentEvent("other", evt)); + evt.RunId = "remote-next"; + Assert.False(state.ShouldSuppressCompletedAgentEvent("main", evt)); + } + + [Theory] + [InlineData(false, "start")] + [InlineData(false, "result")] + [InlineData(true, "start")] + [InlineData(true, "result")] + public void Lifecycle_CompletedToolRepairIsAllowedOnlyForNonAbortedRun(bool aborted, string phase) + { + var state = new ChatLifecycleState(); + state.StartRun("main", "run"); + if (aborted) + state.BeginAbort("main", hadActiveTurn: true); + Assert.False(state.ShouldDropTerminal("main", "run", [], true, out _)); + state.RemoveAbortedRun("run"); + var evt = new AgentEventInfo + { + SessionKey = "main", RunId = "run", Stream = "tool", + Data = JsonSerializer.SerializeToElement(new { phase, name = "tool", toolCallId = "tool-1", result = "completed" }), + }; + Assert.Equal(aborted, state.ShouldSuppressCompletedAgentEvent("main", evt, canReconcile: true)); + Assert.True(state.ShouldSuppressCompletedAgentEvent("main", evt, canReconcile: false)); + } + + [Fact] + public void Lifecycle_OnlyMostRecentCompletedRunCanSupplyTrailingFinal() + { + var state = new ChatLifecycleState(); + Assert.False(state.ShouldDropTerminal("main", "old", [], false, out _)); + Assert.False(state.ShouldDropTerminal("main", "new", [], false, out _)); + Assert.True(state.ShouldSuppressChatMessage("main", "old", true, [], turnActive: false)); + Assert.False(state.ShouldSuppressChatMessage("main", "new", true, [], turnActive: false)); + Assert.True(state.ShouldSuppressChatMessage("main", "new", false, [], turnActive: false)); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData(" ")] + public void Lifecycle_MissingRunIdRetainsLegacyAdmission(string? runId) + { + var state = new ChatLifecycleState(); + state.StartRun("main", "current"); + Assert.False(state.ShouldSuppressChatMessage("main", runId, true, [], turnActive: true)); + state.BeginAbort("main", hadActiveTurn: true); + Assert.True(state.ShouldSuppressChatMessage("main", runId, true, [], turnActive: true)); + } + + [Fact] + public void Reset_ExplicitIgnoredRunCannotBorrowCurrentRunTimestampAdmission() + { + var state = new ChatResetState(); + state.AddIgnoredRun("main", "old"); + var ignored = state.EvaluateChatMessage( + "main", "assistant", "late output", timestampMs: long.MaxValue, + hasPendingLocalEcho: false, activeRunId: "new", runId: "old"); + var current = state.EvaluateChatMessage( + "main", "assistant", "current output", timestampMs: long.MaxValue, + hasPendingLocalEcho: false, activeRunId: "new", runId: "new"); + var otherThread = state.EvaluateChatMessage( + "other", "assistant", "other output", timestampMs: long.MaxValue, + hasPendingLocalEcho: false, runId: "old"); + Assert.True(ignored.Drop); + Assert.False(current.Drop); + Assert.False(otherThread.Drop); + } + + private static AgentEventInfo Lifecycle(string runId, string phase, long timestamp = 0) => new() + { + SessionKey = "main", + RunId = runId, + Stream = "lifecycle", + Ts = timestamp, + Data = JsonSerializer.SerializeToElement(new { phase }), + }; +} diff --git a/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs b/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs index 4d0c0f331..4e6b5b811 100644 --- a/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs +++ b/tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs @@ -2547,6 +2547,80 @@ public async Task SendMessageAsync_RejectsEmptyMessage() await Assert.ThrowsAsync(() => provider.SendMessageAsync("main", " ")); } + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ChatMessageReceived_CancelledRunLateFinal_CannotEndANewerActiveRun( + bool terminalCleanupBeforeNewRun) + { + var (bridge, provider, snapshots, notifications) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + bridge.IsConnected = true; + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "old-run", Status = "started" }); + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "new-run", Status = "started" }); + await provider.LoadAsync(); + await provider.SendMessageAsync("main", "First question"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "old-run")); + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"Old partial"}""", runId: "old-run")); + await provider.StopResponseAsync("main"); + Assert.Contains("old-run", bridge.AbortedRunIds); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + + if (terminalCleanupBeforeNewRun) + { + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"STALE idle delta"}""", runId: "old-run")); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "old-run", Role = "assistant", + Text = "STALE idle final", State = "final", + }); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "old-run")); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + Assert.DoesNotContain(snapshots[^1].Timelines["main"].Entries, e => e.Text.Contains("STALE", StringComparison.Ordinal)); + } + + await provider.SendMessageAsync("main", "Second question"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "new-run")); + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"New partial"}""", runId: "new-run")); + Assert.True(snapshots[^1].Timelines["main"].TurnActive); + notifications.Clear(); + + var beforeLateEvents = snapshots[^1].Timelines["main"]; + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"STALE delta"}""", runId: "old-run")); + Assert.True(snapshots[^1].Timelines["main"].TurnActive); + Assert.DoesNotContain(snapshots[^1].Timelines["main"].Entries, + e => e.Text.Contains("STALE", StringComparison.Ordinal)); + Assert.Equal(beforeLateEvents, snapshots[^1].Timelines["main"]); + var beforeLateFinal = snapshots[^1].Timelines["main"]; + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "old-run", Role = "assistant", + Text = "STALE final", State = "final", + }); + + Assert.True(snapshots[^1].Timelines["main"].TurnActive, + "A cancelled run's late final ended the new turn."); + Assert.Equal(beforeLateFinal, snapshots[^1].Timelines["main"]); + Assert.DoesNotContain(notifications, n => n.Kind == ChatProviderNotificationKind.TurnComplete); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "old-run")); + Assert.Equal(beforeLateFinal, snapshots[^1].Timelines["main"]); + + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":" Still active after old events."}""", runId: "new-run")); + Assert.True(snapshots[^1].Timelines["main"].TurnActive); + Assert.Contains(snapshots[^1].Timelines["main"].Entries, + e => e.Text == "New partial Still active after old events."); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "new-run", Role = "assistant", + Text = "New final", State = "final", + }); + var completed = snapshots[^1].Timelines["main"]; + Assert.False(completed.TurnActive); + Assert.Contains(completed.Entries, e => e.Kind == ChatTimelineItemKind.Assistant && e.Text == "New final"); + Assert.DoesNotContain(completed.Entries, e => e.Text.Contains("STALE", StringComparison.Ordinal)); + Assert.Single(notifications, n => n.Kind == ChatProviderNotificationKind.TurnComplete); + } + [Fact] public async Task ChatMessageReceived_FinalAssistant_AppendsAssistantEntry() { @@ -2570,6 +2644,128 @@ public async Task ChatMessageReceived_FinalAssistant_AppendsAssistantEntry() Assert.Contains(notifications, n => n.Kind == ChatProviderNotificationKind.TurnComplete); } + [Theory] + [InlineData("final", false)] + [InlineData("final", true)] + [InlineData("lifecycle", false)] + [InlineData("lifecycle", true)] + public async Task ChatMessageReceived_CompletedRun_CannotReplaceNewTurn( + string completion, + bool newLifecycleStarted) + { + var (bridge, provider, snapshots, notifications) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + bridge.IsConnected = true; + bridge.SendResults.Enqueue(new ChatSendResult { RunId = "new-run", Status = "started" }); + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "old-run")); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "old-run", Role = "assistant", + Text = "Old answer", State = completion == "final" ? "final" : "delta", + }); + if (completion == "lifecycle") + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "old-run")); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + + await provider.SendMessageAsync("main", "Next question"); + if (newLifecycleStarted) + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "new-run")); + var beforeStale = snapshots[^1].Timelines["main"]; + Assert.True(beforeStale.TurnActive); + notifications.Clear(); + foreach (var staleRun in new[] { "old-run", "unrelated-run" }) + { + foreach (var frameState in new[] { "delta", "final" }) + { + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = staleRun, Role = "assistant", + Text = "STALE output", State = frameState, + }); + Assert.Equal(beforeStale, snapshots[^1].Timelines["main"]); + } + } + Assert.Empty(notifications); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "new-run", Role = "assistant", + Text = "Current final", State = "final", + }); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + Assert.Contains(snapshots[^1].Timelines["main"].Entries, e => e.Text == "Current final"); + Assert.Single(notifications, n => n.Kind == ChatProviderNotificationKind.TurnComplete); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ChatMessageReceived_FinalAfterLifecycleEnd_OnlyCompletesNonAbortedRun(bool aborted) + { + var (bridge, provider, snapshots, notifications) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run")); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", RunId = "run", Role = "assistant", + Text = "Partial", State = "delta", + }); + if (aborted) + await provider.StopResponseAsync("main"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run")); + notifications.Clear(); + var beforeFinal = snapshots[^1].Timelines["main"]; + var final = new ChatMessageInfo + { + SessionKey = "main", RunId = "run", Role = "assistant", + Text = "Complete answer", State = "final", + }; + bridge.RaiseChat(final); + var afterFinal = snapshots[^1].Timelines["main"]; + Assert.False(afterFinal.TurnActive); + if (aborted) + { + Assert.Equal(beforeFinal, afterFinal); + Assert.Empty(notifications); + } + else + { + Assert.Contains(afterFinal.Entries, e => e.Text == "Complete answer"); + Assert.Single(notifications, n => n.Kind == ChatProviderNotificationKind.TurnComplete); + } + notifications.Clear(); + bridge.RaiseChat(final); + Assert.Equal(afterFinal, snapshots[^1].Timelines["main"]); + Assert.Empty(notifications); + } + + [Fact] + public async Task ChatMessageReceived_RunCorrelation_PreservesOtherSessionsAndRemoteFinal() + { + var (bridge, provider, snapshots, _) = CreateProvider(new[] + { + MainSession(), new SessionInfo { Key = "other", DisplayName = "Other" }, + }); + await using var lifetime = provider; + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "main-run")); + var main = snapshots[^1].Timelines["main"]; + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "other", RunId = "remote-run", Role = "assistant", + Text = "Remote final without lifecycle", State = "final", + }); + Assert.Equal(main, snapshots[^1].Timelines["main"]); + Assert.False(snapshots[^1].Timelines["other"].TurnActive); + Assert.Contains(snapshots[^1].Timelines["other"].Entries, e => e.Text == "Remote final without lifecycle"); + bridge.RaiseChat(new ChatMessageInfo + { + SessionKey = "main", Role = "assistant", Text = "Legacy final", State = "final", + }); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + } + [Fact] public async Task ChatMessageReceived_AssistantNoReply_IsSuppressed() { @@ -9744,6 +9940,97 @@ public async Task AgentEvent_ChildBeforeTurnEndAndLateParentMaterializesTerminal Assert.Equal(ChatToolCallStatus.Success, entry.ToolResult); } + [Theory] + [InlineData("tool", """{"phase":"start","name":"stale","toolCallId":"unknown"}""")] + [InlineData("tool", """{"phase":"result","toolCallId":"unknown","result":"stale"}""")] + [InlineData("tool", """{"phase":"error","toolCallId":"unknown","message":"stale"}""")] + [InlineData("item", """{"phase":"start","kind":"command","name":"system.run","toolCallId":"unknown","itemId":"command:unknown"}""")] + public async Task AgentEvent_CompletedUnknownToolCannotMutateTimeline(string stream, string data) + { + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "old")); + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"Completed answer"}""", runId: "old")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "old")); + var idle = snapshots[^1].Timelines["main"]; + bridge.RaiseAgent(MakeAgentEvent(stream, data, runId: "old")); + Assert.Equal(idle, snapshots[^1].Timelines["main"]); + + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "new")); + bridge.RaiseAgent(MakeAgentEvent("tool", + """{"phase":"start","name":"current","toolCallId":"unknown"}""", runId: "new")); + var current = snapshots[^1].Timelines["main"]; + bridge.RaiseAgent(MakeAgentEvent(stream, data, runId: "old")); + Assert.Equal(current, snapshots[^1].Timelines["main"]); + } + + [Theory] + [InlineData("chat", false)] + [InlineData("chat", true)] + [InlineData("session.message", false)] + [InlineData("session.message", true)] + public async Task ChatMessageReceived_RejectedWireFinalCannotEmitNotification(string eventName, bool legacy) + { + using var identity = new OpenClaw.TestSupport.TempDirectory(); + using var client = new OpenClawGatewayClient("ws://127.0.0.1:1", "test-token", identityPath: identity.Path); + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + var notifications = new List(); + client.ChatMessageReceived += (_, message) => bridge.RaiseChat(message); + client.NotificationReceived += (_, notification) => notifications.Add(notification); + var processMessage = typeof(OpenClawGatewayClient).GetMethod( + "ProcessMessage", System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic)!; + void Final(string runId, string text) + { + var payload = new Dictionary + { + ["sessionKey"] = "main", ["runId"] = runId, ["state"] = "final", + }; + if (legacy) + { + payload["role"] = "assistant"; + payload["text"] = text; + } + else + payload["message"] = new { role = "assistant", content = text }; + processMessage.Invoke(client, [JsonSerializer.Serialize(new { type = "event", @event = eventName, payload })]); + } + + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "old")); + await provider.StopResponseAsync("main"); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "new")); + bridge.RaiseAgent(MakeAgentEvent("assistant", """{"delta":"Current partial"}""", runId: "new")); + var current = snapshots[^1].Timelines["main"]; + Final("old", "Stale final"); + Assert.Equal(current, snapshots[^1].Timelines["main"]); + Assert.Empty(notifications); + Final("new", "Current complete answer"); + Assert.Equal("Current complete answer", Assert.Single(notifications).FullMessage); + Final("new", "Current complete answer"); + Assert.Single(notifications); + } + + [Fact] + public async Task AgentEvent_KnownCompletedChildUpdateDoesNotReactivateBeforeLateParent() + { + var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() }); + await using var lifetime = provider; + await provider.LoadAsync(); + bridge.RaiseAgent(MakeAgentEvent("item", + """{"phase":"start","kind":"command","name":"system.run","title":"Bash","itemId":"command:tool-1","toolCallId":"tool-1"}""", runId: "run")); + bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run")); + bridge.RaiseAgent(MakeAgentEvent("item", + """{"phase":"update","kind":"command","name":"system.run","title":"Bash","meta":"late display","itemId":"command:tool-1","toolCallId":"tool-1"}""", runId: "run")); + Assert.False(snapshots[^1].Timelines["main"].TurnActive); + bridge.RaiseAgent(MakeAgentEvent("item", + """{"phase":"start","kind":"tool","name":"system.run","itemId":"tool:tool-1","toolCallId":"tool-1"}""", runId: "run")); + var timeline = snapshots[^1].Timelines["main"]; + Assert.False(timeline.TurnActive); + Assert.Equal(ChatToolCallStatus.Interrupted, Assert.Single(timeline.Entries).ToolResult); + } + [Fact] public async Task AgentEvent_LateLegacyParentUpsertsResolvedCacheGeneration() {