use serde::Serialize;
use serde_json::json;
use onlyne_proto::lifecycle::{
self, AgentPhase, DeliveryPhase, Observation, RecoveryPhase, ResourcePhase, Version,
};
use onlyne_store::session::{SessionLedger, SessionRecord, VersionedSession};
use super::bridge::{Bridge, initial_observation};
use super::fault::{DEFAULT_ISOLATE_AFTER, DEFAULT_TERMINATE_AFTER};
#[derive(Debug, Clone, Default, Serialize)]
pub struct FaultOutcome {
pub fault_id: Option<i64>,
}
impl FaultOutcome {
pub(super) fn recorded(&self) -> bool {
self.fault_id.is_some()
}
}
fn tag<T: Serialize>(value: &T) -> anyhow::Result<String> {
match serde_json::to_value(value)? {
serde_json::Value::String(text) => Ok(text),
other => anyhow::bail!("lifecycle state must serialize to a string: {other}"),
}
}
fn counter(value: u64) -> i64 {
i64::try_from(value).unwrap_or(i64::MAX)
}
pub(super) fn now_unix() -> i64 {
chrono::Utc::now().timestamp()
}
pub(super) fn short(task_id: &str) -> &str {
&task_id[..8.min(task_id.len())]
}
pub fn to_versioned(
obs: &Observation,
backend_ref: &str,
desired_json: &str,
) -> anyhow::Result<VersionedSession> {
Ok(VersionedSession {
agent_state: tag(&obs.agent)?,
delivery_state: tag(&obs.delivery)?,
resource_state: tag(&obs.resource)?,
recovery_substate: tag(&obs.recovery)?,
desired_json: desired_json.to_string(),
observed_json: serde_json::to_string(obs)?,
generation: counter(obs.version.generation),
seq: counter(obs.version.seq),
backend_ref: backend_ref.to_string(),
mismatch_count: counter(u64::from(obs.mismatch_count)),
updated_at: now_unix(),
})
}
pub(super) fn backend_ref_json(
bridge: &Bridge,
task_id: &str,
row: Option<&SessionRecord>,
) -> String {
if let Some(session) = bridge.live.lock().get(task_id) {
if let Ok(json) = serde_json::to_string(session) {
return json;
}
}
if let Some(row) = row {
let stored = row.backend_ref.trim();
if !stored.is_empty() && stored != "{}" {
return stored.to_string();
}
}
"{}".to_string()
}
pub fn stored_observation(ledger: &dyn SessionLedger, row: Option<&SessionRecord>) -> Observation {
let Some(row) = row else {
return initial_observation();
};
match serde_json::from_str::<Observation>(&row.observed_json) {
Ok(obs) if lifecycle::is_legal(&obs) => Observation {
version: Version::new(
if row.generation > 0 {
row.generation as u64
} else {
obs.version.generation
},
row.seq.max(0) as u64,
),
..obs
},
Ok(obs) => corrupt_observation(ledger, row, &format!("illegal tuple {obs:?}")),
Err(err) => corrupt_observation(ledger, row, &err.to_string()),
}
}
fn corrupt_observation(
ledger: &dyn SessionLedger,
row: &SessionRecord,
reason: &str,
) -> Observation {
tracing::warn!(
task = %row.task_id,
reason,
"session row is corrupt; rebuilt the reducer watermark from its columns"
);
ledger.note_alert(format!(
"session {} observed_json is corrupt",
short(&row.task_id)
));
ledger.emit(
"lifecycle_corrupt",
json!({"task_id": row.task_id, "reason": reason}),
);
Observation::build(
Version::new(
if row.generation > 0 {
row.generation as u64
} else {
1
},
row.seq.max(0) as u64,
),
true,
DEFAULT_ISOLATE_AFTER,
DEFAULT_TERMINATE_AFTER,
row.mismatch_count.clamp(0, i64::from(u32::MAX)) as u32,
AgentPhase::Booting,
DeliveryPhase::NoIntent,
ResourcePhase::Detached,
RecoveryPhase::NoRecovery,
)
}