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