use super::*;
use super::projection::{note_verdict, sync_session};
use super::retire::on_recycled;
use super::settle::{SettleAuthority, on_out};
use super::state::{ControlWord, note_beat};
use super::transport::{serves_ending, serves_task};
use super::turn_end::on_turn_end;
pub async fn on_control(state: &DispatchState, op: &ControlOp) -> Result<bool> {
let task_id = op.task_id();
let held = state.holds_task(task_id);
if !held {
tracing::warn!(
op = op.name(),
task = %task_id,
"a control command naming no session this role holds was refused; nothing was applied"
);
return Ok(false);
}
match op {
ControlOp::Recycle { reason, .. } => {
state.owe_controlled_settle(task_id, ControlWord::Recycle, Instant::now());
state.recycle_plugin(task_id, reason, None).await;
on_recycled(state, task_id, crate::backend::CloseReason::Operator)?;
}
ControlOp::Cancel { reason, .. } => {
state.owe_controlled_settle(task_id, ControlWord::Cancel, Instant::now());
state
.recycle_plugin(task_id, reason, Some(Outcome::Cancelled))
.await;
on_recycled(state, task_id, crate::backend::CloseReason::Cancelled)?;
}
ControlOp::Probe { .. } => {
if state.probe_plugin(task_id).await {
sync_session(state, task_id).await?;
}
}
ControlOp::Snapshot { .. } => sync_session(state, task_id).await?,
ControlOp::Focus { .. } => {
let (backend, session) = {
let inner = state.inner.lock();
let session = inner
.sessions
.values()
.find(|slot| slot.task_id.as_deref() == Some(task_id))
.map(|slot| slot.session.clone());
(inner.backend.clone(), session)
};
match session {
Some(session_ref) => {
if let Err(error) = backend.focus(&session_ref) {
tracing::warn!(error = %error, task = %task_id, "focus refused");
if let Err(fault_error) = on_plugin_report(
state,
None,
Report::Fault {
task_id: Some(task_id.to_string()),
session_id: None,
generation: None,
seq: None,
kind: "focus".into(),
reason: error.to_string(),
desired: None,
observed: None,
},
)
.await
{
tracing::warn!(
error = %fault_error,
task = %task_id,
"focus refused"
);
}
}
}
None => {
tracing::warn!(task = %task_id, "focus has no live session");
}
}
}
}
Ok(true)
}
pub async fn on_plugin_report(
state: &DispatchState,
from: Option<&AdapterIo>,
report: Report,
) -> Result<()> {
let subject = task_id_of(&report).to_string();
let touched = match report {
Report::Ready { task_id, .. } => {
let verdict = {
let inner = state.inner.lock();
feed_ready(&inner.bridge, &inner.store, &task_id)?
};
note_verdict(&verdict, &task_id).is_some()
}
Report::Heartbeat {
task_id,
generation: _,
seq,
observed,
..
} => {
let (touched, ended) = {
let mut inner = state.inner.lock();
if from.is_some_and(|io| !serves_task(&inner, &task_id, io)) {
note_beat(&mut inner, &task_id, Instant::now());
tracing::warn!(
task = %task_id,
"a beat arrived on a connection that is not this task's transport: \
the liveness stamp is refreshed and no state is applied"
);
(false, false)
} else {
let row = inner.store.get_session(&task_id).ok().flatten();
let stored = stored_observation(&inner.store, row.as_ref());
let generation = stored.version.generation;
let task_state = inner
.store
.task(&task_id)
.ok()
.flatten()
.map(|record| record.task_state);
if !inner.stall.note_beat_seq(&task_id, generation, seq) {
note_beat(&mut inner, &task_id, Instant::now());
tracing::debug!(
task = %task_id,
generation,
seq,
"a heartbeat at or below the watermark was taken as a replay; it refreshes liveness only"
);
return Ok(());
}
let (ended, verdict) = match serde_json::from_value::<Observation>(observed) {
Ok(body) => {
let composed = compose_observation(&stored, body, task_state);
let ended = stored.agent == AgentPhase::Running
&& composed.agent == AgentPhase::Idle
&& task_state == Some(TaskState::Pending)
&& composed.recovery == RecoveryPhase::IdleWaiting;
let verdict = apply_persist(
&inner.bridge,
&inner.store,
&task_id,
&LifecycleEvent::Heartbeat {
v: Version::new(generation, seq),
body: composed,
},
)?;
(ended, verdict)
}
Err(error) => {
tracing::warn!(task = %task_id, error = %error, "heartbeat carries no readable observation; liveness only");
note_beat(&mut inner, &task_id, Instant::now());
(false, Verdict::Ignored(IgnoredReason::NoOp))
}
};
note_verdict(&verdict, &task_id);
let touched = match &verdict {
Verdict::Applied(_) => {
note_beat(&mut inner, &task_id, Instant::now());
inner.stall.note_applied(&task_id, Instant::now());
true
}
Verdict::Ignored(IgnoredReason::NoOp) => {
note_beat(&mut inner, &task_id, Instant::now());
true
}
Verdict::Ignored(_) | Verdict::Rejected(_) => false,
};
(touched, ended && !matches!(verdict, Verdict::Rejected(_)))
}
};
if ended {
on_turn_end(state, &task_id, None).await?;
}
touched
}
Report::Complete {
task_id,
outcome,
head,
details,
..
} => {
if !ending_is_authorised(state, from, &task_id) {
tracing::warn!(
task = %task_id,
"a completion from a connection serving no such session was refused"
);
return Ok(());
}
let asked = if state.take_controlled_settle(&task_id) {
SettleAuthority::ControlDriven
} else {
SettleAuthority::PluginReport
};
on_out(state, &task_id, outcome, head, details, asked).await?;
false
}
Report::Fault {
task_id: Some(task_id),
kind,
reason,
..
} => {
if !ending_is_authorised(state, from, &task_id) {
tracing::warn!(
task = %task_id,
kind = %kind,
"a fault from a connection serving no such session was refused"
);
return Ok(());
}
let inner = state.inner.lock();
crate::reconcile::record_fault(&inner.store, &task_id, &kind, "plugin", &reason)?;
false
}
Report::Fault { task_id: None, .. } => false,
};
if touched {
sync_session(state, &subject).await?;
}
Ok(())
}
fn ending_is_authorised(state: &DispatchState, from: Option<&AdapterIo>, task_id: &str) -> bool {
let Some(io) = from else {
return true;
};
let inner = state.inner.lock();
serves_ending(&inner, task_id, io)
}
fn compose_observation(
client: &Observation,
mut body: Observation,
task_state: Option<TaskState>,
) -> Observation {
body.delivery = client.delivery;
body.recovery = client.recovery;
body.generation_live = client.generation_live;
body.isolate_after = client.isolate_after;
body.terminate_after = client.terminate_after;
body.mismatch_count = client.mismatch_count;
if body.agent == AgentPhase::Ready && body.delivery == DeliveryPhase::Exhausted {
body.delivery = DeliveryPhase::NoIntent;
}
if body.agent == AgentPhase::Booting && body.delivery == DeliveryPhase::Accepted {
body.delivery = DeliveryPhase::Pending;
}
let held = match body.recovery {
RecoveryPhase::IdleWaiting | RecoveryPhase::IdleFault => body.agent == AgentPhase::Idle,
RecoveryPhase::Draining => matches!(body.agent, AgentPhase::Idle | AgentPhase::Running),
RecoveryPhase::NoRecovery => true,
};
if !held {
body.recovery = RecoveryPhase::NoRecovery;
}
if body.agent == AgentPhase::Idle
&& task_state == Some(TaskState::Pending)
&& body.delivery != DeliveryPhase::Accepted
&& body.recovery == RecoveryPhase::NoRecovery
{
body.recovery = RecoveryPhase::IdleWaiting;
}
body
}
fn task_id_of(report: &Report) -> &str {
match report {
Report::Ready { task_id, .. } | Report::Heartbeat { task_id, .. } => task_id,
Report::Complete { task_id, .. } => task_id,
Report::Fault { .. } => "",
}
}
#[cfg(test)]
mod tests;