Skip to content

Commit 4decdf6

Browse files
samzongsteipete
andauthored
[Fix] Deliver restart recovery replies (#86089)
* fix(agents): deliver restart recovery replies * fix(auto-reply): import session entry updater * test(auto-reply): use current embedded agent mock * test(feishu): refresh typed account fixture --------- Co-authored-by: Peter Steinberger <[email protected]>
1 parent ac0fb97 commit 4decdf6

12 files changed

Lines changed: 1207 additions & 107 deletions

src/agents/agent-command.live-model-switch.test.ts

Lines changed: 535 additions & 64 deletions
Large diffs are not rendered by default.

src/agents/agent-command.ts

Lines changed: 203 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,11 @@ import {
2121
registerAgentRunContext,
2222
} from "../infra/agent-events.js";
2323
import { formatErrorMessage } from "../infra/errors.js";
24+
import {
25+
resolveAgentDeliveryPlan,
26+
resolveAgentOutboundTarget,
27+
} from "../infra/outbound/agent-delivery.js";
28+
import { resolveMessageChannelSelection } from "../infra/outbound/channel-selection.js";
2429
import { buildOutboundSessionContext } from "../infra/outbound/session-context.js";
2530
import { parseStrictNonNegativeInteger } from "../infra/parse-finite-number.js";
2631
import { createSubsystemLogger } from "../logging/subsystem.js";
@@ -47,7 +52,15 @@ import type { getRemoteSkillEligibility } from "../skills/runtime/remote.js";
4752
import type { resolveReusableWorkspaceSkillSnapshot } from "../skills/runtime/session-snapshot.js";
4853
import { createTrajectoryRuntimeRecorder } from "../trajectory/runtime.js";
4954
import { resolveUserPath } from "../utils.js";
50-
import { resolveMessageChannel } from "../utils/message-channel.js";
55+
import {
56+
normalizeDeliveryContext,
57+
type DeliveryContext,
58+
} from "../utils/delivery-context.shared.js";
59+
import {
60+
INTERNAL_MESSAGE_CHANNEL,
61+
isDeliverableMessageChannel,
62+
resolveMessageChannel,
63+
} from "../utils/message-channel.js";
5164
import { resolveAgentRuntimeConfig } from "./agent-runtime-config.js";
5265
import {
5366
clearAutoFallbackPrimaryProbeSelection,
@@ -286,6 +299,115 @@ function clearPendingFinalDeliveryFields(entry: SessionEntry, updatedAt: number)
286299
};
287300
}
288301

302+
async function resolveCurrentRunDeliveryContext(params: {
303+
cfg: OpenClawConfig;
304+
opts: AgentCommandOpts;
305+
sessionEntry?: SessionEntry;
306+
}): Promise<DeliveryContext | undefined> {
307+
const { cfg, opts, sessionEntry } = params;
308+
if (opts.deliver !== true) {
309+
return undefined;
310+
}
311+
// Restart recovery only needs durable route fields; final delivery resolves plugin-specific routes.
312+
const deliveryPlan = resolveAgentDeliveryPlan({
313+
sessionEntry,
314+
requestedChannel: opts.replyChannel ?? opts.channel,
315+
explicitTo: opts.replyTo ?? opts.to,
316+
explicitThreadId: opts.threadId,
317+
accountId: opts.replyAccountId ?? opts.accountId,
318+
wantsDelivery: true,
319+
turnSourceChannel: opts.runContext?.messageChannel ?? opts.messageChannel,
320+
turnSourceTo: opts.runContext?.currentChannelId ?? opts.to,
321+
turnSourceAccountId: opts.runContext?.accountId ?? opts.accountId,
322+
turnSourceThreadId: opts.runContext?.currentThreadTs ?? opts.threadId,
323+
});
324+
const explicitChannelHint = normalizeOptionalString(opts.replyChannel ?? opts.channel);
325+
const explicitThreadId =
326+
opts.threadId != null && opts.threadId !== "" ? opts.threadId : undefined;
327+
let effectivePlan = deliveryPlan;
328+
if (deliveryPlan.resolvedChannel === INTERNAL_MESSAGE_CHANNEL && !explicitChannelHint) {
329+
try {
330+
const selection = await resolveMessageChannelSelection({ cfg });
331+
effectivePlan = {
332+
...deliveryPlan,
333+
resolvedChannel: selection.channel,
334+
deliveryTargetMode: deliveryPlan.deliveryTargetMode ?? "implicit",
335+
};
336+
} catch {
337+
return undefined;
338+
}
339+
}
340+
if (!isDeliverableMessageChannel(effectivePlan.resolvedChannel)) {
341+
return undefined;
342+
}
343+
const targetMode =
344+
opts.deliveryTargetMode ??
345+
effectivePlan.deliveryTargetMode ??
346+
(opts.to ? "explicit" : "implicit");
347+
const resolvedTo =
348+
effectivePlan.resolvedTo ??
349+
resolveAgentOutboundTarget({
350+
cfg,
351+
plan: effectivePlan,
352+
targetMode,
353+
validateExplicitTarget: false,
354+
}).resolvedTo;
355+
if (!resolvedTo) {
356+
return undefined;
357+
}
358+
const threadId =
359+
targetMode === "explicit"
360+
? (explicitThreadId ??
361+
(effectivePlan.baseDelivery.threadIdSource === "explicit"
362+
? effectivePlan.resolvedThreadId
363+
: undefined))
364+
: effectivePlan.resolvedThreadId;
365+
return normalizeDeliveryContext({
366+
channel: effectivePlan.resolvedChannel,
367+
to: resolvedTo,
368+
accountId: effectivePlan.resolvedAccountId,
369+
threadId,
370+
});
371+
}
372+
373+
function shouldPersistCurrentRunSessionCleanup(
374+
current: SessionEntry | undefined,
375+
sessionId: string,
376+
): boolean {
377+
return (
378+
current !== undefined && current.sessionId === sessionId && current.abortedLastRun !== true
379+
);
380+
}
381+
382+
function shouldPersistRestartRecoveryContextClaim(
383+
current: SessionEntry | undefined,
384+
sessionId: string,
385+
runId: string,
386+
allowCreate: boolean,
387+
): boolean {
388+
if (!current) {
389+
return allowCreate;
390+
}
391+
if (!shouldPersistCurrentRunSessionCleanup(current, sessionId)) {
392+
return false;
393+
}
394+
return (
395+
current.restartRecoveryDeliveryRunId === undefined ||
396+
current.restartRecoveryDeliveryRunId === runId
397+
);
398+
}
399+
400+
function shouldPersistRestartRecoveryCleanup(
401+
current: SessionEntry | undefined,
402+
sessionId: string,
403+
runId: string,
404+
): boolean {
405+
return (
406+
shouldPersistCurrentRunSessionCleanup(current, sessionId) &&
407+
current?.restartRecoveryDeliveryRunId === runId
408+
);
409+
}
410+
289411
function containsControlCharacters(value: string): boolean {
290412
for (const char of value) {
291413
const code = char.codePointAt(0);
@@ -607,6 +729,8 @@ async function agentCommandInternal(
607729
} = prepared;
608730
const effectiveCwd = cwd ? resolveUserPath(cwd) : workspaceDir;
609731
let sessionEntry = prepared.sessionEntry;
732+
let trackedRestartRecoveryDeliveryContext = false;
733+
let currentRunDeliveryContext: DeliveryContext | undefined;
610734

611735
try {
612736
if (opts.deliver === true) {
@@ -626,6 +750,49 @@ async function agentCommandInternal(
626750
throw acpResolution.error;
627751
}
628752

753+
if (
754+
sessionStore &&
755+
sessionKey &&
756+
!suppressVisibleSessionEffects &&
757+
!isSubagentSessionKey(sessionKey)
758+
) {
759+
const now = Date.now();
760+
const currentStoreEntry = sessionStore[sessionKey];
761+
const allowCreateRestartRecoveryEntry =
762+
currentStoreEntry === undefined && sessionEntry === undefined;
763+
const entry = currentStoreEntry ??
764+
sessionEntry ?? { sessionId, updatedAt: now, sessionStartedAt: now };
765+
currentRunDeliveryContext = await resolveCurrentRunDeliveryContext({
766+
cfg,
767+
opts,
768+
sessionEntry: entry,
769+
});
770+
const next: SessionEntry = {
771+
...entry,
772+
sessionId,
773+
updatedAt: now,
774+
restartRecoveryDeliveryContext: currentRunDeliveryContext,
775+
restartRecoveryDeliveryRunId: currentRunDeliveryContext ? runId : undefined,
776+
};
777+
const persisted = await persistSessionEntry({
778+
sessionStore,
779+
sessionKey,
780+
storePath,
781+
entry: next,
782+
shouldPersist: (current) =>
783+
shouldPersistRestartRecoveryContextClaim(
784+
current,
785+
sessionId,
786+
runId,
787+
allowCreateRestartRecoveryEntry,
788+
),
789+
});
790+
sessionEntry = persisted ?? sessionEntry;
791+
trackedRestartRecoveryDeliveryContext =
792+
Boolean(persisted?.restartRecoveryDeliveryContext) &&
793+
persisted?.restartRecoveryDeliveryRunId === runId;
794+
}
795+
629796
if (!isRawModelRun && acpResolution?.kind === "ready" && sessionKey) {
630797
const attemptExecutionRuntime = await loadAttemptExecutionRuntime();
631798
const startedAt = Date.now();
@@ -1730,16 +1897,18 @@ async function agentCommandInternal(
17301897
...entry,
17311898
pendingFinalDelivery: true,
17321899
pendingFinalDeliveryText: combinedPayload,
1900+
pendingFinalDeliveryContext: currentRunDeliveryContext,
17331901
pendingFinalDeliveryCreatedAt: now,
17341902
updatedAt: now,
17351903
};
1736-
await persistSessionEntry({
1904+
const persisted = await persistSessionEntry({
17371905
sessionStore,
17381906
sessionKey,
17391907
storePath,
17401908
entry: next,
1909+
shouldPersist: (current) => shouldPersistCurrentRunSessionCleanup(current, sessionId),
17411910
});
1742-
sessionEntry = next;
1911+
sessionEntry = persisted ?? sessionEntry;
17431912
}
17441913
}
17451914

@@ -1795,13 +1964,14 @@ async function agentCommandInternal(
17951964
!entry.pendingFinalDeliveryText;
17961965
if (deliveryResult?.deliverySucceeded === true || noPendingTextForThisRun) {
17971966
const next = clearPendingFinalDeliveryFields(entry, Date.now());
1798-
await persistSessionEntry({
1967+
const persisted = await persistSessionEntry({
17991968
sessionStore,
18001969
sessionKey,
18011970
storePath,
18021971
entry: next,
1972+
shouldPersist: (current) => shouldPersistCurrentRunSessionCleanup(current, sessionId),
18031973
});
1804-
sessionEntry = next;
1974+
sessionEntry = persisted ?? sessionEntry;
18051975
}
18061976
}
18071977

@@ -1812,6 +1982,34 @@ async function agentCommandInternal(
18121982
throw error;
18131983
}
18141984
} finally {
1985+
if (trackedRestartRecoveryDeliveryContext && sessionStore && sessionKey) {
1986+
try {
1987+
const entry = sessionStore[sessionKey] ?? sessionEntry;
1988+
if (entry?.restartRecoveryDeliveryContext && entry.restartRecoveryDeliveryRunId === runId) {
1989+
const next: SessionEntry = {
1990+
...entry,
1991+
restartRecoveryDeliveryContext: undefined,
1992+
restartRecoveryDeliveryRunId: undefined,
1993+
updatedAt: Date.now(),
1994+
};
1995+
const persisted = await persistSessionEntry({
1996+
sessionStore,
1997+
sessionKey,
1998+
storePath,
1999+
entry: next,
2000+
shouldPersist: (current) =>
2001+
shouldPersistRestartRecoveryCleanup(current, sessionId, runId),
2002+
});
2003+
sessionEntry = persisted ?? sessionEntry;
2004+
}
2005+
} catch (error) {
2006+
log.warn(
2007+
`failed to clear restart recovery delivery context for ${sessionKey}: ${
2008+
error instanceof Error ? error.message : String(error)
2009+
}`,
2010+
);
2011+
}
2012+
}
18152013
clearAgentRunContext(runId);
18162014
}
18172015
}

0 commit comments

Comments
 (0)