use super::*;
#[test]
fn held_auth_lock_preserves_other_sessions_disconnect_and_accepted_logout() {
for (method, cross_process) in ["status", "auth.status", "auth.logout", "auth.login.start"]
.into_iter()
.flat_map(|method| [false, true].map(|cross_process| (method, cross_process)))
{
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
coordinator
.execution
.set_login_worker(Arc::new(|_| panic!("disconnected login ran")));
let owner = connect(&mut coordinator);
let lock =
crate::persistence::in_process_file_lock(&runtime.config.paths.auth_file, "auth")
.unwrap();
let held = (!cross_process).then(|| lock.lock().unwrap());
let file_lock = cross_process.then(|| {
crate::persistence::CrossProcessFileLock::acquire(&runtime.config.paths.auth_file)
.unwrap()
});
let work = request(
&coordinator,
&owner,
method,
None,
None,
matches!(method, "auth.logout" | "auth.login.start").then_some("auth-work"),
match method {
"auth.logout" => json!({"provider_id":"openai-codex", "confirmed":true}),
"auth.login.start" => json!({"provider_id":"openai-codex"}),
_ => json!({}),
},
);
let (completed, result) = bounded(1);
let worker = std::thread::spawn(move || {
coordinator
.submit(&owner, &serde_json::to_vec(&work).unwrap(), Instant::now())
.unwrap();
if method == "auth.logout" {
assert_eq!(
coordinator.operations.lookup(
&coordinator.instance,
&coordinator.instance,
"auth-work"
)["state"],
"accepted"
);
}
coordinator.disconnect(&owner);
let other = connect(&mut coordinator);
let session = create(&mut coordinator, &other);
claim(&mut coordinator, &other, &session);
completed.send(coordinator).unwrap();
});
let progress = result.recv_timeout(Duration::from_secs(2));
drop(held);
drop(file_lock);
worker.join().unwrap();
let mut coordinator = progress.expect("auth lock blocked disconnect or another session");
run_until(&mut coordinator, |state| {
state.auth_work.is_none() && state.login.is_none()
});
if method == "auth.logout" {
let outcome = coordinator.operations.lookup(
&coordinator.instance,
&coordinator.instance,
"auth-work",
);
assert_eq!(outcome["state"], "terminal", "{outcome}");
assert_eq!(outcome["result"]["status"], "completed", "{outcome}");
assert_eq!(
crate::config::read_auth_store(&runtime.config.paths)
.unwrap()
.provider_generation("openai-codex"),
1
);
}
if method == "auth.login.start" {
let outcome = coordinator.operations.lookup(
&coordinator.instance,
&coordinator.instance,
"auth-work",
);
assert_eq!(outcome["result"]["status"], "cancelled", "{outcome}");
}
}
}
#[test]
fn held_auth_lock_does_not_block_login_cancel_or_turn_terminal_cleanup() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let (started, ready) = bounded(1);
coordinator.execution.set_login_worker(Arc::new(move |job| {
started.send(()).unwrap();
while !job.cancel.load(Ordering::Acquire) {
std::thread::sleep(Duration::from_millis(1));
}
"cancelled"
}));
let (finish, terminal) = bounded(1);
coordinator.execution.set_turn_worker(Arc::new(move |job| {
terminal.recv().unwrap();
job.sender
.send(TurnWorkerMessage::Terminal {
turn_id: job.turn_id.clone(),
event: ServiceEvent::turn_terminal(
job.request_id,
job.session.id().into(),
job.turn_id,
TurnTerminalStatus::Cancelled,
String::new(),
),
})
.unwrap();
}));
let owner = connect(&mut coordinator);
let other = connect(&mut coordinator);
let session = create(&mut coordinator, &other);
let grant = claim(&mut coordinator, &other, &session);
let start = request(
&coordinator,
&other,
"turn.start",
Some(&session),
Some(grant.clone()),
Some("turn"),
json!({"prompt":"hello"}),
);
let turn = invoke(&mut coordinator, &other, start);
let login = request(
&coordinator,
&owner,
"auth.login.start",
None,
None,
Some("login"),
json!({"provider_id":"openai-codex"}),
);
let login = invoke(&mut coordinator, &owner, login);
ready.recv_timeout(Duration::from_secs(2)).unwrap();
let lock =
crate::persistence::in_process_file_lock(&runtime.config.paths.auth_file, "auth").unwrap();
let held = lock.lock().unwrap();
let (completed, result) = bounded(1);
let worker = std::thread::spawn(move || {
let cancel = request(
&coordinator,
&owner,
"auth.login.cancel",
None,
None,
Some("cancel-login"),
json!({"login_id":login["payload"]["login_id"]}),
);
assert!(invoke(&mut coordinator, &owner, cancel)["error"].is_null());
let logout = request(
&coordinator,
&owner,
"auth.logout",
None,
None,
Some("logout"),
json!({"provider_id":"openai-codex", "confirmed":true}),
);
coordinator
.submit(
&owner,
&serde_json::to_vec(&logout).unwrap(),
Instant::now(),
)
.unwrap();
coordinator.disconnect(&owner);
let cancel = request(
&coordinator,
&other,
"turn.cancel",
Some(&session),
Some(grant),
Some("cancel-turn"),
json!({"turn_id":turn["payload"]["turn_id"]}),
);
coordinator
.submit(
&other,
&serde_json::to_vec(&cancel).unwrap(),
Instant::now(),
)
.unwrap();
finish.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "idle"
});
assert!(coordinator.login.is_some());
assert!(coordinator.auth_work.is_some());
completed.send(coordinator).unwrap();
});
let progress = result.recv_timeout(Duration::from_secs(2));
drop(held);
worker.join().unwrap();
let mut coordinator = progress.expect("auth cancellation blocked turn cleanup");
run_until(&mut coordinator, |state| {
state.login.is_none() && state.auth_work.is_none()
});
let outcome =
coordinator
.operations
.lookup(&coordinator.instance, &coordinator.instance, "logout");
assert_eq!(outcome["result"]["status"], "completed", "{outcome}");
assert_eq!(
crate::config::read_auth_store(&runtime.config.paths)
.unwrap()
.provider_generation("openai-codex"),
1
);
}