use super::*;
use super::outbound::{store_ack, transport_envelope};
use super::projection::{note_verdict, sync_session};
use super::retire::{release_locked, retire_idle_locked};
use super::state::{
DispatchInner, DispatchState, slot_key_named, slot_key_serving_task, slot_task,
};
use super::transport::{names_session, serves_session};
use onlyne_proto::ErrorCode;
use onlyne_proto::adapter::HandoffArgs;
pub const SETTLE_WITHOUT_TURN: &str = "settle_without_turn";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum SettleAuthority {
PluginReport,
ClientOwned,
ControlDriven,
}
fn turn_recorded(inner: &DispatchInner, task_id: &str) -> (bool, String) {
let Ok(Some(row)) = inner.store.get_session(task_id) else {
return (false, "no session row".to_string());
};
let agent = stored_observation(&inner.store, Some(&row)).agent;
(
matches!(agent, AgentState::Running | AgentState::Idle),
row.agent_state,
)
}
pub async fn on_out(
state: &DispatchState,
task_id: &str,
outcome: Outcome,
head: Option<String>,
head_kind: Option<&str>,
handoffs: &[Handoff],
asked: SettleAuthority,
) -> Result<()> {
if asked == SettleAuthority::PluginReport {
let inner = state.inner.lock();
let (turn, phase) = turn_recorded(&inner, task_id);
if !turn {
let reason = format!(
"no turn ran: the agent phase this client holds for the session reads {phase}"
);
onlyne_session::record_fault(
&inner.store,
task_id,
SETTLE_WITHOUT_TURN,
"client",
&reason,
)?;
tracing::warn!(
task = %task_id,
?outcome,
phase = %phase,
"a completion arrived for a session that never ran a turn; the task stays open"
);
return Ok(());
}
}
let settled = {
let mut inner = state.inner.lock();
let verdict = settle(&inner.bridge, &inner.store, task_id)?;
if !inner.store.settle_task(task_id, task_state_of(outcome))? {
tracing::warn!(
task = %task_id,
?outcome,
"a second verdict arrived for a settled task; the first one stands"
);
note_verdict(&verdict, task_id);
release_locked(&mut inner, task_id, None)?;
None
} else {
inner
.store
.put_out_head(task_id, head.as_deref().unwrap_or(""))?;
let slot =
slot_key_serving_task(&inner, task_id).and_then(|key| inner.sessions.get_mut(&key));
let origin = slot.as_ref().and_then(|slot| slot.origin.clone());
let causality = slot.as_ref().map(|slot| slot.causality.clone());
let msg_id = slot.and_then(|slot| slot.msg_id.take());
if let Some(msg_id) = msg_id {
store_ack(
&inner,
AckArgs {
msg_id,
op_id: None,
accepted: true,
reason: None,
},
);
}
release_locked(&mut inner, task_id, None)?;
let held = inner.held_handoffs.remove(task_id);
Some((
verdict,
completion_envelope(
&inner.role,
origin,
task_id,
head.as_deref(),
causality.as_ref(),
),
inner.role.clone(),
causality,
held,
))
}
};
let Some((verdict, receipt, role, causality, held)) = settled else {
return sync_session(state, task_id).await;
};
note_verdict(&verdict, task_id);
let routed = merged_handoffs(
handoffs,
held.as_deref(),
head.as_deref().unwrap_or_default(),
);
let parent = causality.unwrap_or_else(|| Causality::root(task_id));
let denied = handoff::route(
state,
&role,
&parent,
head_kind,
head.as_deref().unwrap_or_default(),
&routed,
)
.await;
record_denials(state, task_id, &denied)?;
retire_revived(state, task_id).await;
if let Some(envelope) = receipt {
transport_envelope(state, &envelope).await?;
}
sync_session(state, task_id).await
}
fn merged_handoffs<'a>(
own: &'a [Handoff],
held: Option<&'a [Handoff]>,
head: &str,
) -> Cow<'a, [Handoff]> {
let held: &[Handoff] = match held {
Some(held) if !held.is_empty() => held,
_ => return Cow::Borrowed(own),
};
let mut order: Vec<String> = Vec::new();
let mut segments: HashMap<String, Vec<String>> = HashMap::new();
for (marker, group) in [("[retry]", own), ("[zombie]", held)] {
for handoff in group {
let line = format!("{marker} {}", handoff.text_or(head));
if !segments.contains_key(handoff.to_role.as_str()) {
order.push(handoff.to_role.clone());
}
segments
.entry(handoff.to_role.clone())
.or_default()
.push(line);
}
}
Cow::Owned(
order
.into_iter()
.map(|to_role| Handoff {
text: Some(
segments
.remove(to_role.as_str())
.unwrap_or_default()
.join("\n"),
),
to_role,
})
.collect(),
)
}
async fn retire_revived(state: &DispatchState, task_id: &str) {
let leaving = {
let mut inner = state.inner.lock();
let mut leaving: Vec<AdapterIo> = Vec::new();
let mut silent: Vec<String> = Vec::new();
for (session_id, io, _) in inner.revived.iter() {
if inner.in_frame.iter().any(|busy| busy.same_connection(io)) {
continue;
}
let reaches = slot_key_named(&inner, session_id)
.and_then(|key| {
inner
.sessions
.get(&key)
.map(|slot| slot_task(slot) == task_id)
})
.unwrap_or_else(|| session_id == task_id);
if reaches {
leaving.push(io.clone());
silent.push(session_id.clone());
}
}
inner
.revived
.retain(|(session_id, _, _)| !silent.contains(session_id));
let silenced: Vec<String> = inner
.sessions
.iter()
.filter(|(_, slot)| slot.read_only && slot_task(slot) == task_id)
.map(|(key, _)| key.clone())
.collect();
for key in silenced {
let Some(slot) = inner.sessions.get(&key).cloned() else {
continue;
};
if slot.payload.is_some() {
tracing::warn!(
session = %key,
task = %task_id,
"a read-only session retires with a payload it was never handed"
);
}
let served: Vec<String> = inner
.transports
.keys()
.filter(|served| names_session(&key, &slot, served))
.cloned()
.collect();
for session_id in served {
inner.transports.remove(&session_id);
}
if let Some(current) = inner.sessions.get_mut(&key) {
current.task_id = None;
current.ready = false;
current.read_only = false;
current.dropped_at = None;
}
retire_idle_locked(&mut inner, &key, onlyne_session::CloseReason::Replaced);
}
leaving
};
for io in leaving {
let notice = AdapterMsg::Host(HostOp::Bye(onlyne_proto::ByeNotice {
reason: "the session that took this task answered for yours".into(),
}));
if let Err(error) = io.notify(notice).await {
tracing::debug!(error = %error, "the read-only connection had already left");
}
}
}
fn record_denials(state: &DispatchState, task_id: &str, denied: &[Denial]) -> Result<()> {
if denied.is_empty() {
return Ok(());
}
let inner = state.inner.lock();
for refusal in denied {
tracing::warn!(
task = %task_id,
to_role = %refusal.to_role,
error = %refusal.reason,
"handoff denied"
);
inner.store.append_event(
"handoff_denied",
&serde_json::json!({
"task_id": task_id,
"to_role": refusal.to_role,
"text": refusal.text,
"error": refusal.reason,
}),
)?;
onlyne_session::record_fault(
&inner.store,
task_id,
"handoff_denied",
"acp",
&format!("{}: {}", refusal.to_role, refusal.reason),
)?;
}
Ok(())
}
fn completion_envelope(
role: &str,
origin: Option<Principal>,
task_id: &str,
head: Option<&str>,
causality: Option<&Causality>,
) -> Option<Envelope> {
let origin = origin?;
let body = Body::text(head.unwrap_or_default());
let mut causality = causality.cloned().unwrap_or_default();
causality.task = task_id.to_string();
causality.parent_task = None;
causality.reply_to = None;
causality.attempt = 0;
new_envelope(
MsgKind::Completion,
Principal::role(role),
origin,
body,
Some(causality),
)
.ok()
}
impl DispatchState {
pub fn plugin_handoff(&self, io: &AdapterIo, args: HandoffArgs) -> ResBody {
let (role, parent) = {
let inner = self.inner.lock();
let found = slot_key_serving_task(&inner, &args.task_id).and_then(|key| {
inner
.sessions
.get(&key)
.map(|slot| (key, slot.causality.clone()))
});
let Some((key, parent)) = found else {
return ResBody::err(
ErrorCode::Invalid,
format!("no session serves task {}", args.task_id),
Some("task_id".into()),
);
};
if !serves_session(&inner, &key, io) {
return ResBody::err(
ErrorCode::Invalid,
format!("this connection does not serve task {}", args.task_id),
Some("task_id".into()),
);
}
(inner.role.clone(), parent)
};
let (envelope, child) =
match handoff::relay(&role, &parent, &args.to, &args.text, args.image) {
Ok(built) => built,
Err(message) => return ResBody::err(ErrorCode::Invalid, message, None),
};
let queued = match self.plugin_send(io, &envelope) {
Ok(queued) => queued,
Err(error) => return ResBody::err(ErrorCode::Internal, error.to_string(), None),
};
ResBody::ok(serde_json::json!({
"task_id": child.task,
"hop": child.hop,
"queued": true,
"op_id": queued["op_id"],
}))
}
}
#[cfg(test)]
mod tests;