返回 CodeWhale
manager.rs
根目录 / crates / tui / src / fleet / manager.rs
1 //! Local-first fleet manager loop and operator controls.
2 //!
3 //! This module is intentionally ledger-first: the first manager can run in the
4 //! foreground and coordinate logical local workers while later host adapters
5 //! add real process and SSH execution behind the same records.
6
7 #![allow(dead_code)]
8
9 use std::collections::{BTreeMap, BTreeSet};
10 use std::fs::OpenOptions;
11 use std::io::ErrorKind;
12 use std::path::{Path, PathBuf};
13 use std::time::Duration;
14
15 use anyhow::{Context, Result, anyhow, bail};
16 use chrono::{DateTime, SecondsFormat, Utc};
17 use codewhale_protocol::fleet::*;
18 use serde_json::Value;
19 use uuid::Uuid;
20
21 use super::executor::{
22 FleetExecutor, FleetExecutorAttempt, FleetWorkerReportedRoute, FleetWorkerTerminalEvent,
23 authority_envelope_for_worker, build_worker_exec_command_with_launch_spec,
24 };
25 use super::host::FleetHostErrorKind;
26 use super::ledger::{
27 FleetEventReplayError, FleetLedger, FleetLedgerState, FleetTaskLedgerStatus, FleetTaskState,
28 };
29 use super::scheduler::{FleetScheduler, FleetSchedulerPolicy};
30 use super::task_spec::{
31 FleetTaskSpecDocument, FleetTaskVerificationInput, load_task_spec_document,
32 prepare_verification_receipt, validate_task_spec_document, verify_task_result,
33 };
34 use super::worker_runtime;
35 use crate::config::Config;
36 use crate::tools::subagent::{AgentWorkerSpec, SharedSubAgentManager, SubAgentManager};
37
38 const DEFAULT_STALE_AFTER_SECONDS: u64 = 300;
39
40 pub struct FleetManager {
41 workspace: PathBuf,
42 ledger: FleetLedger,
43 stale_after: Duration,
44 exec_config: codewhale_config::FleetExecConfig,
45 /// `[fleet]` table used to build the agent roster for dispatch
46 /// (#fleet-roster cutover (v0.8.67)). Defaults keep built-in + workspace
47 /// members resolvable even when the caller has no parsed config.
48 fleet_config: codewhale_config::FleetConfigToml,
49 /// Optional sub-agent manager for headless worker execution.
50 /// When set, fleet workers spawn real sub-agents; when None,
51 /// the manager falls back to local simulation.
52 sub_agent_manager: Option<SharedSubAgentManager>,
53 /// The live session route — the operator's model. Workers whose task and
54 /// roster profile pin no model inherit this instead of `"auto"`, so the
55 /// model the user picked in `/model` is the model that runs the fleet
56 /// (matching the `/fleet roster` operator row). `None` keeps the legacy
57 /// `"auto"` fallback for headless callers with no session.
58 session_model: Option<String>,
59 /// Live provider-route authority used to mint truthful Fleet receipts.
60 /// Kept out of Debug because it may contain credentials.
61 route_config: Option<Config>,
62 }
63
64 impl std::fmt::Debug for FleetManager {
65 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66 f.debug_struct("FleetManager")
67 .field("workspace", &self.workspace)
68 .field("ledger", &self.ledger)
69 .field("stale_after", &self.stale_after)
70 .field("exec_config", &self.exec_config)
71 .field(
72 "sub_agent_manager",
73 &self
74 .sub_agent_manager
75 .as_ref()
76 .map(|_| "SharedSubAgentManager"),
77 )
78 .finish()
79 }
80 }
81
82 #[derive(Debug, Clone)]
83 pub struct FleetRunReport {
84 pub run_id: FleetRunId,
85 pub task_count: usize,
86 pub leased: usize,
87 pub queued: usize,
88 pub worker_ids: Vec<String>,
89 /// Non-blocking dispatch warnings about a brief/profile mismatch.
90 pub warnings: Vec<String>,
91 }
92
93 /// Product identity captured with a managed Fleet run.
94 ///
95 /// CLI task-spec runs predate these fields and use the default descriptor.
96 /// Runtime API callers provide all three values so target and Workflow
97 /// identity remain durable and inspectable after the creating client exits.
98 #[derive(Debug, Clone, Default)]
99 pub struct ManagedFleetRunDescriptor {
100 pub target: Option<FleetRuntimeTarget>,
101 pub workflow: Option<FleetWorkflowDescriptor>,
102 pub roles: Vec<String>,
103 }
104
105 /// Durable restart transition plus the execution context a caller must drive.
106 ///
107 /// Restarting is intentionally split from execution so a live foreground
108 /// manager can observe the transition. Standalone callers must pass this
109 /// context to [`FleetManager::run_to_completion`] before exiting.
110 #[derive(Debug, Clone)]
111 pub struct FleetRestartReport {
112 pub run_id: FleetRunId,
113 pub max_workers: usize,
114 pub inspection: FleetWorkerInspection,
115 }
116
117 #[derive(Debug, Clone, Default)]
118 pub struct FleetTickReport {
119 pub leased: usize,
120 pub heartbeats: usize,
121 }
122
123 #[derive(Debug, Clone, Default)]
124 pub struct FleetExecutorTickReport {
125 pub started: usize,
126 pub events: usize,
127 pub terminals: usize,
128 }
129
130 /// Typed Fleet control-plane refusals.
131 ///
132 /// These are *state conflicts*, not transient backend faults, and the control
133 /// surface classifies them by type. Matching on the message text instead made
134 /// the classification depend on prose that any refactor could silently change.
135 #[derive(Debug, thiserror::Error)]
136 pub enum FleetControlError {
137 /// The exact worker exists (or does not) but has no leased task to cancel.
138 #[error("worker {worker_id} has no running fleet task")]
139 NoActiveTask { worker_id: String },
140 /// No durable run with that exact id in this workspace's ledger.
141 #[error("no fleet run with id {run_id}")]
142 UnknownRun { run_id: String },
143 }
144
145 #[derive(Debug, Clone, Default)]
146 pub struct FleetStatusSnapshot {
147 pub runs: usize,
148 pub queued: usize,
149 pub running: usize,
150 pub completed: usize,
151 pub partial: usize,
152 pub failed: usize,
153 pub restarted: usize,
154 pub escalated: usize,
155 pub transport_failed: usize,
156 pub task_failed: usize,
157 pub verifier_failed: usize,
158 pub cancelled: usize,
159 pub stale: usize,
160 pub workers: BTreeMap<String, FleetWorkerStatus>,
161 }
162
163 /// Outcome of resuming a fleet run from durable ledger state after a manager
164 /// restart. The counts reflect the reconciliation pass; `status` is the
165 /// post-resume inspectable snapshot.
166 #[derive(Debug, Clone)]
167 pub struct FleetResumeReport {
168 pub run_id: FleetRunId,
169 /// Orphaned in-flight leases detected as stale and reclaimed.
170 pub reclaimed_stale: usize,
171 /// Stale leases retried within their retry budget.
172 pub restarted: usize,
173 /// Stale leases that exhausted their retry budget and were failed.
174 pub failed: usize,
175 /// Escalation alerts emitted for exhausted tasks.
176 pub escalated: usize,
177 /// Inspectable run status after the resume pass.
178 pub status: FleetStatusSnapshot,
179 }
180
181 #[derive(Debug, Clone)]
182 pub struct FleetWorkerInspection {
183 pub worker_id: String,
184 pub status: FleetWorkerStatus,
185 pub current_run_id: Option<FleetRunId>,
186 pub current_task_id: Option<String>,
187 pub objective: Option<String>,
188 pub role: Option<String>,
189 pub host: Option<String>,
190 pub latest_heartbeat_at: Option<String>,
191 pub latest_event: Option<FleetWorkerEvent>,
192 pub artifacts: Vec<FleetArtifactRef>,
193 pub receipt_summary: Option<String>,
194 pub last_error: Option<String>,
195 pub alert_state: Option<String>,
196 /// Lightweight projection from the sub-agent worker runtime.
197 /// Populated when a sub-agent manager is attached.
198 pub runtime_state: Option<FleetWorkerRuntimeProjection>,
199 }
200
201 /// Lightweight TUI projection of a headless sub-agent worker's current state.
202 ///
203 /// Derived from the sub-agent manager's `AgentWorkerRecord`.
204 #[derive(Debug, Clone)]
205 pub struct FleetWorkerRuntimeProjection {
206 /// Sub-agent lifecycle status (Queued, Starting, Running, Completed, etc.)
207 pub agent_status: String,
208 /// Steps taken so far (tool calls + model turns)
209 pub steps_taken: u32,
210 /// Latest human-readable message from the worker
211 pub latest_message: Option<String>,
212 /// Error message if the worker failed
213 pub error: Option<String>,
214 /// Result summary if the worker completed
215 pub result_summary: Option<String>,
216 /// Whether the worker has a sub-agent session running
217 pub has_session: bool,
218 }
219
220 #[derive(Debug, Clone)]
221 struct FleetExecutorTaskContext {
222 entry: FleetInboxEntry,
223 task_spec: FleetTaskSpec,
224 worker_id: String,
225 }
226
227 impl FleetManager {
228 pub fn open(workspace: impl AsRef<Path>) -> Result<Self> {
229 let workspace = workspace.as_ref().to_path_buf();
230 let ledger = FleetLedger::open(&workspace)?;
231 Ok(Self {
232 workspace,
233 ledger,
234 stale_after: Duration::from_secs(DEFAULT_STALE_AFTER_SECONDS),
235 exec_config: codewhale_config::FleetExecConfig::default(),
236 fleet_config: codewhale_config::FleetConfigToml::default(),
237 sub_agent_manager: None,
238 session_model: None,
239 route_config: None,
240 })
241 }
242
243 /// Adopt the active session route as the run-level model: whatever the
244 /// user selected in `/model` becomes the operator, and workers without a
245 /// task/profile model pin inherit it. Empty and `"auto"` values are
246 /// ignored so the resolver default keeps applying.
247 pub fn with_session_model(mut self, model: impl Into<String>) -> Self {
248 let model = model.into();
249 let trimmed = model.trim();
250 if !trimmed.is_empty() && !trimmed.eq_ignore_ascii_case("auto") {
251 self.session_model = Some(trimmed.to_string());
252 }
253 self
254 }
255
256 pub fn with_route_config(mut self, config: Config) -> Self {
257 self.route_config = Some(config);
258 self
259 }
260
261 /// The run-level model handed to worker-spec resolution: the session
262 /// model when one was adopted, else the legacy `"auto"` sentinel.
263 fn run_model(&self) -> &str {
264 self.session_model.as_deref().unwrap_or("auto")
265 }
266
267 pub fn with_stale_after(mut self, stale_after: Duration) -> Self {
268 self.stale_after = stale_after;
269 self
270 }
271
272 /// Apply fleet headless-worker execution policy from config.
273 pub fn with_exec_config(mut self, exec_config: codewhale_config::FleetExecConfig) -> Self {
274 self.exec_config = exec_config;
275 self
276 }
277
278 /// Apply the parsed `[fleet]` table so `[fleet.profiles]` members join
279 /// the dispatch roster (#fleet-roster cutover (v0.8.67)).
280 pub fn with_fleet_config(mut self, fleet_config: codewhale_config::FleetConfigToml) -> Self {
281 self.fleet_config = fleet_config;
282 self
283 }
284
285 /// Merged agent roster (built-ins + `[fleet.profiles]` + workspace files)
286 /// used everywhere a task references an `agent_profile` id.
287 fn agent_roster(&self) -> crate::fleet::roster::FleetRoster {
288 crate::fleet::roster::FleetRoster::load(&self.fleet_config, &self.workspace)
289 }
290
291 /// Attach a sub-agent manager so fleet workers can spawn real headless agents.
292 pub fn with_sub_agent_manager(mut self, mgr: SharedSubAgentManager) -> Self {
293 self.sub_agent_manager = Some(mgr);
294 self
295 }
296
297 pub fn ledger_path(&self) -> &Path {
298 self.ledger.path()
299 }
300
301 fn manager_lock_path(&self, run_id: &FleetRunId) -> PathBuf {
302 self.workspace
303 .join(".codewhale")
304 .join("fleet")
305 .join(format!("manager-{}.lock", safe_path_segment(&run_id.0)))
306 }
307
308 pub fn rebuild_state(&self) -> Result<FleetLedgerState> {
309 self.ledger.rebuild_state()
310 }
311
312 pub fn replay_events(
313 &self,
314 run_id: &FleetRunId,
315 after: Option<&str>,
316 limit: usize,
317 ) -> std::result::Result<FleetEventReplay, FleetEventReplayError> {
318 self.ledger.replay_events(run_id, after, limit)
319 }
320
321 pub fn load_task_spec(path: &Path) -> Result<FleetTaskSpecDocument> {
322 load_task_spec_document(path)
323 }
324
325 pub fn create_run_from_task_spec_path(
326 &self,
327 path: &Path,
328 max_workers: usize,
329 ) -> Result<FleetRunReport> {
330 let doc = Self::load_task_spec(path)?;
331 self.create_run(doc, max_workers)
332 }
333
334 pub fn create_run(
335 &self,
336 doc: FleetTaskSpecDocument,
337 max_workers: usize,
338 ) -> Result<FleetRunReport> {
339 let mut prepared = self.create_queued_run(doc, max_workers)?;
340 let started = self.start_run(&prepared.run_id)?;
341 prepared.leased = started.leased;
342 prepared.queued = started.queued;
343 Ok(prepared)
344 }
345
346 /// Validate and durably create a queued run without launching work.
347 ///
348 /// Managed clients use this as the first half of an explicit two-step
349 /// launch gate. The CLI's historical `fleet run` path calls `create_run`,
350 /// which immediately invokes [`Self::start_run`] to preserve compatibility.
351 pub fn create_queued_run(
352 &self,
353 doc: FleetTaskSpecDocument,
354 max_workers: usize,
355 ) -> Result<FleetRunReport> {
356 self.create_queued_run_with_descriptor(
357 doc,
358 max_workers,
359 ManagedFleetRunDescriptor::default(),
360 )
361 }
362
363 pub fn create_queued_run_with_descriptor(
364 &self,
365 mut doc: FleetTaskSpecDocument,
366 max_workers: usize,
367 descriptor: ManagedFleetRunDescriptor,
368 ) -> Result<FleetRunReport> {
369 validate_task_spec_document(&doc)?;
370 // The single funnel: `create_run` and `create_queued_run` both land
371 // here, so counting at either of those would double-count a plain
372 // `fleet run`. Counted after validation, so a rejected spec is not a
373 // dispatch.
374 codewhale_telemetry::session_counters().bump(codewhale_telemetry::Counter::FleetDispatch);
375 worker_runtime::canonicalize_fleet_task_roles(&mut doc.tasks);
376 let roster = self.agent_roster();
377 worker_runtime::validate_task_agent_profiles(&doc.tasks, roster.members())?;
378 worker_runtime::validate_fleet_task_routes(
379 &doc.tasks,
380 roster.members(),
381 self.session_model(),
382 self.route_config.as_ref(),
383 )?;
384 let warnings = doc
385 .tasks
386 .iter()
387 .filter_map(|task| {
388 worker_runtime::network_posture_warning_for_task(
389 task,
390 roster.members(),
391 self.session_model(),
392 )
393 })
394 .collect::<Vec<_>>();
395 let max_workers = max_workers.clamp(1, 128);
396 let run_id = FleetRunId::from(format!(
397 "fleet-{}",
398 &Uuid::new_v4().simple().to_string()[..8]
399 ));
400 let now = timestamp();
401 if doc.workers.is_empty() {
402 doc.workers = default_local_workers(&run_id, max_workers);
403 }
404 let run = FleetRun {
405 id: run_id.clone(),
406 name: doc.name.unwrap_or_else(|| run_id.0.clone()),
407 status: FleetRunStatus::Queued,
408 target: descriptor.target,
409 workflow: descriptor.workflow,
410 roles: descriptor.roles,
411 max_workers: Some(max_workers),
412 task_specs: doc.tasks.clone(),
413 worker_specs: doc.workers.clone(),
414 labels: doc.labels,
415 security_policy: doc.security_policy.clone(),
416 created_at: now.clone(),
417 updated_at: Some(now.clone()),
418 completed_at: None,
419 };
420 self.ledger.create_run(&run)?;
421 for task in &run.task_specs {
422 self.ledger.enqueue(FleetInboxEntry {
423 run_id: run.id.clone(),
424 task_id: task.id.clone(),
425 priority: task_priority(task),
426 enqueued_at: now.clone(),
427 lease_deadline: None,
428 attempts: 0,
429 })?;
430 }
431 let state = self.ledger.rebuild_state()?;
432 let snapshot = self.status_from_state(Some(&run.id), &state);
433 Ok(FleetRunReport {
434 run_id: run.id,
435 task_count: run.task_specs.len(),
436 leased: 0,
437 queued: snapshot.queued,
438 worker_ids: run.worker_specs.iter().map(|w| w.id.clone()).collect(),
439 warnings,
440 })
441 }
442
443 /// Activate one durable queued run without leasing work.
444 ///
445 /// Managed clients use this transition before spawning the executor
446 /// driver, so no partial scheduling failure can strand an unowned lease.
447 /// Terminal runs are never reactivated through this method; the existing
448 /// explicit worker restart control remains separate.
449 pub fn activate_run(&self, run_id: &FleetRunId) -> Result<FleetRunReport> {
450 let state = self.ledger.rebuild_state()?;
451 let run =
452 state
453 .runs
454 .get(&run_id.0)
455 .cloned()
456 .ok_or_else(|| FleetControlError::UnknownRun {
457 run_id: run_id.0.clone(),
458 })?;
459 let lifecycle = state
460 .run_status_overrides
461 .get(&run_id.0)
462 .unwrap_or(&run.status);
463 if matches!(
464 lifecycle,
465 FleetRunStatus::Completed | FleetRunStatus::Failed | FleetRunStatus::Cancelled
466 ) {
467 bail!("fleet run {} is already terminal ({lifecycle:?})", run_id.0);
468 }
469 if !matches!(lifecycle, FleetRunStatus::Running) {
470 self.ledger
471 .update_run_status(run_id, FleetRunStatus::Running, &timestamp())?;
472 }
473 let state = self.ledger.rebuild_state()?;
474 let snapshot = self.status_from_state(Some(run_id), &state);
475 Ok(FleetRunReport {
476 run_id: run.id,
477 task_count: run.task_specs.len(),
478 leased: 0,
479 queued: snapshot.queued,
480 worker_ids: run
481 .worker_specs
482 .iter()
483 .map(|worker| worker.id.clone())
484 .collect(),
485 warnings: Vec::new(),
486 })
487 }
488
489 /// Activate a run and lease its first worker batch for the historical CLI
490 /// path. Managed Runtime callers use [`Self::activate_run`] and start the
491 /// executor driver before it performs scheduling.
492 pub fn start_run(&self, run_id: &FleetRunId) -> Result<FleetRunReport> {
493 let mut report = self.activate_run(run_id)?;
494 let state = self.ledger.rebuild_state()?;
495 let run = state
496 .runs
497 .get(&run_id.0)
498 .ok_or_else(|| FleetControlError::UnknownRun {
499 run_id: run_id.0.clone(),
500 })?;
501 let max_workers = run
502 .max_workers
503 .unwrap_or_else(|| run.worker_specs.len().max(1))
504 .clamp(1, 128);
505 let tick = self.schedule_run(run_id, max_workers)?;
506 self.refresh_run_status(run_id)?;
507 let state = self.ledger.rebuild_state()?;
508 let snapshot = self.status_from_state(Some(run_id), &state);
509 report.leased = tick.leased;
510 report.queued = snapshot.queued;
511 Ok(report)
512 }
513
514 pub fn schedule_run(&self, run_id: &FleetRunId, max_workers: usize) -> Result<FleetTickReport> {
515 self.schedule_run_excluding(run_id, max_workers, &BTreeSet::new())
516 }
517
518 fn schedule_run_excluding(
519 &self,
520 run_id: &FleetRunId,
521 max_workers: usize,
522 unavailable_workers: &BTreeSet<String>,
523 ) -> Result<FleetTickReport> {
524 self.reconcile_coordination_worker_statuses()?;
525 let max_workers = max_workers.clamp(1, 128);
526 let mut report = FleetTickReport::default();
527 let state = self.ledger.rebuild_state()?;
528 let run = state
529 .runs
530 .get(&run_id.0)
531 .cloned()
532 .ok_or_else(|| anyhow!("fleet run {} does not exist", run_id.0))?;
533 let worker_ids = worker_ids_for_run(&run, max_workers);
534
535 for task in active_tasks_for_run(&state, run_id) {
536 if let Some(worker_id) = task.leased_to.as_deref()
537 && worker_ids.iter().any(|id| id == worker_id)
538 {
539 self.ledger.heartbeat(worker_id, &timestamp(), None, None)?;
540 report.heartbeats += 1;
541 }
542 }
543
544 loop {
545 let state = self.ledger.rebuild_state()?;
546 let active_workers = active_workers_for_run(&state, run_id);
547 if active_workers.len() >= max_workers {
548 break;
549 }
550 let Some(worker_id) = worker_ids
551 .iter()
552 .find(|id| {
553 !active_workers.contains(*id) && !unavailable_workers.contains(id.as_str())
554 })
555 .cloned()
556 else {
557 break;
558 };
559 let Some((entry, task_spec)) = next_enqueued_task_for_run(&state, run_id) else {
560 break;
561 };
562 if !self.start_worker_task(&worker_id, &entry, &task_spec, Some(max_workers))? {
563 // Busy coordination or a write-claim contention leaves the
564 // task durably queued. Returning to the driver tick avoids a
565 // tight loop over unchanged state and lets live workers make
566 // progress before scheduling retries.
567 break;
568 }
569 report.leased += 1;
570 }
571
572 self.refresh_run_status(run_id)?;
573 Ok(report)
574 }
575
576 fn reconcile_coordination_worker_statuses(&self) -> Result<()> {
577 let Some(manager) = self.sub_agent_manager.as_ref() else {
578 return Ok(());
579 };
580 let Ok(mut guard) = manager.try_write() else {
581 // The next scheduler tick retries before leasing more work.
582 return Ok(());
583 };
584 let state = self.ledger.rebuild_state()?;
585 for record in guard.list_worker_records() {
586 let current = state
587 .tasks
588 .values()
589 .filter(|task| task.entry.run_id.0 == record.spec.run_id)
590 .filter(|task| task.leased_to.as_deref() == Some(record.spec.worker_id.as_str()))
591 .max_by_key(|task| task.lifecycle_seq);
592 let Some(current) = current else {
593 continue;
594 };
595 let (status, status_label) = match current.status {
596 FleetTaskLedgerStatus::Enqueued => {
597 (crate::tools::subagent::AgentWorkerStatus::Queued, "queued")
598 }
599 FleetTaskLedgerStatus::Leased => (
600 crate::tools::subagent::AgentWorkerStatus::Running,
601 "running",
602 ),
603 FleetTaskLedgerStatus::Completed => (
604 crate::tools::subagent::AgentWorkerStatus::Completed,
605 "completed",
606 ),
607 FleetTaskLedgerStatus::Failed => {
608 (crate::tools::subagent::AgentWorkerStatus::Failed, "failed")
609 }
610 FleetTaskLedgerStatus::Cancelled => (
611 crate::tools::subagent::AgentWorkerStatus::Cancelled,
612 "cancelled",
613 ),
614 };
615 guard.project_external_worker_status(
616 &record.spec.worker_id,
617 status,
618 Some(format!(
619 "Fleet task {} is {}",
620 current.entry.task_id, status_label
621 )),
622 );
623 }
624 Ok(())
625 }
626
627 pub fn status(&self) -> Result<FleetStatusSnapshot> {
628 let state = self.ledger.rebuild_state()?;
629 Ok(self.status_from_state(None, &state))
630 }
631
632 pub fn run_status(&self, run_id: &FleetRunId) -> Result<FleetStatusSnapshot> {
633 let state = self.ledger.rebuild_state()?;
634 Ok(self.status_from_state(Some(run_id), &state))
635 }
636
637 pub fn run_has_open_work(&self, run_id: &FleetRunId) -> Result<bool> {
638 let status = self.run_status(run_id)?;
639 Ok(status.queued + status.running + status.stale > 0)
640 }
641
642 /// Resume a run from durable ledger state after a manager restart.
643 ///
644 /// A crashed or detached manager can leave in-flight tasks `Leased` to
645 /// workers whose processes are gone. Resume rebuilds run state from the
646 /// ledger, reconciles those orphaned/stale leases through the shared
647 /// scheduler recovery semantics (retry within budget, else fail and
648 /// escalate), records every decision durably, and returns an inspectable
649 /// status. It launches no new work and does not re-process tasks that
650 /// already reached a terminal state, so it is safe to call repeatedly.
651 pub fn resume_run(&self, run_id: &FleetRunId) -> Result<FleetResumeReport> {
652 self.resume_run_at(run_id, Utc::now())
653 }
654
655 /// Resume reconciliation at an explicit instant. This is the deterministic
656 /// seam behind `resume_run`'s wall clock: stale detection compares the
657 /// last heartbeat against `now`.
658 pub(crate) fn resume_run_at(
659 &self,
660 run_id: &FleetRunId,
661 now: DateTime<Utc>,
662 ) -> Result<FleetResumeReport> {
663 // Reuse the shared scheduler recovery engine over the same ledger so
664 // resume and steady-state supervision converge on one store and one
665 // retry/escalation policy. The manager's `stale_after` becomes the
666 // scheduler's heartbeat timeout so both surfaces agree on staleness.
667 let policy = FleetSchedulerPolicy {
668 heartbeat_timeout: self.stale_after,
669 ..FleetSchedulerPolicy::default()
670 };
671 let mut scheduler = FleetScheduler::open(&self.workspace, policy)?;
672 scheduler.set_now(now);
673 // Keep the lock order coordination -> ledger. A restart generation is
674 // durably prepared by the callback before the scheduler publishes the
675 // replacement lease, and an exact one-generation-ahead preparation is
676 // safe to consume after a process crash.
677 let mut coordination_guard = match &self.sub_agent_manager {
678 Some(manager) => Some(
679 manager
680 .try_write()
681 .map_err(|_| anyhow!("Fleet coordination state is busy; retry resume"))?,
682 ),
683 None => None,
684 };
685 let report = scheduler.resume_run_with_restart_callback(
686 run_id,
687 |state, task, task_spec, worker_id| {
688 if let Some(guard) = coordination_guard.as_mut() {
689 self.prepare_registered_restart_generation(
690 guard, state, task, task_spec, worker_id,
691 )?;
692 }
693 Ok(())
694 },
695 )?;
696 let status = self.run_status(run_id)?;
697 Ok(FleetResumeReport {
698 run_id: run_id.clone(),
699 reclaimed_stale: report.marked_stale,
700 restarted: report.restarted,
701 failed: report.failed,
702 escalated: report.alerts,
703 status,
704 })
705 }
706
707 pub async fn run_to_completion(
708 &self,
709 run_id: &FleetRunId,
710 max_workers: usize,
711 executor: &mut FleetExecutor,
712 codewhale_binary: &str,
713 model: Option<&str>,
714 tick_interval: Duration,
715 ) -> Result<FleetStatusSnapshot> {
716 let max_workers = max_workers.clamp(1, 128);
717 let manager_lock_path = self.manager_lock_path(run_id);
718 if let Some(parent) = manager_lock_path.parent() {
719 std::fs::create_dir_all(parent)
720 .with_context(|| format!("creating fleet manager lock dir {}", parent.display()))?;
721 }
722 let lock_file = OpenOptions::new()
723 .create(true)
724 .truncate(false)
725 .read(true)
726 .write(true)
727 .open(&manager_lock_path)
728 .with_context(|| {
729 format!("opening fleet manager lock {}", manager_lock_path.display())
730 })?;
731 let mut manager_lock = fd_lock::RwLock::new(lock_file);
732 let standby_interval = tick_interval
733 .min(Duration::from_millis(100))
734 .max(Duration::from_millis(10));
735 let mut observed_owner = false;
736 let _manager_guard = loop {
737 match manager_lock.try_write() {
738 Ok(guard) => {
739 if observed_owner {
740 if self.run_has_open_work(run_id)? {
741 bail!(
742 "fleet manager for run {} exited with open work; wait for stale reconciliation before resuming",
743 run_id.0
744 );
745 }
746 return self.run_status(run_id);
747 }
748 break guard;
749 }
750 Err(err) if err.kind() == ErrorKind::WouldBlock => {
751 // Another process owns this run. Wait for it to finish,
752 // but never treat lock release as permission to relaunch
753 // its unchanged leased attempts: an orphan child may still
754 // be alive after a crash. Stale reconciliation owns that
755 // recovery/generation transition.
756 observed_owner = true;
757 if !self.run_has_open_work(run_id)? {
758 return self.run_status(run_id);
759 }
760 tokio::time::sleep(standby_interval).await;
761 }
762 Err(err) => {
763 return Err(err).with_context(|| {
764 format!(
765 "locking fleet manager ownership {}",
766 manager_lock_path.display()
767 )
768 });
769 }
770 }
771 };
772 loop {
773 // A terminal ledger update can race the foreground host process.
774 // Do not lease new work onto a logical worker until its executor
775 // handle has been observed and forgotten below.
776 let unavailable_workers = executor.worker_ids().into_iter().collect();
777 let scheduling_error = self
778 .schedule_run_excluding(run_id, max_workers, &unavailable_workers)
779 .err();
780 self.drive_executor_tick(run_id, executor, codewhale_binary, model)?;
781 self.refresh_run_status(run_id)?;
782 if let Some(error) = scheduling_error {
783 if executor.worker_ids().is_empty() {
784 return Err(error).with_context(|| {
785 format!(
786 "scheduling Fleet run {} after draining owned workers",
787 run_id.0
788 )
789 });
790 }
791 tracing::warn!(
792 run_id = %run_id.0,
793 error = %error,
794 "Fleet scheduling paused while already-leased workers continue"
795 );
796 }
797 // A separate `fleet interrupt` process can make the ledger
798 // terminal while this manager still owns a live host child. Keep
799 // driving until the executor has observed that cancellation and
800 // stopped every tracked process.
801 if !self.run_has_open_work(run_id)? && executor.worker_ids().is_empty() {
802 return self.run_status(run_id);
803 }
804 tokio::time::sleep(tick_interval).await;
805 }
806 }
807
808 pub fn drive_executor_tick(
809 &self,
810 run_id: &FleetRunId,
811 executor: &mut FleetExecutor,
812 codewhale_binary: &str,
813 model: Option<&str>,
814 ) -> Result<FleetExecutorTickReport> {
815 let mut report = FleetExecutorTickReport::default();
816 report.started += self.start_leased_workers(run_id, executor, codewhale_binary, model)?;
817
818 for worker_id in executor.worker_ids() {
819 let tracked_attempt = executor.tracked_attempt(&worker_id);
820 if let Some(attempt) = tracked_attempt.as_ref()
821 && self
822 .executor_task_context_for_attempt(&worker_id, attempt)?
823 .is_none()
824 {
825 // The ledger advanced to another attempt (restart), or made
826 // this attempt terminal (cancel/stop), while this process was
827 // still alive. The executor owns the host handle, so fence and
828 // reap the old process without publishing any event against the
829 // replacement generation.
830 executor.stop_worker(&worker_id)?;
831 executor.forget_worker(&worker_id);
832 report.terminals += 1;
833 continue;
834 }
835 if tracked_attempt.is_none()
836 && let Some(_task) = self.cancelled_executor_task_context(&worker_id)?
837 {
838 // Cancellation is ledgered by an out-of-process control
839 // command. Only this executor owns the host process handle,
840 // so it must enforce the terminal state before returning from
841 // the foreground manager loop. Do not ingest output produced
842 // after cancellation; publish one final authoritative event
843 // after the process is actually stopped instead.
844 executor.stop_worker(&worker_id)?;
845 executor.forget_worker(&worker_id);
846 report.terminals += 1;
847 continue;
848 }
849
850 for payload in executor.drain_events(&worker_id) {
851 // The subprocess exit is the task-completion authority. Stream
852 // `done` / `error` lines are useful progress signals, but
853 // appending them as terminal ledger events before the process
854 // exits would free the logical worker too early.
855 if is_terminal_payload(&payload) {
856 continue;
857 }
858 let task = if let Some(attempt) = tracked_attempt.as_ref() {
859 self.executor_task_context_for_attempt(&worker_id, attempt)?
860 } else {
861 self.executor_task_context(&worker_id)?
862 };
863 let Some(task) = task else {
864 continue;
865 };
866 if self
867 .ledger
868 .append_event_if_leased(
869 &task.entry.run_id,
870 &worker_id,
871 &task.entry.task_id,
872 task.entry.attempts,
873 &timestamp(),
874 payload,
875 )?
876 .is_none()
877 {
878 continue;
879 }
880 self.ledger
881 .heartbeat(&worker_id, &timestamp(), None, None)?;
882 report.events += 1;
883 }
884
885 if let Some(terminal) = executor.poll_terminal_with_status(&worker_id) {
886 let task = if let Some(attempt) = tracked_attempt.as_ref() {
887 self.executor_task_context_for_attempt(&worker_id, attempt)?
888 } else {
889 self.executor_task_context(&worker_id)?
890 };
891 let Some(task) = task else {
892 executor.forget_worker(&worker_id);
893 continue;
894 };
895 if self.record_task_outcome(&task, terminal)? {
896 report.terminals += 1;
897 }
898 executor.forget_worker(&worker_id);
899 }
900 }
901
902 self.refresh_run_status(run_id)?;
903 Ok(report)
904 }
905
906 pub fn inspect_worker(&self, worker_id: &str) -> Result<FleetWorkerInspection> {
907 let state = self.ledger.rebuild_state()?;
908 let latest_event = latest_event_for_worker(&state, worker_id).cloned();
909 let current = active_task_for_worker(&state, worker_id)
910 .or_else(|| latest_task_for_worker(&state, worker_id));
911 let current_run_id = current.as_ref().map(|task| task.entry.run_id.clone());
912 let current_task_id = current.as_ref().map(|task| task.entry.task_id.clone());
913 let (objective, role) = current
914 .as_ref()
915 .and_then(|task| task_spec_for_state(&state, task))
916 .map(|task_spec| {
917 (
918 task_spec.objective.or(task_spec.description),
919 task_spec
920 .worker
921 .and_then(|worker| worker.role)
922 .map(|role| super::profile::canonical_public_role_name(role.trim())),
923 )
924 })
925 .unwrap_or((None, None));
926 let host = current_run_id
927 .as_ref()
928 .and_then(|run_id| worker_host_for_run(&state, run_id, worker_id));
929 let artifacts = state
930 .artifact_events
931 .values()
932 .filter(|event| event.worker_id == worker_id)
933 .filter_map(|event| match &event.payload {
934 FleetWorkerEventPayload::Artifact(artifact) => Some(artifact.clone()),
935 _ => None,
936 })
937 .chain(
938 state
939 .receipts
940 .values()
941 .filter(|receipt| receipt.worker_id == worker_id)
942 .flat_map(|receipt| receipt.artifacts.clone()),
943 )
944 .collect();
945 let receipt_summary = latest_receipt_for_worker(&state, worker_id).map(receipt_summary);
946 let last_error = latest_error_for_worker(&state, worker_id);
947 let status = state
948 .workers
949 .get(worker_id)
950 .cloned()
951 .unwrap_or(FleetWorkerStatus::Unknown);
952 let latest_heartbeat_at = state
953 .heartbeats
954 .get(worker_id)
955 .map(|heartbeat| heartbeat.timestamp.clone());
956 let alert_state = latest_alert_for_worker(&state, worker_id);
957
958 // Enrich only a live lease with its in-memory worker projection. A
959 // terminal durable task always wins over a lagging runtime record.
960 let runtime_state = current
961 .as_ref()
962 .filter(|task| task.status == FleetTaskLedgerStatus::Leased)
963 .and(self.sub_agent_manager.as_ref())
964 .and_then(|mgr| {
965 mgr.try_read()
966 .ok()
967 .and_then(|guard| guard.get_worker_record(worker_id))
968 .map(|record| FleetWorkerRuntimeProjection {
969 agent_status: format!("{:?}", record.status).to_lowercase(),
970 steps_taken: record.steps_taken,
971 latest_message: record.latest_message,
972 error: record.error,
973 result_summary: record.result_summary,
974 has_session: !matches!(
975 record.status,
976 crate::tools::subagent::AgentWorkerStatus::Completed
977 | crate::tools::subagent::AgentWorkerStatus::Failed
978 | crate::tools::subagent::AgentWorkerStatus::Cancelled
979 ),
980 })
981 });
982
983 Ok(FleetWorkerInspection {
984 worker_id: worker_id.to_string(),
985 status,
986 current_run_id,
987 current_task_id,
988 objective,
989 role,
990 host,
991 latest_heartbeat_at,
992 latest_event,
993 artifacts,
994 receipt_summary,
995 last_error,
996 alert_state,
997 runtime_state,
998 })
999 }
1000
1001 pub fn interrupt_worker(&self, worker_id: &str) -> Result<FleetWorkerInspection> {
1002 let state = self.ledger.rebuild_state()?;
1003 let Some(task) = active_task_for_worker(&state, worker_id) else {
1004 return Err(FleetControlError::NoActiveTask {
1005 worker_id: worker_id.to_string(),
1006 }
1007 .into());
1008 };
1009 let cancelled = self.ledger.cancel_task_if_active(
1010 &task.entry.run_id,
1011 &task.entry.task_id,
1012 Some(worker_id),
1013 &timestamp(),
1014 Some("operator"),
1015 Some("operator"),
1016 )?;
1017 if !cancelled {
1018 bail!("worker {worker_id} no longer has that running fleet task");
1019 }
1020 self.refresh_run_status(&task.entry.run_id)?;
1021 self.inspect_worker(worker_id)
1022 }
1023
1024 pub fn restart_worker(&self, worker_id: &str) -> Result<FleetRestartReport> {
1025 let state = self.ledger.rebuild_state()?;
1026 let Some(task) = active_task_for_worker(&state, worker_id)
1027 .or_else(|| latest_task_for_worker(&state, worker_id))
1028 else {
1029 bail!("worker {worker_id} has no fleet task to restart");
1030 };
1031 let run = state
1032 .runs
1033 .get(&task.entry.run_id.0)
1034 .ok_or_else(|| anyhow!("fleet run {} does not exist", task.entry.run_id.0))?;
1035 let max_workers = run
1036 .max_workers
1037 .unwrap_or_else(|| run.worker_specs.len().max(1))
1038 .clamp(1, 128);
1039 let mut coordination_guard = match &self.sub_agent_manager {
1040 Some(manager) => {
1041 let Ok(guard) = manager.try_write() else {
1042 bail!("Fleet worker {worker_id} coordination state is busy; retry restart");
1043 };
1044 Some(guard)
1045 }
1046 None => None,
1047 };
1048 let now = timestamp();
1049 let latest_seq = state
1050 .latest_seq
1051 .get(&event_key(
1052 worker_id,
1053 &task.entry.run_id.0,
1054 &task.entry.task_id,
1055 ))
1056 .copied()
1057 .unwrap_or(0);
1058 let heartbeat_at = state
1059 .heartbeats
1060 .get(worker_id)
1061 .map(|heartbeat| heartbeat.timestamp.as_str());
1062 let restarted = self.ledger.restart_task_if_unchanged_with_callback(
1063 &task.entry.run_id,
1064 &task.entry.task_id,
1065 worker_id,
1066 task.status,
1067 task.entry.attempts,
1068 latest_seq,
1069 heartbeat_at,
1070 &now,
1071 None,
1072 task.entry.attempts,
1073 || {
1074 if let Some(guard) = coordination_guard.as_mut() {
1075 let task_spec = run
1076 .task_specs
1077 .iter()
1078 .find(|spec| spec.id == task.entry.task_id)
1079 .ok_or_else(|| {
1080 anyhow!("fleet task {} does not exist", task.entry.task_id)
1081 })?;
1082 self.prepare_registered_restart_generation(
1083 guard, &state, task, task_spec, worker_id,
1084 )?;
1085 }
1086 Ok(())
1087 },
1088 )?;
1089 if !restarted {
1090 bail!("worker {worker_id} task changed before it could be restarted");
1091 }
1092 self.ledger
1093 .update_run_status(&task.entry.run_id, FleetRunStatus::Running, &timestamp())?;
1094 Ok(FleetRestartReport {
1095 run_id: task.entry.run_id.clone(),
1096 max_workers,
1097 inspection: self.inspect_worker(worker_id)?,
1098 })
1099 }
1100
1101 /// Prepare or consume the exact durable launch generation for one Fleet
1102 /// retry. The persisted one-generation-ahead record is the prepare marker:
1103 /// it is not launchable while the ledger remains on the old attempt, but a
1104 /// retry after a crash may validate and consume it idempotently.
1105 fn prepare_registered_restart_generation(
1106 &self,
1107 coordination: &mut SubAgentManager,
1108 state: &FleetLedgerState,
1109 task: &FleetTaskState,
1110 task_spec: &FleetTaskSpec,
1111 worker_id: &str,
1112 ) -> Result<()> {
1113 let run = state
1114 .runs
1115 .get(&task.entry.run_id.0)
1116 .ok_or_else(|| anyhow!("fleet run {} does not exist", task.entry.run_id.0))?;
1117 let worker_spec = run
1118 .worker_specs
1119 .iter()
1120 .find(|worker| worker.id == worker_id)
1121 .cloned()
1122 .unwrap_or_else(|| default_local_worker(worker_id));
1123 let cwd = resolve_task_cwd(&self.workspace, task_spec)?;
1124 validate_task_cwd_for_host(&self.workspace, &worker_spec.host, &cwd)?;
1125 let roster = self.agent_roster();
1126 let expected_current = bind_fleet_launch_attempt(
1127 worker_runtime::apply_exec_hardening(
1128 worker_runtime::fleet_task_to_worker_spec_with_profiles(
1129 worker_id,
1130 &task.entry.run_id.0,
1131 task_spec,
1132 &worker_spec,
1133 self.run_model(),
1134 &cwd,
1135 &self.workspace,
1136 roster.members(),
1137 None,
1138 )?,
1139 &self.exec_config,
1140 ),
1141 task.entry.attempts,
1142 );
1143 let record = coordination
1144 .get_worker_record(worker_id)
1145 .ok_or_else(|| anyhow!("Fleet worker {worker_id} has no registered launch spec"))?;
1146 let current_generation = task.entry.attempts.max(1);
1147 let next_generation = current_generation
1148 .checked_add(1)
1149 .ok_or_else(|| anyhow!("Fleet worker {worker_id} exhausted launch generations"))?;
1150 let registered_generation = record
1151 .spec
1152 .launch_manifest
1153 .as_ref()
1154 .map(|manifest| manifest.generation)
1155 .ok_or_else(|| anyhow!("Fleet worker {worker_id} has no persisted launch manifest"))?;
1156
1157 match registered_generation {
1158 generation if generation == current_generation => {
1159 validate_registered_launch_spec(&record.spec, &expected_current)?;
1160 coordination
1161 .advance_registered_worker_generation(bind_fleet_launch_attempt(
1162 record.spec,
1163 next_generation,
1164 ))
1165 .map_err(anyhow::Error::msg)?;
1166 }
1167 generation if generation == next_generation => {
1168 let mut normalized = record.spec;
1169 normalized
1170 .launch_manifest
1171 .as_mut()
1172 .expect("prepared launch manifest checked above")
1173 .generation = current_generation;
1174 validate_registered_launch_spec(&normalized, &expected_current)?;
1175 }
1176 generation => {
1177 bail!(
1178 "Fleet worker {worker_id} persisted launch generation {generation} does not match ledger attempt {current_generation} or its prepared retry {next_generation}"
1179 );
1180 }
1181 }
1182 Ok(())
1183 }
1184
1185 pub fn stop_all(&self) -> Result<usize> {
1186 let state = self.ledger.rebuild_state()?;
1187 let now = timestamp();
1188 let mut affected_runs = BTreeSet::new();
1189 let mut stopped = 0usize;
1190 for task in state.tasks.values() {
1191 if !matches!(
1192 task.status,
1193 FleetTaskLedgerStatus::Enqueued | FleetTaskLedgerStatus::Leased
1194 ) {
1195 continue;
1196 }
1197 if !self.ledger.cancel_task_if_active(
1198 &task.entry.run_id,
1199 &task.entry.task_id,
1200 None,
1201 &now,
1202 Some("stop_all"),
1203 Some("operator"),
1204 )? {
1205 continue;
1206 }
1207 affected_runs.insert(task.entry.run_id.0.clone());
1208 stopped += 1;
1209 }
1210 for run_id in affected_runs {
1211 self.ledger.update_run_status(
1212 &FleetRunId::from(run_id),
1213 FleetRunStatus::Cancelled,
1214 &timestamp(),
1215 )?;
1216 }
1217 Ok(stopped)
1218 }
1219
1220 pub fn stop_run(&self, run_id: &FleetRunId) -> Result<usize> {
1221 let state = self.ledger.rebuild_state()?;
1222 if !state.runs.contains_key(&run_id.0) {
1223 bail!("fleet run {} does not exist", run_id.0);
1224 }
1225 let now = timestamp();
1226 let mut stopped = 0usize;
1227 for task in state
1228 .tasks
1229 .values()
1230 .filter(|task| task.entry.run_id == *run_id)
1231 {
1232 if !matches!(
1233 task.status,
1234 FleetTaskLedgerStatus::Enqueued | FleetTaskLedgerStatus::Leased
1235 ) {
1236 continue;
1237 }
1238 if !self.ledger.cancel_task_if_active(
1239 &task.entry.run_id,
1240 &task.entry.task_id,
1241 None,
1242 &now,
1243 Some("stop_run"),
1244 Some("operator"),
1245 )? {
1246 continue;
1247 }
1248 stopped += 1;
1249 }
1250 self.ledger
1251 .update_run_status(run_id, FleetRunStatus::Cancelled, &timestamp())?;
1252 Ok(stopped)
1253 }
1254
1255 fn start_worker_task(
1256 &self,
1257 worker_id: &str,
1258 entry: &FleetInboxEntry,
1259 task_spec: &FleetTaskSpec,
1260 max_active_for_run: Option<usize>,
1261 ) -> Result<bool> {
1262 let run = self
1263 .ledger
1264 .rebuild_state()?
1265 .runs
1266 .get(&entry.run_id.0)
1267 .cloned()
1268 .ok_or_else(|| anyhow!("fleet run {} does not exist", entry.run_id.0))?;
1269 let worker_spec = run
1270 .worker_specs
1271 .iter()
1272 .find(|worker| worker.id == worker_id)
1273 .cloned()
1274 .unwrap_or_else(|| default_local_worker(worker_id));
1275 let worker_workspace = resolve_task_cwd(&self.workspace, task_spec)?;
1276 validate_task_cwd_for_host(&self.workspace, &worker_spec.host, &worker_workspace)?;
1277 let roster = self.agent_roster();
1278 let sub_agent_worker = bind_fleet_launch_attempt(
1279 worker_runtime::apply_exec_hardening(
1280 worker_runtime::fleet_task_to_worker_spec_with_profiles(
1281 worker_id,
1282 &entry.run_id.0,
1283 task_spec,
1284 &worker_spec,
1285 self.run_model(),
1286 &worker_workspace,
1287 &self.workspace,
1288 roster.members(),
1289 None,
1290 )?,
1291 &self.exec_config,
1292 ),
1293 entry.attempts.saturating_add(1),
1294 );
1295 authority_envelope_for_worker(&sub_agent_worker, task_spec)?;
1296 let log_artifact = self.write_log_artifact(&entry.run_id, worker_id, task_spec)?;
1297 // Hold the coordination manager from pure preflight through the
1298 // ledger's pre-commit projection callback and any append compensation. This keeps a concurrent
1299 // agent/Fleet registration from invalidating the overlap decision in
1300 // between, while a busy manager simply leaves the task queued for the
1301 // next scheduler tick.
1302 let mut coordination_guard = match &self.sub_agent_manager {
1303 Some(manager) => {
1304 let Ok(mut guard) = manager.try_write() else {
1305 return Ok(false);
1306 };
1307 if let Err(error) = guard.preflight_worker_coordination(&sub_agent_worker) {
1308 tracing::debug!(
1309 worker_id,
1310 run_id = %entry.run_id.0,
1311 task_id = %entry.task_id,
1312 error = %error,
1313 "Fleet worker coordination is not currently admissible"
1314 );
1315 return Ok(false);
1316 }
1317 Some(guard)
1318 }
1319 None => None,
1320 };
1321 let registration_snapshot = coordination_guard
1322 .as_ref()
1323 .map(|guard| guard.coordination_registration_snapshot());
1324 let mut registration_succeeded = false;
1325 let now = timestamp();
1326 let start_result = self.ledger.start_task_if_enqueued(
1327 &entry.run_id,
1328 &entry.task_id,
1329 worker_id,
1330 &now,
1331 None,
1332 max_active_for_run,
1333 vec![
1334 FleetWorkerEventPayload::Leased {
1335 lease_expires_at: None,
1336 },
1337 FleetWorkerEventPayload::Starting,
1338 FleetWorkerEventPayload::Artifact(log_artifact),
1339 FleetWorkerEventPayload::Running,
1340 ],
1341 || {
1342 // Registration shares the durable transition lock, so a
1343 // cancellation cannot win between claim and projection setup.
1344 if let Some(guard) = coordination_guard.as_mut() {
1345 guard
1346 .register_worker_with_coordination(sub_agent_worker)
1347 .map_err(anyhow::Error::msg)?;
1348 registration_succeeded = true;
1349 }
1350 Ok(())
1351 },
1352 );
1353 let started = match start_result {
1354 Ok(started) => started,
1355 Err(start_error) => {
1356 if registration_succeeded
1357 && let (Some(guard), Some(snapshot)) =
1358 (coordination_guard.as_mut(), registration_snapshot)
1359 && let Err(rollback_error) =
1360 guard.restore_coordination_registration_snapshot(snapshot)
1361 {
1362 return Err(anyhow!("{start_error:#}; additionally {rollback_error}"));
1363 }
1364 return Err(start_error);
1365 }
1366 };
1367 if !started {
1368 return Ok(false);
1369 }
1370
1371 Ok(true)
1372 }
1373
1374 fn start_leased_workers(
1375 &self,
1376 run_id: &FleetRunId,
1377 executor: &mut FleetExecutor,
1378 codewhale_binary: &str,
1379 model: Option<&str>,
1380 ) -> Result<usize> {
1381 let state = self.ledger.rebuild_state()?;
1382 let run = state
1383 .runs
1384 .get(&run_id.0)
1385 .cloned()
1386 .ok_or_else(|| anyhow!("fleet run {} does not exist", run_id.0))?;
1387 let roster = self.agent_roster();
1388 let mut started = 0usize;
1389 for task in active_tasks_for_run(&state, run_id) {
1390 let Some(worker_id) = task.leased_to.as_deref() else {
1391 continue;
1392 };
1393 if executor.is_tracking(worker_id) {
1394 continue;
1395 }
1396 let Some(task_spec) = run
1397 .task_specs
1398 .iter()
1399 .find(|spec| spec.id == task.entry.task_id)
1400 .cloned()
1401 else {
1402 continue;
1403 };
1404 let worker_spec = run
1405 .worker_specs
1406 .iter()
1407 .find(|worker| worker.id == worker_id)
1408 .cloned()
1409 .unwrap_or_else(|| default_local_worker(worker_id));
1410 let coordination_record = if let Some(manager) = self.sub_agent_manager.as_ref() {
1411 let Ok(guard) = manager.try_read() else {
1412 continue;
1413 };
1414 Some(guard.get_worker_record(worker_id))
1415 } else {
1416 None
1417 };
1418 let preparation = (|| -> Result<_> {
1419 let cwd = resolve_task_cwd(&self.workspace, &task_spec)?;
1420 validate_task_cwd_for_host(&self.workspace, &worker_spec.host, &cwd)?;
1421 let expected_launch_spec = bind_fleet_launch_attempt(
1422 worker_runtime::apply_exec_hardening(
1423 worker_runtime::fleet_task_to_worker_spec_with_profiles(
1424 worker_id,
1425 &run_id.0,
1426 &task_spec,
1427 &worker_spec,
1428 self.run_model(),
1429 &cwd,
1430 &self.workspace,
1431 roster.members(),
1432 None,
1433 )?,
1434 &self.exec_config,
1435 ),
1436 task.entry.attempts,
1437 );
1438 let launch_spec = match coordination_record {
1439 Some(Some(record)) => {
1440 validate_registered_launch_spec(&record.spec, &expected_launch_spec)?;
1441 record.spec
1442 }
1443 Some(None) => {
1444 bail!("Fleet worker {worker_id} has no coordination-registered launch spec")
1445 }
1446 None => expected_launch_spec,
1447 };
1448 let command = build_worker_exec_command_with_launch_spec(
1449 codewhale_binary,
1450 &task_spec,
1451 &launch_spec,
1452 &self.exec_config,
1453 model,
1454 roster.members(),
1455 )?;
1456 Ok((cwd, command))
1457 })();
1458 let attempt = FleetExecutorAttempt {
1459 run_id: task.entry.run_id.clone(),
1460 task_id: task.entry.task_id.clone(),
1461 attempt: task.entry.attempts,
1462 };
1463 let (cwd, command) = match preparation {
1464 Ok(prepared) => prepared,
1465 Err(err) => {
1466 let task = FleetExecutorTaskContext {
1467 entry: task.entry.clone(),
1468 task_spec,
1469 worker_id: worker_id.to_string(),
1470 };
1471 let terminal = FleetWorkerTerminalEvent {
1472 payload: FleetWorkerEventPayload::Failed {
1473 reason: format!("worker launch preparation failed: {err:#}"),
1474 recoverable: false,
1475 },
1476 exit_code: None,
1477 tail_payloads: Vec::new(),
1478 reported_route: None,
1479 requires_reported_route: false,
1480 };
1481 let _ = self.record_task_outcome(&task, terminal)?;
1482 continue;
1483 }
1484 };
1485 match executor.start_worker_attempt_on_host(
1486 worker_id,
1487 &worker_spec.host,
1488 command,
1489 Some(cwd),
1490 attempt,
1491 ) {
1492 Ok(handle) => {
1493 let artifact = self.host_log_artifact(&handle.log_path);
1494 if self
1495 .ledger
1496 .append_event_if_leased(
1497 run_id,
1498 worker_id,
1499 &task.entry.task_id,
1500 task.entry.attempts,
1501 &timestamp(),
1502 FleetWorkerEventPayload::Artifact(artifact),
1503 )?
1504 .is_none()
1505 {
1506 executor.stop_worker(worker_id)?;
1507 executor.forget_worker(worker_id);
1508 continue;
1509 }
1510 started += 1;
1511 }
1512 Err(err) => {
1513 let recoverable = matches!(err.kind, FleetHostErrorKind::Retryable);
1514 let task = FleetExecutorTaskContext {
1515 entry: task.entry.clone(),
1516 task_spec,
1517 worker_id: worker_id.to_string(),
1518 };
1519 let terminal = FleetWorkerTerminalEvent {
1520 payload: FleetWorkerEventPayload::Failed {
1521 reason: err.message,
1522 recoverable,
1523 },
1524 exit_code: None,
1525 tail_payloads: Vec::new(),
1526 reported_route: None,
1527 requires_reported_route: false,
1528 };
1529 let _ = self.record_task_outcome(&task, terminal)?;
1530 }
1531 }
1532 }
1533 Ok(started)
1534 }
1535
1536 fn executor_task_context(&self, worker_id: &str) -> Result<Option<FleetExecutorTaskContext>> {
1537 let state = self.ledger.rebuild_state()?;
1538 let Some(task) = active_task_for_worker(&state, worker_id)
1539 .or_else(|| latest_task_for_worker(&state, worker_id))
1540 else {
1541 return Ok(None);
1542 };
1543 let Some(run) = state.runs.get(&task.entry.run_id.0) else {
1544 return Ok(None);
1545 };
1546 let Some(task_spec) = run
1547 .task_specs
1548 .iter()
1549 .find(|spec| spec.id == task.entry.task_id)
1550 .cloned()
1551 else {
1552 return Ok(None);
1553 };
1554 Ok(Some(FleetExecutorTaskContext {
1555 entry: task.entry.clone(),
1556 task_spec,
1557 worker_id: worker_id.to_string(),
1558 }))
1559 }
1560
1561 fn executor_task_context_for_attempt(
1562 &self,
1563 worker_id: &str,
1564 attempt: &FleetExecutorAttempt,
1565 ) -> Result<Option<FleetExecutorTaskContext>> {
1566 let state = self.ledger.rebuild_state()?;
1567 let key = task_key(&attempt.run_id.0, &attempt.task_id);
1568 let Some(task) = state.tasks.get(&key) else {
1569 return Ok(None);
1570 };
1571 if task.status != FleetTaskLedgerStatus::Leased
1572 || task.leased_to.as_deref() != Some(worker_id)
1573 || task.entry.attempts != attempt.attempt
1574 {
1575 return Ok(None);
1576 }
1577 let Some(run) = state.runs.get(&attempt.run_id.0) else {
1578 return Ok(None);
1579 };
1580 let Some(task_spec) = run
1581 .task_specs
1582 .iter()
1583 .find(|spec| spec.id == attempt.task_id)
1584 .cloned()
1585 else {
1586 return Ok(None);
1587 };
1588 Ok(Some(FleetExecutorTaskContext {
1589 entry: task.entry.clone(),
1590 task_spec,
1591 worker_id: worker_id.to_string(),
1592 }))
1593 }
1594
1595 fn cancelled_executor_task_context(
1596 &self,
1597 worker_id: &str,
1598 ) -> Result<Option<FleetExecutorTaskContext>> {
1599 let state = self.ledger.rebuild_state()?;
1600 let Some(task) = latest_task_for_worker(&state, worker_id) else {
1601 return Ok(None);
1602 };
1603 if task.status != FleetTaskLedgerStatus::Cancelled {
1604 return Ok(None);
1605 }
1606 let Some(run) = state.runs.get(&task.entry.run_id.0) else {
1607 return Ok(None);
1608 };
1609 let Some(task_spec) = run
1610 .task_specs
1611 .iter()
1612 .find(|spec| spec.id == task.entry.task_id)
1613 .cloned()
1614 else {
1615 return Ok(None);
1616 };
1617 Ok(Some(FleetExecutorTaskContext {
1618 entry: task.entry.clone(),
1619 task_spec,
1620 worker_id: worker_id.to_string(),
1621 }))
1622 }
1623
1624 fn record_task_outcome(
1625 &self,
1626 task: &FleetExecutorTaskContext,
1627 terminal: FleetWorkerTerminalEvent,
1628 ) -> Result<bool> {
1629 let state = self.ledger.rebuild_state()?;
1630 let key = task_key(&task.entry.run_id.0, &task.entry.task_id);
1631 let Some(current) = state.tasks.get(&key) else {
1632 return Ok(false);
1633 };
1634 if current.status != FleetTaskLedgerStatus::Leased
1635 || current.leased_to.as_deref() != Some(task.worker_id.as_str())
1636 || current.entry.attempts != task.entry.attempts
1637 {
1638 return Ok(false);
1639 }
1640
1641 let FleetWorkerTerminalEvent {
1642 payload,
1643 exit_code,
1644 tail_payloads,
1645 reported_route,
1646 requires_reported_route,
1647 } = terminal;
1648 let (receipt_result, failure_kind, exit_code) = task_receipt_outcome(&payload, exit_code);
1649 let terminal_completed = matches!(&payload, FleetWorkerEventPayload::Completed { .. });
1650 let expected_terminal_status = match &payload {
1651 FleetWorkerEventPayload::Completed { .. } => FleetTaskLedgerStatus::Completed,
1652 FleetWorkerEventPayload::Failed { .. } => FleetTaskLedgerStatus::Failed,
1653 FleetWorkerEventPayload::Cancelled { .. } => FleetTaskLedgerStatus::Cancelled,
1654 _ => bail!("fleet executor outcome must contain a terminal worker event"),
1655 };
1656 for tail_payload in tail_payloads {
1657 if is_terminal_payload(&tail_payload) {
1658 continue;
1659 }
1660 if self
1661 .ledger
1662 .append_event_if_leased(
1663 &task.entry.run_id,
1664 &task.worker_id,
1665 &task.entry.task_id,
1666 task.entry.attempts,
1667 &timestamp(),
1668 tail_payload,
1669 )?
1670 .is_none()
1671 {
1672 return Ok(false);
1673 }
1674 }
1675 let artifacts = self.task_artifacts_for_receipt(
1676 &task.entry.run_id,
1677 &task.entry.task_id,
1678 &task.worker_id,
1679 )?;
1680 // A terminal worker report is the sole authority for provider/model
1681 // actually used. Never re-resolve those fields through manager-local
1682 // config: remote workers may intentionally run a different config.
1683 // A headless worker that omits or malforms the terminal route fails
1684 // closed to no actual route. Pre-launch transport/simulated paths have
1685 // no process evidence by design, so they retain the explicitly labeled
1686 // intent route rather than pretending it was observed.
1687 let resolved_route = match (reported_route.as_ref(), requires_reported_route) {
1688 (Some(reported_route), _) => {
1689 self.resolve_reported_task_route(&task.task_spec, reported_route)
1690 }
1691 (None, true) => None,
1692 (None, false) => self.resolve_task_route(&task.task_spec),
1693 };
1694 let effective_permissions = self.resolve_task_effective_permissions(task);
1695 let verification_input = FleetTaskVerificationInput {
1696 run_id: task.entry.run_id.clone(),
1697 task_id: task.entry.task_id.clone(),
1698 worker_id: task.worker_id.clone(),
1699 attempt: task.entry.attempts,
1700 exit_code,
1701 artifacts,
1702 resolved_route,
1703 effective_permissions,
1704 };
1705 let receipt = if task.task_spec.scorer.is_some() || terminal_completed {
1706 let verification =
1707 verify_task_result(&self.workspace, &task.task_spec, &verification_input);
1708 prepare_verification_receipt(&self.workspace, &verification_input, verification)?
1709 } else {
1710 FleetReceipt {
1711 run_id: task.entry.run_id.clone(),
1712 task_id: task.entry.task_id.clone(),
1713 worker_id: task.worker_id.clone(),
1714 attempt: Some(task.entry.attempts),
1715 terminal_seq: None,
1716 completed_at: timestamp(),
1717 result: receipt_result,
1718 failure_kind,
1719 artifacts: verification_input.artifacts,
1720 score: None,
1721 resolved_route: verification_input.resolved_route,
1722 effective_permissions: verification_input.effective_permissions,
1723 }
1724 };
1725 let final_status = (matches!(
1726 receipt.result,
1727 FleetTaskResult::Fail | FleetTaskResult::Timeout
1728 ) && expected_terminal_status != FleetTaskLedgerStatus::Failed)
1729 .then_some(FleetTaskLedgerStatus::Failed);
1730 Ok(self
1731 .ledger
1732 .finalize_task_attempt_if_leased(
1733 &task.entry.run_id,
1734 &task.worker_id,
1735 &task.entry.task_id,
1736 task.entry.attempts,
1737 &timestamp(),
1738 payload,
1739 final_status,
1740 receipt,
1741 )?
1742 .is_some())
1743 }
1744
1745 /// Resolve the route snapshot to persist on a task's receipt (#3154).
1746 ///
1747 /// Loads the merged agent roster so role/loadout intent composes the same
1748 /// way as the worker-spec path, then mints a secret-free route candidate via
1749 /// the hermetic resolver bridge. Returns `None` (never a fabricated route)
1750 /// when resolution is unavailable.
1751 fn resolve_task_route(&self, task_spec: &FleetTaskSpec) -> Option<FleetResolvedRoute> {
1752 let roster = self.agent_roster();
1753 worker_runtime::resolve_fleet_route_with_config(
1754 task_spec,
1755 roster.members(),
1756 self.session_model(),
1757 self.route_config.as_ref(),
1758 )
1759 }
1760
1761 fn resolve_reported_task_route(
1762 &self,
1763 task_spec: &FleetTaskSpec,
1764 reported_route: &FleetWorkerReportedRoute,
1765 ) -> Option<FleetResolvedRoute> {
1766 let roster = self.agent_roster();
1767 worker_runtime::resolve_fleet_route_from_worker_report(
1768 task_spec,
1769 roster.members(),
1770 self.session_model(),
1771 &reported_route.provider,
1772 reported_route.provider_exact_id.as_deref(),
1773 &reported_route.model,
1774 )
1775 }
1776
1777 /// The adopted session route, if any — the operator's model.
1778 fn session_model(&self) -> Option<&str> {
1779 self.session_model.as_deref()
1780 }
1781
1782 /// Resolve the effective worker authority to persist on a task's receipt
1783 /// (#3211). This mirrors Fleet worker registration and applies exec
1784 /// hardening before snapshotting the runtime profile. Failures degrade to
1785 /// `None` so receipt writing never widens or fabricates authority.
1786 fn resolve_task_effective_permissions(
1787 &self,
1788 task: &FleetExecutorTaskContext,
1789 ) -> Option<FleetEffectivePermissions> {
1790 let state = self.ledger.rebuild_state().ok()?;
1791 let run = state.runs.get(&task.entry.run_id.0)?;
1792 let worker_spec = run
1793 .worker_specs
1794 .iter()
1795 .find(|worker| worker.id == task.worker_id)
1796 .cloned()
1797 .unwrap_or_else(|| default_local_worker(&task.worker_id));
1798 let roster = self.agent_roster();
1799 let worker = worker_runtime::fleet_task_to_worker_spec_with_profiles(
1800 &task.worker_id,
1801 &task.entry.run_id.0,
1802 &task.task_spec,
1803 &worker_spec,
1804 self.run_model(),
1805 &self.workspace,
1806 &self.workspace,
1807 roster.members(),
1808 None,
1809 )
1810 .ok()?;
1811 let worker = bind_fleet_launch_attempt(
1812 worker_runtime::apply_exec_hardening(worker, &self.exec_config),
1813 task.entry.attempts,
1814 );
1815 Some(worker_runtime::fleet_effective_permissions_for_task(
1816 &task.task_spec,
1817 roster.members(),
1818 &worker,
1819 ))
1820 }
1821
1822 fn task_artifacts_for_receipt(
1823 &self,
1824 run_id: &FleetRunId,
1825 task_id: &str,
1826 worker_id: &str,
1827 ) -> Result<Vec<FleetArtifactRef>> {
1828 let state = self.ledger.rebuild_state()?;
1829 Ok(state
1830 .artifact_events
1831 .values()
1832 .filter(|event| {
1833 event.run_id == *run_id && event.task_id == task_id && event.worker_id == worker_id
1834 })
1835 .filter_map(|event| match &event.payload {
1836 FleetWorkerEventPayload::Artifact(artifact) => {
1837 Some(self.refresh_artifact_size(artifact.clone()))
1838 }
1839 _ => None,
1840 })
1841 .collect())
1842 }
1843
1844 fn refresh_artifact_size(&self, mut artifact: FleetArtifactRef) -> FleetArtifactRef {
1845 let path = if artifact.path.is_absolute() {
1846 artifact.path.clone()
1847 } else {
1848 self.workspace.join(&artifact.path)
1849 };
1850 artifact.size_bytes = std::fs::metadata(path).ok().map(|meta| meta.len());
1851 artifact
1852 }
1853
1854 fn host_log_artifact(&self, path: &Path) -> FleetArtifactRef {
1855 let rel_path = path
1856 .strip_prefix(&self.workspace)
1857 .map(Path::to_path_buf)
1858 .unwrap_or_else(|_| path.to_path_buf());
1859 let size_bytes = std::fs::metadata(path).ok().map(|meta| meta.len());
1860 FleetArtifactRef {
1861 kind: FleetArtifactKind::Log,
1862 path: rel_path,
1863 checksum: None,
1864 mime_type: Some("application/x-ndjson".to_string()),
1865 size_bytes,
1866 }
1867 }
1868
1869 fn append_worker_event(
1870 &self,
1871 run_id: &FleetRunId,
1872 worker_id: &str,
1873 task_id: &str,
1874 payload: FleetWorkerEventPayload,
1875 ) -> Result<FleetWorkerEvent> {
1876 self.ledger
1877 .append_event_next_seq(run_id, worker_id, task_id, &timestamp(), payload)
1878 }
1879
1880 fn write_log_artifact(
1881 &self,
1882 run_id: &FleetRunId,
1883 worker_id: &str,
1884 task_spec: &FleetTaskSpec,
1885 ) -> Result<FleetArtifactRef> {
1886 let rel_path = PathBuf::from(".codewhale")
1887 .join("fleet")
1888 .join(safe_path_segment(&run_id.0))
1889 .join(safe_path_segment(&task_spec.id))
1890 .join(format!("{}.log", safe_path_segment(worker_id)));
1891 let abs_path = self.workspace.join(&rel_path);
1892 if let Some(parent) = abs_path.parent() {
1893 std::fs::create_dir_all(parent)
1894 .with_context(|| format!("creating fleet artifact dir {}", parent.display()))?;
1895 }
1896 let contents = format!(
1897 "run_id={}\ntask_id={}\ntask_name={}\nworker_id={}\nstatus=started\n",
1898 run_id.0, task_spec.id, task_spec.name, worker_id
1899 );
1900 std::fs::write(&abs_path, contents)
1901 .with_context(|| format!("writing fleet worker log {}", abs_path.display()))?;
1902 let size_bytes = std::fs::metadata(&abs_path).ok().map(|m| m.len());
1903 Ok(FleetArtifactRef {
1904 kind: FleetArtifactKind::Log,
1905 path: rel_path,
1906 checksum: None,
1907 mime_type: Some("text/plain".to_string()),
1908 size_bytes,
1909 })
1910 }
1911
1912 fn refresh_run_status(&self, run_id: &FleetRunId) -> Result<()> {
1913 let state = self.ledger.rebuild_state()?;
1914 let mut has_queued = false;
1915 let mut has_running = false;
1916 let mut has_failed = false;
1917 let mut has_cancelled = false;
1918 let mut has_tasks = false;
1919 for task in state
1920 .tasks
1921 .values()
1922 .filter(|task| task.entry.run_id == *run_id)
1923 {
1924 has_tasks = true;
1925 match task.status {
1926 FleetTaskLedgerStatus::Enqueued => has_queued = true,
1927 FleetTaskLedgerStatus::Leased => has_running = true,
1928 FleetTaskLedgerStatus::Failed => has_failed = true,
1929 FleetTaskLedgerStatus::Cancelled => has_cancelled = true,
1930 FleetTaskLedgerStatus::Completed => {}
1931 }
1932 }
1933 let status = if !has_tasks {
1934 FleetRunStatus::Completed
1935 } else if has_queued || has_running {
1936 FleetRunStatus::Running
1937 } else if has_failed {
1938 FleetRunStatus::Failed
1939 } else if has_cancelled {
1940 FleetRunStatus::Cancelled
1941 } else {
1942 FleetRunStatus::Completed
1943 };
1944 self.ledger
1945 .update_run_status(run_id, status, &timestamp())
1946 .context("updating fleet run status")
1947 }
1948
1949 fn status_from_state(
1950 &self,
1951 run_filter: Option<&FleetRunId>,
1952 state: &FleetLedgerState,
1953 ) -> FleetStatusSnapshot {
1954 let mut snapshot = FleetStatusSnapshot {
1955 runs: state.runs.len(),
1956 workers: state.workers.clone(),
1957 ..FleetStatusSnapshot::default()
1958 };
1959 for task in state.tasks.values() {
1960 if run_filter.is_some_and(|run_id| task.entry.run_id != *run_id) {
1961 continue;
1962 }
1963 match task.status {
1964 FleetTaskLedgerStatus::Enqueued => snapshot.queued += 1,
1965 FleetTaskLedgerStatus::Leased => {
1966 if self.task_is_stale(task, state) {
1967 snapshot.stale += 1;
1968 } else {
1969 snapshot.running += 1;
1970 }
1971 }
1972 FleetTaskLedgerStatus::Completed => snapshot.completed += 1,
1973 FleetTaskLedgerStatus::Failed => snapshot.failed += 1,
1974 FleetTaskLedgerStatus::Cancelled => snapshot.cancelled += 1,
1975 }
1976 }
1977 for receipt in state.receipts.values() {
1978 if run_filter.is_some_and(|run_id| receipt.run_id != *run_id) {
1979 continue;
1980 }
1981 if receipt.result == FleetTaskResult::Partial {
1982 snapshot.partial += 1;
1983 }
1984 match &receipt.failure_kind {
1985 Some(FleetTaskFailureKind::Transport) => snapshot.transport_failed += 1,
1986 Some(FleetTaskFailureKind::Task) => snapshot.task_failed += 1,
1987 Some(FleetTaskFailureKind::Verifier) => snapshot.verifier_failed += 1,
1988 None => {}
1989 }
1990 }
1991 snapshot.restarted = state
1992 .restarted_events
1993 .values()
1994 .filter(|event| run_filter.is_none_or(|run_id| event.run_id == *run_id))
1995 .count();
1996 snapshot.escalated = state
1997 .escalated_events
1998 .values()
1999 .filter(|event| run_filter.is_none_or(|run_id| event.run_id == *run_id))
2000 .count();
2001 snapshot
2002 }
2003
2004 fn task_is_stale(&self, task: &FleetTaskState, state: &FleetLedgerState) -> bool {
2005 let Some(worker_id) = task.leased_to.as_deref() else {
2006 return true;
2007 };
2008 let Some(heartbeat) = state.heartbeats.get(worker_id) else {
2009 return true;
2010 };
2011 let Ok(last) = DateTime::parse_from_rfc3339(&heartbeat.timestamp) else {
2012 return true;
2013 };
2014 let age = Utc::now().signed_duration_since(last.with_timezone(&Utc));
2015 age.to_std()
2016 .is_ok_and(|duration| duration > self.stale_after)
2017 }
2018 }
2019
2020 fn default_local_workers(run_id: &FleetRunId, max_workers: usize) -> Vec<FleetWorkerSpec> {
2021 (1..=max_workers)
2022 .map(|index| {
2023 default_local_worker_with_name(&format!("{}-local-{}", run_id.0, index), index)
2024 })
2025 .collect()
2026 }
2027
2028 fn default_local_worker_with_name(worker_id: &str, index: usize) -> FleetWorkerSpec {
2029 FleetWorkerSpec {
2030 id: worker_id.to_string(),
2031 name: format!("Local worker {index}"),
2032 host: FleetHostSpec::Local,
2033 trust_level: Some(FleetTrustLevel::Local),
2034 labels: BTreeMap::new(),
2035 capabilities: vec!["local".to_string()],
2036 max_concurrent_tasks: Some(1),
2037 }
2038 }
2039
2040 fn default_local_worker(worker_id: &str) -> FleetWorkerSpec {
2041 FleetWorkerSpec {
2042 id: worker_id.to_string(),
2043 name: worker_id.to_string(),
2044 host: FleetHostSpec::Local,
2045 trust_level: Some(FleetTrustLevel::Local),
2046 labels: BTreeMap::new(),
2047 capabilities: vec!["local".to_string()],
2048 max_concurrent_tasks: Some(1),
2049 }
2050 }
2051
2052 fn worker_ids_for_run(run: &FleetRun, max_workers: usize) -> Vec<String> {
2053 run.worker_specs
2054 .iter()
2055 .take(max_workers)
2056 .map(|worker| worker.id.clone())
2057 .collect()
2058 }
2059
2060 fn active_workers_for_run(state: &FleetLedgerState, run_id: &FleetRunId) -> BTreeSet<String> {
2061 active_tasks_for_run(state, run_id)
2062 .filter_map(|task| task.leased_to.clone())
2063 .collect()
2064 }
2065
2066 fn active_tasks_for_run<'a>(
2067 state: &'a FleetLedgerState,
2068 run_id: &'a FleetRunId,
2069 ) -> impl Iterator<Item = &'a FleetTaskState> {
2070 state.tasks.values().filter(move |task| {
2071 task.entry.run_id == *run_id && matches!(task.status, FleetTaskLedgerStatus::Leased)
2072 })
2073 }
2074
2075 fn active_task_for_worker<'a>(
2076 state: &'a FleetLedgerState,
2077 worker_id: &str,
2078 ) -> Option<&'a FleetTaskState> {
2079 state.tasks.values().find(|task| {
2080 task.leased_to.as_deref() == Some(worker_id)
2081 && matches!(task.status, FleetTaskLedgerStatus::Leased)
2082 })
2083 }
2084
2085 fn latest_task_for_worker<'a>(
2086 state: &'a FleetLedgerState,
2087 worker_id: &str,
2088 ) -> Option<&'a FleetTaskState> {
2089 state
2090 .tasks
2091 .values()
2092 .filter(|task| task.leased_to.as_deref() == Some(worker_id))
2093 .max_by_key(|task| task.completed_at.as_deref().or(task.leased_at.as_deref()))
2094 }
2095
2096 fn next_enqueued_task_for_run(
2097 state: &FleetLedgerState,
2098 run_id: &FleetRunId,
2099 ) -> Option<(FleetInboxEntry, FleetTaskSpec)> {
2100 let run = state.runs.get(&run_id.0)?;
2101 let task = state
2102 .tasks
2103 .values()
2104 .filter(|task| {
2105 task.entry.run_id == *run_id && matches!(task.status, FleetTaskLedgerStatus::Enqueued)
2106 })
2107 .min_by_key(|task| {
2108 (
2109 task.entry.priority,
2110 task.entry.enqueued_at.clone(),
2111 task.entry.task_id.clone(),
2112 )
2113 })?;
2114 let task_spec = run
2115 .task_specs
2116 .iter()
2117 .find(|spec| spec.id == task.entry.task_id)
2118 .cloned()?;
2119 Some((task.entry.clone(), task_spec))
2120 }
2121
2122 fn task_spec_for_state(state: &FleetLedgerState, task: &FleetTaskState) -> Option<FleetTaskSpec> {
2123 state
2124 .runs
2125 .get(&task.entry.run_id.0)?
2126 .task_specs
2127 .iter()
2128 .find(|spec| spec.id == task.entry.task_id)
2129 .cloned()
2130 }
2131
2132 fn worker_host_for_run(
2133 state: &FleetLedgerState,
2134 run_id: &FleetRunId,
2135 worker_id: &str,
2136 ) -> Option<String> {
2137 let run = state.runs.get(&run_id.0)?;
2138 let worker = run
2139 .worker_specs
2140 .iter()
2141 .find(|worker| worker.id == worker_id)?;
2142 Some(host_label(&worker.host))
2143 }
2144
2145 fn host_label(host: &FleetHostSpec) -> String {
2146 match host {
2147 FleetHostSpec::Local => "local".to_string(),
2148 FleetHostSpec::Ssh { host, .. } => format!("ssh:{host}"),
2149 FleetHostSpec::Docker { image, .. } => format!("docker:{image}"),
2150 }
2151 }
2152
2153 fn latest_event_for_worker<'a>(
2154 state: &'a FleetLedgerState,
2155 worker_id: &str,
2156 ) -> Option<&'a FleetWorkerEvent> {
2157 state
2158 .latest_events
2159 .values()
2160 .filter(|event| event.worker_id == worker_id)
2161 .max_by_key(|event| event.seq)
2162 }
2163
2164 fn latest_alert_for_worker(state: &FleetLedgerState, worker_id: &str) -> Option<String> {
2165 state
2166 .escalated_events
2167 .values()
2168 .filter(|event| event.worker_id == worker_id)
2169 .filter_map(|event| match &event.payload {
2170 FleetWorkerEventPayload::Escalated { channel, alert_id } => Some((
2171 event.seq,
2172 alert_id
2173 .as_ref()
2174 .map(|alert_id| format!("escalated via {channel} alert_id={alert_id}"))
2175 .unwrap_or_else(|| format!("escalated via {channel}")),
2176 )),
2177 _ => None,
2178 })
2179 .max_by_key(|(seq, _)| *seq)
2180 .map(|(_, message)| message)
2181 }
2182
2183 fn latest_receipt_for_worker<'a>(
2184 state: &'a FleetLedgerState,
2185 worker_id: &str,
2186 ) -> Option<&'a FleetReceipt> {
2187 state
2188 .receipts
2189 .values()
2190 .filter(|receipt| receipt.worker_id == worker_id)
2191 .max_by_key(|receipt| &receipt.completed_at)
2192 }
2193
2194 fn receipt_summary(receipt: &FleetReceipt) -> String {
2195 let result = match receipt.result {
2196 FleetTaskResult::Pass => "pass",
2197 FleetTaskResult::Partial => "partial",
2198 FleetTaskResult::Fail => "fail",
2199 FleetTaskResult::Skip => "skip",
2200 FleetTaskResult::Timeout => "timeout",
2201 };
2202 let mut summary = format!("result={result}");
2203 if let Some(kind) = &receipt.failure_kind {
2204 let kind = match kind {
2205 FleetTaskFailureKind::Transport => "transport",
2206 FleetTaskFailureKind::Task => "task",
2207 FleetTaskFailureKind::Verifier => "verifier",
2208 };
2209 summary.push_str(&format!(" failure_kind={kind}"));
2210 }
2211 if let Some(notes) = receipt
2212 .score
2213 .as_ref()
2214 .and_then(|score| score.notes.as_deref())
2215 .filter(|notes| !notes.trim().is_empty())
2216 {
2217 summary.push_str(&format!(" notes={notes}"));
2218 }
2219 summary
2220 }
2221
2222 fn latest_error_for_worker(state: &FleetLedgerState, worker_id: &str) -> Option<String> {
2223 state
2224 .latest_events
2225 .values()
2226 .filter(|event| event.worker_id == worker_id)
2227 .filter_map(|event| match &event.payload {
2228 FleetWorkerEventPayload::Failed { reason, .. } => {
2229 Some((event.seq, format!("failed: {reason}")))
2230 }
2231 FleetWorkerEventPayload::Cancelled { cancelled_by } => Some((
2232 event.seq,
2233 cancelled_by
2234 .as_ref()
2235 .map(|by| format!("cancelled by {by}"))
2236 .unwrap_or_else(|| "cancelled".to_string()),
2237 )),
2238 FleetWorkerEventPayload::Interrupted { signal } => Some((
2239 event.seq,
2240 signal
2241 .as_ref()
2242 .map(|signal| format!("interrupted by {signal}"))
2243 .unwrap_or_else(|| "interrupted".to_string()),
2244 )),
2245 FleetWorkerEventPayload::Stale { last_heartbeat_at } => Some((
2246 event.seq,
2247 last_heartbeat_at
2248 .as_ref()
2249 .map(|ts| format!("stale since {ts}"))
2250 .unwrap_or_else(|| "stale".to_string()),
2251 )),
2252 _ => None,
2253 })
2254 .max_by_key(|(seq, _)| *seq)
2255 .map(|(_, message)| message)
2256 }
2257
2258 fn task_priority(task: &FleetTaskSpec) -> i32 {
2259 task.metadata
2260 .get("priority")
2261 .and_then(Value::as_i64)
2262 .and_then(|value| i32::try_from(value).ok())
2263 .unwrap_or(0)
2264 }
2265
2266 fn resolve_task_cwd(workspace: &Path, task: &FleetTaskSpec) -> Result<PathBuf> {
2267 let Some(root) = task
2268 .workspace
2269 .as_ref()
2270 .and_then(|workspace| workspace.root.as_ref())
2271 else {
2272 return crate::tools::spec::resolve_strict_authority_path(
2273 &crate::tools::ToolContext::new(workspace.to_path_buf()),
2274 ".",
2275 )
2276 .map_err(anyhow::Error::new);
2277 };
2278 crate::tools::spec::resolve_strict_authority_path(
2279 &crate::tools::ToolContext::new(workspace.to_path_buf()),
2280 &root.to_string_lossy(),
2281 )
2282 .map_err(anyhow::Error::new)
2283 }
2284
2285 fn bind_fleet_launch_attempt(mut spec: AgentWorkerSpec, attempt: u32) -> AgentWorkerSpec {
2286 // The outer machine-readable cap is workspace-relative and cannot yet be
2287 // intersected into a grandchild's narrower launch context. Fleet workers
2288 // are therefore truthful leaves in v0.9.1: the nested-agent surface is
2289 // disabled for the authority-bound subprocess.
2290 spec.max_spawn_depth = 0;
2291 spec.runtime_profile.max_spawn_depth = 0;
2292 if let Some(manifest) = spec.launch_manifest.as_mut() {
2293 manifest.generation = attempt.max(1);
2294 manifest.profile.max_spawn_depth = 0;
2295 }
2296 spec
2297 }
2298
2299 fn validate_registered_launch_spec(
2300 registered: &AgentWorkerSpec,
2301 expected: &AgentWorkerSpec,
2302 ) -> Result<()> {
2303 let Some(registered_manifest) = registered.launch_manifest.as_ref() else {
2304 bail!(
2305 "Fleet worker {} has no persisted launch manifest",
2306 registered.worker_id
2307 );
2308 };
2309 if registered_manifest.prompt != registered.objective
2310 || !registered_prompt_matches_expected(&registered.objective, &expected.objective)
2311 {
2312 bail!(
2313 "Fleet worker {} has an inconsistent persisted prompt",
2314 registered.worker_id
2315 );
2316 }
2317
2318 // Coordination may append a bounded decision projection to the prompt.
2319 // Every identity, route, permission, workspace, scope, and attempt field
2320 // must otherwise match a fresh derivation from this exact leased task.
2321 let mut registered_identity = registered.clone();
2322 let mut expected_identity = expected.clone();
2323 registered_identity.objective.clear();
2324 expected_identity.objective.clear();
2325 if let Some(manifest) = registered_identity.launch_manifest.as_mut() {
2326 manifest.prompt.clear();
2327 }
2328 if let Some(manifest) = expected_identity.launch_manifest.as_mut() {
2329 manifest.prompt.clear();
2330 }
2331 if registered_identity != expected_identity {
2332 bail!(
2333 "Fleet worker {} persisted launch spec does not match the exact task lease and attempt",
2334 registered.worker_id
2335 );
2336 }
2337 Ok(())
2338 }
2339
2340 fn registered_prompt_matches_expected(registered: &str, expected: &str) -> bool {
2341 const HEADER: &str = "Accepted coordination decisions relevant to this child (bounded):\n";
2342 if registered == expected {
2343 return true;
2344 }
2345 let Some(projection) = registered
2346 .strip_prefix(expected)
2347 .and_then(|suffix| suffix.strip_prefix("\n\n"))
2348 else {
2349 return false;
2350 };
2351 let Some(lines) = projection.strip_prefix(HEADER) else {
2352 return false;
2353 };
2354 !lines.is_empty()
2355 && projection.len() <= 4096
2356 && lines.lines().count() <= 8
2357 && lines
2358 .lines()
2359 .all(|line| line.starts_with("- ") && line.len() <= 512)
2360 }
2361
2362 fn validate_task_cwd_for_host(
2363 workspace: &Path,
2364 host: &FleetHostSpec,
2365 task_cwd: &Path,
2366 ) -> Result<()> {
2367 if !matches!(host, FleetHostSpec::Ssh { .. }) {
2368 return Ok(());
2369 }
2370 let workspace_root = crate::tools::spec::resolve_strict_authority_path(
2371 &crate::tools::ToolContext::new(workspace.to_path_buf()),
2372 ".",
2373 )
2374 .map_err(anyhow::Error::new)?;
2375 if task_cwd != workspace_root {
2376 bail!(
2377 "SSH Fleet workers do not yet support nested workspace.root values; task cwd '{}' cannot be mapped safely beneath the remote working_directory",
2378 task_cwd.display()
2379 );
2380 }
2381 Ok(())
2382 }
2383
2384 fn task_receipt_outcome(
2385 payload: &FleetWorkerEventPayload,
2386 exit_code: Option<i32>,
2387 ) -> (FleetTaskResult, Option<FleetTaskFailureKind>, Option<i32>) {
2388 match payload {
2389 FleetWorkerEventPayload::Completed {
2390 exit_code: payload_exit_code,
2391 ..
2392 } => (
2393 FleetTaskResult::Pass,
2394 None,
2395 exit_code.or(*payload_exit_code),
2396 ),
2397 FleetWorkerEventPayload::Cancelled { .. } => (FleetTaskResult::Skip, None, exit_code),
2398 FleetWorkerEventPayload::Failed { .. } => {
2399 let failure_kind = if exit_code.is_none() {
2400 FleetTaskFailureKind::Transport
2401 } else {
2402 FleetTaskFailureKind::Task
2403 };
2404 (FleetTaskResult::Fail, Some(failure_kind), exit_code)
2405 }
2406 _ => (FleetTaskResult::Partial, None, exit_code),
2407 }
2408 }
2409
2410 fn is_terminal_payload(payload: &FleetWorkerEventPayload) -> bool {
2411 matches!(
2412 payload,
2413 FleetWorkerEventPayload::Completed { .. }
2414 | FleetWorkerEventPayload::Failed { .. }
2415 | FleetWorkerEventPayload::Cancelled { .. }
2416 | FleetWorkerEventPayload::Interrupted { .. }
2417 )
2418 }
2419
2420 fn task_key(run_id: &str, task_id: &str) -> String {
2421 format!("{run_id}:{task_id}")
2422 }
2423
2424 fn event_key(worker_id: &str, run_id: &str, task_id: &str) -> String {
2425 format!("{worker_id}:{run_id}:{task_id}")
2426 }
2427
2428 fn timestamp() -> String {
2429 Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true)
2430 }
2431
2432 fn safe_path_segment(value: &str) -> String {
2433 value
2434 .chars()
2435 .map(|ch| {
2436 if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
2437 ch
2438 } else {
2439 '_'
2440 }
2441 })
2442 .collect()
2443 }
2444
2445 #[cfg(test)]
2446 mod tests {
2447 use super::*;
2448 use serde_json::json;
2449 use tempfile::TempDir;
2450
2451 fn task(id: &str) -> FleetTaskSpec {
2452 FleetTaskSpec {
2453 id: id.to_string(),
2454 name: id.to_string(),
2455 description: None,
2456 objective: Some(format!("Complete {id}")),
2457 instructions: format!("do {id}"),
2458 worker: Some(FleetTaskWorkerProfile {
2459 agent_profile: None,
2460 role: Some("reviewer".to_string()),
2461 loadout: None,
2462 model_class: None,
2463 model: None,
2464 tool_profile: Some("read-only".to_string()),
2465 tools: Vec::new(),
2466 capabilities: Vec::new(),
2467 }),
2468 workspace: None,
2469 input_files: Vec::new(),
2470 context: Vec::new(),
2471 budget: None,
2472 tags: Vec::new(),
2473 expected_artifacts: vec![FleetArtifactKind::Log],
2474 scorer: None,
2475 retry_policy: None,
2476 alert_policy: None,
2477 timeout_seconds: None,
2478 metadata: BTreeMap::new(),
2479 }
2480 }
2481
2482 #[test]
2483 fn ssh_workers_fail_closed_for_nested_task_roots() {
2484 let tmp = TempDir::new().unwrap();
2485 std::fs::create_dir(tmp.path().join("nested")).unwrap();
2486 let mut nested = task("nested");
2487 nested.workspace = Some(FleetWorkspaceRequirements {
2488 root: Some(PathBuf::from("nested")),
2489 ..FleetWorkspaceRequirements::default()
2490 });
2491 let nested_cwd = resolve_task_cwd(tmp.path(), &nested).unwrap();
2492 let ssh = FleetHostSpec::Ssh {
2493 host: "builder.example.test".to_string(),
2494 port: None,
2495 user: None,
2496 identity: None,
2497 known_hosts: None,
2498 host_key_fingerprint: None,
2499 working_directory: Some(PathBuf::from("/srv/codewhale")),
2500 env_allowlist: Vec::new(),
2501 codewhale_binary: Some("/usr/local/bin/codewhale".to_string()),
2502 };
2503
2504 let error = validate_task_cwd_for_host(tmp.path(), &ssh, &nested_cwd)
2505 .expect_err("nested SSH task roots must fail closed");
2506 assert!(error.to_string().contains("cannot be mapped safely"));
2507
2508 let root_cwd = resolve_task_cwd(tmp.path(), &task("root")).unwrap();
2509 validate_task_cwd_for_host(tmp.path(), &ssh, &root_cwd).unwrap();
2510 validate_task_cwd_for_host(tmp.path(), &FleetHostSpec::Local, &nested_cwd).unwrap();
2511 }
2512
2513 #[test]
2514 fn ssh_nested_root_failure_commits_neither_lease_nor_coordination_record() {
2515 let tmp = TempDir::new().unwrap();
2516 std::fs::create_dir(tmp.path().join("nested")).unwrap();
2517 let coordination =
2518 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2519 let manager = FleetManager::open(tmp.path())
2520 .unwrap()
2521 .with_sub_agent_manager(coordination.clone());
2522 let mut nested = task("nested");
2523 nested.worker = Some(FleetTaskWorkerProfile {
2524 agent_profile: None,
2525 role: Some("reviewer".to_string()),
2526 loadout: None,
2527 model_class: None,
2528 model: None,
2529 tool_profile: Some("read-only".to_string()),
2530 tools: Vec::new(),
2531 capabilities: Vec::new(),
2532 });
2533 nested.workspace = Some(FleetWorkspaceRequirements {
2534 root: Some(PathBuf::from("nested")),
2535 ..FleetWorkspaceRequirements::default()
2536 });
2537 let worker = FleetWorkerSpec {
2538 id: "ssh-worker".to_string(),
2539 name: "SSH worker".to_string(),
2540 host: FleetHostSpec::Ssh {
2541 host: "builder.example.test".to_string(),
2542 port: None,
2543 user: None,
2544 identity: None,
2545 known_hosts: None,
2546 host_key_fingerprint: None,
2547 working_directory: Some(PathBuf::from("/srv/codewhale")),
2548 env_allowlist: Vec::new(),
2549 codewhale_binary: Some("/usr/local/bin/codewhale".to_string()),
2550 },
2551 trust_level: None,
2552 labels: BTreeMap::new(),
2553 capabilities: Vec::new(),
2554 max_concurrent_tasks: Some(1),
2555 };
2556 let error = manager
2557 .create_run(
2558 FleetTaskSpecDocument {
2559 name: Some("nested SSH".to_string()),
2560 labels: BTreeMap::new(),
2561 security_policy: None,
2562 workers: vec![worker],
2563 tasks: vec![nested],
2564 },
2565 1,
2566 )
2567 .expect_err("nested SSH root must fail before leasing");
2568 assert!(error.to_string().contains("cannot be mapped safely"));
2569
2570 let state = manager.rebuild_state().unwrap();
2571 let task = state.tasks.values().next().expect("queued task remains");
2572 assert_eq!(task.status, FleetTaskLedgerStatus::Enqueued);
2573 assert_eq!(task.entry.attempts, 0);
2574 assert!(task.leased_to.is_none());
2575 assert!(
2576 coordination
2577 .try_read()
2578 .unwrap()
2579 .list_worker_records()
2580 .is_empty()
2581 );
2582 }
2583
2584 #[test]
2585 fn ledger_append_failure_rolls_back_coordination_registration() {
2586 let tmp = TempDir::new().unwrap();
2587 let coordination =
2588 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2589 let manager = FleetManager::open(tmp.path())
2590 .unwrap()
2591 .with_sub_agent_manager(coordination.clone());
2592 let mut spec = task("read-only");
2593 spec.worker = Some(FleetTaskWorkerProfile {
2594 agent_profile: None,
2595 role: Some("reviewer".to_string()),
2596 loadout: None,
2597 model_class: None,
2598 model: None,
2599 tool_profile: Some("read-only".to_string()),
2600 tools: Vec::new(),
2601 capabilities: Vec::new(),
2602 });
2603 manager.ledger.fail_next_start_append_after_callback();
2604
2605 let error = manager
2606 .create_run(
2607 FleetTaskSpecDocument {
2608 name: Some("forced rollback".to_string()),
2609 labels: BTreeMap::new(),
2610 security_policy: None,
2611 workers: Vec::new(),
2612 tasks: vec![spec],
2613 },
2614 1,
2615 )
2616 .expect_err("forced append failure");
2617 assert!(error.to_string().contains("forced Fleet ledger append"));
2618
2619 let state = manager.rebuild_state().unwrap();
2620 let task = state.tasks.values().next().expect("queued task remains");
2621 assert_eq!(task.status, FleetTaskLedgerStatus::Enqueued);
2622 assert_eq!(task.entry.attempts, 0);
2623 assert!(task.leased_to.is_none());
2624 let guard = coordination.try_read().unwrap();
2625 assert!(guard.list_worker_records().is_empty());
2626 assert!(guard.coordination_snapshot().write_claims.is_empty());
2627 }
2628
2629 #[test]
2630 fn invalid_worker_identity_is_rejected_before_run_journal_creation() {
2631 let tmp = TempDir::new().unwrap();
2632 let manager = FleetManager::open(tmp.path()).unwrap();
2633 let error = manager
2634 .create_run(
2635 FleetTaskSpecDocument {
2636 name: Some("invalid identity".to_string()),
2637 labels: BTreeMap::new(),
2638 security_policy: None,
2639 workers: vec![resume_worker_spec("worker\r\nforged")],
2640 tasks: vec![task("task-a")],
2641 },
2642 1,
2643 )
2644 .expect_err("multiline worker identity must fail before journaling");
2645 assert!(
2646 error
2647 .to_string()
2648 .contains("worker id must be a simple ASCII token")
2649 );
2650
2651 let state = manager.rebuild_state().unwrap();
2652 assert!(state.runs.is_empty());
2653 assert!(state.tasks.is_empty());
2654 }
2655
2656 #[test]
2657 fn queued_creation_waits_for_explicit_idempotent_start() {
2658 let tmp = TempDir::new().unwrap();
2659 let coordination =
2660 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2661 let manager = FleetManager::open(tmp.path())
2662 .unwrap()
2663 .with_sub_agent_manager(coordination.clone());
2664 let report = manager
2665 .create_queued_run_with_descriptor(
2666 FleetTaskSpecDocument {
2667 name: Some("managed launch gate".to_string()),
2668 labels: BTreeMap::new(),
2669 security_policy: None,
2670 workers: Vec::new(),
2671 tasks: vec![task("task-a")],
2672 },
2673 1,
2674 ManagedFleetRunDescriptor {
2675 target: Some(FleetRuntimeTarget::ThisComputer),
2676 workflow: Some(FleetWorkflowDescriptor {
2677 id: "managed-launch".to_string(),
2678 kind: FleetWorkflowKind::Parallel,
2679 }),
2680 roles: vec!["reviewer".to_string()],
2681 },
2682 )
2683 .unwrap();
2684
2685 let queued = manager.rebuild_state().unwrap();
2686 assert_eq!(queued.runs[&report.run_id.0].status, FleetRunStatus::Queued);
2687 assert_eq!(
2688 queued.tasks[&task_key(&report.run_id.0, "task-a")].status,
2689 FleetTaskLedgerStatus::Enqueued
2690 );
2691 assert!(
2692 coordination
2693 .try_read()
2694 .unwrap()
2695 .list_worker_records()
2696 .is_empty()
2697 );
2698
2699 let activated = manager.activate_run(&report.run_id).unwrap();
2700 assert_eq!(activated.leased, 0);
2701 let activated_state = manager.rebuild_state().unwrap();
2702 assert_eq!(
2703 activated_state.tasks[&task_key(&report.run_id.0, "task-a")].status,
2704 FleetTaskLedgerStatus::Enqueued
2705 );
2706 assert!(
2707 coordination
2708 .try_read()
2709 .unwrap()
2710 .list_worker_records()
2711 .is_empty()
2712 );
2713
2714 let started = manager.start_run(&report.run_id).unwrap();
2715 assert_eq!(started.leased, 1);
2716 let running = manager.rebuild_state().unwrap();
2717 let task = &running.tasks[&task_key(&report.run_id.0, "task-a")];
2718 assert_eq!(task.status, FleetTaskLedgerStatus::Leased);
2719 assert_eq!(task.entry.attempts, 1);
2720 assert_eq!(
2721 running.run_status_overrides[&report.run_id.0],
2722 FleetRunStatus::Running
2723 );
2724 assert_eq!(
2725 coordination.try_read().unwrap().list_worker_records().len(),
2726 1
2727 );
2728
2729 let repeated = manager.start_run(&report.run_id).unwrap();
2730 assert_eq!(repeated.leased, 0);
2731 let repeated_state = manager.rebuild_state().unwrap();
2732 assert_eq!(
2733 repeated_state.tasks[&task_key(&report.run_id.0, "task-a")]
2734 .entry
2735 .attempts,
2736 1
2737 );
2738 assert_eq!(
2739 coordination.try_read().unwrap().list_worker_records().len(),
2740 1
2741 );
2742 }
2743
2744 #[test]
2745 fn busy_coordination_yields_without_spinning_or_leasing() {
2746 let tmp = TempDir::new().unwrap();
2747 let coordination =
2748 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2749 let manager = FleetManager::open(tmp.path())
2750 .unwrap()
2751 .with_sub_agent_manager(coordination.clone());
2752 let prepared = manager
2753 .create_queued_run(
2754 FleetTaskSpecDocument {
2755 name: Some("busy coordination".to_string()),
2756 labels: BTreeMap::new(),
2757 security_policy: None,
2758 workers: Vec::new(),
2759 tasks: vec![task("task-a")],
2760 },
2761 1,
2762 )
2763 .unwrap();
2764
2765 let guard = coordination.try_write().unwrap();
2766 let blocked = manager.start_run(&prepared.run_id).unwrap();
2767 assert_eq!(blocked.leased, 0);
2768 let blocked_state = manager.rebuild_state().unwrap();
2769 assert_eq!(
2770 blocked_state.tasks[&task_key(&prepared.run_id.0, "task-a")].status,
2771 FleetTaskLedgerStatus::Enqueued
2772 );
2773
2774 drop(guard);
2775 let retried = manager.start_run(&prepared.run_id).unwrap();
2776 assert_eq!(retried.leased, 1);
2777 }
2778
2779 #[test]
2780 fn cross_run_write_contention_leaves_later_work_queued_until_claim_releases() {
2781 let tmp = TempDir::new().unwrap();
2782 let coordination =
2783 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2784 let manager = FleetManager::open(tmp.path())
2785 .unwrap()
2786 .with_sub_agent_manager(coordination);
2787 let write_task = |id: &str| {
2788 let mut task = task(id);
2789 task.worker = Some(FleetTaskWorkerProfile {
2790 agent_profile: None,
2791 role: Some("builder".to_string()),
2792 loadout: None,
2793 model_class: None,
2794 model: None,
2795 tool_profile: Some("explicit".to_string()),
2796 tools: vec!["apply_patch".to_string()],
2797 capabilities: Vec::new(),
2798 });
2799 task.workspace = Some(FleetWorkspaceRequirements {
2800 writable_paths: vec![PathBuf::from("src")],
2801 ..FleetWorkspaceRequirements::default()
2802 });
2803 task
2804 };
2805 let prepare = |name: &str, task_id: &str| {
2806 manager
2807 .create_queued_run(
2808 FleetTaskSpecDocument {
2809 name: Some(name.to_string()),
2810 labels: BTreeMap::new(),
2811 security_policy: None,
2812 workers: Vec::new(),
2813 tasks: vec![write_task(task_id)],
2814 },
2815 1,
2816 )
2817 .unwrap()
2818 };
2819 let first = prepare("first writer", "write-a");
2820 let second = prepare("second writer", "write-b");
2821
2822 assert_eq!(manager.start_run(&first.run_id).unwrap().leased, 1);
2823 let blocked = manager.start_run(&second.run_id).unwrap();
2824 assert_eq!(blocked.leased, 0);
2825 assert_eq!(
2826 manager.rebuild_state().unwrap().tasks[&task_key(&second.run_id.0, "write-b")].status,
2827 FleetTaskLedgerStatus::Enqueued
2828 );
2829
2830 assert_eq!(manager.stop_run(&first.run_id).unwrap(), 1);
2831 assert_eq!(manager.start_run(&second.run_id).unwrap().leased, 1);
2832 }
2833
2834 #[test]
2835 fn restored_lease_without_launch_record_fails_durably_instead_of_poisoning_ticks() {
2836 let tmp = TempDir::new().unwrap();
2837 let coordination =
2838 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 4);
2839 let empty_snapshot = coordination
2840 .try_read()
2841 .unwrap()
2842 .coordination_registration_snapshot();
2843 let manager = FleetManager::open(tmp.path())
2844 .unwrap()
2845 .with_sub_agent_manager(coordination.clone());
2846 let mut spec = task("task-a");
2847 spec.worker = Some(FleetTaskWorkerProfile {
2848 agent_profile: None,
2849 role: Some("reviewer".to_string()),
2850 loadout: None,
2851 model_class: None,
2852 model: None,
2853 tool_profile: Some("read-only".to_string()),
2854 tools: Vec::new(),
2855 capabilities: Vec::new(),
2856 });
2857 let report = manager
2858 .create_run(
2859 FleetTaskSpecDocument {
2860 name: Some("restored lease".to_string()),
2861 labels: BTreeMap::new(),
2862 security_policy: None,
2863 workers: Vec::new(),
2864 tasks: vec![spec],
2865 },
2866 1,
2867 )
2868 .unwrap();
2869 coordination
2870 .try_write()
2871 .unwrap()
2872 .restore_coordination_registration_snapshot(empty_snapshot)
2873 .unwrap();
2874
2875 let mut executor = FleetExecutor::new(tmp.path());
2876 manager
2877 .drive_executor_tick(&report.run_id, &mut executor, "unused-codewhale", None)
2878 .expect("missing restored launch state must become a durable task failure");
2879 manager
2880 .drive_executor_tick(&report.run_id, &mut executor, "unused-codewhale", None)
2881 .expect("the next scheduler tick must not remain poisoned");
2882
2883 let state = manager.rebuild_state().unwrap();
2884 let key = task_key(&report.run_id.0, "task-a");
2885 assert_eq!(state.tasks[&key].status, FleetTaskLedgerStatus::Failed);
2886 assert_eq!(state.receipts[&key].result, FleetTaskResult::Fail);
2887 assert!(
2888 latest_error_for_worker(&state, &report.worker_ids[0])
2889 .is_some_and(|error| error.contains("no coordination-registered launch spec"))
2890 );
2891 }
2892
2893 fn read_only_launch_spec(workspace: &Path, task_id: &str, attempt: u32) -> AgentWorkerSpec {
2894 let mut spec = task(task_id);
2895 spec.worker = Some(FleetTaskWorkerProfile {
2896 agent_profile: None,
2897 role: Some("reviewer".to_string()),
2898 loadout: None,
2899 model_class: None,
2900 model: None,
2901 tool_profile: Some("read-only".to_string()),
2902 tools: Vec::new(),
2903 capabilities: Vec::new(),
2904 });
2905 bind_fleet_launch_attempt(
2906 worker_runtime::fleet_task_to_worker_spec_with_profiles(
2907 "worker-1",
2908 "run-1",
2909 &spec,
2910 &default_local_worker("worker-1"),
2911 "auto",
2912 workspace,
2913 workspace,
2914 &[],
2915 None,
2916 )
2917 .unwrap(),
2918 attempt,
2919 )
2920 }
2921
2922 #[test]
2923 fn persisted_launch_identity_is_bound_to_task_and_attempt() {
2924 let tmp = TempDir::new().unwrap();
2925 let task_a = read_only_launch_spec(tmp.path(), "task-a", 1);
2926 let task_b = read_only_launch_spec(tmp.path(), "task-b", 1);
2927 let task_b_retry = read_only_launch_spec(tmp.path(), "task-b", 2);
2928
2929 assert_eq!(task_b.max_spawn_depth, 0);
2930 assert_eq!(task_b.runtime_profile.max_spawn_depth, 0);
2931 assert_eq!(
2932 task_b
2933 .launch_manifest
2934 .as_ref()
2935 .unwrap()
2936 .profile
2937 .max_spawn_depth,
2938 0
2939 );
2940 validate_registered_launch_spec(&task_b, &task_b).unwrap();
2941 assert!(validate_registered_launch_spec(&task_a, &task_b).is_err());
2942 assert!(validate_registered_launch_spec(&task_b, &task_b_retry).is_err());
2943
2944 let mut projected = task_b.clone();
2945 projected.objective.push_str(
2946 "\n\nAccepted coordination decisions relevant to this child (bounded):\n- api v1 [decision-1] owner=planner: keep scope bounded",
2947 );
2948 projected.launch_manifest.as_mut().unwrap().prompt = projected.objective.clone();
2949 validate_registered_launch_spec(&projected, &task_b)
2950 .expect("a bounded coordination prompt projection preserves launch identity");
2951
2952 let mut corrupt = task_b.clone();
2953 corrupt.objective.push_str("\n\narbitrary stale prompt");
2954 corrupt.launch_manifest.as_mut().unwrap().prompt = corrupt.objective.clone();
2955 assert!(validate_registered_launch_spec(&corrupt, &task_b).is_err());
2956 }
2957
2958 #[test]
2959 fn with_session_model_adopts_route_and_ignores_auto_or_empty() {
2960 let tmp = TempDir::new().unwrap();
2961
2962 // No session: legacy auto sentinel.
2963 let manager = FleetManager::open(tmp.path()).unwrap();
2964 assert_eq!(manager.run_model(), "auto");
2965 assert_eq!(manager.session_model(), None);
2966
2967 // The session route becomes the run model — the operator's model.
2968 let manager = FleetManager::open(tmp.path())
2969 .unwrap()
2970 .with_session_model("deepseek-v4-pro");
2971 assert_eq!(manager.run_model(), "deepseek-v4-pro");
2972 assert_eq!(manager.session_model(), Some("deepseek-v4-pro"));
2973
2974 // "auto" and empty/whitespace inputs keep the resolver default.
2975 for noop in ["auto", "AUTO", "", " "] {
2976 let manager = FleetManager::open(tmp.path())
2977 .unwrap()
2978 .with_session_model(noop);
2979 assert_eq!(manager.run_model(), "auto");
2980 assert_eq!(manager.session_model(), None);
2981 }
2982 }
2983
2984 fn task_spec_file(dir: &TempDir, tasks: Vec<FleetTaskSpec>) -> PathBuf {
2985 let path = dir.path().join("fleet-tasks.json");
2986 let doc = json!({
2987 "name": "manager smoke",
2988 "tasks": tasks,
2989 });
2990 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
2991 path
2992 }
2993
2994 /// Read the process RSS (Resident Set Size) in kilobytes from
2995 /// `/proc/self/status`. Returns `None` when the file is unavailable
2996 /// (non-Linux) or the `VmRSS` field is missing.
2997 #[cfg(target_os = "linux")]
2998 fn rss_kb() -> Option<u64> {
2999 let status = std::fs::read_to_string("/proc/self/status").ok()?;
3000 status
3001 .lines()
3002 .find(|line| line.starts_with("VmRSS:"))
3003 .and_then(|line| line.split_whitespace().nth(1))
3004 .and_then(|v| v.parse().ok())
3005 }
3006
3007 #[cfg(unix)]
3008 fn fake_codewhale(dir: &TempDir, body: &str) -> PathBuf {
3009 use std::os::unix::fs::PermissionsExt;
3010
3011 let path = dir.path().join("fake-codewhale");
3012 std::fs::write(&path, body).unwrap();
3013 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
3014 permissions.set_mode(0o755);
3015 std::fs::set_permissions(&path, permissions).unwrap();
3016 path
3017 }
3018
3019 #[cfg(unix)]
3020 fn complete_with_fake_codewhale(
3021 manager: &FleetManager,
3022 run_id: &FleetRunId,
3023 max_workers: usize,
3024 binary: &Path,
3025 ) -> FleetStatusSnapshot {
3026 let rt = tokio::runtime::Runtime::new().unwrap();
3027 let mut executor = FleetExecutor::new(&manager.workspace);
3028 rt.block_on(async {
3029 manager
3030 .run_to_completion(
3031 run_id,
3032 max_workers,
3033 &mut executor,
3034 &binary.display().to_string(),
3035 None,
3036 Duration::from_millis(10),
3037 )
3038 .await
3039 .unwrap()
3040 })
3041 }
3042
3043 const RESUME_T0: &str = "2026-06-13T01:00:00Z";
3044
3045 fn role_task_with_retry(id: &str, role: &str, max_attempts: u32) -> FleetTaskSpec {
3046 let mut spec = task(id);
3047 spec.worker = Some(FleetTaskWorkerProfile {
3048 agent_profile: None,
3049 role: Some(role.to_string()),
3050 loadout: None,
3051 model: None,
3052 model_class: None,
3053 tool_profile: None,
3054 tools: Vec::new(),
3055 capabilities: Vec::new(),
3056 });
3057 spec.retry_policy = Some(FleetRetryPolicy {
3058 max_attempts,
3059 ..FleetRetryPolicy::default()
3060 });
3061 spec
3062 }
3063
3064 fn resume_worker_spec(id: &str) -> FleetWorkerSpec {
3065 FleetWorkerSpec {
3066 id: id.to_string(),
3067 name: id.to_string(),
3068 host: FleetHostSpec::Local,
3069 trust_level: Some(FleetTrustLevel::Local),
3070 labels: BTreeMap::new(),
3071 capabilities: vec!["local".to_string()],
3072 max_concurrent_tasks: Some(1),
3073 }
3074 }
3075
3076 fn resume_now(offset_secs: i64) -> DateTime<Utc> {
3077 DateTime::parse_from_rfc3339(RESUME_T0)
3078 .unwrap()
3079 .with_timezone(&Utc)
3080 + chrono::Duration::seconds(offset_secs)
3081 }
3082
3083 /// Seed the durable ledger with the state a crashed manager would leave: a
3084 /// running run whose `completed` task ids finished with receipts, and whose
3085 /// `orphaned` (task_id, worker_id) pairs are still `Leased` to workers that
3086 /// last heartbeat at `heartbeat_ts` — stale once the resume clock advances
3087 /// past `stale_after`.
3088 fn seed_crashed_run(
3089 ledger: &FleetLedger,
3090 run_id: &FleetRunId,
3091 tasks: &[FleetTaskSpec],
3092 workers: &[FleetWorkerSpec],
3093 completed: &[&str],
3094 orphaned: &[(&str, &str)],
3095 heartbeat_ts: &str,
3096 ) {
3097 ledger
3098 .create_run(&FleetRun {
3099 id: run_id.clone(),
3100 name: "resume smoke".to_string(),
3101 status: FleetRunStatus::Running,
3102 target: None,
3103 workflow: None,
3104 roles: Vec::new(),
3105 max_workers: Some(workers.len().max(1)),
3106 task_specs: tasks.to_vec(),
3107 worker_specs: workers.to_vec(),
3108 labels: BTreeMap::new(),
3109 security_policy: None,
3110 created_at: heartbeat_ts.to_string(),
3111 updated_at: Some(heartbeat_ts.to_string()),
3112 completed_at: None,
3113 })
3114 .unwrap();
3115 for spec in tasks {
3116 ledger
3117 .enqueue(FleetInboxEntry {
3118 run_id: run_id.clone(),
3119 task_id: spec.id.clone(),
3120 priority: 0,
3121 enqueued_at: heartbeat_ts.to_string(),
3122 lease_deadline: None,
3123 attempts: 0,
3124 })
3125 .unwrap();
3126 }
3127 for (idx, &task_id) in completed.iter().enumerate() {
3128 let worker_id = format!("done-worker-{idx}");
3129 ledger
3130 .lease_task(run_id, task_id, &worker_id, heartbeat_ts, None)
3131 .unwrap();
3132 ledger
3133 .mark_task_terminal_status(
3134 run_id,
3135 task_id,
3136 Some(worker_id.as_str()),
3137 heartbeat_ts,
3138 FleetTaskLedgerStatus::Completed,
3139 )
3140 .unwrap();
3141 ledger
3142 .record_receipt(FleetReceipt {
3143 run_id: run_id.clone(),
3144 task_id: task_id.to_string(),
3145 worker_id,
3146 attempt: Some(1),
3147 terminal_seq: None,
3148 completed_at: heartbeat_ts.to_string(),
3149 result: FleetTaskResult::Pass,
3150 failure_kind: None,
3151 artifacts: Vec::new(),
3152 score: None,
3153 resolved_route: None,
3154 effective_permissions: None,
3155 })
3156 .unwrap();
3157 }
3158 for &(task_id, worker_id) in orphaned {
3159 ledger
3160 .lease_task(run_id, task_id, worker_id, heartbeat_ts, None)
3161 .unwrap();
3162 ledger
3163 .heartbeat(worker_id, heartbeat_ts, None, None)
3164 .unwrap();
3165 }
3166 }
3167
3168 #[test]
3169 fn fleet_resume_reconciles_orphaned_lease_and_retries_within_budget() {
3170 let tmp = TempDir::new().unwrap();
3171 let ledger = FleetLedger::open(tmp.path()).unwrap();
3172 let run_id = FleetRunId::from("resume-run");
3173 // Three roles, three workers; scout and verifier finished, builder is
3174 // orphaned mid-flight (its worker stopped heartbeating at the crash).
3175 let tasks = vec![
3176 role_task_with_retry("scout-1", "read-only", 3),
3177 role_task_with_retry("build-1", "builder", 3),
3178 role_task_with_retry("verify-1", "smoke-runner", 3),
3179 ];
3180 let workers = vec![
3181 resume_worker_spec("w-scout"),
3182 resume_worker_spec("w-build"),
3183 resume_worker_spec("w-verify"),
3184 ];
3185 seed_crashed_run(
3186 &ledger,
3187 &run_id,
3188 &tasks,
3189 &workers,
3190 &["scout-1", "verify-1"],
3191 &[("build-1", "w-build")],
3192 RESUME_T0,
3193 );
3194
3195 // Restart: a fresh manager over the same workspace resumes from ledger.
3196 let manager = FleetManager::open(tmp.path())
3197 .unwrap()
3198 .with_stale_after(Duration::from_secs(5));
3199 let outcome = manager.resume_run_at(&run_id, resume_now(30)).unwrap();
3200
3201 assert_eq!(
3202 outcome.reclaimed_stale, 1,
3203 "orphaned builder lease detected stale"
3204 );
3205 assert_eq!(outcome.restarted, 1, "builder retried within budget");
3206 assert_eq!(outcome.failed, 0);
3207 assert_eq!(outcome.escalated, 0);
3208 assert_eq!(
3209 outcome.status.completed, 2,
3210 "pre-crash completions preserved"
3211 );
3212 assert_eq!(outcome.status.restarted, 1);
3213
3214 let state = manager.rebuild_state().unwrap();
3215 assert_eq!(state.receipts.len(), 2, "pre-crash receipts survive resume");
3216 let builder = &state.tasks["resume-run:build-1"];
3217 assert_eq!(builder.status, FleetTaskLedgerStatus::Leased);
3218 assert_eq!(builder.entry.attempts, 2, "retry leased a second attempt");
3219
3220 let text = std::fs::read_to_string(manager.ledger_path()).unwrap();
3221 assert!(
3222 text.contains("\"state\":\"stale\""),
3223 "stale event durably recorded"
3224 );
3225 assert!(
3226 text.contains("\"state\":\"restarted\""),
3227 "restart durably recorded"
3228 );
3229 }
3230
3231 #[test]
3232 fn fleet_resume_exhausted_retry_fails_and_escalates_idempotently() {
3233 let tmp = TempDir::new().unwrap();
3234 let ledger = FleetLedger::open(tmp.path()).unwrap();
3235 let run_id = FleetRunId::from("resume-run");
3236 let mut builder = role_task_with_retry("build-1", "builder", 1);
3237 builder.alert_policy = Some(FleetAlertPolicy {
3238 events: vec![FleetAlertEventClass::RestartExhausted],
3239 channels: vec![FleetAlertChannel::Slack {
3240 webhook: FleetAlertEndpoint::inline("https://hooks.slack.invalid/secret"),
3241 }],
3242 after_attempts: Some(1),
3243 after_minutes_stale: Some(1),
3244 });
3245 let tasks = vec![builder];
3246 let workers = vec![resume_worker_spec("w-build")];
3247 seed_crashed_run(
3248 &ledger,
3249 &run_id,
3250 &tasks,
3251 &workers,
3252 &[],
3253 &[("build-1", "w-build")],
3254 RESUME_T0,
3255 );
3256
3257 let manager = FleetManager::open(tmp.path())
3258 .unwrap()
3259 .with_stale_after(Duration::from_secs(5));
3260 let outcome = manager.resume_run_at(&run_id, resume_now(30)).unwrap();
3261
3262 assert_eq!(outcome.reclaimed_stale, 1);
3263 assert_eq!(outcome.restarted, 0);
3264 assert_eq!(outcome.failed, 1, "exhausted retry budget fails the task");
3265 assert_eq!(
3266 outcome.escalated, 1,
3267 "exhaustion escalates per alert policy"
3268 );
3269 assert_eq!(outcome.status.failed, 1);
3270 assert_eq!(outcome.status.escalated, 1);
3271
3272 let text = std::fs::read_to_string(manager.ledger_path()).unwrap();
3273 assert!(text.contains("\"state\":\"failed\""));
3274 assert!(text.contains("\"record\":\"alert_sent\""));
3275 assert_eq!(manager.rebuild_state().unwrap().escalated_events.len(), 1);
3276 assert!(
3277 !text.contains("hooks.slack.invalid/secret"),
3278 "secret webhook redacted in ledger"
3279 );
3280
3281 // Resuming again must not resurrect or re-escalate a terminal failure.
3282 let again = manager.resume_run_at(&run_id, resume_now(30)).unwrap();
3283 assert_eq!(again.reclaimed_stale, 0);
3284 assert_eq!(again.failed, 0);
3285 assert_eq!(again.escalated, 0);
3286 assert_eq!(
3287 manager.run_status(&run_id).unwrap().escalated,
3288 1,
3289 "no duplicate escalation across resumes"
3290 );
3291 }
3292
3293 #[test]
3294 fn fleet_resume_retry_is_idempotent_at_same_instant() {
3295 let tmp = TempDir::new().unwrap();
3296 let ledger = FleetLedger::open(tmp.path()).unwrap();
3297 let run_id = FleetRunId::from("resume-run");
3298 let tasks = vec![role_task_with_retry("build-1", "builder", 3)];
3299 let workers = vec![resume_worker_spec("w-build")];
3300 seed_crashed_run(
3301 &ledger,
3302 &run_id,
3303 &tasks,
3304 &workers,
3305 &[],
3306 &[("build-1", "w-build")],
3307 RESUME_T0,
3308 );
3309
3310 let manager = FleetManager::open(tmp.path())
3311 .unwrap()
3312 .with_stale_after(Duration::from_secs(5));
3313 let first = manager.resume_run_at(&run_id, resume_now(30)).unwrap();
3314 assert_eq!(first.restarted, 1);
3315
3316 // Re-leased at the resume instant, the task is no longer stale, so a
3317 // second resume at the same instant is a no-op (no double retry).
3318 let second = manager.resume_run_at(&run_id, resume_now(30)).unwrap();
3319 assert_eq!(second.reclaimed_stale, 0);
3320 assert_eq!(second.restarted, 0);
3321 assert_eq!(
3322 manager.rebuild_state().unwrap().tasks["resume-run:build-1"]
3323 .entry
3324 .attempts,
3325 2,
3326 "attempts did not double on the second resume"
3327 );
3328 }
3329
3330 #[test]
3331 fn fleet_resume_uses_wall_clock_for_stale_detection() {
3332 let tmp = TempDir::new().unwrap();
3333 let ledger = FleetLedger::open(tmp.path()).unwrap();
3334 let run_id = FleetRunId::from("resume-run");
3335 let tasks = vec![role_task_with_retry("build-1", "builder", 3)];
3336 let workers = vec![resume_worker_spec("w-build")];
3337 // Heartbeat an hour in the past so it is reliably stale under the real
3338 // wall clock used by the production `resume_run` entrypoint.
3339 let stale_ts = (Utc::now() - chrono::Duration::seconds(3600))
3340 .to_rfc3339_opts(SecondsFormat::Secs, true);
3341 seed_crashed_run(
3342 &ledger,
3343 &run_id,
3344 &tasks,
3345 &workers,
3346 &[],
3347 &[("build-1", "w-build")],
3348 &stale_ts,
3349 );
3350
3351 let manager = FleetManager::open(tmp.path())
3352 .unwrap()
3353 .with_stale_after(Duration::from_secs(5));
3354 let outcome = manager.resume_run(&run_id).unwrap();
3355
3356 assert_eq!(outcome.reclaimed_stale, 1);
3357 assert_eq!(outcome.restarted, 1);
3358 }
3359
3360 #[test]
3361 fn fleet_manager_creates_run_and_starts_workers_up_to_cap() {
3362 let tmp = TempDir::new().unwrap();
3363 let manager = FleetManager::open(tmp.path()).unwrap();
3364 let path = task_spec_file(&tmp, vec![task("task-a"), task("task-b"), task("task-c")]);
3365
3366 let report = manager.create_run_from_task_spec_path(&path, 2).unwrap();
3367
3368 assert_eq!(report.task_count, 3);
3369 assert_eq!(report.leased, 2);
3370 assert_eq!(report.queued, 1);
3371 assert_eq!(report.worker_ids.len(), 2);
3372 let status = manager.run_status(&report.run_id).unwrap();
3373 assert_eq!(status.queued, 1);
3374 assert_eq!(status.running, 2);
3375 assert_eq!(status.completed, 0);
3376 }
3377
3378 #[test]
3379 fn fleet_manager_rejects_unknown_agent_profile_before_run_creation() {
3380 let tmp = TempDir::new().unwrap();
3381 let manager = FleetManager::open(tmp.path()).unwrap();
3382 let mut task = task("task-a");
3383 task.worker = Some(FleetTaskWorkerProfile {
3384 role: None,
3385 agent_profile: Some("missing".to_string()),
3386 loadout: None,
3387 model_class: None,
3388 model: None,
3389 tool_profile: None,
3390 tools: Vec::new(),
3391 capabilities: Vec::new(),
3392 });
3393 let doc = FleetTaskSpecDocument {
3394 name: Some("profile guard".to_string()),
3395 labels: BTreeMap::new(),
3396 security_policy: None,
3397 workers: Vec::new(),
3398 tasks: vec![task],
3399 };
3400
3401 let err = manager
3402 .create_run(doc, 1)
3403 .expect_err("unknown agent profile must reject the run");
3404
3405 assert!(
3406 err.to_string()
3407 .contains("references unknown agent profile \"missing\"")
3408 );
3409 assert!(manager.ledger.rebuild_state().unwrap().runs.is_empty());
3410 }
3411
3412 #[test]
3413 fn fleet_manager_inspect_exposes_heartbeat_artifacts_and_errors() {
3414 let tmp = TempDir::new().unwrap();
3415 let manager = FleetManager::open(tmp.path()).unwrap();
3416 let path = task_spec_file(&tmp, vec![task("task-a")]);
3417 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3418 let worker_id = &report.worker_ids[0];
3419
3420 let inspection = manager.inspect_worker(worker_id).unwrap();
3421 assert_eq!(inspection.status, FleetWorkerStatus::Busy);
3422 assert_eq!(inspection.current_task_id.as_deref(), Some("task-a"));
3423 assert!(inspection.latest_heartbeat_at.is_some());
3424 assert_eq!(inspection.artifacts.len(), 1);
3425 assert!(inspection.last_error.is_none());
3426
3427 let inspection = manager.interrupt_worker(worker_id).unwrap();
3428 assert_eq!(inspection.status, FleetWorkerStatus::Online);
3429 assert_eq!(
3430 inspection.last_error.as_deref(),
3431 Some("cancelled by operator")
3432 );
3433 let status = manager.run_status(&report.run_id).unwrap();
3434 assert_eq!(status.cancelled, 1);
3435 }
3436
3437 #[test]
3438 fn fleet_manager_inspect_canonicalizes_advisory_role_aliases() {
3439 for alias in ["oracle", "advisor"] {
3440 let tmp = TempDir::new().unwrap();
3441 let manager = FleetManager::open(tmp.path()).unwrap();
3442 let path = task_spec_file(&tmp, vec![role_task_with_retry("advice", alias, 1)]);
3443 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3444
3445 let inspection = manager.inspect_worker(&report.worker_ids[0]).unwrap();
3446 assert_eq!(
3447 inspection.role.as_deref(),
3448 Some("consultant"),
3449 "inspection must not emit compatibility alias {alias}"
3450 );
3451 let state = manager.rebuild_state().unwrap();
3452 let persisted_role = state
3453 .runs
3454 .get(&report.run_id.0)
3455 .and_then(|run| run.task_specs[0].worker.as_ref())
3456 .and_then(|worker| worker.role.as_deref());
3457 assert_eq!(
3458 persisted_role,
3459 Some("consultant"),
3460 "new durable task must not persist compatibility alias {alias}"
3461 );
3462 }
3463 }
3464
3465 #[cfg(unix)]
3466 fn skip_if_process_table_unavailable() -> bool {
3467 !crate::fleet::host::process_table_inspection_available()
3468 }
3469
3470 #[cfg(unix)]
3471 #[test]
3472 fn separate_manager_interrupt_stops_live_worker_and_stays_terminal() {
3473 if skip_if_process_table_unavailable() {
3474 return;
3475 }
3476 let tmp = TempDir::new().unwrap();
3477 let manager = FleetManager::open(tmp.path()).unwrap();
3478 let controller = FleetManager::open(tmp.path()).unwrap();
3479 let path = task_spec_file(&tmp, vec![task("task-a"), task("task-b")]);
3480 let pid_path = tmp.path().join("live-worker.pid");
3481 let first_worker_marker = tmp.path().join("first-worker-started");
3482 let stopped_marker = tmp.path().join("first-worker-stopped");
3483 let fake = fake_codewhale(
3484 &tmp,
3485 &format!(
3486 r#"#!/bin/sh
3487 if [ -e '{first_worker_marker}' ]; then
3488 printf '{{"type":"content","content":"second task"}}\n'
3489 exit 0
3490 fi
3491 touch '{first_worker_marker}'
3492 printf '%s' "$$" > '{}'
3493 printf '{{"type":"content","content":"running"}}\n'
3494 trap 'touch "{stopped_marker}"; exit 0' INT TERM
3495 sleep 30
3496 "#,
3497 pid_path.display(),
3498 first_worker_marker = first_worker_marker.display(),
3499 stopped_marker = stopped_marker.display(),
3500 ),
3501 );
3502 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3503 let worker_id = report.worker_ids[0].clone();
3504 let mut executor = FleetExecutor::new(&manager.workspace);
3505 let rt = tokio::runtime::Runtime::new().unwrap();
3506
3507 let (status, interrupted) = rt.block_on(async {
3508 tokio::time::timeout(Duration::from_secs(15), async {
3509 tokio::join!(
3510 async {
3511 manager
3512 .run_to_completion(
3513 &report.run_id,
3514 1,
3515 &mut executor,
3516 &fake.display().to_string(),
3517 None,
3518 Duration::from_millis(10),
3519 )
3520 .await
3521 .unwrap()
3522 },
3523 async {
3524 tokio::time::timeout(Duration::from_secs(10), async {
3525 while !pid_path.is_file() {
3526 tokio::time::sleep(Duration::from_millis(5)).await;
3527 }
3528 })
3529 .await
3530 .expect("fake worker never started");
3531 controller.interrupt_worker(&worker_id).unwrap()
3532 }
3533 )
3534 })
3535 .await
3536 .expect("Fleet cancellation did not beat the worker's natural exit")
3537 });
3538
3539 assert_eq!(interrupted.status, FleetWorkerStatus::Online);
3540 assert_eq!(status.cancelled, 1);
3541 assert_eq!(status.completed, 1);
3542 assert_eq!(status.running, 0);
3543 assert!(executor.worker_ids().is_empty());
3544 assert!(
3545 stopped_marker.is_file(),
3546 "cancelled Fleet worker did not observe the stop signal"
3547 );
3548
3549 let inspection = controller.inspect_worker(&worker_id).unwrap();
3550 assert_eq!(inspection.status, FleetWorkerStatus::Online);
3551 let state = controller.rebuild_state().unwrap();
3552 let task_key = task_key(&report.run_id.0, "task-a");
3553 assert_eq!(
3554 state.tasks[&task_key].status,
3555 FleetTaskLedgerStatus::Cancelled
3556 );
3557 let event_key = event_key(&worker_id, &report.run_id.0, "task-a");
3558 assert!(matches!(
3559 &state.latest_events[&event_key].payload,
3560 FleetWorkerEventPayload::Cancelled { .. }
3561 ));
3562 }
3563
3564 #[cfg(unix)]
3565 #[test]
3566 fn live_restart_fences_old_process_and_only_attempt_two_completes() {
3567 if skip_if_process_table_unavailable() {
3568 return;
3569 }
3570 let tmp = TempDir::new().unwrap();
3571 let manager = FleetManager::open(tmp.path()).unwrap();
3572 let controller = FleetManager::open(tmp.path()).unwrap();
3573 let path = task_spec_file(&tmp, vec![task("task-a")]);
3574 let first_worker_marker = tmp.path().join("first-attempt-started");
3575 let stopped_marker = tmp.path().join("first-attempt-stopped");
3576 let fake = fake_codewhale(
3577 &tmp,
3578 &format!(
3579 r#"#!/bin/sh
3580 if [ -e '{first_worker_marker}' ]; then
3581 printf '{{"type":"content","content":"attempt two"}}\n'
3582 exit 0
3583 fi
3584 touch '{first_worker_marker}'
3585 printf '{{"type":"content","content":"attempt one still running"}}\n'
3586 trap 'touch "{stopped_marker}"; exit 0' INT TERM
3587 while :; do sleep 1; done
3588 "#,
3589 first_worker_marker = first_worker_marker.display(),
3590 stopped_marker = stopped_marker.display(),
3591 ),
3592 );
3593 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3594 let worker_id = report.worker_ids[0].clone();
3595 let mut executor = FleetExecutor::new(&manager.workspace);
3596 let rt = tokio::runtime::Runtime::new().unwrap();
3597
3598 let status = rt.block_on(async {
3599 tokio::time::timeout(Duration::from_secs(15), async {
3600 tokio::join!(
3601 async {
3602 manager
3603 .run_to_completion(
3604 &report.run_id,
3605 1,
3606 &mut executor,
3607 &fake.display().to_string(),
3608 None,
3609 Duration::from_millis(10),
3610 )
3611 .await
3612 .unwrap()
3613 },
3614 async {
3615 tokio::time::timeout(Duration::from_secs(10), async {
3616 while !first_worker_marker.is_file() {
3617 tokio::time::sleep(Duration::from_millis(5)).await;
3618 }
3619 })
3620 .await
3621 .expect("first Fleet attempt never started");
3622 controller.restart_worker(&worker_id).unwrap();
3623 }
3624 )
3625 .0
3626 })
3627 .await
3628 .expect("restarted Fleet task did not finish")
3629 });
3630
3631 assert_eq!(status.completed, 1);
3632 assert_eq!(status.running, 0);
3633 assert_eq!(status.restarted, 1);
3634 assert!(executor.worker_ids().is_empty());
3635 assert!(
3636 stopped_marker.is_file(),
3637 "the restarted attempt's old host process was not stopped"
3638 );
3639 let state = controller.rebuild_state().unwrap();
3640 let task_key = task_key(&report.run_id.0, "task-a");
3641 assert_eq!(state.tasks[&task_key].entry.attempts, 2);
3642 assert_eq!(
3643 state.tasks[&task_key].status,
3644 FleetTaskLedgerStatus::Completed
3645 );
3646 let receipt = &state.receipts[&task_key];
3647 assert_eq!(receipt.attempt, Some(2));
3648 assert!(receipt.terminal_seq.is_some());
3649 assert_eq!(receipt.result, FleetTaskResult::Partial);
3650 }
3651
3652 #[test]
3653 fn fleet_manager_restart_and_stop_all_are_ledgered() {
3654 let tmp = TempDir::new().unwrap();
3655 let manager = FleetManager::open(tmp.path()).unwrap();
3656 let path = task_spec_file(&tmp, vec![task("task-a"), task("task-b")]);
3657 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3658 let worker_id = &report.worker_ids[0];
3659
3660 manager.interrupt_worker(worker_id).unwrap();
3661 let restart = manager.restart_worker(worker_id).unwrap();
3662 assert_eq!(restart.run_id, report.run_id);
3663 assert_eq!(restart.max_workers, 1);
3664 assert_eq!(restart.inspection.status, FleetWorkerStatus::Busy);
3665 let status = manager.run_status(&report.run_id).unwrap();
3666 assert_eq!(status.running, 1);
3667 assert_eq!(status.queued, 1);
3668
3669 let stopped = manager.stop_all().unwrap();
3670 assert_eq!(stopped, 2);
3671 let status = manager.run_status(&report.run_id).unwrap();
3672 assert_eq!(status.cancelled, 2);
3673 assert_eq!(status.running, 0);
3674 }
3675
3676 #[cfg(unix)]
3677 #[test]
3678 fn standalone_restart_drives_replacement_attempt_to_terminal_receipt() {
3679 let tmp = TempDir::new().unwrap();
3680 let coordination =
3681 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3682 let manager = FleetManager::open(tmp.path())
3683 .unwrap()
3684 .with_sub_agent_manager(coordination.clone());
3685 let path = task_spec_file(&tmp, vec![task("task-a")]);
3686 let marker = tmp.path().join("replacement-attempt-ran");
3687 let fake = fake_codewhale(
3688 &tmp,
3689 &format!(
3690 r#"#!/bin/sh
3691 touch '{}'
3692 printf '{{"type":"content","content":"replacement attempt"}}\n'
3693 exit 0
3694 "#,
3695 marker.display()
3696 ),
3697 );
3698 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3699 let worker_id = &report.worker_ids[0];
3700
3701 manager.interrupt_worker(worker_id).unwrap();
3702 let restart = manager.restart_worker(worker_id).unwrap();
3703 let mut executor = FleetExecutor::new(&manager.workspace);
3704 let rt = tokio::runtime::Runtime::new().unwrap();
3705 let status = rt
3706 .block_on(async {
3707 tokio::time::timeout(
3708 Duration::from_secs(5),
3709 manager.run_to_completion(
3710 &restart.run_id,
3711 restart.max_workers,
3712 &mut executor,
3713 &fake.display().to_string(),
3714 None,
3715 Duration::from_millis(10),
3716 ),
3717 )
3718 .await
3719 })
3720 .expect("standalone restart left a ghost running task")
3721 .unwrap();
3722
3723 assert!(
3724 marker.is_file(),
3725 "replacement worker process never launched"
3726 );
3727 assert_eq!(status.completed, 1);
3728 assert_eq!(status.running, 0);
3729 assert_eq!(status.restarted, 1);
3730 let state = manager.rebuild_state().unwrap();
3731 let key = task_key(&report.run_id.0, "task-a");
3732 assert_eq!(state.tasks[&key].entry.attempts, 2);
3733 assert_eq!(state.tasks[&key].status, FleetTaskLedgerStatus::Completed);
3734 assert_eq!(state.receipts[&key].attempt, Some(2));
3735 assert!(state.receipts[&key].terminal_seq.is_some());
3736 let generation = coordination
3737 .try_read()
3738 .unwrap()
3739 .get_worker_record(worker_id)
3740 .and_then(|record| record.spec.launch_manifest)
3741 .map(|manifest| manifest.generation);
3742 assert_eq!(generation, Some(2));
3743 }
3744
3745 #[test]
3746 fn prepared_restart_survives_reload_and_commits_once() {
3747 let tmp = TempDir::new().unwrap();
3748 let coordination =
3749 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3750 let manager = FleetManager::open(tmp.path())
3751 .unwrap()
3752 .with_sub_agent_manager(coordination.clone());
3753 let path = task_spec_file(&tmp, vec![task("task-a")]);
3754 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3755 let worker_id = report.worker_ids[0].clone();
3756
3757 manager.ledger.fail_next_restart_append_after_callback();
3758 let error = manager
3759 .restart_worker(&worker_id)
3760 .expect_err("restart append failpoint must leave a durable preparation");
3761 assert!(
3762 error
3763 .to_string()
3764 .contains("forced Fleet restart ledger append failure"),
3765 "{error:#}"
3766 );
3767 assert_eq!(
3768 manager.rebuild_state().unwrap().tasks[&task_key(&report.run_id.0, "task-a")]
3769 .entry
3770 .attempts,
3771 1,
3772 "the replacement lease was not published"
3773 );
3774 let prepared_spec = coordination
3775 .try_read()
3776 .unwrap()
3777 .get_worker_record(&worker_id)
3778 .unwrap()
3779 .spec;
3780 assert_eq!(
3781 prepared_spec.launch_manifest.as_ref().unwrap().generation,
3782 2,
3783 "generation two is the durable prepare marker"
3784 );
3785
3786 drop(manager);
3787 drop(coordination);
3788
3789 let reloaded_coordination =
3790 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3791 let reloaded = FleetManager::open(tmp.path())
3792 .unwrap()
3793 .with_sub_agent_manager(reloaded_coordination.clone());
3794 reloaded
3795 .restart_worker(&worker_id)
3796 .expect("reload must idempotently consume the prepared generation");
3797
3798 let state = reloaded.rebuild_state().unwrap();
3799 assert_eq!(
3800 state.tasks[&task_key(&report.run_id.0, "task-a")]
3801 .entry
3802 .attempts,
3803 2
3804 );
3805 let committed_spec = reloaded_coordination
3806 .try_read()
3807 .unwrap()
3808 .get_worker_record(&worker_id)
3809 .unwrap()
3810 .spec;
3811 assert_eq!(committed_spec, prepared_spec, "only generation may change");
3812 let ledger_text = std::fs::read_to_string(reloaded.ledger_path()).unwrap();
3813 assert_eq!(
3814 ledger_text.matches("\"state\":\"restarted\"").count(),
3815 1,
3816 "the prepared retry commits exactly one restart event"
3817 );
3818 }
3819
3820 #[test]
3821 fn prepared_restart_with_corrupt_depth_fails_closed_after_reload() {
3822 let tmp = TempDir::new().unwrap();
3823 let coordination =
3824 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3825 let manager = FleetManager::open(tmp.path())
3826 .unwrap()
3827 .with_sub_agent_manager(coordination.clone());
3828 let path = task_spec_file(&tmp, vec![task("task-a")]);
3829 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3830 let worker_id = report.worker_ids[0].clone();
3831
3832 manager.ledger.fail_next_restart_append_after_callback();
3833 manager
3834 .restart_worker(&worker_id)
3835 .expect_err("failpoint leaves generation two prepared");
3836 {
3837 let mut guard = coordination.try_write().unwrap();
3838 let mut corrupt = guard.get_worker_record(&worker_id).unwrap().spec;
3839 corrupt.max_spawn_depth = corrupt.max_spawn_depth.saturating_add(1);
3840 guard
3841 .replace_registered_worker_spec_for_test(corrupt)
3842 .unwrap();
3843 }
3844 drop(manager);
3845 drop(coordination);
3846
3847 let reloaded_coordination =
3848 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3849 let reloaded = FleetManager::open(tmp.path())
3850 .unwrap()
3851 .with_sub_agent_manager(reloaded_coordination);
3852 let error = reloaded
3853 .restart_worker(&worker_id)
3854 .expect_err("prepared authority corruption must fail before ledger commit");
3855 assert!(
3856 error
3857 .to_string()
3858 .contains("does not match the exact task lease")
3859 );
3860 let state = reloaded.rebuild_state().unwrap();
3861 assert_eq!(
3862 state.tasks[&task_key(&report.run_id.0, "task-a")]
3863 .entry
3864 .attempts,
3865 1
3866 );
3867 assert!(state.restarted_events.is_empty());
3868 }
3869
3870 #[cfg(unix)]
3871 #[test]
3872 fn resume_stale_registered_worker_advances_generation_and_launches() {
3873 let tmp = TempDir::new().unwrap();
3874 let coordination =
3875 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3876 let manager = FleetManager::open(tmp.path())
3877 .unwrap()
3878 .with_stale_after(Duration::from_secs(5))
3879 .with_sub_agent_manager(coordination.clone());
3880 let path = task_spec_file(&tmp, vec![role_task_with_retry("task-a", "reviewer", 3)]);
3881 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3882 let worker_id = report.worker_ids[0].clone();
3883 let marker = tmp.path().join("resumed-attempt-ran");
3884 let fake = fake_codewhale(
3885 &tmp,
3886 &format!(
3887 r#"#!/bin/sh
3888 touch '{}'
3889 printf '{{"type":"content","content":"resumed attempt"}}\n'
3890 exit 0
3891 "#,
3892 marker.display()
3893 ),
3894 );
3895
3896 drop(manager);
3897 drop(coordination);
3898
3899 let reloaded_coordination =
3900 crate::tools::subagent::new_shared_subagent_manager(tmp.path().to_path_buf(), 2);
3901 let reloaded = FleetManager::open(tmp.path())
3902 .unwrap()
3903 .with_stale_after(Duration::from_secs(5))
3904 .with_sub_agent_manager(reloaded_coordination.clone());
3905 let resumed = reloaded
3906 .resume_run_at(&report.run_id, Utc::now() + chrono::Duration::minutes(10))
3907 .unwrap();
3908 assert_eq!(resumed.reclaimed_stale, 1);
3909 assert_eq!(resumed.restarted, 1);
3910 assert_eq!(
3911 reloaded.rebuild_state().unwrap().tasks[&task_key(&report.run_id.0, "task-a")]
3912 .entry
3913 .attempts,
3914 2
3915 );
3916 assert_eq!(
3917 reloaded_coordination
3918 .try_read()
3919 .unwrap()
3920 .get_worker_record(&worker_id)
3921 .unwrap()
3922 .spec
3923 .launch_manifest
3924 .unwrap()
3925 .generation,
3926 2
3927 );
3928
3929 let status = complete_with_fake_codewhale(&reloaded, &report.run_id, 1, &fake);
3930 assert!(marker.is_file(), "the recovered replacement never launched");
3931 assert_eq!(status.completed, 1);
3932 assert_eq!(status.restarted, 1);
3933 let state = reloaded.rebuild_state().unwrap();
3934 assert_eq!(
3935 state.receipts[&task_key(&report.run_id.0, "task-a")].attempt,
3936 Some(2)
3937 );
3938 }
3939
3940 #[cfg(unix)]
3941 #[test]
3942 fn concurrent_manager_loops_launch_each_attempt_once() {
3943 let tmp = TempDir::new().unwrap();
3944 let manager = FleetManager::open(tmp.path()).unwrap();
3945 let standby = FleetManager::open(tmp.path()).unwrap();
3946 let path = task_spec_file(&tmp, vec![task("task-a")]);
3947 let starts = tmp.path().join("worker-starts");
3948 let fake = fake_codewhale(
3949 &tmp,
3950 &format!(
3951 r#"#!/bin/sh
3952 printf 'started\n' >> '{}'
3953 sleep 1
3954 printf '{{"type":"content","content":"done"}}\n'
3955 exit 0
3956 "#,
3957 starts.display()
3958 ),
3959 );
3960 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
3961 let mut primary_executor = FleetExecutor::new(&manager.workspace);
3962 let mut standby_executor = FleetExecutor::new(&manager.workspace);
3963 let binary = fake.display().to_string();
3964 let rt = tokio::runtime::Runtime::new().unwrap();
3965
3966 let (primary_status, standby_status) = rt
3967 .block_on(async {
3968 tokio::time::timeout(Duration::from_secs(5), async {
3969 tokio::join!(
3970 manager.run_to_completion(
3971 &report.run_id,
3972 1,
3973 &mut primary_executor,
3974 &binary,
3975 None,
3976 Duration::from_millis(10),
3977 ),
3978 standby.run_to_completion(
3979 &report.run_id,
3980 1,
3981 &mut standby_executor,
3982 &binary,
3983 None,
3984 Duration::from_millis(10),
3985 )
3986 )
3987 })
3988 .await
3989 })
3990 .expect("competing Fleet managers did not converge");
3991
3992 assert_eq!(primary_status.unwrap().completed, 1);
3993 assert_eq!(standby_status.unwrap().completed, 1);
3994 let starts = std::fs::read_to_string(starts).unwrap();
3995 assert_eq!(
3996 starts.lines().count(),
3997 1,
3998 "competing managers launched the same leased attempt more than once"
3999 );
4000 }
4001
4002 #[cfg(unix)]
4003 #[test]
4004 fn fleet_manager_can_record_completed_local_smoke_tasks() {
4005 let tmp = TempDir::new().unwrap();
4006 let manager = FleetManager::open(tmp.path()).unwrap();
4007 let path = task_spec_file(&tmp, vec![task("task-a"), task("task-b"), task("task-c")]);
4008 let fake = fake_codewhale(
4009 &tmp,
4010 r#"#!/bin/sh
4011 printf '{"type":"tool_use","name":"read_file","id":"fake","input":{}}\n'
4012 printf '{"type":"done"}\n'
4013 exit 0
4014 "#,
4015 );
4016
4017 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4018
4019 assert_eq!(report.leased, 1);
4020 assert_eq!(report.queued, 2);
4021 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4022 assert_eq!(status.completed, 3);
4023 assert_eq!(status.running, 0);
4024 let state = manager.ledger.rebuild_state().unwrap();
4025 assert_eq!(state.receipts.len(), 3);
4026 }
4027
4028 #[test]
4029 fn fleet_task_spec_sample_launches_independent_worker_tasks() {
4030 let tmp = TempDir::new().unwrap();
4031 let manager = FleetManager::open(tmp.path()).unwrap();
4032 let path = task_spec_file(
4033 &tmp,
4034 vec![
4035 task("release-triage"),
4036 task("risk-review"),
4037 task("docs-check"),
4038 ],
4039 );
4040
4041 let report = manager.create_run_from_task_spec_path(&path, 2).unwrap();
4042
4043 assert_eq!(report.task_count, 3);
4044 assert_eq!(report.leased, 2);
4045 assert_eq!(report.queued, 1);
4046 assert_ne!(report.worker_ids[0], report.worker_ids[1]);
4047 let state = manager.ledger.rebuild_state().unwrap();
4048 assert!(
4049 state
4050 .tasks
4051 .contains_key(&format!("{}:release-triage", report.run_id.0))
4052 );
4053 assert!(
4054 state
4055 .tasks
4056 .contains_key(&format!("{}:risk-review", report.run_id.0))
4057 );
4058 assert!(
4059 state
4060 .tasks
4061 .contains_key(&format!("{}:docs-check", report.run_id.0))
4062 );
4063 }
4064
4065 #[cfg(unix)]
4066 #[test]
4067 fn fleet_task_spec_local_scorer_records_receipt_artifact() {
4068 let tmp = TempDir::new().unwrap();
4069 let manager = FleetManager::open(tmp.path()).unwrap();
4070 let mut completed = task("task-a");
4071 completed.scorer = Some(FleetScorerSpec::ExitCode);
4072 let path = task_spec_file(&tmp, vec![completed]);
4073 let fake = fake_codewhale(
4074 &tmp,
4075 r#"#!/bin/sh
4076 printf '{"type":"done"}\n'
4077 exit 0
4078 "#,
4079 );
4080
4081 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4082 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4083
4084 assert_eq!(status.completed, 1);
4085 assert_eq!(status.failed, 0);
4086 assert_eq!(status.partial, 0);
4087 let state = manager.ledger.rebuild_state().unwrap();
4088 let receipt = &state.receipts[&format!("{}:task-a", report.run_id.0)];
4089 assert_eq!(receipt.result, FleetTaskResult::Pass);
4090 assert_eq!(receipt.failure_kind, None);
4091 assert_eq!(
4092 receipt.resolved_route, None,
4093 "a headless worker with no terminal route report must fail closed"
4094 );
4095 assert!(receipt.score.as_ref().unwrap().value > 0.99);
4096 assert!(
4097 receipt
4098 .artifacts
4099 .iter()
4100 .any(|artifact| matches!(artifact.kind, FleetArtifactKind::Receipt))
4101 );
4102 }
4103
4104 #[cfg(unix)]
4105 #[test]
4106 fn fleet_task_spec_unscored_zero_exit_records_partial_receipt() {
4107 let tmp = TempDir::new().unwrap();
4108 let manager = FleetManager::open(tmp.path()).unwrap();
4109 let path = task_spec_file(&tmp, vec![task("task-a")]);
4110 let fake = fake_codewhale(
4111 &tmp,
4112 r#"#!/bin/sh
4113 printf '{"type":"done"}\n'
4114 exit 0
4115 "#,
4116 );
4117
4118 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4119 let worker_id = report.worker_ids[0].clone();
4120 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4121
4122 assert_eq!(status.completed, 1);
4123 assert_eq!(status.partial, 1);
4124 assert_eq!(status.failed, 0);
4125 let state = manager.ledger.rebuild_state().unwrap();
4126 let receipt = &state.receipts[&format!("{}:task-a", report.run_id.0)];
4127 assert_eq!(receipt.result, FleetTaskResult::Partial);
4128 assert_eq!(receipt.failure_kind, None);
4129 assert!(
4130 receipt
4131 .score
4132 .as_ref()
4133 .and_then(|score| score.notes.as_deref())
4134 .unwrap_or_default()
4135 .contains("no verifiable output")
4136 );
4137 assert!(
4138 receipt
4139 .artifacts
4140 .iter()
4141 .any(|artifact| matches!(artifact.kind, FleetArtifactKind::Receipt))
4142 );
4143 let inspection = manager.inspect_worker(&worker_id).unwrap();
4144 let summary = inspection.receipt_summary.as_deref().unwrap_or_default();
4145 assert!(summary.contains("result=partial"));
4146 assert!(summary.contains("no verifiable output"));
4147 }
4148
4149 #[cfg(unix)]
4150 #[test]
4151 fn fleet_task_spec_unscored_worker_error_records_failed_receipt() {
4152 let tmp = TempDir::new().unwrap();
4153 let manager = FleetManager::open(tmp.path()).unwrap();
4154 let path = task_spec_file(&tmp, vec![task("task-a")]);
4155 let fake = fake_codewhale(
4156 &tmp,
4157 r#"#!/bin/sh
4158 printf '{"type":"error","error":"tool failed"}\n'
4159 exit 7
4160 "#,
4161 );
4162
4163 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4164 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4165
4166 assert_eq!(status.completed, 0);
4167 assert_eq!(status.partial, 0);
4168 assert_eq!(status.failed, 1);
4169 assert_eq!(status.task_failed, 1);
4170 let state = manager.ledger.rebuild_state().unwrap();
4171 let receipt = &state.receipts[&format!("{}:task-a", report.run_id.0)];
4172 assert_eq!(receipt.result, FleetTaskResult::Fail);
4173 assert_eq!(receipt.failure_kind, Some(FleetTaskFailureKind::Task));
4174 }
4175
4176 #[cfg(unix)]
4177 #[test]
4178 fn fleet_task_spec_status_distinguishes_failure_sources() {
4179 let tmp = TempDir::new().unwrap();
4180 let manager = FleetManager::open(tmp.path()).unwrap();
4181 let mut task_failed = task("a-task-failure");
4182 task_failed.scorer = Some(FleetScorerSpec::ExitCode);
4183 task_failed.instructions = "task-failure".to_string();
4184 let mut transport = task("b-transport-failure");
4185 transport.scorer = Some(FleetScorerSpec::ExitCode);
4186 let mut verifier_failed = task("c-verifier-failure");
4187 verifier_failed.scorer = Some(FleetScorerSpec::RegexMatch {
4188 path: PathBuf::from("missing.log"),
4189 pattern: "[".to_string(),
4190 });
4191 let fake = fake_codewhale(
4192 &tmp,
4193 r#"#!/bin/sh
4194 case "$*" in
4195 *task-failure*)
4196 printf '{"type":"error","error":"task failed"}\n'
4197 exit 7
4198 ;;
4199 *)
4200 printf '{"type":"done"}\n'
4201 exit 0
4202 ;;
4203 esac
4204 "#,
4205 );
4206 let doc = FleetTaskSpecDocument {
4207 name: Some("failure source smoke".to_string()),
4208 labels: BTreeMap::new(),
4209 security_policy: None,
4210 workers: vec![
4211 default_local_worker("local-task"),
4212 FleetWorkerSpec {
4213 id: "docker-transport".to_string(),
4214 name: "Docker transport".to_string(),
4215 host: FleetHostSpec::Docker {
4216 image: "fake".to_string(),
4217 args: Vec::new(),
4218 },
4219 trust_level: Some(FleetTrustLevel::Sandbox),
4220 labels: BTreeMap::new(),
4221 capabilities: vec![],
4222 max_concurrent_tasks: Some(1),
4223 },
4224 default_local_worker("local-verifier"),
4225 ],
4226 tasks: vec![task_failed, transport, verifier_failed],
4227 };
4228
4229 let report = manager.create_run(doc, 3).unwrap();
4230 let status = complete_with_fake_codewhale(&manager, &report.run_id, 3, &fake);
4231
4232 assert_eq!(status.failed, 3);
4233 assert_eq!(status.transport_failed, 1);
4234 assert_eq!(status.task_failed, 1);
4235 assert_eq!(status.verifier_failed, 1);
4236 assert_eq!(status.running, 0);
4237 }
4238
4239 #[cfg(unix)]
4240 #[test]
4241 fn remote_terminal_route_x_wins_manager_config_y_and_receipt_is_secret_free() {
4242 let tmp = TempDir::new().unwrap();
4243 let manager_config = Config {
4244 provider: Some("manager-y".to_string()),
4245 providers: Some(crate::config::ProvidersConfig {
4246 custom: std::collections::HashMap::from([(
4247 "manager-y".to_string(),
4248 crate::config::ProviderConfig {
4249 kind: Some("openai-compatible".to_string()),
4250 base_url: Some("https://manager-y.invalid/v1".to_string()),
4251 model: Some("manager-model-y".to_string()),
4252 api_key: Some("sk-manager-y-must-not-leak".to_string()),
4253 ..Default::default()
4254 },
4255 )]),
4256 ..Default::default()
4257 }),
4258 ..Default::default()
4259 };
4260 let manager = FleetManager::open(tmp.path())
4261 .unwrap()
4262 .with_session_model("manager-model-y")
4263 .with_route_config(manager_config);
4264 let fake = fake_codewhale(
4265 &tmp,
4266 r#"#!/bin/sh
4267 printf '%s\n' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model-x","base_url":"https://remote-x.invalid/v1","api_key":"sk-remote-x-must-not-leak"}}'
4268 printf '%s\n' '{"type":"done"}'
4269 "#,
4270 );
4271 let report = manager
4272 .create_run(
4273 FleetTaskSpecDocument {
4274 name: Some("remote config drift".to_string()),
4275 labels: BTreeMap::new(),
4276 security_policy: None,
4277 workers: vec![],
4278 tasks: vec![task("route-drift")],
4279 },
4280 1,
4281 )
4282 .unwrap();
4283
4284 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4285 assert_eq!(status.completed, 1);
4286 let state = manager.ledger.rebuild_state().unwrap();
4287 let receipt = &state.receipts[&format!("{}:route-drift", report.run_id.0)];
4288 let route = receipt
4289 .resolved_route
4290 .as_ref()
4291 .expect("terminal-reported route receipt");
4292 assert_eq!(route.provider_id, "remote-x");
4293 assert_eq!(route.provider_exact_id.as_deref(), Some("remote-x"));
4294 assert_eq!(route.provider_kind, "custom");
4295 assert_eq!(route.wire_model_id, "worker-model-x");
4296 assert_eq!(route.canonical_model, None);
4297 assert_eq!(route.protocol, "unreported");
4298 assert_eq!(route.source, "worker_terminal_metadata");
4299
4300 let output = serde_json::to_string(receipt).unwrap().to_ascii_lowercase();
4301 for forbidden in [
4302 "manager-y",
4303 "manager-model-y",
4304 "base_url",
4305 "https://",
4306 "api_key",
4307 "sk-manager-y-must-not-leak",
4308 "sk-remote-x-must-not-leak",
4309 ] {
4310 assert!(
4311 !output.contains(forbidden),
4312 "receipt leaked manager drift or secret field {forbidden:?}: {output}"
4313 );
4314 }
4315 }
4316
4317 #[cfg(unix)]
4318 #[test]
4319 fn headless_terminal_with_invalid_route_metadata_does_not_fall_back_to_manager_config() {
4320 let tmp = TempDir::new().unwrap();
4321 let manager = FleetManager::open(tmp.path())
4322 .unwrap()
4323 .with_session_model("manager-model-y")
4324 .with_route_config(Config {
4325 provider: Some("manager-y".to_string()),
4326 providers: Some(crate::config::ProvidersConfig {
4327 custom: std::collections::HashMap::from([(
4328 "manager-y".to_string(),
4329 crate::config::ProviderConfig {
4330 kind: Some("openai-compatible".to_string()),
4331 base_url: Some("https://manager-y.invalid/v1".to_string()),
4332 model: Some("manager-model-y".to_string()),
4333 api_key: Some("sk-manager-y-must-not-leak".to_string()),
4334 ..Default::default()
4335 },
4336 )]),
4337 ..Default::default()
4338 }),
4339 ..Default::default()
4340 });
4341 let fake = fake_codewhale(
4342 &tmp,
4343 r#"#!/bin/sh
4344 printf '%s\n' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"deepseek","provider_id":"custom-x","model":"deepseek-v4-pro"}}'
4345 printf '%s\n' '{"type":"done"}'
4346 "#,
4347 );
4348 let report = manager
4349 .create_run(
4350 FleetTaskSpecDocument {
4351 name: Some("invalid terminal route".to_string()),
4352 labels: BTreeMap::new(),
4353 security_policy: None,
4354 workers: vec![],
4355 tasks: vec![task("invalid-route")],
4356 },
4357 1,
4358 )
4359 .unwrap();
4360
4361 let status = complete_with_fake_codewhale(&manager, &report.run_id, 1, &fake);
4362 assert_eq!(status.completed, 1);
4363 let state = manager.ledger.rebuild_state().unwrap();
4364 let receipt = &state.receipts[&format!("{}:invalid-route", report.run_id.0)];
4365 assert_eq!(receipt.resolved_route, None);
4366 let output = serde_json::to_string(receipt).unwrap().to_ascii_lowercase();
4367 assert!(!output.contains("manager-y"));
4368 assert!(!output.contains("manager-model-y"));
4369 assert!(!output.contains("sk-manager-y-must-not-leak"));
4370 }
4371
4372 #[cfg(unix)]
4373 #[test]
4374 fn fleet_smoke_runs_three_roles_ten_tasks_with_receipts_and_failure() {
4375 let tmp = TempDir::new().unwrap();
4376 let manager = FleetManager::open(tmp.path()).unwrap();
4377 let fake = fake_codewhale(
4378 &tmp,
4379 r#"#!/bin/sh
4380 case "$*" in
4381 *intentional-failure*)
4382 printf '{"type":"tool_use","name":"exec_shell","id":"fail","input":{}}\n'
4383 printf '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"deepseek","model":"deepseek-v4-pro"}}\n'
4384 printf '{"type":"error","error":"intentional failure"}\n'
4385 exit 7
4386 ;;
4387 *)
4388 printf '{"type":"tool_use","name":"read_file","id":"ok","input":{}}\n'
4389 printf '{"type":"content","delta":"ok"}\n'
4390 printf '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"deepseek","model":"deepseek-v4-pro"}}\n'
4391 printf '{"type":"done"}\n'
4392 exit 0
4393 ;;
4394 esac
4395 "#,
4396 );
4397 let smoke_task = |id: &str, role: &str, tools: Vec<&str>, marker: &str| {
4398 let mut task = task(id);
4399 task.name = format!("{role} {id}");
4400 task.objective = Some(format!("{role} smoke task {id}"));
4401 task.instructions = format!("run deterministic fleet smoke lane {marker}");
4402 task.worker = Some(FleetTaskWorkerProfile {
4403 role: Some(role.to_string()),
4404 agent_profile: None,
4405 loadout: None,
4406 model_class: None,
4407 model: None,
4408 tool_profile: Some("explicit".to_string()),
4409 tools: tools.into_iter().map(str::to_string).collect(),
4410 capabilities: vec!["local-smoke".to_string()],
4411 });
4412 if role == "builder" {
4413 task.workspace = Some(FleetWorkspaceRequirements {
4414 writable_paths: vec![PathBuf::from(".codewhale/fleet")],
4415 ..FleetWorkspaceRequirements::default()
4416 });
4417 }
4418 task.expected_artifacts = vec![FleetArtifactKind::Log, FleetArtifactKind::Receipt];
4419 task.scorer = Some(FleetScorerSpec::ExitCode);
4420 task.retry_policy = Some(FleetRetryPolicy {
4421 max_attempts: 1,
4422 ..Default::default()
4423 });
4424 task
4425 };
4426 let tasks = vec![
4427 smoke_task("scout-1", "scout", vec!["read_file", "grep_files"], "ok"),
4428 smoke_task(
4429 "builder-1",
4430 "builder",
4431 vec!["read_file", "apply_patch"],
4432 "ok",
4433 ),
4434 smoke_task(
4435 "verifier-1",
4436 "verifier",
4437 vec!["exec_shell", "read_file"],
4438 "ok",
4439 ),
4440 smoke_task("scout-2", "scout", vec!["read_file", "grep_files"], "ok"),
4441 smoke_task(
4442 "builder-2",
4443 "builder",
4444 vec!["read_file", "apply_patch"],
4445 "ok",
4446 ),
4447 smoke_task(
4448 "verifier-2",
4449 "verifier",
4450 vec!["exec_shell", "read_file"],
4451 "ok",
4452 ),
4453 smoke_task("scout-3", "scout", vec!["read_file", "grep_files"], "ok"),
4454 smoke_task(
4455 "builder-3",
4456 "builder",
4457 vec!["read_file", "apply_patch"],
4458 "ok",
4459 ),
4460 smoke_task(
4461 "verifier-3",
4462 "verifier",
4463 vec!["exec_shell", "read_file"],
4464 "ok",
4465 ),
4466 smoke_task(
4467 "verifier-4-fail",
4468 "verifier",
4469 vec!["exec_shell", "read_file"],
4470 "intentional-failure",
4471 ),
4472 ];
4473
4474 let report = manager
4475 .create_run(
4476 FleetTaskSpecDocument {
4477 name: Some("fleet route parity smoke".to_string()),
4478 labels: BTreeMap::from([("issue".to_string(), "3166".to_string())]),
4479 security_policy: Some(FleetSecurityPolicy {
4480 default_trust_level: FleetTrustLevel::Local,
4481 ..Default::default()
4482 }),
4483 workers: vec![],
4484 tasks,
4485 },
4486 3,
4487 )
4488 .unwrap();
4489
4490 assert_eq!(report.task_count, 10);
4491 assert_eq!(report.worker_ids.len(), 3);
4492 assert_eq!(report.leased, 3);
4493 assert_eq!(report.queued, 7);
4494
4495 // #3885 / item 3: baseline RSS before workers run.
4496 #[cfg(target_os = "linux")]
4497 let rss_baseline_kb = rss_kb();
4498
4499 let status = complete_with_fake_codewhale(&manager, &report.run_id, 3, &fake);
4500
4501 // #3885 / item 3: RSS after fleet completes — logged so memory
4502 // regressions produce numbers, not user reports.
4503 #[cfg(target_os = "linux")]
4504 {
4505 let rss_after_kb = rss_kb();
4506 eprintln!(
4507 "[fleet-smoke rss] baseline={} kB, after_run={} kB, delta={} kB",
4508 rss_baseline_kb.map_or("n/a".to_string(), |v| v.to_string()),
4509 rss_after_kb.map_or("n/a".to_string(), |v| v.to_string()),
4510 match (rss_baseline_kb, rss_after_kb) {
4511 (Some(b), Some(a)) => (a as i64 - b as i64).to_string(),
4512 _ => "n/a".to_string(),
4513 }
4514 );
4515 }
4516 assert_eq!(status.completed, 9);
4517 assert_eq!(status.failed, 1);
4518 assert_eq!(status.task_failed, 1);
4519 assert_eq!(status.partial, 0);
4520 assert_eq!(status.running, 0);
4521 assert_eq!(status.queued, 0);
4522
4523 let state = manager.ledger.rebuild_state().unwrap();
4524 let run = &state.runs[&report.run_id.0];
4525 let roles = run
4526 .task_specs
4527 .iter()
4528 .filter_map(|task| task.worker.as_ref()?.role.clone())
4529 .collect::<BTreeSet<_>>();
4530 assert_eq!(
4531 roles,
4532 BTreeSet::from([
4533 "builder".to_string(),
4534 "scout".to_string(),
4535 "verifier".to_string()
4536 ])
4537 );
4538 assert_eq!(state.receipts.len(), 10);
4539
4540 // #3166 scope #10: every receipt persists a resolved-route snapshot
4541 // (#3154) with non-empty provider/wire-model, a role, and the resolver
4542 // source — and the serialized receipt leaks no credential material.
4543 for (key, receipt) in &state.receipts {
4544 let route = receipt
4545 .resolved_route
4546 .as_ref()
4547 .unwrap_or_else(|| panic!("receipt {key} should carry a resolved route"));
4548 assert!(
4549 !route.provider_id.is_empty(),
4550 "receipt {key} resolved-route provider_id must be non-empty"
4551 );
4552 assert!(
4553 !route.wire_model_id.is_empty(),
4554 "receipt {key} resolved-route wire_model_id must be non-empty"
4555 );
4556 assert!(
4557 route.role.as_deref().is_some_and(|role| !role.is_empty()),
4558 "receipt {key} resolved-route should record a role"
4559 );
4560 assert_eq!(
4561 route.source, "worker_terminal_metadata",
4562 "receipt {key} resolved-route source must be the worker terminal"
4563 );
4564 assert!(
4565 route
4566 .model_route
4567 .as_deref()
4568 .is_some_and(|route| !route.is_empty()),
4569 "receipt {key} resolved-route should record the model route seam"
4570 );
4571 assert!(
4572 route
4573 .role_source
4574 .as_deref()
4575 .is_some_and(|source| !source.is_empty()),
4576 "receipt {key} resolved-route should record role source"
4577 );
4578 assert!(
4579 route
4580 .model_source
4581 .as_deref()
4582 .is_some_and(|source| !source.is_empty()),
4583 "receipt {key} resolved-route should record model source"
4584 );
4585 let permissions = receipt
4586 .effective_permissions
4587 .as_ref()
4588 .unwrap_or_else(|| panic!("receipt {key} should carry effective permissions"));
4589 assert_eq!(
4590 permissions.source, "worker_runtime_profile",
4591 "receipt {key} permissions source must be the worker runtime profile"
4592 );
4593 assert!(
4594 permissions.background,
4595 "receipt {key} should record background worker execution"
4596 );
4597 assert_eq!(
4598 permissions.tool_scope, "explicit",
4599 "receipt {key} should preserve explicit tool scope"
4600 );
4601 assert!(
4602 !permissions.tools.is_empty(),
4603 "receipt {key} should record explicit tool names"
4604 );
4605 match route.role.as_deref() {
4606 Some("builder") => {
4607 assert!(permissions.write, "builder receipt {key} should write");
4608 assert_eq!(permissions.shell, "full");
4609 }
4610 Some("scout") => {
4611 assert!(
4612 !permissions.write,
4613 "scout receipt {key} must stay read-only"
4614 );
4615 // Scout/reviewer lanes now carry the recon posture:
4616 // network reach + bounded verification surface, so the
4617 // receipt records full shell authority (raw shell still
4618 // requires write and stays denied by the clamp).
4619 assert_eq!(permissions.shell, "full");
4620 }
4621 Some("verifier") => {
4622 assert!(
4623 !permissions.write,
4624 "verifier receipt {key} must stay read-only"
4625 );
4626 assert_eq!(permissions.shell, "full");
4627 }
4628 role => panic!("unexpected receipt role for {key}: {role:?}"),
4629 }
4630
4631 let receipt_json = serde_json::to_string(receipt).unwrap();
4632 let haystack = receipt_json.to_ascii_lowercase();
4633 for needle in [
4634 "api_key",
4635 "apikey",
4636 "api-key",
4637 "authorization",
4638 "bearer ",
4639 "auth_token",
4640 "auth-token",
4641 "password",
4642 "credential",
4643 "sk-ant-",
4644 "sk-proj-",
4645 "sk-or-",
4646 "secret",
4647 ] {
4648 assert!(
4649 !haystack.contains(needle),
4650 "receipt {key} JSON must not contain secret marker {needle:?}: {receipt_json}"
4651 );
4652 }
4653 }
4654
4655 let failed_receipt = &state.receipts[&format!("{}:verifier-4-fail", report.run_id.0)];
4656 assert_eq!(failed_receipt.result, FleetTaskResult::Fail);
4657 assert_eq!(
4658 failed_receipt.failure_kind,
4659 Some(FleetTaskFailureKind::Task)
4660 );
4661 assert!(failed_receipt.artifacts.iter().any(|artifact| {
4662 matches!(artifact.kind, FleetArtifactKind::Log)
4663 && artifact.mime_type.as_deref() == Some("application/x-ndjson")
4664 && artifact.size_bytes.unwrap_or_default() > 0
4665 }));
4666 assert!(
4667 failed_receipt
4668 .artifacts
4669 .iter()
4670 .any(|artifact| matches!(artifact.kind, FleetArtifactKind::Receipt))
4671 );
4672
4673 for worker_id in &report.worker_ids {
4674 let inspection = manager.inspect_worker(worker_id).unwrap();
4675 assert_eq!(inspection.status, FleetWorkerStatus::Online);
4676 assert!(inspection.latest_heartbeat_at.is_some());
4677 assert!(
4678 inspection.receipt_summary.is_some(),
4679 "{worker_id} should expose latest receipt summary"
4680 );
4681 assert!(
4682 inspection.artifacts.iter().any(|artifact| matches!(
4683 artifact.kind,
4684 FleetArtifactKind::Log | FleetArtifactKind::Receipt
4685 )),
4686 "{worker_id} should expose artifact refs"
4687 );
4688 }
4689 }
4690
4691 #[test]
4692 fn fleet_status_counts_restarted_and_escalated_events() {
4693 let tmp = TempDir::new().unwrap();
4694 let manager = FleetManager::open(tmp.path()).unwrap();
4695 let path = task_spec_file(&tmp, vec![task("task-a")]);
4696 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4697 let worker_id = &report.worker_ids[0];
4698
4699 manager.restart_worker(worker_id).unwrap();
4700 manager
4701 .append_worker_event(
4702 &report.run_id,
4703 worker_id,
4704 "task-a",
4705 FleetWorkerEventPayload::Escalated {
4706 channel: "slack".to_string(),
4707 alert_id: None,
4708 },
4709 )
4710 .unwrap();
4711
4712 let status = manager.run_status(&report.run_id).unwrap();
4713 assert_eq!(status.restarted, 1);
4714 assert_eq!(status.escalated, 1);
4715
4716 manager.ledger.compact().unwrap();
4717 let status = manager.run_status(&report.run_id).unwrap();
4718 assert_eq!(status.restarted, 1);
4719 assert_eq!(status.escalated, 1);
4720 }
4721
4722 #[test]
4723 fn fleet_status_inspect_exposes_task_context_host_and_alert() {
4724 let tmp = TempDir::new().unwrap();
4725 let manager = FleetManager::open(tmp.path()).unwrap();
4726 let mut contextual = task("task-a");
4727 contextual.objective = Some("Review the release ledger".to_string());
4728 contextual.worker = Some(FleetTaskWorkerProfile {
4729 agent_profile: None,
4730 role: Some("reviewer".to_string()),
4731 loadout: None,
4732 model_class: None,
4733 model: None,
4734 tool_profile: Some("read-only".to_string()),
4735 tools: vec!["git".to_string()],
4736 capabilities: vec!["rust".to_string()],
4737 });
4738 let path = task_spec_file(&tmp, vec![contextual]);
4739 let report = manager.create_run_from_task_spec_path(&path, 1).unwrap();
4740 let worker_id = &report.worker_ids[0];
4741 manager
4742 .append_worker_event(
4743 &report.run_id,
4744 worker_id,
4745 "task-a",
4746 FleetWorkerEventPayload::Escalated {
4747 channel: "pagerduty".to_string(),
4748 alert_id: Some("alert-1".to_string()),
4749 },
4750 )
4751 .unwrap();
4752
4753 let inspection = manager.inspect_worker(worker_id).unwrap();
4754
4755 assert_eq!(
4756 inspection.objective.as_deref(),
4757 Some("Review the release ledger")
4758 );
4759 assert_eq!(inspection.role.as_deref(), Some("reviewer"));
4760 assert_eq!(inspection.host.as_deref(), Some("local"));
4761 assert_eq!(
4762 inspection.alert_state.as_deref(),
4763 Some("escalated via pagerduty alert_id=alert-1")
4764 );
4765 }
4766
4767 #[test]
4768 fn fleet_dogfood_smoke_run_two_local_workers_two_tasks() {
4769 let tmp = TempDir::new().unwrap();
4770 let workspace = tmp.path().join("repo");
4771 std::fs::create_dir_all(&workspace).unwrap();
4772 // Create a minimal Cargo.toml so the cargo-check task can succeed.
4773 std::fs::write(
4774 workspace.join("Cargo.toml"),
4775 "[package]\nname = \"smoke\"\nversion = \"0.1.0\"\nedition = \"2021\"\n",
4776 )
4777 .unwrap();
4778 std::fs::create_dir_all(workspace.join("src")).unwrap();
4779 std::fs::write(
4780 workspace.join("src").join("lib.rs"),
4781 "pub fn answer() -> u8 { 42 }\n",
4782 )
4783 .unwrap();
4784
4785 let tasks = vec![
4786 FleetTaskSpec {
4787 id: "check".to_string(),
4788 name: "check".to_string(),
4789 description: None,
4790 objective: Some("cargo check".to_string()),
4791 instructions: "run cargo check and report result".to_string(),
4792 worker: Some(FleetTaskWorkerProfile {
4793 agent_profile: None,
4794 role: Some("release-checker".to_string()),
4795 loadout: None,
4796 model_class: None,
4797 model: None,
4798 tool_profile: Some("read-only".to_string()),
4799 tools: vec!["cargo".to_string()],
4800 capabilities: vec!["rust".to_string()],
4801 }),
4802 workspace: Some(FleetWorkspaceRequirements {
4803 root: None,
4804 required_files: vec![PathBuf::from("Cargo.toml")],
4805 writable_paths: vec![PathBuf::from(".codewhale/fleet")],
4806 environment: Some(FleetEnvironmentRequirements {
4807 required: vec!["PATH".to_string()],
4808 allowlist: vec![],
4809 }),
4810 }),
4811 input_files: vec![],
4812 context: vec![],
4813 budget: None,
4814 tags: vec!["smoke".to_string()],
4815 expected_artifacts: vec![FleetArtifactKind::Log, FleetArtifactKind::Receipt],
4816 scorer: Some(FleetScorerSpec::ExitCode),
4817 retry_policy: Some(FleetRetryPolicy {
4818 max_attempts: 1,
4819 ..Default::default()
4820 }),
4821 alert_policy: None,
4822 timeout_seconds: Some(60),
4823 metadata: BTreeMap::new(),
4824 },
4825 FleetTaskSpec {
4826 id: "review".to_string(),
4827 name: "review".to_string(),
4828 description: None,
4829 objective: Some("review source".to_string()),
4830 instructions: "read src/lib.rs and report findings".to_string(),
4831 worker: Some(FleetTaskWorkerProfile {
4832 agent_profile: None,
4833 role: Some("reviewer".to_string()),
4834 loadout: None,
4835 model_class: None,
4836 model: None,
4837 tool_profile: Some("read-only".to_string()),
4838 tools: vec!["cargo".to_string()],
4839 capabilities: vec!["rust".to_string()],
4840 }),
4841 workspace: Some(FleetWorkspaceRequirements {
4842 root: None,
4843 required_files: vec![],
4844 writable_paths: vec![],
4845 environment: Some(FleetEnvironmentRequirements {
4846 required: vec!["PATH".to_string()],
4847 allowlist: vec![],
4848 }),
4849 }),
4850 input_files: vec![],
4851 context: vec![],
4852 budget: None,
4853 tags: vec!["smoke".to_string()],
4854 expected_artifacts: vec![FleetArtifactKind::Log, FleetArtifactKind::Receipt],
4855 scorer: None,
4856 retry_policy: Some(FleetRetryPolicy {
4857 max_attempts: 1,
4858 ..Default::default()
4859 }),
4860 alert_policy: None,
4861 timeout_seconds: Some(60),
4862 metadata: BTreeMap::new(),
4863 },
4864 ];
4865
4866 let manager = FleetManager::open(&workspace).unwrap();
4867 let report = manager
4868 .create_run(
4869 FleetTaskSpecDocument {
4870 name: Some("dogfood smoke".to_string()),
4871 labels: BTreeMap::new(),
4872 security_policy: Some(FleetSecurityPolicy {
4873 default_trust_level: FleetTrustLevel::Local,
4874 ..Default::default()
4875 }),
4876 workers: vec![],
4877 tasks,
4878 },
4879 2,
4880 )
4881 .unwrap();
4882
4883 assert_eq!(report.task_count, 2);
4884 assert!(!report.worker_ids.is_empty());
4885 assert_eq!(report.worker_ids.len(), 2);
4886 // After immediate scheduling, tasks may already be leased,
4887 // so queued+running should total 2.
4888 let status = manager.run_status(&report.run_id).unwrap();
4889 assert_eq!(status.queued + status.running, 2);
4890 }
4891
4892 #[test]
4893 fn fleet_security_policy_propagates_from_task_spec_document_to_run() {
4894 let tmp = TempDir::new().unwrap();
4895 let manager = FleetManager::open(tmp.path()).unwrap();
4896 // Rewrite the spec file with a security_policy block.
4897 let doc = serde_json::json!({
4898 "name": "secure smoke",
4899 "tasks": [{
4900 "id": "task-a",
4901 "name": "task-a",
4902 "instructions": "report ok",
4903 "worker": {"role": "reviewer", "tool_profile": "read-only"},
4904 "expected_artifacts": ["log"]
4905 }],
4906 "security_policy": {
4907 "default_trust_level": "local",
4908 "allowed_secrets": [{"key": "GH_TOKEN", "source": "env"}],
4909 "max_trust_level": "remote_verified",
4910 "require_identity_verification": true
4911 }
4912 });
4913 let spec_path = tmp.path().join("secure-tasks.json");
4914 std::fs::write(&spec_path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
4915
4916 let report = manager
4917 .create_run_from_task_spec_path(&spec_path, 1)
4918 .unwrap();
4919
4920 let state = manager.ledger.rebuild_state().unwrap();
4921 let run = state.runs.get(&report.run_id.0).unwrap();
4922 let policy = run.security_policy.as_ref().unwrap();
4923 assert_eq!(policy.default_trust_level, FleetTrustLevel::Local);
4924 assert_eq!(policy.allowed_secrets.len(), 1);
4925 assert_eq!(policy.allowed_secrets[0].key, "GH_TOKEN");
4926 assert_eq!(policy.max_trust_level, FleetTrustLevel::RemoteVerified);
4927 assert!(policy.require_identity_verification);
4928 }
4929 }
4930
4930 lines RUST