use car_memgine::MemgineEngine;
use car_server_core::{run_dispatch, ServerState, ServerStateConfig};
use futures::{SinkExt, StreamExt};
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>>;
const AUTH_TOKEN: &str = "callback-state-auth-token-0123456789abcdef";
const HOST_TOKEN: &str = "callback-state-host-token-0123456789abcdef";
fn loopback_state(journal_dir: std::path::PathBuf) -> Arc<ServerState> {
let engine = Arc::new(Mutex::new(MemgineEngine::new(None)));
let config = ServerStateConfig::new(journal_dir).with_shared_memgine(engine);
let state = Arc::new(ServerState::with_config(config));
state
.install_auth_token(AUTH_TOKEN.to_string())
.expect("install ordinary auth token");
state
.install_host_token(HOST_TOKEN.to_string())
.expect("install host token");
state
}
async fn spawn_dispatcher(state: Arc<ServerState>) -> SocketAddr {
let listener = TcpListener::bind(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 0)))
.await
.expect("bind loopback");
let address = listener.local_addr().expect("dispatcher address");
tokio::spawn(async move {
loop {
let (stream, peer) = match listener.accept().await {
Ok(accepted) => accepted,
Err(_) => return,
};
let state = state.clone();
tokio::spawn(async move {
let websocket = match accept_async(stream).await {
Ok(websocket) => websocket,
Err(_) => return,
};
let (write, read) = websocket.split();
let _ = run_dispatch(read, Box::pin(write), peer.to_string(), state).await;
});
}
});
address
}
async fn connect(address: SocketAddr) -> Ws {
connect_async(format!("ws://{address}"))
.await
.expect("connect websocket")
.0
}
async fn send(ws: &mut Ws, id: &str, method: &str, params: 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 websocket request");
}
async fn next_json(ws: &mut Ws) -> serde_json::Value {
loop {
match ws
.next()
.await
.expect("websocket frame")
.expect("valid frame")
{
Message::Text(text) => return serde_json::from_str(&text).expect("JSON text frame"),
Message::Ping(_) | Message::Pong(_) => continue,
other => panic!("unexpected websocket frame: {other:?}"),
}
}
}
async fn call(ws: &mut Ws, id: &str, method: &str, params: serde_json::Value) -> serde_json::Value {
send(ws, id, method, params).await;
loop {
let message = next_json(ws).await;
if message.get("id").and_then(serde_json::Value::as_str) == Some(id) {
return message;
}
}
}
async fn authenticate_with_callback_state(ws: &mut Ws, host: bool, callback_state: bool) {
let params = if host {
serde_json::json!({"host_token": HOST_TOKEN})
} else {
serde_json::json!({"token": AUTH_TOKEN})
};
let auth = call(ws, "auth", "session.auth", params).await;
assert!(auth.get("error").is_none(), "auth failed: {auth}");
let mut required_capabilities = vec![car_proto::RUNS_PAGINATION_CAPABILITY];
if callback_state {
required_capabilities.push(car_proto::TOOLS_CALLBACK_STATE_CAPABILITY);
}
let handshake = call(
ws,
"handshake",
"server.handshake",
serde_json::json!({
"protocol_version": car_proto::PROTOCOL_VERSION,
"required_capabilities": required_capabilities
}),
)
.await;
assert_eq!(
handshake["result"]["protocol_version"],
car_proto::PROTOCOL_VERSION,
"protocol handshake: {handshake}"
);
let negotiated = handshake["result"]["negotiated_capabilities"]
.as_array()
.expect("negotiated capabilities");
assert_eq!(
negotiated
.iter()
.any(|capability| capability == car_proto::TOOLS_CALLBACK_STATE_CAPABILITY),
callback_state,
"callback-state negotiation: {handshake}"
);
}
async fn authenticate(ws: &mut Ws, host: bool) {
authenticate_with_callback_state(ws, host, true).await;
}
async fn register_newsroom_tools(ws: &mut Ws) {
let response = call(
ws,
"register",
"tools.register",
serde_json::json!([
{
"name": "report_source",
"description": "produce one source report",
"parameters": {"type":"object"},
"returns": {
"type":"object",
"properties":{"report_id":{"type":"string"}},
"required":["report_id"],
"additionalProperties":false
}
},
{
"name": "edit_report",
"description": "consume a filed report",
"parameters": {"type":"object"},
"returns": {
"type":"object",
"properties":{"edition":{"type":"string"}},
"required":["edition"],
"additionalProperties":false
}
}
]),
)
.await;
assert!(
response.get("error").is_none(),
"register tools: {response}"
);
}
fn newsroom_proposal() -> serde_json::Value {
serde_json::json!({
"proposal": {
"id": "proposal-newsroom-callback-state",
"source": "callback-state-test",
"actions": [
{
"id": "reporter-calendar",
"type": "tool_call",
"tool": "report_source",
"parameters": {},
"expected_effects": {"report.calendar": true},
"state_dependencies": [],
"max_retries": 0,
"failure_behavior": "abort"
},
{
"id": "editor",
"type": "tool_call",
"tool": "edit_report",
"parameters": {},
"expected_effects": {},
"state_dependencies": ["report.calendar"],
"max_retries": 0,
"failure_behavior": "abort"
}
]
}
})
}
fn one_action_proposal(
proposal_id: &str,
action_id: &str,
expected_effects: serde_json::Value,
) -> serde_json::Value {
serde_json::json!({
"proposal": {
"id": proposal_id,
"source": "callback-state-negative-test",
"actions": [{
"id": action_id,
"type": "tool_call",
"tool": "report_source",
"parameters": {},
"expected_effects": expected_effects,
"state_dependencies": [],
"max_retries": 0,
"failure_behavior": "abort"
}]
}
})
}
async fn submit_with_callback_frame(
ws: &mut Ws,
request_id: &str,
proposal: serde_json::Value,
callback_frame: impl FnOnce(&str) -> serde_json::Value,
) -> serde_json::Value {
send(ws, request_id, "proposal.submit", proposal).await;
let mut callback_frame = Some(callback_frame);
loop {
let message = next_json(ws).await;
if message.get("method").and_then(serde_json::Value::as_str) == Some("tools.execute") {
let callback_id = message["id"].as_str().expect("callback id");
let frame = callback_frame
.take()
.expect("one callback expected for one-action proposal")(
callback_id
);
ws.send(Message::Text(frame.to_string().into()))
.await
.expect("send callback response");
continue;
}
if message.get("id").and_then(serde_json::Value::as_str) == Some(request_id) {
return message;
}
}
}
fn assert_failed_without_state(
response: &serde_json::Value,
runtime: &Arc<car_engine::Runtime>,
absent_keys: &[&str],
) {
let result = &response["result"]["results"][0];
assert_eq!(result["status"], "failed", "{response}");
assert_eq!(result["state_changes"], serde_json::json!({}), "{response}");
for key in absent_keys {
assert_eq!(runtime.state.get(key), None, "{key} must not be written");
}
}
#[tokio::test]
async fn unnegotiated_v3_preserves_legacy_raw_results_without_runtime_writes() {
let temp = TempDir::new().expect("temporary state");
let state = loopback_state(temp.path().join("journals"));
let address = spawn_dispatcher(state.clone()).await;
let mut harness = connect(address).await;
authenticate_with_callback_state(&mut harness, false, false).await;
register_newsroom_tools(&mut harness).await;
let effectful = submit_with_callback_frame(
&mut harness,
"submit-legacy-effectful",
one_action_proposal(
"proposal-legacy-effectful",
"action-legacy-effectful",
serde_json::json!({"report.legacy":true}),
),
|callback_id| {
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":{"report_id":"legacy-effectful"}
})
},
)
.await;
assert_eq!(
effectful["result"]["results"][0]["status"], "succeeded",
"an unnegotiated v3 host keeps legacy raw-output behavior: {effectful}"
);
assert_eq!(
effectful["result"]["results"][0]["output"],
serde_json::json!({"report_id":"legacy-effectful"})
);
assert_eq!(
effectful["result"]["results"][0]["state_changes"],
serde_json::json!({})
);
let raw_envelope_shape = serde_json::json!({
"output":{"report_id":"nested-output"},
"state_changes":{"this":"is ordinary tool output"}
});
let expected_raw_output = raw_envelope_shape.clone();
let no_effects = submit_with_callback_frame(
&mut harness,
"submit-legacy-envelope-shape",
one_action_proposal(
"proposal-legacy-envelope-shape",
"action-legacy-envelope-shape",
serde_json::json!({}),
),
move |callback_id| {
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":raw_envelope_shape
})
},
)
.await;
assert_eq!(
no_effects["result"]["results"][0]["status"], "succeeded",
"legacy exact-envelope-shaped output remains raw: {no_effects}"
);
assert_eq!(
no_effects["result"]["results"][0]["output"], expected_raw_output,
"an unnegotiated client must never have ordinary output unwrapped"
);
let sessions: Vec<_> = state.sessions.lock().await.values().cloned().collect();
let runtime = sessions
.into_iter()
.find(|session| !session.client_id.is_empty())
.expect("legacy harness session")
.runtime
.clone();
assert_eq!(runtime.state.get("report.legacy"), None);
assert_eq!(runtime.state.get("this"), None);
}
#[tokio::test]
async fn authenticated_callback_state_unblocks_dependent_action_and_journals_actual_mutation() {
let temp = TempDir::new().expect("temporary state");
let state = loopback_state(temp.path().join("journals"));
let address = spawn_dispatcher(state.clone()).await;
let mut harness = connect(address).await;
authenticate(&mut harness, false).await;
register_newsroom_tools(&mut harness).await;
let started = call(
&mut harness,
"start",
"runs.start",
serde_json::json!({"agent_id":"daily-newsroom","intent":"publish the morning edition"}),
)
.await;
let run_id = started["result"]["run_id"]
.as_str()
.expect("run id")
.to_string();
let mut subscriber = connect(address).await;
authenticate(&mut subscriber, true).await;
let subscription = call(
&mut subscriber,
"subscribe",
"runs.subscribe",
serde_json::json!({"run_id":run_id,"cursor":0,"limit":100}),
)
.await;
assert_eq!(subscription["result"]["cursor"], 0, "{subscription}");
send(
&mut harness,
"submit",
"proposal.submit",
newsroom_proposal(),
)
.await;
let proposal = loop {
let message = next_json(&mut harness).await;
if message.get("method").and_then(serde_json::Value::as_str) == Some("tools.execute") {
let callback_id = message["id"].as_str().expect("callback request id");
let action_id = message["params"]["action_id"]
.as_str()
.expect("callback action id");
let result = match action_id {
"reporter-calendar" => serde_json::json!({
"output": {"report_id":"calendar-2026-09-01"},
"state_changes": {
"report.calendar": {
"status":"healthy",
"receipt_digest":"sha256:calendar-receipt"
}
}
}),
"editor" => serde_json::json!({"edition":"2026-09-01-v1"}),
other => panic!("unexpected action callback: {other}"),
};
harness
.send(Message::Text(
serde_json::json!({"jsonrpc":"2.0","id":callback_id,"result":result})
.to_string()
.into(),
))
.await
.expect("reply to tool callback");
continue;
}
if message.get("id").and_then(serde_json::Value::as_str) == Some("submit") {
break message;
}
};
assert!(
proposal.get("error").is_none(),
"proposal submit: {proposal}"
);
let results = proposal["result"]["results"]
.as_array()
.expect("proposal results");
assert_eq!(results.len(), 2, "{proposal}");
assert!(
results.iter().all(|result| result["status"] == "succeeded"),
"reporter state must make the editor runnable: {proposal}"
);
let sessions: Vec<_> = state.sessions.lock().await.values().cloned().collect();
let mut harness_session = None;
for session in sessions {
if session.current_run_id.lock().await.as_deref() == Some(run_id.as_str()) {
harness_session = Some(session);
break;
}
}
let session = harness_session.expect("the harness session owns the active run");
let stored = session
.runtime
.state
.get("report.calendar")
.expect("callback mutation reached the StateStore");
assert_eq!(stored["status"], "healthy");
let log = session.runtime.log.lock().await;
let state_event = log
.events()
.iter()
.find(|event| {
event.kind == car_eventlog::EventKind::StateChanged
&& event.action_id.as_deref() == Some("reporter-calendar")
})
.expect("reporter state_changed event");
assert_eq!(
state_event.data["runtime_state_mutations"]["report.calendar"]["op"],
"set"
);
assert_eq!(
state_event.data["runtime_state_mutations"]["report.calendar"]["value"],
stored
);
}
#[tokio::test]
async fn invalid_callback_envelopes_tool_errors_and_wrong_session_never_write_state() {
let temp = TempDir::new().expect("temporary state");
let state = loopback_state(temp.path().join("journals"));
let address = spawn_dispatcher(state.clone()).await;
let mut harness = connect(address).await;
authenticate(&mut harness, false).await;
register_newsroom_tools(&mut harness).await;
let started = call(
&mut harness,
"negative-start",
"runs.start",
serde_json::json!({"agent_id":"daily-newsroom","intent":"reject invalid callbacks"}),
)
.await;
let run_id = started["result"]["run_id"]
.as_str()
.expect("run id")
.to_string();
let sessions: Vec<_> = state.sessions.lock().await.values().cloned().collect();
let mut harness_session = None;
for session in sessions {
if session.current_run_id.lock().await.as_deref() == Some(run_id.as_str()) {
harness_session = Some(session);
break;
}
}
let runtime = harness_session.expect("harness session").runtime.clone();
let cases = [
(
"missing-key",
serde_json::json!({"report.missing":true}),
serde_json::json!({"output":{"report_id":"r1"},"state_changes":{}}),
vec!["report.missing"],
),
(
"extra-key",
serde_json::json!({"report.extra":true}),
serde_json::json!({
"output":{"report_id":"r2"},
"state_changes":{"report.extra":{},"report.unexpected":{}}
}),
vec!["report.extra", "report.unexpected"],
),
(
"wrong-key",
serde_json::json!({"report.expected":true}),
serde_json::json!({
"output":{"report_id":"r3"},
"state_changes":{"report.wrong":{}}
}),
vec!["report.expected", "report.wrong"],
),
(
"missing-envelope",
serde_json::json!({"report.envelope":true}),
serde_json::json!({"report_id":"r4"}),
vec!["report.envelope"],
),
(
"extra-envelope-field",
serde_json::json!({"report.envelope-extra":true}),
serde_json::json!({
"output":{"report_id":"r4-extra"},
"state_changes":{"report.envelope-extra":{}},
"action_id":"action-envelope-extra"
}),
vec!["report.envelope-extra"],
),
(
"invalid-output-schema",
serde_json::json!({"report.schema":true}),
serde_json::json!({
"output":{"report_id":42},
"state_changes":{"report.schema":{}}
}),
vec!["report.schema"],
),
(
"invalid-state-value",
serde_json::json!({"report.state-value":true}),
serde_json::json!({
"output":{"report_id":"r5"},
"state_changes":{"report.state-value":9007199254740992_i64}
}),
vec!["report.state-value"],
),
(
"no-effects-nonempty",
serde_json::json!({}),
serde_json::json!({
"output":{"report_id":"r6"},
"state_changes":{"report.forbidden":{}}
}),
vec!["report.forbidden"],
),
];
for (name, expected_effects, callback_result, absent_keys) in cases {
let response = submit_with_callback_frame(
&mut harness,
&format!("submit-{name}"),
one_action_proposal(
&format!("proposal-{name}"),
&format!("action-{name}"),
expected_effects,
),
move |callback_id| {
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":callback_result
})
},
)
.await;
assert_failed_without_state(&response, &runtime, &absent_keys);
}
let empty_changes = submit_with_callback_frame(
&mut harness,
"submit-empty-changes",
one_action_proposal(
"proposal-empty-changes",
"action-empty-changes",
serde_json::json!({}),
),
|callback_id| {
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":{
"output":{"report_id":"empty-changes"},
"state_changes":{}
}
})
},
)
.await;
assert_eq!(
empty_changes["result"]["results"][0]["status"], "succeeded",
"an action with no expected effects may return an empty mutation envelope: {empty_changes}"
);
assert_eq!(
empty_changes["result"]["results"][0]["state_changes"],
serde_json::json!({}),
"empty mutation envelopes cannot fabricate writes: {empty_changes}"
);
let tool_error = submit_with_callback_frame(
&mut harness,
"submit-tool-error",
one_action_proposal(
"proposal-tool-error",
"action-tool-error",
serde_json::json!({"report.error":true}),
),
|callback_id| {
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"error":{"code":-32000,"message":"reporter failed"}
})
},
)
.await;
assert_failed_without_state(&tool_error, &runtime, &["report.error"]);
let mut wrong_session = connect(address).await;
authenticate(&mut wrong_session, false).await;
send(
&mut harness,
"submit-wrong-session",
"proposal.submit",
one_action_proposal(
"proposal-wrong-session",
"action-wrong-session",
serde_json::json!({"report.session":true}),
),
)
.await;
loop {
let message = next_json(&mut harness).await;
if message.get("method").and_then(serde_json::Value::as_str) == Some("tools.execute") {
let callback_id = message["id"].as_str().expect("callback id");
wrong_session
.send(Message::Text(
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":{
"output":{"report_id":"wrong-session"},
"state_changes":{"report.session":{}}
}
})
.to_string()
.into(),
))
.await
.expect("send response on wrong session");
harness
.send(Message::Text(
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"error":{"code":-32000,"message":"right session rejects tool"}
})
.to_string()
.into(),
))
.await
.expect("fail the real callback");
continue;
}
if message.get("id").and_then(serde_json::Value::as_str) == Some("submit-wrong-session") {
assert_failed_without_state(&message, &runtime, &["report.session"]);
break;
}
}
}
#[tokio::test]
async fn swapped_action_callbacks_and_duplicate_response_do_not_cross_or_repeat_writes() {
let temp = TempDir::new().expect("temporary state");
let state = loopback_state(temp.path().join("journals"));
let address = spawn_dispatcher(state.clone()).await;
let mut harness = connect(address).await;
authenticate(&mut harness, false).await;
register_newsroom_tools(&mut harness).await;
let started = call(
&mut harness,
"correlation-start",
"runs.start",
serde_json::json!({"agent_id":"daily-newsroom","intent":"bind callback correlation"}),
)
.await;
let run_id = started["result"]["run_id"]
.as_str()
.expect("run id")
.to_string();
let sessions: Vec<_> = state.sessions.lock().await.values().cloned().collect();
let mut runtime = None;
for session in sessions {
if session.current_run_id.lock().await.as_deref() == Some(run_id.as_str()) {
runtime = Some(session.runtime.clone());
break;
}
}
let runtime = runtime.expect("harness runtime");
let swapped_proposal = serde_json::json!({
"proposal": {
"id":"proposal-swapped",
"source":"callback-state-negative-test",
"actions":[
{"id":"action-left","type":"tool_call","tool":"report_source","parameters":{},"expected_effects":{"report.left":true},"failure_behavior":"skip"},
{"id":"action-right","type":"tool_call","tool":"report_source","parameters":{},"expected_effects":{"report.right":true},"failure_behavior":"skip"}
]
}
});
send(
&mut harness,
"submit-swapped",
"proposal.submit",
swapped_proposal,
)
.await;
let mut callbacks = Vec::new();
let swapped = loop {
let message = next_json(&mut harness).await;
if message.get("method").and_then(serde_json::Value::as_str) == Some("tools.execute") {
callbacks.push((
message["id"].as_str().expect("callback id").to_string(),
message["params"]["action_id"]
.as_str()
.expect("action id")
.to_string(),
));
if callbacks.len() == 2 {
for (callback_id, action_id) in &callbacks {
let wrong_key = if action_id == "action-left" {
"report.right"
} else {
"report.left"
};
harness
.send(Message::Text(
serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":{
"output":{"report_id":format!("wrong-for-{action_id}")},
"state_changes":{(wrong_key):{}}
}
})
.to_string()
.into(),
))
.await
.expect("send swapped callback");
}
}
continue;
}
if message.get("id").and_then(serde_json::Value::as_str) == Some("submit-swapped") {
break message;
}
};
assert!(
swapped["result"]["results"]
.as_array()
.expect("swapped results")
.iter()
.all(|result| {
result["status"] == "skipped" && result["state_changes"] == serde_json::json!({})
}),
"{swapped}"
);
assert_eq!(runtime.state.get("report.left"), None);
assert_eq!(runtime.state.get("report.right"), None);
send(
&mut harness,
"submit-duplicate",
"proposal.submit",
one_action_proposal(
"proposal-duplicate",
"action-duplicate",
serde_json::json!({"report.duplicate":true}),
),
)
.await;
let duplicate_response = loop {
let message = next_json(&mut harness).await;
if message.get("method").and_then(serde_json::Value::as_str) == Some("tools.execute") {
let callback_id = message["id"].as_str().expect("callback id");
let response = serde_json::json!({
"jsonrpc":"2.0",
"id":callback_id,
"result":{
"output":{"report_id":"duplicate"},
"state_changes":{"report.duplicate":{"receipt":"one"}}
}
});
harness
.send(Message::Text(response.to_string().into()))
.await
.expect("send first response");
harness
.send(Message::Text(response.to_string().into()))
.await
.expect("send duplicate response");
continue;
}
if message.get("id").and_then(serde_json::Value::as_str) == Some("submit-duplicate") {
break message;
}
};
assert_eq!(
duplicate_response["result"]["results"][0]["status"], "succeeded",
"{duplicate_response}"
);
let transitions: Vec<_> = runtime
.state
.transitions()
.into_iter()
.filter(|transition| transition.key == "report.duplicate")
.collect();
assert_eq!(
transitions.len(),
1,
"duplicate callback must not write twice"
);
assert_eq!(transitions[0].action_id, "action-duplicate");
}