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}