use super::*;
fn queued_messages(coordinator: &Coordinator, connection: &str) -> Vec<Value> {
coordinator.connections[connection]
.queue
.iter()
.map(|record| serde_json::from_str(record).unwrap())
.collect()
}
#[test]
fn claim_snapshot_and_terminal_delivery_cover_both_sequence_boundaries() {
for terminal_before_claim in [false, true] {
let temp = TempDir::new().unwrap();
std::fs::write(temp.path().join("evidence.txt"), "synthetic evidence").unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
let (release, wait) = bounded(1);
let (cleanup, cleanup_wait) = bounded(1);
install_tool_worker(
&mut coordinator,
wait,
Arc::new(AtomicUsize::new(0)),
cleanup_wait,
);
let first = connect(&mut coordinator);
let session = create(&mut coordinator, &first);
let grant = claim(&mut coordinator, &first, &session);
let start = request(
&coordinator,
&first,
"turn.start",
Some(&session),
Some(grant),
Some("boundary-turn"),
json!({"prompt":"Read the synthetic fixture"}),
);
let accepted = invoke(&mut coordinator, &first, start);
assert_eq!(accepted["payload"]["status"], "accepted");
coordinator.disconnect(&first);
release.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "cleaning"
});
if terminal_before_claim {
cleanup.send(()).unwrap();
run_until(&mut coordinator, |state| {
!state.actors.contains_key(&session)
});
}
let second = connect(&mut coordinator);
let claim_request = request(
&coordinator,
&second,
"session.claim",
Some(&session),
None,
Some("boundary-claim"),
json!({}),
);
coordinator
.submit(
&second,
&serde_json::to_vec(&claim_request).unwrap(),
Instant::now(),
)
.unwrap();
if !terminal_before_claim {
cleanup.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "idle"
});
}
coordinator.tick(Instant::now());
let messages = queued_messages(&coordinator, &second);
let reply = &messages[0];
assert_eq!(reply["request_id"], claim_request.request_id);
assert!(reply["error"].is_null(), "{reply}");
let snapshot = &reply["payload"]["snapshot"];
let boundary = snapshot["live_sequence"].as_u64().unwrap();
let terminal = &coordinator.terminals[&session];
if terminal_before_claim {
assert_eq!(
messages.len(),
1,
"settled terminal must not be redelivered"
);
assert_eq!(snapshot["phase"], "idle");
assert!(snapshot["turn"].is_null());
assert_eq!(snapshot["terminal"], terminal.payload);
assert_eq!(boundary, terminal.sequence);
} else {
assert_eq!(
messages.len(),
2,
"claim must deliver exactly one later terminal"
);
assert_eq!(snapshot["phase"], "cleaning");
assert!(snapshot["terminal"].is_null());
let event = &messages[1];
assert_eq!(event["event"], "turn.terminal");
assert_eq!(event["live_sequence"], boundary + 1);
assert_eq!(event["live_sequence"], terminal.sequence);
assert_eq!(event["grant_generation"], reply["payload"]["generation"]);
assert_eq!(event["payload"]["turn_id"], snapshot["turn"]["turn_id"]);
assert_eq!(
event["payload"]["assistant_text"],
"before tool\n\nafter tool"
);
assert_eq!(event["payload"]["status"], "completed");
}
coordinator.disconnect(&second);
}
}
#[test]
fn committed_login_survives_initiating_disconnect_and_lookup_reports_cleanup() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let paths = runtime.config.paths.clone();
let worker_paths = paths.clone();
let mut coordinator = Coordinator::new(runtime).unwrap();
let (committed, commit_observed) = bounded(1);
let (release, wait) = bounded(1);
coordinator.execution.set_login_worker(Arc::new(move |job| {
let generation = crate::config::read_auth_store(&worker_paths)
.unwrap()
.provider_generation("openai-codex");
crate::login::complete_service_codex_login(
&worker_paths,
generation,
&job.cancel,
|| {
Ok(crate::config::NormalizedToken {
access: "synthetic-access".into(),
refresh: Some("synthetic-refresh".into()),
expires: Some(i64::MAX),
account_id: "synthetic-account".into(),
})
},
|| {
committed.send(Arc::clone(&job.cancel)).unwrap();
wait.recv_timeout(Duration::from_secs(5)).unwrap();
},
)
}));
let owner = connect(&mut coordinator);
let start = request(
&coordinator,
&owner,
"auth.login.start",
None,
None,
Some("committed-login"),
json!({"provider_id":"openai-codex"}),
);
assert!(invoke(&mut coordinator, &owner, start)["error"].is_null());
let cancel = commit_observed
.recv_timeout(Duration::from_secs(2))
.unwrap();
let expected_credentials = crate::config::AuthProviderRecord::OAuth {
access: "synthetic-access".into(),
refresh: Some("synthetic-refresh".into()),
expires: Some(i64::MAX),
account_id: Some("synthetic-account".into()),
};
assert_eq!(
crate::config::read_auth(&paths)
.unwrap()
.providers
.get("openai-codex"),
Some(&expected_credentials),
);
coordinator.disconnect(&owner);
run_until(&mut coordinator, |_| cancel.load(Ordering::Relaxed));
assert_eq!(
coordinator.operations.lookup(
&coordinator.instance,
&coordinator.instance,
"committed-login"
)["state"],
"accepted",
"commit alone does not establish cleanup"
);
release.send(()).unwrap();
run_until(&mut coordinator, |state| state.login.is_none());
let observer = connect(&mut coordinator);
let lookup = request(
&coordinator,
&observer,
"operation.lookup",
None,
None,
None,
json!({"target_instance_id":coordinator.instance,"operation_id":"committed-login"}),
);
let reply = invoke(&mut coordinator, &observer, lookup);
assert!(reply["error"].is_null(), "{reply}");
assert_eq!(reply["payload"]["state"], "terminal");
assert_eq!(reply["payload"]["result"]["status"], "completed");
assert_eq!(reply["payload"]["result"]["cleanup_complete"], true);
assert!(reply["payload"]["error"].is_null());
assert_eq!(
crate::config::read_auth(&paths)
.unwrap()
.providers
.get("openai-codex"),
Some(&expected_credentials),
"disconnect after commit must not remove the stored login",
);
}
#[test]
fn standalone_writer_waits_for_daemon_cleanup_and_control_release() {
let temp = TempDir::new().unwrap();
std::fs::write(temp.path().join("evidence.txt"), "synthetic evidence").unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let (release, wait) = bounded(1);
let (cleanup, cleanup_wait) = bounded(1);
install_tool_worker(
&mut coordinator,
wait,
Arc::new(AtomicUsize::new(0)),
cleanup_wait,
);
let owner = connect(&mut coordinator);
let session = create(&mut coordinator, &owner);
let separate = create(&mut coordinator, &owner);
let disk = runtime.session_manager.open_existing(&session).unwrap();
let separate_disk = runtime.session_manager.open_existing(&separate).unwrap();
let grant = claim(&mut coordinator, &owner, &session);
assert!(
disk.try_frontend_writer().unwrap().is_none(),
"idle claim owns writer"
);
let separate_writer = separate_disk.try_frontend_writer().unwrap().unwrap();
let start = request(
&coordinator,
&owner,
"turn.start",
Some(&session),
Some(grant),
Some("writer-handoff"),
json!({"prompt":"Read the synthetic fixture"}),
);
assert_eq!(
invoke(&mut coordinator, &owner, start)["payload"]["status"],
"accepted"
);
assert_eq!(coordinator.actors[&session].phase, "running");
assert!(disk.try_frontend_writer().unwrap().is_none());
release.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "cleaning"
});
assert!(disk.try_frontend_writer().unwrap().is_none());
cleanup.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "idle"
});
assert!(
disk.try_frontend_writer().unwrap().is_none(),
"cleanup does not revoke control"
);
coordinator.disconnect(&owner);
assert!(!coordinator.actors.contains_key(&session));
let standalone_writer = disk.try_frontend_writer().unwrap().unwrap();
assert!(separate_disk.try_frontend_writer().unwrap().is_none());
drop(standalone_writer);
drop(separate_writer);
}