Skip to main content

onlyne_store/
session.rs

1//! The session mirror's row shape and the port its writers go through.
2//!
3//! [`SessionRecord`], [`VersionedSession`], and [`FaultRecord`] mirror the
4//! columns the `sessions` and `faults` tables carry in both databases, and
5//! [`SessionLedger`] is the seam the session reducer's persistence bridge writes
6//! them through. This crate implements the port; the client's reconcile module
7//! drives it.
8//!
9//! A stored session row answers for a session, which is the table's key in both
10//! databases; the delivery a session serves is a binding of its own. The port
11//! stays delivery-addressed because that is how the reducer holds a session —
12//! one tuple per delivery it opened — and the implementation resolves the
13//! delivery to its session through the binding.
14
15use serde::{Deserialize, Serialize};
16
17/// The only persistence surface the session bridge needs. The write gate is
18/// monotonic: `upsert_session` applies a row only when its `(generation, seq)` is
19/// strictly newer than the stored watermark, and reports whether the row changed.
20pub trait SessionLedger: Send + Sync {
21    /// Load the session row serving one delivery.
22    fn get_session(&self, task_id: &str) -> anyhow::Result<Option<SessionRecord>>;
23    /// Monotonic session upsert, which also opens the delivery's binding.
24    /// Returns true when the row changed.
25    fn upsert_session(&self, task_id: &str, version: &VersionedSession) -> anyhow::Result<bool>;
26    /// Whether the task itself is tracked, so a stray report cannot conjure a row.
27    fn task_is_known(&self, task_id: &str) -> anyhow::Result<bool>;
28    /// Attempt count for a task, for the fault audit entry.
29    fn task_attempt(&self, task_id: &str) -> anyhow::Result<i64>;
30    /// Existing faults for one task, for the `(task_id, kind, generation)` dedupe.
31    fn list_faults(&self, task_id: &str) -> anyhow::Result<Vec<FaultRecord>>;
32    /// Append one fault row. Returns its id.
33    fn insert_fault(&self, fault: &FaultRecord) -> anyhow::Result<i64>;
34    /// Observe an emitted event. The bridge never reads back from this hook.
35    fn emit(&self, kind: &str, data: serde_json::Value);
36    /// Observe an operator-visible alert line.
37    fn note_alert(&self, line: String);
38}
39
40/// One stored session tuple, as the ledger columns hold it. No public view is
41/// among them: the projection is derived by whoever asks, from these columns
42/// plus the task state that caller owns.
43#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
44pub struct SessionRecord {
45    /// The row's own key. A session serves one delivery at a time and keeps its
46    /// own id while that binding moves, so the two are separate facts: this one
47    /// names the session, and `task_id` names the delivery the caller asked
48    /// about.
49    pub session_id: String,
50    /// The delivery this session serves, as the caller named it. Not the row's
51    /// key: the session id is.
52    pub task_id: String,
53    pub agent_state: String,
54    pub delivery_state: String,
55    pub resource_state: String,
56    pub recovery_substate: String,
57    pub desired_json: String,
58    pub observed_json: String,
59    pub generation: i64,
60    pub seq: i64,
61    pub backend_ref: String,
62    pub mismatch_count: i64,
63    pub updated_at: i64,
64}
65
66/// One projected row ready for the ledger write gate.
67#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
68pub struct VersionedSession {
69    pub agent_state: String,
70    pub delivery_state: String,
71    pub resource_state: String,
72    pub recovery_substate: String,
73    pub desired_json: String,
74    pub observed_json: String,
75    pub generation: i64,
76    pub seq: i64,
77    pub backend_ref: String,
78    pub mismatch_count: i64,
79    pub updated_at: i64,
80}
81
82/// One fault queue row.
83#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
84pub struct FaultRecord {
85    pub id: i64,
86    pub task_id: String,
87    pub session_id: String,
88    pub generation: i64,
89    pub seq: i64,
90    pub desired_json: String,
91    pub observed_json: String,
92    pub intent: String,
93    pub attempt: i64,
94    pub backend_ref: String,
95    pub kind: String,
96    pub reason: String,
97    pub state: String,
98    pub created_at: i64,
99}