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_session;
pub async fn on_control(state: &DispatchState, op: &ControlOp) -> Result<bool> {
let task_id = op.task_id();
let held = state.holds_task(task_id);
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, onlyne_session::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, onlyne_session::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(held)
}
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 mut inner = state.inner.lock();
if from.is_some_and(|io| !serves_session(&inner, &task_id, io)) {
note_beat(&mut inner, &task_id, Instant::now());
tracing::warn!(
task = %task_id,
"a beat from a connection held read-only refreshes the liveness stamp and applies no state"
);
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);
let verdict = match serde_json::from_value::<Observation>(observed) {
Ok(body) => apply_persist(
&inner.bridge,
&inner.store,
&task_id,
&LifecycleEvent::Heartbeat {
v: Version::new(generation, seq),
body: compose_observation(&stored, body, task_state),
},
)?,
Err(error) => {
tracing::warn!(task = %task_id, error = %error, "heartbeat carries no readable observation; liveness only");
note_beat(&mut inner, &task_id, Instant::now());
Verdict::Ignored(IgnoredReason::NoOp)
}
};
note_verdict(&verdict, &task_id);
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());
inner
.store
.bump_session_version(&task_id, generation, seq)?
}
Verdict::Ignored(_) | Verdict::Rejected(_) => false,
}
}
}
Report::Complete {
task_id,
outcome,
head,
..
} => {
let asked = if state.take_controlled_settle(&task_id) {
SettleAuthority::ControlDriven
} else {
SettleAuthority::PluginReport
};
on_out(state, &task_id, outcome, head, None, &[], asked).await?;
false
}
Report::Fault {
task_id: Some(task_id),
kind,
reason,
..
} => {
let inner = state.inner.lock();
onlyne_session::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 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 == AgentState::Ready && body.delivery == DeliveryState::Exhausted {
body.delivery = DeliveryState::None;
}
if body.agent == AgentState::Booting && body.delivery == DeliveryState::Accepted {
body.delivery = DeliveryState::Pending;
}
let held = match body.recovery {
RecoveryState::IdleWaiting | RecoveryState::IdleFault => body.agent == AgentState::Idle,
RecoveryState::Draining => matches!(body.agent, AgentState::Idle | AgentState::Running),
RecoveryState::None => true,
};
if !held {
body.recovery = RecoveryState::None;
}
if body.agent == AgentState::Idle
&& task_state == Some(TaskState::Pending)
&& body.delivery != DeliveryState::Accepted
&& body.recovery == RecoveryState::None
{
body.recovery = RecoveryState::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;