mod auth_work;
mod entrypoint;
mod operations;
mod wire;
pub use entrypoint::PersistentService;
#[cfg(test)]
mod tests;
use super::{
coordinator::ServiceCoordinator,
protocol::{ServiceErrorCode as Code, ServiceMessage},
runtime::ServiceRuntime,
};
use operations::Operations;
use serde_json::{Value, json};
use std::{
collections::{BTreeMap, HashMap, VecDeque},
sync::Arc,
time::{Duration, Instant},
};
use wire::{Control, Request};
struct Connection {
initialized: bool,
accepted_at: Instant,
requests: HashMap<String, usize>,
queue: VecDeque<String>,
queued_bytes: usize,
}
struct Actor {
controller: Option<(String, Control)>,
live_sequence: u64,
turn: Option<Value>,
terminal: Option<Value>,
phase: &'static str,
delivery_dropped: u64,
}
struct RetainedTerminal {
operation_id: String,
payload: Value,
sequence: u64,
}
struct PendingRequest {
connection: String,
request: Request,
accepted: bool,
}
struct PendingTurnPreparation {
connection: String,
request: Request,
deadline: Instant,
worker: std::thread::JoinHandle<Result<Arc<ServiceRuntime>, Code>>,
}
struct Coordinator {
runtime: Arc<ServiceRuntime>,
execution: ServiceCoordinator,
instance: String,
workspace: String,
state_root: String,
unix_transport: bool,
generation: u64,
connections: HashMap<String, Connection>,
actors: BTreeMap<String, Actor>,
operations: Operations,
terminals: HashMap<String, RetainedTerminal>,
pending: HashMap<String, PendingRequest>,
turn_preparation: Option<PendingTurnPreparation>,
auth_work: Option<auth_work::PendingAuthWork>,
cleaning: HashMap<String, super::protocol::ServiceEvent>,
login: Option<(String, String)>,
}
impl Coordinator {
fn new(runtime: Arc<ServiceRuntime>) -> anyhow::Result<Self> {
use sha2::{Digest, Sha256};
let identity = |path: &std::path::Path| {
let canonical = path.canonicalize()?;
Ok::<_, std::io::Error>(crate::hex::lower_hex(Sha256::digest(
canonical.as_os_str().as_encoded_bytes(),
)))
};
let workspace = identity(&runtime.cwd)?;
let state_root = identity(&runtime.config.paths.root)?;
Ok(Self {
execution: ServiceCoordinator::persistent(Arc::clone(&runtime)),
runtime,
instance: uuid::Uuid::new_v4().to_string(),
workspace,
state_root,
unix_transport: false,
generation: 0,
connections: HashMap::new(),
actors: BTreeMap::new(),
operations: Operations::new(),
terminals: HashMap::new(),
pending: HashMap::new(),
turn_preparation: None,
auth_work: None,
cleaning: HashMap::new(),
login: None,
})
}
fn is_idle(&self) -> bool {
self.connections.is_empty()
&& self.actors.is_empty()
&& self.cleaning.is_empty()
&& self.operations.admitted() == 0
&& self.pending.is_empty()
&& self.turn_preparation.is_none()
&& self.auth_work.is_none()
&& self.login.is_none()
&& !self.execution.configuration_busy()
}
fn connect(&mut self, now: Instant) -> Result<String, Code> {
if self.connections.len() >= 16 {
return Err(Code::LimitExceeded);
}
let id = uuid::Uuid::new_v4().to_string();
self.connections.insert(
id.clone(),
Connection {
initialized: false,
accepted_at: now,
requests: HashMap::new(),
queue: VecDeque::new(),
queued_bytes: 0,
},
);
Ok(id)
}
fn disconnect(&mut self, connection: &str) {
self.connections.remove(connection);
let sessions: Vec<_> = self
.actors
.iter()
.filter(|(_, actor)| {
actor
.controller
.as_ref()
.is_some_and(|(owner, _)| owner == connection)
})
.map(|(id, _)| id.clone())
.collect();
for id in sessions {
self.detach(&id);
}
if self
.login
.as_ref()
.is_some_and(|(owner, _)| owner == connection)
{
self.execution.cancel_login();
}
}
fn detach(&mut self, session: &str) {
if let Some(actor) = self.actors.get_mut(session) {
if actor.controller.take().is_some() {
self.generation += 1;
}
if actor.phase == "idle" {
self.actors.remove(session);
}
}
self.execution.detach_session(session);
}
fn queue(&mut self, connection: &str, message: Value, activity: bool) -> bool {
let Some(client) = self.connections.get_mut(connection) else {
return false;
};
let mut encoded = message.to_string();
encoded.push('\n');
if encoded.len() > super::protocol::MAX_RECORD_BYTES
|| super::protocol::payload_is_bounded(&message["payload"]).is_err()
|| client.queue.len() >= 64
|| client.queued_bytes + encoded.len() > 262_144
{
if !activity {
self.disconnect(connection);
}
return false;
}
client.queued_bytes += encoded.len();
client.queue.push_back(encoded);
true
}
fn reject_decoding(&mut self, connection: &str, bytes: &[u8], code: Code) -> Result<(), Code> {
let mut reply = wire::decoding_error(&self.instance, connection, bytes, code);
let client = self
.connections
.get_mut(connection)
.ok_or(Code::StaleConnection)?;
if let Some(id) = reply["request_id"].as_str() {
if let Some(count) = client.requests.get_mut(id) {
*count += 1;
reply["error"] = wire::error(Code::DuplicateRequestId);
} else if client.requests.len() < 32 {
client.requests.insert(id.to_owned(), 1);
} else {
return Err(Code::LimitExceeded);
}
}
self.queue(connection, reply, false);
Ok(())
}
#[cfg(test)]
fn submit(&mut self, connection: &str, bytes: &[u8], now: Instant) -> Result<(), Code> {
self.submit_until(connection, bytes, now, now + Duration::from_secs(30))
}
fn submit_until(
&mut self,
connection: &str,
bytes: &[u8],
now: Instant,
deadline: Instant,
) -> Result<(), Code> {
let request = match wire::decode(bytes) {
Ok(request) => request,
Err(code) => return self.reject_decoding(connection, bytes, code),
};
if !wire::valid_id(&request.request_id) {
return self.reject_decoding(connection, bytes, Code::InvalidRequest);
}
let client = self
.connections
.get_mut(connection)
.ok_or(Code::StaleConnection)?;
if let Some(reservations) = client.requests.get_mut(&request.request_id) {
*reservations += 1;
let reply = wire::response(
&self.instance,
connection,
&request,
Err(Code::DuplicateRequestId),
);
self.queue(connection, reply, false);
return Ok(());
}
if client.requests.len() >= 32 {
return Err(Code::LimitExceeded);
}
client.requests.insert(request.request_id.clone(), 1);
let result = self.admit(connection, &request, now);
if let Some(pending) = &mut self.turn_preparation
&& pending.connection == connection
&& pending.request.request_id == request.request_id
{
pending.deadline = pending.deadline.min(deadline);
}
if let Some(result) = result {
self.queue(
connection,
wire::response(&self.instance, connection, &request, result),
false,
);
}
Ok(())
}
fn admit(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
if let Err(code) = self.validate_connection(connection, request) {
return Some(Err(code));
}
if request.method == "initialize" {
return Some(self.initialize(connection, request));
}
if request.mutation() {
if let Err(code) = self.operations.reserve(request, now) {
return Some(Err(code));
}
if let Err(code) = self.validate_grant(connection, request) {
self.settle_rejection(request, code, now);
return Some(Err(code));
}
}
let result = self.route(connection, request, now);
if let Some(ref result) = result {
self.settle_immediate(request, result, now);
}
result
}
fn validate_connection(&self, connection: &str, request: &Request) -> Result<(), Code> {
request.validate()?;
let client = self
.connections
.get(connection)
.ok_or(Code::StaleConnection)?;
if request.method == "initialize" {
if client.initialized {
return Err(Code::AlreadyInitialized);
}
if request.instance_id.is_some() || request.connection_id.is_some() {
return Err(Code::InvalidRequest);
}
} else {
if !client.initialized {
return Err(Code::NotInitialized);
}
if request.instance_id.as_deref() != Some(&self.instance) {
return Err(Code::StaleInstance);
}
if request.connection_id.as_deref() != Some(connection) {
return Err(Code::StaleConnection);
}
}
Ok(())
}
fn validate_grant(&self, connection: &str, request: &Request) -> Result<(), Code> {
if !request.needs_control() {
return Ok(());
}
let controller = request
.session_id
.as_ref()
.and_then(|id| self.actors.get(id))
.and_then(|actor| actor.controller.as_ref());
if !controller.is_some_and(|(owner, control)| {
owner == connection && Some(control) == request.control.as_ref()
}) {
return Err(Code::StaleGrant);
}
Ok(())
}
fn initialize(&mut self, connection: &str, request: &Request) -> Result<Value, Code> {
let params: super::protocol::InitializeParams =
serde_json::from_value(request.payload.clone()).map_err(|_| Code::InvalidPayload)?;
if !params.supported_protocol_versions.contains(&2) {
return Err(Code::UnsupportedVersion);
}
if params.supported_protocol_versions.len() > 8 || params.requested_capabilities.len() > 8 {
return Err(Code::LimitExceeded);
}
if params
.requested_capabilities
.iter()
.any(|name| !wire::OPERATIONS.contains(&name.as_str()))
{
return Err(Code::UnsupportedCapability);
}
self.connections
.get_mut(connection)
.ok_or(Code::StaleConnection)?
.initialized = true;
Ok(
json!({"protocol_version":2,"server_name":"magi-code","server_version":env!("CARGO_PKG_VERSION"),
"instance_id":self.instance,"connection_id":connection,"workspace_id":self.workspace,"state_root_id":self.state_root,
"limits":limits(),"capabilities":capabilities(self.unix_transport)}),
)
}
fn route(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
let synchronous = match request.method.as_str() {
"capabilities" => empty(request).map(|()| capabilities(self.unix_transport)),
"status" | "auth.status" | "auth.logout" => {
return self.start_auth_work(connection, request, now);
}
"auth.login.start" if self.auth_work.is_some() => Err(Code::AuthBusy),
"session.active" => empty(request).map(|()| {
json!({"sessions":self.actors.iter().map(|(id, actor)| json!({
"session_id":id,"controlled":actor.controller.is_some(),"phase":actor.phase,
"turn_id":actor.turn.as_ref().map(|turn| &turn["turn_id"])
})).collect::<Vec<_>>()})
}),
"session.claim" => self.claim(connection, request),
"session.detach" => empty(request).map(|()| {
self.detach(request.session_id.as_deref().unwrap_or_default());
json!({"status":"detached"})
}),
"operation.lookup" => self.lookup(request),
"turn.start" | "turn.cancel" => return self.turn_request(connection, request, now),
"auth.login.callback" | "auth.login.cancel"
if !self
.login
.as_ref()
.is_some_and(|(owner, _)| owner == connection) =>
{
Err(Code::UnknownLogin)
}
_ => return self.legacy_request(connection, request, now),
};
Some(synchronous)
}
fn lookup(&self, request: &Request) -> Result<Value, Code> {
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct Lookup {
target_instance_id: String,
operation_id: String,
}
let params: Lookup =
serde_json::from_value(request.payload.clone()).map_err(|_| Code::InvalidPayload)?;
if !wire::valid_id(¶ms.target_instance_id) || !wire::valid_id(¶ms.operation_id) {
return Err(Code::InvalidPayload);
}
Ok(self.operations.lookup(
&self.instance,
¶ms.target_instance_id,
¶ms.operation_id,
))
}
fn snapshot(&self, session: &str, actor: &Actor) -> Value {
json!({"instance_id":self.instance,"session_id":session,"live_sequence":actor.live_sequence,
"phase":actor.phase,"turn":actor.turn,"terminal":actor.terminal,"activity_replay_available":false})
}
fn claim(&mut self, connection: &str, request: &Request) -> Result<Value, Code> {
empty(request)?;
let session = request.session_id.as_deref().ok_or(Code::InvalidPayload)?;
if self
.actors
.get(session)
.is_some_and(|actor| actor.controller.is_some())
{
return Err(Code::SessionControlled);
}
let controlled = self
.actors
.values()
.filter(|actor| actor.controller.is_some())
.count() as u64;
if self.generation > wire::MAX_COUNTER - controlled - 2 {
return Err(Code::LimitExceeded);
}
self.execution.claim_session(session)?;
self.generation += 1;
let control = Control {
grant_id: uuid::Uuid::new_v4().to_string(),
generation: self.generation,
};
let retained = self.terminals.get(session);
let actor = self
.actors
.entry(session.to_owned())
.or_insert_with(|| Actor {
controller: None,
live_sequence: retained.map_or(0, |terminal| terminal.sequence),
turn: None,
terminal: retained.map(|terminal| terminal.payload.clone()),
phase: "idle",
delivery_dropped: 0,
});
actor.controller = Some((connection.to_owned(), control.clone()));
let snapshot = self.snapshot(session, &self.actors[session]);
Ok(json!({"grant_id":control.grant_id,"generation":control.generation,"snapshot":snapshot}))
}
fn turn_request(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
let session = request.session_id.as_deref().unwrap_or_default();
let actor = self.actors.get(session)?;
if request.method == "turn.start" && actor.phase != "idle" {
return Some(Err(Code::SessionBusy));
}
if request.method == "turn.cancel"
&& (actor.phase != "running"
|| actor
.turn
.as_ref()
.is_none_or(|turn| turn["turn_id"] != request.payload["turn_id"]))
{
return Some(Err(Code::UnknownTurn));
}
if request.method == "turn.start" {
return self.prepare_turn(connection, request, now);
}
self.legacy_request(connection, request, now)
}
fn prepare_turn(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
if self.turn_preparation.is_some() || self.execution.configuration_busy() {
return Some(Err(Code::ConfigurationBusy));
}
let runtime = Arc::clone(&self.runtime);
let worker = match std::thread::Builder::new()
.name("magi-turn-preparation".into())
.spawn(move || {
runtime
.capture_turn_settings()
.map(Arc::new)
.map_err(|_| Code::InternalError)
}) {
Ok(worker) => worker,
Err(_) => return Some(Err(Code::InternalError)),
};
self.turn_preparation = Some(PendingTurnPreparation {
connection: connection.to_owned(),
request: request.clone(),
deadline: now + Duration::from_secs(30),
worker,
});
None
}
fn finish_turn_preparation(&mut self, now: Instant) {
if !self
.turn_preparation
.as_ref()
.is_some_and(|pending| pending.worker.is_finished())
{
return;
}
let pending = self.turn_preparation.take().expect("finished preparation");
let prepared = pending.worker.join().unwrap_or(Err(Code::InternalError));
let result = self
.validate_prepared_turn(&pending.connection, &pending.request, pending.deadline, now)
.and(prepared);
let response = match result {
Ok(runtime) => self.legacy_request_with_runtime(
&pending.connection,
&pending.request,
now,
Some(runtime),
),
Err(code) => Some(Err(code)),
};
if let Some(result) = response {
self.settle_immediate(&pending.request, &result, now);
self.queue(
&pending.connection,
wire::response(
&self.instance,
&pending.connection,
&pending.request,
result,
),
false,
);
}
}
fn validate_prepared_turn(
&self,
connection: &str,
request: &Request,
deadline: Instant,
now: Instant,
) -> Result<(), Code> {
if now >= deadline {
return Err(Code::RequestTimeout);
}
self.validate_connection(connection, request)?;
self.validate_grant(connection, request)?;
let actor = self
.actors
.get(request.session_id.as_deref().unwrap_or_default())
.ok_or(Code::StaleGrant)?;
if actor.phase != "idle" {
return Err(Code::SessionBusy);
}
let operation = self.operations.lookup(
&self.instance,
&self.instance,
request.operation_id.as_deref().unwrap_or_default(),
);
if operation["state"] != "in_progress" {
return Err(Code::OperationAlreadyKnown);
}
Ok(())
}
fn legacy_request(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
self.legacy_request_with_runtime(connection, request, now, None)
}
fn legacy_request_with_runtime(
&mut self,
connection: &str,
request: &Request,
now: Instant,
prepared_runtime: Option<Arc<ServiceRuntime>>,
) -> Option<Result<Value, Code>> {
if request.method == "config.set" && self.turn_preparation.is_some() {
return Some(Err(Code::ConfigurationBusy));
}
let internal_id = uuid::Uuid::new_v4().to_string();
let mut legacy = request.legacy(&request.method);
legacy.request_id = internal_id.clone();
let outbound = if request.method == "turn.start" {
let live_sequence = request
.session_id
.as_ref()
.and_then(|id| self.actors.get(id))
.map_or(wire::MAX_COUNTER, |actor| actor.live_sequence);
if live_sequence >= wire::MAX_COUNTER - 1 {
return Some(Err(Code::LimitExceeded));
}
self.execution.dispatch_prepared_turn_request(
legacy,
wire::MAX_COUNTER - live_sequence - 1,
prepared_runtime,
)
} else {
self.execution.dispatch_request(legacy)
};
let mut correlation = request.clone();
correlation.payload = Value::Null;
self.pending.insert(
internal_id.clone(),
PendingRequest {
connection: connection.to_owned(),
request: correlation,
accepted: outbound.messages().is_empty(),
},
);
if outbound.messages().is_empty() {
if let Some(operation) = &request.operation_id {
self.operations
.update(operation, "accepted", Value::Null, Value::Null, now);
}
return None;
}
self.consume_legacy(outbound.messages(), now);
None
}
fn consume_legacy(&mut self, messages: &[ServiceMessage], now: Instant) {
for message in messages {
match message {
ServiceMessage::Response(response) => {
let Some(id) = &response.request_id else {
continue;
};
let Some(pending) = self.pending.remove(id) else {
continue;
};
let result = if let Some(error) = &response.error {
Err(error.code)
} else {
serde_json::to_value(&response.payload).map_err(|_| Code::InternalError)
};
self.accept_legacy(&pending, &result, now);
if pending.accepted
&& let Err(code) = result
{
if let Some(operation) = &pending.request.operation_id {
self.operations.update(
operation,
"terminal",
json!({"status":"failed"}),
wire::error(code),
now,
);
}
} else {
self.settle_immediate(&pending.request, &result, now);
}
let mut result = result;
if let Ok(payload) = &mut result {
if pending.request.method == "turn.start"
&& let Some(fields) = payload.as_object_mut()
{
fields.remove("session_id");
}
if pending.request.method == "session.create" {
*payload = json!({"session_id":payload["session_id"]});
}
if pending.request.method == "turn.cancel" {
payload["session_id"] = json!(pending.request.session_id);
}
}
self.queue(
&pending.connection,
wire::response(
&self.instance,
&pending.connection,
&pending.request,
result,
),
false,
);
}
ServiceMessage::Event(event) => {
if event.event.starts_with("auth.") {
self.auth_event(event, now);
} else {
self.turn_event(event, now);
}
}
}
}
}
fn accept_legacy(
&mut self,
pending: &PendingRequest,
result: &Result<Value, Code>,
now: Instant,
) {
let Ok(payload) = result else { return };
let request = &pending.request;
let operation = request.operation_id.as_deref().unwrap_or_default();
match request.method.as_str() {
"turn.start" => {
let Some(session) = request.session_id.as_ref() else {
return;
};
let Some(actor) = self.actors.get_mut(session) else {
return;
};
actor.phase = "running";
actor.terminal = None;
actor.delivery_dropped = 0;
actor.turn = Some(
json!({"turn_id":payload["turn_id"],"operation_id":operation,"sequence":0,
"assistant_text":"","activity_dropped":0,"replay_required":false}),
);
if let Some(turn_id) = payload["turn_id"].as_str() {
self.execution.release_turn_request(turn_id);
}
self.operations.update(
operation,
"accepted",
json!({"turn_id":payload["turn_id"]}),
Value::Null,
now,
);
}
"auth.login.start" => {
self.login = Some((pending.connection.clone(), operation.to_owned()));
self.operations.update(
operation,
"accepted",
json!({"login_id":payload["login_id"]}),
Value::Null,
now,
);
}
"session.create" => {
if let Some(session) = payload["session_id"].as_str() {
self.execution.detach_session(session);
}
}
_ => {}
}
}
fn settle_rejection(&mut self, request: &Request, code: Code, now: Instant) {
if let Some(operation) = &request.operation_id {
self.operations
.update(operation, "rejected", Value::Null, wire::error(code), now);
}
}
fn settle_immediate(&mut self, request: &Request, result: &Result<Value, Code>, now: Instant) {
let Some(operation) = &request.operation_id else {
return;
};
match result {
Err(code)
if (request.method == "session.create" && *code == Code::SessionUnavailable)
|| (request.method == "auth.logout" && *code == Code::InternalError) =>
{
let mut settlement = json!({"status":"failed"});
if request.method == "session.create" {
settlement["session_id"] = Value::Null;
}
self.operations
.update(operation, "terminal", settlement, wire::error(*code), now);
}
Err(code) => self.settle_rejection(request, *code, now),
Ok(_) if matches!(request.method.as_str(), "turn.start" | "auth.login.start") => {}
Ok(payload) => {
let mut settlement = json!({"status":"completed"});
if request.method == "session.create" {
settlement["session_id"] = payload["session_id"].clone();
}
self.operations
.update(operation, "terminal", settlement, Value::Null, now);
}
}
}
fn turn_event(&mut self, event: &super::protocol::ServiceEvent, now: Instant) {
let Some(session) = &event.session_id else {
return;
};
let Some(actor) = self.actors.get_mut(session) else {
return;
};
let Some(turn) = actor.turn.as_mut() else {
return;
};
if turn["turn_id"] != event.payload["turn_id"] {
return;
}
let operation = turn["operation_id"].as_str().unwrap_or_default().to_owned();
let sequence = event.payload["sequence"].as_u64().unwrap_or(0);
let previous = turn["sequence"].as_u64().unwrap_or(0);
actor.live_sequence = actor
.live_sequence
.saturating_add(sequence.saturating_sub(previous).max(1));
turn["sequence"] = json!(sequence);
turn["activity_dropped"] = json!(
turn["activity_dropped"].as_u64().unwrap_or(0)
+ sequence.saturating_sub(previous).saturating_sub(1)
);
if event.event == "turn.assistant_delta" {
let mut text = turn["assistant_text"]
.as_str()
.unwrap_or_default()
.to_owned();
text.push_str(event.payload["text"].as_str().unwrap_or_default());
turn["assistant_text"] = json!(text);
}
let activity = event.event == "turn.activity";
let mut payload = event.payload.clone();
if event.event == "turn.terminal" {
payload["activity_dropped"] =
json!(payload["activity_dropped"].as_u64().unwrap_or(0) + actor.delivery_dropped);
let mut terminal = payload.clone();
terminal["operation_id"] = json!(operation);
actor.terminal = Some(terminal.clone());
actor.turn = None;
actor.phase = "idle";
self.terminals.insert(
session.clone(),
RetainedTerminal {
operation_id: operation.clone(),
payload: terminal,
sequence: actor.live_sequence,
},
);
self.operations.update(&operation, "terminal", json!({"turn_id":payload["turn_id"],"status":payload["status"],"persistence":payload["persistence"]}), payload["error"].clone(), now);
}
let route = actor.controller.clone();
let live_sequence = actor.live_sequence;
let delivered = if let Some((owner, control)) = route {
let message = json!({"protocol_version":2,"kind":"event","instance_id":self.instance,"connection_id":owner,
"event_id":uuid::Uuid::new_v4().to_string(),"session_id":session,"operation_id":operation,
"grant_generation":control.generation,"live_sequence":live_sequence,"event":event.event,"payload":payload});
self.queue(&owner, message, activity)
} else {
false
};
if activity
&& !delivered
&& let Some(actor) = self.actors.get_mut(session)
{
actor.delivery_dropped += 1;
if let Some(turn) = &mut actor.turn {
turn["activity_dropped"] =
json!(turn["activity_dropped"].as_u64().unwrap_or(0) + 1);
}
}
if event.event == "turn.terminal"
&& self
.actors
.get(session)
.is_some_and(|actor| actor.controller.is_none())
{
self.actors.remove(session);
self.execution.detach_session(session);
}
}
fn auth_event(&mut self, event: &super::protocol::ServiceEvent, now: Instant) {
let Some((owner, operation)) = self.login.clone() else {
return;
};
if event.event == "auth.login.terminal" {
let status = match event.payload["state"].as_str() {
Some("exchange_succeeded") => "completed",
Some("cancelled") => "cancelled",
_ => "failed",
};
let error = if status == "failed" {
wire::error(Code::InternalError)
} else {
Value::Null
};
self.operations.update(&operation, "terminal", json!({"login_id":event.payload["login_id"],"status":status,"cleanup_complete":true}), error, now);
self.login = None;
}
self.queue(&owner, json!({"protocol_version":2,"kind":"event","instance_id":self.instance,"connection_id":owner,
"event_id":uuid::Uuid::new_v4().to_string(),"session_id":null,"operation_id":operation,
"grant_generation":null,"live_sequence":null,"event":event.event,"payload":event.payload}), false);
}
fn tick(&mut self, now: Instant) {
self.finish_turn_preparation(now);
self.finish_auth_work(now);
self.operations.expire(now);
self.terminals
.retain(|_, terminal| self.operations.retained(&terminal.operation_id));
for actor in self.actors.values_mut() {
if actor.terminal.as_ref().is_some_and(|terminal| {
!self
.operations
.retained(terminal["operation_id"].as_str().unwrap_or_default())
}) {
actor.terminal = None;
}
}
let expired: Vec<_> = self
.connections
.iter()
.filter(|(_, client)| {
!client.initialized
&& now.saturating_duration_since(client.accepted_at) >= Duration::from_secs(5)
})
.map(|(id, _)| id.clone())
.collect();
for id in expired {
self.disconnect(&id);
}
for _ in 0..128 {
let Ok(message) = self.execution.worker_receiver().try_recv() else {
break;
};
if let super::turns::TurnWorkerMessage::Persisting { turn_id } = &message {
if let Some(actor) = self.actors.values_mut().find(|actor| {
actor
.turn
.as_ref()
.is_some_and(|turn| turn["turn_id"] == *turn_id)
}) {
actor.phase = "persisting";
}
continue;
}
if let Some(output) = self.execution.worker_output(message) {
if let Some(turn_id) = output.terminal_turn_id {
for message in output.outbound.messages() {
if let ServiceMessage::Event(event) = message {
if let Some(actor) = event
.session_id
.as_ref()
.and_then(|id| self.actors.get_mut(id))
{
actor.phase = "cleaning";
if let Some(turn) = actor.turn.as_mut() {
turn["assistant_text"] =
event.payload["assistant_text"].clone();
}
}
self.cleaning.insert(turn_id.clone(), event.clone());
}
}
} else {
self.consume_legacy(output.outbound.messages(), now);
}
}
}
let finished: Vec<_> = self
.cleaning
.keys()
.filter(|id| self.execution.worker_finished(id))
.cloned()
.collect();
for id in finished {
self.execution.finish_worker_output(&id);
if let Some(event) = self.cleaning.remove(&id) {
self.turn_event(&event, now);
}
}
for _ in 0..16 {
let Ok(message) = self.execution.auth_receiver().try_recv() else {
break;
};
if let Some(output) = self.execution.auth_output(message) {
self.consume_legacy(output.messages(), now);
}
}
}
}
fn empty(request: &Request) -> Result<(), Code> {
serde_json::from_value::<super::protocol::EmptyParams>(request.payload.clone())
.map(|_| ())
.map_err(|_| Code::InvalidPayload)
}
fn limits() -> Value {
let mut limits = serde_json::to_value(super::protocol::ProtocolLimits::current())
.expect("static limits serialize");
let additional = json!({"max_connections":16,"handshake_timeout_ms":5000,"max_inflight_requests":32,
"request_timeout_ms":30000,"max_active_turns":16,"max_operation_entries":256,"operation_retention_ms":600000,
"operation_tombstone_ms":600000,"max_snapshot_bytes":24576,"max_queued_records":64,"max_queued_bytes":262144,
"write_timeout_ms":5000,"max_auth_operations":1,"max_config_operations":1,"idle_shutdown_ms":60000,"startup_grace_ms":10000});
limits
.as_object_mut()
.expect("limits object")
.extend(additional.as_object().expect("limits object").clone());
limits
}
fn capabilities(unix_transport: bool) -> Value {
let transports: &[&str] = if unix_transport { &["unix"] } else { &[] };
json!({"protocol_versions":[2],"transports":transports,"operations":wire::OPERATIONS,"limits":limits(),
"events":["turn.started","turn.assistant_delta","turn.activity","turn.terminal","auth.login.progress","auth.login.terminal"],
"activity":super::activity::ActivityCapabilities::current()})
}