Skip to content

Commit 8943b29

Browse files
calebedenCopilot
andcommitted
Harden queued message follow-ups
Clear stale queued-drain guards when queue state is emptied, clean up idle local-run tracking after cancel, and document ephemeral queue identifiers. Add targeted queue regressions for stale drain guards, remote backfill after cancel, and deferred in-flight cancel behavior. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
1 parent 1cdb7e0 commit 8943b29

3 files changed

Lines changed: 168 additions & 3 deletions

File tree

src/OpenClaw.Tray.WinUI/Chat/OpenClawChatDataProvider.cs

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -607,6 +607,7 @@ private async Task DispatchQueuedSendAsync(
607607
if (!sendStillCurrent)
608608
return;
609609

610+
Logger.Warn($"[Queue] chat.send failed threadId='{threadId}' queuedMessageId='{request.Id}' sendRunId='{request.SendRunId}': {ex.Message}");
610611
// Surface as an error in the timeline + notification, while the
611612
// failed queue card keeps the attempted text visible for retry/edit.
612613
Publish(failureSnapshot!);
@@ -1651,6 +1652,7 @@ public ValueTask DisposeAsync()
16511652
_localInlineApprovals.Clear();
16521653
_queuedMessages.Clear();
16531654
_queuedSendRequests.Clear();
1655+
_queuedDrainScheduledThreads.Clear();
16541656
_queuedMessageIdsByRunId.Clear();
16551657
_terminalRunIdsByThread.Clear();
16561658
_localSentTexts.Clear();
@@ -1752,6 +1754,7 @@ private void OnStatusChanged(object? sender, ConnectionStatus status)
17521754
_localSentTexts.Clear();
17531755
_queuedMessages.Clear();
17541756
_queuedSendRequests.Clear();
1757+
_queuedDrainScheduledThreads.Clear();
17551758
_assistantFallbackPromotedThreads.Clear();
17561759
_queuedMessageIdsByRunId.Clear();
17571760
_terminalRunIdsByThread.Clear();
@@ -3032,7 +3035,11 @@ private bool RemoveQueuedMessageLocked(string threadId, string messageId)
30323035
RemoveQueuedSendRequestLocked(threadId, messageId);
30333036
}
30343037
if (list.Count == 0)
3038+
{
30353039
_queuedMessages.Remove(threadId);
3040+
ClearQueuedDrainScheduleLocked(threadId);
3041+
ClearLocallyInitiatedIfIdleLocked(threadId);
3042+
}
30363043
return removed;
30373044
}
30383045

@@ -3053,10 +3060,17 @@ private bool CancelQueuedMessageLocked(string threadId, string messageId)
30533060
RemoveQueuedRunMappingByMessageIdLocked(threadId, messageId);
30543061
RemoveQueuedSendRequestLocked(threadId, messageId);
30553062
if (list.Count == 0)
3063+
{
30563064
_queuedMessages.Remove(threadId);
3065+
ClearQueuedDrainScheduleLocked(threadId);
3066+
ClearLocallyInitiatedIfIdleLocked(threadId);
3067+
}
30573068
return true;
30583069
}
30593070

3071+
private void ClearQueuedDrainScheduleLocked(string threadId)
3072+
=> _queuedDrainScheduledThreads.Remove(threadId);
3073+
30603074
private bool PromoteQueuedMessageLocked(
30613075
string threadId,
30623076
string messageId,
@@ -3088,10 +3102,25 @@ private bool PromoteQueuedMessageLocked(
30883102
_assistantFallbackPromotedThreads.Add(threadId);
30893103
RemoveQueuedSendRequestLocked(threadId, messageId);
30903104
if (list.Count == 0)
3105+
{
30913106
_queuedMessages.Remove(threadId);
3107+
ClearQueuedDrainScheduleLocked(threadId);
3108+
}
30923109
return true;
30933110
}
30943111

3112+
private void ClearLocallyInitiatedIfIdleLocked(string threadId)
3113+
{
3114+
if (_activeRunIds.ContainsKey(threadId))
3115+
return;
3116+
if (_timelines.TryGetValue(threadId, out var timeline) && timeline.TurnActive)
3117+
return;
3118+
if (HasPendingQueuedMessagesLocked(threadId))
3119+
return;
3120+
3121+
_locallyInitiatedThreads.Remove(threadId);
3122+
}
3123+
30953124
private bool ReconcileQueuedMessageEchoLocked(
30963125
string threadId,
30973126
string messageId,
@@ -4672,6 +4701,7 @@ private ResetClearPersistence ClearThreadHistoryAfterResetLocked(string threadId
46724701
_localSentTexts.Remove(threadId);
46734702
_queuedMessages.Remove(threadId);
46744703
_queuedSendRequests.Remove(threadId);
4704+
ClearQueuedDrainScheduleLocked(threadId);
46754705
_queuedMessageIdsByRunId.Remove(threadId);
46764706
_terminalRunIdsByThread.Remove(threadId);
46774707
_assistantFallbackPromotedThreads.Remove(threadId);

src/OpenClaw.WinNode.Cli/skill.md

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -401,13 +401,15 @@ List native chat outgoing queue entries.
401401
When `threadId`/`sessionKey` is omitted this returns all queued threads. Returns:
402402
`{ defaultThreadId, requestedThreadId, totalCount, selectedThread, threads }`
403403
where each thread has `{ threadId, count, messages }`, and each message has
404-
`{ id, text, createdAt, sendState, errorText, canCancel }`.
404+
`{ id, text, createdAt, sendState, errorText, canCancel }`. Queue message IDs
405+
are ephemeral UI/provider IDs; read them from the current `app.chat.queue.list`
406+
or `app.chat.snapshot` response and do not cache them across app lifetimes.
405407

406408
### app.chat.queue.cancel
407409
Cancel/remove one native chat outgoing queue entry before it is sent.
408410
```
409411
{
410-
"queuedMessageId": "string", // required; from app.chat.queue.list or app.chat.snapshot queue
412+
"queuedMessageId": "string", // required; from the current app.chat.queue.list or app.chat.snapshot queue
411413
"threadId": "string" // required; "sessionKey" alias also accepted
412414
}
413415
```

tests/OpenClaw.Tray.Tests/OpenClawChatDataProviderTests.cs

Lines changed: 134 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2815,6 +2815,29 @@ public async Task CancelQueuedMessageAsync_RemovesFailedQueuedCard()
28152815
Assert.False(snapshots[^1].QueuedMessagesByThread?.ContainsKey("main") == true);
28162816
}
28172817

2818+
[Fact]
2819+
public async Task CancelQueuedMessageAsync_RemovingLastQueuedMessageClearsStaleDrainGuard()
2820+
{
2821+
var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() });
2822+
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-active", Status = "started" });
2823+
await provider.LoadAsync();
2824+
bridge.RaiseStatus(ConnectionStatus.Connected);
2825+
2826+
await provider.SendMessageAsync("main", "active");
2827+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-active"));
2828+
await provider.SendMessageAsync("main", "queued");
2829+
2830+
var queued = Assert.Single(GetQueuedMessages(snapshots[^1], "main"));
2831+
var scheduled = GetQueuedDrainScheduledThreads(provider);
2832+
scheduled.Add("main");
2833+
2834+
var canceled = await provider.CancelQueuedMessageAsync("main", queued.Id);
2835+
2836+
Assert.True(canceled);
2837+
Assert.DoesNotContain("main", scheduled);
2838+
Assert.Empty(GetQueuedMessages(snapshots[^1], "main"));
2839+
}
2840+
28182841
[Fact]
28192842
public async Task CancelQueuedMessageAsync_ReturnsFalseForSendingQueuedCard()
28202843
{
@@ -2884,6 +2907,48 @@ public async Task CancelQueuedMessageAsync_DoesNotTurnActiveLocalRunIntoRemoteRu
28842907
Assert.Equal(0, historyCalls);
28852908
}
28862909

2910+
[Fact]
2911+
public async Task CancelQueuedMessageAsync_LastQueuedAfterLifecycleEndAllowsNextRemoteRunBackfill()
2912+
{
2913+
var historyCalls = 0;
2914+
var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() });
2915+
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-active", Status = "started" });
2916+
bridge.HistoryBehavior = _ =>
2917+
{
2918+
historyCalls++;
2919+
return Task.FromResult(new ChatHistoryInfo
2920+
{
2921+
SessionKey = "main",
2922+
Messages = new[]
2923+
{
2924+
new ChatMessageInfo { SessionKey = "main", Role = "user", Text = "remote prompt" },
2925+
},
2926+
});
2927+
};
2928+
await provider.LoadAsync();
2929+
bridge.RaiseStatus(ConnectionStatus.Connected);
2930+
2931+
await provider.SendMessageAsync("main", "active");
2932+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-active"));
2933+
await provider.SendMessageAsync("main", "queued");
2934+
var queued = Assert.Single(GetQueuedMessages(snapshots[^1], "main"));
2935+
2936+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"end"}""", runId: "run-active"));
2937+
var canceled = await provider.CancelQueuedMessageAsync("main", queued.Id);
2938+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-remote"));
2939+
2940+
Assert.True(canceled);
2941+
await WaitForConditionAsync(() =>
2942+
historyCalls > 0 &&
2943+
snapshots[^1].Timelines["main"].Entries.Any(entry =>
2944+
entry.Kind == ChatTimelineItemKind.User && entry.Text == "remote prompt"));
2945+
Assert.True(historyCalls > 0);
2946+
Assert.Contains(snapshots[^1].Timelines["main"].Entries, entry =>
2947+
entry.Kind == ChatTimelineItemKind.User && entry.Text == "remote prompt");
2948+
Assert.DoesNotContain(snapshots[^1].Timelines["main"].Entries, entry =>
2949+
entry.Kind == ChatTimelineItemKind.User && entry.Text == "queued");
2950+
}
2951+
28872952
[Fact]
28882953
public async Task QueuedSend_LifecycleStartBeforeAck_PromotesByIdempotencyKey()
28892954
{
@@ -2988,6 +3053,48 @@ public async Task QueuedSend_InFlightAckWithoutLifecycle_RequeuesAndRetriesSameI
29883053
e.Kind == ChatTimelineItemKind.User && e.Text == "Hello");
29893054
}
29903055

3056+
[Fact]
3057+
public async Task CancelQueuedMessageAsync_DeferredInFlightRetryRemovesQueuedCardWithoutAbortOrResend()
3058+
{
3059+
var (bridge, provider, snapshots, _) = CreateProvider(new[] { MainSession() });
3060+
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-1", Status = "started" });
3061+
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-2", Status = "in_flight" });
3062+
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-2", Status = "started" });
3063+
await provider.LoadAsync();
3064+
bridge.RaiseStatus(ConnectionStatus.Connected);
3065+
3066+
await provider.SendMessageAsync("main", "first");
3067+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-1"));
3068+
await provider.SendMessageAsync("main", "deferred");
3069+
3070+
bridge.RaiseChat(new ChatMessageInfo
3071+
{
3072+
SessionKey = "main",
3073+
Role = "assistant",
3074+
Text = "first response",
3075+
State = "final",
3076+
});
3077+
await WaitForConditionAsync(() =>
3078+
{
3079+
var queued = GetQueuedMessages(snapshots[^1], "main");
3080+
return bridge.SentMessages.Count == 2 &&
3081+
queued.Count == 1 &&
3082+
queued[0].Text == "deferred" &&
3083+
queued[0].SendState == ChatQueuedMessageSendState.Queued;
3084+
});
3085+
3086+
var deferred = Assert.Single(GetQueuedMessages(snapshots[^1], "main"));
3087+
var canceled = await provider.CancelQueuedMessageAsync("main", deferred.Id);
3088+
await Task.Delay(250);
3089+
3090+
Assert.True(canceled);
3091+
Assert.Empty(GetQueuedMessages(snapshots[^1], "main"));
3092+
Assert.Equal(new[] { "first", "deferred" }, bridge.SentMessages);
3093+
Assert.Empty(bridge.AbortedRunIds);
3094+
Assert.DoesNotContain(snapshots[^1].Timelines["main"].Entries, entry =>
3095+
entry.Kind == ChatTimelineItemKind.User && entry.Text == "deferred");
3096+
}
3097+
29913098
[Fact]
29923099
public async Task QueuedSend_InFlightAckThenLifecycleBeforeRetry_PromotesWithoutResend()
29933100
{
@@ -2996,6 +3103,23 @@ public async Task QueuedSend_InFlightAckThenLifecycleBeforeRetry_PromotesWithout
29963103
bridge.SendResults.Enqueue(new ChatSendResult { RunId = "run-2", Status = "in_flight" });
29973104
await provider.LoadAsync();
29983105
bridge.RaiseStatus(ConnectionStatus.Connected);
3106+
var lifecycleRaisedBeforeRetry = false;
3107+
provider.Changed += (_, e) =>
3108+
{
3109+
if (lifecycleRaisedBeforeRetry || bridge.SentMessages.Count != 2)
3110+
return;
3111+
3112+
var queued = GetQueuedMessages(e.Snapshot, "main");
3113+
if (queued.Count != 1 ||
3114+
queued[0].Text != "Hello" ||
3115+
queued[0].SendState != ChatQueuedMessageSendState.Queued)
3116+
{
3117+
return;
3118+
}
3119+
3120+
lifecycleRaisedBeforeRetry = true;
3121+
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-2"));
3122+
};
29993123

30003124
await provider.SendMessageAsync("main", "first");
30013125
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-1"));
@@ -3017,7 +3141,6 @@ await WaitForConditionAsync(() =>
30173141
queued[0].SendState == ChatQueuedMessageSendState.Queued;
30183142
});
30193143

3020-
bridge.RaiseAgent(MakeAgentEvent("lifecycle", """{"phase":"start"}""", runId: "run-2"));
30213144
await WaitForConditionAsync(() =>
30223145
GetQueuedMessages(snapshots[^1], "main").Count == 0 &&
30233146
snapshots[^1].Timelines["main"].Entries.Count(e =>
@@ -3031,6 +3154,7 @@ await WaitForConditionAsync(() =>
30313154
});
30323155
await Task.Delay(250);
30333156

3157+
Assert.True(lifecycleRaisedBeforeRetry);
30343158
Assert.Equal(2, bridge.SentMessages.Count);
30353159
Assert.Empty(GetQueuedMessages(snapshots[^1], "main"));
30363160
Assert.Single(snapshots[^1].Timelines["main"].Entries, e =>
@@ -6935,6 +7059,15 @@ private static IReadOnlyList<ChatQueuedMessage> GetQueuedMessages(ChatDataSnapsh
69357059
? queued
69367060
: Array.Empty<ChatQueuedMessage>();
69377061

7062+
private static ISet<string> GetQueuedDrainScheduledThreads(OpenClawChatDataProvider provider)
7063+
{
7064+
var field = typeof(OpenClawChatDataProvider).GetField(
7065+
"_queuedDrainScheduledThreads",
7066+
System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic);
7067+
Assert.NotNull(field);
7068+
return Assert.IsAssignableFrom<ISet<string>>(field.GetValue(provider));
7069+
}
7070+
69387071
private static bool HasFailedQueuedMessage(ChatDataSnapshot snapshot, string threadId, string text) =>
69397072
GetQueuedMessages(snapshot, threadId).Any(message =>
69407073
message.Text == text &&

0 commit comments

Comments
 (0)