1use 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
16pub struct SqliteMailStore {
21 db: Db,
22}
23
24impl SqliteMailStore {
25 pub fn new(db: Db) -> Self {
31 Self { db }
32 }
33
34 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 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 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 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 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 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
643enum RawInsertOutcome {
651 Inserted,
652 Deduplicated(String),
653}
654
655struct 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#[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#[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
711fn 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
743fn 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
758fn 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
800fn 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
816fn 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
838fn address_text(address: &Address) -> String {
852 address.to_string()
853}
854
855#[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}