car_engine/executor.rs
1//! Core execution engine — propose → validate → execute → commit.
2
3use crate::cache::ResultCache;
4use crate::capabilities::CapabilitySet;
5use crate::checkpoint::Checkpoint;
6use crate::rate_limit::{RateLimit, RateLimiter};
7use car_eventlog::{EventKind, EventLog, SpanStatus};
8use car_ir::{
9 build_dag, AcceptedProposalPreimage, Action, ActionProposal, ActionResult, ActionStatus,
10 ActionType, CostSummary, FailureBehavior, ProposalLineageEntry, ProposalLineageStatus,
11 ProposalResult, StateMutation, ToolFailure, ToolFailureClassification, ToolSchema,
12};
13use car_policy::{ApprovalDecision, ApprovalLedger, ApprovalRecord, PermissionTier, PolicyEngine};
14use car_state::{RestoreDurability, StateStore};
15use car_validator::validate_action;
16use serde_json::Value;
17use sha2::{Digest, Sha256};
18use std::collections::HashMap;
19use std::sync::Arc;
20use std::time::Duration;
21use tokio::sync::{Mutex as TokioMutex, RwLock as TokioRwLock};
22use tokio::time::timeout;
23use tracing::instrument;
24use uuid::Uuid;
25
26/// Retry backoff constants.
27const RETRY_BASE_DELAY_MS: u64 = 100;
28const RETRY_BACKOFF_FACTOR: u64 = 2;
29const ROLLBACK_WARNING: &str =
30 "proposal aborted; state effects were rolled back; external effects may remain and were not undone";
31const PLAN_FALLBACK_ROLLBACK_WARNING: &str =
32 "planning candidate rejected; state effects were rolled back; external effects may remain and were not undone";
33const ROLLBACK_DURABILITY_ERROR: &str =
34 "durable state rollback failed before publication; in-memory state and idempotency entries were preserved";
35const ROLLBACK_DURABILITY_UNKNOWN: &str =
36 "state rollback was published but parent-directory durability is unknown; idempotency entries were invalidated for exact retry";
37
38// ---------------------------------------------------------------------------
39// Replan types — failure recovery via model callback
40// ---------------------------------------------------------------------------
41
42/// Callback trait for replanning failed proposals.
43/// Implement this to let the runtime ask the model for an alternative plan
44/// when a proposal aborts.
45#[async_trait::async_trait]
46pub trait ReplanCallback: Send + Sync {
47 async fn replan(&self, ctx: &ReplanContext) -> Result<ActionProposal, String>;
48}
49
50/// Context provided to the replan callback so the model can generate an alternative.
51#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
52pub struct ReplanContext {
53 /// Original proposal ID.
54 pub proposal_id: String,
55 /// Which replan attempt this is (1-indexed; attempt 1 = first replan after initial failure).
56 pub attempt: u32,
57 /// Actions that failed and caused the abort.
58 pub failed_actions: Vec<FailedActionSummary>,
59 /// Action IDs that succeeded before the abort (now rolled back).
60 pub completed_action_ids: Vec<String>,
61 /// State snapshot after rollback.
62 pub state_snapshot: HashMap<String, Value>,
63 /// How many replans remain.
64 pub replans_remaining: u32,
65 /// Original proposal source (model name, agent, etc.).
66 pub original_source: String,
67 /// Total actions in the original proposal.
68 pub original_action_count: usize,
69 /// The original goal/context from the proposal (for generating alternatives).
70 pub original_context: HashMap<String, Value>,
71}
72
73/// Summary of a failed action, included in ReplanContext.
74#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
75pub struct FailedActionSummary {
76 pub action_id: String,
77 pub tool: Option<String>,
78 pub error: String,
79 pub parameters: HashMap<String, Value>,
80}
81
82/// Configuration for the replan loop.
83#[derive(Debug, Clone)]
84pub struct ReplanConfig {
85 /// Maximum number of replan attempts. 0 = disabled (default).
86 pub max_replans: u32,
87 /// Delay in milliseconds between replan attempts. Prevents burning through
88 /// attempts instantly if the model returns garbage fast. 0 = no delay.
89 pub delay_ms: u64,
90 /// If true, replan proposals are scored via car-planner's verify() before
91 /// execution. Proposals with errors are rejected without executing.
92 /// Prevents the engine from running a worse plan than the one that failed.
93 pub verify_before_execute: bool,
94 /// When true, validator/policy/capability rejections (ActionStatus::Rejected)
95 /// also trigger rollback + replan, not just runtime Failed actions.
96 /// Default false (conservative: preserves prior abort-only-on-Failed behavior).
97 pub replan_on_rejected: bool,
98}
99
100/// How the runtime treats a proposal that conflicts with the current
101/// shared state on a pre-execution transactional check (survey §4.3/§5.2.4
102/// — `car_verify::check_transaction` against the versioned `StateStore`).
103#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
104pub enum TransactionCheckMode {
105 /// Don't run the check (default — preserves prior behavior exactly).
106 #[default]
107 Off,
108 /// Run it and emit any conflicts as `TransactionConflict` telemetry, but
109 /// execute anyway.
110 Warn,
111 /// Run it; if any conflict is found, reject the proposal without
112 /// executing (the conflicting actions become `Rejected` results).
113 Strict,
114}
115
116impl Default for ReplanConfig {
117 fn default() -> Self {
118 Self {
119 max_replans: 0,
120 delay_ms: 0,
121 verify_before_execute: true,
122 replan_on_rejected: false,
123 }
124 }
125}
126
127/// Trait for tool execution. Implement this to provide tools to the runtime.
128///
129/// In-process: implement directly with function calls.
130/// Daemon mode: implement by sending JSON-RPC to the client.
131#[async_trait::async_trait]
132pub trait ToolExecutor: Send + Sync {
133 async fn execute(&self, tool: &str, params: &Value) -> Result<Value, String>;
134
135 /// Variant that also carries the originating proposal `Action.id`
136 /// and its `timeout_ms` budget.
137 ///
138 /// `action_id`: WS-based executors (`car-server-core::WsToolExecutor`)
139 /// use it so the daemon-initiated `tools.execute` request to the client
140 /// carries the same id the host's process-wide handler is keyed on
141 /// — without this round-trip the host can't disambiguate concurrent
142 /// callbacks for the same tool (Parslee-ai/car-releases#43 follow-up).
143 ///
144 /// `timeout_ms`: the action's per-call budget. WS executors MUST bound
145 /// their callback wait by this (falling back to a default when `None`)
146 /// so the daemon→host wait and the executor's own action deadline stay
147 /// coordinated — otherwise a hardcoded inner wait reaps a call the outer
148 /// deadline still permits (Parslee-ai/car#259). In-process executors
149 /// that don't need either can keep the default forward to [`execute`](crate::executor::ToolExecutor::execute).
150 async fn execute_with_action(
151 &self,
152 tool: &str,
153 params: &Value,
154 _action_id: &str,
155 _timeout_ms: Option<u64>,
156 ) -> Result<Value, String> {
157 self.execute(tool, params).await
158 }
159
160 /// Variant that also carries the Runtime execution session and the retry
161 /// attempt. Executors that retain safety-relevant state across calls (such
162 /// as a read-before-edit ledger) override this to keep that state isolated
163 /// by session. Existing executors retain their current behavior through the
164 /// default delegation.
165 ///
166 /// `attempt` is **1-based**: `1` on the first try, `2` on the first retry.
167 /// WS executors surface it on the `tools.execute` payload so a host can
168 /// tell which retry it is serving — `action_id` cannot, being neither
169 /// unique across attempts nor varying between them.
170 ///
171 /// It is threaded from `execute_with_retry`'s own counter rather than
172 /// synthesized here. It was previously hardcoded to `1` at the one place
173 /// that put it on the wire, which made the field a constant and any join
174 /// built on it silently degenerate (Parslee-ai/car#928).
175 async fn execute_with_action_in_session(
176 &self,
177 tool: &str,
178 params: &Value,
179 action_id: &str,
180 timeout_ms: Option<u64>,
181 _session_id: Option<&str>,
182 _attempt: u32,
183 ) -> Result<Value, String> {
184 self.execute_with_action(tool, params, action_id, timeout_ms)
185 .await
186 }
187
188 /// Execute one action while allowing a transport to return runtime-observed
189 /// state mutations alongside its ordinary tool output.
190 ///
191 /// The default preserves every direct/internal executor: its existing
192 /// output is wrapped with no mutations, and declared `expected_effects`
193 /// remain assertions rather than writes. A transport that overrides this
194 /// method must validate its mutation envelope against `expected_effects`
195 /// before returning. The Runtime repeats the key-set and I-JSON checks and
196 /// applies accepted changes through its own [`StateStore`].
197 async fn execute_with_action_state_in_session(
198 &self,
199 tool: &str,
200 params: &Value,
201 action_id: &str,
202 timeout_ms: Option<u64>,
203 session_id: Option<&str>,
204 attempt: u32,
205 _expected_effects: &HashMap<String, Value>,
206 _return_schema: Option<&Value>,
207 ) -> Result<ToolExecution, String> {
208 self.execute_with_action_in_session(
209 tool, params, action_id, timeout_ms, session_id, attempt,
210 )
211 .await
212 .map(ToolExecution::output_only)
213 }
214
215 /// Execute one action with typed failure evidence.
216 ///
217 /// This is a new default-delegating layer rather than a signature change
218 /// to the existing execution methods. Existing executors therefore retain
219 /// their exact behavior: every legacy string error becomes an ordinary,
220 /// non-terminal failure. Executors opt into fail-stop behavior only by
221 /// overriding this method and returning [`ToolFailure::terminal`].
222 async fn execute_classified(
223 &self,
224 tool: &str,
225 params: &Value,
226 action_id: &str,
227 timeout_ms: Option<u64>,
228 session_id: Option<&str>,
229 attempt: u32,
230 expected_effects: &HashMap<String, Value>,
231 return_schema: Option<&Value>,
232 ) -> Result<ToolExecution, ToolFailure> {
233 self.execute_with_action_state_in_session(
234 tool,
235 params,
236 action_id,
237 timeout_ms,
238 session_id,
239 attempt,
240 expected_effects,
241 return_schema,
242 )
243 .await
244 .map_err(ToolFailure::from)
245 }
246
247 /// Streaming entry point for detached invocation modes (C2). Start
248 /// the tool and return a channel of [`car_ir::ToolStreamChunk`]s; the
249 /// runtime drains it into the per-runtime handle registry while the
250 /// DAG proceeds. End the stream with a terminal chunk (`done` /
251 /// `error`); dropping the sender without one is reported as failure.
252 /// A cooperative executor should stop work when the receiver returned
253 /// here is dropped (that's what cancellation looks like from its side).
254 ///
255 /// Default: unsupported — existing one-shot executors compile and
256 /// behave unchanged; a detached action against them is rejected with
257 /// this error.
258 async fn execute_stream(
259 &self,
260 tool: &str,
261 _params: &Value,
262 _action_id: &str,
263 ) -> Result<tokio::sync::mpsc::Receiver<car_ir::ToolStreamChunk>, String> {
264 Err(format!(
265 "tool '{tool}': this executor does not support streaming/long-running invocation"
266 ))
267 }
268}
269
270/// Result of a tool dispatch before the Runtime commits observed state.
271#[derive(Debug, Clone, PartialEq)]
272pub struct ToolExecution {
273 pub output: Value,
274 pub state_changes: HashMap<String, Value>,
275}
276
277impl ToolExecution {
278 pub fn output_only(output: Value) -> Self {
279 Self {
280 output,
281 state_changes: HashMap::new(),
282 }
283 }
284}
285
286/// Deterministic key for idempotency deduplication.
287/// Deterministic key for idempotency deduplication. Carries the tenant
288/// dimension (linus review): without it, tenant A's cached result for
289/// an identical idempotent action was served to tenant B — a
290/// cross-tenant data leak through the dedup cache. Unscoped executions
291/// keep their historical keys (empty tenant segment).
292fn idempotency_key(action: &Action, scope: Option<&crate::scope::RuntimeScope>) -> String {
293 let sorted: std::collections::BTreeMap<_, _> = action.parameters.iter().collect();
294 let params = serde_json::to_string(&sorted).unwrap_or_default();
295 let tenant = scope.and_then(|s| s.tenant_id.as_deref()).unwrap_or("");
296 format!(
297 "{}:{}:{}:{}",
298 tenant,
299 serde_json::to_string(&action.action_type).unwrap_or_default(),
300 action.tool.as_deref().unwrap_or(""),
301 params
302 )
303}
304
305/// Exact normal-serde proposal preimage plus its lowercase RFC 8785/JCS
306/// SHA-256 identity. The shared inference helper is the repository's settled
307/// canonicalizer; journal identity must not drift onto a second JSON sorter.
308fn proposal_journal_identity(proposal: &ActionProposal) -> Result<(Value, String), String> {
309 let preimage = serde_json::to_value(proposal)
310 .map_err(|error| format!("proposal serialization failed: {error}"))?;
311 let canonical = car_inference::catalog_identity::canonical_json(&preimage)?;
312 let digest = format!("{:x}", Sha256::digest(canonical.as_bytes()));
313 Ok((preimage, digest))
314}
315
316fn proposal_digest(proposal: &ActionProposal) -> Result<String, String> {
317 proposal_journal_identity(proposal).map(|(_, digest)| digest)
318}
319
320fn proposal_rejection_boundary_data(
321 proposal: &ActionProposal,
322 proposal_digest: Option<&str>,
323 reason: &str,
324) -> HashMap<String, Value> {
325 let mut data: HashMap<String, Value> = [
326 ("affected_actions".to_string(), Value::Array(Vec::new())),
327 (
328 "rolled_back_changes".to_string(),
329 Value::Object(serde_json::Map::new()),
330 ),
331 (
332 "changes_semantics".to_string(),
333 Value::from("proposal_rejected_no_state_commit"),
334 ),
335 ("attempted".to_string(), Value::from(false)),
336 ("stage".to_string(), Value::from("proposal_rejection")),
337 ("rejection_reason".to_string(), Value::from(reason)),
338 ]
339 .into();
340 if let Ok(preimage) = serde_json::to_value(proposal) {
341 data.insert("proposal".to_string(), preimage);
342 }
343 if let Some(digest) = proposal_digest {
344 data.insert("proposal_digest".to_string(), Value::from(digest));
345 }
346 data
347}
348
349fn proposal_lineage_entry(
350 proposal: &ActionProposal,
351 generation: u32,
352 status: ProposalLineageStatus,
353 rejection_reason: Option<String>,
354) -> ProposalLineageEntry {
355 ProposalLineageEntry {
356 generation,
357 proposal_id: proposal.id.clone(),
358 proposal_digest: proposal_digest(proposal).ok(),
359 status,
360 rejection_reason,
361 }
362}
363
364fn finalize_proposal_result(
365 mut result: ProposalResult,
366 original_proposal_id: &str,
367 lineage: &[ProposalLineageEntry],
368 accepted_proposal_preimages: &[AcceptedProposalPreimage],
369) -> ProposalResult {
370 result.original_proposal_id = original_proposal_id.to_string();
371 result.replan_lineage = lineage.to_vec();
372 result.accepted_proposal_preimages = accepted_proposal_preimages.to_vec();
373 result
374}
375
376/// Validate the proposal-local action identity boundary used by DAG execution,
377/// action receipts, results, and state-transition attribution. A retry reuses
378/// one admitted action id across attempts; two declared actions may not share
379/// an id because every downstream join would become ambiguous.
380pub fn validate_proposal_action_ids(proposal: &ActionProposal) -> Result<(), String> {
381 let mut seen = std::collections::HashSet::with_capacity(proposal.actions.len());
382 for action in &proposal.actions {
383 if !seen.insert(action.id.as_str()) {
384 return Err(format!("duplicate action id '{}'", action.id));
385 }
386 }
387 Ok(())
388}
389
390fn validate_proposal_retry_limits(proposal: &ActionProposal) -> Result<(), String> {
391 for action in &proposal.actions {
392 if action.failure_behavior == FailureBehavior::Retry
393 && action.max_retries.checked_add(1).is_none()
394 {
395 return Err(format!(
396 "action '{}' retry attempt count overflows u32",
397 action.id
398 ));
399 }
400 }
401 Ok(())
402}
403
404fn requires_integrity_rollback(result: &ActionResult) -> bool {
405 result.status == ActionStatus::Failed
406 && result.error.as_deref().is_some_and(|error| {
407 error.starts_with("tool output failed JCS/I-JSON validation after dispatch")
408 })
409}
410
411fn record_rollback_durability_error(
412 results: &mut [ActionResult],
413 prefix: &str,
414 error: &str,
415) -> String {
416 let detail = format!("{prefix}: {error}");
417 if let Some(result) = results
418 .iter_mut()
419 .find(|result| result.status != ActionStatus::Succeeded)
420 {
421 result.error = Some(match result.error.take() {
422 Some(existing) => format!("{existing}; {detail}"),
423 None => detail.clone(),
424 });
425 }
426 detail
427}
428
429fn rejected_result(action_id: &str, error: String) -> ActionResult {
430 ActionResult {
431 action_id: action_id.to_string(),
432 status: ActionStatus::Rejected,
433 output: None,
434 error: Some(error),
435 terminal: false,
436 state_changes: HashMap::new(),
437 rolled_back: false,
438 duration_ms: None,
439 timestamp: chrono::Utc::now(),
440 }
441}
442
443/// Capture only the state keys relevant to an action (state_dependencies + expected_effects).
444/// Returns an empty map if the action declares no relevant keys (e.g., tool calls
445/// that don't interact with state).
446fn snapshot_relevant_keys(
447 state: &car_state::StateStore,
448 action: &Action,
449) -> HashMap<String, Value> {
450 let mut keys: std::collections::HashSet<&str> = std::collections::HashSet::new();
451 for dep in &action.state_dependencies {
452 keys.insert(dep.as_str());
453 }
454 for key in action.expected_effects.keys() {
455 keys.insert(key.as_str());
456 }
457 // For state_write actions, capture the key being written
458 if action.action_type == ActionType::StateWrite {
459 if let Some(key) = action.parameters.get("key").and_then(|v| v.as_str()) {
460 keys.insert(key);
461 }
462 }
463
464 if keys.is_empty() {
465 // No declared state interaction — skip snapshot (empty map)
466 return HashMap::new();
467 }
468
469 keys.iter()
470 .filter_map(|&k| state.get(k).map(|v| (k.to_string(), v)))
471 .collect()
472}
473
474fn skipped_result(action_id: &str, reason: &str) -> ActionResult {
475 ActionResult {
476 action_id: action_id.to_string(),
477 status: ActionStatus::Skipped,
478 output: None,
479 error: Some(reason.to_string()),
480 terminal: false,
481 state_changes: HashMap::new(),
482 rolled_back: false,
483 duration_ms: None,
484 timestamp: chrono::Utc::now(),
485 }
486}
487
488fn insert_tool_event_provenance(
489 data: &mut HashMap<String, Value>,
490 action: &Action,
491 source: Option<car_ir::ToolSourceKind>,
492) {
493 if let Some(tool) = action.tool.as_deref() {
494 data.insert("tool".to_string(), Value::from(tool));
495 data.insert(
496 "tool_source".to_string(),
497 Value::from(
498 source
499 .unwrap_or(car_ir::ToolSourceKind::UserDefined)
500 .as_str(),
501 ),
502 );
503 }
504}
505
506/// SHA-256 over the RFC 8785/JCS rendering of only the action parameters.
507///
508/// The accepted proposal has already crossed `proposal_journal_identity`, so
509/// canonicalization cannot fail here without violating that admission
510/// invariant. Keeping only the digest on outcome events lets a later join find
511/// the raw parameters in `ProposalReceived` without copying them into every
512/// failure record.
513fn action_params_digest(action: &Action) -> String {
514 // Serialize the map directly: RFC 8785/JCS sorts object keys
515 // recursively at render time (see `canonical_json` and the
516 // `canonical_json_sorts_object_keys_recursively` test), so the map's
517 // iteration order does not affect the digest and cloning into a
518 // `Value::Object` would be pure per-action allocation overhead.
519 let canonical = car_inference::catalog_identity::canonical_json(&action.parameters)
520 .expect("accepted action parameters must remain RFC 8785/JCS canonicalizable");
521 format!("{:x}", Sha256::digest(canonical.as_bytes()))
522}
523
524fn insert_action_outcome_signal(
525 data: &mut HashMap<String, Value>,
526 action: &Action,
527 params_digest: &str,
528 error_class: Option<&str>,
529) {
530 data.insert(
531 "params_digest".to_string(),
532 Value::from(params_digest.to_string()),
533 );
534 data.insert(
535 "expected_effects".to_string(),
536 serde_json::to_value(&action.expected_effects)
537 .expect("action expected_effects must serialize as JSON"),
538 );
539 if let Some(error_class) = error_class {
540 data.insert("error_class".to_string(), Value::from(error_class));
541 }
542}
543
544/// Normalize execution failures into a stable, deliberately low-cardinality
545/// event-log vocabulary. Validation and policy checks that happen before
546/// dispatch remain `ActionRejected`; this maps failures reached after an
547/// action began executing.
548pub(crate) fn action_error_class(
549 action: &Action,
550 error: &str,
551 engine_timeout: bool,
552) -> &'static str {
553 // Classification is by substring over prose and is a published wire
554 // contract (docs/agent-ir-spec.md `error_class`). Each literal is pinned
555 // to its producer below; changing a producer's wording must update the
556 // matching literal (and `action_failure_error_class_mapping_is_stable`).
557 let error = error.to_ascii_lowercase();
558 // "callback timed out" — the engine deadline itself is typed via
559 // `engine_timeout`, so this literal only covers the detached-dispatch
560 // producer `tool '{tool}' callback timed out ({s}s)`
561 // (car-server-core/src/session.rs, `reap_detached_calls`/wait path).
562 // A tool faithfully relaying a *downstream* timeout in its own message
563 // also lands here; accepted drift, low-cardinality vocabulary.
564 if engine_timeout || error.contains("callback timed out") {
565 "timeout"
566 } else if error.starts_with("denied by policy:") || error.starts_with("rejected by policy:") {
567 // Exact producers return this prefix unwrapped from the tool's
568 // guard check: car-server-core/src/assistant/executor.rs and
569 // car-server-core/src/coder/shell_tool.rs
570 // (`format!("denied by policy: {reason}")`). An executor that
571 // *wraps* the denial (e.g. `tool 'x' failed: denied by policy: …`)
572 // deliberately falls through to `tool_error` — see the wrapped
573 // case in `action_failure_error_class_mapping_is_stable`.
574 "rejected_by_policy"
575 } else if error.contains("jcs/i-json validation")
576 || error.contains("output validation:")
577 || error.contains("invalid return json schema")
578 || error.contains("must return exact envelope {output,state_changes}")
579 || error.contains("state_changes must be an object")
580 || error.contains("state_changes keys do not match")
581 || error.contains("state_changes serialization failed")
582 {
583 "validation"
584 } else if action.tool.is_some() {
585 "tool_error"
586 } else {
587 "unknown"
588 }
589}
590
591fn action_outcome_data(
592 stage: &str,
593 reason: &str,
594 attempted: bool,
595 attempt: Option<u32>,
596) -> HashMap<String, Value> {
597 let mut data: HashMap<String, Value> = [
598 ("stage".to_string(), Value::from(stage)),
599 ("reason".to_string(), Value::from(reason)),
600 ("attempted".to_string(), Value::from(attempted)),
601 ]
602 .into();
603 if let Some(attempt) = attempt {
604 data.insert("attempt".to_string(), Value::from(attempt));
605 }
606 data
607}
608
609/// Prefix on the `error` field of an `ActionResult` that distinguishes
610/// "the user pulled the plug" from "earlier abort cascaded." The
611/// `ActionStatus` itself is `Skipped` in both cases (introducing a
612/// new variant ripples through every IR consumer + FFI binding); the
613/// prefix lets callers like the A2A bridge tell the cases apart
614/// without string-matching a magic literal.
615pub const CANCELED_PREFIX: &str = "canceled: ";
616
617/// Result for an action that didn't run because the proposal was
618/// cancelled mid-flight. Distinct from `Skipped` so callers can tell
619/// "didn't run because of an earlier abort" from "didn't run because
620/// the user pulled the plug."
621fn canceled_result(action_id: &str, reason: &str) -> ActionResult {
622 ActionResult {
623 action_id: action_id.to_string(),
624 status: ActionStatus::Skipped,
625 output: None,
626 error: Some(format!("{}{}", CANCELED_PREFIX, reason)),
627 terminal: false,
628 state_changes: HashMap::new(),
629 rolled_back: false,
630 duration_ms: None,
631 timestamp: chrono::Utc::now(),
632 }
633}
634
635/// Format a tool result for feeding back to a model.
636pub fn format_tool_result(result: &ActionResult) -> String {
637 match result.status {
638 ActionStatus::Succeeded => match &result.output {
639 Some(v) => serde_json::to_string(v).unwrap_or_else(|_| v.to_string()),
640 None => String::new(),
641 },
642 ActionStatus::Rejected => format!("[REJECTED] {}", result.error.as_deref().unwrap_or("")),
643 ActionStatus::Failed => format!("[FAILED] {}", result.error.as_deref().unwrap_or("")),
644 _ => format!(
645 "[{:?}] {}",
646 result.status,
647 result.error.as_deref().unwrap_or("")
648 ),
649 }
650}
651
652/// Names the shape [`format_tool_result_for_model`] gives a model, so a
653/// measurement can record what its model read. Change it whenever that shape
654/// changes: results measured under different formats are not comparable.
655pub const MODEL_OBSERVATION_FORMAT: &str = "text-v1";
656
657/// Fields that say how a call went; [`render_tool_output`] puts them first.
658const CONTROL_FIELDS: [&str; 6] = [
659 "truncated",
660 "error",
661 "ok",
662 "status",
663 "exit_code",
664 "timed_out",
665];
666
667/// A tool result's output as a model should read it: outcome fields first
668/// (`truncated`, `error`, `ok`, `status`, `exit_code`, `timed_out`), one line
669/// each; then other scalars as `key: value` lines; then each long or
670/// multi-line string (a file's content, a command's output) as a raw block
671/// headed `key (N bytes):`, and each array of records as `key (N records):`
672/// with one line per record. A cut at the budget loses the tail of the bulk,
673/// never the verdict.
674///
675/// Journals, receipts and events keep the JSON; only the model's copy changes.
676/// The model used to be handed the JSON too — so a file arrived as one line of
677/// `\n`, `\t` and `\"` escapes. On claude-haiku-4-5, 23% of native `edit_file`
678/// calls (42 of 182) missed with "old_text not found", against 1 of 60 for
679/// Claude Code's `Edit` on the same model and tasks; the misses carried escapes
680/// copied out of the JSON (`\"\"\"` for a docstring's `"""`). Claude Code's
681/// `Read` returns plain numbered lines; this is that shape.
682///
683/// The text reads back exactly with [`parse_rendered_output`], apart from the
684/// outcome fields (shown on one line) and a record list (read back as text):
685/// blocks carry their length, so no content line can be mistaken for a field,
686/// and a string that would read back as another JSON type (`"123"`, `"null"`)
687/// is written quoted.
688pub fn render_tool_output(v: &Value) -> String {
689 let map = match v {
690 Value::Object(map) => map,
691 Value::String(s) => return s.clone(),
692 other => return other.to_string(),
693 };
694 let mut head = Vec::new();
695 let mut lines = Vec::new();
696 let mut blocks = Vec::new();
697 // Sorted, so every build shows the same bytes: serde_json's map order
698 // depends on whether `preserve_order` is enabled anywhere in the graph.
699 let mut fields: Vec<_> = map.iter().collect();
700 fields.sort_by(|a, b| a.0.cmp(b.0));
701 for (key, value) in fields {
702 match value {
703 // Outcome fields lead, on one line each, whatever their size: a cut
704 // at the budget must never take the verdict with it.
705 _ if CONTROL_FIELDS.contains(&key.as_str()) => {
706 // Readable, not escaped: a traceback in `error` shows as one
707 // line with its breaks turned into spaces, keeping its start
708 // and its end (where the exception is) when it is long.
709 let line = match value {
710 Value::String(s) => scalar_text(&Value::String(s.replace('\n', " "))),
711 other => scalar_text(other),
712 };
713 head.push(format!(
714 "{}: {}",
715 key_text(key),
716 clip_keeping_ends(&line, CONTROL_FIELD_BYTES)
717 ));
718 }
719 Value::String(s) if s.contains('\n') || s.len() > 200 => {
720 blocks.push(format!("{} ({} bytes):\n{s}", key_text(key), s.len()))
721 }
722 Value::Array(items) if !items.is_empty() && items.iter().all(Value::is_object) => {
723 let records: Vec<String> = items.iter().map(render_record).collect();
724 blocks.push(format!(
725 "{} ({} records):\n{}",
726 key_text(key),
727 records.len(),
728 records.join("\n")
729 ));
730 }
731 other => lines.push(format!("{}: {}", key_text(key), scalar_text(other))),
732 }
733 }
734 head.extend(lines);
735 head.extend(blocks);
736 head.join("\n")
737}
738
739/// A scalar as it appears after `key: `. A string is written raw unless raw
740/// text would read back as something else — JSON of another type, a quoted
741/// string, or surrounding whitespace — in which case it is written as a JSON
742/// string. Everything else is compact JSON.
743fn scalar_text(value: &Value) -> String {
744 match value {
745 Value::String(s)
746 if s.trim() != s
747 || s.is_empty()
748 || s.contains('\n')
749 || s.ends_with("):")
750 || serde_json::from_str::<Value>(s).is_ok() =>
751 {
752 serde_json::to_string(s).unwrap_or_else(|_| s.clone())
753 }
754 Value::String(s) => s.clone(),
755 other => other.to_string(),
756 }
757}
758
759/// How much of an outcome field a model is shown on its one line.
760const CONTROL_FIELD_BYTES: usize = 500;
761
762/// A key as written before `: ` or ` (N bytes):` — raw unless raw text would
763/// not read back as this key, then as a JSON string.
764fn key_text(key: &str) -> String {
765 if key.is_empty()
766 || key.contains(char::is_whitespace)
767 || key.ends_with(':')
768 || key.contains(": ")
769 || key.contains(" (")
770 || key.contains('\n')
771 || key.starts_with('"')
772 || key.starts_with('[')
773 || key.starts_with('…')
774 {
775 serde_json::to_string(key).unwrap_or_else(|_| key.to_string())
776 } else {
777 key.to_string()
778 }
779}
780
781/// Read a key written by [`key_text`] off the front of `line`, returning it
782/// and what follows it.
783fn read_key(line: &str) -> Option<(String, &str)> {
784 if line.starts_with('"') {
785 let mut stream = serde_json::Deserializer::from_str(line).into_iter::<String>();
786 let key = stream.next()?.ok()?;
787 return Some((key, &line[stream.byte_offset()..]));
788 }
789 let end = line.find(": ").or_else(|| line.find(" ("))?;
790 let key = &line[..end];
791 // `key_text` quotes any key with whitespace or a leading `[`/`…`, so a
792 // raw one without them is a key; `[FAILED] tool 'x': …` is not.
793 (!key.is_empty()
794 && !key.contains(char::is_whitespace)
795 && !key.starts_with('[')
796 && !key.starts_with('…'))
797 .then(|| (key.to_string(), &line[end..]))
798}
799
800/// `s` in at most about `max` bytes, keeping its first third and its end,
801/// with the number of bytes dropped between them.
802fn clip_keeping_ends(s: &str, max: usize) -> String {
803 if s.len() <= max {
804 return s.to_string();
805 }
806 let mut head = max / 3;
807 while !s.is_char_boundary(head) {
808 head -= 1;
809 }
810 let mut tail = s.len() - (max - head);
811 while !s.is_char_boundary(tail) {
812 tail += 1;
813 }
814 format!("{}…[{} bytes]…{}", &s[..head], tail - head, &s[tail..])
815}
816
817/// A tool result's fields from whichever form it was stored in: the JSON a
818/// receipt holds, or the rendered text a model was shown (a transcript). Any
819/// JSON parses as itself; rendered text reads back through
820/// [`parse_rendered_output`]; `None` for anything else, such as a `[FAILED] …`
821/// line.
822pub fn parse_tool_output(text: &str) -> Option<Value> {
823 serde_json::from_str::<Value>(text)
824 .ok()
825 .or_else(|| parse_rendered_output(text))
826}
827
828/// Read [`render_tool_output`]'s text back into a JSON object — for consumers
829/// that need fields (an outcome, a receipt's structured redaction) when only
830/// the model's copy survives, such as a transcript replayed without its
831/// recorded `ok`. Exact for scalars and blocks; an outcome field reads back as
832/// the one line the model saw (breaks as spaces, long ones clipped), and a
833/// record list reads back as its text. A block cut short by a budget reads
834/// back as what is left of it, and parsing stops at the first line that is not
835/// a field, which is where a budget cut the text. `None` unless the first line
836/// is a field — a `[FAILED] …` line, a refusal or prose is not rendered output.
837pub fn parse_rendered_output(text: &str) -> Option<Value> {
838 fn block_header(line: &str) -> Option<(String, usize, bool)> {
839 let (key, rest) = read_key(line)?;
840 let count = rest.strip_prefix(" (")?.strip_suffix("):")?;
841 let (n, bytes) = if let Some(n) = count.strip_suffix(" bytes") {
842 (n, true)
843 } else {
844 (count.strip_suffix(" records")?, false)
845 };
846 Some((key, n.parse().ok()?, bytes))
847 }
848 let mut map = serde_json::Map::new();
849 let mut rest = text;
850 let mut seen_block = false;
851 while !rest.is_empty() {
852 let (line, after) = rest.split_once('\n').unwrap_or((rest, ""));
853 if let Some((key, n, bytes)) = block_header(line) {
854 let (value, remaining) = if bytes {
855 let mut end = n.min(after.len());
856 while !after.is_char_boundary(end) {
857 end -= 1;
858 }
859 (
860 &after[..end],
861 after[end..].strip_prefix('\n').unwrap_or(&after[end..]),
862 )
863 } else {
864 let mut end = 0;
865 let mut taken = 0;
866 while taken < n && end < after.len() {
867 end = after[end..].find('\n').map_or(after.len(), |i| end + i + 1);
868 taken += 1;
869 }
870 (after[..end].trim_end_matches('\n'), &after[end..])
871 };
872 // First wins: the renderer writes each key once, so a second is
873 // text a cut spliced in, and must not replace an outcome field.
874 map.entry(key)
875 .or_insert_with(|| Value::String(value.to_string()));
876 seen_block = true;
877 rest = remaining;
878 continue;
879 }
880 match read_key(line) {
881 Some((key, raw)) if !seen_block && raw.starts_with(": ") => {
882 let raw = &raw[2..];
883 let value = serde_json::from_str::<Value>(raw)
884 .unwrap_or_else(|_| Value::String(raw.to_string()));
885 map.entry(key).or_insert(value);
886 }
887 // The first line decides what this is: a `[FAILED] …` line, a
888 // refusal or prose is not rendered output and is not read as
889 // fields. After a field, an unrecognised line is where a budget
890 // cut the text (mid-line, or through a marker spliced into a
891 // block): keep what came before it — above all the outcome
892 // fields, which the renderer writes first.
893 _ if map.is_empty() => return None,
894 _ => break,
895 }
896 rest = after;
897 }
898 (!map.is_empty()).then_some(Value::Object(map))
899}
900
901/// One record of an array field (a grep match, a directory entry) on one line,
902/// its string values raw.
903fn render_record(item: &Value) -> String {
904 let Value::Object(fields) = item else {
905 return item.to_string();
906 };
907 let mut fields: Vec<_> = fields.iter().collect();
908 fields.sort_by(|a, b| a.0.cmp(b.0));
909 fields
910 .into_iter()
911 .map(|(k, v)| {
912 let k = k.replace('\n', "\\n");
913 match v {
914 Value::String(s) => format!("{k}={}", s.replace('\n', "\\n")),
915 other => format!("{k}={other}"),
916 }
917 })
918 .collect::<Vec<_>>()
919 .join(" ")
920}
921
922/// The phrase [`truncate_keeping_ends`] puts in its elision marker, for
923/// readers that count cut observations.
924pub const ELIDED_MARKER: &str = "bytes elided from the middle";
925
926/// Fit `s` in `max` bytes keeping its first third and its end, with a marker
927/// saying how much of the middle was dropped. For command output: the verdict
928/// a test run or build prints comes last, so a head-only cut loses exactly the
929/// line the model needs.
930pub fn truncate_keeping_ends(s: &str, max: usize) -> String {
931 if s.len() <= max {
932 return s.to_string();
933 }
934 if max < 512 {
935 // Too small for a head, a tail and the marker: keep the end, the part
936 // that carries the verdict.
937 let mut start = s.len() - max;
938 while !s.is_char_boundary(start) {
939 start += 1;
940 }
941 return s[start..].to_string();
942 }
943 let budget = max.saturating_sub(128);
944 let mut head = budget / 3;
945 while !s.is_char_boundary(head) {
946 head -= 1;
947 }
948 let mut tail = s.len() - (budget - head);
949 while !s.is_char_boundary(tail) {
950 tail += 1;
951 }
952 format!(
953 "{}\n…[{} {ELIDED_MARKER} of {}; re-run with narrower output to see them]…\n{}",
954 &s[..head],
955 tail - head,
956 s.len(),
957 &s[tail..]
958 )
959}
960
961/// Whether a tool's result is command output, whose end matters most: the
962/// tools every CAR agent loop truncates with [`truncate_keeping_ends`].
963pub fn is_command_tool(name: &str) -> bool {
964 matches!(name, "shell" | "run_command")
965}
966
967/// [`format_tool_result`] as a model should read it: a successful result's
968/// output goes through [`render_tool_output`] instead of being serialized as
969/// JSON; failures keep their `[FAILED]` / `[REJECTED]` prefix.
970pub fn format_tool_result_for_model(result: &ActionResult) -> String {
971 match (&result.status, &result.output) {
972 (ActionStatus::Succeeded, Some(v)) => render_tool_output(v),
973 _ => format_tool_result(result),
974 }
975}
976
977/// Budget constraints for proposal execution.
978#[derive(Debug, Clone)]
979pub struct CostBudget {
980 pub max_tool_calls: Option<u32>,
981 pub max_duration_ms: Option<f64>,
982 pub max_actions: Option<u32>,
983}
984
985/// Common Agent Runtime — deterministic execution layer.
986///
987/// Lock ordering discipline (never hold multiple simultaneously, never hold sync locks across .await):
988/// 1. capabilities (RwLock, read-only during execution)
989/// 2. tools (RwLock, read-only during execution)
990/// 3. policies (RwLock, read-only during execution)
991/// 4. session_policies (RwLock to find the per-session engine, then read inner Arc)
992/// 5. cost_budget (RwLock, read-only during execution)
993/// 6. log (TokioMutex, acquired/released per event)
994/// 7. tool_executor (TokioMutex, clone Arc and drop before await)
995/// 8. idempotency_cache (TokioMutex, acquired/released per check)
996///
997/// StateStore uses parking_lot::Mutex (sync) — NEVER hold across .await points.
998pub struct Runtime {
999 pub state: Arc<StateStore>,
1000 pub tools: Arc<TokioRwLock<HashMap<String, ToolSchema>>>,
1001 pub policies: Arc<TokioRwLock<PolicyEngine>>,
1002 /// Per-session policy registries. Hosts that multiplex multiple
1003 /// concurrent agent sessions over a single `Runtime` (IDE-style
1004 /// frontends with per-project rules, multi-tenant servers) call
1005 /// [`Runtime::open_session`] to mint an id, then
1006 /// [`Runtime::register_policy_in_session`] to attach session-scoped
1007 /// rules. Validation under that session walks the global registry
1008 /// AND the session's — both must pass. Sessions can deny what
1009 /// global allows; sessions cannot allow what global denies.
1010 /// Closing a session drops its registry and any closures it holds.
1011 /// See `docs/proposals/per-session-policy-scoping.md`.
1012 pub session_policies: Arc<TokioRwLock<HashMap<String, Arc<TokioRwLock<PolicyEngine>>>>>,
1013 pub log: Arc<TokioMutex<EventLog>>,
1014 pub rate_limiter: Arc<RateLimiter>,
1015 pub result_cache: Arc<ResultCache>,
1016 /// Per-execution-session read ledgers backing the read-before-edit / staleness guard on
1017 /// the built-in file tools (H1/F4-remainder, audit 2026-07-06). Threaded into
1018 /// `agent_basics::execute_with_ledger` on the built-in fallback path so a raw
1019 /// Runtime session (no configured executor) still requires reading a file
1020 /// before editing or overwriting it. Configured executors
1021 /// (coder/assistant/bench) own their own ledger, so this one only governs the
1022 /// built-in fallback path.
1023 read_ledgers: crate::agent_basics::SessionReadLedgers,
1024 tool_executor: TokioMutex<Option<Arc<dyn ToolExecutor>>>,
1025 idempotency_cache: TokioMutex<HashMap<String, ActionResult>>,
1026 cost_budget: TokioRwLock<Option<CostBudget>>,
1027 capabilities: TokioRwLock<Option<CapabilitySet>>,
1028 inference_engine: Option<Arc<car_inference::InferenceEngine>>,
1029 /// Optional outbound transport backing the `messaging.send` built-in.
1030 /// `None` by default — a runtime with no sink refuses `messaging.send`
1031 /// outright rather than falling through to the host executor, and never
1032 /// advertises the tool. Attach one with
1033 /// [`Runtime::with_message_sink`]; see [`crate::messaging`] for why
1034 /// human-directed messaging belongs inside the governed chain at all.
1035 message_sink: Option<Arc<dyn crate::messaging::MessageSink>>,
1036 /// Optional memgine for skill learning (auto-distillation after execution).
1037 memgine: Option<Arc<TokioMutex<car_memgine::MemgineEngine>>>,
1038 /// Whether to auto-distill skills after each proposal execution.
1039 auto_distill: bool,
1040 /// Optional trajectory store for persisting execution traces.
1041 trajectory_store: Option<Arc<car_memgine::TrajectoryStore>>,
1042 /// Optional replan callback for failure recovery.
1043 replan_callback: TokioMutex<Option<Arc<dyn ReplanCallback>>>,
1044 /// Replan configuration.
1045 replan_config: TokioRwLock<ReplanConfig>,
1046 /// Pre-execution transactional conflict check mode (survey §4.3/§5.2.4).
1047 /// `Off` by default. When `Warn`/`Strict`, each proposal is checked
1048 /// against the versioned shared state before execution; `Strict`
1049 /// rejects on conflict. See [`Runtime::set_transaction_check_mode`].
1050 transaction_check: TokioRwLock<TransactionCheckMode>,
1051 /// The live harness operating config the Evolution Agent tunes (survey
1052 /// §3.5). `None` until one is installed via [`Runtime::set_harness_config`]
1053 /// — so default behavior is byte-identical to before (no cap on
1054 /// per-action retries, built-in backoff). Once installed, the runtime
1055 /// *reads* it: `max_retries` caps per-action retry budgets and
1056 /// `retry_backoff_ms` sets the inter-attempt delay; the setter also maps
1057 /// `planning_max_replans` onto the replan config. This is what makes an
1058 /// applied [`car_memgine::HarnessConfigPatch`] take effect.
1059 harness_config: TokioRwLock<Option<car_memgine::HarnessConfig>>,
1060 /// Canonical tool registry (optional — new code should use this).
1061 pub registry: Arc<crate::registry::ToolRegistry>,
1062 /// The environment the agent's side-effecting built-in tools act within.
1063 /// Defaults to [`crate::substrate::LocalSubstrate`] (host fs/process), which
1064 /// reproduces the historic `agent_basics` host behavior byte-for-byte.
1065 /// Bind a different environment (e.g. a VM via `McpSubstrate`) with
1066 /// [`Runtime::with_substrate`] / [`Runtime::set_substrate`]. `calculate`
1067 /// stays pure and never consults the substrate.
1068 substrate: TokioRwLock<Arc<dyn crate::substrate::Substrate>>,
1069 /// Proposal-admission gates — the pre-execution safety seam (EPIC A /
1070 /// task A1). Each registered [`crate::admission::AdmissionGate`] runs
1071 /// during admission, before any action executes; a proposal that any
1072 /// gate blocks (or escalates to approval) is refused. Empty by default,
1073 /// so a runtime that registers no gates behaves exactly as before.
1074 /// Individual gates (information-flow, concurrency, blocking-policy)
1075 /// are layered on via [`Runtime::register_admission_gate`].
1076 admission_gates: TokioRwLock<Vec<Arc<dyn crate::admission::AdmissionGate>>>,
1077 /// Runtime taint provenance for the VIGIL intent gate
1078 /// ([`crate::taint::TaintLedger`]). Installed by
1079 /// [`Runtime::install_intent_gate`] and by nothing else: with no intent
1080 /// gate there is no ledger, so a runtime that never configures VIGIL
1081 /// pays no per-action cost and behaves exactly as before. When present,
1082 /// each successful action records which state keys its result wrote and
1083 /// whether that result was tainted, and the gate reads it back at
1084 /// admission to mark actions that READ a tainted key — the cross-
1085 /// proposal (replan) and trusted-tool-laundering cases the proposal-
1086 /// local dependency DAG cannot see.
1087 taint_ledger: TokioRwLock<Option<Arc<crate::taint::TaintLedger>>>,
1088 /// Durable human-in-the-loop approval ledger (EPIC A / A7). When set,
1089 /// an admission gate's `NeedsApproval` verdict is resolved against this
1090 /// ledger by fingerprint: a prior Approved decision admits the
1091 /// proposal, a Rejected decision blocks it, an unseen one stays pending
1092 /// (fail-closed). `None` by default — escalations fail closed with an
1093 /// explanatory reason until a ledger is installed.
1094 approval_ledger: TokioRwLock<Option<ApprovalLedger>>,
1095 /// Optional JSONL journal backing the idempotency cache (EPIC A / C3).
1096 /// `None` by default — the cache is in-memory and lost on restart. When
1097 /// set via [`Runtime::set_idempotency_cache_path`], idempotent results
1098 /// are persisted and reloaded so a crash-restart doesn't re-execute a
1099 /// completed idempotent action (avoiding duplicate side effects).
1100 idempotency_journal: TokioRwLock<Option<std::path::PathBuf>>,
1101 /// Detached tool invocations (C2): a `ToolCall` with a
1102 /// `streaming`/`long_running` invocation mode is registered here and
1103 /// its handle returned as the action's output while the DAG proceeds.
1104 /// Chunks are drained via [`Runtime::tool_poll`], cancelled via
1105 /// [`Runtime::tool_cancel`], and fanned out to
1106 /// [`Runtime::subscribe_tool_events`] subscribers.
1107 pub tool_handles: Arc<crate::tool_handles::ToolHandleRegistry>,
1108}
1109
1110/// A durable idempotency-cache record (C3). `result: None` is a tombstone
1111/// recorded when a cached entry is invalidated by a rollback.
1112#[derive(serde::Serialize, serde::Deserialize)]
1113struct IdempotencyEntry {
1114 key: String,
1115 result: Option<ActionResult>,
1116}
1117
1118impl Runtime {
1119 pub fn new() -> Self {
1120 Self {
1121 state: Arc::new(StateStore::new()),
1122 tools: Arc::new(TokioRwLock::new(HashMap::new())),
1123 policies: Arc::new(TokioRwLock::new(PolicyEngine::new())),
1124 session_policies: Arc::new(TokioRwLock::new(HashMap::new())),
1125 log: Arc::new(TokioMutex::new(EventLog::new())),
1126 rate_limiter: Arc::new(RateLimiter::new()),
1127 result_cache: Arc::new(ResultCache::new()),
1128 read_ledgers: crate::agent_basics::SessionReadLedgers::new(),
1129 tool_executor: TokioMutex::new(None),
1130 idempotency_cache: TokioMutex::new(HashMap::new()),
1131 cost_budget: TokioRwLock::new(None),
1132 capabilities: TokioRwLock::new(None),
1133 inference_engine: None,
1134 message_sink: None,
1135 memgine: None,
1136 auto_distill: false,
1137 trajectory_store: None,
1138 replan_callback: TokioMutex::new(None),
1139 replan_config: TokioRwLock::new(ReplanConfig::default()),
1140 transaction_check: TokioRwLock::new(TransactionCheckMode::Off),
1141 harness_config: TokioRwLock::new(None),
1142 registry: Arc::new(crate::registry::ToolRegistry::new()),
1143 substrate: TokioRwLock::new(Arc::new(crate::substrate::LocalSubstrate::new())),
1144 admission_gates: TokioRwLock::new(Vec::new()),
1145 taint_ledger: TokioRwLock::new(None),
1146 approval_ledger: TokioRwLock::new(None),
1147 idempotency_journal: TokioRwLock::new(None),
1148 tool_handles: Arc::new(crate::tool_handles::ToolHandleRegistry::new()),
1149 }
1150 }
1151
1152 /// Create a runtime with shared state, event log, and policies.
1153 /// Each runtime gets its own tool set, executor, and idempotency cache.
1154 pub fn with_shared(
1155 state: Arc<StateStore>,
1156 log: Arc<TokioMutex<EventLog>>,
1157 policies: Arc<TokioRwLock<PolicyEngine>>,
1158 ) -> Self {
1159 Self {
1160 state,
1161 tools: Arc::new(TokioRwLock::new(HashMap::new())),
1162 policies,
1163 // Session-policy registries are per-runtime — sharing them
1164 // across embedders that share global policies would defeat
1165 // the isolation point. Hosts that genuinely want shared
1166 // sessions should drive them through one shared Runtime.
1167 session_policies: Arc::new(TokioRwLock::new(HashMap::new())),
1168 log,
1169 rate_limiter: Arc::new(RateLimiter::new()),
1170 result_cache: Arc::new(ResultCache::new()),
1171 read_ledgers: crate::agent_basics::SessionReadLedgers::new(),
1172 tool_executor: TokioMutex::new(None),
1173 idempotency_cache: TokioMutex::new(HashMap::new()),
1174 cost_budget: TokioRwLock::new(None),
1175 capabilities: TokioRwLock::new(None),
1176 inference_engine: None,
1177 message_sink: None,
1178 memgine: None,
1179 auto_distill: false,
1180 trajectory_store: None,
1181 replan_callback: TokioMutex::new(None),
1182 replan_config: TokioRwLock::new(ReplanConfig::default()),
1183 transaction_check: TokioRwLock::new(TransactionCheckMode::Off),
1184 harness_config: TokioRwLock::new(None),
1185 registry: Arc::new(crate::registry::ToolRegistry::new()),
1186 substrate: TokioRwLock::new(Arc::new(crate::substrate::LocalSubstrate::new())),
1187 admission_gates: TokioRwLock::new(Vec::new()),
1188 taint_ledger: TokioRwLock::new(None),
1189 approval_ledger: TokioRwLock::new(None),
1190 idempotency_journal: TokioRwLock::new(None),
1191 tool_handles: Arc::new(crate::tool_handles::ToolHandleRegistry::new()),
1192 }
1193 }
1194
1195 // ─── Session policy lifecycle ───────────────────────────────────
1196 //
1197 // Per-session policy scoping. Policies registered against a
1198 // session apply to proposals executed under that session id (via
1199 // [`Self::execute_with_session`] / [`Self::execute_with_session_and_cancel`]).
1200 // Policies registered globally always apply, on top of any
1201 // session-scoped layer. See `docs/proposals/per-session-policy-scoping.md`.
1202
1203 /// Mint a new session id and pre-register an empty policy engine
1204 /// under it. Hosts call this once per concurrent agent context
1205 /// (an IDE project window, a multi-tenant client, etc.) and pair
1206 /// it with [`Self::close_session`] when the context ends.
1207 ///
1208 /// Returns the opaque id to pass to subsequent
1209 /// [`Self::register_policy_in_session`] / [`Self::execute_with_session`]
1210 /// calls. Ids are UUIDs so collisions across concurrent calls
1211 /// don't matter.
1212 pub async fn open_session(&self) -> String {
1213 let id = Uuid::new_v4().to_string();
1214 let mut sessions = self.session_policies.write().await;
1215 sessions.insert(id.clone(), Arc::new(TokioRwLock::new(PolicyEngine::new())));
1216 self.read_ledgers.ledger_for(Some(&id));
1217 id
1218 }
1219
1220 /// Drop the session and every policy scoped to it. Returns true
1221 /// if a session by that id existed; false if it didn't (already
1222 /// closed, never opened, etc.). Idempotent in effect — closing a
1223 /// missing session is a no-op the caller is free to ignore.
1224 pub async fn close_session(&self, session_id: &str) -> bool {
1225 let mut sessions = self.session_policies.write().await;
1226 let removed = sessions.remove(session_id).is_some();
1227 if removed {
1228 self.read_ledgers.remove(session_id);
1229 }
1230 removed
1231 }
1232
1233 /// Register a policy under a specific session id. The policy
1234 /// applies only when a proposal is executed under that session;
1235 /// proposals executed without a session (the default) only see
1236 /// global policies.
1237 ///
1238 /// Returns `Err(...)` if the session is unknown — callers either
1239 /// forgot to call [`Self::open_session`] or are using a
1240 /// stale/closed id.
1241 pub async fn register_policy_in_session(
1242 &self,
1243 session_id: &str,
1244 name: &str,
1245 check: car_policy::PolicyCheck,
1246 description: &str,
1247 ) -> Result<(), String> {
1248 let engine = {
1249 let sessions = self.session_policies.read().await;
1250 sessions
1251 .get(session_id)
1252 .cloned()
1253 .ok_or_else(|| format!("unknown session id '{session_id}'"))?
1254 };
1255 let mut engine = engine.write().await;
1256 engine.register(name, check, description);
1257 Ok(())
1258 }
1259
1260 /// [`Self::register_policy_in_session`] for a check that forbids `tool`
1261 /// outright, so `PolicyEngine::blanket_denied_tools` can read it back.
1262 /// Separate method rather than an extra parameter: the existing signature
1263 /// is public and crosses the bindings, and this is purely additive.
1264 pub async fn register_tool_deny_in_session(
1265 &self,
1266 session_id: &str,
1267 name: &str,
1268 tool: &str,
1269 check: car_policy::PolicyCheck,
1270 description: &str,
1271 ) -> Result<(), String> {
1272 let engine = {
1273 let sessions = self.session_policies.read().await;
1274 sessions
1275 .get(session_id)
1276 .cloned()
1277 .ok_or_else(|| format!("unknown session id '{session_id}'"))?
1278 };
1279 let mut engine = engine.write().await;
1280 engine.register_tool_deny(name, tool, check, description);
1281 Ok(())
1282 }
1283
1284 /// Remove a policy by name. `session_id` targets a session's policy set;
1285 /// `None` targets the global one.
1286 ///
1287 /// Returns how many policies were dropped, or `Err` if `session_id` names a
1288 /// session that doesn't exist. Session-scoped policies could always be
1289 /// dropped wholesale by [`Self::close_session`], but a *global* policy had
1290 /// no removal path at all — once registered it lived until the process
1291 /// exited, so a mistyped or over-broad rule could only be cleared by
1292 /// restarting the daemon (Parslee-ai/car#623).
1293 pub async fn unregister_policy(
1294 &self,
1295 name: &str,
1296 session_id: Option<&str>,
1297 ) -> Result<usize, String> {
1298 match session_id {
1299 Some(sid) => {
1300 let engine = {
1301 let sessions = self.session_policies.read().await;
1302 sessions
1303 .get(sid)
1304 .cloned()
1305 .ok_or_else(|| format!("unknown session id '{sid}'"))?
1306 };
1307 let mut engine = engine.write().await;
1308 Ok(engine.unregister(name))
1309 }
1310 None => {
1311 let mut engine = self.policies.write().await;
1312 Ok(engine.unregister(name))
1313 }
1314 }
1315 }
1316
1317 /// Registered policies as `(name, description)`. `session_id` lists a
1318 /// session's set; `None` lists the global one. Without this a client could
1319 /// register a policy but never ask what was in force, so a rejection could
1320 /// not be explained beyond its single message (Parslee-ai/car#623).
1321 pub async fn list_policies(
1322 &self,
1323 session_id: Option<&str>,
1324 ) -> Result<Vec<(String, String)>, String> {
1325 match session_id {
1326 Some(sid) => {
1327 let engine = {
1328 let sessions = self.session_policies.read().await;
1329 sessions
1330 .get(sid)
1331 .cloned()
1332 .ok_or_else(|| format!("unknown session id '{sid}'"))?
1333 };
1334 let engine = engine.read().await;
1335 Ok(engine.policy_details())
1336 }
1337 None => {
1338 let engine = self.policies.read().await;
1339 Ok(engine.policy_details())
1340 }
1341 }
1342 }
1343
1344 /// Load declarative deny rules from a project's `.car/policies/`
1345 /// directory and register them on the global policy engine (EPIC A /
1346 /// task A2).
1347 ///
1348 /// `car_dir` is a `.car` directory; this looks for
1349 /// `car_dir/policies/*.toml`. **The caller chooses the directory and
1350 /// there is no walk-up here** — this joins the path it is given, once.
1351 /// The two production callers pick differently: the daemon passes
1352 /// `$HOME/.car`, and the assistant passes its working directory's
1353 /// `.car` (`--dir`, else cwd). So a rule file at a repository root does
1354 /// not govern a `car do` run started from a subdirectory. Say which
1355 /// directory you mean at the call site rather than assuming discovery.
1356 ///
1357 /// A missing directory is not an error (returns 0). A malformed rule
1358 /// file *is* an error — a dropped security rule must surface loudly.
1359 /// Returns the number of rules registered.
1360 ///
1361 /// These rules are additive on top of any code-registered policies, and
1362 /// every one of them is a prohibition — though `allow_tool_param` states
1363 /// its prohibition as an allowlist, denying everything about its tool
1364 /// that it does not name. Once A9 makes policy violations blocking at
1365 /// admission, a matching action refuses the proposal.
1366 pub async fn load_project_policies(
1367 &self,
1368 car_dir: impl AsRef<std::path::Path>,
1369 ) -> Result<usize, car_policy::PolicyLoadError> {
1370 let dir = car_dir.as_ref().join("policies");
1371 let rules = car_policy::load_policy_dir(&dir)?;
1372 // `PolicyRules::len` counts every kind. Summing a hand-picked subset
1373 // here is what made this under-report once already — it covered the
1374 // three kinds that existed at the time and was not revisited when
1375 // more landed, so a file full of allowlists reported zero rules
1376 // loaded while `apply` registered all of them.
1377 let count = rules.len();
1378 let mut engine = self.policies.write().await;
1379 rules.apply(&mut engine);
1380 Ok(count)
1381 }
1382
1383 /// Load information-flow tool labels from a project's `.car` directory
1384 /// and register the information-flow admission gate (EPIC A / A3+A4).
1385 ///
1386 /// Reads `car_dir/tool-labels.json` (merged over built-in defaults) and
1387 /// registers an [`crate::flow::InformationFlowGate`] so every admitted
1388 /// proposal is checked for data exfiltration (blocked) and forbidden
1389 /// tool orderings (escalated to approval). A missing labels file is
1390 /// fine — the built-in defaults still mark the network tools as sinks.
1391 /// A malformed file is a loud error.
1392 pub async fn install_information_flow_gate(
1393 &self,
1394 car_dir: impl AsRef<std::path::Path>,
1395 ) -> Result<(), crate::flow::FlowLoadError> {
1396 let config = crate::flow::load_tool_labels(car_dir)?;
1397 let gate = Arc::new(crate::flow::InformationFlowGate::new(config));
1398 self.register_admission_gate(gate).await;
1399 Ok(())
1400 }
1401
1402 /// True if a session with this id is currently open. Mostly for
1403 /// tests and FFI surface validation — production code should
1404 /// trust the id it just opened.
1405 pub async fn session_exists(&self, session_id: &str) -> bool {
1406 self.session_policies.read().await.contains_key(session_id)
1407 }
1408
1409 /// Attach a local inference engine. Registers `infer`, `embed`, `classify`
1410 /// as built-in tools with real implementations.
1411 pub fn with_inference(mut self, engine: Arc<car_inference::InferenceEngine>) -> Self {
1412 self.inference_engine = Some(engine);
1413 // Register inference tool schemas (non-async init, use try_lock).
1414 // Keep the canonical registry populated too: execution-event provenance
1415 // is read from ToolEntry.source rather than inferred from a tool name.
1416 //
1417 // Contention contract: this runs at construction time on a private
1418 // runtime, so both try_locks are expected to succeed. If either fails
1419 // the two stores would disagree (validator reads `tools`, provenance
1420 // reads the registry first) — loud in debug, logged in release.
1421 if let Ok(mut tools) = self.tools.try_write() {
1422 for schema in car_inference::service::all_schemas() {
1423 let name = schema.name.clone();
1424 let entry = crate::registry::ToolEntry::builtin(schema);
1425 tools.insert(name.clone(), entry.schema.clone());
1426 if !self.registry.try_register(entry) {
1427 tracing::warn!(
1428 tool = %name,
1429 "with_inference: canonical registry contention; \
1430 schema map registered the inference tool but the \
1431 registry did not — provenance falls back to the map"
1432 );
1433 }
1434 }
1435 } else {
1436 debug_assert!(
1437 false,
1438 "with_inference must populate both stores; schema-map lock was contended \
1439 at construction time — inference tools stay unregistered"
1440 );
1441 tracing::warn!(
1442 "with_inference: schema map lock contended at construction; \
1443 inference tools were NOT registered"
1444 );
1445 }
1446 self
1447 }
1448
1449 /// Attach an outbound message sink, making `messaging.send` a real tool on
1450 /// this runtime.
1451 ///
1452 /// Registering the schema here — and ONLY here — is deliberate. A tool the
1453 /// runtime cannot execute must not be advertised to the model: without a
1454 /// sink the dispatch arm can only return an error, and a model that keeps
1455 /// seeing "message the human" in its tool list will keep trying to use it.
1456 /// So sink and schema arrive together, and neither exists alone.
1457 ///
1458 /// The entry is `AskUser` with side effects: reaching a human is
1459 /// irreversible, so the conservative default is the right one, and hosts
1460 /// that have their own consent model can re-register with a different
1461 /// permission. `category = "messaging"` groups it for the capability
1462 /// surfaces that filter by category.
1463 ///
1464 /// Follows [`Self::with_inference`]'s non-async registration pattern
1465 /// (`try_write` during construction, while the locks are uncontended) —
1466 /// making the builders async would break every `Runtime::new().with_…()`
1467 /// chain in the workspace for no gain.
1468 pub fn with_message_sink(mut self, sink: Arc<dyn crate::messaging::MessageSink>) -> Self {
1469 self.message_sink = Some(sink);
1470
1471 let entry = crate::registry::ToolEntry::builtin(car_ir::builtins::messaging_send())
1472 .with_permission(crate::registry::ToolPermission::AskUser)
1473 .with_side_effects(true)
1474 .with_category("messaging");
1475 let schema = entry.schema.clone();
1476
1477 // Honour the schema's declared rate limit, the way the async
1478 // `register_tool_schema` path does — otherwise the backstop on
1479 // outbound human messaging would silently not exist.
1480 if let Some(ref rl) = schema.rate_limit {
1481 self.rate_limiter.try_set_limit(
1482 &schema.name,
1483 RateLimit {
1484 max_calls: rl.max_calls,
1485 interval_secs: rl.interval_secs,
1486 },
1487 );
1488 }
1489 self.registry.try_register(entry);
1490 // The legacy schema map is what `validate_action` reads, so the
1491 // validator can only check `messaging.send` parameters once it lands
1492 // here.
1493 if let Ok(mut tools) = self.tools.try_write() {
1494 tools.insert(schema.name.clone(), schema);
1495 }
1496 self
1497 }
1498
1499 /// Attach a memgine for automatic skill learning after execution.
1500 /// When `auto_distill` is true, execution traces are automatically distilled
1501 /// into skills and domains are evolved when underperforming.
1502 pub fn with_learning(
1503 mut self,
1504 memgine: Arc<TokioMutex<car_memgine::MemgineEngine>>,
1505 auto_distill: bool,
1506 ) -> Self {
1507 self.memgine = Some(memgine);
1508 self.auto_distill = auto_distill;
1509 self
1510 }
1511
1512 /// Attach a memgine with auto-distillation enabled (recommended default).
1513 pub fn with_memgine(self, memgine: Arc<TokioMutex<car_memgine::MemgineEngine>>) -> Self {
1514 self.with_learning(memgine, true)
1515 }
1516
1517 /// Attach a trajectory store for persisting execution traces.
1518 pub fn with_trajectory_store(mut self, store: Arc<car_memgine::TrajectoryStore>) -> Self {
1519 self.trajectory_store = Some(store);
1520 self
1521 }
1522
1523 /// The attached trajectory store, if any.
1524 ///
1525 /// Writing traces was always the point of the store; reading them back is
1526 /// what makes them a feedback signal rather than an audit log.
1527 pub fn trajectory_store(&self) -> Option<&Arc<car_memgine::TrajectoryStore>> {
1528 self.trajectory_store.as_ref()
1529 }
1530
1531 /// Per-tool success rates observed over the last `window_days`, or `None`
1532 /// when no trajectory store is attached.
1533 ///
1534 /// Dispatch-conditional (see
1535 /// [`ToolFeedback::dispatched_from_trajectories`](car_planner::ToolFeedback::dispatched_from_trajectories))
1536 /// — rejected and skipped actions are excluded, because the tool never ran
1537 /// and a consumer that models rejection separately would otherwise count
1538 /// it twice.
1539 ///
1540 /// The window exists because a success rate is a claim about how a tool
1541 /// behaves *now*. An API that was broken for a week six months ago and has
1542 /// been fixed since should not still be dragging its own rate down, and an
1543 /// unbounded history makes the rate progressively less responsive to
1544 /// exactly the recent change an operator is trying to see. 30 days is the
1545 /// suggested default at the call sites.
1546 ///
1547 /// Reads from disk on every call — parsing is bounded by the window (day
1548 /// files outside it are skipped by filename, never opened), and the callers
1549 /// are interactive rather than hot-path. If that stops being true, cache
1550 /// here rather than at each call site.
1551 pub fn tool_feedback(&self, window_days: u32) -> Option<car_planner::ToolFeedback> {
1552 let store = self.trajectory_store.as_ref()?;
1553 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(window_days));
1554 let trajectories = store.load_since(cutoff);
1555 Some(car_planner::ToolFeedback::dispatched_from_trajectories(
1556 &trajectories,
1557 ))
1558 }
1559
1560 pub fn with_executor(self, executor: Arc<dyn ToolExecutor>) -> Self {
1561 // Use try_lock for non-async init context. Safe because we just created the mutex.
1562 if let Ok(mut guard) = self.tool_executor.try_lock() {
1563 *guard = Some(executor);
1564 }
1565 self
1566 }
1567
1568 /// Set a tool executor for the next execute() call.
1569 /// Used by NAPI bindings where executor varies per call.
1570 pub async fn set_executor(&self, executor: Arc<dyn ToolExecutor>) {
1571 *self.tool_executor.lock().await = Some(executor);
1572 }
1573
1574 /// Bind the execution substrate the side-effecting built-in tools
1575 /// (`read_file`/`write_file`/`edit_file`/`list_dir`/`find_files`/
1576 /// `grep_files`) act within (builder). Defaults to
1577 /// [`crate::substrate::LocalSubstrate`]; bind e.g.
1578 /// [`crate::substrate::McpSubstrate`] to make those tools hit a VM.
1579 /// `calculate` stays pure and ignores the substrate.
1580 pub fn with_substrate(self, substrate: Arc<dyn crate::substrate::Substrate>) -> Self {
1581 // Use try_write for non-async init context. Safe because we just created the lock.
1582 if let Ok(mut guard) = self.substrate.try_write() {
1583 *guard = substrate;
1584 }
1585 self
1586 }
1587
1588 /// Set the execution substrate at runtime.
1589 pub async fn set_substrate(&self, substrate: Arc<dyn crate::substrate::Substrate>) {
1590 {
1591 *self.substrate.write().await = substrate;
1592 }
1593 self.read_ledgers.clear();
1594 }
1595
1596 /// Clone the currently bound substrate.
1597 pub async fn substrate(&self) -> Arc<dyn crate::substrate::Substrate> {
1598 self.substrate.read().await.clone()
1599 }
1600
1601 pub fn with_event_log(mut self, log: EventLog) -> Self {
1602 self.log = Arc::new(TokioMutex::new(log));
1603 self
1604 }
1605
1606 /// Bind an event log the caller already holds a handle to.
1607 ///
1608 /// [`Self::with_event_log`] takes ownership, which leaves no way for
1609 /// anything outside the runtime to read what was recorded. That was fine
1610 /// while the log was observability-only, but a model-callable events
1611 /// surface has to run inside a `ToolExecutor` — and the executor is
1612 /// constructed *before* the runtime that would own it, so the two can only
1613 /// meet through a handle created first and shared into both
1614 /// (Parslee-ai/car#815).
1615 pub fn with_shared_event_log(mut self, log: Arc<TokioMutex<EventLog>>) -> Self {
1616 self.log = log;
1617 self
1618 }
1619
1620 /// The runtime's event log handle, for callers that need to read the record
1621 /// the runtime is writing.
1622 pub fn event_log_handle(&self) -> Arc<TokioMutex<EventLog>> {
1623 Arc::clone(&self.log)
1624 }
1625
1626 /// Attach a replan callback for failure recovery (builder).
1627 pub fn with_replan(self, callback: Arc<dyn ReplanCallback>, config: ReplanConfig) -> Self {
1628 if let Ok(mut guard) = self.replan_callback.try_lock() {
1629 *guard = Some(callback);
1630 }
1631 if let Ok(mut guard) = self.replan_config.try_write() {
1632 *guard = config;
1633 }
1634 self
1635 }
1636
1637 /// Set a replan callback at runtime.
1638 pub async fn set_replan_callback(&self, callback: Arc<dyn ReplanCallback>) {
1639 *self.replan_callback.lock().await = Some(callback);
1640 }
1641
1642 /// Set replan configuration at runtime.
1643 pub async fn set_replan_config(&self, config: ReplanConfig) {
1644 *self.replan_config.write().await = config;
1645 }
1646
1647 /// Set the pre-execution transactional conflict-check mode (survey
1648 /// §4.3/§5.2.4). `Off` (default) preserves prior behavior; `Warn`
1649 /// records conflicts as telemetry; `Strict` rejects a conflicting
1650 /// proposal before executing it.
1651 pub async fn set_transaction_check_mode(&self, mode: TransactionCheckMode) {
1652 *self.transaction_check.write().await = mode;
1653 }
1654
1655 /// Register a proposal-admission gate (EPIC A / task A1).
1656 ///
1657 /// Gates run during admission, before any action executes, in the
1658 /// order they were registered. Each is a verified pre-execution safety
1659 /// check — information-flow (A4), concurrency (A5), blocking-policy
1660 /// (A9) — that can block a proposal or escalate it to human approval.
1661 /// Registering no gates leaves behavior unchanged.
1662 pub async fn register_admission_gate(&self, gate: Arc<dyn crate::admission::AdmissionGate>) {
1663 self.admission_gates.write().await.push(gate);
1664 }
1665
1666 /// Register `gate`, **replacing** any already-registered gate answering
1667 /// to the same [`crate::admission::AdmissionGate::name`] instead of
1668 /// appending a second one. For gates that are *reconfigured* rather than
1669 /// layered — the VIGIL intent gate is reloaded whenever an operator
1670 /// re-reads `.car/intent.json` — appending is a bug, not a no-op: gate
1671 /// aggregation is fail-closed, so the stale gate's verdict still wins
1672 /// while it holds state (a [`crate::taint::TaintLedger`]) nothing writes
1673 /// to any more, and two gates with the same name emit two
1674 /// `AdmissionGateDecision` events under that name.
1675 ///
1676 /// The replaced gate keeps its **position**, so a reconfigure never
1677 /// silently reorders admission relative to the other gates. Any further
1678 /// gate carrying the same name is a stale duplicate from an earlier
1679 /// registration and is dropped, so exactly one gate answers to the name.
1680 async fn replace_admission_gate(&self, gate: Arc<dyn crate::admission::AdmissionGate>) {
1681 let name = gate.name().to_string();
1682 let mut gates = self.admission_gates.write().await;
1683 match gates.iter().position(|g| g.name() == name) {
1684 Some(i) => {
1685 gates[i] = gate;
1686 let mut idx = 0usize;
1687 gates.retain(|g| {
1688 let keep = idx == i || g.name() != name;
1689 idx += 1;
1690 keep
1691 });
1692 }
1693 None => gates.push(gate),
1694 }
1695 }
1696
1697 /// Remove all registered admission gates (primarily for tests and
1698 /// reconfiguration). After this, proposal admission reverts to the
1699 /// transactional pre-check only.
1700 pub async fn clear_admission_gates(&self) {
1701 self.admission_gates.write().await.clear();
1702 }
1703
1704 /// The number of currently-registered admission gates.
1705 pub async fn admission_gate_count(&self) -> usize {
1706 self.admission_gates.read().await.len()
1707 }
1708
1709 /// The `name()` of every registered admission gate, in registration order.
1710 ///
1711 /// Prefer this over [`Self::admission_gate_count`] when asserting that a
1712 /// *particular* gate is installed: a bare count couples the assertion to
1713 /// every other gate the runtime happens to register, so adding one breaks
1714 /// unrelated tests that only cared about their own.
1715 pub async fn admission_gate_names(&self) -> Vec<String> {
1716 self.admission_gates
1717 .read()
1718 .await
1719 .iter()
1720 .map(|g| g.name().to_string())
1721 .collect()
1722 }
1723
1724 /// Register the VIGIL intent gate (arXiv 2601.05755 — the live
1725 /// verify-before-commit call-site). Installs
1726 /// [`crate::intent_gate::IntentGate`] as an admission gate: a
1727 /// forbidden capability or a tool-stream-influenced out-of-intent
1728 /// action hard-rejects the proposal; untainted drift escalates to
1729 /// the durable approval flow (A7) by content-bound fingerprint.
1730 /// Replans are covered automatically (gates re-run on every
1731 /// replanned proposal).
1732 ///
1733 /// This is also where the runtime's [`crate::taint::TaintLedger`] is
1734 /// installed. From here on, every successful action records whether its
1735 /// result was tainted and which state keys it wrote, and the gate marks
1736 /// any incoming action that READS a tainted key as tool-stream-
1737 /// influenced. That closes the two holes the static `untrusted_tools`
1738 /// list leaves open: a trusted tool laundering attacker-controlled
1739 /// content out of state, and a replanned proposal that carries no
1740 /// dependency edge back to the poisoned action. Installing no intent
1741 /// gate installs no ledger, so nothing changes for a runtime that
1742 /// doesn't configure VIGIL.
1743 ///
1744 /// **Calling this again REPLACES the installed gate; it does not add a
1745 /// second one.** Reinstalling is a real operation — an operator editing
1746 /// `.car/intent.json` and re-running
1747 /// [`Runtime::install_intent_gate_from_project`] lands here — and each
1748 /// call builds a *fresh* ledger that becomes the only one the executor
1749 /// writes to. Appending would leave the previous gate registered while
1750 /// holding an orphaned ledger, and because gate aggregation is
1751 /// fail-closed that frozen gate's verdict would still win: every key
1752 /// tainted before the reload would stay tainted forever, no trusted
1753 /// overwrite could clear it, and legitimate out-of-intent work would
1754 /// hard-refuse with no approval path. So the runtime holds exactly one
1755 /// intent gate, and it is always the one bound to the live ledger. The
1756 /// replacement keeps the previous gate's position in the admission
1757 /// order, so a reload never silently reorders the other gates.
1758 pub async fn install_intent_gate(&self, config: crate::intent_gate::IntentGateConfig) {
1759 let ledger = Arc::new(crate::taint::TaintLedger::new(
1760 config.untrusted_tools.iter().cloned().collect(),
1761 ));
1762 *self.taint_ledger.write().await = Some(Arc::clone(&ledger));
1763 self.replace_admission_gate(Arc::new(crate::intent_gate::IntentGate::with_taint(
1764 config, ledger,
1765 )))
1766 .await;
1767 }
1768
1769 /// The runtime's taint provenance ledger, when an intent gate is
1770 /// installed (see [`Runtime::install_intent_gate`]). `None` otherwise —
1771 /// the ledger is strictly opt-in, alongside VIGIL.
1772 pub async fn taint_ledger(&self) -> Option<Arc<crate::taint::TaintLedger>> {
1773 self.taint_ledger.read().await.clone()
1774 }
1775
1776 /// Load `.car/intent.json` from `car_dir` and install the intent
1777 /// gate when present. Absent file → Ok(false) (opt-in, ungated);
1778 /// malformed file → loud error, never a silently-ungated session.
1779 pub async fn install_intent_gate_from_project(
1780 &self,
1781 car_dir: impl AsRef<std::path::Path>,
1782 ) -> Result<bool, crate::intent_gate::IntentLoadError> {
1783 match crate::intent_gate::load_intent_config(car_dir)? {
1784 Some(cfg) => {
1785 self.install_intent_gate(cfg).await;
1786 Ok(true)
1787 }
1788 None => Ok(false),
1789 }
1790 }
1791
1792 /// Register the skill deployment-tier ceiling gate (EPIC A / A8).
1793 ///
1794 /// Requires a memgine to be attached (skills live there). When a
1795 /// proposal names its driving skill in `context["skill"]`, the gate
1796 /// caps the proposal's actions at that skill's persisted
1797 /// `deployment_tier`, escalating an over-ceiling action to the durable
1798 /// approval flow (A7). Returns false if no memgine is attached.
1799 pub async fn install_skill_ceiling_gate(&self) -> bool {
1800 match &self.memgine {
1801 Some(mem) => {
1802 let gate = Arc::new(crate::skill_ceiling::SkillCeilingGate::new(mem.clone()));
1803 self.register_admission_gate(gate).await;
1804 true
1805 }
1806 None => false,
1807 }
1808 }
1809
1810 /// Drain buffered chunks + status for a detached tool invocation (C2).
1811 /// `None` for an unknown or fully-consumed handle. See
1812 /// [`crate::tool_handles::ToolHandleRegistry::poll`] for the
1813 /// consume-on-terminal contract.
1814 pub async fn tool_poll(&self, handle_id: &str) -> Option<crate::tool_handles::ToolPollResult> {
1815 self.tool_handles.poll(handle_id).await
1816 }
1817
1818 /// Request cancellation of a detached tool invocation (C2): fires the
1819 /// handle's cancel token (dropping the executor's chunk receiver) and
1820 /// seals its status as `cancelled` unless already terminal. Returns
1821 /// false for an unknown handle.
1822 pub async fn tool_cancel(&self, handle_id: &str) -> bool {
1823 self.tool_handles.cancel(handle_id).await
1824 }
1825
1826 /// Subscribe to the live [`car_ir::ToolStreamEvent`] fanout for all
1827 /// detached tool invocations on this runtime (C2). The WS layer
1828 /// forwards these as `tools.stream.event` notifications.
1829 pub fn subscribe_tool_events(
1830 &self,
1831 ) -> tokio::sync::broadcast::Receiver<car_ir::ToolStreamEvent> {
1832 self.tool_handles.subscribe()
1833 }
1834
1835 /// Enable tamper-evident hash chaining on the event log (EPIC A / A9).
1836 /// Every event appended from now on is linked to its predecessor by a
1837 /// content hash, so an after-the-fact edit to a chained event — or an
1838 /// interior deletion/reordering — is detectable via
1839 /// [`Runtime::verify_event_log_chain`]. Truncation at either end of the
1840 /// log (dropping a prefix or suffix wholesale) is NOT detectable — the
1841 /// chain has no anchored head hash and no trusted tail witness; that is
1842 /// out of scope until the chain head is anchored. Opt-in: existing logs
1843 /// stay byte-identical until enabled.
1844 pub async fn enable_event_log_hash_chaining(&self) {
1845 self.log.lock().await.enable_hash_chaining();
1846 }
1847
1848 /// Verify the event log's tamper-evidence chain (EPIC A / A9). Returns
1849 /// `Ok(n)` for `n` verified chained events, or `Err(index)` naming the
1850 /// first event whose hash/linkage doesn't match — the point of an
1851 /// interior edit, deletion, or reordering. Head/tail truncation is not
1852 /// detectable (no anchored head hash; the first chained event's
1853 /// `prev_hash` is taken on trust) — see `EventLog::verify_chain`.
1854 pub async fn verify_event_log_chain(&self) -> Result<usize, usize> {
1855 self.log.lock().await.verify_chain()
1856 }
1857
1858 /// Install a durable HITL approval ledger backed by a JSONL journal
1859 /// (EPIC A / A7). Loads any existing decisions so approvals survive a
1860 /// restart, then resolves future admission-gate `NeedsApproval`
1861 /// verdicts against it. The canonical daemon path is
1862 /// `~/.car/approvals.jsonl`.
1863 pub async fn set_approval_ledger_path(
1864 &self,
1865 path: impl Into<std::path::PathBuf>,
1866 ) -> std::io::Result<()> {
1867 let ledger = ApprovalLedger::with_journal(path.into())?;
1868 // Surface journal corruption (or a torn concurrent write) instead of
1869 // silently trusting a partial ledger (review A7 — the doc on
1870 // `skipped_on_load` promises callers surface it).
1871 if ledger.skipped_on_load() > 0 {
1872 tracing::warn!(
1873 skipped = ledger.skipped_on_load(),
1874 "approval ledger journal had unparseable lines skipped on load — \
1875 the ledger may be missing decisions"
1876 );
1877 }
1878 *self.approval_ledger.write().await = Some(ledger);
1879 Ok(())
1880 }
1881
1882 /// Use an in-memory approval ledger (no persistence) — primarily for
1883 /// tests and ephemeral runtimes.
1884 pub async fn set_approval_ledger_in_memory(&self) {
1885 *self.approval_ledger.write().await = Some(ApprovalLedger::new());
1886 }
1887
1888 /// Record a human approval for an admission fingerprint (EPIC A / A7).
1889 /// A subsequently re-submitted proposal whose escalation matches this
1890 /// fingerprint is admitted without asking again.
1891 pub async fn approve_admission(
1892 &self,
1893 fingerprint: &str,
1894 reviewer: &str,
1895 reason: &str,
1896 ) -> Result<(), String> {
1897 self.record_admission_decision(fingerprint, ApprovalDecision::Approved, reviewer, reason)
1898 .await
1899 }
1900
1901 /// Record a human rejection for an admission fingerprint (EPIC A / A7).
1902 /// A proposal whose escalation matches a rejected fingerprint is
1903 /// blocked outright.
1904 pub async fn reject_admission(
1905 &self,
1906 fingerprint: &str,
1907 reviewer: &str,
1908 reason: &str,
1909 ) -> Result<(), String> {
1910 self.record_admission_decision(fingerprint, ApprovalDecision::Rejected, reviewer, reason)
1911 .await
1912 }
1913
1914 async fn record_admission_decision(
1915 &self,
1916 fingerprint: &str,
1917 decision: ApprovalDecision,
1918 reviewer: &str,
1919 reason: &str,
1920 ) -> Result<(), String> {
1921 {
1922 let mut guard = self.approval_ledger.write().await;
1923 let ledger = guard
1924 .as_mut()
1925 .ok_or_else(|| "no approval ledger installed".to_string())?;
1926 ledger
1927 .record(ApprovalRecord {
1928 fingerprint: fingerprint.to_string(),
1929 // Admission escalations aren't tier-classified; record the
1930 // most restrictive tier so the decision reads as "elevated".
1931 required_tier: PermissionTier::FullAccess,
1932 decision,
1933 reviewer: reviewer.to_string(),
1934 reason: reason.to_string(),
1935 evidence: None,
1936 decided_at: chrono::Utc::now().to_rfc3339(),
1937 })
1938 // A journal write failure means the decision is NOT durable —
1939 // surface it instead of emitting a false ApprovalRecorded
1940 // audit event (review A7).
1941 .map_err(|e| format!("failed to persist approval decision: {e}"))?;
1942 }
1943 let approval = match decision {
1944 ApprovalDecision::Approved => "approved",
1945 ApprovalDecision::Rejected => "rejected",
1946 };
1947 let mut log = self.log.lock().await;
1948 log.append(
1949 EventKind::ApprovalRecorded,
1950 None,
1951 None,
1952 [
1953 ("fingerprint".to_string(), Value::from(fingerprint)),
1954 ("approval".to_string(), Value::from(approval)),
1955 ("reviewer".to_string(), Value::from(reviewer)),
1956 ("reason".to_string(), Value::from(reason)),
1957 ]
1958 .into(),
1959 );
1960 Ok(())
1961 }
1962
1963 /// Back the idempotency cache with a durable JSONL journal (EPIC A / C3).
1964 ///
1965 /// Loads any existing entries into the in-memory cache, then persists
1966 /// future idempotent results (and rollback invalidations, as
1967 /// tombstones) to the journal. After a crash-restart, re-submitting a
1968 /// completed idempotent action returns the cached result instead of
1969 /// re-executing it — preventing duplicate external side effects. The
1970 /// canonical daemon path is `~/.car/idempotency.jsonl`. Returns the
1971 /// number of live entries loaded.
1972 pub async fn set_idempotency_cache_path(
1973 &self,
1974 path: impl Into<std::path::PathBuf>,
1975 ) -> std::io::Result<usize> {
1976 let path = path.into();
1977 if let Some(parent) = path.parent() {
1978 let _ = std::fs::create_dir_all(parent);
1979 }
1980 // Replay the journal (last entry per key wins; a tombstone removes).
1981 let mut loaded: HashMap<String, ActionResult> = HashMap::new();
1982 if path.exists() {
1983 let content = std::fs::read_to_string(&path)?;
1984 for line in content.lines() {
1985 if line.trim().is_empty() {
1986 continue;
1987 }
1988 if let Ok(entry) = serde_json::from_str::<IdempotencyEntry>(line) {
1989 match entry.result {
1990 Some(r) => {
1991 loaded.insert(entry.key, r);
1992 }
1993 None => {
1994 loaded.remove(&entry.key);
1995 }
1996 }
1997 }
1998 }
1999 }
2000 {
2001 let mut cache = self.idempotency_cache.lock().await;
2002 for (k, v) in loaded.iter() {
2003 cache.entry(k.clone()).or_insert_with(|| v.clone());
2004 }
2005 }
2006 let count = loaded.len();
2007 *self.idempotency_journal.write().await = Some(path);
2008 Ok(count)
2009 }
2010
2011 /// Append an idempotency record to the journal if one is configured.
2012 /// `result = None` writes a tombstone (invalidation).
2013 async fn journal_idempotency(&self, key: &str, result: Option<&ActionResult>) {
2014 let path = {
2015 let guard = self.idempotency_journal.read().await;
2016 guard.clone()
2017 };
2018 let Some(path) = path else {
2019 return;
2020 };
2021 let entry = IdempotencyEntry {
2022 key: key.to_string(),
2023 result: result.cloned(),
2024 };
2025 if let Ok(mut line) = serde_json::to_string(&entry) {
2026 line.push('\n');
2027 use std::io::Write;
2028 // Persistence failure must be VISIBLE (linus review): a
2029 // silently-dropped journal write means an idempotent action
2030 // re-executes its side effects after a restart while the
2031 // operator believes it's covered. Fail-open (the in-memory
2032 // cache still dedups this process) but loudly.
2033 let write = std::fs::OpenOptions::new()
2034 .create(true)
2035 .append(true)
2036 .open(&path)
2037 .and_then(|mut f| f.write_all(line.as_bytes()));
2038 if let Err(e) = write {
2039 tracing::warn!(
2040 path = %path.display(),
2041 error = %e,
2042 "idempotency journal write failed — durable dedup is NOT covering this result"
2043 );
2044 }
2045 }
2046 }
2047
2048 /// Look up the current decision for an admission fingerprint, if any.
2049 pub async fn admission_decision(&self, fingerprint: &str) -> Option<ApprovalDecision> {
2050 let guard = self.approval_ledger.read().await;
2051 guard
2052 .as_ref()
2053 .and_then(|l| l.lookup(fingerprint).map(|r| r.decision))
2054 }
2055
2056 /// Verify a model's tool-use claims against the runtime's own execution
2057 /// receipts (EPIC A / A6 — arXiv 2603.10060). The runtime owns tool
2058 /// execution and logs it, so it holds unforgeable ground truth: this
2059 /// projects [`car_eventlog::tool_receipts::ToolReceipt`]s from the event
2060 /// log and cross-checks the supplied claims, catching a fabricated tool
2061 /// reference, a misstated result count, or a false "found nothing".
2062 ///
2063 /// `proposal_id` scopes the cross-check window to a single proposal's
2064 /// events (pass the proposal whose response the claims came from) — a
2065 /// claim is never judged against another run's receipts. The check is
2066 /// retention-coherent (review A6): when the log has trimmed events and
2067 /// the window can't be proven complete (unscoped, or the proposal's
2068 /// `ProposalReceived` marker — which precedes every receipt of that
2069 /// proposal — was itself evicted), a claim without a receipt comes back
2070 /// in `ReceiptReport::ungroundable` ("window evicted") instead of being
2071 /// mis-flagged `fabricated_tool_reference`.
2072 ///
2073 /// Returns the [`car_eventlog::tool_receipts::ReceiptReport`]; when it is not grounded, a
2074 /// `ToolReceiptHallucination` event is emitted so the caller's
2075 /// verdict→action loop (reject/flag the response) is auditable.
2076 /// Deterministic, zero-inference. Claims arrive structured — CAR's
2077 /// thesis is that intent is structured IR, so a caller extracts claims
2078 /// from the model's tool_calls / IR rather than regexing prose.
2079 pub async fn verify_tool_receipts(
2080 &self,
2081 claims: &[car_eventlog::tool_receipts::ToolClaim],
2082 proposal_id: Option<&str>,
2083 ) -> car_eventlog::tool_receipts::ReceiptReport {
2084 let (receipts, window_complete) = {
2085 let log = self.log.lock().await;
2086 let complete = if log.trimmed_events() == 0 {
2087 // Nothing was ever evicted — the window is complete whether
2088 // or not it is scoped.
2089 true
2090 } else {
2091 match proposal_id {
2092 // A proposal's window is complete iff its ProposalReceived
2093 // marker survived retention: every receipt of the proposal
2094 // was appended after it, so if the marker is retained, so
2095 // are the receipts.
2096 Some(pid) => log.events().iter().any(|e| {
2097 e.kind == EventKind::ProposalReceived
2098 && e.proposal_id.as_deref() == Some(pid)
2099 }),
2100 // Unscoped check over a trimmed log: unknowable.
2101 None => false,
2102 }
2103 };
2104 (
2105 car_eventlog::tool_receipts::receipts_from_events_scoped(log.events(), proposal_id),
2106 complete,
2107 )
2108 };
2109 let report = car_eventlog::tool_receipts::verify_tool_claims_windowed(
2110 claims,
2111 &receipts,
2112 window_complete,
2113 );
2114 if !report.grounded {
2115 let mut log = self.log.lock().await;
2116 log.append(
2117 EventKind::ToolReceiptHallucination,
2118 None,
2119 proposal_id,
2120 [
2121 (
2122 "count".to_string(),
2123 Value::from(report.hallucinations.len()),
2124 ),
2125 (
2126 "hallucinations".to_string(),
2127 serde_json::to_value(&report.hallucinations).unwrap_or_default(),
2128 ),
2129 ]
2130 .into(),
2131 );
2132 }
2133 report
2134 }
2135
2136 /// Run every registered admission gate against a proposal and fold
2137 /// their verdicts into a single [`crate::admission::AdmissionDecision`].
2138 ///
2139 /// Each gate's outcome is recorded as an `AdmissionGateDecision` event
2140 /// so a denial is attributable. Aggregation is conjunctive and
2141 /// fail-closed: the proposal is admitted only if every gate allowed it.
2142 /// Returns an admit decision immediately when no gates are registered
2143 /// (zero overhead on the common path).
2144 async fn run_admission_gates(
2145 &self,
2146 proposal: &ActionProposal,
2147 session_id: Option<&str>,
2148 scope: Option<&crate::scope::RuntimeScope>,
2149 ) -> crate::admission::AdmissionDecision {
2150 use crate::admission::{AdmissionDecision, GateContext};
2151
2152 let gates = self.admission_gates.read().await;
2153 if gates.is_empty() {
2154 return AdmissionDecision::admit();
2155 }
2156
2157 // One consistent snapshot for every gate this pass.
2158 let (state, versions) = self.state.versioned_snapshot();
2159 let ctx = GateContext {
2160 session_id,
2161 scope,
2162 state: &state,
2163 versions: &versions,
2164 };
2165
2166 let mut decision = AdmissionDecision::admit();
2167 for gate in gates.iter() {
2168 let outcome = gate.check(proposal, &ctx).await;
2169 // Audit every gate decision (allow included) so the trail shows
2170 // which gates ran, not just which objected.
2171 let mut props: HashMap<String, Value> = HashMap::new();
2172 props.insert("gate".to_string(), Value::from(gate.name()));
2173 // Pre-execution proposal admission (vs the multi-agent
2174 // commit barrier, which emits the same event kind with
2175 // phase:"commit_barrier").
2176 props.insert("phase".to_string(), Value::from("admission"));
2177 props.insert("decision".to_string(), Value::from(outcome.label()));
2178 match &outcome {
2179 crate::admission::GateOutcome::Allow => {}
2180 crate::admission::GateOutcome::Reject { blocked, reason } => {
2181 props.insert("reason".to_string(), Value::from(reason.clone()));
2182 props.insert(
2183 "blocked".to_string(),
2184 serde_json::to_value(blocked).unwrap_or_default(),
2185 );
2186 }
2187 crate::admission::GateOutcome::NeedsApproval {
2188 actions,
2189 fingerprint,
2190 reason,
2191 } => {
2192 props.insert("reason".to_string(), Value::from(reason.clone()));
2193 props.insert(
2194 "blocked".to_string(),
2195 serde_json::to_value(actions).unwrap_or_default(),
2196 );
2197 props.insert("fingerprint".to_string(), Value::from(fingerprint.clone()));
2198 }
2199 }
2200 {
2201 let mut log = self.log.lock().await;
2202 log.append(
2203 EventKind::AdmissionGateDecision,
2204 None,
2205 Some(&proposal.id),
2206 props,
2207 );
2208 }
2209 decision.absorb(gate.name(), outcome);
2210 }
2211 decision
2212 }
2213
2214 /// Install a harness operating config — the live end of the Evolution
2215 /// Agent loop (survey §3.5). After the meta-agent's `HarnessConfig::apply`
2216 /// produces a governed, regression-gated config, hand it here to take
2217 /// effect: `max_retries`/`retry_backoff_ms` drive the per-action retry
2218 /// loop, and `planning_max_replans` is mapped onto the replan budget.
2219 /// Every `HarnessConfig` knob is consumed here — none is aspirational.
2220 pub async fn set_harness_config(&self, cfg: car_memgine::HarnessConfig) {
2221 self.replan_config.write().await.max_replans = cfg.planning_max_replans;
2222 *self.harness_config.write().await = Some(cfg);
2223 }
2224
2225 /// The current harness operating config, if one has been installed.
2226 pub async fn harness_config(&self) -> Option<car_memgine::HarnessConfig> {
2227 self.harness_config.read().await.clone()
2228 }
2229
2230 /// Atomically read-modify-write the harness operating config under ONE
2231 /// write lock (installing the default first when none is set). The
2232 /// get→mutate→set alternative is a lost-update race when two requests
2233 /// mutate concurrently — the daemon's `evolution.run` harness-apply path
2234 /// uses this instead (kernel review S3). Keeps the same replan-budget
2235 /// propagation as [`Self::set_harness_config`].
2236 pub async fn update_harness_config<R>(
2237 &self,
2238 f: impl FnOnce(&mut car_memgine::HarnessConfig) -> R,
2239 ) -> R {
2240 let mut guard = self.harness_config.write().await;
2241 let cfg = guard.get_or_insert_with(car_memgine::HarnessConfig::default);
2242 let out = f(cfg);
2243 let max_replans = cfg.planning_max_replans;
2244 drop(guard);
2245 self.replan_config.write().await.max_replans = max_replans;
2246 out
2247 }
2248
2249 /// Run the pre-execution transactional check against the current
2250 /// versioned shared state. Returns the conflicting action ids when the
2251 /// mode is `Strict` and conflicts exist (so the caller can reject those
2252 /// actions); always emits `TransactionConflict` telemetry for each
2253 /// conflict found. Empty/`None` means "proceed".
2254 async fn transaction_precheck(
2255 &self,
2256 proposal: &ActionProposal,
2257 ) -> std::collections::HashSet<String> {
2258 let mode = *self.transaction_check.read().await;
2259 if mode == TransactionCheckMode::Off {
2260 return std::collections::HashSet::new();
2261 }
2262 let (state, versions) = self.state.versioned_snapshot();
2263 let report = car_verify::check_transaction(proposal, &versions, Some(&state));
2264 if report.consistent {
2265 return std::collections::HashSet::new();
2266 }
2267 let mut blocked = std::collections::HashSet::new();
2268 let mut log = self.log.lock().await;
2269 for c in &report.conflicts {
2270 for aid in &c.actions {
2271 blocked.insert(aid.clone());
2272 }
2273 log.append(
2274 EventKind::TransactionConflict,
2275 c.actions.first().map(|s| s.as_str()),
2276 Some(&proposal.id),
2277 [
2278 // Use the serde representation (snake_case: write_write
2279 // / read_write / stale_assumption) the .d.ts/.pyi/doc
2280 // contract is written against — NOT Debug, which would
2281 // emit "writewrite" (neo review #4).
2282 (
2283 "kind".to_string(),
2284 serde_json::to_value(c.kind).unwrap_or_default(),
2285 ),
2286 ("key".to_string(), Value::from(c.key.clone())),
2287 (
2288 "actions".to_string(),
2289 serde_json::to_value(&c.actions).unwrap_or_default(),
2290 ),
2291 (
2292 "explanation".to_string(),
2293 Value::from(c.explanation.clone()),
2294 ),
2295 ("resolution".to_string(), Value::from(c.resolution.clone())),
2296 ]
2297 .into(),
2298 );
2299 }
2300 drop(log);
2301 // Only Strict blocks execution; Warn records and proceeds.
2302 match mode {
2303 TransactionCheckMode::Strict => blocked,
2304 _ => std::collections::HashSet::new(),
2305 }
2306 }
2307
2308 /// Register a tool with just a name (backward compatible).
2309 pub async fn register_tool(&self, name: &str) {
2310 let schema = ToolSchema {
2311 name: name.to_string(),
2312 source: car_ir::ToolSourceKind::UserDefined,
2313 description: String::new(),
2314 parameters: serde_json::Value::Object(Default::default()),
2315 returns: None,
2316 idempotent: false,
2317 cache_ttl_secs: None,
2318 rate_limit: None,
2319 };
2320 self.register_tool_schema(schema).await;
2321 }
2322
2323 /// Register a tool with full schema.
2324 pub async fn register_tool_schema(&self, schema: ToolSchema) {
2325 // Auto-configure cache if schema specifies it
2326 if let Some(ttl) = schema.cache_ttl_secs {
2327 self.result_cache.enable_caching(&schema.name, ttl).await;
2328 }
2329 // Auto-configure rate limit if schema specifies it
2330 if let Some(ref rl) = schema.rate_limit {
2331 self.rate_limiter
2332 .set_limit(
2333 &schema.name,
2334 RateLimit {
2335 max_calls: rl.max_calls,
2336 interval_secs: rl.interval_secs,
2337 },
2338 )
2339 .await;
2340 }
2341 self.tools.write().await.insert(schema.name.clone(), schema);
2342 }
2343
2344 /// Register a tool via the canonical registry.
2345 /// This is the preferred way to register tools — it updates both the
2346 /// registry and the legacy tools HashMap for backward compatibility.
2347 pub async fn register_tool_entry(&self, mut entry: crate::registry::ToolEntry) {
2348 // The registry's source is authoritative for both event provenance and
2349 // the public tools.list/schema view, including literal-built entries.
2350 entry.schema.source = entry.source.kind();
2351 let schema = entry.schema.clone();
2352 self.registry.register(entry).await;
2353 self.register_tool_schema(schema).await;
2354 }
2355
2356 /// Remove a tool from both the canonical registry and the legacy
2357 /// `tools` schema map, so the model no longer sees it and the
2358 /// validator no longer accepts it. Used when a remote MCP connector
2359 /// tool is disabled or its connector is removed. Returns true if the
2360 /// tool was present in either store.
2361 pub async fn unregister_tool(&self, name: &str) -> bool {
2362 let removed_entry = self.registry.remove(name).await.is_some();
2363 let removed_schema = self.tools.write().await.remove(name).is_some();
2364 removed_entry || removed_schema
2365 }
2366
2367 /// Register CAR's built-in agent utility stdlib.
2368 ///
2369 /// This is an opt-in convenience layer for common local-file and text tools.
2370 /// Existing runtimes remain unchanged until this is called.
2371 pub async fn register_agent_basics(&self) {
2372 for entry in crate::agent_basics::entries() {
2373 self.register_tool_entry(entry).await;
2374 }
2375 }
2376
2377 /// Get all registered tool schemas (for model prompt generation).
2378 pub async fn tool_schemas(&self) -> Vec<ToolSchema> {
2379 self.tools.read().await.values().cloned().collect()
2380 }
2381
2382 /// Set a cost budget that limits proposal execution.
2383 pub async fn set_cost_budget(&self, budget: CostBudget) {
2384 *self.cost_budget.write().await = Some(budget);
2385 }
2386
2387 /// Set per-agent capability permissions that restrict tools, state keys, and action count.
2388 pub async fn set_capabilities(&self, caps: CapabilitySet) {
2389 *self.capabilities.write().await = Some(caps);
2390 }
2391
2392 /// Set a per-tool rate limit (token bucket).
2393 ///
2394 /// `max_calls` tokens are available per `interval_secs` window.
2395 /// When the bucket is empty, `dispatch()` applies backpressure by
2396 /// waiting until a token refills.
2397 pub async fn set_rate_limit(&self, tool: &str, max_calls: u32, interval_secs: f64) {
2398 self.rate_limiter
2399 .set_limit(
2400 tool,
2401 RateLimit {
2402 max_calls,
2403 interval_secs,
2404 },
2405 )
2406 .await;
2407 }
2408
2409 /// Enable cross-proposal result caching for a tool with a TTL in seconds.
2410 pub async fn enable_tool_cache(&self, tool: &str, ttl_secs: u64) {
2411 self.result_cache.enable_caching(tool, ttl_secs).await;
2412 }
2413
2414 /// Execute a proposal with automatic replanning on failure.
2415 ///
2416 /// If a `ReplanCallback` is registered and `max_replans > 0`, the runtime
2417 /// will catch abort failures, roll back state, ask the model for an
2418 /// alternative proposal via the callback, and re-execute. This transforms
2419 /// "execute-and-hope" into "execute-and-recover."
2420 ///
2421 /// If no callback is registered or `max_replans == 0`, behaves identically
2422 /// to a single `execute_inner()` call (zero overhead, fully backward compatible).
2423 #[instrument(
2424 name = "proposal.execute",
2425 skip_all,
2426 fields(
2427 proposal_id = %proposal.id,
2428 action_count = proposal.actions.len(),
2429 )
2430 )]
2431 pub async fn execute(&self, proposal: &ActionProposal) -> ProposalResult {
2432 // Forward to the cancel-aware variant with a never-cancelled
2433 // token. Existing callers see no behaviour change.
2434 let token = tokio_util::sync::CancellationToken::new();
2435 self.execute_with_cancel(proposal, &token).await
2436 }
2437
2438 /// Execute a proposal scoped to a specific session id.
2439 ///
2440 /// Validation walks the global policy registry plus the session's
2441 /// own registry — both must pass for an action to run. Session
2442 /// policies can deny what global allows; they cannot allow what
2443 /// global denies (validation is conjunctive).
2444 ///
2445 /// Returns the same [`ProposalResult`] shape as [`Self::execute`].
2446 /// Errors with an action-level rejection if the session id is
2447 /// unknown — callers should check via [`Self::session_exists`] or
2448 /// trust an id they minted via [`Self::open_session`].
2449 pub async fn execute_with_session(
2450 &self,
2451 proposal: &ActionProposal,
2452 session_id: &str,
2453 ) -> ProposalResult {
2454 let token = tokio_util::sync::CancellationToken::new();
2455 self.execute_with_session_and_cancel(proposal, session_id, &token)
2456 .await
2457 }
2458
2459 /// Combined session-scoped + cancellable execute. The session id
2460 /// is passed verbatim to the per-action policy check; the cancel
2461 /// token behaves identically to [`Self::execute_with_cancel`].
2462 pub async fn execute_with_session_and_cancel(
2463 &self,
2464 proposal: &ActionProposal,
2465 session_id: &str,
2466 cancel: &tokio_util::sync::CancellationToken,
2467 ) -> ProposalResult {
2468 self.execute_with_optional_session(proposal, Some(session_id), None, None, cancel)
2469 .await
2470 }
2471
2472 /// Execute an active-run proposal while requiring every accepted replan
2473 /// to retain the authenticated proposal id already claimed by the server.
2474 pub async fn execute_with_session_and_stable_replan_id(
2475 &self,
2476 proposal: &ActionProposal,
2477 session_id: &str,
2478 ) -> ProposalResult {
2479 let token = tokio_util::sync::CancellationToken::new();
2480 self.execute_with_optional_session(
2481 proposal,
2482 Some(session_id),
2483 None,
2484 Some(&proposal.id),
2485 &token,
2486 )
2487 .await
2488 }
2489
2490 /// Execute a proposal with cooperative cancellation.
2491 ///
2492 /// The runtime checks `token.is_cancelled()` at each DAG level
2493 /// boundary. When set, every action that hadn't yet started runs
2494 /// is reported as `Skipped` with `error = "canceled: ..."` so
2495 /// callers can distinguish "user pulled the plug" from "earlier
2496 /// abort cascaded." Actions already in flight continue to
2497 /// completion — tool calls dispatched to user-provided executors
2498 /// can't be safely interrupted from the engine.
2499 ///
2500 /// The CAR A2A bridge uses this so `tasks/cancel` produces a
2501 /// `ProposalResult` with clean partial state rather than relying
2502 /// on `JoinHandle::abort` to interrupt mid-await (which leaves
2503 /// no record of which actions actually ran).
2504 ///
2505 /// **FFI exposure:** this method is intentionally not surfaced
2506 /// through the NAPI / PyO3 / `car-server-core` JSON-RPC bindings.
2507 /// Those consumers (Node, Python, WebSocket) don't currently
2508 /// expose long-running async-task surfaces that need
2509 /// cancellation; the bridge is the lone consumer. When a binding
2510 /// gains a long-running task surface, the path is clear: add a
2511 /// per-binding token registry keyed by some caller-provided id,
2512 /// expose `cancelExecution(id)` / `cancel_execution(id)` /
2513 /// `proposal.cancel { id }`, and have the runtime call
2514 /// `execute_with_cancel` with the matching token. Skipping that
2515 /// today avoids speculative API surface that bloats bindings
2516 /// without a consumer.
2517 pub async fn execute_with_cancel(
2518 &self,
2519 proposal: &ActionProposal,
2520 cancel: &tokio_util::sync::CancellationToken,
2521 ) -> ProposalResult {
2522 self.execute_with_optional_session(proposal, None, None, None, cancel)
2523 .await
2524 }
2525
2526 pub async fn execute_with_stable_replan_id(&self, proposal: &ActionProposal) -> ProposalResult {
2527 let token = tokio_util::sync::CancellationToken::new();
2528 self.execute_with_optional_session(proposal, None, None, Some(&proposal.id), &token)
2529 .await
2530 }
2531
2532 /// Execute a proposal with an attached [`RuntimeScope`](crate::scope::RuntimeScope)
2533 /// (Parslee-ai/car#187 phase 3).
2534 ///
2535 /// Same contract as [`Self::execute`] plus a per-execution
2536 /// identity surface — typically built by the car-a2a dispatcher
2537 /// from the verified `Identity` and cooperative `a2a_caller`
2538 /// metadata on the inbound `ActionProposal`. The scope is
2539 /// recorded on the event log so downstream audit / log analysis
2540 /// can see which caller / tenant issued each action.
2541 ///
2542 /// **What this enforces today**: scope is captured + logged.
2543 /// Memgine queries and state-store ops still hit global
2544 /// namespaces — those follow-ups are tracked under #187.
2545 /// Tool / policy code that needs per-tenant behaviour right now
2546 /// should keep reading `proposal.context["a2a_caller_verified"]`
2547 /// directly (the phase 1 / 2 surface).
2548 pub async fn execute_scoped(
2549 &self,
2550 proposal: &ActionProposal,
2551 scope: &crate::scope::RuntimeScope,
2552 ) -> ProposalResult {
2553 let token = tokio_util::sync::CancellationToken::new();
2554 self.execute_scoped_with_cancel(proposal, scope, &token)
2555 .await
2556 }
2557
2558 /// Combined scoped + cancellable execute. Mirrors the shape of
2559 /// [`Self::execute_with_session_and_cancel`] for symmetry — both
2560 /// add a side-channel (session id / scope) on top of the
2561 /// cancellable form.
2562 pub async fn execute_scoped_with_cancel(
2563 &self,
2564 proposal: &ActionProposal,
2565 scope: &crate::scope::RuntimeScope,
2566 cancel: &tokio_util::sync::CancellationToken,
2567 ) -> ProposalResult {
2568 self.execute_with_optional_session(proposal, None, Some(scope), None, cancel)
2569 .await
2570 }
2571
2572 pub async fn execute_scoped_with_stable_replan_id(
2573 &self,
2574 proposal: &ActionProposal,
2575 scope: &crate::scope::RuntimeScope,
2576 ) -> ProposalResult {
2577 let token = tokio_util::sync::CancellationToken::new();
2578 self.execute_with_optional_session(proposal, None, Some(scope), Some(&proposal.id), &token)
2579 .await
2580 }
2581
2582 /// Internal entry point that backs both
2583 /// [`Self::execute_with_cancel`] (no session) and
2584 /// [`Self::execute_with_session_and_cancel`]. Holds the replan
2585 /// loop and threads `session_id` into the per-action validation
2586 /// path so session-scoped policies stack on top of global ones.
2587 async fn execute_with_optional_session(
2588 &self,
2589 proposal: &ActionProposal,
2590 session_id: Option<&str>,
2591 scope: Option<&crate::scope::RuntimeScope>,
2592 required_replan_proposal_id: Option<&str>,
2593 cancel: &tokio_util::sync::CancellationToken,
2594 ) -> ProposalResult {
2595 // The shared StateStore has one proposal transaction at a time. The
2596 // guard spans validation, replanning, and rollback so distinct Runtime
2597 // facades over the same store cannot interleave state snapshots.
2598 let _proposal_execution = self.state.lock_proposal_execution().await;
2599 self.execute_with_optional_session_already_guarded(
2600 proposal,
2601 session_id,
2602 scope,
2603 required_replan_proposal_id,
2604 cancel,
2605 )
2606 .await
2607 }
2608
2609 /// Execute while the caller retains this StateStore's proposal-execution
2610 /// guard. `plan_and_execute` uses this path so one guard covers its initial
2611 /// snapshot, ranking, every candidate, and the final commit/rollback.
2612 async fn execute_with_optional_session_already_guarded(
2613 &self,
2614 proposal: &ActionProposal,
2615 session_id: Option<&str>,
2616 scope: Option<&crate::scope::RuntimeScope>,
2617 required_replan_proposal_id: Option<&str>,
2618 cancel: &tokio_util::sync::CancellationToken,
2619 ) -> ProposalResult {
2620 // Establish canonical identity before any other rejection so an
2621 // undigested lineage entry can only mean the proposal was genuinely
2622 // impossible to represent as JCS/I-JSON. No journal row is emitted.
2623 let original_proposal_digest = match proposal_digest(proposal) {
2624 Ok(digest) => digest,
2625 Err(error) => {
2626 self.log.lock().await.append(
2627 EventKind::StateRollback,
2628 None,
2629 Some(&proposal.id),
2630 proposal_rejection_boundary_data(proposal, None, &error),
2631 );
2632 let result = ProposalResult::for_proposal(
2633 proposal,
2634 proposal
2635 .actions
2636 .iter()
2637 .map(|action| rejected_result(&action.id, error.clone()))
2638 .collect(),
2639 CostSummary::default(),
2640 );
2641 return finalize_proposal_result(
2642 result,
2643 &proposal.id,
2644 &[ProposalLineageEntry {
2645 generation: 0,
2646 proposal_id: proposal.id.clone(),
2647 proposal_digest: None,
2648 status: ProposalLineageStatus::Rejected,
2649 rejection_reason: Some(error),
2650 }],
2651 &[],
2652 );
2653 }
2654 };
2655
2656 // Proposal-local action ids are the join key for DAG scheduling,
2657 // result rows, receipts, and StateTransition attribution. Reject a
2658 // duplicate before scope/admission/transaction logging or any state
2659 // mutation. Retried attempts remain valid reuse of one admitted id.
2660 if let Err(error) = validate_proposal_action_ids(proposal) {
2661 self.log.lock().await.append(
2662 EventKind::StateRollback,
2663 None,
2664 Some(&proposal.id),
2665 proposal_rejection_boundary_data(proposal, Some(&original_proposal_digest), &error),
2666 );
2667 let result = ProposalResult::for_proposal(
2668 proposal,
2669 proposal
2670 .actions
2671 .iter()
2672 .map(|action| rejected_result(&action.id, error.clone()))
2673 .collect(),
2674 CostSummary::default(),
2675 );
2676 return finalize_proposal_result(
2677 result,
2678 &proposal.id,
2679 &[proposal_lineage_entry(
2680 proposal,
2681 0,
2682 ProposalLineageStatus::Rejected,
2683 Some(error),
2684 )],
2685 &[],
2686 );
2687 }
2688 if let Err(error) = validate_proposal_retry_limits(proposal) {
2689 self.log.lock().await.append(
2690 EventKind::StateRollback,
2691 None,
2692 Some(&proposal.id),
2693 proposal_rejection_boundary_data(proposal, Some(&original_proposal_digest), &error),
2694 );
2695 let result = ProposalResult::for_proposal(
2696 proposal,
2697 proposal
2698 .actions
2699 .iter()
2700 .map(|action| rejected_result(&action.id, error.clone()))
2701 .collect(),
2702 CostSummary::default(),
2703 );
2704 return finalize_proposal_result(
2705 result,
2706 &proposal.id,
2707 &[proposal_lineage_entry(
2708 proposal,
2709 0,
2710 ProposalLineageStatus::Rejected,
2711 Some(error),
2712 )],
2713 &[],
2714 );
2715 }
2716 let config = self.replan_config.read().await.clone();
2717 let mut current_proposal = proposal.clone();
2718 let mut attempt: u32 = 0;
2719 let mut lineage = vec![proposal_lineage_entry(
2720 proposal,
2721 0,
2722 ProposalLineageStatus::Accepted,
2723 None,
2724 )];
2725
2726 // Phase 3 foundation (Parslee-ai/car#187): record the scope
2727 // on the event log so audit / log analysis can correlate
2728 // actions to the caller / tenant that triggered them. Only
2729 // logged when at least one identity field is set — keeps
2730 // the existing in-process call sites free of noise.
2731 if let Some(s) = scope {
2732 if !s.is_unscoped() {
2733 let mut props: HashMap<String, Value> = HashMap::new();
2734 if let Some(cid) = &s.caller_id {
2735 props.insert("caller_id".to_string(), Value::from(cid.as_str()));
2736 }
2737 if let Some(tid) = &s.tenant_id {
2738 props.insert("tenant_id".to_string(), Value::from(tid.as_str()));
2739 }
2740 if !s.claims.is_empty() {
2741 if let Ok(claims_json) = serde_json::to_value(&s.claims) {
2742 props.insert("claims".to_string(), claims_json);
2743 }
2744 }
2745 let mut log = self.log.lock().await;
2746 log.append(EventKind::SessionScope, None, Some(&proposal.id), props);
2747 }
2748 }
2749
2750 // Pre-execution transactional conflict check (survey §4.3/§5.2.4).
2751 // Off by default; Warn records conflicts; Strict rejects the whole
2752 // proposal before any action runs, since a transactional conflict is
2753 // a property of the action *set* against current state, not an
2754 // isolated action. This is an advisory planning-time gate, not a
2755 // substitute for per-action execution-time validation. Replanned
2756 // proposals are re-checked inside the loop (at the replan quality
2757 // gate below). A pre-execution rejection here deliberately produces
2758 // no execution trajectory — nothing ran; the emitted
2759 // `TransactionConflict` events are the audit record.
2760 let blocked = self.transaction_precheck(¤t_proposal).await;
2761 if !blocked.is_empty() {
2762 let reason = "transactional conflict with current shared state (strict mode); \
2763 see TransactionConflict events for details and resolution"
2764 .to_string();
2765 let results = current_proposal
2766 .actions
2767 .iter()
2768 .map(|a| rejected_result(&a.id, reason.clone()))
2769 .collect();
2770 self.log.lock().await.append(
2771 EventKind::StateRollback,
2772 None,
2773 Some(¤t_proposal.id),
2774 proposal_rejection_boundary_data(
2775 ¤t_proposal,
2776 Some(&original_proposal_digest),
2777 &reason,
2778 ),
2779 );
2780 lineage[0].status = ProposalLineageStatus::Rejected;
2781 lineage[0].rejection_reason = Some(reason);
2782 return finalize_proposal_result(
2783 ProposalResult::for_proposal(proposal, results, Default::default()),
2784 &proposal.id,
2785 &lineage,
2786 &[],
2787 );
2788 }
2789
2790 // Pre-execution admission gates (EPIC A / task A1 — the safety
2791 // seam). Runs every registered AdmissionGate against the proposal
2792 // before any action executes. Aggregation is conjunctive and
2793 // fail-closed: a proposal any gate blocks (or escalates to
2794 // approval) does not run. Like the transactional pre-check above, a
2795 // rejection here produces no execution trajectory — the emitted
2796 // AdmissionGateDecision events are the audit record. No gates
2797 // registered → zero overhead, identical behavior.
2798 let admission = self
2799 .run_admission_gates(¤t_proposal, session_id, scope)
2800 .await;
2801 if !admission.admitted {
2802 let gate = admission.deciding_gate.as_deref().unwrap_or("admission");
2803 let base_reason = admission
2804 .reason
2805 .clone()
2806 .unwrap_or_else(|| "blocked by admission gate".to_string());
2807 // Resolve approval escalations against the durable ledger (A7).
2808 // A hard `Reject` from ANY gate is never overridable — the
2809 // ledger is not consulted at all in that case (an old approval
2810 // for one gate's escalation must not steamroll another gate's
2811 // deny). Otherwise EVERY escalation must resolve to Approved,
2812 // each by its own fingerprint: one Rejected fingerprint blocks,
2813 // one unseen fingerprint stays pending (fail-closed).
2814 let mut approved = false;
2815 let reason = if admission.hard_rejected {
2816 format!("{base_reason} (gate: {gate})")
2817 } else if admission.needs_approval() {
2818 let mut blocking_reason = None;
2819 for esc in &admission.escalations {
2820 match self.admission_decision(&esc.fingerprint).await {
2821 Some(ApprovalDecision::Approved) => continue,
2822 Some(ApprovalDecision::Rejected) => {
2823 blocking_reason = Some(format!(
2824 "rejected by operator (gate: {}; fingerprint: {})",
2825 esc.gate, esc.fingerprint
2826 ));
2827 break;
2828 }
2829 None => {
2830 blocking_reason = Some(format!(
2831 "requires human approval (gate: {}; {}); \
2832 approve fingerprint '{}' then re-run",
2833 esc.gate, esc.reason, esc.fingerprint
2834 ));
2835 break;
2836 }
2837 }
2838 }
2839 match blocking_reason {
2840 None => {
2841 approved = true;
2842 // Audit every durable approval taking effect.
2843 let mut log = self.log.lock().await;
2844 for esc in &admission.escalations {
2845 log.append(
2846 EventKind::ApprovalRecorded,
2847 None,
2848 Some(&proposal.id),
2849 [
2850 (
2851 "fingerprint".to_string(),
2852 Value::from(esc.fingerprint.as_str()),
2853 ),
2854 ("gate".to_string(), Value::from(esc.gate.as_str())),
2855 ("approval".to_string(), Value::from("approved")),
2856 ("applied".to_string(), Value::from(true)),
2857 ]
2858 .into(),
2859 );
2860 }
2861 String::new()
2862 }
2863 Some(r) => r,
2864 }
2865 } else {
2866 format!("{base_reason} (gate: {gate})")
2867 };
2868 if approved {
2869 // Escalation cleared by a durable approval — proceed to
2870 // execution as if admitted.
2871 } else {
2872 let results = current_proposal
2873 .actions
2874 .iter()
2875 .map(|a| {
2876 // Name the specific reason on the offending actions; a
2877 // generic note on the rest (the whole proposal is held,
2878 // since a safety hazard is a property of the set).
2879 if admission.blocked.is_empty() || admission.blocked.contains(&a.id) {
2880 rejected_result(&a.id, reason.clone())
2881 } else {
2882 rejected_result(
2883 &a.id,
2884 format!("proposal blocked by admission gate: {gate}"),
2885 )
2886 }
2887 })
2888 .collect();
2889 self.log.lock().await.append(
2890 EventKind::StateRollback,
2891 None,
2892 Some(¤t_proposal.id),
2893 proposal_rejection_boundary_data(
2894 ¤t_proposal,
2895 Some(&original_proposal_digest),
2896 &reason,
2897 ),
2898 );
2899 lineage[0].status = ProposalLineageStatus::Rejected;
2900 lineage[0].rejection_reason = Some(reason.clone());
2901 let result = finalize_proposal_result(
2902 ProposalResult::for_proposal(proposal, results, Default::default()),
2903 &proposal.id,
2904 &lineage,
2905 &[],
2906 );
2907 // Record the trajectory even though nothing dispatched. This
2908 // return is *before* the execution loop's persist calls, so
2909 // without this an admission rejection is invisible to the
2910 // trajectory store — and the per-tool success rates
2911 // `verify.monte_carlo` derives from it would silently skew
2912 // optimistic, counting only calls that got far enough to run.
2913 // The state map is empty because no action mutated anything.
2914 if let Some(err) = self.persist_trajectory(
2915 proposal,
2916 ¤t_proposal,
2917 &result,
2918 car_memgine::TrajectoryOutcome::Failed,
2919 0,
2920 &HashMap::new(),
2921 ) {
2922 tracing::warn!(
2923 error = %err,
2924 proposal_id = %proposal.id,
2925 "failed to persist trajectory for admission-rejected proposal"
2926 );
2927 }
2928 return finalize_proposal_result(result, &proposal.id, &lineage, &[]);
2929 }
2930 }
2931
2932 let mut accepted_proposal_preimages = vec![AcceptedProposalPreimage {
2933 generation: 0,
2934 proposal_digest: original_proposal_digest,
2935 proposal: proposal.clone(),
2936 }];
2937
2938 loop {
2939 let (result, state_before_map) = self
2940 .execute_inner_with_cancel(¤t_proposal, Some(cancel), session_id, scope)
2941 .await;
2942
2943 // Statuses that trigger rollback + replan: always runtime Failed,
2944 // and (opt-in) validator/policy/capability Rejected.
2945 let replan_triggers = |s: &ActionStatus| {
2946 *s == ActionStatus::Failed
2947 || (config.replan_on_rejected && *s == ActionStatus::Rejected)
2948 };
2949
2950 // A terminal tool failure ends this engine execution. It may not
2951 // enter the retry loop above or the proposal-replan loop here;
2952 // daemon-owned session halting is layered on top by the caller.
2953 let terminal_failure = result.results.iter().any(|result| result.terminal);
2954 // Check if we aborted
2955 let aborted = result.results.iter().any(|r| replan_triggers(&r.status));
2956 let rollback_durability_failed = result.results.iter().any(|result| {
2957 result.error.as_deref().is_some_and(|error| {
2958 error.contains(ROLLBACK_DURABILITY_ERROR)
2959 || error.contains(ROLLBACK_DURABILITY_UNKNOWN)
2960 })
2961 });
2962 if !aborted
2963 || terminal_failure
2964 || attempt >= config.max_replans
2965 || rollback_durability_failed
2966 {
2967 if aborted && !terminal_failure && attempt > 0 && !rollback_durability_failed {
2968 // Exhausted all replan attempts
2969 let mut log = self.log.lock().await;
2970 log.append(
2971 EventKind::ReplanExhausted,
2972 None,
2973 Some(&proposal.id),
2974 [("attempts".to_string(), Value::from(attempt))].into(),
2975 );
2976 }
2977
2978 // Persist trajectory
2979 let outcome = if !aborted {
2980 if attempt > 0 {
2981 car_memgine::TrajectoryOutcome::ReplanSuccess
2982 } else {
2983 car_memgine::TrajectoryOutcome::Success
2984 }
2985 } else if attempt > 0 && !terminal_failure {
2986 car_memgine::TrajectoryOutcome::ReplanExhausted
2987 } else {
2988 car_memgine::TrajectoryOutcome::Failed
2989 };
2990 if let Some(err) = self.persist_trajectory(
2991 proposal,
2992 ¤t_proposal,
2993 &result,
2994 outcome,
2995 attempt,
2996 &state_before_map,
2997 ) {
2998 let mut log = self.log.lock().await;
2999 log.append(
3000 EventKind::ActionFailed,
3001 None,
3002 Some(&proposal.id),
3003 [(
3004 "trajectory_persist_error".to_string(),
3005 Value::from(err.as_str()),
3006 )]
3007 .into(),
3008 );
3009 }
3010
3011 return finalize_proposal_result(
3012 result,
3013 &proposal.id,
3014 &lineage,
3015 &accepted_proposal_preimages,
3016 );
3017 }
3018
3019 // Get replan callback (clone Arc, drop lock immediately)
3020 let callback = {
3021 let guard = self.replan_callback.lock().await;
3022 guard.clone()
3023 };
3024 let Some(callback) = callback else {
3025 // No callback registered — persist trajectory and return
3026 if let Some(err) = self.persist_trajectory(
3027 proposal,
3028 ¤t_proposal,
3029 &result,
3030 car_memgine::TrajectoryOutcome::Failed,
3031 attempt,
3032 &state_before_map,
3033 ) {
3034 let mut log = self.log.lock().await;
3035 log.append(
3036 EventKind::ActionFailed,
3037 None,
3038 Some(&proposal.id),
3039 [(
3040 "trajectory_persist_error".to_string(),
3041 Value::from(err.as_str()),
3042 )]
3043 .into(),
3044 );
3045 }
3046 return finalize_proposal_result(
3047 result,
3048 &proposal.id,
3049 &lineage,
3050 &accepted_proposal_preimages,
3051 );
3052 };
3053
3054 // Build the stable failure context once. Candidate rejection is a
3055 // planning-only loop below: it consumes a replan generation but
3056 // never returns to `execute_inner_with_cancel`, so a failed
3057 // proposal with external effects cannot be dispatched twice.
3058 let failed_actions: Vec<FailedActionSummary> = result
3059 .results
3060 .iter()
3061 .filter(|r| replan_triggers(&r.status) && !r.rolled_back)
3062 .map(|r| {
3063 let action = current_proposal
3064 .actions
3065 .iter()
3066 .find(|a| a.id == r.action_id);
3067 FailedActionSummary {
3068 action_id: r.action_id.clone(),
3069 tool: action.and_then(|a| a.tool.clone()),
3070 error: r.error.clone().unwrap_or_default(),
3071 parameters: action.map(|a| a.parameters.clone()).unwrap_or_default(),
3072 }
3073 })
3074 .collect();
3075
3076 let completed_action_ids: Vec<String> = result
3077 .results
3078 .iter()
3079 .filter(|r| r.rolled_back)
3080 .map(|r| r.action_id.clone())
3081 .collect();
3082
3083 'replan_candidates: loop {
3084 if attempt >= config.max_replans {
3085 let mut log = self.log.lock().await;
3086 log.append(
3087 EventKind::ReplanExhausted,
3088 None,
3089 Some(&proposal.id),
3090 [("attempts".to_string(), Value::from(attempt))].into(),
3091 );
3092 drop(log);
3093 if let Some(err) = self.persist_trajectory(
3094 proposal,
3095 ¤t_proposal,
3096 &result,
3097 car_memgine::TrajectoryOutcome::ReplanExhausted,
3098 attempt,
3099 &state_before_map,
3100 ) {
3101 self.log.lock().await.append(
3102 EventKind::ActionFailed,
3103 None,
3104 Some(&proposal.id),
3105 [("trajectory_persist_error".to_string(), Value::from(err))].into(),
3106 );
3107 }
3108 return finalize_proposal_result(
3109 result,
3110 &proposal.id,
3111 &lineage,
3112 &accepted_proposal_preimages,
3113 );
3114 }
3115
3116 let generation = attempt + 1;
3117 let ctx = ReplanContext {
3118 proposal_id: proposal.id.clone(),
3119 attempt: generation,
3120 failed_actions: failed_actions.clone(),
3121 completed_action_ids: completed_action_ids.clone(),
3122 state_snapshot: self.state.snapshot(),
3123 replans_remaining: config.max_replans.saturating_sub(generation),
3124 original_source: proposal.source.clone(),
3125 original_action_count: proposal.actions.len(),
3126 original_context: proposal.context.clone(),
3127 };
3128
3129 // Backoff delay between replan attempts
3130 if config.delay_ms > 0 {
3131 tokio::time::sleep(Duration::from_millis(config.delay_ms)).await;
3132 }
3133
3134 // Log replan attempt
3135 {
3136 let mut log = self.log.lock().await;
3137 log.append(
3138 EventKind::ReplanAttempted,
3139 None,
3140 Some(&proposal.id),
3141 [
3142 ("attempt".to_string(), Value::from(generation)),
3143 (
3144 "failed_count".to_string(),
3145 Value::from(ctx.failed_actions.len()),
3146 ),
3147 ]
3148 .into(),
3149 );
3150 // Deep-telemetry breadcrumbs (§3.5.1): record the fork the
3151 // harness took (replan vs accept vs abandon) and the
3152 // approach it discarded, so failure-mode diagnosis can see
3153 // the path *not* taken, not just the path taken.
3154 let rejected_ids: Vec<String> = ctx
3155 .failed_actions
3156 .iter()
3157 .map(|f| f.action_id.clone())
3158 .collect();
3159 log.append(
3160 EventKind::BranchDecision,
3161 None,
3162 Some(&proposal.id),
3163 [
3164 ("branch".to_string(), Value::from("replan")),
3165 (
3166 "reason".to_string(),
3167 Value::from("actions failed; requesting a revised plan"),
3168 ),
3169 ("attempt".to_string(), Value::from(generation)),
3170 ]
3171 .into(),
3172 );
3173 log.append(
3174 EventKind::AlternativeRejected,
3175 None,
3176 Some(&proposal.id),
3177 [
3178 (
3179 "alternative".to_string(),
3180 serde_json::to_value(&rejected_ids).unwrap_or_default(),
3181 ),
3182 (
3183 "reason".to_string(),
3184 Value::from("plan superseded by replan after action failure"),
3185 ),
3186 ]
3187 .into(),
3188 );
3189 }
3190
3191 // Call the model for a new plan
3192 match callback.replan(&ctx).await {
3193 Ok(new_proposal) => {
3194 attempt = generation;
3195 let (new_proposal_preimage, new_proposal_digest) =
3196 match proposal_journal_identity(&new_proposal) {
3197 Ok(identity) => identity,
3198 Err(error) => {
3199 lineage.push(ProposalLineageEntry {
3200 generation,
3201 proposal_id: new_proposal.id.clone(),
3202 proposal_digest: None,
3203 status: ProposalLineageStatus::Rejected,
3204 rejection_reason: Some(error.clone()),
3205 });
3206 self.log.lock().await.append(
3207 EventKind::ReplanRejected,
3208 None,
3209 Some(&proposal.id),
3210 [
3211 ("errors".to_string(), Value::from(error)),
3212 ("attempt".to_string(), Value::from(generation)),
3213 ]
3214 .into(),
3215 );
3216 continue 'replan_candidates;
3217 }
3218 };
3219 if required_replan_proposal_id
3220 .is_some_and(|required| new_proposal.id != required)
3221 {
3222 let reason = format!(
3223 "replan proposal id `{}` does not retain authenticated proposal id `{}`",
3224 new_proposal.id,
3225 required_replan_proposal_id.expect("checked above")
3226 );
3227 lineage.push(ProposalLineageEntry {
3228 generation,
3229 proposal_id: new_proposal.id.clone(),
3230 proposal_digest: Some(new_proposal_digest),
3231 status: ProposalLineageStatus::Rejected,
3232 rejection_reason: Some(reason.clone()),
3233 });
3234 self.log.lock().await.append(
3235 EventKind::ReplanRejected,
3236 None,
3237 Some(&proposal.id),
3238 [
3239 ("errors".to_string(), Value::from(reason)),
3240 ("attempt".to_string(), Value::from(generation)),
3241 ]
3242 .into(),
3243 );
3244 continue 'replan_candidates;
3245 }
3246 if let Err(error) = validate_proposal_action_ids(&new_proposal) {
3247 lineage.push(ProposalLineageEntry {
3248 generation,
3249 proposal_id: new_proposal.id.clone(),
3250 proposal_digest: Some(new_proposal_digest),
3251 status: ProposalLineageStatus::Rejected,
3252 rejection_reason: Some(error.clone()),
3253 });
3254 let mut log = self.log.lock().await;
3255 log.append(
3256 EventKind::ReplanRejected,
3257 None,
3258 Some(&proposal.id),
3259 [
3260 (
3261 "errors".to_string(),
3262 Value::from(format!("invalid proposal: {error}")),
3263 ),
3264 ("attempt".to_string(), Value::from(generation)),
3265 ]
3266 .into(),
3267 );
3268 continue 'replan_candidates;
3269 }
3270
3271 if let Err(error) = validate_proposal_retry_limits(&new_proposal) {
3272 lineage.push(ProposalLineageEntry {
3273 generation,
3274 proposal_id: new_proposal.id.clone(),
3275 proposal_digest: Some(new_proposal_digest),
3276 status: ProposalLineageStatus::Rejected,
3277 rejection_reason: Some(error.clone()),
3278 });
3279 self.log.lock().await.append(
3280 EventKind::ReplanRejected,
3281 None,
3282 Some(&proposal.id),
3283 [
3284 ("errors".to_string(), Value::from(error)),
3285 ("attempt".to_string(), Value::from(generation)),
3286 ]
3287 .into(),
3288 );
3289 continue 'replan_candidates;
3290 }
3291
3292 // Quality gate: verify replan proposal before executing
3293 if config.verify_before_execute {
3294 let current_state = self.state.snapshot();
3295 // Verify against the full registered schemas
3296 // (not just names) so a replan with a bad
3297 // parameter type / missing required field is
3298 // rejected here, per register_tool_schema's
3299 // contract (car-releases#56). The read guard is
3300 // held across the synchronous verify call.
3301 // Same helper the admission gate uses, so the two
3302 // verification points cannot disagree about what is
3303 // fatal. They previously did: this path blocked on the
3304 // loop heuristic and on state-dependent findings that
3305 // the gate treats as advisory, and it passed the tool
3306 // map through unconditionally — so an embedder with no
3307 // registered schemas burned its whole replan budget on
3308 // "unregistered tool" for every call.
3309 let blocking = {
3310 let tools_guard = self.tools.read().await;
3311 crate::verify_gate::blocking_errors(
3312 &new_proposal,
3313 Some(¤t_state),
3314 &tools_guard,
3315 100,
3316 )
3317 };
3318 if !blocking.is_empty() {
3319 let error_msgs: Vec<String> =
3320 blocking.iter().map(|i| i.message.clone()).collect();
3321 let reason = error_msgs.join("; ");
3322 lineage.push(ProposalLineageEntry {
3323 generation,
3324 proposal_id: new_proposal.id.clone(),
3325 proposal_digest: Some(new_proposal_digest),
3326 status: ProposalLineageStatus::Rejected,
3327 rejection_reason: Some(reason.clone()),
3328 });
3329 let mut log = self.log.lock().await;
3330 log.append(
3331 EventKind::ReplanRejected,
3332 None,
3333 Some(&proposal.id),
3334 [
3335 ("errors".to_string(), Value::from(reason)),
3336 ("attempt".to_string(), Value::from(generation)),
3337 ]
3338 .into(),
3339 );
3340 // Don't execute a broken replan — count as failed attempt
3341 continue 'replan_candidates;
3342 }
3343 }
3344
3345 // Transactional re-check of the replan against the now-
3346 // mutated shared state (neo review #1): a replan runs
3347 // after earlier actions changed state, so it's exactly
3348 // where a fresh stale-assumption/write conflict appears.
3349 // A Strict conflict rejects this replan attempt via the
3350 // same bounded ReplanRejected path (capped by
3351 // max_replans), never an infinite loop.
3352 let replan_blocked = self.transaction_precheck(&new_proposal).await;
3353 if !replan_blocked.is_empty() {
3354 let reason = "transactional conflict with current state".to_string();
3355 lineage.push(ProposalLineageEntry {
3356 generation,
3357 proposal_id: new_proposal.id.clone(),
3358 proposal_digest: Some(new_proposal_digest),
3359 status: ProposalLineageStatus::Rejected,
3360 rejection_reason: Some(reason.clone()),
3361 });
3362 let mut log = self.log.lock().await;
3363 log.append(
3364 EventKind::ReplanRejected,
3365 None,
3366 Some(&proposal.id),
3367 [
3368 ("errors".to_string(), Value::from(reason)),
3369 ("attempt".to_string(), Value::from(generation)),
3370 ]
3371 .into(),
3372 );
3373 drop(log);
3374 continue 'replan_candidates;
3375 }
3376
3377 // Admission gates re-run on EVERY replanned proposal
3378 // (linus review C-2): the replan callback is exactly
3379 // where injected/adversarial content reshapes a plan,
3380 // so a proposal that was clean at first admission must
3381 // not smuggle a hazardous replan past the gates. Same
3382 // fail-closed contract as first admission — a hard
3383 // reject or any unapproved escalation rejects this
3384 // replan attempt (bounded by max_replans).
3385 let replan_admission = self
3386 .run_admission_gates(&new_proposal, session_id, scope)
3387 .await;
3388 if !replan_admission.admitted {
3389 let mut cleared = !replan_admission.hard_rejected
3390 && replan_admission.needs_approval();
3391 if cleared {
3392 for esc in &replan_admission.escalations {
3393 if self.admission_decision(&esc.fingerprint).await
3394 != Some(ApprovalDecision::Approved)
3395 {
3396 cleared = false;
3397 break;
3398 }
3399 }
3400 }
3401 if !cleared {
3402 let gate = replan_admission
3403 .deciding_gate
3404 .as_deref()
3405 .unwrap_or("admission");
3406 let why = replan_admission
3407 .reason
3408 .clone()
3409 .unwrap_or_else(|| "blocked by admission gate".to_string());
3410 let reason =
3411 format!("admission gate blocked replan (gate: {gate}; {why})");
3412 lineage.push(ProposalLineageEntry {
3413 generation,
3414 proposal_id: new_proposal.id.clone(),
3415 proposal_digest: Some(new_proposal_digest),
3416 status: ProposalLineageStatus::Rejected,
3417 rejection_reason: Some(reason.clone()),
3418 });
3419 let mut log = self.log.lock().await;
3420 log.append(
3421 EventKind::ReplanRejected,
3422 None,
3423 Some(&proposal.id),
3424 [
3425 ("errors".to_string(), Value::from(reason)),
3426 ("attempt".to_string(), Value::from(generation)),
3427 ]
3428 .into(),
3429 );
3430 drop(log);
3431 continue 'replan_candidates;
3432 }
3433 }
3434
3435 lineage.push(ProposalLineageEntry {
3436 generation,
3437 proposal_id: new_proposal.id.clone(),
3438 proposal_digest: Some(new_proposal_digest.clone()),
3439 status: ProposalLineageStatus::Accepted,
3440 rejection_reason: None,
3441 });
3442 accepted_proposal_preimages.push(AcceptedProposalPreimage {
3443 generation,
3444 proposal_digest: new_proposal_digest.clone(),
3445 proposal: new_proposal.clone(),
3446 });
3447
3448 // Log accepted proposal
3449 {
3450 let mut log = self.log.lock().await;
3451 log.append(
3452 EventKind::ReplanProposalReceived,
3453 None,
3454 Some(&proposal.id),
3455 [
3456 ("attempt".to_string(), Value::from(generation)),
3457 (
3458 "proposal_id".to_string(),
3459 Value::from(new_proposal.id.as_str()),
3460 ),
3461 (
3462 "proposal_digest".to_string(),
3463 Value::from(new_proposal_digest),
3464 ),
3465 ("proposal".to_string(), new_proposal_preimage),
3466 (
3467 "new_action_count".to_string(),
3468 Value::from(new_proposal.actions.len()),
3469 ),
3470 ]
3471 .into(),
3472 );
3473 }
3474 current_proposal = new_proposal;
3475 break 'replan_candidates;
3476 }
3477 Err(e) => {
3478 // Replan callback itself failed — log and return original failure
3479 let mut log = self.log.lock().await;
3480 log.append(
3481 EventKind::ReplanExhausted,
3482 None,
3483 Some(&proposal.id),
3484 [
3485 ("reason".to_string(), Value::from("callback_error")),
3486 ("error".to_string(), Value::from(e.as_str())),
3487 ("attempt".to_string(), Value::from(generation)),
3488 ]
3489 .into(),
3490 );
3491 if let Some(err) = self.persist_trajectory(
3492 proposal,
3493 ¤t_proposal,
3494 &result,
3495 car_memgine::TrajectoryOutcome::Failed,
3496 attempt,
3497 &state_before_map,
3498 ) {
3499 log.append(
3500 EventKind::ActionFailed,
3501 None,
3502 Some(&proposal.id),
3503 [(
3504 "trajectory_persist_error".to_string(),
3505 Value::from(err.as_str()),
3506 )]
3507 .into(),
3508 );
3509 }
3510 return finalize_proposal_result(
3511 result,
3512 &proposal.id,
3513 &lineage,
3514 &accepted_proposal_preimages,
3515 );
3516 }
3517 }
3518 }
3519 }
3520 }
3521
3522 /// Persist a trajectory to the store if configured.
3523 fn persist_trajectory(
3524 &self,
3525 proposal: &ActionProposal,
3526 current_proposal: &ActionProposal,
3527 result: &ProposalResult,
3528 outcome: car_memgine::TrajectoryOutcome,
3529 attempt: u32,
3530 state_before_map: &HashMap<String, HashMap<String, Value>>,
3531 ) -> Option<String> {
3532 let store = self.trajectory_store.as_ref()?;
3533
3534 let trace_events: Vec<car_memgine::TraceEvent> = result
3535 .results
3536 .iter()
3537 .map(|r| {
3538 let kind = match r.status {
3539 ActionStatus::Succeeded => "action_succeeded",
3540 ActionStatus::Failed => "action_failed",
3541 ActionStatus::Rejected => "action_rejected",
3542 ActionStatus::Skipped => "action_skipped",
3543 _ => "unknown",
3544 };
3545 let tool = current_proposal
3546 .actions
3547 .iter()
3548 .find(|a| a.id == r.action_id)
3549 .and_then(|a| a.tool.clone());
3550 let reward = match r.status {
3551 ActionStatus::Succeeded => Some(1.0),
3552 ActionStatus::Failed => Some(0.0),
3553 ActionStatus::Rejected => Some(0.0),
3554 ActionStatus::Skipped => None,
3555 _ => None,
3556 };
3557 car_memgine::TraceEvent {
3558 kind: kind.to_string(),
3559 action_id: Some(r.action_id.clone()),
3560 tool,
3561 data: r
3562 .error
3563 .as_ref()
3564 .map(|e| serde_json::json!({"error": e}))
3565 .unwrap_or(serde_json::json!({})),
3566 duration_ms: r.duration_ms,
3567 state_before: state_before_map.get(&r.action_id).cloned(),
3568 state_after: if !r.state_changes.is_empty() {
3569 Some(r.state_changes.clone())
3570 } else {
3571 None
3572 },
3573 reward,
3574 }
3575 })
3576 .collect();
3577
3578 let trajectory = car_memgine::Trajectory {
3579 proposal_id: proposal.id.clone(),
3580 source: proposal.source.clone(),
3581 action_count: current_proposal.actions.len(),
3582 events: trace_events,
3583 outcome,
3584 timestamp: chrono::Utc::now(),
3585 duration_ms: result.cost.total_duration_ms,
3586 replan_attempts: attempt,
3587 };
3588
3589 match store.append(&trajectory) {
3590 Ok(()) => None,
3591 Err(e) => Some(e.to_string()),
3592 }
3593 }
3594
3595 async fn rollback_failed_plan_candidate(
3596 &self,
3597 proposal: &ActionProposal,
3598 result: &mut ProposalResult,
3599 pre_plan_snapshot: &HashMap<String, Value>,
3600 pre_plan_transitions: usize,
3601 pre_plan_idempotency_keys: &std::collections::HashSet<String>,
3602 ) -> Result<(), String> {
3603 let executed_proposal = result.final_proposal.as_ref().unwrap_or(proposal);
3604 let executed_proposal_digest = result
3605 .accepted_proposal_preimages
3606 .last()
3607 .map(|accepted| accepted.proposal_digest.clone())
3608 .unwrap_or_else(|| {
3609 proposal_digest(executed_proposal)
3610 .expect("an executed planning candidate already passed proposal identity")
3611 });
3612 let observed_changes_by_action: std::collections::BTreeMap<String, HashMap<String, Value>> =
3613 result
3614 .results
3615 .iter()
3616 .filter(|action_result| {
3617 action_result.status == ActionStatus::Succeeded
3618 && !action_result.state_changes.is_empty()
3619 })
3620 .map(|action_result| {
3621 (
3622 action_result.action_id.clone(),
3623 action_result.state_changes.clone(),
3624 )
3625 })
3626 .collect();
3627 let affected_actions: Vec<String> = observed_changes_by_action.keys().cloned().collect();
3628
3629 // The planning transaction owns the StateStore proposal guard, so this
3630 // exact snapshot/count restore cannot erase another proposal's commit.
3631 let durability = match self
3632 .state
3633 .restore(pre_plan_snapshot.clone(), pre_plan_transitions)
3634 {
3635 Ok(durability) => durability,
3636 Err(error) => {
3637 let detail = record_rollback_durability_error(
3638 &mut result.results,
3639 ROLLBACK_DURABILITY_ERROR,
3640 &error.to_string(),
3641 );
3642 self.log.lock().await.append(
3643 EventKind::ActionFailed,
3644 None,
3645 Some(&proposal.id),
3646 [
3647 ("error".to_string(), Value::from(detail.clone())),
3648 ("attempted".to_string(), Value::from(true)),
3649 ("publication_succeeded".to_string(), Value::from(false)),
3650 ("rollback_succeeded".to_string(), Value::from(false)),
3651 ("stage".to_string(), Value::from("plan_fallback_rollback")),
3652 ]
3653 .into(),
3654 );
3655 return Err(detail);
3656 }
3657 };
3658 let durability_error = match &durability {
3659 RestoreDurability::Durable => None,
3660 RestoreDurability::DurabilityUnknown { error } => Some(error.clone()),
3661 };
3662
3663 {
3664 let mut log = self.log.lock().await;
3665 log.append(
3666 EventKind::StateSnapshot,
3667 None,
3668 Some(&proposal.id),
3669 [(
3670 "state".to_string(),
3671 serde_json::to_value(pre_plan_snapshot).unwrap_or_default(),
3672 )]
3673 .into(),
3674 );
3675 log.append(
3676 EventKind::StateRollback,
3677 None,
3678 Some(&proposal.id),
3679 [
3680 (
3681 "rolled_back_to".to_string(),
3682 Value::from("pre-plan snapshot"),
3683 ),
3684 (
3685 "affected_actions".to_string(),
3686 serde_json::to_value(&affected_actions).unwrap_or_default(),
3687 ),
3688 (
3689 "rolled_back_changes".to_string(),
3690 serde_json::to_value(&observed_changes_by_action).unwrap_or_default(),
3691 ),
3692 (
3693 "changes_semantics".to_string(),
3694 Value::from("rolled_back_failed_planning_candidate"),
3695 ),
3696 (
3697 "proposal_digest".to_string(),
3698 Value::from(executed_proposal_digest),
3699 ),
3700 ("attempted".to_string(), Value::from(true)),
3701 (
3702 "durability_unknown".to_string(),
3703 Value::from(durability_error.is_some()),
3704 ),
3705 ("publication_succeeded".to_string(), Value::from(true)),
3706 ("rollback_succeeded".to_string(), Value::from(true)),
3707 ("stage".to_string(), Value::from("plan_fallback_rollback")),
3708 ]
3709 .into(),
3710 );
3711 }
3712
3713 // Only invalidate entries introduced by this planning transaction.
3714 // Pre-existing keys may have been read by a candidate, but their
3715 // committed result predates the pre-plan snapshot and must survive.
3716 let mut invalidated = std::collections::BTreeSet::new();
3717 {
3718 let mut cache = self.idempotency_cache.lock().await;
3719 for action_result in &result.results {
3720 if action_result.status != ActionStatus::Succeeded {
3721 continue;
3722 }
3723 let Some(action) = executed_proposal
3724 .actions
3725 .iter()
3726 .find(|action| action.id == action_result.action_id)
3727 else {
3728 continue;
3729 };
3730 if !action.idempotent || action.invocation_mode.is_detached() {
3731 continue;
3732 }
3733 let key = idempotency_key(action, None);
3734 if !pre_plan_idempotency_keys.contains(&key) && cache.remove(&key).is_some() {
3735 invalidated.insert(key);
3736 }
3737 }
3738 }
3739 for key in invalidated {
3740 self.journal_idempotency(&key, None).await;
3741 }
3742
3743 // Same reasoning as the abort path above: a rejected planning
3744 // candidate loses its state effects, not the record of what ran. This
3745 // result is observable — `plan_and_execute` returns it directly on a
3746 // rollback durability error, and otherwise keeps it as
3747 // `first_failure` (car#1157). `rolled_back` is the machine-readable
3748 // fact; the warning remains descriptive and must never drive logic.
3749 for action_result in &mut result.results {
3750 if action_result.status == ActionStatus::Succeeded {
3751 action_result.rolled_back = true;
3752 action_result.error = Some(PLAN_FALLBACK_ROLLBACK_WARNING.to_string());
3753 action_result.state_changes.clear();
3754 }
3755 }
3756 if let Some(error) = durability_error {
3757 return Err(record_rollback_durability_error(
3758 &mut result.results,
3759 ROLLBACK_DURABILITY_UNKNOWN,
3760 &error,
3761 ));
3762 }
3763 Ok(())
3764 }
3765
3766 /// Score N candidate proposals, execute the best valid one, fall back to
3767 /// next-best on failure. Combines car-planner scoring with engine execution.
3768 ///
3769 /// Returns the result from whichever proposal was executed (best or fallback).
3770 /// If all candidates fail verification, returns an error result for the first.
3771 pub async fn plan_and_execute(
3772 &self,
3773 candidates: &[ActionProposal],
3774 planner_config: Option<car_planner::PlannerConfig>,
3775 feedback: Option<&car_planner::ToolFeedback>,
3776 ) -> ProposalResult {
3777 if candidates.is_empty() {
3778 return ProposalResult::new("empty", vec![], car_ir::CostSummary::default());
3779 }
3780
3781 // One transaction guard covers ranking's state snapshot, every
3782 // candidate attempt and rollback, and the accepted candidate's commit.
3783 // Candidate execution calls the already-guarded path to avoid recursive
3784 // acquisition of the non-reentrant mutex.
3785 let _proposal_execution = self.state.lock_proposal_execution().await;
3786
3787 // Score all candidates
3788 let planner = car_planner::Planner::new(planner_config.unwrap_or_default());
3789 let tools_guard = self.tools.read().await;
3790 let tool_names: std::collections::HashSet<String> = tools_guard.keys().cloned().collect();
3791 drop(tools_guard);
3792
3793 let pre_plan_snapshot = self.state.snapshot();
3794 let pre_plan_transitions = self.state.transition_count();
3795 let pre_plan_idempotency_keys: std::collections::HashSet<String> = self
3796 .idempotency_cache
3797 .lock()
3798 .await
3799 .keys()
3800 .cloned()
3801 .collect();
3802 let ranked = planner.rank_with_feedback(
3803 candidates,
3804 Some(&pre_plan_snapshot),
3805 Some(&tool_names),
3806 feedback,
3807 );
3808
3809 // Try each valid candidate in score order
3810 let mut first_failure: Option<ProposalResult> = None;
3811 for scored in &ranked {
3812 if !scored.valid {
3813 continue;
3814 }
3815
3816 let proposal = &candidates[scored.index];
3817 let cancel = tokio_util::sync::CancellationToken::new();
3818 let mut result = self
3819 .execute_with_optional_session_already_guarded(proposal, None, None, None, &cancel)
3820 .await;
3821
3822 if result.all_succeeded() {
3823 return result;
3824 }
3825
3826 if self
3827 .rollback_failed_plan_candidate(
3828 proposal,
3829 &mut result,
3830 &pre_plan_snapshot,
3831 pre_plan_transitions,
3832 &pre_plan_idempotency_keys,
3833 )
3834 .await
3835 .is_err()
3836 {
3837 return result;
3838 }
3839
3840 tracing::info!(
3841 proposal_id = %proposal.id,
3842 score = scored.score,
3843 "plan_and_execute: proposal failed, trying next candidate"
3844 );
3845
3846 if first_failure.is_none() {
3847 first_failure = Some(result);
3848 }
3849 }
3850
3851 // Return the first failure result (don't re-execute — avoids duplicate side effects)
3852 first_failure.unwrap_or_else(|| {
3853 ProposalResult::for_proposal(&candidates[0], vec![], car_ir::CostSummary::default())
3854 })
3855 }
3856
3857 /// Execute a single proposal through the runtime loop (no replanning).
3858 /// Returns (result, state_before_map) where state_before_map has per-action snapshots.
3859 ///
3860 /// `session_id`, when `Some`, scopes per-action policy validation
3861 /// to the named session in addition to global policies. Both
3862 /// layers must pass for an action to run; the session layer cannot
3863 /// loosen what global denies.
3864 async fn execute_inner_with_cancel(
3865 &self,
3866 proposal: &ActionProposal,
3867 cancel: Option<&tokio_util::sync::CancellationToken>,
3868 session_id: Option<&str>,
3869 scope: Option<&crate::scope::RuntimeScope>,
3870 ) -> (ProposalResult, HashMap<String, HashMap<String, Value>>) {
3871 // Identity is established before CAR admits the proposal or emits any
3872 // action transition. In particular, an I-JSON-incompatible number may
3873 // not produce a journal row whose mandatory JCS digest is missing.
3874 let (proposal_preimage, proposal_digest) = match proposal_journal_identity(proposal) {
3875 Ok(identity) => identity,
3876 Err(error) => {
3877 let results = proposal
3878 .actions
3879 .iter()
3880 .map(|action| rejected_result(&action.id, error.clone()))
3881 .collect();
3882 self.log.lock().await.append(
3883 EventKind::StateRollback,
3884 None,
3885 Some(&proposal.id),
3886 proposal_rejection_boundary_data(proposal, None, &error),
3887 );
3888 return (
3889 ProposalResult::for_proposal(proposal, results, CostSummary::default()),
3890 HashMap::new(),
3891 );
3892 }
3893 };
3894
3895 // Whether validator/policy/capability rejections should count toward
3896 // the abort-and-replan path (default false). Read once up front so the
3897 // per-action loop doesn't take the replan_config lock repeatedly (and
3898 // never inside join_all). Independent lock from log/tools/policies/
3899 // capabilities, so no lock-ordering conflict.
3900 let replan_on_rejected = self.replan_config.read().await.replan_on_rejected;
3901
3902 // Generate trace_id for this proposal execution
3903 let trace_id = Uuid::new_v4().to_string();
3904
3905 // Begin root span for proposal execution
3906 let root_span_id = {
3907 let mut log = self.log.lock().await;
3908 log.begin_span(
3909 "proposal.execute",
3910 &trace_id,
3911 None,
3912 [("proposal_id".to_string(), Value::from(proposal.id.as_str()))].into(),
3913 )
3914 };
3915
3916 // Log proposal received
3917 {
3918 let mut log = self.log.lock().await;
3919 log.append(
3920 EventKind::ProposalReceived,
3921 None,
3922 Some(&proposal.id),
3923 [
3924 ("source".to_string(), Value::from(proposal.source.as_str())),
3925 (
3926 "action_count".to_string(),
3927 Value::from(proposal.actions.len()),
3928 ),
3929 ("proposal".to_string(), proposal_preimage.clone()),
3930 (
3931 "proposal_digest".to_string(),
3932 Value::from(proposal_digest.clone()),
3933 ),
3934 ]
3935 .into(),
3936 );
3937 }
3938
3939 // Capability check: max_actions budget for entire proposal
3940 {
3941 let caps = self.capabilities.read().await;
3942 if let Some(ref cap) = *caps {
3943 if !cap.actions_within_budget(proposal.actions.len() as u32) {
3944 let reason = format!(
3945 "capability denied: proposal has {} actions, max allowed is {:?}",
3946 proposal.actions.len(),
3947 cap.max_actions
3948 );
3949 let mut action_results = Vec::new();
3950 let mut log = self.log.lock().await;
3951 for action in &proposal.actions {
3952 log.append(
3953 EventKind::ActionRejected,
3954 Some(&action.id),
3955 Some(&proposal.id),
3956 action_outcome_data("proposal_capability", &reason, false, None),
3957 );
3958 action_results.push(rejected_result(&action.id, reason.clone()));
3959 }
3960 log.append(
3961 EventKind::StateRollback,
3962 None,
3963 Some(&proposal.id),
3964 proposal_rejection_boundary_data(proposal, Some(&proposal_digest), &reason),
3965 );
3966 drop(log);
3967 return (
3968 ProposalResult::for_proposal(
3969 proposal,
3970 action_results,
3971 CostSummary::default(),
3972 ),
3973 HashMap::new(),
3974 );
3975 }
3976 }
3977 }
3978
3979 // Snapshot for rollback. When the proposal carries a tenant scope,
3980 // snapshot only that tenant's namespace so a rollback can't clobber
3981 // concurrent tenants' state (EPIC E / E2). No tenant → unchanged
3982 // full snapshot.
3983 let rollback_tenant: Option<&str> = scope.and_then(|s| s.tenant_id.as_deref());
3984 let snapshot = match rollback_tenant {
3985 Some(_) => self.state.snapshot_scoped(rollback_tenant),
3986 None => self.state.snapshot(),
3987 };
3988 let transition_count = self.state.transition_count();
3989
3990 let mut results: Vec<ActionResult> = Vec::new();
3991 // Per-action state snapshots captured before execution (for TraceEvent.state_before).
3992 let mut state_before_map: HashMap<String, HashMap<String, Value>> = HashMap::new();
3993 let mut aborted = false;
3994 let mut budget_exceeded = false;
3995 let mut total_retries: u32 = 0;
3996
3997 // Running cost counters for budget enforcement
3998 let mut running_tool_calls: u32 = 0;
3999 let mut running_actions: u32 = 0;
4000 let mut running_duration_ms: f64 = 0.0;
4001
4002 // Snapshot the budget once
4003 let budget = self.cost_budget.read().await.clone();
4004
4005 // Build DAG
4006 let levels = build_dag(&proposal.actions);
4007
4008 let mut canceled = false;
4009 for level in &levels {
4010 // Cooperative cancellation check at the level boundary.
4011 // Actions already in flight aren't interrupted (we can't
4012 // safely cancel a tool call dispatched to a user-provided
4013 // executor), but every action that hadn't started runs
4014 // is recorded as canceled with a clear reason.
4015 if !canceled {
4016 if let Some(token) = cancel {
4017 if token.is_cancelled() {
4018 canceled = true;
4019 }
4020 }
4021 }
4022 if canceled {
4023 for &idx in level {
4024 let action = &proposal.actions[idx];
4025 let reason = format!("{CANCELED_PREFIX}cancellation requested by caller");
4026 self.log.lock().await.append(
4027 EventKind::ActionSkipped,
4028 Some(&action.id),
4029 Some(&proposal.id),
4030 action_outcome_data("cancellation", &reason, false, None),
4031 );
4032 results.push(canceled_result(
4033 &action.id,
4034 "cancellation requested by caller",
4035 ));
4036 }
4037 continue;
4038 }
4039 if aborted || budget_exceeded {
4040 let skip_reason = if budget_exceeded {
4041 "cost budget exceeded"
4042 } else {
4043 "skipped due to earlier abort"
4044 };
4045 for &idx in level {
4046 let action = &proposal.actions[idx];
4047 let stage = if budget_exceeded {
4048 "cost_budget"
4049 } else {
4050 "dependency_abort"
4051 };
4052 self.log.lock().await.append(
4053 EventKind::ActionSkipped,
4054 Some(&action.id),
4055 Some(&proposal.id),
4056 action_outcome_data(stage, skip_reason, false, None),
4057 );
4058 results.push(skipped_result(&action.id, skip_reason));
4059 }
4060 continue;
4061 }
4062
4063 // Check if any action in this level has ABORT behavior
4064 let has_abort = level
4065 .iter()
4066 .any(|&i| proposal.actions[i].failure_behavior == FailureBehavior::Abort);
4067
4068 if level.len() == 1 || has_abort {
4069 // Sequential execution
4070 for &idx in level {
4071 if aborted || budget_exceeded {
4072 let skip_reason = if budget_exceeded {
4073 "cost budget exceeded"
4074 } else {
4075 "skipped due to abort"
4076 };
4077 let action = &proposal.actions[idx];
4078 let stage = if budget_exceeded {
4079 "cost_budget"
4080 } else {
4081 "dependency_abort"
4082 };
4083 self.log.lock().await.append(
4084 EventKind::ActionSkipped,
4085 Some(&action.id),
4086 Some(&proposal.id),
4087 action_outcome_data(stage, skip_reason, false, None),
4088 );
4089 results.push(skipped_result(&action.id, skip_reason));
4090 continue;
4091 }
4092
4093 // Budget check before execution
4094 if let Some(ref b) = budget {
4095 if let Some(max) = b.max_actions {
4096 if running_actions >= max {
4097 budget_exceeded = true;
4098 self.log.lock().await.append(
4099 EventKind::ActionSkipped,
4100 Some(&proposal.actions[idx].id),
4101 Some(&proposal.id),
4102 action_outcome_data(
4103 "cost_budget",
4104 "cost budget exceeded",
4105 false,
4106 None,
4107 ),
4108 );
4109 results.push(skipped_result(
4110 &proposal.actions[idx].id,
4111 "cost budget exceeded",
4112 ));
4113 continue;
4114 }
4115 }
4116 if let Some(max) = b.max_tool_calls {
4117 if proposal.actions[idx].action_type == ActionType::ToolCall
4118 && running_tool_calls >= max
4119 {
4120 budget_exceeded = true;
4121 self.log.lock().await.append(
4122 EventKind::ActionSkipped,
4123 Some(&proposal.actions[idx].id),
4124 Some(&proposal.id),
4125 action_outcome_data(
4126 "cost_budget",
4127 "cost budget exceeded",
4128 false,
4129 None,
4130 ),
4131 );
4132 results.push(skipped_result(
4133 &proposal.actions[idx].id,
4134 "cost budget exceeded",
4135 ));
4136 continue;
4137 }
4138 }
4139 if let Some(max) = b.max_duration_ms {
4140 if running_duration_ms >= max {
4141 budget_exceeded = true;
4142 self.log.lock().await.append(
4143 EventKind::ActionSkipped,
4144 Some(&proposal.actions[idx].id),
4145 Some(&proposal.id),
4146 action_outcome_data(
4147 "cost_budget",
4148 "cost budget exceeded",
4149 false,
4150 None,
4151 ),
4152 );
4153 results.push(skipped_result(
4154 &proposal.actions[idx].id,
4155 "cost budget exceeded",
4156 ));
4157 continue;
4158 }
4159 }
4160 }
4161
4162 state_before_map.insert(
4163 proposal.actions[idx].id.clone(),
4164 snapshot_relevant_keys(&self.state, &proposal.actions[idx]),
4165 );
4166 let (ar, action_retries) = self
4167 .process_action(
4168 &proposal.actions[idx],
4169 &proposal.id,
4170 &trace_id,
4171 &root_span_id,
4172 session_id,
4173 scope,
4174 )
4175 .await;
4176 total_retries += action_retries;
4177
4178 // Update running counters
4179 if ar.status == ActionStatus::Succeeded
4180 && proposal.actions[idx].action_type == ActionType::ToolCall
4181 {
4182 running_tool_calls += 1;
4183 }
4184 if ar.status != ActionStatus::Skipped {
4185 running_actions += 1;
4186 }
4187 if let Some(d) = ar.duration_ms {
4188 running_duration_ms += d;
4189 }
4190
4191 if ar.terminal
4192 || requires_integrity_rollback(&ar)
4193 || ((ar.status == ActionStatus::Failed
4194 || (replan_on_rejected && ar.status == ActionStatus::Rejected))
4195 && proposal.actions[idx].failure_behavior == FailureBehavior::Abort)
4196 {
4197 aborted = true;
4198 }
4199 results.push(ar);
4200 }
4201 } else {
4202 // Concurrent execution via futures::join_all
4203 // Snapshot only relevant keys per action (all see same pre-level state)
4204 for &idx in level {
4205 state_before_map.insert(
4206 proposal.actions[idx].id.clone(),
4207 snapshot_relevant_keys(&self.state, &proposal.actions[idx]),
4208 );
4209 }
4210 let futs: Vec<_> = level
4211 .iter()
4212 .map(|&idx| {
4213 self.process_action(
4214 &proposal.actions[idx],
4215 &proposal.id,
4216 &trace_id,
4217 &root_span_id,
4218 session_id,
4219 scope,
4220 )
4221 })
4222 .collect();
4223 let level_results = futures::future::join_all(futs).await;
4224
4225 for (i, (ar, action_retries)) in level_results.into_iter().enumerate() {
4226 let idx = level[i];
4227 total_retries += action_retries;
4228 if ar.status == ActionStatus::Succeeded
4229 && proposal.actions[idx].action_type == ActionType::ToolCall
4230 {
4231 running_tool_calls += 1;
4232 }
4233 if ar.status != ActionStatus::Skipped {
4234 running_actions += 1;
4235 }
4236 if let Some(d) = ar.duration_ms {
4237 running_duration_ms += d;
4238 }
4239 if ar.terminal || requires_integrity_rollback(&ar) {
4240 aborted = true;
4241 }
4242 results.push(ar);
4243 }
4244 }
4245 }
4246
4247 let observed_changes_by_action: std::collections::BTreeMap<String, HashMap<String, Value>> =
4248 proposal
4249 .actions
4250 .iter()
4251 .filter_map(|action| {
4252 results
4253 .iter()
4254 .find(|result| {
4255 result.action_id == action.id
4256 && result.status == ActionStatus::Succeeded
4257 && !result.state_changes.is_empty()
4258 })
4259 .map(|result| (action.id.clone(), result.state_changes.clone()))
4260 })
4261 .collect();
4262 let affected_actions: Vec<String> = observed_changes_by_action.keys().cloned().collect();
4263
4264 // Handle rollback. Per-action StateChanged rows are provisional until
4265 // this proposal-level transaction boundary is durably recorded.
4266 if aborted {
4267 // Tenant-scoped rollback when the proposal is tenant-scoped
4268 // (EPIC E / E2), else the full restore.
4269 let rollback = match rollback_tenant {
4270 Some(_) => {
4271 self.state
4272 .restore_scoped(rollback_tenant, snapshot.clone(), transition_count)
4273 }
4274 None => self.state.restore(snapshot.clone(), transition_count),
4275 };
4276
4277 let rollback_durability = match rollback {
4278 Ok(durability) => Some(durability),
4279 Err(error) => {
4280 let detail = record_rollback_durability_error(
4281 &mut results,
4282 ROLLBACK_DURABILITY_ERROR,
4283 &error.to_string(),
4284 );
4285 self.log.lock().await.append(
4286 EventKind::ActionFailed,
4287 None,
4288 Some(&proposal.id),
4289 [
4290 ("error".to_string(), Value::from(detail)),
4291 ("attempted".to_string(), Value::from(true)),
4292 ("publication_succeeded".to_string(), Value::from(false)),
4293 ("in_memory_state_preserved".to_string(), Value::from(true)),
4294 (
4295 "idempotency_entries_preserved".to_string(),
4296 Value::from(true),
4297 ),
4298 ("rollback_succeeded".to_string(), Value::from(false)),
4299 ("stage".to_string(), Value::from("proposal_rollback")),
4300 ]
4301 .into(),
4302 );
4303 None
4304 }
4305 };
4306
4307 if let Some(durability) = rollback_durability {
4308 let durability_error = match &durability {
4309 RestoreDurability::Durable => None,
4310 RestoreDurability::DurabilityUnknown { error } => Some(error.clone()),
4311 };
4312 if let Some(error) = &durability_error {
4313 record_rollback_durability_error(
4314 &mut results,
4315 ROLLBACK_DURABILITY_UNKNOWN,
4316 error,
4317 );
4318 }
4319 let mut log = self.log.lock().await;
4320 log.append(
4321 EventKind::StateSnapshot,
4322 None,
4323 Some(&proposal.id),
4324 [(
4325 "state".to_string(),
4326 serde_json::to_value(&snapshot).unwrap_or_default(),
4327 )]
4328 .into(),
4329 );
4330 log.append(
4331 EventKind::StateRollback,
4332 None,
4333 Some(&proposal.id),
4334 [
4335 (
4336 "rolled_back_to".to_string(),
4337 Value::from("pre-proposal snapshot"),
4338 ),
4339 (
4340 "affected_actions".to_string(),
4341 serde_json::to_value(&affected_actions).unwrap_or_default(),
4342 ),
4343 (
4344 "rolled_back_changes".to_string(),
4345 serde_json::to_value(&observed_changes_by_action).unwrap_or_default(),
4346 ),
4347 (
4348 "changes_semantics".to_string(),
4349 Value::from("rolled_back_provisional_state_mutations"),
4350 ),
4351 (
4352 "proposal_digest".to_string(),
4353 Value::from(proposal_digest.clone()),
4354 ),
4355 ("attempted".to_string(), Value::from(true)),
4356 (
4357 "durability_unknown".to_string(),
4358 Value::from(durability_error.is_some()),
4359 ),
4360 ("publication_succeeded".to_string(), Value::from(true)),
4361 ("rollback_succeeded".to_string(), Value::from(true)),
4362 ("stage".to_string(), Value::from("proposal_rollback")),
4363 ]
4364 .into(),
4365 );
4366 drop(log);
4367
4368 // Clear idempotency cache for rolled-back actions only after
4369 // the state journal has durably accepted the rollback.
4370 let mut invalidated: Vec<String> = Vec::new();
4371 {
4372 let mut cache = self.idempotency_cache.lock().await;
4373 for r in &results {
4374 if r.status == ActionStatus::Succeeded {
4375 for action in &proposal.actions {
4376 if action.id == r.action_id && action.idempotent {
4377 let key = idempotency_key(action, scope);
4378 cache.remove(&key);
4379 invalidated.push(key);
4380 }
4381 }
4382 }
4383 }
4384 }
4385 // Record tombstones so the durable journal (C3) doesn't
4386 // resurrect a rolled-back result on reload.
4387 for key in &invalidated {
4388 self.journal_idempotency(key, None).await;
4389 }
4390 // The proposal did not commit, so no action's state changes
4391 // survive. Each action's STATUS still records whether that
4392 // action executed; `rolled_back` independently records whether
4393 // that successful execution preceded this rollback. Folding
4394 // the second fact into ActionStatus erased which action
4395 // aborted the proposal and marked `failed` precisely the
4396 // actions whose external effects DID happen and may remain.
4397 // The warning is human-readable context, never control flow
4398 // (car#1157).
4399 for result in &mut results {
4400 if result.status == ActionStatus::Succeeded {
4401 result.rolled_back = true;
4402 result.error = Some(ROLLBACK_WARNING.to_string());
4403 result.state_changes.clear();
4404 }
4405 }
4406 }
4407 } else {
4408 self.log.lock().await.append(
4409 EventKind::StateCommitted,
4410 None,
4411 Some(&proposal.id),
4412 [
4413 (
4414 "affected_actions".to_string(),
4415 serde_json::to_value(&affected_actions).unwrap_or_default(),
4416 ),
4417 (
4418 "committed_changes".to_string(),
4419 serde_json::to_value(&observed_changes_by_action).unwrap_or_default(),
4420 ),
4421 (
4422 "changes_semantics".to_string(),
4423 Value::from("committed_provisional_state_mutations"),
4424 ),
4425 ("proposal_digest".to_string(), Value::from(proposal_digest)),
4426 ("attempted".to_string(), Value::from(true)),
4427 ("stage".to_string(), Value::from("proposal_commit")),
4428 ]
4429 .into(),
4430 );
4431 }
4432
4433 // Sort results to match original action order
4434 let action_order: HashMap<String, usize> = proposal
4435 .actions
4436 .iter()
4437 .enumerate()
4438 .map(|(i, a)| (a.id.clone(), i))
4439 .collect();
4440 results.sort_by_key(|r| {
4441 action_order
4442 .get(&r.action_id)
4443 .copied()
4444 .unwrap_or(usize::MAX)
4445 });
4446
4447 // Compute cost summary from results
4448 let mut cost = CostSummary::default();
4449 for r in &results {
4450 let action = action_order
4451 .get(&r.action_id)
4452 .and_then(|&i| proposal.actions.get(i));
4453 match r.status {
4454 ActionStatus::Succeeded => {
4455 cost.actions_executed += 1;
4456 if let Some(a) = action {
4457 if a.action_type == ActionType::ToolCall {
4458 cost.tool_calls += 1;
4459 }
4460 }
4461 }
4462 // A failed action ran — the tool was invoked and errored — so it
4463 // belongs in `actions_executed`. A rejected one was blocked by
4464 // the validator or a policy and never started; counting it as
4465 // executed reported work that never happened (car#624).
4466 ActionStatus::Failed => {
4467 cost.actions_executed += 1;
4468 }
4469 ActionStatus::Rejected => {
4470 cost.actions_rejected += 1;
4471 }
4472 ActionStatus::Skipped => {
4473 cost.actions_skipped += 1;
4474 }
4475 _ => {}
4476 }
4477 if let Some(d) = r.duration_ms {
4478 cost.total_duration_ms += d;
4479 }
4480 }
4481
4482 // Set retries from inline counter
4483 cost.retries = total_retries;
4484
4485 // End root span — Ok if no abort, Error if aborted
4486 {
4487 let span_status = if aborted {
4488 SpanStatus::Error
4489 } else {
4490 SpanStatus::Ok
4491 };
4492 let mut log = self.log.lock().await;
4493 log.end_span(&root_span_id, span_status);
4494 }
4495
4496 let proposal_result = ProposalResult::for_proposal(proposal, results, cost);
4497
4498 // Post-execution: auto-distill skills from this execution trace
4499 if self.auto_distill {
4500 if let Some(ref memgine) = self.memgine {
4501 // Convert results to TraceEvents for distillation
4502 let trace_events: Vec<car_memgine::TraceEvent> = proposal_result
4503 .results
4504 .iter()
4505 .map(|r| {
4506 let kind = match r.status {
4507 ActionStatus::Succeeded => "action_succeeded",
4508 ActionStatus::Failed => "action_failed",
4509 ActionStatus::Rejected => "action_rejected",
4510 ActionStatus::Skipped => "action_skipped",
4511 _ => "unknown",
4512 };
4513 // Find the matching action to get the tool name
4514 let tool = proposal
4515 .actions
4516 .iter()
4517 .find(|a| a.id == r.action_id)
4518 .and_then(|a| a.tool.clone());
4519 let mut data = serde_json::Map::new();
4520 if let Some(ref e) = r.error {
4521 data.insert("error".into(), Value::from(e.as_str()));
4522 }
4523 if let Some(ref o) = r.output {
4524 data.insert("output".into(), o.clone());
4525 }
4526 car_memgine::TraceEvent {
4527 kind: kind.to_string(),
4528 action_id: Some(r.action_id.clone()),
4529 tool,
4530 data: Value::Object(data),
4531 duration_ms: r.duration_ms,
4532 reward: match r.status {
4533 ActionStatus::Succeeded => Some(1.0),
4534 ActionStatus::Failed | ActionStatus::Rejected => Some(0.0),
4535 _ => None,
4536 },
4537 ..Default::default()
4538 }
4539 })
4540 .collect();
4541
4542 let mut engine = memgine.lock().await;
4543 let skills = engine.distill_skills(&trace_events).await;
4544 if !skills.is_empty() {
4545 let count = skills.len();
4546 // Validation-gated ingest (SkillOpt-inspired). Distilled
4547 // skills are NOT trusted straight into the active pool; each
4548 // enters as a PROVISIONAL candidate on trial and must prove
4549 // itself (or beat its incumbent) before the promotion gate in
4550 // consolidate() makes it Active. Tenant is threaded through so
4551 // candidates are stamped for the calling tenant's namespace.
4552 let tenant = scope.and_then(|s| s.tenant_id.as_deref());
4553 let provisional = engine.ingest_provisional_candidates(&skills, tenant);
4554
4555 // Log the distillation event
4556 let mut log = self.log.lock().await;
4557 log.append(
4558 EventKind::SkillDistilled,
4559 None,
4560 Some(&proposal_result.proposal_id),
4561 [
4562 ("skills_count".to_string(), Value::from(count)),
4563 ("provisional_ingested".to_string(), Value::from(provisional)),
4564 (
4565 "skill_names".to_string(),
4566 Value::from(
4567 skills.iter().map(|s| s.name.as_str()).collect::<Vec<_>>(),
4568 ),
4569 ),
4570 ]
4571 .into(),
4572 );
4573
4574 // Check if any domains need evolution
4575 let threshold = engine.evolution_threshold();
4576 let domains = engine.domains_needing_evolution(threshold);
4577 for domain in &domains {
4578 // Collect failed events for this domain
4579 let failed: Vec<car_memgine::TraceEvent> = trace_events
4580 .iter()
4581 .filter(|e| {
4582 matches!(e.kind.as_str(), "action_failed" | "action_rejected")
4583 })
4584 .cloned()
4585 .collect();
4586 if !failed.is_empty() {
4587 let evolved = engine.evolve_skills(&failed, domain).await;
4588 if !evolved.is_empty() {
4589 log.append(
4590 EventKind::EvolutionTriggered,
4591 None,
4592 Some(&proposal_result.proposal_id),
4593 [
4594 ("domain".to_string(), Value::from(domain.as_str())),
4595 ("new_skills".to_string(), Value::from(evolved.len())),
4596 ]
4597 .into(),
4598 );
4599 }
4600 }
4601 }
4602 }
4603 }
4604 }
4605
4606 (proposal_result, state_before_map)
4607 }
4608
4609 /// Process a single action: capability → validate → policy → idempotency → execute.
4610 /// (Idempotency dedup runs AFTER the deny checks — review C3: a durable
4611 /// cached result must never outlive a tool's revocation.)
4612 /// Returns (ActionResult, retries_count).
4613 async fn process_action(
4614 &self,
4615 action: &Action,
4616 proposal_id: &str,
4617 trace_id: &str,
4618 parent_span_id: &str,
4619 session_id: Option<&str>,
4620 scope: Option<&crate::scope::RuntimeScope>,
4621 ) -> (ActionResult, u32) {
4622 // Derive action type name for span naming
4623 let action_type_name = serde_json::to_string(&action.action_type)
4624 .unwrap_or_default()
4625 .trim_matches('"')
4626 .to_string();
4627 let span_name = format!("action.{}", action_type_name);
4628
4629 // Begin child span for this action
4630 let action_span_id = {
4631 let mut attrs: HashMap<String, Value> = HashMap::new();
4632 attrs.insert("action_id".to_string(), Value::from(action.id.as_str()));
4633 if let Some(ref tool) = action.tool {
4634 attrs.insert("tool".to_string(), Value::from(tool.as_str()));
4635 }
4636 let mut log = self.log.lock().await;
4637 log.begin_span(&span_name, trace_id, Some(parent_span_id), attrs)
4638 };
4639
4640 // Execute the action pipeline and capture result
4641 let (result, retries) = self
4642 .process_action_inner(action, proposal_id, session_id, scope)
4643 .await;
4644
4645 // End action span based on result status
4646 let span_status = match result.status {
4647 ActionStatus::Succeeded => SpanStatus::Ok,
4648 ActionStatus::Failed | ActionStatus::Rejected => SpanStatus::Error,
4649 _ => SpanStatus::Unset,
4650 };
4651 {
4652 let mut log = self.log.lock().await;
4653 log.end_span(&action_span_id, span_status);
4654 }
4655
4656 (result, retries)
4657 }
4658
4659 /// Inner action processing: idempotency -> validate -> policy -> execute.
4660 /// Returns (ActionResult, retries_count).
4661 #[instrument(
4662 name = "action.process",
4663 skip_all,
4664 fields(
4665 action_id = %action.id,
4666 action_type = ?action.action_type,
4667 tool = action.tool.as_deref().unwrap_or("none"),
4668 )
4669 )]
4670 async fn process_action_inner(
4671 &self,
4672 action: &Action,
4673 proposal_id: &str,
4674 session_id: Option<&str>,
4675 scope: Option<&crate::scope::RuntimeScope>,
4676 ) -> (ActionResult, u32) {
4677 // Capability check
4678 {
4679 let caps = self.capabilities.read().await;
4680 if let Some(ref cap) = *caps {
4681 // Check tool capability for ToolCall actions
4682 if action.action_type == ActionType::ToolCall {
4683 if let Some(ref tool_name) = action.tool {
4684 if !cap.tool_allowed(tool_name) {
4685 let reason =
4686 format!("capability denied: tool '{}' not allowed", tool_name);
4687 let mut log = self.log.lock().await;
4688 log.append(
4689 EventKind::ActionRejected,
4690 Some(&action.id),
4691 Some(proposal_id),
4692 action_outcome_data("capability", &reason, false, None),
4693 );
4694 return (rejected_result(&action.id, reason), 0);
4695 }
4696 }
4697 }
4698
4699 // Check state key capability for StateWrite/StateRead actions
4700 if action.action_type == ActionType::StateWrite
4701 || action.action_type == ActionType::StateRead
4702 {
4703 if let Some(key) = action.parameters.get("key").and_then(|v| v.as_str()) {
4704 if !cap.state_key_allowed(key) {
4705 let reason =
4706 format!("capability denied: state key '{}' not allowed", key);
4707 let mut log = self.log.lock().await;
4708 log.append(
4709 EventKind::ActionRejected,
4710 Some(&action.id),
4711 Some(proposal_id),
4712 action_outcome_data("capability", &reason, false, None),
4713 );
4714 return (rejected_result(&action.id, reason), 0);
4715 }
4716 }
4717 }
4718 }
4719 }
4720
4721 // Validate
4722 let tools = self.tools.read().await;
4723 let validation = validate_action(action, &self.state, &tools);
4724 drop(tools);
4725
4726 if !validation.valid() {
4727 let error = validation
4728 .errors
4729 .iter()
4730 .map(|e| e.reason.as_str())
4731 .collect::<Vec<_>>()
4732 .join("; ");
4733 let mut log = self.log.lock().await;
4734 log.append(
4735 EventKind::ActionRejected,
4736 Some(&action.id),
4737 Some(proposal_id),
4738 action_outcome_data("validation", &error, false, None),
4739 );
4740 return (rejected_result(&action.id, error), 0);
4741 }
4742
4743 // Policy check — global registry plus, when the proposal is
4744 // executed under a session, that session's registry. Both
4745 // layers must pass for the action to proceed; session
4746 // policies are additive deny rules — they can deny what
4747 // global allows but cannot allow what global denies.
4748 {
4749 let mut violations = {
4750 let policies = self.policies.read().await;
4751 policies.check(action, &self.state)
4752 };
4753 if let Some(sid) = session_id {
4754 // Snapshot the per-session engine handle out from under
4755 // the outer registry lock so the inner check holds
4756 // only the engine's own RwLock — preserves the
4757 // documented lock-ordering discipline.
4758 let session_engine = {
4759 let sessions = self.session_policies.read().await;
4760 sessions.get(sid).cloned()
4761 };
4762 if let Some(engine) = session_engine {
4763 let engine = engine.read().await;
4764 violations.extend(engine.check(action, &self.state));
4765 } else {
4766 let error = format!(
4767 "unknown session id '{sid}' — open one via Runtime::open_session before executing under a session"
4768 );
4769 // Unknown session id — refuse the action rather
4770 // than silently fall back to global-only. A
4771 // proposal submitted under a closed session
4772 // shouldn't run with looser rules than the caller
4773 // intended.
4774 let mut log = self.log.lock().await;
4775 log.append(
4776 EventKind::PolicyViolation,
4777 Some(&action.id),
4778 Some(proposal_id),
4779 action_outcome_data("policy", &error, false, None),
4780 );
4781 return (rejected_result(&action.id, error), 0);
4782 }
4783 }
4784 if !violations.is_empty() {
4785 let error = violations
4786 .iter()
4787 .map(|v| format!("policy '{}': {}", v.policy_name, v.reason))
4788 .collect::<Vec<_>>()
4789 .join("; ");
4790 let mut log = self.log.lock().await;
4791 log.append(
4792 EventKind::PolicyViolation,
4793 Some(&action.id),
4794 Some(proposal_id),
4795 action_outcome_data("policy", &error, false, None),
4796 );
4797 return (rejected_result(&action.id, error), 0);
4798 }
4799 }
4800
4801 // Idempotency check — deliberately AFTER capability + policy
4802 // (linus review C3): with the durable journal, a cached result
4803 // consulted first would let a tool that is denied TODAY serve
4804 // yesterday's cached result after every restart, forever. Deny
4805 // rules must win over dedup.
4806 if action.idempotent && !action.invocation_mode.is_detached() {
4807 let key = idempotency_key(action, scope);
4808 let cache = self.idempotency_cache.lock().await;
4809 if let Some(cached) = cache.get(&key) {
4810 let mut log = self.log.lock().await;
4811 log.append(
4812 EventKind::ActionDeduplicated,
4813 Some(&action.id),
4814 Some(proposal_id),
4815 [
4816 (
4817 "cached_action_id".to_string(),
4818 Value::from(cached.action_id.as_str()),
4819 ),
4820 ("attempted".to_string(), Value::from(false)),
4821 ("stage".to_string(), Value::from("idempotency_cache")),
4822 ]
4823 .into(),
4824 );
4825 return (
4826 ActionResult {
4827 action_id: action.id.clone(),
4828 status: cached.status.clone(),
4829 output: cached.output.clone(),
4830 error: cached.error.clone(),
4831 terminal: cached.terminal,
4832 // The cached action committed these mutations in an
4833 // earlier proposal. This proposal did not observe or
4834 // apply them, so its actual-change evidence is empty.
4835 state_changes: HashMap::new(),
4836 rolled_back: false,
4837 duration_ms: Some(0.0),
4838 timestamp: chrono::Utc::now(),
4839 },
4840 0,
4841 );
4842 }
4843 }
4844
4845 // Validated
4846 {
4847 let mut log = self.log.lock().await;
4848 log.append(
4849 EventKind::ActionValidated,
4850 Some(&action.id),
4851 Some(proposal_id),
4852 HashMap::new(),
4853 );
4854 }
4855
4856 // Execute with retry
4857 let (result, retries) = self
4858 .execute_with_retry(action, proposal_id, session_id, scope)
4859 .await;
4860
4861 // Cache idempotent results. Detached invocations are exempt
4862 // (linus review D4): their Succeeded output is a live tool
4863 // HANDLE, exactly as uncacheable here as in the result cache —
4864 // deduping would return a stale handle instead of starting the
4865 // tool, and journaling would replay a handle id that doesn't
4866 // exist in a fresh registry after restart.
4867 if action.idempotent
4868 && !action.invocation_mode.is_detached()
4869 && result.status == ActionStatus::Succeeded
4870 {
4871 let key = idempotency_key(action, scope);
4872 {
4873 let mut cache = self.idempotency_cache.lock().await;
4874 cache.insert(key.clone(), result.clone());
4875 }
4876 // Persist to the durable journal (C3) so the result survives a
4877 // restart and the action isn't re-executed.
4878 self.journal_idempotency(&key, Some(&result)).await;
4879 }
4880
4881 tracing::info!(
4882 status = ?result.status,
4883 duration_ms = result.duration_ms,
4884 "action completed"
4885 );
4886
4887 (result, retries)
4888 }
4889
4890 /// Execute with retry logic and timeout.
4891 /// Returns (ActionResult, retries_count).
4892 async fn execute_with_retry(
4893 &self,
4894 action: &Action,
4895 proposal_id: &str,
4896 session_id: Option<&str>,
4897 scope: Option<&crate::scope::RuntimeScope>,
4898 ) -> (ActionResult, u32) {
4899 // The harness config (Evolution Agent, §3.5), when installed, caps
4900 // the per-action retry budget and sets the inter-attempt backoff —
4901 // this is where an applied retry-config mutation actually takes
4902 // effect. `None` (the default) preserves prior behavior exactly.
4903 let hc = self.harness_config.read().await.clone();
4904 let retry_cap = hc.as_ref().map(|c| c.max_retries).unwrap_or(u32::MAX);
4905 let backoff_override = hc.as_ref().map(|c| c.retry_backoff_ms).filter(|&v| v > 0);
4906 let max_attempts = if action.failure_behavior == FailureBehavior::Retry {
4907 match action.max_retries.min(retry_cap).checked_add(1) {
4908 Some(max_attempts) => max_attempts,
4909 None => {
4910 let reason =
4911 format!("action '{}' retry attempt count overflows u32", action.id);
4912 self.log.lock().await.append(
4913 EventKind::ActionRejected,
4914 Some(&action.id),
4915 Some(proposal_id),
4916 action_outcome_data("retry_admission", &reason, false, None),
4917 );
4918 return (rejected_result(&action.id, reason), 0);
4919 }
4920 }
4921 } else {
4922 1
4923 };
4924
4925 let mut last_error: Option<String> = None;
4926 let mut retries: u32 = 0;
4927 let tool_source = self.tool_source_for_action(action).await;
4928 let params_digest = action_params_digest(action);
4929
4930 for attempt in 0..max_attempts {
4931 if attempt > 0 {
4932 retries += 1;
4933 // Base delay from the harness config when installed (an
4934 // applied retry-config mutation), else the built-in default.
4935 let base = backoff_override.unwrap_or(RETRY_BASE_DELAY_MS);
4936 let delay = base * RETRY_BACKOFF_FACTOR.pow(attempt - 1);
4937 tokio::time::sleep(Duration::from_millis(delay)).await;
4938 let mut log = self.log.lock().await;
4939 log.append(
4940 EventKind::ActionRetrying,
4941 Some(&action.id),
4942 Some(proposal_id),
4943 [("attempt".to_string(), Value::from(attempt + 1))].into(),
4944 );
4945 }
4946
4947 {
4948 let mut data: HashMap<String, Value> =
4949 [("attempt".to_string(), Value::from(attempt + 1))].into();
4950 insert_tool_event_provenance(&mut data, action, tool_source);
4951 let mut log = self.log.lock().await;
4952 log.append(
4953 EventKind::ActionExecuting,
4954 Some(&action.id),
4955 Some(proposal_id),
4956 data,
4957 );
4958 }
4959
4960 let start = std::time::Instant::now();
4961 let transitions_before = self.state.transition_count();
4962
4963 // Execute with optional timeout. Keep the timeout bit typed until
4964 // event classification rather than trying to recover it from the
4965 // human-readable error later.
4966 let (exec_result, engine_timeout) = if let Some(timeout_ms) = action.timeout_ms {
4967 match timeout(
4968 Duration::from_millis(timeout_ms),
4969 self.dispatch(action, session_id, scope, attempt + 1),
4970 )
4971 .await
4972 {
4973 Ok(result) => (result, false),
4974 Err(_) => (
4975 Err(ToolFailure::ordinary(format!(
4976 "action timed out after {}ms",
4977 timeout_ms
4978 ))),
4979 true,
4980 ),
4981 }
4982 } else {
4983 (
4984 self.dispatch(action, session_id, scope, attempt + 1).await,
4985 false,
4986 )
4987 };
4988
4989 let duration_ms = start.elapsed().as_secs_f64() * 1000.0;
4990
4991 match exec_result {
4992 Ok(output) => {
4993 if let Err(error) = car_inference::catalog_identity::canonical_json(&output) {
4994 let reason = format!(
4995 "tool output failed JCS/I-JSON validation after dispatch; external effect may have occurred and was not undone: {error}"
4996 );
4997 let mut data = action_outcome_data(
4998 "output_validation",
4999 &reason,
5000 true,
5001 Some(attempt + 1),
5002 );
5003 data.insert(
5004 "external_effect_status".to_string(),
5005 Value::from("may_have_occurred_not_undone"),
5006 );
5007 insert_tool_event_provenance(&mut data, action, tool_source);
5008 insert_action_outcome_signal(
5009 &mut data,
5010 action,
5011 ¶ms_digest,
5012 Some("validation"),
5013 );
5014 if action.tool.is_some() {
5015 data.insert("ok".to_string(), Value::from(false));
5016 }
5017 self.log.lock().await.append(
5018 EventKind::ActionFailed,
5019 Some(&action.id),
5020 Some(proposal_id),
5021 data,
5022 );
5023 return (
5024 ActionResult {
5025 action_id: action.id.clone(),
5026 status: ActionStatus::Failed,
5027 output: None,
5028 error: Some(reason),
5029 terminal: false,
5030 state_changes: HashMap::new(),
5031 rolled_back: false,
5032 duration_ms: Some(duration_ms),
5033 timestamp: chrono::Utc::now(),
5034 },
5035 retries,
5036 );
5037 }
5038
5039 // Snapshot the changes made by dispatch before applying
5040 // any result accounting. Deletion and set-to-null are
5041 // intentionally tagged differently on the authenticated
5042 // wire; a bare JSON null cannot represent both.
5043 let runtime_state_changes: HashMap<String, Value> = self
5044 .state
5045 .transitions_since(transitions_before)
5046 .into_iter()
5047 // Actions at the same DAG level execute concurrently,
5048 // so the shared transition tail can contain a sibling's
5049 // writes. StateTransition.action_id is assigned at the
5050 // mutation site and is the authoritative attribution
5051 // boundary for this per-action journal field.
5052 .filter(|transition| transition.action_id == action.id)
5053 .map(|transition| {
5054 (
5055 transition.key,
5056 StateMutation::from_new_value(transition.new_value).encode(),
5057 )
5058 })
5059 .collect();
5060
5061 // Declared expected effects are plan assertions only.
5062 // They must never be applied as if the runtime observed
5063 // them, nor merged into committed result evidence.
5064 let tenant_for_effects = scope.and_then(|s| s.tenant_id.as_deref());
5065 let state_changes = runtime_state_changes.clone();
5066
5067 // Record this result's provenance for the VIGIL intent
5068 // gate (`crate::taint`) — the only point in the system
5069 // that knows which SPECIFIC results were tainted. Only
5070 // on success: a failed action commits no effects, so it
5071 // changes no taint. No intent gate ⇒ no ledger ⇒ one
5072 // `Option` read and nothing else.
5073 let ledger = self.taint_ledger.read().await.clone();
5074 if let Some(ledger) = ledger {
5075 ledger
5076 .record_result(
5077 tenant_for_effects,
5078 action,
5079 state_changes.keys().cloned(),
5080 )
5081 .await;
5082 }
5083
5084 let mut log = self.log.lock().await;
5085 // Record the tool name + result count so the tool-receipt
5086 // verifier (A6) can project ground-truth receipts from the
5087 // log and catch tool-use hallucinations.
5088 let mut succ_data: HashMap<String, Value> = HashMap::new();
5089 succ_data.insert("attempt".to_string(), Value::from(attempt + 1));
5090 succ_data.insert("attempted".to_string(), Value::from(true));
5091 succ_data.insert("stage".to_string(), Value::from("dispatch"));
5092 insert_tool_event_provenance(&mut succ_data, action, tool_source);
5093 insert_action_outcome_signal(&mut succ_data, action, ¶ms_digest, None);
5094 if action.tool.is_some() {
5095 succ_data.insert("ok".to_string(), Value::from(true));
5096 if let Value::Array(arr) = &output {
5097 succ_data
5098 .insert("result_count".to_string(), Value::from(arr.len() as u64));
5099 }
5100 }
5101 // Record duration through the standardized metric path
5102 // (§3.5.1) so it feeds `EventLog::metrics_totals` via the
5103 // same contract as inference token metrics — one path,
5104 // not a coincidentally-matching `"duration_ms"` string.
5105 log.append_metered(
5106 EventKind::ActionSucceeded,
5107 Some(&action.id),
5108 Some(proposal_id),
5109 succ_data,
5110 car_eventlog::Metrics::latency(duration_ms),
5111 );
5112
5113 if !action.expected_effects.is_empty() || !runtime_state_changes.is_empty() {
5114 log.append(
5115 EventKind::StateChanged,
5116 Some(&action.id),
5117 Some(proposal_id),
5118 [
5119 (
5120 "declared_expected_effects".to_string(),
5121 serde_json::to_value(&action.expected_effects)
5122 .unwrap_or_default(),
5123 ),
5124 (
5125 "runtime_state_mutations".to_string(),
5126 serde_json::to_value(&runtime_state_changes)
5127 .unwrap_or_default(),
5128 ),
5129 (
5130 "changes".to_string(),
5131 serde_json::to_value(&state_changes).unwrap_or_default(),
5132 ),
5133 (
5134 "changes_semantics".to_string(),
5135 Value::from("provisional_state_mutations"),
5136 ),
5137 ("attempt".to_string(), Value::from(attempt + 1)),
5138 ("attempted".to_string(), Value::from(true)),
5139 ("stage".to_string(), Value::from("state_effect")),
5140 ]
5141 .into(),
5142 );
5143 }
5144
5145 return (
5146 ActionResult {
5147 action_id: action.id.clone(),
5148 status: ActionStatus::Succeeded,
5149 output: Some(output),
5150 error: None,
5151 terminal: false,
5152 state_changes,
5153 rolled_back: false,
5154 duration_ms: Some(duration_ms),
5155 timestamp: chrono::Utc::now(),
5156 },
5157 retries,
5158 );
5159 }
5160 Err(failure) => {
5161 let terminal = match failure.classification {
5162 ToolFailureClassification::Ordinary => false,
5163 ToolFailureClassification::Terminal => true,
5164 };
5165 let error = failure.message;
5166 last_error = Some(error.clone());
5167 let mut log = self.log.lock().await;
5168 let mut fail_data: HashMap<String, Value> = [
5169 ("error".to_string(), Value::from(error.as_str())),
5170 ("attempt".to_string(), Value::from(attempt + 1)),
5171 ("attempted".to_string(), Value::from(true)),
5172 ("stage".to_string(), Value::from("dispatch")),
5173 ]
5174 .into();
5175 // Tool name + ok=false so the receipt verifier (A6) records
5176 // that the tool *executed* (a failed call still ran).
5177 insert_tool_event_provenance(&mut fail_data, action, tool_source);
5178 insert_action_outcome_signal(
5179 &mut fail_data,
5180 action,
5181 ¶ms_digest,
5182 Some(action_error_class(action, &error, engine_timeout)),
5183 );
5184 if action.tool.is_some() {
5185 fail_data.insert("ok".to_string(), Value::from(false));
5186 }
5187 if terminal {
5188 fail_data.insert("terminal".to_string(), Value::from(true));
5189 }
5190 log.append(
5191 EventKind::ActionFailed,
5192 Some(&action.id),
5193 Some(proposal_id),
5194 fail_data,
5195 );
5196 drop(log);
5197
5198 if terminal {
5199 return (
5200 ActionResult {
5201 action_id: action.id.clone(),
5202 status: ActionStatus::Failed,
5203 output: None,
5204 error: Some(error),
5205 terminal: true,
5206 rolled_back: false,
5207 state_changes: HashMap::new(),
5208 duration_ms: Some(duration_ms),
5209 timestamp: chrono::Utc::now(),
5210 },
5211 retries,
5212 );
5213 }
5214 }
5215 }
5216 }
5217
5218 // All attempts exhausted
5219 if action.failure_behavior == FailureBehavior::Skip {
5220 let reason = last_error.as_deref().unwrap_or("all attempts exhausted");
5221 self.log.lock().await.append(
5222 EventKind::ActionSkipped,
5223 Some(&action.id),
5224 Some(proposal_id),
5225 action_outcome_data("failure_behavior", reason, true, Some(max_attempts)),
5226 );
5227 return (skipped_result(&action.id, reason), retries);
5228 }
5229
5230 (
5231 ActionResult {
5232 action_id: action.id.clone(),
5233 status: ActionStatus::Failed,
5234 output: None,
5235 error: last_error,
5236 terminal: false,
5237 state_changes: HashMap::new(),
5238 rolled_back: false,
5239 duration_ms: None,
5240 timestamp: chrono::Utc::now(),
5241 },
5242 retries,
5243 )
5244 }
5245
5246 /// Resolve the stable source category for an action's tool. Canonical
5247 /// ToolEntry provenance wins; the ToolSchema fallback keeps legacy direct
5248 /// registrations observable, and the final default covers old callers that
5249 /// supplied a tool field on a non-tool action.
5250 async fn tool_source_for_action(&self, action: &Action) -> Option<car_ir::ToolSourceKind> {
5251 let tool = action.tool.as_deref()?;
5252 if let Some(entry) = self.registry.get(tool).await {
5253 return Some(entry.source.kind());
5254 }
5255 Some(
5256 self.tools
5257 .read()
5258 .await
5259 .get(tool)
5260 .map(|schema| schema.source)
5261 .unwrap_or(car_ir::ToolSourceKind::UserDefined),
5262 )
5263 }
5264
5265 /// Dispatch an action to the appropriate handler.
5266 ///
5267 /// `scope` is the per-execution caller / tenant surface
5268 /// (Parslee-ai/car#187 phase 3). When the scope carries a
5269 /// tenant id, state R/W operations (`StateWrite`, `StateRead`,
5270 /// `Assertion`) route through `StateStore::scoped(tenant_id)`
5271 /// so distinct tenants can't see each other's keys. Unscoped
5272 /// proposals get the legacy flat-namespace behaviour
5273 /// automatically.
5274 async fn dispatch(
5275 &self,
5276 action: &Action,
5277 session_id: Option<&str>,
5278 scope: Option<&crate::scope::RuntimeScope>,
5279 attempt: u32,
5280 ) -> Result<Value, ToolFailure> {
5281 match action.action_type {
5282 ActionType::ToolCall => {
5283 let tool_name = action.tool.as_deref().ok_or("tool_call has no tool")?;
5284 let params = Value::Object(
5285 action
5286 .parameters
5287 .iter()
5288 .map(|(k, v)| (k.clone(), v.clone()))
5289 .collect(),
5290 );
5291
5292 // Detached invocation modes (C2): start the tool via the
5293 // configured executor's streaming entry point, register a
5294 // handle, and return immediately — the DAG must not block
5295 // on a streaming/long-running tool. Deliberately BEFORE
5296 // the result cache (a handle is a live invocation, never a
5297 // cacheable value) but still behind the rate limiter.
5298 // Built-ins don't stream; a detached call requires a
5299 // configured executor that implements execute_stream.
5300 if action.invocation_mode.is_detached() {
5301 self.rate_limiter.acquire(tool_name).await;
5302 let configured = {
5303 let guard = self.tool_executor.lock().await;
5304 guard.as_ref().cloned()
5305 };
5306 let executor = configured.ok_or_else(|| {
5307 format!("tool '{tool_name}': detached invocation requires a tool executor")
5308 })?;
5309 let rx = executor
5310 .execute_stream(tool_name, ¶ms, &action.id)
5311 .await?;
5312 let (handle, cancel) = self.tool_handles.register(tool_name, &action.id).await;
5313 crate::tool_handles::spawn_drain(
5314 self.tool_handles.clone(),
5315 handle.id.clone(),
5316 rx,
5317 cancel,
5318 );
5319 return Ok(serde_json::json!({
5320 "tool_handle": handle.id,
5321 "status": "running",
5322 }));
5323 }
5324
5325 // Check cross-proposal result cache.
5326 if let Some(cached) = self.result_cache.get(tool_name, ¶ms).await {
5327 return Ok(cached);
5328 }
5329
5330 // Apply rate limit backpressure before executing.
5331 self.rate_limiter.acquire(tool_name).await;
5332
5333 // Try built-in inference tools first when the inference engine is available.
5334 if matches!(
5335 tool_name,
5336 "infer" | "infer.grounded" | "embed" | "classify" | "transcribe" | "synthesize"
5337 ) {
5338 if let Some(ref engine) = self.inference_engine {
5339 // For "infer.grounded" or "infer" with memgine available,
5340 // build context from memory and attach it to the request.
5341 let params = {
5342 let should_ground =
5343 tool_name == "infer.grounded" || tool_name == "infer";
5344 if should_ground {
5345 if let Some(ref memgine) = self.memgine {
5346 if let Some(prompt) =
5347 params.get("prompt").and_then(|v| v.as_str())
5348 {
5349 let ctx = {
5350 let mut m = memgine.lock().await;
5351 m.build_context(prompt)
5352 };
5353 if !ctx.is_empty() {
5354 let mut p = params.clone();
5355 if let Some(obj) = p.as_object_mut() {
5356 obj.insert("context".to_string(), Value::from(ctx));
5357 }
5358 p
5359 } else {
5360 params
5361 }
5362 } else {
5363 params
5364 }
5365 } else {
5366 params
5367 }
5368 } else {
5369 params
5370 }
5371 };
5372
5373 // Route "infer.grounded" to "infer" for the service layer
5374 let effective_tool = if tool_name == "infer.grounded" {
5375 "infer"
5376 } else {
5377 tool_name
5378 };
5379 let result =
5380 car_inference::service::execute_tool(engine, effective_tool, ¶ms)
5381 .await
5382 .map_err(|e| e.to_string());
5383
5384 if let Ok(ref value) = result {
5385 self.result_cache
5386 .put(tool_name, ¶ms, value.clone())
5387 .await;
5388 }
5389
5390 return result.map_err(ToolFailure::from);
5391 }
5392 }
5393
5394 // Built-in memory consolidation tool.
5395 if tool_name == "memory.consolidate" {
5396 if let Some(ref memgine) = self.memgine {
5397 let report = {
5398 let mut m = memgine.lock().await;
5399 m.consolidate().await
5400 };
5401 // Log the consolidation event
5402 {
5403 let mut log = self.log.lock().await;
5404 log.append(
5405 EventKind::Consolidated,
5406 None,
5407 None,
5408 [
5409 (
5410 "expired_pruned".to_string(),
5411 Value::from(report.expired_pruned),
5412 ),
5413 (
5414 "superseded_gc".to_string(),
5415 Value::from(report.superseded_gc),
5416 ),
5417 (
5418 "stale_embeddings_removed".to_string(),
5419 Value::from(report.stale_embeddings_removed),
5420 ),
5421 (
5422 "nodes_embedded".to_string(),
5423 Value::from(report.nodes_embedded),
5424 ),
5425 (
5426 "domains_evolved".to_string(),
5427 Value::from(report.domains_evolved.clone()),
5428 ),
5429 ("total_nodes".to_string(), Value::from(report.total_nodes)),
5430 ("total_edges".to_string(), Value::from(report.total_edges)),
5431 ]
5432 .into(),
5433 );
5434 // Per-candidate gate telemetry (SkillOpt-inspired):
5435 // one event per promotion/rejection so the skill
5436 // bank's evolution is auditable.
5437 for key in &report.candidates_promoted {
5438 log.append(
5439 EventKind::CandidatePromoted,
5440 None,
5441 None,
5442 [("candidate".to_string(), Value::from(key.as_str()))].into(),
5443 );
5444 }
5445 for key in &report.candidates_rejected {
5446 log.append(
5447 EventKind::CandidateRejected,
5448 None,
5449 None,
5450 [("candidate".to_string(), Value::from(key.as_str()))].into(),
5451 );
5452 }
5453 }
5454 return Ok(serde_json::to_value(&report).unwrap_or(Value::Null));
5455 } else {
5456 return Err(
5457 "memory.consolidate requires memgine (attach with with_learning)"
5458 .into(),
5459 );
5460 }
5461 }
5462
5463 // Built-in outbound human messaging. Sits here, alongside the
5464 // other built-ins — downstream of `validate_action` and the
5465 // policy check in the caller, and downstream of the result
5466 // cache read and `rate_limiter.acquire` above — so a message
5467 // to a human traverses exactly the same governed chain
5468 // (validator → policy → rate limit → eventlog) as any other
5469 // side effect, instead of the hand-rolled transport each agent
5470 // used to carry.
5471 if tool_name == "messaging.send" {
5472 // No fall-through to the host executor when unconfigured.
5473 // Falling through would let an ungoverned host transport
5474 // answer the call — the exact path this built-in exists to
5475 // close — so an unconfigured runtime refuses instead.
5476 let Some(ref sink) = self.message_sink else {
5477 return Err("messaging.send: messaging is not configured on this \
5478 runtime — attach a message sink via \
5479 Runtime::with_message_sink"
5480 .into());
5481 };
5482 let msg = crate::messaging::OutboundMessage::from_tool_params(¶ms)?;
5483 let receipt = sink.send(&msg).await?;
5484 // Not cached: the schema declares no cache TTL, and a
5485 // cached send would swallow a second, genuinely wanted
5486 // message with identical text.
5487 return serde_json::to_value(&receipt).map_err(|error| {
5488 ToolFailure::ordinary(format!("messaging.send: serialize receipt: {error}"))
5489 });
5490 }
5491
5492 // Prefer a configured tool_executor for any tool it claims to handle.
5493 // Fall through to agent_basics only when the configured executor is absent
5494 // or explicitly returns "unknown tool" — this prevents agent_basics' built-in
5495 // read_file/write_file (which resolve paths via std::env::current_dir) from
5496 // silently overriding an executor that carries its own working_dir.
5497 let configured = {
5498 let guard = self.tool_executor.lock().await;
5499 guard.as_ref().cloned()
5500 };
5501
5502 if let Some(ref executor) = configured {
5503 let return_schema = self
5504 .tools
5505 .read()
5506 .await
5507 .get(tool_name)
5508 .and_then(|schema| schema.returns.clone());
5509 let result = executor
5510 .execute_classified(
5511 tool_name,
5512 ¶ms,
5513 &action.id,
5514 action.timeout_ms,
5515 session_id,
5516 attempt,
5517 &action.expected_effects,
5518 return_schema.as_ref(),
5519 )
5520 .await;
5521 let fall_through = matches!(
5522 &result,
5523 Err(error) if error.classification == ToolFailureClassification::Ordinary
5524 && error.message.starts_with("unknown tool")
5525 );
5526 if !fall_through {
5527 match result {
5528 Ok(execution) => {
5529 if !execution.state_changes.is_empty() {
5530 let expected_keys: std::collections::BTreeSet<&str> = action
5531 .expected_effects
5532 .keys()
5533 .map(String::as_str)
5534 .collect();
5535 let actual_keys: std::collections::BTreeSet<&str> = execution
5536 .state_changes
5537 .keys()
5538 .map(String::as_str)
5539 .collect();
5540 if expected_keys != actual_keys {
5541 return Err(format!(
5542 "tool '{}' callback state_changes keys do not match action '{}': expected {:?}, got {:?}",
5543 tool_name, action.id, expected_keys, actual_keys
5544 )
5545 .into());
5546 }
5547 car_inference::catalog_identity::canonical_json(&execution.output)
5548 .map_err(|error| {
5549 format!(
5550 "tool output failed JCS/I-JSON validation after dispatch; external effect may have occurred and was not undone: {error}"
5551 )
5552 })?;
5553 let changes_value = serde_json::to_value(&execution.state_changes)
5554 .map_err(|error| {
5555 format!(
5556 "tool '{}' callback state_changes serialization failed: {error}",
5557 tool_name
5558 )
5559 })?;
5560 car_inference::catalog_identity::canonical_json(&changes_value)
5561 .map_err(|error| {
5562 format!(
5563 "tool '{}' callback state_changes failed JCS/I-JSON validation: {error}",
5564 tool_name
5565 )
5566 })?;
5567
5568 // Every rejection above happens before the
5569 // first write. Apply in deterministic key order,
5570 // attributed to the active action, through ONE
5571 // batch mutation boundary: a concurrent reader,
5572 // journal replay, and restart recovery observe
5573 // the complete old state or the complete new
5574 // state — never a prefix of this callback's
5575 // keys (Parslee-ai/car#1140).
5576 let mut changes: Vec<(String, Value)> = execution
5577 .state_changes
5578 .iter()
5579 .map(|(key, value)| (key.clone(), value.clone()))
5580 .collect();
5581 changes.sort_by(|(left, _), (right, _)| left.cmp(right));
5582 let tenant = scope.and_then(|scope| scope.tenant_id.as_deref());
5583 self.state.scoped(tenant).set_batch(changes, &action.id);
5584 }
5585 self.result_cache
5586 .put(tool_name, ¶ms, execution.output.clone())
5587 .await;
5588 return Ok(execution.output);
5589 }
5590 Err(error) => return Err(error),
5591 }
5592 }
5593 }
5594
5595 // Built-in commodity tools resolve against the bound substrate
5596 // (default: LocalSubstrate → historic host behavior). `calculate`
5597 // stays pure inside agent_basics and ignores the substrate.
5598 let substrate = self.substrate.read().await.clone();
5599 let read_ledger = self.read_ledgers.ledger_for(session_id);
5600 if let Some(result) = crate::agent_basics::execute_with_ledger(
5601 &substrate,
5602 &read_ledger,
5603 tool_name,
5604 ¶ms,
5605 )
5606 .await
5607 {
5608 if let Ok(ref value) = result {
5609 self.result_cache
5610 .put(tool_name, ¶ms, value.clone())
5611 .await;
5612 }
5613 return result.map_err(ToolFailure::from);
5614 }
5615
5616 Err(format!("no handler for tool '{}'", tool_name).into())
5617 }
5618 ActionType::StateWrite => {
5619 let key = action
5620 .parameters
5621 .get("key")
5622 .and_then(|v| v.as_str())
5623 .ok_or("state_write requires 'key' parameter")?;
5624 let value = action
5625 .parameters
5626 .get("value")
5627 .cloned()
5628 .unwrap_or(Value::Null);
5629 let tenant = scope.and_then(|s| s.tenant_id.as_deref());
5630 self.state.scoped(tenant).set(key, value, &action.id);
5631 Ok(Value::from(format!("written: {}", key)))
5632 }
5633 ActionType::StateRead => {
5634 let key = action
5635 .parameters
5636 .get("key")
5637 .and_then(|v| v.as_str())
5638 .ok_or("state_read requires 'key' parameter")?;
5639 let tenant = scope.and_then(|s| s.tenant_id.as_deref());
5640 Ok(self.state.scoped(tenant).get(key).unwrap_or(Value::Null))
5641 }
5642 ActionType::Assertion => {
5643 let key = action
5644 .parameters
5645 .get("key")
5646 .and_then(|v| v.as_str())
5647 .ok_or("assertion requires 'key' parameter")?;
5648 let expected = action
5649 .parameters
5650 .get("expected")
5651 .cloned()
5652 .unwrap_or(Value::Null);
5653 let tenant = scope.and_then(|s| s.tenant_id.as_deref());
5654 let actual = self.state.scoped(tenant).get(key).unwrap_or(Value::Null);
5655 if actual != expected {
5656 Err(format!(
5657 "assertion failed: state['{}'] = {:?}, expected {:?}",
5658 key, actual, expected
5659 )
5660 .into())
5661 } else {
5662 Ok(serde_json::json!({"asserted": key, "value": actual}))
5663 }
5664 }
5665 }
5666 }
5667
5668 // --- Checkpoint and resume ---
5669
5670 /// Save a checkpoint of the current runtime state.
5671 pub async fn save_checkpoint(&self) -> Checkpoint {
5672 let state = self.state.snapshot();
5673 let tools: Vec<String> = self.tools.read().await.keys().cloned().collect();
5674 let log = self.log.lock().await;
5675 let events: Vec<Value> = log
5676 .events()
5677 .iter()
5678 .map(|e| serde_json::to_value(e).unwrap_or_default())
5679 .collect();
5680
5681 Checkpoint {
5682 checkpoint_id: Uuid::new_v4().to_string(),
5683 created_at: chrono::Utc::now(),
5684 state,
5685 events,
5686 tools,
5687 metadata: HashMap::new(),
5688 }
5689 }
5690
5691 /// Save checkpoint to a JSON file.
5692 pub async fn save_checkpoint_to_file(&self, path: &str) -> Result<(), String> {
5693 let checkpoint = self.save_checkpoint().await;
5694 let json = serde_json::to_string_pretty(&checkpoint)
5695 .map_err(|e| format!("serialize error: {}", e))?;
5696 tokio::fs::write(path, json)
5697 .await
5698 .map_err(|e| format!("write error: {}", e))?;
5699 Ok(())
5700 }
5701
5702 /// Load a checkpoint from a JSON file and restore state.
5703 pub async fn load_checkpoint_from_file(&self, path: &str) -> Result<Checkpoint, String> {
5704 let json = tokio::fs::read_to_string(path)
5705 .await
5706 .map_err(|e| format!("read error: {}", e))?;
5707 let checkpoint: Checkpoint =
5708 serde_json::from_str(&json).map_err(|e| format!("deserialize error: {}", e))?;
5709 self.restore_checkpoint(&checkpoint).await;
5710 Ok(checkpoint)
5711 }
5712
5713 /// Restore runtime state from a checkpoint.
5714 pub async fn restore_checkpoint(&self, checkpoint: &Checkpoint) {
5715 // Replace state completely — don't merge, don't create synthetic transitions
5716 self.state.replace_all(checkpoint.state.clone());
5717 // Clear idempotency cache — stale results from pre-checkpoint execution
5718 // must not bypass validation/policy on the restored state
5719 self.idempotency_cache.lock().await.clear();
5720 // Restore tools (as name-only schemas; full schemas are not persisted in checkpoint)
5721 let mut tools = self.tools.write().await;
5722 tools.clear();
5723 for tool_name in &checkpoint.tools {
5724 let schema = ToolSchema {
5725 name: tool_name.clone(),
5726 source: car_ir::ToolSourceKind::UserDefined,
5727 description: String::new(),
5728 parameters: serde_json::Value::Object(Default::default()),
5729 returns: None,
5730 idempotent: false,
5731 cache_ttl_secs: None,
5732 rate_limit: None,
5733 };
5734 tools.insert(tool_name.clone(), schema);
5735 }
5736 }
5737
5738 /// Register a subprocess tool and set up the subprocess executor.
5739 /// If no executor exists, creates a new SubprocessToolExecutor.
5740 /// If one already exists, creates a new SubprocessToolExecutor with the
5741 /// existing executor as fallback.
5742 pub async fn register_subprocess_tool(
5743 &self,
5744 name: &str,
5745 tool: crate::subprocess::SubprocessTool,
5746 ) {
5747 use crate::subprocess::SubprocessToolExecutor;
5748
5749 let schema = ToolSchema {
5750 name: name.to_string(),
5751 source: car_ir::ToolSourceKind::Subprocess,
5752 description: format!("Subprocess tool: {}", tool.command),
5753 parameters: serde_json::Value::Object(Default::default()),
5754 returns: None,
5755 idempotent: false,
5756 cache_ttl_secs: None,
5757 rate_limit: None,
5758 };
5759 self.register_tool_entry(
5760 crate::registry::ToolEntry::new(schema)
5761 .with_source(crate::registry::ToolSource::Subprocess),
5762 )
5763 .await;
5764
5765 let mut guard = self.tool_executor.lock().await;
5766 let mut executor = match guard.take() {
5767 Some(existing) => {
5768 let mut sub = SubprocessToolExecutor::new();
5769 sub = sub.with_fallback(existing);
5770 sub
5771 }
5772 None => SubprocessToolExecutor::new(),
5773 };
5774 executor.register(name, tool);
5775 *guard = Some(std::sync::Arc::new(executor));
5776 }
5777}
5778
5779impl Default for Runtime {
5780 fn default() -> Self {
5781 Self::new()
5782 }
5783}
5784
5785#[cfg(test)]
5786mod timeout_integration_tests {
5787 //! Integration coverage for the #259/#262 timeout coordination the unit
5788 //! tests in `car-server-core` / `car-ffi-common` only exercise as pure
5789 //! selection helpers (#266 item 4): that the executor's per-action deadline
5790 //! reaps a slow dispatch first, and that `action.timeout_ms` flows
5791 //! end-to-end onto `execute_with_action`.
5792 use super::*;
5793 use car_ir::ActionProposal;
5794 use std::sync::atomic::{AtomicU64, Ordering};
5795
5796 /// A tool executor that (a) records the `timeout_ms` it was handed and
5797 /// (b) sleeps `delay_ms` before returning, so a test can prove the
5798 /// executor's own `timeout(action.timeout_ms, dispatch)` reaps a call that
5799 /// outlives its budget.
5800 struct RecordingExecutor {
5801 seen_timeout_ms: Arc<AtomicU64>,
5802 delay_ms: u64,
5803 }
5804
5805 #[async_trait::async_trait]
5806 impl ToolExecutor for RecordingExecutor {
5807 async fn execute(&self, _tool: &str, _params: &Value) -> Result<Value, String> {
5808 Ok(Value::Null)
5809 }
5810 async fn execute_with_action(
5811 &self,
5812 _tool: &str,
5813 _params: &Value,
5814 _action_id: &str,
5815 timeout_ms: Option<u64>,
5816 ) -> Result<Value, String> {
5817 // `u64::MAX` sentinel = "None was passed".
5818 self.seen_timeout_ms
5819 .store(timeout_ms.unwrap_or(u64::MAX), Ordering::SeqCst);
5820 tokio::time::sleep(Duration::from_millis(self.delay_ms)).await;
5821 Ok(serde_json::json!({ "ok": true }))
5822 }
5823 }
5824
5825 fn one_tool_proposal(timeout_ms: Option<u64>) -> ActionProposal {
5826 let mut action = serde_json::json!({
5827 "id": "a0",
5828 "type": "tool_call",
5829 "tool": "slow",
5830 "parameters": {},
5831 "dependencies": [],
5832 });
5833 if let Some(ms) = timeout_ms {
5834 action["timeout_ms"] = serde_json::json!(ms);
5835 }
5836 serde_json::from_value(serde_json::json!({
5837 "source": "test",
5838 "actions": [action],
5839 }))
5840 .expect("proposal deserializes")
5841 }
5842
5843 #[tokio::test(start_paused = true)]
5844 async fn action_timeout_ms_reaches_executor_and_reaps_slow_dispatch() {
5845 let seen = Arc::new(AtomicU64::new(0));
5846 let rt = Runtime::new();
5847 rt.register_tool("slow").await;
5848 rt.set_executor(Arc::new(RecordingExecutor {
5849 seen_timeout_ms: seen.clone(),
5850 delay_ms: 2_000, // outlives the 100ms budget below
5851 }))
5852 .await;
5853
5854 // 100ms budget against a 2s dispatch: the executor's own
5855 // `timeout(action.timeout_ms, dispatch)` must reap it (#262 — the
5856 // executor is the authority), and the budget must have reached
5857 // `execute_with_action` (the #262 harness→action propagation, proven
5858 // end-to-end here rather than only via the pure helper).
5859 let result = rt.execute(&one_tool_proposal(Some(100))).await;
5860 assert_eq!(
5861 seen.load(Ordering::SeqCst),
5862 100,
5863 "budget must reach the executor"
5864 );
5865 let action_result = &result.results[0];
5866 assert!(
5867 action_result
5868 .error
5869 .as_deref()
5870 .unwrap_or("")
5871 .contains("timed out"),
5872 "slow dispatch must be reaped by the action deadline: {:?}",
5873 action_result.error
5874 );
5875 }
5876
5877 /// An executor that records every `attempt` it is handed and fails until
5878 /// the Nth, so a test can watch the counter advance across real retries.
5879 struct AttemptRecordingExecutor {
5880 seen: Arc<tokio::sync::Mutex<Vec<u32>>>,
5881 succeed_on: u32,
5882 }
5883
5884 #[async_trait::async_trait]
5885 impl ToolExecutor for AttemptRecordingExecutor {
5886 async fn execute(&self, _tool: &str, _params: &Value) -> Result<Value, String> {
5887 Ok(Value::Null)
5888 }
5889 async fn execute_with_action_in_session(
5890 &self,
5891 _tool: &str,
5892 _params: &Value,
5893 _action_id: &str,
5894 _timeout_ms: Option<u64>,
5895 _session_id: Option<&str>,
5896 attempt: u32,
5897 ) -> Result<Value, String> {
5898 self.seen.lock().await.push(attempt);
5899 if attempt >= self.succeed_on {
5900 Ok(serde_json::json!({ "ok": true }))
5901 } else {
5902 Err("transient".to_string())
5903 }
5904 }
5905 }
5906
5907 /// Parslee-ai/car#928 — the retry counter the executor receives must be the
5908 /// engine's real one, and it must advance.
5909 ///
5910 /// The WS executor hardcoded `attempt: 1` onto the `tools.execute` payload,
5911 /// so the field a host would build a retry-disambiguating join on never
5912 /// varied. Asserting on the sequence the executor actually observes is the
5913 /// property that matters: a frame test can only prove the payload carries
5914 /// whatever it was handed, not that the engine hands it the truth.
5915 #[tokio::test]
5916 async fn the_executor_sees_the_engines_real_attempt_sequence() {
5917 let seen = Arc::new(tokio::sync::Mutex::new(Vec::new()));
5918 let rt = Runtime::new();
5919 rt.register_tool("flaky").await;
5920 rt.set_executor(Arc::new(AttemptRecordingExecutor {
5921 seen: seen.clone(),
5922 succeed_on: 3,
5923 }))
5924 .await;
5925
5926 let proposal: ActionProposal = serde_json::from_value(serde_json::json!({
5927 "source": "test",
5928 "actions": [{
5929 "id": "a0",
5930 "type": "tool_call",
5931 "tool": "flaky",
5932 "parameters": {},
5933 "dependencies": [],
5934 "failure_behavior": "retry",
5935 "max_retries": 3,
5936 }],
5937 }))
5938 .expect("proposal deserializes");
5939
5940 let result = rt.execute(&proposal).await;
5941 assert_eq!(
5942 result.results[0].status,
5943 ActionStatus::Succeeded,
5944 "third attempt succeeds: {:?}",
5945 result.results[0].error
5946 );
5947 assert_eq!(
5948 *seen.lock().await,
5949 vec![1, 2, 3],
5950 "1-based and advancing — a constant here is the #928 bug"
5951 );
5952 }
5953
5954 #[tokio::test(start_paused = true)]
5955 async fn no_budget_applies_no_executor_deadline() {
5956 // The `None` path: the executor applies NO deadline, so a dispatch that
5957 // outlives any per-attempt budget still completes (the WS callback wait
5958 // is the sole bound on that path, exercised in car-server-core). Tokio's
5959 // paused clock advances the 300ms dispatch without a wall-clock wait.
5960 let seen = Arc::new(AtomicU64::new(0));
5961 let rt = Runtime::new();
5962 rt.register_tool("slow").await;
5963 rt.set_executor(Arc::new(RecordingExecutor {
5964 seen_timeout_ms: seen.clone(),
5965 delay_ms: 300,
5966 }))
5967 .await;
5968
5969 let result = rt.execute(&one_tool_proposal(None)).await;
5970 assert_eq!(
5971 seen.load(Ordering::SeqCst),
5972 u64::MAX,
5973 "None budget must be forwarded as None"
5974 );
5975 assert_eq!(result.results[0].status, ActionStatus::Succeeded);
5976 }
5977}
5978
5979#[cfg(test)]
5980mod model_facing_result_tests {
5981 use super::*;
5982
5983 #[test]
5984 fn a_model_reads_a_successful_output_as_text_and_a_failure_with_its_prefix() {
5985 let ok = ActionResult {
5986 action_id: "a".into(),
5987 status: ActionStatus::Succeeded,
5988 output: Some(serde_json::json!({ "path": "a.py", "content": "x = \"1\"\ny = 2" })),
5989 error: None,
5990 terminal: false,
5991 state_changes: HashMap::new(),
5992 rolled_back: false,
5993 duration_ms: None,
5994 timestamp: chrono::Utc::now(),
5995 };
5996 assert_eq!(
5997 format_tool_result_for_model(&ok),
5998 "path: a.py\ncontent (13 bytes):\nx = \"1\"\ny = 2"
5999 );
6000 // The JSON form stays what it was for every other consumer.
6001 assert!(format_tool_result(&ok).starts_with('{'));
6002 let failed = ActionResult {
6003 status: ActionStatus::Failed,
6004 output: None,
6005 error: Some("boom".into()),
6006 ..ok
6007 };
6008 assert_eq!(format_tool_result_for_model(&failed), "[FAILED] boom");
6009 }
6010
6011 #[test]
6012 fn rendered_output_reads_back_exactly() {
6013 // Shapes tools actually produce, including content that looks like
6014 // fields and strings that look like other JSON types.
6015 for original in [
6016 serde_json::json!({
6017 "path": "config/locales/en.yml",
6018 "content": "en:\n x: 1\nerror:\n not_found: Missing\n",
6019 "total_lines": 4,
6020 }),
6021 serde_json::json!({
6022 "status": 200,
6023 "url": "https://example.com/a",
6024 "headers": {"content-type": "text/plain", "authorization": "Bearer abc"},
6025 "body": "hello\nstatus:\nall good\nzebra:\n",
6026 }),
6027 serde_json::json!({
6028 "exit_code": 1,
6029 "output": "make: ***\nzebra:\nx",
6030 "timed_out": false,
6031 }),
6032 serde_json::json!({
6033 "error": "null",
6034 "flag": "true",
6035 "id": "123456789012345678901234",
6036 "name": "123",
6037 "content-type": "text/plain",
6038 "empty": "",
6039 "padded": " x ",
6040 "quoted": "\"q\"",
6041 }),
6042 ] {
6043 let text = render_tool_output(&original);
6044 assert_eq!(
6045 parse_rendered_output(&text).as_ref(),
6046 Some(&original),
6047 "{text}"
6048 );
6049 }
6050 assert_eq!(parse_rendered_output("[FAILED] boom"), None);
6051 assert_eq!(parse_rendered_output("plain text"), None);
6052 // Runtime failure and refusal lines are not fields, even with `: ` in them.
6053 for line in [
6054 "[FAILED] tool 'read_file': no such file",
6055 "[REJECTED] policy 'deny_shell': shell denied",
6056 "[FAILED] canceled: user stop",
6057 "Tool call denied: not allowed\nsecond line without a key",
6058 ] {
6059 assert_eq!(parse_rendered_output(line), None, "{line}");
6060 }
6061 // A shell result cut in the middle by the budget keeps its verdict:
6062 // the stderr block claims bytes that now run through the marker.
6063 let shell = serde_json::json!({
6064 "exit_code": 1,
6065 "stderr": "warning: x\n".repeat(700),
6066 "stdout": "test a ... FAILED\n".repeat(2500),
6067 });
6068 let cut = truncate_keeping_ends(&render_tool_output(&shell), 16 * 1024);
6069 assert_eq!(parse_rendered_output(&cut).unwrap()["exit_code"], 1);
6070 // A block spliced into a cut tail cannot replace an outcome field.
6071 let spliced = "ok: false\nout (10 bytes):\nab\nok (4 bytes):\ntrue";
6072 assert_eq!(parse_rendered_output(spliced).unwrap()["ok"], false);
6073 // … and a cut mid-line after the fields keeps them.
6074 let cut = "exit_code: 1\nstdo\n…[truncated: 10 of 20 bytes]…";
6075 assert_eq!(parse_rendered_output(cut).unwrap()["exit_code"], 1);
6076 // Keys and values that look like the format's own syntax.
6077 let odd = serde_json::json!({
6078 "a:": "x\ny",
6079 "a: b": "c",
6080 "note": "wrote (3 bytes):",
6081 " padded": 1,
6082 "k (x": "y\nz",
6083 });
6084 let text = render_tool_output(&odd);
6085 assert_eq!(parse_rendered_output(&text).as_ref(), Some(&odd), "{text}");
6086 // A block cut by a budget reads back as what is left of it.
6087 let cut = "exit_code: 1\noutput (100 bytes):\nfirst lines only";
6088 assert_eq!(
6089 parse_rendered_output(cut).unwrap()["output"],
6090 "first lines only"
6091 );
6092 }
6093
6094 #[test]
6095 fn outcome_fields_lead_and_stay_on_one_line() {
6096 let out = render_tool_output(&serde_json::json!({
6097 "content": "x\ny",
6098 "error": format!("{}\nsecond line", "e".repeat(600)),
6099 "exit_code": 1,
6100 "matches": [{"path": "a"}],
6101 "truncated": true,
6102 }));
6103 let lines: Vec<&str> = out.lines().collect();
6104 assert!(lines[0].starts_with("error: "), "{out}");
6105 assert!(
6106 lines[0].len() <= 530 && lines[0].contains(" bytes]…"),
6107 "{out}"
6108 );
6109 assert!(
6110 lines[0].ends_with("second line"),
6111 "the end of a long error survives: {out}"
6112 );
6113 assert_eq!(lines[1], "exit_code: 1");
6114 assert_eq!(lines[2], "truncated: true");
6115 assert!(
6116 out.ends_with("content (3 bytes):\nx\ny\nmatches (1 records):\npath=a"),
6117 "{out}"
6118 );
6119 }
6120
6121 #[test]
6122 fn long_command_output_keeps_its_verdict() {
6123 let out = format!(
6124 "exit_code: 1\noutput:\n{}\n1 failed, 99 passed",
6125 "é.".repeat(20_000)
6126 );
6127 let cut = truncate_keeping_ends(&out, 4096);
6128 assert!(cut.len() <= 4096, "{}", cut.len());
6129 assert!(cut.starts_with("exit_code: 1"));
6130 assert!(cut.ends_with("1 failed, 99 passed"));
6131 assert!(cut.contains("bytes elided from the middle"));
6132 assert_eq!(truncate_keeping_ends("short", 4096), "short");
6133 assert!(truncate_keeping_ends(&out, 4096).contains(ELIDED_MARKER));
6134 let tiny = truncate_keeping_ends(&"x".repeat(1000), 64);
6135 assert!(tiny.len() <= 64, "{}", tiny.len());
6136 }
6137}