1use std::collections::HashMap;
9use std::path::Path;
10use std::sync::Mutex;
11
12use chacha20poly1305::aead::{Aead, KeyInit};
13use chacha20poly1305::{XChaCha20Poly1305, XNonce};
14use rusqlite::{params, Connection, OptionalExtension, Row};
15use rusqlite_migration::{Migrations, M};
16
17use crate::control::ControlOp;
18use crate::error::{Error, Result};
19use crate::faults::Faults;
20use crate::invite::Invite;
21use crate::keys::{Kind, PublicKey};
22use crate::limits::Limits;
23use crate::message::{Content, DataClass, Message, MessageType, GENESIS_PREV};
24use crate::room::{Member, Role, Room};
25
26const MIGRATIONS: &[M<'static>] = &[
27 M::up(
28 r#"
29 CREATE TABLE rooms (
30 id TEXT PRIMARY KEY,
31 name TEXT NOT NULL UNIQUE,
32 about TEXT NOT NULL DEFAULT '',
33 owner TEXT NOT NULL,
34 created TEXT NOT NULL,
35 retention_days INTEGER,
36 class TEXT NOT NULL DEFAULT 'internal',
37 paused INTEGER NOT NULL DEFAULT 0,
38 hold INTEGER NOT NULL DEFAULT 0,
39 home_node TEXT NOT NULL,
40 home_hints TEXT NOT NULL DEFAULT 'null'
41 );
42 CREATE TABLE members (
43 room_id TEXT NOT NULL REFERENCES rooms(id) ON DELETE CASCADE,
44 name TEXT NOT NULL,
45 key TEXT NOT NULL,
46 kind TEXT NOT NULL,
47 role TEXT NOT NULL,
48 node TEXT,
49 granted_by TEXT NOT NULL,
50 expires_at TEXT,
51 joined_at TEXT NOT NULL,
52 last_seen TEXT,
53 muted INTEGER NOT NULL DEFAULT 0,
54 revoked INTEGER NOT NULL DEFAULT 0,
55 profile TEXT NOT NULL DEFAULT 'null',
56 identity TEXT,
57 PRIMARY KEY (room_id, name)
58 );
59 CREATE UNIQUE INDEX members_key ON members(room_id, key);
60 CREATE TABLE messages (
61 room_id TEXT NOT NULL REFERENCES rooms(id) ON DELETE CASCADE,
62 seq INTEGER NOT NULL,
63 id TEXT NOT NULL UNIQUE,
64 from_name TEXT NOT NULL,
65 type TEXT NOT NULL,
66 ts TEXT NOT NULL,
67 prev TEXT NOT NULL,
68 chain_hash TEXT NOT NULL,
69 envelope TEXT NOT NULL,
70 PRIMARY KEY (room_id, seq)
71 );
72 CREATE INDEX messages_room_from_ts ON messages(room_id, from_name, ts);
73 CREATE INDEX messages_room_ts ON messages(room_id, ts);
74 CREATE TABLE contents (
75 msg_id TEXT PRIMARY KEY REFERENCES messages(id) ON DELETE CASCADE,
76 body TEXT NOT NULL,
77 deleted INTEGER NOT NULL DEFAULT 0
78 );
79 CREATE TABLE bookmarks (
80 room_id TEXT NOT NULL,
81 reader TEXT NOT NULL,
82 seq INTEGER NOT NULL DEFAULT 0,
83 PRIMARY KEY (room_id, reader)
84 );
85 CREATE TABLE outbox (
86 msg_id TEXT PRIMARY KEY,
87 room_id TEXT NOT NULL,
88 message TEXT NOT NULL,
89 created TEXT NOT NULL,
90 attempts INTEGER NOT NULL DEFAULT 0,
91 last_error TEXT
92 );
93 CREATE TABLE invites (
94 nonce TEXT PRIMARY KEY,
95 room_id TEXT NOT NULL,
96 name TEXT NOT NULL,
97 created TEXT NOT NULL,
98 expires TEXT NOT NULL,
99 used_at TEXT,
100 used_by TEXT
101 );
102 "#,
103 ),
104 M::up(
105 r#"
106 ALTER TABLE rooms ADD COLUMN closed INTEGER NOT NULL DEFAULT 0;
107 ALTER TABLE messages ADD COLUMN reply_to TEXT;
108 UPDATE messages SET reply_to = json_extract(envelope, '$.reply_to');
109 CREATE INDEX messages_reply_to ON messages(room_id, reply_to);
110 CREATE TABLE approvals_used (
111 msg_id TEXT PRIMARY KEY,
112 action_hash TEXT NOT NULL,
113 used_at TEXT NOT NULL
114 );
115 "#,
116 ),
117 M::up(
120 r#"
121 ALTER TABLE messages ADD COLUMN received TEXT;
122 UPDATE messages SET received = ts;
123 CREATE INDEX messages_room_from_received ON messages(room_id, from_name, received);
124 CREATE INDEX messages_room_received ON messages(room_id, received);
125 "#,
126 ),
127 M::up(
131 r#"
132 ALTER TABLE outbox ADD COLUMN sender TEXT NOT NULL DEFAULT '';
133 ALTER TABLE outbox ADD COLUMN state TEXT NOT NULL DEFAULT 'pending';
134 ALTER TABLE outbox ADD COLUMN retry_at TEXT;
135 ALTER TABLE outbox ADD COLUMN reason_class TEXT;
136 ALTER TABLE outbox ADD COLUMN reason_code INTEGER;
137 ALTER TABLE outbox ADD COLUMN unknown_code INTEGER;
138 ALTER TABLE outbox ADD COLUMN unknown_since TEXT;
139 ALTER TABLE outbox ADD COLUMN unknown_count INTEGER NOT NULL DEFAULT 0;
140 ALTER TABLE outbox ADD COLUMN updated TEXT;
141 ALTER TABLE outbox ADD COLUMN sig TEXT;
142 UPDATE outbox SET
143 sender = CASE WHEN json_valid(message)
144 THEN COALESCE(json_extract(message, '$.from'), '') ELSE '' END,
145 sig = CASE WHEN json_valid(message) THEN json_extract(message, '$.sig') END;
146 CREATE INDEX outbox_room_state ON outbox(room_id, state, sender, created, msg_id);
147 CREATE TABLE meta (
148 key TEXT PRIMARY KEY,
149 value TEXT NOT NULL
150 );
151 "#,
152 ),
153 M::up(
157 r#"
158 CREATE TABLE deliveries (
159 room_id TEXT NOT NULL,
160 reader TEXT NOT NULL,
161 seq INTEGER NOT NULL,
162 state TEXT NOT NULL,
163 token TEXT UNIQUE,
164 lease_until TEXT,
165 retry_at TEXT,
166 attempt INTEGER NOT NULL DEFAULT 0,
167 updated TEXT NOT NULL,
168 PRIMARY KEY (room_id, reader, seq)
169 );
170 "#,
171 ),
172 M::up(
176 r#"
177 CREATE TABLE spends (
178 approve_id TEXT PRIMARY KEY,
179 room_id TEXT NOT NULL,
180 action_hash TEXT NOT NULL,
181 op_id TEXT NOT NULL,
182 spender TEXT NOT NULL,
183 node TEXT NOT NULL,
184 at TEXT NOT NULL,
185 audit_seq INTEGER NOT NULL
186 );
187 CREATE TABLE pending_spends (
188 room_id TEXT NOT NULL,
189 action_hash TEXT NOT NULL,
190 spender TEXT NOT NULL,
191 op_id TEXT NOT NULL,
192 created TEXT NOT NULL,
193 PRIMARY KEY (room_id, action_hash, spender)
194 );
195 "#,
196 ),
197 M::up(
202 r#"
203 CREATE TABLE file_refs (
204 msg_id TEXT NOT NULL REFERENCES messages(id) ON DELETE CASCADE,
205 room_id TEXT NOT NULL,
206 hash TEXT NOT NULL,
207 name TEXT NOT NULL,
208 size INTEGER NOT NULL,
209 type TEXT NOT NULL DEFAULT '',
210 from_name TEXT NOT NULL,
211 to_name TEXT,
212 PRIMARY KEY (msg_id, hash)
213 );
214 CREATE INDEX file_refs_hash ON file_refs(room_id, hash);
215 CREATE TABLE files (
216 room_id TEXT NOT NULL,
217 hash TEXT NOT NULL,
218 size INTEGER NOT NULL,
219 kind TEXT NOT NULL,
220 uploader TEXT NOT NULL,
221 created TEXT NOT NULL,
222 PRIMARY KEY (room_id, hash)
223 );
224 CREATE INDEX files_uploader ON files(room_id, uploader, created);
225 CREATE TABLE file_fetches (
226 room_id TEXT NOT NULL,
227 hash TEXT NOT NULL,
228 member TEXT NOT NULL,
229 at TEXT NOT NULL,
230 PRIMARY KEY (room_id, hash, member)
231 );
232 "#,
233 ),
234];
235
236const SEAL_BATCH: usize = 100;
238
239pub struct Store {
241 conn: Mutex<Connection>,
242 cipher: Option<XChaCha20Poly1305>,
246 pub faults: Faults,
248}
249
250const ENC_PREFIX: &str = "enc1:";
251
252impl std::fmt::Debug for Store {
253 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
254 f.write_str("Store")
255 }
256}
257
258#[derive(Debug, Clone)]
260pub struct LocalMember {
261 pub member: Member,
262 pub identity: String,
264}
265
266#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
268#[serde(rename_all = "lowercase")]
269pub enum OutboxState {
270 Pending,
272 Waiting,
274 Failed,
276 Quarantined,
278 Dropped,
280}
281
282impl OutboxState {
283 pub fn as_str(&self) -> &'static str {
284 match self {
285 OutboxState::Pending => "pending",
286 OutboxState::Waiting => "waiting",
287 OutboxState::Failed => "failed",
288 OutboxState::Quarantined => "quarantined",
289 OutboxState::Dropped => "dropped",
290 }
291 }
292}
293
294impl std::str::FromStr for OutboxState {
295 type Err = Error;
296 fn from_str(s: &str) -> Result<Self> {
297 Ok(match s {
298 "pending" => OutboxState::Pending,
299 "waiting" => OutboxState::Waiting,
300 "failed" => OutboxState::Failed,
301 "quarantined" => OutboxState::Quarantined,
302 "dropped" => OutboxState::Dropped,
303 other => return Err(Error::Invalid(format!("unknown outbox state {other}"))),
304 })
305 }
306}
307
308#[derive(Debug, Clone)]
310pub struct OutboxEntry {
311 pub msg_id: String,
312 pub room_id: String,
313 pub sender: String,
315 pub state: OutboxState,
316 pub created: String,
317 pub attempts: u32,
318 pub retry_at: Option<String>,
319 pub reason: Option<String>,
321 pub reason_class: Option<String>,
323 pub reason_code: Option<i32>,
324 pub updated: Option<String>,
325 pub message: Option<Message>,
328 pub unreadable: bool,
330}
331
332#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
333pub struct OutboxCounts {
334 pub pending: u64,
335 pub waiting: u64,
336 pub failed: u64,
337 pub quarantined: u64,
338 pub dropped: u64,
339}
340
341impl OutboxCounts {
342 pub fn queued(&self) -> u64 {
344 self.pending + self.waiting
345 }
346}
347
348#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
350#[serde(rename_all = "lowercase")]
351pub enum DeliveryState {
352 Leased,
355 Acked,
358 Delayed,
360 Quarantined,
363 Replay,
365}
366
367impl DeliveryState {
368 pub fn as_str(&self) -> &'static str {
369 match self {
370 DeliveryState::Leased => "leased",
371 DeliveryState::Acked => "acked",
372 DeliveryState::Delayed => "delayed",
373 DeliveryState::Quarantined => "quarantined",
374 DeliveryState::Replay => "replay",
375 }
376 }
377
378 fn parse(s: &str) -> DeliveryState {
379 match s {
380 "acked" => DeliveryState::Acked,
381 "delayed" => DeliveryState::Delayed,
382 "quarantined" => DeliveryState::Quarantined,
383 "replay" => DeliveryState::Replay,
384 _ => DeliveryState::Leased,
385 }
386 }
387}
388
389#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
392pub struct WakePending {
393 pub count: u64,
394 pub newest_seq: u64,
395 pub newest_id: String,
396 pub looped: u64,
398}
399
400#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
402pub struct Delivery {
403 pub room_id: String,
404 pub reader: String,
405 pub seq: u64,
406 pub state: DeliveryState,
407 #[serde(default, skip_serializing_if = "Option::is_none")]
408 pub token: Option<String>,
409 #[serde(default, skip_serializing_if = "Option::is_none")]
410 pub lease_until: Option<String>,
411 #[serde(default, skip_serializing_if = "Option::is_none")]
412 pub retry_at: Option<String>,
413 pub attempt: u32,
415 pub updated: String,
416}
417
418#[derive(Debug, Clone)]
420pub enum Settle {
421 Ack,
423 Renew { lease_until: String },
425 Nack { retry_at: String },
427}
428
429#[derive(Debug, Clone)]
432pub struct SpendRequest {
433 pub room_id: String,
434 pub approve_id: String,
435 pub action_hash: String,
436 pub op_id: String,
438 pub spender: String,
440 pub node: String,
442 pub now: String,
444}
445
446#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
448pub struct SpendRecord {
449 pub approve_id: String,
450 pub room_id: String,
451 pub action_hash: String,
452 pub op_id: String,
453 pub spender: String,
454 pub node: String,
455 pub at: String,
456 pub audit_seq: u64,
458}
459
460#[derive(Debug, Clone, PartialEq, Eq)]
462pub enum SpendOutcome {
463 Spent(SpendRecord),
465 AlreadyYours(SpendRecord),
467 SpentByAnother,
469}
470
471type Candidate = (
474 i64,
475 String,
476 String,
477 Option<String>,
478 Option<String>,
479 Option<String>,
480);
481
482fn wanted(from: &str, kind: &str, reader: &str) -> bool {
485 from != reader && kind != "system" && kind != "control"
486}
487
488fn seconds_between(a: &str, b: &str) -> i64 {
490 match (
491 chrono::DateTime::parse_from_rfc3339(a),
492 chrono::DateTime::parse_from_rfc3339(b),
493 ) {
494 (Ok(a), Ok(b)) => (b - a).num_seconds(),
495 _ => 0,
496 }
497}
498
499impl Store {
500 pub fn open(path: &Path) -> Result<Self> {
503 Self::open_with_key(path, None)
504 }
505
506 pub fn open_with_key(path: &Path, key: Option<[u8; 32]>) -> Result<Self> {
513 let store = Self::open_unsealed(path, key)?;
514 store.seal_legacy_outbox(SEAL_BATCH)?;
515 Ok(store)
516 }
517
518 #[doc(hidden)]
521 pub fn open_unsealed(path: &Path, key: Option<[u8; 32]>) -> Result<Self> {
522 if let Some(parent) = path.parent() {
523 std::fs::create_dir_all(parent)?;
524 }
525 let conn = Connection::open(path)?;
526 Self::init(conn, key)
527 }
528
529 pub fn open_memory() -> Result<Self> {
531 Self::init(Connection::open_in_memory()?, None)
532 }
533
534 pub fn open_memory_with_key(key: [u8; 32]) -> Result<Self> {
536 Self::init(Connection::open_in_memory()?, Some(key))
537 }
538
539 fn init(mut conn: Connection, key: Option<[u8; 32]>) -> Result<Self> {
540 conn.pragma_update(None, "journal_mode", "WAL")?;
541 conn.pragma_update(None, "synchronous", "FULL")?;
542 conn.pragma_update(None, "foreign_keys", "ON")?;
543 conn.pragma_update(None, "secure_delete", "ON")?;
547 Migrations::from_slice(MIGRATIONS).to_latest(&mut conn)?;
548 Ok(Store {
549 conn: Mutex::new(conn),
550 cipher: key.map(|k| XChaCha20Poly1305::new((&k).into())),
551 faults: Faults::default(),
552 })
553 }
554
555 #[doc(hidden)]
558 pub fn cap_size_for_test(&self, extra: u32) -> Result<()> {
559 let conn = self.lock();
560 let pages: i64 = conn.query_row("PRAGMA page_count", [], |r| r.get(0))?;
561 conn.pragma_update(None, "max_page_count", pages + extra as i64)?;
562 Ok(())
563 }
564
565 pub fn checkpoint(&self) -> Result<()> {
569 let conn = self.lock();
570 conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |_| Ok(()))?;
571 Ok(())
572 }
573
574 pub fn seal_legacy_outbox(&self, batch: usize) -> Result<u64> {
577 if self.cipher.is_none() {
578 return Ok(0);
579 }
580 let mut sealed = 0u64;
581 loop {
582 let mut conn = self.lock();
583 let tx = conn.transaction()?;
584 let rows: Vec<(String, String)> = {
585 let mut stmt = tx.prepare(
586 "SELECT msg_id, message FROM outbox
587 WHERE message != '' AND message NOT LIKE 'enc1:%' LIMIT ?1",
588 )?;
589 let rows = stmt.query_map(params![batch as i64], |r| Ok((r.get(0)?, r.get(1)?)))?;
590 rows.collect::<rusqlite::Result<Vec<_>>>()?
591 };
592 if rows.is_empty() {
593 tx.execute(
594 "INSERT INTO meta (key, value) VALUES ('outbox_sealed', ?1)
595 ON CONFLICT(key) DO UPDATE SET value=excluded.value",
596 params![crate::message::now_ts()],
597 )?;
598 tx.commit()?;
599 return Ok(sealed);
600 }
601 for (id, plain) in rows {
602 self.faults.check("migration.row")?;
603 tx.execute(
604 "UPDATE outbox SET message=?2 WHERE msg_id=?1",
605 params![id, self.seal(&plain)?],
606 )?;
607 sealed += 1;
608 }
609 tx.commit()?;
610 }
611 }
612
613 pub fn is_encrypted(&self) -> bool {
615 self.cipher.is_some()
616 }
617
618 fn seal(&self, body: &str) -> Result<String> {
619 match &self.cipher {
620 None => Ok(body.to_string()),
621 Some(c) => {
622 let nonce_bytes: [u8; 24] = rand::random();
623 let nonce = XNonce::from(nonce_bytes);
624 let ct = c
625 .encrypt(&nonce, body.as_bytes())
626 .map_err(|_| Error::Other("encrypt failed".into()))?;
627 let mut out = nonce_bytes.to_vec();
628 out.extend_from_slice(&ct);
629 Ok(format!(
630 "{ENC_PREFIX}{}",
631 data_encoding::BASE64.encode(&out)
632 ))
633 }
634 }
635 }
636
637 fn unseal(&self, stored: &str) -> Option<String> {
638 let Some(b64) = stored.strip_prefix(ENC_PREFIX) else {
639 return Some(stored.to_string());
640 };
641 let c = self.cipher.as_ref()?;
642 let raw = data_encoding::BASE64.decode(b64.as_bytes()).ok()?;
643 if raw.len() < 24 {
644 return None;
645 }
646 let (n, ct) = raw.split_at(24);
647 let nonce = XNonce::from(<[u8; 24]>::try_from(n).ok()?);
648 let pt = c.decrypt(&nonce, ct).ok()?;
649 String::from_utf8(pt).ok()
650 }
651
652 fn lock(&self) -> std::sync::MutexGuard<'_, Connection> {
653 self.conn.lock().unwrap_or_else(|p| p.into_inner())
654 }
655
656 pub fn create_room(&self, room: &Room) -> Result<()> {
659 let conn = self.lock();
660 conn.execute(
661 "INSERT INTO rooms (id, name, about, owner, created, retention_days, class, paused, hold, home_node, home_hints)
662 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
663 params![
664 room.id,
665 room.name,
666 room.about,
667 room.owner.to_string(),
668 room.created,
669 room.retention_days,
670 room.class.as_str(),
671 room.paused as i32,
672 room.hold as i32,
673 room.home_node,
674 serde_json::to_string(&room.home_hints)?,
675 ],
676 )
677 .map_err(|e| match e {
678 rusqlite::Error::SqliteFailure(f, _) if f.code == rusqlite::ErrorCode::ConstraintViolation => {
679 Error::NameTaken(format!("a room named {} already exists here", room.name))
680 }
681 other => Error::Db(other),
682 })?;
683 Ok(())
684 }
685
686 pub fn update_room(&self, room: &Room) -> Result<()> {
687 let conn = self.lock();
688 conn.execute(
689 "UPDATE rooms SET about=?2, retention_days=?3, class=?4, paused=?5, hold=?6, home_node=?7, home_hints=?8, closed=?9 WHERE id=?1",
690 params![
691 room.id,
692 room.about,
693 room.retention_days,
694 room.class.as_str(),
695 room.paused as i32,
696 room.hold as i32,
697 room.home_node,
698 serde_json::to_string(&room.home_hints)?,
699 room.closed as i32,
700 ],
701 )?;
702 Ok(())
703 }
704
705 fn row_to_room(row: &Row<'_>) -> rusqlite::Result<Room> {
706 let owner: String = row.get("owner")?;
707 let class: String = row.get("class")?;
708 let hints: String = row.get("home_hints")?;
709 Ok(Room {
710 id: row.get("id")?,
711 name: row.get("name")?,
712 about: row.get("about")?,
713 owner: owner.parse().map_err(|_| rusqlite::Error::InvalidQuery)?,
714 created: row.get("created")?,
715 retention_days: row.get("retention_days")?,
716 class: class.parse().unwrap_or(DataClass::Internal),
717 paused: row.get::<_, i32>("paused")? != 0,
718 hold: row.get::<_, i32>("hold")? != 0,
719 closed: row.get::<_, i32>("closed").unwrap_or(0) != 0,
720 home_node: row.get("home_node")?,
721 home_hints: serde_json::from_str(&hints).unwrap_or(serde_json::Value::Null),
722 })
723 }
724
725 pub fn room_by_name(&self, name: &str) -> Result<Option<Room>> {
726 let conn = self.lock();
727 Ok(conn
728 .query_row(
729 "SELECT * FROM rooms WHERE name=?1",
730 params![name],
731 Self::row_to_room,
732 )
733 .optional()?)
734 }
735
736 pub fn room_by_id(&self, id: &str) -> Result<Option<Room>> {
737 let conn = self.lock();
738 Ok(conn
739 .query_row(
740 "SELECT * FROM rooms WHERE id=?1",
741 params![id],
742 Self::row_to_room,
743 )
744 .optional()?)
745 }
746
747 pub fn room(&self, name_or_id: &str) -> Result<Room> {
749 if let Some(r) = self.room_by_name(name_or_id)? {
750 return Ok(r);
751 }
752 if let Some(r) = self.room_by_id(name_or_id)? {
753 return Ok(r);
754 }
755 Err(Error::NotInRoom(name_or_id.to_string()))
756 }
757
758 pub fn list_rooms(&self) -> Result<Vec<Room>> {
759 let conn = self.lock();
760 let mut stmt = conn.prepare("SELECT * FROM rooms ORDER BY created")?;
761 let rows = stmt.query_map([], Self::row_to_room)?;
762 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
763 }
764
765 pub fn upsert_member(&self, m: &Member, identity: Option<&str>) -> Result<()> {
770 let conn = self.lock();
771 conn.execute(
772 "INSERT INTO members (room_id, name, key, kind, role, node, granted_by, expires_at, joined_at, last_seen, muted, revoked, profile, identity)
773 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
774 ON CONFLICT(room_id, name) DO UPDATE SET
775 key=excluded.key, kind=excluded.kind, role=excluded.role, node=excluded.node,
776 granted_by=excluded.granted_by, expires_at=excluded.expires_at,
777 last_seen=COALESCE(excluded.last_seen, members.last_seen),
778 muted=excluded.muted, revoked=excluded.revoked, profile=excluded.profile,
779 identity=COALESCE(excluded.identity, members.identity)",
780 params![
781 m.room_id,
782 m.name,
783 m.key.to_string(),
784 m.kind.to_string(),
785 m.role.as_str(),
786 m.node,
787 m.granted_by.to_string(),
788 m.expires_at,
789 m.joined_at,
790 m.last_seen,
791 m.muted as i32,
792 m.revoked as i32,
793 serde_json::to_string(&m.profile)?,
794 identity,
795 ],
796 )?;
797 Ok(())
798 }
799
800 fn row_to_member(row: &Row<'_>) -> rusqlite::Result<Member> {
801 let key: String = row.get("key")?;
802 let kind: String = row.get("kind")?;
803 let role: String = row.get("role")?;
804 let granted_by: String = row.get("granted_by")?;
805 let profile: String = row.get("profile")?;
806 Ok(Member {
807 room_id: row.get("room_id")?,
808 name: row.get("name")?,
809 key: key.parse().map_err(|_| rusqlite::Error::InvalidQuery)?,
810 kind: kind.parse().unwrap_or(Kind::Agent),
811 role: role.parse().unwrap_or(Role::Observer),
812 node: row.get("node")?,
813 granted_by: granted_by
814 .parse()
815 .map_err(|_| rusqlite::Error::InvalidQuery)?,
816 expires_at: row.get("expires_at")?,
817 joined_at: row.get("joined_at")?,
818 last_seen: row.get("last_seen")?,
819 muted: row.get::<_, i32>("muted")? != 0,
820 revoked: row.get::<_, i32>("revoked")? != 0,
821 profile: serde_json::from_str(&profile).unwrap_or(serde_json::Value::Null),
822 })
823 }
824
825 pub fn member_by_name(&self, room_id: &str, name: &str) -> Result<Option<Member>> {
826 let conn = self.lock();
827 Ok(conn
828 .query_row(
829 "SELECT * FROM members WHERE room_id=?1 AND name=?2",
830 params![room_id, name],
831 Self::row_to_member,
832 )
833 .optional()?)
834 }
835
836 pub fn member_by_key(&self, room_id: &str, key: &PublicKey) -> Result<Option<Member>> {
837 let conn = self.lock();
838 Ok(conn
839 .query_row(
840 "SELECT * FROM members WHERE room_id=?1 AND key=?2",
841 params![room_id, key.to_string()],
842 Self::row_to_member,
843 )
844 .optional()?)
845 }
846
847 pub fn members(&self, room_id: &str) -> Result<Vec<Member>> {
848 let conn = self.lock();
849 let mut stmt =
850 conn.prepare("SELECT * FROM members WHERE room_id=?1 ORDER BY joined_at, name")?;
851 let rows = stmt.query_map(params![room_id], Self::row_to_member)?;
852 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
853 }
854
855 pub fn local_members(&self, room_id: &str) -> Result<Vec<LocalMember>> {
857 let conn = self.lock();
858 let mut stmt = conn.prepare(
859 "SELECT * FROM members WHERE room_id=?1 AND identity IS NOT NULL ORDER BY joined_at",
860 )?;
861 let rows = stmt.query_map(params![room_id], |row| {
862 Ok(LocalMember {
863 member: Self::row_to_member(row)?,
864 identity: row.get("identity")?,
865 })
866 })?;
867 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
868 }
869
870 pub fn local_member(&self, room_id: &str, identity: &str) -> Result<Option<LocalMember>> {
872 let conn = self.lock();
873 Ok(conn
874 .query_row(
875 "SELECT * FROM members WHERE room_id=?1 AND identity=?2",
876 params![room_id, identity],
877 |row| {
878 Ok(LocalMember {
879 member: Self::row_to_member(row)?,
880 identity: row.get("identity")?,
881 })
882 },
883 )
884 .optional()?)
885 }
886
887 pub fn set_member_role(
888 &self,
889 room_id: &str,
890 name: &str,
891 role: Role,
892 expires_at: Option<&str>,
893 ) -> Result<bool> {
894 let conn = self.lock();
895 let n = conn.execute(
896 "UPDATE members SET role=?3, expires_at=?4 WHERE room_id=?1 AND name=?2",
897 params![room_id, name, role.as_str(), expires_at],
898 )?;
899 Ok(n == 1)
900 }
901
902 pub fn set_member_muted(&self, room_id: &str, name: &str, muted: bool) -> Result<bool> {
903 let conn = self.lock();
904 let n = conn.execute(
905 "UPDATE members SET muted=?3 WHERE room_id=?1 AND name=?2",
906 params![room_id, name, muted as i32],
907 )?;
908 Ok(n == 1)
909 }
910
911 pub fn set_member_revoked(&self, room_id: &str, name: &str) -> Result<bool> {
912 let conn = self.lock();
913 let n = conn.execute(
914 "UPDATE members SET revoked=1, node=NULL WHERE room_id=?1 AND name=?2",
915 params![room_id, name],
916 )?;
917 Ok(n == 1)
918 }
919
920 pub fn clear_member_node(&self, room_id: &str, name: &str) -> Result<()> {
923 let conn = self.lock();
924 conn.execute(
925 "UPDATE members SET node=NULL WHERE room_id=?1 AND name=?2",
926 params![room_id, name],
927 )?;
928 Ok(())
929 }
930
931 pub fn rename_room(&self, room_id: &str, new_name: &str) -> Result<()> {
932 let conn = self.lock();
933 conn.execute(
934 "UPDATE rooms SET name=?2 WHERE id=?1",
935 params![room_id, new_name],
936 )?;
937 Ok(())
938 }
939
940 pub fn touch_member(
942 &self,
943 room_id: &str,
944 name: &str,
945 node: Option<&str>,
946 now: &str,
947 ) -> Result<()> {
948 let conn = self.lock();
949 conn.execute(
950 "UPDATE members SET last_seen=?3, node=COALESCE(?4, node) WHERE room_id=?1 AND name=?2",
951 params![room_id, name, now, node],
952 )?;
953 Ok(())
954 }
955
956 pub fn chain_head(&self, room_id: &str) -> Result<(u64, String)> {
961 let conn = self.lock();
962 Self::chain_head_in(&conn, room_id)
963 }
964
965 fn chain_head_in(conn: &Connection, room_id: &str) -> Result<(u64, String)> {
966 let head: Option<(i64, String)> = conn
967 .query_row(
968 "SELECT seq, chain_hash FROM messages WHERE room_id=?1 ORDER BY seq DESC LIMIT 1",
969 params![room_id],
970 |row| Ok((row.get(0)?, row.get(1)?)),
971 )
972 .optional()?;
973 Ok(match head {
974 Some((seq, hash)) => (seq as u64, hash),
975 None => (0, GENESIS_PREV.to_string()),
976 })
977 }
978
979 pub fn append(&self, msg: &Message) -> Result<bool> {
983 let mut conn = self.lock();
984 let tx = conn.transaction()?;
985 let written = self.append_in(&tx, msg)?;
986 tx.commit()?;
987 Ok(written)
988 }
989
990 fn append_in(&self, conn: &Connection, msg: &Message) -> Result<bool> {
991 if !msg.is_sequenced() {
992 return Err(Error::Invalid("message has no seq/prev yet".into()));
993 }
994 let exists: Option<i64> = conn
995 .query_row(
996 "SELECT seq FROM messages WHERE id=?1",
997 params![msg.id],
998 |r| r.get(0),
999 )
1000 .optional()?;
1001 if exists.is_some() {
1002 Self::unqueue_in(conn, msg)?;
1006 return Ok(false);
1007 }
1008 let (last_seq, last_hash) = Self::chain_head_in(conn, &msg.room)?;
1009 if msg.seq != last_seq + 1 || msg.prev != last_hash {
1010 return Err(Error::Invalid(format!(
1011 "chain break in room {}: got seq {} prev {}, expected seq {} prev {}",
1012 msg.room,
1013 msg.seq,
1014 &msg.prev[..msg.prev.len().min(16)],
1015 last_seq + 1,
1016 &last_hash[..16]
1017 )));
1018 }
1019 conn.execute(
1020 "INSERT INTO messages (room_id, seq, id, from_name, type, ts, prev, chain_hash, envelope, reply_to, received)
1021 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
1022 params![
1023 msg.room,
1024 msg.seq as i64,
1025 msg.id,
1026 msg.from,
1027 msg.kind.as_str(),
1028 msg.ts,
1029 msg.prev,
1030 msg.chain_hash(),
1031 serde_json::to_string(&msg.envelope())?,
1032 msg.reply_to,
1033 crate::message::now_ts(),
1034 ],
1035 )?;
1036 if msg.tombstone {
1037 conn.execute(
1038 "INSERT INTO contents (msg_id, body, deleted) VALUES (?1, '', 1)",
1039 params![msg.id],
1040 )?;
1041 } else {
1042 conn.execute(
1043 "INSERT INTO contents (msg_id, body, deleted) VALUES (?1, ?2, 0)",
1044 params![msg.id, self.seal(&serde_json::to_string(&msg.content())?)?],
1045 )?;
1046 for f in crate::files::refs(&msg.data).unwrap_or_default() {
1049 conn.execute(
1050 "INSERT OR IGNORE INTO file_refs (msg_id, room_id, hash, name, size, type, from_name, to_name)
1051 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1052 params![msg.id, msg.room, f.id, f.name, f.size as i64, f.mime, msg.from, msg.to],
1053 )?;
1054 }
1055 }
1056 Self::unqueue_in(conn, msg)?;
1059 Ok(true)
1060 }
1061
1062 fn unqueue_in(conn: &Connection, msg: &Message) -> Result<()> {
1065 conn.execute(
1066 "DELETE FROM outbox WHERE msg_id=?1 AND (sig IS NULL OR sig=?2)",
1067 params![msg.id, msg.sig],
1068 )?;
1069 Ok(())
1070 }
1071
1072 pub fn sequence_and_append(&self, msg: &mut Message) -> Result<bool> {
1079 let mut conn = self.lock();
1080 let tx = conn.transaction()?;
1081 let exists: Option<i64> = tx
1082 .query_row(
1083 "SELECT seq FROM messages WHERE id=?1",
1084 params![msg.id],
1085 |r| r.get(0),
1086 )
1087 .optional()?;
1088 if exists.is_some() {
1089 return Ok(false);
1090 }
1091 let (last_seq, last_hash) = Self::chain_head_in(&tx, &msg.room)?;
1092 msg.sequence(last_seq + 1, &last_hash);
1093 self.append_in(&tx, msg)?;
1094 tx.commit()?;
1095 Ok(true)
1096 }
1097
1098 pub fn sequence_and_append_control(&self, msg: &mut Message, op: &ControlOp) -> Result<bool> {
1102 let mut conn = self.lock();
1103 let tx = conn.transaction()?;
1104 let exists: Option<i64> = tx
1105 .query_row(
1106 "SELECT seq FROM messages WHERE id=?1",
1107 params![msg.id],
1108 |r| r.get(0),
1109 )
1110 .optional()?;
1111 if exists.is_some() {
1112 return Ok(false);
1113 }
1114 let (last_seq, last_hash) = Self::chain_head_in(&tx, &msg.room)?;
1115 msg.sequence(last_seq + 1, &last_hash);
1116 self.append_in(&tx, msg)?;
1117 Self::apply_control_in(&tx, &msg.room, op)?;
1118 tx.commit()?;
1119 Ok(true)
1120 }
1121
1122 fn apply_control_in(conn: &Connection, room_id: &str, op: &ControlOp) -> Result<()> {
1123 match op {
1124 ControlOp::Grant { name, role, until } => {
1125 conn.execute(
1126 "UPDATE members SET role=?3, expires_at=?4 WHERE room_id=?1 AND name=?2",
1127 params![room_id, name, role.as_str(), until],
1128 )?;
1129 }
1130 ControlOp::Pause | ControlOp::Resume => {
1131 conn.execute(
1132 "UPDATE rooms SET paused=?2 WHERE id=?1",
1133 params![room_id, matches!(op, ControlOp::Pause) as i32],
1134 )?;
1135 }
1136 ControlOp::Mute { name } | ControlOp::Unmute { name } => {
1137 conn.execute(
1138 "UPDATE members SET muted=?3 WHERE room_id=?1 AND name=?2",
1139 params![room_id, name, matches!(op, ControlOp::Mute { .. }) as i32],
1140 )?;
1141 }
1142 ControlOp::Revoke { name } => {
1143 conn.execute(
1144 "UPDATE members SET revoked=1, node=NULL WHERE room_id=?1 AND name=?2",
1145 params![room_id, name],
1146 )?;
1147 }
1148 ControlOp::Hold { on } => {
1149 conn.execute(
1150 "UPDATE rooms SET hold=?2 WHERE id=?1",
1151 params![room_id, *on as i32],
1152 )?;
1153 }
1154 ControlOp::Rotated { .. } => {}
1156 }
1157 Ok(())
1158 }
1159
1160 fn row_to_spend(row: &Row<'_>) -> rusqlite::Result<SpendRecord> {
1161 Ok(SpendRecord {
1162 approve_id: row.get("approve_id")?,
1163 room_id: row.get("room_id")?,
1164 action_hash: row.get("action_hash")?,
1165 op_id: row.get("op_id")?,
1166 spender: row.get("spender")?,
1167 node: row.get("node")?,
1168 at: row.get("at")?,
1169 audit_seq: row.get::<_, i64>("audit_seq")? as u64,
1170 })
1171 }
1172
1173 pub fn spend_approve(&self, req: &SpendRequest, audit: &mut Message) -> Result<SpendOutcome> {
1191 let mut conn = self.lock();
1192 let tx = conn.transaction()?;
1193 let prior = tx
1194 .query_row(
1195 "SELECT * FROM spends WHERE approve_id=?1",
1196 params![req.approve_id],
1197 Self::row_to_spend,
1198 )
1199 .optional()?;
1200 if let Some(rec) = prior {
1201 if rec.op_id == req.op_id && rec.spender == req.spender && rec.room_id == req.room_id {
1202 return Ok(SpendOutcome::AlreadyYours(rec));
1203 }
1204 return Ok(SpendOutcome::SpentByAnother);
1205 }
1206 let legacy: Option<String> = tx
1208 .query_row(
1209 "SELECT msg_id FROM approvals_used WHERE msg_id=?1",
1210 params![req.approve_id],
1211 |r| r.get(0),
1212 )
1213 .optional()?;
1214 if legacy.is_some() {
1215 return Ok(SpendOutcome::SpentByAnother);
1216 }
1217 let room = tx
1218 .query_row(
1219 "SELECT * FROM rooms WHERE id=?1",
1220 params![req.room_id],
1221 Self::row_to_room,
1222 )
1223 .optional()?
1224 .ok_or_else(|| Error::NotInRoom(req.room_id.clone()))?;
1225 if room.closed {
1226 return Err(Error::Denied(format!("room {} was rotated", room.name)));
1227 }
1228 if room.paused {
1229 return Err(Error::RoomPaused(room.name));
1230 }
1231 let member = |name: &str| -> Result<Option<Member>> {
1232 Ok(tx
1233 .query_row(
1234 "SELECT * FROM members WHERE room_id=?1 AND name=?2",
1235 params![req.room_id, name],
1236 Self::row_to_member,
1237 )
1238 .optional()?)
1239 };
1240 let spender = member(&req.spender)?
1241 .ok_or_else(|| Error::Denied(format!("{} is not a member", req.spender)))?;
1242 if spender.revoked || spender.muted || spender.is_expired(&req.now) {
1243 return Err(Error::Denied(format!(
1244 "{} may not act on approves: revoked, muted or expired",
1245 spender.name
1246 )));
1247 }
1248 if spender.role == Role::Observer {
1249 return Err(Error::Denied(format!(
1250 "{} is an observer and may not act on approves",
1251 spender.name
1252 )));
1253 }
1254 let approve = tx
1255 .query_row(
1256 "SELECT m.envelope, c.body, c.deleted FROM messages m
1257 LEFT JOIN contents c ON c.msg_id = m.id WHERE m.id=?1 AND m.room_id=?2",
1258 params![req.approve_id, req.room_id],
1259 |r| self.row_to_message(r),
1260 )
1261 .optional()?
1262 .ok_or_else(|| Error::Denied(format!("no approve {} in this room", req.approve_id)))?;
1263 let invalid = |why: &str| Err(Error::Denied(format!("approve {}: {why}", approve.id)));
1264 if approve.kind != MessageType::Approve {
1265 return invalid("is not an approve");
1266 }
1267 if approve.action_hash.as_deref() != Some(req.action_hash.as_str()) {
1268 return invalid("is for a different action");
1269 }
1270 match &approve.expires {
1271 Some(exp) if exp.as_str() > req.now.as_str() => {}
1272 _ => return invalid("has expired"),
1273 }
1274 if approve.ts.as_str() > req.now.as_str() {
1275 return invalid("is dated in the future");
1276 }
1277 let approver = member(&approve.from)?
1278 .ok_or_else(|| Error::Denied(format!("approver {} is not a member", approve.from)))?;
1279 if approver.kind != Kind::Human {
1280 return invalid("was not signed by a human key");
1281 }
1282 approver
1283 .check_may_send(MessageType::Approve, approver.key == room.owner, &req.now)
1284 .map_err(|e| Error::Denied(format!("approve {}: {e}", approve.id)))?;
1285 approve.verify(&approver.key)?;
1286
1287 let (last_seq, last_hash) = Self::chain_head_in(&tx, &req.room_id)?;
1288 audit.sequence(last_seq + 1, &last_hash);
1289 self.append_in(&tx, audit)?;
1290 tx.execute(
1291 "INSERT INTO spends (approve_id, room_id, action_hash, op_id, spender, node, at, audit_seq)
1292 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1293 params![
1294 req.approve_id,
1295 req.room_id,
1296 req.action_hash,
1297 req.op_id,
1298 req.spender,
1299 req.node,
1300 req.now,
1301 audit.seq as i64
1302 ],
1303 )?;
1304 tx.execute(
1305 "INSERT OR IGNORE INTO approvals_used (msg_id, action_hash, used_at) VALUES (?1, ?2, ?3)",
1306 params![req.approve_id, req.action_hash, req.now],
1307 )?;
1308 self.faults.check("spend.commit")?;
1309 tx.commit()?;
1310 Ok(SpendOutcome::Spent(SpendRecord {
1311 approve_id: req.approve_id.clone(),
1312 room_id: req.room_id.clone(),
1313 action_hash: req.action_hash.clone(),
1314 op_id: req.op_id.clone(),
1315 spender: req.spender.clone(),
1316 node: req.node.clone(),
1317 at: req.now.clone(),
1318 audit_seq: audit.seq,
1319 }))
1320 }
1321
1322 pub fn spend_of(&self, approve_id: &str) -> Result<Option<SpendRecord>> {
1324 let conn = self.lock();
1325 Ok(conn
1326 .query_row(
1327 "SELECT * FROM spends WHERE approve_id=?1",
1328 params![approve_id],
1329 Self::row_to_spend,
1330 )
1331 .optional()?)
1332 }
1333
1334 pub fn pending_spend(
1338 &self,
1339 room_id: &str,
1340 action_hash: &str,
1341 spender: &str,
1342 fresh: &str,
1343 now: &str,
1344 ) -> Result<String> {
1345 let conn = self.lock();
1346 conn.execute(
1347 "INSERT OR IGNORE INTO pending_spends (room_id, action_hash, spender, op_id, created)
1348 VALUES (?1, ?2, ?3, ?4, ?5)",
1349 params![room_id, action_hash, spender, fresh, now],
1350 )?;
1351 Ok(conn.query_row(
1352 "SELECT op_id FROM pending_spends WHERE room_id=?1 AND action_hash=?2 AND spender=?3",
1353 params![room_id, action_hash, spender],
1354 |r| r.get(0),
1355 )?)
1356 }
1357
1358 pub fn pending_spend_done(
1360 &self,
1361 room_id: &str,
1362 action_hash: &str,
1363 spender: &str,
1364 ) -> Result<()> {
1365 let conn = self.lock();
1366 conn.execute(
1367 "DELETE FROM pending_spends WHERE room_id=?1 AND action_hash=?2 AND spender=?3",
1368 params![room_id, action_hash, spender],
1369 )?;
1370 Ok(())
1371 }
1372
1373 fn row_to_message(&self, row: &Row<'_>) -> rusqlite::Result<Message> {
1374 let envelope: String = row.get("envelope")?;
1375 let body: Option<String> = row.get("body")?;
1376 let deleted: Option<i32> = row.get("deleted")?;
1377 let mut v: serde_json::Value =
1378 serde_json::from_str(&envelope).map_err(|_| rusqlite::Error::InvalidQuery)?;
1379 let (content, tombstone) = match (body, deleted) {
1380 (Some(b), Some(0)) => match self.unseal(&b) {
1381 Some(plain) => (serde_json::from_str(&plain).unwrap_or_default(), false),
1382 None => (Content::default(), true),
1383 },
1384 _ => (Content::default(), true),
1385 };
1386 if let Some(obj) = v.as_object_mut() {
1387 if tombstone {
1388 obj.insert("tombstone".into(), serde_json::Value::Bool(true));
1389 } else {
1390 obj.remove("content_hash");
1391 }
1392 obj.insert("text".into(), serde_json::Value::String(content.text));
1393 obj.insert(
1394 "action".into(),
1395 serde_json::to_value(content.action).unwrap_or(serde_json::Value::Null),
1396 );
1397 obj.insert("data".into(), content.data);
1398 }
1399 serde_json::from_value(v).map_err(|_| rusqlite::Error::InvalidQuery)
1400 }
1401
1402 pub fn messages_after(&self, room_id: &str, after: u64, limit: u32) -> Result<Vec<Message>> {
1404 let conn = self.lock();
1405 let mut stmt = conn.prepare(
1406 "SELECT m.envelope, c.body, c.deleted FROM messages m
1407 LEFT JOIN contents c ON c.msg_id = m.id
1408 WHERE m.room_id=?1 AND m.seq>?2 ORDER BY m.seq LIMIT ?3",
1409 )?;
1410 let rows = stmt.query_map(params![room_id, after as i64, limit], |r| {
1411 self.row_to_message(r)
1412 })?;
1413 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1414 }
1415
1416 pub fn messages_from_ts(&self, room_id: &str, since: &str, limit: u32) -> Result<Vec<Message>> {
1418 let conn = self.lock();
1419 let mut stmt = conn.prepare(
1420 "SELECT m.envelope, c.body, c.deleted FROM messages m
1421 LEFT JOIN contents c ON c.msg_id = m.id
1422 WHERE m.room_id=?1 AND m.ts>=?2 ORDER BY m.seq LIMIT ?3",
1423 )?;
1424 let rows = stmt.query_map(params![room_id, since, limit], |r| self.row_to_message(r))?;
1425 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1426 }
1427
1428 pub fn message_by_id(&self, id: &str) -> Result<Option<Message>> {
1429 let conn = self.lock();
1430 Ok(conn
1431 .query_row(
1432 "SELECT m.envelope, c.body, c.deleted FROM messages m
1433 LEFT JOIN contents c ON c.msg_id = m.id WHERE m.id=?1",
1434 params![id],
1435 |r| self.row_to_message(r),
1436 )
1437 .optional()?)
1438 }
1439
1440 pub fn replies_to(&self, room_id: &str, msg_id: &str) -> Result<Vec<Message>> {
1442 let conn = self.lock();
1443 let mut stmt = conn.prepare(
1444 "SELECT m.envelope, c.body, c.deleted FROM messages m
1445 LEFT JOIN contents c ON c.msg_id = m.id
1446 WHERE m.room_id=?1 AND m.reply_to=?2 ORDER BY m.seq",
1447 )?;
1448 let rows = stmt.query_map(params![room_id, msg_id], |r| self.row_to_message(r))?;
1449 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1450 }
1451
1452 pub fn messages_with_trace(&self, room_id: &str, trace: &str) -> Result<Vec<Message>> {
1454 let conn = self.lock();
1455 let mut stmt = conn.prepare(
1456 "SELECT m.envelope, c.body, c.deleted FROM messages m
1457 LEFT JOIN contents c ON c.msg_id = m.id
1458 WHERE m.room_id=?1 AND json_extract(m.envelope, '$.trace')=?2
1459 ORDER BY m.seq LIMIT 5000",
1460 )?;
1461 let rows = stmt.query_map(params![room_id, trace], |r| self.row_to_message(r))?;
1462 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1463 }
1464
1465 pub fn claim_holder(&self, room_id: &str, task_id: &str) -> Result<Option<String>> {
1468 let mut holder: Option<String> = None;
1469 for m in self.replies_to(room_id, task_id)? {
1470 match m.kind {
1471 MessageType::Claim if holder.is_none() => holder = Some(m.from.clone()),
1472 MessageType::Release if holder.as_deref() == Some(m.from.as_str()) => holder = None,
1473 _ => {}
1474 }
1475 }
1476 Ok(holder)
1477 }
1478
1479 pub fn approvals_for(&self, room_id: &str, action_hash: &str) -> Result<Vec<Message>> {
1482 self.approvals_for_action(room_id, action_hash, false)
1483 }
1484
1485 pub fn approvals_for_action(
1486 &self,
1487 room_id: &str,
1488 action_hash: &str,
1489 with_spent: bool,
1490 ) -> Result<Vec<Message>> {
1491 let conn = self.lock();
1492 let mut stmt = conn.prepare(
1493 "SELECT m.envelope, c.body, c.deleted,
1494 m.id IN (SELECT msg_id FROM approvals_used) AS spent
1495 FROM messages m
1496 LEFT JOIN contents c ON c.msg_id = m.id
1497 WHERE m.room_id=?1 AND m.type='approve'
1498 AND json_extract(m.envelope, '$.action_hash')=?2
1499 AND (?3 OR m.id NOT IN (SELECT msg_id FROM approvals_used))
1500 ORDER BY spent, m.seq",
1501 )?;
1502 let rows = stmt.query_map(params![room_id, action_hash, with_spent], |r| {
1503 self.row_to_message(r)
1504 })?;
1505 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1506 }
1507
1508 pub fn approval_use(&self, msg_id: &str, action_hash: &str, now: &str) -> Result<bool> {
1510 let conn = self.lock();
1511 let n = conn.execute(
1512 "INSERT OR IGNORE INTO approvals_used (msg_id, action_hash, used_at) VALUES (?1, ?2, ?3)",
1513 params![msg_id, action_hash, now],
1514 )?;
1515 Ok(n == 1)
1516 }
1517
1518 pub fn tombstone_before(&self, room_id: &str, cutoff: &str) -> Result<u64> {
1522 let conn = self.lock();
1523 let n = conn.execute(
1524 "UPDATE contents SET body='', deleted=1
1525 WHERE deleted=0 AND msg_id IN (SELECT id FROM messages WHERE room_id=?1 AND ts<?2)",
1526 params![room_id, cutoff],
1527 )?;
1528 conn.execute(
1529 "DELETE FROM file_refs WHERE msg_id IN (SELECT id FROM messages WHERE room_id=?1 AND ts<?2)",
1530 params![room_id, cutoff],
1531 )?;
1532 Ok(n as u64)
1533 }
1534
1535 pub fn tombstone_message(&self, msg_id: &str) -> Result<bool> {
1537 let conn = self.lock();
1538 let n = conn.execute(
1539 "UPDATE contents SET body='', deleted=1 WHERE msg_id=?1 AND deleted=0",
1540 params![msg_id],
1541 )?;
1542 conn.execute("DELETE FROM file_refs WHERE msg_id=?1", params![msg_id])?;
1543 Ok(n == 1)
1544 }
1545
1546 pub fn message_count(&self, room_id: &str) -> Result<u64> {
1547 let conn = self.lock();
1548 let n: i64 = conn.query_row(
1549 "SELECT COUNT(*) FROM messages WHERE room_id=?1",
1550 params![room_id],
1551 |r| r.get(0),
1552 )?;
1553 Ok(n as u64)
1554 }
1555
1556 pub fn bookmark(&self, room_id: &str, reader: &str) -> Result<u64> {
1559 let conn = self.lock();
1560 let seq: Option<i64> = conn
1561 .query_row(
1562 "SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2",
1563 params![room_id, reader],
1564 |r| r.get(0),
1565 )
1566 .optional()?;
1567 Ok(seq.unwrap_or(0) as u64)
1568 }
1569
1570 pub fn set_bookmark(&self, room_id: &str, reader: &str, seq: u64) -> Result<()> {
1572 let conn = self.lock();
1573 conn.execute(
1574 "INSERT INTO bookmarks (room_id, reader, seq) VALUES (?1, ?2, ?3)
1575 ON CONFLICT(room_id, reader) DO UPDATE SET seq=MAX(bookmarks.seq, excluded.seq)",
1576 params![room_id, reader, seq as i64],
1577 )?;
1578 Ok(())
1579 }
1580
1581 fn row_to_delivery(row: &Row<'_>) -> rusqlite::Result<Delivery> {
1596 let state: String = row.get("state")?;
1597 Ok(Delivery {
1598 room_id: row.get("room_id")?,
1599 reader: row.get("reader")?,
1600 seq: row.get::<_, i64>("seq")? as u64,
1601 state: DeliveryState::parse(&state),
1602 token: row.get("token")?,
1603 lease_until: row.get("lease_until")?,
1604 retry_at: row.get("retry_at")?,
1605 attempt: row.get::<_, i64>("attempt")? as u32,
1606 updated: row.get("updated")?,
1607 })
1608 }
1609
1610 fn bookmark_in(conn: &Connection, room_id: &str, reader: &str) -> Result<u64> {
1611 let seq: Option<i64> = conn
1612 .query_row(
1613 "SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2",
1614 params![room_id, reader],
1615 |r| r.get(0),
1616 )
1617 .optional()?;
1618 Ok(seq.unwrap_or(0) as u64)
1619 }
1620
1621 fn advance_bookmark_in(conn: &Connection, room_id: &str, reader: &str) -> Result<u64> {
1625 let start = Self::bookmark_in(conn, room_id, reader)?;
1626 let mut bm = start;
1627 'pages: loop {
1628 let mut stmt = conn.prepare_cached(
1629 "SELECT m.seq, m.from_name, m.type, d.state FROM messages m
1630 LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
1631 WHERE m.room_id=?1 AND m.seq>?3 ORDER BY m.seq LIMIT 500",
1632 )?;
1633 let rows: Vec<(i64, String, String, Option<String>)> = stmt
1634 .query_map(params![room_id, reader, bm as i64], |r| {
1635 Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
1636 })?
1637 .collect::<rusqlite::Result<_>>()?;
1638 let n = rows.len();
1639 for (seq, from, kind, state) in rows {
1640 let settled = !wanted(&from, &kind, reader)
1641 || matches!(state.as_deref(), Some("acked") | Some("quarantined"));
1642 if !settled {
1643 break 'pages;
1644 }
1645 bm = seq as u64;
1646 }
1647 if n < 500 {
1648 break;
1649 }
1650 }
1651 if bm > start {
1652 conn.execute(
1653 "INSERT INTO bookmarks (room_id, reader, seq) VALUES (?1, ?2, ?3)
1654 ON CONFLICT(room_id, reader) DO UPDATE SET seq=MAX(bookmarks.seq, excluded.seq)",
1655 params![room_id, reader, bm as i64],
1656 )?;
1657 }
1658 Ok(bm)
1659 }
1660
1661 pub fn delivery_lease(
1664 &self,
1665 room_id: &str,
1666 reader: &str,
1667 now: &str,
1668 lease_until: &str,
1669 max_attempts: u32,
1670 ) -> Result<Option<(Message, Delivery)>> {
1671 let mut conn = self.lock();
1672 let tx = conn.transaction()?;
1673 tx.execute(
1675 "UPDATE deliveries SET state='quarantined', token=NULL, updated=?3
1676 WHERE room_id=?1 AND reader=?2 AND state='leased' AND lease_until<=?3 AND attempt>=?4",
1677 params![room_id, reader, now, max_attempts],
1678 )?;
1679 let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
1680 let mut seq: Option<i64> = tx
1683 .query_row(
1684 "SELECT seq FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq<=?3 AND (
1685 state='replay'
1686 OR (state='leased' AND lease_until<=?4)
1687 OR (state='delayed' AND retry_at<=?4))
1688 ORDER BY seq LIMIT 1",
1689 params![room_id, reader, bm as i64, now],
1690 |r| r.get(0),
1691 )
1692 .optional()?;
1693 let mut after = bm as i64;
1694 while seq.is_none() {
1695 let rows: Vec<Candidate> = {
1696 let mut stmt = tx.prepare_cached(
1697 "SELECT m.seq, m.from_name, m.type, d.state, d.lease_until, d.retry_at
1698 FROM messages m
1699 LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
1700 WHERE m.room_id=?1 AND m.seq>?3 ORDER BY m.seq LIMIT 500",
1701 )?;
1702 let rows = stmt.query_map(params![room_id, reader, after], |r| {
1703 Ok((
1704 r.get(0)?,
1705 r.get(1)?,
1706 r.get(2)?,
1707 r.get(3)?,
1708 r.get(4)?,
1709 r.get(5)?,
1710 ))
1711 })?;
1712 rows.collect::<rusqlite::Result<_>>()?
1713 };
1714 if rows.is_empty() {
1715 break;
1716 }
1717 for (s, from, kind, state, lease_until, retry_at) in &rows {
1718 after = *s;
1719 if !wanted(from, kind, reader) {
1720 continue;
1721 }
1722 let free = match state.as_deref() {
1723 None | Some("replay") => true,
1724 Some("leased") => lease_until.as_deref().is_some_and(|t| t <= now),
1725 Some("delayed") => retry_at.as_deref().is_some_and(|t| t <= now),
1726 _ => false,
1727 };
1728 if free {
1729 seq = Some(*s);
1730 break;
1731 }
1732 }
1733 }
1734 let Some(seq) = seq else {
1735 tx.commit()?;
1736 return Ok(None);
1737 };
1738 let token = format!(
1739 "d_{}",
1740 data_encoding::HEXLOWER.encode(&rand::random::<[u8; 16]>())
1741 );
1742 tx.execute(
1743 "INSERT INTO deliveries (room_id, reader, seq, state, token, lease_until, attempt, updated)
1744 VALUES (?1, ?2, ?3, 'leased', ?4, ?5, 1, ?6)
1745 ON CONFLICT(room_id, reader, seq) DO UPDATE SET
1746 state='leased', token=excluded.token, lease_until=excluded.lease_until,
1747 retry_at=NULL, attempt=deliveries.attempt+1, updated=excluded.updated",
1748 params![room_id, reader, seq, token, lease_until, now],
1749 )?;
1750 self.faults.check("delivery.lease")?;
1751 let msg = tx.query_row(
1752 "SELECT m.envelope, c.body, c.deleted FROM messages m
1753 LEFT JOIN contents c ON c.msg_id = m.id WHERE m.room_id=?1 AND m.seq=?2",
1754 params![room_id, seq],
1755 |r| self.row_to_message(r),
1756 )?;
1757 let d = tx.query_row(
1758 "SELECT * FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq=?3",
1759 params![room_id, reader, seq],
1760 Self::row_to_delivery,
1761 )?;
1762 tx.commit()?;
1763 Ok(Some((msg, d)))
1764 }
1765
1766 pub fn delivery_by_token(&self, token: &str) -> Result<Option<Delivery>> {
1768 let conn = self.lock();
1769 Ok(conn
1770 .query_row(
1771 "SELECT * FROM deliveries WHERE token=?1",
1772 params![token],
1773 Self::row_to_delivery,
1774 )
1775 .optional()?)
1776 }
1777
1778 pub fn delivery_settle(
1785 &self,
1786 token: &str,
1787 how: &Settle,
1788 now: &str,
1789 max_attempts: u32,
1790 ) -> Result<Delivery> {
1791 let mut conn = self.lock();
1792 let tx = conn.transaction()?;
1793 let d = tx
1794 .query_row(
1795 "SELECT * FROM deliveries WHERE token=?1",
1796 params![token],
1797 Self::row_to_delivery,
1798 )
1799 .optional()?
1800 .ok_or_else(|| {
1801 Error::Denied(
1802 "that delivery is no longer yours: it was settled, or its lease ran out and \
1803 it was handed out again under a new token"
1804 .into(),
1805 )
1806 })?;
1807 let key = params![d.room_id, d.reader, d.seq as i64, now];
1808 match (d.state, how) {
1809 (DeliveryState::Acked, Settle::Ack) => {}
1810 (DeliveryState::Leased, Settle::Ack) => {
1811 tx.execute(
1812 "UPDATE deliveries SET state='acked', lease_until=NULL, updated=?4
1813 WHERE room_id=?1 AND reader=?2 AND seq=?3",
1814 key,
1815 )?;
1816 }
1817 (DeliveryState::Leased, Settle::Renew { lease_until }) => {
1818 tx.execute(
1819 "UPDATE deliveries SET lease_until=?5, updated=?4
1820 WHERE room_id=?1 AND reader=?2 AND seq=?3",
1821 params![d.room_id, d.reader, d.seq as i64, now, lease_until],
1822 )?;
1823 }
1824 (DeliveryState::Leased, Settle::Nack { retry_at }) => {
1825 let quarantine = d.attempt >= max_attempts;
1826 tx.execute(
1827 "UPDATE deliveries SET state=?5, token=NULL, lease_until=NULL, retry_at=?6, updated=?4
1828 WHERE room_id=?1 AND reader=?2 AND seq=?3",
1829 params![
1830 d.room_id,
1831 d.reader,
1832 d.seq as i64,
1833 now,
1834 if quarantine { "quarantined" } else { "delayed" },
1835 if quarantine { None } else { Some(retry_at) }
1836 ],
1837 )?;
1838 }
1839 (state, _) => {
1840 return Err(Error::Denied(format!(
1841 "that delivery is {} and can only be acked again",
1842 state.as_str()
1843 )))
1844 }
1845 }
1846 Self::advance_bookmark_in(&tx, &d.room_id, &d.reader)?;
1847 let day_ago = chrono::DateTime::parse_from_rfc3339(now)
1850 .map(|t| {
1851 (t - chrono::Duration::days(1)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1852 })
1853 .unwrap_or_default();
1854 tx.execute(
1855 "DELETE FROM deliveries WHERE room_id=?1 AND reader=?2 AND state='acked'
1856 AND updated<?3 AND seq<=(SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2)",
1857 params![d.room_id, d.reader, day_ago],
1858 )?;
1859 let out = tx.query_row(
1860 "SELECT * FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq=?3",
1861 params![d.room_id, d.reader, d.seq as i64],
1862 Self::row_to_delivery,
1863 )?;
1864 tx.commit()?;
1865 Ok(out)
1866 }
1867
1868 pub fn delivery_ack_seqs(
1872 &self,
1873 room_id: &str,
1874 reader: &str,
1875 seqs: &[u64],
1876 now: &str,
1877 ) -> Result<u64> {
1878 let mut conn = self.lock();
1879 let tx = conn.transaction()?;
1880 for seq in seqs {
1881 tx.execute(
1882 "INSERT INTO deliveries (room_id, reader, seq, state, attempt, updated)
1883 VALUES (?1, ?2, ?3, 'acked', 1, ?4)
1884 ON CONFLICT(room_id, reader, seq) DO UPDATE SET
1885 state='acked', token=NULL, lease_until=NULL, retry_at=NULL, updated=?4
1886 WHERE deliveries.state != 'quarantined'",
1887 params![room_id, reader, *seq as i64, now],
1888 )?;
1889 }
1890 let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
1891 tx.commit()?;
1892 Ok(bm)
1893 }
1894
1895 pub fn delivery_replay(
1897 &self,
1898 room_id: &str,
1899 reader: &str,
1900 seq: u64,
1901 now: &str,
1902 ) -> Result<bool> {
1903 let conn = self.lock();
1904 let n = conn.execute(
1905 "UPDATE deliveries SET state='replay', token=NULL, lease_until=NULL, retry_at=NULL,
1906 attempt=0, updated=?4
1907 WHERE room_id=?1 AND reader=?2 AND seq=?3 AND state='quarantined'",
1908 params![room_id, reader, seq as i64, now],
1909 )?;
1910 Ok(n == 1)
1911 }
1912
1913 pub fn delivery_peek(
1918 &self,
1919 room_id: &str,
1920 reader: &str,
1921 now: &str,
1922 limit: u32,
1923 ) -> Result<Vec<Message>> {
1924 let bm = {
1925 let mut conn = self.lock();
1926 let tx = conn.transaction()?;
1927 let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
1928 tx.commit()?;
1929 bm
1930 };
1931 let conn = self.lock();
1932 let mut stmt = conn.prepare(
1933 "SELECT m.envelope, c.body, c.deleted FROM messages m
1934 LEFT JOIN contents c ON c.msg_id = m.id
1935 LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
1936 WHERE m.room_id=?1 AND m.seq>?3
1937 AND m.from_name != ?2 AND m.type NOT IN ('system', 'control')
1938 AND (d.state IS NULL OR d.state IN ('leased', 'replay')
1939 OR (d.state='delayed' AND d.retry_at<=?4))
1940 ORDER BY m.seq LIMIT ?5",
1941 )?;
1942 let rows = stmt.query_map(params![room_id, reader, bm as i64, now, limit], |r| {
1943 self.row_to_message(r)
1944 })?;
1945 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1946 }
1947
1948 pub fn deliveries(&self, room_id: &str, reader: Option<&str>) -> Result<Vec<Delivery>> {
1951 let conn = self.lock();
1952 let mut stmt = conn.prepare(
1953 "SELECT * FROM deliveries WHERE room_id=?1 AND (?2 IS NULL OR reader=?2)
1954 AND state != 'acked' ORDER BY reader, seq",
1955 )?;
1956 let rows = stmt.query_map(params![room_id, reader], Self::row_to_delivery)?;
1957 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
1958 }
1959
1960 pub fn delivery_next_due(
1963 &self,
1964 room_id: &str,
1965 reader: &str,
1966 now: &str,
1967 ) -> Result<Option<String>> {
1968 let conn = self.lock();
1969 Ok(conn.query_row(
1970 "SELECT MIN(t) FROM (
1971 SELECT lease_until AS t FROM deliveries
1972 WHERE room_id=?1 AND reader=?2 AND state='leased' AND lease_until>?3
1973 UNION ALL
1974 SELECT retry_at AS t FROM deliveries
1975 WHERE room_id=?1 AND reader=?2 AND state='delayed' AND retry_at>?3)",
1976 params![room_id, reader, now],
1977 |r| r.get(0),
1978 )?)
1979 }
1980
1981 pub fn wake_pending(
1987 &self,
1988 room_id: &str,
1989 reader: &str,
1990 now: &str,
1991 trace_limit: u64,
1992 ) -> Result<WakePending> {
1993 let conn = self.lock();
1994 let bm = Self::bookmark_in(&conn, room_id, reader)?;
1995 let mut stmt = conn.prepare_cached(
1996 "SELECT m.seq, m.id, json_extract(m.envelope, '$.trace') FROM messages m
1997 LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
1998 WHERE m.room_id=?1 AND m.from_name != ?2 AND m.type NOT IN ('system', 'control')
1999 AND ((d.state IS NULL AND m.seq > ?3)
2000 OR d.state='replay'
2001 OR (d.state='leased' AND d.lease_until<=?4)
2002 OR (d.state='delayed' AND d.retry_at<=?4))
2003 ORDER BY m.seq LIMIT 10000",
2004 )?;
2005 let rows: Vec<(i64, String, Option<String>)> = stmt
2006 .query_map(params![room_id, reader, bm as i64, now], |r| {
2007 Ok((r.get(0)?, r.get(1)?, r.get(2)?))
2008 })?
2009 .collect::<rusqlite::Result<_>>()?;
2010 let mut sent_on: HashMap<String, u64> = HashMap::new();
2011 let mut out = WakePending::default();
2012 for (seq, id, trace) in rows {
2013 if let Some(t) = trace.filter(|t| !t.is_empty()) {
2014 let n = match sent_on.get(&t) {
2015 Some(n) => *n,
2016 None => {
2017 let n: i64 = conn.query_row(
2018 "SELECT COUNT(*) FROM messages WHERE room_id=?1 AND from_name=?2
2019 AND json_extract(envelope, '$.trace')=?3",
2020 params![room_id, reader, t],
2021 |r| r.get(0),
2022 )?;
2023 sent_on.insert(t.clone(), n as u64);
2024 n as u64
2025 }
2026 };
2027 if n > trace_limit {
2028 out.looped += 1;
2029 continue;
2030 }
2031 }
2032 out.count += 1;
2033 out.newest_seq = seq as u64;
2034 out.newest_id = id;
2035 }
2036 Ok(out)
2037 }
2038
2039 pub fn outbox_add(&self, msg: &Message) -> Result<()> {
2053 let conn = self.lock();
2054 let sealed = self.seal(&serde_json::to_string(msg)?)?;
2055 self.faults.check("outbox.add")?;
2056 conn.execute(
2057 "INSERT OR IGNORE INTO outbox (msg_id, room_id, message, created, sender, state, updated, sig)
2058 VALUES (?1, ?2, ?3, ?4, ?5, 'pending', ?6, ?7)",
2059 params![
2060 msg.id,
2061 msg.room,
2062 sealed,
2063 msg.ts,
2064 msg.from,
2065 crate::message::now_ts(),
2066 msg.sig
2067 ],
2068 )?;
2069 Ok(())
2070 }
2071
2072 fn row_to_outbox(&self, row: &Row<'_>) -> rusqlite::Result<OutboxEntry> {
2073 let state: String = row.get("state")?;
2074 let stored: String = row.get("message")?;
2075 let message = if stored.is_empty() {
2076 None
2077 } else {
2078 self.unseal(&stored)
2079 .and_then(|plain| serde_json::from_str::<Message>(&plain).ok())
2080 };
2081 Ok(OutboxEntry {
2082 msg_id: row.get("msg_id")?,
2083 room_id: row.get("room_id")?,
2084 sender: row.get("sender")?,
2085 state: state.parse().unwrap_or(OutboxState::Pending),
2086 created: row.get("created")?,
2087 attempts: row.get::<_, i64>("attempts")? as u32,
2088 retry_at: row.get("retry_at")?,
2089 reason: row.get("last_error")?,
2090 reason_class: row.get("reason_class")?,
2091 reason_code: row.get("reason_code")?,
2092 updated: row.get("updated")?,
2093 unreadable: !stored.is_empty() && message.is_none(),
2094 message,
2095 })
2096 }
2097
2098 pub fn outbox_get(&self, msg_id: &str) -> Result<Option<OutboxEntry>> {
2100 let conn = self.lock();
2101 Ok(conn
2102 .query_row(
2103 "SELECT * FROM outbox WHERE msg_id=?1",
2104 params![msg_id],
2105 |r| self.row_to_outbox(r),
2106 )
2107 .optional()?)
2108 }
2109
2110 pub fn outbox_lanes(&self, room_id: &str) -> Result<Vec<(String, Vec<OutboxEntry>)>> {
2115 let conn = self.lock();
2116 let mut stmt = conn.prepare(
2117 "SELECT * FROM outbox WHERE room_id=?1 AND state IN ('pending', 'waiting')
2118 ORDER BY sender, rowid",
2119 )?;
2120 let rows = stmt.query_map(params![room_id], |r| self.row_to_outbox(r))?;
2121 let mut lanes: Vec<(String, Vec<OutboxEntry>)> = Vec::new();
2122 for e in rows {
2123 let e = e?;
2124 match lanes.last_mut() {
2125 Some((sender, lane)) if *sender == e.sender => lane.push(e),
2126 _ => lanes.push((e.sender.clone(), vec![e])),
2127 }
2128 }
2129 Ok(lanes)
2130 }
2131
2132 pub fn outbox_entries(&self, room_id: Option<&str>) -> Result<Vec<OutboxEntry>> {
2135 let conn = self.lock();
2136 let mut stmt =
2137 conn.prepare("SELECT * FROM outbox WHERE (?1 IS NULL OR room_id=?1) ORDER BY rowid")?;
2138 let rows = stmt.query_map(params![room_id], |r| self.row_to_outbox(r))?;
2139 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
2140 }
2141
2142 fn outbox_set(&self, sql: &str, args: &[&dyn rusqlite::ToSql]) -> Result<bool> {
2143 let conn = self.lock();
2144 self.faults.check("outbox.state")?;
2145 Ok(conn.execute(sql, args)? == 1)
2146 }
2147
2148 pub fn outbox_wait(
2152 &self,
2153 msg_id: &str,
2154 class: &str,
2155 reason: &str,
2156 code: Option<i32>,
2157 retry_at: &str,
2158 now: &str,
2159 ) -> Result<bool> {
2160 self.outbox_set(
2161 "UPDATE outbox SET state='waiting', attempts=attempts+1, retry_at=?3,
2162 reason_class=?4, last_error=?5, reason_code=?6, updated=?2
2163 WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
2164 &[&msg_id, &now, &retry_at, &class, &reason, &code],
2165 )
2166 }
2167
2168 pub fn outbox_fail(
2171 &self,
2172 msg_id: &str,
2173 reason: &str,
2174 code: Option<i32>,
2175 now: &str,
2176 ) -> Result<bool> {
2177 self.outbox_set(
2178 "UPDATE outbox SET state='failed', attempts=attempts+1, retry_at=NULL,
2179 reason_class='refused', last_error=?3, reason_code=?4, updated=?2
2180 WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
2181 &[&msg_id, &now, &reason, &code],
2182 )
2183 }
2184
2185 #[allow(clippy::too_many_arguments)]
2190 pub fn outbox_unknown(
2191 &self,
2192 msg_id: &str,
2193 reason: &str,
2194 code: Option<i32>,
2195 retry_at: &str,
2196 now: &str,
2197 min_count: u32,
2198 min_secs: i64,
2199 ) -> Result<OutboxState> {
2200 let mut conn = self.lock();
2201 let tx = conn.transaction()?;
2202 let row: Option<(Option<i32>, Option<String>, i64)> = tx
2203 .query_row(
2204 "SELECT unknown_code, unknown_since, unknown_count FROM outbox
2205 WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
2206 params![msg_id],
2207 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
2208 )
2209 .optional()?;
2210 let Some((last_code, since, count)) = row else {
2211 return Err(Error::Invalid(format!(
2212 "{msg_id} is not waiting to be sent"
2213 )));
2214 };
2215 let (since, count) = match (last_code == code, since) {
2216 (true, Some(since)) => (since, count as u32 + 1),
2217 _ => (now.to_string(), 1),
2218 };
2219 let held_for = seconds_between(&since, now);
2220 let state = if count >= min_count && held_for >= min_secs {
2221 OutboxState::Quarantined
2222 } else {
2223 OutboxState::Waiting
2224 };
2225 self.faults.check("outbox.state")?;
2226 tx.execute(
2227 "UPDATE outbox SET state=?2, attempts=attempts+1, retry_at=?3,
2228 reason_class='unknown', last_error=?4, reason_code=?5,
2229 unknown_code=?5, unknown_since=?6, unknown_count=?7, updated=?8
2230 WHERE msg_id=?1",
2231 params![
2232 msg_id,
2233 state.as_str(),
2234 if state == OutboxState::Waiting {
2235 Some(retry_at)
2236 } else {
2237 None
2238 },
2239 reason,
2240 code,
2241 since,
2242 count,
2243 now
2244 ],
2245 )?;
2246 tx.commit()?;
2247 Ok(state)
2248 }
2249
2250 pub fn outbox_retry(&self, msg_id: &str, now: &str) -> Result<bool> {
2253 self.outbox_set(
2254 "UPDATE outbox SET state='pending', attempts=0, retry_at=NULL,
2255 unknown_code=NULL, unknown_since=NULL, unknown_count=0, updated=?2
2256 WHERE msg_id=?1 AND state IN ('failed', 'quarantined') AND message != ''",
2257 &[&msg_id, &now],
2258 )
2259 }
2260
2261 pub fn outbox_drop(&self, msg_id: &str, now: &str) -> Result<bool> {
2264 self.outbox_set(
2265 "UPDATE outbox SET state='dropped', message='', retry_at=NULL, updated=?2
2266 WHERE msg_id=?1 AND state IN ('failed', 'quarantined')",
2267 &[&msg_id, &now],
2268 )
2269 }
2270
2271 pub fn outbox_wake(&self, room_id: &str, class: &str) -> Result<u64> {
2274 let conn = self.lock();
2275 let n = conn.execute(
2276 "UPDATE outbox SET retry_at=NULL
2277 WHERE room_id=?1 AND state='waiting' AND reason_class=?2",
2278 params![room_id, class],
2279 )?;
2280 Ok(n as u64)
2281 }
2282
2283 pub fn outbox_remove(&self, msg_id: &str) -> Result<()> {
2286 let conn = self.lock();
2287 conn.execute(
2288 "DELETE FROM outbox WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
2289 params![msg_id],
2290 )?;
2291 Ok(())
2292 }
2293
2294 pub fn outbox_count(&self, room_id: &str) -> Result<u64> {
2296 Ok(self.outbox_counts(room_id)?.queued())
2297 }
2298
2299 pub fn outbox_counts(&self, room_id: &str) -> Result<OutboxCounts> {
2300 let conn = self.lock();
2301 let mut stmt =
2302 conn.prepare("SELECT state, COUNT(*) FROM outbox WHERE room_id=?1 GROUP BY state")?;
2303 let rows = stmt.query_map(params![room_id], |r| {
2304 Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?))
2305 })?;
2306 let mut c = OutboxCounts::default();
2307 for row in rows {
2308 let (state, n) = row?;
2309 let n = n as u64;
2310 match state.parse().unwrap_or(OutboxState::Pending) {
2311 OutboxState::Pending => c.pending += n,
2312 OutboxState::Waiting => c.waiting += n,
2313 OutboxState::Failed => c.failed += n,
2314 OutboxState::Quarantined => c.quarantined += n,
2315 OutboxState::Dropped => c.dropped += n,
2316 }
2317 }
2318 Ok(c)
2319 }
2320
2321 pub fn invite_record(&self, inv: &Invite) -> Result<()> {
2324 let conn = self.lock();
2325 conn.execute(
2326 "INSERT OR IGNORE INTO invites (nonce, room_id, name, created, expires) VALUES (?1, ?2, ?3, ?4, ?5)",
2327 params![inv.nonce, inv.room_id, inv.name, inv.created, inv.expires],
2328 )?;
2329 Ok(())
2330 }
2331
2332 pub fn invite_use(&self, nonce: &str, used_by: &str, now: &str) -> Result<bool> {
2335 let conn = self.lock();
2336 let n = conn.execute(
2337 "UPDATE invites SET used_at=?2, used_by=?3 WHERE nonce=?1 AND used_at IS NULL",
2338 params![nonce, now, used_by],
2339 )?;
2340 Ok(n == 1)
2341 }
2342
2343 fn count_since(
2346 conn: &Connection,
2347 room_id: &str,
2348 from: Option<&str>,
2349 since: &str,
2350 ) -> Result<u32> {
2351 let n: i64 = match from {
2352 Some(f) => conn.query_row(
2353 "SELECT COUNT(*) FROM messages WHERE room_id=?1 AND from_name=?2 AND received>=?3",
2354 params![room_id, f, since],
2355 |r| r.get(0),
2356 )?,
2357 None => conn.query_row(
2358 "SELECT COUNT(*) FROM messages WHERE room_id=?1 AND received>=?2",
2359 params![room_id, since],
2360 |r| r.get(0),
2361 )?,
2362 };
2363 Ok(n as u32)
2364 }
2365
2366 pub fn check_limits(&self, room_id: &str, from: &str, limits: &Limits) -> Result<bool> {
2370 let conn = self.lock();
2371 let now = chrono::Utc::now();
2372 let minute_ago =
2373 (now - chrono::Duration::minutes(1)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
2374 let day_start = now
2375 .date_naive()
2376 .and_hms_opt(0, 0, 0)
2377 .expect("midnight")
2378 .and_utc()
2379 .to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
2380 let per_minute = Self::count_since(&conn, room_id, Some(from), &minute_ago)?;
2381 if per_minute >= limits.per_minute_per_sender {
2382 let oldest: Option<String> = conn.query_row(
2385 "SELECT MIN(received) FROM messages WHERE room_id=?1 AND from_name=?2 AND received>=?3",
2386 params![room_id, from, minute_ago],
2387 |r| r.get(0),
2388 )?;
2389 let wait = oldest
2390 .map(|o| {
2391 60 - seconds_between(
2392 &o,
2393 &now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
2394 )
2395 })
2396 .unwrap_or(60)
2397 .clamp(1, 60) as u64;
2398 return Err(Error::OverBudget(
2399 format!(
2400 "{from} sent {per_minute} messages in the last minute; the limit is {}",
2401 limits.per_minute_per_sender
2402 ),
2403 Some(wait),
2404 ));
2405 }
2406 let today = Self::count_since(&conn, room_id, None, &day_start)?;
2407 if today >= limits.daily_per_room {
2408 let midnight = now
2409 .date_naive()
2410 .and_hms_opt(0, 0, 0)
2411 .expect("midnight")
2412 .and_utc()
2413 + chrono::Duration::days(1);
2414 let wait = (midnight - now).num_seconds().max(1) as u64;
2415 return Err(Error::OverBudget(
2416 format!(
2417 "room has used its daily budget of {} messages",
2418 limits.daily_per_room
2419 ),
2420 Some(wait),
2421 ));
2422 }
2423 let alert_at = limits.daily_per_room as u64 * limits.burst_alert_percent as u64 / 100;
2424 Ok(today as u64 + 1 == alert_at)
2425 }
2426}
2427
2428#[derive(Debug, Clone, PartialEq, Eq)]
2430pub struct StoredFile {
2431 pub room_id: String,
2432 pub hash: String,
2433 pub size: u64,
2434 pub kind: String,
2435 pub uploader: String,
2436 pub created: String,
2437}
2438
2439#[derive(Debug, Clone, PartialEq, Eq)]
2441pub struct FileRefRow {
2442 pub msg_id: String,
2443 pub room_id: String,
2444 pub file: crate::files::FileRef,
2445 pub from: String,
2446 pub to: Option<String>,
2447}
2448
2449pub const ORPHAN_FILE_SECS: i64 = 24 * 3600;
2452
2453impl Store {
2454 fn row_to_ref(r: &Row) -> rusqlite::Result<FileRefRow> {
2457 Ok(FileRefRow {
2458 msg_id: r.get("msg_id")?,
2459 room_id: r.get("room_id")?,
2460 file: crate::files::FileRef {
2461 id: r.get("hash")?,
2462 name: r.get("name")?,
2463 size: r.get::<_, i64>("size")? as u64,
2464 mime: r.get("type")?,
2465 },
2466 from: r.get("from_name")?,
2467 to: r.get("to_name")?,
2468 })
2469 }
2470
2471 pub fn file_refs(&self, room_id: &str, hash: &str) -> Result<Vec<FileRefRow>> {
2473 let conn = self.lock();
2474 let mut stmt = conn.prepare(
2475 "SELECT f.* FROM file_refs f JOIN messages m ON m.id = f.msg_id
2476 WHERE f.room_id=?1 AND f.hash=?2 ORDER BY m.seq",
2477 )?;
2478 let rows = stmt.query_map(params![room_id, hash], Self::row_to_ref)?;
2479 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
2480 }
2481
2482 pub fn file_refs_like(&self, room_id: &str, prefix: &str) -> Result<Vec<FileRefRow>> {
2485 let conn = self.lock();
2486 let mut stmt = conn.prepare(
2487 "SELECT f.* FROM file_refs f JOIN messages m ON m.id = f.msg_id
2488 WHERE f.room_id=?1 AND substr(f.hash, 1, length(?2))=?2 ORDER BY f.hash, m.seq",
2489 )?;
2490 let rows = stmt.query_map(params![room_id, prefix], Self::row_to_ref)?;
2491 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
2492 }
2493
2494 pub fn file_refs_of(&self, msg_id: &str) -> Result<Vec<FileRefRow>> {
2496 let conn = self.lock();
2497 let mut stmt = conn.prepare("SELECT * FROM file_refs WHERE msg_id=?1 ORDER BY rowid")?;
2498 let rows = stmt.query_map(params![msg_id], Self::row_to_ref)?;
2499 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
2500 }
2501
2502 fn row_to_file(r: &Row) -> rusqlite::Result<StoredFile> {
2503 Ok(StoredFile {
2504 room_id: r.get("room_id")?,
2505 hash: r.get("hash")?,
2506 size: r.get::<_, i64>("size")? as u64,
2507 kind: r.get("kind")?,
2508 uploader: r.get("uploader")?,
2509 created: r.get("created")?,
2510 })
2511 }
2512
2513 pub fn file_stored(&self, room_id: &str, hash: &str) -> Result<Option<StoredFile>> {
2515 let conn = self.lock();
2516 Ok(conn
2517 .query_row(
2518 "SELECT * FROM files WHERE room_id=?1 AND hash=?2",
2519 params![room_id, hash],
2520 Self::row_to_file,
2521 )
2522 .optional()?)
2523 }
2524
2525 pub fn files_stored(&self, room_id: &str) -> Result<Vec<StoredFile>> {
2527 let conn = self.lock();
2528 let mut stmt = conn.prepare("SELECT * FROM files WHERE room_id=?1 ORDER BY created")?;
2529 let rows = stmt.query_map(params![room_id], Self::row_to_file)?;
2530 Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
2531 }
2532
2533 pub fn file_add(&self, f: &StoredFile) -> Result<bool> {
2536 let conn = self.lock();
2537 let n = conn.execute(
2538 "INSERT OR IGNORE INTO files (room_id, hash, size, kind, uploader, created)
2539 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2540 params![
2541 f.room_id,
2542 f.hash,
2543 f.size as i64,
2544 f.kind,
2545 f.uploader,
2546 f.created
2547 ],
2548 )?;
2549 Ok(n == 1)
2550 }
2551
2552 pub fn file_bytes(&self, room_id: &str, uploader: &str, since: &str) -> Result<(u64, u64)> {
2555 let conn = self.lock();
2556 let room: i64 = conn.query_row(
2557 "SELECT COALESCE(SUM(size), 0) FROM files WHERE room_id=?1",
2558 params![room_id],
2559 |r| r.get(0),
2560 )?;
2561 let mine: i64 = conn.query_row(
2562 "SELECT COALESCE(SUM(size), 0) FROM files WHERE room_id=?1 AND uploader=?2 AND created>=?3",
2563 params![room_id, uploader, since],
2564 |r| r.get(0),
2565 )?;
2566 Ok((room as u64, mine as u64))
2567 }
2568
2569 pub fn file_fetched(&self, room_id: &str, hash: &str, member: &str, now: &str) -> Result<()> {
2571 let conn = self.lock();
2572 conn.execute(
2573 "INSERT OR IGNORE INTO file_fetches (room_id, hash, member, at) VALUES (?1, ?2, ?3, ?4)",
2574 params![room_id, hash, member, now],
2575 )?;
2576 Ok(())
2577 }
2578
2579 pub fn file_sweep(
2586 &self,
2587 room_id: &str,
2588 now: chrono::DateTime<chrono::Utc>,
2589 keep_days: u64,
2590 max_days: u64,
2591 ) -> Result<Vec<String>> {
2592 let ts =
2593 |t: chrono::DateTime<chrono::Utc>| t.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
2594 let days = |d: u64| chrono::Duration::days(d.min(100_000) as i64);
2595 let too_old = ts(now - days(max_days));
2596 let orphan = ts(now - chrono::Duration::seconds(ORPHAN_FILE_SECS));
2597 let settled = ts(now - days(keep_days));
2598 let members: Vec<String> = self
2599 .members(room_id)?
2600 .into_iter()
2601 .filter(|m| !m.revoked)
2602 .map(|m| m.name)
2603 .collect();
2604 let mut due = Vec::new();
2605 for f in self.files_stored(room_id)? {
2606 let refs = self.file_refs(room_id, &f.hash)?;
2607 let go = if f.created < too_old {
2608 true
2609 } else if refs.is_empty() {
2610 f.created < orphan
2611 } else {
2612 let mut audience: Vec<String> = Vec::new();
2613 for r in &refs {
2614 match &r.to {
2615 Some(to) => audience.push(to.clone()),
2616 None => audience.extend(members.iter().filter(|m| **m != r.from).cloned()),
2617 }
2618 }
2619 audience.sort();
2620 audience.dedup();
2621 let fetched: HashMap<String, String> = {
2622 let conn = self.lock();
2623 let mut stmt = conn.prepare(
2624 "SELECT member, at FROM file_fetches WHERE room_id=?1 AND hash=?2",
2625 )?;
2626 let rows = stmt.query_map(params![room_id, f.hash], |r| {
2627 Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
2628 })?;
2629 rows.collect::<rusqlite::Result<HashMap<_, _>>>()?
2630 };
2631 let mut last = f.created.clone();
2632 let mut all = true;
2633 for who in &audience {
2634 match fetched.get(who) {
2635 Some(at) => last = last.max(at.clone()),
2636 None => all = false,
2637 }
2638 }
2639 all && last < settled
2640 };
2641 if go {
2642 due.push(f.hash);
2643 }
2644 }
2645 let conn = self.lock();
2646 for h in &due {
2647 conn.execute(
2648 "DELETE FROM files WHERE room_id=?1 AND hash=?2",
2649 params![room_id, h],
2650 )?;
2651 conn.execute(
2652 "DELETE FROM file_fetches WHERE room_id=?1 AND hash=?2",
2653 params![room_id, h],
2654 )?;
2655 }
2656 Ok(due)
2657 }
2658}
2659
2660#[cfg(test)]
2661mod tests {
2662 use super::*;
2663 use crate::keys::{Identity, Kind};
2664 use crate::message::{Draft, MessageType};
2665
2666 fn room(owner: &Identity) -> Room {
2667 Room {
2668 id: "r_test".into(),
2669 name: "ops".into(),
2670 about: "".into(),
2671 owner: owner.public(),
2672 created: "2026-01-01T00:00:00Z".into(),
2673 retention_days: None,
2674 class: DataClass::Internal,
2675 paused: false,
2676 hold: false,
2677 closed: false,
2678 home_node: "node".into(),
2679 home_hints: serde_json::Value::Null,
2680 }
2681 }
2682
2683 fn msg(owner: &Identity, text: &str) -> Message {
2684 Message::new(
2685 Draft {
2686 room: "r_test".into(),
2687 from: "haris".into(),
2688 text: text.into(),
2689 kind: Some(MessageType::Chat),
2690 ..Default::default()
2691 },
2692 owner,
2693 )
2694 .unwrap()
2695 }
2696
2697 #[test]
2698 fn wake_pending_counts_what_waits_and_nothing_else() {
2699 let owner = Identity::generate("haris", Kind::Human);
2700 let s = Store::open_memory().unwrap();
2701 s.create_room(&room(&owner)).unwrap();
2702 let put = |from: &str, kind: MessageType, trace: Option<&str>| {
2703 let mut m = Message::new(
2704 Draft {
2705 room: "r_test".into(),
2706 from: from.into(),
2707 text: "hi".into(),
2708 kind: Some(kind),
2709 trace: trace.map(String::from),
2710 ..Default::default()
2711 },
2712 &owner,
2713 )
2714 .unwrap();
2715 s.sequence_and_append(&mut m).unwrap();
2716 m
2717 };
2718 let now = "2026-09-27T12:00:00Z";
2719 for _ in 0..3 {
2720 put("haris", MessageType::Task, None);
2721 }
2722 put("bob", MessageType::Chat, None);
2724 let newest = put("haris", MessageType::System, None);
2725 let p = s.wake_pending("r_test", "bob", now, 5).unwrap();
2726 assert_eq!((p.count, p.newest_seq), (3, 3));
2727 assert_ne!(p.newest_id, newest.id);
2728
2729 let (leased, _) = s
2732 .delivery_lease("r_test", "bob", now, "2026-09-27T12:10:00Z", 5)
2733 .unwrap()
2734 .unwrap();
2735 assert_eq!(leased.seq, 1);
2736 assert_eq!(s.wake_pending("r_test", "bob", now, 5).unwrap().count, 2);
2737 assert_eq!(
2738 s.wake_pending("r_test", "bob", "2026-09-27T12:11:00Z", 5)
2739 .unwrap()
2740 .count,
2741 3
2742 );
2743
2744 for _ in 0..6 {
2746 put("bob", MessageType::Reply, Some("ping-pong"));
2747 }
2748 put("haris", MessageType::Reply, Some("ping-pong"));
2749 put("haris", MessageType::Task, Some("fresh"));
2750 let p = s.wake_pending("r_test", "bob", now, 5).unwrap();
2751 assert_eq!((p.count, p.looped), (3, 1));
2752 }
2753
2754 #[test]
2755 fn rooms_and_members() {
2756 let owner = Identity::generate("haris", Kind::Human);
2757 let s = Store::open_memory().unwrap();
2758 s.create_room(&room(&owner)).unwrap();
2759 assert!(matches!(
2760 s.create_room(&room(&owner)),
2761 Err(Error::NameTaken(_))
2762 ));
2763 assert_eq!(s.room("ops").unwrap().id, "r_test");
2764 assert!(matches!(s.room("nope"), Err(Error::NotInRoom(_))));
2765 let m = Member {
2766 room_id: "r_test".into(),
2767 name: "haris".into(),
2768 key: owner.public(),
2769 kind: Kind::Human,
2770 role: Role::Approver,
2771 node: None,
2772 granted_by: owner.public(),
2773 expires_at: None,
2774 joined_at: "2026-01-01T00:00:00Z".into(),
2775 last_seen: None,
2776 muted: false,
2777 revoked: false,
2778 profile: serde_json::Value::Null,
2779 };
2780 s.upsert_member(&m, Some("default")).unwrap();
2781 assert_eq!(s.members("r_test").unwrap().len(), 1);
2782 assert_eq!(s.local_members("r_test").unwrap()[0].identity, "default");
2783 assert_eq!(
2784 s.member_by_key("r_test", &owner.public())
2785 .unwrap()
2786 .unwrap()
2787 .name,
2788 "haris"
2789 );
2790 }
2791
2792 #[test]
2793 fn chain_append_dedup_and_bookmarks() {
2794 let owner = Identity::generate("haris", Kind::Human);
2795 let s = Store::open_memory().unwrap();
2796 s.create_room(&room(&owner)).unwrap();
2797 assert_eq!(
2798 s.chain_head("r_test").unwrap(),
2799 (0, GENESIS_PREV.to_string())
2800 );
2801 let mut a = msg(&owner, "one");
2802 s.sequence_and_append(&mut a).unwrap();
2803 assert_eq!(a.seq, 1);
2804 assert_eq!(a.prev, GENESIS_PREV);
2805 let mut b = msg(&owner, "two");
2806 s.sequence_and_append(&mut b).unwrap();
2807 assert_eq!(b.seq, 2);
2808 assert_eq!(b.prev, a.chain_hash());
2809 assert!(!s.append(&b).unwrap());
2811 let mut c = msg(&owner, "three");
2813 c.sequence(3, "sha256:bad");
2814 assert!(s.append(&c).is_err());
2815 c.sequence(3, &b.chain_hash());
2816 assert!(s.append(&c).unwrap());
2817 let all = s.messages_after("r_test", 0, 100).unwrap();
2819 assert_eq!(all, vec![a.clone(), b.clone(), c.clone()]);
2820 assert!(all[0].verify(&owner.public()).is_ok());
2821 assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
2823 s.set_bookmark("r_test", "haris", 2).unwrap();
2824 s.set_bookmark("r_test", "haris", 1).unwrap();
2825 assert_eq!(s.bookmark("r_test", "haris").unwrap(), 2);
2826 assert_eq!(s.bookmark("r_test", "bob").unwrap(), 0);
2827 assert_eq!(s.messages_after("r_test", 2, 100).unwrap(), vec![c]);
2828 assert_eq!(s.message_count("r_test").unwrap(), 3);
2830 }
2831
2832 #[test]
2833 fn outbox_and_invites() {
2834 let owner = Identity::generate("haris", Kind::Human);
2835 let s = Store::open_memory().unwrap();
2836 s.create_room(&room(&owner)).unwrap();
2837 let m = msg(&owner, "queued");
2838 s.outbox_add(&m).unwrap();
2839 s.outbox_add(&m).unwrap();
2840 assert_eq!(s.outbox_count("r_test").unwrap(), 1);
2841 let lanes = s.outbox_lanes("r_test").unwrap();
2842 assert_eq!(lanes.len(), 1);
2843 assert_eq!(lanes[0].0, "haris");
2844 assert_eq!(lanes[0].1[0].message.as_ref(), Some(&m));
2845 s.outbox_remove(&m.id).unwrap();
2846 assert_eq!(s.outbox_count("r_test").unwrap(), 0);
2847
2848 let inv = Invite::create(
2849 crate::invite::InviteSpec {
2850 room_id: "r_test".into(),
2851 room_name: "ops".into(),
2852 name: "bob".into(),
2853 kind: Kind::Agent,
2854 role: Role::TaskGiver,
2855 home_node: "node".into(),
2856 home_hints: serde_json::Value::Null,
2857 for_node: None,
2858 ttl_hours: None,
2859 },
2860 &owner,
2861 )
2862 .unwrap();
2863 s.invite_record(&inv).unwrap();
2864 assert!(s.invite_use(&inv.nonce, "key", "now").unwrap());
2865 assert!(!s.invite_use(&inv.nonce, "key2", "now").unwrap());
2866 assert!(!s.invite_use("i_unknown", "key", "now").unwrap());
2867 }
2868
2869 #[test]
2870 fn limits() {
2871 let owner = Identity::generate("haris", Kind::Human);
2872 let s = Store::open_memory().unwrap();
2873 s.create_room(&room(&owner)).unwrap();
2874 let limits = Limits {
2875 per_minute_per_sender: 2,
2876 daily_per_room: 3,
2877 burst_alert_percent: 100,
2878 };
2879 s.check_limits("r_test", "haris", &limits).unwrap();
2880 let mut a = msg(&owner, "1");
2881 s.sequence_and_append(&mut a).unwrap();
2882 let mut b = msg(&owner, "2");
2883 s.sequence_and_append(&mut b).unwrap();
2884 assert!(matches!(
2885 s.check_limits("r_test", "haris", &limits),
2886 Err(Error::OverBudget(..))
2887 ));
2888 assert!(s.check_limits("r_test", "bob", &limits).unwrap());
2891 }
2892
2893 #[test]
2894 fn appending_the_same_message_twice_stores_it_once() {
2895 let owner = Identity::generate("haris", Kind::Human);
2896 let s = Store::open_memory().unwrap();
2897 s.create_room(&room(&owner)).unwrap();
2898 let mut a = msg(&owner, "1");
2899 let mut again = a.clone();
2900 assert!(s.sequence_and_append(&mut a).unwrap());
2901 assert!(!s.sequence_and_append(&mut again).unwrap());
2903 assert_eq!(again.seq, 0);
2904 assert_eq!(s.message_count("r_test").unwrap(), 1);
2905 }
2906
2907 #[test]
2908 fn a_backdated_ts_still_counts_against_the_limits() {
2909 let owner = Identity::generate("haris", Kind::Human);
2910 let s = Store::open_memory().unwrap();
2911 s.create_room(&room(&owner)).unwrap();
2912 let limits = Limits {
2913 per_minute_per_sender: 2,
2914 daily_per_room: 100,
2915 burst_alert_percent: 100,
2916 };
2917 for text in ["1", "2"] {
2918 let mut m = msg(&owner, text);
2919 m.ts = "2000-01-01T00:00:00Z".into();
2920 s.sequence_and_append(&mut m).unwrap();
2921 }
2922 assert!(matches!(
2923 s.check_limits("r_test", "haris", &limits),
2924 Err(Error::OverBudget(..))
2925 ));
2926 }
2927
2928 fn msg_from(who: &Identity, from: &str, text: &str) -> Message {
2929 Message::new(
2930 Draft {
2931 room: "r_test".into(),
2932 from: from.into(),
2933 text: text.into(),
2934 kind: Some(MessageType::Chat),
2935 ..Default::default()
2936 },
2937 who,
2938 )
2939 .unwrap()
2940 }
2941
2942 fn temp_db(tag: &str) -> std::path::PathBuf {
2943 static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
2944 let dir = std::env::temp_dir().join(format!(
2945 "diavlos-store-{tag}-{}-{}-{}",
2946 std::process::id(),
2947 NEXT.fetch_add(1, std::sync::atomic::Ordering::SeqCst),
2948 std::time::SystemTime::now()
2949 .duration_since(std::time::UNIX_EPOCH)
2950 .unwrap()
2951 .as_nanos()
2952 ));
2953 std::fs::create_dir_all(&dir).unwrap();
2954 dir.join("diavlos.db")
2955 }
2956
2957 const T0: &str = "2026-01-01T00:00:00Z";
2958
2959 #[test]
2960 fn the_outbox_accounts_for_every_message() {
2961 let owner = Identity::generate("haris", Kind::Human);
2962 let s = Store::open_memory().unwrap();
2963 s.create_room(&room(&owner)).unwrap();
2964 let msgs: Vec<Message> = (0..6)
2965 .map(|i| {
2966 msg_from(
2967 &owner,
2968 if i % 2 == 0 { "haris" } else { "bob" },
2969 &format!("m{i}"),
2970 )
2971 })
2972 .collect();
2973 for m in &msgs {
2974 s.outbox_add(m).unwrap();
2975 }
2976 let lanes = s.outbox_lanes("r_test").unwrap();
2978 assert_eq!(lanes.len(), 2);
2979 assert!(lanes.iter().all(|(_, l)| l.len() == 3));
2980
2981 let mut delivered = msgs[0].clone();
2983 s.sequence_and_append(&mut delivered).unwrap();
2984 s.outbox_wait(
2985 &msgs[1].id,
2986 "transport",
2987 "offline",
2988 Some(3),
2989 "2026-01-01T00:01:00Z",
2990 T0,
2991 )
2992 .unwrap();
2993 s.outbox_fail(&msgs[2].id, "denied", Some(6), T0).unwrap();
2994 s.outbox_fail(&msgs[3].id, "denied", Some(6), T0).unwrap();
2995 s.outbox_drop(&msgs[3].id, T0).unwrap();
2996 for i in 0..5 {
2997 let now = format!("2026-01-01T0{i}:00:00Z");
2998 s.outbox_unknown(&msgs[4].id, "huh", Some(1), &now, &now, 5, 3600)
2999 .unwrap();
3000 }
3001 let c = s.outbox_counts("r_test").unwrap();
3002 assert_eq!(
3003 c,
3004 OutboxCounts {
3005 pending: 1,
3006 waiting: 1,
3007 failed: 1,
3008 quarantined: 1,
3009 dropped: 1
3010 }
3011 );
3012 let in_chain = s.messages_after("r_test", 0, 100).unwrap().len() as u64;
3014 assert_eq!(
3015 in_chain + c.pending + c.waiting + c.failed + c.quarantined + c.dropped,
3016 msgs.len() as u64
3017 );
3018 let dropped = s.outbox_get(&msgs[3].id).unwrap().unwrap();
3020 assert_eq!(dropped.state, OutboxState::Dropped);
3021 assert!(dropped.message.is_none());
3022 assert!(s.outbox_retry(&msgs[4].id, T0).unwrap());
3024 let back = s.outbox_get(&msgs[4].id).unwrap().unwrap();
3025 assert_eq!(back.state, OutboxState::Pending);
3026 assert_eq!(back.message.as_ref(), Some(&msgs[4]));
3027 assert!(!s.outbox_retry(&msgs[5].id, T0).unwrap());
3029 assert!(!s.outbox_drop(&msgs[5].id, T0).unwrap());
3030 assert!(!s.outbox_retry(&msgs[3].id, T0).unwrap());
3031 }
3032
3033 #[test]
3034 fn an_unknown_answer_is_quarantined_only_after_both_count_and_time() {
3035 let owner = Identity::generate("haris", Kind::Human);
3036 let s = Store::open_memory().unwrap();
3037 s.create_room(&room(&owner)).unwrap();
3038 let m = msg(&owner, "x");
3039 s.outbox_add(&m).unwrap();
3040 let at = |mins: u32| format!("2026-01-01T{:02}:{:02}:00Z", mins / 60, mins % 60);
3041 for i in 0..20 {
3043 let st = s
3044 .outbox_unknown(&m.id, "huh", Some(1), &at(i), &at(i), 5, 3600)
3045 .unwrap();
3046 assert_eq!(st, OutboxState::Waiting);
3047 }
3048 let st = s
3050 .outbox_unknown(&m.id, "other", Some(9), &at(61), &at(61), 5, 3600)
3051 .unwrap();
3052 assert_eq!(st, OutboxState::Waiting);
3053 for i in [90, 120, 180] {
3055 let st = s
3056 .outbox_unknown(&m.id, "other", Some(9), &at(i), &at(i), 5, 3600)
3057 .unwrap();
3058 assert_eq!(st, OutboxState::Waiting);
3059 }
3060 let st = s
3062 .outbox_unknown(&m.id, "other", Some(9), &at(200), &at(200), 5, 3600)
3063 .unwrap();
3064 assert_eq!(st, OutboxState::Quarantined);
3065 assert_eq!(s.outbox_count("r_test").unwrap(), 0);
3066 assert_eq!(s.outbox_counts("r_test").unwrap().quarantined, 1);
3067 }
3068
3069 #[test]
3070 fn queued_messages_are_sealed_and_old_plain_ones_get_sealed_too() {
3071 let owner = Identity::generate("haris", Kind::Human);
3072 let key: [u8; 32] = rand::random();
3073 let canary = "CANARY-4f1b9e-queued-text";
3074 let path = temp_db("seal");
3075
3076 let ids: Vec<String> = {
3078 let s = Store::open_with_key(&path, None).unwrap();
3079 s.create_room(&room(&owner)).unwrap();
3080 (0..3)
3081 .map(|i| {
3082 let m = msg(&owner, &format!("{canary} {i}"));
3083 s.outbox_add(&m).unwrap();
3084 m.id
3085 })
3086 .collect()
3087 };
3088
3089 {
3092 let s = Store::open_unsealed(&path, Some(key)).unwrap();
3093 s.faults.arm_after("migration.row", 2);
3094 assert!(s.seal_legacy_outbox(2).is_err());
3095 let entries = s.outbox_entries(None).unwrap();
3096 assert_eq!(entries.len(), 3);
3097 assert!(
3098 entries.iter().all(|e| e.message.is_some()),
3099 "plain and sealed both read"
3100 );
3101 }
3102
3103 {
3105 let s = Store::open_with_key(&path, Some(key)).unwrap();
3106 let entries = s.outbox_entries(None).unwrap();
3107 assert_eq!(entries.len(), 3);
3108 for (e, id) in entries.iter().zip(&ids) {
3109 assert_eq!(&e.msg_id, id);
3110 assert!(e.message.as_ref().unwrap().text.starts_with(canary));
3111 }
3112 s.outbox_add(&msg(&owner, &format!("{canary} new")))
3114 .unwrap();
3115 s.outbox_fail(&ids[0], "no", Some(6), T0).unwrap();
3117 s.outbox_fail(&ids[1], "no", Some(6), T0).unwrap();
3118 s.outbox_drop(&ids[1], T0).unwrap();
3119 s.checkpoint().unwrap();
3120 }
3121 let mut bytes = std::fs::read(&path).unwrap();
3122 bytes.extend(std::fs::read(path.with_extension("db-wal")).unwrap_or_default());
3123 let text = String::from_utf8_lossy(&bytes);
3124 assert!(
3125 !text.contains(canary),
3126 "plain queued text left in the database files"
3127 );
3128
3129 let s = Store::open_with_key(&path, None).unwrap();
3131 let e = s.outbox_get(&ids[2]).unwrap().unwrap();
3132 assert!(e.unreadable && e.message.is_none());
3133 let _ = std::fs::remove_dir_all(path.parent().unwrap());
3134 }
3135
3136 #[test]
3137 fn a_full_disk_refuses_the_message_and_leaves_nothing_half_written() {
3138 let owner = Identity::generate("haris", Kind::Human);
3139 let s = Store::open_memory_with_key(rand::random()).unwrap();
3140 s.create_room(&room(&owner)).unwrap();
3141 s.cap_size_for_test(4).unwrap();
3142 let big = "x".repeat(3000);
3143 let mut kept = Vec::new();
3144 let err = loop {
3145 let m = msg(&owner, &big);
3146 match s.outbox_add(&m) {
3147 Ok(()) => kept.push(m.id),
3148 Err(e) => break (e, m.id),
3149 }
3150 assert!(kept.len() < 1000, "never filled up");
3151 };
3152 match &err.0 {
3153 Error::Db(rusqlite::Error::SqliteFailure(f, _)) => {
3154 assert_eq!(f.code, rusqlite::ErrorCode::DiskFull)
3155 }
3156 other => panic!("expected a full disk, got {other}"),
3157 }
3158 assert!(s.outbox_get(&err.1).unwrap().is_none());
3160 for id in &kept {
3161 assert!(s.outbox_get(id).unwrap().unwrap().message.is_some());
3162 }
3163 let before = s.message_count("r_test").unwrap();
3165 let mut m = msg(&owner, &big);
3166 assert!(s.sequence_and_append(&mut m).is_err());
3167 assert_eq!(s.message_count("r_test").unwrap(), before);
3168 }
3169
3170 #[test]
3171 fn a_queued_copy_of_a_message_already_in_the_chain_is_cleared() {
3172 let owner = Identity::generate("haris", Kind::Human);
3173 let s = Store::open_memory().unwrap();
3174 s.create_room(&room(&owner)).unwrap();
3175 let mut m = msg(&owner, "stored at the home, answer lost");
3176 s.sequence_and_append(&mut m).unwrap();
3177 s.outbox_add(&m).unwrap();
3179 assert!(!s.append(&m).unwrap(), "not stored twice");
3180 assert!(s.outbox_get(&m.id).unwrap().is_none());
3181 let mut other = msg(&owner, "something else");
3183 other.id = m.id.clone();
3184 s.outbox_add(&other).unwrap();
3185 let _ = s.append(&m);
3186 assert!(s.outbox_get(&m.id).unwrap().is_some());
3187 }
3188
3189 fn at(secs: i64) -> String {
3190 (chrono::DateTime::parse_from_rfc3339(T0).unwrap() + chrono::Duration::seconds(secs))
3191 .to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
3192 }
3193
3194 fn owed(n: usize) -> (Store, Vec<u64>) {
3196 let owner = Identity::generate("haris", Kind::Human);
3197 let s = Store::open_memory().unwrap();
3198 s.create_room(&room(&owner)).unwrap();
3199 let mut seqs = Vec::new();
3200 for i in 0..n {
3201 let mut m = msg_from(&owner, "bob", &format!("m{i}"));
3202 s.sequence_and_append(&mut m).unwrap();
3203 seqs.push(m.seq);
3204 }
3205 (s, seqs)
3206 }
3207
3208 fn lease(s: &Store, now: i64) -> Option<Delivery> {
3209 s.delivery_lease("r_test", "haris", &at(now), &at(now + 600), 5)
3210 .unwrap()
3211 .map(|(_, d)| d)
3212 }
3213
3214 fn ack(s: &Store, d: &Delivery, now: i64) -> Result<Delivery> {
3215 s.delivery_settle(d.token.as_deref().unwrap(), &Settle::Ack, &at(now), 5)
3216 }
3217
3218 #[test]
3219 fn the_bookmark_never_passes_a_message_still_owed() {
3220 let (s, seqs) = owed(3);
3221 let a = lease(&s, 0).unwrap();
3222 let b = lease(&s, 0).unwrap();
3223 let c = lease(&s, 0).unwrap();
3224 assert_eq!((a.seq, b.seq, c.seq), (seqs[0], seqs[1], seqs[2]));
3225 assert!(lease(&s, 0).is_none(), "all three are out");
3226 ack(&s, &c, 1).unwrap();
3228 assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
3229 ack(&s, &a, 1).unwrap();
3230 assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
3231 ack(&s, &b, 1).unwrap();
3232 assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[2]);
3233 ack(&s, &b, 2).unwrap();
3235 }
3236
3237 #[test]
3238 fn a_lease_that_runs_out_is_handed_out_again_and_the_old_token_is_refused() {
3239 let (s, seqs) = owed(1);
3240 let first = lease(&s, 0).unwrap();
3241 assert_eq!(first.attempt, 1);
3242 assert!(lease(&s, 599).is_none(), "still leased");
3243 let second = lease(&s, 600).unwrap();
3244 assert_eq!(second.seq, seqs[0]);
3245 assert_eq!(second.attempt, 2);
3246 assert_ne!(first.token, second.token);
3247 assert!(matches!(ack(&s, &first, 700), Err(Error::Denied(_))));
3249 let renew = Settle::Renew {
3250 lease_until: at(5000),
3251 };
3252 assert!(s
3253 .delivery_settle(first.token.as_deref().unwrap(), &renew, &at(700), 5)
3254 .is_err());
3255 assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
3256 ack(&s, &second, 700).unwrap();
3257 assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
3258 }
3259
3260 #[test]
3261 fn a_late_ack_counts_if_nobody_was_handed_it_since() {
3262 let (s, seqs) = owed(1);
3263 let d = lease(&s, 0).unwrap();
3264 ack(&s, &d, 5000).unwrap();
3265 assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
3266 }
3267
3268 #[test]
3269 fn renew_holds_it_and_nack_hands_it_back_later() {
3270 let (s, _) = owed(1);
3271 let d = lease(&s, 0).unwrap();
3272 let t = d.token.clone().unwrap();
3273 s.delivery_settle(
3274 &t,
3275 &Settle::Renew {
3276 lease_until: at(2000),
3277 },
3278 &at(500),
3279 5,
3280 )
3281 .unwrap();
3282 assert!(lease(&s, 1500).is_none(), "renewed past the first lease");
3283 s.delivery_settle(&t, &Settle::Nack { retry_at: at(3000) }, &at(1600), 5)
3284 .unwrap();
3285 assert!(lease(&s, 2999).is_none(), "not before retry_at");
3286 let again = lease(&s, 3000).unwrap();
3287 assert_eq!(again.attempt, 2);
3288 }
3289
3290 #[test]
3291 fn five_strikes_quarantine_it_the_lane_moves_on_and_replay_brings_it_back() {
3292 let (s, seqs) = owed(2);
3293 let mut now = 0;
3294 for i in 1..=5 {
3295 let d = lease(&s, now).unwrap();
3296 assert_eq!((d.seq, d.attempt), (seqs[0], i));
3297 now += 600;
3298 }
3299 let next = lease(&s, now).unwrap();
3301 assert_eq!(next.seq, seqs[1]);
3302 let q = s.deliveries("r_test", Some("haris")).unwrap();
3303 assert!(q
3304 .iter()
3305 .any(|d| d.seq == seqs[0] && d.state == DeliveryState::Quarantined));
3306 ack(&s, &next, now).unwrap();
3307 assert_eq!(
3308 s.bookmark("r_test", "haris").unwrap(),
3309 seqs[1],
3310 "quarantine counts as settled"
3311 );
3312 assert!(s
3314 .delivery_replay("r_test", "haris", seqs[0], &at(now))
3315 .unwrap());
3316 let back = lease(&s, now).unwrap();
3317 assert_eq!((back.seq, back.attempt), (seqs[0], 1));
3318 ack(&s, &back, now).unwrap();
3319 assert!(lease(&s, now).is_none());
3320 }
3321
3322 #[test]
3323 fn own_and_housekeeping_messages_settle_by_themselves() {
3324 let owner = Identity::generate("haris", Kind::Human);
3325 let s = Store::open_memory().unwrap();
3326 s.create_room(&room(&owner)).unwrap();
3327 let mut mine = msg(&owner, "from me");
3328 s.sequence_and_append(&mut mine).unwrap();
3329 let mut sys = msg(&owner, "joined");
3330 sys.kind = MessageType::System;
3331 s.sequence_and_append(&mut sys).unwrap();
3332 let mut theirs = msg_from(&owner, "bob", "for you");
3333 s.sequence_and_append(&mut theirs).unwrap();
3334 let d = lease(&s, 0).unwrap();
3335 assert_eq!(d.seq, theirs.seq);
3336 assert_eq!(s.bookmark("r_test", "haris").unwrap(), sys.seq);
3338 }
3339
3340 #[test]
3341 fn peek_shows_what_is_owed_without_taking_it() {
3342 let (s, seqs) = owed(3);
3343 let d = lease(&s, 0).unwrap();
3344 s.delivery_settle(
3345 d.token.as_deref().unwrap(),
3346 &Settle::Nack { retry_at: at(100) },
3347 &at(1),
3348 5,
3349 )
3350 .unwrap();
3351 let peeked: Vec<u64> = s
3352 .delivery_peek("r_test", "haris", &at(2), 10)
3353 .unwrap()
3354 .iter()
3355 .map(|m| m.seq)
3356 .collect();
3357 assert_eq!(peeked, vec![seqs[1], seqs[2]], "the delayed one is not due");
3358 let later: Vec<u64> = s
3359 .delivery_peek("r_test", "haris", &at(100), 10)
3360 .unwrap()
3361 .iter()
3362 .map(|m| m.seq)
3363 .collect();
3364 assert_eq!(later, seqs);
3365 assert_eq!(lease(&s, 100).unwrap().seq, seqs[0]);
3367 }
3368
3369 #[test]
3370 fn read_and_ack_settles_exactly_what_was_read() {
3371 let (s, seqs) = owed(3);
3372 s.delivery_ack_seqs("r_test", "haris", &[seqs[1]], &at(0))
3373 .unwrap();
3374 assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
3375 s.delivery_ack_seqs("r_test", "haris", &[seqs[0]], &at(0))
3376 .unwrap();
3377 assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[1]);
3378 assert_eq!(lease(&s, 0).unwrap().seq, seqs[2]);
3379 }
3380
3381 #[test]
3382 fn a_crash_after_the_lease_is_written_hands_it_out_again_later() {
3383 let owner = Identity::generate("haris", Kind::Human);
3384 let path = temp_db("lease");
3385 let seq = {
3386 let s = Store::open_with_key(&path, None).unwrap();
3387 s.create_room(&room(&owner)).unwrap();
3388 let mut m = msg_from(&owner, "bob", "work");
3389 s.sequence_and_append(&mut m).unwrap();
3390 lease(&s, 0).unwrap();
3392 m.seq
3393 };
3394 let s = Store::open_with_key(&path, None).unwrap();
3395 assert!(lease(&s, 10).is_none(), "the lease survived the restart");
3396 let again = lease(&s, 600).unwrap();
3397 assert_eq!((again.seq, again.attempt), (seq, 2));
3398 s.faults.arm("delivery.lease");
3400 assert!(s
3401 .delivery_lease("r_test", "haris", &at(1300), &at(1900), 5)
3402 .is_err());
3403 let d = s.deliveries("r_test", Some("haris")).unwrap();
3404 assert_eq!(d[0].attempt, 2);
3405 let _ = std::fs::remove_dir_all(path.parent().unwrap());
3406 }
3407
3408 #[test]
3409 fn a_spend_and_its_audit_event_land_together_or_not_at_all() {
3410 let owner = Identity::generate("haris", Kind::Human);
3411 let bot = Identity::generate("bot", Kind::Agent);
3412 let s = Store::open_memory().unwrap();
3413 s.create_room(&room(&owner)).unwrap();
3414 let member = |name: &str, id: &Identity, kind: Kind, role: Role| Member {
3415 room_id: "r_test".into(),
3416 name: name.into(),
3417 key: id.public(),
3418 kind,
3419 role,
3420 node: None,
3421 granted_by: owner.public(),
3422 expires_at: None,
3423 joined_at: T0.into(),
3424 last_seen: None,
3425 muted: false,
3426 revoked: false,
3427 profile: serde_json::Value::Null,
3428 };
3429 s.upsert_member(&member("haris", &owner, Kind::Human, Role::Approver), None)
3430 .unwrap();
3431 s.upsert_member(&member("bot", &bot, Kind::Agent, Role::TaskGiver), None)
3432 .unwrap();
3433 let action: crate::message::Action =
3434 serde_json::from_value(serde_json::json!({"verb":"deploy","target":"api","params":{}}))
3435 .unwrap();
3436 let mut q = Message::new(
3437 Draft {
3438 room: "r_test".into(),
3439 from: "bot".into(),
3440 text: "deploy?".into(),
3441 kind: Some(MessageType::Question),
3442 action: Some(action.clone()),
3443 ..Default::default()
3444 },
3445 &bot,
3446 )
3447 .unwrap();
3448 s.sequence_and_append(&mut q).unwrap();
3449 let mut a = Message::new(
3450 Draft {
3451 room: "r_test".into(),
3452 from: "haris".into(),
3453 text: "approved".into(),
3454 kind: Some(MessageType::Approve),
3455 reply_to: Some(q.id.clone()),
3456 action_hash: Some(action.hash()),
3457 expires: Some("2999-01-01T00:00:00Z".into()),
3458 once: Some(true),
3459 ..Default::default()
3460 },
3461 &owner,
3462 )
3463 .unwrap();
3464 s.sequence_and_append(&mut a).unwrap();
3465 let req = SpendRequest {
3466 room_id: "r_test".into(),
3467 approve_id: a.id.clone(),
3468 action_hash: action.hash(),
3469 op_id: "op-1".into(),
3470 spender: "bot".into(),
3471 node: "n1".into(),
3472 now: crate::message::now_ts(),
3473 };
3474 let audit = || msg(&owner, "spent");
3475 let before = s.message_count("r_test").unwrap();
3476
3477 s.faults.arm("spend.commit");
3479 assert!(s.spend_approve(&req, &mut audit()).is_err());
3480 assert!(s.spend_of(&a.id).unwrap().is_none());
3481 assert_eq!(s.message_count("r_test").unwrap(), before);
3482
3483 let SpendOutcome::Spent(rec) = s.spend_approve(&req, &mut audit()).unwrap() else {
3486 panic!("not spent")
3487 };
3488 assert_eq!(s.message_count("r_test").unwrap(), before + 1);
3489 assert_eq!(
3490 s.spend_approve(&req, &mut audit()).unwrap(),
3491 SpendOutcome::AlreadyYours(rec.clone())
3492 );
3493 let other = SpendRequest {
3494 op_id: "op-2".into(),
3495 ..req.clone()
3496 };
3497 assert_eq!(
3498 s.spend_approve(&other, &mut audit()).unwrap(),
3499 SpendOutcome::SpentByAnother
3500 );
3501 assert_eq!(s.message_count("r_test").unwrap(), before + 1);
3502 }
3503 #[test]
3504 fn files_are_indexed_kept_and_swept() {
3505 let owner = Identity::generate("haris", Kind::Human);
3506 let bot = Identity::generate("bot", Kind::Agent);
3507 let s = Store::open_memory().unwrap();
3508 s.create_room(&room(&owner)).unwrap();
3509 for (name, id) in [("haris", &owner), ("bot", &bot)] {
3510 s.upsert_member(
3511 &Member {
3512 room_id: "r_test".into(),
3513 name: name.into(),
3514 key: id.public(),
3515 kind: Kind::Agent,
3516 role: Role::TaskGiver,
3517 node: None,
3518 granted_by: owner.public(),
3519 expires_at: None,
3520 joined_at: T0.into(),
3521 last_seen: None,
3522 muted: false,
3523 revoked: false,
3524 profile: serde_json::Value::Null,
3525 },
3526 None,
3527 )
3528 .unwrap();
3529 }
3530 let hash = format!("sha256:{}", "a".repeat(64));
3531 let orphan = format!("sha256:{}", "b".repeat(64));
3532 let mut m = Message::new(
3533 Draft {
3534 room: "r_test".into(),
3535 from: "haris".into(),
3536 text: "log".into(),
3537 kind: Some(MessageType::Chat),
3538 data: serde_json::json!({"files": [{"id": hash, "name": "a.log", "size": 3, "type": "text/plain"}]}),
3539 ..Default::default()
3540 },
3541 &owner,
3542 )
3543 .unwrap();
3544 assert!(s.sequence_and_append(&mut m).unwrap());
3545 let refs = s.file_refs("r_test", &hash).unwrap();
3546 assert_eq!(refs.len(), 1);
3547 assert_eq!(refs[0].file.name, "a.log");
3548 assert_eq!(s.file_refs_of(&m.id).unwrap()[0].from, "haris");
3549
3550 let day0 = chrono::DateTime::parse_from_rfc3339(T0)
3551 .unwrap()
3552 .with_timezone(&chrono::Utc);
3553 let at = |d: i64| day0 + chrono::Duration::days(d);
3554 let ts = |d: i64| at(d).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
3555 for h in [&hash, &orphan] {
3556 s.file_add(&StoredFile {
3557 room_id: "r_test".into(),
3558 hash: h.clone(),
3559 size: 3,
3560 kind: "text".into(),
3561 uploader: "haris".into(),
3562 created: ts(0),
3563 })
3564 .unwrap();
3565 }
3566 assert_eq!(s.file_bytes("r_test", "haris", &ts(0)).unwrap(), (6, 6));
3567 assert_eq!(s.file_bytes("r_test", "bot", &ts(0)).unwrap(), (6, 0));
3568 assert_eq!(
3570 s.file_sweep("r_test", at(2), 7, 30).unwrap(),
3571 vec![orphan.clone()]
3572 );
3573 assert!(s.file_sweep("r_test", at(20), 7, 30).unwrap().is_empty());
3574 s.file_fetched("r_test", &hash, "bot", &ts(20)).unwrap();
3576 assert!(s.file_sweep("r_test", at(26), 7, 30).unwrap().is_empty());
3577 assert_eq!(
3578 s.file_sweep("r_test", at(28), 7, 30).unwrap(),
3579 vec![hash.clone()]
3580 );
3581 assert!(s.file_stored("r_test", &hash).unwrap().is_none());
3582
3583 s.file_add(&StoredFile {
3585 room_id: "r_test".into(),
3586 hash: hash.clone(),
3587 size: 3,
3588 kind: "text".into(),
3589 uploader: "haris".into(),
3590 created: ts(0),
3591 })
3592 .unwrap();
3593 assert_eq!(
3594 s.file_sweep("r_test", at(31), 7, 30).unwrap(),
3595 vec![hash.clone()]
3596 );
3597 assert!(s.tombstone_message(&m.id).unwrap());
3598 assert!(s.file_refs("r_test", &hash).unwrap().is_empty());
3599 }
3600}