Skip to content

Commit e85a541

Browse files
fix(auto-reply): harden pending final delivery recovery
1 parent 4ff5e4d commit e85a541

10 files changed

Lines changed: 421 additions & 24 deletions

src/agents/main-session-restart-recovery.test.ts

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -616,6 +616,51 @@ describe("main-session-restart-recovery", () => {
616616
);
617617
});
618618

619+
it("resumes Slack DM pending finals through the native DM channel when stored target is user-scoped", async () => {
620+
const sessionsDir = await makeSessionsDir();
621+
const pendingPayload = "The final answer is 42.";
622+
await writeStore(sessionsDir, {
623+
"agent:main:slack:direct:u039p2yjr29": {
624+
sessionId: "main-session",
625+
updatedAt: Date.now() - 10_000,
626+
status: "running",
627+
abortedLastRun: true,
628+
pendingFinalDelivery: true,
629+
pendingFinalDeliveryText: pendingPayload,
630+
pendingFinalDeliveryContext: {
631+
channel: "slack",
632+
to: "user:U039P2YJR29",
633+
accountId: "son-of-anton",
634+
},
635+
pendingFinalDeliveryCreatedAt: Date.now() - 5_000,
636+
origin: {
637+
provider: "slack",
638+
chatType: "direct",
639+
to: "user:U039P2YJR29",
640+
nativeChannelId: "D0ACL8LRTJP",
641+
accountId: "son-of-anton",
642+
},
643+
},
644+
});
645+
await writeTranscript(sessionsDir, "main-session", [
646+
{ role: "user", content: "calculate the answer" },
647+
{ role: "assistant", content: [{ type: "toolCall", id: "call-1", name: "calc" }] },
648+
{ role: "toolResult", content: "42" },
649+
]);
650+
651+
const result = await recoverRestartAbortedMainSessions({ stateDir: tmpDir });
652+
653+
expect(result).toEqual({ recovered: 1, failed: 0, skipped: 0 });
654+
expect(firstGatewayParams()).toMatchObject({
655+
deliver: true,
656+
bestEffortDeliver: true,
657+
channel: "slack",
658+
to: "channel:D0ACL8LRTJP",
659+
accountId: "son-of-anton",
660+
});
661+
expect(firstGatewayParams().message).toContain(pendingPayload);
662+
});
663+
619664
it("sanitizes durable pending final delivery payloads before resume prompts", async () => {
620665
const sessionsDir = await makeSessionsDir();
621666
const pendingPayload = [

src/agents/main-session-restart-recovery.ts

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,10 @@ import crypto from "node:crypto";
66
import fs from "node:fs";
77
import path from "node:path";
88
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
9-
import { sanitizePendingFinalDeliveryText } from "../auto-reply/reply/pending-final-delivery.js";
9+
import {
10+
resolveSlackDirectPendingFinalDeliveryContext,
11+
sanitizePendingFinalDeliveryText,
12+
} from "../auto-reply/reply/pending-final-delivery.js";
1013
import { resolveStateDir } from "../config/paths.js";
1114
import {
1215
type SessionEntry,
@@ -369,11 +372,22 @@ function resolveRestartRecoveryDeliveryContext(params: {
369372
) {
370373
return undefined;
371374
}
372-
return {
373-
...deliveryContext,
374-
channel,
375-
to,
376-
};
375+
return (
376+
resolveSlackDirectPendingFinalDeliveryContext({
377+
context: {
378+
...deliveryContext,
379+
channel,
380+
to,
381+
},
382+
nativeChannelId: params.entry.origin?.nativeChannelId,
383+
chatType: params.entry.origin?.chatType,
384+
directUserTarget: params.entry.origin?.to,
385+
}) ?? {
386+
...deliveryContext,
387+
channel,
388+
to,
389+
}
390+
);
377391
}
378392

379393
async function resumeMainSession(params: {

src/auto-reply/reply/agent-runner.runreplyagent.e2e.test.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,9 @@ const state = vi.hoisted(() => ({
3434
queueEmbeddedAgentMessageMock: vi.fn(),
3535
runEmbeddedAgentMock: vi.fn(),
3636
}));
37+
const diagnosticMocks = vi.hoisted(() => ({
38+
markDiagnosticSessionPendingFinalDelivery: vi.fn(),
39+
}));
3740

3841
function countMatching<T>(items: readonly T[], predicate: (item: T) => boolean): number {
3942
let count = 0;
@@ -111,6 +114,16 @@ vi.mock("./queue.js", () => ({
111114
refreshQueuedFollowupSession: vi.fn(),
112115
scheduleFollowupDrain: vi.fn(),
113116
}));
117+
vi.mock("../../logging/diagnostic.js", async () => {
118+
const actual = await vi.importActual<typeof import("../../logging/diagnostic.js")>(
119+
"../../logging/diagnostic.js",
120+
);
121+
return {
122+
...actual,
123+
markDiagnosticSessionPendingFinalDelivery:
124+
diagnosticMocks.markDiagnosticSessionPendingFinalDelivery,
125+
};
126+
});
114127

115128
beforeAll(async () => {
116129
// Avoid attributing the initial agent-runner import cost to the first test case.
@@ -134,6 +147,7 @@ beforeEach(() => {
134147
});
135148
state.queueEmbeddedAgentMessageMock.mockReset();
136149
state.queueEmbeddedAgentMessageMock.mockReturnValue(false);
150+
diagnosticMocks.markDiagnosticSessionPendingFinalDelivery.mockReset();
137151
vi.mocked(enqueueFollowupRun).mockClear();
138152
vi.mocked(refreshQueuedFollowupSession).mockClear();
139153
vi.mocked(scheduleFollowupDrain).mockClear();
@@ -517,6 +531,11 @@ describe("runReplyAgent pending final delivery capture", () => {
517531
const stored = await readStoredMainSession(storePath);
518532
expect(stored.pendingFinalDelivery).toBe(true);
519533
expect(stored.pendingFinalDeliveryText).toBe("visible final");
534+
expect(stored.pendingFinalDeliveryIntentId).toBe("msg");
535+
expect(diagnosticMocks.markDiagnosticSessionPendingFinalDelivery).toHaveBeenCalledWith({
536+
sessionKey: "main",
537+
pending: true,
538+
});
520539
});
521540

522541
it("persists auto-reply delivery context for restart recovery", async () => {
@@ -568,6 +587,7 @@ describe("runReplyAgent pending final delivery capture", () => {
568587
accountId: "work",
569588
threadId: "1503645939964055592",
570589
});
590+
expect(stored.pendingFinalDeliveryIntentId).toBe("1503645939964055592");
571591
expect(stored.restartRecoveryDeliveryContext).toBeUndefined();
572592
expect(stored.restartRecoveryDeliveryRunId).toBeUndefined();
573593
});

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

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ import {
4040
} from "../../infra/diagnostic-trace-context.js";
4141
import { measureDiagnosticsTimelineSpan } from "../../infra/diagnostics-timeline.js";
4242
import { enqueueSystemEvent } from "../../infra/system-events.js";
43+
import { markDiagnosticSessionPendingFinalDelivery } from "../../logging/diagnostic.js";
4344
import { CommandLaneClearedError, GatewayDrainingError } from "../../process/command-queue.js";
4445
import { shouldPreserveUserFacingSessionStateForInputProvenance } from "../../sessions/input-provenance.js";
4546
import { resolveSendPolicy } from "../../sessions/send-policy.js";
@@ -93,7 +94,10 @@ import { resolveEffectiveReplyRoute } from "./effective-reply-route.js";
9394
import { createFollowupRunner } from "./followup-runner.js";
9495
import { REPLY_RUN_STILL_SHUTTING_DOWN_TEXT } from "./get-reply-run-queue.js";
9596
import { resolveOriginMessageProvider, resolveOriginMessageTo } from "./origin-routing.js";
96-
import { sanitizePendingFinalDeliveryText } from "./pending-final-delivery.js";
97+
import {
98+
resolveSlackDirectPendingFinalDeliveryContext,
99+
sanitizePendingFinalDeliveryText,
100+
} from "./pending-final-delivery.js";
97101
import { drainPendingToolTasks } from "./pending-tool-task-drain.js";
98102
import { readPostCompactionContext } from "./post-compaction-context.js";
99103
import {
@@ -196,13 +200,21 @@ function resolveReplyRunDeliveryContext(params: {
196200
normalizeOptionalString(
197201
parseSessionThreadInfoFast(params.sessionCtx.SessionKey ?? params.sessionKey).threadId,
198202
);
199-
return normalizeDeliveryContext({
203+
const context = normalizeDeliveryContext({
200204
...resolveEffectiveReplyRoute({
201205
ctx: params.sessionCtx,
202206
entry: params.sessionEntry,
203207
}),
204208
threadId,
205209
});
210+
return resolveSlackDirectPendingFinalDeliveryContext({
211+
context,
212+
nativeChannelId: normalizeOptionalString(params.sessionCtx.NativeChannelId),
213+
chatType: normalizeOptionalString(params.sessionCtx.ChatType),
214+
directUserTarget: normalizeOptionalString(
215+
params.sessionCtx.OriginatingTo ?? params.sessionCtx.To,
216+
),
217+
});
206218
}
207219

208220
function hasNonEmptyStringArray(value: unknown): boolean {
@@ -2334,6 +2346,9 @@ export async function runReplyAgent(params: {
23342346
});
23352347
}
23362348
const pendingText = sourceReplyPolicy.suppressDelivery ? "" : finalDeliveryText;
2349+
const pendingFinalDeliveryIntentId = normalizeOptionalString(
2350+
sessionCtx.MessageSidFull ?? sessionCtx.MessageSid,
2351+
);
23372352
const agentId = followupRun.run.agentId;
23382353
const heartbeatAgentCfg = agentId ? resolveAgentConfig(cfg, agentId)?.heartbeat : undefined;
23392354
const heartbeatAckMaxChars = Math.max(
@@ -2369,10 +2384,12 @@ export async function runReplyAgent(params: {
23692384
pendingFinalDelivery: true,
23702385
pendingFinalDeliveryText: resolvedPendingText,
23712386
pendingFinalDeliveryContext,
2387+
pendingFinalDeliveryIntentId,
23722388
pendingFinalDeliveryCreatedAt: Date.now(),
23732389
updatedAt: Date.now(),
23742390
},
23752391
});
2392+
markDiagnosticSessionPendingFinalDelivery({ sessionKey, pending: true });
23762393
}
23772394
}
23782395

src/auto-reply/reply/dispatch-from-config.reply-dispatch.test.ts

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,59 @@ describe("dispatchReplyFromConfig reply_dispatch hook", () => {
263263
expect(sessionStoreMocks.currentEntry?.pendingFinalDeliveryIntentId).toBeUndefined();
264264
});
265265

266+
it("replays old Slack user-scoped pending finals through the concrete DM channel", async () => {
267+
hookMocks.runner.hasHooks.mockReturnValue(false);
268+
sessionStoreMocks.currentEntry = {
269+
sessionKey: "agent:test:session",
270+
pendingFinalDelivery: true,
271+
pendingFinalDeliveryText: "lost final answer",
272+
pendingFinalDeliveryCreatedAt: 1,
273+
pendingFinalDeliveryContext: {
274+
channel: "slack",
275+
to: "user:U039P2YJR29",
276+
accountId: "son-of-anton",
277+
},
278+
pendingFinalDeliveryIntentId: "1780718204.607",
279+
origin: {
280+
provider: "slack",
281+
chatType: "direct",
282+
to: "user:U039P2YJR29",
283+
nativeChannelId: "D0ACL8LRTJP",
284+
accountId: "son-of-anton",
285+
},
286+
};
287+
sessionStoreMocks.resolveSessionStoreEntry.mockReturnValue({
288+
existing: sessionStoreMocks.currentEntry,
289+
});
290+
const dispatcher = createDispatcher();
291+
const replyResolver = vi.fn(async () => ({ text: "regenerated answer" }));
292+
293+
const result = await dispatchReplyFromConfig({
294+
ctx: createHookCtx({
295+
Provider: "slack",
296+
Surface: "slack",
297+
OriginatingChannel: "slack",
298+
OriginatingTo: "user:U039P2YJR29",
299+
AccountId: "son-of-anton",
300+
To: "user:U039P2YJR29",
301+
ChatType: "direct",
302+
NativeChannelId: "D0ACL8LRTJP",
303+
MessageSid: "1780718204.607",
304+
}),
305+
cfg: emptyConfig,
306+
dispatcher,
307+
replyResolver,
308+
});
309+
310+
expect(result.queuedFinal).toBe(true);
311+
expect(dispatcher.sendFinalReply).toHaveBeenCalledOnce();
312+
expect(vi.mocked(dispatcher.sendFinalReply).mock.calls[0]?.[0]).toMatchObject({
313+
text: "lost final answer",
314+
});
315+
expect(replyResolver).not.toHaveBeenCalled();
316+
expect(sessionStoreMocks.currentEntry?.pendingFinalDelivery).toBeUndefined();
317+
});
318+
266319
it("does not replay pending final delivery when the current route differs", async () => {
267320
hookMocks.runner.hasHooks.mockReturnValue(false);
268321
sessionStoreMocks.currentEntry = {
@@ -303,6 +356,53 @@ describe("dispatchReplyFromConfig reply_dispatch hook", () => {
303356
expect(sessionStoreMocks.currentEntry?.pendingFinalDeliveryText).toBe("lost final answer");
304357
});
305358

359+
it("includes replayed pending final accounting when reply_dispatch handles the turn", async () => {
360+
hookMocks.runner.hasHooks.mockImplementation(
361+
(hookName?: string) => hookName === "reply_dispatch",
362+
);
363+
hookMocks.runner.runReplyDispatch.mockResolvedValue({
364+
handled: true,
365+
queuedFinal: false,
366+
counts: { tool: 0, block: 0, final: 0 },
367+
});
368+
sessionStoreMocks.currentEntry = {
369+
sessionKey: "agent:test:session",
370+
pendingFinalDelivery: true,
371+
pendingFinalDeliveryText: "lost final answer",
372+
pendingFinalDeliveryCreatedAt: 1,
373+
pendingFinalDeliveryContext: {
374+
channel: "slack",
375+
to: "D0ACL8LRTJP",
376+
accountId: "son-of-anton",
377+
},
378+
pendingFinalDeliveryIntentId: "previous-message",
379+
};
380+
sessionStoreMocks.resolveSessionStoreEntry.mockReturnValue({
381+
existing: sessionStoreMocks.currentEntry,
382+
});
383+
const dispatcher = createDispatcher();
384+
385+
const result = await dispatchReplyFromConfig({
386+
ctx: createHookCtx({
387+
Provider: "slack",
388+
Surface: "slack",
389+
OriginatingChannel: "slack",
390+
OriginatingTo: "D0ACL8LRTJP",
391+
AccountId: "son-of-anton",
392+
To: "D0ACL8LRTJP",
393+
ChatType: "direct",
394+
MessageSid: "new-message",
395+
}),
396+
cfg: emptyConfig,
397+
dispatcher,
398+
replyResolver: async () => ({ text: "plugin handled" }),
399+
});
400+
401+
expect(result.queuedFinal).toBe(true);
402+
expect(result.counts.final).toBe(0);
403+
expect(dispatcher.sendFinalReply).toHaveBeenCalledOnce();
404+
});
405+
306406
it("preserves pending final delivery when final dispatch fails", async () => {
307407
hookMocks.runner.hasHooks.mockReturnValue(false);
308408
sessionStoreMocks.currentEntry = {

0 commit comments

Comments
 (0)