Skip to main content

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(&current_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(&current_proposal.id),
2774                proposal_rejection_boundary_data(
2775                    &current_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(&current_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(&current_proposal.id),
2893                    proposal_rejection_boundary_data(
2894                        &current_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                    &current_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(&current_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                    &current_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                    &current_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                        &current_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(&current_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                            &current_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                            &params_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, &params_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                        &params_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, &params, &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, &params).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, &params)
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, &params, 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(&params)?;
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                            &params,
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, &params, 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                    &params,
5605                )
5606                .await
5607                {
5608                    if let Ok(ref value) = result {
5609                        self.result_cache
5610                            .put(tool_name, &params, 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}