use super::*;
use super::projection::stored_task_state;
use super::settle::{SettleAuthority, on_out};
use super::state::{DispatchInner, slot_key_serving_task};
use super::transport::NUDGE_TEXT;
pub const TURN_END_WITHOUT_COMPLETE: &str = "turn_end_without_complete";
pub const DELIVERY_BLOCKED: &str = "delivery_blocked";
pub const HANDOFF: &str = "handoff";
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
enum TurnEnd {
#[default]
Fresh,
Nudged,
Blocked,
}
#[derive(Debug, Default)]
pub(super) struct TurnEndWatch {
states: HashMap<String, TurnEnd>,
}
impl TurnEndWatch {
fn of(&self, task_id: &str) -> TurnEnd {
self.states.get(task_id).copied().unwrap_or_default()
}
fn set(&mut self, task_id: &str, state: TurnEnd) {
self.states.insert(task_id.to_string(), state);
}
fn forget(&mut self, task_id: &str) {
self.states.remove(task_id);
}
fn tasks(&self) -> impl Iterator<Item = &String> + '_ {
self.states.keys()
}
}
fn prune(inner: &mut DispatchInner) {
let stale: Vec<String> = inner
.turn_end
.tasks()
.filter(|task_id| {
!inner
.sessions
.values()
.any(|slot| slot.task_id.as_deref() == Some(task_id.as_str()))
|| stored_task_state(inner, task_id) != TaskState::Pending
})
.cloned()
.collect();
for task_id in stale {
inner.turn_end.forget(&task_id);
}
}
enum Nudge {
Plugin,
Backend {
backend: Arc<dyn SessionBackend>,
session: SessionRef,
},
}
enum Step {
Quiet,
Nudge(Nudge),
Settle(&'static str),
}
pub async fn on_turn_end(
state: &DispatchState,
task_id: &str,
closing: Option<String>,
) -> Result<()> {
let (step, ending) = {
let mut inner = state.inner.lock();
decide(&mut inner, task_id)?
};
if let Some(op) = ending {
state.enqueue_op(&op)?;
}
match step {
Step::Quiet => Ok(()),
Step::Nudge(Nudge::Plugin) => {
if state.nudge_plugin(task_id).await {
Ok(())
} else {
settle(
state,
task_id,
closing,
"the nudge did not reach the plugin",
)
.await
}
}
Step::Nudge(Nudge::Backend { backend, session }) => {
state.feed_turn_started(task_id);
match backend.nudge(&session, task_id, NUDGE_TEXT) {
Ok(()) => Ok(()),
Err(error) => {
settle(
state,
task_id,
closing,
&format!("the nudge did not reach the agent: {error}"),
)
.await
}
}
}
Step::Settle(why) => settle(state, task_id, closing, why).await,
}
}
async fn settle(
state: &DispatchState,
task_id: &str,
closing: Option<String>,
why: &str,
) -> Result<()> {
let op = {
let mut inner = state.inner.lock();
inner.turn_end.set(task_id, TurnEnd::Blocked);
let session_id =
slot_key_serving_task(&inner, task_id).unwrap_or_else(|| task_id.to_string());
record_blocked(&inner, task_id, &session_id, why)
};
state.enqueue_op(&op)?;
tracing::info!(task = %task_id, reason = why, "a delivery settled blocked at its turn end");
on_out(
state,
task_id,
Outcome::Blocked,
closing,
None,
SettleAuthority::ClientOwned,
)
.await
}
fn decide(inner: &mut DispatchInner, task_id: &str) -> Result<(Step, Option<ClientOp>)> {
prune(inner);
let Some(key) = slot_key_serving_task(inner, task_id) else {
tracing::debug!(task = %task_id, "a turn ended for a task no session of this role serves");
return Ok((Step::Quiet, None));
};
let session = inner
.sessions
.get(&key)
.filter(|slot| !slot.read_only && !slot.suspended)
.map(|slot| slot.session.clone());
let Some(session) = session else {
tracing::debug!(task = %task_id, session = %key, "a turn ended for a session that serves nothing");
return Ok((Step::Quiet, None));
};
if stored_task_state(inner, task_id) != TaskState::Pending {
inner.turn_end.forget(task_id);
return Ok((Step::Quiet, None));
}
let step = match inner.turn_end.of(task_id) {
TurnEnd::Fresh => {
let target = if inner.backend.self_driven() {
Nudge::Backend {
backend: inner.backend.clone(),
session,
}
} else {
Nudge::Plugin
};
inner.turn_end.set(task_id, TurnEnd::Nudged);
Step::Nudge(target)
}
TurnEnd::Nudged => Step::Settle("a second turn ended without a completion"),
TurnEnd::Blocked => return Ok((Step::Quiet, None)),
};
let op = record_ending(inner, task_id, &key, matches!(step, Step::Nudge(_)));
Ok((step, Some(op)))
}
pub fn record_handoff(
state: &DispatchState,
task_id: &str,
to_role: &str,
hop: u32,
text: &str,
) -> ClientOp {
let inner = state.inner.lock();
let mut payload = serde_json::json!({
"task_id": task_id,
"role": inner.role.as_str(),
"to_role": to_role,
"hop": hop,
"text": text,
});
if let Some(session_id) = slot_key_serving_task(&inner, task_id) {
payload["session_id"] = serde_json::Value::String(session_id);
}
drop(inner);
publish(HANDOFF, payload)
}
fn record_ending(inner: &DispatchInner, task_id: &str, session_id: &str, nudge: bool) -> ClientOp {
publish(
TURN_END_WITHOUT_COMPLETE,
serde_json::json!({
"task_id": task_id,
"session_id": session_id,
"role": inner.role.as_str(),
"nudge": nudge,
}),
)
}
fn record_blocked(inner: &DispatchInner, task_id: &str, session_id: &str, why: &str) -> ClientOp {
publish(
DELIVERY_BLOCKED,
serde_json::json!({
"task_id": task_id,
"session_id": session_id,
"role": inner.role.as_str(),
"reason": why,
}),
)
}
fn publish(class: &str, payload: serde_json::Value) -> ClientOp {
ClientOp::PublishEvent(onlyne_proto::PublishEventArgs {
class: class.to_string(),
payload,
})
}