Skip to main content

mail4agent_core/
store.rs

1//! The mailbox's own storage boundary: the [`MailStore`] trait plus
2//! [`InMemoryStore`], the implementation this crate's tests run against.
3//!
4//! **This crate does not write SQLite.** Persistence is a separate task; what
5//! belongs here is the shape of the boundary, chosen so a SQLite
6//! implementation can satisfy it with one transaction per mutating method --
7//! each mutating method below corresponds to exactly one engine-level
8//! mutation (register a participant, send a message, record an ack, ...),
9//! never a fragment of one, so nothing needs two round trips to stay
10//! consistent.
11//!
12//! Every method returns `Result<_, StoreError>`: a store is a real
13//! dependency (disk, a lock, a connection) that can fail independently of
14//! anything a caller did wrong, and a fallible implementation cannot satisfy
15//! an infallible trait without panicking -- unacceptable in a service whose
16//! job is to stay reachable. [`crate::MailboxEngine`] maps a [`StoreError`]
17//! into [`mail4agent_api::MailError::StoreUnavailable`], naming only the
18//! failing operation; the [`StoreError`] itself is logged via
19//! `tracing::error!` and never crosses the wire (see the engine's
20//! `store_unavailable` helper).
21
22use std::collections::{BTreeSet, HashMap};
23
24use mail4agent_api::{Ack, Address, Message, MessageId, ParticipantId, RoomId, SessionCard, SessionId};
25use thiserror::Error;
26
27/// The SHA-256 digest of a participant's secret. The engine stores only
28/// this, never the secret itself -- see `MailboxEngine::authenticate`'s doc
29/// comment for why an exact-digest index is safe to look up directly rather
30/// than scanned.
31pub type SecretDigest = [u8; 32];
32
33/// An error from a [`MailStore`] implementation itself -- I/O, a lock, a
34/// corrupt row -- as distinct from a domain refusal
35/// ([`mail4agent_api::MailError`]). Carries a message meant for
36/// `tracing::error!`, never for a caller: see the module doc comment.
37#[derive(Debug, Error)]
38#[error("{0}")]
39pub struct StoreError(pub String);
40
41impl StoreError {
42    pub fn new(message: impl Into<String>) -> Self {
43        Self(message.into())
44    }
45}
46
47/// A registered participant, as the mailbox's own registry holds it.
48///
49/// `label` is display metadata only. `secret_digest` is the SHA-256 digest
50/// of the participant's bearer secret -- the plaintext is generated once, on
51/// registration or rotation, handed back to the caller, and never stored.
52/// `may_send` and `may_read` gate [`crate::MailboxEngine::send`] and the
53/// read-side operations respectively; `operator` additionally lets a
54/// participant read any *named* address one at a time (see
55/// `crate::MailboxEngine::inbox_of`), never the whole mailbox in one call.
56#[derive(Clone, Debug, Eq, PartialEq)]
57pub struct ParticipantRecord {
58    pub label: Option<String>,
59    pub secret_digest: SecretDigest,
60    pub may_send: bool,
61    pub may_read: bool,
62    pub operator: bool,
63    /// The URL the mailbox POSTs a [`mail4agent_api::DeliveryNotification`]
64    /// to when mail arrives for this account or any of its sessions --
65    /// `None` until `crate::MailboxEngine::set_listener` registers one.
66    /// Structural validation (bounded, loopback-only) happens once, in
67    /// that engine method, before this is ever written -- this store keeps
68    /// whatever string it is handed, the same way `label` does.
69    pub listener_url: Option<String>,
70}
71
72/// A room the mailbox tracks membership for. Membership is explicit and
73/// mailbox-owned -- never recomputed from a foreign graph (see
74/// `mail4agent/CLAUDE.md` and
75/// `mailbox-service-extraction-and-signed-session-identity-2026-09-16.md`
76/// ยง5b: this is what makes a room's readability immune to something
77/// unrelated growing too large elsewhere).
78#[derive(Clone, Debug, Eq, PartialEq)]
79pub struct RoomRecord {
80    pub created_at_unix_ms: u64,
81    pub members: BTreeSet<ParticipantId>,
82}
83
84/// A directory-listable summary of a registered participant: enough to
85/// list it, and no more. Deliberately excludes [`ParticipantRecord`]'s
86/// `secret_digest`, `may_send`, `may_read` and `operator` -- a directory
87/// answers "who exists", never "what may they do" or anything that would
88/// help forge one, and a type that structurally has no `secret_digest`
89/// field cannot leak one even by accident, regardless of what
90/// [`MailStore::list_participants`]'s implementation does internally. See
91/// `mail4agent_api::DirectoryEntry`, the wire type this is assembled into.
92#[derive(Clone, Debug, Eq, PartialEq)]
93pub struct ParticipantSummary {
94    pub id: ParticipantId,
95    pub label: Option<String>,
96}
97
98/// A directory-listable summary of a room: its id and current membership,
99/// as raw fact -- not yet filtered through any one caller's point of view.
100/// `crate::MailboxEngine::directory` is what turns "who is a member" into
101/// "is the caller a member".
102#[derive(Clone, Debug, Eq, PartialEq)]
103pub struct RoomSummary {
104    pub id: RoomId,
105    pub members: BTreeSet<ParticipantId>,
106}
107
108/// One live session, as the mailbox's own registry holds it. Carries no
109/// secret and no permission bits of its own -- a session authenticates
110/// through its account's bearer secret plus kernel attestation of the
111/// calling process (`mail4agent-attest`, outside this crate entirely), and
112/// borrows its account's `may_send`/`may_read`/`operator` bits rather than
113/// carrying its own (see `crate::MailboxEngine::resolve_identity`). This is
114/// "a participant record gains a kind" made a type-system fact rather than
115/// a runtime tag: [`ParticipantRecord`] is the account kind, this is the
116/// session kind, and which one a lookup returns says which kind it found.
117#[derive(Clone, Debug, Eq, PartialEq)]
118pub struct SessionRecord {
119    pub account: ParticipantId,
120    pub card: SessionCard,
121    pub last_seen_unix_ms: u64,
122}
123
124/// What [`MailStore::insert_message`] did. Distinguishes a genuinely new
125/// message from a retry recognised by its idempotency key, so the engine can
126/// return the *original* [`mail4agent_api::SendResponse`] without needing a
127/// second read.
128#[derive(Clone, Debug, Eq, PartialEq)]
129pub enum InsertMessageOutcome {
130    /// The message was new and is now stored.
131    Inserted,
132    /// A prior send from the same participant already used this idempotency
133    /// key; nothing was created, and this is the id of that prior message.
134    Deduplicated { message_id: MessageId },
135}
136
137/// The mailbox's persistence boundary. One mutating method per engine-level
138/// mutation (see the module doc comment on why); reads are split finely
139/// enough that each maps onto a single indexed SQL query rather than a
140/// linear scan.
141pub trait MailStore {
142    /// Registers a new participant. The caller (the engine) has already
143    /// confirmed no participant is registered under this id.
144    fn register_participant(&mut self, id: ParticipantId, record: ParticipantRecord) -> Result<(), StoreError>;
145
146    /// Removes a participant's registration entirely, including its place
147    /// in the secret-digest index. Room memberships naming this id are left
148    /// as-is: the id can never authenticate again without a fresh
149    /// registration, so a stale membership entry is inert, not a leak.
150    fn deregister_participant(&mut self, id: &ParticipantId) -> Result<(), StoreError>;
151
152    /// Replaces a participant's stored secret digest -- the shared
153    /// mechanism behind both revoking and rotating a secret (see
154    /// `crate::MailboxEngine::rotate_participant_secret`).
155    fn set_participant_secret_digest(&mut self, id: &ParticipantId, digest: SecretDigest) -> Result<(), StoreError>;
156
157    /// Replaces a participant's registered delivery-listener URL --
158    /// `Some` to register or replace one, `None` to remove it. The caller
159    /// (the engine) has already confirmed `id` is registered and, on
160    /// `Some`, already validated `url`'s shape.
161    fn set_listener_url(&mut self, id: &ParticipantId, url: Option<String>) -> Result<(), StoreError>;
162
163    fn get_participant(&self, id: &ParticipantId) -> Result<Option<ParticipantRecord>, StoreError>;
164
165    /// Looks a participant up by the exact digest of a presented secret.
166    /// See `crate::MailboxEngine::authenticate` for why this is an index
167    /// lookup, not a scan.
168    fn find_participant_by_digest(
169        &self,
170        digest: &SecretDigest,
171    ) -> Result<Option<(ParticipantId, ParticipantRecord)>, StoreError>;
172
173    /// Creates a room with no members. The caller has already confirmed no
174    /// room is registered under this id.
175    fn create_room(&mut self, id: RoomId, created_at_unix_ms: u64) -> Result<(), StoreError>;
176
177    /// Idempotent: adding an existing member is a no-op.
178    fn add_room_member(&mut self, room: &RoomId, participant: ParticipantId) -> Result<(), StoreError>;
179
180    /// Idempotent: removing a non-member is a no-op.
181    fn remove_room_member(&mut self, room: &RoomId, participant: &ParticipantId) -> Result<(), StoreError>;
182
183    fn get_room(&self, id: &RoomId) -> Result<Option<RoomRecord>, StoreError>;
184
185    /// Stores `message` unless `idempotency` names a (sender address, key)
186    /// pair already recorded against an earlier message, in which case
187    /// nothing is created and that earlier message's id is returned. One
188    /// call, one transaction: a SQLite implementation satisfies this with
189    /// an insert under a `UNIQUE (sender, key)` constraint (or an
190    /// equivalent check-and-insert within one transaction), never a
191    /// separate read-then-write pair that could race under concurrent
192    /// callers. The sender is an [`Address`] rather than a
193    /// [`ParticipantId`] so that a retry from one specific session is
194    /// deduplicated against that session, not against every session of its
195    /// account.
196    fn insert_message(
197        &mut self,
198        message: Message,
199        idempotency: Option<(Address, String)>,
200    ) -> Result<InsertMessageOutcome, StoreError>;
201
202    fn get_message(&self, id: &MessageId) -> Result<Option<Message>, StoreError>;
203
204    /// Messages addressed to exactly `to` (a [`Address::Direct`] account
205    /// address or a [`Address::Session`] one), no older than
206    /// `since_unix_ms`. Never `Address::Room` -- room mail is
207    /// [`Self::room_messages_since`], keyed by [`RoomId`] rather than by a
208    /// full address, since it is never gated on the reader's own identity
209    /// the way this method's result is.
210    fn messages_to_since(&self, to: &Address, since_unix_ms: u64) -> Result<Vec<Message>, StoreError>;
211
212    /// Messages addressed to `room`, no older than `since_unix_ms`. Not
213    /// gated on membership -- the caller (the engine) decides who may see
214    /// the result.
215    fn room_messages_since(&self, room: &RoomId, since_unix_ms: u64) -> Result<Vec<Message>, StoreError>;
216
217    /// Every room `participant` currently belongs to. Membership stays on
218    /// the account: a session looks its account's rooms up through this
219    /// same method, it does not have a membership set of its own (see
220    /// [`SessionRecord`]).
221    fn rooms_containing(&self, participant: &ParticipantId) -> Result<Vec<RoomId>, StoreError>;
222
223    /// Records an acknowledgement, or returns the one already on file for
224    /// this `(message_id, reader)` pair unchanged. One call, one
225    /// transaction, so two concurrent acks of the same message by the same
226    /// reader cannot both "win" with different timestamps. `reader` is an
227    /// [`Address`] so a session's ack is tracked separately from its
228    /// account's and from its sibling sessions', the way Matrix scopes a
229    /// read marker to a `(user_id, device_id)` pair.
230    fn record_ack(&mut self, ack: Ack) -> Result<Ack, StoreError>;
231
232    fn get_ack(&self, message_id: &MessageId, reader: &Address) -> Result<Option<Ack>, StoreError>;
233
234    /// Every registered participant, for the mailbox's own directory
235    /// (`crate::MailboxEngine::directory`). Returns the full set, always
236    /// -- this mailbox is a small, local directory, not a paginated
237    /// social graph. **Never returns a secret digest or a permission
238    /// bit**: see [`ParticipantSummary`]'s own doc comment for why the
239    /// return type itself rules that out.
240    fn list_participants(&self) -> Result<Vec<ParticipantSummary>, StoreError>;
241
242    /// Every room the mailbox tracks, with its current membership, for
243    /// the same directory. Also the full set, always, for the same reason.
244    fn list_rooms(&self) -> Result<Vec<RoomSummary>, StoreError>;
245
246    /// Looks a session up by its id. `None` until
247    /// `crate::MailboxEngine::ensure_session` has registered it at least
248    /// once.
249    fn get_session(&self, session: &SessionId) -> Result<Option<SessionRecord>, StoreError>;
250
251    /// Inserts or wholesale-replaces the record for `session`. The engine,
252    /// not this trait, is responsible for merging a fresh reading into an
253    /// existing record before calling this -- see
254    /// `crate::MailboxEngine::ensure_session` and `::set_declared`, the
255    /// only two callers, and the only two ways a [`SessionRecord`] ever
256    /// changes.
257    fn upsert_session(&mut self, session: SessionId, record: SessionRecord) -> Result<(), StoreError>;
258
259    /// Every session currently registered under `account`, for the
260    /// mailbox's own directory (`crate::MailboxEngine::directory`), which
261    /// nests them under their account the way Matrix nests devices under a
262    /// `user_id`.
263    fn sessions_of(&self, account: &ParticipantId) -> Result<Vec<(SessionId, SessionRecord)>, StoreError>;
264}
265
266/// An in-memory [`MailStore`], used by this crate's own tests. Not meant for
267/// production use: nothing here survives a process restart, and every
268/// method always succeeds -- there is no disk, lock or connection here to
269/// fail.
270#[derive(Default)]
271pub struct InMemoryStore {
272    participants: HashMap<ParticipantId, ParticipantRecord>,
273    digest_index: HashMap<SecretDigest, ParticipantId>,
274    rooms: HashMap<RoomId, RoomRecord>,
275    messages: HashMap<MessageId, Message>,
276    idempotency: HashMap<(Address, String), MessageId>,
277    acks: HashMap<(MessageId, Address), Ack>,
278    sessions: HashMap<SessionId, SessionRecord>,
279}
280
281impl MailStore for InMemoryStore {
282    fn register_participant(&mut self, id: ParticipantId, record: ParticipantRecord) -> Result<(), StoreError> {
283        self.digest_index.insert(record.secret_digest, id.clone());
284        self.participants.insert(id, record);
285        Ok(())
286    }
287
288    fn deregister_participant(&mut self, id: &ParticipantId) -> Result<(), StoreError> {
289        if let Some(record) = self.participants.remove(id) {
290            self.digest_index.remove(&record.secret_digest);
291        }
292        Ok(())
293    }
294
295    fn set_participant_secret_digest(&mut self, id: &ParticipantId, digest: SecretDigest) -> Result<(), StoreError> {
296        if let Some(record) = self.participants.get_mut(id) {
297            self.digest_index.remove(&record.secret_digest);
298            record.secret_digest = digest;
299            self.digest_index.insert(digest, id.clone());
300        }
301        Ok(())
302    }
303
304    fn get_participant(&self, id: &ParticipantId) -> Result<Option<ParticipantRecord>, StoreError> {
305        Ok(self.participants.get(id).cloned())
306    }
307
308    fn set_listener_url(&mut self, id: &ParticipantId, url: Option<String>) -> Result<(), StoreError> {
309        if let Some(record) = self.participants.get_mut(id) {
310            record.listener_url = url;
311        }
312        Ok(())
313    }
314
315    fn find_participant_by_digest(
316        &self,
317        digest: &SecretDigest,
318    ) -> Result<Option<(ParticipantId, ParticipantRecord)>, StoreError> {
319        let Some(id) = self.digest_index.get(digest) else {
320            return Ok(None);
321        };
322        Ok(self.participants.get(id).map(|record| (id.clone(), record.clone())))
323    }
324
325    fn create_room(&mut self, id: RoomId, created_at_unix_ms: u64) -> Result<(), StoreError> {
326        self.rooms.insert(id, RoomRecord { created_at_unix_ms, members: BTreeSet::new() });
327        Ok(())
328    }
329
330    fn add_room_member(&mut self, room: &RoomId, participant: ParticipantId) -> Result<(), StoreError> {
331        if let Some(record) = self.rooms.get_mut(room) {
332            record.members.insert(participant);
333        }
334        Ok(())
335    }
336
337    fn remove_room_member(&mut self, room: &RoomId, participant: &ParticipantId) -> Result<(), StoreError> {
338        if let Some(record) = self.rooms.get_mut(room) {
339            record.members.remove(participant);
340        }
341        Ok(())
342    }
343
344    fn get_room(&self, id: &RoomId) -> Result<Option<RoomRecord>, StoreError> {
345        Ok(self.rooms.get(id).cloned())
346    }
347
348    fn insert_message(
349        &mut self,
350        message: Message,
351        idempotency: Option<(Address, String)>,
352    ) -> Result<InsertMessageOutcome, StoreError> {
353        if let Some(key) = &idempotency {
354            if let Some(existing) = self.idempotency.get(key) {
355                return Ok(InsertMessageOutcome::Deduplicated { message_id: existing.clone() });
356            }
357        }
358        let message_id = message.message_id.clone();
359        self.messages.insert(message_id.clone(), message);
360        if let Some(key) = idempotency {
361            self.idempotency.insert(key, message_id);
362        }
363        Ok(InsertMessageOutcome::Inserted)
364    }
365
366    fn get_message(&self, id: &MessageId) -> Result<Option<Message>, StoreError> {
367        Ok(self.messages.get(id).cloned())
368    }
369
370    fn messages_to_since(&self, to: &Address, since_unix_ms: u64) -> Result<Vec<Message>, StoreError> {
371        Ok(self
372            .messages
373            .values()
374            .filter(|message| message.created_at_unix_ms >= since_unix_ms && &message.to == to)
375            .cloned()
376            .collect())
377    }
378
379    fn room_messages_since(&self, room: &RoomId, since_unix_ms: u64) -> Result<Vec<Message>, StoreError> {
380        Ok(self
381            .messages
382            .values()
383            .filter(|message| {
384                message.created_at_unix_ms >= since_unix_ms
385                    && matches!(&message.to, Address::Room { room: r } if r == room)
386            })
387            .cloned()
388            .collect())
389    }
390
391    fn rooms_containing(&self, participant: &ParticipantId) -> Result<Vec<RoomId>, StoreError> {
392        Ok(self
393            .rooms
394            .iter()
395            .filter(|(_, record)| record.members.contains(participant))
396            .map(|(id, _)| id.clone())
397            .collect())
398    }
399
400    fn record_ack(&mut self, ack: Ack) -> Result<Ack, StoreError> {
401        Ok(self.acks.entry((ack.message_id.clone(), ack.reader.clone())).or_insert(ack).clone())
402    }
403
404    fn get_ack(&self, message_id: &MessageId, reader: &Address) -> Result<Option<Ack>, StoreError> {
405        Ok(self.acks.get(&(message_id.clone(), reader.clone())).cloned())
406    }
407
408    fn list_participants(&self) -> Result<Vec<ParticipantSummary>, StoreError> {
409        Ok(self
410            .participants
411            .iter()
412            .map(|(id, record)| ParticipantSummary { id: id.clone(), label: record.label.clone() })
413            .collect())
414    }
415
416    fn list_rooms(&self) -> Result<Vec<RoomSummary>, StoreError> {
417        Ok(self
418            .rooms
419            .iter()
420            .map(|(id, record)| RoomSummary { id: id.clone(), members: record.members.clone() })
421            .collect())
422    }
423
424    fn get_session(&self, session: &SessionId) -> Result<Option<SessionRecord>, StoreError> {
425        Ok(self.sessions.get(session).cloned())
426    }
427
428    fn upsert_session(&mut self, session: SessionId, record: SessionRecord) -> Result<(), StoreError> {
429        self.sessions.insert(session, record);
430        Ok(())
431    }
432
433    fn sessions_of(&self, account: &ParticipantId) -> Result<Vec<(SessionId, SessionRecord)>, StoreError> {
434        Ok(self
435            .sessions
436            .iter()
437            .filter(|(_, record)| &record.account == account)
438            .map(|(id, record)| (id.clone(), record.clone()))
439            .collect())
440    }
441}