use onlyne_proto::lifecycle::{
DeliveryPhase, IgnoredReason, LifecycleEvent, Observation, RecoveryPhase, Verdict,
};
use super::bridge::{Bridge, apply_at_next, apply_from_stored};
use onlyne_store::session::SessionLedger;
pub fn feed_created(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
if ledger.get_session(task_id)?.is_some() {
tracing::debug!(task = %task_id, "session row already exists; created is a no-op");
return Ok(Verdict::Ignored(IgnoredReason::NoOp));
}
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::Created { v })
}
pub fn feed_dispatched(bridge: &Bridge, ledger: &dyn SessionLedger, task_id: &str) {
if let Err(err) = feed_created(bridge, ledger, task_id) {
tracing::warn!(task = %task_id, error = %err, "could not seed the session row");
}
match feed_resource_attached(bridge, ledger, task_id) {
Ok(verdict) => {
if matches!(verdict, Verdict::Rejected(_)) {
tracing::warn!(task = %task_id, ?verdict, "resource attach refused for a dispatched session");
}
}
Err(err) => {
tracing::warn!(task = %task_id, error = %err, "could not record the resource attach")
}
}
}
pub fn feed_resource_attached(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| {
LifecycleEvent::ResourceAttach { v }
})
}
pub fn feed_ready(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::Ready { v })
}
pub fn feed_turn_started(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::TurnStarted {
v,
})
}
pub fn feed_turn_ended(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::TurnEnded { v })
}
pub fn feed_intent_receipt(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::IntentReceipt {
v,
})
}
pub fn feed_resource_closed(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| {
LifecycleEvent::ResourceClosed { v }
})
}
pub fn feed_agent_gone(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::AgentGone { v })
}
pub fn feed_suspended(
bridge: &Bridge,
ledger: &dyn SessionLedger,
session_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, session_id, |v| LifecycleEvent::Suspend {
v,
})
}
pub fn feed_resumed(
bridge: &Bridge,
ledger: &dyn SessionLedger,
session_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, session_id, |v| LifecycleEvent::Resume { v })
}
pub fn feed_fail(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::Fail { v })
}
pub fn feed_cancel(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::Cancel { v })
}
pub fn feed_mismatch(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| {
LifecycleEvent::ReconcileMismatch { v }
})
}
pub fn feed_reconcile_ok(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::ReconcileOk {
v,
})
}
pub fn feed_delivered(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_at_next(bridge, ledger, task_id, |v| LifecycleEvent::Complete { v })?;
settle(bridge, ledger, task_id)
}
pub fn settle(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
) -> anyhow::Result<Verdict> {
apply_from_stored(bridge, ledger, task_id, |current, version| {
LifecycleEvent::Heartbeat {
v: version,
body: settle_body(current),
}
})
}
fn settle_body(obs: &Observation) -> Observation {
let recovery = match obs.recovery {
RecoveryPhase::Draining | RecoveryPhase::IdleWaiting => RecoveryPhase::NoRecovery,
other => other,
};
Observation::build(
obs.version,
obs.generation_live,
obs.isolate_after,
obs.terminate_after,
obs.mismatch_count,
obs.agent,
DeliveryPhase::Accepted,
obs.resource,
recovery,
)
.with_host(obs.host.clone())
}