use super::*;
use super::outbound::send_frame;
use super::state::{DispatchInner, DispatchState};
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,
reply_to,
..
} => Report::Complete {
task_id,
outcome,
head,
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 = wire_lifecycle(project(
phase(&row.agent_state, AgentState::Booting),
phase(&row.delivery_state, DeliveryState::None),
phase(&row.resource_state, ResourceState::Detached),
phase(&row.recovery_substate, RecoveryState::None),
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,
}
}
fn wire_lifecycle(lifecycle: PublicLifecycle) -> Lifecycle {
match lifecycle {
PublicLifecycle::Created => Lifecycle::Created,
PublicLifecycle::Working => Lifecycle::Working,
PublicLifecycle::Idle => Lifecycle::Idle,
PublicLifecycle::Exited => Lifecycle::Exited,
}
}
pub fn task_state_of(outcome: Outcome) -> TaskState {
match outcome {
Outcome::Done => TaskState::Done,
Outcome::Failed => TaskState::Failed,
Outcome::Cancelled => TaskState::Cancelled,
}
}
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),
}
}
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, task_id: &str) -> Result<()> {
let (row, task_state) = {
let inner = state.inner.lock();
(
inner.store.get_session(task_id)?,
stored_task_state(&inner, task_id),
)
};
let Some(row) = row else { return Ok(()) };
let projection = projection_of(&row, task_state);
send_frame(
state,
ClientOp::Report(Report::Heartbeat {
task_id: row.task_id.clone(),
session_id: row.task_id.clone(),
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),
}),
)
.await
}
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"
),
}
}