How Scheduled Prompts Work with Prime Agent Worker Recovery: Complete Technical Guide
Prime Agent persists scheduled prompts across worker crashes by combining a durable cron job store with a JSON-L recovery journal that records in-flight operations and replays them on restart.
The Prime Agent codebase implements a robust daemon mode that handles recurring AI prompts (often called heartbeat or cron jobs) with full crash recovery. When running as a supervised worker, the system guarantees that scheduled prompts survive process restarts without duplication or data loss. This article examines the exact mechanism using source code from the PrimeIntellect-ai/prime-agent repository.
Architecture Overview: Three Layers of Durability
Prime Agent's scheduled prompt system operates across three coordinated layers:
| Layer | Component | Persistence Strategy |
|---|---|---|
| Schedule Definition | AgentCronJobStore |
JSON file on disk (cron-jobs-path) |
| Execution State | WorkerRecoveryJournal |
Append-only JSON-L file (worker.recovery.jsonl) |
| Runtime Coordination | AgentCronScheduler |
In-memory queue with disk-backed recovery |
This design separates what should run (the schedule) from what is currently running (the recovery state).
How Scheduled Prompts Are Created
Scheduled prompts enter the system through RPC commands. In packages/coding-agent/src/modes/daemon/daemon-mode.ts, the createCronJobForState method handles the conversion from client request to durable job.
The method signature and core logic appear around lines 1844–1856:
private createCronJobForState(
state: ActiveSessionState,
schedule: string,
prompt: string,
): AgentCronJob {
const job: AgentCronJob = {
id: randomUUID(),
schedule: parseAgentCronSchedule(schedule),
prompt,
createdAt: Date.now(),
};
this.cronJobStore.add(job);
this.recordWorkerRecoveryState(state, "create_job");
return job;
}
The AgentCronJob interface is defined in packages/coding-agent/src/core/cron-jobs.ts (lines 45–62). It stores:
id: UUID for deduplicationschedule: Parsed cron expression with interval calculationprompt: The actual text sent to the AI sessioncreatedAt: Timestamp for debugging
The Cron Scheduler: Driving Execution
The AgentCronScheduler class in packages/coding-agent/src/core/cron-jobs.ts (lines 228–371) evaluates all registered jobs every second. When a job becomes due, it invokes runCronJob via callback.
The scheduler maintains no internal persistence—it relies entirely on the AgentCronJobStore for job definitions and the WorkerRecoveryJournal for execution state.
// From cron-jobs.ts - simplified scheduler loop
public start(): void {
this.interval = setInterval(() => {
const now = Date.now();
for (const job of this.store.getAll()) {
if (this.isDue(job, now)) {
this.onRunJob(job); // Callback to daemon-mode.ts
}
}
}, 1000);
}
Recording Worker Recovery State
The critical bridge between scheduled prompts and crash recovery is recordWorkerRecoveryState. This method in daemon-mode.ts (around lines 7271–7285) creates a durable record before any state-changing operation.
private recordWorkerRecoveryState(
state: ActiveSessionState,
operation: string,
): void {
if (!this.recoveryJournal) return;
const record: WorkerRecoveryRecord = {
timestamp: Date.now(),
activeSessionId: state.activeSessionId,
sessionId: state.runtime.session.sessionId,
sessionFile: state.runtime.session.sessionFile,
busy: true,
operation,
};
this.recoveryJournal.record(record);
}
Key parameters in WorkerRecoveryRecord:
busy:truewhen work is in-progress,falsewhen completeoperation: String identifier like"run_job","session_end","create_job"sessionId: Links the record to a specific AI sessionsessionFile: Path to session state for full reconstruction
The Recovery Journal: JSON-L Implementation
The WorkerRecoveryJournal class in packages/coding-agent/src/modes/daemon/worker-recovery-journal.ts implements an append-only log with automatic compaction.
Recording entries (lines 66–84):
public record(entry: WorkerRecoveryRecord): void {
const line = JSON.stringify(entry) + '\n';
fs.appendFileSync(this.path, line);
// Track latest state per session
this.latestBySession.set(entry.sessionId, entry);
// Compact if all entries are idle
if (this.allEntriesIdle()) {
this.compact();
}
}
Reading latest state (lines 91–93):
public readLatest(): Map<string, WorkerRecoveryRecord> {
return this.latestBySession;
}
The readLatest method returns a map from sessionId to most recent record, enabling O(1) lookup of recovery state on restart.
Complete Recovery Flow: From Crash to Resume
When a worker process crashes and restarts, this sequence ensures scheduled prompts resume correctly:
-
Supervisor restarts worker with
DAEMON_WORKER_RECOVERY_JOURNAL_ENVpointing to the journal file -
Daemon constructor (
daemon-mode.tslines 40–44) initializes:const journalPath = process.env.DAEMON_WORKER_RECOVERY_JOURNAL_ENV; if (journalPath) { this.recoveryJournal = new WorkerRecoveryJournal(journalPath); this.recoveryJournal.open(); } -
Load cron jobs from
AgentCronJobStoredisk file—schedules are fully restored -
Replay recovery journal by iterating
recoveryJournal.readLatest():- Jobs with
busy: falseare considered complete - Jobs with
busy: trueare re-queued for execution
- Jobs with
-
Scheduler starts, picking up where execution left off
Exactly-Once Semantics for In-Flight Prompts
The busy flag prevents duplicate execution of scheduled prompts that were mid-flight during a crash. Consider this scenario:
| Time | Event | Journal State |
|---|---|---|
| T0 | Cron job becomes due | — |
| T1 | recordWorkerRecoveryState(..., "run_job") called |
busy: true |
| T2 | Prompt sent to AI session | busy: true |
| T3 | CRASH — process terminates | busy: true (persisted) |
| T4 | Worker restarts, reads journal | busy: true detected |
| T5 | Job re-queued, executes again | busy: true (new record) |
| T6 | Completion recorded | busy: false |
The compacted journal would show only the final busy: false record after successful completion.
Code Example: End-to-End Scheduled Prompt with Recovery
// Client code: Request a recurring prompt
import { RPCClient } from 'prime-agent';
const client = new RPCClient({ socketPath: '/tmp/prime-agent.sock' });
await client.addSchedule('*/5 * * * *', 'Review code changes and suggest improvements');
// Daemon-mode.ts: Handling the schedule creation
class DaemonMode {
public async handleAddSchedule(
state: ActiveSessionState,
scheduleExpr: string,
promptText: string,
): Promise<void> {
// 1. Record intent to create job
this.recordWorkerRecoveryState(state, "create_job");
// 2. Create durable cron job
const job = this.createCronJobForState(state, scheduleExpr, promptText);
// 3. Schedule persists to disk
await this.cronJobStore.flush();
// 4. Confirm success
return { jobId: job.id };
}
private async runCronJob(job: AgentCronJob): Promise<void> {
const state = this.getOrCreateSessionState();
// Critical: Record before execution
this.recordWorkerRecoveryState(state, "run_job");
try {
await this.sendPromptToSession(state, job.prompt);
// Record completion
this.recordWorkerRecoveryState(state, "job_complete");
} catch (err) {
// Busy remains true; will be retried on recovery
throw err;
}
}
}
Key Files and Their Responsibilities
| File Path | Core Responsibility |
|---|---|
packages/coding-agent/src/modes/daemon/daemon-mode.ts |
RPC handling, cron job lifecycle, recovery coordination |
packages/coding-agent/src/core/cron-jobs.ts |
AgentCronJob definition, AgentCronScheduler, schedule parsing |
packages/coding-agent/src/modes/daemon/worker-recovery-journal.ts |
JSON-L persistence, compaction, replay logic |
packages/coding-agent/src/core/cron-job-store.ts |
On-disk storage for job definitions |
packages/coding-agent/src/modes/rpc/rpc-client.ts |
Client-side API for addSchedule, setHeartbeat |
Performance Characteristics
- Journal append: O(1) — single
fs.appendFileSyncper state change - Compaction: O(n) where n = number of sessions, triggered only when all idle
- Recovery read: O(n) on startup, results cached in
Mapfor O(1) lookups - Scheduler evaluation: O(m) per second where m = number of scheduled jobs
The journal file size remains bounded because compaction collapses completed operations.
Summary
- Scheduled prompts in Prime Agent use
AgentCronJobdefinitions stored inAgentCronJobStorefor durability across restarts. - Worker recovery relies on
WorkerRecoveryJournal, an append-only JSON-L log that recordsbusystate for every operation. - The
recordWorkerRecoveryStatemethod indaemon-mode.tscreates recovery records before any state-changing operation, ensuring crash consistency. - On restart, the daemon reads the latest recovery records via
readLatest()and replays anybusy: trueoperations, providing exactly-once execution for in-flight prompts. - Automatic compaction keeps the recovery journal small while preserving necessary state.
Frequently Asked Questions
What happens to scheduled prompts if the worker crashes during execution?
The partially-executed prompt is recorded with busy: true in the recovery journal. When the worker restarts, readLatest() returns this record, and the daemon re-queues the job for execution. The busy flag ensures the prompt runs to completion without being lost.
How does Prime Agent prevent duplicate scheduled prompts after recovery?
The recovery journal uses the busy flag to distinguish between completed and in-flight work. Only records with busy: true are replayed. Additionally, each AgentCronJob has a unique id, and the scheduler deduplicates jobs during initialization.
Where is the recovery journal stored and how is it configured?
The journal path is passed via the DAEMON_WORKER_RECOVERY_JOURNAL_ENV environment variable, typically set to a file like /var/run/prime-agent/worker.recovery.jsonl. The WorkerRecoveryJournal constructor opens this path for append-only writing.
Can scheduled prompts survive a full system reboot?
Yes. The AgentCronJobStore persists job definitions to a JSON file on disk (configured via cron-jobs-path). The recovery journal is also file-based. Both survive process termination and system reboots, though the reboot itself must not corrupt the filesystem.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →