use super::*;
use super::outbound::{store_ack, transport_envelope};
use super::projection::{note_verdict, phase, sync_session};
use super::retire::{PendingClose, close_retired, 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 = phase(&row.agent_state, AgentPhase::Booting);
(
matches!(agent, AgentPhase::Running | AgentPhase::Idle),
row.agent_state,
)
}
pub async fn on_out(
state: &DispatchState,
task_id: &str,
outcome: Outcome,
head: Option<String>,
details: Option<String>,
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}"
);
crate::reconcile::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, session_id) = {
let mut inner = state.inner.lock();
if take_verdict(&inner, task_id, outcome)? {
inner
.store
.put_out_head(task_id, head.as_deref().unwrap_or(""))?;
let key = slot_key_serving_task(&inner, task_id);
let slot = key.as_deref().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());
let session_id = key.unwrap_or_else(|| task_id.to_string());
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)?;
(
Some((completion_envelope(
&inner.role,
origin,
task_id,
head.as_deref(),
details.as_deref(),
causality.as_ref(),
),)),
session_id,
)
} else {
tracing::warn!(
task = %task_id,
?outcome,
"a second verdict arrived for a settled task; the first one stands"
);
let session_id =
slot_key_serving_task(&inner, task_id).unwrap_or_else(|| task_id.to_string());
release_locked(&mut inner, task_id, None)?;
(None, session_id)
}
};
let Some((receipt,)) = settled else {
return sync_session(state, &session_id).await;
};
retire_revived(state, task_id).await;
if let Some(envelope) = receipt {
transport_envelope(state, &envelope).await?;
}
sync_session(state, &session_id).await
}
fn take_verdict(inner: &DispatchInner, task_id: &str, outcome: Outcome) -> Result<bool> {
let drain = inner
.store
.get_session(task_id)?
.map(|_| settle(&inner.bridge, &inner.store, task_id))
.transpose()?;
let first = inner.store.settle_task(task_id, task_state_of(outcome))?;
if let Some(verdict) = drain {
note_verdict(&verdict, task_id);
}
Ok(first)
}
async fn retire_revived(state: &DispatchState, task_id: &str) {
let (leaving, pending) = {
let mut inner = state.inner.lock();
let mut leaving: Vec<AdapterIo> = Vec::new();
let mut pending: Vec<PendingClose> = 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());
}
}
inner
.revived
.retain(|(_, revived, _)| !leaving.iter().any(|io| io.same_connection(revived)));
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,
crate::backend::CloseReason::Replaced,
&mut pending,
);
}
(leaving, pending)
};
close_retired(pending);
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 completion_envelope(
role: &str,
origin: Option<Principal>,
task_id: &str,
head: Option<&str>,
details: Option<&str>,
causality: Option<&Causality>,
) -> Option<Envelope> {
let origin = origin?;
let body = Body {
text: Some(details.or(head).unwrap_or_default().to_string()),
head: head.map(str::to_string),
image: None,
};
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),
};
let op =
super::turn_end::record_handoff(self, &args.task_id, &args.to, child.hop, &args.text);
if let Err(error) = self.enqueue_op(&op) {
tracing::warn!(
task = %args.task_id,
to = %args.to,
error = %error,
"the handoff event was not queued"
);
}
ResBody::ok(serde_json::json!({
"task_id": child.task,
"hop": child.hop,
"queued": true,
"op_id": queued["op_id"],
}))
}
}