返回 CodeWhale
executor.rs
根目录 / crates / tui / src / fleet / executor.rs
1 //! Fleet executor — runs a Fleet worker as a real `codewhale exec` subprocess.
2 //!
3 //! A Fleet worker IS a headless `codewhale exec` run. There is no separate
4 //! "Fleet worker" execution engine: the sub-agent runtime, full tool surface,
5 //! and recursion depth all come from the one `codewhale exec` runtime, so
6 //! Fleets and sub-agents are one substrate (not two moving targets).
7 //!
8 //! This module is the bridge:
9 //! - [`build_worker_exec_command`] turns a `FleetTaskSpec` + `FleetExecConfig`
10 //! into the `codewhale [route flags] exec --output-format stream-json …`
11 //! argv that a host adapter ([`super::host`]) launches locally or over SSH.
12 //! - [`map_exec_stream_line`] maps one stream-json line emitted by that worker
13 //! into a [`FleetWorkerEventPayload`] for the durable ledger, so the ledger
14 //! persists the worker's own event vocabulary instead of a simulated one.
15 //! - [`classify_worker_exit`] turns the process exit into a terminal event.
16 //!
17 //! The TUI/CLI/Runtime API observe the ledger's compact event stream — they
18 //! never render a child session, which is what keeps the orchestrator light at
19 //! high fanout.
20
21 #![allow(dead_code)]
22
23 use anyhow::Result;
24 use codewhale_config::FleetExecConfig;
25 use codewhale_protocol::fleet::{FleetHostSpec, FleetTaskSpec, FleetWorkerEventPayload};
26
27 use super::host::{FleetHostAdapter, FleetWorkerCommand};
28 use super::profile::AgentProfile;
29 use super::task_spec::FleetWorkerFinalAnswer;
30 use super::worker_runtime::{
31 fleet_task_prompt, fleet_task_prompt_with_profiles, fleet_worker_launch_reasoning_effort,
32 fleet_worker_launch_route,
33 };
34 use crate::tools::spec::{
35 ToolAuthorityEnvelope, ToolMutationAuthority, ToolShellAuthority, ToolVerificationAuthority,
36 };
37 use crate::tools::subagent::{AgentWorkerSpec, FleetRole};
38
39 #[derive(Clone, Copy, Default)]
40 struct WorkerExecRoute<'a> {
41 model: Option<&'a str>,
42 provider: Option<&'a str>,
43 reasoning_effort: Option<&'a str>,
44 }
45
46 #[derive(Clone, Copy, Default)]
47 struct WorkerExecLimits {
48 max_turns: Option<u32>,
49 max_tool_calls: Option<u32>,
50 }
51
52 /// Resolve the executable used for Fleet worker subprocesses.
53 ///
54 /// Kept here so every long-lived surface (CLI and Runtime API) launches the
55 /// same configured worker binary instead of silently diverging.
56 pub fn configured_codewhale_binary() -> String {
57 std::env::var("CODEWHALE_FLEET_CODEWHALE_BINARY")
58 .ok()
59 .map(|value| value.trim().to_string())
60 .filter(|value| !value.is_empty())
61 .unwrap_or_else(|| "codewhale".to_string())
62 }
63
64 /// Build the `codewhale exec` argv that runs a fleet task headlessly.
65 ///
66 /// `--auto` is always passed: a headless worker has no human to approve tool
67 /// calls, so it runs with full (policy-gated) tool access. `--output-format
68 /// stream-json` makes the worker emit the NDJSON event stream this module
69 /// parses. A worker launched with the v0.9.1 machine-readable outer authority
70 /// cap is a truthful leaf (`max_spawn_depth = 0`): the nested-agent surface is
71 /// disabled until authority scopes can be intersected across
72 /// process/workspace boundaries.
73 ///
74 /// Secrets are NEVER placed on the argv: provider credentials are resolved by
75 /// the worker process from its own config/keyring exactly like an interactive
76 /// run. The host adapter additionally refuses secret-bearing env keys. The
77 /// `--provider` flag threaded by [`build_worker_exec_command_with_profiles`] is
78 /// a non-secret provider *identifier* only (#4093) — the worker still resolves
79 /// that provider's credentials from its own env/config, so this invariant
80 /// holds.
81 pub fn build_worker_exec_command(
82 codewhale_binary: &str,
83 task_spec: &FleetTaskSpec,
84 exec_config: &FleetExecConfig,
85 model: Option<&str>,
86 ) -> FleetWorkerCommand {
87 let max_turns = effective_task_max_turns(task_spec, exec_config);
88 let max_tool_calls = task_max_tool_calls(task_spec);
89 build_worker_exec_command_from_prompt(
90 codewhale_binary,
91 fleet_task_prompt(task_spec),
92 exec_config,
93 WorkerExecRoute {
94 model,
95 ..WorkerExecRoute::default()
96 },
97 None,
98 WorkerExecLimits {
99 max_turns,
100 max_tool_calls,
101 },
102 )
103 }
104
105 /// Build a worker command after resolving workspace Fleet profile input.
106 ///
107 /// The launched subprocess runs on the worker's RESOLVED route, not blindly on
108 /// the run-level session model (#4093 AC #4): the per-worker model+provider are
109 /// resolved from the task's agent profile via the same explicit-only path the
110 /// receipt uses ([`fleet_worker_launch_route`]). A worker whose profile pins
111 /// provider B thus launches on provider B's model even when the parent session
112 /// is on provider A. Workers with no profile-bound provider fall back to the
113 /// run-level model and emit no `--provider`, so the worker keeps its own
114 /// session default (today's behavior, unchanged).
115 pub fn build_worker_exec_command_with_profiles(
116 codewhale_binary: &str,
117 task_spec: &FleetTaskSpec,
118 exec_config: &FleetExecConfig,
119 model: Option<&str>,
120 agent_profiles: &[AgentProfile],
121 ) -> Result<FleetWorkerCommand> {
122 let (worker_model, worker_provider) =
123 fleet_worker_launch_route(task_spec, agent_profiles, model.unwrap_or_default());
124 let worker_reasoning_effort = fleet_worker_launch_reasoning_effort(task_spec, agent_profiles);
125 let max_turns = effective_task_max_turns(task_spec, exec_config);
126 let max_tool_calls = task_max_tool_calls(task_spec);
127 Ok(build_worker_exec_command_from_prompt(
128 codewhale_binary,
129 fleet_task_prompt_with_profiles(task_spec, agent_profiles)?,
130 exec_config,
131 WorkerExecRoute {
132 model: Some(worker_model.as_str()),
133 provider: worker_provider.as_deref(),
134 reasoning_effort: worker_reasoning_effort.as_deref(),
135 },
136 None,
137 WorkerExecLimits {
138 max_turns,
139 max_tool_calls,
140 },
141 ))
142 }
143
144 /// Build the exact Fleet subprocess command from the coordination-registered
145 /// worker spec. Unlike the compatibility helpers above, production dispatch
146 /// uses the projected objective and carries a machine-readable outer authority
147 /// envelope into the child process.
148 pub fn build_worker_exec_command_with_launch_spec(
149 codewhale_binary: &str,
150 task_spec: &FleetTaskSpec,
151 launch_spec: &AgentWorkerSpec,
152 exec_config: &FleetExecConfig,
153 model: Option<&str>,
154 agent_profiles: &[AgentProfile],
155 ) -> Result<FleetWorkerCommand> {
156 let (worker_model, worker_provider) =
157 fleet_worker_launch_route(task_spec, agent_profiles, model.unwrap_or_default());
158 let worker_reasoning_effort = fleet_worker_launch_reasoning_effort(task_spec, agent_profiles);
159 let authority = authority_envelope_for_worker(launch_spec, task_spec)?;
160 // Production dispatch receives an already-hardened launch spec from the
161 // Fleet manager. Its positive max_steps value has therefore already been
162 // intersected with a positive FleetExecConfig.max_turns; zero remains the
163 // explicit unbounded sentinel.
164 let max_turns = (launch_spec.max_steps > 0).then_some(launch_spec.max_steps);
165 let max_tool_calls = task_max_tool_calls(task_spec);
166 Ok(build_worker_exec_command_from_prompt(
167 codewhale_binary,
168 launch_spec.objective.clone(),
169 exec_config,
170 WorkerExecRoute {
171 model: Some(worker_model.as_str()),
172 provider: worker_provider.as_deref(),
173 reasoning_effort: worker_reasoning_effort.as_deref(),
174 },
175 Some(&authority),
176 WorkerExecLimits {
177 max_turns,
178 max_tool_calls,
179 },
180 ))
181 }
182
183 fn task_max_steps(task_spec: &FleetTaskSpec) -> Option<u32> {
184 task_spec
185 .budget
186 .as_ref()
187 .and_then(|budget| budget.max_steps)
188 .filter(|max_steps| *max_steps > 0)
189 }
190
191 fn task_max_tool_calls(task_spec: &FleetTaskSpec) -> Option<u32> {
192 task_spec
193 .budget
194 .as_ref()
195 .and_then(|budget| budget.max_tool_calls)
196 .filter(|max_tool_calls| *max_tool_calls > 0)
197 }
198
199 fn effective_task_max_turns(
200 task_spec: &FleetTaskSpec,
201 exec_config: &FleetExecConfig,
202 ) -> Option<u32> {
203 match (task_max_steps(task_spec), exec_config.max_turns) {
204 (Some(task_max_steps), fleet_max_turns) if fleet_max_turns > 0 => {
205 Some(task_max_steps.min(fleet_max_turns))
206 }
207 (Some(task_max_steps), _) => Some(task_max_steps),
208 (None, fleet_max_turns) if fleet_max_turns > 0 => Some(fleet_max_turns),
209 (None, _) => None,
210 }
211 }
212
213 pub(crate) fn authority_envelope_for_worker(
214 spec: &AgentWorkerSpec,
215 task_spec: &FleetTaskSpec,
216 ) -> Result<ToolAuthorityEnvelope> {
217 let (authority, writable_roots, writable_files, coordination_contracts) =
218 if spec.runtime_profile.permissions.write {
219 let manifest = spec.launch_manifest.as_ref().ok_or_else(|| {
220 anyhow::anyhow!(
221 "write-capable Fleet worker '{}' has no launch manifest",
222 spec.worker_id
223 )
224 })?;
225 (
226 ToolMutationAuthority::ScopedWrite,
227 super::worker_runtime::fleet_runtime_write_roots(task_spec)?,
228 manifest.writable_files.clone(),
229 manifest.coordination_contracts.clone(),
230 )
231 } else {
232 (
233 ToolMutationAuthority::ReadOnly,
234 Vec::new(),
235 Vec::new(),
236 Vec::new(),
237 )
238 };
239 let shell = if authority == ToolMutationAuthority::ReadOnly
240 && matches!(
241 &spec.agent_type,
242 FleetRole::Scout | FleetRole::Reviewer | FleetRole::Planner
243 )
244 && spec.runtime_profile.shell.allows_shell()
245 {
246 // Scout, Reviewer, and Planner get one classifier-bounded
247 // foreground shell. Verifiers keep the dedicated Run surface,
248 // consultants remain shell-less, and writers retain the historical
249 // subprocess contract.
250 ToolShellAuthority::ReadOnly
251 } else {
252 ToolShellAuthority::None
253 };
254 let verification = if authority == ToolMutationAuthority::ReadOnly
255 && matches!(&spec.agent_type, FleetRole::Verifier)
256 && spec.runtime_profile.shell.allows_shell()
257 {
258 ToolVerificationAuthority::Bounded
259 } else {
260 ToolVerificationAuthority::None
261 };
262 ToolAuthorityEnvelope {
263 schema_version: 1,
264 owner: spec.worker_id.clone(),
265 authority,
266 network_access: Some(spec.runtime_profile.permissions.network),
267 shell,
268 verification,
269 writable_roots,
270 writable_files,
271 coordination_contracts,
272 }
273 .normalized()
274 .map_err(anyhow::Error::msg)
275 }
276
277 fn build_worker_exec_command_from_prompt(
278 codewhale_binary: &str,
279 task_prompt: String,
280 exec_config: &FleetExecConfig,
281 route: WorkerExecRoute<'_>,
282 authority: Option<&ToolAuthorityEnvelope>,
283 limits: WorkerExecLimits,
284 ) -> FleetWorkerCommand {
285 let mut args: Vec<String> = Vec::new();
286
287 // The canonical `codewhale` dispatcher owns these route overrides as
288 // global flags and deliberately rejects them after `exec`. Keep them in
289 // front of the subcommand so Fleet commands work through the installed
290 // dispatcher as well as when a host points directly at `codewhale-tui`.
291 if let Some(model) = route.model.map(str::trim).filter(|m| !m.is_empty()) {
292 args.push("--model".to_string());
293 args.push(model.to_string());
294 }
295
296 // Non-secret provider identifier only (#4093): the worker resolves the
297 // provider's credentials from its own env/config. Emitted ONLY when the
298 // worker's profile explicitly pins a provider, so profile-less workers keep
299 // their own session default exactly as before.
300 if let Some(provider) = route.provider.map(str::trim).filter(|p| !p.is_empty()) {
301 args.push("--provider".to_string());
302 args.push(provider.to_string());
303 }
304
305 args.extend([
306 "exec".to_string(),
307 "--auto".to_string(),
308 "--output-format".to_string(),
309 "stream-json".to_string(),
310 // R7: the worker shuts itself down when the manager dies (stdin EOF).
311 "--parent-death-watch".to_string(),
312 ]);
313
314 // Non-secret thinking tier only (#4137). This is profile metadata and
315 // follows the same explicit-only policy as provider: omit it when the
316 // worker profile inherits the session/default reasoning setting.
317 if let Some(reasoning_effort) = route
318 .reasoning_effort
319 .map(str::trim)
320 .filter(|e| !e.is_empty())
321 {
322 args.push("--reasoning-effort".to_string());
323 args.push(reasoning_effort.to_string());
324 }
325
326 if !exec_config.allowed_tools.is_empty() {
327 args.push("--allowed-tools".to_string());
328 args.push(exec_config.allowed_tools.join(","));
329 }
330 if !exec_config.disallowed_tools.is_empty() {
331 args.push("--disallowed-tools".to_string());
332 args.push(exec_config.disallowed_tools.join(","));
333 }
334 if let Some(max_turns) = limits.max_turns {
335 args.push("--max-turns".to_string());
336 args.push(max_turns.to_string());
337 }
338 if let Some(max_tool_calls) = limits.max_tool_calls {
339 args.push("--max-tool-calls".to_string());
340 args.push(max_tool_calls.to_string());
341 }
342 if !exec_config.append_system_prompt.trim().is_empty() {
343 // One `--flag=value` argument, so the policy is always the value: as a
344 // separate argument, text that reads like one of exec's own options
345 // (`--hooks`) is parsed as that option and the worker fails to start.
346 args.push(format!(
347 "--append-system-prompt={}",
348 exec_config.append_system_prompt
349 ));
350 }
351
352 if let Some(authority) = authority {
353 args.push("--tool-authority-json".to_string());
354 args.push(
355 serde_json::to_string(authority)
356 .expect("validated Fleet tool authority envelope must serialize"),
357 );
358 }
359
360 // The composed task prompt is the final positional argument.
361 args.push(task_prompt);
362
363 FleetWorkerCommand::new(codewhale_binary.to_string(), args)
364 }
365
366 /// Map one `codewhale exec` stream-json line into a fleet ledger event.
367 ///
368 /// Returns `None` for lines that don't correspond to a worker lifecycle
369 /// transition (e.g. `session_capture`, `metadata`). The exec event schema is
370 /// `{"type": "...", ...}` (see `ExecStreamEvent` in `main.rs`).
371 pub fn map_exec_stream_line(line: &str) -> Option<FleetWorkerEventPayload> {
372 let value: serde_json::Value = serde_json::from_str(line.trim()).ok()?;
373 map_exec_stream_value(&value)
374 }
375
376 /// [`map_exec_stream_line`] on an already-parsed line, so the incremental
377 /// stream reader parses each frame exactly once.
378 fn map_exec_stream_value(value: &serde_json::Value) -> Option<FleetWorkerEventPayload> {
379 match value.get("type").and_then(serde_json::Value::as_str)? {
380 "tool_use" => {
381 let tool = value
382 .get("name")
383 .and_then(serde_json::Value::as_str)
384 .unwrap_or("tool")
385 .to_string();
386 let call_id = value
387 .get("id")
388 .and_then(serde_json::Value::as_str)
389 .map(str::to_string);
390 Some(FleetWorkerEventPayload::RunningTool { tool, call_id })
391 }
392 "workflow_event" => Some(FleetWorkerEventPayload::WorkflowEvent {
393 workflow_run_id: value.get("run_id")?.as_str()?.to_string(),
394 event: value.get("event")?.clone(),
395 }),
396 // Streaming model output / tool results mean the worker is alive and
397 // making progress; surface a coarse Running heartbeat.
398 "content" | "tool_result" => Some(FleetWorkerEventPayload::Running),
399 // Per-step usage receipts feed the run-level accumulator behind the
400 // fleet usage ceiling (R6, #5567); a malformed line still counts as
401 // liveness.
402 "turn_usage" => {
403 let tokens = |field: &str| {
404 value
405 .get(field)
406 .and_then(serde_json::Value::as_u64)
407 .unwrap_or(0)
408 };
409 let input_tokens = tokens("input_tokens");
410 let output_tokens = tokens("output_tokens");
411 if input_tokens == 0 && output_tokens == 0 {
412 Some(FleetWorkerEventPayload::Running)
413 } else {
414 Some(FleetWorkerEventPayload::UsageReport {
415 input_tokens,
416 output_tokens,
417 })
418 }
419 }
420 "done" => Some(FleetWorkerEventPayload::Completed {
421 exit_code: Some(0),
422 summary: None,
423 }),
424 "error" => {
425 let reason = value
426 .get("error")
427 .and_then(serde_json::Value::as_str)
428 .unwrap_or("worker reported an error")
429 .to_string();
430 Some(FleetWorkerEventPayload::Failed {
431 reason,
432 recoverable: false,
433 })
434 }
435 _ => None,
436 }
437 }
438
439 #[derive(Debug)]
440 enum ParsedTerminalRoute {
441 NotTerminal,
442 Valid(FleetWorkerReportedRoute),
443 Invalid,
444 }
445
446 /// The `meta` object of a terminal exec receipt, or `None` for every other
447 /// stream line.
448 fn exec_terminal_meta(
449 value: &serde_json::Value,
450 ) -> Option<&serde_json::Map<String, serde_json::Value>> {
451 if value.get("type").and_then(serde_json::Value::as_str) != Some("metadata") {
452 return None;
453 }
454 let meta = value.get("meta").and_then(serde_json::Value::as_object)?;
455 (meta.get("receipt_kind").and_then(serde_json::Value::as_str) == Some("terminal"))
456 .then_some(meta)
457 }
458
459 /// The worker's visible final answer from a terminal exec receipt: the
460 /// emitter already bounded and redacted `visible_final_answer_excerpt`, and
461 /// `visible_final_answer_chars` is the real pre-truncation length.
462 fn parse_exec_terminal_final_answer(value: &serde_json::Value) -> Option<FleetWorkerFinalAnswer> {
463 let meta = exec_terminal_meta(value)?;
464 let excerpt = meta
465 .get("visible_final_answer_excerpt")
466 .and_then(serde_json::Value::as_str)?
467 .trim();
468 if excerpt.is_empty() {
469 return None;
470 }
471 let chars = meta
472 .get("visible_final_answer_chars")
473 .and_then(serde_json::Value::as_u64)
474 .and_then(|chars| usize::try_from(chars).ok())
475 .unwrap_or_else(|| excerpt.chars().count());
476 Some(FleetWorkerFinalAnswer {
477 excerpt: crate::exec_stream_final_answer_excerpt(excerpt),
478 chars,
479 })
480 }
481
482 /// Parse one allowlisted, secret-free route identity from terminal exec
483 /// metadata. Once a line declares itself as a terminal receipt, malformed
484 /// route fields are distinct from ordinary non-terminal stream noise so a
485 /// prior valid record cannot survive contradictory evidence.
486 fn parse_exec_terminal_route(value: &serde_json::Value) -> ParsedTerminalRoute {
487 let Some(meta) = exec_terminal_meta(value) else {
488 return ParsedTerminalRoute::NotTerminal;
489 };
490
491 let route = (|| {
492 let provider = meta.get("provider")?.as_str()?.trim();
493 let model = meta.get("model")?.as_str()?.trim();
494 if provider.is_empty() || model.is_empty() {
495 return None;
496 }
497 let provider_kind = crate::config::ProviderKind::parse(provider)?;
498 let provider_exact_id = match meta.get("provider_id") {
499 None => None,
500 Some(value) => {
501 let id = value.as_str()?.trim();
502 if id.is_empty() {
503 return None;
504 }
505 Some(id.to_string())
506 }
507 };
508 if provider_exact_id.is_some() && provider_kind != crate::config::ProviderKind::Custom {
509 return None;
510 }
511 Some(FleetWorkerReportedRoute {
512 provider: provider.to_string(),
513 provider_exact_id,
514 model: model.to_string(),
515 })
516 })();
517
518 route.map_or(ParsedTerminalRoute::Invalid, ParsedTerminalRoute::Valid)
519 }
520
521 #[cfg(test)]
522 fn map_exec_terminal_route(line: &str) -> Option<FleetWorkerReportedRoute> {
523 let value: serde_json::Value = serde_json::from_str(line).ok()?;
524 match parse_exec_terminal_route(&value) {
525 ParsedTerminalRoute::Valid(route) => Some(route),
526 ParsedTerminalRoute::NotTerminal | ParsedTerminalRoute::Invalid => None,
527 }
528 }
529
530 /// Classify a worker process exit into a terminal fleet event.
531 ///
532 /// `stopped` means the operator stopped the worker (cancellation), which takes
533 /// precedence over the exit code.
534 pub fn classify_worker_exit(exit_code: Option<i32>, stopped: bool) -> FleetWorkerEventPayload {
535 if stopped {
536 return FleetWorkerEventPayload::Cancelled { cancelled_by: None };
537 }
538 match exit_code {
539 Some(0) => FleetWorkerEventPayload::Completed {
540 exit_code: Some(0),
541 summary: None,
542 },
543 Some(code) => FleetWorkerEventPayload::Failed {
544 reason: format!("worker exited with code {code}"),
545 recoverable: true,
546 },
547 None => FleetWorkerEventPayload::Failed {
548 reason: "worker exited without a status code".to_string(),
549 recoverable: true,
550 },
551 }
552 }
553
554 /// Drives fleet workers as real `codewhale exec` subprocesses on the local
555 /// host, incrementally draining each worker's stream-json output into fleet
556 /// ledger events.
557 ///
558 /// The caller (the `codewhale fleet run` loop / `FleetManager`) owns the
559 /// ledger; the executor owns the OS process boundary and the incremental log
560 /// parse. Because the worker is a separate process, its heavy runtime/tool
561 /// construction never touches the orchestrator — the parent only ingests a
562 /// compact event stream, which is what keeps it light at high fanout.
563 pub struct FleetExecutor {
564 workspace: std::path::PathBuf,
565 sessions_dir: Option<std::path::PathBuf>,
566 adapter: super::host::LocalProcessFleetHostAdapter,
567 ssh_adapters: std::collections::BTreeMap<String, super::host::SshFleetHostAdapter>,
568 streams: std::collections::BTreeMap<String, WorkerStream>,
569 }
570
571 /// Durable lease identity owned by one concrete host process.
572 #[derive(Debug, Clone, PartialEq, Eq)]
573 pub struct FleetExecutorAttempt {
574 pub run_id: codewhale_protocol::fleet::FleetRunId,
575 pub task_id: String,
576 pub attempt: u32,
577 }
578
579 struct WorkerStream {
580 log_path: std::path::PathBuf,
581 host: WorkerStreamHost,
582 attempt: Option<FleetExecutorAttempt>,
583 offset: u64,
584 // Keep incomplete stream frames as bytes. Decoding each read separately
585 // corrupts valid UTF-8 when a multibyte code point crosses a read boundary.
586 pending: Vec<u8>,
587 terminal: bool,
588 terminal_route: TerminalRouteEvidence,
589 /// When this worker process was started, for per-task wall-clock limits (R5).
590 started_at: std::time::Instant,
591 /// The worker's visible final answer from its terminal exec receipt. This
592 /// is the task's deliverable for report/summary work that produces no
593 /// file artifact; surfaced as `Completed.summary` and in the receipt note
594 /// so receipts stop reporting "no verifiable output" for a worker that
595 /// wrote a full report. Bounded by the emitter, so nothing accumulates
596 /// here.
597 final_answer: Option<FleetWorkerFinalAnswer>,
598 /// Saved exec session id reported by the worker's `session_capture` event.
599 /// Resolving it via `GET /v1/sessions/{id}` yields the full transcript
600 /// (the worker's final assistant reply).
601 saved_session_id: Option<String>,
602 /// Parent-owned capture identity in the Runtime's existing session store.
603 session_capture: Option<(String, std::path::PathBuf)>,
604 }
605
606 impl WorkerStream {
607 /// Observe one raw stream-json frame: record route evidence, the final
608 /// answer, and the saved-session id, and map it to a ledger payload. The
609 /// frame is parsed exactly once.
610 fn observe_line(&mut self, line: &[u8]) -> Option<FleetWorkerEventPayload> {
611 let Ok(line) = std::str::from_utf8(line) else {
612 // stream-json is a UTF-8 contract. Never accept a lossy-decoded route
613 // receipt: replacement characters could turn corrupt provider/model
614 // bytes into apparently valid provenance.
615 self.terminal_route.observe(ParsedTerminalRoute::Invalid);
616 return None;
617 };
618 let value: serde_json::Value = serde_json::from_str(line.trim()).ok()?;
619 self.terminal_route
620 .observe(parse_exec_terminal_route(&value));
621 if exec_terminal_meta(&value).is_some() {
622 self.final_answer = parse_exec_terminal_final_answer(&value);
623 }
624 if value.get("type").and_then(serde_json::Value::as_str) == Some("session_capture")
625 && let Some(id) = value
626 .get("saved_session_id")
627 .and_then(serde_json::Value::as_str)
628 .map(str::trim)
629 .filter(|id| !id.is_empty())
630 {
631 if self
632 .session_capture
633 .as_ref()
634 .is_some_and(|(expected, _)| expected == id)
635 {
636 self.saved_session_id = Some(id.to_string());
637 } else {
638 // A worker log cannot select an unrelated local transcript.
639 self.session_capture = None;
640 self.saved_session_id = None;
641 }
642 }
643 map_exec_stream_value(&value)
644 }
645 }
646
647 #[derive(Debug, Clone, Default)]
648 enum TerminalRouteEvidence {
649 #[default]
650 Missing,
651 Valid(FleetWorkerReportedRoute),
652 InvalidOrAmbiguous,
653 }
654
655 impl TerminalRouteEvidence {
656 fn observe(&mut self, parsed: ParsedTerminalRoute) {
657 match parsed {
658 ParsedTerminalRoute::NotTerminal => {}
659 ParsedTerminalRoute::Invalid => *self = Self::InvalidOrAmbiguous,
660 ParsedTerminalRoute::Valid(route) => {
661 *self = if matches!(&*self, Self::Missing) {
662 Self::Valid(route)
663 } else {
664 // The stream contract emits exactly one terminal receipt.
665 // Any second record, even an identical one, is ambiguous
666 // provenance and must permanently fail closed.
667 Self::InvalidOrAmbiguous
668 };
669 }
670 }
671 }
672
673 fn reported_route(&self) -> Option<&FleetWorkerReportedRoute> {
674 match self {
675 Self::Valid(route) => Some(route),
676 Self::Missing | Self::InvalidOrAmbiguous => None,
677 }
678 }
679 }
680
681 enum WorkerStreamHost {
682 Local,
683 Ssh(String),
684 }
685
686 #[derive(Debug, Clone)]
687 pub struct FleetWorkerReportedRoute {
688 pub provider: String,
689 pub provider_exact_id: Option<String>,
690 pub model: String,
691 }
692
693 #[derive(Debug, Clone)]
694 pub struct FleetWorkerTerminalEvent {
695 pub payload: FleetWorkerEventPayload,
696 pub exit_code: Option<i32>,
697 /// Non-terminal payloads discovered by the mandatory post-exit drain.
698 pub tail_payloads: Vec<FleetWorkerEventPayload>,
699 pub reported_route: Option<FleetWorkerReportedRoute>,
700 /// The worker's visible final answer from its terminal exec receipt,
701 /// whatever the outcome: a worker that fails after writing most of a
702 /// report keeps the text on its receipt.
703 pub final_answer: Option<FleetWorkerFinalAnswer>,
704 /// Saved exec session id reported by the worker's `session_capture` event,
705 /// when one was persisted on completion.
706 pub saved_session_id: Option<String>,
707 /// A real headless exec process must report its actual route. Callers use
708 /// this bit to distinguish a missing/invalid report (fail closed) from
709 /// pre-launch or simulated paths that only have declared route intent.
710 pub requires_reported_route: bool,
711 }
712
713 impl FleetExecutor {
714 pub fn new(workspace: impl AsRef<std::path::Path>) -> Self {
715 let workspace = workspace.as_ref().to_path_buf();
716 Self {
717 adapter: super::host::LocalProcessFleetHostAdapter::new(&workspace),
718 workspace,
719 sessions_dir: None,
720 ssh_adapters: std::collections::BTreeMap::new(),
721 streams: std::collections::BTreeMap::new(),
722 }
723 }
724
725 /// Share the caller's SessionManager directory; never create a second store.
726 pub fn with_sessions_dir(mut self, sessions_dir: std::path::PathBuf) -> Self {
727 self.sessions_dir = Some(sessions_dir);
728 self
729 }
730
731 /// Start a worker process and begin tracking its event stream.
732 pub fn start_worker(
733 &mut self,
734 worker_id: &str,
735 command: FleetWorkerCommand,
736 cwd: Option<std::path::PathBuf>,
737 ) -> super::host::FleetHostResult<super::host::FleetWorkerHandle> {
738 self.start_worker_on_host(worker_id, &FleetHostSpec::Local, command, cwd)
739 }
740
741 /// Start a worker on the requested fleet host.
742 pub fn start_worker_on_host(
743 &mut self,
744 worker_id: &str,
745 host: &FleetHostSpec,
746 command: FleetWorkerCommand,
747 cwd: Option<std::path::PathBuf>,
748 ) -> super::host::FleetHostResult<super::host::FleetWorkerHandle> {
749 self.start_worker_on_host_inner(worker_id, host, command, cwd, None)
750 }
751
752 /// Start the concrete process for one exact durable Fleet lease.
753 pub fn start_worker_attempt_on_host(
754 &mut self,
755 worker_id: &str,
756 host: &FleetHostSpec,
757 command: FleetWorkerCommand,
758 cwd: Option<std::path::PathBuf>,
759 attempt: FleetExecutorAttempt,
760 ) -> super::host::FleetHostResult<super::host::FleetWorkerHandle> {
761 self.start_worker_on_host_inner(worker_id, host, command, cwd, Some(attempt))
762 }
763
764 fn start_worker_on_host_inner(
765 &mut self,
766 worker_id: &str,
767 host: &FleetHostSpec,
768 command: FleetWorkerCommand,
769 cwd: Option<std::path::PathBuf>,
770 attempt: Option<FleetExecutorAttempt>,
771 ) -> super::host::FleetHostResult<super::host::FleetWorkerHandle> {
772 let mut request = super::host::FleetWorkerStartRequest::new(worker_id, command);
773 request.cwd = cwd;
774 let session_capture = if matches!(host, FleetHostSpec::Local) {
775 self.sessions_dir.as_ref().map(|dir| {
776 let id = uuid::Uuid::new_v4().to_string();
777 for (key, value) in [
778 ("CODEWHALE_FLEET_CAPTURE_ID", id.clone()),
779 (
780 "CODEWHALE_FLEET_CAPTURE_DIR",
781 dir.to_string_lossy().into_owned(),
782 ),
783 ] {
784 request.env.insert(key.to_string(), value);
785 request.env_allowlist.insert(key.to_string());
786 }
787 // Preserve an explicit local account/config root as well.
788 if let Some(home) = std::env::var_os("CODEWHALE_HOME") {
789 request
790 .env
791 .insert("CODEWHALE_HOME".into(), home.to_string_lossy().into_owned());
792 request.env_allowlist.insert("CODEWHALE_HOME".into());
793 }
794 (id, dir.clone())
795 })
796 } else {
797 // SSH transcripts live remotely and have no local retrieval link.
798 None
799 };
800 let (handle, host) = match host {
801 FleetHostSpec::Local => {
802 let handle = self.adapter.start_worker(request)?;
803 (handle, WorkerStreamHost::Local)
804 }
805 FleetHostSpec::Ssh { .. } => {
806 let config = super::host::SshFleetHostConfig::from_host_spec(host)?;
807 let key = worker_id.to_string();
808 let adapter = self.ssh_adapters.entry(key.clone()).or_insert(
809 super::host::SshFleetHostAdapter::new(&self.workspace, config)?,
810 );
811 let handle = adapter.start_worker(request)?;
812 (handle, WorkerStreamHost::Ssh(key))
813 }
814 FleetHostSpec::Docker { image, .. } => {
815 return Err(super::host::FleetHostError {
816 kind: super::host::FleetHostErrorKind::Configuration,
817 message: format!("docker Fleet workers are not wired yet (image {image})"),
818 });
819 }
820 };
821 self.streams.insert(
822 worker_id.to_string(),
823 WorkerStream {
824 log_path: handle.log_path.clone(),
825 host,
826 attempt,
827 offset: 0,
828 pending: Vec::new(),
829 terminal: false,
830 terminal_route: TerminalRouteEvidence::default(),
831 started_at: std::time::Instant::now(),
832 final_answer: None,
833 saved_session_id: None,
834 session_capture,
835 },
836 );
837 Ok(handle)
838 }
839
840 pub fn is_tracking(&self, worker_id: &str) -> bool {
841 self.streams.contains_key(worker_id)
842 }
843
844 pub fn worker_ids(&self) -> Vec<String> {
845 self.streams.keys().cloned().collect()
846 }
847
848 pub fn tracked_attempt(&self, worker_id: &str) -> Option<FleetExecutorAttempt> {
849 self.streams
850 .get(worker_id)
851 .and_then(|stream| stream.attempt.clone())
852 }
853
854 /// Wall-clock time this worker process has been running (R5). `None` when
855 /// the worker is not tracked or its exit was already observed: a process
856 /// that has exited is not running, so a deadline must not rewrite the
857 /// outcome that exit produced as a timeout.
858 #[must_use]
859 pub fn worker_running_for(&self, worker_id: &str) -> Option<std::time::Duration> {
860 self.streams
861 .get(worker_id)
862 .filter(|stream| !stream.terminal)
863 .map(|stream| stream.started_at.elapsed())
864 }
865
866 /// Stop a tracked worker at the host boundary.
867 ///
868 /// Operator controls run in a separate process from the foreground Fleet
869 /// manager, so they communicate cancellation through the durable ledger.
870 /// The manager calls this method after observing that terminal state; the
871 /// executor is the only owner that can reliably reach the live local/SSH
872 /// adapter handle.
873 pub fn stop_worker(&mut self, worker_id: &str) -> Result<()> {
874 let ssh_key = match self.streams.get(worker_id).map(|stream| &stream.host) {
875 Some(WorkerStreamHost::Local) => None,
876 Some(WorkerStreamHost::Ssh(key)) => Some(key.clone()),
877 None => return Ok(()),
878 };
879 if let Some(key) = ssh_key {
880 let adapter = self.ssh_adapters.get_mut(&key).ok_or_else(|| {
881 anyhow::anyhow!("tracked SSH Fleet worker {worker_id} has no host adapter")
882 })?;
883 adapter.stop_worker(worker_id)?;
884 } else {
885 self.adapter.stop_worker(worker_id)?;
886 }
887 Ok(())
888 }
889
890 /// Stop tracking a terminal worker so the scheduler can reuse the same
891 /// logical worker id for the next queued task.
892 pub fn forget_worker(&mut self, worker_id: &str) {
893 let Some(stream) = self.streams.remove(worker_id) else {
894 return;
895 };
896 match stream.host {
897 WorkerStreamHost::Local => {
898 let _ = self.adapter.cleanup_worker(worker_id);
899 }
900 WorkerStreamHost::Ssh(key) => {
901 if let Some(adapter) = self.ssh_adapters.get_mut(&key) {
902 let _ = adapter.cleanup_worker(worker_id);
903 }
904 self.ssh_adapters.remove(&key);
905 }
906 }
907 }
908
909 /// Maximum bytes drained from a worker log per call, and the cap for the
910 /// buffered partial line. Event lines are small; the remainder stays for
911 /// the next drain.
912 const MAX_DRAIN_BYTES: u64 = 1024 * 1024;
913
914 /// Read any newly-written stream-json lines for a worker and map them to
915 /// fleet ledger events. Safe to call repeatedly; only new bytes are parsed,
916 /// and a trailing partial line is buffered until its newline arrives.
917 pub fn drain_events(&mut self, worker_id: &str) -> Vec<FleetWorkerEventPayload> {
918 let Some(stream) = self.streams.get_mut(worker_id) else {
919 return Vec::new();
920 };
921 let mut events = Vec::new();
922 let Ok(mut file) = std::fs::File::open(&stream.log_path) else {
923 return events;
924 };
925 use std::io::{Read, Seek, SeekFrom};
926 if file.seek(SeekFrom::Start(stream.offset)).is_err() {
927 return events;
928 }
929 let mut buf = Vec::new();
930 if let Ok(read) = file.take(Self::MAX_DRAIN_BYTES).read_to_end(&mut buf) {
931 stream.offset += read as u64;
932 stream.pending.extend_from_slice(&buf);
933 while let Some(idx) = stream.pending.iter().position(|byte| *byte == b'\n') {
934 let line: Vec<u8> = stream.pending.drain(..=idx).collect();
935 if let Some(event) = stream.observe_line(&line) {
936 events.push(event);
937 }
938 }
939 // Whatever remains has no newline; drop it rather than buffering
940 // a newline-free flood forever.
941 if stream.pending.len() as u64 > Self::MAX_DRAIN_BYTES {
942 tracing::debug!(
943 worker_id,
944 dropped_bytes = stream.pending.len(),
945 "fleet drain dropped a newline-free flood exceeding the per-call budget"
946 );
947 stream.pending.clear();
948 }
949 }
950 events
951 }
952
953 /// Poll the worker process; once it exits, return the terminal event exactly
954 /// once. Returns `None` while the worker is still running or already
955 /// finalized.
956 pub fn poll_terminal(&mut self, worker_id: &str) -> Option<FleetWorkerEventPayload> {
957 self.poll_terminal_with_status(worker_id)
958 .map(|event| event.payload)
959 }
960
961 /// Poll the worker process and include the raw exit code for receipt
962 /// verification.
963 pub fn poll_terminal_with_status(
964 &mut self,
965 worker_id: &str,
966 ) -> Option<FleetWorkerTerminalEvent> {
967 if self.streams.get(worker_id).is_none_or(|s| s.terminal) {
968 return None;
969 }
970 let status = match self.streams.get(worker_id).map(|s| &s.host)? {
971 WorkerStreamHost::Local => self.adapter.read_status(worker_id).ok()?,
972 WorkerStreamHost::Ssh(key) => self
973 .ssh_adapters
974 .get_mut(key)
975 .and_then(|adapter| adapter.read_status(worker_id).ok())?,
976 };
977 let mut terminal = match status.state {
978 super::host::FleetHostWorkerState::Running
979 | super::host::FleetHostWorkerState::Draining
980 | super::host::FleetHostWorkerState::Unknown => return None,
981 super::host::FleetHostWorkerState::Stopped => {
982 classify_worker_exit(status.exit_code, true)
983 }
984 super::host::FleetHostWorkerState::Exited
985 | super::host::FleetHostWorkerState::Failed => {
986 classify_worker_exit(status.exit_code, false)
987 }
988 };
989 // Once status is terminal the worker can no longer append. Drain one
990 // final time before snapshotting route evidence so metadata written
991 // between the scheduler's ordinary drain and this status poll cannot
992 // be lost when the worker is forgotten.
993 let mut tail_payloads = self.drain_events(worker_id);
994 let stream = self.streams.get_mut(worker_id)?;
995 let trailing_line = std::mem::take(&mut stream.pending);
996 if trailing_line.iter().any(|byte| !byte.is_ascii_whitespace())
997 && let Some(payload) = stream.observe_line(&trailing_line)
998 {
999 tail_payloads.push(payload);
1000 }
1001 stream.terminal = true;
1002 let final_answer = stream.final_answer.take();
1003 // Surface the visible final answer on a successful completion so
1004 // report/summary tasks (no scorer, no file artifact) show their
1005 // deliverable instead of "no verifiable output". `Failed` has no
1006 // summary slot; the receipt keeps the text via `final_answer`.
1007 if let (Some(answer), FleetWorkerEventPayload::Completed { summary, .. }) =
1008 (final_answer.as_ref(), &mut terminal)
1009 {
1010 *summary = Some(answer.excerpt.clone());
1011 }
1012 Some(FleetWorkerTerminalEvent {
1013 payload: terminal,
1014 exit_code: status.exit_code,
1015 tail_payloads,
1016 reported_route: stream.terminal_route.reported_route().cloned(),
1017 final_answer,
1018 saved_session_id: stream.saved_session_id.take().filter(|id| {
1019 stream
1020 .session_capture
1021 .as_ref()
1022 .is_some_and(|(expected, dir)| {
1023 id == expected
1024 && crate::session_manager::SessionManager::new(dir.clone())
1025 .and_then(|manager| manager.load_session(id))
1026 .is_ok()
1027 })
1028 }),
1029 requires_reported_route: true,
1030 })
1031 }
1032
1033 /// True once every started worker has reached a terminal state.
1034 pub fn all_terminal(&self) -> bool {
1035 !self.streams.is_empty() && self.streams.values().all(|s| s.terminal)
1036 }
1037 }
1038
1039 #[cfg(test)]
1040 mod tests {
1041 use super::*;
1042 use codewhale_config::{
1043 FleetDelegationHints, FleetLoadout, FleetProfile, FleetProfilePermissions, FleetRole,
1044 FleetSlot,
1045 };
1046 use codewhale_protocol::fleet::{
1047 FleetHostSpec, FleetTaskBudget, FleetTaskSpec, FleetTaskWorkerProfile, FleetWorkerSpec,
1048 FleetWorkspaceRequirements,
1049 };
1050 use std::collections::BTreeMap;
1051 use tempfile::TempDir;
1052
1053 fn task(instructions: &str) -> FleetTaskSpec {
1054 FleetTaskSpec {
1055 id: "t1".to_string(),
1056 name: "Smoke".to_string(),
1057 description: None,
1058 objective: Some("prove it runs".to_string()),
1059 instructions: instructions.to_string(),
1060 worker: Some(FleetTaskWorkerProfile {
1061 agent_profile: None,
1062 role: Some("reviewer".to_string()),
1063 loadout: None,
1064 model_class: None,
1065 model: None,
1066 tool_profile: Some("read-only".to_string()),
1067 tools: vec![],
1068 capabilities: vec![],
1069 }),
1070 workspace: None,
1071 input_files: vec![],
1072 context: vec![],
1073 budget: None,
1074 tags: vec![],
1075 expected_artifacts: vec![],
1076 scorer: None,
1077 retry_policy: None,
1078 alert_policy: None,
1079 timeout_seconds: None,
1080 metadata: BTreeMap::new(),
1081 }
1082 }
1083
1084 fn agent_profile(id: &str, role: &str, instructions: &str) -> AgentProfile {
1085 AgentProfile {
1086 native_preset: None,
1087 id: id.to_string(),
1088 display_name: Some(format!("{role} profile")),
1089 description: Some(format!("{role} description")),
1090 requires: Vec::new(),
1091 profile: FleetProfile {
1092 slot: FleetSlot::from_name(role),
1093 role: FleetRole {
1094 name: role.to_string(),
1095 description: None,
1096 instructions: Some(instructions.to_string()),
1097 },
1098 loadout: FleetLoadout::Inherit,
1099 model: None,
1100 provider: None,
1101 reasoning_effort: None,
1102 permissions: FleetProfilePermissions::default(),
1103 delegation: FleetDelegationHints::default(),
1104 },
1105 source: std::path::PathBuf::from(format!("{id}.toml")),
1106 origin: crate::fleet::roster::ProfileOrigin::Workspace,
1107 plugin_authority: None,
1108 }
1109 }
1110
1111 fn launch_spec(task: &FleetTaskSpec, workspace: &std::path::Path) -> AgentWorkerSpec {
1112 let worker = FleetWorkerSpec {
1113 id: "worker-1".to_string(),
1114 name: "Worker 1".to_string(),
1115 host: FleetHostSpec::Local,
1116 trust_level: None,
1117 labels: BTreeMap::new(),
1118 capabilities: Vec::new(),
1119 max_concurrent_tasks: Some(1),
1120 };
1121 crate::fleet::worker_runtime::fleet_task_to_worker_spec_with_profiles(
1122 "worker-1",
1123 "run-1",
1124 task,
1125 &worker,
1126 "auto",
1127 workspace,
1128 workspace,
1129 &[],
1130 None,
1131 )
1132 .unwrap()
1133 }
1134
1135 fn track_test_stream(
1136 executor: &mut FleetExecutor,
1137 worker_id: &str,
1138 log_path: std::path::PathBuf,
1139 ) {
1140 executor.streams.insert(
1141 worker_id.to_string(),
1142 WorkerStream {
1143 log_path,
1144 host: WorkerStreamHost::Local,
1145 attempt: None,
1146 offset: 0,
1147 pending: Vec::new(),
1148 terminal: false,
1149 terminal_route: TerminalRouteEvidence::default(),
1150 started_at: std::time::Instant::now(),
1151 final_answer: None,
1152 saved_session_id: None,
1153 session_capture: None,
1154 },
1155 );
1156 }
1157
1158 fn append_test_stream(path: &std::path::Path, bytes: &[u8]) {
1159 use std::io::Write as _;
1160
1161 std::fs::OpenOptions::new()
1162 .append(true)
1163 .open(path)
1164 .unwrap()
1165 .write_all(bytes)
1166 .unwrap();
1167 }
1168
1169 #[test]
1170 fn worker_command_is_a_headless_codewhale_exec_run() {
1171 let exec = FleetExecConfig::default();
1172 let cmd = build_worker_exec_command("codewhale", &task("read the file"), &exec, None);
1173 assert_eq!(cmd.program, "codewhale");
1174 assert_eq!(cmd.args[0], "exec");
1175 assert!(cmd.args.contains(&"--auto".to_string()));
1176 // stream-json so the executor can ingest the worker's event stream.
1177 let joined = cmd.args.join(" ");
1178 assert!(joined.contains("--output-format stream-json"));
1179 // The task instructions ride in the positional prompt (last arg).
1180 assert!(cmd.args.last().unwrap().contains("read the file"));
1181 }
1182
1183 #[test]
1184 fn worker_command_threads_exec_hardening_flags() {
1185 let exec = FleetExecConfig {
1186 allowed_tools: vec!["read_file".to_string(), "grep_files".to_string()],
1187 disallowed_tools: vec!["exec_shell".to_string()],
1188 max_turns: 40,
1189 append_system_prompt: "never push to main".to_string(),
1190 ..FleetExecConfig::default()
1191 };
1192 let cmd = build_worker_exec_command("codewhale", &task("audit"), &exec, Some("glm-5.1"));
1193 let exec_idx = cmd
1194 .args
1195 .iter()
1196 .position(|arg| arg == "exec")
1197 .expect("worker command must contain exec");
1198 let model_idx = cmd
1199 .args
1200 .iter()
1201 .position(|arg| arg == "--model")
1202 .expect("worker command must contain --model");
1203 assert!(
1204 model_idx < exec_idx,
1205 "global --model must precede exec: {:?}",
1206 cmd.args
1207 );
1208 let joined = cmd.args.join(" ");
1209 assert!(joined.contains("--model glm-5.1"));
1210 assert!(joined.contains("--allowed-tools read_file,grep_files"));
1211 assert!(joined.contains("--disallowed-tools exec_shell"));
1212 assert!(joined.contains("--max-turns 40"));
1213 assert!(
1214 cmd.args
1215 .iter()
1216 .any(|a| a == "--append-system-prompt=never push to main")
1217 );
1218 }
1219
1220 #[test]
1221 fn worker_command_threads_positive_task_budgets_and_caps_steps() {
1222 let mut task = task("audit");
1223 task.budget = Some(FleetTaskBudget {
1224 max_steps: Some(75),
1225 max_tool_calls: Some(11),
1226 ..FleetTaskBudget::default()
1227 });
1228 let exec = FleetExecConfig {
1229 max_turns: 40,
1230 ..FleetExecConfig::default()
1231 };
1232
1233 let cmd = build_worker_exec_command("codewhale", &task, &exec, None);
1234 let max_turns_idx = cmd
1235 .args
1236 .iter()
1237 .position(|arg| arg == "--max-turns")
1238 .expect("positive max_steps must reach exec");
1239 let max_tool_calls_idx = cmd
1240 .args
1241 .iter()
1242 .position(|arg| arg == "--max-tool-calls")
1243 .expect("positive max_tool_calls must reach exec");
1244
1245 assert_eq!(cmd.args[max_turns_idx + 1], "40");
1246 assert_eq!(cmd.args[max_tool_calls_idx + 1], "11");
1247 }
1248
1249 #[test]
1250 fn production_worker_command_uses_hardened_launch_spec_steps() {
1251 let tmp = TempDir::new().unwrap();
1252 let mut task = task("audit");
1253 task.budget = Some(FleetTaskBudget {
1254 max_steps: Some(75),
1255 max_tool_calls: Some(11),
1256 ..FleetTaskBudget::default()
1257 });
1258 let mut launch_spec = launch_spec(&task, tmp.path());
1259 // The manager owns this hardening step. A different value here proves
1260 // production argv comes from the registered spec, not a second budget
1261 // projection from the task document.
1262 launch_spec.max_steps = 13;
1263
1264 let cmd = build_worker_exec_command_with_launch_spec(
1265 "codewhale",
1266 &task,
1267 &launch_spec,
1268 &FleetExecConfig {
1269 max_turns: 40,
1270 ..FleetExecConfig::default()
1271 },
1272 None,
1273 &[],
1274 )
1275 .unwrap();
1276 let max_turns_idx = cmd
1277 .args
1278 .iter()
1279 .position(|arg| arg == "--max-turns")
1280 .expect("hardened launch max_steps must reach exec");
1281 let max_tool_calls_idx = cmd
1282 .args
1283 .iter()
1284 .position(|arg| arg == "--max-tool-calls")
1285 .expect("task max_tool_calls must reach exec");
1286
1287 assert_eq!(cmd.args[max_turns_idx + 1], "13");
1288 assert_eq!(cmd.args[max_tool_calls_idx + 1], "11");
1289 }
1290
1291 #[test]
1292 fn worker_command_threads_agent_profile_prompt() {
1293 let mut task = task("audit");
1294 task.worker.as_mut().unwrap().agent_profile = Some("reviewer".to_string());
1295 let cmd = build_worker_exec_command_with_profiles(
1296 "codewhale",
1297 &task,
1298 &FleetExecConfig::default(),
1299 None,
1300 &[agent_profile(
1301 "reviewer",
1302 "reviewer",
1303 "Focus on defects, regressions, and missing tests.",
1304 )],
1305 )
1306 .unwrap();
1307 let prompt = cmd.args.last().unwrap();
1308
1309 assert!(prompt.contains("Fleet profile: reviewer"));
1310 assert!(prompt.contains("Focus on defects, regressions, and missing tests."));
1311 }
1312
1313 #[test]
1314 fn consultant_launch_spec_carries_read_only_authority_with_network_reads() {
1315 let tmp = TempDir::new().unwrap();
1316 let mut task = task("advise on the release candidate");
1317 task.worker.as_mut().unwrap().role = Some("consultant".to_string());
1318 let launch_spec = launch_spec(&task, tmp.path());
1319 assert_eq!(
1320 launch_spec.agent_type,
1321 crate::tools::subagent::FleetRole::Consultant
1322 );
1323 // Counsel only: never writes the workspace, but may read the web.
1324 assert!(!launch_spec.runtime_profile.permissions.write);
1325 assert!(launch_spec.runtime_profile.permissions.network);
1326
1327 let cmd = build_worker_exec_command_with_launch_spec(
1328 "codewhale",
1329 &task,
1330 &launch_spec,
1331 &FleetExecConfig::default(),
1332 None,
1333 &[],
1334 )
1335 .unwrap();
1336
1337 assert_eq!(cmd.args.last(), Some(&launch_spec.objective));
1338 let authority_index = cmd
1339 .args
1340 .iter()
1341 .position(|arg| arg == "--tool-authority-json")
1342 .expect("launch command must carry machine-readable authority");
1343 let authority = ToolAuthorityEnvelope::from_json(&cmd.args[authority_index + 1]).unwrap();
1344 assert_eq!(authority.owner, "worker-1");
1345 assert_eq!(authority.authority, ToolMutationAuthority::ReadOnly);
1346 assert_eq!(authority.network_access, Some(true));
1347 assert_eq!(authority.shell, ToolShellAuthority::None);
1348 assert_eq!(authority.verification, ToolVerificationAuthority::None);
1349 assert!(authority.writable_roots.is_empty());
1350 assert!(authority.writable_files.is_empty());
1351 assert!(authority.coordination_contracts.is_empty());
1352 }
1353
1354 #[test]
1355 fn scout_reviewer_and_planner_launches_carry_only_read_only_shell_authority() {
1356 let tmp = TempDir::new().unwrap();
1357 for role in ["scout", "reviewer", "planner"] {
1358 let mut task = task("inspect repository and GitHub state");
1359 task.worker.as_mut().unwrap().role = Some(role.to_string());
1360 let launch_spec = launch_spec(&task, tmp.path());
1361 let cmd = build_worker_exec_command_with_launch_spec(
1362 "codewhale",
1363 &task,
1364 &launch_spec,
1365 &FleetExecConfig::default(),
1366 None,
1367 &[],
1368 )
1369 .unwrap();
1370 let authority_index = cmd
1371 .args
1372 .iter()
1373 .position(|arg| arg == "--tool-authority-json")
1374 .expect("launch command must carry machine-readable authority");
1375 let authority =
1376 ToolAuthorityEnvelope::from_json(&cmd.args[authority_index + 1]).unwrap();
1377
1378 assert_eq!(authority.authority, ToolMutationAuthority::ReadOnly);
1379 assert_eq!(authority.network_access, Some(true));
1380 assert_eq!(authority.shell, ToolShellAuthority::ReadOnly);
1381 assert_eq!(authority.verification, ToolVerificationAuthority::None);
1382 }
1383 }
1384
1385 #[test]
1386 fn verifier_launch_carries_only_bounded_verification_process_authority() {
1387 let tmp = TempDir::new().unwrap();
1388 let mut task = task("run the focused release checks");
1389 task.worker.as_mut().unwrap().role = Some("verifier".to_string());
1390 let mut launch_spec = launch_spec(&task, tmp.path());
1391 let authority = authority_envelope_for_worker(&launch_spec, &task).unwrap();
1392
1393 assert_eq!(authority.authority, ToolMutationAuthority::ReadOnly);
1394 assert_eq!(authority.shell, ToolShellAuthority::None);
1395 assert_eq!(authority.verification, ToolVerificationAuthority::Bounded);
1396
1397 launch_spec.runtime_profile.shell = crate::worker_profile::ShellPolicy::None;
1398 let shellless = authority_envelope_for_worker(&launch_spec, &task).unwrap();
1399 assert_eq!(
1400 shellless.verification,
1401 ToolVerificationAuthority::None,
1402 "the parent shell ceiling also removes bounded verification process authority"
1403 );
1404 }
1405
1406 #[test]
1407 fn launch_spec_command_preserves_exact_write_scope() {
1408 let tmp = TempDir::new().unwrap();
1409 let mut task = task("edit the bounded source tree");
1410 let worker = task.worker.as_mut().unwrap();
1411 worker.role = Some("implementer".to_string());
1412 worker.tool_profile = None;
1413 task.workspace = Some(FleetWorkspaceRequirements {
1414 writable_paths: vec![std::path::PathBuf::from("src")],
1415 ..FleetWorkspaceRequirements::default()
1416 });
1417 let launch_spec = launch_spec(&task, tmp.path());
1418
1419 let cmd = build_worker_exec_command_with_launch_spec(
1420 "codewhale",
1421 &task,
1422 &launch_spec,
1423 &FleetExecConfig::default(),
1424 None,
1425 &[],
1426 )
1427 .unwrap();
1428 let authority_index = cmd
1429 .args
1430 .iter()
1431 .position(|arg| arg == "--tool-authority-json")
1432 .expect("launch command must carry machine-readable authority");
1433 let authority = ToolAuthorityEnvelope::from_json(&cmd.args[authority_index + 1]).unwrap();
1434
1435 assert_eq!(authority.authority, ToolMutationAuthority::ScopedWrite);
1436 assert_eq!(authority.network_access, Some(true));
1437 assert_eq!(authority.shell, ToolShellAuthority::None);
1438 assert_eq!(authority.verification, ToolVerificationAuthority::None);
1439 assert_eq!(authority.writable_roots, ["src"]);
1440 assert!(authority.writable_files.is_empty());
1441 assert!(authority.coordination_contracts.is_empty());
1442 assert_eq!(cmd.args.last(), Some(&launch_spec.objective));
1443 }
1444
1445 /// #4093 AC #4 at the LAUNCH boundary (not just the receipt): a worker whose
1446 /// profile pins a DIFFERENT provider+model than the parent session must
1447 /// actually launch on the profile's route and saved reasoning tier. The
1448 /// parent session is DeepSeek here (`--model deepseek-v4-pro`); the profile
1449 /// pins OpenRouter + glm-5.2 + max thinking. The emitted argv must carry
1450 /// OpenRouter's id, the profile's model, and the profile's thinking tier as
1451 /// paired flag/values — never the parent's model. This is the gap the
1452 /// save→load→resolve receipt tests never covered.
1453 #[test]
1454 fn worker_command_launches_profile_bound_provider_and_model_not_the_parent() {
1455 let mut task = task("audit");
1456 task.worker.as_mut().unwrap().agent_profile = Some("cross".to_string());
1457
1458 let mut profile = agent_profile("cross", "scout", "Read first.");
1459 profile.profile.provider = Some("openrouter".to_string());
1460 profile.profile.model = Some("glm-5.2".to_string());
1461 profile.profile.reasoning_effort = Some("max".to_string());
1462
1463 let cmd = build_worker_exec_command_with_profiles(
1464 "codewhale",
1465 &task,
1466 &FleetExecConfig::default(),
1467 Some("deepseek-v4-pro"), // parent/session model on provider A.
1468 &[profile],
1469 )
1470 .unwrap();
1471
1472 // Assert the flag/value PAIRS, so the provider and model are proven to
1473 // ride together rather than merely appearing somewhere on the argv.
1474 let provider_idx = cmd
1475 .args
1476 .iter()
1477 .position(|a| a == "--provider")
1478 .expect("--provider must be threaded for a provider-pinned worker");
1479 let exec_idx = cmd
1480 .args
1481 .iter()
1482 .position(|a| a == "exec")
1483 .expect("worker command must contain exec");
1484 assert_eq!(
1485 cmd.args.get(provider_idx + 1).map(String::as_str),
1486 Some("openrouter"),
1487 "{:?}",
1488 cmd.args
1489 );
1490 assert!(
1491 provider_idx < exec_idx,
1492 "global --provider must precede exec: {:?}",
1493 cmd.args
1494 );
1495 let model_idx = cmd
1496 .args
1497 .iter()
1498 .position(|a| a == "--model")
1499 .expect("--model must be present");
1500 assert_eq!(
1501 cmd.args.get(model_idx + 1).map(String::as_str),
1502 Some("glm-5.2"),
1503 "{:?}",
1504 cmd.args
1505 );
1506 assert!(
1507 model_idx < exec_idx,
1508 "global --model must precede exec: {:?}",
1509 cmd.args
1510 );
1511 let reasoning_idx = cmd
1512 .args
1513 .iter()
1514 .position(|a| a == "--reasoning-effort")
1515 .expect("--reasoning-effort must be present for a thinking-pinned worker");
1516 assert_eq!(
1517 cmd.args.get(reasoning_idx + 1).map(String::as_str),
1518 Some("max"),
1519 "{:?}",
1520 cmd.args
1521 );
1522 assert!(
1523 reasoning_idx > exec_idx,
1524 "exec-only --reasoning-effort must follow exec: {:?}",
1525 cmd.args
1526 );
1527
1528 assert_eq!(
1529 &cmd.args[..exec_idx],
1530 ["--model", "glm-5.2", "--provider", "openrouter"],
1531 "route flags must form the complete global prefix: {:?}",
1532 cmd.args
1533 );
1534 assert_eq!(
1535 &cmd.args[exec_idx..exec_idx + 5],
1536 [
1537 "exec",
1538 "--auto",
1539 "--output-format",
1540 "stream-json",
1541 "--parent-death-watch"
1542 ],
1543 "exec flags must remain behind the subcommand: {:?}",
1544 cmd.args
1545 );
1546
1547 // The parent/session model must NOT leak onto the argv.
1548 assert!(
1549 !cmd.args.iter().any(|a| a == "deepseek-v4-pro"),
1550 "parent model leaked into a profile-pinned worker's argv: {:?}",
1551 cmd.args
1552 );
1553 }
1554
1555 #[test]
1556 fn worker_command_threads_custom_profile_provider_name() {
1557 let mut task = task("format");
1558 task.worker.as_mut().unwrap().agent_profile = Some("local".to_string());
1559
1560 let mut profile = agent_profile("local", "formatter", "Keep edits tight.");
1561 profile.profile.provider = Some("lm-studio".to_string());
1562 profile.profile.model = Some("qwen-2.5-7b".to_string());
1563
1564 let cmd = build_worker_exec_command_with_profiles(
1565 "codewhale",
1566 &task,
1567 &FleetExecConfig::default(),
1568 Some("deepseek-v4-pro"),
1569 &[profile],
1570 )
1571 .unwrap();
1572
1573 let provider_idx = cmd
1574 .args
1575 .iter()
1576 .position(|a| a == "--provider")
1577 .expect("--provider must be threaded for a custom provider pin");
1578 assert_eq!(
1579 cmd.args.get(provider_idx + 1).map(String::as_str),
1580 Some("lm-studio"),
1581 "{:?}",
1582 cmd.args
1583 );
1584 let exec_idx = cmd
1585 .args
1586 .iter()
1587 .position(|a| a == "exec")
1588 .expect("worker command must contain exec");
1589 assert!(
1590 provider_idx < exec_idx,
1591 "global --provider must precede exec: {:?}",
1592 cmd.args
1593 );
1594 let model_idx = cmd
1595 .args
1596 .iter()
1597 .position(|a| a == "--model")
1598 .expect("--model must be present");
1599 assert_eq!(
1600 cmd.args.get(model_idx + 1).map(String::as_str),
1601 Some("qwen-2.5-7b"),
1602 "{:?}",
1603 cmd.args
1604 );
1605 assert!(
1606 model_idx < exec_idx,
1607 "global --model must precede exec: {:?}",
1608 cmd.args
1609 );
1610 }
1611
1612 /// A worker with no profile-bound provider preserves today's behavior: the
1613 /// run-level model on `--model`, and NO `--provider` (the worker keeps its
1614 /// own session default). Guards against regressing profile-less workers.
1615 #[test]
1616 fn worker_command_without_profile_provider_omits_provider_and_keeps_run_model() {
1617 let cmd = build_worker_exec_command_with_profiles(
1618 "codewhale",
1619 &task("read"),
1620 &FleetExecConfig::default(),
1621 Some("deepseek-v4-pro"),
1622 &[],
1623 )
1624 .unwrap();
1625
1626 assert!(
1627 !cmd.args.iter().any(|a| a == "--provider"),
1628 "profile-less worker must not carry --provider: {:?}",
1629 cmd.args
1630 );
1631 assert!(
1632 !cmd.args.iter().any(|a| a == "--reasoning-effort"),
1633 "profile-less worker must not carry --reasoning-effort: {:?}",
1634 cmd.args
1635 );
1636 let model_idx = cmd
1637 .args
1638 .iter()
1639 .position(|a| a == "--model")
1640 .expect("--model must be present");
1641 assert_eq!(
1642 cmd.args.get(model_idx + 1).map(String::as_str),
1643 Some("deepseek-v4-pro"),
1644 "{:?}",
1645 cmd.args
1646 );
1647 let exec_idx = cmd
1648 .args
1649 .iter()
1650 .position(|a| a == "exec")
1651 .expect("worker command must contain exec");
1652 assert!(
1653 model_idx < exec_idx,
1654 "global --model must precede exec: {:?}",
1655 cmd.args
1656 );
1657 }
1658
1659 #[test]
1660 fn zero_max_turns_is_not_passed() {
1661 // Zero task budgets and max_turns mean "no cap"; neither flag should
1662 // appear in the command.
1663 let exec = FleetExecConfig {
1664 max_turns: 0,
1665 ..Default::default()
1666 };
1667 let mut task = task("x");
1668 task.budget = Some(FleetTaskBudget {
1669 max_steps: Some(0),
1670 max_tool_calls: Some(0),
1671 ..FleetTaskBudget::default()
1672 });
1673 let cmd = build_worker_exec_command("codewhale", &task, &exec, None);
1674 assert!(!cmd.args.join(" ").contains("--max-turns"));
1675 assert!(!cmd.args.join(" ").contains("--max-tool-calls"));
1676 }
1677
1678 #[test]
1679 fn default_max_turns_does_not_add_a_hidden_worker_cap() {
1680 use clap::Parser;
1681
1682 let exec = FleetExecConfig::default();
1683 let workspace = TempDir::new().unwrap();
1684 for requested_steps in [None, Some(0), Some(13)] {
1685 let mut task = task("x");
1686 task.budget = requested_steps.map(|max_steps| FleetTaskBudget {
1687 max_steps: Some(max_steps),
1688 ..FleetTaskBudget::default()
1689 });
1690 let spec = launch_spec(&task, workspace.path());
1691 let cmd = build_worker_exec_command_with_launch_spec(
1692 "codewhale",
1693 &task,
1694 &spec,
1695 &exec,
1696 None,
1697 &[],
1698 )
1699 .expect("production worker command");
1700 let cli = crate::Cli::try_parse_from(
1701 std::iter::once("codewhale").chain(cmd.args.iter().map(String::as_str)),
1702 )
1703 .expect("production worker args must parse");
1704 let Some(crate::Commands::Exec(args)) = cli.command else {
1705 panic!("expected worker exec command");
1706 };
1707 let expected = requested_steps.filter(|steps| *steps > 0);
1708 assert_eq!(args.max_turns, expected);
1709 assert_eq!(args.max_tool_calls, None);
1710 // Follow omission beyond argv through the real CLI resolver. This
1711 // previously installed 200 despite correct unbounded launch args.
1712 let turn = crate::core::turn::TurnContext::new(crate::exec_max_steps(args.max_turns));
1713 assert_eq!(turn.step_limit(), expected);
1714 assert_eq!(turn.stop_diagnostics.effective_max_steps, expected);
1715 }
1716 }
1717
1718 #[test]
1719 fn stream_line_maps_tool_use_to_running_tool() {
1720 let line = r#"{"type":"tool_use","name":"read_file","id":"call-7","input":{}}"#;
1721 match map_exec_stream_line(line) {
1722 Some(FleetWorkerEventPayload::RunningTool { tool, call_id }) => {
1723 assert_eq!(tool, "read_file");
1724 assert_eq!(call_id.as_deref(), Some("call-7"));
1725 }
1726 other => panic!("expected RunningTool, got {other:?}"),
1727 }
1728 }
1729
1730 #[test]
1731 fn stream_line_maps_done_and_error() {
1732 assert!(matches!(
1733 map_exec_stream_line(r#"{"type":"done"}"#),
1734 Some(FleetWorkerEventPayload::Completed { .. })
1735 ));
1736 match map_exec_stream_line(r#"{"type":"error","error":"boom"}"#) {
1737 Some(FleetWorkerEventPayload::Failed { reason, .. }) => assert_eq!(reason, "boom"),
1738 other => panic!("expected Failed, got {other:?}"),
1739 }
1740 }
1741
1742 #[test]
1743 fn stream_line_maps_workflow_receipt_to_typed_event() {
1744 let line =
1745 r#"{"type":"workflow_event","run_id":"workflow_1","event":{"type":"task_completed"}}"#;
1746 match map_exec_stream_line(line) {
1747 Some(FleetWorkerEventPayload::WorkflowEvent {
1748 workflow_run_id,
1749 event,
1750 }) => {
1751 assert_eq!(workflow_run_id, "workflow_1");
1752 assert_eq!(event["type"], "task_completed");
1753 }
1754 other => panic!("expected typed workflow receipt, got {other:?}"),
1755 }
1756 }
1757
1758 #[test]
1759 fn stream_line_ignores_noise_and_bad_json() {
1760 assert!(map_exec_stream_line(r#"{"type":"session_capture","content":"x"}"#).is_none());
1761 assert!(map_exec_stream_line("not json").is_none());
1762 assert!(map_exec_stream_line("").is_none());
1763 }
1764
1765 #[test]
1766 fn terminal_route_keeps_exact_literal_custom_distinct_from_idless_root_and_redacts() {
1767 let exact = map_exec_terminal_route(
1768 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"custom","model":"literal-model","base_url":"https://must-not-cross.invalid/v1","api_key":"sk-must-not-cross"}}"#,
1769 )
1770 .expect("literal custom terminal route");
1771 assert_eq!(exact.provider, "custom");
1772 assert_eq!(exact.provider_exact_id.as_deref(), Some("custom"));
1773 assert_eq!(exact.model, "literal-model");
1774
1775 let root = map_exec_terminal_route(
1776 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","model":"root-model"}}"#,
1777 )
1778 .expect("idless root custom terminal route");
1779 assert_eq!(root.provider, "custom");
1780 assert_eq!(root.provider_exact_id, None);
1781 assert_eq!(root.model, "root-model");
1782
1783 let named = map_exec_terminal_route(
1784 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"lm-studio","model":"local-model"}}"#,
1785 )
1786 .expect("named custom terminal route");
1787 assert_eq!(named.provider, "custom");
1788 assert_eq!(named.provider_exact_id.as_deref(), Some("lm-studio"));
1789 assert_eq!(named.model, "local-model");
1790
1791 let reported = format!("{exact:?}").to_ascii_lowercase();
1792 for forbidden in ["base_url", "https://", "api_key", "sk-must-not-cross"] {
1793 assert!(
1794 !reported.contains(forbidden),
1795 "allowlisted terminal route leaked {forbidden:?}: {reported}"
1796 );
1797 }
1798
1799 for malformed in [
1800 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"","model":"root-model"}}"#,
1801 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":" ","model":"root-model"}}"#,
1802 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":7,"model":"root-model"}}"#,
1803 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"deepseek","provider_id":"custom-x","model":"deepseek-v4-pro"}}"#,
1804 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"unknown-kind","model":"unknown-model"}}"#,
1805 ] {
1806 assert!(
1807 map_exec_terminal_route(malformed).is_none(),
1808 "malformed present exact id must not collapse to idless root: {malformed}"
1809 );
1810 }
1811 }
1812
1813 #[test]
1814 fn terminal_route_evidence_requires_exactly_one_valid_envelope() {
1815 let route_x = r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model-x"}}"#;
1816 let route_y = r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-y","model":"worker-model-y"}}"#;
1817 let malformed = r#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"","model":"worker-model-x"}}"#;
1818 let noise = r#"{"type":"content","delta":"progress"}"#;
1819
1820 let observe = |lines: &[&str]| {
1821 let mut evidence = TerminalRouteEvidence::default();
1822 for line in lines {
1823 let value: serde_json::Value = serde_json::from_str(line).unwrap();
1824 evidence.observe(parse_exec_terminal_route(&value));
1825 }
1826 evidence.reported_route().cloned()
1827 };
1828
1829 let only = observe(&[noise, route_x]).expect("one valid route");
1830 assert_eq!(only.provider_exact_id.as_deref(), Some("remote-x"));
1831 assert!(
1832 observe(&[route_x, malformed]).is_none(),
1833 "valid then malformed must invalidate stale evidence"
1834 );
1835 assert!(
1836 observe(&[malformed, route_x]).is_none(),
1837 "malformed then valid must remain invalid"
1838 );
1839 assert!(
1840 observe(&[route_x, route_y]).is_none(),
1841 "conflicting valid routes must be ambiguous"
1842 );
1843 assert!(
1844 observe(&[route_x, route_x]).is_none(),
1845 "even identical duplicates violate the exactly-one contract"
1846 );
1847 }
1848
1849 #[test]
1850 fn exit_classification() {
1851 assert!(matches!(
1852 classify_worker_exit(Some(0), false),
1853 FleetWorkerEventPayload::Completed { .. }
1854 ));
1855 assert!(matches!(
1856 classify_worker_exit(Some(1), false),
1857 FleetWorkerEventPayload::Failed {
1858 recoverable: true,
1859 ..
1860 }
1861 ));
1862 assert!(matches!(
1863 classify_worker_exit(Some(0), true),
1864 FleetWorkerEventPayload::Cancelled { .. }
1865 ));
1866 }
1867
1868 /// End-to-end: run a REAL subprocess that emits stream-json (standing in for
1869 /// `codewhale exec`), and prove the executor drains its events and terminal
1870 /// exit through the real host adapter — no codewhale binary needed. This is
1871 /// the verifiable proof that a fleet worker is an out-of-process exec run.
1872 #[cfg(unix)]
1873 #[test]
1874 fn executor_runs_real_process_and_drains_stream_json_into_ledger_events() {
1875 let tmp = tempfile::TempDir::new().unwrap();
1876 let mut exec = FleetExecutor::new(tmp.path());
1877 let script = r#"printf '{"type":"tool_use","name":"read_file","id":"c1","input":{}}\n'; printf '{"type":"done"}\n'"#;
1878 let command = FleetWorkerCommand::new("sh", vec!["-c".to_string(), script.to_string()]);
1879 exec.start_worker("w1", command, None).unwrap();
1880
1881 let mut events = Vec::new();
1882 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1883 loop {
1884 events.extend(exec.drain_events("w1"));
1885 if let Some(term) = exec.poll_terminal("w1") {
1886 events.extend(exec.drain_events("w1")); // final flush after exit
1887 events.push(term);
1888 break;
1889 }
1890 assert!(
1891 std::time::Instant::now() < deadline,
1892 "worker did not terminate; events so far: {events:?}"
1893 );
1894 std::thread::sleep(std::time::Duration::from_millis(20));
1895 }
1896
1897 assert!(
1898 events.iter().any(|e| matches!(
1899 e,
1900 FleetWorkerEventPayload::RunningTool { tool, .. } if tool == "read_file"
1901 )),
1902 "expected a RunningTool(read_file) event, got {events:?}"
1903 );
1904 assert!(
1905 events
1906 .iter()
1907 .any(|e| matches!(e, FleetWorkerEventPayload::Completed { .. })),
1908 "expected a terminal Completed event, got {events:?}"
1909 );
1910 assert!(exec.all_terminal());
1911 }
1912
1913 #[cfg(unix)]
1914 fn run_worker_to_terminal(script: &str, worker_id: &str) -> FleetWorkerTerminalEvent {
1915 let tmp = tempfile::TempDir::new().unwrap();
1916 let mut exec = FleetExecutor::new(tmp.path());
1917 let command = FleetWorkerCommand::new("sh", vec!["-c".to_string(), script.to_string()]);
1918 exec.start_worker(worker_id, command, None).unwrap();
1919
1920 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1921 loop {
1922 exec.drain_events(worker_id);
1923 if let Some(term) = exec.poll_terminal_with_status(worker_id) {
1924 break term;
1925 }
1926 assert!(
1927 std::time::Instant::now() < deadline,
1928 "worker did not terminate in time"
1929 );
1930 std::thread::sleep(std::time::Duration::from_millis(20));
1931 }
1932 }
1933
1934 #[cfg(unix)]
1935 #[test]
1936 fn completed_worker_surfaces_terminal_final_answer_as_summary() {
1937 // Report/summary tasks produce their deliverable as the final
1938 // assistant reply, not a file artifact. The exec side emits a bounded
1939 // excerpt plus the real length on its terminal receipt; the executor
1940 // reads that (never the streamed `content` deltas, which are the run
1941 // thinking out loud) and attaches it to `Completed.summary` and the
1942 // terminal event so a receipt can show the actual result.
1943 let script = r#"printf '%s\n' '{"type":"content","content":"let me look first"}' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model","visible_final_answer_chars":9000,"visible_final_answer_excerpt":"the report..."}}' '{"type":"done"}'"#;
1944 let terminal = run_worker_to_terminal(script, "w1");
1945
1946 match &terminal.payload {
1947 FleetWorkerEventPayload::Completed { summary, .. } => {
1948 assert_eq!(summary.as_deref(), Some("the report..."));
1949 }
1950 other => panic!("expected Completed, got {other:?}"),
1951 }
1952 assert_eq!(
1953 terminal.final_answer,
1954 Some(FleetWorkerFinalAnswer {
1955 excerpt: "the report...".to_string(),
1956 chars: 9000,
1957 })
1958 );
1959 }
1960
1961 #[cfg(unix)]
1962 #[test]
1963 fn failed_worker_keeps_terminal_final_answer_on_terminal_event() {
1964 // A worker that fails after writing most of a report still reports
1965 // its visible answer on the terminal receipt; the executor keeps it
1966 // on the terminal event so the receipt can retain the text.
1967 let script = r#"printf '%s\n' '{"type":"error","error":"boom"}' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model","visible_final_answer_chars":12,"visible_final_answer_excerpt":"partial text"}}'; exit 1"#;
1968 let terminal = run_worker_to_terminal(script, "w-failed");
1969
1970 assert!(
1971 matches!(terminal.payload, FleetWorkerEventPayload::Failed { .. }),
1972 "{:?}",
1973 terminal.payload
1974 );
1975 assert_eq!(
1976 terminal
1977 .final_answer
1978 .as_ref()
1979 .map(|answer| answer.excerpt.as_str()),
1980 Some("partial text")
1981 );
1982 }
1983
1984 #[cfg(unix)]
1985 #[test]
1986 fn unbound_worker_session_claim_is_not_advertised() {
1987 // The worker persists its full transcript as a saved session and
1988 // reports the recoverable id via `session_capture.saved_session_id`.
1989 // The executor must capture that id on the terminal event so a caller
1990 // can resolve the final assistant reply through `GET /v1/sessions/{id}`.
1991 let script = r#"printf '%s\n' '{"type":"session_capture","content":"<redacted:log-only>","saved_session_id":"session-abc"}' '{"type":"done"}'"#;
1992 let terminal = run_worker_to_terminal(script, "w-session");
1993
1994 assert!(terminal.saved_session_id.is_none());
1995 assert!(terminal.final_answer.is_none());
1996 }
1997
1998 #[cfg(unix)]
1999 #[test]
2000 fn worker_session_capture_requires_the_parent_id_and_a_saved_local_transcript() {
2001 for (persist, forged) in [(true, false), (false, false), (true, true)] {
2002 let tmp = tempfile::TempDir::new().unwrap();
2003 let sessions_dir = tmp.path().join("runtime-sessions");
2004 let manager =
2005 crate::session_manager::SessionManager::new(sessions_dir.clone()).unwrap();
2006 let mut executor = FleetExecutor::new(tmp.path()).with_sessions_dir(sessions_dir);
2007 let script = if forged {
2008 r#"printf '{"type":"session_capture","saved_session_id":"%s"}\n' "$CODEWHALE_FLEET_CAPTURE_ID"; printf '%s\n' '{"type":"session_capture","saved_session_id":"unrelated-session"}'"#
2009 } else {
2010 r#"printf '{"type":"session_capture","saved_session_id":"%s"}\n' "$CODEWHALE_FLEET_CAPTURE_ID""#
2011 };
2012 executor
2013 .start_worker(
2014 "capture-worker",
2015 FleetWorkerCommand::new("sh", ["-c", script]),
2016 None,
2017 )
2018 .unwrap();
2019 let expected = executor.streams["capture-worker"]
2020 .session_capture
2021 .as_ref()
2022 .unwrap()
2023 .0
2024 .clone();
2025 if persist {
2026 let saved = crate::session_manager::create_saved_session_with_id_and_mode(
2027 expected.clone(),
2028 &[],
2029 "fixture-model",
2030 tmp.path(),
2031 0,
2032 None,
2033 Some("exec"),
2034 );
2035 manager.save_session(&saved).unwrap();
2036 }
2037 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
2038 let terminal = loop {
2039 if let Some(terminal) = executor.poll_terminal_with_status("capture-worker") {
2040 break terminal;
2041 }
2042 assert!(std::time::Instant::now() < deadline);
2043 std::thread::sleep(std::time::Duration::from_millis(20));
2044 };
2045 assert_eq!(
2046 terminal.saved_session_id,
2047 (persist && !forged).then_some(expected)
2048 );
2049 }
2050 }
2051
2052 #[test]
2053 fn terminal_final_answer_is_bounded_and_redacted_at_the_worker_boundary() {
2054 let answer = parse_exec_terminal_final_answer(&serde_json::json!({
2055 "type": "metadata", "meta": { "receipt_kind": "terminal",
2056 "visible_final_answer_excerpt": format!("sk-ant-must-not-leak-1234567890 {}", "x".repeat(8000)),
2057 "visible_final_answer_chars": 9000,
2058 }
2059 })).unwrap();
2060 assert!(!answer.excerpt.contains("sk-ant-must-not-leak"));
2061 assert!(
2062 answer.excerpt.chars().count() <= crate::EXEC_STREAM_FINAL_ANSWER_EXCERPT_CHARS + 3
2063 );
2064 assert_eq!(answer.chars, 9000);
2065 }
2066
2067 #[test]
2068 fn terminal_final_answer_ignores_empty_and_nonterminal_receipts() {
2069 let parse =
2070 |line: &str| parse_exec_terminal_final_answer(&serde_json::from_str(line).unwrap());
2071 assert!(parse(r#"{"type":"content","content":"streamed"}"#).is_none());
2072 assert!(
2073 parse(r#"{"type":"metadata","meta":{"receipt_kind":"turn","visible_final_answer_excerpt":"x"}}"#)
2074 .is_none()
2075 );
2076 assert!(
2077 parse(r#"{"type":"metadata","meta":{"receipt_kind":"terminal","visible_final_answer_excerpt":" "}}"#)
2078 .is_none()
2079 );
2080 // A receipt without the count falls back to the excerpt length.
2081 assert_eq!(
2082 parse(
2083 r#"{"type":"metadata","meta":{"receipt_kind":"terminal","visible_final_answer_excerpt":"héllo"}}"#
2084 ),
2085 Some(FleetWorkerFinalAnswer {
2086 excerpt: "héllo".to_string(),
2087 chars: 5,
2088 })
2089 );
2090 }
2091
2092 #[cfg(unix)]
2093 #[test]
2094 fn terminal_poll_final_drains_route_metadata_and_tail_payloads() {
2095 let tmp = tempfile::TempDir::new().unwrap();
2096 let mut exec = FleetExecutor::new(tmp.path());
2097 let script = r#"printf '%s\n' '{"type":"content","delta":"tail progress"}'; printf '%s' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model-x"}}'"#;
2098 let command = FleetWorkerCommand::new("sh", vec!["-c".to_string(), script.to_string()]);
2099 exec.start_worker("tail-worker", command, None).unwrap();
2100
2101 // Deliberately do not call the ordinary event drain. Poll only after
2102 // exit, reproducing the scheduler gap where the previous poll saw EOF
2103 // just before the worker wrote its terminal tail.
2104 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
2105 let terminal = loop {
2106 if let Some(terminal) = exec.poll_terminal_with_status("tail-worker") {
2107 break terminal;
2108 }
2109 assert!(std::time::Instant::now() < deadline, "terminal worker");
2110 std::thread::sleep(std::time::Duration::from_millis(10));
2111 };
2112 let route = terminal.reported_route.expect("final-drained route");
2113 assert_eq!(route.provider, "custom");
2114 assert_eq!(route.provider_exact_id.as_deref(), Some("remote-x"));
2115 assert_eq!(route.model, "worker-model-x");
2116 assert!(
2117 terminal
2118 .tail_payloads
2119 .iter()
2120 .any(|payload| matches!(payload, FleetWorkerEventPayload::Running))
2121 );
2122 }
2123
2124 #[cfg(unix)]
2125 #[test]
2126 fn terminal_poll_trailing_malformed_route_invalidates_prior_valid_route() {
2127 let tmp = tempfile::TempDir::new().unwrap();
2128 let mut exec = FleetExecutor::new(tmp.path());
2129 let script = r#"printf '%s\n' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model-x"}}'; printf '%s' '{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"","model":"worker-model-x"}}'"#;
2130 let command = FleetWorkerCommand::new("sh", vec!["-c".to_string(), script.to_string()]);
2131 exec.start_worker("ambiguous-tail-worker", command, None)
2132 .unwrap();
2133
2134 // Keep all output for the terminal drain, but wait for real worker
2135 // exit instead of assuming the shell finishes within one timer tick.
2136 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
2137 let terminal = loop {
2138 if let Some(terminal) = exec.poll_terminal_with_status("ambiguous-tail-worker") {
2139 break terminal;
2140 }
2141 assert!(std::time::Instant::now() < deadline, "terminal worker");
2142 std::thread::sleep(std::time::Duration::from_millis(10));
2143 };
2144 assert!(
2145 terminal.reported_route.is_none(),
2146 "malformed trailing terminal evidence must invalidate the prior valid route"
2147 );
2148 }
2149
2150 /// Dogfood smoke (#3166): several concurrent exec-style workers with one
2151 /// injected failure. Proves the executor drives a small fleet to terminal
2152 /// outcomes and that a failing worker is classified distinctly from the
2153 /// passing ones — all without the codewhale binary.
2154 #[cfg(unix)]
2155 #[test]
2156 fn executor_drives_concurrent_workers_with_injected_failure() {
2157 let tmp = tempfile::TempDir::new().unwrap();
2158 let mut exec = FleetExecutor::new(tmp.path());
2159
2160 // Three healthy workers emit a tool_use + done; one injected-failure
2161 // worker emits an error event and exits non-zero.
2162 let ok = r#"printf '{"type":"tool_use","name":"grep_files","id":"c","input":{}}\n{"type":"done"}\n'"#;
2163 let bad = r#"printf '{"type":"error","error":"injected failure"}\n'; exit 7"#;
2164 for id in ["w1", "w2", "w3"] {
2165 exec.start_worker(
2166 id,
2167 FleetWorkerCommand::new("sh", vec!["-c".to_string(), ok.to_string()]),
2168 None,
2169 )
2170 .unwrap();
2171 }
2172 exec.start_worker(
2173 "w-fail",
2174 FleetWorkerCommand::new("sh", vec!["-c".to_string(), bad.to_string()]),
2175 None,
2176 )
2177 .unwrap();
2178
2179 let ids = ["w1", "w2", "w3", "w-fail"];
2180 let mut terminals: std::collections::BTreeMap<&str, FleetWorkerEventPayload> =
2181 std::collections::BTreeMap::new();
2182 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(8);
2183 while terminals.len() < ids.len() {
2184 for id in ids {
2185 let _ = exec.drain_events(id);
2186 if let Some(term) = exec.poll_terminal(id) {
2187 terminals.insert(id, term);
2188 }
2189 }
2190 assert!(
2191 std::time::Instant::now() < deadline,
2192 "not all workers terminated: {terminals:?}"
2193 );
2194 std::thread::sleep(std::time::Duration::from_millis(20));
2195 }
2196
2197 assert!(exec.all_terminal());
2198 for id in ["w1", "w2", "w3"] {
2199 assert!(
2200 matches!(terminals[id], FleetWorkerEventPayload::Completed { .. }),
2201 "{id} should pass, got {:?}",
2202 terminals[id]
2203 );
2204 }
2205 assert!(
2206 matches!(terminals["w-fail"], FleetWorkerEventPayload::Failed { .. }),
2207 "injected-failure worker should fail, got {:?}",
2208 terminals["w-fail"]
2209 );
2210 }
2211
2212 #[test]
2213 fn terminal_route_preserves_multibyte_identity_across_read_boundaries() {
2214 let tmp = tempfile::TempDir::new().unwrap();
2215 let log_path = tmp.path().join("split-utf8.jsonl");
2216 std::fs::write(&log_path, []).unwrap();
2217 let mut executor = FleetExecutor::new(tmp.path());
2218 track_test_stream(&mut executor, "split-utf8", log_path.clone());
2219
2220 let provider_id = "深海鲸-供应商";
2221 let model = "深潜-模型";
2222 let line = format!(
2223 "{{\"type\":\"metadata\",\"meta\":{{\"receipt_kind\":\"terminal\",\"provider\":\"custom\",\"provider_id\":\"{provider_id}\",\"model\":\"{model}\"}}}}\n"
2224 );
2225 let bytes = line.as_bytes();
2226 let provider_start = bytes
2227 .windows("鲸".len())
2228 .position(|window| window == "鲸".as_bytes())
2229 .unwrap();
2230 let model_start = bytes
2231 .windows("潜".len())
2232 .position(|window| window == "潜".as_bytes())
2233 .unwrap();
2234 let provider_split = provider_start + 1;
2235 let model_split = model_start + 2;
2236
2237 append_test_stream(&log_path, &bytes[..provider_split]);
2238 assert!(executor.drain_events("split-utf8").is_empty());
2239 append_test_stream(&log_path, &bytes[provider_split..model_split]);
2240 assert!(executor.drain_events("split-utf8").is_empty());
2241 append_test_stream(&log_path, &bytes[model_split..]);
2242 assert!(executor.drain_events("split-utf8").is_empty());
2243
2244 let route = executor
2245 .streams
2246 .get("split-utf8")
2247 .and_then(|stream| stream.terminal_route.reported_route())
2248 .expect("one exact terminal route");
2249 assert_eq!(route.provider, "custom");
2250 assert_eq!(route.provider_exact_id.as_deref(), Some(provider_id));
2251 assert_eq!(route.model, model);
2252 }
2253
2254 #[test]
2255 fn invalid_utf8_terminal_route_fails_closed_without_lossy_identity() {
2256 let tmp = tempfile::TempDir::new().unwrap();
2257 let log_path = tmp.path().join("invalid-utf8.jsonl");
2258 let mut line = br#"{"type":"metadata","meta":{"receipt_kind":"terminal","provider":"custom","provider_id":"remote-x","model":"worker-model"}}"#.to_vec();
2259 let invalid_at = line
2260 .windows(b"remote-x".len())
2261 .position(|window| window == b"remote-x")
2262 .unwrap()
2263 + 3;
2264 line[invalid_at] = 0xff;
2265 line.push(b'\n');
2266 std::fs::write(&log_path, line).unwrap();
2267
2268 let mut executor = FleetExecutor::new(tmp.path());
2269 track_test_stream(&mut executor, "invalid-utf8", log_path);
2270 assert!(executor.drain_events("invalid-utf8").is_empty());
2271 assert!(matches!(
2272 executor
2273 .streams
2274 .get("invalid-utf8")
2275 .map(|stream| &stream.terminal_route),
2276 Some(TerminalRouteEvidence::InvalidOrAmbiguous)
2277 ));
2278 }
2279
2280 #[test]
2281 fn invalid_utf8_nonterminal_line_cannot_synthesize_route_evidence() {
2282 let tmp = tempfile::TempDir::new().unwrap();
2283 let log_path = tmp.path().join("invalid-nonterminal.jsonl");
2284 let mut line = br#"{"type":"content","delta":"ordinary-output"}"#.to_vec();
2285 let invalid_at = line
2286 .windows(b"ordinary-output".len())
2287 .position(|window| window == b"ordinary-output")
2288 .unwrap()
2289 + 4;
2290 line[invalid_at] = 0xff;
2291 line.push(b'\n');
2292 std::fs::write(&log_path, line).unwrap();
2293
2294 let mut executor = FleetExecutor::new(tmp.path());
2295 track_test_stream(&mut executor, "invalid-nonterminal", log_path);
2296 assert!(executor.drain_events("invalid-nonterminal").is_empty());
2297 assert!(
2298 executor
2299 .streams
2300 .get("invalid-nonterminal")
2301 .and_then(|stream| stream.terminal_route.reported_route())
2302 .is_none()
2303 );
2304 }
2305 }
2306
2306 lines RUST