From 13da4e7f9b5ec0578cce427522fd785f0d924cb0 Mon Sep 17 00:00:00 2001 From: Seth Date: Thu, 25 Jun 2026 14:00:38 -0700 Subject: [PATCH 1/4] Defer heartbeat cron jobs while sessions are active --- packages/coding-agent/src/core/cron-jobs.ts | 16 +++++ .../src/modes/daemon/daemon-mode.ts | 13 ++-- packages/coding-agent/test/cron-jobs.test.ts | 64 +++++++++++++++++++ .../coding-agent/test/daemon-mode.test.ts | 38 +++++++---- 4 files changed, 111 insertions(+), 20 deletions(-) diff --git a/packages/coding-agent/src/core/cron-jobs.ts b/packages/coding-agent/src/core/cron-jobs.ts index 95bf322bac..0ec3d4458a 100644 --- a/packages/coding-agent/src/core/cron-jobs.ts +++ b/packages/coding-agent/src/core/cron-jobs.ts @@ -58,6 +58,12 @@ export interface AgentCronSchedulerHooks { onError?: (job: AgentCronJob, error: unknown) => void; } +export interface HeartbeatCronSessionActivity { + isStreaming: boolean; + isBashRunning: boolean; + pendingMessageCount: number; +} + interface CronJobsFile { jobs?: unknown; } @@ -814,6 +820,16 @@ function consumeLeadingEverySchedule(text: string): { interval: string; rest: st }; } +export function isHeartbeatCronJob(job: AgentCronJob): boolean { + return job.source === "heartbeat" || job.source === "rlm_heartbeat"; +} + +export function shouldDeferHeartbeatCronJob(job: AgentCronJob, activity: HeartbeatCronSessionActivity): boolean { + return ( + isHeartbeatCronJob(job) && (activity.isStreaming || activity.isBashRunning || activity.pendingMessageCount > 0) + ); +} + function nextCronRunAfter(expression: string, after: Date): Date { const fields = parseCronExpression(expression); const candidate = new Date(after.getTime()); diff --git a/packages/coding-agent/src/modes/daemon/daemon-mode.ts b/packages/coding-agent/src/modes/daemon/daemon-mode.ts index 0461285d2a..aae2cd36ce 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-mode.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-mode.ts @@ -22,7 +22,9 @@ import { type AgentHeartbeatUpdateAction, createAgentHeartbeatToolDefinitions, DEFAULT_HEARTBEAT_SCHEDULE, + isHeartbeatCronJob, normalizeHeartbeatSchedule, + shouldDeferHeartbeatCronJob, } from "../../core/cron-jobs.js"; import type { CreateRlmSubagentRuntimeOptions, @@ -471,11 +473,10 @@ export class AgentDaemon { if (!state) { return; } - const followUpQueueKey = isHeartbeatCronJob(job) ? `heartbeat:${job.id}` : undefined; - if (followUpQueueKey && (state.runtime.session.isStreaming || state.runtime.session.pendingMessageCount > 0)) { - const didQueue = await state.runtime.session.followUp(job.prompt, undefined, { queueKey: followUpQueueKey }); - return didQueue ? undefined : "skipped"; + if (shouldDeferHeartbeatCronJob(job, state.runtime.session)) { + return "skipped"; } + const followUpQueueKey = isHeartbeatCronJob(job) ? `heartbeat:${job.id}` : undefined; if (!followUpQueueKey && (state.runtime.session.isStreaming || state.runtime.session.pendingMessageCount > 0)) { await state.runtime.session.followUp(job.prompt); return; @@ -1670,10 +1671,6 @@ export class AgentDaemon { } } -function isHeartbeatCronJob(job: AgentCronJob): boolean { - return job.source === "heartbeat" || job.source === "rlm_heartbeat"; -} - function serializeSavedSessionInfo(session: SessionInfo): DaemonSavedSessionInfo { return { path: session.path, diff --git a/packages/coding-agent/test/cron-jobs.test.ts b/packages/coding-agent/test/cron-jobs.test.ts index 0f355051e3..95df22ee5e 100644 --- a/packages/coding-agent/test/cron-jobs.test.ts +++ b/packages/coding-agent/test/cron-jobs.test.ts @@ -9,6 +9,7 @@ import { createAgentHeartbeatToolDefinitions, parseAgentCronSchedule, parseHeartbeatCommand, + shouldDeferHeartbeatCronJob, } from "../src/core/cron-jobs.js"; const start = new Date("2026-01-01T12:34:00.000Z"); @@ -752,6 +753,69 @@ describe("AgentCronScheduler", () => { }); }); +describe("shouldDeferHeartbeatCronJob", () => { + const baseJob: AgentCronJob = { + id: "job-1", + status: "active", + activeSessionId: "active-1", + sessionId: "session-1", + sessionFile: "/tmp/session.jsonl", + cwd: "/tmp/project", + prompt: "check progress", + schedule: { kind: "interval", expression: "every 5m", intervalMs: 300_000 }, + createdAt: "2026-01-01T12:34:00.000Z", + updatedAt: "2026-01-01T12:34:00.000Z", + nextRunAt: "2026-01-01T12:39:00.000Z", + runCount: 0, + }; + + it("defers user and RLM heartbeats while the target session is working", () => { + for (const source of ["heartbeat", "rlm_heartbeat"] as const) { + const job = { ...baseJob, source }; + + expect( + shouldDeferHeartbeatCronJob(job, { + isStreaming: true, + isBashRunning: false, + pendingMessageCount: 0, + }), + ).toBe(true); + expect( + shouldDeferHeartbeatCronJob(job, { + isStreaming: false, + isBashRunning: true, + pendingMessageCount: 0, + }), + ).toBe(true); + expect( + shouldDeferHeartbeatCronJob(job, { + isStreaming: false, + isBashRunning: false, + pendingMessageCount: 1, + }), + ).toBe(true); + } + }); + + it("allows heartbeats when the target session is idle", () => { + expect( + shouldDeferHeartbeatCronJob( + { ...baseJob, source: "heartbeat" }, + { isStreaming: false, isBashRunning: false, pendingMessageCount: 0 }, + ), + ).toBe(false); + }); + + it("does not defer ordinary cron jobs", () => { + expect( + shouldDeferHeartbeatCronJob( + { ...baseJob, source: "cron" }, + { isStreaming: true, isBashRunning: true, pendingMessageCount: 2 }, + ), + ).toBe(false); + }); +}); + describe("createAgentHeartbeatToolDefinitions", () => { it("exposes only read-only user heartbeat inspection to the model", () => { const tools = createAgentHeartbeatToolDefinitions({ diff --git a/packages/coding-agent/test/daemon-mode.test.ts b/packages/coding-agent/test/daemon-mode.test.ts index 28514dbe4b..39634272a1 100644 --- a/packages/coding-agent/test/daemon-mode.test.ts +++ b/packages/coding-agent/test/daemon-mode.test.ts @@ -181,7 +181,7 @@ describe("daemon mode helpers", () => { } }); - it("queues busy heartbeat cron jobs with a per-job coalescing key", async () => { + it("defers busy heartbeat cron jobs instead of queueing a follow-up", async () => { const daemon = new AgentDaemon("/tmp/prime-agent-test.sock", { defaultSessionConfig: { agentDir: "/tmp/prime-agent-test-agent", @@ -197,6 +197,7 @@ describe("daemon mode helpers", () => { runtime: ActiveSessionState["runtime"] & { session: { isStreaming: boolean; + isBashRunning: boolean; pendingMessageCount: number; prompt: typeof prompt; followUp: typeof followUp; @@ -205,6 +206,7 @@ describe("daemon mode helpers", () => { }; state.runtime.session = { isStreaming: true, + isBashRunning: false, pendingMessageCount: 0, prompt, followUp, @@ -216,17 +218,20 @@ describe("daemon mode helpers", () => { ).sessions.set(state.activeSessionId, state); const runCronJob = ( daemon as unknown as { - runCronJob(job: AgentCronJob): Promise; + runCronJob(job: AgentCronJob): Promise<"skipped" | undefined>; } ).runCronJob.bind(daemon); - await runCronJob(makeCronJob({ id: "heartbeat-1", source: "heartbeat", activeSessionId: state.activeSessionId })); + const result = await runCronJob( + makeCronJob({ id: "heartbeat-1", source: "heartbeat", activeSessionId: state.activeSessionId }), + ); - expect(followUp).toHaveBeenCalledWith("heartbeat prompt", undefined, { queueKey: "heartbeat:heartbeat-1" }); + expect(result).toBe("skipped"); + expect(followUp).not.toHaveBeenCalled(); expect(prompt).not.toHaveBeenCalled(); }); - it("uses separate queue keys for separate RLM heartbeat cron jobs", async () => { + it("defers separate RLM heartbeat cron jobs while the session is busy", async () => { const daemon = new AgentDaemon("/tmp/prime-agent-test.sock", { defaultSessionConfig: { agentDir: "/tmp/prime-agent-test-agent", @@ -242,6 +247,7 @@ describe("daemon mode helpers", () => { runtime: ActiveSessionState["runtime"] & { session: { isStreaming: boolean; + isBashRunning: boolean; pendingMessageCount: number; prompt: typeof prompt; followUp: typeof followUp; @@ -250,6 +256,7 @@ describe("daemon mode helpers", () => { }; state.runtime.session = { isStreaming: true, + isBashRunning: false, pendingMessageCount: 0, prompt, followUp, @@ -261,19 +268,24 @@ describe("daemon mode helpers", () => { ).sessions.set(state.activeSessionId, state); const runCronJob = ( daemon as unknown as { - runCronJob(job: AgentCronJob): Promise; + runCronJob(job: AgentCronJob): Promise<"skipped" | undefined>; } ).runCronJob.bind(daemon); - await runCronJob(makeCronJob({ id: "rlm-1", source: "rlm_heartbeat", activeSessionId: state.activeSessionId })); - await runCronJob(makeCronJob({ id: "rlm-2", source: "rlm_heartbeat", activeSessionId: state.activeSessionId })); + const first = await runCronJob( + makeCronJob({ id: "rlm-1", source: "rlm_heartbeat", activeSessionId: state.activeSessionId }), + ); + const second = await runCronJob( + makeCronJob({ id: "rlm-2", source: "rlm_heartbeat", activeSessionId: state.activeSessionId }), + ); - expect(followUp).toHaveBeenNthCalledWith(1, "heartbeat prompt", undefined, { queueKey: "heartbeat:rlm-1" }); - expect(followUp).toHaveBeenNthCalledWith(2, "heartbeat prompt", undefined, { queueKey: "heartbeat:rlm-2" }); + expect(first).toBe("skipped"); + expect(second).toBe("skipped"); + expect(followUp).not.toHaveBeenCalled(); expect(prompt).not.toHaveBeenCalled(); }); - it("skips duplicate queued heartbeat cron jobs", async () => { + it("does not enqueue another heartbeat when one is already pending", async () => { const daemon = new AgentDaemon("/tmp/prime-agent-test.sock", { defaultSessionConfig: { agentDir: "/tmp/prime-agent-test-agent", cwd: "/tmp" }, createRuntime: async () => { @@ -284,6 +296,7 @@ describe("daemon mode helpers", () => { runtime: ActiveSessionState["runtime"] & { session: { isStreaming: boolean; + isBashRunning: boolean; pendingMessageCount: number; prompt: ReturnType; followUp: ReturnType; @@ -294,6 +307,7 @@ describe("daemon mode helpers", () => { const removeQueuedFollowUp = vi.fn(() => true); state.runtime.session = { isStreaming: true, + isBashRunning: false, pendingMessageCount: 1, prompt: vi.fn(), followUp, @@ -305,7 +319,7 @@ describe("daemon mode helpers", () => { ).runCronJob(makeCronJob({ id: "heartbeat-1", source: "heartbeat", activeSessionId: state.activeSessionId })); expect(result).toBe("skipped"); - expect(followUp).toHaveBeenCalledWith("heartbeat prompt", undefined, { queueKey: "heartbeat:heartbeat-1" }); + expect(followUp).not.toHaveBeenCalled(); expect(removeQueuedFollowUp).not.toHaveBeenCalled(); }); From 5a73a9baaa202ac60df19fe4fcf262690c7e550f Mon Sep 17 00:00:00 2001 From: Kevin Thomas Date: Mon, 29 Jun 2026 18:47:10 +0000 Subject: [PATCH 2/4] remove dead heartbeat follow-up plumbing and surface deferred skips --- packages/coding-agent/src/core/cron-jobs.ts | 11 +++++++++-- .../coding-agent/src/modes/daemon/daemon-mode.ts | 12 +++++------- packages/coding-agent/test/cron-jobs.test.ts | 1 + 3 files changed, 15 insertions(+), 9 deletions(-) diff --git a/packages/coding-agent/src/core/cron-jobs.ts b/packages/coding-agent/src/core/cron-jobs.ts index 0ec3d4458a..457b97e0e9 100644 --- a/packages/coding-agent/src/core/cron-jobs.ts +++ b/packages/coding-agent/src/core/cron-jobs.ts @@ -33,6 +33,7 @@ export interface AgentCronJob { updatedAt: string; nextRunAt?: string; lastRunAt?: string; + lastSkippedAt?: string; lastError?: string; runCount: number; } @@ -489,7 +490,12 @@ export class AgentCronJobStore { return job; } const nextRunAt = nextRunAtForSchedule(job.schedule, now); - updated = { ...job, nextRunAt: nextRunAt?.toISOString(), updatedAt: now.toISOString() }; + updated = { + ...job, + nextRunAt: nextRunAt?.toISOString(), + lastSkippedAt: now.toISOString(), + updatedAt: now.toISOString(), + }; return updated; }); if (updated) { @@ -766,7 +772,8 @@ export function formatAgentCronJob(job: AgentCronJob): string { const preview = job.prompt.replace(/\s+/g, " ").slice(0, 80); const error = job.lastError ? ` error=${job.lastError}` : ""; const label = job.label ? ` label="${job.label}"` : ""; - return `${job.id} ${job.status}${label} next=${next} last=${last} runs=${job.runCount} schedule="${job.schedule.expression}" prompt="${preview}"${error}`; + const skipped = job.lastSkippedAt ? ` skipped=${new Date(job.lastSkippedAt).toLocaleString()}` : ""; + return `${job.id} ${job.status}${label} next=${next} last=${last}${skipped} runs=${job.runCount} schedule="${job.schedule.expression}" prompt="${preview}"${error}`; } export function createAgentHeartbeatToolDefinitions(controller: AgentCronToolController): ToolDefinition[] { diff --git a/packages/coding-agent/src/modes/daemon/daemon-mode.ts b/packages/coding-agent/src/modes/daemon/daemon-mode.ts index aae2cd36ce..a5276b65d3 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-mode.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-mode.ts @@ -476,16 +476,14 @@ export class AgentDaemon { if (shouldDeferHeartbeatCronJob(job, state.runtime.session)) { return "skipped"; } - const followUpQueueKey = isHeartbeatCronJob(job) ? `heartbeat:${job.id}` : undefined; - if (!followUpQueueKey && (state.runtime.session.isStreaming || state.runtime.session.pendingMessageCount > 0)) { + if ( + !isHeartbeatCronJob(job) && + (state.runtime.session.isStreaming || state.runtime.session.pendingMessageCount > 0) + ) { await state.runtime.session.followUp(job.prompt); return; } - await state.runtime.session.prompt(job.prompt, { - streamingBehavior: state.runtime.session.isStreaming ? "followUp" : undefined, - followUpQueueKey, - source: "rpc", - }); + await state.runtime.session.prompt(job.prompt, { source: "rpc" }); } private createCronJobForState(state: ActiveSessionState, schedule: string, prompt: string): AgentCronJob { diff --git a/packages/coding-agent/test/cron-jobs.test.ts b/packages/coding-agent/test/cron-jobs.test.ts index 95df22ee5e..fcb2d1b431 100644 --- a/packages/coding-agent/test/cron-jobs.test.ts +++ b/packages/coding-agent/test/cron-jobs.test.ts @@ -703,6 +703,7 @@ describe("AgentCronScheduler", () => { id: job.id, status: "active", nextRunAt: "2026-01-01T12:45:00.000Z", + lastSkippedAt: "2026-01-01T12:40:00.000Z", runCount: 0, }); expect(store.getHeartbeat("active-1")).not.toHaveProperty("lastRunAt"); From f4501f8caec2ce5ca2cfcaaf0cd8430def51fb3e Mon Sep 17 00:00:00 2001 From: Kevin Thomas Date: Mon, 29 Jun 2026 18:50:49 +0000 Subject: [PATCH 3/4] pass followUp streaming behavior on cron prompt to survive mid-call stream start --- .../src/modes/daemon/daemon-mode.ts | 5 ++- .../coding-agent/test/daemon-mode.test.ts | 37 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/packages/coding-agent/src/modes/daemon/daemon-mode.ts b/packages/coding-agent/src/modes/daemon/daemon-mode.ts index a5276b65d3..e113a14175 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-mode.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-mode.ts @@ -483,7 +483,10 @@ export class AgentDaemon { await state.runtime.session.followUp(job.prompt); return; } - await state.runtime.session.prompt(job.prompt, { source: "rpc" }); + await state.runtime.session.prompt(job.prompt, { + streamingBehavior: "followUp", + source: "rpc", + }); } private createCronJobForState(state: ActiveSessionState, schedule: string, prompt: string): AgentCronJob { diff --git a/packages/coding-agent/test/daemon-mode.test.ts b/packages/coding-agent/test/daemon-mode.test.ts index 39634272a1..77d3897b6e 100644 --- a/packages/coding-agent/test/daemon-mode.test.ts +++ b/packages/coding-agent/test/daemon-mode.test.ts @@ -353,6 +353,43 @@ describe("daemon mode helpers", () => { expect(prompt).not.toHaveBeenCalled(); }); + it("prompts idle sessions with a followUp streaming behavior to survive a mid-call stream start", async () => { + const daemon = new AgentDaemon("/tmp/prime-agent-test.sock", { + defaultSessionConfig: { agentDir: "/tmp/prime-agent-test-agent", cwd: "/tmp" }, + createRuntime: async () => { + throw new Error("unexpected runtime creation"); + }, + }); + const prompt = vi.fn(async () => {}); + const followUp = vi.fn(async () => true); + const state = makeState("active-1") as ActiveSessionState & { + runtime: ActiveSessionState["runtime"] & { + session: { + isStreaming: boolean; + isBashRunning: boolean; + pendingMessageCount: number; + prompt: typeof prompt; + followUp: typeof followUp; + }; + }; + }; + state.runtime.session = { + isStreaming: false, + isBashRunning: false, + pendingMessageCount: 0, + prompt, + followUp, + } as never; + (daemon as unknown as { sessions: Map }).sessions.set(state.activeSessionId, state); + + await (daemon as unknown as { runCronJob(job: AgentCronJob): Promise<"skipped" | undefined> }).runCronJob( + makeCronJob({ id: "cron-1", source: "cron", activeSessionId: state.activeSessionId }), + ); + + expect(prompt).toHaveBeenCalledWith("heartbeat prompt", { streamingBehavior: "followUp", source: "rpc" }); + expect(followUp).not.toHaveBeenCalled(); + }); + it("removes queued heartbeat follow-ups when a heartbeat is cleared", async () => { const tempDir = mkdtempSync(join(tmpdir(), "prime-agent-daemon-heartbeat-clear-")); try { From da4ff041b7770a0e4d5720a2e69c5863e6c1e94c Mon Sep 17 00:00:00 2001 From: Kevin Thomas Date: Mon, 29 Jun 2026 18:59:41 +0000 Subject: [PATCH 4/4] coalesce raced heartbeat follow-ups with a per-job queue key --- .../src/modes/daemon/daemon-mode.ts | 1 + .../coding-agent/test/daemon-mode.test.ts | 47 ++++++++++++++++++- 2 files changed, 47 insertions(+), 1 deletion(-) diff --git a/packages/coding-agent/src/modes/daemon/daemon-mode.ts b/packages/coding-agent/src/modes/daemon/daemon-mode.ts index e113a14175..8e985845ae 100644 --- a/packages/coding-agent/src/modes/daemon/daemon-mode.ts +++ b/packages/coding-agent/src/modes/daemon/daemon-mode.ts @@ -485,6 +485,7 @@ export class AgentDaemon { } await state.runtime.session.prompt(job.prompt, { streamingBehavior: "followUp", + followUpQueueKey: isHeartbeatCronJob(job) ? `heartbeat:${job.id}` : undefined, source: "rpc", }); } diff --git a/packages/coding-agent/test/daemon-mode.test.ts b/packages/coding-agent/test/daemon-mode.test.ts index 77d3897b6e..87ba96faa7 100644 --- a/packages/coding-agent/test/daemon-mode.test.ts +++ b/packages/coding-agent/test/daemon-mode.test.ts @@ -382,11 +382,56 @@ describe("daemon mode helpers", () => { } as never; (daemon as unknown as { sessions: Map }).sessions.set(state.activeSessionId, state); + await (daemon as unknown as { runCronJob(job: AgentCronJob): Promise<"skipped" | undefined> }).runCronJob( + makeCronJob({ id: "heartbeat-1", source: "heartbeat", activeSessionId: state.activeSessionId }), + ); + + expect(prompt).toHaveBeenCalledWith("heartbeat prompt", { + streamingBehavior: "followUp", + followUpQueueKey: "heartbeat:heartbeat-1", + source: "rpc", + }); + expect(followUp).not.toHaveBeenCalled(); + }); + + it("prompts idle generic cron jobs without a heartbeat coalescing key", async () => { + const daemon = new AgentDaemon("/tmp/prime-agent-test.sock", { + defaultSessionConfig: { agentDir: "/tmp/prime-agent-test-agent", cwd: "/tmp" }, + createRuntime: async () => { + throw new Error("unexpected runtime creation"); + }, + }); + const prompt = vi.fn(async () => {}); + const followUp = vi.fn(async () => true); + const state = makeState("active-1") as ActiveSessionState & { + runtime: ActiveSessionState["runtime"] & { + session: { + isStreaming: boolean; + isBashRunning: boolean; + pendingMessageCount: number; + prompt: typeof prompt; + followUp: typeof followUp; + }; + }; + }; + state.runtime.session = { + isStreaming: false, + isBashRunning: false, + pendingMessageCount: 0, + prompt, + followUp, + } as never; + (daemon as unknown as { sessions: Map }).sessions.set(state.activeSessionId, state); + await (daemon as unknown as { runCronJob(job: AgentCronJob): Promise<"skipped" | undefined> }).runCronJob( makeCronJob({ id: "cron-1", source: "cron", activeSessionId: state.activeSessionId }), ); - expect(prompt).toHaveBeenCalledWith("heartbeat prompt", { streamingBehavior: "followUp", source: "rpc" }); + expect(prompt).toHaveBeenCalledWith("heartbeat prompt", { + streamingBehavior: "followUp", + followUpQueueKey: undefined, + source: "rpc", + }); expect(followUp).not.toHaveBeenCalled(); });