返回 CodeWhale
exec_agent.rs
根目录 / crates / tui / src / exec_agent.rs
1 //! Non-interactive exec agent assembly: the `run_exec_agent` pipeline
2 //! that resolves the CLI route, builds the engine configuration, spawns
3 //! the engine, and drives the exec output stream to completion.
4 //!
5 //! Extracted verbatim from `lib.rs` (#5586, the issue's prescribed
6 //! engine-config-assembly cut). The two functions were crate-private in
7 //! the root and are `pub(crate)` here purely so the root's glob re-export
8 //! keeps the dispatch site and tests resolving unchanged.
9
10 use super::*;
11 use crate::core::ops::TurnSpec;
12
13 /// Resolve the headless `exec` model-step ceiling.
14 ///
15 /// Omission leaves model steps uncapped. Clap rejects `--max-turns 0`;
16 /// explicit positive values retain the documented finite range.
17 pub(crate) fn exec_max_steps(max_turns: Option<u32>) -> u32 {
18 crate::core::engine::turn_budget::resolve_max_model_steps(max_turns)
19 }
20
21 /// Model-step ceiling for a plain (zero-tool) `exec` run without
22 /// `--max-turns`. Its only extra steps are output-limit continuations, which
23 /// have no progress signal of their own; without this a model stuck at the
24 /// output limit would be re-asked, with growing history, until the turn wall
25 /// clock (#6510 review).
26 pub(crate) const ONE_SHOT_DEFAULT_MAX_STEPS: u32 = 8;
27
28 /// Default-denied tools for headless `exec`, on top of the operator's own
29 /// `--disallowed-tools` flag.
30 ///
31 /// A headless run has no responder for `request_user_input`, so offering the
32 /// tool can only stall the run until the turn wall clock, or forever with
33 /// `[tools] user_input_timeout_seconds = 0`. Withholding it is the default
34 /// form of the operator workaround (`--disallowed-tools request_user_input`):
35 /// the model reports the tool absent and finishes instead of parking. This
36 /// stays unconditional: there is no channel on which a one-shot CLI run
37 /// could answer, so advertising the tool cannot work.
38 pub(crate) fn exec_disallowed_tools(disallowed_tools: Option<Vec<String>>) -> Option<Vec<String>> {
39 use crate::core::engine::tool_catalog::REQUEST_USER_INPUT_NAME;
40 let mut disallowed = disallowed_tools.unwrap_or_default();
41 if !disallowed
42 .iter()
43 .any(|tool| tool.as_str() == REQUEST_USER_INPUT_NAME)
44 {
45 disallowed.push(REQUEST_USER_INPUT_NAME.to_string());
46 }
47 Some(disallowed)
48 }
49
50 type ExecSettlementProbe = std::pin::Pin<
51 Box<dyn std::future::Future<Output = Result<crate::core::ops::SubAgentSettlement>> + Send>,
52 >;
53
54 /// Read the existing Engine stream through the one-shot host's final boundary.
55 /// Successful parent receipts remain pending while admitted children or their
56 /// completion inbox can still produce another normal Engine turn. This owns
57 /// only the deferred output receipt, never child execution or a turn loop.
58 pub(crate) struct ExecAgentEvents {
59 handle: crate::core::engine::EngineHandle,
60 deadline: tokio::time::Instant,
61 terminal: Option<crate::core::events::Event>,
62 probe: Option<ExecSettlementProbe>,
63 next_probe_at: tokio::time::Instant,
64 in_flight_usage: codewhale_models::Usage,
65 }
66
67 impl ExecAgentEvents {
68 pub(crate) fn new(handle: crate::core::engine::EngineHandle, deadline: Instant) -> Self {
69 Self {
70 handle,
71 deadline: deadline.into(),
72 terminal: None,
73 probe: None,
74 next_probe_at: tokio::time::Instant::now(),
75 in_flight_usage: codewhale_models::Usage::default(),
76 }
77 }
78
79 fn stop_settlement(
80 &mut self,
81 status: crate::core::events::TurnOutcomeStatus,
82 error: String,
83 ) -> crate::core::events::Event {
84 use crate::core::events::Event;
85 self.handle
86 .cancel_with_reason(crate::core::engine::CancelReason::External);
87 // Cancellation is out of band; shutdown remains in the existing
88 // Engine mailbox so it also cancels detached session children.
89 let _ = self.handle.try_send(crate::core::ops::Op::Shutdown);
90 self.probe = None;
91 let mut terminal = self.terminal.take().expect("pending parent receipt");
92 if let Event::TurnComplete {
93 status: terminal_status,
94 error: terminal_error,
95 usage,
96 parent_route_usage,
97 routed_usage_dropped_records,
98 ..
99 } = &mut terminal
100 {
101 *terminal_status = status;
102 *terminal_error = Some(error);
103 crate::core::turn::add_usage_to(usage, &self.in_flight_usage);
104 crate::core::turn::add_usage_to(parent_route_usage, &self.in_flight_usage);
105 // The host cannot prove usage settlement after abandoning the
106 // inbox. Keep reported usage and explicitly mark coverage partial.
107 *routed_usage_dropped_records = routed_usage_dropped_records.saturating_add(1);
108 }
109 self.in_flight_usage = codewhale_models::Usage::default();
110 terminal
111 }
112
113 pub(crate) async fn next(&mut self) -> Option<crate::core::events::Event> {
114 use crate::core::events::{Event, TurnOutcomeStatus};
115 // Keep the streamed event on the stack instead of allocating another
116 // box for every token merely to equalize the two small control arms.
117 #[allow(clippy::large_enum_variant)]
118 enum Input {
119 Event(Option<Event>),
120 Probe(Result<crate::core::ops::SubAgentSettlement>),
121 Poll,
122 }
123 loop {
124 if matches!(
125 self.terminal,
126 Some(Event::TurnComplete { status, .. }) if status != TurnOutcomeStatus::Completed
127 ) {
128 return self.terminal.take();
129 }
130 let settling = self.terminal.is_some();
131 if settling && self.handle.is_cancelled() {
132 return Some(self.stop_settlement(
133 TurnOutcomeStatus::Interrupted,
134 "Headless exec cancelled while settling children; recorded usage is partial."
135 .to_string(),
136 ));
137 }
138 if settling && tokio::time::Instant::now() >= self.deadline {
139 return Some(self.stop_settlement(
140 TurnOutcomeStatus::Failed,
141 "Headless exec wall-clock budget exhausted while settling children; recorded usage is partial."
142 .to_string(),
143 ));
144 }
145 if settling && self.probe.is_none() && tokio::time::Instant::now() >= self.next_probe_at
146 {
147 let handle = self.handle.clone();
148 self.probe = Some(Box::pin(
149 async move { handle.get_subagent_settlement().await },
150 ));
151 }
152 let probing = self.probe.is_some();
153 let wake_at = self.deadline.min(if probing {
154 tokio::time::Instant::now() + Duration::from_millis(250)
155 } else {
156 self.next_probe_at
157 });
158 let input = {
159 let mut events = self.handle.rx_event.write().await;
160 tokio::select! {
161 biased;
162 // Drain queued SessionUpdated/TurnComplete events before
163 // accepting the later actor-owned idle receipt.
164 event = events.recv() => Input::Event(event),
165 result = async { self.probe.as_mut().expect("active probe").await }, if probing => Input::Probe(result),
166 () = tokio::time::sleep_until(wake_at), if settling => Input::Poll,
167 }
168 };
169 match input {
170 Input::Poll => {}
171 Input::Probe(Ok(snapshot)) if snapshot.is_settled() => {
172 self.probe = None;
173 return self.terminal.take();
174 }
175 Input::Probe(Ok(_)) => {
176 self.probe = None;
177 self.next_probe_at = tokio::time::Instant::now() + Duration::from_millis(250);
178 }
179 Input::Probe(Err(error)) => {
180 return Some(self.stop_settlement(
181 TurnOutcomeStatus::Failed,
182 format!(
183 "Cannot verify child settlement: {error}; recorded usage is partial."
184 ),
185 ));
186 }
187 Input::Event(None) if settling => {
188 return Some(self.stop_settlement(
189 TurnOutcomeStatus::Failed,
190 "Engine event channel closed before child settlement; recorded usage is partial."
191 .to_string(),
192 ));
193 }
194 Input::Event(None) => return None,
195 Input::Event(Some(mut event)) => {
196 match &mut event {
197 Event::TurnComplete {
198 usage,
199 parent_route_usage,
200 routed_usage_dropped_records,
201 status,
202 error,
203 ..
204 } => {
205 if let Some(Event::TurnComplete {
206 usage: prior_usage,
207 parent_route_usage: prior_parent_usage,
208 routed_usage_dropped_records: prior_dropped,
209 ..
210 }) = self.terminal.take()
211 {
212 crate::core::turn::add_usage_to(usage, &prior_usage);
213 crate::core::turn::add_usage_to(
214 parent_route_usage,
215 &prior_parent_usage,
216 );
217 *routed_usage_dropped_records =
218 routed_usage_dropped_records.saturating_add(prior_dropped);
219 }
220 self.in_flight_usage = codewhale_models::Usage::default();
221 if *status == TurnOutcomeStatus::Completed && error.is_none() {
222 self.terminal = Some(event);
223 self.next_probe_at = tokio::time::Instant::now();
224 continue;
225 }
226 }
227 Event::TurnUsage { usage, .. } => {
228 crate::core::turn::add_usage_to(&mut self.in_flight_usage, usage);
229 }
230 Event::Error { envelope, .. }
231 if settling && exec_error_event_is_fatal(envelope) =>
232 {
233 let terminal = self.stop_settlement(
234 TurnOutcomeStatus::Failed,
235 format!(
236 "{}; child settlement stopped and recorded usage is partial.",
237 envelope.message
238 ),
239 );
240 self.terminal = Some(terminal);
241 }
242 _ => {}
243 }
244 return Some(event);
245 }
246 }
247 }
248 }
249 }
250
251 /// Attach the durable automation store headless `exec` inspects.
252 ///
253 /// Headless exec builds its catalog from the same tool surface the TUI and the
254 /// Runtime host do, so it advertises `automation` and `send_later` whether or
255 /// not the store behind them is attached. Left unattached, every call failed
256 /// "AutomationManager is not attached" — the tool was real and the service was
257 /// missing. This opens the same store those two hosts open: a shared directory
258 /// guarded per transaction by its own file locks (`AutomationManager::open`),
259 /// so it adds no second store, no scheduler, and no second scheduling
260 /// authority.
261 ///
262 /// What exec deliberately does not take is the Runtime's task-execution lease.
263 /// That lease is exclusive (`TaskExecutionLease::new`) and a one-shot host must
264 /// neither contend with it nor recover work from the process that holds it. So
265 /// the manager returned here is *unbound*: inspection works, and everything
266 /// that would promise dispatch is refused by the tool's own admission check
267 /// (`tools::automation::require_dispatch_owner`) rather than persisting a
268 /// schedule nothing would honor.
269 ///
270 /// Fleet worker subprocesses get nothing, keeping the narrowed envelope they
271 /// were launched with alongside the empty plugin registry and disabled
272 /// subagents. A store that cannot be opened is reported, never swallowed:
273 /// both other hosts fail startup on it, so exec does too instead of
274 /// advertising an automation surface it silently cannot serve.
275 pub(crate) fn exec_automation_services(
276 fleet_authority_active: bool,
277 ) -> Result<Option<crate::automation_manager::SharedAutomationManager>> {
278 if fleet_authority_active {
279 return Ok(None);
280 }
281 let service = crate::automation_manager::AutomationManager::default_location()
282 .context("open the automation store for headless exec")?;
283 Ok(Some(std::sync::Arc::new(tokio::sync::Mutex::new(service))))
284 }
285
286 /// Printed to stderr when a tool-less one-shot `exec` answer contained
287 /// tool-call markup that the engine stripped from the visible output.
288 pub(crate) const ONE_SHOT_TOOL_CALL_NOTICE: &str = "codewhale exec: the model tried to call a tool, but this run offers none, so the tool call was removed from the answer. Re-run with --auto (or --allowed-tools) to let it use tools.";
289
290 #[allow(clippy::too_many_arguments)]
291 pub(crate) async fn run_exec_agent(
292 config: &Config,
293 model: &str,
294 prompt: &str,
295 workspace: PathBuf,
296 max_subagents: usize,
297 auto_approve: bool,
298 allow_sandbox_elevation: bool,
299 explicit_sandbox: Option<&str>,
300 trust_mode: bool,
301 json_output: bool,
302 resume_session: Option<session_manager::SavedSession>,
303 force_configured_route: bool,
304 output_format: ExecOutputFormat,
305 max_turns: u32,
306 max_tool_calls: Option<u32>,
307 allowed_tools: Option<Vec<String>>,
308 disallowed_tools: Option<Vec<String>>,
309 append_system_prompt: Option<String>,
310 tool_authority_json: Option<String>,
311 exec_hooks_enabled: bool,
312 plugin_registry: std::sync::Arc<crate::plugins::PluginRegistry>,
313 // #6510: plain `exec` (no tool-surface flag). The caller passes an empty
314 // allowlist; this also skips workspace snapshots, LSP and the automation
315 // store, and writes the one-shot `--json` receipt shape.
316 one_shot: bool,
317 ) -> Result<()> {
318 use crate::compaction::CompactionConfig;
319 use crate::core::engine::{EngineConfig, spawn_engine};
320 use crate::core::events::Event;
321 use crate::core::ops::Op;
322 use crate::tools::plan::new_shared_plan_state;
323 use crate::tools::todo::new_shared_todo_list;
324 use codewhale_config::AppMode;
325 use codewhale_execpolicy::ApprovalMode;
326
327 ignore_sigpipe_for_headless_exec();
328
329 // Withhold `request_user_input`; a headless run has no responder.
330 let disallowed_tools = exec_disallowed_tools(disallowed_tools);
331
332 // Headless exec registers the model-facing notify tool too. Project the
333 // final merged config before tool setup so `off`, quiet/category gates,
334 // and explicit `always` are truthful outside the interactive TUI. With no
335 // focus-reporting channel, fail closed to focused; only explicit `always`
336 // may authorize a headless desktop notification.
337 let terminal = crate::host_terminal::host();
338 terminal.set_terminal_focused(true);
339 terminal.apply_notification_settings(&config.notifications_config());
340
341 validate_exec_tool_authority_resume(tool_authority_json.as_deref(), resume_session.is_some())?;
342 let fleet_authority = tool_authority_json
343 .as_deref()
344 .map(crate::tools::spec::ToolAuthorityEnvelope::from_json)
345 .transpose()
346 .map_err(anyhow::Error::msg)?;
347 let fleet_authority_active = fleet_authority.is_some();
348 let outer_network_access = fleet_authority
349 .as_ref()
350 .and_then(|authority| authority.network_access);
351 let outer_shell_authority = fleet_authority
352 .as_ref()
353 .map(|authority| authority.shell)
354 .unwrap_or_default();
355 if let Some(envelope) = fleet_authority {
356 crate::tools::spec::install_process_tool_authority(envelope).map_err(anyhow::Error::msg)?;
357 }
358
359 let fleet_capture = match (
360 std::env::var("CODEWHALE_FLEET_CAPTURE_ID").ok(),
361 std::env::var_os("CODEWHALE_FLEET_CAPTURE_DIR"),
362 ) {
363 (Some(id), Some(dir)) => {
364 uuid::Uuid::parse_str(&id).context("invalid Fleet session capture id")?;
365 anyhow::ensure!(
366 resume_session.is_none(),
367 "Fleet capture cannot resume a session"
368 );
369 let manager = SessionManager::new(PathBuf::from(dir))?;
370 match manager.load_session(&id) {
371 Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
372 _ => anyhow::bail!("Fleet session capture id is already in use or unavailable"),
373 }
374 Some((id, manager))
375 }
376 (None, None) => None,
377 _ => anyhow::bail!("incomplete Fleet session capture destination"),
378 };
379
380 let route = resolve_cli_exec_route(config, model, prompt, force_configured_route).await?;
381 let execution_config = config_for_cli_route(config, &route)?;
382 let auto_model = route.auto_model;
383 let effective_identity = execution_config
384 .active_provider_identity()
385 .map_err(anyhow::Error::msg)?;
386 let effective_provider = effective_identity.provider;
387 let effective_model = route.model;
388 let validated_route = crate::route_runtime::resolve_runtime_route_for_identity(
389 &execution_config,
390 &effective_identity,
391 Some(&effective_model),
392 )
393 .map_err(anyhow::Error::msg)?
394 .validate_for(crate::route_runtime::RouteErrorSurface::Headless)
395 .map_err(anyhow::Error::msg)?;
396 let effective_identity = validated_route.identity.clone();
397 let effective_provider_name = effective_identity.key.to_string();
398 let (effective_provider_kind, effective_stream_provider_id) =
399 exec_stream_provider_route(&validated_route.identity);
400 let route_source = if auto_model {
401 "auto_resolver"
402 } else {
403 "explicit_or_configured"
404 }
405 .to_string();
406 let exec_started = Instant::now();
407 let prompt_sha256 = format!("sha256:{}", crate::hashing::sha256_hex(prompt.as_bytes()));
408 let binary_sha256 = current_binary_sha256();
409 let approval_posture = if auto_approve { "auto_tools" } else { "ask" }.to_string();
410 let sandbox_posture = explicit_sandbox.unwrap_or("configured_default").to_string();
411 let active_route_limits =
412 crate::route_budget::known_route_limits(validated_route.candidate.limits());
413 let max_subagents = if max_subagents
414 == config.max_subagents_for_provider(
415 &config
416 .active_provider_identity()
417 .map_err(anyhow::Error::msg)?,
418 ) {
419 execution_config
420 .max_subagents_for_provider(&effective_identity)
421 .clamp(1, MAX_SUBAGENTS)
422 } else {
423 max_subagents
424 };
425 // A FIXED model with `--reasoning-effort auto` (the exact shape a Fleet
426 // worker subprocess launches with: `--model <exact> --reasoning-effort
427 // auto`) is still Auto. `auto_model` is a *model* decision and is false
428 // here, so deriving the auto flag from it left this path both raw and
429 // non-auto: the literal string `"auto"` travelled to the engine while the
430 // receipt claimed no Auto was in play.
431 let reasoning_effort_auto = route.auto_controls_reasoning;
432 // Resolve Auto against this run's prompt at the CLI boundary, exactly like
433 // the interactive launch path does, so the tier the engine (and the
434 // receipt below) sees is concrete.
435 let effective_reasoning_effort = route.reasoning_effort.and_then(|effort| {
436 cli_reasoning_effort_value_for_prompt(&execution_config, &effective_model, effort)
437 });
438
439 let settings = crate::settings::Settings::load().unwrap_or_default();
440 let auto_compact_enabled = if crate::settings::Settings::auto_compact_explicitly_configured() {
441 settings.auto_compact
442 } else {
443 crate::route_budget::auto_compact_default_for_route(
444 effective_provider,
445 &effective_model,
446 active_route_limits,
447 )
448 };
449 let compaction = CompactionConfig {
450 enabled: auto_compact_enabled,
451 model: effective_model.clone(),
452 effective_context_window: Some(crate::route_budget::route_context_window_tokens(
453 effective_provider,
454 &effective_model,
455 active_route_limits,
456 )),
457 token_threshold: crate::route_budget::compaction_threshold_for_route_at_percent(
458 effective_provider,
459 &effective_model,
460 active_route_limits,
461 settings.auto_compact_threshold_percent,
462 ),
463 summary_instructions: execution_config.compaction_summary_instructions(),
464 retained_user_message_tokens: execution_config.compaction_retained_user_message_tokens(),
465 ..Default::default()
466 };
467
468 let network_policy = exec_network_policy(&execution_config, outer_network_access);
469
470 let lsp_config = (!fleet_authority_active && !one_shot)
471 .then(|| {
472 execution_config
473 .lsp
474 .clone()
475 .map(crate::config::LspConfigToml::into_runtime)
476 })
477 .flatten();
478 let mut engine_features = execution_config.features();
479 apply_fleet_engine_feature_caps(
480 &mut engine_features,
481 fleet_authority_active,
482 outer_network_access,
483 outer_shell_authority,
484 );
485 if crate::core::allowlist_is_native_file_and_shell_only(allowed_tools.as_deref()) {
486 engine_features.disable(crate::features::Feature::Mcp);
487 }
488 let engine_plugin_registry = if fleet_authority_active {
489 std::sync::Arc::new(crate::plugins::PluginRegistry::empty(&workspace))
490 } else {
491 plugin_registry
492 };
493 // `exec --hooks` (#6099) is the operator's explicit opt-in: headless runs
494 // fire no hooks by default. When armed, the executor is the same one the
495 // TUI builds — global config, reviewed plugin snapshots, then trusted
496 // project `.codewhale/hooks.toml` — so `tool_call_before` can still deny
497 // and `shell_env` still applies. It is shared with the engine config, the
498 // turn's SendMessage op (which re-installs it into the engine), and the
499 // tool runtime services. Fleet workers never opt in: the narrowed
500 // authority envelope does not carry the operator's hook set into a child.
501 let exec_hook_executor = (exec_hooks_enabled && !fleet_authority_active).then(|| {
502 let hooks_config = crate::hooks::HooksConfig::load_with_project_and_plugins(
503 execution_config.hooks_config(),
504 &workspace,
505 Some(engine_plugin_registry.as_ref()),
506 );
507 std::sync::Arc::new(crate::hooks::HookExecutor::new(
508 hooks_config,
509 workspace.clone(),
510 ))
511 });
512 let exec_allow_shell = crate::tools::spec::fleet_exec_shell_enabled(
513 fleet_authority_active,
514 outer_shell_authority,
515 disallowed_tools.as_deref(),
516 ) || (!fleet_authority_active
517 && (auto_approve || execution_config.allow_shell()));
518 let persist_services_enabled = cfg!(unix)
519 && !fleet_authority_active
520 && exec_allow_shell
521 && explicit_sandbox
522 .is_some_and(|sandbox| sandbox.eq_ignore_ascii_case("danger-full-access"));
523 let exec_shell_manager = crate::tools::shell::new_shared_shell_manager(workspace.clone());
524 let exec_automations = exec_automation_services(fleet_authority_active || one_shot)?;
525 let runtime_services = crate::tools::spec::RuntimeToolServices {
526 shell_manager: Some(exec_shell_manager.clone()),
527 persist_services_enabled,
528 automations: exec_automations,
529 media_originals_dir: crate::media_originals::default_store_dir(),
530 hook_executor: exec_hook_executor.clone(),
531 ..crate::tools::spec::RuntimeToolServices::default()
532 };
533
534 let engine_config = EngineConfig {
535 model: effective_model.clone(),
536 active_route_limits,
537 workspace: workspace.clone(),
538 session_id: fleet_capture.as_ref().map(|(id, _)| id.clone()),
539 subagent_state_root: None,
540 plugin_registry: Some(std::sync::Arc::clone(&engine_plugin_registry)),
541 allow_shell: exec_allow_shell,
542 trust_mode,
543 notes_path: execution_config.notes_path(),
544 mcp_config_path: execution_config.mcp_config_path(),
545 // Non-interactive exec has no user-level MCP OAuth callback
546 // overrides; the loopback default applies.
547 mcp_oauth_callback_port: None,
548 mcp_oauth_callback_url: None,
549 skills_dir: execution_config.skills_dir(),
550 skills_discovery_mode: crate::skills::SkillDiscoveryMode::from_config(
551 &execution_config.skills_config(),
552 ),
553 instructions: {
554 let mut instrs: Vec<crate::prompts::InstructionSource> = execution_config
555 .instructions_paths()
556 .into_iter()
557 .map(Into::into)
558 .collect();
559 if let Some(ref extra) = append_system_prompt {
560 instrs.push(crate::prompts::InstructionSource::Inline {
561 name: "cli:append-system-prompt".into(),
562 content: extra.clone(),
563 });
564 }
565 instrs
566 },
567 project_context_pack_enabled: execution_config.project_context_pack_enabled(),
568 translation_enabled: false,
569 max_steps: max_turns,
570 max_subagents,
571 max_admitted_subagents: execution_config
572 .max_admitted_subagents_for_provider(&effective_identity)
573 .max(max_subagents),
574 launch_concurrency: execution_config.launch_concurrency_for_provider(&effective_identity),
575 subagents_enabled: !fleet_authority_active
576 && execution_config.subagents_enabled_for_provider(&effective_identity),
577 features: engine_features,
578 auto_review_policy: execution_config.auto_review_policy(),
579 compaction: compaction.clone(),
580 todos: new_shared_todo_list(),
581 plan_state: new_shared_plan_state(),
582 goal_state: crate::tools::goal::new_shared_goal_state(),
583 max_spawn_depth: if fleet_authority_active {
584 0
585 } else {
586 execution_config.subagent_max_spawn_depth_for_provider(&effective_identity)
587 },
588 network_policy,
589 snapshots_enabled: !fleet_authority_active
590 && !one_shot
591 && execution_config.snapshots_config().enabled,
592 snapshots_max_workspace_bytes: execution_config
593 .snapshots_config()
594 .max_workspace_gb
595 .saturating_mul(1024 * 1024 * 1024),
596 // No host here records snapshot receipts.
597 record_restore_points: false,
598 lsp_config,
599 runtime_services,
600 subagent_model_overrides: execution_config.subagent_model_overrides(),
601 fleet_roster: std::sync::Arc::new(crate::fleet::identity::load_effective_roster(
602 &execution_config.fleet_config(),
603 &workspace,
604 Some(engine_plugin_registry.as_ref()),
605 )),
606 subagent_api_timeout: std::time::Duration::from_secs(
607 execution_config.subagent_api_timeout_secs_for_provider(&effective_identity),
608 ),
609 stream_chunk_timeout: std::time::Duration::from_secs(
610 execution_config.stream_chunk_timeout_secs(),
611 ),
612 turn_wall_clock: execution_config.turn_wall_clock(),
613 stream_max_content_bytes: execution_config.stream_max_content_bytes(),
614 stream_max_duration: execution_config.stream_max_duration(),
615 stream_retry_limits: execution_config.stream_retry_limits(),
616 stream_open_timeout: execution_config.stream_open_timeout(),
617 subagent_heartbeat_timeout: std::time::Duration::from_secs(
618 execution_config.subagent_heartbeat_timeout_secs_for_provider(&effective_identity),
619 ),
620 prefer_bwrap: execution_config.prefer_bwrap.unwrap_or(false),
621 bwrap_extensions: crate::sandbox::BwrapMountExtensions {
622 read_only_roots: execution_config.bwrap_ro_roots.clone(),
623 device_roots: execution_config.bwrap_dev_roots.clone(),
624 },
625 read_denylist: execution_config.read_denylist(),
626 memory_enabled: execution_config.memory_enabled(),
627 memory_path: execution_config.memory_path(),
628 speech_output_dir: execution_config.speech_output_dir(),
629 vision_config: execution_config.vision_model_config(),
630 strict_tool_mode: execution_config.strict_tool_mode.unwrap_or(false),
631 goal_objective: None,
632 goal_token_budget: None,
633 goal_status: crate::tools::goal::GoalStatus::Active,
634 goal_max_continuations: execution_config.goal_max_continuations(),
635 goal_continuation_delay_seconds: execution_config.goal_continuation_delay_seconds(),
636 goal_enforce_token_budget: execution_config.goal_enforce_token_budget(),
637 reasoning_only_max_reprompts: execution_config.reasoning_only_max_reprompts(),
638 reasoning_only_reprompt_message: Some(
639 execution_config
640 .reasoning_only_reprompt_message()
641 .to_string(),
642 ),
643 allowed_tools: allowed_tools.clone(),
644 disallowed_tools: disallowed_tools.clone(),
645 max_tool_calls,
646 hook_executor: exec_hook_executor.clone(),
647 locale_tag: codewhale_localization::resolve_locale(&settings.locale)
648 .tag()
649 .to_string(),
650 workshop: {
651 crate::tools::large_output_router::WorkshopConfig::install_active(
652 config.workshop.as_ref(),
653 );
654 config.workshop.clone()
655 },
656 search_provider: execution_config.search_provider(),
657 search_api_key: execution_config
658 .search
659 .as_ref()
660 .and_then(|s| s.api_key.clone()),
661 search_native: execution_config.search_native(),
662 search_base_url: execution_config
663 .search
664 .as_ref()
665 .and_then(|s| s.base_url.clone()),
666 tools_always_load: if fleet_authority_active {
667 std::collections::HashSet::new()
668 } else {
669 execution_config.tools_always_load()
670 },
671 user_input_limits: execution_config.user_input_limits(),
672 user_input_timeout: execution_config.user_input_timeout(),
673 goal_max_steps: None,
674 tools: if fleet_authority_active {
675 None
676 } else {
677 execution_config.tools.clone()
678 },
679 verbosity: execution_config.verbosity.clone(),
680 workspace_follow_symlinks: settings.workspace_follow_symlinks,
681 exec_policy_engine: execution_config.exec_policy_engine.clone(),
682 terminal_chrome_enabled: false,
683 advisor_config: execution_config
684 .advisor
685 .as_ref()
686 .map(crate::tools::subagent::AdvisorConfig::from_toml)
687 .unwrap_or_else(crate::tools::subagent::AdvisorConfig::disabled),
688 };
689
690 let engine_handle = spawn_engine(engine_config, &execution_config);
691 // The Full Access posture travels in the op's auto_approve/approval_mode
692 // fields; modes no longer carry permission.
693 let mode = AppMode::Agent;
694
695 let resuming_session = resume_session.is_some();
696 let mut loaded_session_id = None;
697 if let Some(saved) = resume_session {
698 let saved_id = saved.metadata.id.clone();
699 if saved.metadata.workspace != workspace && output_format == ExecOutputFormat::Text {
700 eprintln!(
701 "Warning: session {} was created in a different workspace ({}). Resuming anyway.",
702 truncate_id(&saved_id),
703 saved.metadata.workspace.display(),
704 );
705 }
706
707 engine_handle
708 .send(Op::SyncSession {
709 session_id: Some(saved_id.clone()),
710 messages: saved.messages,
711 system_prompt: saved.system_prompt.map(SystemPrompt::Text),
712 system_prompt_override: false,
713 model: saved.metadata.model,
714 workspace: saved.metadata.workspace,
715 mode,
716 })
717 .await?;
718 loaded_session_id = Some(saved_id.clone());
719 if output_format == ExecOutputFormat::Text && !json_output {
720 eprintln!("{}", exec_resumed_session_line(&saved_id));
721 }
722 }
723
724 // Lifecycle outbox (`[lifecycle_outbox]`): headless `codewhale exec`
725 // gets the same turn boundaries as the interactive TUI. Disabled
726 // (all emits no-op) when the config has no path.
727 let lifecycle_outbox = config
728 .lifecycle_outbox
729 .as_ref()
730 .map(|outbox| {
731 codewhale_hooks::LifecycleOutbox::new(
732 outbox.path.clone(),
733 outbox.webhook_url.clone(),
734 outbox.webhook_token.clone(),
735 )
736 })
737 .unwrap_or_else(codewhale_hooks::LifecycleOutbox::disabled);
738 // Wall clock for the outbox `turn_end` duration. `exec` never receives
739 // a TurnStarted engine event, so the start is marked at the same
740 // `Op::SendMessage` boundary where `turn_start` is emitted below.
741 let exec_turn_started_at = Instant::now();
742
743 engine_handle
744 .send(Op::SendMessage(TurnSpec {
745 max_output_tokens: None,
746 content: prompt.to_string(),
747 images: Vec::new(),
748 mode,
749 route: Box::new(validated_route.into_resolved()),
750 compaction: Box::new(compaction.clone()),
751 initial_routed_usage: Box::default(),
752 goal_objective: None,
753 goal_token_budget: None,
754 goal_status: crate::tools::goal::GoalStatus::Active,
755 allowed_tools: allowed_tools.clone(),
756 dynamic_tools: Vec::new(),
757 hook_executor: exec_hook_executor.clone(),
758 reasoning_effort: effective_reasoning_effort,
759 reasoning_effort_auto,
760 auto_model,
761 allow_shell: auto_approve || execution_config.allow_shell(),
762 trust_mode,
763 auto_approve,
764 translation_enabled: false,
765 approval_mode: if auto_approve {
766 ApprovalMode::Bypass
767 } else {
768 execution_config
769 .approval_policy
770 .as_deref()
771 .and_then(ApprovalMode::from_config_value)
772 .unwrap_or_default()
773 },
774 verbosity: execution_config.verbosity.clone(),
775 provenance: crate::core::ops::UserInputProvenance::ExternalUser,
776 // Headless exec does not correlate submissions.
777 submission_id: None,
778 }))
779 .await?;
780
781 // Lifecycle outbox: the clean headless turn-start boundary. `exec` has
782 // no TurnStarted engine event; the message submission above is exactly
783 // where the engine begins the turn. No-op when the feature is disabled.
784 lifecycle_outbox.emit(codewhale_hooks::LifecycleEvent {
785 event: "turn_start".to_string(),
786 kind: "turn.started".to_string(),
787 thread_id: loaded_session_id.clone().unwrap_or_default(),
788 turn_id: None,
789 item_id: None,
790 payload: serde_json::json!({
791 "model": codewhale_hooks::bounded_text(
792 &effective_model,
793 codewhale_hooks::OUTBOX_DETAIL_MAX_CHARS,
794 ),
795 "workspace": workspace.display().to_string(),
796 }),
797 });
798
799 let mut summary = ExecSummary {
800 mode: if one_shot { "one-shot" } else { "agent" }.to_string(),
801 provider: effective_provider_name.clone(),
802 model: effective_model.clone(),
803 prompt: prompt.to_string(),
804 ..ExecSummary::default()
805 };
806 let can_elevate_sandbox =
807 exec_sandbox_elevation_authorized(allow_sandbox_elevation, explicit_sandbox);
808 let mut sandbox_denied = false;
809 let mut approval_required = false;
810 let mut tool_error_seen = false;
811 let mut last_error_category = None;
812 let mut reported_sandbox_contract = false;
813
814 let mut should_persist_session =
815 resuming_session || output_format == ExecOutputFormat::StreamJson;
816 let mut latest_session_id = loaded_session_id;
817 let mut latest_messages: Arc<Vec<Message>> = Arc::new(Vec::new());
818 let mut latest_system_prompt: Option<SystemPrompt> = None;
819 let mut latest_model = effective_model;
820 let mut latest_workspace = workspace.clone();
821 let mut tool_starts: HashMap<String, (Instant, String)> = HashMap::new();
822 let mut turn_usage_seq: u32 = 0;
823 // None means no actual terminal request snapshot was observed. A known
824 // zero must remain distinguishable from that missing receipt.
825 let mut observed_retry_count = None;
826 let mut settled_usage: Option<codewhale_models::Usage> = None;
827
828 let mut ends_with_newline = false;
829 // One absolute host deadline includes every autonomous child fan-in turn;
830 // child-specific shorter deadlines remain enforced by their runtime.
831 // The default wall clock is unbounded (`Duration::MAX`); a century
832 // stands in for "never" without overflowing `Instant`.
833 let exec_deadline = exec_turn_started_at
834 .checked_add(execution_config.turn_wall_clock())
835 .unwrap_or_else(|| exec_turn_started_at + Duration::from_secs(100 * 365 * 86_400));
836 let mut events = ExecAgentEvents::new(engine_handle.clone(), exec_deadline);
837 loop {
838 let Some(event) = events.next().await else {
839 break;
840 };
841
842 match event {
843 Event::MessageDelta { content, .. } => {
844 summary.output.push_str(&content);
845 if output_format == ExecOutputFormat::StreamJson {
846 emit_exec_stream_event(&ExecStreamEvent::Content { content })?;
847 } else if !json_output {
848 write_exec_stdout(&content)?;
849 }
850 ends_with_newline = summary.output.ends_with('\n');
851 }
852 Event::MessageComplete { .. }
853 if output_format == ExecOutputFormat::Text
854 && !json_output
855 && !ends_with_newline =>
856 {
857 write_exec_stdout("\n")?;
858 }
859 Event::ThinkingDelta { .. } => {
860 // Exec stream-json intentionally omits reasoning deltas; the
861 // TUI transcript retains its existing Activity Detail surface.
862 }
863 Event::ToolProjectionWarning {
864 provider,
865 omitted_tool_names,
866 omitted_tool_count,
867 } if !json_output => {
868 eprintln!(
869 "{}",
870 crate::core::events::tool_projection_warning_message(
871 &provider,
872 &omitted_tool_names,
873 omitted_tool_count,
874 )
875 );
876 }
877 Event::ToolCallStarted {
878 id, name, input, ..
879 } => {
880 let started_at = chrono::Utc::now().to_rfc3339();
881 tool_starts.insert(id.clone(), (Instant::now(), started_at.clone()));
882 if output_format == ExecOutputFormat::StreamJson {
883 emit_exec_stream_event(&ExecStreamEvent::ToolUse {
884 name,
885 id,
886 input,
887 started_at,
888 })?;
889 } else if !json_output {
890 let summary = summarize_tool_args(&input);
891 if let Some(summary) = summary {
892 eprintln!("tool: {name} ({summary})");
893 } else {
894 eprintln!("tool: {name}");
895 }
896 }
897 }
898 Event::ToolCallComplete {
899 id, name, result, ..
900 } => {
901 let (duration_ms, started_at) = tool_starts
902 .remove(&id)
903 .map(|(started, timestamp)| {
904 (
905 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
906 timestamp,
907 )
908 })
909 .unwrap_or_else(|| (0, chrono::Utc::now().to_rfc3339()));
910 let receipt_name = name.clone();
911 match result {
912 Ok(output) => {
913 tool_error_seen |= !output.success;
914 summary.tools.push(ExecToolEntry {
915 name: name.clone(),
916 success: output.success,
917 output: output.content.clone(),
918 });
919 if output_format == ExecOutputFormat::StreamJson {
920 emit_exec_stream_event(&ExecStreamEvent::ToolResult {
921 id,
922 name: receipt_name,
923 output: output.content,
924 status: if output.success {
925 "success".to_string()
926 } else {
927 "error".to_string()
928 },
929 started_at,
930 completed_at: chrono::Utc::now().to_rfc3339(),
931 duration_ms,
932 side_effect_status: output
933 .metadata
934 .as_ref()
935 .and_then(|metadata| metadata.get("side_effect_status"))
936 .and_then(serde_json::Value::as_str)
937 .unwrap_or("unknown")
938 .to_string(),
939 error_category: (!output.success).then(|| {
940 output
941 .metadata
942 .as_ref()
943 .and_then(|metadata| metadata.get("error_category"))
944 .and_then(serde_json::Value::as_str)
945 .unwrap_or("tool_reported_failure")
946 .to_string()
947 }),
948 truncated: output
949 .metadata
950 .as_ref()
951 .and_then(|metadata| metadata.get("truncated"))
952 .and_then(serde_json::Value::as_bool),
953 artifact: tool_artifact_receipt(output.metadata.as_ref()),
954 result_metadata: output.metadata,
955 })?;
956 } else if !json_output {
957 if name == "exec_shell" && !output.content.trim().is_empty() {
958 eprintln!("tool {name} completed");
959 eprintln!(
960 "--- stdout/stderr ---\n{}\n---------------------",
961 output.content
962 );
963 } else {
964 eprintln!(
965 "tool {name} completed: {}",
966 summarize_tool_output(&output.content)
967 );
968 }
969 }
970 }
971 Err(err) => {
972 tool_error_seen = true;
973 let error_text = err.to_string();
974 summary.tools.push(ExecToolEntry {
975 name: name.clone(),
976 success: false,
977 output: error_text.clone(),
978 });
979 if output_format == ExecOutputFormat::StreamJson {
980 emit_exec_stream_event(&ExecStreamEvent::ToolResult {
981 id,
982 name: receipt_name,
983 output: error_text,
984 status: "error".to_string(),
985 started_at,
986 completed_at: chrono::Utc::now().to_rfc3339(),
987 duration_ms,
988 side_effect_status: "not_started_or_unknown".to_string(),
989 error_category: Some(tool_error_receipt_category(&err).to_string()),
990 truncated: None,
991 artifact: None,
992 result_metadata: None,
993 })?;
994 } else if !json_output {
995 eprintln!("tool {name} failed: {err}");
996 }
997 }
998 }
999 }
1000 Event::AgentSpawned { id, prompt, .. }
1001 if output_format == ExecOutputFormat::Text && !json_output =>
1002 {
1003 eprintln!("sub-agent {id} spawned: {}", summarize_tool_output(&prompt));
1004 }
1005 Event::AgentProgress { id, status, .. }
1006 if output_format == ExecOutputFormat::Text && !json_output =>
1007 {
1008 eprintln!("sub-agent {id}: {status}");
1009 }
1010 Event::AgentComplete {
1011 id,
1012 result,
1013 outcome,
1014 ..
1015 } if output_format == ExecOutputFormat::Text && !json_output => {
1016 eprintln!(
1017 "sub-agent {id} {}: {}",
1018 outcome
1019 .as_ref()
1020 .map(crate::tools::subagent::subagent_status_name)
1021 .unwrap_or("settled (outcome unconfirmed)"),
1022 summarize_tool_output(&result)
1023 );
1024 }
1025 Event::AgentSpawned {
1026 id,
1027 parent_run_id,
1028 spawn_depth,
1029 model,
1030 route_source,
1031 ..
1032 } if output_format == ExecOutputFormat::StreamJson => {
1033 emit_exec_stream_event(&ExecStreamEvent::AgentSpawned {
1034 id,
1035 model,
1036 spawn_depth,
1037 parent_run_id,
1038 route_source,
1039 })?;
1040 }
1041 Event::AgentSpawned { .. }
1042 | Event::AgentProgress { .. }
1043 | Event::AgentComplete { .. } => {}
1044 Event::WorkflowUi { run_id, event, .. }
1045 if output_format == ExecOutputFormat::StreamJson =>
1046 {
1047 emit_exec_stream_event(&ExecStreamEvent::WorkflowEvent { run_id, event })?;
1048 }
1049 // Headless runs have no person at the prompt: the run's flags
1050 // (the posture) answer every request.
1051 Event::ApprovalRequired {
1052 id,
1053 approval_force_prompt,
1054 ..
1055 } => {
1056 // An exact user decision (including extension-sourced shell
1057 // and network calls) cannot be supplied by a headless posture.
1058 if auto_approve && !approval_force_prompt {
1059 let _ = engine_handle
1060 .approve_tool_call_by(id, crate::approval_log::ApprovalDecider::Posture)
1061 .await;
1062 } else {
1063 approval_required = true;
1064 let _ = engine_handle
1065 .deny_tool_call_by(id, crate::approval_log::ApprovalDecider::Posture)
1066 .await;
1067 }
1068 }
1069 Event::ElevationRequired {
1070 tool_id,
1071 tool_name,
1072 denial_reason,
1073 ..
1074 } => {
1075 if can_elevate_sandbox {
1076 let policy = crate::sandbox::SandboxPolicy::DangerFullAccess;
1077 let _ = engine_handle
1078 .retry_tool_with_policy_by(
1079 tool_id,
1080 policy,
1081 crate::approval_log::ApprovalDecider::Posture,
1082 )
1083 .await;
1084 } else {
1085 sandbox_denied = true;
1086 approval_required = true;
1087 summary.outcomes.push(ExecOutcome {
1088 kind: "sandbox_denied".to_string(),
1089 outcome: "approval_required".to_string(),
1090 tool_name: tool_name.clone(),
1091 reason: denial_reason.clone(),
1092 });
1093 if !reported_sandbox_contract {
1094 eprintln!(
1095 "sandbox denied {tool_name}: {denial_reason}; --auto approves tools but does not elevate sandbox access — use --sandbox danger-full-access or --allow-sandbox-elevation to opt in"
1096 );
1097 reported_sandbox_contract = true;
1098 }
1099 if output_format == ExecOutputFormat::StreamJson {
1100 emit_exec_stream_event(&ExecStreamEvent::SandboxDenied {
1101 tool_id: tool_id.clone(),
1102 tool_name,
1103 reason: denial_reason,
1104 outcome: "approval_required".to_string(),
1105 })?;
1106 }
1107 let _ = engine_handle
1108 .deny_tool_call_by(tool_id, crate::approval_log::ApprovalDecider::Posture)
1109 .await;
1110 }
1111 }
1112 Event::Error {
1113 envelope,
1114 recoverable: _,
1115 } => {
1116 // Only a non-recoverable envelope may force the run summary
1117 // into failure. Recoverable warnings (stream-stall notices,
1118 // transient retry noise) are still streamed for visibility,
1119 // but the terminal TurnComplete event carries the
1120 // authoritative turn outcome — letting a warning set
1121 // `summary.error` here would exit an otherwise-successful
1122 // `exec` run non-zero.
1123 if exec_error_event_is_fatal(&envelope) {
1124 last_error_category = Some(envelope.category);
1125 summary.error_category = Some(envelope.category.to_string());
1126 summary.error = Some(envelope.message.clone());
1127 }
1128 if output_format == ExecOutputFormat::StreamJson {
1129 emit_exec_stream_event(&ExecStreamEvent::Error {
1130 error: envelope.message,
1131 })?;
1132 } else if !json_output {
1133 eprintln!("error: {}", envelope.message);
1134 }
1135 }
1136 Event::TurnUsage {
1137 usage, duration_ms, ..
1138 } => {
1139 if output_format == ExecOutputFormat::StreamJson {
1140 turn_usage_seq = turn_usage_seq.saturating_add(1);
1141 emit_exec_stream_event(&ExecStreamEvent::TurnUsage {
1142 turn: turn_usage_seq,
1143 input_tokens: usage.input_tokens,
1144 output_tokens: usage.output_tokens,
1145 reasoning_tokens: usage.reasoning_tokens,
1146 prompt_cache_hit_tokens: usage.prompt_cache_hit_tokens,
1147 prompt_cache_miss_tokens: usage.prompt_cache_miss_tokens,
1148 prompt_cache_write_tokens: usage.prompt_cache_write_tokens,
1149 reasoning_replay_tokens: usage.reasoning_replay_tokens,
1150 duration_ms,
1151 })?;
1152 }
1153 }
1154 Event::TurnComplete {
1155 status,
1156 error,
1157 usage,
1158 tool_catalog,
1159 ..
1160 } => {
1161 let (terminal_status, terminal_error) = (status, error);
1162 settled_usage = Some(usage.clone());
1163 #[cfg(unix)]
1164 let (mut terminal_status, mut terminal_error) = (terminal_status, terminal_error);
1165 if matches!(
1166 terminal_status,
1167 crate::core::events::TurnOutcomeStatus::Completed
1168 ) && terminal_error.is_none()
1169 {
1170 #[cfg(unix)]
1171 match exec_shell_manager.lock() {
1172 Ok(mut manager) => match manager.commit_persistent_services() {
1173 Ok(receipts) => {
1174 for receipt in &receipts {
1175 if output_format == ExecOutputFormat::StreamJson {
1176 emit_exec_stream_event(
1177 &ExecStreamEvent::ServiceReleased {
1178 task_id: receipt.task_id.clone(),
1179 pid: receipt.pid,
1180 process_group_id: receipt.process_group_id,
1181 ownership: receipt.ownership.clone(),
1182 },
1183 )?;
1184 } else if !json_output {
1185 eprintln!(
1186 "persistent service released: {} pid={} pgid={} ownership={}",
1187 receipt.task_id,
1188 receipt.pid,
1189 receipt.process_group_id,
1190 receipt.ownership
1191 );
1192 }
1193 }
1194 summary.released_services.extend(receipts);
1195 }
1196 Err(error) => {
1197 manager.abort_persistent_services();
1198 terminal_status = crate::core::events::TurnOutcomeStatus::Failed;
1199 terminal_error = Some(format!(
1200 "Persistent service ownership transfer failed: {error}"
1201 ));
1202 }
1203 },
1204 Err(_) => {
1205 terminal_status = crate::core::events::TurnOutcomeStatus::Failed;
1206 terminal_error = Some(
1207 "Persistent service ownership transfer failed: shell manager lock poisoned"
1208 .to_string(),
1209 );
1210 }
1211 }
1212 } else if let Ok(mut manager) = exec_shell_manager.lock() {
1213 manager.abort_persistent_services();
1214 }
1215 summary.status = Some(format!("{terminal_status:?}").to_lowercase());
1216 if terminal_error.is_some() {
1217 summary.error = terminal_error;
1218 }
1219 if sandbox_denied
1220 && summary.error.is_none()
1221 && matches!(
1222 terminal_status,
1223 crate::core::events::TurnOutcomeStatus::Failed
1224 )
1225 {
1226 summary.error = Some(
1227 "exec turn failed after sandbox denial; explicit sandbox elevation was not authorized"
1228 .to_string(),
1229 );
1230 }
1231 // Lifecycle outbox: the clean headless turn-end boundary.
1232 // `terminal_status` is authoritative here — persistent-service
1233 // handoff failures above already demoted it to Failed, and
1234 // `summary.error` includes the sandbox-denial augmentation.
1235 // No-op when the feature is disabled.
1236 {
1237 let outbox_status = format!("{terminal_status:?}").to_lowercase();
1238 let kind = match terminal_status {
1239 crate::core::events::TurnOutcomeStatus::Completed => "turn.completed",
1240 crate::core::events::TurnOutcomeStatus::Failed => "turn.failed",
1241 crate::core::events::TurnOutcomeStatus::Interrupted => "turn.interrupted",
1242 };
1243 lifecycle_outbox.emit(codewhale_hooks::LifecycleEvent {
1244 event: "turn_end".to_string(),
1245 kind: kind.to_string(),
1246 thread_id: latest_session_id.clone().unwrap_or_default(),
1247 turn_id: None,
1248 item_id: None,
1249 payload: serde_json::json!({
1250 "status": outbox_status,
1251 "duration_ms": exec_turn_started_at.elapsed().as_millis() as u64,
1252 "workspace": latest_workspace.display().to_string(),
1253 "error": summary.error.as_deref().map(|message| {
1254 codewhale_hooks::bounded_text(
1255 message,
1256 codewhale_hooks::OUTBOX_DETAIL_MAX_CHARS,
1257 )
1258 }),
1259 }),
1260 });
1261 }
1262 if last_error_category.is_none() {
1263 last_error_category = summary
1264 .error
1265 .as_deref()
1266 .map(crate::error_taxonomy::classify_error_message);
1267 summary.error_category =
1268 last_error_category.map(|category| category.to_string());
1269 }
1270 let termination_reason = crate::core::termination::classify_turn_termination(
1271 terminal_status,
1272 last_error_category,
1273 tool_error_seen,
1274 approval_required,
1275 );
1276 summary.termination_reason = Some(termination_reason.as_str().to_string());
1277 // State the exit class here rather than inferring it later
1278 // from the process exit code: `Canceled` exits 130, the same
1279 // value the SIGINT path uses, so a code-based derivation would
1280 // report every Esc-cancelled turn as a signal. A no-op unless
1281 // this process was armed.
1282 if !termination_reason.is_success() {
1283 codewhale_telemetry::set_exit_class(codewhale_telemetry::ExitClass::Error);
1284 }
1285 let saved_session_id = if should_persist_session && !latest_messages.is_empty() {
1286 match persist_exec_session(
1287 &latest_messages,
1288 &latest_model,
1289 PersistedProviderRoute {
1290 kind: effective_identity.persisted_kind(),
1291 id: effective_identity.persisted_id(),
1292 },
1293 &latest_workspace,
1294 &latest_system_prompt,
1295 latest_session_id.as_deref(),
1296 u64::from(usage.input_tokens) + u64::from(usage.output_tokens),
1297 fleet_capture.as_ref().map(|(_, manager)| manager),
1298 ) {
1299 Ok(id) => {
1300 if output_format == ExecOutputFormat::Text && !json_output {
1301 eprintln!("{}", exec_saved_session_line(&id));
1302 }
1303 Some(id)
1304 }
1305 Err(err) => {
1306 if output_format == ExecOutputFormat::Text && !json_output {
1307 eprintln!("warning: failed to save exec session: {err}");
1308 }
1309 None
1310 }
1311 }
1312 } else {
1313 None
1314 };
1315 if output_format == ExecOutputFormat::StreamJson {
1316 if let Some(id) = saved_session_id.as_ref() {
1317 emit_exec_stream_event(&ExecStreamEvent::SessionCapture {
1318 content: exec_stream_session_ref(id),
1319 saved_session_id: id.clone(),
1320 })?;
1321 }
1322 // Resolved output ceiling and its provenance, surfaced so a
1323 // wrong ceiling is visible in the receipt rather than
1324 // requiring packet capture.
1325 let codewhale_max_output_tokens =
1326 crate::route_budget::effective_max_output_tokens_for_route(
1327 effective_provider,
1328 &latest_model,
1329 active_route_limits,
1330 );
1331 let codewhale_max_output_tokens_source =
1332 crate::route_budget::output_ceiling_source(
1333 effective_provider,
1334 &latest_model,
1335 )
1336 .as_str();
1337 // The deliverable is the final assistant reply of the
1338 // session, not the cumulative stream output: a
1339 // multi-step turn streams pre-tool commentary first,
1340 // and that commentary is not part of the answer.
1341 let final_answer = exec_stream_final_answer_text(
1342 &latest_messages,
1343 !summary.output.trim().is_empty(),
1344 )
1345 .unwrap_or_default();
1346 emit_exec_stream_event(&ExecStreamEvent::Metadata {
1347 meta: Box::new(ExecStreamMeta {
1348 receipt_kind: "terminal",
1349 provider: effective_provider_kind.clone(),
1350 provider_id: effective_stream_provider_id.clone(),
1351 model: latest_model.clone(),
1352 route_source: route_source.clone(),
1353 input_tokens: Some(usage.input_tokens),
1354 output_tokens: Some(usage.output_tokens),
1355 prompt_cache_hit_tokens: usage.prompt_cache_hit_tokens,
1356 prompt_cache_miss_tokens: usage.prompt_cache_miss_tokens,
1357 prompt_cache_write_tokens: usage.prompt_cache_write_tokens,
1358 reasoning_tokens: usage.reasoning_tokens,
1359 codewhale_max_output_tokens: Some(codewhale_max_output_tokens),
1360 codewhale_max_output_tokens_source: Some(
1361 codewhale_max_output_tokens_source,
1362 ),
1363 duration_ms: u64::try_from(exec_started.elapsed().as_millis())
1364 .unwrap_or(u64::MAX),
1365 retry_count: observed_retry_count,
1366 approval_posture: approval_posture.clone(),
1367 sandbox_posture: sandbox_posture.clone(),
1368 binary_sha256: binary_sha256.clone(),
1369 config_sha256: None,
1370 prompt_sha256: prompt_sha256.clone(),
1371 tool_catalog_sha256: tool_catalog.as_ref().and_then(|catalog| {
1372 serde_json::to_vec(catalog).ok().map(|bytes| {
1373 format!("sha256:{}", crate::hashing::sha256_hex(&bytes))
1374 })
1375 }),
1376 input_analysis: exec_stream_input_analysis(
1377 &latest_messages,
1378 latest_system_prompt.as_ref(),
1379 ),
1380 visible_final_answer_chars: final_answer.chars().count(),
1381 visible_final_answer_excerpt: exec_stream_final_answer_excerpt(
1382 &final_answer,
1383 ),
1384 resume_command: saved_session_id
1385 .as_deref()
1386 .map(exec_stream_resume_hint)
1387 .unwrap_or_default(),
1388 session_id: saved_session_id
1389 .as_deref()
1390 .map(exec_stream_session_ref)
1391 .unwrap_or_default(),
1392 workspace: latest_workspace.display().to_string(),
1393 message_count: latest_messages.len(),
1394 status: summary.status.clone(),
1395 termination_reason: summary.termination_reason.clone(),
1396 error_category: summary.error_category.clone(),
1397 error: summary.error.clone(),
1398 }),
1399 })?;
1400 emit_exec_stream_event(&ExecStreamEvent::Done)?;
1401 }
1402 let _ =
1403 tokio::time::timeout(Duration::from_secs(2), engine_handle.send(Op::Shutdown))
1404 .await;
1405 break;
1406 }
1407 Event::CompactionStarted { .. } => {
1408 // The Engine writes recovery artifacts under its session ID.
1409 // Keep the owning session discoverable even in text output.
1410 should_persist_session = true;
1411 }
1412 Event::SessionUpdated {
1413 session_id,
1414 messages,
1415 system_prompt,
1416 model,
1417 workspace,
1418 } => {
1419 latest_session_id = Some(session_id);
1420 latest_messages = messages;
1421 latest_system_prompt = system_prompt;
1422 latest_model = model;
1423 latest_workspace = workspace;
1424 }
1425 // A tool-less one-shot run has no tool channel, so a model that
1426 // still tries to call a tool writes the call as text. The engine
1427 // strips that markup from the answer; say why the answer is short
1428 // and how to give the model tools, instead of exiting on nothing.
1429 Event::Status { message }
1430 if one_shot && message == crate::core::engine::FAKE_WRAPPER_NOTICE =>
1431 {
1432 eprintln!("{ONE_SHOT_TOOL_CALL_NOTICE}");
1433 }
1434 // #3027: surface the engine's max-steps notice in text mode so a
1435 // --max-turns run that stops early says why instead of going quiet.
1436 Event::Status { message }
1437 if output_format == ExecOutputFormat::Text
1438 && !json_output
1439 && message.contains("Maximum model steps") =>
1440 {
1441 eprintln!("{message}");
1442 }
1443 Event::ToolRequestSnapshot { snapshot } => {
1444 observed_retry_count =
1445 accumulate_exec_retry_count(observed_retry_count, snapshot.terminal.as_ref());
1446 }
1447 Event::Status { message } => {
1448 if let Some(receipt) = exec_retry_status(&message) {
1449 if output_format == ExecOutputFormat::StreamJson {
1450 emit_exec_stream_event(&receipt)?;
1451 } else if output_format == ExecOutputFormat::Text && !json_output {
1452 eprintln!("{message}");
1453 }
1454 }
1455 }
1456 _ => {}
1457 }
1458 }
1459
1460 if summary.status.is_none() {
1461 if let Ok(mut manager) = exec_shell_manager.lock() {
1462 manager.abort_persistent_services();
1463 }
1464 let error = summary.error.clone().unwrap_or_else(|| {
1465 "Engine event channel closed before a terminal turn receipt".to_string()
1466 });
1467 let category = last_error_category
1468 .unwrap_or_else(|| crate::error_taxonomy::classify_error_message(&error));
1469 let termination_reason = crate::core::termination::classify_turn_termination(
1470 crate::core::events::TurnOutcomeStatus::Failed,
1471 Some(category),
1472 tool_error_seen,
1473 approval_required,
1474 );
1475 summary.status = Some("failed".to_string());
1476 summary.error_category = Some(category.to_string());
1477 summary.termination_reason = Some(termination_reason.as_str().to_string());
1478 summary.error = Some(error.clone());
1479 // Lifecycle outbox: the engine channel closed before a terminal
1480 // turn receipt. Every emitted `turn_start` still gets its matching
1481 // `turn_end` so a supervisor never sees an orphaned in-progress
1482 // turn. No-op when the feature is disabled.
1483 lifecycle_outbox.emit(codewhale_hooks::LifecycleEvent {
1484 event: "turn_end".to_string(),
1485 kind: "turn.failed".to_string(),
1486 thread_id: latest_session_id.clone().unwrap_or_default(),
1487 turn_id: None,
1488 item_id: None,
1489 payload: serde_json::json!({
1490 "status": "failed",
1491 "duration_ms": exec_turn_started_at.elapsed().as_millis() as u64,
1492 "workspace": latest_workspace.display().to_string(),
1493 "error": codewhale_hooks::bounded_text(
1494 &error,
1495 codewhale_hooks::OUTBOX_DETAIL_MAX_CHARS,
1496 ),
1497 }),
1498 });
1499 if output_format == ExecOutputFormat::StreamJson {
1500 emit_exec_stream_event(&ExecStreamEvent::Error { error })?;
1501 }
1502 }
1503
1504 // Drain the terminal receipt before either returning or taking the explicit
1505 // retryable-failure process exit below. Outbox failures cannot change the
1506 // authoritative turn outcome.
1507 if let Err(error) = lifecycle_outbox.flush(Duration::from_secs(2)).await {
1508 tracing::warn!(target: "lifecycle_outbox", %error, "exec lifecycle outbox did not drain before exit");
1509 }
1510
1511 if one_shot {
1512 summary.record_one_shot_outcome(settled_usage);
1513 }
1514 if json_output {
1515 write_exec_stdout(&format!("{}\n", serde_json::to_string_pretty(&summary)?))?;
1516 }
1517
1518 if let Some(error) = summary.error.as_ref()
1519 && !error.trim().is_empty()
1520 {
1521 // Distinguish retryable infrastructure failures (provider/transport,
1522 // after all in-session retries are exhausted) from genuine task
1523 // failures so supervisors and bench harnesses can tell them apart at
1524 // the process level without parsing the stream. Genuine failures
1525 // keep the historical `bail!` → exit 1 path.
1526 let exit_code = exec_failure_exit_code(summary.error_category.as_deref());
1527 // The final line always carries the message: automation greps it and
1528 // a caller may keep only the last stderr line, even when the stream
1529 // already printed the same error above.
1530 if exit_code != 1 {
1531 eprintln!("Error: exec turn failed: {error}");
1532 let _ = io::stdout().flush();
1533 std::process::exit(exit_code);
1534 }
1535 bail!("exec turn failed: {error}");
1536 }
1537
1538 if matches!(
1539 summary.status.as_deref(),
1540 Some("failed" | "canceled" | "interrupted")
1541 ) {
1542 let status = summary.status.as_deref().unwrap_or("unknown");
1543 bail!("exec turn ended with status {status}");
1544 }
1545
1546 Ok(())
1547 }
1548
1549 fn exec_retry_status(message: &str) -> Option<ExecStreamEvent> {
1550 crate::core::events::is_retry_status_receipt(message).then(|| ExecStreamEvent::Status {
1551 message: message.to_string(),
1552 })
1553 }
1554
1555 fn accumulate_exec_retry_count(
1556 previous: Option<u32>,
1557 terminal: Option<&crate::tool_inspection::TurnStopDiagnostics>,
1558 ) -> Option<u32> {
1559 let Some(terminal) = terminal.filter(|facts| facts.status.is_some()) else {
1560 return previous;
1561 };
1562 let retries = terminal
1563 .transport_retries
1564 .saturating_add(terminal.transparent_stream_retries)
1565 .saturating_add(terminal.stream_resumes)
1566 .saturating_add(terminal.empty_stop_retries)
1567 .saturating_add(terminal.reasoning_only_reprompts);
1568 Some(previous.unwrap_or(0).saturating_add(retries))
1569 }
1570
1571 #[cfg(test)]
1572 mod tests {
1573 use super::{ExecAgentEvents, exec_automation_services, exec_disallowed_tools};
1574
1575 use crate::core::engine::mock_engine_handle;
1576 use crate::core::engine::tool_catalog::REQUEST_USER_INPUT_NAME;
1577 use crate::core::events::{Event, TurnOutcomeStatus};
1578 use crate::core::ops::{Op, SubAgentSettlement};
1579 use codewhale_models::Usage;
1580 use std::time::{Duration, Instant};
1581
1582 fn completed_parent(input_tokens: u32) -> Event {
1583 let usage = Usage {
1584 input_tokens,
1585 ..Usage::default()
1586 };
1587 Event::TurnComplete {
1588 usage: usage.clone(),
1589 parent_route_usage: usage,
1590 routed_usage_dropped_records: 0,
1591 status: TurnOutcomeStatus::Completed,
1592 error: None,
1593 tool_catalog: None,
1594 base_url: None,
1595 }
1596 }
1597
1598 async fn reply_to_probe(
1599 operations: &mut tokio::sync::mpsc::Receiver<Op>,
1600 snapshot: SubAgentSettlement,
1601 ) {
1602 let Op::GetSubAgentSettlement { tx } = operations.recv().await.expect("host probe") else {
1603 panic!("host must not shut down while child work remains");
1604 };
1605 tx.lock().unwrap().take().unwrap().send(snapshot).unwrap();
1606 }
1607
1608 #[test]
1609 fn headless_exec_withholds_request_user_input_without_a_responder() {
1610 // No responder exists on a one-shot CLI run, so the tool is
1611 // withheld by default rather than offered and stalled on.
1612 let disallowed = exec_disallowed_tools(None).expect("withhold list");
1613 assert!(
1614 disallowed
1615 .iter()
1616 .any(|tool| tool.as_str() == REQUEST_USER_INPUT_NAME),
1617 "request_user_input must be withheld by default: {disallowed:?}"
1618 );
1619 // An operator-passed entry is kept exactly once, not duplicated.
1620 let disallowed = exec_disallowed_tools(Some(vec![REQUEST_USER_INPUT_NAME.to_string()]))
1621 .expect("withhold list");
1622 assert_eq!(
1623 disallowed
1624 .iter()
1625 .filter(|tool| tool.as_str() == REQUEST_USER_INPUT_NAME)
1626 .count(),
1627 1
1628 );
1629 }
1630
1631 #[tokio::test]
1632 async fn headless_success_waits_for_children_workflow_phases_and_parent_fan_in() {
1633 let mut engine = mock_engine_handle();
1634 let mut events = ExecAgentEvents::new(
1635 engine.handle.clone(),
1636 Instant::now() + Duration::from_secs(3),
1637 );
1638 engine.tx_event.send(completed_parent(11)).await.unwrap();
1639 let actor = async {
1640 reply_to_probe(
1641 &mut engine.rx_op,
1642 SubAgentSettlement {
1643 running_children: 1,
1644 running_workflows: 1,
1645 pending_completions: 0,
1646 },
1647 )
1648 .await;
1649 reply_to_probe(
1650 &mut engine.rx_op,
1651 SubAgentSettlement {
1652 running_children: 0,
1653 running_workflows: 1,
1654 pending_completions: 0,
1655 },
1656 )
1657 .await;
1658 reply_to_probe(
1659 &mut engine.rx_op,
1660 SubAgentSettlement {
1661 running_children: 0,
1662 running_workflows: 0,
1663 pending_completions: 1,
1664 },
1665 )
1666 .await;
1667 engine
1668 .tx_event
1669 .send(Event::MessageDelta {
1670 content: "child findings reviewed".into(),
1671 index: 0,
1672 })
1673 .await
1674 .unwrap();
1675 engine.tx_event.send(completed_parent(7)).await.unwrap();
1676 reply_to_probe(&mut engine.rx_op, SubAgentSettlement::default()).await;
1677 };
1678 let host = async {
1679 assert!(
1680 matches!(events.next().await, Some(Event::MessageDelta { content, .. }) if content == "child findings reviewed")
1681 );
1682 let Some(Event::TurnComplete { usage, status, .. }) = events.next().await else {
1683 panic!("settled parent receipt")
1684 };
1685 assert_eq!(status, TurnOutcomeStatus::Completed);
1686 assert_eq!(
1687 usage.input_tokens, 18,
1688 "both parent turns are accounted once"
1689 );
1690 };
1691 tokio::time::timeout(Duration::from_secs(4), async { tokio::join!(actor, host) })
1692 .await
1693 .unwrap();
1694 assert!(
1695 !engine.handle.is_cancelled(),
1696 "ordinary success must not cancel children"
1697 );
1698 assert!(
1699 engine.rx_op.try_recv().is_err(),
1700 "the event reader does not send early Shutdown"
1701 );
1702 }
1703
1704 #[tokio::test]
1705 async fn headless_child_settlement_deadline_bounds_a_stalled_engine_probe() {
1706 let mut engine = mock_engine_handle();
1707 let mut events = ExecAgentEvents::new(
1708 engine.handle.clone(),
1709 Instant::now() + Duration::from_millis(30),
1710 );
1711 engine.tx_event.send(completed_parent(13)).await.unwrap();
1712 let event = tokio::time::timeout(Duration::from_secs(1), events.next())
1713 .await
1714 .unwrap();
1715 let Some(Event::TurnComplete {
1716 status,
1717 error,
1718 usage,
1719 routed_usage_dropped_records,
1720 ..
1721 }) = event
1722 else {
1723 panic!("bounded failure receipt")
1724 };
1725 assert_eq!(status, TurnOutcomeStatus::Failed);
1726 assert!(error.unwrap().contains("wall-clock budget exhausted"));
1727 assert_eq!(usage.input_tokens, 13);
1728 assert_eq!(routed_usage_dropped_records, 1);
1729 assert!(engine.handle.is_cancelled());
1730 let mut shutdown = false;
1731 while let Ok(op) = engine.rx_op.try_recv() {
1732 shutdown |= matches!(op, Op::Shutdown);
1733 }
1734 assert!(shutdown, "shutdown must also cancel detached session tasks");
1735 }
1736
1737 #[tokio::test]
1738 async fn headless_failed_or_interrupted_parent_skips_child_settlement() {
1739 for terminal_status in [TurnOutcomeStatus::Failed, TurnOutcomeStatus::Interrupted] {
1740 let mut engine = mock_engine_handle();
1741 let mut events = ExecAgentEvents::new(
1742 engine.handle.clone(),
1743 Instant::now() + Duration::from_secs(30),
1744 );
1745 let mut terminal = completed_parent(3);
1746 if let Event::TurnComplete { status, .. } = &mut terminal {
1747 *status = terminal_status;
1748 }
1749 engine.tx_event.send(terminal).await.unwrap();
1750 assert!(
1751 matches!(events.next().await, Some(Event::TurnComplete { status, .. }) if status == terminal_status)
1752 );
1753 assert!(
1754 engine.rx_op.try_recv().is_err(),
1755 "failure cannot admit another settling turn"
1756 );
1757 }
1758 }
1759
1760 #[tokio::test]
1761 async fn headless_fatal_fan_in_error_cancels_before_releasing_terminal_receipt() {
1762 let engine = mock_engine_handle();
1763 let mut events = ExecAgentEvents::new(
1764 engine.handle.clone(),
1765 Instant::now() + Duration::from_secs(30),
1766 );
1767 engine.tx_event.send(completed_parent(5)).await.unwrap();
1768 engine
1769 .tx_event
1770 .send(Event::error(crate::error_taxonomy::ErrorEnvelope::fatal(
1771 "fan-in route unavailable",
1772 )))
1773 .await
1774 .unwrap();
1775 assert!(matches!(events.next().await, Some(Event::Error { .. })));
1776 assert!(engine.handle.is_cancelled());
1777 assert!(
1778 matches!(events.next().await, Some(Event::TurnComplete { status: TurnOutcomeStatus::Failed, error: Some(error), .. }) if error.contains("fan-in route unavailable"))
1779 );
1780 }
1781
1782 #[tokio::test]
1783 async fn headless_cancel_during_child_wait_returns_interrupted() {
1784 let mut engine = mock_engine_handle();
1785 let mut events = ExecAgentEvents::new(
1786 engine.handle.clone(),
1787 Instant::now() + Duration::from_secs(30),
1788 );
1789 engine.tx_event.send(completed_parent(5)).await.unwrap();
1790 let cancel = async {
1791 reply_to_probe(
1792 &mut engine.rx_op,
1793 SubAgentSettlement {
1794 running_children: 1,
1795 running_workflows: 0,
1796 pending_completions: 0,
1797 },
1798 )
1799 .await;
1800 engine.handle.cancel();
1801 };
1802 let host = async {
1803 assert!(matches!(
1804 events.next().await,
1805 Some(Event::TurnComplete {
1806 status: TurnOutcomeStatus::Interrupted,
1807 ..
1808 })
1809 ));
1810 };
1811 tokio::time::timeout(Duration::from_secs(1), async { tokio::join!(cancel, host) })
1812 .await
1813 .unwrap();
1814 }
1815
1816 #[tokio::test]
1817 async fn headless_closed_engine_cannot_reuse_an_earlier_success_receipt() {
1818 let mut engine = mock_engine_handle();
1819 let mut events = ExecAgentEvents::new(
1820 engine.handle.clone(),
1821 Instant::now() + Duration::from_secs(30),
1822 );
1823 engine.tx_event.send(completed_parent(5)).await.unwrap();
1824 engine.close_event_stream();
1825 assert!(
1826 matches!(events.next().await, Some(Event::TurnComplete { status: TurnOutcomeStatus::Failed, error: Some(error), .. }) if error.contains("channel closed"))
1827 );
1828 }
1829
1830 /// The reproduced defect: headless exec advertised `automation` while
1831 /// attaching no store, so every call — including the read-only `list` and
1832 /// `read` — failed "AutomationManager is not attached". Exec must attach
1833 /// the same durable store the TUI and the Runtime open.
1834 #[test]
1835 fn headless_exec_attaches_the_shared_automation_store() {
1836 let _lock = crate::test_support::lock_test_env();
1837 let tmp = tempfile::TempDir::new().expect("tempdir");
1838 // SAFETY: serialised by lock_test_env.
1839 unsafe {
1840 std::env::set_var("CODEWHALE_AUTOMATIONS_DIR", tmp.path());
1841 }
1842 let attached = exec_automation_services(false).expect("open store");
1843 // SAFETY: cleanup under the same lock.
1844 unsafe {
1845 std::env::remove_var("CODEWHALE_AUTOMATIONS_DIR");
1846 }
1847 let attached = attached.expect("exec attaches the automation store");
1848 let manager = attached.blocking_lock();
1849 // Reads work against the shared store...
1850 assert!(
1851 manager.list_automations().is_ok(),
1852 "an attached store must serve inspection"
1853 );
1854 // ...while the exec host stays outside the Runtime's exclusive
1855 // task-execution lease, so it claims no dispatch ownership.
1856 assert!(
1857 manager.execution_scope().is_none(),
1858 "a one-shot host must not claim an execution scope"
1859 );
1860 }
1861
1862 /// A Fleet worker keeps the narrowed envelope it was launched with.
1863 #[test]
1864 fn fleet_workers_get_no_automation_store() {
1865 assert!(
1866 exec_automation_services(true)
1867 .expect("no store to open")
1868 .is_none()
1869 );
1870 }
1871
1872 /// A store that cannot be opened is reported, not swallowed into a silent
1873 /// "not attached" at the first tool call.
1874 #[test]
1875 fn an_unopenable_store_fails_loudly() {
1876 let _lock = crate::test_support::lock_test_env();
1877 let tmp = tempfile::NamedTempFile::new().expect("temp file");
1878 // A regular file cannot host the store's directories.
1879 let blocked = tmp.path().join("automations");
1880 // SAFETY: serialised by lock_test_env.
1881 unsafe {
1882 std::env::set_var("CODEWHALE_AUTOMATIONS_DIR", &blocked);
1883 }
1884 let result = exec_automation_services(false);
1885 // SAFETY: cleanup under the same lock.
1886 unsafe {
1887 std::env::remove_var("CODEWHALE_AUTOMATIONS_DIR");
1888 }
1889 let err = result.expect_err("opening the store must fail");
1890 assert!(
1891 format!("{err:#}").contains("automation store for headless exec"),
1892 "the failure must name what could not be opened: {err:#}"
1893 );
1894 }
1895
1896 #[test]
1897 fn exec_retry_receipts_keep_the_existing_jsonl_schema_and_text() {
1898 for message in [
1899 "Retry attempt: transport 1/2; upstream 503; waiting 0.00s",
1900 "Retry recovery: transport request recovered after 1 retries",
1901 "Retry exhaustion: stream-resume stopped after 2 retries; stream interrupted",
1902 "Retry stopped: transparent stream completion was not observed",
1903 ] {
1904 let event = super::exec_retry_status(message).expect("retry-only projection");
1905 let value = crate::exec_stream_value(&event).unwrap();
1906 assert_eq!(value["type"], "status");
1907 assert_eq!(value["message"], message);
1908 assert_eq!(value["schema"], "codewhale.exec-stream");
1909 assert_eq!(value["schema_version"], 1);
1910 }
1911 assert!(super::exec_retry_status("Executing tools sequentially").is_none());
1912 assert!(super::exec_retry_status("Goal set; starting goal work.").is_none());
1913 }
1914
1915 #[test]
1916 fn exec_retry_count_requires_terminal_facts_and_keeps_unknown_distinct_from_zero() {
1917 use crate::tool_inspection::TurnStopDiagnostics;
1918 assert_eq!(super::accumulate_exec_retry_count(None, None), None);
1919 assert_eq!(
1920 super::accumulate_exec_retry_count(None, Some(&TurnStopDiagnostics::default())),
1921 None
1922 );
1923 let zero = TurnStopDiagnostics {
1924 status: Some(TurnOutcomeStatus::Completed),
1925 ..Default::default()
1926 };
1927 assert_eq!(
1928 super::accumulate_exec_retry_count(None, Some(&zero)),
1929 Some(0)
1930 );
1931 let retries = TurnStopDiagnostics {
1932 status: Some(TurnOutcomeStatus::Failed),
1933 transport_retries: 2,
1934 stream_resumes: 3,
1935 transparent_stream_retries: 1,
1936 empty_stop_retries: 1,
1937 reasoning_only_reprompts: 2,
1938 ..Default::default()
1939 };
1940 assert_eq!(
1941 super::accumulate_exec_retry_count(Some(0), Some(&retries)),
1942 Some(9)
1943 );
1944 assert_eq!(super::accumulate_exec_retry_count(Some(9), None), Some(9));
1945 assert_eq!(
1946 super::accumulate_exec_retry_count(Some(u32::MAX), Some(&retries)),
1947 Some(u32::MAX)
1948 );
1949 }
1950 }
1951
1951 lines RUST