@@ -21,6 +21,11 @@ import {
2121 registerAgentRunContext ,
2222} from "../infra/agent-events.js" ;
2323import { 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" ;
2429import { buildOutboundSessionContext } from "../infra/outbound/session-context.js" ;
2530import { parseStrictNonNegativeInteger } from "../infra/parse-finite-number.js" ;
2631import { createSubsystemLogger } from "../logging/subsystem.js" ;
@@ -47,7 +52,15 @@ import type { getRemoteSkillEligibility } from "../skills/runtime/remote.js";
4752import type { resolveReusableWorkspaceSkillSnapshot } from "../skills/runtime/session-snapshot.js" ;
4853import { createTrajectoryRuntimeRecorder } from "../trajectory/runtime.js" ;
4954import { 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" ;
5164import { resolveAgentRuntimeConfig } from "./agent-runtime-config.js" ;
5265import {
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+
289411function 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