use super::*;
use super::env::missing_capability;
use super::idle::{resumable, suspend_locked};
use super::outbound::queue_outbound_locked;
use super::retire::{
PendingClose, close_retired, keeps_idle, retire_idle_locked, stored_close_reason,
};
use super::state::{
DispatchInner, DispatchState, FrameGuard, SessionSlot, has_attached_transport,
rebase_generation, slot_key_named, slot_key_serving_task, slot_task,
};
pub const NUDGE_TEXT: &str =
"If this task is finished, report it with onlyne_complete; if something is missing, say what.";
pub(super) fn names_session(key: &str, slot: &SessionSlot, session_id: &str) -> bool {
key == session_id
|| slot.session.task_id == session_id
|| slot.task_id.as_deref() == Some(session_id)
}
fn attached_to_other(inner: &DispatchInner, key: &str, slot: &SessionSlot, io: &AdapterIo) -> bool {
inner.transports.iter().any(|(session_id, (live, _))| {
!live.same_connection(io) && names_session(key, slot, session_id)
})
}
pub(super) fn serves_session(inner: &DispatchInner, session_id: &str, io: &AdapterIo) -> bool {
let Some(key) = slot_key_named(inner, session_id) else {
return false;
};
let Some(slot) = inner.sessions.get(&key) else {
return false;
};
slot_is_served_by(inner, &key, slot, io)
}
fn slot_is_served_by(inner: &DispatchInner, key: &str, slot: &SessionSlot, io: &AdapterIo) -> bool {
inner.transports.iter().any(|(served, (transport, _))| {
transport.same_connection(io) && names_session(key, slot, served)
}) || tools_bound_to(inner, key, io)
}
fn tools_bound_to(inner: &DispatchInner, key: &str, io: &AdapterIo) -> bool {
inner
.tools_mounts
.iter()
.any(|(held, bound)| held == key && bound.same_connection(io))
}
pub(super) fn serves_task(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
inner.sessions.iter().any(|(key, slot)| {
names_session(key, slot, task_id) && slot_is_served_by(inner, key, slot, io)
})
}
pub(super) fn is_bound_transport(inner: &DispatchInner, io: &AdapterIo) -> bool {
inner
.transports
.values()
.any(|(transport, _)| transport.same_connection(io))
}
pub(super) fn is_revived_connection(inner: &DispatchInner, io: &AdapterIo) -> bool {
inner
.revived
.iter()
.any(|(_, revived, _)| revived.same_connection(io))
}
fn held_for_task(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
inner.revived.iter().any(|(name, revived, _)| {
revived.same_connection(io)
&& slot_key_named(inner, name)
.and_then(|key| inner.sessions.get(&key))
.is_some_and(|slot| slot_task(slot) == task_id)
})
}
pub(super) fn serves_ending(inner: &DispatchInner, task_id: &str, io: &AdapterIo) -> bool {
serves_task(inner, task_id, io)
|| held_for_task(inner, task_id, io)
|| (inner
.control_settles
.iter()
.any(|note| note.task_id == task_id)
&& is_bound_transport(inner, io))
}
pub(super) fn held_read_only(inner: &DispatchInner, key: &str, slot: &SessionSlot) -> bool {
slot.read_only
&& inner
.revived
.iter()
.any(|(name, _, _)| names_session(key, slot, name))
}
pub(super) fn record_revived_connection(
inner: &mut DispatchInner,
session_id: &str,
io: AdapterIo,
capabilities: Vec<Capability>,
) {
if !is_revived_connection(inner, &io) {
inner
.revived
.push((session_id.to_string(), io, capabilities));
}
}
fn promote_held_connection(inner: &mut DispatchInner, key: &str) {
let attached = inner
.sessions
.get(key)
.is_some_and(|slot| has_attached_transport(inner, key, slot));
if attached {
return;
}
let held = inner.revived.iter().position(|(name, _, _)| {
slot_key_named(inner, name).is_some_and(|held_key| held_key == key)
});
let Some(index) = held else { return };
let (name, io, capabilities) = inner.revived.remove(index);
if let Some(slot) = inner.sessions.get_mut(key) {
slot.read_only = false;
}
tracing::info!(
session = %name,
"a held connection takes the session its predecessor left"
);
attach_transport_locked(inner, &name, io, capabilities);
}
pub(super) fn note_binding_locked(
inner: &mut DispatchInner,
session_id: &str,
io: &AdapterIo,
) -> bool {
if is_revived_connection(inner, io) {
return false;
}
let Some(key) = slot_key_named(inner, session_id) else {
return true;
};
let Some(slot) = inner.sessions.get(&key) else {
return true;
};
let task = slot_task(slot);
let taken = attached_to_other(inner, &key, slot, io);
let moved_on = !taken
&& inner.sessions.iter().any(|(other, other_slot)| {
*other != key
&& slot_task(other_slot) == task
&& attached_to_other(inner, other, other_slot, io)
});
let revived = taken || moved_on;
let already_serving = slot_is_served_by(inner, &key, slot, io);
let returning = !revived && slot.dropped_at.is_some() && slot.ready && slot.task_id.is_some();
if let Some(slot) = inner.sessions.get_mut(&key) {
if revived {
slot.read_only = moved_on;
} else {
slot.dropped_at = None;
slot.read_only = false;
if !already_serving {
slot.last_beat = Some(Instant::now());
}
}
}
if revived {
tracing::warn!(
session = %session_id,
task = %task,
"a plugin mounted a session this client already serves; it is held read-only"
);
return false;
}
if returning {
rebase_returned_reporter(inner, &key);
}
true
}
fn rebase_returned_reporter(inner: &mut DispatchInner, key: &str) {
let Some(slot) = inner.sessions.get(key) else {
return;
};
let task_id = slot.session.task_id.clone();
let verdict = rebase_generation(inner, &task_id, |stored| {
let mut body = stored.clone();
body.agent = AgentPhase::Booting;
body.recovery = RecoveryPhase::NoRecovery;
if body.delivery == DeliveryPhase::Accepted {
body.delivery = DeliveryPhase::Pending;
}
body
});
match verdict {
Ok(Some(Verdict::Applied(next))) => tracing::info!(
task = %task_id,
generation = next.version.generation,
"a returning plugin's watermark was rebased onto a new generation"
),
Ok(Some(verdict)) => tracing::warn!(
task = %task_id,
?verdict,
"the returning plugin's watermark was not rebased"
),
Ok(None) => tracing::warn!(
task = %task_id,
"the returning plugin's row was not there to rebase"
),
Err(error) => tracing::warn!(
task = %task_id,
error = %error,
"the returning plugin's watermark was not rebased"
),
}
}
fn attach_transport_locked(
inner: &mut DispatchInner,
session_id: &str,
io: AdapterIo,
capabilities: Vec<Capability>,
) -> bool {
if !note_binding_locked(inner, session_id, &io) {
record_revived_connection(inner, session_id, io, capabilities);
return false;
}
inner
.transports
.insert(session_id.to_string(), (io, capabilities));
true
}
impl DispatchState {
pub fn hold_frame(&self, io: &AdapterIo) -> FrameGuard<'_> {
self.inner.lock().in_frame.push(io.clone());
FrameGuard {
state: self,
io: io.clone(),
}
}
pub fn bind_adapter(&self, session_id: &str, io: AdapterIo, capabilities: Vec<Capability>) {
attach_transport_locked(&mut self.inner.lock(), session_id, io, capabilities);
}
pub fn attach_msg_id(&self, task_id: &str, msg_id: &str) {
let mut inner = self.inner.lock();
let Some(key) = slot_key_serving_task(&inner, task_id) else {
return;
};
if let Some(slot) = inner.sessions.get_mut(&key) {
slot.msg_id = Some(msg_id.to_string());
}
}
pub fn plugin_send(&self, io: &AdapterIo, envelope: &Envelope) -> Result<serde_json::Value> {
let mut inner = self.inner.lock();
send_is_authorised(&inner, envelope)?;
let op_id = queue_outbound_locked(&mut inner, envelope)?;
if let Some(key) = super::guards::session_key_of_connection(&inner, io) {
super::guards::record_delivery(&mut inner, &key, &envelope.to);
}
Ok(serde_json::json!({"queued": true, "op_id": op_id}))
}
pub fn park_transport(&self, io: AdapterIo, capabilities: Vec<Capability>) {
let mut inner = self.inner.lock();
let waiting = inner
.parked
.iter()
.position(|(parked, _)| parked.same_connection(&io));
match waiting {
Some(index) => inner.parked[index].1 = capabilities,
None => inner.parked.push((io, capabilities)),
}
}
pub(super) fn claim_parked_transport(
&self,
session_id: &str,
) -> Option<(AdapterIo, Vec<Capability>)> {
let mut inner = self.inner.lock();
if inner.parked.is_empty() {
return None;
}
let (io, capabilities) = inner.parked.remove(0);
if note_binding_locked(&mut inner, session_id, &io) {
inner
.transports
.insert(session_id.to_string(), (io.clone(), capabilities.clone()));
return Some((io, capabilities));
}
tracing::warn!(
session = %session_id,
"the longest-waiting agent took nothing from this session; it waits for the next"
);
inner.parked.push((io, capabilities));
None
}
pub fn stand_transport(&self, io: AdapterIo, capabilities: Vec<Capability>) {
let mut inner = self.inner.lock();
let standing = inner
.standing
.iter()
.position(|(held, _)| held.same_connection(&io));
match standing {
Some(index) => inner.standing[index].1 = capabilities,
None => inner.standing.push((io, capabilities)),
}
}
pub(crate) fn hosted_session_ready(
&self,
session_id: &str,
opened: &onlyne_proto::OpenedArgs,
io: AdapterIo,
capabilities: Vec<Capability>,
) -> bool {
let mut inner = self.inner.lock();
let Some(slot) = inner.sessions.get_mut(session_id) else {
return false;
};
if !opened.session_id.is_empty() {
slot.session.backend_ref = serde_json::Value::String(opened.session_id.clone());
}
slot.resume_handle = opened.resume_handle.clone();
inner
.transports
.insert(session_id.to_string(), (io, capabilities));
true
}
pub fn session_transport(&self, session_id: &str) -> Option<(AdapterIo, Vec<Capability>)> {
let inner = self.inner.lock();
if let Some(transport) = inner.transports.get(session_id) {
return Some(transport.clone());
}
let key = inner
.sessions
.iter()
.find(|(key, slot)| names_session(key, slot, session_id))
.map(|(key, _)| key.clone())?;
inner.transports.get(&key).cloned()
}
pub async fn recycle_plugin(&self, task_id: &str, reason: &str, outcome: Option<Outcome>) {
let Some((io, capabilities)) = self.session_transport(task_id) else {
return;
};
if missing_capability(&capabilities, Capability::Recycle) {
tracing::debug!(
task = %task_id,
"plugin does not implement recycle; the host closes the resource"
);
return;
}
let args = RecycleArgs {
task_id: task_id.to_string(),
reason: reason.to_string(),
outcome,
};
if let Err(error) = io.notify(AdapterMsg::Host(HostOp::Recycle(args))).await {
tracing::warn!(error = %error, task = %task_id, "recycle frame did not reach the plugin");
}
}
pub async fn probe_plugin(&self, task_id: &str) -> bool {
let Some((io, _)) = self.session_transport(task_id) else {
tracing::warn!(task = %task_id, "probe found no plugin connection to ask");
return false;
};
let request = serde_json::json!({"task_id": task_id});
if let Err(error) = io.notify(AdapterMsg::Host(HostOp::Probe(request))).await {
tracing::warn!(error = %error, task = %task_id, "probe frame did not reach the plugin");
return false;
}
true
}
pub async fn nudge_plugin(&self, task_id: &str) -> bool {
let Some((io, capabilities)) = self.session_transport(task_id) else {
tracing::debug!(
task = %task_id,
"nudge found no plugin connection to hand the sentence to"
);
return false;
};
if missing_capability(&capabilities, Capability::Inject) {
tracing::debug!(
task = %task_id,
"plugin does not implement inject; the turn end settles the delivery"
);
return false;
}
let frame = AdapterMsg::Host(HostOp::Nudge {
task_id: task_id.to_string(),
text: NUDGE_TEXT.to_string(),
});
if let Err(error) = io.notify(frame).await {
tracing::warn!(error = %error, task = %task_id, "nudge frame did not reach the plugin");
return false;
}
true
}
pub fn release_connection(
&self,
session_id: Option<&str>,
io: &AdapterIo,
graceful_detach: bool,
) -> Vec<String> {
let mut inner = self.inner.lock();
let revived_connection = {
let before = inner.revived.len();
inner
.revived
.retain(|(_, revived, _)| !revived.same_connection(io));
before != inner.revived.len()
};
inner
.parked
.retain(|(parked, _)| !parked.same_connection(io));
let released: Vec<String> = inner
.transports
.iter()
.filter(|(session, (transport, _))| {
transport.same_connection(io)
&& session_id.is_none_or(|mounted| mounted == session.as_str())
})
.map(|(session, _)| session.clone())
.collect();
let served_tasks: Vec<String> = released
.iter()
.map(|session| {
inner
.sessions
.iter()
.find(|(key, slot)| names_session(key, slot, session))
.map(|(_, slot)| {
slot.task_id
.clone()
.unwrap_or_else(|| slot.session.task_id.clone())
})
.unwrap_or_else(|| session.clone())
})
.collect();
for task_id in served_tasks {
inner.stall.forget(&task_id);
}
for session in &released {
inner.transports.remove(session);
}
if !revived_connection {
let now = Instant::now();
for session in &released {
if let Some((_, slot)) = inner.sessions.iter_mut().find(|(key, slot)| {
names_session(key, slot, session)
&& (!graceful_detach || slot.task_id.is_some())
}) {
slot.dropped_at = Some(now);
}
}
}
for session in &released {
if let Some(key) = slot_key_named(&inner, session) {
promote_held_connection(&mut inner, &key);
}
}
let mut retired: Vec<String> = Vec::new();
let mut pending: Vec<PendingClose> = Vec::new();
if graceful_detach {
let idle: Vec<String> = released
.iter()
.filter_map(|session| {
inner
.sessions
.iter()
.find(|(key, slot)| {
names_session(key, slot, session) && slot.task_id.is_none()
})
.map(|(key, _)| key.clone())
})
.collect();
for key in idle {
let Some(task_id) = inner
.sessions
.get(&key)
.map(|slot| slot.session.task_id.clone())
else {
continue;
};
if keeps_idle(&inner, &key) && resumable(&inner, &key) {
if suspend_locked(&mut inner, &key, &mut pending) {
retired.push(task_id);
}
continue;
}
let reason = inner
.sessions
.get(&key)
.and_then(|slot| stored_close_reason(&inner, &slot.session.task_id))
.unwrap_or(crate::backend::CloseReason::Completed);
if retire_idle_locked(&mut inner, &key, reason, &mut pending) {
retired.push(task_id);
}
}
}
drop(inner);
close_retired(pending);
retired
}
}
fn send_is_authorised(inner: &DispatchInner, envelope: &Envelope) -> Result<()> {
if envelope.kind == MsgKind::Control {
tracing::warn!(
role = %inner.role,
to = %envelope.to,
"a plugin send naming a control command was refused"
);
return Err(anyhow!("this connection may not queue a control command"));
}
if let Some(from) = envelope.from.role_name()
&& from != inner.role
{
tracing::warn!(
role = %inner.role,
from,
"a plugin send written as another role was refused"
);
return Err(anyhow!("this role may not send as {from}"));
}
Ok(())
}
#[cfg(test)]
mod tests;