Skip to main content

lash_sqlite_store/process_registry/
support.rs

1use super::*;
2
3pub(super) fn process_status_label(record: &ProcessRecord) -> &'static str {
4    record.status.label()
5}
6
7impl SqliteProcessRegistry {
8    /// Open a process registry whose terminal-retention prune removes the two
9    /// process-owned session stores from `session_store_root` before the process
10    /// row. The root is required and explicit; no sibling-directory convention
11    /// is inferred.
12    pub async fn open(
13        path: &Path,
14        session_store_root: impl Into<PathBuf>,
15    ) -> tokio_rusqlite::Result<Self> {
16        Self::open_with_clock(path, Arc::new(lash_core::SystemClock), session_store_root).await
17    }
18
19    pub async fn open_with_clock(
20        path: &Path,
21        clock: Arc<dyn lash_core::Clock>,
22        session_store_root: impl Into<PathBuf>,
23    ) -> tokio_rusqlite::Result<Self> {
24        Self::open_configured(path, clock, session_store_root.into()).await
25    }
26
27    async fn open_configured(
28        path: &Path,
29        clock: Arc<dyn lash_core::Clock>,
30        process_session_store_root: PathBuf,
31    ) -> tokio_rusqlite::Result<Self> {
32        let conn = SqliteConnection::open(path).await?;
33        ensure_process_schema(&conn).await?;
34        apply_pragmas(&conn, StoreBacking::File).await?;
35        Ok(Self {
36            conn,
37            clock,
38            process_session_store_root: Some(process_session_store_root),
39        })
40    }
41
42    pub async fn memory() -> tokio_rusqlite::Result<Self> {
43        Self::memory_with_clock(Arc::new(lash_core::SystemClock)).await
44    }
45
46    pub async fn memory_with_clock(
47        clock: Arc<dyn lash_core::Clock>,
48    ) -> tokio_rusqlite::Result<Self> {
49        let conn = SqliteConnection::open_in_memory().await?;
50        ensure_process_schema(&conn).await?;
51        apply_pragmas(&conn, StoreBacking::Memory).await?;
52        Ok(Self {
53            conn,
54            clock,
55            process_session_store_root: None,
56        })
57    }
58
59    pub(crate) fn load_process_conn(
60        conn: &Connection,
61        process_id: &str,
62    ) -> Result<Option<ProcessRecord>, lash_core::PluginError> {
63        let json: Option<String> = conn
64            .query_row(
65                "SELECT record_json FROM processes WHERE process_id = ?1",
66                params![process_id],
67                |row| row.get(0),
68            )
69            .optional()
70            .map_err(process_sqlite_error)?;
71        json.map(|json| serde_json::from_str(&json).map_err(process_decode_error))
72            .transpose()
73    }
74
75    pub(crate) fn save_process_conn(
76        conn: &Connection,
77        record: &ProcessRecord,
78    ) -> Result<(), lash_core::PluginError> {
79        let change_seq = Self::next_change_seq_conn(conn)?;
80        conn.execute(
81            "UPDATE processes
82             SET updated_at_ms = ?2, change_seq = ?3, status = ?4, record_json = ?5
83             WHERE process_id = ?1",
84            params![
85                record.id.as_str(),
86                record.updated_at_ms as i64,
87                change_seq as i64,
88                process_status_label(record),
89                process_encode_json(record)?
90            ],
91        )
92        .map_err(process_sqlite_error)?;
93        Ok(())
94    }
95
96    pub(super) fn next_change_seq_conn(conn: &Connection) -> Result<u64, lash_core::PluginError> {
97        conn.execute(
98            "UPDATE process_change_clock
99             SET current_seq = current_seq + 1
100             WHERE singleton = 1",
101            [],
102        )
103        .map_err(process_sqlite_error)?;
104        conn.query_row(
105            "SELECT current_seq FROM process_change_clock WHERE singleton = 1",
106            [],
107            |row| row.get::<_, i64>(0),
108        )
109        .map(|seq| seq as u64)
110        .map_err(process_sqlite_error)
111    }
112
113    pub(crate) fn load_event_by_key_conn(
114        conn: &Connection,
115        process_id: &str,
116        replay_key: &str,
117    ) -> Result<Option<(String, ProcessEvent)>, lash_core::PluginError> {
118        let row: Option<(String, String)> = conn
119            .query_row(
120                "SELECT payload_hash, event_json
121                 FROM process_events
122                 WHERE process_id = ?1 AND idempotency_key = ?2",
123                params![process_id, replay_key],
124                |row| Ok((row.get(0)?, row.get(1)?)),
125            )
126            .optional()
127            .map_err(process_sqlite_error)?;
128        row.map(|(hash, json)| {
129            serde_json::from_str(&json)
130                .map(|event| (hash, event))
131                .map_err(process_decode_error)
132        })
133        .transpose()
134    }
135
136    pub(crate) fn load_process_lease_conn(
137        conn: &Connection,
138        process_id: &str,
139    ) -> Result<Option<ProcessLease>, lash_core::PluginError> {
140        conn.query_row(
141            "SELECT lease_owner_id, lease_token, lease_fencing_token,
142                    lease_claimed_at_ms, lease_expires_at_ms,
143                    lease_owner_incarnation_id, lease_owner_liveness_json
144             FROM process_leases
145             WHERE process_id = ?1",
146            params![process_id],
147            |row| {
148                let owner_id: Option<String> = row.get(0)?;
149                let lease_token: Option<String> = row.get(1)?;
150                let incarnation_id: Option<String> = row.get(5)?;
151                let liveness_json: Option<String> = row.get(6)?;
152                let (Some(owner_id), Some(lease_token)) = (owner_id, lease_token) else {
153                    return Ok(None);
154                };
155                Ok(Some(ProcessLease {
156                    schema_version: PROCESS_LEASE_SCHEMA_VERSION,
157                    process_id: process_id.to_string(),
158                    owner: process_lease_owner_from_columns(
159                        owner_id,
160                        incarnation_id,
161                        liveness_json,
162                    ),
163                    lease_token,
164                    fencing_token: row.get::<_, i64>(2)? as u64,
165                    claimed_at_epoch_ms: row.get::<_, i64>(3)? as u64,
166                    expires_at_epoch_ms: row.get::<_, i64>(4)? as u64,
167                }))
168            },
169        )
170        .optional()
171        .map(|lease| lease.flatten())
172        .map_err(process_sqlite_error)
173    }
174
175    /// Insert-or-replace the persisted lease row for `process_id` with a fresh
176    /// lease owned by `owner` at `fencing_token`.
177    pub(super) fn acquire_process_lease_conn(
178        conn: &Connection,
179        process_id: &str,
180        owner: &LeaseOwnerIdentity,
181        fencing_token: u64,
182        now: u64,
183        lease_ttl_ms: u64,
184    ) -> Result<ProcessLease, lash_core::PluginError> {
185        let lease = ProcessLease {
186            schema_version: PROCESS_LEASE_SCHEMA_VERSION,
187            process_id: process_id.to_string(),
188            owner: owner.clone(),
189            lease_token: format!(
190                "{:x}",
191                Sha256::digest(
192                    format!(
193                        "{process_id}:{}:{}:{now}:{fencing_token}",
194                        owner.owner_id, owner.incarnation_id
195                    )
196                    .as_bytes()
197                )
198            ),
199            fencing_token,
200            claimed_at_epoch_ms: now,
201            expires_at_epoch_ms: now.saturating_add(lease_ttl_ms),
202        };
203        conn.execute(
204            "INSERT INTO process_leases (
205                process_id, lease_owner_id, lease_owner_incarnation_id,
206                lease_owner_liveness_json, lease_token, lease_fencing_token,
207                lease_claimed_at_ms, lease_expires_at_ms
208             )
209             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
210             ON CONFLICT(process_id) DO UPDATE SET
211                lease_owner_id = excluded.lease_owner_id,
212                lease_owner_incarnation_id = excluded.lease_owner_incarnation_id,
213                lease_owner_liveness_json = excluded.lease_owner_liveness_json,
214                lease_token = excluded.lease_token,
215                lease_fencing_token = excluded.lease_fencing_token,
216                lease_claimed_at_ms = excluded.lease_claimed_at_ms,
217                lease_expires_at_ms = excluded.lease_expires_at_ms",
218            params![
219                lease.process_id.as_str(),
220                lease.owner.owner_id.as_str(),
221                lease.owner.incarnation_id.as_str(),
222                encode_process_lease_liveness(&lease.owner.liveness)?,
223                lease.lease_token.as_str(),
224                lease.fencing_token as i64,
225                lease.claimed_at_epoch_ms as i64,
226                lease.expires_at_epoch_ms as i64,
227            ],
228        )
229        .map_err(process_sqlite_error)?;
230        Ok(lease)
231    }
232
233    pub(super) fn list_grants_for_scope_conn(
234        conn: &Connection,
235        session_scope: &SessionScope,
236        live_only: bool,
237    ) -> Result<Vec<ProcessHandleGrantEntry>, lash_core::PluginError> {
238        let session_scope_id = session_scope.id();
239        let status_clause = if live_only {
240            "AND p.status = 'running'"
241        } else {
242            ""
243        };
244        let mut stmt = conn
245            .prepare(&format!(
246                "SELECT g.process_id, g.descriptor_json, p.record_json
247                 FROM process_handle_grants g
248                 JOIN processes p ON p.process_id = g.process_id
249                 WHERE g.scope_id = ?1 {status_clause}
250                 ORDER BY g.process_id ASC"
251            ))
252            .map_err(process_sqlite_error)?;
253        let rows = stmt
254            .query_map(params![session_scope_id.as_str()], |row| {
255                Ok((
256                    row.get::<_, String>(0)?,
257                    row.get::<_, String>(1)?,
258                    row.get::<_, String>(2)?,
259                ))
260            })
261            .map_err(process_sqlite_error)?;
262        let mut entries = Vec::new();
263        for row in rows {
264            let (process_id, descriptor_json, record_json) = row.map_err(process_sqlite_error)?;
265            let descriptor: ProcessHandleDescriptor =
266                serde_json::from_str(&descriptor_json).map_err(process_decode_error)?;
267            let record: ProcessRecord =
268                serde_json::from_str(&record_json).map_err(process_decode_error)?;
269            entries.push((
270                ProcessHandleGrant {
271                    session_id: session_scope.session_id.clone(),
272                    process_id,
273                    descriptor,
274                },
275                record,
276            ));
277        }
278        Ok(entries)
279    }
280}
281
282/// Map a `Result<T, PluginError>` produced by a synchronous transaction body to
283/// a [`TxOutcome`]: commit on success, roll back on logical error. Both arms
284/// carry the inner `Result` back so the caller recovers the value or the
285/// `PluginError` after the transaction resolves.
286pub(crate) fn tx_outcome<T>(
287    result: Result<T, lash_core::PluginError>,
288) -> TxOutcome<Result<T, lash_core::PluginError>> {
289    match result {
290        Ok(value) => TxOutcome::Commit(Ok(value)),
291        Err(err) => TxOutcome::Rollback(Err(err)),
292    }
293}