Skip to content

Commit 887aa2c

Browse files
committed
fix(heartbeat): bound pending final delivery replay
1 parent e3b77d6 commit 887aa2c

2 files changed

Lines changed: 172 additions & 49 deletions

File tree

src/auto-reply/reply/get-reply.fast-path.test.ts

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -315,6 +315,86 @@ describe("getReplyFromConfig fast test bootstrap", () => {
315315
expect(stored.pendingFinalDeliveryAttemptCount).toBe(1);
316316
});
317317

318+
it("clears stale heartbeat pending delivery after the replay attempt limit", async () => {
319+
const home = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-heartbeat-pending-limit-"));
320+
const storePath = path.join(home, "sessions.json");
321+
const sessionKey = "agent:main:telegram:123";
322+
await fs.writeFile(
323+
storePath,
324+
JSON.stringify({
325+
[sessionKey]: {
326+
sessionId: "pending-limit",
327+
updatedAt: Date.now(),
328+
pendingFinalDelivery: true,
329+
pendingFinalDeliveryText: "HEARTBEAT_OK short",
330+
pendingFinalDeliveryCreatedAt: Date.now(),
331+
pendingFinalDeliveryAttemptCount: 10,
332+
pendingFinalDeliveryLastError: null,
333+
},
334+
}),
335+
"utf8",
336+
);
337+
const cfg = withFastReplyConfig({
338+
agents: {
339+
defaults: {
340+
model: "openai/gpt-5.5",
341+
workspace: home,
342+
heartbeat: { ackMaxChars: 0 },
343+
},
344+
},
345+
session: { store: storePath },
346+
} as OpenClawConfig);
347+
348+
await expect(
349+
getReplyFromConfig(buildGetReplyCtx(), { isHeartbeat: true }, cfg),
350+
).resolves.toEqual({ text: "ok" });
351+
352+
const stored = JSON.parse(await fs.readFile(storePath, "utf8"))[sessionKey];
353+
expect(stored.pendingFinalDelivery).toBeUndefined();
354+
expect(stored.pendingFinalDeliveryText).toBeUndefined();
355+
expect(stored.pendingFinalDeliveryAttemptCount).toBeUndefined();
356+
});
357+
358+
it("clears stale heartbeat pending delivery after the replay TTL", async () => {
359+
const home = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-heartbeat-pending-ttl-"));
360+
const storePath = path.join(home, "sessions.json");
361+
const sessionKey = "agent:main:telegram:123";
362+
await fs.writeFile(
363+
storePath,
364+
JSON.stringify({
365+
[sessionKey]: {
366+
sessionId: "pending-expired",
367+
updatedAt: Date.now(),
368+
pendingFinalDelivery: true,
369+
pendingFinalDeliveryText: "HEARTBEAT_OK short",
370+
pendingFinalDeliveryCreatedAt: Date.now() - 24 * 60 * 60 * 1000 - 1,
371+
pendingFinalDeliveryAttemptCount: 2,
372+
pendingFinalDeliveryLastError: null,
373+
},
374+
}),
375+
"utf8",
376+
);
377+
const cfg = withFastReplyConfig({
378+
agents: {
379+
defaults: {
380+
model: "openai/gpt-5.5",
381+
workspace: home,
382+
heartbeat: { ackMaxChars: 0 },
383+
},
384+
},
385+
session: { store: storePath },
386+
} as OpenClawConfig);
387+
388+
await expect(
389+
getReplyFromConfig(buildGetReplyCtx(), { isHeartbeat: true }, cfg),
390+
).resolves.toEqual({ text: "ok" });
391+
392+
const stored = JSON.parse(await fs.readFile(storePath, "utf8"))[sessionKey];
393+
expect(stored.pendingFinalDelivery).toBeUndefined();
394+
expect(stored.pendingFinalDeliveryText).toBeUndefined();
395+
expect(stored.pendingFinalDeliveryAttemptCount).toBeUndefined();
396+
});
397+
318398
it("sanitizes stale heartbeat pending delivery before replay", async () => {
319399
const home = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-heartbeat-pending-sanitize-"));
320400
const storePath = path.join(home, "sessions.json");

src/auto-reply/reply/get-reply.ts

Lines changed: 92 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,9 @@ import { createTypingController } from "./typing.js";
5656

5757
type 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+
5962
function 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+
80105
const 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

Comments
 (0)