use super::*;
#[test]
fn v2_turn_payload_rejects_options_and_unknown_fields_before_execution() {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
coordinator
.execution
.set_turn_worker(Arc::new(|_| panic!("invalid turn executed")));
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
for payload in [
json!({"prompt":"hello", "options":{}}),
json!({"prompt":"hello", "extra":true}),
] {
let start = request(
&coordinator,
&connection,
"turn.start",
Some(&session),
Some(grant.clone()),
Some(&uuid::Uuid::new_v4().to_string()),
payload,
);
assert_eq!(
invoke(&mut coordinator, &connection, start)["error"]["code"],
"invalid_payload"
);
assert!(coordinator.turn_preparation.is_none());
}
assert!(
serde_json::from_value::<crate::service::protocol::TurnStartParams>(
json!({"prompt":"hello", "options":{}})
)
.is_ok()
);
}
#[test]
fn held_settings_lock_allows_cancel_disconnect_and_terminal_drain() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let (terminal, finish) = bounded(1);
coordinator.execution.set_turn_worker(Arc::new(move |job| {
finish.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 running = create(&mut coordinator, &owner);
let waiting = create(&mut coordinator, &other);
let running_grant = claim(&mut coordinator, &owner, &running);
let waiting_grant = claim(&mut coordinator, &other, &waiting);
let start = request(
&coordinator,
&owner,
"turn.start",
Some(&running),
Some(running_grant.clone()),
Some("running"),
json!({"prompt":"hello"}),
);
let accepted = invoke(&mut coordinator, &owner, start);
let turn_id = accepted["payload"]["turn_id"].clone();
coordinator
.connections
.get_mut(&owner)
.unwrap()
.queue
.clear();
coordinator
.connections
.get_mut(&owner)
.unwrap()
.queued_bytes = 0;
let lock =
crate::persistence::CrossProcessFileLock::acquire(&runtime.config.paths.settings_file)
.unwrap();
let start = request(
&coordinator,
&other,
"turn.start",
Some(&waiting),
Some(waiting_grant),
Some("waiting"),
json!({"prompt":"hello"}),
);
coordinator
.submit(&other, &serde_json::to_vec(&start).unwrap(), Instant::now())
.unwrap();
assert!(coordinator.turn_preparation.is_some());
let mut duplicate = start.clone();
duplicate.request_id = "duplicate".into();
assert_eq!(
invoke(&mut coordinator, &other, duplicate)["error"]["code"],
"operation_already_known"
);
let cancel = request(
&coordinator,
&owner,
"turn.cancel",
Some(&running),
Some(running_grant),
Some("cancel"),
json!({"turn_id":turn_id}),
);
assert!(invoke(&mut coordinator, &owner, cancel)["error"].is_null());
coordinator.disconnect(&other);
terminal.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&running].phase == "idle"
});
assert!(coordinator.turn_preparation.is_some());
drop(lock);
run_until(&mut coordinator, |state| state.turn_preparation.is_none());
assert_eq!(
coordinator
.operations
.lookup(&coordinator.instance, &coordinator.instance, "waiting")["error"]["code"],
"stale_connection"
);
}
#[test]
fn prepared_turn_rechecks_grant_deadline_and_settings_failure() {
for failure in ["grant", "deadline", "settings"] {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
coordinator
.execution
.set_turn_worker(Arc::new(|_| panic!("rejected turn executed")));
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
let lock =
crate::persistence::CrossProcessFileLock::acquire(&runtime.config.paths.settings_file)
.unwrap();
let start = request(
&coordinator,
&connection,
"turn.start",
Some(&session),
Some(grant),
Some("prepare"),
json!({"prompt":"hello"}),
);
coordinator
.submit(
&connection,
&serde_json::to_vec(&start).unwrap(),
Instant::now(),
)
.unwrap();
let expected = match failure {
"grant" => {
coordinator.detach(&session);
claim(&mut coordinator, &connection, &session);
"stale_grant"
}
"deadline" => {
coordinator.turn_preparation.as_mut().unwrap().deadline = Instant::now();
"request_timeout"
}
_ => {
std::fs::write(&runtime.config.paths.settings_file, "{secret-marker").unwrap();
"internal_error"
}
};
drop(lock);
run_until(&mut coordinator, |state| state.turn_preparation.is_none());
let reply = coordinator.connections[&connection].queue.back().unwrap();
assert_eq!(
serde_json::from_str::<Value>(reply).unwrap()["error"]["code"],
expected
);
assert!(!reply.contains("secret-marker"));
assert!(
coordinator
.actors
.values()
.all(|actor| actor.phase == "idle")
);
}
}
#[test]
fn command_loop_serves_disconnect_while_settings_capture_waits() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
let start = request(
&coordinator,
&connection,
"turn.start",
Some(&session),
Some(grant),
Some("blocked"),
json!({"prompt":"hello"}),
);
let lock =
crate::persistence::CrossProcessFileLock::acquire(&runtime.config.paths.settings_file)
.unwrap();
let service = PersistentService::spawn(coordinator).unwrap();
let (release, wait) = bounded(1);
let lock_owner = std::thread::spawn(move || {
let _ = wait.recv_timeout(Duration::from_secs(3));
drop(lock);
});
let began = Instant::now();
service
.submit(&connection, &serde_json::to_vec(&start).unwrap())
.unwrap();
service.disconnect(&connection).unwrap();
let responsive = began.elapsed() < Duration::from_secs(1);
let _ = release.send(());
lock_owner.join().unwrap();
drop(service);
assert!(responsive, "settings preparation blocked the command loop");
}
#[test]
fn configuration_writes_and_turn_preparation_cannot_overtake_each_other() {
for settings_first in [false, true] {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
let start = request(
&coordinator,
&connection,
"turn.start",
Some(&session),
Some(grant),
Some("turn"),
json!({"prompt":"hello"}),
);
let settings = request(
&coordinator,
&connection,
"config.set",
None,
None,
Some("settings"),
json!({"scope":"global", "fast":true}),
);
let lock =
crate::persistence::CrossProcessFileLock::acquire(&runtime.config.paths.settings_file)
.unwrap();
let (first, second) = if settings_first {
(settings, start)
} else {
(start, settings)
};
coordinator
.submit(
&connection,
&serde_json::to_vec(&first).unwrap(),
Instant::now(),
)
.unwrap();
assert_eq!(
invoke(&mut coordinator, &connection, second)["error"]["code"],
"configuration_busy"
);
coordinator.disconnect(&connection);
drop(lock);
run_until(&mut coordinator, |state| {
state.turn_preparation.is_none() && state.pending.is_empty()
});
assert_eq!(
crate::config::read_settings(&runtime.config.paths)
.unwrap()
.fast
.enabled,
settings_first
);
}
}