use super::*;
#[test]
fn ordinary_v2_cannot_create_executors_or_submit_injected_turns() {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
coordinator.unix_transport = true;
let connection = connect(&mut coordinator);
let create_executor = request(
&coordinator,
&connection,
"executor.create",
None,
None,
Some("executor"),
json!({}),
);
assert_eq!(
invoke(&mut coordinator, &connection, create_executor)["error"]["code"],
"unsupported_capability"
);
for method in ["executor.rotate", "executor.rotate_confirm"] {
let rotation = request(
&coordinator,
&connection,
method,
None,
None,
Some(method),
json!({}),
);
assert_eq!(
invoke(&mut coordinator, &connection, rotation)["error"]["code"],
"unsupported_capability"
);
}
let capabilities = request(
&coordinator,
&connection,
"capabilities",
None,
None,
None,
json!({}),
);
let capabilities = invoke(&mut coordinator, &connection, capabilities);
assert!(
capabilities["payload"]
.get("application_callbacks")
.is_none()
);
for method in ["executor.rotate", "executor.rotate_confirm"] {
assert!(
!capabilities["payload"]["operations"]
.as_array()
.unwrap()
.contains(&json!(method))
);
}
assert!(
!capabilities["payload"]["events"]
.as_array()
.unwrap()
.contains(&json!("turn.application_tool_call"))
);
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":"do not start","application":{}}),
);
assert_eq!(
invoke(&mut coordinator, &connection, start)["error"]["code"],
"unsupported_capability"
);
assert!(coordinator.turn_preparation.is_none());
assert_eq!(coordinator.actors[&session].phase, "idle");
}
#[test]
fn session_registry_retains_admitted_snapshots_and_releases_capacity_on_drop() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = super::super::super::application::ApplicationState::default();
let mut revision = Value::Null;
let mut resource_revision = Value::Null;
let mut pins = Vec::new();
for index in 0..256 {
let registered = state.registry.route("session", "session.resources.register", &json!({
"expected_revision":resource_revision,"resource":{"kind":"skill","name":"instructions","description":"Captured","source":{"kind":"inline","content":format!("revision {index}")}}
}), &runtime).unwrap();
revision = registered["registry_revision"].clone();
resource_revision = registered["entry"]["revision_id"].clone();
pins.push(state.capture_session("connection", "session", &json!({"expected_registry_revision":revision,"executor_id":null,"executor_generation":null,"tool_authorizations":[]})).unwrap());
}
let replacement = json!({"expected_revision":resource_revision,"resource":{"kind":"skill","name":"instructions","description":"Captured","source":{"kind":"inline","content":"replacement"}}});
assert_eq!(
state.registry.route(
"session",
"session.resources.register",
&replacement,
&runtime
),
Err(Code::LimitExceeded)
);
state
.registry
.route(
"session",
"session.resources.clear",
&json!({"expected_registry_revision":revision}),
&runtime,
)
.unwrap();
assert_eq!(pins[0].skill.as_deref(), Some("revision 0"));
let fresh = json!({"expected_revision":null,"resource":{"kind":"skill","name":"instructions","description":"Captured","source":{"kind":"inline","content":"fresh"}}});
assert_eq!(
state
.registry
.route("other", "session.resources.register", &fresh, &runtime),
Err(Code::LimitExceeded)
);
drop(pins);
let registered = state
.registry
.route("other", "session.resources.register", &fresh, &runtime)
.unwrap();
let captured = state.capture_session("connection", "other", &json!({"expected_registry_revision":registered["registry_revision"],"executor_id":null,"executor_generation":null,"tool_authorizations":[]})).unwrap();
assert_eq!(captured.skill.as_deref(), Some("fresh"));
}
#[test]
fn maximum_multi_tool_captures_share_instructions_until_callback_cleanup() {
use super::super::super::application::ApplicationState;
use std::sync::Mutex;
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let state = Arc::new(Mutex::new(ApplicationState::default()));
let mut locked = state.lock().unwrap();
let issued = locked
.route(
"connection",
"executor.create",
None,
&json!({}),
Instant::now(),
)
.unwrap();
locked
.route(
"connection",
"executor.confirm",
None,
&json!({
"executor_id":issued["executor_id"],"resume_secret":issued["resume_secret"]
}),
Instant::now(),
)
.unwrap();
let schema =
json!({"type":"object","properties":{},"required":[],"additionalProperties":false});
let mut revision = Value::Null;
let mut authorizations = Vec::new();
let mut skill_revision = Value::Null;
for index in 0..16 {
let registered = locked
.registry
.route(
"session",
"session.resources.register",
&json!({
"expected_revision":null,"resource":{"kind":"tool",
"name":format!("client_tool_{index}"),"description":"Captured callback",
"source":{"kind":"inline","content":json!({"input_schema":schema,"output_schema":schema}).to_string()}
}
}),
&runtime,
)
.unwrap();
revision = registered["registry_revision"].clone();
authorizations.push(authorization(®istered["entry"]));
}
let mut callbacks = Vec::new();
let mut instruction_allocations = Vec::new();
for index in 0..60 {
let registered = locked.registry.route("session", "session.resources.register", &json!({
"expected_revision":skill_revision,"resource":{"kind":"skill","name":"instructions",
"description":"Captured","source":{"kind":"inline","content":"x".repeat(8_000)}}
}), &runtime).unwrap();
revision = registered["registry_revision"].clone();
skill_revision = registered["entry"]["revision_id"].clone();
let captured = locked.capture_session("connection", "session", &json!({
"expected_registry_revision":revision,"executor_id":issued["executor_id"],"executor_generation":1,"tool_authorizations":authorizations
})).unwrap();
let skill = captured.skill.as_ref().unwrap();
assert_eq!(Arc::strong_count(skill), 18);
instruction_allocations.push(Arc::downgrade(skill));
drop(locked);
let tools = captured.tools(
Arc::clone(&state),
"session".into(),
format!("turn-{index}"),
);
assert_eq!(tools.definitions.len(), 16);
assert_eq!(Arc::strong_count(skill), 34);
callbacks.push(tools);
drop(captured);
locked = state.lock().unwrap();
}
locked
.registry
.route(
"session",
"session.resources.clear",
&json!({
"expected_registry_revision":revision
}),
&runtime,
)
.unwrap();
let replacement = json!({"expected_revision":null,"resource":{
"kind":"skill","name":"replacement","description":"","source":{"kind":"inline","content":"new"}
}});
for index in 0..4 {
locked
.registry
.route(
&format!("other-{index}"),
"session.resources.register",
&replacement,
&runtime,
)
.unwrap();
}
assert_eq!(
locked.registry.route(
"blocked",
"session.resources.register",
&replacement,
&runtime
),
Err(Code::LimitExceeded)
);
for index in 0..60 {
locked.finish_turn(&format!("turn-{index}"));
}
assert!(
instruction_allocations
.iter()
.all(|skill| skill.upgrade().is_some())
);
assert_eq!(
locked.registry.route(
"blocked",
"session.resources.register",
&replacement,
&runtime
),
Err(Code::LimitExceeded)
);
drop(callbacks);
assert!(
instruction_allocations
.iter()
.all(|skill| skill.upgrade().is_none())
);
locked
.registry
.route(
"blocked",
"session.resources.register",
&replacement,
&runtime,
)
.unwrap();
}
#[test]
fn session_registry_rejects_collisions_stale_revisions_and_combined_skill_overflow() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = super::super::super::application::ApplicationState::default();
for name in ["bash", "mcp_external", "Bad_name"] {
assert_eq!(
state.registry.route(
"session",
"session.resources.register",
&json!({"expected_revision":null,
"resource":{"kind":"skill","name":name,"description":"","source":{"kind":"inline","content":"instructions"}}}),
&runtime
),
Err(Code::InvalidPayload)
);
}
let registered = state
.registry
.route(
"session",
"session.resources.register",
&json!({"expected_revision":null,
"resource":{"kind":"skill","name":"first","description":"","source":{"kind":"inline","content":"x".repeat(8192)}}}),
&runtime,
)
.unwrap();
let second = json!({"expected_revision":null,
"resource":{"kind":"skill","name":"second","description":"","source":{"kind":"inline","content":"y".repeat(8192)}}});
assert_eq!(
state
.registry
.route("session", "session.resources.register", &second, &runtime),
Err(Code::LimitExceeded)
);
assert_eq!(
state.registry.route(
"session",
"session.resources.clear",
&json!({"expected_registry_revision":"stale"}),
&runtime
),
Err(Code::Conflict)
);
let listed = state
.registry
.route("session", "session.resources.list", &json!({}), &runtime)
.unwrap();
assert_eq!(listed["entries"].as_array().unwrap().len(), 1);
assert_eq!(listed["registry_revision"], registered["registry_revision"]);
assert!(!listed.to_string().contains(&"x".repeat(100)));
}
#[test]
fn disabled_and_discovered_skills_reserve_names_at_registration_and_admission() {
let temp = TempDir::new().unwrap();
let mut runtime = (*runtime(&temp)).clone();
let mut state = super::super::super::application::ApplicationState::default();
let resource = |name: &str| {
json!({"expected_revision":null,
"resource":{"kind":"skill","name":name,"description":"","source":{"kind":"inline","content":"instructions"}}})
};
let registered = state
.registry
.route(
"session",
"session.resources.register",
&resource("later_disabled"),
&runtime,
)
.unwrap();
let captured = state
.capture_session(
"connection",
"session",
&json!({
"expected_registry_revision":registered["registry_revision"],
"executor_id":null,"executor_generation":null,"tool_authorizations":[]}),
)
.unwrap();
runtime
.discovered_skill_names
.push("Discovered_Skill".into());
runtime.settings.skills.disabled = vec!["Later_Disabled".into(), "Absent_Skill".into()];
for name in ["discovered_skill", "later_disabled", "absent_skill"] {
assert_eq!(
state.registry.route(
"other",
"session.resources.register",
&resource(name),
&runtime
),
Err(Code::InvalidPayload)
);
}
assert_eq!(
captured.validate_catalog(&runtime),
Err(Code::InvalidPayload)
);
}
#[test]
fn application_resource_profile_rejects_retired_tokens_and_non_unix_transports() {
for (unix, profiles, accepted) in [
(false, json!(["application_resources_v1"]), false),
(true, json!(["application_callbacks_inline_v1"]), false),
(true, json!(["application_callbacks_session_v1"]), false),
(
true,
json!([
"application_resources_v1",
"application_callbacks_inline_v1"
]),
false,
),
(true, json!(["application_resources_v1"]), true),
] {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
coordinator.unix_transport = unix;
let connection = coordinator.connect(Instant::now()).unwrap();
let init = request(
&coordinator,
&connection,
"initialize",
None,
None,
None,
json!({"supported_protocol_versions":[2],"requested_capabilities":profiles}),
);
let response = invoke(&mut coordinator, &connection, init);
assert_eq!(response["error"].is_null(), accepted, "{response}");
if accepted {
assert_eq!(
response["payload"]["capabilities"]["negotiated_capabilities"],
json!(["application_resources_v1"])
);
let capabilities_request = request(
&coordinator,
&connection,
"capabilities",
None,
None,
None,
json!({}),
);
let capabilities = invoke(&mut coordinator, &connection, capabilities_request);
assert_eq!(capabilities["payload"], response["payload"]["capabilities"]);
assert_eq!(
response["payload"]["capabilities"]["application_callbacks"]["max_tools"],
16
);
assert_eq!(
response["payload"]["capabilities"]["application_callbacks"]["registry"],
true
);
} else {
assert_eq!(response["error"]["code"], "unsupported_capability");
assert!(!coordinator.connections[&connection].initialized);
}
}
}
#[test]
fn session_registry_global_resource_count_includes_removed_pins() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = super::super::super::application::ApplicationState::default();
let mut revision = Value::Null;
let mut resource_revision = Value::Null;
for index in 0..16 {
let registered = state.registry.route("session", "session.resources.register", &json!({"expected_revision":null,
"resource":{"kind":"skill","name":format!("skill_{index}"),"description":"","source":{"kind":"inline","content":"text"}}}), &runtime).unwrap();
revision = registered["registry_revision"].clone();
if index == 0 {
resource_revision = registered["entry"]["revision_id"].clone();
}
}
let mut pins = Vec::new();
for index in 0..64 {
pins.push(state.capture_session("connection", "session", &json!({"expected_registry_revision":revision,"executor_id":null,"executor_generation":null,"tool_authorizations":[]})).unwrap());
let result = state.registry.route("session", "session.resources.register", &json!({"expected_revision":resource_revision,
"resource":{"kind":"skill","name":"skill_0","description":"","source":{"kind":"inline","content":format!("revision {index}")}}}), &runtime);
if index == 63 {
assert_eq!(result, Err(Code::LimitExceeded));
} else {
let registered = result.unwrap();
revision = registered["registry_revision"].clone();
resource_revision = registered["entry"]["revision_id"].clone();
}
}
state
.registry
.route(
"session",
"session.resources.clear",
&json!({"expected_registry_revision":revision}),
&runtime,
)
.unwrap();
let fresh = json!({"expected_revision":null,"resource":{"kind":"skill","name":"fresh","description":"","source":{"kind":"inline","content":"new"}}});
assert_eq!(
state
.registry
.route("other", "session.resources.register", &fresh, &runtime),
Err(Code::LimitExceeded)
);
drop(pins);
assert!(
state
.registry
.route("other", "session.resources.register", &fresh, &runtime)
.is_ok()
);
}
#[test]
fn registry_idle_expiry_uses_exact_deadline_and_reclaim_resets_timer() {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
claim(&mut coordinator, &connection, &session);
coordinator.application.lock().unwrap().registry.route(&session, "session.resources.register", &json!({
"expected_revision":null,"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"original"}}
}), &coordinator.runtime).unwrap();
let now = Instant::now();
let lifetime = Duration::from_secs(900);
coordinator.tick(now + lifetime);
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
coordinator.detach(&session, now + lifetime);
assert!(!coordinator.actors.contains_key(&session));
let lease = coordinator
.runtime
.session_manager
.open_existing(&session)
.unwrap();
let writer = lease
.try_daemon_writer()
.unwrap()
.expect("registry holds no writer lease");
drop(writer);
let reclaim = request(
&coordinator,
&connection,
"session.claim",
Some(&session),
None,
Some("reclaim"),
json!({}),
);
assert!(
coordinator
.admit(
&connection,
&reclaim,
now + lifetime * 2 - Duration::from_nanos(1)
)
.unwrap()
.is_ok()
);
coordinator.tick(now + lifetime * 3);
let idle = now + lifetime * 3;
coordinator.detach(&session, idle);
coordinator.tick(idle + lifetime - Duration::from_nanos(1));
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
let reclaim_at_expiry = request(
&coordinator,
&connection,
"session.claim",
Some(&session),
None,
Some("reclaim_at_expiry"),
json!({}),
);
assert!(
coordinator
.admit(&connection, &reclaim_at_expiry, idle + lifetime)
.unwrap()
.is_ok()
);
assert!(
!coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
}
#[test]
fn registry_expiry_defers_for_old_snapshot_work_until_cleanup() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = super::super::super::application::ApplicationState::default();
let registered = state.registry.route("session", "session.resources.register", &json!({
"expected_revision":null,"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"original"}}
}), &runtime).unwrap();
let captured = state.capture_session("connection", "session", &json!({
"expected_registry_revision":registered["registry_revision"],"executor_id":null,"executor_generation":null,"tool_authorizations":[]
})).unwrap();
state.registry.route("session", "session.resources.register", &json!({
"expected_revision":registered["entry"]["revision_id"],"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"replacement"}}
}), &runtime).unwrap();
let now = Instant::now();
let lifetime = Duration::from_secs(900);
state.registry.expire_idle(now, |_| true);
state.registry.expire_idle(now + lifetime * 2, |_| true);
assert!(state.registry.nonempty("session"));
assert_eq!(captured.skill.as_deref(), Some("original"));
drop(captured);
state.registry.expire_idle(now + lifetime * 2, |_| false);
state
.registry
.expire_idle(now + lifetime * 3 - Duration::from_nanos(1), |_| false);
assert!(state.registry.nonempty("session"));
state.registry.expire_idle(now + lifetime * 3, |_| false);
assert!(!state.registry.nonempty("session"));
}
#[test]
fn coordinator_defers_registry_expiry_through_disconnected_preparation() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let connection = connect(&mut coordinator);
coordinator
.connections
.get_mut(&connection)
.unwrap()
.application_profile = true;
coordinator
.connections
.get_mut(&connection)
.unwrap()
.application_profile = true;
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
let registered = coordinator.application.lock().unwrap().registry.route(&session, "session.resources.register", &json!({
"expected_revision":null,"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"original"}}
}), &runtime).unwrap();
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","application":{"expected_registry_revision":registered["registry_revision"],"executor_id":null,"executor_generation":null,"tool_authorizations":[]}
}),
);
coordinator
.submit(
&connection,
&serde_json::to_vec(&start).unwrap(),
Instant::now(),
)
.unwrap();
assert!(coordinator.turn_preparation.is_some());
coordinator.disconnect(&connection);
let now = Instant::now();
coordinator.tick(now + Duration::from_secs(1800));
assert!(coordinator.turn_preparation.is_some());
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
drop(lock);
run_until(&mut coordinator, |state| state.turn_preparation.is_none());
let settled = Instant::now();
coordinator.tick(settled + Duration::from_secs(899));
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
coordinator.tick(settled + Duration::from_secs(900));
assert!(
!coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
}
#[test]
fn coordinator_defers_registry_expiry_through_disconnected_worker_cleanup() {
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 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("active"),
json!({"prompt":"Read the fixture"}),
);
assert_eq!(
invoke(&mut coordinator, &connection, start)["payload"]["status"],
"accepted"
);
coordinator.application.lock().unwrap().registry.route(&session, "session.resources.register", &json!({
"expected_revision":null,"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"later turn"}}
}), &runtime).unwrap();
coordinator.disconnect(&connection);
coordinator.tick(Instant::now() + Duration::from_secs(1800));
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
release.send(()).unwrap();
run_until(&mut coordinator, |state| {
state.actors[&session].phase == "cleaning"
});
coordinator.tick(Instant::now() + Duration::from_secs(3600));
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
cleanup.send(()).unwrap();
run_until(&mut coordinator, |state| {
!state.actors.contains_key(&session)
});
let settled = Instant::now();
coordinator.tick(settled + Duration::from_secs(600));
assert!(
!coordinator.terminals.contains_key(&session),
"registry retention must not extend outcome evidence"
);
assert!(
coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
coordinator.tick(settled + Duration::from_secs(900));
assert!(
!coordinator
.application
.lock()
.unwrap()
.registry
.nonempty(&session)
);
}
#[test]
fn evidence_reads_require_profile_and_controller_and_do_not_extend_outcomes() {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
let connection = connect(&mut coordinator);
let session = create(&mut coordinator, &connection);
let control = claim(&mut coordinator, &connection, &session);
let read = request(
&coordinator,
&connection,
"turn.application_evidence",
Some(&session),
Some(control.clone()),
None,
json!({"turn_id":"turn"}),
);
assert_eq!(
invoke(&mut coordinator, &connection, read.clone())["error"]["code"],
"unsupported_capability"
);
coordinator
.connections
.get_mut(&connection)
.unwrap()
.application_profile = true;
let now = Instant::now();
let start = request(
&coordinator,
&connection,
"turn.start",
Some(&session),
Some(control),
Some("evidence-operation"),
json!({"prompt":"test"}),
);
coordinator.operations.reserve(&start, now).unwrap();
coordinator.operations.update(
"evidence-operation",
"terminal",
Value::Null,
Value::Null,
now,
);
{
let mut state = coordinator.application.lock().unwrap();
state
.reserve_evidence(
"evidence-operation",
&session,
&json!({"profile":"application_resources_v1","resources":[],"tools":[]}),
)
.unwrap();
state.bind_evidence("evidence-operation", Some("turn"));
}
assert!(
coordinator
.admit(&connection, &read, now + Duration::from_secs(599))
.unwrap()
.is_ok()
);
assert_eq!(
coordinator
.admit(&connection, &read, now + Duration::from_secs(600))
.unwrap(),
Err(Code::UnknownTurn)
);
}
#[test]
fn source_capture_revalidates_control_deadline_and_registry_before_publication() {
for failure in ["control", "revision", "deadline"] {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut coordinator = Coordinator::new(Arc::clone(&runtime)).unwrap();
let connection = connect(&mut coordinator);
coordinator
.connections
.get_mut(&connection)
.unwrap()
.application_profile = true;
let session = create(&mut coordinator, &connection);
let grant = claim(&mut coordinator, &connection, &session);
let registration = request(
&coordinator,
&connection,
"session.resources.register",
Some(&session),
Some(grant),
Some("capture"),
json!({"expected_revision":null,
"resource":{"kind":"skill","name":"instructions","description":"captured","source":{"kind":"inline","content":"staged instructions"}}}),
);
coordinator
.submit(
&connection,
&serde_json::to_vec(®istration).unwrap(),
Instant::now(),
)
.unwrap();
assert!(coordinator.source_capture.is_some());
match failure {
"control" => coordinator.disconnect(&connection),
"revision" => {
coordinator.application.lock().unwrap().registry.route(&session, "session.resources.register", &json!({"expected_revision":null,
"resource":{"kind":"skill","name":"instructions","description":"","source":{"kind":"inline","content":"published instead"}}}), &runtime).unwrap();
}
"deadline" => coordinator.source_capture.as_mut().unwrap().deadline = Instant::now(),
_ => unreachable!(),
}
run_until(&mut coordinator, |state| state.source_capture.is_none());
let listed = coordinator
.application
.lock()
.unwrap()
.registry
.route(&session, "session.resources.list", &json!({}), &runtime)
.unwrap();
assert!(
listed["entries"]
.as_array()
.unwrap()
.iter()
.all(|entry| failure == "revision" || entry["name"] != "instructions")
);
let outcome =
coordinator
.operations
.lookup(&coordinator.instance, &coordinator.instance, "capture");
assert_eq!(outcome["state"], "rejected", "{outcome}");
}
}
#[test]
fn source_preparation_has_no_admission_attribution_and_skill_only_acceptance_has_no_executor() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = crate::service::application::ApplicationState::default();
let payload = json!({"expected_revision":null,"resource":{
"kind":"skill","name":"instructions","description":"captured",
"source":{"kind":"inline","content":"Use captured instructions"}}});
let prepared = crate::service::application::prepare_source(&payload, &runtime).unwrap();
let published = state
.registry
.publish_source(
"session",
prepared,
&runtime,
json!({"operation_id":"registration"}),
)
.unwrap();
let captured = state
.capture_session(
"connection",
"session",
&json!({
"expected_registry_revision":published["registry_revision"],
"executor_id":null,"executor_generation":null,"tool_authorizations":[]
}),
)
.unwrap();
assert!(captured.manifest().get("accepted_controller").is_none());
let controller = json!({"instance_id":"instance","connection_id":"connection",
"operation_id":"start","grant_generation":3});
let accepted = captured.accept_controller(controller.clone());
let manifest = accepted.manifest();
assert_eq!(manifest["accepted_controller"], controller);
assert_eq!(manifest["accepted_executor"], Value::Null);
state
.registry
.route(
"session",
"session.resources.clear",
&json!({
"expected_registry_revision":published["registry_revision"]}),
&runtime,
)
.unwrap();
assert_eq!(accepted.manifest(), manifest);
}
fn authorization(entry: &Value) -> Value {
json!({"resource_id":entry["resource_id"],"revision_id":entry["revision_id"],"revision_sha256":entry["revision_sha256"]})
}
#[test]
fn resource_mutations_compare_resource_revisions_and_keep_empty_registry_identity() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = crate::service::application::ApplicationState::default();
let empty = state
.registry
.route("session", "session.resources.list", &json!({}), &runtime)
.unwrap();
assert!(
empty["registry_revision"]
.as_str()
.is_some_and(wire::valid_id)
);
assert_eq!(
empty,
state
.registry
.route("session", "session.resources.list", &json!({}), &runtime)
.unwrap()
);
let resource = |expected: Value, kind: &str, name: &str| {
let schema =
json!({"type":"object","properties":{},"required":[],"additionalProperties":false});
let content = if kind == "skill" {
"captured instructions".into()
} else {
json!({"input_schema":schema,"output_schema":schema}).to_string()
};
json!({"expected_revision":expected,"resource":{"kind":kind,"name":name,"description":"captured","source":{"kind":"inline","content":content}}})
};
let first = state
.registry
.route(
"session",
"session.resources.register",
&resource(Value::Null, "tool", "alpha"),
&runtime,
)
.unwrap();
let skill = state
.registry
.route(
"session",
"session.resources.register",
&resource(Value::Null, "skill", "zulu"),
&runtime,
)
.unwrap();
let listed = state
.registry
.route("session", "session.resources.list", &json!({}), &runtime)
.unwrap();
assert_eq!(listed["entries"], json!([skill["entry"], first["entry"]]));
assert_ne!(empty["registry_revision"], first["registry_revision"]);
assert_eq!(
state.registry.route(
"session",
"session.resources.register",
&resource(Value::Null, "tool", "alpha"),
&runtime
),
Err(Code::Conflict)
);
assert_eq!(
state.registry.route(
"session",
"session.resources.register",
&resource(first["entry"]["revision_id"].clone(), "skill", "alpha"),
&runtime
),
Err(Code::Conflict)
);
let replacement = state
.registry
.route(
"session",
"session.resources.register",
&resource(first["entry"]["revision_id"].clone(), "tool", "alpha"),
&runtime,
)
.unwrap();
assert_eq!(
replacement["entry"]["resource_id"],
first["entry"]["resource_id"]
);
assert_ne!(
replacement["entry"]["revision_id"],
first["entry"]["revision_id"]
);
assert_eq!(
replacement["entry"]["revision_sha256"],
first["entry"]["revision_sha256"]
);
assert_eq!(state.registry.route("session", "session.resources.unregister", &json!({"resource_id":first["entry"]["resource_id"],"expected_revision":first["entry"]["revision_id"]}), &runtime), Err(Code::Conflict));
let removed = state.registry.route("session", "session.resources.unregister", &json!({"resource_id":replacement["entry"]["resource_id"],"expected_revision":replacement["entry"]["revision_id"]}), &runtime).unwrap();
assert_eq!(
removed["removed_revision"],
replacement["entry"]["revision_id"]
);
assert_eq!(removed["resource_id"], replacement["entry"]["resource_id"]);
assert_eq!(state.registry.route("session", "session.resources.unregister", &json!({"resource_id":replacement["entry"]["resource_id"],"expected_revision":replacement["entry"]["revision_id"]}), &runtime), Err(Code::ResourceNotFound));
let recreated = state
.registry
.route(
"session",
"session.resources.register",
&resource(Value::Null, "tool", "alpha"),
&runtime,
)
.unwrap();
assert_ne!(
recreated["entry"]["resource_id"],
replacement["entry"]["resource_id"]
);
let cleared = state
.registry
.route(
"session",
"session.resources.clear",
&json!({"expected_registry_revision":recreated["registry_revision"]}),
&runtime,
)
.unwrap();
assert_eq!(cleared["removed_count"], 2);
assert_ne!(cleared["registry_revision"], recreated["registry_revision"]);
let cleared_again = state
.registry
.route(
"session",
"session.resources.clear",
&json!({"expected_registry_revision":cleared["registry_revision"]}),
&runtime,
)
.unwrap();
assert_eq!(cleared_again["removed_count"], 0);
assert_eq!(
cleared_again["registry_revision"],
cleared["registry_revision"]
);
let selection = json!({"expected_registry_revision":cleared["registry_revision"],"executor_id":null,"executor_generation":null,"tool_authorizations":[]});
assert!(
state
.capture_session("connection", "session", &selection)
.is_ok()
);
let mut missing = selection;
missing
.as_object_mut()
.unwrap()
.remove("tool_authorizations");
assert!(matches!(
state.capture_session("connection", "session", &missing),
Err(Code::InvalidPayload)
));
}
#[test]
fn empty_registry_table_is_bounded_and_unclaimed_empty_entries_are_removed() {
let temp = TempDir::new().unwrap();
let runtime = runtime(&temp);
let mut state = crate::service::application::ApplicationState::default();
for index in 0..256 {
state
.registry
.route(
&format!("session-{index}"),
"session.resources.list",
&json!({}),
&runtime,
)
.unwrap();
}
assert_eq!(
state
.registry
.route("overflow", "session.resources.list", &json!({}), &runtime),
Err(Code::LimitExceeded)
);
state
.registry
.expire_idle(Instant::now(), |session| session != "session-0");
assert!(
state
.registry
.route("overflow", "session.resources.list", &json!({}), &runtime)
.is_ok()
);
}