use car_memgine::MemgineEngine;
use car_proto::RunTermination;
use car_server_core::{run_dispatch, ServerState, ServerStateConfig};
use futures::{SinkExt, StreamExt};
use sha2::{Digest, Sha256};
use std::io::Write;
use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4};
use std::sync::Arc;
use tempfile::TempDir;
use tokio::net::TcpListener;
use tokio::sync::Mutex;
use tokio_tungstenite::{accept_async, connect_async, tungstenite::Message};
type Ws =
tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
fn loopback_state(journal_dir: std::path::PathBuf) -> Arc<ServerState> {
let engine = Arc::new(Mutex::new(MemgineEngine::new(None)));
let cfg = ServerStateConfig::new(journal_dir).with_shared_memgine(engine);
Arc::new(ServerState::with_config(cfg))
}
async fn spawn_dispatcher(state: Arc<ServerState>) -> SocketAddr {
let listener = TcpListener::bind(SocketAddr::V4(SocketAddrV4::new(
Ipv4Addr::new(127, 0, 0, 1),
0,
)))
.await
.expect("bind loopback");
let addr = listener.local_addr().expect("local_addr");
tokio::spawn(async move {
let (stream, peer) = listener.accept().await.expect("accept");
let ws = accept_async(stream).await.expect("ws handshake");
let (write, read) = ws.split();
let _ = run_dispatch(read, Box::pin(write), peer.to_string(), state).await;
});
addr
}
async fn spawn_multi_dispatcher(state: Arc<ServerState>) -> SocketAddr {
let listener = TcpListener::bind(SocketAddr::V4(SocketAddrV4::new(
Ipv4Addr::new(127, 0, 0, 1),
0,
)))
.await
.expect("bind loopback");
let addr = listener.local_addr().expect("local_addr");
tokio::spawn(async move {
while let Ok((stream, peer)) = listener.accept().await {
let state = state.clone();
tokio::spawn(async move {
let ws = accept_async(stream).await.expect("ws handshake");
let (write, read) = ws.split();
let _ = run_dispatch(read, Box::pin(write), peer.to_string(), state).await;
});
}
});
addr
}
async fn call_raw(
ws: &mut Ws,
id: &str,
method: &str,
params: serde_json::Value,
) -> serde_json::Value {
ws.send(Message::Text(
serde_json::json!({ "jsonrpc": "2.0", "id": id, "method": method, "params": params })
.to_string()
.into(),
))
.await
.expect("send request");
let text = ws
.next()
.await
.expect("response frame")
.expect("response frame ok")
.into_text()
.expect("text response")
.to_string();
serde_json::from_str(&text).expect("parse response")
}
async fn call_ok(
ws: &mut Ws,
id: &str,
method: &str,
params: serde_json::Value,
) -> serde_json::Value {
let resp = call_raw(ws, id, method, params).await;
assert!(
resp.get("error").is_none(),
"{method} should succeed; full response: {resp}"
);
resp["result"].clone()
}
async fn only_session_current_run(state: &Arc<ServerState>) -> Option<String> {
let session = only_session(state).await;
let current = session.current_run_id.lock().await.clone();
current
}
async fn only_session(state: &Arc<ServerState>) -> Arc<car_server_core::session::ClientSession> {
let sessions = state.sessions.lock().await;
sessions
.values()
.next()
.expect("one connected session")
.clone()
}
async fn wait_for_run_termination(
state: &Arc<ServerState>,
run_id: &str,
) -> car_server_core::session::RunMeta {
tokio::time::timeout(
car_server_core::session::RUN_DISCONNECT_CLEANUP_TIMEOUT,
async {
loop {
if let Some(meta) = state.run_meta(run_id).await {
if meta.termination.is_some() {
break meta;
}
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
},
)
.await
.unwrap_or_else(|_| {
panic!("run `{run_id}` did not become terminal within the disconnect cleanup contract")
})
}
fn journal_kind_count(path: &std::path::Path, kind: &str) -> usize {
std::fs::read_to_string(path)
.unwrap()
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["kind"] == kind)
.count()
}
fn success_outcome() -> serde_json::Value {
serde_json::json!({
"status": "success",
"summary": "did the thing",
"evidence": [],
"metrics": {
"turns": 1, "tool_calls": 1, "duration_ms": 12.0,
"retries": 0, "actions_succeeded": 1, "actions_failed": 0
},
"timestamp": chrono::Utc::now().to_rfc3339()
})
}
fn assert_retry_safe_durability_unknown(response: &serde_json::Value) {
let message = response["error"]["message"].as_str().unwrap_or_default();
assert!(
message.contains("critical journal durability is unknown")
&& message.contains("retry the exact event safely"),
"critical lifecycle uncertainty must be explicit and retry-safe: {response}"
);
}
async fn assert_only_intended_start_reservation(
state: &Arc<ServerState>,
intended_run_id: &str,
rejected_fresh_run_id: &str,
) {
assert!(
state.run_lifecycle_state(intended_run_id).await.is_some(),
"the retry-owned start reservation must remain"
);
assert_eq!(
state.run_lifecycle_state(rejected_fresh_run_id).await,
None,
"a fresh key rejected while the current start is unfinished must not leak a reservation"
);
}
#[tokio::test]
async fn runs_start_returns_unique_ids_for_sequential_starts() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let r1 = call_ok(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "first goal" }),
)
.await;
let r2 = call_ok(
&mut ws,
"s2",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "second goal" }),
)
.await;
let id1 = r1["run_id"].as_str().unwrap();
let id2 = r2["run_id"].as_str().unwrap();
assert!(!id1.is_empty() && !id2.is_empty());
assert_ne!(id1, id2, "two sequential starts must mint distinct run_ids");
assert_eq!(r1["agent_id"], "agent-a");
}
#[tokio::test]
async fn idempotency_key_dedups_runs_start() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let r1 = call_ok(
&mut ws,
"k1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "occurrence", "idempotency_key": "occ-42" }),
)
.await;
assert_eq!(r1["run_id"], "occ-42", "the key becomes the run_id");
assert_eq!(r1["agent_id"], "agent-a");
let r2 = call_ok(
&mut ws,
"k2",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "occurrence", "idempotency_key": "occ-42" }),
)
.await;
assert_eq!(r2["run_id"], "occ-42", "same key returns the same run");
assert_eq!(r2["agent_id"], "agent-a");
assert_eq!(
state.run_store.list_runs("agent-a").len(),
1,
"an idempotent re-start must not open a duplicate run"
);
let r3 = call_ok(
&mut ws,
"k3",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "different", "idempotency_key": "occ-99" }),
)
.await;
assert_eq!(r3["run_id"], "occ-99");
assert_eq!(
state.run_store.list_runs("agent-a").len(),
2,
"a fresh key is a fresh run"
);
}
async fn assert_changed_start_payload_is_idempotency_conflict(
changed_field: &str,
first_start_is_durability_unknown: bool,
) {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let journal_failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(journal_failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let run_id = format!(
"{}-{}",
if first_start_is_durability_unknown {
"unknown"
} else {
"committed"
},
changed_field
);
let original = serde_json::json!({
"agent_id": "agent-retry-identity",
"intent": "original intent",
"outcome_description": "original outcome",
"idempotency_key": run_id
});
if first_start_is_durability_unknown {
journal_failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
}
let first = call_raw(&mut ws, "original-start", "runs.start", original.clone()).await;
if first_start_is_durability_unknown {
assert_retry_safe_durability_unknown(&first);
} else {
assert!(
first.get("error").is_none(),
"original committed start failed: {first}"
);
}
let session = only_session(&state).await;
let run_path = state
.run_store
.root()
.join("agent-retry-identity")
.join(format!("{run_id}.jsonl"));
let journal_path = journal_dir.join(format!("{}.jsonl", session.client_id));
let run_bytes_before = std::fs::read(&run_path).unwrap();
let journal_bytes_before = std::fs::read(&journal_path).unwrap();
let mut changed = original.clone();
changed[changed_field] = serde_json::Value::String(format!("changed {changed_field}"));
let rejected = call_raw(&mut ws, "changed-retry", "runs.start", changed).await;
assert_eq!(
rejected["error"]["code"],
car_proto::RUN_OWNERSHIP_CONFLICT_ERROR_CODE,
"a changed occurrence payload must be an explicit idempotency conflict: {rejected}"
);
let message = rejected["error"]["message"].as_str().unwrap_or_default();
assert!(
message.contains("idempotency conflict") && message.contains(changed_field),
"the conflict must name the changed occurrence field: {rejected}"
);
assert_eq!(
std::fs::read(&run_path).unwrap(),
run_bytes_before,
"a conflicting retry must not alter the durable run trace"
);
assert_eq!(
std::fs::read(&journal_path).unwrap(),
journal_bytes_before,
"a conflicting retry must not alter durable journal bytes"
);
let meta = state.run_meta(&run_id).await.unwrap();
assert_eq!(meta.intent, "original intent");
assert_eq!(
meta.outcome_description.as_deref(),
Some("original outcome")
);
let exact = call_ok(&mut ws, "exact-retry", "runs.start", original).await;
assert_eq!(exact["run_id"], run_id);
assert_eq!(
std::fs::read(&run_path).unwrap(),
run_bytes_before,
"the exact retry must reuse the original durable RunStarted row"
);
assert_eq!(
std::fs::read(&journal_path).unwrap(),
journal_bytes_before,
"the exact retry must reconcile the original journal row without duplication"
);
}
#[tokio::test]
async fn committed_start_retry_with_changed_intent_is_idempotency_conflict() {
assert_changed_start_payload_is_idempotency_conflict("intent", false).await;
}
#[tokio::test]
async fn committed_start_retry_with_changed_outcome_description_is_idempotency_conflict() {
assert_changed_start_payload_is_idempotency_conflict("outcome_description", false).await;
}
#[tokio::test]
async fn durability_unknown_start_retry_with_changed_intent_is_idempotency_conflict() {
assert_changed_start_payload_is_idempotency_conflict("intent", true).await;
}
#[tokio::test]
async fn durability_unknown_start_retry_with_changed_outcome_description_is_idempotency_conflict() {
assert_changed_start_payload_is_idempotency_conflict("outcome_description", true).await;
}
#[tokio::test]
async fn stale_idempotency_key_rejection_keeps_current_run_active() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
call_ok(
&mut ws,
"old-start",
"runs.start",
serde_json::json!({
"agent_id": "agent-a",
"intent": "old",
"idempotency_key": "old-terminal"
}),
)
.await;
call_ok(
&mut ws,
"old-complete",
"runs.complete",
serde_json::json!({"run_id": "old-terminal", "outcome": success_outcome()}),
)
.await;
call_ok(
&mut ws,
"live-start",
"runs.start",
serde_json::json!({
"agent_id": "agent-a",
"intent": "live",
"idempotency_key": "still-live"
}),
)
.await;
let rejected = call_raw(
&mut ws,
"stale-start",
"runs.start",
serde_json::json!({
"agent_id": "agent-a",
"intent": "must not replace live",
"idempotency_key": "old-terminal"
}),
)
.await;
assert!(rejected.get("error").is_some(), "stale key must fail");
assert_eq!(
only_session_current_run(&state).await.as_deref(),
Some("still-live")
);
assert!(
state
.run_meta("still-live")
.await
.expect("live run remains registered")
.termination
.is_none(),
"rejecting a stale target must not terminalize the current live run"
);
}
#[tokio::test]
async fn current_run_id_is_set_before_runs_start_responds() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let r = call_ok(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "do it" }),
)
.await;
let run_id = r["run_id"].as_str().unwrap().to_string();
let current = only_session_current_run(&state).await;
assert_eq!(
current.as_deref(),
Some(run_id.as_str()),
"current_run_id must be set before runs.start responds"
);
}
#[tokio::test]
async fn runs_complete_records_terminal_outcome_status() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "do it" }),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let ack = call_ok(
&mut ws,
"c1",
"runs.complete",
serde_json::json!({ "run_id": run_id, "outcome": success_outcome() }),
)
.await;
assert_eq!(ack["ok"], true);
assert_eq!(ack["run_id"], run_id);
let meta = state.run_meta(&run_id).await.expect("run recorded");
match meta.termination {
Some(RunTermination::Outcome { status, .. }) => {
assert_eq!(status, car_ir::OutcomeStatus::Success);
}
other => panic!("expected Outcome termination, got {other:?}"),
}
assert_eq!(only_session_current_run(&state).await, None);
}
#[tokio::test]
async fn completed_then_closed_run_is_not_incomplete() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "do it" }),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
call_ok(
&mut ws,
"c1",
"runs.complete",
serde_json::json!({ "run_id": run_id, "outcome": success_outcome() }),
)
.await;
ws.close(None).await.ok();
drop(ws);
tokio::time::sleep(std::time::Duration::from_millis(600)).await;
let meta = state.run_meta(&run_id).await.expect("run still recorded");
match meta.termination {
Some(RunTermination::Outcome { status, .. }) => {
assert_eq!(
status,
car_ir::OutcomeStatus::Success,
"a completed run must stay Success after close, not be raced to Incomplete"
);
}
other => panic!("expected the run to remain terminal Outcome, got {other:?}"),
}
}
#[tokio::test]
async fn mid_run_drop_yields_incomplete() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-a", "intent": "do it" }),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let client_id = started["client_id"].as_str().unwrap().to_string();
let session = only_session(&state).await;
ws.close(None).await.ok();
drop(ws);
let meta = wait_for_run_termination(&state, &run_id).await;
assert!(
matches!(meta.termination, Some(RunTermination::Incomplete)),
"a mid-run drop must record Incomplete, got {:?}",
meta.termination
);
assert!(meta.pending_terminal.is_none());
let trace = state.run_store.get_run_trace(&run_id).unwrap();
assert!(matches!(
trace.last(),
Some(car_proto::RunRecord::Ended(car_proto::RunEnded {
termination: RunTermination::Incomplete,
..
}))
));
assert_eq!(
journal_kind_count(
&tmp.path()
.join("journals")
.join(format!("{client_id}.jsonl")),
"run_completed"
),
1,
"disconnect must durably journal exactly one terminal event"
);
assert!(
session
.runtime
.event_log_handle()
.lock()
.await
.active_run_binding()
.is_none(),
"terminal acknowledgement must clear the journal binding"
);
assert!(
session.current_run_id.lock().await.is_none(),
"terminal acknowledgement must clear the session's current run"
);
}
#[tokio::test]
async fn disconnect_retries_retry_safe_terminal_uncertainty_before_clearing_state() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let journal_failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(journal_failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"start",
"runs.start",
serde_json::json!({
"agent_id": "agent-disconnect-retry",
"intent": "survive an ambiguous disconnect terminal",
"idempotency_key": "disconnect-terminal-retry"
}),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let client_id = started["client_id"].as_str().unwrap().to_string();
let session = only_session(&state).await;
journal_failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
ws.close(None).await.ok();
drop(ws);
let meta = wait_for_run_termination(&state, &run_id).await;
assert!(matches!(meta.termination, Some(RunTermination::Incomplete)));
assert!(
meta.pending_terminal.is_none(),
"the exact retry must commit the prepared terminal transaction"
);
let trace = state.run_store.get_run_trace(&run_id).unwrap();
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Ended(_)))
.count(),
1,
"the retry must preserve one durable RunEnded preimage"
);
assert_eq!(
journal_kind_count(
&journal_dir.join(format!("{client_id}.jsonl")),
"run_completed"
),
1,
"the exact journal retry must suppress a duplicate terminal row"
);
assert!(
session
.runtime
.event_log_handle()
.lock()
.await
.active_run_binding()
.is_none(),
"cleanup must clear the binding only after terminal acknowledgement"
);
assert!(session.current_run_id.lock().await.is_none());
}
#[tokio::test]
async fn disconnect_preserves_outcome_unknown_execution_quarantine() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"unknown-start",
"runs.start",
serde_json::json!({ "agent_id": "agent-unknown", "intent": "do not guess" }),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let client_id = started["client_id"].as_str().unwrap().to_string();
let original = serde_json::json!({
"id": "unknown-proposal",
"source": "disconnect-crash-test",
"actions": []
});
let typed_original: car_ir::ActionProposal = serde_json::from_value(original.clone()).unwrap();
let normal_original = serde_json::to_value(&typed_original).unwrap();
let canonical = car_inference::catalog_identity::canonical_json(&normal_original).unwrap();
state
.run_store
.write_execution_marker(&car_server_core::run_store::ProposalExecutionMarker {
run_id: run_id.clone(),
client_id,
requested_policy_session_id: None,
policy_session_id: None,
original_proposal_id: "unknown-proposal".to_string(),
original_submission: original.clone(),
original_proposal: normal_original,
proposal_digest: format!("{:x}", Sha256::digest(canonical.as_bytes())),
})
.unwrap();
ws.close(None).await.ok();
drop(ws);
tokio::time::sleep(std::time::Duration::from_millis(600)).await;
let meta = state
.run_meta(&run_id)
.await
.expect("run remains registered");
assert!(
meta.termination.is_none(),
"disconnect must not fabricate Incomplete while an execution outcome is unknown: {:?}",
meta.termination
);
assert!(state.run_store.execution_marker(&run_id).unwrap().is_some());
assert!(
state
.run_store
.get_run_trace(&run_id)
.unwrap()
.iter()
.all(|record| !matches!(record, car_proto::RunRecord::Ended(_))),
"durable trace must remain unterminated while outcome is unknown"
);
}
fn write_truncated_proposal_sidecar(car_root: &std::path::Path, kind: &str, run_id: &str) {
let root = car_root.join(kind);
car_secrets::ensure_private_dir(&root).unwrap();
let key = format!("{:x}", Sha256::digest(run_id.as_bytes()));
let mut file = car_secrets::create_private_file(&root.join(format!("{key}.json"))).unwrap();
file.write_all(b"{").unwrap();
file.sync_all().unwrap();
}
fn write_proposal_sidecar_value(
car_root: &std::path::Path,
kind: &str,
run_id: &str,
value: &impl serde::Serialize,
) {
let root = car_root.join(kind);
car_secrets::ensure_private_dir(&root).unwrap();
let key = format!("{:x}", Sha256::digest(run_id.as_bytes()));
let mut file = car_secrets::create_private_file(&root.join(format!("{key}.json"))).unwrap();
serde_json::to_writer(&mut file, value).unwrap();
file.sync_all().unwrap();
}
async fn assert_corrupt_proposal_state_quarantines_every_live_gate(kind: &str) {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"corrupt-start",
"runs.start",
serde_json::json!({ "agent_id": "agent-corrupt", "intent": "fail closed" }),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
write_truncated_proposal_sidecar(tmp.path(), kind, &run_id);
let submit = call_raw(
&mut ws,
"corrupt-submit",
"proposal.submit",
serde_json::json!({
"proposal": {"id":"must-not-run", "source":"test", "actions":[]}
}),
)
.await;
assert!(
submit["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("durability state is unreadable"),
"corrupt {kind} must block proposal.submit before runtime: {submit}"
);
let complete = call_raw(
&mut ws,
"corrupt-complete",
"runs.complete",
serde_json::json!({"run_id": run_id, "outcome": success_outcome()}),
)
.await;
assert!(
complete["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("durability state is unreadable"),
"corrupt {kind} must block runs.complete: {complete}"
);
let replacement = call_raw(
&mut ws,
"corrupt-replacement",
"runs.start",
serde_json::json!({
"agent_id":"agent-corrupt",
"intent":"must not replace",
"idempotency_key":"corrupt-replacement-run"
}),
)
.await;
assert!(
replacement["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("durability state is unreadable"),
"corrupt {kind} must block replacement: {replacement}"
);
assert!(
state
.run_lifecycle_state("corrupt-replacement-run")
.await
.is_none(),
"failed-closed replacement must not reserve a new run"
);
ws.close(None).await.ok();
drop(ws);
tokio::time::sleep(std::time::Duration::from_millis(600)).await;
let meta = state.run_meta(&run_id).await.unwrap();
assert!(
meta.termination.is_none(),
"corrupt {kind} must block disconnect Incomplete fabrication"
);
}
#[tokio::test]
async fn corrupt_execution_marker_quarantines_every_live_gate() {
assert_corrupt_proposal_state_quarantines_every_live_gate("proposal-execution").await;
}
#[tokio::test]
async fn corrupt_finalization_outbox_quarantines_every_live_gate() {
assert_corrupt_proposal_state_quarantines_every_live_gate("proposal-finalization").await;
}
#[tokio::test]
async fn runs_start_rejects_unresolvable_agent_but_name_fallback_records() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let restore = std::env::var("CAR_AGENT_ID").ok();
std::env::remove_var("CAR_AGENT_ID");
let rejected = call_raw(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "intent": "orphan run" }),
)
.await;
assert!(
rejected.get("error").is_some(),
"runs.start with no resolvable agent_id must be rejected; got {rejected}"
);
let ok = call_ok(
&mut ws,
"s2",
"runs.start",
serde_json::json!({ "agent_name": "Bulldozer Agent", "intent": "one-shot run" }),
)
.await;
let agent_id = ok["agent_id"].as_str().unwrap();
assert_eq!(
agent_id, "name:bulldozer-agent",
"agent_name must synthesize a deterministic name-derived id"
);
assert!(!ok["run_id"].as_str().unwrap().is_empty());
if let Some(v) = restore {
std::env::set_var("CAR_AGENT_ID", v);
}
}
async fn bind_only_session(state: &Arc<ServerState>, agent_id: &str) {
let session = {
let sessions = state.sessions.lock().await;
sessions
.values()
.next()
.expect("one connected session")
.clone()
};
*session.agent_id.lock().await = Some(agent_id.to_string());
}
#[tokio::test]
async fn bound_session_cannot_forge_run_under_other_agent() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started_a = call_ok(
&mut ws,
"s0",
"runs.start",
serde_json::json!({ "agent_id": "agent-A", "intent": "warm up" }),
)
.await;
assert_eq!(started_a["agent_id"], "agent-A");
bind_only_session(&state, "agent-A").await;
let forged = call_raw(
&mut ws,
"s1",
"runs.start",
serde_json::json!({ "agent_id": "agent-B", "intent": "forge under B" }),
)
.await;
assert!(
forged.get("error").is_some(),
"a session bound to A must NOT be able to start a run as B; got {forged}"
);
assert!(
state.run_store.list_runs("agent-B").is_empty(),
"no run may be attributed to agent-B by the A-bound session"
);
let matching = call_ok(
&mut ws,
"s2",
"runs.start",
serde_json::json!({ "agent_id": "agent-A", "intent": "legit, explicit" }),
)
.await;
assert_eq!(matching["agent_id"], "agent-A");
let derived = call_ok(
&mut ws,
"s3",
"runs.start",
serde_json::json!({ "intent": "legit, derived" }),
)
.await;
assert_eq!(
derived["agent_id"], "agent-A",
"an unspecified agent_id on a bound session derives the bound agent"
);
}
#[tokio::test]
async fn concurrent_sessions_reserve_one_global_idempotency_owner() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_multi_dispatcher(state.clone()).await;
let (mut left, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let (mut right, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let params = serde_json::json!({
"agent_id": "agent-race",
"intent": "one global occurrence",
"idempotency_key": "global-occurrence"
});
let (left_response, right_response) = tokio::join!(
call_raw(&mut left, "left", "runs.start", params.clone()),
call_raw(&mut right, "right", "runs.start", params),
);
let responses = [&left_response, &right_response];
assert_eq!(
responses
.iter()
.filter(|response| response.get("result").is_some())
.count(),
1,
"exactly one client owns the absent key: {responses:?}"
);
let conflict = responses
.iter()
.find(|response| response.get("error").is_some())
.unwrap();
assert_eq!(
conflict["error"]["code"],
car_proto::RUN_OWNERSHIP_CONFLICT_ERROR_CODE
);
assert!(conflict["error"]["message"]
.as_str()
.unwrap()
.starts_with(car_proto::RUN_OWNERSHIP_CONFLICT_MESSAGE_PREFIX));
let trace = state.run_store.get_run_trace("global-occurrence").unwrap();
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Started(_)))
.count(),
1
);
let winning_client = responses
.iter()
.find_map(|response| response["result"]["client_id"].as_str())
.unwrap();
match trace.first().unwrap() {
car_proto::RunRecord::Started(started) => {
assert_eq!(started.client_id.as_deref(), Some(winning_client))
}
other => panic!("expected RunStarted, got {other:?}"),
}
}
#[tokio::test]
async fn start_fsync_failure_is_not_acked_and_exact_retry_finishes_once() {
let tmp = TempDir::new().unwrap();
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let params = serde_json::json!({
"agent_id": "agent-durable",
"intent": "retry start",
"idempotency_key": "durable-start"
});
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Fsync);
let failed = call_raw(&mut ws, "start-fail", "runs.start", params.clone()).await;
assert!(
failed.get("error").is_some(),
"fsync failure was acknowledged: {failed}"
);
assert_eq!(
only_session_current_run(&state).await.as_deref(),
Some("durable-start"),
"binding remains for exact retry"
);
let fresh = call_raw(
&mut ws,
"start-fresh-while-pending",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "must not reserve behind unfinished start",
"idempotency_key": "fresh-after-store-fsync"
}),
)
.await;
assert!(
fresh.get("error").is_some(),
"unfinished start was replaced: {fresh}"
);
assert_only_intended_start_reservation(&state, "durable-start", "fresh-after-store-fsync")
.await;
let retried = call_ok(&mut ws, "start-retry", "runs.start", params).await;
assert_eq!(retried["run_id"], "durable-start");
let trace = state.run_store.get_run_trace("durable-start").unwrap();
assert_eq!(trace.len(), 1, "retry must not duplicate RunStarted");
}
#[tokio::test]
async fn start_write_rejection_releases_reservation_for_same_or_fresh_retry() {
for exact_retry in [true, false] {
let tmp = TempDir::new().unwrap();
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let rejected_params = serde_json::json!({
"agent_id": "agent-rejected",
"intent": "retry rejected start",
"idempotency_key": "rejected-start"
});
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Write);
let failed = call_raw(
&mut ws,
"start-rejected",
"runs.start",
rejected_params.clone(),
)
.await;
assert!(
failed.get("error").is_some(),
"RunStore rejection was acknowledged: {failed}"
);
assert_eq!(
state.run_lifecycle_state("rejected-start").await,
None,
"a preaccept rejection must release its reservation"
);
assert_eq!(only_session_current_run(&state).await, None);
let retry_params = if exact_retry {
rejected_params
} else {
serde_json::json!({
"agent_id": "agent-rejected",
"intent": "fresh start after proven rejection",
"idempotency_key": "fresh-after-store-rejection"
})
};
let expected_run_id = if exact_retry {
"rejected-start"
} else {
"fresh-after-store-rejection"
};
let retried = call_ok(&mut ws, "start-retry", "runs.start", retry_params).await;
assert_eq!(retried["run_id"], expected_run_id);
assert!(state.run_lifecycle_state(expected_run_id).await.is_some());
if !exact_retry {
assert_eq!(state.run_lifecycle_state("rejected-start").await, None);
}
let trace = state.run_store.get_run_trace(expected_run_id).unwrap();
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Started(_)))
.count(),
1,
"the accepted retry must create one occurrence"
);
}
}
#[tokio::test]
async fn journal_binding_rejection_releases_same_and_fresh_candidate_reservations() {
let tmp = TempDir::new().unwrap();
let state = loopback_state(tmp.path().join("journals"));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let session = only_session(&state).await;
session
.bind_run_journal("blocking-binding")
.await
.expect("install deterministic conflicting journal binding");
for run_id in ["journal-rejected-start", "fresh-after-journal-rejection"] {
let rejected = call_raw(
&mut ws,
run_id,
"runs.start",
serde_json::json!({
"agent_id": "agent-journal-rejected",
"intent": format!("rejected before journal admission: {run_id}"),
"idempotency_key": run_id
}),
)
.await;
assert!(
rejected["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("journal is already bound"),
"journal preaccept rejection must be explicit: {rejected}"
);
assert_eq!(
state.run_lifecycle_state(run_id).await,
None,
"journal rejection must release candidate `{run_id}`"
);
}
assert_eq!(only_session_current_run(&state).await, None);
session
.clear_run_journal_binding("blocking-binding")
.await
.expect("clear deterministic blocker");
let retried = call_ok(
&mut ws,
"journal-rejected-exact-retry",
"runs.start",
serde_json::json!({
"agent_id": "agent-journal-rejected",
"intent": "rejected before journal admission: journal-rejected-start",
"idempotency_key": "journal-rejected-start"
}),
)
.await;
assert_eq!(retried["run_id"], "journal-rejected-start");
assert_eq!(
state
.run_store
.get_run_trace("journal-rejected-start")
.unwrap()
.len(),
1
);
assert_eq!(
state
.run_lifecycle_state("fresh-after-journal-rejection")
.await,
None
);
}
#[tokio::test]
async fn start_first_use_parent_sync_failure_is_not_acked_and_exact_retry_finishes_once() {
let tmp = TempDir::new().unwrap();
let failures = car_secrets::PrivatePathDurabilityFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_private_path_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let params = serde_json::json!({
"agent_id": "agent-private-first-use",
"intent": "retry first-use directory durability",
"idempotency_key": "private-first-use-start"
});
failures.fail_next(car_secrets::PrivatePathDurabilityFailurePoint::ParentDirectorySync);
let failed = call_raw(&mut ws, "start-fail", "runs.start", params.clone()).await;
assert!(
failed.get("error").is_some(),
"parent directory sync failure was acknowledged: {failed}"
);
assert!(
state
.run_store
.get_run_trace("private-first-use-start")
.is_none(),
"failed first-use storage must not expose a RunStarted receipt"
);
let retried = call_ok(&mut ws, "start-retry", "runs.start", params).await;
assert_eq!(retried["run_id"], "private-first-use-start");
let trace = state
.run_store
.get_run_trace("private-first-use-start")
.unwrap();
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Started(_)))
.count(),
1,
"retry must persist exactly one RunStarted"
);
}
#[tokio::test]
async fn start_journal_fsync_failure_is_not_acked_and_exact_retry_finishes_once() {
let tmp = TempDir::new().unwrap();
let failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let params = serde_json::json!({
"agent_id": "agent-durable",
"intent": "retry journal start",
"idempotency_key": "durable-journal-start"
});
failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
let failed = call_raw(&mut ws, "start-fail", "runs.start", params.clone()).await;
assert!(
failed.get("error").is_some(),
"journal fsync failure was acknowledged: {failed}"
);
assert_retry_safe_durability_unknown(&failed);
let fresh = call_raw(
&mut ws,
"start-fresh-while-journal-unknown",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "must not reserve behind unknown journal durability",
"idempotency_key": "fresh-after-journal-unknown"
}),
)
.await;
assert!(
fresh.get("error").is_some(),
"unknown start was replaced: {fresh}"
);
assert_only_intended_start_reservation(
&state,
"durable-journal-start",
"fresh-after-journal-unknown",
)
.await;
let retried = call_ok(&mut ws, "start-retry", "runs.start", params).await;
let journal = std::fs::read_to_string(
tmp.path()
.join("journals")
.join(format!("{}.jsonl", retried["client_id"].as_str().unwrap())),
)
.unwrap();
assert_eq!(
journal
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["kind"] == "run_started")
.count(),
1
);
assert_eq!(
state
.run_store
.get_run_trace("durable-journal-start")
.unwrap()
.len(),
1
);
}
#[tokio::test]
async fn fresh_start_reconciles_runstore_rejected_pending_terminal_before_reserving() {
let tmp = TempDir::new().unwrap();
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
call_ok(
&mut ws,
"start-current",
"runs.start",
serde_json::json!({
"agent_id": "agent-terminal-retry",
"intent": "finish exact pending terminal",
"idempotency_key": "pending-store-terminal"
}),
)
.await;
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Write);
let failed = call_raw(
&mut ws,
"complete-rejected",
"runs.complete",
serde_json::json!({
"run_id": "pending-store-terminal",
"outcome": success_outcome()
}),
)
.await;
assert!(
failed.get("error").is_some(),
"terminal rejection was acked: {failed}"
);
let fresh = call_ok(
&mut ws,
"start-after-terminal-rejection",
"runs.start",
serde_json::json!({
"agent_id": "agent-terminal-retry",
"intent": "start only after exact pending terminal",
"idempotency_key": "after-store-terminal"
}),
)
.await;
assert_eq!(fresh["run_id"], "after-store-terminal");
assert_eq!(
state
.run_lifecycle_state("pending-store-terminal")
.await
.map(|state| (state.1, state.3)),
Some((true, false)),
"the old exact terminal must commit before the fresh reservation becomes active"
);
assert_eq!(
only_session_current_run(&state).await.as_deref(),
Some("after-store-terminal")
);
}
#[tokio::test]
async fn fresh_start_reconciles_durability_unknown_terminal_before_reserving() {
let tmp = TempDir::new().unwrap();
let failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
call_ok(
&mut ws,
"start-current",
"runs.start",
serde_json::json!({
"agent_id": "agent-terminal-unknown",
"intent": "finish unknown terminal",
"idempotency_key": "pending-journal-terminal"
}),
)
.await;
failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
let failed = call_raw(
&mut ws,
"complete-unknown",
"runs.complete",
serde_json::json!({
"run_id": "pending-journal-terminal",
"outcome": success_outcome()
}),
)
.await;
assert_retry_safe_durability_unknown(&failed);
let fresh = call_ok(
&mut ws,
"start-after-terminal-unknown",
"runs.start",
serde_json::json!({
"agent_id": "agent-terminal-unknown",
"intent": "start only after exact unknown terminal",
"idempotency_key": "after-journal-terminal"
}),
)
.await;
assert_eq!(fresh["run_id"], "after-journal-terminal");
assert_eq!(
state
.run_lifecycle_state("pending-journal-terminal")
.await
.map(|state| (state.1, state.3)),
Some((true, false)),
"the unknown old terminal must resolve exactly before the fresh run starts"
);
assert_eq!(
only_session_current_run(&state).await.as_deref(),
Some("after-journal-terminal")
);
}
#[tokio::test]
async fn start_write_failure_disconnect_releases_absent_occurrence_for_reconnect() {
let tmp = TempDir::new().unwrap();
let journal_root = tmp.path().join("journals");
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_root.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(failures.clone()),
));
let addr = spawn_multi_dispatcher(state.clone()).await;
let (mut owner, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let params = serde_json::json!({
"agent_id": "agent-undurable",
"intent": "release an occurrence that never reached disk",
"idempotency_key": "undurable-start"
});
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Write);
let failed = call_raw(&mut owner, "start-fail", "runs.start", params.clone()).await;
assert!(
failed.get("error").is_some(),
"write failure was acknowledged: {failed}"
);
let owner_client_id = state
.sessions
.lock()
.await
.keys()
.next()
.expect("owning connection is registered")
.clone();
owner.close(None).await.ok();
tokio::time::sleep(std::time::Duration::from_millis(600)).await;
assert_eq!(
state.run_lifecycle_state("undurable-start").await,
None,
"disconnect must release a reservation with no durable start surface"
);
assert!(
state.run_store.get_run_trace("undurable-start").is_none(),
"an empty file from the failed write must not strand the global key"
);
assert!(
!journal_root
.join(format!("{owner_client_id}.jsonl"))
.exists(),
"RunStore rejection precedes journal admission, so no old-client start surface exists"
);
let (mut replacement, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(&mut replacement, "start-reconnect", "runs.start", params).await;
assert_eq!(started["run_id"], "undurable-start");
}
async fn assert_durable_failed_start_disconnects_to_historical_terminal(
state: Arc<ServerState>,
addr: SocketAddr,
journal_root: &std::path::Path,
run_id: &str,
) {
let lifecycle = tokio::time::timeout(std::time::Duration::from_secs(3), async {
loop {
if let Some(lifecycle) = state.run_lifecycle_state(run_id).await {
if lifecycle.1 {
break lifecycle;
}
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.expect("durable failed start must terminalize within three seconds");
assert!(
lifecycle.1,
"durable failed start must be terminal after disconnect"
);
assert!(
lifecycle.2,
"historical terminal requires a committed RunStarted"
);
assert!(
!lifecycle.3,
"historical terminal must not retain a pending transaction"
);
let trace = state.run_store.get_run_trace(run_id).unwrap();
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Started(_)))
.count(),
1,
"disconnect reconciliation must not duplicate RunStarted"
);
assert_eq!(
trace
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Ended(_)))
.count(),
1,
"disconnect reconciliation must leave one orphan-safe terminal"
);
assert!(matches!(
trace.last(),
Some(car_proto::RunRecord::Ended(car_proto::RunEnded {
termination: RunTermination::Incomplete,
..
}))
));
let client_id = match trace.first().unwrap() {
car_proto::RunRecord::Started(started) => started.client_id.as_deref().unwrap(),
other => panic!("expected RunStarted, got {other:?}"),
};
let journal = std::fs::read_to_string(journal_root.join(format!("{client_id}.jsonl")))
.expect("reconciled client journal");
let lifecycle_kinds: Vec<String> = journal
.lines()
.map(|line| serde_json::from_str::<serde_json::Value>(line).unwrap())
.filter_map(|row| {
row["kind"]
.as_str()
.filter(|kind| *kind == "run_started" || *kind == "run_completed")
.map(str::to_string)
})
.collect();
assert_eq!(
lifecycle_kinds,
["run_started", "run_completed"],
"reconciliation must preserve one ordered journal bracket"
);
let (mut replacement, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let conflict = call_raw(
&mut replacement,
"owned-key",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "cannot adopt a historical occurrence",
"idempotency_key": run_id
}),
)
.await;
assert_eq!(
conflict["error"]["code"],
car_proto::RUN_OWNERSHIP_CONFLICT_ERROR_CODE,
"a new socket cannot adopt the original durable idempotency key"
);
let fresh = call_ok(
&mut replacement,
"fresh-key",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "liveness after historical reconciliation",
"idempotency_key": format!("{run_id}-fresh")
}),
)
.await;
assert_eq!(fresh["run_id"], format!("{run_id}-fresh"));
}
#[tokio::test]
async fn start_store_fsync_failure_disconnects_to_one_historical_occurrence() {
let tmp = TempDir::new().unwrap();
let journal_root = tmp.path().join("journals");
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_root.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(failures.clone()),
));
let addr = spawn_multi_dispatcher(state.clone()).await;
let (mut owner, _) = connect_async(format!("ws://{addr}")).await.unwrap();
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Fsync);
let failed = call_raw(
&mut owner,
"start-fail",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "recover a written RunStarted",
"idempotency_key": "store-fsync-start"
}),
)
.await;
assert!(
failed.get("error").is_some(),
"fsync failure was acknowledged: {failed}"
);
owner.close(None).await.ok();
drop(owner);
assert_durable_failed_start_disconnects_to_historical_terminal(
state,
addr,
&journal_root,
"store-fsync-start",
)
.await;
}
#[tokio::test]
async fn start_journal_fsync_failure_disconnects_to_one_historical_occurrence() {
let tmp = TempDir::new().unwrap();
let journal_root = tmp.path().join("journals");
let failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_root.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(failures.clone()),
));
let addr = spawn_multi_dispatcher(state.clone()).await;
let (mut owner, _) = connect_async(format!("ws://{addr}")).await.unwrap();
failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
let failed = call_raw(
&mut owner,
"start-fail",
"runs.start",
serde_json::json!({
"agent_id": "agent-durable",
"intent": "recover a written journal boundary",
"idempotency_key": "journal-fsync-start"
}),
)
.await;
assert!(
failed.get("error").is_some(),
"journal fsync failure was acknowledged: {failed}"
);
owner.close(None).await.ok();
drop(owner);
assert_durable_failed_start_disconnects_to_historical_terminal(
state,
addr,
&journal_root,
"journal-fsync-start",
)
.await;
}
#[tokio::test]
async fn terminal_journal_fsync_failure_retries_same_digest_once() {
let tmp = TempDir::new().unwrap();
let journal_failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(tmp.path().join("journals"))
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(journal_failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"start",
"runs.start",
serde_json::json!({"agent_id":"agent-durable","intent":"retry complete"}),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let outcome = success_outcome();
journal_failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
let failed = call_raw(
&mut ws,
"complete-fail",
"runs.complete",
serde_json::json!({"run_id":run_id,"outcome":outcome.clone()}),
)
.await;
assert!(
failed.get("error").is_some(),
"journal fsync failure was acknowledged: {failed}"
);
assert_retry_safe_durability_unknown(&failed);
assert_eq!(
only_session_current_run(&state).await.as_deref(),
Some(run_id.as_str())
);
let run_path = state
.run_store
.root()
.join("agent-durable")
.join(format!("{run_id}.jsonl"));
let journal_path = tmp
.path()
.join("journals")
.join(format!("{}.jsonl", started["client_id"].as_str().unwrap()));
let run_bytes_before_conflict = std::fs::read(&run_path).unwrap();
let journal_bytes_before_conflict = std::fs::read(&journal_path).unwrap();
let mut changed_outcome = outcome.clone();
changed_outcome["summary"] = serde_json::Value::String("different outcome".to_string());
let conflict = call_raw(
&mut ws,
"complete-conflict",
"runs.complete",
serde_json::json!({"run_id":run_id,"outcome":changed_outcome}),
)
.await;
assert!(
conflict["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("different pending terminal transaction"),
"a non-exact terminal retry must be rejected: {conflict}"
);
assert_eq!(
std::fs::read(&run_path).unwrap(),
run_bytes_before_conflict,
"a conflicting terminal retry must not change durable run bytes"
);
assert_eq!(
std::fs::read(&journal_path).unwrap(),
journal_bytes_before_conflict,
"a conflicting terminal retry must not change durable journal bytes"
);
let completed = call_ok(
&mut ws,
"complete-retry",
"runs.complete",
serde_json::json!({"run_id":run_id,"outcome":outcome}),
)
.await;
let digest = completed["completion_digest"].as_str().unwrap();
let trace = state.run_store.get_run_trace(&run_id).unwrap();
let terminals: Vec<_> = trace
.iter()
.filter_map(|row| match row {
car_proto::RunRecord::Ended(ended) => Some(ended),
_ => None,
})
.collect();
assert_eq!(terminals.len(), 1);
assert_eq!(terminals[0].completion_digest.as_deref(), Some(digest));
let journal = std::fs::read_to_string(journal_path).unwrap();
assert_eq!(
journal
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["kind"] == "run_completed")
.count(),
1
);
}
#[tokio::test]
async fn rejected_terminal_retry_with_changed_outcome_preserves_exact_preimage() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let run_store_failures = car_server_core::run_store::RunStoreFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_run_store_failures(run_store_failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let started = call_ok(
&mut ws,
"start",
"runs.start",
serde_json::json!({
"agent_id": "agent-rejected-terminal",
"intent": "retain one terminal preimage",
"idempotency_key": "rejected-terminal"
}),
)
.await;
let outcome = success_outcome();
run_store_failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::Write);
let failed = call_raw(
&mut ws,
"complete-fail",
"runs.complete",
serde_json::json!({"run_id":"rejected-terminal","outcome":outcome.clone()}),
)
.await;
assert!(
failed["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("RunEnded durability failed"),
"the rejected terminal must remain unacknowledged: {failed}"
);
let run_path = state
.run_store
.root()
.join("agent-rejected-terminal")
.join("rejected-terminal.jsonl");
let journal_path =
journal_dir.join(format!("{}.jsonl", started["client_id"].as_str().unwrap()));
let run_bytes_before_conflict = std::fs::read(&run_path).unwrap();
let journal_bytes_before_conflict = std::fs::read(&journal_path).unwrap();
let mut changed_outcome = outcome.clone();
changed_outcome["summary"] = serde_json::Value::String("different outcome".to_string());
let conflict = call_raw(
&mut ws,
"complete-conflict",
"runs.complete",
serde_json::json!({"run_id":"rejected-terminal","outcome":changed_outcome}),
)
.await;
assert!(
conflict["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("different pending terminal transaction"),
"a rejected terminal still owns its exact retry identity: {conflict}"
);
assert_eq!(std::fs::read(&run_path).unwrap(), run_bytes_before_conflict);
assert_eq!(
std::fs::read(&journal_path).unwrap(),
journal_bytes_before_conflict
);
let completed = call_ok(
&mut ws,
"complete-exact",
"runs.complete",
serde_json::json!({"run_id":"rejected-terminal","outcome":outcome}),
)
.await;
assert_eq!(completed["run_id"], "rejected-terminal");
assert_eq!(
state
.run_store
.get_run_trace("rejected-terminal")
.unwrap()
.iter()
.filter(|row| matches!(row, car_proto::RunRecord::Ended(_)))
.count(),
1,
"the exact retry must finish one terminal transaction"
);
assert_eq!(journal_kind_count(&journal_path, "run_completed"), 1);
}
#[tokio::test]
async fn proposal_terminal_durability_unknown_preserves_retry_state_until_exact_retry() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let journal_failures = car_eventlog::JournalFailureInjector::default();
let state = Arc::new(ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(journal_failures.clone()),
));
let addr = spawn_dispatcher(state.clone()).await;
let (mut ws, _) = connect_async(format!("ws://{addr}")).await.unwrap();
let policy_session_id = call_ok(
&mut ws,
"proposal-policy-open",
"session.policy.open",
serde_json::json!({}),
)
.await["session_id"]
.as_str()
.unwrap()
.to_string();
let started = call_ok(
&mut ws,
"proposal-start",
"runs.start",
serde_json::json!({
"agent_id": "agent-proposal-durability",
"intent": "retry one proposal terminal"
}),
)
.await;
let run_id = started["run_id"].as_str().unwrap().to_string();
let client_id = started["client_id"].as_str().unwrap().to_string();
let proposal = serde_json::json!({
"id": "proposal-durability-unknown",
"source": "run-lifecycle-test",
"actions": []
});
let params = serde_json::json!({
"session_id": policy_session_id,
"proposal": proposal
});
journal_failures.fail_next(car_eventlog::JournalFailurePoint::Fsync);
let failed = call_raw(
&mut ws,
"proposal-terminal-fail",
"proposal.submit",
params.clone(),
)
.await;
let failure_message = failed["error"]["message"].as_str().unwrap_or_default();
assert!(
failure_message.contains("proposal finalization durability is unknown")
&& failure_message.contains("retry the exact proposal.submit safely"),
"proposal uncertainty must name the retry-safe operation: {failed}"
);
let pending = state
.run_store
.pending_proposal(&run_id)
.unwrap()
.expect("durability-unknown terminal keeps its finalization outbox");
assert_eq!(
pending.policy_session_id.as_deref(),
Some(policy_session_id.as_str())
);
assert!(
state.run_store.execution_marker(&run_id).unwrap().is_some(),
"durability-unknown terminal keeps its execution guard"
);
assert!(
state
.run_store
.completed_proposal(
&run_id,
&client_id,
Some(&policy_session_id),
¶ms["proposal"],
)
.unwrap()
.is_none(),
"a completion receipt must not precede journal acknowledgement"
);
let session = state
.sessions
.lock()
.await
.values()
.next()
.expect("one connected session")
.clone();
let event_log = session.runtime.event_log_handle();
let binding = event_log.lock().await.active_run_binding().map(
|(bound_run, bound_client, bound_policy)| {
(
bound_run.to_string(),
bound_client.to_string(),
bound_policy.map(str::to_string),
)
},
);
assert_eq!(
binding,
Some((
run_id.clone(),
client_id.clone(),
Some(policy_session_id.clone())
)),
"durability-unknown terminal must retain the run and policy binding"
);
let retried = call_ok(
&mut ws,
"proposal-terminal-retry",
"proposal.submit",
params.clone(),
)
.await;
assert_eq!(retried["proposal_id"], "proposal-durability-unknown");
assert!(state.run_store.pending_proposal(&run_id).unwrap().is_none());
assert!(state.run_store.execution_marker(&run_id).unwrap().is_none());
assert!(
state
.run_store
.completed_proposal(
&run_id,
&client_id,
Some(&policy_session_id),
¶ms["proposal"],
)
.unwrap()
.is_some(),
"acknowledged exact retry writes the completion receipt"
);
assert_eq!(
event_log
.lock()
.await
.active_run_binding()
.and_then(|(_, _, policy)| policy),
None,
"acknowledged exact retry clears the policy binding"
);
let journal = std::fs::read_to_string(journal_dir.join(format!("{client_id}.jsonl"))).unwrap();
assert_eq!(
journal
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["kind"] == "proposal_completed")
.count(),
1,
"exact retry must not duplicate the proposal terminal"
);
}
#[tokio::test]
async fn restart_reconciles_missing_critical_journal_from_durable_run_preimages() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let started = car_proto::RunStarted {
run_id: "restart-run".to_string(),
client_id: Some("restart-client".to_string()),
agent_id: "agent-restart".to_string(),
intent: "recover outbox".to_string(),
outcome_description: Some("journal catches durable run".to_string()),
started_at: chrono::Utc::now(),
};
store.write_started(&started).unwrap();
let termination = car_proto::RunTermination::Incomplete;
let canonical = car_inference::catalog_identity::canonical_json(&termination).unwrap();
let digest = format!("{:x}", Sha256::digest(canonical.as_bytes()));
let ended = car_proto::RunEnded {
run_id: started.run_id.clone(),
client_id: started.client_id.clone(),
agent_id: started.agent_id.clone(),
termination,
completion_digest: Some(digest.clone()),
ended_at: chrono::Utc::now(),
};
store.write_ended(&ended).unwrap();
assert!(
!journal_dir.exists(),
"precondition: critical journal is missing"
);
let _restarted = ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None)))),
);
let rows: Vec<serde_json::Value> =
std::fs::read_to_string(journal_dir.join("restart-client.jsonl"))
.unwrap()
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
let kinds: Vec<_> = rows
.iter()
.map(|row| row["kind"].as_str().unwrap())
.collect();
assert_eq!(kinds, vec!["run_started", "run_completed"]);
for row in &rows {
assert_eq!(row["run_id"], "restart-run");
assert_eq!(row["client_id"], "restart-client");
}
assert_eq!(rows[1]["data"]["completion_digest"], digest);
assert_eq!(
rows[1]["data"]["termination"],
serde_json::json!({"kind":"incomplete"})
);
}
#[test]
fn startup_reconciliation_timeout_fails_before_start_and_preserves_quarantine() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let (started, marker, pending) = proposal_finalization_fixture(
"startup-stalled-run",
"startup-stalled-client",
Some("startup-stalled-policy"),
Some("startup-stalled-policy"),
);
store.write_started(&started).unwrap();
store.write_execution_marker(&marker).unwrap();
store.write_pending_proposal(&pending).unwrap();
let failures = car_eventlog::JournalFailureInjector::default();
failures.fail_next(car_eventlog::JournalFailurePoint::HoldAcknowledgement);
let acknowledgement_timeout = std::time::Duration::from_millis(20);
let config = ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None))))
.with_journal_failures(failures.clone())
.with_startup_reconciliation_acknowledgement_timeout(acknowledgement_timeout);
let (result_tx, result_rx) = std::sync::mpsc::sync_channel(1);
let started_wait = std::time::Instant::now();
let startup = std::thread::spawn(move || {
let result = ServerState::try_with_config(config).map(|_| ());
result_tx.send(result).unwrap();
});
let result = result_rx
.recv_timeout(std::time::Duration::from_secs(5))
.expect("the 20ms startup acknowledgement bound must return without journal release");
let elapsed = started_wait.elapsed();
startup.join().unwrap();
assert!(
elapsed >= acknowledgement_timeout,
"startup returned before its configured acknowledgement bound: {elapsed:?}"
);
assert_eq!(
result.unwrap_err(),
"startup reconciliation failed for run startup-stalled-run run_started: critical journal durability is unknown; retry the exact event safely: journal writer did not acknowledge within 20ms"
);
assert_eq!(
failures.held_acknowledgement_count(),
1,
"the permitted stalled-journal double must still own the unacknowledged receipt"
);
assert!(
store.pending_proposal(&started.run_id).unwrap().is_some(),
"startup timeout must preserve the durable finalization outbox"
);
assert!(
store.execution_marker(&started.run_id).unwrap().is_some(),
"startup timeout must preserve the outcome-unknown execution quarantine"
);
assert!(
store
.completed_proposal(
&started.run_id,
started.client_id.as_deref().unwrap(),
Some("startup-stalled-policy"),
&pending.original_submission,
)
.unwrap()
.is_none(),
"a timed-out acknowledgement must never claim a completed response"
);
let _retried = ServerState::try_with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None)))),
)
.expect("a fresh startup must reconcile the exact durable rows");
assert!(store.pending_proposal(&started.run_id).unwrap().is_none());
assert!(store.execution_marker(&started.run_id).unwrap().is_none());
}
#[test]
fn completed_receipt_migration_checkpoint_prevents_permanent_rescan() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let failures = car_server_core::run_store::RunStoreFailureInjector::default();
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir)
.with_failure_injector(failures.clone());
let (started, marker, pending) = proposal_finalization_fixture(
"migration-checkpoint-run",
"migration-checkpoint-client",
None,
None,
);
store.write_started(&started).unwrap();
store.write_execution_marker(&marker).unwrap();
store.write_pending_proposal(&pending).unwrap();
let _receipt = store.write_completed_proposal(&pending).unwrap();
failures.fail_next(car_server_core::run_store::RunStoreFailurePoint::MarkerUnlink);
assert!(store.reconcile_completed_proposal_migration().is_err());
assert!(store.execution_marker(&started.run_id).unwrap().is_some());
assert!(
!tmp.path()
.join("proposal-completed-index/owner-index-migration.json")
.exists(),
"checkpoint must be last so failed cleanup remains retryable"
);
store
.reconcile_completed_proposal_migration()
.expect("the first migration must validate, index, and clean legacy receipts");
assert!(store.execution_marker(&started.run_id).unwrap().is_none());
assert!(store.pending_proposal(&started.run_id).unwrap().is_none());
let unrelated = tmp.path().join("proposal-completed/unrelated");
std::fs::create_dir_all(&unrelated).unwrap();
std::fs::write(unrelated.join("corrupt.json"), b"not-json").unwrap();
store
.reconcile_completed_proposal_migration()
.expect("a completed migration must not enumerate permanent receipts again");
}
fn proposal_finalization_fixture(
run_id: &str,
client_id: &str,
requested_policy_session_id: Option<&str>,
policy_session_id: Option<&str>,
) -> (
car_proto::RunStarted,
car_server_core::run_store::ProposalExecutionMarker,
car_server_core::run_store::PendingProposalFinalization,
) {
let started = car_proto::RunStarted {
run_id: run_id.to_string(),
client_id: Some(client_id.to_string()),
agent_id: "agent-provenance".to_string(),
intent: "bind pending finalization provenance".to_string(),
outcome_description: None,
started_at: chrono::Utc::now(),
};
let raw = serde_json::json!({
"id": format!("proposal-{run_id}"),
"source": "provenance-test",
"actions": [],
"caller_extension": {"kept": true}
});
let proposal: car_ir::ActionProposal = serde_json::from_value(raw.clone()).unwrap();
let proposal_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&proposal)
.unwrap()
.as_bytes()
)
);
let proposal_result = car_ir::ProposalResult {
proposal_id: proposal.id.clone(),
original_proposal_id: proposal.id.clone(),
final_proposal: Some(proposal.clone()),
replan_lineage: vec![car_ir::ProposalLineageEntry {
generation: 0,
proposal_id: proposal.id.clone(),
proposal_digest: Some(proposal_digest.clone()),
status: car_ir::ProposalLineageStatus::Accepted,
rejection_reason: None,
}],
accepted_proposal_preimages: vec![car_ir::AcceptedProposalPreimage {
generation: 0,
proposal_digest: proposal_digest.clone(),
proposal: proposal.clone(),
}],
results: vec![],
cost: car_ir::CostSummary::default(),
};
let result_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&proposal_result)
.unwrap()
.as_bytes()
)
);
let pending = car_server_core::run_store::PendingProposalFinalization {
run_id: started.run_id.clone(),
client_id: client_id.to_string(),
requested_policy_session_id: requested_policy_session_id.map(str::to_string),
policy_session_id: policy_session_id.map(str::to_string),
original_proposal_id: proposal.id.clone(),
final_proposal_id: proposal.id.clone(),
original_submission: raw.clone(),
original_proposal: proposal.clone(),
final_proposal: proposal.clone(),
accepted_proposal_preimages: vec![car_server_core::run_store::AcceptedProposalPreimage {
generation: 0,
proposal: proposal.clone(),
}],
proposal_result,
result_digest,
};
let marker = car_server_core::run_store::ProposalExecutionMarker {
run_id: started.run_id.clone(),
client_id: client_id.to_string(),
requested_policy_session_id: requested_policy_session_id.map(str::to_string),
policy_session_id: policy_session_id.map(str::to_string),
original_proposal_id: proposal.id.clone(),
original_submission: raw,
original_proposal: serde_json::to_value(proposal).unwrap(),
proposal_digest,
};
(started, marker, pending)
}
fn mutate_pending_original_proposal(
pending: &mut car_server_core::run_store::PendingProposalFinalization,
) {
pending.original_submission["source"] = serde_json::json!("pending-self-claim");
pending.original_proposal.source = "pending-self-claim".to_string();
pending.final_proposal = pending.original_proposal.clone();
pending.accepted_proposal_preimages[0].proposal = pending.original_proposal.clone();
pending.proposal_result.final_proposal = Some(pending.original_proposal.clone());
let proposal_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&pending.original_proposal)
.unwrap()
.as_bytes()
)
);
pending.proposal_result.replan_lineage[0].proposal_digest = Some(proposal_digest);
pending.result_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&pending.proposal_result)
.unwrap()
.as_bytes()
)
);
}
#[tokio::test]
async fn restart_replays_exact_pending_proposal_result_without_a_runtime_dispatch() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let started = car_proto::RunStarted {
run_id: "restart-proposal-run".to_string(),
client_id: Some("restart-proposal-client".to_string()),
agent_id: "agent-restart-proposal".to_string(),
intent: "recover proposal result".to_string(),
outcome_description: None,
started_at: chrono::Utc::now(),
};
store.write_started(&started).unwrap();
let proposal_timestamp = chrono::Utc::now();
let original_proposal: car_ir::ActionProposal = serde_json::from_value(serde_json::json!({
"id": "proposal-original",
"source": "restart-test",
"actions": [{
"id": "original-action",
"type": "tool_call",
"tool": "drive_cli",
"parameters": {"prompt": "original parameters"}
}],
"timestamp": proposal_timestamp,
"context": {}
}))
.unwrap();
let final_proposal: car_ir::ActionProposal = serde_json::from_value(serde_json::json!({
"id": "proposal-final",
"source": "restart-replan",
"actions": [{
"id": "final-action",
"type": "tool_call",
"tool": "drive_cli",
"parameters": {"prompt": "authenticated final parameters"}
}],
"timestamp": proposal_timestamp,
"context": {}
}))
.unwrap();
let original_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&original_proposal)
.unwrap()
.as_bytes()
)
);
let final_digest = format!(
"{:x}",
Sha256::digest(
car_inference::catalog_identity::canonical_json(&final_proposal)
.unwrap()
.as_bytes()
)
);
let proposal_result = car_ir::ProposalResult {
proposal_id: final_proposal.id.clone(),
original_proposal_id: original_proposal.id.clone(),
final_proposal: Some(final_proposal.clone()),
replan_lineage: vec![
car_ir::ProposalLineageEntry {
generation: 0,
proposal_id: original_proposal.id.clone(),
proposal_digest: Some(original_digest.clone()),
status: car_ir::ProposalLineageStatus::Accepted,
rejection_reason: None,
},
car_ir::ProposalLineageEntry {
generation: 1,
proposal_id: final_proposal.id.clone(),
proposal_digest: Some(final_digest.clone()),
status: car_ir::ProposalLineageStatus::Accepted,
rejection_reason: None,
},
car_ir::ProposalLineageEntry {
generation: 2,
proposal_id: "proposal-rejected-after-final".to_string(),
proposal_digest: Some("d".repeat(64)),
status: car_ir::ProposalLineageStatus::Rejected,
rejection_reason: Some("replan candidate failed quality gate".to_string()),
},
],
accepted_proposal_preimages: vec![
car_ir::AcceptedProposalPreimage {
generation: 0,
proposal_digest: original_digest.clone(),
proposal: original_proposal.clone(),
},
car_ir::AcceptedProposalPreimage {
generation: 1,
proposal_digest: final_digest,
proposal: final_proposal.clone(),
},
],
results: vec![car_ir::ActionResult {
action_id: "final-action".to_string(),
status: car_ir::ActionStatus::Succeeded,
output: Some(serde_json::json!({"ok": true})),
error: None,
state_changes: std::collections::HashMap::new(),
duration_ms: Some(12.0),
timestamp: chrono::Utc::now(),
}],
cost: car_ir::CostSummary::default(),
};
let canonical = car_inference::catalog_identity::canonical_json(&proposal_result).unwrap();
let digest = format!("{:x}", Sha256::digest(canonical.as_bytes()));
let pending = car_server_core::run_store::PendingProposalFinalization {
run_id: started.run_id.clone(),
client_id: started.client_id.clone().unwrap(),
requested_policy_session_id: Some("policy-restart".to_string()),
policy_session_id: Some("policy-restart".to_string()),
original_proposal_id: "proposal-original".to_string(),
final_proposal_id: "proposal-final".to_string(),
original_submission: serde_json::json!({
"id": "proposal-original",
"source": "restart-test",
"actions": [{
"id": "original-action",
"type": "tool_call",
"tool": "drive_cli",
"parameters": {"prompt": "original parameters"}
}]
}),
original_proposal: original_proposal.clone(),
final_proposal: final_proposal.clone(),
accepted_proposal_preimages: vec![
car_server_core::run_store::AcceptedProposalPreimage {
generation: 0,
proposal: original_proposal.clone(),
},
car_server_core::run_store::AcceptedProposalPreimage {
generation: 1,
proposal: final_proposal,
},
],
proposal_result: proposal_result.clone(),
result_digest: digest.clone(),
};
let expected_terminal_data = serde_json::to_value(pending.event_data()).unwrap();
store
.write_execution_marker(&car_server_core::run_store::ProposalExecutionMarker {
run_id: pending.run_id.clone(),
client_id: pending.client_id.clone(),
requested_policy_session_id: pending.requested_policy_session_id.clone(),
policy_session_id: pending.policy_session_id.clone(),
original_proposal_id: pending.original_proposal_id.clone(),
original_submission: pending.original_submission.clone(),
original_proposal: serde_json::to_value(&pending.original_proposal).unwrap(),
proposal_digest: original_digest,
})
.unwrap();
store.write_pending_proposal(&pending).unwrap();
assert!(!journal_dir.exists(), "precondition: no journal exists");
let _restarted = ServerState::with_config(
ServerStateConfig::new(journal_dir.clone())
.with_shared_memgine(Arc::new(Mutex::new(MemgineEngine::new(None)))),
);
let rows: Vec<serde_json::Value> =
std::fs::read_to_string(journal_dir.join("restart-proposal-client.jsonl"))
.unwrap()
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
let kinds: Vec<_> = rows
.iter()
.map(|row| row["kind"].as_str().unwrap())
.collect();
assert_eq!(
kinds,
vec!["run_started", "proposal_completed", "run_completed"]
);
let terminal = &rows[1];
assert_eq!(terminal["run_id"], started.run_id);
assert_eq!(terminal["client_id"], pending.client_id);
assert_eq!(terminal["policy_session_id"], "policy-restart");
assert_eq!(terminal["proposal_id"], pending.final_proposal_id);
assert_eq!(
terminal["data"], expected_terminal_data,
"restart reconstruction must emit the same terminal bytes as fresh finalization"
);
assert_eq!(
terminal["data"]["proposal_result"],
serde_json::to_value(&proposal_result).unwrap()
);
assert_eq!(terminal["data"]["result_digest"], digest);
assert_eq!(
terminal["data"]["original_submission_id"],
pending.original_proposal_id
);
assert_eq!(
terminal["data"]["replan_lineage"][2]["status"], "rejected",
"a rejected tail remains in lineage while the most recent accepted proposal is final"
);
assert_eq!(terminal["data"]["final_proposal_id"], "proposal-final");
let turns: Vec<_> = store
.get_run_trace(&started.run_id)
.unwrap()
.into_iter()
.filter_map(|record| match record {
car_proto::RunRecord::Turn(turn) => Some(turn),
_ => None,
})
.collect();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].proposal_id.as_deref(), Some("proposal-final"));
assert_eq!(turns[0].action_id.as_deref(), Some("final-action"));
assert_eq!(
turns[0].parameters,
serde_json::json!({"prompt": "authenticated final parameters"})
);
assert!(store.pending_proposal(&started.run_id).unwrap().is_none());
assert!(store.execution_marker(&started.run_id).unwrap().is_none());
}
#[tokio::test]
async fn restart_quarantines_every_pending_provenance_mismatch_without_stamping_or_cleanup() {
for mismatch in [
"client",
"requested_policy",
"authenticated_policy",
"run",
"original_proposal",
"raw_null",
"raw_array",
"raw_missing_actions",
"raw_mismatched_actions",
] {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let run_id = format!("provenance-{mismatch}");
let client_id = format!("client-{mismatch}");
let (requested, authenticated) = if mismatch == "requested_policy" {
(Some("requested-policy"), None)
} else {
(Some("live-policy"), Some("live-policy"))
};
let (started, marker, mut pending) =
proposal_finalization_fixture(&run_id, &client_id, requested, authenticated);
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
store.write_started(&started).unwrap();
store.write_execution_marker(&marker).unwrap();
match mismatch {
"client" => pending.client_id = "pending-self-claimed-client".to_string(),
"requested_policy" => {
pending.requested_policy_session_id = Some("other-requested-policy".to_string())
}
"authenticated_policy" => pending.policy_session_id = None,
"run" => pending.run_id = "pending-self-claimed-run".to_string(),
"original_proposal" => mutate_pending_original_proposal(&mut pending),
"raw_null" => pending.original_submission = serde_json::Value::Null,
"raw_array" => pending.original_submission = serde_json::json!([]),
"raw_missing_actions" => {
pending.original_submission = serde_json::json!({
"id": pending.original_proposal_id.clone(),
"source": pending.original_proposal.source.clone()
})
}
"raw_mismatched_actions" => {
pending.original_submission["actions"] = serde_json::json!([{
"id": "not-the-accepted-action",
"type": "state_read",
"parameters": {"key": "untrusted"}
}])
}
_ => unreachable!(),
}
write_proposal_sidecar_value(tmp.path(), "proposal-finalization", &run_id, &pending);
let _restarted = loopback_state(journal_dir.clone());
let journal_rows: Vec<serde_json::Value> = if journal_dir.exists() {
std::fs::read_dir(&journal_dir)
.unwrap()
.filter_map(Result::ok)
.filter(|entry| {
entry.path().extension().and_then(|ext| ext.to_str()) == Some("jsonl")
})
.flat_map(|entry| {
std::fs::read_to_string(entry.path())
.unwrap()
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect::<Vec<_>>()
})
.collect()
} else {
Vec::new()
};
assert!(
journal_rows
.iter()
.all(|row| row["kind"] != "proposal_completed"),
"{mismatch} mismatch must never stamp proposal_completed: {journal_rows:?}"
);
assert!(
store
.get_run_trace(&run_id)
.unwrap()
.iter()
.all(|record| !matches!(record, car_proto::RunRecord::Ended(_))),
"{mismatch} mismatch must preserve the run as outcome-unknown"
);
let key = format!("{:x}", Sha256::digest(run_id.as_bytes()));
assert!(
tmp.path()
.join("proposal-execution")
.join(format!("{key}.json"))
.exists(),
"{mismatch} mismatch must preserve the execution marker"
);
assert!(
tmp.path()
.join("proposal-finalization")
.join(format!("{key}.json"))
.exists(),
"{mismatch} mismatch must preserve the finalization outbox"
);
}
}
#[tokio::test]
async fn restart_quarantines_execution_in_progress_without_an_exact_result() {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let started = car_proto::RunStarted {
run_id: "outcome-unknown-run".to_string(),
client_id: Some("outcome-unknown-client".to_string()),
agent_id: "agent-outcome-unknown".to_string(),
intent: "never guess whether effect happened".to_string(),
outcome_description: None,
started_at: chrono::Utc::now(),
};
store.write_started(&started).unwrap();
let original = serde_json::json!({
"id": "outcome-unknown-proposal",
"source": "crash-test",
"actions": [{
"id": "possibly-irreversible",
"type": "tool_call",
"tool": "external_side_effect",
"parameters": {}
}]
});
let typed_original: car_ir::ActionProposal = serde_json::from_value(original.clone()).unwrap();
let normal_original = serde_json::to_value(&typed_original).unwrap();
let original_canonical =
car_inference::catalog_identity::canonical_json(&normal_original).unwrap();
let marker = car_server_core::run_store::ProposalExecutionMarker {
run_id: started.run_id.clone(),
client_id: started.client_id.clone().unwrap(),
requested_policy_session_id: None,
policy_session_id: None,
original_proposal_id: "outcome-unknown-proposal".to_string(),
original_submission: original.clone(),
original_proposal: normal_original,
proposal_digest: format!("{:x}", Sha256::digest(original_canonical.as_bytes())),
};
store.write_execution_marker(&marker).unwrap();
let restarted = loopback_state(journal_dir.clone());
assert!(
store.execution_marker(&started.run_id).unwrap().is_some(),
"startup must preserve the outcome-unknown quarantine"
);
let trace = store.get_run_trace(&started.run_id).unwrap();
assert!(
!trace
.iter()
.any(|row| matches!(row, car_proto::RunRecord::Ended(_))),
"outcome unknown must not be fabricated as Incomplete"
);
let rows: Vec<serde_json::Value> =
std::fs::read_to_string(journal_dir.join("outcome-unknown-client.jsonl"))
.unwrap()
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
assert_eq!(
rows.iter()
.map(|row| row["kind"].as_str().unwrap())
.collect::<Vec<_>>(),
vec!["run_started"],
"startup may restore the bracket but cannot invent a proposal/run terminal"
);
let addr = spawn_dispatcher(restarted).await;
let mut ws = connect_async(format!("ws://{addr}")).await.unwrap().0;
let rejected = call_raw(
&mut ws,
"outcome-unknown-reopen",
"runs.start",
serde_json::json!({
"agent_id": started.agent_id,
"intent": started.intent,
"idempotency_key": started.run_id
}),
)
.await;
assert!(
rejected["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("proposal execution outcome unknown"),
"restart must refuse automatic redispatch/reopen: {rejected}"
);
}
async fn assert_corrupt_sidecar_quarantined_across_restart(kind: &str) {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let started = car_proto::RunStarted {
run_id: format!("restart-corrupt-{kind}"),
client_id: Some(format!("restart-corrupt-client-{kind}")),
agent_id: "agent-restart-corrupt".to_string(),
intent: "preserve unreadable proposal durability state".to_string(),
outcome_description: None,
started_at: chrono::Utc::now(),
};
store.write_started(&started).unwrap();
write_truncated_proposal_sidecar(tmp.path(), kind, &started.run_id);
let restarted = loopback_state(journal_dir.clone());
let trace = store.get_run_trace(&started.run_id).unwrap();
assert!(
trace
.iter()
.all(|record| !matches!(record, car_proto::RunRecord::Ended(_))),
"startup must not fabricate a terminal for corrupt {kind}"
);
let addr = spawn_dispatcher(restarted).await;
let mut ws = connect_async(format!("ws://{addr}")).await.unwrap().0;
let reopen = call_raw(
&mut ws,
"restart-corrupt-reopen",
"runs.start",
serde_json::json!({
"agent_id": started.agent_id,
"intent": started.intent,
"idempotency_key": started.run_id
}),
)
.await;
assert!(
reopen["error"]["message"]
.as_str()
.unwrap_or_default()
.contains("durability state is unreadable"),
"restart must preserve corrupt {kind} quarantine: {reopen}"
);
}
#[tokio::test]
async fn restart_quarantines_corrupt_execution_marker() {
assert_corrupt_sidecar_quarantined_across_restart("proposal-execution").await;
}
#[tokio::test]
async fn restart_quarantines_corrupt_finalization_outbox() {
assert_corrupt_sidecar_quarantined_across_restart("proposal-finalization").await;
}
#[tokio::test]
async fn restart_quarantines_invalid_typed_outbox_and_lineage() {
for invalid_lineage in [false, true] {
let tmp = TempDir::new().unwrap();
let journal_dir = tmp.path().join("journals");
let store = car_server_core::run_store::RunStore::from_journal_dir(&journal_dir);
let run_id = if invalid_lineage {
"invalid-lineage-run"
} else {
"invalid-typed-result-run"
};
let started = car_proto::RunStarted {
run_id: run_id.to_string(),
client_id: Some(format!("client-{run_id}")),
agent_id: "agent-invalid-outbox".to_string(),
intent: "reject untrusted outbox".to_string(),
outcome_description: None,
started_at: chrono::Utc::now(),
};
store.write_started(&started).unwrap();
let proposal = serde_json::json!({
"id":"invalid-outbox-proposal",
"source":"test",
"actions":[],
"timestamp":"2026-08-28T10:00:00Z",
"context":{}
});
let mut result = serde_json::json!({
"proposal_id":"invalid-outbox-proposal",
"original_proposal_id":"invalid-outbox-proposal",
"final_proposal": proposal,
"replan_lineage":[{
"generation":0,
"proposal_id":"invalid-outbox-proposal",
"proposal_digest":"a".repeat(64),
"status":"accepted"
}],
"results":[],
"cost":{
"tool_calls":0,
"actions_executed":0,
"actions_rejected":0,
"actions_skipped":0,
"total_duration_ms":0.0,
"retries":0
}
});
if invalid_lineage {
result["replan_lineage"][0]["proposal_digest"] = serde_json::json!("A".repeat(64));
} else {
result.as_object_mut().unwrap().remove("cost");
}
let outbox = serde_json::json!({
"run_id": started.run_id,
"client_id": started.client_id,
"original_proposal_id":"invalid-outbox-proposal",
"final_proposal_id":"invalid-outbox-proposal",
"original_submission": {"id":"invalid-outbox-proposal","source":"test","actions":[]},
"original_proposal": proposal,
"final_proposal": proposal,
"accepted_proposal_preimages":[{"generation":0,"proposal":proposal}],
"proposal_result": result,
"result_digest":"b".repeat(64)
});
let root = tmp.path().join("proposal-finalization");
car_secrets::ensure_private_dir(&root).unwrap();
let key = format!("{:x}", Sha256::digest(run_id.as_bytes()));
let mut file = car_secrets::create_private_file(&root.join(format!("{key}.json"))).unwrap();
serde_json::to_writer(&mut file, &outbox).unwrap();
file.sync_all().unwrap();
let _restarted = loopback_state(journal_dir);
let trace = store.get_run_trace(run_id).unwrap();
assert!(
trace
.iter()
.all(|record| !matches!(record, car_proto::RunRecord::Ended(_))),
"invalid active-v3 outbox must remain quarantined (lineage={invalid_lineage})"
);
assert!(store.pending_proposal(run_id).is_err());
}
}