use super::*;
#[test]
fn accepted_manifest_survives_projection_takeover_and_expires_with_terminal() {
let temp = TempDir::new().unwrap();
let mut coordinator = Coordinator::new(runtime(&temp)).unwrap();
let first = connect(&mut coordinator);
coordinator
.connections
.get_mut(&first)
.unwrap()
.application_profile = true;
let session = create(&mut coordinator, &first);
let control = claim(&mut coordinator, &first, &session);
let now = Instant::now();
let mut start = request(
&coordinator,
&first,
"turn.start",
Some(&session),
Some(control),
Some("manifest-turn"),
json!({"prompt":"test"}),
);
coordinator.operations.reserve(&start, now).unwrap();
let manifest = {
let mut state = coordinator.application.lock().unwrap();
let listed = state
.registry
.route(
&session,
"session.resources.list",
&json!({}),
&coordinator.runtime,
)
.unwrap();
let captured = state
.capture_session(
&first,
&session,
&json!({
"expected_registry_revision":listed["registry_revision"],
"executor_id":null,"executor_generation":null,"tool_authorizations":[]
}),
)
.unwrap()
.accept_controller(json!({"connection_id":first}));
let manifest = captured.manifest();
assert_eq!(manifest["resources"], json!([]));
assert_eq!(manifest["tools"], json!([]));
assert_eq!(manifest.get("accepted_executor"), Some(&Value::Null));
manifest
};
start.payload = json!({"application_manifest":manifest});
let pending = PendingRequest {
connection: first.clone(),
request: start,
accepted: false,
};
coordinator.accept_execution_result(
&pending,
&Ok(json!({"turn_id":"turn","application_manifest":manifest})),
now,
);
let lookup = request(
&coordinator,
&first,
"operation.lookup",
None,
None,
None,
json!({"target_instance_id":coordinator.instance,"operation_id":"manifest-turn"}),
);
assert_eq!(
invoke(&mut coordinator, &first, lookup)["payload"]["result"]["application_manifest"],
manifest
);
coordinator.disconnect(&first);
let second = connect(&mut coordinator);
coordinator
.connections
.get_mut(&second)
.unwrap()
.application_profile = true;
let claim_request = request(
&coordinator,
&second,
"session.claim",
Some(&session),
None,
Some("manifest-claim"),
json!({}),
);
let claimed = invoke(&mut coordinator, &second, claim_request);
assert_eq!(
claimed["payload"]["snapshot"]["turn"]["application_manifest"],
manifest
);
let mut event = ServiceEvent::turn_terminal(
"request".into(),
session.clone(),
"turn".into(),
TurnTerminalStatus::Completed,
"accepted text".into(),
);
event.payload["sequence"] = json!(1);
coordinator.turn_event(&event, now);
let terminal: Value =
serde_json::from_str(coordinator.connections[&second].queue.back().unwrap()).unwrap();
assert_eq!(terminal["payload"]["application_manifest"], manifest);
assert_eq!(
coordinator.terminals[&session].payload["application_manifest"],
manifest
);
let client = coordinator.connections.get_mut(&second).unwrap();
client.queue.clear();
client.queued_bytes = 0;
for selected in [false, true] {
coordinator
.connections
.get_mut(&second)
.unwrap()
.application_profile = selected;
let lookup = request(
&coordinator,
&second,
"operation.lookup",
None,
None,
None,
json!({"target_instance_id":coordinator.instance,"operation_id":"manifest-turn"}),
);
let reply = invoke(&mut coordinator, &second, lookup);
assert_eq!(
reply["payload"]["result"].get("application_manifest"),
selected.then_some(&manifest)
);
}
coordinator.tick(now + Duration::from_secs(599));
assert_eq!(
coordinator.terminals[&session].payload["application_manifest"],
manifest
);
coordinator.tick(now + Duration::from_secs(600));
assert!(!coordinator.terminals.contains_key(&session));
assert!(coordinator.actors[&session].terminal.is_none());
assert_eq!(
coordinator.operations.lookup(
&coordinator.instance,
&coordinator.instance,
"manifest-turn"
)["state"],
"expired"
);
}
#[test]
fn manifest_admission_boundary_reserves_escaped_claim_and_64_kib_terminal_frames() {
let base = json!({"resources":["", ""]});
let overhead = serde_json::to_vec(&json!({"application_manifest":base}))
.unwrap()
.len();
let mut manifest = base;
manifest["resources"] = json!(["x".repeat(9_000), "x".repeat(9_000 - overhead)]);
let budget = crate::service::protocol::persistent_assistant_budget(Some(&manifest)).unwrap();
assert_eq!(budget, 2_000);
let mut excessive = manifest.clone();
excessive["resources"][0] = json!(format!("{}x", excessive["resources"][0].as_str().unwrap()));
assert_eq!(
crate::service::protocol::persistent_assistant_budget(Some(&excessive)),
Err(Code::LimitExceeded)
);
let identity = "\0".repeat(128);
let text = "\"".repeat((budget - 2) / 2);
let mut terminal = ServiceEvent::turn_terminal(
identity.clone(),
identity.clone(),
"t".repeat(36),
TurnTerminalStatus::Failed,
text,
);
terminal.payload["application_manifest"] = manifest;
terminal.payload["operation_id"] = json!(identity);
terminal.payload["sequence"] = json!(u64::MAX);
terminal.payload["activity_dropped"] = json!(u64::MAX);
let snapshot = json!({"instance_id":"i".repeat(36),"session_id":identity,"live_sequence":u64::MAX,"phase":"idle","turn":null,"terminal":terminal.payload,"activity_replay_available":false});
assert!(serde_json::to_vec(&snapshot).unwrap().len() <= 24_576);
let claim = json!({"protocol_version":2,"kind":"response","instance_id":"i".repeat(36),"connection_id":"c".repeat(36),"request_id":identity,"session_id":identity,"operation_id":identity,"method":"session.claim","payload":{"grant_id":"g".repeat(36),"generation":wire::MAX_COUNTER,"snapshot":snapshot},"error":null});
let event = json!({"protocol_version":2,"kind":"event","instance_id":"i".repeat(36),"connection_id":"c".repeat(36),"event_id":"e".repeat(36),"session_id":identity,"operation_id":identity,"grant_generation":wire::MAX_COUNTER,"live_sequence":u64::MAX,"event":"turn.terminal","payload":terminal.payload});
for frame in [claim, event] {
assert!(serde_json::to_vec(&frame).unwrap().len() < 64 * 1024);
crate::service::protocol::payload_is_bounded(&frame["payload"]).unwrap();
}
}