@@ -56,6 +56,9 @@ import { createTypingController } from "./typing.js";
5656
5757type ResetCommandAction = "new" | "reset" ;
5858
59+ const PENDING_FINAL_DELIVERY_MAX_HEARTBEAT_ATTEMPTS = 10 ;
60+ const PENDING_FINAL_DELIVERY_HEARTBEAT_TTL_MS = 24 * 60 * 60 * 1000 ;
61+
5962function classifyHeartbeatPendingFinalDelivery ( text : string , ackMaxChars : number ) {
6063 const stripped = stripHeartbeatToken ( text , {
6164 mode : "heartbeat" ,
@@ -77,6 +80,28 @@ function resolveHeartbeatAckMaxChars(cfg: OpenClawConfig, agentId: string): numb
7780 ) ;
7881}
7982
83+ function resolveHeartbeatPendingFinalDeliveryGiveUp ( args : {
84+ attemptCount : number ;
85+ now : number ;
86+ createdAt ?: number ;
87+ } ) {
88+ const ageMs =
89+ typeof args . createdAt === "number" ? Math . max ( 0 , args . now - args . createdAt ) : undefined ;
90+ if ( args . attemptCount > PENDING_FINAL_DELIVERY_MAX_HEARTBEAT_ATTEMPTS ) {
91+ return {
92+ reason : "retry-limit" as const ,
93+ ageMs,
94+ } ;
95+ }
96+ if ( typeof ageMs === "number" && ageMs > PENDING_FINAL_DELIVERY_HEARTBEAT_TTL_MS ) {
97+ return {
98+ reason : "expiry" as const ,
99+ ageMs,
100+ } ;
101+ }
102+ return null ;
103+ }
104+
80105const sessionResetModelRuntimeLoader = createLazyImportLoader (
81106 ( ) => import ( "./session-reset-model.runtime.js" ) ,
82107) ;
@@ -412,6 +437,34 @@ export async function getReplyFromConfig(
412437
413438 if ( sessionEntry ?. pendingFinalDelivery && sessionEntry . pendingFinalDeliveryText ) {
414439 const text = sanitizePendingFinalDeliveryText ( sessionEntry . pendingFinalDeliveryText ) ;
440+ const clearPendingFinalDelivery = async ( ) => {
441+ sessionEntry . pendingFinalDelivery = undefined ;
442+ sessionEntry . pendingFinalDeliveryText = undefined ;
443+ sessionEntry . pendingFinalDeliveryCreatedAt = undefined ;
444+ sessionEntry . pendingFinalDeliveryLastAttemptAt = undefined ;
445+ sessionEntry . pendingFinalDeliveryAttemptCount = undefined ;
446+ sessionEntry . pendingFinalDeliveryLastError = undefined ;
447+ sessionEntry . pendingFinalDeliveryContext = undefined ;
448+ if ( sessionKey && sessionStore ) {
449+ sessionStore [ sessionKey ] = sessionEntry ;
450+ }
451+ if ( sessionKey && storePath ) {
452+ const { updateSessionStoreEntry } = await import ( "../../config/sessions.js" ) ;
453+ await updateSessionStoreEntry ( {
454+ storePath,
455+ sessionKey,
456+ update : async ( ) => ( {
457+ pendingFinalDelivery : undefined ,
458+ pendingFinalDeliveryText : undefined ,
459+ pendingFinalDeliveryCreatedAt : undefined ,
460+ pendingFinalDeliveryLastAttemptAt : undefined ,
461+ pendingFinalDeliveryAttemptCount : undefined ,
462+ pendingFinalDeliveryLastError : undefined ,
463+ pendingFinalDeliveryContext : undefined ,
464+ } ) ,
465+ } ) ;
466+ }
467+ } ;
415468
416469 // If it's a heartbeat, we definitely want to try delivering the lost reply now.
417470 // If it's a user message, we deliver the lost reply first, then continue.
@@ -422,59 +475,49 @@ export async function getReplyFromConfig(
422475 resolveHeartbeatAckMaxChars ( cfg , agentId ) ,
423476 ) ;
424477 if ( heartbeatPending . shouldClear ) {
425- sessionEntry . pendingFinalDelivery = undefined ;
426- sessionEntry . pendingFinalDeliveryText = undefined ;
427- sessionEntry . pendingFinalDeliveryCreatedAt = undefined ;
428- sessionEntry . pendingFinalDeliveryLastAttemptAt = undefined ;
429- sessionEntry . pendingFinalDeliveryAttemptCount = undefined ;
430- sessionEntry . pendingFinalDeliveryLastError = undefined ;
431- sessionEntry . pendingFinalDeliveryContext = undefined ;
432- if ( sessionKey && sessionStore ) {
433- sessionStore [ sessionKey ] = sessionEntry ;
434- }
435- if ( sessionKey && storePath ) {
436- const { updateSessionStoreEntry } = await import ( "../../config/sessions.js" ) ;
437- await updateSessionStoreEntry ( {
438- storePath,
439- sessionKey,
440- update : async ( ) => ( {
441- pendingFinalDelivery : undefined ,
442- pendingFinalDeliveryText : undefined ,
443- pendingFinalDeliveryCreatedAt : undefined ,
444- pendingFinalDeliveryLastAttemptAt : undefined ,
445- pendingFinalDeliveryAttemptCount : undefined ,
446- pendingFinalDeliveryLastError : undefined ,
447- pendingFinalDeliveryContext : undefined ,
448- } ) ,
449- } ) ;
450- }
478+ await clearPendingFinalDelivery ( ) ;
451479 } else {
452480 const updatedAt = Date . now ( ) ;
453481 const attemptCount = ( sessionEntry . pendingFinalDeliveryAttemptCount ?? 0 ) + 1 ;
454- sessionEntry . pendingFinalDeliveryLastAttemptAt = updatedAt ;
455- sessionEntry . pendingFinalDeliveryAttemptCount = attemptCount ;
456- sessionEntry . pendingFinalDeliveryLastError = null ;
457- const replayText = sanitizePendingFinalDeliveryText ( heartbeatPending . replayText ) ;
458- sessionEntry . pendingFinalDeliveryText = replayText ;
459- sessionEntry . updatedAt = updatedAt ;
460- if ( sessionKey && sessionStore ) {
461- sessionStore [ sessionKey ] = sessionEntry ;
462- }
463- if ( sessionKey && storePath ) {
464- const { updateSessionStoreEntry } = await import ( "../../config/sessions.js" ) ;
465- await updateSessionStoreEntry ( {
466- storePath,
467- sessionKey,
468- update : async ( ) => ( {
469- pendingFinalDeliveryText : replayText ,
470- pendingFinalDeliveryLastAttemptAt : updatedAt ,
471- pendingFinalDeliveryAttemptCount : attemptCount ,
472- pendingFinalDeliveryLastError : null ,
473- updatedAt,
474- } ) ,
475- } ) ;
482+ const giveUp = resolveHeartbeatPendingFinalDeliveryGiveUp ( {
483+ attemptCount,
484+ now : updatedAt ,
485+ createdAt : sessionEntry . pendingFinalDeliveryCreatedAt ,
486+ } ) ;
487+ if ( giveUp ) {
488+ const replayAgeSuffix = typeof giveUp . ageMs === "number" ? ` ageMs=${ giveUp . ageMs } ` : "" ;
489+ defaultRuntime . log (
490+ `[warn] Clearing stale heartbeat pending final delivery (${ giveUp . reason } ) after ${ attemptCount } attempts${ replayAgeSuffix } ${
491+ sessionKey ? ` for session ${ sessionKey } ` : ""
492+ } .`,
493+ ) ;
494+ await clearPendingFinalDelivery ( ) ;
495+ } else {
496+ sessionEntry . pendingFinalDeliveryLastAttemptAt = updatedAt ;
497+ sessionEntry . pendingFinalDeliveryAttemptCount = attemptCount ;
498+ sessionEntry . pendingFinalDeliveryLastError = null ;
499+ const replayText = sanitizePendingFinalDeliveryText ( heartbeatPending . replayText ) ;
500+ sessionEntry . pendingFinalDeliveryText = replayText ;
501+ sessionEntry . updatedAt = updatedAt ;
502+ if ( sessionKey && sessionStore ) {
503+ sessionStore [ sessionKey ] = sessionEntry ;
504+ }
505+ if ( sessionKey && storePath ) {
506+ const { updateSessionStoreEntry } = await import ( "../../config/sessions.js" ) ;
507+ await updateSessionStoreEntry ( {
508+ storePath,
509+ sessionKey,
510+ update : async ( ) => ( {
511+ pendingFinalDeliveryText : replayText ,
512+ pendingFinalDeliveryLastAttemptAt : updatedAt ,
513+ pendingFinalDeliveryAttemptCount : attemptCount ,
514+ pendingFinalDeliveryLastError : null ,
515+ updatedAt,
516+ } ) ,
517+ } ) ;
518+ }
519+ return { text : replayText } ;
476520 }
477- return { text : replayText } ;
478521 }
479522 }
480523 }
0 commit comments