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