use super::*;
pub(in crate::service) const EVIDENCE_METHOD: &str = "turn.application_evidence";
const MAX_EVIDENCE_TURNS: usize = 256;
const MAX_EVIDENCE_BYTES: usize = 30 * 1024;
pub(super) struct TurnEvidence {
session: String,
turn: Option<String>,
manifest: Value,
calls: BTreeMap<String, Value>,
}
impl Call {
pub(super) fn pending_evidence(&self, id: &str, now: Instant) -> Value {
let mut evidence = json!({
"call_id":id,"session_id":self.session,"turn_id":self.turn,
"provider_call_id":self.provider_call_id,"revision_sha256":self.revision,
"status":if self.dispatched {"pending"} else {"waiting_for_executor"},
"external_effects":if self.dispatched {"uncertain"} else {"not_dispatched"},
"wait_kind":if self.dispatched {"execution"} else {"availability"},
"remaining_ms":self.deadline.saturating_duration_since(now).as_millis() as u64
});
if let Some(identity) = &self.resource_identity {
identity.add_to(&mut evidence);
}
evidence
}
}
fn reserve_waiting_state(call: &mut Value) {
call["status"] = json!("waiting_for_executor");
call["external_effects"] = json!("not_dispatched");
call["wait_kind"] = json!("availability");
call["remaining_ms"] = json!(600_000);
}
impl TurnEvidence {
fn payload(&self) -> Value {
json!({"turn_id":self.turn,"application_manifest":self.manifest,
"calls":self.calls.values().collect::<Vec<_>>()})
}
fn reserved_bytes(&self) -> usize {
let mut payload = self.payload();
payload["turn_id"] = json!("\u{0001}".repeat(128));
for call in payload["calls"].as_array_mut().expect("calls array") {
reserve_waiting_state(call);
}
payload.to_string().len()
}
}
impl ApplicationState {
pub(super) fn reserve_resume_capacity(&self, id: &str, candidate: &Call) -> Result<(), Code> {
let now = Instant::now();
let calls: Vec<_> = self
.calls
.iter()
.filter(|(_, call)| call.executor == candidate.executor && call.outcome.is_none())
.map(|(id, call)| (id.as_str(), call))
.chain(std::iter::once((id, candidate)))
.map(|(id, call)| {
let mut evidence = call.pending_evidence(id, now);
reserve_waiting_state(&mut evidence);
evidence
})
.collect();
let payload = json!({"executor_id":candidate.executor,
"executor_generation":9_007_199_254_740_991u64,"pending_calls":calls});
if payload.to_string().len() > MAX_EVIDENCE_BYTES {
return Err(Code::LimitExceeded);
}
Ok(())
}
pub fn reserve_evidence(
&mut self,
operation: &str,
session: &str,
manifest: &Value,
) -> Result<(), Code> {
if self.evidence.len() >= MAX_EVIDENCE_TURNS {
return Err(Code::LimitExceeded);
}
let evidence = TurnEvidence {
session: session.into(),
turn: None,
manifest: manifest.clone(),
calls: BTreeMap::new(),
};
if evidence.reserved_bytes() > MAX_EVIDENCE_BYTES {
return Err(Code::LimitExceeded);
}
self.evidence.insert(operation.into(), evidence);
Ok(())
}
pub fn bind_evidence(&mut self, operation: &str, turn: Option<&str>) {
if let Some(turn) = turn {
if let Some(evidence) = self.evidence.get_mut(operation) {
evidence.turn = Some(turn.into());
}
} else {
self.evidence.remove(operation);
}
}
pub fn retain_evidence(&mut self, retained: impl Fn(&str) -> bool) {
self.evidence.retain(|operation, _| retained(operation));
}
pub fn read_evidence(&self, session: &str, payload: &Value) -> Result<Value, Code> {
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct Read {
turn_id: String,
}
let request: Read = decode(payload)?;
let evidence = self
.evidence
.values()
.find(|evidence| {
evidence.session == session && evidence.turn.as_deref() == Some(&request.turn_id)
})
.ok_or(Code::UnknownTurn)?;
let mut payload = evidence.payload();
let now = Instant::now();
for entry in payload["calls"].as_array_mut().expect("calls array") {
let pending = entry["call_id"]
.as_str()
.and_then(|id| self.calls.get(id))
.filter(|call| call.outcome.is_none());
if let Some(call) = pending {
let current =
call.pending_evidence(entry["call_id"].as_str().expect("call id"), now);
for field in ["status", "wait_kind", "remaining_ms", "external_effects"] {
entry[field] = current[field].clone();
}
} else {
entry["wait_kind"] = Value::Null;
entry["remaining_ms"] = Value::Null;
}
}
Ok(payload)
}
pub(super) fn reserve_call_evidence(
&mut self,
session: &str,
turn: &str,
id: &str,
attribution: Value,
) -> Result<(), Code> {
let Some(evidence) = self.evidence.values_mut().find(|evidence| {
evidence.session == session
&& (evidence.turn.is_none() || evidence.turn.as_deref() == Some(turn))
}) else {
return Ok(());
};
if evidence.calls.len() >= MAX_CALLS || evidence.calls.contains_key(id) {
return Err(Code::LimitExceeded);
}
evidence.turn = Some(turn.into());
evidence.calls.insert(id.into(), attribution);
if evidence.reserved_bytes() > MAX_EVIDENCE_BYTES {
evidence.calls.remove(id);
return Err(Code::LimitExceeded);
}
Ok(())
}
pub(super) fn update_call_evidence(&mut self, id: &str, status: &str, effects: &str) {
for evidence in self.evidence.values_mut() {
if let Some(call) = evidence.calls.get_mut(id) {
call["status"] = json!(status);
call["external_effects"] = json!(effects);
break;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn evidence_capacity_reserves_all_status_growth_and_releases_only_expired_turns() {
let mut state = ApplicationState::default();
let manifest = json!({"profile":PROFILE,"tools":(0..16).map(|index|
json!({"name":format!("tool_{index}"),"executor_id":"e".repeat(36),
"executor_generation":9_007_199_254_740_991u64,"revision_sha256":"f".repeat(64)}))
.collect::<Vec<_>>()});
state
.reserve_evidence("operation", "session", &manifest)
.unwrap();
for index in 0..32 {
let id = format!("call_{index}");
state.reserve_call_evidence("session", "turn", &id, json!({
"call_id":id,"provider_call_id":"\u{0001}".repeat(128),"name":"n".repeat(48),
"executor_id":"e".repeat(36),"executor_generation":9_007_199_254_740_991u64,
"revision_sha256":"f".repeat(64),"status":"pending","external_effects":"uncertain"
})).ok();
}
let accepted = state
.read_evidence("session", &json!({"turn_id":"turn"}))
.unwrap();
let calls = accepted["calls"].as_array().unwrap();
assert!(!calls.is_empty());
assert!(
calls.len() < 32,
"escaped identities must hit byte capacity first"
);
for call in calls {
state.update_call_evidence(
call["call_id"].as_str().unwrap(),
"cancelled",
"not_dispatched",
);
}
let final_payload = state
.read_evidence("session", &json!({"turn_id":"turn"}))
.unwrap();
assert!(final_payload.to_string().len() <= MAX_EVIDENCE_BYTES);
crate::service::protocol::payload_is_bounded(&final_payload).unwrap();
assert_eq!(
final_payload["calls"].as_array().unwrap().len(),
calls.len()
);
assert_eq!(
state.read_evidence("foreign", &json!({"turn_id":"turn"})),
Err(Code::UnknownTurn)
);
for index in 1..MAX_EVIDENCE_TURNS {
state
.reserve_evidence(&index.to_string(), "other", &manifest)
.unwrap();
}
assert_eq!(
state.reserve_evidence("overflow", "other", &manifest),
Err(Code::LimitExceeded)
);
state.retain_evidence(|operation| operation != "operation");
assert_eq!(
state.read_evidence("session", &json!({"turn_id":"turn"})),
Err(Code::UnknownTurn)
);
state
.reserve_evidence("replacement", "other", &manifest)
.unwrap();
}
}