Skip to content

Commit 3acef2b

Browse files
committed
fix(auto-reply): preserve queued undelivered finals
1 parent 974adc9 commit 3acef2b

6 files changed

Lines changed: 195 additions & 127 deletions

File tree

extensions/qa-channel/src/channel-actions.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -132,7 +132,7 @@ export const qaChannelMessageActions: ChannelMessageActionAdapter = {
132132
threadId,
133133
replyToId: readStringParam(params, "replyTo") ?? readStringParam(params, "replyToId"),
134134
});
135-
return jsonResult({ message });
135+
return jsonResult({ message, messageId: message.id });
136136
}
137137
case "thread-create": {
138138
const channelId =

extensions/qa-channel/src/channel.test.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -652,8 +652,12 @@ describe("qa-channel plugin", () => {
652652
message: "hello from action",
653653
},
654654
});
655-
const payload = extractToolPayload(result) as { message: { text: string } };
655+
const payload = extractToolPayload(result) as {
656+
message: { id: string; text: string };
657+
messageId?: string;
658+
};
656659
expect(payload.message.text).toBe("hello from action");
660+
expect(payload.messageId).toBe(payload.message.id);
657661

658662
const outbound = await state.waitFor({
659663
kind: "message-text",

src/auto-reply/reply/agent-runner.ts

Lines changed: 3 additions & 88 deletions
Original file line numberDiff line numberDiff line change
@@ -36,17 +36,8 @@ import {
3636
resolveSessionPluginTraceLines,
3737
type SessionEntry,
3838
} from "../../config/sessions.js";
39-
import {
40-
loadSessionEntry,
41-
persistSessionTranscriptTurn,
42-
updateSessionEntry,
43-
} from "../../config/sessions/session-accessor.js";
39+
import { loadSessionEntry, updateSessionEntry } from "../../config/sessions/session-accessor.js";
4440
import { parseSessionThreadInfoFast } from "../../config/sessions/thread-info.js";
45-
import {
46-
MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_CUSTOM_TYPE,
47-
MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_NOTICE,
48-
type MessageToolOnlyUndeliveredFinalNoticeDetails,
49-
} from "../../config/sessions/undelivered-final-notice.js";
5041
import type { TypingMode } from "../../config/types.js";
5142
import { resolveSessionTranscriptCandidates } from "../../gateway/session-utils.fs.js";
5243
import { logVerbose } from "../../globals.js";
@@ -83,7 +74,7 @@ import {
8374
} from "../reply-payload.js";
8475
import type { OriginatingChannelType, TemplateContext } from "../templating.js";
8576
import type { VerboseLevel } from "../thinking.js";
86-
import { isSilentReplyText, SILENT_REPLY_TOKEN } from "../tokens.js";
77+
import { SILENT_REPLY_TOKEN } from "../tokens.js";
8778
import type { GetReplyOptions, ReplyPayload } from "../types.js";
8879
import {
8980
buildEmptyInteractiveReplyPayload,
@@ -153,6 +144,7 @@ import {
153144
} from "./stranded-reply-recovery.js";
154145
import { createTypingSignaler } from "./typing-mode.js";
155146
import type { TypingController } from "./typing.js";
147+
import { persistMessageToolOnlyUndeliveredFinalNotice } from "./undelivered-final-notice.js";
156148

157149
const BLOCK_REPLY_SEND_TIMEOUT_MS = 15_000;
158150
const RESTART_LIFECYCLE_REPLY_TEXT =
@@ -291,83 +283,6 @@ async function hasMessagingToolDeliveredFinalTextToCurrentSourceRoute(params: {
291283
return hasDeliveredFinalText({ finalText: params.finalText, sentTexts: routeSentTexts });
292284
}
293285

294-
async function persistMessageToolOnlyUndeliveredFinalNotice(params: {
295-
cfg: OpenClawConfig;
296-
sessionEntry?: SessionEntry;
297-
sessionStore?: Record<string, SessionEntry>;
298-
sessionId: string;
299-
expectedLifecycleRevision?: string;
300-
sessionKey?: string;
301-
storePath?: string;
302-
sessionAgentId?: string;
303-
threadId?: string | number;
304-
workspaceDir: string;
305-
sourceReplyDeliveryMode?: string;
306-
sendPolicyDenied: boolean;
307-
finalTextDeliveredToCurrentSourceRoute: boolean;
308-
finalText: string;
309-
}): Promise<void> {
310-
if (
311-
params.sourceReplyDeliveryMode !== "message_tool_only" ||
312-
params.sendPolicyDenied ||
313-
params.finalTextDeliveredToCurrentSourceRoute
314-
) {
315-
return;
316-
}
317-
const trimmed = params.finalText.trim();
318-
if (!trimmed || isSilentReplyText(trimmed)) {
319-
return;
320-
}
321-
const sessionKey = params.sessionKey?.trim();
322-
const sessionId = params.sessionId.trim();
323-
if (!sessionKey || !sessionId) {
324-
return;
325-
}
326-
327-
await persistSessionTranscriptTurn(
328-
{
329-
sessionId,
330-
sessionKey,
331-
sessionEntry: params.sessionEntry,
332-
sessionStore: params.sessionStore,
333-
storePath: params.storePath,
334-
agentId: params.sessionAgentId,
335-
threadId: params.threadId,
336-
},
337-
{
338-
config: params.cfg,
339-
cwd: params.workspaceDir,
340-
messages: [
341-
{
342-
message: {
343-
role: "custom",
344-
customType: MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_CUSTOM_TYPE,
345-
content: MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_NOTICE,
346-
display: false,
347-
details: {
348-
sourceReplyDeliveryMode: "message_tool_only",
349-
delivered: false,
350-
finalTextLength: trimmed.length,
351-
} satisfies MessageToolOnlyUndeliveredFinalNoticeDetails,
352-
timestamp: Date.now(),
353-
},
354-
},
355-
],
356-
publishWhen: "when-appended",
357-
touchSessionEntry: true,
358-
updateMode: "file-only",
359-
...(params.storePath
360-
? {
361-
expectedSessionId: sessionId,
362-
...(params.expectedLifecycleRevision
363-
? { expectedLifecycleRevision: params.expectedLifecycleRevision }
364-
: {}),
365-
}
366-
: {}),
367-
},
368-
);
369-
}
370-
371286
function resolveReplyRunDeliveryContext(params: {
372287
cfg: OpenClawConfig;
373288
sessionCtx: TemplateContext;

src/auto-reply/reply/followup-runner.test.ts

Lines changed: 74 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vite
88
import { setCliSessionBinding } from "../../agents/cli-session.js";
99
import type { OpenClawConfig } from "../../config/config.js";
1010
import type { SessionEntry } from "../../config/sessions/types.js";
11+
import { MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_CUSTOM_TYPE } from "../../config/sessions/undelivered-final-notice.js";
1112
import {
1213
createUserTurnTranscriptRecorder,
1314
type PersistedUserTurnMessage,
@@ -58,6 +59,7 @@ const FOLLOWUP_TEST_QUEUES = new Map<
5859
>();
5960
const FOLLOWUP_TEST_SESSION_STORES = new Map<string, Record<string, SessionEntry>>();
6061
const FOLLOWUP_TEST_SESSION_STORE_PATHS = new Set<string>();
62+
let inspectEnqueueFollowupRunForTest: (() => void) | undefined;
6163

6264
function debugFollowupTest(message: string): void {
6365
if (!FOLLOWUP_DEBUG) {
@@ -242,6 +244,7 @@ function enqueueFollowupRunForFollowupTest(
242244
_restartIfIdle?: unknown,
243245
options?: { position?: "tail" | "front" },
244246
): boolean {
247+
inspectEnqueueFollowupRunForTest?.();
245248
if (options?.position === "front") {
246249
run.protectFromQueueOverflow = true;
247250
}
@@ -387,6 +390,7 @@ async function loadFreshFollowupRunnerModuleForTest() {
387390
release: async () => {},
388391
})),
389392
resolveSessionLockMaxHoldFromTimeout: vi.fn(() => 1),
393+
resolveSessionWriteLockOptions: vi.fn(() => ({})),
390394
}));
391395
vi.doMock("../../agents/embedded-agent.js", () => ({
392396
abortEmbeddedAgentRun: vi.fn(async () => false),
@@ -604,6 +608,7 @@ beforeEach(() => {
604608
clearFollowupQueue("main");
605609
FOLLOWUP_TEST_QUEUES.clear();
606610
FOLLOWUP_TEST_SESSION_STORES.clear();
611+
inspectEnqueueFollowupRunForTest = undefined;
607612
});
608613

609614
afterEach(() => {
@@ -5750,46 +5755,82 @@ describe("createFollowupRunner messaging delivery and dedupe", () => {
57505755
expect(onBlockReply).not.toHaveBeenCalled();
57515756
});
57525757

5753-
it("enqueues a one-shot recovery retry for substantive message-tool-only queued followup finals", async () => {
5758+
it("persists queued message-tool-only undelivered finals before enqueueing recovery", async () => {
57545759
const finalText =
57555760
"Here is the answer the queued user asked for. It includes enough detail to be a visible response, and it has another sentence so the substantive-final detector treats it as a real reply.";
57565761
const parentOnComplete = vi.fn();
57575762
const parentLifecycle = { onComplete: parentOnComplete };
57585763
const queued = baseQueuedRun("discord");
5759-
const { onBlockReply } = await runMessagingCase({
5760-
agentResult: {
5761-
payloads: [{ text: finalText }],
5762-
meta: { finalAssistantVisibleText: finalText },
5763-
},
5764-
queued: {
5765-
...queued,
5766-
originatingChannel: "discord",
5767-
originatingTo: "channel:C1",
5768-
queuedLifecycle: parentLifecycle,
5769-
run: {
5770-
...queued.run,
5771-
sourceReplyDeliveryMode: "message_tool_only",
5764+
const sessionFile = path.join(tmpdir(), "openclaw-followup-stranded-session.jsonl");
5765+
const storePath = path.join(tmpdir(), "openclaw-followup-stranded-store.json");
5766+
const sessionEntry: SessionEntry = {
5767+
sessionId: queued.run.sessionId,
5768+
sessionFile,
5769+
updatedAt: Date.now(),
5770+
};
5771+
const sessionStore = { main: sessionEntry };
5772+
await fs.writeFile(
5773+
sessionFile,
5774+
`${JSON.stringify({
5775+
type: "session",
5776+
id: sessionEntry.sessionId,
5777+
timestamp: new Date().toISOString(),
5778+
cwd: queued.run.workspaceDir,
5779+
})}\n`,
5780+
"utf8",
5781+
);
5782+
let transcriptAtEnqueue = "";
5783+
inspectEnqueueFollowupRunForTest = () => {
5784+
transcriptAtEnqueue = fsSync.existsSync(sessionFile)
5785+
? fsSync.readFileSync(sessionFile, "utf8")
5786+
: "";
5787+
};
5788+
5789+
try {
5790+
const { onBlockReply } = await runMessagingCase({
5791+
agentResult: {
5792+
payloads: [{ text: finalText }],
5793+
meta: { finalAssistantVisibleText: finalText },
57725794
},
5773-
} as FollowupRun,
5774-
});
5795+
runnerOverrides: {
5796+
sessionEntry,
5797+
sessionStore,
5798+
sessionKey: "main",
5799+
storePath,
5800+
},
5801+
queued: {
5802+
...queued,
5803+
originatingChannel: "discord",
5804+
originatingTo: "channel:C1",
5805+
queuedLifecycle: parentLifecycle,
5806+
run: {
5807+
...queued.run,
5808+
sessionFile,
5809+
sourceReplyDeliveryMode: "message_tool_only",
5810+
},
5811+
} as FollowupRun,
5812+
});
57755813

5776-
expect(onBlockReply).not.toHaveBeenCalled();
5777-
expect(routeReplyMock).not.toHaveBeenCalled();
5778-
const retry = FOLLOWUP_TEST_QUEUES.get("main")?.items[0];
5779-
expect(retry?.summaryLine).toBe("stranded-reply-retry");
5780-
expect(retry?.strandedReplyRetry).toBe(true);
5781-
expect(retry?.disableCollectBatching).toBe(true);
5782-
expect(retry?.protectFromQueueOverflow).toBe(true);
5783-
expect(retry?.transcriptPrompt).toBeUndefined();
5784-
expect(retry?.userTurnTranscriptRecorder).toBeUndefined();
5785-
expect(retry?.currentInboundContext).toBeUndefined();
5786-
expect(retry?.run.suppressNextUserMessagePersistence).toBe(true);
5787-
expect(retry?.run.sourceReplyDeliveryMode).toBe("message_tool_only");
5788-
expect(retry?.prompt).toContain("message(action=send)");
5789-
expect(retry?.prompt).toContain(finalText);
5790-
// System retry detaches from the client turn lifecycle; parent completion owns onComplete once.
5791-
expect(retry?.queuedLifecycle).toBeUndefined();
5792-
expect(parentOnComplete).toHaveBeenCalledTimes(1);
5814+
expect(transcriptAtEnqueue).toContain(MESSAGE_TOOL_ONLY_UNDELIVERED_FINAL_CUSTOM_TYPE);
5815+
expect(onBlockReply).not.toHaveBeenCalled();
5816+
expect(routeReplyMock).not.toHaveBeenCalled();
5817+
const retry = FOLLOWUP_TEST_QUEUES.get("main")?.items[0];
5818+
expect(retry?.summaryLine).toBe("stranded-reply-retry");
5819+
expect(retry?.strandedReplyRetry).toBe(true);
5820+
expect(retry?.disableCollectBatching).toBe(true);
5821+
expect(retry?.protectFromQueueOverflow).toBe(true);
5822+
expect(retry?.transcriptPrompt).toBeUndefined();
5823+
expect(retry?.userTurnTranscriptRecorder).toBeUndefined();
5824+
expect(retry?.currentInboundContext).toBeUndefined();
5825+
expect(retry?.run.suppressNextUserMessagePersistence).toBe(true);
5826+
expect(retry?.run.sourceReplyDeliveryMode).toBe("message_tool_only");
5827+
expect(retry?.prompt).toContain("message(action=send)");
5828+
expect(retry?.prompt).toContain(finalText);
5829+
expect(retry?.queuedLifecycle).toBeUndefined();
5830+
expect(parentOnComplete).toHaveBeenCalledTimes(1);
5831+
} finally {
5832+
await fs.rm(sessionFile, { force: true });
5833+
}
57935834
});
57945835

57955836
it("excludes raw trace and status payloads from queued stranded recovery prompts", async () => {

src/auto-reply/reply/followup-runner.ts

Lines changed: 26 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -129,6 +129,7 @@ import {
129129
} from "./stranded-reply-recovery.js";
130130
import { createTypingSignaler } from "./typing-mode.js";
131131
import type { TypingController } from "./typing.js";
132+
import { persistMessageToolOnlyUndeliveredFinalNotice } from "./undelivered-final-notice.js";
132133

133134
type EmbeddedAgentRunResult = Awaited<ReturnType<typeof runEmbeddedAgent>>;
134135

@@ -1687,15 +1688,36 @@ export function createFollowupRunner(params: {
16871688
typeof runResult.meta?.finalAssistantVisibleText === "string"
16881689
? normalizeAssistantFinalDeliveryText(runResult.meta.finalAssistantVisibleText)
16891690
: "";
1691+
const successfulSourceReplyDelivery = hasSuccessfulFollowupSourceReplyDelivery({
1692+
didDeliverSourceReplyViaMessageTool: runResult.didDeliverSourceReplyViaMessageTool,
1693+
messagingToolSourceReplyPayloads: runResult.messagingToolSourceReplyPayloads,
1694+
});
1695+
try {
1696+
await persistMessageToolOnlyUndeliveredFinalNotice({
1697+
cfg: runtimeConfig,
1698+
sessionEntry: activeSessionEntry,
1699+
sessionStore,
1700+
sessionId: activeSessionEntry?.sessionId ?? run.sessionId,
1701+
expectedLifecycleRevision: activeSessionEntry?.lifecycleRevision,
1702+
sessionKey: replySessionKey,
1703+
storePath,
1704+
sessionAgentId: run.agentId,
1705+
threadId: queued.originatingThreadId,
1706+
workspaceDir: run.workspaceDir,
1707+
sourceReplyDeliveryMode: sourceReplyPolicy.sourceReplyDeliveryMode,
1708+
sendPolicyDenied: sourceReplyPolicy.sendPolicyDenied,
1709+
finalTextDeliveredToCurrentSourceRoute: successfulSourceReplyDelivery,
1710+
finalText: assistantFinalText,
1711+
});
1712+
} catch (error) {
1713+
logVerbose(`failed to persist undelivered final notice: ${String(error)}`);
1714+
}
16901715
const isStrandedReply =
16911716
queued.currentInboundEventKind !== "room_event" &&
16921717
shouldWarnAboutPrivateMessageToolFinal({
16931718
sourceReplyDeliveryMode: sourceReplyPolicy.sourceReplyDeliveryMode,
16941719
sendPolicyDenied: sourceReplyPolicy.sendPolicyDenied,
1695-
successfulSourceReplyDelivery: hasSuccessfulFollowupSourceReplyDelivery({
1696-
didDeliverSourceReplyViaMessageTool: runResult.didDeliverSourceReplyViaMessageTool,
1697-
messagingToolSourceReplyPayloads: runResult.messagingToolSourceReplyPayloads,
1698-
}),
1720+
successfulSourceReplyDelivery,
16991721
finalText: assistantFinalText,
17001722
});
17011723
if (!isStrandedReply) {

0 commit comments

Comments
 (0)