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