返回 CodeWhale
ledger.rs
根目录 / crates / tui / src / fleet / ledger.rs
1 //! Durable fleet inbox and run ledger.
2 //!
3 //! Stores fleet state as append-only JSONL so the manager can survive
4 //! restarts and reconstruct queue/worker state by replaying records.
5 //! Artifacts are referenced by bounded metadata; large payloads live on disk
6 //! and are never embedded in the ledger.
7
8 #![allow(dead_code)]
9
10 use std::collections::BTreeMap;
11 use std::collections::HashMap;
12 #[cfg(test)]
13 use std::fs::OpenOptions;
14 use std::sync::{Mutex, OnceLock};
15
16 use super::files::{WorkspaceFile, same_file};
17 use std::io::{BufRead, Read, Seek, SeekFrom, Write};
18 use std::path::{Path, PathBuf};
19
20 use anyhow::{Context, Result, bail};
21 use codewhale_protocol::fleet::*;
22 use serde::{Deserialize, Serialize};
23 use serde_json::{Value, json};
24 use tokio::sync::Notify;
25
26 const FLEET_DIR: &str = ".codewhale";
27 const FLEET_LEDGER_FILE: &str = "fleet.jsonl";
28 const FLEET_LEDGER_LOCK_FILE: &str = "fleet.lock";
29
30 /// Path of the ledger file under `workspace`, mirroring [`FleetLedger::open`].
31 pub(crate) fn fleet_ledger_path(workspace: &Path) -> PathBuf {
32 workspace.join(FLEET_DIR).join(FLEET_LEDGER_FILE)
33 }
34
35 /// Process-wide append wakes, one [`Notify`] per ledger file (#6211 R7b).
36 /// Ledger managers open per operation, so instances cannot share a channel
37 /// handle; the registry joins differently-spelled opens of one file by
38 /// canonicalized path. A wake means "re-poll with your cursor" — the cursor
39 /// stays the source of truth, so a missed or spurious wake only costs one
40 /// extra poll, never lost or duplicated events.
41 static FLEET_LEDGER_WAKES: OnceLock<Mutex<HashMap<PathBuf, std::sync::Arc<Notify>>>> =
42 OnceLock::new();
43
44 fn fleet_ledger_wake_key(ledger_path: &Path) -> PathBuf {
45 if let Ok(canonical) = ledger_path.canonicalize() {
46 return canonical;
47 }
48 // The ledger file may not exist yet at subscribe time; canonicalize the
49 // parent so both spellings still meet. A fully missing tree falls back
50 // to the raw path on both sides.
51 match ledger_path
52 .parent()
53 .and_then(|parent| parent.canonicalize().ok())
54 .zip(ledger_path.file_name())
55 {
56 Some((parent, name)) => parent.join(name),
57 None => ledger_path.to_path_buf(),
58 }
59 }
60
61 /// Subscribe to append wakes for the ledger at `ledger_path`. Callers must
62 /// still poll on a fallback interval: the registry only sees in-process
63 /// appends, and a wake can race the subscriber's last read.
64 pub(crate) fn subscribe_fleet_ledger_appends(ledger_path: &Path) -> std::sync::Arc<Notify> {
65 let key = fleet_ledger_wake_key(ledger_path);
66 let mut wakes = FLEET_LEDGER_WAKES
67 .get_or_init(|| Mutex::new(HashMap::new()))
68 .lock()
69 .expect("fleet ledger wake registry poisoned");
70 wakes
71 .entry(key)
72 .or_insert_with(|| std::sync::Arc::new(Notify::new()))
73 .clone()
74 }
75
76 fn notify_fleet_ledger_append(ledger_path: &Path) {
77 let key = fleet_ledger_wake_key(ledger_path);
78 let notify = FLEET_LEDGER_WAKES
79 .get()
80 .and_then(|wakes| wakes.lock().ok())
81 .and_then(|wakes| wakes.get(&key).cloned());
82 if let Some(notify) = notify {
83 notify.notify_waiters();
84 }
85 }
86
87 fn inline_secret_assignment_pattern() -> &'static regex::Regex {
88 static PATTERN: std::sync::OnceLock<regex::Regex> = std::sync::OnceLock::new();
89 PATTERN.get_or_init(|| {
90 regex::Regex::new(
91 r#"(?i)(?:"|')?\b(api[_-]?key|apikey|secret|token|password|passwd|authorization|auth[_-]?token|access[_-]?key|client[_-]?secret|private[_-]?key)\b(?:"|')?\s*([:=])\s*(?:"[^"]*"|'[^']*'|[^\s,;}\]]+)"#,
92 )
93 .expect("Fleet inline secret redaction pattern must compile")
94 })
95 }
96
97 fn bearer_secret_pattern() -> &'static regex::Regex {
98 static PATTERN: std::sync::OnceLock<regex::Regex> = std::sync::OnceLock::new();
99 PATTERN.get_or_init(|| {
100 regex::Regex::new(r"(?i)\bbearer\s+[A-Za-z0-9._~+/=-]{8,}")
101 .expect("Fleet bearer secret redaction pattern must compile")
102 })
103 }
104
105 /// A single append-only record in the fleet ledger.
106 #[derive(Debug, Clone, Serialize, Deserialize)]
107 #[serde(tag = "record", rename_all = "snake_case")]
108 pub enum FleetLedgerRecord {
109 /// Replay generation rotated whenever compaction replaces durable history.
110 /// Legacy ledgers omit this record and use the stable `legacy` generation
111 /// until their first compaction.
112 ReplayEpoch {
113 epoch: String,
114 },
115 RunCreated {
116 // Boxed: FleetRun is by far the largest payload; boxing keeps the enum
117 // small (clippy::large_enum_variant). Serde treats Box<T> as T.
118 run: Box<FleetRun>,
119 },
120 RunStatusChanged {
121 run_id: FleetRunId,
122 status: FleetRunStatus,
123 timestamp: String,
124 },
125 TaskEnqueued {
126 entry: FleetInboxEntry,
127 },
128 TaskLeased {
129 run_id: FleetRunId,
130 task_id: String,
131 worker_id: String,
132 leased_at: String,
133 lease_expires_at: Option<String>,
134 },
135 TaskCompletedOrFailed {
136 run_id: FleetRunId,
137 task_id: String,
138 worker_id: String,
139 timestamp: String,
140 #[serde(default = "default_terminal_task_status")]
141 status: FleetTaskLedgerStatus,
142 },
143 /// Exact owner lifecycle high-water mark retained across compaction. Raw
144 /// lifecycle records remain the normal source; this checkpoint prevents a
145 /// compacted multi-lease history from reusing lower graph idempotency keys.
146 TaskLifecycleCheckpoint {
147 run_id: FleetRunId,
148 task_id: String,
149 lifecycle_seq: u64,
150 },
151 /// One crash-atomic terminal transition and its attempt-fenced receipt.
152 /// Keeping these values in one JSONL record prevents a process crash from
153 /// publishing a terminal task without the receipt that explains it.
154 TaskAttemptFinalized {
155 event: FleetWorkerEvent,
156 #[serde(default, skip_serializing_if = "Option::is_none")]
157 final_status: Option<FleetTaskLedgerStatus>,
158 receipt: Box<FleetReceipt>,
159 },
160 EventAppended {
161 event: FleetWorkerEvent,
162 },
163 /// Durable per-task event high-water mark retained by compaction even when
164 /// the highest raw event was intentionally excluded from projections (for
165 /// example, late progress after cancellation).
166 EventSequenceCheckpoint {
167 run_id: FleetRunId,
168 worker_id: String,
169 task_id: String,
170 seq: u64,
171 },
172 Heartbeat {
173 worker_id: String,
174 timestamp: String,
175 #[serde(default)]
176 cpu_percent: Option<f32>,
177 #[serde(default)]
178 memory_mb: Option<u64>,
179 },
180 ReceiptRecorded {
181 // Boxed for the same reason as RunCreated: FleetReceipt is the largest
182 // variant payload. Serde treats Box<T> as T.
183 receipt: Box<FleetReceipt>,
184 },
185 AlertSent {
186 run_id: FleetRunId,
187 task_id: String,
188 channel: String,
189 timestamp: String,
190 /// Present on attempt-fenced restart-exhaustion alerts. Older records
191 /// omit these fields and remain replayable as audit-only entries.
192 #[serde(default, skip_serializing_if = "Option::is_none")]
193 worker_id: Option<String>,
194 #[serde(default, skip_serializing_if = "Option::is_none")]
195 attempt: Option<u32>,
196 #[serde(default, skip_serializing_if = "Option::is_none")]
197 seq: Option<u64>,
198 },
199 }
200
201 /// Reconstructed fleet state after replaying the ledger.
202 #[derive(Debug, Clone, Default)]
203 pub struct FleetLedgerState {
204 pub runs: BTreeMap<String, FleetRun>,
205 pub run_status_overrides: BTreeMap<String, FleetRunStatus>,
206 /// Tasks keyed by run_id:task_id.
207 pub tasks: BTreeMap<String, FleetTaskState>,
208 /// Worker status by worker_id.
209 pub workers: BTreeMap<String, FleetWorkerStatus>,
210 /// Latest heartbeat by worker_id.
211 pub heartbeats: BTreeMap<String, FleetHeartbeatState>,
212 /// Latest event seq per worker_id:task_id.
213 pub latest_seq: BTreeMap<String, u64>,
214 /// Structured owner for each sequence key, used to write compaction
215 /// checkpoints without parsing delimiter-bearing ids back out of a string.
216 pub(crate) sequence_owners: BTreeMap<String, FleetEventSequenceOwner>,
217 /// Latest event envelope per worker_id:run_id:task_id.
218 pub latest_events: BTreeMap<String, FleetWorkerEvent>,
219 /// Artifact events keyed by worker_id:run_id:task_id:path.
220 pub artifact_events: BTreeMap<String, FleetWorkerEvent>,
221 /// Restart events keyed by worker_id:run_id:task_id.
222 pub restarted_events: BTreeMap<String, FleetWorkerEvent>,
223 /// Escalation events keyed by worker_id:run_id:task_id.
224 pub escalated_events: BTreeMap<String, FleetWorkerEvent>,
225 /// Completed receipts by run_id:task_id.
226 pub receipts: BTreeMap<String, FleetReceipt>,
227 /// Durable alert deliveries keyed by run/task/attempt/channel.
228 pub(crate) alerts: BTreeMap<(String, String, Option<u32>, String), FleetLedgerAlert>,
229 /// Accumulated worker usage per run_id (R6, #5567), folded from
230 /// `UsageReport` events during replay.
231 pub run_usage: BTreeMap<String, FleetRunUsage>,
232 }
233
234 /// Run-level usage accumulator (R6, #5567).
235 #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
236 pub struct FleetRunUsage {
237 pub input_tokens: u64,
238 pub output_tokens: u64,
239 }
240
241 impl FleetRunUsage {
242 #[must_use]
243 pub fn total_tokens(&self) -> u64 {
244 self.input_tokens.saturating_add(self.output_tokens)
245 }
246 }
247
248 /// Run-level budget alert identity (R6, #5567): one alert per run, not per
249 /// task, so the dedupe key uses a reserved task slot.
250 pub(crate) const BUDGET_ALERT_TASK: &str = "__run__";
251 pub(crate) const BUDGET_ALERT_CHANNEL: &str = "budget_exceeded";
252
253 impl FleetLedgerState {
254 /// Accumulated worker usage for one run; zero when nothing reported.
255 #[must_use]
256 pub fn run_usage(&self, run_id: &FleetRunId) -> FleetRunUsage {
257 self.run_usage.get(&run_id.0).copied().unwrap_or_default()
258 }
259
260 /// `(total, ceiling)` when the run declares a usage ceiling and the
261 /// accumulated total has reached it; `None` for unbounded runs.
262 #[must_use]
263 pub fn run_ceiling_breached(&self, run_id: &FleetRunId) -> Option<(u64, u64)> {
264 let ceiling = self
265 .runs
266 .get(&run_id.0)?
267 .usage_ceiling
268 .as_ref()?
269 .max_total_tokens;
270 let total = self.run_usage(run_id).total_tokens();
271 (total >= ceiling).then_some((total, ceiling))
272 }
273 }
274
275 #[derive(Debug, Clone)]
276 pub struct FleetTaskState {
277 pub entry: FleetInboxEntry,
278 pub status: FleetTaskLedgerStatus,
279 /// Monotonic owner sequence reconstructed from durable lifecycle records.
280 pub lifecycle_seq: u64,
281 pub leased_to: Option<String>,
282 pub leased_at: Option<String>,
283 pub completed_at: Option<String>,
284 }
285
286 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
287 #[serde(rename_all = "snake_case")]
288 pub enum FleetTaskLedgerStatus {
289 Enqueued,
290 Leased,
291 Completed,
292 Failed,
293 Cancelled,
294 }
295
296 fn default_terminal_task_status() -> FleetTaskLedgerStatus {
297 FleetTaskLedgerStatus::Completed
298 }
299
300 #[derive(Debug, Clone)]
301 pub struct FleetHeartbeatState {
302 pub timestamp: String,
303 pub cpu_percent: Option<f32>,
304 pub memory_mb: Option<u64>,
305 }
306
307 #[derive(Debug, Clone)]
308 pub(crate) struct FleetLedgerAlert {
309 run_id: FleetRunId,
310 task_id: String,
311 channel: String,
312 timestamp: String,
313 worker_id: Option<String>,
314 attempt: Option<u32>,
315 seq: Option<u64>,
316 }
317
318 #[derive(Debug, Clone)]
319 pub(crate) struct FleetEventSequenceOwner {
320 run_id: FleetRunId,
321 worker_id: String,
322 task_id: String,
323 }
324
325 /// Precise failure boundary for managed-client replay.
326 #[derive(Debug, thiserror::Error)]
327 pub enum FleetEventReplayError {
328 #[error("fleet run {run_id} does not exist")]
329 UnknownRun { run_id: String },
330 #[error("fleet event cursor is no longer available for run {run_id}")]
331 CursorUnavailable { run_id: String },
332 #[error("failed to read fleet event history: {message}")]
333 Storage { message: String },
334 }
335
336 /// Append-only JSONL ledger for fleet runs.
337 #[derive(Debug)]
338 pub struct FleetLedger {
339 ledger_path: PathBuf,
340 ledger_file: WorkspaceFile,
341 lock_file: WorkspaceFile,
342 original_lock: std::fs::File,
343 /// Stable advisory-lock inode shared by every manager for this workspace.
344 ///
345 /// The ledger itself is replaced during compaction, so locking
346 /// `fleet.jsonl` would leave appenders holding a lock on the old inode.
347 lock_path: PathBuf,
348 #[cfg(test)]
349 fail_start_append_after_callback: std::sync::atomic::AtomicBool,
350 #[cfg(test)]
351 fail_restart_append_after_callback: std::sync::atomic::AtomicBool,
352 }
353
354 impl FleetLedger {
355 /// Open (or create) the ledger under `workspace/.codewhale/fleet.jsonl`.
356 pub fn open(workspace: &Path) -> Result<Self> {
357 let relative = Path::new(FLEET_DIR).join(FLEET_LEDGER_FILE);
358 let ledger_path = workspace.join(&relative);
359 let ledger_file = WorkspaceFile::open(workspace, &relative, true)?;
360 ledger_file
361 .open_update(true, true)
362 .with_context(|| format!("creating fleet ledger {}", ledger_path.display()))?;
363 let lock_path = workspace.join(FLEET_DIR).join(FLEET_LEDGER_LOCK_FILE);
364 let lock_file = ledger_file.sibling(FLEET_LEDGER_LOCK_FILE)?;
365 let original_lock = lock_file
366 .open_update(true, false)
367 .with_context(|| format!("creating fleet ledger lock {}", lock_path.display()))?;
368 Ok(Self {
369 ledger_file,
370 lock_file,
371 original_lock,
372 ledger_path,
373 lock_path,
374 #[cfg(test)]
375 fail_start_append_after_callback: std::sync::atomic::AtomicBool::new(false),
376 #[cfg(test)]
377 fail_restart_append_after_callback: std::sync::atomic::AtomicBool::new(false),
378 })
379 }
380
381 pub fn path(&self) -> &Path {
382 &self.ledger_path
383 }
384
385 #[cfg(test)]
386 pub(crate) fn fail_next_start_append_after_callback(&self) {
387 self.fail_start_append_after_callback
388 .store(true, std::sync::atomic::Ordering::SeqCst);
389 }
390
391 #[cfg(test)]
392 pub(crate) fn fail_next_restart_append_after_callback(&self) {
393 self.fail_restart_append_after_callback
394 .store(true, std::sync::atomic::Ordering::SeqCst);
395 }
396
397 fn open_lock_file(&self) -> Result<std::fs::File> {
398 let file = self
399 .lock_file
400 .open_update(false, false)
401 .with_context(|| format!("opening fleet ledger lock {}", self.lock_path.display()))?;
402 if !same_file(&file, &self.original_lock)? {
403 bail!("Fleet ledger lock was replaced; reopen the workspace before continuing");
404 }
405 Ok(file)
406 }
407
408 fn with_read_lock<T>(&self, action: impl FnOnce() -> Result<T>) -> Result<T> {
409 let lock_file = self.open_lock_file()?;
410 let lock = fd_lock::RwLock::new(lock_file);
411 let _guard = lock
412 .read()
413 .with_context(|| format!("read-locking fleet ledger {}", self.ledger_path.display()))?;
414 action()
415 }
416
417 fn with_write_lock<T>(&self, action: impl FnOnce() -> Result<T>) -> Result<T> {
418 let lock_file = self.open_lock_file()?;
419 let mut lock = fd_lock::RwLock::new(lock_file);
420 let _guard = lock.write().with_context(|| {
421 format!("write-locking fleet ledger {}", self.ledger_path.display())
422 })?;
423 action()
424 }
425
426 /// Append a single record without rewriting existing ledger contents.
427 fn append_record_unlocked(&self, record: &FleetLedgerRecord) -> Result<()> {
428 self.append_records_unlocked(std::slice::from_ref(record))
429 }
430
431 /// Append a transaction's records with one write while the caller holds
432 /// the workspace ledger lock. Cancellation uses this to publish its
433 /// interrupted + cancelled pair without another manager allocating a
434 /// sequence number between them.
435 fn append_records_unlocked(&self, records: &[FleetLedgerRecord]) -> Result<()> {
436 let mut lines = String::new();
437 for record in records {
438 lines.push_str(
439 &serde_json::to_string(record).context("serializing fleet ledger record")?,
440 );
441 lines.push('\n');
442 }
443 let mut file = self
444 .ledger_file
445 .open_update(true, true)
446 .with_context(|| format!("opening fleet ledger {}", self.ledger_path.display()))?;
447 // A process can die after writing only part of its final JSON record.
448 // O_APPEND would otherwise concatenate the next valid record directly
449 // onto that unterminated tail, causing replay to discard both. Preserve
450 // the forensic fragment but quarantine it as its own malformed line.
451 let len = file
452 .metadata()
453 .with_context(|| {
454 format!(
455 "reading fleet ledger metadata {}",
456 self.ledger_path.display()
457 )
458 })?
459 .len();
460 if len > 0 {
461 file.seek(SeekFrom::End(-1)).with_context(|| {
462 format!("seeking fleet ledger tail {}", self.ledger_path.display())
463 })?;
464 let mut tail = [0_u8; 1];
465 file.read_exact(&mut tail).with_context(|| {
466 format!("reading fleet ledger tail {}", self.ledger_path.display())
467 })?;
468 if tail[0] != b'\n' {
469 file.write_all(b"\n").with_context(|| {
470 format!(
471 "quarantining fleet ledger tail {}",
472 self.ledger_path.display()
473 )
474 })?;
475 }
476 }
477 file.write_all(lines.as_bytes())
478 .with_context(|| format!("appending fleet ledger {}", self.ledger_path.display()))?;
479 file.flush()
480 .with_context(|| format!("flushing fleet ledger {}", self.ledger_path.display()))?;
481 file.sync_data()
482 .with_context(|| format!("syncing fleet ledger {}", self.ledger_path.display()))?;
483 // Wake SSE subscribers after the bytes are durable (#6211 R7b). Every
484 // record kind funnels through here, so a wake only promises "something
485 // changed" — subscribers re-poll with their cursor.
486 notify_fleet_ledger_append(&self.ledger_path);
487 Ok(())
488 }
489
490 /// Append a single record while excluding compaction and other writers.
491 fn append_record(&self, record: &FleetLedgerRecord) -> Result<()> {
492 self.with_write_lock(|| self.append_record_unlocked(record))
493 }
494
495 pub fn create_run(&self, run: &FleetRun) -> Result<()> {
496 self.append_record(&FleetLedgerRecord::RunCreated {
497 run: Box::new(sanitize_run_for_ledger(run)),
498 })
499 }
500
501 pub fn update_run_status(
502 &self,
503 run_id: &FleetRunId,
504 status: FleetRunStatus,
505 timestamp: &str,
506 ) -> Result<()> {
507 self.append_record(&FleetLedgerRecord::RunStatusChanged {
508 run_id: run_id.clone(),
509 status,
510 timestamp: timestamp.to_string(),
511 })
512 }
513
514 pub fn enqueue(&self, entry: FleetInboxEntry) -> Result<()> {
515 self.append_record(&FleetLedgerRecord::TaskEnqueued { entry })
516 }
517
518 /// Mark a task as leased to a worker.
519 ///
520 /// Production transitions must use one of the compare-and-set helpers
521 /// below. This raw append remains available only to focused replay tests.
522 #[cfg(test)]
523 pub fn lease_task(
524 &self,
525 run_id: &FleetRunId,
526 task_id: &str,
527 worker_id: &str,
528 leased_at: &str,
529 lease_expires_at: Option<&str>,
530 ) -> Result<()> {
531 self.append_record(&FleetLedgerRecord::TaskLeased {
532 run_id: run_id.clone(),
533 task_id: task_id.to_string(),
534 worker_id: worker_id.to_string(),
535 leased_at: leased_at.to_string(),
536 lease_expires_at: lease_expires_at.map(String::from),
537 })
538 }
539
540 /// Atomically lease one queued task to an idle logical worker.
541 ///
542 /// The queue snapshot, task/worker eligibility checks, run-capacity check,
543 /// and durable lease append share one cross-process lock. A competing
544 /// manager therefore observes either the queued task or the winning lease,
545 /// never the same stale snapshot followed by a second lease.
546 pub fn lease_task_if_enqueued(
547 &self,
548 run_id: &FleetRunId,
549 task_id: &str,
550 worker_id: &str,
551 leased_at: &str,
552 lease_expires_at: Option<&str>,
553 max_active_for_run: Option<usize>,
554 ) -> Result<bool> {
555 self.lease_task_if_enqueued_with_events(
556 run_id,
557 task_id,
558 worker_id,
559 leased_at,
560 lease_expires_at,
561 max_active_for_run,
562 Vec::new(),
563 false,
564 || Ok(()),
565 )
566 }
567
568 /// Atomically claim a queued task and publish its initial lifecycle.
569 ///
570 /// The callback runs before the durable transaction while the same ledger
571 /// lock is held. Fleet uses it to synchronously persist the exact launch
572 /// projection; a callback failure leaves the task Enqueued. Callers must
573 /// compensate the projection if the subsequent ledger append fails.
574 #[allow(clippy::too_many_arguments)]
575 pub fn start_task_if_enqueued(
576 &self,
577 run_id: &FleetRunId,
578 task_id: &str,
579 worker_id: &str,
580 leased_at: &str,
581 lease_expires_at: Option<&str>,
582 max_active_for_run: Option<usize>,
583 initial_events: Vec<FleetWorkerEventPayload>,
584 on_started: impl FnOnce() -> Result<()>,
585 ) -> Result<bool> {
586 self.lease_task_if_enqueued_with_events(
587 run_id,
588 task_id,
589 worker_id,
590 leased_at,
591 lease_expires_at,
592 max_active_for_run,
593 initial_events,
594 true,
595 on_started,
596 )
597 }
598
599 #[allow(clippy::too_many_arguments)]
600 fn lease_task_if_enqueued_with_events(
601 &self,
602 run_id: &FleetRunId,
603 task_id: &str,
604 worker_id: &str,
605 leased_at: &str,
606 lease_expires_at: Option<&str>,
607 max_active_for_run: Option<usize>,
608 initial_events: Vec<FleetWorkerEventPayload>,
609 record_heartbeat: bool,
610 on_started: impl FnOnce() -> Result<()>,
611 ) -> Result<bool> {
612 self.with_write_lock(move || {
613 let state = self.rebuild_state_unlocked()?;
614 let key = task_key(&run_id.0, task_id);
615 let Some(task) = state.tasks.get(&key) else {
616 return Ok(false);
617 };
618 if task.status != FleetTaskLedgerStatus::Enqueued {
619 return Ok(false);
620 }
621 if state.tasks.values().any(|candidate| {
622 candidate.status == FleetTaskLedgerStatus::Leased
623 && candidate.leased_to.as_deref() == Some(worker_id)
624 }) {
625 return Ok(false);
626 }
627 if max_active_for_run.is_some_and(|max_active| {
628 state
629 .tasks
630 .values()
631 .filter(|candidate| candidate.entry.run_id == *run_id)
632 .filter(|candidate| candidate.status == FleetTaskLedgerStatus::Leased)
633 .count()
634 >= max_active
635 }) {
636 return Ok(false);
637 }
638
639 // Budget ceiling (R6, #5567): a breached run admits nothing new.
640 // Already-leased tasks finish; the once-only alert is recorded by
641 // the usage-event path.
642 if state.run_ceiling_breached(run_id).is_some() {
643 return Ok(false);
644 }
645
646 let mut records =
647 Vec::with_capacity(1 + initial_events.len() + usize::from(record_heartbeat));
648 records.push(FleetLedgerRecord::TaskLeased {
649 run_id: run_id.clone(),
650 task_id: task_id.to_string(),
651 worker_id: worker_id.to_string(),
652 leased_at: leased_at.to_string(),
653 lease_expires_at: lease_expires_at.map(String::from),
654 });
655 let event_key = event_key(worker_id, &run_id.0, task_id);
656 let mut next_seq = state.latest_seq.get(&event_key).copied().unwrap_or(0) + 1;
657 for payload in initial_events {
658 records.push(FleetLedgerRecord::EventAppended {
659 event: FleetWorkerEvent {
660 seq: next_seq,
661 run_id: run_id.clone(),
662 worker_id: worker_id.to_string(),
663 task_id: task_id.to_string(),
664 timestamp: leased_at.to_string(),
665 payload,
666 extra: BTreeMap::new(),
667 },
668 });
669 next_seq = next_seq.saturating_add(1);
670 }
671 if record_heartbeat {
672 records.push(FleetLedgerRecord::Heartbeat {
673 worker_id: worker_id.to_string(),
674 timestamp: leased_at.to_string(),
675 cpu_percent: None,
676 memory_mb: None,
677 });
678 }
679 // Persist the exact launch/coordination projection before making
680 // the lease visible. A failed callback leaves the task Enqueued
681 // with attempt 0 rather than a Running lease that cannot launch.
682 on_started()?;
683 #[cfg(test)]
684 if self
685 .fail_start_append_after_callback
686 .swap(false, std::sync::atomic::Ordering::SeqCst)
687 {
688 bail!("forced Fleet ledger append failure after coordination callback");
689 }
690 self.append_records_unlocked(&records)?;
691 Ok(true)
692 })
693 }
694
695 #[cfg(test)]
696 pub fn mark_task_terminal_status(
697 &self,
698 run_id: &FleetRunId,
699 task_id: &str,
700 worker_id: Option<&str>,
701 timestamp: &str,
702 status: FleetTaskLedgerStatus,
703 ) -> Result<()> {
704 self.append_record(&FleetLedgerRecord::TaskCompletedOrFailed {
705 run_id: run_id.clone(),
706 task_id: task_id.to_string(),
707 worker_id: worker_id.unwrap_or_default().to_string(),
708 timestamp: timestamp.to_string(),
709 status,
710 })
711 }
712
713 /// Change a task's terminal projection only if the exact attempt and event
714 /// observed by the verifier are still current. A concurrent restart changes
715 /// both the attempt counter and lifecycle sequence, so late verification
716 /// can never fail the replacement attempt.
717 #[allow(clippy::too_many_arguments)]
718 pub fn mark_task_terminal_status_if_unchanged(
719 &self,
720 run_id: &FleetRunId,
721 task_id: &str,
722 worker_id: &str,
723 expected_status: FleetTaskLedgerStatus,
724 expected_attempts: u32,
725 expected_latest_seq: u64,
726 timestamp: &str,
727 status: FleetTaskLedgerStatus,
728 ) -> Result<bool> {
729 self.with_write_lock(|| {
730 let state = self.rebuild_state_unlocked()?;
731 let key = task_key(&run_id.0, task_id);
732 let Some(task) = state.tasks.get(&key) else {
733 return Ok(false);
734 };
735 let latest_seq = state
736 .latest_seq
737 .get(&event_key(worker_id, &run_id.0, task_id))
738 .copied()
739 .unwrap_or(0);
740 if task.status != expected_status
741 || task.leased_to.as_deref() != Some(worker_id)
742 || task.entry.attempts != expected_attempts
743 || latest_seq != expected_latest_seq
744 {
745 return Ok(false);
746 }
747 self.append_record_unlocked(&FleetLedgerRecord::TaskCompletedOrFailed {
748 run_id: run_id.clone(),
749 task_id: task_id.to_string(),
750 worker_id: worker_id.to_string(),
751 timestamp: timestamp.to_string(),
752 status,
753 })?;
754 Ok(true)
755 })
756 }
757
758 pub fn append_event(&self, event: FleetWorkerEvent) -> Result<()> {
759 let usage_receipt = matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. });
760 let run_id = event.run_id.clone();
761 let timestamp = event.timestamp.clone();
762 self.append_record(&FleetLedgerRecord::EventAppended { event })?;
763 if usage_receipt {
764 self.record_budget_exceeded_alert_once(&run_id, &timestamp)?;
765 }
766 Ok(())
767 }
768
769 /// Allocate and append the next event sequence under one ledger lock.
770 pub fn append_event_next_seq(
771 &self,
772 run_id: &FleetRunId,
773 worker_id: &str,
774 task_id: &str,
775 timestamp: &str,
776 payload: FleetWorkerEventPayload,
777 ) -> Result<FleetWorkerEvent> {
778 let event = self.with_write_lock(|| {
779 let state = self.rebuild_state_unlocked()?;
780 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
781 self.append_record_unlocked(&FleetLedgerRecord::EventAppended {
782 event: event.clone(),
783 })?;
784 Ok(event)
785 })?;
786 if matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. }) {
787 self.record_budget_exceeded_alert_once(run_id, timestamp)?;
788 }
789 Ok(event)
790 }
791
792 /// Append progress only while this exact worker still owns the live lease.
793 /// Host startup and stream draining use this guard so output produced after
794 /// an out-of-process cancellation cannot become durable task progress.
795 pub fn append_event_if_leased(
796 &self,
797 run_id: &FleetRunId,
798 worker_id: &str,
799 task_id: &str,
800 expected_attempts: u32,
801 timestamp: &str,
802 payload: FleetWorkerEventPayload,
803 ) -> Result<Option<FleetWorkerEvent>> {
804 if matches!(
805 &payload,
806 FleetWorkerEventPayload::Completed { .. }
807 | FleetWorkerEventPayload::Failed { .. }
808 | FleetWorkerEventPayload::Cancelled { .. }
809 ) {
810 bail!("conditional progress append does not accept terminal worker events");
811 }
812 let appended = self.with_write_lock(|| {
813 let state = self.rebuild_state_unlocked()?;
814 let key = task_key(&run_id.0, task_id);
815 let Some(task) = state.tasks.get(&key) else {
816 return Ok(None);
817 };
818 if task.status != FleetTaskLedgerStatus::Leased
819 || task.leased_to.as_deref() != Some(worker_id)
820 || task.entry.attempts != expected_attempts
821 {
822 return Ok(None);
823 }
824 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
825 self.append_record_unlocked(&FleetLedgerRecord::EventAppended {
826 event: event.clone(),
827 })?;
828 Ok(Some(event))
829 });
830 if let Ok(Some(event)) = &appended
831 && matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. })
832 {
833 self.record_budget_exceeded_alert_once(run_id, timestamp)?;
834 }
835 appended
836 }
837
838 /// Append a non-terminal scheduler event only if the complete lease
839 /// version observed during policy evaluation is still current.
840 #[allow(clippy::too_many_arguments)]
841 pub fn append_event_if_lease_unchanged(
842 &self,
843 run_id: &FleetRunId,
844 worker_id: &str,
845 task_id: &str,
846 expected_attempts: u32,
847 expected_latest_seq: u64,
848 expected_heartbeat_at: Option<&str>,
849 timestamp: &str,
850 payload: FleetWorkerEventPayload,
851 ) -> Result<Option<FleetWorkerEvent>> {
852 if matches!(
853 &payload,
854 FleetWorkerEventPayload::Completed { .. }
855 | FleetWorkerEventPayload::Failed { .. }
856 | FleetWorkerEventPayload::Cancelled { .. }
857 ) {
858 bail!("conditional progress append does not accept terminal worker events");
859 }
860 let appended = self.with_write_lock(|| {
861 let state = self.rebuild_state_unlocked()?;
862 let key = task_key(&run_id.0, task_id);
863 let Some(task) = state.tasks.get(&key) else {
864 return Ok(None);
865 };
866 let latest_seq = state
867 .latest_seq
868 .get(&event_key(worker_id, &run_id.0, task_id))
869 .copied()
870 .unwrap_or(0);
871 let heartbeat_at = state
872 .heartbeats
873 .get(worker_id)
874 .map(|heartbeat| heartbeat.timestamp.as_str());
875 if task.status != FleetTaskLedgerStatus::Leased
876 || task.leased_to.as_deref() != Some(worker_id)
877 || task.entry.attempts != expected_attempts
878 || latest_seq != expected_latest_seq
879 || heartbeat_at != expected_heartbeat_at
880 {
881 return Ok(None);
882 }
883 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
884 self.append_record_unlocked(&FleetLedgerRecord::EventAppended {
885 event: event.clone(),
886 })?;
887 Ok(Some(event))
888 });
889 if let Ok(Some(event)) = &appended
890 && matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. })
891 {
892 self.record_budget_exceeded_alert_once(run_id, timestamp)?;
893 }
894 appended
895 }
896
897 /// Append a terminal worker event only while the expected task lease is
898 /// still live. This is the completion side of the same compare-and-set used
899 /// by cancellation: whichever terminal transition acquires the ledger lock
900 /// first wins, and the loser cannot overwrite the task or mint a receipt.
901 pub fn append_terminal_event_if_leased(
902 &self,
903 run_id: &FleetRunId,
904 worker_id: &str,
905 task_id: &str,
906 expected_attempts: u32,
907 timestamp: &str,
908 payload: FleetWorkerEventPayload,
909 ) -> Result<Option<FleetWorkerEvent>> {
910 if !matches!(
911 &payload,
912 FleetWorkerEventPayload::Completed { .. }
913 | FleetWorkerEventPayload::Failed { .. }
914 | FleetWorkerEventPayload::Cancelled { .. }
915 ) {
916 bail!("conditional terminal append requires a terminal worker event");
917 }
918 let appended = self.with_write_lock(|| {
919 let state = self.rebuild_state_unlocked()?;
920 let key = task_key(&run_id.0, task_id);
921 let Some(task) = state.tasks.get(&key) else {
922 return Ok(None);
923 };
924 if task.status != FleetTaskLedgerStatus::Leased
925 || task.leased_to.as_deref() != Some(worker_id)
926 || task.entry.attempts != expected_attempts
927 {
928 return Ok(None);
929 }
930 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
931 self.append_record_unlocked(&FleetLedgerRecord::EventAppended {
932 event: event.clone(),
933 })?;
934 Ok(Some(event))
935 });
936 if let Ok(Some(event)) = &appended
937 && matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. })
938 {
939 self.record_budget_exceeded_alert_once(run_id, timestamp)?;
940 }
941 appended
942 }
943
944 /// Finalize one exact process attempt and its receipt as one JSONL record.
945 /// A restart increments `attempts`, so a late process or verifier cannot
946 /// terminalize or publish evidence for its replacement.
947 #[allow(clippy::too_many_arguments)]
948 pub fn finalize_task_attempt_if_leased(
949 &self,
950 run_id: &FleetRunId,
951 worker_id: &str,
952 task_id: &str,
953 expected_attempts: u32,
954 timestamp: &str,
955 payload: FleetWorkerEventPayload,
956 final_status: Option<FleetTaskLedgerStatus>,
957 mut receipt: FleetReceipt,
958 ) -> Result<Option<FleetWorkerEvent>> {
959 if !matches!(
960 &payload,
961 FleetWorkerEventPayload::Completed { .. }
962 | FleetWorkerEventPayload::Failed { .. }
963 | FleetWorkerEventPayload::Cancelled { .. }
964 ) {
965 bail!("attempt finalization requires a terminal worker event");
966 }
967 if final_status.is_some_and(|status| {
968 !matches!(
969 status,
970 FleetTaskLedgerStatus::Completed
971 | FleetTaskLedgerStatus::Failed
972 | FleetTaskLedgerStatus::Cancelled
973 )
974 }) {
975 bail!("attempt finalization status must be terminal");
976 }
977 if receipt.run_id != *run_id || receipt.task_id != task_id || receipt.worker_id != worker_id
978 {
979 bail!("attempt receipt identity does not match its terminal event");
980 }
981 if receipt
982 .attempt
983 .is_some_and(|attempt| attempt != expected_attempts)
984 {
985 bail!("attempt receipt generation does not match its lease");
986 }
987 self.with_write_lock(|| {
988 let state = self.rebuild_state_unlocked()?;
989 let key = task_key(&run_id.0, task_id);
990 let Some(task) = state.tasks.get(&key) else {
991 return Ok(None);
992 };
993 if task.status != FleetTaskLedgerStatus::Leased
994 || task.leased_to.as_deref() != Some(worker_id)
995 || task.entry.attempts != expected_attempts
996 {
997 return Ok(None);
998 }
999 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
1000 receipt.attempt = Some(expected_attempts);
1001 receipt.terminal_seq = Some(event.seq);
1002 self.append_record_unlocked(&FleetLedgerRecord::TaskAttemptFinalized {
1003 event: event.clone(),
1004 final_status,
1005 receipt: Box::new(receipt),
1006 })?;
1007 Ok(Some(event))
1008 })
1009 }
1010
1011 /// Append a scheduler-owned terminal event only if the stale lease version
1012 /// it evaluated is still exact. In particular, a fresh heartbeat arriving
1013 /// after stale detection invalidates exhaustion/failure of that worker.
1014 #[allow(clippy::too_many_arguments)]
1015 pub fn append_terminal_event_if_lease_unchanged(
1016 &self,
1017 run_id: &FleetRunId,
1018 worker_id: &str,
1019 task_id: &str,
1020 expected_attempts: u32,
1021 expected_latest_seq: u64,
1022 expected_heartbeat_at: Option<&str>,
1023 timestamp: &str,
1024 payload: FleetWorkerEventPayload,
1025 ) -> Result<Option<FleetWorkerEvent>> {
1026 if !matches!(
1027 &payload,
1028 FleetWorkerEventPayload::Completed { .. }
1029 | FleetWorkerEventPayload::Failed { .. }
1030 | FleetWorkerEventPayload::Cancelled { .. }
1031 ) {
1032 bail!("conditional terminal append requires a terminal worker event");
1033 }
1034 let appended = self.with_write_lock(|| {
1035 let state = self.rebuild_state_unlocked()?;
1036 let key = task_key(&run_id.0, task_id);
1037 let Some(task) = state.tasks.get(&key) else {
1038 return Ok(None);
1039 };
1040 let latest_seq = state
1041 .latest_seq
1042 .get(&event_key(worker_id, &run_id.0, task_id))
1043 .copied()
1044 .unwrap_or(0);
1045 let heartbeat_at = state
1046 .heartbeats
1047 .get(worker_id)
1048 .map(|heartbeat| heartbeat.timestamp.as_str());
1049 if task.status != FleetTaskLedgerStatus::Leased
1050 || task.leased_to.as_deref() != Some(worker_id)
1051 || task.entry.attempts != expected_attempts
1052 || latest_seq != expected_latest_seq
1053 || heartbeat_at != expected_heartbeat_at
1054 {
1055 return Ok(None);
1056 }
1057 let event = next_worker_event(&state, run_id, worker_id, task_id, timestamp, payload);
1058 self.append_record_unlocked(&FleetLedgerRecord::EventAppended {
1059 event: event.clone(),
1060 })?;
1061 Ok(Some(event))
1062 });
1063 if let Ok(Some(event)) = &appended
1064 && matches!(&event.payload, FleetWorkerEventPayload::UsageReport { .. })
1065 {
1066 self.record_budget_exceeded_alert_once(run_id, timestamp)?;
1067 }
1068 appended
1069 }
1070
1071 /// Atomically restart the exact task attempt observed by a manager.
1072 ///
1073 /// Status, worker ownership, attempt count, latest event sequence, and
1074 /// heartbeat are the transition version. Any intervening cancellation,
1075 /// completion, progress event, heartbeat, or competing restart makes this
1076 /// compare-and-set lose without appending a second lease.
1077 #[allow(clippy::too_many_arguments)]
1078 pub fn restart_task_if_unchanged(
1079 &self,
1080 run_id: &FleetRunId,
1081 task_id: &str,
1082 worker_id: &str,
1083 expected_status: FleetTaskLedgerStatus,
1084 expected_attempts: u32,
1085 expected_latest_seq: u64,
1086 expected_heartbeat_at: Option<&str>,
1087 leased_at: &str,
1088 lease_expires_at: Option<&str>,
1089 restart_count: u32,
1090 ) -> Result<bool> {
1091 self.restart_task_if_unchanged_with_callback(
1092 run_id,
1093 task_id,
1094 worker_id,
1095 expected_status,
1096 expected_attempts,
1097 expected_latest_seq,
1098 expected_heartbeat_at,
1099 leased_at,
1100 lease_expires_at,
1101 restart_count,
1102 || Ok(()),
1103 )
1104 }
1105
1106 /// Restart the observed attempt while persisting its replacement launch
1107 /// identity before the new lease becomes visible. This mirrors
1108 /// `start_task_if_enqueued`: a failed callback leaves the old attempt
1109 /// untouched, so no leased generation can exist without a matching
1110 /// durable launch manifest.
1111 #[allow(clippy::too_many_arguments)]
1112 pub fn restart_task_if_unchanged_with_callback(
1113 &self,
1114 run_id: &FleetRunId,
1115 task_id: &str,
1116 worker_id: &str,
1117 expected_status: FleetTaskLedgerStatus,
1118 expected_attempts: u32,
1119 expected_latest_seq: u64,
1120 expected_heartbeat_at: Option<&str>,
1121 leased_at: &str,
1122 lease_expires_at: Option<&str>,
1123 restart_count: u32,
1124 on_restarted: impl FnOnce() -> Result<()>,
1125 ) -> Result<bool> {
1126 self.with_write_lock(move || {
1127 let state = self.rebuild_state_unlocked()?;
1128 let key = task_key(&run_id.0, task_id);
1129 let Some(task) = state.tasks.get(&key) else {
1130 return Ok(false);
1131 };
1132 let lifecycle_key = event_key(worker_id, &run_id.0, task_id);
1133 let latest_seq = state.latest_seq.get(&lifecycle_key).copied().unwrap_or(0);
1134 let heartbeat_at = state
1135 .heartbeats
1136 .get(worker_id)
1137 .map(|heartbeat| heartbeat.timestamp.as_str());
1138 if task.status != expected_status
1139 || task.leased_to.as_deref() != Some(worker_id)
1140 || task.entry.attempts != expected_attempts
1141 || latest_seq != expected_latest_seq
1142 || heartbeat_at != expected_heartbeat_at
1143 {
1144 return Ok(false);
1145 }
1146 if state.tasks.iter().any(|(candidate_key, candidate)| {
1147 candidate_key != &key
1148 && candidate.status == FleetTaskLedgerStatus::Leased
1149 && candidate.leased_to.as_deref() == Some(worker_id)
1150 }) {
1151 return Ok(false);
1152 }
1153
1154 let restarted = FleetWorkerEvent {
1155 seq: latest_seq.saturating_add(1),
1156 run_id: run_id.clone(),
1157 worker_id: worker_id.to_string(),
1158 task_id: task_id.to_string(),
1159 timestamp: leased_at.to_string(),
1160 payload: FleetWorkerEventPayload::Restarted { restart_count },
1161 extra: BTreeMap::new(),
1162 };
1163 let running = FleetWorkerEvent {
1164 seq: restarted.seq.saturating_add(1),
1165 run_id: run_id.clone(),
1166 worker_id: worker_id.to_string(),
1167 task_id: task_id.to_string(),
1168 timestamp: leased_at.to_string(),
1169 payload: FleetWorkerEventPayload::Running,
1170 extra: BTreeMap::new(),
1171 };
1172 on_restarted()?;
1173 #[cfg(test)]
1174 if self
1175 .fail_restart_append_after_callback
1176 .swap(false, std::sync::atomic::Ordering::SeqCst)
1177 {
1178 bail!("forced Fleet restart ledger append failure after coordination callback");
1179 }
1180 self.append_records_unlocked(&[
1181 FleetLedgerRecord::TaskLeased {
1182 run_id: run_id.clone(),
1183 task_id: task_id.to_string(),
1184 worker_id: worker_id.to_string(),
1185 leased_at: leased_at.to_string(),
1186 lease_expires_at: lease_expires_at.map(String::from),
1187 },
1188 FleetLedgerRecord::EventAppended { event: restarted },
1189 FleetLedgerRecord::EventAppended { event: running },
1190 FleetLedgerRecord::Heartbeat {
1191 worker_id: worker_id.to_string(),
1192 timestamp: leased_at.to_string(),
1193 cpu_percent: None,
1194 memory_mb: None,
1195 },
1196 ])?;
1197 Ok(true)
1198 })
1199 }
1200
1201 /// Atomically cancel an active task, optionally requiring one exact lease.
1202 ///
1203 /// `expected_worker_id = Some` is a compare-and-set for an operator command
1204 /// targeting a specific live worker. `None` expresses run-wide cancellation
1205 /// and accepts either a queued task or whichever worker leased it before
1206 /// this transaction acquired the lock.
1207 pub fn cancel_task_if_active(
1208 &self,
1209 run_id: &FleetRunId,
1210 task_id: &str,
1211 expected_worker_id: Option<&str>,
1212 timestamp: &str,
1213 signal: Option<&str>,
1214 cancelled_by: Option<&str>,
1215 ) -> Result<bool> {
1216 self.with_write_lock(|| {
1217 let state = self.rebuild_state_unlocked()?;
1218 let key = task_key(&run_id.0, task_id);
1219 let Some(task) = state.tasks.get(&key) else {
1220 return Ok(false);
1221 };
1222 if !matches!(
1223 task.status,
1224 FleetTaskLedgerStatus::Enqueued | FleetTaskLedgerStatus::Leased
1225 ) {
1226 return Ok(false);
1227 }
1228 if expected_worker_id.is_some_and(|expected| {
1229 task.status != FleetTaskLedgerStatus::Leased
1230 || task.leased_to.as_deref() != Some(expected)
1231 }) {
1232 return Ok(false);
1233 }
1234
1235 if let Some(worker_id) = task.leased_to.as_deref() {
1236 let event_key = event_key(worker_id, &run_id.0, task_id);
1237 let first_seq = state.latest_seq.get(&event_key).copied().unwrap_or(0) + 1;
1238 let mut records = Vec::with_capacity(2);
1239 let mut next_seq = first_seq;
1240 if let Some(signal) = signal {
1241 records.push(FleetLedgerRecord::EventAppended {
1242 event: FleetWorkerEvent {
1243 seq: next_seq,
1244 run_id: run_id.clone(),
1245 worker_id: worker_id.to_string(),
1246 task_id: task_id.to_string(),
1247 timestamp: timestamp.to_string(),
1248 payload: FleetWorkerEventPayload::Interrupted {
1249 signal: Some(signal.to_string()),
1250 },
1251 extra: BTreeMap::new(),
1252 },
1253 });
1254 next_seq += 1;
1255 }
1256 records.push(FleetLedgerRecord::EventAppended {
1257 event: FleetWorkerEvent {
1258 seq: next_seq,
1259 run_id: run_id.clone(),
1260 worker_id: worker_id.to_string(),
1261 task_id: task_id.to_string(),
1262 timestamp: timestamp.to_string(),
1263 payload: FleetWorkerEventPayload::Cancelled {
1264 cancelled_by: cancelled_by.map(str::to_string),
1265 },
1266 extra: BTreeMap::new(),
1267 },
1268 });
1269 self.append_records_unlocked(&records)?;
1270 } else {
1271 self.append_record_unlocked(&FleetLedgerRecord::TaskCompletedOrFailed {
1272 run_id: run_id.clone(),
1273 task_id: task_id.to_string(),
1274 worker_id: String::new(),
1275 timestamp: timestamp.to_string(),
1276 status: FleetTaskLedgerStatus::Cancelled,
1277 })?;
1278 }
1279 Ok(true)
1280 })
1281 }
1282
1283 pub fn heartbeat(
1284 &self,
1285 worker_id: &str,
1286 timestamp: &str,
1287 cpu_percent: Option<f32>,
1288 memory_mb: Option<u64>,
1289 ) -> Result<()> {
1290 self.append_record(&FleetLedgerRecord::Heartbeat {
1291 worker_id: worker_id.to_string(),
1292 timestamp: timestamp.to_string(),
1293 cpu_percent,
1294 memory_mb,
1295 })
1296 }
1297
1298 pub fn record_receipt(&self, receipt: FleetReceipt) -> Result<()> {
1299 self.append_record(&FleetLedgerRecord::ReceiptRecorded {
1300 receipt: Box::new(receipt),
1301 })
1302 }
1303
1304 /// Record the run's budget-exceeded alert exactly once and pause the
1305 /// run (R6, #5567). Returns true only for the recording transition; a
1306 /// non-breached or already-alerted run is a no-op.
1307 pub fn record_budget_exceeded_alert_once(
1308 &self,
1309 run_id: &FleetRunId,
1310 timestamp: &str,
1311 ) -> Result<bool> {
1312 self.with_write_lock(|| {
1313 let state = self.rebuild_state_unlocked()?;
1314 let Some((total, ceiling)) = state.run_ceiling_breached(run_id) else {
1315 return Ok(false);
1316 };
1317 let key = alert_key(run_id, BUDGET_ALERT_TASK, None, BUDGET_ALERT_CHANNEL);
1318 if state.alerts.contains_key(&key) {
1319 return Ok(false);
1320 }
1321 tracing::warn!(
1322 target: "fleet",
1323 run = %run_id.0,
1324 total,
1325 ceiling,
1326 "fleet run crossed its usage ceiling; pausing run and refusing new admissions"
1327 );
1328 self.append_record_unlocked(&FleetLedgerRecord::AlertSent {
1329 run_id: run_id.clone(),
1330 task_id: BUDGET_ALERT_TASK.to_string(),
1331 channel: BUDGET_ALERT_CHANNEL.to_string(),
1332 timestamp: timestamp.to_string(),
1333 worker_id: None,
1334 attempt: None,
1335 seq: None,
1336 })?;
1337 self.append_record_unlocked(&FleetLedgerRecord::RunStatusChanged {
1338 run_id: run_id.clone(),
1339 status: FleetRunStatus::Paused,
1340 timestamp: timestamp.to_string(),
1341 })?;
1342 Ok(true)
1343 })
1344 }
1345
1346 pub fn record_alert(
1347 &self,
1348 run_id: &FleetRunId,
1349 task_id: &str,
1350 channel: &str,
1351 timestamp: &str,
1352 ) -> Result<()> {
1353 self.append_record(&FleetLedgerRecord::AlertSent {
1354 run_id: run_id.clone(),
1355 task_id: task_id.to_string(),
1356 channel: channel.to_string(),
1357 timestamp: timestamp.to_string(),
1358 worker_id: None,
1359 attempt: None,
1360 seq: None,
1361 })
1362 }
1363
1364 /// Record one restart-exhaustion alert for one exact failed attempt.
1365 /// The audit marker and inspectable escalation are represented by one JSONL
1366 /// record, so competing schedulers and crash recovery cannot duplicate it.
1367 #[allow(clippy::too_many_arguments)]
1368 pub fn record_failed_attempt_alert_once(
1369 &self,
1370 run_id: &FleetRunId,
1371 task_id: &str,
1372 worker_id: &str,
1373 expected_attempts: u32,
1374 channel_label: &str,
1375 channel_key: &str,
1376 timestamp: &str,
1377 ) -> Result<bool> {
1378 self.with_write_lock(|| {
1379 let state = self.rebuild_state_unlocked()?;
1380 let task_key = task_key(&run_id.0, task_id);
1381 let Some(task) = state.tasks.get(&task_key) else {
1382 return Ok(false);
1383 };
1384 if task.status != FleetTaskLedgerStatus::Failed
1385 || task.leased_to.as_deref() != Some(worker_id)
1386 || task.entry.attempts != expected_attempts
1387 {
1388 return Ok(false);
1389 }
1390 let exact_alert_key = alert_key(run_id, task_id, Some(expected_attempts), channel_key);
1391 // Pre-attempt ledgers stored only the display label. Replay infers
1392 // their attempt, so suppress that same delivery without treating a
1393 // legacy attempt as a permanent block on future retries.
1394 let legacy_alert_key =
1395 alert_key(run_id, task_id, Some(expected_attempts), channel_label);
1396 if state.alerts.contains_key(&exact_alert_key)
1397 || state.alerts.contains_key(&legacy_alert_key)
1398 {
1399 return Ok(false);
1400 }
1401 let lifecycle_key = event_key(worker_id, &run_id.0, task_id);
1402 let seq = state
1403 .latest_seq
1404 .get(&lifecycle_key)
1405 .copied()
1406 .unwrap_or(0)
1407 .saturating_add(1);
1408 self.append_record_unlocked(&FleetLedgerRecord::AlertSent {
1409 run_id: run_id.clone(),
1410 task_id: task_id.to_string(),
1411 channel: channel_key.to_string(),
1412 timestamp: timestamp.to_string(),
1413 worker_id: Some(worker_id.to_string()),
1414 attempt: Some(expected_attempts),
1415 seq: Some(seq),
1416 })?;
1417 Ok(true)
1418 })
1419 }
1420
1421 /// Replay the ledger and reconstruct current state. Malformed or partial
1422 /// lines are skipped so an interrupted write cannot corrupt earlier state.
1423 pub fn rebuild_state(&self) -> Result<FleetLedgerState> {
1424 self.with_read_lock(|| self.rebuild_state_unlocked())
1425 }
1426
1427 /// Read one bounded page from the durable Fleet transition history.
1428 ///
1429 /// A cursor is the opaque digest of a ledger transition, not a global
1430 /// worker sequence. That distinction matters because worker `seq` values
1431 /// restart for each `(worker, task)` lifecycle. Recent cursors survive a
1432 /// process restart and normal appends. Ledger compaction may intentionally
1433 /// discard old history; in that case callers receive `CursorUnavailable`
1434 /// and must reload the current run projection instead of silently skipping
1435 /// an unknown gap.
1436 pub fn replay_events(
1437 &self,
1438 run_id: &FleetRunId,
1439 after: Option<&str>,
1440 limit: usize,
1441 ) -> std::result::Result<FleetEventReplay, FleetEventReplayError> {
1442 let (run_exists, all_events, history_compacted) = self
1443 .with_read_lock(|| self.scan_runtime_events_unlocked(run_id))
1444 .map_err(|error| FleetEventReplayError::Storage {
1445 message: error.to_string(),
1446 })?;
1447 if !run_exists {
1448 return Err(FleetEventReplayError::UnknownRun {
1449 run_id: run_id.0.clone(),
1450 });
1451 }
1452
1453 let limit = limit.clamp(1, 1_000);
1454 let (events, has_more, history_truncated) = if let Some(after) = after {
1455 let Some(position) = all_events.iter().position(|event| event.cursor == after) else {
1456 return Err(FleetEventReplayError::CursorUnavailable {
1457 run_id: run_id.0.clone(),
1458 });
1459 };
1460 let remaining = &all_events[position.saturating_add(1)..];
1461 (
1462 remaining.iter().take(limit).cloned().collect::<Vec<_>>(),
1463 remaining.len() > limit,
1464 false,
1465 )
1466 } else {
1467 let start = all_events.len().saturating_sub(limit);
1468 (
1469 all_events[start..].to_vec(),
1470 false,
1471 start > 0 || history_compacted,
1472 )
1473 };
1474 let next_cursor = events.last().map(|event| event.cursor.clone());
1475 Ok(FleetEventReplay {
1476 run_id: run_id.clone(),
1477 events,
1478 has_more,
1479 history_truncated,
1480 next_cursor,
1481 })
1482 }
1483
1484 fn scan_runtime_events_unlocked(
1485 &self,
1486 run_filter: &FleetRunId,
1487 ) -> Result<(bool, Vec<FleetRuntimeEvent>, bool)> {
1488 let mut state = FleetLedgerState::default();
1489 let mut events = Vec::new();
1490 let mut replay_epoch = "legacy".to_string();
1491 let mut history_compacted = false;
1492 let file = match self.ledger_file.open_file() {
1493 Ok(file) => file,
1494 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1495 return Ok((false, events, history_compacted));
1496 }
1497 Err(error) => return Err(error).context("Opening Fleet ledger for replay"),
1498 };
1499 let reader = std::io::BufReader::new(file);
1500 for (line_no, line) in reader.lines().enumerate() {
1501 let line = match line {
1502 Ok(line) => line,
1503 Err(error) => {
1504 tracing::warn!(
1505 "fleet ledger line {} unreadable during event replay: {}",
1506 line_no + 1,
1507 error
1508 );
1509 continue;
1510 }
1511 };
1512 if line.trim().is_empty() {
1513 continue;
1514 }
1515 let record = match serde_json::from_str::<FleetLedgerRecord>(&line) {
1516 Ok(record) => record,
1517 Err(error) => {
1518 tracing::warn!(
1519 "fleet ledger line {} parse error during event replay (skipping): {}",
1520 line_no + 1,
1521 error
1522 );
1523 continue;
1524 }
1525 };
1526 if let FleetLedgerRecord::ReplayEpoch { epoch } = &record {
1527 replay_epoch.clone_from(epoch);
1528 history_compacted = true;
1529 apply_record(&mut state, record);
1530 continue;
1531 }
1532 for event in runtime_events_from_record(&record, &state, line_no, &replay_epoch) {
1533 if event.run_id == *run_filter {
1534 events.push(event);
1535 }
1536 }
1537 apply_record(&mut state, record);
1538 }
1539 Ok((
1540 state.runs.contains_key(&run_filter.0),
1541 events,
1542 history_compacted,
1543 ))
1544 }
1545
1546 fn rebuild_state_unlocked(&self) -> Result<FleetLedgerState> {
1547 let mut state = FleetLedgerState::default();
1548 let file = match self.ledger_file.open_file() {
1549 Ok(file) => file,
1550 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(state),
1551 Err(error) => return Err(error).context("Opening Fleet ledger for state rebuild"),
1552 };
1553 let reader = std::io::BufReader::new(file);
1554 for (line_no, line) in reader.lines().enumerate() {
1555 let line = match line {
1556 Ok(l) => l,
1557 Err(err) => {
1558 tracing::warn!("fleet ledger line {} unreadable: {}", line_no + 1, err);
1559 continue;
1560 }
1561 };
1562 if line.trim().is_empty() {
1563 continue;
1564 }
1565 let record: FleetLedgerRecord = match serde_json::from_str(&line) {
1566 Ok(r) => r,
1567 Err(err) => {
1568 tracing::warn!(
1569 "fleet ledger line {} parse error (skipping): {}",
1570 line_no + 1,
1571 err
1572 );
1573 continue;
1574 }
1575 };
1576 apply_record(&mut state, record);
1577 }
1578 Ok(state)
1579 }
1580
1581 /// Claim the next available inbox task for `worker_id`. Returns the
1582 /// enqueued entry and appends a lease record.
1583 pub fn claim_next(
1584 &self,
1585 worker_id: &str,
1586 _worker_capabilities: &[String],
1587 timestamp: &str,
1588 ) -> Result<Option<FleetInboxEntry>> {
1589 self.with_write_lock(|| {
1590 let state = self.rebuild_state_unlocked()?;
1591 if state.tasks.values().any(|candidate| {
1592 candidate.status == FleetTaskLedgerStatus::Leased
1593 && candidate.leased_to.as_deref() == Some(worker_id)
1594 }) {
1595 return Ok(None);
1596 }
1597 // Find oldest enqueued task whose task spec (if known) matches worker
1598 // capabilities. For now, tasks without specs match everything.
1599 let candidate = state
1600 .tasks
1601 .values()
1602 .filter(|t| matches!(t.status, FleetTaskLedgerStatus::Enqueued))
1603 .map(|t| &t.entry)
1604 .min_by_key(|e| (e.priority, e.enqueued_at.clone()))
1605 .cloned();
1606 let Some(entry) = candidate else {
1607 return Ok(None);
1608 };
1609 self.append_record_unlocked(&FleetLedgerRecord::TaskLeased {
1610 run_id: entry.run_id.clone(),
1611 task_id: entry.task_id.clone(),
1612 worker_id: worker_id.to_string(),
1613 leased_at: timestamp.to_string(),
1614 lease_expires_at: None,
1615 })?;
1616 Ok(Some(entry))
1617 })
1618 }
1619
1620 /// Compact the ledger by rewriting only the records needed to reconstruct
1621 /// current state. This truncates history but preserves run/task/event
1622 /// metadata and receipts.
1623 pub fn compact(&self) -> Result<()> {
1624 self.compact_with_snapshot_hook(|| {})
1625 }
1626
1627 fn compact_with_snapshot_hook(&self, after_snapshot: impl FnOnce()) -> Result<()> {
1628 self.with_write_lock(|| {
1629 let state = self.rebuild_state_unlocked()?;
1630 after_snapshot();
1631 let mut lines = vec![serde_json::to_string(&FleetLedgerRecord::ReplayEpoch {
1632 epoch: uuid::Uuid::new_v4().simple().to_string(),
1633 })?];
1634 let mut terminal_lines = Vec::new();
1635 let mut lifecycle_lines = Vec::new();
1636 for run in state.runs.values() {
1637 lines.push(serde_json::to_string(&FleetLedgerRecord::RunCreated {
1638 run: Box::new(run.clone()),
1639 })?);
1640 if let Some(status) = state.run_status_overrides.get(&run.id.0) {
1641 lines.push(serde_json::to_string(
1642 &FleetLedgerRecord::RunStatusChanged {
1643 run_id: run.id.clone(),
1644 status: status.clone(),
1645 timestamp: run.updated_at.clone().unwrap_or_default(),
1646 },
1647 )?);
1648 }
1649 }
1650 for task in state.tasks.values() {
1651 let mut enqueued_entry = task.entry.clone();
1652 if task.leased_to.is_some() {
1653 // Replaying the retained TaskLeased record increments the
1654 // attempt counter, so checkpoint the value immediately
1655 // before that transition instead of inventing an attempt
1656 // on every compaction.
1657 enqueued_entry.attempts = enqueued_entry.attempts.saturating_sub(1);
1658 enqueued_entry.lease_deadline = None;
1659 }
1660 lines.push(serde_json::to_string(&FleetLedgerRecord::TaskEnqueued {
1661 entry: enqueued_entry,
1662 })?);
1663 if let Some(worker) = &task.leased_to {
1664 lines.push(serde_json::to_string(&FleetLedgerRecord::TaskLeased {
1665 run_id: task.entry.run_id.clone(),
1666 task_id: task.entry.task_id.clone(),
1667 worker_id: worker.clone(),
1668 leased_at: task.leased_at.clone().unwrap_or_default(),
1669 lease_expires_at: task.entry.lease_deadline.clone(),
1670 })?);
1671 }
1672 if matches!(
1673 task.status,
1674 FleetTaskLedgerStatus::Completed
1675 | FleetTaskLedgerStatus::Failed
1676 | FleetTaskLedgerStatus::Cancelled
1677 ) {
1678 terminal_lines.push(serde_json::to_string(
1679 &FleetLedgerRecord::TaskCompletedOrFailed {
1680 run_id: task.entry.run_id.clone(),
1681 task_id: task.entry.task_id.clone(),
1682 worker_id: task.leased_to.clone().unwrap_or_default(),
1683 timestamp: task.completed_at.clone().unwrap_or_default(),
1684 status: task.status,
1685 },
1686 )?);
1687 }
1688 lifecycle_lines.push(serde_json::to_string(
1689 &FleetLedgerRecord::TaskLifecycleCheckpoint {
1690 run_id: task.entry.run_id.clone(),
1691 task_id: task.entry.task_id.clone(),
1692 lifecycle_seq: task.lifecycle_seq,
1693 },
1694 )?);
1695 }
1696 for alert in state.alerts.values() {
1697 lines.push(serde_json::to_string(&FleetLedgerRecord::AlertSent {
1698 run_id: alert.run_id.clone(),
1699 task_id: alert.task_id.clone(),
1700 channel: alert.channel.clone(),
1701 timestamp: alert.timestamp.clone(),
1702 worker_id: alert.worker_id.clone(),
1703 attempt: alert.attempt,
1704 seq: alert.seq,
1705 })?);
1706 }
1707 let mut compacted_events = BTreeMap::new();
1708 for event in state.latest_events.values() {
1709 compacted_events.insert(compact_event_key(event), event.clone());
1710 }
1711 for event in state.artifact_events.values() {
1712 compacted_events.insert(compact_event_key(event), event.clone());
1713 }
1714 for event in state.restarted_events.values() {
1715 compacted_events.insert(compact_event_key(event), event.clone());
1716 }
1717 for event in state.escalated_events.values() {
1718 compacted_events.insert(compact_event_key(event), event.clone());
1719 }
1720 let mut compacted_events = compacted_events.into_values().collect::<Vec<_>>();
1721 compacted_events.sort_by(|left, right| {
1722 left.worker_id
1723 .cmp(&right.worker_id)
1724 .then_with(|| left.run_id.0.cmp(&right.run_id.0))
1725 .then_with(|| left.task_id.cmp(&right.task_id))
1726 .then_with(|| left.seq.cmp(&right.seq))
1727 });
1728 for event in compacted_events {
1729 lines.push(serde_json::to_string(&FleetLedgerRecord::EventAppended {
1730 event,
1731 })?);
1732 }
1733 for (key, seq) in &state.latest_seq {
1734 let Some(owner) = state.sequence_owners.get(key) else {
1735 continue;
1736 };
1737 lines.push(serde_json::to_string(
1738 &FleetLedgerRecord::EventSequenceCheckpoint {
1739 run_id: owner.run_id.clone(),
1740 worker_id: owner.worker_id.clone(),
1741 task_id: owner.task_id.clone(),
1742 seq: *seq,
1743 },
1744 )?);
1745 }
1746 for (worker_id, heartbeat) in &state.heartbeats {
1747 lines.push(serde_json::to_string(&FleetLedgerRecord::Heartbeat {
1748 worker_id: worker_id.clone(),
1749 timestamp: heartbeat.timestamp.clone(),
1750 cpu_percent: heartbeat.cpu_percent,
1751 memory_mb: heartbeat.memory_mb,
1752 })?);
1753 }
1754 // Terminal status is the final task-state projection. In
1755 // particular, a verifier may override a successful process exit to
1756 // Failed. Emit these records after retained worker events so replay
1757 // cannot let an earlier Completed event erase that override.
1758 lines.extend(terminal_lines);
1759 // Lifecycle checkpoints follow every reconstructed task-state
1760 // transition so replay ends at the exact pre-compaction sequence.
1761 lines.extend(lifecycle_lines);
1762 for receipt in state.receipts.values() {
1763 lines.push(serde_json::to_string(
1764 &FleetLedgerRecord::ReceiptRecorded {
1765 receipt: Box::new(receipt.clone()),
1766 },
1767 )?);
1768 }
1769 let mut contents = lines.join("\n");
1770 if !contents.is_empty() {
1771 contents.push('\n');
1772 }
1773 self.ledger_file
1774 .replace(contents.as_bytes())
1775 .with_context(|| {
1776 format!(
1777 "atomically compacting Fleet ledger {}",
1778 self.ledger_path.display()
1779 )
1780 })?;
1781 Ok(())
1782 })
1783 }
1784 }
1785
1786 fn runtime_events_from_record(
1787 record: &FleetLedgerRecord,
1788 state: &FleetLedgerState,
1789 ledger_ordinal: usize,
1790 replay_epoch: &str,
1791 ) -> Vec<FleetRuntimeEvent> {
1792 match record {
1793 FleetLedgerRecord::ReplayEpoch { .. } => Vec::new(),
1794 FleetLedgerRecord::RunCreated { run } => vec![FleetRuntimeEvent {
1795 cursor: fleet_runtime_cursor(&run.id, "run_created", ledger_ordinal, replay_epoch),
1796 event: "fleet.run.created".to_string(),
1797 run_id: run.id.clone(),
1798 worker_id: None,
1799 task_id: None,
1800 timestamp: Some(run.created_at.clone()),
1801 worker_seq: None,
1802 payload: json!({
1803 "target": run.target,
1804 "workflow": run.workflow,
1805 "roles": run.roles,
1806 "task_count": run.task_specs.len(),
1807 "worker_count": run.worker_specs.len(),
1808 }),
1809 }],
1810 FleetLedgerRecord::RunStatusChanged {
1811 run_id,
1812 status,
1813 timestamp,
1814 } => vec![FleetRuntimeEvent {
1815 cursor: fleet_runtime_cursor(run_id, "run_status", ledger_ordinal, replay_epoch),
1816 event: "fleet.run.status_changed".to_string(),
1817 run_id: run_id.clone(),
1818 worker_id: None,
1819 task_id: None,
1820 timestamp: Some(timestamp.clone()),
1821 worker_seq: None,
1822 payload: json!({ "status": status }),
1823 }],
1824 FleetLedgerRecord::TaskEnqueued { entry } => vec![FleetRuntimeEvent {
1825 cursor: fleet_runtime_cursor(
1826 &entry.run_id,
1827 "task_enqueued",
1828 ledger_ordinal,
1829 replay_epoch,
1830 ),
1831 event: "fleet.task.enqueued".to_string(),
1832 run_id: entry.run_id.clone(),
1833 worker_id: None,
1834 task_id: Some(entry.task_id.clone()),
1835 timestamp: Some(entry.enqueued_at.clone()),
1836 worker_seq: None,
1837 payload: json!({
1838 "priority": entry.priority,
1839 "attempts": entry.attempts,
1840 }),
1841 }],
1842 FleetLedgerRecord::TaskLeased {
1843 run_id,
1844 task_id,
1845 worker_id,
1846 leased_at,
1847 lease_expires_at,
1848 } => vec![FleetRuntimeEvent {
1849 cursor: fleet_runtime_cursor(run_id, "task_leased", ledger_ordinal, replay_epoch),
1850 event: "fleet.task.leased".to_string(),
1851 run_id: run_id.clone(),
1852 worker_id: Some(worker_id.clone()),
1853 task_id: Some(task_id.clone()),
1854 timestamp: Some(leased_at.clone()),
1855 worker_seq: None,
1856 payload: json!({ "lease_expires_at": lease_expires_at }),
1857 }],
1858 FleetLedgerRecord::TaskCompletedOrFailed {
1859 run_id,
1860 task_id,
1861 worker_id,
1862 timestamp,
1863 status,
1864 } => vec![FleetRuntimeEvent {
1865 cursor: fleet_runtime_cursor(run_id, "task_terminal", ledger_ordinal, replay_epoch),
1866 event: "fleet.task.terminal".to_string(),
1867 run_id: run_id.clone(),
1868 worker_id: (!worker_id.is_empty()).then(|| worker_id.clone()),
1869 task_id: Some(task_id.clone()),
1870 timestamp: Some(timestamp.clone()),
1871 worker_seq: None,
1872 payload: json!({ "status": status }),
1873 }],
1874 FleetLedgerRecord::TaskAttemptFinalized { event, receipt, .. } => vec![
1875 runtime_worker_event("worker", ledger_ordinal, replay_epoch, event),
1876 runtime_receipt_event("receipt", ledger_ordinal, replay_epoch, receipt),
1877 ],
1878 FleetLedgerRecord::EventAppended { event } => {
1879 vec![runtime_worker_event(
1880 "worker",
1881 ledger_ordinal,
1882 replay_epoch,
1883 event,
1884 )]
1885 }
1886 FleetLedgerRecord::Heartbeat {
1887 worker_id,
1888 timestamp,
1889 cpu_percent,
1890 memory_mb,
1891 } => active_task_for_replay_worker(state, worker_id)
1892 .map(|task| FleetRuntimeEvent {
1893 cursor: fleet_runtime_cursor(
1894 &task.entry.run_id,
1895 "heartbeat",
1896 ledger_ordinal,
1897 replay_epoch,
1898 ),
1899 event: "fleet.worker.heartbeat".to_string(),
1900 run_id: task.entry.run_id.clone(),
1901 worker_id: Some(worker_id.clone()),
1902 task_id: Some(task.entry.task_id.clone()),
1903 timestamp: Some(timestamp.clone()),
1904 worker_seq: None,
1905 payload: json!({
1906 "cpu_percent": cpu_percent,
1907 "memory_mb": memory_mb,
1908 }),
1909 })
1910 .into_iter()
1911 .collect(),
1912 FleetLedgerRecord::ReceiptRecorded { receipt } => {
1913 vec![runtime_receipt_event(
1914 "receipt",
1915 ledger_ordinal,
1916 replay_epoch,
1917 receipt,
1918 )]
1919 }
1920 FleetLedgerRecord::AlertSent {
1921 run_id,
1922 task_id,
1923 timestamp,
1924 worker_id,
1925 attempt,
1926 ..
1927 } => vec![FleetRuntimeEvent {
1928 cursor: fleet_runtime_cursor(run_id, "alert", ledger_ordinal, replay_epoch),
1929 event: "fleet.alert.sent".to_string(),
1930 run_id: run_id.clone(),
1931 worker_id: worker_id.clone(),
1932 task_id: Some(task_id.clone()),
1933 timestamp: Some(timestamp.clone()),
1934 worker_seq: None,
1935 payload: json!({ "attempt": attempt }),
1936 }],
1937 FleetLedgerRecord::TaskLifecycleCheckpoint { .. }
1938 | FleetLedgerRecord::EventSequenceCheckpoint { .. } => Vec::new(),
1939 }
1940 }
1941
1942 fn active_task_for_replay_worker<'a>(
1943 state: &'a FleetLedgerState,
1944 worker_id: &str,
1945 ) -> Option<&'a FleetTaskState> {
1946 state.tasks.values().find(|task| {
1947 task.status == FleetTaskLedgerStatus::Leased && task.leased_to.as_deref() == Some(worker_id)
1948 })
1949 }
1950
1951 fn runtime_worker_event(
1952 cursor_kind: &str,
1953 ledger_ordinal: usize,
1954 replay_epoch: &str,
1955 event: &FleetWorkerEvent,
1956 ) -> FleetRuntimeEvent {
1957 let (event_name, payload) = privacy_bounded_worker_payload(&event.payload);
1958 FleetRuntimeEvent {
1959 cursor: fleet_runtime_cursor(&event.run_id, cursor_kind, ledger_ordinal, replay_epoch),
1960 event: event_name,
1961 run_id: event.run_id.clone(),
1962 worker_id: Some(event.worker_id.clone()),
1963 task_id: Some(event.task_id.clone()),
1964 timestamp: Some(event.timestamp.clone()),
1965 worker_seq: Some(event.seq),
1966 payload,
1967 }
1968 }
1969
1970 fn runtime_receipt_event(
1971 cursor_kind: &str,
1972 ledger_ordinal: usize,
1973 replay_epoch: &str,
1974 receipt: &FleetReceipt,
1975 ) -> FleetRuntimeEvent {
1976 let score = receipt.score.as_ref().map(|score| {
1977 json!({
1978 "value": score.value,
1979 "max": score.max,
1980 })
1981 });
1982 FleetRuntimeEvent {
1983 cursor: fleet_runtime_cursor(&receipt.run_id, cursor_kind, ledger_ordinal, replay_epoch),
1984 event: "fleet.task.receipt_recorded".to_string(),
1985 run_id: receipt.run_id.clone(),
1986 worker_id: Some(receipt.worker_id.clone()),
1987 task_id: Some(receipt.task_id.clone()),
1988 timestamp: Some(receipt.completed_at.clone()),
1989 worker_seq: receipt.terminal_seq,
1990 payload: json!({
1991 "attempt": receipt.attempt,
1992 "result": receipt.result,
1993 "failure_kind": receipt.failure_kind,
1994 "score": score,
1995 "artifact_kinds": receipt.artifacts.iter().map(|artifact| &artifact.kind).collect::<Vec<_>>(),
1996 }),
1997 }
1998 }
1999
2000 fn privacy_bounded_worker_payload(payload: &FleetWorkerEventPayload) -> (String, Value) {
2001 let state = match payload {
2002 FleetWorkerEventPayload::Queued => "queued",
2003 FleetWorkerEventPayload::Leased { .. } => "leased",
2004 FleetWorkerEventPayload::Starting => "starting",
2005 FleetWorkerEventPayload::Running => "running",
2006 FleetWorkerEventPayload::ModelWait { .. } => "model_wait",
2007 FleetWorkerEventPayload::RunningTool { .. } => "running_tool",
2008 FleetWorkerEventPayload::WorkflowEvent { .. } => "workflow_event",
2009 FleetWorkerEventPayload::Heartbeat { .. } => "heartbeat",
2010 FleetWorkerEventPayload::UsageReport { .. } => "usage_report",
2011 FleetWorkerEventPayload::Artifact(_) => "artifact",
2012 FleetWorkerEventPayload::Completed { .. } => "completed",
2013 FleetWorkerEventPayload::Failed { .. } => "failed",
2014 FleetWorkerEventPayload::Cancelled { .. } => "cancelled",
2015 FleetWorkerEventPayload::Interrupted { .. } => "interrupted",
2016 FleetWorkerEventPayload::Stale { .. } => "stale",
2017 FleetWorkerEventPayload::Restarted { .. } => "restarted",
2018 FleetWorkerEventPayload::Escalated { .. } => "escalated",
2019 };
2020 let value = match payload {
2021 FleetWorkerEventPayload::Leased { lease_expires_at } => {
2022 json!({ "state": state, "lease_expires_at": lease_expires_at })
2023 }
2024 FleetWorkerEventPayload::ModelWait { model } => {
2025 json!({ "state": state, "model": model })
2026 }
2027 FleetWorkerEventPayload::RunningTool { tool, .. } => {
2028 json!({ "state": state, "tool": tool })
2029 }
2030 FleetWorkerEventPayload::WorkflowEvent {
2031 workflow_run_id,
2032 event,
2033 } => json!({
2034 "state": state,
2035 "workflow_run_id": workflow_run_id,
2036 "event_type": event.get("type").and_then(Value::as_str),
2037 }),
2038 FleetWorkerEventPayload::Heartbeat {
2039 cpu_percent,
2040 memory_mb,
2041 } => json!({
2042 "state": state,
2043 "cpu_percent": cpu_percent,
2044 "memory_mb": memory_mb,
2045 }),
2046 FleetWorkerEventPayload::UsageReport {
2047 input_tokens,
2048 output_tokens,
2049 } => json!({
2050 "state": state,
2051 "input_tokens": input_tokens,
2052 "output_tokens": output_tokens,
2053 }),
2054 FleetWorkerEventPayload::Artifact(artifact) => json!({
2055 "state": state,
2056 "kind": artifact.kind,
2057 "mime_type": artifact.mime_type,
2058 "size_bytes": artifact.size_bytes,
2059 }),
2060 FleetWorkerEventPayload::Completed { exit_code, .. } => {
2061 json!({ "state": state, "exit_code": exit_code })
2062 }
2063 FleetWorkerEventPayload::Failed {
2064 reason,
2065 recoverable,
2066 } => json!({
2067 "state": state,
2068 "reason": redact_fleet_event_text(reason),
2069 "recoverable": recoverable,
2070 }),
2071 FleetWorkerEventPayload::Stale { last_heartbeat_at } => {
2072 json!({ "state": state, "last_heartbeat_at": last_heartbeat_at })
2073 }
2074 FleetWorkerEventPayload::Restarted { restart_count } => {
2075 json!({ "state": state, "restart_count": restart_count })
2076 }
2077 FleetWorkerEventPayload::Escalated { channel, .. } => {
2078 json!({ "state": state, "channel": channel })
2079 }
2080 FleetWorkerEventPayload::Queued
2081 | FleetWorkerEventPayload::Starting
2082 | FleetWorkerEventPayload::Running
2083 | FleetWorkerEventPayload::Cancelled { .. }
2084 | FleetWorkerEventPayload::Interrupted { .. } => json!({ "state": state }),
2085 };
2086 (format!("fleet.worker.{state}"), value)
2087 }
2088
2089 fn redact_fleet_event_text(value: &str) -> String {
2090 // The shared redactor deliberately treats one whole line as an assignment.
2091 // Failure diagnostics often prefix an inline assignment (`provider failed:
2092 // api_key=...`), so make a second token-level pass before exposing the
2093 // bounded preview. Whitespace is normalized because this is a status
2094 // summary, not the forensic worker log.
2095 let redacted = bearer_secret_pattern()
2096 .replace_all(value, "Bearer [redacted]")
2097 .into_owned();
2098 let redacted = inline_secret_assignment_pattern()
2099 .replace_all(&redacted, "$1$2[redacted]")
2100 .into_owned();
2101 let redacted = codewhale_config::persistence::redact_secrets(&redacted);
2102 let redacted = redacted
2103 .split_whitespace()
2104 .map(codewhale_config::persistence::redact_secrets)
2105 .collect::<Vec<_>>()
2106 .join(" ");
2107 let mut chars = redacted.chars();
2108 let preview = chars.by_ref().take(1_000).collect::<String>();
2109 if chars.next().is_some() {
2110 format!("{preview}...")
2111 } else {
2112 preview
2113 }
2114 }
2115
2116 fn fleet_runtime_cursor(
2117 run_id: &FleetRunId,
2118 kind: &str,
2119 ledger_ordinal: usize,
2120 replay_epoch: &str,
2121 ) -> String {
2122 // Cursor material is deliberately metadata-only: hashing a raw ledger
2123 // record would let a caller correlate or guess secret-bearing diagnostic
2124 // text. The generation rotates on compaction, while the physical JSONL
2125 // ordinal remains stable across ordinary appends and process restarts.
2126 let material = format!(
2127 "fleet-event-v1\0{replay_epoch}\0{}\0{ledger_ordinal}\0{kind}",
2128 run_id.0
2129 );
2130 format!(
2131 "fev1_{}_{}",
2132 crate::hashing::sha256_hex(material.as_bytes()),
2133 kind
2134 )
2135 }
2136
2137 fn task_key(run_id: &str, task_id: &str) -> String {
2138 format!("{run_id}:{task_id}")
2139 }
2140
2141 fn event_key(worker_id: &str, run_id: &str, task_id: &str) -> String {
2142 format!("{worker_id}:{run_id}:{task_id}")
2143 }
2144
2145 fn alert_key(
2146 run_id: &FleetRunId,
2147 task_id: &str,
2148 attempt: Option<u32>,
2149 channel: &str,
2150 ) -> (String, String, Option<u32>, String) {
2151 (
2152 run_id.0.clone(),
2153 task_id.to_string(),
2154 attempt,
2155 channel.to_string(),
2156 )
2157 }
2158
2159 fn alert_channel_label(channel_key: &str) -> &str {
2160 channel_key
2161 .rsplit_once('#')
2162 .filter(|(_, ordinal)| ordinal.parse::<usize>().is_ok())
2163 .map_or(channel_key, |(label, _)| label)
2164 }
2165
2166 fn next_worker_event(
2167 state: &FleetLedgerState,
2168 run_id: &FleetRunId,
2169 worker_id: &str,
2170 task_id: &str,
2171 timestamp: &str,
2172 payload: FleetWorkerEventPayload,
2173 ) -> FleetWorkerEvent {
2174 let key = event_key(worker_id, &run_id.0, task_id);
2175 FleetWorkerEvent {
2176 seq: state.latest_seq.get(&key).copied().unwrap_or(0) + 1,
2177 run_id: run_id.clone(),
2178 worker_id: worker_id.to_string(),
2179 task_id: task_id.to_string(),
2180 timestamp: timestamp.to_string(),
2181 payload,
2182 extra: BTreeMap::new(),
2183 }
2184 }
2185
2186 fn compact_event_key(event: &FleetWorkerEvent) -> String {
2187 format!(
2188 "{}:{}:{}:{}",
2189 event.worker_id, event.run_id.0, event.task_id, event.seq
2190 )
2191 }
2192
2193 fn mark_task_terminal(
2194 state: &mut FleetLedgerState,
2195 run_id: &FleetRunId,
2196 task_id: &str,
2197 worker_id: &str,
2198 timestamp: &str,
2199 status: FleetTaskLedgerStatus,
2200 ) {
2201 let key = task_key(&run_id.0, task_id);
2202 if let Some(task) = state.tasks.get_mut(&key) {
2203 if task.status != status {
2204 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
2205 }
2206 task.status = status;
2207 if !worker_id.is_empty() {
2208 task.leased_to = Some(worker_id.to_string());
2209 }
2210 task.completed_at = Some(timestamp.to_string());
2211 }
2212 }
2213
2214 fn artifact_event_key(event: &FleetWorkerEvent, artifact: &FleetArtifactRef) -> String {
2215 format!(
2216 "{}:{}:{}:{}",
2217 event.worker_id,
2218 event.run_id.0,
2219 event.task_id,
2220 artifact.path.display()
2221 )
2222 }
2223
2224 fn apply_record(state: &mut FleetLedgerState, record: FleetLedgerRecord) {
2225 match record {
2226 FleetLedgerRecord::ReplayEpoch { .. } => {}
2227 FleetLedgerRecord::RunCreated { run } => {
2228 state.runs.insert(run.id.0.clone(), *run);
2229 }
2230 FleetLedgerRecord::RunStatusChanged {
2231 run_id,
2232 status,
2233 timestamp: _,
2234 } => {
2235 state.run_status_overrides.insert(run_id.0, status);
2236 }
2237 FleetLedgerRecord::TaskEnqueued { entry } => {
2238 let key = task_key(&entry.run_id.0, &entry.task_id);
2239 state.tasks.entry(key).or_insert_with(|| FleetTaskState {
2240 entry,
2241 status: FleetTaskLedgerStatus::Enqueued,
2242 lifecycle_seq: 1,
2243 leased_to: None,
2244 leased_at: None,
2245 completed_at: None,
2246 });
2247 }
2248 FleetLedgerRecord::TaskLeased {
2249 run_id,
2250 task_id,
2251 worker_id,
2252 leased_at,
2253 lease_expires_at,
2254 } => {
2255 let key = task_key(&run_id.0, &task_id);
2256 if let Some(task) = state.tasks.get_mut(&key) {
2257 if task.status != FleetTaskLedgerStatus::Leased
2258 || task.leased_to.as_deref() != Some(worker_id.as_str())
2259 || task.leased_at.as_deref() != Some(leased_at.as_str())
2260 {
2261 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
2262 }
2263 task.status = FleetTaskLedgerStatus::Leased;
2264 task.leased_to = Some(worker_id);
2265 task.leased_at = Some(leased_at);
2266 task.entry.lease_deadline = lease_expires_at;
2267 task.entry.attempts = task.entry.attempts.saturating_add(1);
2268 }
2269 }
2270 FleetLedgerRecord::TaskCompletedOrFailed {
2271 run_id,
2272 task_id,
2273 worker_id,
2274 timestamp,
2275 status,
2276 } => {
2277 mark_task_terminal(state, &run_id, &task_id, &worker_id, &timestamp, status);
2278 }
2279 FleetLedgerRecord::TaskLifecycleCheckpoint {
2280 run_id,
2281 task_id,
2282 lifecycle_seq,
2283 } => {
2284 if let Some(task) = state.tasks.get_mut(&task_key(&run_id.0, &task_id)) {
2285 task.lifecycle_seq = task.lifecycle_seq.max(lifecycle_seq.max(1));
2286 }
2287 }
2288 FleetLedgerRecord::TaskAttemptFinalized {
2289 event,
2290 final_status,
2291 receipt,
2292 } => {
2293 let run_id = event.run_id.clone();
2294 let task_id = event.task_id.clone();
2295 let worker_id = event.worker_id.clone();
2296 let timestamp = event.timestamp.clone();
2297 apply_record(state, FleetLedgerRecord::EventAppended { event });
2298 if let Some(status) = final_status {
2299 mark_task_terminal(state, &run_id, &task_id, &worker_id, &timestamp, status);
2300 }
2301 let key = task_key(&receipt.run_id.0, &receipt.task_id);
2302 state.receipts.insert(key, *receipt);
2303 }
2304 FleetLedgerRecord::EventAppended { event } => {
2305 let latest_event_key = event_key(&event.worker_id, &event.run_id.0, &event.task_id);
2306 state.sequence_owners.insert(
2307 latest_event_key.clone(),
2308 FleetEventSequenceOwner {
2309 run_id: event.run_id.clone(),
2310 worker_id: event.worker_id.clone(),
2311 task_id: event.task_id.clone(),
2312 },
2313 );
2314 let task_is_terminal = state
2315 .tasks
2316 .get(&task_key(&event.run_id.0, &event.task_id))
2317 .is_some_and(|task| {
2318 matches!(
2319 task.status,
2320 FleetTaskLedgerStatus::Completed
2321 | FleetTaskLedgerStatus::Failed
2322 | FleetTaskLedgerStatus::Cancelled
2323 )
2324 });
2325 let regresses_terminal_state = task_is_terminal
2326 && matches!(
2327 &event.payload,
2328 FleetWorkerEventPayload::Leased { .. }
2329 | FleetWorkerEventPayload::Artifact(_)
2330 | FleetWorkerEventPayload::ModelWait { .. }
2331 | FleetWorkerEventPayload::RunningTool { .. }
2332 | FleetWorkerEventPayload::WorkflowEvent { .. }
2333 | FleetWorkerEventPayload::Heartbeat { .. }
2334 | FleetWorkerEventPayload::UsageReport { .. }
2335 | FleetWorkerEventPayload::Starting
2336 | FleetWorkerEventPayload::Running
2337 | FleetWorkerEventPayload::Stale { .. }
2338 | FleetWorkerEventPayload::Restarted { .. }
2339 | FleetWorkerEventPayload::Interrupted { .. }
2340 );
2341 if state
2342 .latest_seq
2343 .get(&latest_event_key)
2344 .copied()
2345 .is_none_or(|seq| event.seq > seq)
2346 {
2347 state.latest_seq.insert(latest_event_key.clone(), event.seq);
2348 if !regresses_terminal_state {
2349 state.latest_events.insert(latest_event_key, event.clone());
2350 }
2351 }
2352 if let FleetWorkerEventPayload::Artifact(artifact) = &event.payload {
2353 state
2354 .artifact_events
2355 .insert(artifact_event_key(&event, artifact), event.clone());
2356 }
2357 if matches!(&event.payload, FleetWorkerEventPayload::Restarted { .. }) {
2358 state.restarted_events.insert(
2359 event_key(&event.worker_id, &event.run_id.0, &event.task_id),
2360 event.clone(),
2361 );
2362 }
2363 if matches!(&event.payload, FleetWorkerEventPayload::Escalated { .. }) {
2364 state.escalated_events.insert(
2365 event_key(&event.worker_id, &event.run_id.0, &event.task_id),
2366 event.clone(),
2367 );
2368 }
2369 // Usage always counts, even when a late receipt lands after the
2370 // task went terminal — the provider billed it either way.
2371 if let FleetWorkerEventPayload::UsageReport {
2372 input_tokens,
2373 output_tokens,
2374 } = &event.payload
2375 {
2376 let usage = state.run_usage.entry(event.run_id.0.clone()).or_default();
2377 usage.input_tokens = usage.input_tokens.saturating_add(*input_tokens);
2378 usage.output_tokens = usage.output_tokens.saturating_add(*output_tokens);
2379 }
2380 // Derive worker status from lifecycle events. A late stream event
2381 // must never resurrect a terminal task after an out-of-process
2382 // cancellation raced the foreground executor's final drain.
2383 if regresses_terminal_state {
2384 return;
2385 }
2386 match &event.payload {
2387 FleetWorkerEventPayload::Leased { .. }
2388 | FleetWorkerEventPayload::Restarted { .. }
2389 | FleetWorkerEventPayload::ModelWait { .. }
2390 | FleetWorkerEventPayload::RunningTool { .. }
2391 | FleetWorkerEventPayload::WorkflowEvent { .. }
2392 | FleetWorkerEventPayload::Heartbeat { .. }
2393 | FleetWorkerEventPayload::UsageReport { .. }
2394 | FleetWorkerEventPayload::Starting
2395 | FleetWorkerEventPayload::Running => {
2396 state
2397 .workers
2398 .insert(event.worker_id.clone(), FleetWorkerStatus::Busy);
2399 }
2400 FleetWorkerEventPayload::Interrupted { .. } => {
2401 state
2402 .workers
2403 .insert(event.worker_id.clone(), FleetWorkerStatus::Draining);
2404 }
2405 FleetWorkerEventPayload::Stale { .. } => {
2406 state
2407 .workers
2408 .insert(event.worker_id.clone(), FleetWorkerStatus::Unhealthy);
2409 }
2410 FleetWorkerEventPayload::Completed { .. } => {
2411 mark_task_terminal(
2412 state,
2413 &event.run_id,
2414 &event.task_id,
2415 &event.worker_id,
2416 &event.timestamp,
2417 FleetTaskLedgerStatus::Completed,
2418 );
2419 state
2420 .workers
2421 .insert(event.worker_id.clone(), FleetWorkerStatus::Online);
2422 }
2423 FleetWorkerEventPayload::Failed { .. } => {
2424 mark_task_terminal(
2425 state,
2426 &event.run_id,
2427 &event.task_id,
2428 &event.worker_id,
2429 &event.timestamp,
2430 FleetTaskLedgerStatus::Failed,
2431 );
2432 state
2433 .workers
2434 .insert(event.worker_id.clone(), FleetWorkerStatus::Online);
2435 }
2436 FleetWorkerEventPayload::Cancelled { .. } => {
2437 mark_task_terminal(
2438 state,
2439 &event.run_id,
2440 &event.task_id,
2441 &event.worker_id,
2442 &event.timestamp,
2443 FleetTaskLedgerStatus::Cancelled,
2444 );
2445 state
2446 .workers
2447 .insert(event.worker_id.clone(), FleetWorkerStatus::Online);
2448 }
2449 _ => {}
2450 }
2451 }
2452 FleetLedgerRecord::EventSequenceCheckpoint {
2453 run_id,
2454 worker_id,
2455 task_id,
2456 seq,
2457 } => {
2458 let key = event_key(&worker_id, &run_id.0, &task_id);
2459 state.sequence_owners.insert(
2460 key.clone(),
2461 FleetEventSequenceOwner {
2462 run_id,
2463 worker_id,
2464 task_id,
2465 },
2466 );
2467 state
2468 .latest_seq
2469 .entry(key)
2470 .and_modify(|current| *current = (*current).max(seq))
2471 .or_insert(seq);
2472 }
2473 FleetLedgerRecord::Heartbeat {
2474 worker_id,
2475 timestamp,
2476 cpu_percent,
2477 memory_mb,
2478 } => {
2479 state.heartbeats.insert(
2480 worker_id.clone(),
2481 FleetHeartbeatState {
2482 timestamp,
2483 cpu_percent,
2484 memory_mb,
2485 },
2486 );
2487 if state
2488 .workers
2489 .get(&worker_id)
2490 .cloned()
2491 .unwrap_or(FleetWorkerStatus::Unknown)
2492 != FleetWorkerStatus::Busy
2493 {
2494 state.workers.insert(worker_id, FleetWorkerStatus::Online);
2495 }
2496 }
2497 FleetLedgerRecord::ReceiptRecorded { receipt } => {
2498 let key = task_key(&receipt.run_id.0, &receipt.task_id);
2499 state.receipts.insert(key, *receipt);
2500 }
2501 FleetLedgerRecord::AlertSent {
2502 run_id,
2503 task_id,
2504 channel,
2505 timestamp,
2506 worker_id,
2507 attempt,
2508 seq,
2509 } => {
2510 // Old AlertSent records had no attempt. They were appended after
2511 // terminalization, so the replay projection at this point carries
2512 // the exact attempt that delivered them. Normalize immediately;
2513 // compaction then upgrades the durable record too.
2514 let attempt = attempt.or_else(|| {
2515 state
2516 .tasks
2517 .get(&task_key(&run_id.0, &task_id))
2518 .map(|task| task.entry.attempts)
2519 .filter(|attempt| *attempt > 0)
2520 });
2521 let key = alert_key(&run_id, &task_id, attempt, &channel);
2522 let is_new = !state.alerts.contains_key(&key);
2523 state.alerts.entry(key).or_insert_with(|| FleetLedgerAlert {
2524 run_id: run_id.clone(),
2525 task_id: task_id.clone(),
2526 channel: channel.clone(),
2527 timestamp: timestamp.clone(),
2528 worker_id: worker_id.clone(),
2529 attempt,
2530 seq,
2531 });
2532 if is_new && let (Some(worker_id), Some(seq)) = (worker_id, seq) {
2533 apply_record(
2534 state,
2535 FleetLedgerRecord::EventAppended {
2536 event: FleetWorkerEvent {
2537 seq,
2538 run_id,
2539 worker_id,
2540 task_id,
2541 timestamp,
2542 payload: FleetWorkerEventPayload::Escalated {
2543 channel: alert_channel_label(&channel).to_string(),
2544 alert_id: None,
2545 },
2546 extra: BTreeMap::new(),
2547 },
2548 },
2549 );
2550 }
2551 }
2552 }
2553 }
2554
2555 fn sanitize_run_for_ledger(run: &FleetRun) -> FleetRun {
2556 let mut run = run.clone();
2557 for task in &mut run.task_specs {
2558 if let Some(policy) = &mut task.alert_policy {
2559 for channel in &mut policy.channels {
2560 match channel {
2561 FleetAlertChannel::Slack { webhook } => {
2562 webhook.url = webhook.url.as_ref().map(|_| "<redacted>".to_string());
2563 }
2564 FleetAlertChannel::Webhook { endpoint } => {
2565 *endpoint = FleetAlertEndpoint {
2566 url: endpoint.url.as_ref().map(|_| "<redacted>".to_string()),
2567 url_ref: endpoint
2568 .url_ref
2569 .as_ref()
2570 .map(|_| FleetSecretRef::new("<redacted>")),
2571 secret_ref: endpoint
2572 .secret_ref
2573 .as_ref()
2574 .map(|_| FleetSecretRef::new("<redacted>")),
2575 };
2576 }
2577 FleetAlertChannel::PagerDuty { routing_key, .. } => {
2578 *routing_key = "<redacted>".to_string();
2579 }
2580 }
2581 }
2582 }
2583 }
2584 run
2585 }
2586
2587 #[cfg(test)]
2588 mod tests {
2589 use super::*;
2590 use anyhow::anyhow;
2591 use std::sync::{
2592 Arc, Barrier,
2593 atomic::{AtomicBool, Ordering},
2594 mpsc,
2595 };
2596 use std::thread;
2597 use std::time::Duration;
2598 use tempfile::TempDir;
2599
2600 #[tokio::test]
2601 async fn ledger_append_wakes_subscribers() {
2602 let workspace = TempDir::new().unwrap();
2603 let ledger = FleetLedger::open(workspace.path()).unwrap();
2604 let appends = subscribe_fleet_ledger_appends(&fleet_ledger_path(workspace.path()));
2605 let notified = appends.notified();
2606 tokio::pin!(notified);
2607 let appended = tokio::task::spawn_blocking(move || {
2608 ledger.create_run(&sample_run("run-1")).unwrap();
2609 });
2610 tokio::time::timeout(Duration::from_secs(5), &mut notified)
2611 .await
2612 .expect("append should wake subscribers");
2613 appended.await.unwrap();
2614 }
2615
2616 #[tokio::test]
2617 async fn ledger_wake_is_scoped_to_its_file() {
2618 let first = TempDir::new().unwrap();
2619 let second = TempDir::new().unwrap();
2620 let ledger = FleetLedger::open(first.path()).unwrap();
2621 let appends = subscribe_fleet_ledger_appends(&fleet_ledger_path(second.path()));
2622 ledger.create_run(&sample_run("run-1")).unwrap();
2623 assert!(
2624 tokio::time::timeout(Duration::from_millis(100), appends.notified())
2625 .await
2626 .is_err(),
2627 "an append must not wake another ledger's subscribers"
2628 );
2629 }
2630
2631 /// The workspace lock sidecar is created through the confined, private
2632 /// file helper: owner-only, never umask-widened.
2633 #[cfg(unix)]
2634 #[test]
2635 fn ledger_lock_file_is_owner_only() {
2636 use std::os::unix::fs::PermissionsExt as _;
2637 let workspace = TempDir::new().unwrap();
2638 let ledger = FleetLedger::open(workspace.path()).unwrap();
2639 let mode = |path: &Path| std::fs::metadata(path).unwrap().permissions().mode() & 0o777;
2640 assert_eq!(mode(&ledger.lock_path), 0o600);
2641 assert_eq!(mode(ledger.path()), 0o600);
2642 }
2643
2644 #[test]
2645 fn ledger_rejects_replaced_lock_identity() {
2646 let workspace = TempDir::new().unwrap();
2647 let ledger = FleetLedger::open(workspace.path()).unwrap();
2648 std::fs::remove_file(&ledger.lock_path).unwrap();
2649 std::fs::write(&ledger.lock_path, b"").unwrap();
2650 assert!(ledger.rebuild_state().is_err());
2651 assert!(ledger.compact().is_err());
2652 }
2653
2654 #[test]
2655 fn ledger_rejects_hard_linked_ledger_and_lock() {
2656 for name in ["fleet.jsonl", "fleet.lock"] {
2657 let workspace = TempDir::new().unwrap();
2658 let outside = TempDir::new().unwrap();
2659 std::fs::create_dir(workspace.path().join(".codewhale")).unwrap();
2660 let canary = outside.path().join("canary.txt");
2661 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_HARD_LINK_CANARY").unwrap();
2662 std::fs::hard_link(&canary, workspace.path().join(".codewhale").join(name)).unwrap();
2663 assert!(FleetLedger::open(workspace.path()).is_err());
2664 assert_eq!(
2665 std::fs::read(canary).unwrap(),
2666 b"OUTSIDE_SYNTHETIC_HARD_LINK_CANARY"
2667 );
2668 }
2669 }
2670
2671 #[cfg(unix)]
2672 #[test]
2673 fn ledger_rejects_parent_and_final_linked_paths() {
2674 use std::os::unix::fs::symlink;
2675 let workspace = TempDir::new().unwrap();
2676 let outside = TempDir::new().unwrap();
2677 symlink(outside.path(), workspace.path().join(".codewhale")).unwrap();
2678 assert!(FleetLedger::open(workspace.path()).is_err());
2679 assert!(!outside.path().join("fleet.jsonl").exists());
2680 assert!(!outside.path().join("fleet.lock").exists());
2681 for name in ["fleet.jsonl", "fleet.lock"] {
2682 for hard_link in [false, true] {
2683 let workspace = TempDir::new().unwrap();
2684 std::fs::create_dir(workspace.path().join(".codewhale")).unwrap();
2685 let outside = TempDir::new().unwrap();
2686 let canary = outside.path().join("canary.txt");
2687 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_LEDGER_CANARY").unwrap();
2688 let linked = workspace.path().join(".codewhale").join(name);
2689 if hard_link {
2690 std::fs::hard_link(&canary, linked).unwrap();
2691 } else {
2692 symlink(&canary, linked).unwrap();
2693 }
2694 assert!(
2695 FleetLedger::open(workspace.path()).is_err(),
2696 "accepted linked {name}"
2697 );
2698 assert_eq!(
2699 std::fs::read(canary).unwrap(),
2700 b"OUTSIDE_SYNTHETIC_LEDGER_CANARY"
2701 );
2702 }
2703 }
2704 }
2705
2706 #[cfg(unix)]
2707 #[test]
2708 fn ledger_compaction_ignores_preexisting_temporary_links() {
2709 use std::os::unix::fs::symlink;
2710 let workspace = TempDir::new().unwrap();
2711 let outside = TempDir::new().unwrap();
2712 let ledger = FleetLedger::open(workspace.path()).unwrap();
2713 let canary = outside.path().join("canary.txt");
2714 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_COMPACTION_CANARY").unwrap();
2715 symlink(&canary, ledger.path().with_extension(".tmp")).unwrap();
2716 ledger.compact().unwrap();
2717 assert_eq!(
2718 std::fs::read(canary).unwrap(),
2719 b"OUTSIDE_SYNTHETIC_COMPACTION_CANARY"
2720 );
2721 assert!(ledger.rebuild_state().unwrap().runs.is_empty());
2722 assert!(
2723 !std::fs::symlink_metadata(ledger.path())
2724 .unwrap()
2725 .file_type()
2726 .is_symlink()
2727 );
2728 }
2729
2730 #[cfg(unix)]
2731 #[test]
2732 fn ledger_rejects_linked_ledger_and_replaced_lock_after_open() {
2733 use std::os::unix::fs::symlink;
2734 let workspace = TempDir::new().unwrap();
2735 let outside = TempDir::new().unwrap();
2736 let ledger = FleetLedger::open(workspace.path()).unwrap();
2737 let canary = outside.path().join("canary.txt");
2738 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_LEDGER_CANARY").unwrap();
2739 std::fs::remove_file(ledger.path()).unwrap();
2740 symlink(&canary, ledger.path()).unwrap();
2741 assert!(ledger.compact().is_err());
2742 assert!(ledger.rebuild_state().is_err());
2743 assert_eq!(
2744 std::fs::read(&canary).unwrap(),
2745 b"OUTSIDE_SYNTHETIC_LEDGER_CANARY"
2746 );
2747 std::fs::remove_file(ledger.path()).unwrap();
2748 std::fs::write(ledger.path(), b"").unwrap();
2749 std::fs::remove_file(&ledger.lock_path).unwrap();
2750 std::fs::write(&ledger.lock_path, b"").unwrap();
2751 assert!(ledger.compact().is_err());
2752 assert!(ledger.rebuild_state().is_err());
2753 }
2754
2755 #[cfg(unix)]
2756 #[test]
2757 fn ledger_uses_pinned_parent_after_workspace_path_swap() {
2758 use std::os::unix::fs::symlink;
2759 let workspace = TempDir::new().unwrap();
2760 let outside = TempDir::new().unwrap();
2761 let ledger = FleetLedger::open(workspace.path()).unwrap();
2762 std::fs::rename(
2763 workspace.path().join(".codewhale"),
2764 workspace.path().join("retained"),
2765 )
2766 .unwrap();
2767 symlink(outside.path(), workspace.path().join(".codewhale")).unwrap();
2768 ledger.compact().unwrap();
2769 assert!(
2770 std::fs::read(workspace.path().join("retained/fleet.jsonl"))
2771 .unwrap()
2772 .starts_with(b"{\"record\":\"replay_epoch\"")
2773 );
2774 assert!(!outside.path().join("fleet.jsonl").exists());
2775 assert!(!outside.path().join("fleet.lock").exists());
2776 }
2777
2778 fn sample_run(id: &str) -> FleetRun {
2779 FleetRun {
2780 id: FleetRunId::from(id),
2781 name: "smoke".to_string(),
2782 status: FleetRunStatus::Running,
2783 target: None,
2784 workflow: None,
2785 roles: Vec::new(),
2786 max_workers: None,
2787 usage_ceiling: None,
2788 task_specs: vec![],
2789 worker_specs: vec![],
2790 labels: BTreeMap::new(),
2791 security_policy: None,
2792 created_at: "2026-06-12T17:00:00Z".to_string(),
2793 updated_at: None,
2794 completed_at: None,
2795 }
2796 }
2797
2798 fn sample_entry(run_id: &str, task_id: &str) -> FleetInboxEntry {
2799 FleetInboxEntry {
2800 run_id: FleetRunId::from(run_id),
2801 task_id: task_id.to_string(),
2802 priority: 0,
2803 enqueued_at: "2026-06-12T17:00:00Z".to_string(),
2804 lease_deadline: None,
2805 attempts: 0,
2806 }
2807 }
2808
2809 #[test]
2810 fn usage_ceiling_breach_pauses_run_refuses_admissions_and_alerts_once() {
2811 let tmp = TempDir::new().unwrap();
2812 let ledger = FleetLedger::open(tmp.path()).unwrap();
2813 let mut run = sample_run("budget-run");
2814 run.usage_ceiling = Some(codewhale_protocol::fleet::FleetUsageCeiling {
2815 max_total_tokens: 1_000,
2816 });
2817 ledger.create_run(&run).unwrap();
2818 ledger
2819 .enqueue(sample_entry("budget-run", "task-a"))
2820 .unwrap();
2821 ledger
2822 .enqueue(sample_entry("budget-run", "task-b"))
2823 .unwrap();
2824
2825 // Under the ceiling, admission works.
2826 assert!(
2827 ledger
2828 .lease_task_if_enqueued(
2829 &run.id,
2830 "task-a",
2831 "worker-1",
2832 "2026-06-12T17:00:01Z",
2833 None,
2834 None,
2835 )
2836 .unwrap()
2837 );
2838
2839 // Usage receipts accumulate; the second crosses the 1 000 ceiling and
2840 // must pause the run and fire exactly one durable alert.
2841 for (ts, input, output) in [
2842 ("2026-06-12T17:00:02Z", 600, 0),
2843 ("2026-06-12T17:00:03Z", 300, 200),
2844 ] {
2845 ledger
2846 .append_event_next_seq(
2847 &run.id,
2848 "worker-1",
2849 "task-a",
2850 ts,
2851 FleetWorkerEventPayload::UsageReport {
2852 input_tokens: input,
2853 output_tokens: output,
2854 },
2855 )
2856 .unwrap();
2857 }
2858 let budget_alerts = |state: &FleetLedgerState| {
2859 state
2860 .alerts
2861 .keys()
2862 .filter(|(run, task, _, channel)| {
2863 run == "budget-run"
2864 && task == BUDGET_ALERT_TASK
2865 && channel == BUDGET_ALERT_CHANNEL
2866 })
2867 .count()
2868 };
2869 let state = ledger.rebuild_state().unwrap();
2870 assert_eq!(state.run_usage(&run.id).total_tokens(), 1_100);
2871 assert!(state.run_ceiling_breached(&run.id).is_some());
2872 assert_eq!(budget_alerts(&state), 1);
2873 assert_eq!(
2874 state.run_status_overrides["budget-run"],
2875 FleetRunStatus::Paused
2876 );
2877
2878 // A later receipt cannot duplicate the alert…
2879 ledger
2880 .append_event_next_seq(
2881 &run.id,
2882 "worker-1",
2883 "task-a",
2884 "2026-06-12T17:00:04Z",
2885 FleetWorkerEventPayload::UsageReport {
2886 input_tokens: 50,
2887 output_tokens: 0,
2888 },
2889 )
2890 .unwrap();
2891 let state = ledger.rebuild_state().unwrap();
2892 assert_eq!(budget_alerts(&state), 1);
2893
2894 // …and the breached run admits nothing new.
2895 assert!(
2896 !ledger
2897 .lease_task_if_enqueued(
2898 &run.id,
2899 "task-b",
2900 "worker-2",
2901 "2026-06-12T17:00:05Z",
2902 None,
2903 None,
2904 )
2905 .unwrap()
2906 );
2907 }
2908
2909 #[test]
2910 fn fleet_ledger_create_and_rebuild_run() {
2911 let tmp = TempDir::new().unwrap();
2912 let ledger = FleetLedger::open(tmp.path()).unwrap();
2913 let run = sample_run("run-1");
2914 ledger.create_run(&run).unwrap();
2915 ledger
2916 .update_run_status(&run.id, FleetRunStatus::Completed, "2026-06-12T18:00:00Z")
2917 .unwrap();
2918
2919 let state = ledger.rebuild_state().unwrap();
2920 assert_eq!(state.runs.len(), 1);
2921 assert_eq!(
2922 state.run_status_overrides["run-1"],
2923 FleetRunStatus::Completed
2924 );
2925 }
2926
2927 #[test]
2928 fn fleet_event_replay_is_durable_cursor_bounded_and_privacy_safe() {
2929 let tmp = TempDir::new().unwrap();
2930 let ledger = FleetLedger::open(tmp.path()).unwrap();
2931 let mut run = sample_run("managed-run");
2932 run.target = Some(FleetRuntimeTarget::ThisComputer);
2933 run.workflow = Some(FleetWorkflowDescriptor {
2934 id: "release-check".to_string(),
2935 kind: FleetWorkflowKind::Parallel,
2936 });
2937 run.roles = vec!["reviewer".to_string()];
2938 ledger.create_run(&run).unwrap();
2939 ledger
2940 .update_run_status(&run.id, FleetRunStatus::Paused, "2026-06-12T17:00:10Z")
2941 .unwrap();
2942 ledger
2943 .update_run_status(&run.id, FleetRunStatus::Running, "2026-06-12T17:00:20Z")
2944 .unwrap();
2945 ledger
2946 .enqueue(sample_entry("managed-run", "task-a"))
2947 .unwrap();
2948 ledger
2949 .append_event(FleetWorkerEvent {
2950 seq: 1,
2951 run_id: run.id.clone(),
2952 worker_id: "worker-1".to_string(),
2953 task_id: "task-a".to_string(),
2954 timestamp: "2026-06-12T17:01:00Z".to_string(),
2955 payload: FleetWorkerEventPayload::Artifact(FleetArtifactRef {
2956 kind: FleetArtifactKind::Report,
2957 path: PathBuf::from(".codewhale/private/full-report.md"),
2958 checksum: Some("sha256:private-checksum".to_string()),
2959 mime_type: Some("text/markdown".to_string()),
2960 size_bytes: Some(42),
2961 }),
2962 extra: BTreeMap::new(),
2963 })
2964 .unwrap();
2965 ledger
2966 .record_receipt(FleetReceipt {
2967 run_id: run.id.clone(),
2968 task_id: "task-a".to_string(),
2969 worker_id: "worker-1".to_string(),
2970 attempt: Some(1),
2971 terminal_seq: Some(2),
2972 completed_at: "2026-06-12T17:02:00Z".to_string(),
2973 result: FleetTaskResult::Fail,
2974 failure_kind: Some(FleetTaskFailureKind::Task),
2975 artifacts: vec![FleetArtifactRef {
2976 kind: FleetArtifactKind::Report,
2977 path: PathBuf::from(".codewhale/private/full-report.md"),
2978 checksum: Some("sha256:private-checksum".to_string()),
2979 mime_type: Some("text/markdown".to_string()),
2980 size_bytes: Some(42),
2981 }],
2982 score: Some(FleetScore {
2983 value: 0.0,
2984 max: Some(1.0),
2985 notes: Some("verifier note contained super-secret".to_string()),
2986 }),
2987 resolved_route: None,
2988 saved_session_id: None,
2989 effective_permissions: None,
2990 })
2991 .unwrap();
2992 ledger
2993 .append_event(FleetWorkerEvent {
2994 seq: 2,
2995 run_id: run.id.clone(),
2996 worker_id: "worker-1".to_string(),
2997 task_id: "task-a".to_string(),
2998 timestamp: "2026-06-12T17:02:00Z".to_string(),
2999 payload: FleetWorkerEventPayload::Failed {
3000 // Low-entropy fixtures on purpose: realistic tokens trip
3001 // push-time secret scanners (GitGuardian flagged the originals).
3002 reason: "provider failed: api_key = super-secret; bearer sk-aaaaaaaaaaaaaaaa; authorization: Bearer aaaaaaaaaaaaaa"
3003 .to_string(),
3004 recoverable: true,
3005 },
3006 extra: BTreeMap::new(),
3007 })
3008 .unwrap();
3009
3010 let page = ledger.replay_events(&run.id, None, 100).unwrap();
3011 assert_eq!(page.run_id, run.id);
3012 assert!(!page.history_truncated);
3013 assert!(
3014 page.events
3015 .iter()
3016 .any(|event| event.event == "fleet.run.created")
3017 );
3018 assert!(
3019 page.events
3020 .iter()
3021 .any(|event| event.event == "fleet.worker.artifact")
3022 );
3023 assert!(
3024 page.events
3025 .iter()
3026 .any(|event| event.event == "fleet.worker.failed")
3027 );
3028 let encoded = serde_json::to_string(&page).unwrap();
3029 assert!(!encoded.contains("super-secret"));
3030 assert!(!encoded.contains("sk-aaaaaaaaaaaaaaaa"));
3031 assert!(!encoded.contains("aaaaaaaaaaaaaa\""));
3032 assert!(!encoded.contains("private/full-report.md"));
3033 assert!(!encoded.contains("private-checksum"));
3034
3035 let first_cursor = page.events[0].cursor.clone();
3036 drop(ledger);
3037 let reopened = FleetLedger::open(tmp.path()).unwrap();
3038 let after_restart = reopened
3039 .replay_events(&run.id, Some(&first_cursor), 100)
3040 .unwrap();
3041 assert_eq!(after_restart.events, page.events[1..]);
3042 assert!(after_restart.next_cursor.is_some());
3043
3044 let tail = reopened.replay_events(&run.id, None, 2).unwrap();
3045 assert_eq!(tail.events.len(), 2);
3046 assert!(tail.history_truncated);
3047 assert!(!tail.has_more);
3048
3049 assert!(matches!(
3050 reopened.replay_events(&run.id, Some("fev1_missing_worker"), 100),
3051 Err(FleetEventReplayError::CursorUnavailable { .. })
3052 ));
3053
3054 reopened.compact().unwrap();
3055 assert!(
3056 matches!(
3057 reopened.replay_events(&run.id, Some(&first_cursor), 100),
3058 Err(FleetEventReplayError::CursorUnavailable { .. })
3059 ),
3060 "even a retained RunCreated transition must rotate its cursor when compaction deletes intervening history"
3061 );
3062 let compacted = reopened.replay_events(&run.id, None, 100).unwrap();
3063 assert!(compacted.history_truncated);
3064 assert!(
3065 !compacted.events.is_empty(),
3066 "current projection remains replayable after compaction"
3067 );
3068 assert!(
3069 compacted
3070 .events
3071 .iter()
3072 .all(|event| !page.events.iter().any(|old| old.cursor == event.cursor)),
3073 "a compaction epoch must rotate every retained transition cursor"
3074 );
3075 }
3076
3077 #[test]
3078 fn fleet_ledger_enqueue_and_claim() {
3079 let tmp = TempDir::new().unwrap();
3080 let ledger = FleetLedger::open(tmp.path()).unwrap();
3081 ledger.create_run(&sample_run("run-1")).unwrap();
3082 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3083 ledger.enqueue(sample_entry("run-1", "task-b")).unwrap();
3084
3085 let claimed = ledger
3086 .claim_next("worker-1", &[], "2026-06-12T17:01:00Z")
3087 .unwrap();
3088 assert!(claimed.is_some());
3089 let claimed = claimed.unwrap();
3090 assert_eq!(claimed.task_id, "task-a");
3091
3092 let state = ledger.rebuild_state().unwrap();
3093 assert_eq!(state.tasks.len(), 2);
3094 assert_eq!(
3095 state.tasks["run-1:task-a"].status,
3096 FleetTaskLedgerStatus::Leased
3097 );
3098 assert_eq!(
3099 state.tasks["run-1:task-a"].leased_to.as_deref(),
3100 Some("worker-1")
3101 );
3102 assert_eq!(
3103 state.tasks["run-1:task-b"].status,
3104 FleetTaskLedgerStatus::Enqueued
3105 );
3106 }
3107
3108 #[test]
3109 fn concurrent_ledgers_allocate_unique_monotonic_event_sequences() {
3110 const WRITERS: usize = 2;
3111 const EVENTS_PER_WRITER: usize = 8;
3112
3113 let tmp = TempDir::new().unwrap();
3114 let ledger = FleetLedger::open(tmp.path()).unwrap();
3115 ledger.create_run(&sample_run("run-1")).unwrap();
3116 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3117 let root = tmp.path().to_path_buf();
3118 let barrier = Arc::new(Barrier::new(WRITERS));
3119 let handles = (0..WRITERS)
3120 .map(|_| {
3121 let root = root.clone();
3122 let barrier = Arc::clone(&barrier);
3123 thread::spawn(move || {
3124 let ledger = FleetLedger::open(&root).unwrap();
3125 barrier.wait();
3126 for _ in 0..EVENTS_PER_WRITER {
3127 ledger
3128 .append_event_next_seq(
3129 &FleetRunId::from("run-1"),
3130 "worker-1",
3131 "task-a",
3132 "2026-06-12T17:01:00Z",
3133 FleetWorkerEventPayload::Running,
3134 )
3135 .unwrap();
3136 }
3137 })
3138 })
3139 .collect::<Vec<_>>();
3140 for handle in handles {
3141 handle.join().unwrap();
3142 }
3143
3144 let mut sequences = std::fs::read_to_string(ledger.path())
3145 .unwrap()
3146 .lines()
3147 .filter_map(|line| serde_json::from_str::<FleetLedgerRecord>(line).ok())
3148 .filter_map(|record| match record {
3149 FleetLedgerRecord::EventAppended { event }
3150 if event.run_id.0 == "run-1"
3151 && event.worker_id == "worker-1"
3152 && event.task_id == "task-a" =>
3153 {
3154 Some(event.seq)
3155 }
3156 _ => None,
3157 })
3158 .collect::<Vec<_>>();
3159 sequences.sort_unstable();
3160 assert_eq!(
3161 sequences,
3162 (1..=(WRITERS * EVENTS_PER_WRITER) as u64).collect::<Vec<_>>()
3163 );
3164 }
3165
3166 #[test]
3167 fn concurrent_ledgers_claim_one_queued_task_once() {
3168 let tmp = TempDir::new().unwrap();
3169 let ledger = FleetLedger::open(tmp.path()).unwrap();
3170 ledger.create_run(&sample_run("run-1")).unwrap();
3171 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3172 let root = tmp.path().to_path_buf();
3173 let barrier = Arc::new(Barrier::new(2));
3174 let handles = ["worker-a", "worker-b"].map(|worker_id| {
3175 let root = root.clone();
3176 let barrier = Arc::clone(&barrier);
3177 thread::spawn(move || {
3178 let ledger = FleetLedger::open(&root).unwrap();
3179 barrier.wait();
3180 ledger
3181 .claim_next(worker_id, &[], "2026-06-12T17:01:00Z")
3182 .unwrap()
3183 })
3184 });
3185 let claims = handles
3186 .into_iter()
3187 .filter_map(|handle| handle.join().unwrap())
3188 .collect::<Vec<_>>();
3189
3190 assert_eq!(claims.len(), 1);
3191 assert_eq!(claims[0].task_id, "task-a");
3192 let state = ledger.rebuild_state().unwrap();
3193 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 1);
3194 assert_eq!(
3195 state.tasks["run-1:task-a"].status,
3196 FleetTaskLedgerStatus::Leased
3197 );
3198 }
3199
3200 #[test]
3201 fn concurrent_start_and_cancel_leave_one_terminal_task_without_late_progress() {
3202 let tmp = TempDir::new().unwrap();
3203 let ledger = FleetLedger::open(tmp.path()).unwrap();
3204 ledger.create_run(&sample_run("run-1")).unwrap();
3205 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3206 let root = tmp.path().to_path_buf();
3207 let barrier = Arc::new(Barrier::new(2));
3208
3209 let start_root = root.clone();
3210 let start_barrier = Arc::clone(&barrier);
3211 let starter = thread::spawn(move || {
3212 let ledger = FleetLedger::open(&start_root).unwrap();
3213 start_barrier.wait();
3214 ledger
3215 .start_task_if_enqueued(
3216 &FleetRunId::from("run-1"),
3217 "task-a",
3218 "worker-1",
3219 "2026-06-12T17:01:00Z",
3220 None,
3221 Some(1),
3222 vec![
3223 FleetWorkerEventPayload::Leased {
3224 lease_expires_at: None,
3225 },
3226 FleetWorkerEventPayload::Starting,
3227 FleetWorkerEventPayload::Artifact(FleetArtifactRef {
3228 kind: FleetArtifactKind::Log,
3229 path: PathBuf::from(".codewhale/fleet/run-1/task-a/worker-1.log"),
3230 checksum: None,
3231 mime_type: Some("text/plain".to_string()),
3232 size_bytes: Some(0),
3233 }),
3234 FleetWorkerEventPayload::Running,
3235 ],
3236 || Ok(()),
3237 )
3238 .unwrap()
3239 });
3240 let cancel_root = root.clone();
3241 let cancel_barrier = Arc::clone(&barrier);
3242 let canceller = thread::spawn(move || {
3243 let ledger = FleetLedger::open(&cancel_root).unwrap();
3244 cancel_barrier.wait();
3245 ledger
3246 .cancel_task_if_active(
3247 &FleetRunId::from("run-1"),
3248 "task-a",
3249 None,
3250 "2026-06-12T17:02:00Z",
3251 Some("operator"),
3252 Some("operator"),
3253 )
3254 .unwrap()
3255 });
3256
3257 let started = starter.join().unwrap();
3258 assert!(canceller.join().unwrap());
3259 let state = ledger.rebuild_state().unwrap();
3260 assert_eq!(
3261 state.tasks["run-1:task-a"].status,
3262 FleetTaskLedgerStatus::Cancelled
3263 );
3264 assert_eq!(
3265 state.tasks["run-1:task-a"].entry.attempts,
3266 u32::from(started)
3267 );
3268 if started {
3269 assert!(matches!(
3270 state.latest_events["worker-1:run-1:task-a"].payload,
3271 FleetWorkerEventPayload::Cancelled { .. }
3272 ));
3273 }
3274 let artifacts_before = state.artifact_events.len();
3275 assert!(
3276 ledger
3277 .append_event_if_leased(
3278 &FleetRunId::from("run-1"),
3279 "worker-1",
3280 "task-a",
3281 u32::from(started),
3282 "2026-06-12T17:03:00Z",
3283 FleetWorkerEventPayload::Artifact(FleetArtifactRef {
3284 kind: FleetArtifactKind::Log,
3285 path: PathBuf::from("late.log"),
3286 checksum: None,
3287 mime_type: None,
3288 size_bytes: None,
3289 }),
3290 )
3291 .unwrap()
3292 .is_none()
3293 );
3294 assert_eq!(
3295 ledger.rebuild_state().unwrap().artifact_events.len(),
3296 artifacts_before
3297 );
3298 }
3299
3300 #[test]
3301 fn cancelled_queue_does_not_run_start_projection_callback() {
3302 let tmp = TempDir::new().unwrap();
3303 let ledger = FleetLedger::open(tmp.path()).unwrap();
3304 ledger.create_run(&sample_run("run-1")).unwrap();
3305 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3306 assert!(
3307 ledger
3308 .cancel_task_if_active(
3309 &FleetRunId::from("run-1"),
3310 "task-a",
3311 None,
3312 "2026-06-12T17:01:00Z",
3313 None,
3314 Some("operator"),
3315 )
3316 .unwrap()
3317 );
3318 let callback_ran = AtomicBool::new(false);
3319 assert!(
3320 !ledger
3321 .start_task_if_enqueued(
3322 &FleetRunId::from("run-1"),
3323 "task-a",
3324 "worker-1",
3325 "2026-06-12T17:02:00Z",
3326 None,
3327 Some(1),
3328 vec![FleetWorkerEventPayload::Running],
3329 || {
3330 callback_ran.store(true, Ordering::SeqCst);
3331 Ok(())
3332 },
3333 )
3334 .unwrap()
3335 );
3336 assert!(!callback_ran.load(Ordering::SeqCst));
3337 }
3338
3339 #[test]
3340 fn failed_start_projection_leaves_task_unleased_and_without_events() {
3341 let tmp = TempDir::new().unwrap();
3342 let ledger = FleetLedger::open(tmp.path()).unwrap();
3343 ledger.create_run(&sample_run("run-1")).unwrap();
3344 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3345
3346 let error = ledger
3347 .start_task_if_enqueued(
3348 &FleetRunId::from("run-1"),
3349 "task-a",
3350 "worker-1",
3351 "2026-06-12T17:02:00Z",
3352 None,
3353 Some(1),
3354 vec![FleetWorkerEventPayload::Running],
3355 || Err(anyhow!("forced coordination persistence failure")),
3356 )
3357 .expect_err("projection failure must abort the lease");
3358 assert!(error.to_string().contains("forced coordination"));
3359
3360 let state = ledger.rebuild_state().unwrap();
3361 let task = &state.tasks["run-1:task-a"];
3362 assert_eq!(task.status, FleetTaskLedgerStatus::Enqueued);
3363 assert_eq!(task.entry.attempts, 0);
3364 assert!(task.leased_to.is_none());
3365 assert!(state.latest_events.is_empty());
3366 assert!(state.latest_seq.is_empty());
3367 assert!(!state.heartbeats.contains_key("worker-1"));
3368 }
3369
3370 #[test]
3371 fn concurrent_restarts_compare_and_set_one_new_attempt() {
3372 let tmp = TempDir::new().unwrap();
3373 let ledger = FleetLedger::open(tmp.path()).unwrap();
3374 ledger.create_run(&sample_run("run-1")).unwrap();
3375 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3376 assert!(
3377 ledger
3378 .start_task_if_enqueued(
3379 &FleetRunId::from("run-1"),
3380 "task-a",
3381 "worker-1",
3382 "2026-06-12T17:01:00Z",
3383 None,
3384 Some(1),
3385 vec![FleetWorkerEventPayload::Running],
3386 || Ok(()),
3387 )
3388 .unwrap()
3389 );
3390 let state = ledger.rebuild_state().unwrap();
3391 let expected_seq = state.latest_seq["worker-1:run-1:task-a"];
3392 let expected_heartbeat = state.heartbeats["worker-1"].timestamp.clone();
3393 let root = tmp.path().to_path_buf();
3394 let barrier = Arc::new(Barrier::new(2));
3395 let handles = (0..2)
3396 .map(|_| {
3397 let root = root.clone();
3398 let barrier = Arc::clone(&barrier);
3399 let expected_heartbeat = expected_heartbeat.clone();
3400 thread::spawn(move || {
3401 let ledger = FleetLedger::open(&root).unwrap();
3402 barrier.wait();
3403 ledger
3404 .restart_task_if_unchanged(
3405 &FleetRunId::from("run-1"),
3406 "task-a",
3407 "worker-1",
3408 FleetTaskLedgerStatus::Leased,
3409 1,
3410 expected_seq,
3411 Some(&expected_heartbeat),
3412 "2026-06-12T17:02:00Z",
3413 None,
3414 1,
3415 )
3416 .unwrap()
3417 })
3418 })
3419 .collect::<Vec<_>>();
3420 let winners = handles
3421 .into_iter()
3422 .map(|handle| handle.join().unwrap())
3423 .filter(|won| *won)
3424 .count();
3425
3426 assert_eq!(winners, 1);
3427 let state = ledger.rebuild_state().unwrap();
3428 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 2);
3429 assert_eq!(state.latest_seq["worker-1:run-1:task-a"], expected_seq + 2);
3430 }
3431
3432 #[test]
3433 fn stale_verifier_status_cannot_overwrite_restarted_attempt() {
3434 let tmp = TempDir::new().unwrap();
3435 let verifier = FleetLedger::open(tmp.path()).unwrap();
3436 verifier.create_run(&sample_run("run-1")).unwrap();
3437 verifier.enqueue(sample_entry("run-1", "task-a")).unwrap();
3438 assert!(
3439 verifier
3440 .start_task_if_enqueued(
3441 &FleetRunId::from("run-1"),
3442 "task-a",
3443 "worker-1",
3444 "2026-06-12T17:01:00Z",
3445 None,
3446 Some(1),
3447 vec![FleetWorkerEventPayload::Running],
3448 || Ok(()),
3449 )
3450 .unwrap()
3451 );
3452 let terminal = verifier
3453 .append_terminal_event_if_leased(
3454 &FleetRunId::from("run-1"),
3455 "worker-1",
3456 "task-a",
3457 1,
3458 "2026-06-12T17:02:00Z",
3459 FleetWorkerEventPayload::Completed {
3460 exit_code: Some(0),
3461 summary: None,
3462 },
3463 )
3464 .unwrap()
3465 .unwrap();
3466 let state = verifier.rebuild_state().unwrap();
3467 let heartbeat = state.heartbeats["worker-1"].timestamp.clone();
3468 let restarter = FleetLedger::open(tmp.path()).unwrap();
3469 assert!(
3470 restarter
3471 .restart_task_if_unchanged(
3472 &FleetRunId::from("run-1"),
3473 "task-a",
3474 "worker-1",
3475 FleetTaskLedgerStatus::Completed,
3476 1,
3477 terminal.seq,
3478 Some(&heartbeat),
3479 "2026-06-12T17:03:00Z",
3480 None,
3481 1,
3482 )
3483 .unwrap()
3484 );
3485 assert!(
3486 !verifier
3487 .mark_task_terminal_status_if_unchanged(
3488 &FleetRunId::from("run-1"),
3489 "task-a",
3490 "worker-1",
3491 FleetTaskLedgerStatus::Completed,
3492 1,
3493 terminal.seq,
3494 "2026-06-12T17:04:00Z",
3495 FleetTaskLedgerStatus::Failed,
3496 )
3497 .unwrap()
3498 );
3499 let state = verifier.rebuild_state().unwrap();
3500 assert_eq!(
3501 state.tasks["run-1:task-a"].status,
3502 FleetTaskLedgerStatus::Leased
3503 );
3504 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 2);
3505 }
3506
3507 #[test]
3508 fn fresh_heartbeat_invalidates_stale_scheduler_failure() {
3509 let tmp = TempDir::new().unwrap();
3510 let scheduler = FleetLedger::open(tmp.path()).unwrap();
3511 scheduler.create_run(&sample_run("run-1")).unwrap();
3512 scheduler.enqueue(sample_entry("run-1", "task-a")).unwrap();
3513 assert!(
3514 scheduler
3515 .start_task_if_enqueued(
3516 &FleetRunId::from("run-1"),
3517 "task-a",
3518 "worker-1",
3519 "2026-06-12T17:01:00Z",
3520 None,
3521 Some(1),
3522 vec![FleetWorkerEventPayload::Running],
3523 || Ok(()),
3524 )
3525 .unwrap()
3526 );
3527 let stale = scheduler
3528 .append_event_if_leased(
3529 &FleetRunId::from("run-1"),
3530 "worker-1",
3531 "task-a",
3532 1,
3533 "2026-06-12T17:02:00Z",
3534 FleetWorkerEventPayload::Stale {
3535 last_heartbeat_at: Some("2026-06-12T17:01:00Z".to_string()),
3536 },
3537 )
3538 .unwrap()
3539 .unwrap();
3540 let worker = FleetLedger::open(tmp.path()).unwrap();
3541 worker
3542 .heartbeat("worker-1", "2026-06-12T17:02:30Z", None, None)
3543 .unwrap();
3544
3545 assert!(
3546 scheduler
3547 .append_terminal_event_if_lease_unchanged(
3548 &FleetRunId::from("run-1"),
3549 "worker-1",
3550 "task-a",
3551 1,
3552 stale.seq,
3553 Some("2026-06-12T17:01:00Z"),
3554 "2026-06-12T17:03:00Z",
3555 FleetWorkerEventPayload::Failed {
3556 reason: "stale retry budget exhausted".to_string(),
3557 recoverable: false,
3558 },
3559 )
3560 .unwrap()
3561 .is_none()
3562 );
3563 assert_eq!(
3564 scheduler.rebuild_state().unwrap().tasks["run-1:task-a"].status,
3565 FleetTaskLedgerStatus::Leased
3566 );
3567 }
3568
3569 #[test]
3570 fn cancellation_and_completion_are_compare_and_set_terminal_transitions() {
3571 let tmp = TempDir::new().unwrap();
3572 let ledger = FleetLedger::open(tmp.path()).unwrap();
3573 ledger.create_run(&sample_run("run-1")).unwrap();
3574 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3575 assert!(
3576 ledger
3577 .lease_task_if_enqueued(
3578 &FleetRunId::from("run-1"),
3579 "task-a",
3580 "worker-1",
3581 "2026-06-12T17:01:00Z",
3582 None,
3583 Some(1),
3584 )
3585 .unwrap()
3586 );
3587 assert!(
3588 ledger
3589 .cancel_task_if_active(
3590 &FleetRunId::from("run-1"),
3591 "task-a",
3592 Some("worker-1"),
3593 "2026-06-12T17:02:00Z",
3594 Some("operator"),
3595 Some("operator"),
3596 )
3597 .unwrap()
3598 );
3599 assert!(
3600 ledger
3601 .append_terminal_event_if_leased(
3602 &FleetRunId::from("run-1"),
3603 "worker-1",
3604 "task-a",
3605 1,
3606 "2026-06-12T17:03:00Z",
3607 FleetWorkerEventPayload::Completed {
3608 exit_code: Some(0),
3609 summary: None,
3610 },
3611 )
3612 .unwrap()
3613 .is_none()
3614 );
3615 assert!(
3616 !ledger
3617 .cancel_task_if_active(
3618 &FleetRunId::from("run-1"),
3619 "task-a",
3620 Some("worker-1"),
3621 "2026-06-12T17:04:00Z",
3622 Some("operator"),
3623 Some("operator"),
3624 )
3625 .unwrap()
3626 );
3627 assert_eq!(
3628 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].status,
3629 FleetTaskLedgerStatus::Cancelled
3630 );
3631 }
3632
3633 #[test]
3634 fn fleet_ledger_survives_restart() {
3635 let tmp = TempDir::new().unwrap();
3636 {
3637 let ledger = FleetLedger::open(tmp.path()).unwrap();
3638 ledger.create_run(&sample_run("run-1")).unwrap();
3639 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3640 ledger
3641 .lease_task(
3642 &FleetRunId::from("run-1"),
3643 "task-a",
3644 "worker-1",
3645 "2026-06-12T17:01:00Z",
3646 None,
3647 )
3648 .unwrap();
3649 }
3650 // Re-open simulates process restart.
3651 let ledger = FleetLedger::open(tmp.path()).unwrap();
3652 let state = ledger.rebuild_state().unwrap();
3653 assert_eq!(state.runs.len(), 1);
3654 assert_eq!(
3655 state.tasks["run-1:task-a"].status,
3656 FleetTaskLedgerStatus::Leased
3657 );
3658 }
3659
3660 #[test]
3661 fn heartbeat_between_stale_snapshot_and_append_fences_stale_event() {
3662 let tmp = TempDir::new().unwrap();
3663 let scheduler = FleetLedger::open(tmp.path()).unwrap();
3664 scheduler.create_run(&sample_run("run-1")).unwrap();
3665 scheduler.enqueue(sample_entry("run-1", "task-a")).unwrap();
3666 assert!(
3667 scheduler
3668 .start_task_if_enqueued(
3669 &FleetRunId::from("run-1"),
3670 "task-a",
3671 "worker-1",
3672 "2026-06-12T17:01:00Z",
3673 None,
3674 Some(1),
3675 vec![FleetWorkerEventPayload::Running],
3676 || Ok(()),
3677 )
3678 .unwrap()
3679 );
3680 let snapshot = scheduler.rebuild_state().unwrap();
3681 let expected_seq = snapshot.latest_seq["worker-1:run-1:task-a"];
3682 let expected_heartbeat = snapshot.heartbeats["worker-1"].timestamp.clone();
3683
3684 FleetLedger::open(tmp.path())
3685 .unwrap()
3686 .heartbeat("worker-1", "2026-06-12T17:02:30Z", None, None)
3687 .unwrap();
3688
3689 assert!(
3690 scheduler
3691 .append_event_if_lease_unchanged(
3692 &FleetRunId::from("run-1"),
3693 "worker-1",
3694 "task-a",
3695 1,
3696 expected_seq,
3697 Some(&expected_heartbeat),
3698 "2026-06-12T17:03:00Z",
3699 FleetWorkerEventPayload::Stale {
3700 last_heartbeat_at: Some(expected_heartbeat.clone()),
3701 },
3702 )
3703 .unwrap()
3704 .is_none()
3705 );
3706 let state = scheduler.rebuild_state().unwrap();
3707 assert_eq!(
3708 state.tasks["run-1:task-a"].status,
3709 FleetTaskLedgerStatus::Leased
3710 );
3711 assert!(matches!(
3712 state.latest_events["worker-1:run-1:task-a"].payload,
3713 FleetWorkerEventPayload::Running
3714 ));
3715 }
3716
3717 #[test]
3718 fn stale_attempt_cannot_finalize_or_replace_restarted_attempt_receipt() {
3719 let tmp = TempDir::new().unwrap();
3720 let ledger = FleetLedger::open(tmp.path()).unwrap();
3721 let run_id = FleetRunId::from("run-1");
3722 ledger.create_run(&sample_run("run-1")).unwrap();
3723 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3724 assert!(
3725 ledger
3726 .start_task_if_enqueued(
3727 &run_id,
3728 "task-a",
3729 "worker-1",
3730 "2026-06-12T17:01:00Z",
3731 None,
3732 Some(1),
3733 vec![FleetWorkerEventPayload::Running],
3734 || Ok(()),
3735 )
3736 .unwrap()
3737 );
3738 let before_restart = ledger.rebuild_state().unwrap();
3739 assert!(
3740 ledger
3741 .restart_task_if_unchanged(
3742 &run_id,
3743 "task-a",
3744 "worker-1",
3745 FleetTaskLedgerStatus::Leased,
3746 1,
3747 before_restart.latest_seq["worker-1:run-1:task-a"],
3748 Some(&before_restart.heartbeats["worker-1"].timestamp),
3749 "2026-06-12T17:02:00Z",
3750 None,
3751 1,
3752 )
3753 .unwrap()
3754 );
3755
3756 let receipt = |attempt, result| FleetReceipt {
3757 run_id: run_id.clone(),
3758 task_id: "task-a".to_string(),
3759 worker_id: "worker-1".to_string(),
3760 attempt: Some(attempt),
3761 terminal_seq: None,
3762 completed_at: "2026-06-12T17:03:00Z".to_string(),
3763 result,
3764 failure_kind: None,
3765 artifacts: Vec::new(),
3766 score: None,
3767 resolved_route: None,
3768 saved_session_id: None,
3769 effective_permissions: None,
3770 };
3771 assert!(
3772 ledger
3773 .finalize_task_attempt_if_leased(
3774 &run_id,
3775 "worker-1",
3776 "task-a",
3777 1,
3778 "2026-06-12T17:03:00Z",
3779 FleetWorkerEventPayload::Completed {
3780 exit_code: Some(0),
3781 summary: Some("late attempt one".to_string()),
3782 },
3783 None,
3784 receipt(1, FleetTaskResult::Fail),
3785 )
3786 .unwrap()
3787 .is_none()
3788 );
3789 let winning_event = ledger
3790 .finalize_task_attempt_if_leased(
3791 &run_id,
3792 "worker-1",
3793 "task-a",
3794 2,
3795 "2026-06-12T17:04:00Z",
3796 FleetWorkerEventPayload::Completed {
3797 exit_code: Some(0),
3798 summary: Some("attempt two".to_string()),
3799 },
3800 None,
3801 receipt(2, FleetTaskResult::Pass),
3802 )
3803 .unwrap()
3804 .unwrap();
3805
3806 let state = ledger.rebuild_state().unwrap();
3807 let durable = &state.receipts["run-1:task-a"];
3808 assert_eq!(durable.attempt, Some(2));
3809 assert_eq!(durable.terminal_seq, Some(winning_event.seq));
3810 assert_eq!(durable.result, FleetTaskResult::Pass);
3811 assert_eq!(
3812 state.tasks["run-1:task-a"].status,
3813 FleetTaskLedgerStatus::Completed
3814 );
3815 }
3816
3817 #[test]
3818 fn fleet_ledger_quarantines_unterminated_tail_before_next_valid_record() {
3819 let tmp = TempDir::new().unwrap();
3820 let ledger = FleetLedger::open(tmp.path()).unwrap();
3821 ledger.create_run(&sample_run("run-1")).unwrap();
3822 // Simulate a process dying before its trailing newline, then use the
3823 // normal append path. The next record must not be concatenated to and
3824 // lost with the malformed crash tail.
3825 let mut file = OpenOptions::new().append(true).open(ledger.path()).unwrap();
3826 write!(file, "{{\"record\":\"run_created\",\"run\":").unwrap();
3827 file.sync_all().unwrap();
3828 drop(file);
3829 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3830
3831 let state = ledger.rebuild_state().unwrap();
3832 assert_eq!(state.runs.len(), 1);
3833 assert!(state.runs.contains_key("run-1"));
3834 assert!(state.tasks.contains_key("run-1:task-a"));
3835 }
3836
3837 #[test]
3838 fn fleet_ledger_event_and_heartbeat_reconstruct_worker_status() {
3839 let tmp = TempDir::new().unwrap();
3840 let ledger = FleetLedger::open(tmp.path()).unwrap();
3841 ledger.create_run(&sample_run("run-1")).unwrap();
3842 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3843 ledger
3844 .append_event(FleetWorkerEvent {
3845 seq: 1,
3846 run_id: FleetRunId::from("run-1"),
3847 worker_id: "worker-1".to_string(),
3848 task_id: "task-a".to_string(),
3849 timestamp: "2026-06-12T17:01:00Z".to_string(),
3850 payload: FleetWorkerEventPayload::Running,
3851 extra: BTreeMap::new(),
3852 })
3853 .unwrap();
3854 ledger
3855 .heartbeat("worker-1", "2026-06-12T17:02:00Z", Some(12.5), Some(1024))
3856 .unwrap();
3857
3858 let state = ledger.rebuild_state().unwrap();
3859 assert_eq!(state.workers["worker-1"], FleetWorkerStatus::Busy);
3860 assert_eq!(state.heartbeats["worker-1"].cpu_percent, Some(12.5));
3861 }
3862
3863 #[test]
3864 fn fleet_ledger_replays_typed_workflow_receipt_with_distinct_run_ids() {
3865 let tmp = TempDir::new().unwrap();
3866 let ledger = FleetLedger::open(tmp.path()).unwrap();
3867 ledger.create_run(&sample_run("fleet-run-1")).unwrap();
3868 ledger
3869 .enqueue(sample_entry("fleet-run-1", "task-a"))
3870 .unwrap();
3871 ledger
3872 .append_event(FleetWorkerEvent {
3873 seq: 1,
3874 run_id: FleetRunId::from("fleet-run-1"),
3875 worker_id: "worker-1".to_string(),
3876 task_id: "task-a".to_string(),
3877 timestamp: "2026-07-10T00:00:00Z".to_string(),
3878 payload: FleetWorkerEventPayload::WorkflowEvent {
3879 workflow_run_id: "workflow_1".to_string(),
3880 event: serde_json::json!({"type": "task_completed"}),
3881 },
3882 extra: BTreeMap::new(),
3883 })
3884 .unwrap();
3885
3886 let state = ledger.rebuild_state().unwrap();
3887 let event = &state.latest_events["worker-1:fleet-run-1:task-a"];
3888 assert!(matches!(
3889 &event.payload,
3890 FleetWorkerEventPayload::WorkflowEvent {
3891 workflow_run_id,
3892 event,
3893 } if workflow_run_id == "workflow_1" && event["type"] == "task_completed"
3894 ));
3895 let last_line = std::fs::read_to_string(ledger.path())
3896 .unwrap()
3897 .lines()
3898 .last()
3899 .unwrap()
3900 .to_string();
3901 assert_eq!(last_line.matches("\"run_id\"").count(), 1);
3902 assert_eq!(last_line.matches("\"workflow_run_id\"").count(), 1);
3903 }
3904
3905 #[test]
3906 fn fleet_ledger_terminal_events_ignore_late_progress_regressions() {
3907 let tmp = TempDir::new().unwrap();
3908 let ledger = FleetLedger::open(tmp.path()).unwrap();
3909 ledger.create_run(&sample_run("run-1")).unwrap();
3910 ledger
3911 .enqueue(sample_entry("run-1", "task-failed"))
3912 .unwrap();
3913 ledger
3914 .enqueue(sample_entry("run-1", "task-cancelled"))
3915 .unwrap();
3916
3917 ledger
3918 .append_event(FleetWorkerEvent {
3919 seq: 1,
3920 run_id: FleetRunId::from("run-1"),
3921 worker_id: "worker-1".to_string(),
3922 task_id: "task-failed".to_string(),
3923 timestamp: "2026-06-12T17:03:00Z".to_string(),
3924 payload: FleetWorkerEventPayload::Failed {
3925 reason: "test failed".to_string(),
3926 recoverable: false,
3927 },
3928 extra: BTreeMap::new(),
3929 })
3930 .unwrap();
3931 ledger
3932 .append_event(FleetWorkerEvent {
3933 seq: 2,
3934 run_id: FleetRunId::from("run-1"),
3935 worker_id: "worker-2".to_string(),
3936 task_id: "task-cancelled".to_string(),
3937 timestamp: "2026-06-12T17:04:00Z".to_string(),
3938 payload: FleetWorkerEventPayload::Cancelled {
3939 cancelled_by: Some("operator".to_string()),
3940 },
3941 extra: BTreeMap::new(),
3942 })
3943 .unwrap();
3944 // A live worker can flush progress after an out-of-process operator
3945 // command has already made the task terminal. Preserve the raw
3946 // sequence for append ordering without projecting the task or worker
3947 // back to a running state.
3948 ledger
3949 .append_event(FleetWorkerEvent {
3950 seq: 3,
3951 run_id: FleetRunId::from("run-1"),
3952 worker_id: "worker-2".to_string(),
3953 task_id: "task-cancelled".to_string(),
3954 timestamp: "2026-06-12T17:04:01Z".to_string(),
3955 payload: FleetWorkerEventPayload::Running,
3956 extra: BTreeMap::new(),
3957 })
3958 .unwrap();
3959
3960 let state = ledger.rebuild_state().unwrap();
3961 assert_eq!(
3962 state.tasks["run-1:task-failed"].status,
3963 FleetTaskLedgerStatus::Failed
3964 );
3965 assert_eq!(
3966 state.tasks["run-1:task-cancelled"].status,
3967 FleetTaskLedgerStatus::Cancelled
3968 );
3969 assert_eq!(state.workers["worker-2"], FleetWorkerStatus::Online);
3970 assert_eq!(
3971 state.latest_seq["worker-2:run-1:task-cancelled"], 3,
3972 "raw sequence ownership must still advance past ignored progress"
3973 );
3974 assert!(matches!(
3975 state.latest_events["worker-2:run-1:task-cancelled"].payload,
3976 FleetWorkerEventPayload::Cancelled { .. }
3977 ));
3978
3979 ledger.compact().unwrap();
3980 let state = ledger.rebuild_state().unwrap();
3981 assert_eq!(
3982 state.tasks["run-1:task-failed"].status,
3983 FleetTaskLedgerStatus::Failed
3984 );
3985 assert_eq!(
3986 state.tasks["run-1:task-cancelled"].status,
3987 FleetTaskLedgerStatus::Cancelled
3988 );
3989 assert_eq!(state.workers["worker-2"], FleetWorkerStatus::Online);
3990 assert_eq!(
3991 state.latest_seq["worker-2:run-1:task-cancelled"], 3,
3992 "compaction must preserve the ignored event sequence high-water mark"
3993 );
3994 assert!(matches!(
3995 state.latest_events["worker-2:run-1:task-cancelled"].payload,
3996 FleetWorkerEventPayload::Cancelled { .. }
3997 ));
3998 let next = ledger
3999 .append_event_next_seq(
4000 &FleetRunId::from("run-1"),
4001 "worker-2",
4002 "task-cancelled",
4003 "2026-06-12T17:04:02Z",
4004 FleetWorkerEventPayload::Running,
4005 )
4006 .unwrap();
4007 assert_eq!(next.seq, 4, "compaction must never permit sequence reuse");
4008 }
4009
4010 #[test]
4011 fn fleet_ledger_compact_preserves_current_state() {
4012 let tmp = TempDir::new().unwrap();
4013 let ledger = FleetLedger::open(tmp.path()).unwrap();
4014 ledger.create_run(&sample_run("run-1")).unwrap();
4015 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4016 ledger
4017 .lease_task(
4018 &FleetRunId::from("run-1"),
4019 "task-a",
4020 "worker-1",
4021 "2026-06-12T17:01:00Z",
4022 None,
4023 )
4024 .unwrap();
4025 ledger
4026 .append_event(FleetWorkerEvent {
4027 seq: 7,
4028 run_id: FleetRunId::from("run-1"),
4029 worker_id: "worker-1".to_string(),
4030 task_id: "task-a".to_string(),
4031 timestamp: "2026-06-12T17:01:30Z".to_string(),
4032 payload: FleetWorkerEventPayload::Running,
4033 extra: BTreeMap::new(),
4034 })
4035 .unwrap();
4036 ledger
4037 .heartbeat("worker-1", "2026-06-12T17:02:00Z", Some(12.5), Some(1024))
4038 .unwrap();
4039 ledger
4040 .record_receipt(FleetReceipt {
4041 run_id: FleetRunId::from("run-1"),
4042 task_id: "task-a".to_string(),
4043 worker_id: "worker-1".to_string(),
4044 attempt: Some(1),
4045 terminal_seq: None,
4046 completed_at: "2026-06-12T17:03:00Z".to_string(),
4047 result: FleetTaskResult::Pass,
4048 failure_kind: None,
4049 artifacts: vec![],
4050 score: None,
4051 resolved_route: None,
4052 saved_session_id: None,
4053 effective_permissions: None,
4054 })
4055 .unwrap();
4056
4057 let lifecycle_seq_before_compaction =
4058 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].lifecycle_seq;
4059 assert_eq!(
4060 lifecycle_seq_before_compaction, 2,
4061 "enqueue and lease are the two effective owner states"
4062 );
4063
4064 ledger.compact().unwrap();
4065 let contents = std::fs::read_to_string(ledger.path()).unwrap();
4066 assert!(contents.lines().count() >= 5, "{contents}");
4067
4068 let state = ledger.rebuild_state().unwrap();
4069 assert_eq!(state.runs.len(), 1);
4070 assert_eq!(
4071 state.tasks["run-1:task-a"].status,
4072 FleetTaskLedgerStatus::Leased
4073 );
4074 assert_eq!(state.workers["worker-1"], FleetWorkerStatus::Busy);
4075 assert_eq!(state.heartbeats["worker-1"].memory_mb, Some(1024));
4076 assert_eq!(
4077 state.tasks["run-1:task-a"].entry.attempts, 1,
4078 "compaction must not mint a synthetic retry attempt"
4079 );
4080 assert_eq!(
4081 state.tasks["run-1:task-a"].lifecycle_seq, lifecycle_seq_before_compaction,
4082 "compaction must not mint an owner lifecycle transition"
4083 );
4084 assert!(state.latest_seq.values().any(|seq| *seq == 7));
4085 assert_eq!(state.receipts["run-1:task-a"].result, FleetTaskResult::Pass);
4086 }
4087
4088 #[test]
4089 fn fleet_compaction_preserves_multilease_lifecycle_high_water() {
4090 let tmp = TempDir::new().unwrap();
4091 let ledger = FleetLedger::open(tmp.path()).unwrap();
4092 ledger.create_run(&sample_run("run-1")).unwrap();
4093 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4094 ledger
4095 .lease_task(
4096 &FleetRunId::from("run-1"),
4097 "task-a",
4098 "worker-1",
4099 "2026-06-12T17:01:00Z",
4100 None,
4101 )
4102 .unwrap();
4103 ledger
4104 .lease_task(
4105 &FleetRunId::from("run-1"),
4106 "task-a",
4107 "worker-1",
4108 "2026-06-12T17:02:00Z",
4109 None,
4110 )
4111 .unwrap();
4112 let before = ledger.rebuild_state().unwrap();
4113 assert_eq!(before.tasks["run-1:task-a"].lifecycle_seq, 3);
4114
4115 ledger.compact().unwrap();
4116 let after = ledger.rebuild_state().unwrap();
4117 assert_eq!(
4118 after.tasks["run-1:task-a"].lifecycle_seq, 3,
4119 "compaction must not reuse lower Work Graph idempotency keys"
4120 );
4121 }
4122
4123 #[test]
4124 fn compaction_preserves_verifier_failure_override_after_completed_exit() {
4125 let tmp = TempDir::new().unwrap();
4126 let ledger = FleetLedger::open(tmp.path()).unwrap();
4127 let run_id = FleetRunId::from("run-1");
4128 ledger.create_run(&sample_run("run-1")).unwrap();
4129 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4130 assert!(
4131 ledger
4132 .start_task_if_enqueued(
4133 &run_id,
4134 "task-a",
4135 "worker-1",
4136 "2026-06-12T17:01:00Z",
4137 None,
4138 Some(1),
4139 vec![FleetWorkerEventPayload::Running],
4140 || Ok(()),
4141 )
4142 .unwrap()
4143 );
4144 let terminal = ledger
4145 .finalize_task_attempt_if_leased(
4146 &run_id,
4147 "worker-1",
4148 "task-a",
4149 1,
4150 "2026-06-12T17:02:00Z",
4151 FleetWorkerEventPayload::Completed {
4152 exit_code: Some(0),
4153 summary: Some("process succeeded but verification failed".to_string()),
4154 },
4155 Some(FleetTaskLedgerStatus::Failed),
4156 FleetReceipt {
4157 run_id: run_id.clone(),
4158 task_id: "task-a".to_string(),
4159 worker_id: "worker-1".to_string(),
4160 attempt: Some(1),
4161 terminal_seq: None,
4162 completed_at: "2026-06-12T17:02:00Z".to_string(),
4163 result: FleetTaskResult::Fail,
4164 failure_kind: Some(FleetTaskFailureKind::Verifier),
4165 artifacts: Vec::new(),
4166 score: None,
4167 resolved_route: None,
4168 saved_session_id: None,
4169 effective_permissions: None,
4170 },
4171 )
4172 .unwrap()
4173 .unwrap();
4174 assert_eq!(
4175 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].status,
4176 FleetTaskLedgerStatus::Failed
4177 );
4178
4179 ledger.compact().unwrap();
4180 let reopened = FleetLedger::open(tmp.path()).unwrap();
4181 let state = reopened.rebuild_state().unwrap();
4182 assert_eq!(
4183 state.tasks["run-1:task-a"].status,
4184 FleetTaskLedgerStatus::Failed
4185 );
4186 assert_eq!(state.receipts["run-1:task-a"].attempt, Some(1));
4187 assert_eq!(
4188 state.receipts["run-1:task-a"].terminal_seq,
4189 Some(terminal.seq)
4190 );
4191 assert!(matches!(
4192 state.latest_events["worker-1:run-1:task-a"].payload,
4193 FleetWorkerEventPayload::Completed { .. }
4194 ));
4195 }
4196
4197 #[test]
4198 fn legacy_attemptless_alert_suppresses_duplicate_delivery_after_upgrade() {
4199 let tmp = TempDir::new().unwrap();
4200 let ledger = FleetLedger::open(tmp.path()).unwrap();
4201 let run_id = FleetRunId::from("run-1");
4202 ledger.create_run(&sample_run("run-1")).unwrap();
4203 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4204 assert!(
4205 ledger
4206 .start_task_if_enqueued(
4207 &run_id,
4208 "task-a",
4209 "worker-1",
4210 "2026-06-12T17:01:00Z",
4211 None,
4212 Some(1),
4213 vec![FleetWorkerEventPayload::Running],
4214 || Ok(()),
4215 )
4216 .unwrap()
4217 );
4218 ledger
4219 .append_terminal_event_if_leased(
4220 &run_id,
4221 "worker-1",
4222 "task-a",
4223 1,
4224 "2026-06-12T17:02:00Z",
4225 FleetWorkerEventPayload::Failed {
4226 reason: "attempt exhausted before upgrade".to_string(),
4227 recoverable: false,
4228 },
4229 )
4230 .unwrap()
4231 .unwrap();
4232 ledger
4233 .record_alert(&run_id, "task-a", "slack", "2026-06-12T17:02:01Z")
4234 .unwrap();
4235 assert!(
4236 !ledger
4237 .record_failed_attempt_alert_once(
4238 &run_id,
4239 "task-a",
4240 "worker-1",
4241 1,
4242 "slack",
4243 "slack#0",
4244 "2026-06-12T17:03:00Z",
4245 )
4246 .unwrap()
4247 );
4248 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 1);
4249
4250 ledger.compact().unwrap();
4251 assert!(
4252 !ledger
4253 .record_failed_attempt_alert_once(
4254 &run_id,
4255 "task-a",
4256 "worker-1",
4257 1,
4258 "slack",
4259 "slack#0",
4260 "2026-06-12T17:04:00Z",
4261 )
4262 .unwrap()
4263 );
4264 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 1);
4265
4266 let attempt_one = ledger.rebuild_state().unwrap();
4267 assert!(
4268 ledger
4269 .restart_task_if_unchanged(
4270 &run_id,
4271 "task-a",
4272 "worker-1",
4273 FleetTaskLedgerStatus::Failed,
4274 1,
4275 attempt_one.latest_seq["worker-1:run-1:task-a"],
4276 Some(&attempt_one.heartbeats["worker-1"].timestamp),
4277 "2026-06-12T17:05:00Z",
4278 None,
4279 1,
4280 )
4281 .unwrap()
4282 );
4283 ledger
4284 .append_terminal_event_if_leased(
4285 &run_id,
4286 "worker-1",
4287 "task-a",
4288 2,
4289 "2026-06-12T17:06:00Z",
4290 FleetWorkerEventPayload::Failed {
4291 reason: "attempt two also exhausted".to_string(),
4292 recoverable: false,
4293 },
4294 )
4295 .unwrap()
4296 .unwrap();
4297 assert!(
4298 ledger
4299 .record_failed_attempt_alert_once(
4300 &run_id,
4301 "task-a",
4302 "worker-1",
4303 2,
4304 "slack",
4305 "slack#0",
4306 "2026-06-12T17:07:00Z",
4307 )
4308 .unwrap()
4309 );
4310 assert!(
4311 !ledger
4312 .record_failed_attempt_alert_once(
4313 &run_id,
4314 "task-a",
4315 "worker-1",
4316 2,
4317 "slack",
4318 "slack#0",
4319 "2026-06-12T17:08:00Z",
4320 )
4321 .unwrap()
4322 );
4323 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 2);
4324 }
4325
4326 #[test]
4327 fn compaction_lock_keeps_concurrent_append_on_replacement_ledger() {
4328 let tmp = TempDir::new().unwrap();
4329 let ledger = FleetLedger::open(tmp.path()).unwrap();
4330 ledger.create_run(&sample_run("run-1")).unwrap();
4331 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4332
4333 let root = tmp.path().to_path_buf();
4334 let (snapshot_tx, snapshot_rx) = mpsc::sync_channel(0);
4335 let (release_tx, release_rx) = mpsc::sync_channel(0);
4336 let compact_root = root.clone();
4337 let compactor = thread::spawn(move || {
4338 let ledger = FleetLedger::open(&compact_root).unwrap();
4339 ledger
4340 .compact_with_snapshot_hook(|| {
4341 snapshot_tx.send(()).unwrap();
4342 release_rx.recv().unwrap();
4343 })
4344 .unwrap();
4345 });
4346 snapshot_rx
4347 .recv_timeout(Duration::from_secs(5))
4348 .expect("compaction never reached its locked snapshot");
4349
4350 let contender = FleetLedger::open(&root).unwrap();
4351 let lock_file = contender.open_lock_file().unwrap();
4352 let mut lock = fd_lock::RwLock::new(lock_file);
4353 match lock.try_write() {
4354 Err(err) => assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock),
4355 Ok(_) => panic!("compaction snapshot did not retain the Fleet ledger lock"),
4356 }
4357
4358 let (append_started_tx, append_started_rx) = mpsc::sync_channel(0);
4359 let (append_done_tx, append_done_rx) = mpsc::sync_channel(0);
4360 let append_root = root.clone();
4361 let appender = thread::spawn(move || {
4362 let ledger = FleetLedger::open(&append_root).unwrap();
4363 append_started_tx.send(()).unwrap();
4364 let event = ledger
4365 .append_event_next_seq(
4366 &FleetRunId::from("run-1"),
4367 "worker-1",
4368 "task-a",
4369 "2026-06-12T17:01:00Z",
4370 FleetWorkerEventPayload::Running,
4371 )
4372 .unwrap();
4373 append_done_tx.send(event.seq).unwrap();
4374 });
4375 append_started_rx.recv().unwrap();
4376 assert!(
4377 append_done_rx
4378 .recv_timeout(Duration::from_millis(100))
4379 .is_err(),
4380 "append completed while compaction still held the ledger lock"
4381 );
4382
4383 release_tx.send(()).unwrap();
4384 compactor.join().unwrap();
4385 assert_eq!(
4386 append_done_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
4387 1
4388 );
4389 appender.join().unwrap();
4390
4391 let state = ledger.rebuild_state().unwrap();
4392 assert_eq!(state.latest_seq["worker-1:run-1:task-a"], 1);
4393 assert!(matches!(
4394 &state.latest_events["worker-1:run-1:task-a"].payload,
4395 FleetWorkerEventPayload::Running
4396 ));
4397 }
4398
4399 #[test]
4400 fn fleet_ledger_receipt_round_trip() {
4401 let tmp = TempDir::new().unwrap();
4402 let ledger = FleetLedger::open(tmp.path()).unwrap();
4403 let receipt = FleetReceipt {
4404 run_id: FleetRunId::from("run-1"),
4405 task_id: "task-a".to_string(),
4406 worker_id: "worker-1".to_string(),
4407 attempt: None,
4408 terminal_seq: None,
4409 completed_at: "2026-06-12T17:03:00Z".to_string(),
4410 result: FleetTaskResult::Pass,
4411 failure_kind: None,
4412 artifacts: vec![],
4413 score: None,
4414 resolved_route: None,
4415 saved_session_id: None,
4416 effective_permissions: None,
4417 };
4418 ledger.record_receipt(receipt.clone()).unwrap();
4419 let state = ledger.rebuild_state().unwrap();
4420 assert_eq!(state.receipts["run-1:task-a"].result, FleetTaskResult::Pass);
4421 }
4422 }
4423
4423 lines RUST