use super::*;
use super::outbound::send_frame;
use super::state::{DispatchInner, DispatchState, binding_task_state, slot_key_named};
pub fn with_cluster(state: &DispatchState, report: Report) -> Report {
let cluster = state.cluster_ref();
if cluster.is_empty() {
return report;
}
match report {
Report::Ready {
task_id,
session_id,
generation,
seq,
..
} => Report::Ready {
task_id,
session_id,
generation,
seq,
cluster_ref: Some(cluster),
},
Report::Heartbeat {
projection: None,
task_id,
session_id,
generation,
seq,
observed,
..
} => Report::Heartbeat {
task_id,
session_id,
generation,
seq,
observed,
projection: None,
cluster_ref: Some(cluster),
},
Report::Complete {
task_id,
outcome,
head,
details,
files,
reply_to,
cluster_ref: _,
} => Report::Complete {
task_id,
outcome,
head,
details,
files,
reply_to,
cluster_ref: Some(cluster),
},
other => other,
}
}
pub fn projection_of(row: &SessionRecord, task_state: TaskState) -> SessionProjection {
let observed: Option<serde_json::Value> = serde_json::from_str(&row.observed_json).ok();
let lifecycle = project(
phase(&row.agent_state, AgentPhase::Booting),
phase(&row.delivery_state, DeliveryPhase::NoIntent),
phase(&row.resource_state, ResourcePhase::Detached),
phase(&row.recovery_substate, RecoveryPhase::NoRecovery),
task_state,
);
SessionProjection {
lifecycle,
agent: phase(&row.agent_state, AgentPhase::Booting),
delivery: phase(&row.delivery_state, DeliveryPhase::NoIntent),
resource: phase(&row.resource_state, ResourcePhase::Detached),
recovery: phase(&row.recovery_substate, RecoveryPhase::NoRecovery),
outcome: task_outcome_of(task_state),
observed,
}
}
pub fn task_state_of(outcome: Outcome) -> TaskState {
match outcome {
Outcome::Done => TaskState::Done,
Outcome::Failed => TaskState::Failed,
Outcome::Cancelled => TaskState::Cancelled,
Outcome::Blocked => TaskState::Blocked,
}
}
pub fn task_outcome_of(task_state: TaskState) -> Option<Outcome> {
match task_state {
TaskState::Pending => None,
TaskState::Done => Some(Outcome::Done),
TaskState::Failed => Some(Outcome::Failed),
TaskState::Cancelled => Some(Outcome::Cancelled),
TaskState::Blocked => Some(Outcome::Blocked),
}
}
pub(super) fn stored_task_state(inner: &DispatchInner, task_id: &str) -> TaskState {
inner
.store
.task(task_id)
.ok()
.flatten()
.map(|record| record.task_state)
.unwrap_or(TaskState::Pending)
}
pub(super) fn phase<T: serde::de::DeserializeOwned>(word: &str, fallback: T) -> T {
serde_json::from_value(serde_json::Value::String(word.to_string())).unwrap_or(fallback)
}
pub async fn sync_session(state: &DispatchState, session_id: &str) -> Result<()> {
let Some(op) = sync_frame(state, session_id)? else {
return Ok(());
};
send_frame(state, op).await
}
pub fn sync_frame(state: &DispatchState, session_id: &str) -> Result<Option<ClientOp>> {
let (key, row, task_state) = {
let inner = state.inner.lock();
let key = slot_key_named(&inner, session_id);
let task_state = match key.as_deref().and_then(|key| inner.sessions.get(key)) {
Some(slot) => binding_task_state(&inner, slot),
None => match inner.store.get_session(session_id)? {
Some(row) => stored_task_state(&inner, &row.task_id),
None => TaskState::Pending,
},
};
(key, inner.store.get_session(session_id)?, task_state)
};
let Some(row) = row else { return Ok(None) };
let session_id = key.unwrap_or_else(|| row.session_id.clone());
let projection = projection_of(&row, task_state);
Ok(Some(ClientOp::Report(Report::Heartbeat {
task_id: row.task_id.clone(),
session_id,
generation: row.generation.max(0) as u64,
seq: row.seq.max(0) as u64,
observed: projection
.observed
.clone()
.unwrap_or(serde_json::Value::Null),
cluster_ref: None,
projection: Some(projection),
})))
}
pub fn note_verdict(verdict: &Verdict, task_id: &str) -> Option<Version> {
match verdict {
Verdict::Applied(observation) => Some(observation.version),
Verdict::Ignored(reason) => {
tracing::debug!(task = %task_id, ?reason, "lifecycle event ignored");
None
}
Verdict::Rejected(reason) => {
tracing::warn!(task = %task_id, ?reason, "lifecycle event rejected; the ledger kept its state");
None
}
}
}
pub fn note_intent_receipt(state: &DispatchState, task_id: &str) {
let verdict = {
let inner = state.inner.lock();
feed_intent_receipt(&inner.bridge, &inner.store, task_id)
};
match verdict {
Ok(verdict) => {
note_verdict(&verdict, task_id);
}
Err(error) => tracing::debug!(
task = %task_id,
error = %error,
"intent receipt not recorded in the session ledger"
),
}
}