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() {