use std::collections::HashMap;
use parking_lot::{Mutex, MutexGuard};
use serde_json::json;
use crate::backend::SessionRef;
use onlyne_proto::lifecycle::{self, LifecycleEvent, Observation, Verdict, Version};
use super::fault::{DEFAULT_ISOLATE_AFTER, DEFAULT_TERMINATE_AFTER};
use onlyne_store::session::{SessionLedger, SessionRecord};
use super::record::{backend_ref_json, short, stored_observation, to_versioned};
#[derive(Debug, Default)]
pub struct Bridge {
pub(super) live: Mutex<HashMap<String, SessionRef>>,
applying: Mutex<()>,
}
impl Bridge {
pub fn new() -> Self {
Self::default()
}
fn applying(&self) -> MutexGuard<'_, ()> {
self.applying.lock()
}
pub fn track_live(&self, session: SessionRef) {
self.live.lock().insert(session.task_id.clone(), session);
}
pub fn untrack_live(&self, task_id: &str) {
self.live.lock().remove(task_id);
}
}
pub(super) fn initial_observation() -> Observation {
Observation::initial(DEFAULT_ISOLATE_AFTER, DEFAULT_TERMINATE_AFTER)
}
fn version_of(event: &LifecycleEvent) -> Version {
match event {
LifecycleEvent::Created { v }
| LifecycleEvent::Ready { v }
| LifecycleEvent::TurnStarted { v }
| LifecycleEvent::TurnEnded { v }
| LifecycleEvent::Heartbeat { v, .. }
| LifecycleEvent::Complete { v }
| LifecycleEvent::IntentPending { v }
| LifecycleEvent::IntentRetry { v }
| LifecycleEvent::IntentReceipt { v }
| LifecycleEvent::IntentExhausted { v }
| LifecycleEvent::ResourceAttach { v }
| LifecycleEvent::ResourceCloseRequested { v }
| LifecycleEvent::ResourceClosed { v }
| LifecycleEvent::Suspend { v }
| LifecycleEvent::Resume { v }
| LifecycleEvent::AgentGone { v }
| LifecycleEvent::Cancel { v }
| LifecycleEvent::Fail { v }
| LifecycleEvent::ReconcileMismatch { v }
| LifecycleEvent::ReconcileOk { v }
| LifecycleEvent::AdoptNewGeneration { v }
| LifecycleEvent::Supersede { v, .. } => *v,
}
}
fn event_name(event: &LifecycleEvent) -> &'static str {
match event {
LifecycleEvent::Created { .. } => "created",
LifecycleEvent::Ready { .. } => "ready",
LifecycleEvent::TurnStarted { .. } => "turn_started",
LifecycleEvent::TurnEnded { .. } => "turn_ended",
LifecycleEvent::Heartbeat { .. } => "heartbeat",
LifecycleEvent::Complete { .. } => "complete",
LifecycleEvent::IntentPending { .. } => "intent_pending",
LifecycleEvent::IntentRetry { .. } => "intent_retry",
LifecycleEvent::IntentReceipt { .. } => "intent_receipt",
LifecycleEvent::IntentExhausted { .. } => "intent_exhausted",
LifecycleEvent::ResourceAttach { .. } => "resource_attach",
LifecycleEvent::ResourceCloseRequested { .. } => "resource_close_requested",
LifecycleEvent::ResourceClosed { .. } => "resource_closed",
LifecycleEvent::Suspend { .. } => "suspend",
LifecycleEvent::Resume { .. } => "resume",
LifecycleEvent::AgentGone { .. } => "agent_gone",
LifecycleEvent::Cancel { .. } => "cancel",
LifecycleEvent::Fail { .. } => "fail",
LifecycleEvent::ReconcileMismatch { .. } => "reconcile_mismatch",
LifecycleEvent::ReconcileOk { .. } => "reconcile_ok",
LifecycleEvent::AdoptNewGeneration { .. } => "adopt_new_generation",
LifecycleEvent::Supersede { .. } => "supersede",
}
}
pub fn apply_persist(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
event: &LifecycleEvent,
) -> anyhow::Result<Verdict> {
let (verdict, landed) = apply_persist_reported(bridge, ledger, task_id, event)?;
if !landed {
report_lost_write(ledger, task_id, event, &verdict, 1);
}
Ok(verdict)
}
fn apply_persist_reported(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
event: &LifecycleEvent,
) -> anyhow::Result<(Verdict, bool)> {
let _gate = bridge.applying();
let row = ledger.get_session(task_id)?;
apply_round(bridge, ledger, task_id, row.as_ref(), event)
}
fn apply_round(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
row: Option<&SessionRecord>,
event: &LifecycleEvent,
) -> anyhow::Result<(Verdict, bool)> {
if row.is_none() {
let known = ledger.task_is_known(task_id)? || bridge.live.lock().contains_key(task_id);
if !known {
tracing::warn!(
task = %task_id,
event = event_name(event),
"lifecycle event for a session never tracked; nothing persisted"
);
anyhow::bail!("unknown session {task_id}");
}
}
let current = stored_observation(ledger, row);
if row.is_none() && matches!(event, LifecycleEvent::Created { .. }) {
let (seeded, landed) = seed_created(bridge, ledger, task_id, version_of(event))?;
return Ok((Verdict::Applied(seeded), landed));
}
let stamped = if row.is_some() {
stamp_adoption(event, ¤t)
} else {
None
};
if let Some(adopting) = stamped.as_ref() {
tracing::warn!(
task = %task_id,
event = event_name(event),
from = version_of(event).seq,
to = version_of(adopting).seq,
watermark = current.version.seq,
"adoption carried a sequence the stored watermark had passed; advanced it past the gate"
);
}
let event = stamped.as_ref().unwrap_or(event);
let verdict = lifecycle::apply(¤t, event);
let landed = record_verdict(bridge, ledger, task_id, row, event, ¤t, &verdict)?;
Ok((verdict, landed))
}
fn stamp_adoption(event: &LifecycleEvent, current: &Observation) -> Option<LifecycleEvent> {
let v = version_of(event);
let adoption = matches!(
event,
LifecycleEvent::AdoptNewGeneration { .. } | LifecycleEvent::Supersede { .. }
);
if !adoption || v.generation != current.version.generation || v.seq > current.version.seq {
return None;
}
let newer = Version::new(v.generation, current.version.seq.saturating_add(1));
Some(match event {
LifecycleEvent::Supersede {
old_generation_dead,
body,
..
} => LifecycleEvent::Supersede {
v: newer,
old_generation_dead: *old_generation_dead,
body: body.clone(),
},
_ => LifecycleEvent::AdoptNewGeneration { v: newer },
})
}
fn report_lost_write(
ledger: &dyn SessionLedger,
task_id: &str,
event: &LifecycleEvent,
verdict: &Verdict,
attempts: usize,
) {
let v = version_of(event);
tracing::warn!(
task = %task_id,
event = event_name(event),
generation = v.generation,
seq = v.seq,
attempts,
applied = matches!(verdict, Verdict::Applied(_)),
"session transition did not land; the row keeps the older tuple"
);
ledger.note_alert(format!(
"session {} lost its {} write to a newer watermark",
short(task_id),
event_name(event)
));
ledger.emit(
"lifecycle_write_lost",
json!({
"task_id": task_id,
"event": event_name(event),
"generation": v.generation,
"seq": v.seq,
"attempts": attempts,
}),
);
}
fn seed_created(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
version: Version,
) -> anyhow::Result<(Observation, bool)> {
let mut seeded = initial_observation();
seeded.version = version;
let backend_ref = backend_ref_json(bridge, task_id, None);
let desired = serde_json::to_string(&LifecycleEvent::Created { v: version })?;
let stored = to_versioned(&seeded, &backend_ref, &desired)?;
if ledger.upsert_session(task_id, &stored)? {
ledger.emit("lifecycle", transition_payload(task_id, &seeded, "created"));
Ok((seeded, true))
} else {
tracing::warn!(
task = %task_id,
"a newer session row appeared while seeding Created; kept the newer watermark"
);
Ok((seeded, false))
}
}
fn transition_payload(task_id: &str, to: &Observation, event: &str) -> serde_json::Value {
json!({
"task_id": task_id,
"agent": to.agent,
"delivery": to.delivery,
"resource": to.resource,
"recovery": to.recovery,
"generation": to.version.generation,
"seq": to.version.seq,
"event": event,
})
}
#[allow(clippy::too_many_arguments)]
fn record_verdict(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
row: Option<&SessionRecord>,
event: &LifecycleEvent,
current: &Observation,
verdict: &Verdict,
) -> anyhow::Result<bool> {
match verdict {
Verdict::Applied(next) => {
let backend_ref = backend_ref_json(bridge, task_id, row);
let desired = serde_json::to_string(event)?;
let stored = to_versioned(next, &backend_ref, &desired)?;
if !ledger.upsert_session(task_id, &stored)? {
tracing::warn!(
task = %task_id,
event = event_name(event),
generation = next.version.generation,
seq = next.version.seq,
"session write lost to a newer watermark; left the row alone"
);
return Ok(false);
}
tracing::debug!(
task = %task_id,
agent = ?next.agent,
delivery = ?next.delivery,
resource = ?next.resource,
recovery = ?next.recovery,
generation = next.version.generation,
seq = next.version.seq,
"lifecycle transition applied"
);
ledger.emit(
"lifecycle",
transition_payload(task_id, next, event_name(event)),
);
Ok(true)
}
Verdict::Ignored(reason) => {
tracing::debug!(
task = %task_id,
reason = ?reason,
event = event_name(event),
generation = current.version.generation,
seq = current.version.seq,
"lifecycle event ignored; ledger left as-is"
);
Ok(true)
}
Verdict::Rejected(reason) => {
tracing::warn!(
task = %task_id,
reason = ?reason,
event = event_name(event),
generation = current.version.generation,
seq = current.version.seq,
"lifecycle event rejected; ledger left as-is"
);
Ok(true)
}
}
}
pub fn next_version(ledger: &dyn SessionLedger, task_id: &str) -> anyhow::Result<Version> {
let row = ledger.get_session(task_id)?;
let current = stored_observation(ledger, row.as_ref());
Ok(next_after(¤t))
}
fn next_after(current: &Observation) -> Version {
Version::new(
current.version.generation,
current.version.seq.saturating_add(1),
)
}
pub fn apply_at_next(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
make: impl FnMut(Version) -> LifecycleEvent,
) -> anyhow::Result<Verdict> {
let mut make = make;
apply_locally(bridge, ledger, task_id, |_, v| make(v))
}
pub fn apply_from_stored(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
make: impl FnMut(&Observation, Version) -> LifecycleEvent,
) -> anyhow::Result<Verdict> {
apply_locally(bridge, ledger, task_id, make)
}
fn apply_locally(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
mut make: impl FnMut(&Observation, Version) -> LifecycleEvent,
) -> anyhow::Result<Verdict> {
let _gate = bridge.applying();
let mut attempt = 1;
loop {
let row = ledger.get_session(task_id)?;
let current = stored_observation(ledger, row.as_ref());
let event = make(¤t, next_after(¤t));
let (verdict, landed) = apply_round(bridge, ledger, task_id, row.as_ref(), &event)?;
if landed {
return Ok(verdict);
}
if attempt == APPLY_ATTEMPTS {
report_lost_write(ledger, task_id, &event, &verdict, attempt);
return Ok(verdict);
}
tracing::debug!(
task = %task_id,
event = event_name(&event),
attempt,
"write lost to a competing watermark; replaying the round on the fresher row"
);
attempt += 1;
}
}
const APPLY_ATTEMPTS: usize = 4;
pub fn try_feed(
bridge: &Bridge,
ledger: &dyn SessionLedger,
task_id: &str,
make: impl FnMut(Version) -> LifecycleEvent,
) {
if let Err(err) = apply_at_next(bridge, ledger, task_id, make) {
tracing::warn!(
task = %task_id,
error = %err,
"could not record a lifecycle event in the session ledger"
);
}
}