Skip to main content

mail4agent_store_sqlite/
store.rs

1//! [`SqliteMailStore`] -- the [`MailStore`] implementation this crate
2//! exists to provide. See the crate's module doc comment for the
3//! blocking-vs-async discipline every method here follows.
4
5use mail4agent_api::{Ack, Address, Message, MessageId, MessageRef, ParticipantId, RoomId, SessionCard, SessionId};
6use mail4agent_core::{
7    InsertMessageOutcome, MailStore, ParticipantRecord, ParticipantSummary, RoomRecord, RoomSummary, SecretDigest,
8    SessionRecord, StoreError,
9};
10use rusqlite::OptionalExtension;
11use serde::{Deserialize, Serialize};
12
13use crate::db::{Db, DbConfig, MigrationRunner};
14use crate::migrations::migrations;
15
16/// A [`MailStore`] backed by SQLite through this crate's own [`Db`]. Every
17/// mutating method below is exactly one transaction, matching the contract
18/// `mail4agent_core::store`'s module doc comment sets for a persistent
19/// implementation.
20pub struct SqliteMailStore {
21    db: Db,
22}
23
24impl SqliteMailStore {
25    /// Wraps an already-open [`Db`]. The caller is responsible for having
26    /// run [`crate::migrations`] against it first -- typically the same
27    /// `main.rs` boot sequence that opened `db` -- so this constructor
28    /// stays infallible and a daemon can hand in the very `Db` it opened
29    /// itself.
30    pub fn new(db: Db) -> Self {
31        Self { db }
32    }
33
34    /// Opens a fresh in-memory [`Db`] and runs this crate's own migrations
35    /// against it. Convenience for tests and small tools; a real daemon
36    /// wants a file-backed `Db` it built (and migrated) itself and should
37    /// use [`Self::new`] instead.
38    pub fn open_in_memory() -> Result<Self, StoreError> {
39        let db = Db::open(&DbConfig::in_memory())
40            .map_err(|err| StoreError::new(format!("open in-memory db: {err}")))?;
41        db.run_migrations_blocking(MigrationRunner::new(migrations()))
42            .map_err(|err| StoreError::new(format!("run migrations: {err}")))?;
43        Ok(Self { db })
44    }
45
46    /// The underlying [`Db`] handle. Clones share the same connection (see
47    /// [`Db`]'s own doc comment) -- useful when a daemon wants this store
48    /// and some other subsystem sharing one sqlite file.
49    pub fn db(&self) -> Db {
50        self.db.clone()
51    }
52}
53
54impl MailStore for SqliteMailStore {
55    fn register_participant(&mut self, id: ParticipantId, record: ParticipantRecord) -> Result<(), StoreError> {
56        self.db
57            .write_blocking(|conn| {
58                conn.execute(
59                    "INSERT INTO participants (id, label, secret_digest, may_send, may_read, operator, listener_url)
60                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
61                    rusqlite::params![
62                        id.as_str(),
63                        record.label,
64                        record.secret_digest.as_slice(),
65                        record.may_send,
66                        record.may_read,
67                        record.operator,
68                        record.listener_url,
69                    ],
70                )?;
71                Ok(())
72            })
73            .map_err(|err| StoreError::new(format!("register_participant({id}): {err}")))
74    }
75
76    fn set_listener_url(&mut self, id: &ParticipantId, url: Option<String>) -> Result<(), StoreError> {
77        self.db
78            .write_blocking(|conn| {
79                conn.execute(
80                    "UPDATE participants SET listener_url = ?1 WHERE id = ?2",
81                    rusqlite::params![url, id.as_str()],
82                )?;
83                Ok(())
84            })
85            .map_err(|err| StoreError::new(format!("set_listener_url({id}): {err}")))
86    }
87
88    fn deregister_participant(&mut self, id: &ParticipantId) -> Result<(), StoreError> {
89        self.db
90            .write_blocking(|conn| {
91                conn.execute("DELETE FROM participants WHERE id = ?1", rusqlite::params![id.as_str()])?;
92                Ok(())
93            })
94            .map_err(|err| StoreError::new(format!("deregister_participant({id}): {err}")))
95    }
96
97    fn set_participant_secret_digest(&mut self, id: &ParticipantId, digest: SecretDigest) -> Result<(), StoreError> {
98        self.db
99            .write_blocking(|conn| {
100                conn.execute(
101                    "UPDATE participants SET secret_digest = ?1 WHERE id = ?2",
102                    rusqlite::params![digest.as_slice(), id.as_str()],
103                )?;
104                Ok(())
105            })
106            .map_err(|err| StoreError::new(format!("set_participant_secret_digest({id}): {err}")))
107    }
108
109    fn get_participant(&self, id: &ParticipantId) -> Result<Option<ParticipantRecord>, StoreError> {
110        self.db
111            .read_blocking(|conn| {
112                conn.query_row(
113                    "SELECT label, secret_digest, may_send, may_read, operator, listener_url
114                       FROM participants WHERE id = ?1",
115                    rusqlite::params![id.as_str()],
116                    |row| {
117                        Ok(ParticipantRecord {
118                            label: row.get(0)?,
119                            secret_digest: digest_from_row(row, 1)?,
120                            may_send: row.get(2)?,
121                            may_read: row.get(3)?,
122                            operator: row.get(4)?,
123                            listener_url: row.get(5)?,
124                        })
125                    },
126                )
127                .optional()
128            })
129            .map_err(|err| StoreError::new(format!("get_participant({id}): {err}")))
130    }
131
132    fn find_participant_by_digest(
133        &self,
134        digest: &SecretDigest,
135    ) -> Result<Option<(ParticipantId, ParticipantRecord)>, StoreError> {
136        let found = self
137            .db
138            .read_blocking(|conn| {
139                conn.query_row(
140                    "SELECT id, label, may_send, may_read, operator, listener_url
141                       FROM participants WHERE secret_digest = ?1",
142                    rusqlite::params![digest.as_slice()],
143                    |row| {
144                        let id: String = row.get(0)?;
145                        let label: Option<String> = row.get(1)?;
146                        let may_send: bool = row.get(2)?;
147                        let may_read: bool = row.get(3)?;
148                        let operator: bool = row.get(4)?;
149                        let listener_url: Option<String> = row.get(5)?;
150                        Ok((id, label, may_send, may_read, operator, listener_url))
151                    },
152                )
153                .optional()
154            })
155            .map_err(|err| StoreError::new(format!("find_participant_by_digest: {err}")))?;
156
157        let Some((id, label, may_send, may_read, operator, listener_url)) = found else {
158            return Ok(None);
159        };
160        let participant_id = ParticipantId::new(id).map_err(|err| {
161            StoreError::new(format!("find_participant_by_digest: stored participant id failed validation: {err}"))
162        })?;
163        Ok(Some((
164            participant_id,
165            ParticipantRecord { label, secret_digest: *digest, may_send, may_read, operator, listener_url },
166        )))
167    }
168
169    fn create_room(&mut self, id: RoomId, created_at_unix_ms: u64) -> Result<(), StoreError> {
170        let created_at = i64::try_from(created_at_unix_ms)
171            .map_err(|_| StoreError::new(format!("create_room({id}): created_at_unix_ms overflows i64")))?;
172        self.db
173            .write_blocking(|conn| {
174                conn.execute(
175                    "INSERT INTO rooms (id, created_at_unix_ms) VALUES (?1, ?2)",
176                    rusqlite::params![id.as_str(), created_at],
177                )?;
178                Ok(())
179            })
180            .map_err(|err| StoreError::new(format!("create_room({id}): {err}")))
181    }
182
183    fn add_room_member(&mut self, room: &RoomId, participant: ParticipantId) -> Result<(), StoreError> {
184        self.db
185            .write_blocking(|conn| {
186                conn.execute(
187                    "INSERT INTO room_members (room_id, participant_id) VALUES (?1, ?2)
188                     ON CONFLICT (room_id, participant_id) DO NOTHING",
189                    rusqlite::params![room.as_str(), participant.as_str()],
190                )?;
191                Ok(())
192            })
193            .map_err(|err| StoreError::new(format!("add_room_member({room}, {participant}): {err}")))
194    }
195
196    fn remove_room_member(&mut self, room: &RoomId, participant: &ParticipantId) -> Result<(), StoreError> {
197        self.db
198            .write_blocking(|conn| {
199                conn.execute(
200                    "DELETE FROM room_members WHERE room_id = ?1 AND participant_id = ?2",
201                    rusqlite::params![room.as_str(), participant.as_str()],
202                )?;
203                Ok(())
204            })
205            .map_err(|err| StoreError::new(format!("remove_room_member({room}, {participant}): {err}")))
206    }
207
208    fn get_room(&self, id: &RoomId) -> Result<Option<RoomRecord>, StoreError> {
209        let created_at = self
210            .db
211            .read_blocking(|conn| {
212                conn.query_row(
213                    "SELECT created_at_unix_ms FROM rooms WHERE id = ?1",
214                    rusqlite::params![id.as_str()],
215                    |row| row.get::<_, i64>(0),
216                )
217                .optional()
218            })
219            .map_err(|err| StoreError::new(format!("get_room({id}): {err}")))?;
220
221        let Some(created_at) = created_at else {
222            return Ok(None);
223        };
224        let created_at_unix_ms = u64::try_from(created_at)
225            .map_err(|_| StoreError::new(format!("get_room({id}): stored created_at_unix_ms is negative")))?;
226
227        let member_ids: Vec<String> = self
228            .db
229            .read_blocking(|conn| {
230                let mut statement = conn.prepare("SELECT participant_id FROM room_members WHERE room_id = ?1")?;
231                let rows = statement.query_map(rusqlite::params![id.as_str()], |row| row.get(0))?;
232                rows.collect()
233            })
234            .map_err(|err| StoreError::new(format!("get_room({id}): {err}")))?;
235
236        let mut members = std::collections::BTreeSet::new();
237        for member_id in member_ids {
238            let participant_id = ParticipantId::new(member_id).map_err(|err| {
239                StoreError::new(format!("get_room({id}): stored member id failed validation: {err}"))
240            })?;
241            members.insert(participant_id);
242        }
243        Ok(Some(RoomRecord { created_at_unix_ms, members }))
244    }
245
246    fn insert_message(
247        &mut self,
248        message: Message,
249        idempotency: Option<(Address, String)>,
250    ) -> Result<InsertMessageOutcome, StoreError> {
251        let (from_kind, from_participant, from_session) = from_address_columns(&message.from)
252            .map_err(|err| StoreError::new(format!("insert_message({}): {err}", message.message_id)))?;
253        let (to_kind, to_participant, to_session, to_room) = to_address_columns(&message.to);
254        let payload = MessagePayloadWrite {
255            subject: &message.subject,
256            body: &message.body,
257            reply_to: message.reply_to.as_ref().map(MessageId::as_str),
258            correlation: message.correlation.as_deref(),
259            refs: &message.refs,
260        };
261        let payload_json = serde_json::to_string(&payload)
262            .map_err(|err| StoreError::new(format!("insert_message({}): serialize payload: {err}", message.message_id)))?;
263        let created_at_unix_ms = i64::try_from(message.created_at_unix_ms).map_err(|_| {
264            StoreError::new(format!("insert_message({}): created_at_unix_ms overflows i64", message.message_id))
265        })?;
266        // Encoded once, outside the closure, as `Address`'s own `Display`
267        // form -- see [`address_text`] for why that needs no schema change
268        // to hold a session address distinctly from its account's.
269        let idempotency_sender = idempotency.as_ref().map(|(sender, key)| (address_text(sender), key.clone()));
270
271        let raw_outcome = self
272            .db
273            .write_blocking(|conn| -> rusqlite::Result<RawInsertOutcome> {
274                let tx = conn.transaction()?;
275                tx.execute(
276                    "INSERT INTO messages
277                        (message_id, from_participant, from_kind, from_session,
278                         to_kind, to_participant, to_session, to_room, created_at_unix_ms, payload)
279                     VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
280                    rusqlite::params![
281                        message.message_id.as_str(),
282                        from_participant,
283                        from_kind,
284                        from_session,
285                        to_kind,
286                        to_participant,
287                        to_session,
288                        to_room,
289                        created_at_unix_ms,
290                        payload_json,
291                    ],
292                )?;
293
294                let outcome = match &idempotency_sender {
295                    Some((sender, key)) => match tx.execute(
296                        "INSERT INTO idempotency (sender, idempotency_key, message_id) VALUES (?1, ?2, ?3)",
297                        rusqlite::params![sender, key, message.message_id.as_str()],
298                    ) {
299                        Ok(_) => RawInsertOutcome::Inserted,
300                        Err(rusqlite::Error::SqliteFailure(sql_err, _))
301                            if sql_err.code == rusqlite::ErrorCode::ConstraintViolation =>
302                        {
303                            let existing: String = tx.query_row(
304                                "SELECT message_id FROM idempotency WHERE sender = ?1 AND idempotency_key = ?2",
305                                rusqlite::params![sender, key],
306                                |row| row.get(0),
307                            )?;
308                            RawInsertOutcome::Deduplicated(existing)
309                        }
310                        Err(other) => return Err(other),
311                    },
312                    None => RawInsertOutcome::Inserted,
313                };
314
315                // On `Inserted`, commit both statements above. On
316                // `Deduplicated`, `tx` is simply dropped here without a
317                // commit, rolling back the message row this transaction
318                // just wrote -- the idempotency ledger's UNIQUE constraint
319                // already proved another send owns this (sender, key)
320                // pair, so nothing new should persist.
321                if matches!(outcome, RawInsertOutcome::Inserted) {
322                    tx.commit()?;
323                }
324                Ok(outcome)
325            })
326            .map_err(|err| StoreError::new(format!("insert_message({}): {err}", message.message_id)))?;
327
328        match raw_outcome {
329            RawInsertOutcome::Inserted => Ok(InsertMessageOutcome::Inserted),
330            RawInsertOutcome::Deduplicated(existing) => {
331                let message_id = MessageId::new(existing).map_err(|err| {
332                    StoreError::new(format!(
333                        "insert_message({}): stored idempotency row names an invalid message id: {err}",
334                        message.message_id
335                    ))
336                })?;
337                Ok(InsertMessageOutcome::Deduplicated { message_id })
338            }
339        }
340    }
341
342    fn get_message(&self, id: &MessageId) -> Result<Option<Message>, StoreError> {
343        let row = self
344            .db
345            .read_blocking(|conn| {
346                conn.query_row(
347                    "SELECT message_id, from_participant, from_kind, from_session,
348                            to_kind, to_participant, to_session, to_room, created_at_unix_ms, payload
349                       FROM messages WHERE message_id = ?1",
350                    rusqlite::params![id.as_str()],
351                    row_to_stored_message,
352                )
353                .optional()
354            })
355            .map_err(|err| StoreError::new(format!("get_message({id}): {err}")))?;
356        row.map(assemble_message).transpose()
357    }
358
359    fn messages_to_since(&self, to: &Address, since_unix_ms: u64) -> Result<Vec<Message>, StoreError> {
360        let since = i64::try_from(since_unix_ms)
361            .map_err(|_| StoreError::new(format!("messages_to_since({to}): since_unix_ms overflows i64")))?;
362        let rows = match to {
363            Address::Direct { participant } => self
364                .db
365                .read_blocking(|conn| {
366                    let mut statement = conn.prepare(
367                        "SELECT message_id, from_participant, from_kind, from_session,
368                                to_kind, to_participant, to_session, to_room, created_at_unix_ms, payload
369                           FROM messages
370                          WHERE to_kind = 'direct' AND to_participant = ?1 AND created_at_unix_ms >= ?2
371                       ORDER BY created_at_unix_ms ASC",
372                    )?;
373                    let rows =
374                        statement.query_map(rusqlite::params![participant.as_str(), since], row_to_stored_message)?;
375                    rows.collect::<rusqlite::Result<Vec<_>>>()
376                })
377                .map_err(|err| StoreError::new(format!("messages_to_since({to}): {err}")))?,
378            Address::Session { participant, session } => self
379                .db
380                .read_blocking(|conn| {
381                    let mut statement = conn.prepare(
382                        "SELECT message_id, from_participant, from_kind, from_session,
383                                to_kind, to_participant, to_session, to_room, created_at_unix_ms, payload
384                           FROM messages
385                          WHERE to_kind = 'session' AND to_participant = ?1 AND to_session = ?2
386                                AND created_at_unix_ms >= ?3
387                       ORDER BY created_at_unix_ms ASC",
388                    )?;
389                    let rows = statement.query_map(
390                        rusqlite::params![participant.as_str(), session.as_str(), since],
391                        row_to_stored_message,
392                    )?;
393                    rows.collect::<rusqlite::Result<Vec<_>>>()
394                })
395                .map_err(|err| StoreError::new(format!("messages_to_since({to}): {err}")))?,
396            Address::Room { .. } => {
397                return Err(StoreError::new(format!(
398                    "messages_to_since({to}): a room address must use room_messages_since instead"
399                )));
400            }
401        };
402        rows.into_iter().map(assemble_message).collect()
403    }
404
405    fn room_messages_since(&self, room: &RoomId, since_unix_ms: u64) -> Result<Vec<Message>, StoreError> {
406        let since = i64::try_from(since_unix_ms)
407            .map_err(|_| StoreError::new(format!("room_messages_since({room}): since_unix_ms overflows i64")))?;
408        let rows = self
409            .db
410            .read_blocking(|conn| {
411                let mut statement = conn.prepare(
412                    "SELECT message_id, from_participant, from_kind, from_session,
413                            to_kind, to_participant, to_session, to_room, created_at_unix_ms, payload
414                       FROM messages
415                      WHERE to_kind = 'room' AND to_room = ?1 AND created_at_unix_ms >= ?2
416                   ORDER BY created_at_unix_ms ASC",
417                )?;
418                let rows = statement.query_map(rusqlite::params![room.as_str(), since], row_to_stored_message)?;
419                rows.collect::<rusqlite::Result<Vec<_>>>()
420            })
421            .map_err(|err| StoreError::new(format!("room_messages_since({room}): {err}")))?;
422        rows.into_iter().map(assemble_message).collect()
423    }
424
425    fn rooms_containing(&self, participant: &ParticipantId) -> Result<Vec<RoomId>, StoreError> {
426        let room_ids: Vec<String> = self
427            .db
428            .read_blocking(|conn| {
429                let mut statement = conn.prepare("SELECT room_id FROM room_members WHERE participant_id = ?1")?;
430                let rows = statement.query_map(rusqlite::params![participant.as_str()], |row| row.get(0))?;
431                rows.collect()
432            })
433            .map_err(|err| StoreError::new(format!("rooms_containing({participant}): {err}")))?;
434
435        room_ids
436            .into_iter()
437            .map(|id| {
438                RoomId::new(id).map_err(|err| {
439                    StoreError::new(format!("rooms_containing({participant}): stored room id failed validation: {err}"))
440                })
441            })
442            .collect()
443    }
444
445    fn record_ack(&mut self, ack: Ack) -> Result<Ack, StoreError> {
446        let acked_at = i64::try_from(ack.acked_at_unix_ms)
447            .map_err(|_| StoreError::new(format!("record_ack({}, {}): acked_at_unix_ms overflows i64", ack.message_id, ack.reader)))?;
448        let reader = address_text(&ack.reader);
449
450        // `DO UPDATE SET reader = excluded.reader` is a genuine no-op --
451        // `reader` is part of the conflict key, so it never changes -- but
452        // it is what makes SQLite treat this as an upsert rather than a
453        // skipped insert, which is what lets `RETURNING` hand back
454        // whichever row is now on file (the one this call just wrote, or
455        // an earlier ack for the same (message_id, reader) pair) in the
456        // same round trip, atomically: two concurrent acks of the same
457        // message by the same reader read back the very same
458        // `acked_at_unix_ms`, never two different ones.
459        let stored_acked_at: i64 = self
460            .db
461            .write_blocking(|conn| {
462                conn.query_row(
463                    "INSERT INTO acks (message_id, reader, acked_at_unix_ms) VALUES (?1, ?2, ?3)
464                     ON CONFLICT (message_id, reader) DO UPDATE SET reader = excluded.reader
465                     RETURNING acked_at_unix_ms",
466                    rusqlite::params![ack.message_id.as_str(), reader, acked_at],
467                    |row| row.get(0),
468                )
469            })
470            .map_err(|err| StoreError::new(format!("record_ack({}, {}): {err}", ack.message_id, ack.reader)))?;
471
472        let acked_at_unix_ms = u64::try_from(stored_acked_at).map_err(|_| {
473            StoreError::new(format!("record_ack({}, {}): stored acked_at_unix_ms is negative", ack.message_id, ack.reader))
474        })?;
475        Ok(Ack { message_id: ack.message_id, reader: ack.reader, acked_at_unix_ms })
476    }
477
478    fn get_ack(&self, message_id: &MessageId, reader: &Address) -> Result<Option<Ack>, StoreError> {
479        let reader_text = address_text(reader);
480        let row = self
481            .db
482            .read_blocking(|conn| {
483                conn.query_row(
484                    "SELECT acked_at_unix_ms FROM acks WHERE message_id = ?1 AND reader = ?2",
485                    rusqlite::params![message_id.as_str(), reader_text],
486                    |row| row.get::<_, i64>(0),
487                )
488                .optional()
489            })
490            .map_err(|err| StoreError::new(format!("get_ack({message_id}, {reader}): {err}")))?;
491
492        let Some(acked_at) = row else {
493            return Ok(None);
494        };
495        let acked_at_unix_ms = u64::try_from(acked_at)
496            .map_err(|_| StoreError::new(format!("get_ack({message_id}, {reader}): stored acked_at_unix_ms is negative")))?;
497        Ok(Some(Ack { message_id: message_id.clone(), reader: reader.clone(), acked_at_unix_ms }))
498    }
499
500    fn list_participants(&self) -> Result<Vec<ParticipantSummary>, StoreError> {
501        let rows: Vec<(String, Option<String>)> = self
502            .db
503            .read_blocking(|conn| {
504                let mut statement = conn.prepare("SELECT id, label FROM participants ORDER BY id")?;
505                let rows = statement.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?;
506                rows.collect::<rusqlite::Result<Vec<_>>>()
507            })
508            .map_err(|err| StoreError::new(format!("list_participants: {err}")))?;
509
510        rows.into_iter()
511            .map(|(id, label)| {
512                let id = ParticipantId::new(id).map_err(|err| {
513                    StoreError::new(format!("list_participants: stored participant id failed validation: {err}"))
514                })?;
515                Ok(ParticipantSummary { id, label })
516            })
517            .collect()
518    }
519
520    fn list_rooms(&self) -> Result<Vec<RoomSummary>, StoreError> {
521        // One transaction-free read over two statements against the same
522        // connection, not two separate `read_blocking` calls: this keeps
523        // the room list and the membership rows a single consistent
524        // snapshot rather than two reads a concurrent write could land
525        // between.
526        let (room_ids, member_rows): (Vec<String>, Vec<(String, String)>) = self
527            .db
528            .read_blocking(|conn| {
529                let mut room_statement = conn.prepare("SELECT id FROM rooms ORDER BY id")?;
530                let room_ids =
531                    room_statement.query_map([], |row| row.get(0))?.collect::<rusqlite::Result<Vec<String>>>()?;
532
533                let mut member_statement = conn.prepare("SELECT room_id, participant_id FROM room_members")?;
534                let member_rows = member_statement
535                    .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
536                    .collect::<rusqlite::Result<Vec<(String, String)>>>()?;
537
538                Ok((room_ids, member_rows))
539            })
540            .map_err(|err| StoreError::new(format!("list_rooms: {err}")))?;
541
542        let mut members_by_room: std::collections::HashMap<String, std::collections::BTreeSet<ParticipantId>> =
543            std::collections::HashMap::new();
544        for (room_id, participant_id) in member_rows {
545            let participant_id = ParticipantId::new(participant_id).map_err(|err| {
546                StoreError::new(format!("list_rooms: stored member id failed validation: {err}"))
547            })?;
548            members_by_room.entry(room_id).or_default().insert(participant_id);
549        }
550
551        room_ids
552            .into_iter()
553            .map(|id_str| {
554                let members = members_by_room.remove(&id_str).unwrap_or_default();
555                let id = RoomId::new(id_str)
556                    .map_err(|err| StoreError::new(format!("list_rooms: stored room id failed validation: {err}")))?;
557                Ok(RoomSummary { id, members })
558            })
559            .collect()
560    }
561
562    fn get_session(&self, session: &SessionId) -> Result<Option<SessionRecord>, StoreError> {
563        let row = self
564            .db
565            .read_blocking(|conn| {
566                conn.query_row(
567                    "SELECT account, card, last_seen_unix_ms FROM sessions WHERE session_id = ?1",
568                    rusqlite::params![session.as_str()],
569                    |row| {
570                        let account: String = row.get(0)?;
571                        let card: String = row.get(1)?;
572                        let last_seen: i64 = row.get(2)?;
573                        Ok((account, card, last_seen))
574                    },
575                )
576                .optional()
577            })
578            .map_err(|err| StoreError::new(format!("get_session({session}): {err}")))?;
579
580        let Some((account, card_json, last_seen)) = row else {
581            return Ok(None);
582        };
583        let account = ParticipantId::new(account).map_err(|err| {
584            StoreError::new(format!("get_session({session}): stored account failed validation: {err}"))
585        })?;
586        let card: SessionCard = serde_json::from_str(&card_json)
587            .map_err(|err| StoreError::new(format!("get_session({session}): stored card failed to parse: {err}")))?;
588        let last_seen_unix_ms = u64::try_from(last_seen)
589            .map_err(|_| StoreError::new(format!("get_session({session}): stored last_seen_unix_ms is negative")))?;
590        Ok(Some(SessionRecord { account, card, last_seen_unix_ms }))
591    }
592
593    fn upsert_session(&mut self, session: SessionId, record: SessionRecord) -> Result<(), StoreError> {
594        let card_json = serde_json::to_string(&record.card)
595            .map_err(|err| StoreError::new(format!("upsert_session({session}): serialize card: {err}")))?;
596        let last_seen_unix_ms = i64::try_from(record.last_seen_unix_ms)
597            .map_err(|_| StoreError::new(format!("upsert_session({session}): last_seen_unix_ms overflows i64")))?;
598        self.db
599            .write_blocking(|conn| {
600                conn.execute(
601                    "INSERT INTO sessions (session_id, account, card, last_seen_unix_ms) VALUES (?1, ?2, ?3, ?4)
602                     ON CONFLICT (session_id) DO UPDATE SET
603                         account = excluded.account,
604                         card = excluded.card,
605                         last_seen_unix_ms = excluded.last_seen_unix_ms",
606                    rusqlite::params![session.as_str(), record.account.as_str(), card_json, last_seen_unix_ms],
607                )?;
608                Ok(())
609            })
610            .map_err(|err| StoreError::new(format!("upsert_session({session}): {err}")))
611    }
612
613    fn sessions_of(&self, account: &ParticipantId) -> Result<Vec<(SessionId, SessionRecord)>, StoreError> {
614        let rows: Vec<(String, String, i64)> = self
615            .db
616            .read_blocking(|conn| {
617                let mut statement = conn.prepare(
618                    "SELECT session_id, card, last_seen_unix_ms FROM sessions WHERE account = ?1 ORDER BY session_id",
619                )?;
620                let rows = statement
621                    .query_map(rusqlite::params![account.as_str()], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))?;
622                rows.collect::<rusqlite::Result<Vec<_>>>()
623            })
624            .map_err(|err| StoreError::new(format!("sessions_of({account}): {err}")))?;
625
626        rows.into_iter()
627            .map(|(session_id, card_json, last_seen)| {
628                let session_id = SessionId::new(session_id).map_err(|err| {
629                    StoreError::new(format!("sessions_of({account}): stored session id failed validation: {err}"))
630                })?;
631                let card: SessionCard = serde_json::from_str(&card_json).map_err(|err| {
632                    StoreError::new(format!("sessions_of({account}): stored card failed to parse: {err}"))
633                })?;
634                let last_seen_unix_ms = u64::try_from(last_seen).map_err(|_| {
635                    StoreError::new(format!("sessions_of({account}): stored last_seen_unix_ms is negative"))
636                })?;
637                Ok((session_id, SessionRecord { account: account.clone(), card, last_seen_unix_ms }))
638            })
639            .collect()
640    }
641}
642
643/// What [`MailStore::insert_message`]'s own transaction decided, before
644/// its message id has been re-validated into a [`MessageId`]. Kept
645/// separate from [`InsertMessageOutcome`] because the closure passed to
646/// `write_blocking` returns a plain [`rusqlite::Result`] and has no way to
647/// surface a [`StoreError`] from a failed [`MessageId::new`] -- that
648/// validation happens once, after the transaction has already committed
649/// or rolled back.
650enum RawInsertOutcome {
651    Inserted,
652    Deduplicated(String),
653}
654
655/// A message row exactly as its ten columns hold it, before
656/// [`assemble_message`] re-validates each id and parses `payload` back
657/// into the rest of a [`Message`].
658struct StoredMessageRow {
659    message_id: String,
660    from_participant: String,
661    from_kind: String,
662    from_session: Option<String>,
663    to_kind: String,
664    to_participant: Option<String>,
665    to_session: Option<String>,
666    to_room: Option<String>,
667    created_at_unix_ms: i64,
668    payload: String,
669}
670
671fn row_to_stored_message(row: &rusqlite::Row<'_>) -> rusqlite::Result<StoredMessageRow> {
672    Ok(StoredMessageRow {
673        message_id: row.get(0)?,
674        from_participant: row.get(1)?,
675        from_kind: row.get(2)?,
676        from_session: row.get(3)?,
677        to_kind: row.get(4)?,
678        to_participant: row.get(5)?,
679        to_session: row.get(6)?,
680        to_room: row.get(7)?,
681        created_at_unix_ms: row.get(8)?,
682        payload: row.get(9)?,
683    })
684}
685
686/// The JSON shape of `messages.payload`, written by borrowing straight out
687/// of a [`Message`] the caller already owns -- see [`MessagePayloadRead`]
688/// for the owned counterpart a row is parsed back into.
689#[derive(Serialize)]
690struct MessagePayloadWrite<'a> {
691    subject: &'a str,
692    body: &'a str,
693    reply_to: Option<&'a str>,
694    correlation: Option<&'a str>,
695    refs: &'a [MessageRef],
696}
697
698/// The owned counterpart of [`MessagePayloadWrite`], for parsing
699/// `messages.payload` back out of a row. `refs` defaults on decode so a
700/// payload written before some future field addition still deserializes.
701#[derive(Deserialize)]
702struct MessagePayloadRead {
703    subject: String,
704    body: String,
705    reply_to: Option<String>,
706    correlation: Option<String>,
707    #[serde(default)]
708    refs: Vec<MessageRef>,
709}
710
711/// Reassembles a [`Message`] from a [`StoredMessageRow`], re-validating
712/// every id this store itself wrote (defence in depth against a row
713/// hand-edited outside this crate, not against anything a normal call
714/// path can produce).
715fn assemble_message(row: StoredMessageRow) -> Result<Message, StoreError> {
716    let message_id = MessageId::new(row.message_id)
717        .map_err(|err| StoreError::new(format!("stored message id failed validation: {err}")))?;
718    let from = address_from_from_columns(&row.from_kind, row.from_participant, row.from_session)?;
719    let to = address_from_to_columns(&row.to_kind, row.to_participant, row.to_session, row.to_room)?;
720    let created_at_unix_ms = u64::try_from(row.created_at_unix_ms)
721        .map_err(|_| StoreError::new("stored created_at_unix_ms is negative".to_string()))?;
722    let payload: MessagePayloadRead = serde_json::from_str(&row.payload)
723        .map_err(|err| StoreError::new(format!("stored message payload failed to parse: {err}")))?;
724    let reply_to = payload
725        .reply_to
726        .map(MessageId::new)
727        .transpose()
728        .map_err(|err| StoreError::new(format!("stored reply_to failed validation: {err}")))?;
729
730    Ok(Message {
731        message_id,
732        from,
733        to,
734        subject: payload.subject,
735        body: payload.body,
736        reply_to,
737        correlation: payload.correlation,
738        refs: payload.refs,
739        created_at_unix_ms,
740    })
741}
742
743/// Splits a message's `to` [`Address`] into `(to_kind, to_participant,
744/// to_session, to_room)` the way `messages` stores a recipient -- a
745/// session gets its own `to_session` column rather than being packed into
746/// `to_participant` alongside a direct account address, so the two kinds
747/// can never be confused by a query that forgets to check `to_kind` first.
748fn to_address_columns(address: &Address) -> (&'static str, Option<&str>, Option<&str>, Option<&str>) {
749    match address {
750        Address::Direct { participant } => ("direct", Some(participant.as_str()), None, None),
751        Address::Session { participant, session } => {
752            ("session", Some(participant.as_str()), Some(session.as_str()), None)
753        }
754        Address::Room { room } => ("room", None, None, Some(room.as_str())),
755    }
756}
757
758/// The inverse of [`to_address_columns`].
759fn address_from_to_columns(
760    kind: &str,
761    to_participant: Option<String>,
762    to_session: Option<String>,
763    to_room: Option<String>,
764) -> Result<Address, StoreError> {
765    match kind {
766        "direct" => {
767            let participant = to_participant.ok_or_else(|| {
768                StoreError::new("stored message has to_kind = direct but to_participant is NULL".to_string())
769            })?;
770            let participant = ParticipantId::new(participant)
771                .map_err(|err| StoreError::new(format!("stored direct recipient id failed validation: {err}")))?;
772            Ok(Address::Direct { participant })
773        }
774        "session" => {
775            let participant = to_participant.ok_or_else(|| {
776                StoreError::new("stored message has to_kind = session but to_participant is NULL".to_string())
777            })?;
778            let participant = ParticipantId::new(participant).map_err(|err| {
779                StoreError::new(format!("stored session recipient account id failed validation: {err}"))
780            })?;
781            let session = to_session.ok_or_else(|| {
782                StoreError::new("stored message has to_kind = session but to_session is NULL".to_string())
783            })?;
784            let session = SessionId::new(session)
785                .map_err(|err| StoreError::new(format!("stored session recipient id failed validation: {err}")))?;
786            Ok(Address::Session { participant, session })
787        }
788        "room" => {
789            let room = to_room.ok_or_else(|| {
790                StoreError::new("stored message has to_kind = room but to_room is NULL".to_string())
791            })?;
792            let room = RoomId::new(room)
793                .map_err(|err| StoreError::new(format!("stored room recipient id failed validation: {err}")))?;
794            Ok(Address::Room { room })
795        }
796        other => Err(StoreError::new(format!("stored message has unknown to_kind {other:?}"))),
797    }
798}
799
800/// Splits a message's `from` [`Address`] into `(from_kind, from_participant,
801/// from_session)`. Unlike [`to_address_columns`], `from` is never a room --
802/// `Message::validate`'s `validate_participant_address` already rules that
803/// out before a message ever reaches this store -- so a room reaching here
804/// is named as a [`StoreError`] rather than silently coerced into some
805/// other shape.
806fn from_address_columns(address: &Address) -> Result<(&'static str, &str, Option<&str>), StoreError> {
807    match address {
808        Address::Direct { participant } => Ok(("direct", participant.as_str(), None)),
809        Address::Session { participant, session } => Ok(("session", participant.as_str(), Some(session.as_str()))),
810        Address::Room { room } => {
811            Err(StoreError::new(format!("a message's `from` must not be a room address (got \"{room}\")")))
812        }
813    }
814}
815
816/// The inverse of [`from_address_columns`].
817fn address_from_from_columns(
818    kind: &str,
819    from_participant: String,
820    from_session: Option<String>,
821) -> Result<Address, StoreError> {
822    let participant = ParticipantId::new(from_participant)
823        .map_err(|err| StoreError::new(format!("stored sender id failed validation: {err}")))?;
824    match kind {
825        "direct" => Ok(Address::Direct { participant }),
826        "session" => {
827            let session = from_session.ok_or_else(|| {
828                StoreError::new("stored message has from_kind = session but from_session is NULL".to_string())
829            })?;
830            let session = SessionId::new(session)
831                .map_err(|err| StoreError::new(format!("stored sender session id failed validation: {err}")))?;
832            Ok(Address::Session { participant, session })
833        }
834        other => Err(StoreError::new(format!("stored message has unknown from_kind {other:?}"))),
835    }
836}
837
838/// Encodes a participant-identifying [`Address`] (never a room -- see
839/// [`Ack::validate`] and [`Message::validate`], the only two places an
840/// address reaches `acks.reader` or `idempotency.sender`) as the single
841/// TEXT value those two columns already held before this crate's addresses
842/// could name a session. `Address`'s own `Display` -- `"claude"` for an
843/// account, `"claude/s-7f3a..."` for one of its sessions -- is exactly
844/// that value, and the two shapes never collide (see `Address`'s own
845/// `Display`/`FromStr` doc comment: neither `ParticipantId`'s nor
846/// `SessionId`'s charset permits `/`). That is why neither `acks` nor
847/// `idempotency` needed a v2 migration at all: a v1 row's `reader`/`sender`
848/// was already exactly an account's `Display` form, and a session's
849/// `Display` form is simply a new, longer string the same column always
850/// could have held.
851fn address_text(address: &Address) -> String {
852    address.to_string()
853}
854
855/// A tiny local error so [`digest_from_row`] can report a length mismatch
856/// through rusqlite's own `FromSqlConversionFailure` rather than panicking
857/// on the `TryFrom` it uses to size the digest down to `[u8; 32]`.
858#[derive(Debug)]
859struct DigestLengthError(usize);
860
861impl std::fmt::Display for DigestLengthError {
862    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
863        write!(f, "secret_digest column must be exactly 32 bytes, got {}", self.0)
864    }
865}
866
867impl std::error::Error for DigestLengthError {}
868
869fn digest_from_row(row: &rusqlite::Row<'_>, idx: usize) -> rusqlite::Result<SecretDigest> {
870    let bytes: Vec<u8> = row.get(idx)?;
871    let len = bytes.len();
872    <[u8; 32]>::try_from(bytes)
873        .map_err(|_| rusqlite::Error::FromSqlConversionFailure(idx, rusqlite::types::Type::Blob, Box::new(DigestLengthError(len))))
874}