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}