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