返回 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 #[test]
2632 fn ledger_rejects_replaced_lock_identity() {
2633 let workspace = TempDir::new().unwrap();
2634 let ledger = FleetLedger::open(workspace.path()).unwrap();
2635 std::fs::remove_file(&ledger.lock_path).unwrap();
2636 std::fs::write(&ledger.lock_path, b"").unwrap();
2637 assert!(ledger.rebuild_state().is_err());
2638 assert!(ledger.compact().is_err());
2639 }
2640
2641 #[test]
2642 fn ledger_rejects_hard_linked_ledger_and_lock() {
2643 for name in ["fleet.jsonl", "fleet.lock"] {
2644 let workspace = TempDir::new().unwrap();
2645 let outside = TempDir::new().unwrap();
2646 std::fs::create_dir(workspace.path().join(".codewhale")).unwrap();
2647 let canary = outside.path().join("canary.txt");
2648 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_HARD_LINK_CANARY").unwrap();
2649 std::fs::hard_link(&canary, workspace.path().join(".codewhale").join(name)).unwrap();
2650 assert!(FleetLedger::open(workspace.path()).is_err());
2651 assert_eq!(
2652 std::fs::read(canary).unwrap(),
2653 b"OUTSIDE_SYNTHETIC_HARD_LINK_CANARY"
2654 );
2655 }
2656 }
2657
2658 #[cfg(unix)]
2659 #[test]
2660 fn ledger_rejects_parent_and_final_linked_paths() {
2661 use std::os::unix::fs::symlink;
2662 let workspace = TempDir::new().unwrap();
2663 let outside = TempDir::new().unwrap();
2664 symlink(outside.path(), workspace.path().join(".codewhale")).unwrap();
2665 assert!(FleetLedger::open(workspace.path()).is_err());
2666 assert!(!outside.path().join("fleet.jsonl").exists());
2667 assert!(!outside.path().join("fleet.lock").exists());
2668 for name in ["fleet.jsonl", "fleet.lock"] {
2669 for hard_link in [false, true] {
2670 let workspace = TempDir::new().unwrap();
2671 std::fs::create_dir(workspace.path().join(".codewhale")).unwrap();
2672 let outside = TempDir::new().unwrap();
2673 let canary = outside.path().join("canary.txt");
2674 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_LEDGER_CANARY").unwrap();
2675 let linked = workspace.path().join(".codewhale").join(name);
2676 if hard_link {
2677 std::fs::hard_link(&canary, linked).unwrap();
2678 } else {
2679 symlink(&canary, linked).unwrap();
2680 }
2681 assert!(
2682 FleetLedger::open(workspace.path()).is_err(),
2683 "accepted linked {name}"
2684 );
2685 assert_eq!(
2686 std::fs::read(canary).unwrap(),
2687 b"OUTSIDE_SYNTHETIC_LEDGER_CANARY"
2688 );
2689 }
2690 }
2691 }
2692
2693 #[cfg(unix)]
2694 #[test]
2695 fn ledger_compaction_ignores_preexisting_temporary_links() {
2696 use std::os::unix::fs::symlink;
2697 let workspace = TempDir::new().unwrap();
2698 let outside = TempDir::new().unwrap();
2699 let ledger = FleetLedger::open(workspace.path()).unwrap();
2700 let canary = outside.path().join("canary.txt");
2701 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_COMPACTION_CANARY").unwrap();
2702 symlink(&canary, ledger.path().with_extension(".tmp")).unwrap();
2703 ledger.compact().unwrap();
2704 assert_eq!(
2705 std::fs::read(canary).unwrap(),
2706 b"OUTSIDE_SYNTHETIC_COMPACTION_CANARY"
2707 );
2708 assert!(ledger.rebuild_state().unwrap().runs.is_empty());
2709 assert!(
2710 !std::fs::symlink_metadata(ledger.path())
2711 .unwrap()
2712 .file_type()
2713 .is_symlink()
2714 );
2715 }
2716
2717 #[cfg(unix)]
2718 #[test]
2719 fn ledger_rejects_linked_ledger_and_replaced_lock_after_open() {
2720 use std::os::unix::fs::symlink;
2721 let workspace = TempDir::new().unwrap();
2722 let outside = TempDir::new().unwrap();
2723 let ledger = FleetLedger::open(workspace.path()).unwrap();
2724 let canary = outside.path().join("canary.txt");
2725 std::fs::write(&canary, b"OUTSIDE_SYNTHETIC_LEDGER_CANARY").unwrap();
2726 std::fs::remove_file(ledger.path()).unwrap();
2727 symlink(&canary, ledger.path()).unwrap();
2728 assert!(ledger.compact().is_err());
2729 assert!(ledger.rebuild_state().is_err());
2730 assert_eq!(
2731 std::fs::read(&canary).unwrap(),
2732 b"OUTSIDE_SYNTHETIC_LEDGER_CANARY"
2733 );
2734 std::fs::remove_file(ledger.path()).unwrap();
2735 std::fs::write(ledger.path(), b"").unwrap();
2736 std::fs::remove_file(&ledger.lock_path).unwrap();
2737 std::fs::write(&ledger.lock_path, b"").unwrap();
2738 assert!(ledger.compact().is_err());
2739 assert!(ledger.rebuild_state().is_err());
2740 }
2741
2742 #[cfg(unix)]
2743 #[test]
2744 fn ledger_uses_pinned_parent_after_workspace_path_swap() {
2745 use std::os::unix::fs::symlink;
2746 let workspace = TempDir::new().unwrap();
2747 let outside = TempDir::new().unwrap();
2748 let ledger = FleetLedger::open(workspace.path()).unwrap();
2749 std::fs::rename(
2750 workspace.path().join(".codewhale"),
2751 workspace.path().join("retained"),
2752 )
2753 .unwrap();
2754 symlink(outside.path(), workspace.path().join(".codewhale")).unwrap();
2755 ledger.compact().unwrap();
2756 assert!(
2757 std::fs::read(workspace.path().join("retained/fleet.jsonl"))
2758 .unwrap()
2759 .starts_with(b"{\"record\":\"replay_epoch\"")
2760 );
2761 assert!(!outside.path().join("fleet.jsonl").exists());
2762 assert!(!outside.path().join("fleet.lock").exists());
2763 }
2764
2765 fn sample_run(id: &str) -> FleetRun {
2766 FleetRun {
2767 id: FleetRunId::from(id),
2768 name: "smoke".to_string(),
2769 status: FleetRunStatus::Running,
2770 target: None,
2771 workflow: None,
2772 roles: Vec::new(),
2773 max_workers: None,
2774 usage_ceiling: None,
2775 task_specs: vec![],
2776 worker_specs: vec![],
2777 labels: BTreeMap::new(),
2778 security_policy: None,
2779 created_at: "2026-06-12T17:00:00Z".to_string(),
2780 updated_at: None,
2781 completed_at: None,
2782 }
2783 }
2784
2785 fn sample_entry(run_id: &str, task_id: &str) -> FleetInboxEntry {
2786 FleetInboxEntry {
2787 run_id: FleetRunId::from(run_id),
2788 task_id: task_id.to_string(),
2789 priority: 0,
2790 enqueued_at: "2026-06-12T17:00:00Z".to_string(),
2791 lease_deadline: None,
2792 attempts: 0,
2793 }
2794 }
2795
2796 #[test]
2797 fn usage_ceiling_breach_pauses_run_refuses_admissions_and_alerts_once() {
2798 let tmp = TempDir::new().unwrap();
2799 let ledger = FleetLedger::open(tmp.path()).unwrap();
2800 let mut run = sample_run("budget-run");
2801 run.usage_ceiling = Some(codewhale_protocol::fleet::FleetUsageCeiling {
2802 max_total_tokens: 1_000,
2803 });
2804 ledger.create_run(&run).unwrap();
2805 ledger
2806 .enqueue(sample_entry("budget-run", "task-a"))
2807 .unwrap();
2808 ledger
2809 .enqueue(sample_entry("budget-run", "task-b"))
2810 .unwrap();
2811
2812 // Under the ceiling, admission works.
2813 assert!(
2814 ledger
2815 .lease_task_if_enqueued(
2816 &run.id,
2817 "task-a",
2818 "worker-1",
2819 "2026-06-12T17:00:01Z",
2820 None,
2821 None,
2822 )
2823 .unwrap()
2824 );
2825
2826 // Usage receipts accumulate; the second crosses the 1 000 ceiling and
2827 // must pause the run and fire exactly one durable alert.
2828 for (ts, input, output) in [
2829 ("2026-06-12T17:00:02Z", 600, 0),
2830 ("2026-06-12T17:00:03Z", 300, 200),
2831 ] {
2832 ledger
2833 .append_event_next_seq(
2834 &run.id,
2835 "worker-1",
2836 "task-a",
2837 ts,
2838 FleetWorkerEventPayload::UsageReport {
2839 input_tokens: input,
2840 output_tokens: output,
2841 },
2842 )
2843 .unwrap();
2844 }
2845 let budget_alerts = |state: &FleetLedgerState| {
2846 state
2847 .alerts
2848 .keys()
2849 .filter(|(run, task, _, channel)| {
2850 run == "budget-run"
2851 && task == BUDGET_ALERT_TASK
2852 && channel == BUDGET_ALERT_CHANNEL
2853 })
2854 .count()
2855 };
2856 let state = ledger.rebuild_state().unwrap();
2857 assert_eq!(state.run_usage(&run.id).total_tokens(), 1_100);
2858 assert!(state.run_ceiling_breached(&run.id).is_some());
2859 assert_eq!(budget_alerts(&state), 1);
2860 assert_eq!(
2861 state.run_status_overrides["budget-run"],
2862 FleetRunStatus::Paused
2863 );
2864
2865 // A later receipt cannot duplicate the alert…
2866 ledger
2867 .append_event_next_seq(
2868 &run.id,
2869 "worker-1",
2870 "task-a",
2871 "2026-06-12T17:00:04Z",
2872 FleetWorkerEventPayload::UsageReport {
2873 input_tokens: 50,
2874 output_tokens: 0,
2875 },
2876 )
2877 .unwrap();
2878 let state = ledger.rebuild_state().unwrap();
2879 assert_eq!(budget_alerts(&state), 1);
2880
2881 // …and the breached run admits nothing new.
2882 assert!(
2883 !ledger
2884 .lease_task_if_enqueued(
2885 &run.id,
2886 "task-b",
2887 "worker-2",
2888 "2026-06-12T17:00:05Z",
2889 None,
2890 None,
2891 )
2892 .unwrap()
2893 );
2894 }
2895
2896 #[test]
2897 fn fleet_ledger_create_and_rebuild_run() {
2898 let tmp = TempDir::new().unwrap();
2899 let ledger = FleetLedger::open(tmp.path()).unwrap();
2900 let run = sample_run("run-1");
2901 ledger.create_run(&run).unwrap();
2902 ledger
2903 .update_run_status(&run.id, FleetRunStatus::Completed, "2026-06-12T18:00:00Z")
2904 .unwrap();
2905
2906 let state = ledger.rebuild_state().unwrap();
2907 assert_eq!(state.runs.len(), 1);
2908 assert_eq!(
2909 state.run_status_overrides["run-1"],
2910 FleetRunStatus::Completed
2911 );
2912 }
2913
2914 #[test]
2915 fn fleet_event_replay_is_durable_cursor_bounded_and_privacy_safe() {
2916 let tmp = TempDir::new().unwrap();
2917 let ledger = FleetLedger::open(tmp.path()).unwrap();
2918 let mut run = sample_run("managed-run");
2919 run.target = Some(FleetRuntimeTarget::ThisComputer);
2920 run.workflow = Some(FleetWorkflowDescriptor {
2921 id: "release-check".to_string(),
2922 kind: FleetWorkflowKind::Parallel,
2923 });
2924 run.roles = vec!["reviewer".to_string()];
2925 ledger.create_run(&run).unwrap();
2926 ledger
2927 .update_run_status(&run.id, FleetRunStatus::Paused, "2026-06-12T17:00:10Z")
2928 .unwrap();
2929 ledger
2930 .update_run_status(&run.id, FleetRunStatus::Running, "2026-06-12T17:00:20Z")
2931 .unwrap();
2932 ledger
2933 .enqueue(sample_entry("managed-run", "task-a"))
2934 .unwrap();
2935 ledger
2936 .append_event(FleetWorkerEvent {
2937 seq: 1,
2938 run_id: run.id.clone(),
2939 worker_id: "worker-1".to_string(),
2940 task_id: "task-a".to_string(),
2941 timestamp: "2026-06-12T17:01:00Z".to_string(),
2942 payload: FleetWorkerEventPayload::Artifact(FleetArtifactRef {
2943 kind: FleetArtifactKind::Report,
2944 path: PathBuf::from(".codewhale/private/full-report.md"),
2945 checksum: Some("sha256:private-checksum".to_string()),
2946 mime_type: Some("text/markdown".to_string()),
2947 size_bytes: Some(42),
2948 }),
2949 extra: BTreeMap::new(),
2950 })
2951 .unwrap();
2952 ledger
2953 .record_receipt(FleetReceipt {
2954 run_id: run.id.clone(),
2955 task_id: "task-a".to_string(),
2956 worker_id: "worker-1".to_string(),
2957 attempt: Some(1),
2958 terminal_seq: Some(2),
2959 completed_at: "2026-06-12T17:02:00Z".to_string(),
2960 result: FleetTaskResult::Fail,
2961 failure_kind: Some(FleetTaskFailureKind::Task),
2962 artifacts: vec![FleetArtifactRef {
2963 kind: FleetArtifactKind::Report,
2964 path: PathBuf::from(".codewhale/private/full-report.md"),
2965 checksum: Some("sha256:private-checksum".to_string()),
2966 mime_type: Some("text/markdown".to_string()),
2967 size_bytes: Some(42),
2968 }],
2969 score: Some(FleetScore {
2970 value: 0.0,
2971 max: Some(1.0),
2972 notes: Some("verifier note contained super-secret".to_string()),
2973 }),
2974 resolved_route: None,
2975 saved_session_id: None,
2976 effective_permissions: None,
2977 })
2978 .unwrap();
2979 ledger
2980 .append_event(FleetWorkerEvent {
2981 seq: 2,
2982 run_id: run.id.clone(),
2983 worker_id: "worker-1".to_string(),
2984 task_id: "task-a".to_string(),
2985 timestamp: "2026-06-12T17:02:00Z".to_string(),
2986 payload: FleetWorkerEventPayload::Failed {
2987 // Low-entropy fixtures on purpose: realistic tokens trip
2988 // push-time secret scanners (GitGuardian flagged the originals).
2989 reason: "provider failed: api_key = super-secret; bearer sk-aaaaaaaaaaaaaaaa; authorization: Bearer aaaaaaaaaaaaaa"
2990 .to_string(),
2991 recoverable: true,
2992 },
2993 extra: BTreeMap::new(),
2994 })
2995 .unwrap();
2996
2997 let page = ledger.replay_events(&run.id, None, 100).unwrap();
2998 assert_eq!(page.run_id, run.id);
2999 assert!(!page.history_truncated);
3000 assert!(
3001 page.events
3002 .iter()
3003 .any(|event| event.event == "fleet.run.created")
3004 );
3005 assert!(
3006 page.events
3007 .iter()
3008 .any(|event| event.event == "fleet.worker.artifact")
3009 );
3010 assert!(
3011 page.events
3012 .iter()
3013 .any(|event| event.event == "fleet.worker.failed")
3014 );
3015 let encoded = serde_json::to_string(&page).unwrap();
3016 assert!(!encoded.contains("super-secret"));
3017 assert!(!encoded.contains("sk-aaaaaaaaaaaaaaaa"));
3018 assert!(!encoded.contains("aaaaaaaaaaaaaa\""));
3019 assert!(!encoded.contains("private/full-report.md"));
3020 assert!(!encoded.contains("private-checksum"));
3021
3022 let first_cursor = page.events[0].cursor.clone();
3023 drop(ledger);
3024 let reopened = FleetLedger::open(tmp.path()).unwrap();
3025 let after_restart = reopened
3026 .replay_events(&run.id, Some(&first_cursor), 100)
3027 .unwrap();
3028 assert_eq!(after_restart.events, page.events[1..]);
3029 assert!(after_restart.next_cursor.is_some());
3030
3031 let tail = reopened.replay_events(&run.id, None, 2).unwrap();
3032 assert_eq!(tail.events.len(), 2);
3033 assert!(tail.history_truncated);
3034 assert!(!tail.has_more);
3035
3036 assert!(matches!(
3037 reopened.replay_events(&run.id, Some("fev1_missing_worker"), 100),
3038 Err(FleetEventReplayError::CursorUnavailable { .. })
3039 ));
3040
3041 reopened.compact().unwrap();
3042 assert!(
3043 matches!(
3044 reopened.replay_events(&run.id, Some(&first_cursor), 100),
3045 Err(FleetEventReplayError::CursorUnavailable { .. })
3046 ),
3047 "even a retained RunCreated transition must rotate its cursor when compaction deletes intervening history"
3048 );
3049 let compacted = reopened.replay_events(&run.id, None, 100).unwrap();
3050 assert!(compacted.history_truncated);
3051 assert!(
3052 !compacted.events.is_empty(),
3053 "current projection remains replayable after compaction"
3054 );
3055 assert!(
3056 compacted
3057 .events
3058 .iter()
3059 .all(|event| !page.events.iter().any(|old| old.cursor == event.cursor)),
3060 "a compaction epoch must rotate every retained transition cursor"
3061 );
3062 }
3063
3064 #[test]
3065 fn fleet_ledger_enqueue_and_claim() {
3066 let tmp = TempDir::new().unwrap();
3067 let ledger = FleetLedger::open(tmp.path()).unwrap();
3068 ledger.create_run(&sample_run("run-1")).unwrap();
3069 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3070 ledger.enqueue(sample_entry("run-1", "task-b")).unwrap();
3071
3072 let claimed = ledger
3073 .claim_next("worker-1", &[], "2026-06-12T17:01:00Z")
3074 .unwrap();
3075 assert!(claimed.is_some());
3076 let claimed = claimed.unwrap();
3077 assert_eq!(claimed.task_id, "task-a");
3078
3079 let state = ledger.rebuild_state().unwrap();
3080 assert_eq!(state.tasks.len(), 2);
3081 assert_eq!(
3082 state.tasks["run-1:task-a"].status,
3083 FleetTaskLedgerStatus::Leased
3084 );
3085 assert_eq!(
3086 state.tasks["run-1:task-a"].leased_to.as_deref(),
3087 Some("worker-1")
3088 );
3089 assert_eq!(
3090 state.tasks["run-1:task-b"].status,
3091 FleetTaskLedgerStatus::Enqueued
3092 );
3093 }
3094
3095 #[test]
3096 fn concurrent_ledgers_allocate_unique_monotonic_event_sequences() {
3097 const WRITERS: usize = 2;
3098 const EVENTS_PER_WRITER: usize = 8;
3099
3100 let tmp = TempDir::new().unwrap();
3101 let ledger = FleetLedger::open(tmp.path()).unwrap();
3102 ledger.create_run(&sample_run("run-1")).unwrap();
3103 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3104 let root = tmp.path().to_path_buf();
3105 let barrier = Arc::new(Barrier::new(WRITERS));
3106 let handles = (0..WRITERS)
3107 .map(|_| {
3108 let root = root.clone();
3109 let barrier = Arc::clone(&barrier);
3110 thread::spawn(move || {
3111 let ledger = FleetLedger::open(&root).unwrap();
3112 barrier.wait();
3113 for _ in 0..EVENTS_PER_WRITER {
3114 ledger
3115 .append_event_next_seq(
3116 &FleetRunId::from("run-1"),
3117 "worker-1",
3118 "task-a",
3119 "2026-06-12T17:01:00Z",
3120 FleetWorkerEventPayload::Running,
3121 )
3122 .unwrap();
3123 }
3124 })
3125 })
3126 .collect::<Vec<_>>();
3127 for handle in handles {
3128 handle.join().unwrap();
3129 }
3130
3131 let mut sequences = std::fs::read_to_string(ledger.path())
3132 .unwrap()
3133 .lines()
3134 .filter_map(|line| serde_json::from_str::<FleetLedgerRecord>(line).ok())
3135 .filter_map(|record| match record {
3136 FleetLedgerRecord::EventAppended { event }
3137 if event.run_id.0 == "run-1"
3138 && event.worker_id == "worker-1"
3139 && event.task_id == "task-a" =>
3140 {
3141 Some(event.seq)
3142 }
3143 _ => None,
3144 })
3145 .collect::<Vec<_>>();
3146 sequences.sort_unstable();
3147 assert_eq!(
3148 sequences,
3149 (1..=(WRITERS * EVENTS_PER_WRITER) as u64).collect::<Vec<_>>()
3150 );
3151 }
3152
3153 #[test]
3154 fn concurrent_ledgers_claim_one_queued_task_once() {
3155 let tmp = TempDir::new().unwrap();
3156 let ledger = FleetLedger::open(tmp.path()).unwrap();
3157 ledger.create_run(&sample_run("run-1")).unwrap();
3158 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3159 let root = tmp.path().to_path_buf();
3160 let barrier = Arc::new(Barrier::new(2));
3161 let handles = ["worker-a", "worker-b"].map(|worker_id| {
3162 let root = root.clone();
3163 let barrier = Arc::clone(&barrier);
3164 thread::spawn(move || {
3165 let ledger = FleetLedger::open(&root).unwrap();
3166 barrier.wait();
3167 ledger
3168 .claim_next(worker_id, &[], "2026-06-12T17:01:00Z")
3169 .unwrap()
3170 })
3171 });
3172 let claims = handles
3173 .into_iter()
3174 .filter_map(|handle| handle.join().unwrap())
3175 .collect::<Vec<_>>();
3176
3177 assert_eq!(claims.len(), 1);
3178 assert_eq!(claims[0].task_id, "task-a");
3179 let state = ledger.rebuild_state().unwrap();
3180 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 1);
3181 assert_eq!(
3182 state.tasks["run-1:task-a"].status,
3183 FleetTaskLedgerStatus::Leased
3184 );
3185 }
3186
3187 #[test]
3188 fn concurrent_start_and_cancel_leave_one_terminal_task_without_late_progress() {
3189 let tmp = TempDir::new().unwrap();
3190 let ledger = FleetLedger::open(tmp.path()).unwrap();
3191 ledger.create_run(&sample_run("run-1")).unwrap();
3192 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3193 let root = tmp.path().to_path_buf();
3194 let barrier = Arc::new(Barrier::new(2));
3195
3196 let start_root = root.clone();
3197 let start_barrier = Arc::clone(&barrier);
3198 let starter = thread::spawn(move || {
3199 let ledger = FleetLedger::open(&start_root).unwrap();
3200 start_barrier.wait();
3201 ledger
3202 .start_task_if_enqueued(
3203 &FleetRunId::from("run-1"),
3204 "task-a",
3205 "worker-1",
3206 "2026-06-12T17:01:00Z",
3207 None,
3208 Some(1),
3209 vec![
3210 FleetWorkerEventPayload::Leased {
3211 lease_expires_at: None,
3212 },
3213 FleetWorkerEventPayload::Starting,
3214 FleetWorkerEventPayload::Artifact(FleetArtifactRef {
3215 kind: FleetArtifactKind::Log,
3216 path: PathBuf::from(".codewhale/fleet/run-1/task-a/worker-1.log"),
3217 checksum: None,
3218 mime_type: Some("text/plain".to_string()),
3219 size_bytes: Some(0),
3220 }),
3221 FleetWorkerEventPayload::Running,
3222 ],
3223 || Ok(()),
3224 )
3225 .unwrap()
3226 });
3227 let cancel_root = root.clone();
3228 let cancel_barrier = Arc::clone(&barrier);
3229 let canceller = thread::spawn(move || {
3230 let ledger = FleetLedger::open(&cancel_root).unwrap();
3231 cancel_barrier.wait();
3232 ledger
3233 .cancel_task_if_active(
3234 &FleetRunId::from("run-1"),
3235 "task-a",
3236 None,
3237 "2026-06-12T17:02:00Z",
3238 Some("operator"),
3239 Some("operator"),
3240 )
3241 .unwrap()
3242 });
3243
3244 let started = starter.join().unwrap();
3245 assert!(canceller.join().unwrap());
3246 let state = ledger.rebuild_state().unwrap();
3247 assert_eq!(
3248 state.tasks["run-1:task-a"].status,
3249 FleetTaskLedgerStatus::Cancelled
3250 );
3251 assert_eq!(
3252 state.tasks["run-1:task-a"].entry.attempts,
3253 u32::from(started)
3254 );
3255 if started {
3256 assert!(matches!(
3257 state.latest_events["worker-1:run-1:task-a"].payload,
3258 FleetWorkerEventPayload::Cancelled { .. }
3259 ));
3260 }
3261 let artifacts_before = state.artifact_events.len();
3262 assert!(
3263 ledger
3264 .append_event_if_leased(
3265 &FleetRunId::from("run-1"),
3266 "worker-1",
3267 "task-a",
3268 u32::from(started),
3269 "2026-06-12T17:03:00Z",
3270 FleetWorkerEventPayload::Artifact(FleetArtifactRef {
3271 kind: FleetArtifactKind::Log,
3272 path: PathBuf::from("late.log"),
3273 checksum: None,
3274 mime_type: None,
3275 size_bytes: None,
3276 }),
3277 )
3278 .unwrap()
3279 .is_none()
3280 );
3281 assert_eq!(
3282 ledger.rebuild_state().unwrap().artifact_events.len(),
3283 artifacts_before
3284 );
3285 }
3286
3287 #[test]
3288 fn cancelled_queue_does_not_run_start_projection_callback() {
3289 let tmp = TempDir::new().unwrap();
3290 let ledger = FleetLedger::open(tmp.path()).unwrap();
3291 ledger.create_run(&sample_run("run-1")).unwrap();
3292 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3293 assert!(
3294 ledger
3295 .cancel_task_if_active(
3296 &FleetRunId::from("run-1"),
3297 "task-a",
3298 None,
3299 "2026-06-12T17:01:00Z",
3300 None,
3301 Some("operator"),
3302 )
3303 .unwrap()
3304 );
3305 let callback_ran = AtomicBool::new(false);
3306 assert!(
3307 !ledger
3308 .start_task_if_enqueued(
3309 &FleetRunId::from("run-1"),
3310 "task-a",
3311 "worker-1",
3312 "2026-06-12T17:02:00Z",
3313 None,
3314 Some(1),
3315 vec![FleetWorkerEventPayload::Running],
3316 || {
3317 callback_ran.store(true, Ordering::SeqCst);
3318 Ok(())
3319 },
3320 )
3321 .unwrap()
3322 );
3323 assert!(!callback_ran.load(Ordering::SeqCst));
3324 }
3325
3326 #[test]
3327 fn failed_start_projection_leaves_task_unleased_and_without_events() {
3328 let tmp = TempDir::new().unwrap();
3329 let ledger = FleetLedger::open(tmp.path()).unwrap();
3330 ledger.create_run(&sample_run("run-1")).unwrap();
3331 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3332
3333 let error = ledger
3334 .start_task_if_enqueued(
3335 &FleetRunId::from("run-1"),
3336 "task-a",
3337 "worker-1",
3338 "2026-06-12T17:02:00Z",
3339 None,
3340 Some(1),
3341 vec![FleetWorkerEventPayload::Running],
3342 || Err(anyhow!("forced coordination persistence failure")),
3343 )
3344 .expect_err("projection failure must abort the lease");
3345 assert!(error.to_string().contains("forced coordination"));
3346
3347 let state = ledger.rebuild_state().unwrap();
3348 let task = &state.tasks["run-1:task-a"];
3349 assert_eq!(task.status, FleetTaskLedgerStatus::Enqueued);
3350 assert_eq!(task.entry.attempts, 0);
3351 assert!(task.leased_to.is_none());
3352 assert!(state.latest_events.is_empty());
3353 assert!(state.latest_seq.is_empty());
3354 assert!(!state.heartbeats.contains_key("worker-1"));
3355 }
3356
3357 #[test]
3358 fn concurrent_restarts_compare_and_set_one_new_attempt() {
3359 let tmp = TempDir::new().unwrap();
3360 let ledger = FleetLedger::open(tmp.path()).unwrap();
3361 ledger.create_run(&sample_run("run-1")).unwrap();
3362 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3363 assert!(
3364 ledger
3365 .start_task_if_enqueued(
3366 &FleetRunId::from("run-1"),
3367 "task-a",
3368 "worker-1",
3369 "2026-06-12T17:01:00Z",
3370 None,
3371 Some(1),
3372 vec![FleetWorkerEventPayload::Running],
3373 || Ok(()),
3374 )
3375 .unwrap()
3376 );
3377 let state = ledger.rebuild_state().unwrap();
3378 let expected_seq = state.latest_seq["worker-1:run-1:task-a"];
3379 let expected_heartbeat = state.heartbeats["worker-1"].timestamp.clone();
3380 let root = tmp.path().to_path_buf();
3381 let barrier = Arc::new(Barrier::new(2));
3382 let handles = (0..2)
3383 .map(|_| {
3384 let root = root.clone();
3385 let barrier = Arc::clone(&barrier);
3386 let expected_heartbeat = expected_heartbeat.clone();
3387 thread::spawn(move || {
3388 let ledger = FleetLedger::open(&root).unwrap();
3389 barrier.wait();
3390 ledger
3391 .restart_task_if_unchanged(
3392 &FleetRunId::from("run-1"),
3393 "task-a",
3394 "worker-1",
3395 FleetTaskLedgerStatus::Leased,
3396 1,
3397 expected_seq,
3398 Some(&expected_heartbeat),
3399 "2026-06-12T17:02:00Z",
3400 None,
3401 1,
3402 )
3403 .unwrap()
3404 })
3405 })
3406 .collect::<Vec<_>>();
3407 let winners = handles
3408 .into_iter()
3409 .map(|handle| handle.join().unwrap())
3410 .filter(|won| *won)
3411 .count();
3412
3413 assert_eq!(winners, 1);
3414 let state = ledger.rebuild_state().unwrap();
3415 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 2);
3416 assert_eq!(state.latest_seq["worker-1:run-1:task-a"], expected_seq + 2);
3417 }
3418
3419 #[test]
3420 fn stale_verifier_status_cannot_overwrite_restarted_attempt() {
3421 let tmp = TempDir::new().unwrap();
3422 let verifier = FleetLedger::open(tmp.path()).unwrap();
3423 verifier.create_run(&sample_run("run-1")).unwrap();
3424 verifier.enqueue(sample_entry("run-1", "task-a")).unwrap();
3425 assert!(
3426 verifier
3427 .start_task_if_enqueued(
3428 &FleetRunId::from("run-1"),
3429 "task-a",
3430 "worker-1",
3431 "2026-06-12T17:01:00Z",
3432 None,
3433 Some(1),
3434 vec![FleetWorkerEventPayload::Running],
3435 || Ok(()),
3436 )
3437 .unwrap()
3438 );
3439 let terminal = verifier
3440 .append_terminal_event_if_leased(
3441 &FleetRunId::from("run-1"),
3442 "worker-1",
3443 "task-a",
3444 1,
3445 "2026-06-12T17:02:00Z",
3446 FleetWorkerEventPayload::Completed {
3447 exit_code: Some(0),
3448 summary: None,
3449 },
3450 )
3451 .unwrap()
3452 .unwrap();
3453 let state = verifier.rebuild_state().unwrap();
3454 let heartbeat = state.heartbeats["worker-1"].timestamp.clone();
3455 let restarter = FleetLedger::open(tmp.path()).unwrap();
3456 assert!(
3457 restarter
3458 .restart_task_if_unchanged(
3459 &FleetRunId::from("run-1"),
3460 "task-a",
3461 "worker-1",
3462 FleetTaskLedgerStatus::Completed,
3463 1,
3464 terminal.seq,
3465 Some(&heartbeat),
3466 "2026-06-12T17:03:00Z",
3467 None,
3468 1,
3469 )
3470 .unwrap()
3471 );
3472 assert!(
3473 !verifier
3474 .mark_task_terminal_status_if_unchanged(
3475 &FleetRunId::from("run-1"),
3476 "task-a",
3477 "worker-1",
3478 FleetTaskLedgerStatus::Completed,
3479 1,
3480 terminal.seq,
3481 "2026-06-12T17:04:00Z",
3482 FleetTaskLedgerStatus::Failed,
3483 )
3484 .unwrap()
3485 );
3486 let state = verifier.rebuild_state().unwrap();
3487 assert_eq!(
3488 state.tasks["run-1:task-a"].status,
3489 FleetTaskLedgerStatus::Leased
3490 );
3491 assert_eq!(state.tasks["run-1:task-a"].entry.attempts, 2);
3492 }
3493
3494 #[test]
3495 fn fresh_heartbeat_invalidates_stale_scheduler_failure() {
3496 let tmp = TempDir::new().unwrap();
3497 let scheduler = FleetLedger::open(tmp.path()).unwrap();
3498 scheduler.create_run(&sample_run("run-1")).unwrap();
3499 scheduler.enqueue(sample_entry("run-1", "task-a")).unwrap();
3500 assert!(
3501 scheduler
3502 .start_task_if_enqueued(
3503 &FleetRunId::from("run-1"),
3504 "task-a",
3505 "worker-1",
3506 "2026-06-12T17:01:00Z",
3507 None,
3508 Some(1),
3509 vec![FleetWorkerEventPayload::Running],
3510 || Ok(()),
3511 )
3512 .unwrap()
3513 );
3514 let stale = scheduler
3515 .append_event_if_leased(
3516 &FleetRunId::from("run-1"),
3517 "worker-1",
3518 "task-a",
3519 1,
3520 "2026-06-12T17:02:00Z",
3521 FleetWorkerEventPayload::Stale {
3522 last_heartbeat_at: Some("2026-06-12T17:01:00Z".to_string()),
3523 },
3524 )
3525 .unwrap()
3526 .unwrap();
3527 let worker = FleetLedger::open(tmp.path()).unwrap();
3528 worker
3529 .heartbeat("worker-1", "2026-06-12T17:02:30Z", None, None)
3530 .unwrap();
3531
3532 assert!(
3533 scheduler
3534 .append_terminal_event_if_lease_unchanged(
3535 &FleetRunId::from("run-1"),
3536 "worker-1",
3537 "task-a",
3538 1,
3539 stale.seq,
3540 Some("2026-06-12T17:01:00Z"),
3541 "2026-06-12T17:03:00Z",
3542 FleetWorkerEventPayload::Failed {
3543 reason: "stale retry budget exhausted".to_string(),
3544 recoverable: false,
3545 },
3546 )
3547 .unwrap()
3548 .is_none()
3549 );
3550 assert_eq!(
3551 scheduler.rebuild_state().unwrap().tasks["run-1:task-a"].status,
3552 FleetTaskLedgerStatus::Leased
3553 );
3554 }
3555
3556 #[test]
3557 fn cancellation_and_completion_are_compare_and_set_terminal_transitions() {
3558 let tmp = TempDir::new().unwrap();
3559 let ledger = FleetLedger::open(tmp.path()).unwrap();
3560 ledger.create_run(&sample_run("run-1")).unwrap();
3561 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3562 assert!(
3563 ledger
3564 .lease_task_if_enqueued(
3565 &FleetRunId::from("run-1"),
3566 "task-a",
3567 "worker-1",
3568 "2026-06-12T17:01:00Z",
3569 None,
3570 Some(1),
3571 )
3572 .unwrap()
3573 );
3574 assert!(
3575 ledger
3576 .cancel_task_if_active(
3577 &FleetRunId::from("run-1"),
3578 "task-a",
3579 Some("worker-1"),
3580 "2026-06-12T17:02:00Z",
3581 Some("operator"),
3582 Some("operator"),
3583 )
3584 .unwrap()
3585 );
3586 assert!(
3587 ledger
3588 .append_terminal_event_if_leased(
3589 &FleetRunId::from("run-1"),
3590 "worker-1",
3591 "task-a",
3592 1,
3593 "2026-06-12T17:03:00Z",
3594 FleetWorkerEventPayload::Completed {
3595 exit_code: Some(0),
3596 summary: None,
3597 },
3598 )
3599 .unwrap()
3600 .is_none()
3601 );
3602 assert!(
3603 !ledger
3604 .cancel_task_if_active(
3605 &FleetRunId::from("run-1"),
3606 "task-a",
3607 Some("worker-1"),
3608 "2026-06-12T17:04:00Z",
3609 Some("operator"),
3610 Some("operator"),
3611 )
3612 .unwrap()
3613 );
3614 assert_eq!(
3615 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].status,
3616 FleetTaskLedgerStatus::Cancelled
3617 );
3618 }
3619
3620 #[test]
3621 fn fleet_ledger_survives_restart() {
3622 let tmp = TempDir::new().unwrap();
3623 {
3624 let ledger = FleetLedger::open(tmp.path()).unwrap();
3625 ledger.create_run(&sample_run("run-1")).unwrap();
3626 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3627 ledger
3628 .lease_task(
3629 &FleetRunId::from("run-1"),
3630 "task-a",
3631 "worker-1",
3632 "2026-06-12T17:01:00Z",
3633 None,
3634 )
3635 .unwrap();
3636 }
3637 // Re-open simulates process restart.
3638 let ledger = FleetLedger::open(tmp.path()).unwrap();
3639 let state = ledger.rebuild_state().unwrap();
3640 assert_eq!(state.runs.len(), 1);
3641 assert_eq!(
3642 state.tasks["run-1:task-a"].status,
3643 FleetTaskLedgerStatus::Leased
3644 );
3645 }
3646
3647 #[test]
3648 fn heartbeat_between_stale_snapshot_and_append_fences_stale_event() {
3649 let tmp = TempDir::new().unwrap();
3650 let scheduler = FleetLedger::open(tmp.path()).unwrap();
3651 scheduler.create_run(&sample_run("run-1")).unwrap();
3652 scheduler.enqueue(sample_entry("run-1", "task-a")).unwrap();
3653 assert!(
3654 scheduler
3655 .start_task_if_enqueued(
3656 &FleetRunId::from("run-1"),
3657 "task-a",
3658 "worker-1",
3659 "2026-06-12T17:01:00Z",
3660 None,
3661 Some(1),
3662 vec![FleetWorkerEventPayload::Running],
3663 || Ok(()),
3664 )
3665 .unwrap()
3666 );
3667 let snapshot = scheduler.rebuild_state().unwrap();
3668 let expected_seq = snapshot.latest_seq["worker-1:run-1:task-a"];
3669 let expected_heartbeat = snapshot.heartbeats["worker-1"].timestamp.clone();
3670
3671 FleetLedger::open(tmp.path())
3672 .unwrap()
3673 .heartbeat("worker-1", "2026-06-12T17:02:30Z", None, None)
3674 .unwrap();
3675
3676 assert!(
3677 scheduler
3678 .append_event_if_lease_unchanged(
3679 &FleetRunId::from("run-1"),
3680 "worker-1",
3681 "task-a",
3682 1,
3683 expected_seq,
3684 Some(&expected_heartbeat),
3685 "2026-06-12T17:03:00Z",
3686 FleetWorkerEventPayload::Stale {
3687 last_heartbeat_at: Some(expected_heartbeat.clone()),
3688 },
3689 )
3690 .unwrap()
3691 .is_none()
3692 );
3693 let state = scheduler.rebuild_state().unwrap();
3694 assert_eq!(
3695 state.tasks["run-1:task-a"].status,
3696 FleetTaskLedgerStatus::Leased
3697 );
3698 assert!(matches!(
3699 state.latest_events["worker-1:run-1:task-a"].payload,
3700 FleetWorkerEventPayload::Running
3701 ));
3702 }
3703
3704 #[test]
3705 fn stale_attempt_cannot_finalize_or_replace_restarted_attempt_receipt() {
3706 let tmp = TempDir::new().unwrap();
3707 let ledger = FleetLedger::open(tmp.path()).unwrap();
3708 let run_id = FleetRunId::from("run-1");
3709 ledger.create_run(&sample_run("run-1")).unwrap();
3710 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3711 assert!(
3712 ledger
3713 .start_task_if_enqueued(
3714 &run_id,
3715 "task-a",
3716 "worker-1",
3717 "2026-06-12T17:01:00Z",
3718 None,
3719 Some(1),
3720 vec![FleetWorkerEventPayload::Running],
3721 || Ok(()),
3722 )
3723 .unwrap()
3724 );
3725 let before_restart = ledger.rebuild_state().unwrap();
3726 assert!(
3727 ledger
3728 .restart_task_if_unchanged(
3729 &run_id,
3730 "task-a",
3731 "worker-1",
3732 FleetTaskLedgerStatus::Leased,
3733 1,
3734 before_restart.latest_seq["worker-1:run-1:task-a"],
3735 Some(&before_restart.heartbeats["worker-1"].timestamp),
3736 "2026-06-12T17:02:00Z",
3737 None,
3738 1,
3739 )
3740 .unwrap()
3741 );
3742
3743 let receipt = |attempt, result| FleetReceipt {
3744 run_id: run_id.clone(),
3745 task_id: "task-a".to_string(),
3746 worker_id: "worker-1".to_string(),
3747 attempt: Some(attempt),
3748 terminal_seq: None,
3749 completed_at: "2026-06-12T17:03:00Z".to_string(),
3750 result,
3751 failure_kind: None,
3752 artifacts: Vec::new(),
3753 score: None,
3754 resolved_route: None,
3755 saved_session_id: None,
3756 effective_permissions: None,
3757 };
3758 assert!(
3759 ledger
3760 .finalize_task_attempt_if_leased(
3761 &run_id,
3762 "worker-1",
3763 "task-a",
3764 1,
3765 "2026-06-12T17:03:00Z",
3766 FleetWorkerEventPayload::Completed {
3767 exit_code: Some(0),
3768 summary: Some("late attempt one".to_string()),
3769 },
3770 None,
3771 receipt(1, FleetTaskResult::Fail),
3772 )
3773 .unwrap()
3774 .is_none()
3775 );
3776 let winning_event = ledger
3777 .finalize_task_attempt_if_leased(
3778 &run_id,
3779 "worker-1",
3780 "task-a",
3781 2,
3782 "2026-06-12T17:04:00Z",
3783 FleetWorkerEventPayload::Completed {
3784 exit_code: Some(0),
3785 summary: Some("attempt two".to_string()),
3786 },
3787 None,
3788 receipt(2, FleetTaskResult::Pass),
3789 )
3790 .unwrap()
3791 .unwrap();
3792
3793 let state = ledger.rebuild_state().unwrap();
3794 let durable = &state.receipts["run-1:task-a"];
3795 assert_eq!(durable.attempt, Some(2));
3796 assert_eq!(durable.terminal_seq, Some(winning_event.seq));
3797 assert_eq!(durable.result, FleetTaskResult::Pass);
3798 assert_eq!(
3799 state.tasks["run-1:task-a"].status,
3800 FleetTaskLedgerStatus::Completed
3801 );
3802 }
3803
3804 #[test]
3805 fn fleet_ledger_quarantines_unterminated_tail_before_next_valid_record() {
3806 let tmp = TempDir::new().unwrap();
3807 let ledger = FleetLedger::open(tmp.path()).unwrap();
3808 ledger.create_run(&sample_run("run-1")).unwrap();
3809 // Simulate a process dying before its trailing newline, then use the
3810 // normal append path. The next record must not be concatenated to and
3811 // lost with the malformed crash tail.
3812 let mut file = OpenOptions::new().append(true).open(ledger.path()).unwrap();
3813 write!(file, "{{\"record\":\"run_created\",\"run\":").unwrap();
3814 file.sync_all().unwrap();
3815 drop(file);
3816 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3817
3818 let state = ledger.rebuild_state().unwrap();
3819 assert_eq!(state.runs.len(), 1);
3820 assert!(state.runs.contains_key("run-1"));
3821 assert!(state.tasks.contains_key("run-1:task-a"));
3822 }
3823
3824 #[test]
3825 fn fleet_ledger_event_and_heartbeat_reconstruct_worker_status() {
3826 let tmp = TempDir::new().unwrap();
3827 let ledger = FleetLedger::open(tmp.path()).unwrap();
3828 ledger.create_run(&sample_run("run-1")).unwrap();
3829 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
3830 ledger
3831 .append_event(FleetWorkerEvent {
3832 seq: 1,
3833 run_id: FleetRunId::from("run-1"),
3834 worker_id: "worker-1".to_string(),
3835 task_id: "task-a".to_string(),
3836 timestamp: "2026-06-12T17:01:00Z".to_string(),
3837 payload: FleetWorkerEventPayload::Running,
3838 extra: BTreeMap::new(),
3839 })
3840 .unwrap();
3841 ledger
3842 .heartbeat("worker-1", "2026-06-12T17:02:00Z", Some(12.5), Some(1024))
3843 .unwrap();
3844
3845 let state = ledger.rebuild_state().unwrap();
3846 assert_eq!(state.workers["worker-1"], FleetWorkerStatus::Busy);
3847 assert_eq!(state.heartbeats["worker-1"].cpu_percent, Some(12.5));
3848 }
3849
3850 #[test]
3851 fn fleet_ledger_replays_typed_workflow_receipt_with_distinct_run_ids() {
3852 let tmp = TempDir::new().unwrap();
3853 let ledger = FleetLedger::open(tmp.path()).unwrap();
3854 ledger.create_run(&sample_run("fleet-run-1")).unwrap();
3855 ledger
3856 .enqueue(sample_entry("fleet-run-1", "task-a"))
3857 .unwrap();
3858 ledger
3859 .append_event(FleetWorkerEvent {
3860 seq: 1,
3861 run_id: FleetRunId::from("fleet-run-1"),
3862 worker_id: "worker-1".to_string(),
3863 task_id: "task-a".to_string(),
3864 timestamp: "2026-07-10T00:00:00Z".to_string(),
3865 payload: FleetWorkerEventPayload::WorkflowEvent {
3866 workflow_run_id: "workflow_1".to_string(),
3867 event: serde_json::json!({"type": "task_completed"}),
3868 },
3869 extra: BTreeMap::new(),
3870 })
3871 .unwrap();
3872
3873 let state = ledger.rebuild_state().unwrap();
3874 let event = &state.latest_events["worker-1:fleet-run-1:task-a"];
3875 assert!(matches!(
3876 &event.payload,
3877 FleetWorkerEventPayload::WorkflowEvent {
3878 workflow_run_id,
3879 event,
3880 } if workflow_run_id == "workflow_1" && event["type"] == "task_completed"
3881 ));
3882 let last_line = std::fs::read_to_string(ledger.path())
3883 .unwrap()
3884 .lines()
3885 .last()
3886 .unwrap()
3887 .to_string();
3888 assert_eq!(last_line.matches("\"run_id\"").count(), 1);
3889 assert_eq!(last_line.matches("\"workflow_run_id\"").count(), 1);
3890 }
3891
3892 #[test]
3893 fn fleet_ledger_terminal_events_ignore_late_progress_regressions() {
3894 let tmp = TempDir::new().unwrap();
3895 let ledger = FleetLedger::open(tmp.path()).unwrap();
3896 ledger.create_run(&sample_run("run-1")).unwrap();
3897 ledger
3898 .enqueue(sample_entry("run-1", "task-failed"))
3899 .unwrap();
3900 ledger
3901 .enqueue(sample_entry("run-1", "task-cancelled"))
3902 .unwrap();
3903
3904 ledger
3905 .append_event(FleetWorkerEvent {
3906 seq: 1,
3907 run_id: FleetRunId::from("run-1"),
3908 worker_id: "worker-1".to_string(),
3909 task_id: "task-failed".to_string(),
3910 timestamp: "2026-06-12T17:03:00Z".to_string(),
3911 payload: FleetWorkerEventPayload::Failed {
3912 reason: "test failed".to_string(),
3913 recoverable: false,
3914 },
3915 extra: BTreeMap::new(),
3916 })
3917 .unwrap();
3918 ledger
3919 .append_event(FleetWorkerEvent {
3920 seq: 2,
3921 run_id: FleetRunId::from("run-1"),
3922 worker_id: "worker-2".to_string(),
3923 task_id: "task-cancelled".to_string(),
3924 timestamp: "2026-06-12T17:04:00Z".to_string(),
3925 payload: FleetWorkerEventPayload::Cancelled {
3926 cancelled_by: Some("operator".to_string()),
3927 },
3928 extra: BTreeMap::new(),
3929 })
3930 .unwrap();
3931 // A live worker can flush progress after an out-of-process operator
3932 // command has already made the task terminal. Preserve the raw
3933 // sequence for append ordering without projecting the task or worker
3934 // back to a running state.
3935 ledger
3936 .append_event(FleetWorkerEvent {
3937 seq: 3,
3938 run_id: FleetRunId::from("run-1"),
3939 worker_id: "worker-2".to_string(),
3940 task_id: "task-cancelled".to_string(),
3941 timestamp: "2026-06-12T17:04:01Z".to_string(),
3942 payload: FleetWorkerEventPayload::Running,
3943 extra: BTreeMap::new(),
3944 })
3945 .unwrap();
3946
3947 let state = ledger.rebuild_state().unwrap();
3948 assert_eq!(
3949 state.tasks["run-1:task-failed"].status,
3950 FleetTaskLedgerStatus::Failed
3951 );
3952 assert_eq!(
3953 state.tasks["run-1:task-cancelled"].status,
3954 FleetTaskLedgerStatus::Cancelled
3955 );
3956 assert_eq!(state.workers["worker-2"], FleetWorkerStatus::Online);
3957 assert_eq!(
3958 state.latest_seq["worker-2:run-1:task-cancelled"], 3,
3959 "raw sequence ownership must still advance past ignored progress"
3960 );
3961 assert!(matches!(
3962 state.latest_events["worker-2:run-1:task-cancelled"].payload,
3963 FleetWorkerEventPayload::Cancelled { .. }
3964 ));
3965
3966 ledger.compact().unwrap();
3967 let state = ledger.rebuild_state().unwrap();
3968 assert_eq!(
3969 state.tasks["run-1:task-failed"].status,
3970 FleetTaskLedgerStatus::Failed
3971 );
3972 assert_eq!(
3973 state.tasks["run-1:task-cancelled"].status,
3974 FleetTaskLedgerStatus::Cancelled
3975 );
3976 assert_eq!(state.workers["worker-2"], FleetWorkerStatus::Online);
3977 assert_eq!(
3978 state.latest_seq["worker-2:run-1:task-cancelled"], 3,
3979 "compaction must preserve the ignored event sequence high-water mark"
3980 );
3981 assert!(matches!(
3982 state.latest_events["worker-2:run-1:task-cancelled"].payload,
3983 FleetWorkerEventPayload::Cancelled { .. }
3984 ));
3985 let next = ledger
3986 .append_event_next_seq(
3987 &FleetRunId::from("run-1"),
3988 "worker-2",
3989 "task-cancelled",
3990 "2026-06-12T17:04:02Z",
3991 FleetWorkerEventPayload::Running,
3992 )
3993 .unwrap();
3994 assert_eq!(next.seq, 4, "compaction must never permit sequence reuse");
3995 }
3996
3997 #[test]
3998 fn fleet_ledger_compact_preserves_current_state() {
3999 let tmp = TempDir::new().unwrap();
4000 let ledger = FleetLedger::open(tmp.path()).unwrap();
4001 ledger.create_run(&sample_run("run-1")).unwrap();
4002 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4003 ledger
4004 .lease_task(
4005 &FleetRunId::from("run-1"),
4006 "task-a",
4007 "worker-1",
4008 "2026-06-12T17:01:00Z",
4009 None,
4010 )
4011 .unwrap();
4012 ledger
4013 .append_event(FleetWorkerEvent {
4014 seq: 7,
4015 run_id: FleetRunId::from("run-1"),
4016 worker_id: "worker-1".to_string(),
4017 task_id: "task-a".to_string(),
4018 timestamp: "2026-06-12T17:01:30Z".to_string(),
4019 payload: FleetWorkerEventPayload::Running,
4020 extra: BTreeMap::new(),
4021 })
4022 .unwrap();
4023 ledger
4024 .heartbeat("worker-1", "2026-06-12T17:02:00Z", Some(12.5), Some(1024))
4025 .unwrap();
4026 ledger
4027 .record_receipt(FleetReceipt {
4028 run_id: FleetRunId::from("run-1"),
4029 task_id: "task-a".to_string(),
4030 worker_id: "worker-1".to_string(),
4031 attempt: Some(1),
4032 terminal_seq: None,
4033 completed_at: "2026-06-12T17:03:00Z".to_string(),
4034 result: FleetTaskResult::Pass,
4035 failure_kind: None,
4036 artifacts: vec![],
4037 score: None,
4038 resolved_route: None,
4039 saved_session_id: None,
4040 effective_permissions: None,
4041 })
4042 .unwrap();
4043
4044 let lifecycle_seq_before_compaction =
4045 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].lifecycle_seq;
4046 assert_eq!(
4047 lifecycle_seq_before_compaction, 2,
4048 "enqueue and lease are the two effective owner states"
4049 );
4050
4051 ledger.compact().unwrap();
4052 let contents = std::fs::read_to_string(ledger.path()).unwrap();
4053 assert!(contents.lines().count() >= 5, "{contents}");
4054
4055 let state = ledger.rebuild_state().unwrap();
4056 assert_eq!(state.runs.len(), 1);
4057 assert_eq!(
4058 state.tasks["run-1:task-a"].status,
4059 FleetTaskLedgerStatus::Leased
4060 );
4061 assert_eq!(state.workers["worker-1"], FleetWorkerStatus::Busy);
4062 assert_eq!(state.heartbeats["worker-1"].memory_mb, Some(1024));
4063 assert_eq!(
4064 state.tasks["run-1:task-a"].entry.attempts, 1,
4065 "compaction must not mint a synthetic retry attempt"
4066 );
4067 assert_eq!(
4068 state.tasks["run-1:task-a"].lifecycle_seq, lifecycle_seq_before_compaction,
4069 "compaction must not mint an owner lifecycle transition"
4070 );
4071 assert!(state.latest_seq.values().any(|seq| *seq == 7));
4072 assert_eq!(state.receipts["run-1:task-a"].result, FleetTaskResult::Pass);
4073 }
4074
4075 #[test]
4076 fn fleet_compaction_preserves_multilease_lifecycle_high_water() {
4077 let tmp = TempDir::new().unwrap();
4078 let ledger = FleetLedger::open(tmp.path()).unwrap();
4079 ledger.create_run(&sample_run("run-1")).unwrap();
4080 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4081 ledger
4082 .lease_task(
4083 &FleetRunId::from("run-1"),
4084 "task-a",
4085 "worker-1",
4086 "2026-06-12T17:01:00Z",
4087 None,
4088 )
4089 .unwrap();
4090 ledger
4091 .lease_task(
4092 &FleetRunId::from("run-1"),
4093 "task-a",
4094 "worker-1",
4095 "2026-06-12T17:02:00Z",
4096 None,
4097 )
4098 .unwrap();
4099 let before = ledger.rebuild_state().unwrap();
4100 assert_eq!(before.tasks["run-1:task-a"].lifecycle_seq, 3);
4101
4102 ledger.compact().unwrap();
4103 let after = ledger.rebuild_state().unwrap();
4104 assert_eq!(
4105 after.tasks["run-1:task-a"].lifecycle_seq, 3,
4106 "compaction must not reuse lower Work Graph idempotency keys"
4107 );
4108 }
4109
4110 #[test]
4111 fn compaction_preserves_verifier_failure_override_after_completed_exit() {
4112 let tmp = TempDir::new().unwrap();
4113 let ledger = FleetLedger::open(tmp.path()).unwrap();
4114 let run_id = FleetRunId::from("run-1");
4115 ledger.create_run(&sample_run("run-1")).unwrap();
4116 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4117 assert!(
4118 ledger
4119 .start_task_if_enqueued(
4120 &run_id,
4121 "task-a",
4122 "worker-1",
4123 "2026-06-12T17:01:00Z",
4124 None,
4125 Some(1),
4126 vec![FleetWorkerEventPayload::Running],
4127 || Ok(()),
4128 )
4129 .unwrap()
4130 );
4131 let terminal = ledger
4132 .finalize_task_attempt_if_leased(
4133 &run_id,
4134 "worker-1",
4135 "task-a",
4136 1,
4137 "2026-06-12T17:02:00Z",
4138 FleetWorkerEventPayload::Completed {
4139 exit_code: Some(0),
4140 summary: Some("process succeeded but verification failed".to_string()),
4141 },
4142 Some(FleetTaskLedgerStatus::Failed),
4143 FleetReceipt {
4144 run_id: run_id.clone(),
4145 task_id: "task-a".to_string(),
4146 worker_id: "worker-1".to_string(),
4147 attempt: Some(1),
4148 terminal_seq: None,
4149 completed_at: "2026-06-12T17:02:00Z".to_string(),
4150 result: FleetTaskResult::Fail,
4151 failure_kind: Some(FleetTaskFailureKind::Verifier),
4152 artifacts: Vec::new(),
4153 score: None,
4154 resolved_route: None,
4155 saved_session_id: None,
4156 effective_permissions: None,
4157 },
4158 )
4159 .unwrap()
4160 .unwrap();
4161 assert_eq!(
4162 ledger.rebuild_state().unwrap().tasks["run-1:task-a"].status,
4163 FleetTaskLedgerStatus::Failed
4164 );
4165
4166 ledger.compact().unwrap();
4167 let reopened = FleetLedger::open(tmp.path()).unwrap();
4168 let state = reopened.rebuild_state().unwrap();
4169 assert_eq!(
4170 state.tasks["run-1:task-a"].status,
4171 FleetTaskLedgerStatus::Failed
4172 );
4173 assert_eq!(state.receipts["run-1:task-a"].attempt, Some(1));
4174 assert_eq!(
4175 state.receipts["run-1:task-a"].terminal_seq,
4176 Some(terminal.seq)
4177 );
4178 assert!(matches!(
4179 state.latest_events["worker-1:run-1:task-a"].payload,
4180 FleetWorkerEventPayload::Completed { .. }
4181 ));
4182 }
4183
4184 #[test]
4185 fn legacy_attemptless_alert_suppresses_duplicate_delivery_after_upgrade() {
4186 let tmp = TempDir::new().unwrap();
4187 let ledger = FleetLedger::open(tmp.path()).unwrap();
4188 let run_id = FleetRunId::from("run-1");
4189 ledger.create_run(&sample_run("run-1")).unwrap();
4190 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4191 assert!(
4192 ledger
4193 .start_task_if_enqueued(
4194 &run_id,
4195 "task-a",
4196 "worker-1",
4197 "2026-06-12T17:01:00Z",
4198 None,
4199 Some(1),
4200 vec![FleetWorkerEventPayload::Running],
4201 || Ok(()),
4202 )
4203 .unwrap()
4204 );
4205 ledger
4206 .append_terminal_event_if_leased(
4207 &run_id,
4208 "worker-1",
4209 "task-a",
4210 1,
4211 "2026-06-12T17:02:00Z",
4212 FleetWorkerEventPayload::Failed {
4213 reason: "attempt exhausted before upgrade".to_string(),
4214 recoverable: false,
4215 },
4216 )
4217 .unwrap()
4218 .unwrap();
4219 ledger
4220 .record_alert(&run_id, "task-a", "slack", "2026-06-12T17:02:01Z")
4221 .unwrap();
4222 assert!(
4223 !ledger
4224 .record_failed_attempt_alert_once(
4225 &run_id,
4226 "task-a",
4227 "worker-1",
4228 1,
4229 "slack",
4230 "slack#0",
4231 "2026-06-12T17:03:00Z",
4232 )
4233 .unwrap()
4234 );
4235 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 1);
4236
4237 ledger.compact().unwrap();
4238 assert!(
4239 !ledger
4240 .record_failed_attempt_alert_once(
4241 &run_id,
4242 "task-a",
4243 "worker-1",
4244 1,
4245 "slack",
4246 "slack#0",
4247 "2026-06-12T17:04:00Z",
4248 )
4249 .unwrap()
4250 );
4251 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 1);
4252
4253 let attempt_one = ledger.rebuild_state().unwrap();
4254 assert!(
4255 ledger
4256 .restart_task_if_unchanged(
4257 &run_id,
4258 "task-a",
4259 "worker-1",
4260 FleetTaskLedgerStatus::Failed,
4261 1,
4262 attempt_one.latest_seq["worker-1:run-1:task-a"],
4263 Some(&attempt_one.heartbeats["worker-1"].timestamp),
4264 "2026-06-12T17:05:00Z",
4265 None,
4266 1,
4267 )
4268 .unwrap()
4269 );
4270 ledger
4271 .append_terminal_event_if_leased(
4272 &run_id,
4273 "worker-1",
4274 "task-a",
4275 2,
4276 "2026-06-12T17:06:00Z",
4277 FleetWorkerEventPayload::Failed {
4278 reason: "attempt two also exhausted".to_string(),
4279 recoverable: false,
4280 },
4281 )
4282 .unwrap()
4283 .unwrap();
4284 assert!(
4285 ledger
4286 .record_failed_attempt_alert_once(
4287 &run_id,
4288 "task-a",
4289 "worker-1",
4290 2,
4291 "slack",
4292 "slack#0",
4293 "2026-06-12T17:07:00Z",
4294 )
4295 .unwrap()
4296 );
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:08:00Z",
4307 )
4308 .unwrap()
4309 );
4310 assert_eq!(ledger.rebuild_state().unwrap().alerts.len(), 2);
4311 }
4312
4313 #[test]
4314 fn compaction_lock_keeps_concurrent_append_on_replacement_ledger() {
4315 let tmp = TempDir::new().unwrap();
4316 let ledger = FleetLedger::open(tmp.path()).unwrap();
4317 ledger.create_run(&sample_run("run-1")).unwrap();
4318 ledger.enqueue(sample_entry("run-1", "task-a")).unwrap();
4319
4320 let root = tmp.path().to_path_buf();
4321 let (snapshot_tx, snapshot_rx) = mpsc::sync_channel(0);
4322 let (release_tx, release_rx) = mpsc::sync_channel(0);
4323 let compact_root = root.clone();
4324 let compactor = thread::spawn(move || {
4325 let ledger = FleetLedger::open(&compact_root).unwrap();
4326 ledger
4327 .compact_with_snapshot_hook(|| {
4328 snapshot_tx.send(()).unwrap();
4329 release_rx.recv().unwrap();
4330 })
4331 .unwrap();
4332 });
4333 snapshot_rx
4334 .recv_timeout(Duration::from_secs(5))
4335 .expect("compaction never reached its locked snapshot");
4336
4337 let contender = FleetLedger::open(&root).unwrap();
4338 let lock_file = contender.open_lock_file().unwrap();
4339 let mut lock = fd_lock::RwLock::new(lock_file);
4340 match lock.try_write() {
4341 Err(err) => assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock),
4342 Ok(_) => panic!("compaction snapshot did not retain the Fleet ledger lock"),
4343 }
4344
4345 let (append_started_tx, append_started_rx) = mpsc::sync_channel(0);
4346 let (append_done_tx, append_done_rx) = mpsc::sync_channel(0);
4347 let append_root = root.clone();
4348 let appender = thread::spawn(move || {
4349 let ledger = FleetLedger::open(&append_root).unwrap();
4350 append_started_tx.send(()).unwrap();
4351 let event = ledger
4352 .append_event_next_seq(
4353 &FleetRunId::from("run-1"),
4354 "worker-1",
4355 "task-a",
4356 "2026-06-12T17:01:00Z",
4357 FleetWorkerEventPayload::Running,
4358 )
4359 .unwrap();
4360 append_done_tx.send(event.seq).unwrap();
4361 });
4362 append_started_rx.recv().unwrap();
4363 assert!(
4364 append_done_rx
4365 .recv_timeout(Duration::from_millis(100))
4366 .is_err(),
4367 "append completed while compaction still held the ledger lock"
4368 );
4369
4370 release_tx.send(()).unwrap();
4371 compactor.join().unwrap();
4372 assert_eq!(
4373 append_done_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
4374 1
4375 );
4376 appender.join().unwrap();
4377
4378 let state = ledger.rebuild_state().unwrap();
4379 assert_eq!(state.latest_seq["worker-1:run-1:task-a"], 1);
4380 assert!(matches!(
4381 &state.latest_events["worker-1:run-1:task-a"].payload,
4382 FleetWorkerEventPayload::Running
4383 ));
4384 }
4385
4386 #[test]
4387 fn fleet_ledger_receipt_round_trip() {
4388 let tmp = TempDir::new().unwrap();
4389 let ledger = FleetLedger::open(tmp.path()).unwrap();
4390 let receipt = FleetReceipt {
4391 run_id: FleetRunId::from("run-1"),
4392 task_id: "task-a".to_string(),
4393 worker_id: "worker-1".to_string(),
4394 attempt: None,
4395 terminal_seq: None,
4396 completed_at: "2026-06-12T17:03:00Z".to_string(),
4397 result: FleetTaskResult::Pass,
4398 failure_kind: None,
4399 artifacts: vec![],
4400 score: None,
4401 resolved_route: None,
4402 saved_session_id: None,
4403 effective_permissions: None,
4404 };
4405 ledger.record_receipt(receipt.clone()).unwrap();
4406 let state = ledger.rebuild_state().unwrap();
4407 assert_eq!(state.receipts["run-1:task-a"].result, FleetTaskResult::Pass);
4408 }
4409 }
4410
4410 lines RUST