use std::collections::HashMap;
use std::path::Path;
use std::sync::Mutex;
use chacha20poly1305::aead::{Aead, KeyInit};
use chacha20poly1305::{XChaCha20Poly1305, XNonce};
use rusqlite::{params, Connection, OptionalExtension, Row};
use rusqlite_migration::{Migrations, M};
use crate::control::ControlOp;
use crate::error::{Error, Result};
use crate::faults::Faults;
use crate::invite::Invite;
use crate::keys::{Kind, PublicKey};
use crate::limits::Limits;
use crate::message::{Content, DataClass, Message, MessageType, GENESIS_PREV};
use crate::room::{Member, Role, Room};
const MIGRATIONS: &[M<'static>] = &[
M::up(
r#"
CREATE TABLE rooms (
id TEXT PRIMARY KEY,
name TEXT NOT NULL UNIQUE,
about TEXT NOT NULL DEFAULT '',
owner TEXT NOT NULL,
created TEXT NOT NULL,
retention_days INTEGER,
class TEXT NOT NULL DEFAULT 'internal',
paused INTEGER NOT NULL DEFAULT 0,
hold INTEGER NOT NULL DEFAULT 0,
home_node TEXT NOT NULL,
home_hints TEXT NOT NULL DEFAULT 'null'
);
CREATE TABLE members (
room_id TEXT NOT NULL REFERENCES rooms(id) ON DELETE CASCADE,
name TEXT NOT NULL,
key TEXT NOT NULL,
kind TEXT NOT NULL,
role TEXT NOT NULL,
node TEXT,
granted_by TEXT NOT NULL,
expires_at TEXT,
joined_at TEXT NOT NULL,
last_seen TEXT,
muted INTEGER NOT NULL DEFAULT 0,
revoked INTEGER NOT NULL DEFAULT 0,
profile TEXT NOT NULL DEFAULT 'null',
identity TEXT,
PRIMARY KEY (room_id, name)
);
CREATE UNIQUE INDEX members_key ON members(room_id, key);
CREATE TABLE messages (
room_id TEXT NOT NULL REFERENCES rooms(id) ON DELETE CASCADE,
seq INTEGER NOT NULL,
id TEXT NOT NULL UNIQUE,
from_name TEXT NOT NULL,
type TEXT NOT NULL,
ts TEXT NOT NULL,
prev TEXT NOT NULL,
chain_hash TEXT NOT NULL,
envelope TEXT NOT NULL,
PRIMARY KEY (room_id, seq)
);
CREATE INDEX messages_room_from_ts ON messages(room_id, from_name, ts);
CREATE INDEX messages_room_ts ON messages(room_id, ts);
CREATE TABLE contents (
msg_id TEXT PRIMARY KEY REFERENCES messages(id) ON DELETE CASCADE,
body TEXT NOT NULL,
deleted INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE bookmarks (
room_id TEXT NOT NULL,
reader TEXT NOT NULL,
seq INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY (room_id, reader)
);
CREATE TABLE outbox (
msg_id TEXT PRIMARY KEY,
room_id TEXT NOT NULL,
message TEXT NOT NULL,
created TEXT NOT NULL,
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
CREATE TABLE invites (
nonce TEXT PRIMARY KEY,
room_id TEXT NOT NULL,
name TEXT NOT NULL,
created TEXT NOT NULL,
expires TEXT NOT NULL,
used_at TEXT,
used_by TEXT
);
"#,
),
M::up(
r#"
ALTER TABLE rooms ADD COLUMN closed INTEGER NOT NULL DEFAULT 0;
ALTER TABLE messages ADD COLUMN reply_to TEXT;
UPDATE messages SET reply_to = json_extract(envelope, '$.reply_to');
CREATE INDEX messages_reply_to ON messages(room_id, reply_to);
CREATE TABLE approvals_used (
msg_id TEXT PRIMARY KEY,
action_hash TEXT NOT NULL,
used_at TEXT NOT NULL
);
"#,
),
M::up(
r#"
ALTER TABLE messages ADD COLUMN received TEXT;
UPDATE messages SET received = ts;
CREATE INDEX messages_room_from_received ON messages(room_id, from_name, received);
CREATE INDEX messages_room_received ON messages(room_id, received);
"#,
),
M::up(
r#"
ALTER TABLE outbox ADD COLUMN sender TEXT NOT NULL DEFAULT '';
ALTER TABLE outbox ADD COLUMN state TEXT NOT NULL DEFAULT 'pending';
ALTER TABLE outbox ADD COLUMN retry_at TEXT;
ALTER TABLE outbox ADD COLUMN reason_class TEXT;
ALTER TABLE outbox ADD COLUMN reason_code INTEGER;
ALTER TABLE outbox ADD COLUMN unknown_code INTEGER;
ALTER TABLE outbox ADD COLUMN unknown_since TEXT;
ALTER TABLE outbox ADD COLUMN unknown_count INTEGER NOT NULL DEFAULT 0;
ALTER TABLE outbox ADD COLUMN updated TEXT;
ALTER TABLE outbox ADD COLUMN sig TEXT;
UPDATE outbox SET
sender = CASE WHEN json_valid(message)
THEN COALESCE(json_extract(message, '$.from'), '') ELSE '' END,
sig = CASE WHEN json_valid(message) THEN json_extract(message, '$.sig') END;
CREATE INDEX outbox_room_state ON outbox(room_id, state, sender, created, msg_id);
CREATE TABLE meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
"#,
),
M::up(
r#"
CREATE TABLE deliveries (
room_id TEXT NOT NULL,
reader TEXT NOT NULL,
seq INTEGER NOT NULL,
state TEXT NOT NULL,
token TEXT UNIQUE,
lease_until TEXT,
retry_at TEXT,
attempt INTEGER NOT NULL DEFAULT 0,
updated TEXT NOT NULL,
PRIMARY KEY (room_id, reader, seq)
);
"#,
),
M::up(
r#"
CREATE TABLE spends (
approve_id TEXT PRIMARY KEY,
room_id TEXT NOT NULL,
action_hash TEXT NOT NULL,
op_id TEXT NOT NULL,
spender TEXT NOT NULL,
node TEXT NOT NULL,
at TEXT NOT NULL,
audit_seq INTEGER NOT NULL
);
CREATE TABLE pending_spends (
room_id TEXT NOT NULL,
action_hash TEXT NOT NULL,
spender TEXT NOT NULL,
op_id TEXT NOT NULL,
created TEXT NOT NULL,
PRIMARY KEY (room_id, action_hash, spender)
);
"#,
),
M::up(
r#"
CREATE TABLE file_refs (
msg_id TEXT NOT NULL REFERENCES messages(id) ON DELETE CASCADE,
room_id TEXT NOT NULL,
hash TEXT NOT NULL,
name TEXT NOT NULL,
size INTEGER NOT NULL,
type TEXT NOT NULL DEFAULT '',
from_name TEXT NOT NULL,
to_name TEXT,
PRIMARY KEY (msg_id, hash)
);
CREATE INDEX file_refs_hash ON file_refs(room_id, hash);
CREATE TABLE files (
room_id TEXT NOT NULL,
hash TEXT NOT NULL,
size INTEGER NOT NULL,
kind TEXT NOT NULL,
uploader TEXT NOT NULL,
created TEXT NOT NULL,
PRIMARY KEY (room_id, hash)
);
CREATE INDEX files_uploader ON files(room_id, uploader, created);
CREATE TABLE file_fetches (
room_id TEXT NOT NULL,
hash TEXT NOT NULL,
member TEXT NOT NULL,
at TEXT NOT NULL,
PRIMARY KEY (room_id, hash, member)
);
"#,
),
];
const SEAL_BATCH: usize = 100;
pub struct Store {
conn: Mutex<Connection>,
cipher: Option<XChaCha20Poly1305>,
pub faults: Faults,
}
const ENC_PREFIX: &str = "enc1:";
impl std::fmt::Debug for Store {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("Store")
}
}
#[derive(Debug, Clone)]
pub struct LocalMember {
pub member: Member,
pub identity: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum OutboxState {
Pending,
Waiting,
Failed,
Quarantined,
Dropped,
}
impl OutboxState {
pub fn as_str(&self) -> &'static str {
match self {
OutboxState::Pending => "pending",
OutboxState::Waiting => "waiting",
OutboxState::Failed => "failed",
OutboxState::Quarantined => "quarantined",
OutboxState::Dropped => "dropped",
}
}
}
impl std::str::FromStr for OutboxState {
type Err = Error;
fn from_str(s: &str) -> Result<Self> {
Ok(match s {
"pending" => OutboxState::Pending,
"waiting" => OutboxState::Waiting,
"failed" => OutboxState::Failed,
"quarantined" => OutboxState::Quarantined,
"dropped" => OutboxState::Dropped,
other => return Err(Error::Invalid(format!("unknown outbox state {other}"))),
})
}
}
#[derive(Debug, Clone)]
pub struct OutboxEntry {
pub msg_id: String,
pub room_id: String,
pub sender: String,
pub state: OutboxState,
pub created: String,
pub attempts: u32,
pub retry_at: Option<String>,
pub reason: Option<String>,
pub reason_class: Option<String>,
pub reason_code: Option<i32>,
pub updated: Option<String>,
pub message: Option<Message>,
pub unreadable: bool,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct OutboxCounts {
pub pending: u64,
pub waiting: u64,
pub failed: u64,
pub quarantined: u64,
pub dropped: u64,
}
impl OutboxCounts {
pub fn queued(&self) -> u64 {
self.pending + self.waiting
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum DeliveryState {
Leased,
Acked,
Delayed,
Quarantined,
Replay,
}
impl DeliveryState {
pub fn as_str(&self) -> &'static str {
match self {
DeliveryState::Leased => "leased",
DeliveryState::Acked => "acked",
DeliveryState::Delayed => "delayed",
DeliveryState::Quarantined => "quarantined",
DeliveryState::Replay => "replay",
}
}
fn parse(s: &str) -> DeliveryState {
match s {
"acked" => DeliveryState::Acked,
"delayed" => DeliveryState::Delayed,
"quarantined" => DeliveryState::Quarantined,
"replay" => DeliveryState::Replay,
_ => DeliveryState::Leased,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct WakePending {
pub count: u64,
pub newest_seq: u64,
pub newest_id: String,
pub looped: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Delivery {
pub room_id: String,
pub reader: String,
pub seq: u64,
pub state: DeliveryState,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub token: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub lease_until: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retry_at: Option<String>,
pub attempt: u32,
pub updated: String,
}
#[derive(Debug, Clone)]
pub enum Settle {
Ack,
Renew { lease_until: String },
Nack { retry_at: String },
}
#[derive(Debug, Clone)]
pub struct SpendRequest {
pub room_id: String,
pub approve_id: String,
pub action_hash: String,
pub op_id: String,
pub spender: String,
pub node: String,
pub now: String,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct SpendRecord {
pub approve_id: String,
pub room_id: String,
pub action_hash: String,
pub op_id: String,
pub spender: String,
pub node: String,
pub at: String,
pub audit_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SpendOutcome {
Spent(SpendRecord),
AlreadyYours(SpendRecord),
SpentByAnother,
}
type Candidate = (
i64,
String,
String,
Option<String>,
Option<String>,
Option<String>,
);
fn wanted(from: &str, kind: &str, reader: &str) -> bool {
from != reader && kind != "system" && kind != "control"
}
fn seconds_between(a: &str, b: &str) -> i64 {
match (
chrono::DateTime::parse_from_rfc3339(a),
chrono::DateTime::parse_from_rfc3339(b),
) {
(Ok(a), Ok(b)) => (b - a).num_seconds(),
_ => 0,
}
}
impl Store {
pub fn open(path: &Path) -> Result<Self> {
Self::open_with_key(path, None)
}
pub fn open_with_key(path: &Path, key: Option<[u8; 32]>) -> Result<Self> {
let store = Self::open_unsealed(path, key)?;
store.seal_legacy_outbox(SEAL_BATCH)?;
Ok(store)
}
#[doc(hidden)]
pub fn open_unsealed(path: &Path, key: Option<[u8; 32]>) -> Result<Self> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let conn = Connection::open(path)?;
Self::init(conn, key)
}
pub fn open_memory() -> Result<Self> {
Self::init(Connection::open_in_memory()?, None)
}
pub fn open_memory_with_key(key: [u8; 32]) -> Result<Self> {
Self::init(Connection::open_in_memory()?, Some(key))
}
fn init(mut conn: Connection, key: Option<[u8; 32]>) -> Result<Self> {
conn.pragma_update(None, "journal_mode", "WAL")?;
conn.pragma_update(None, "synchronous", "FULL")?;
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.pragma_update(None, "secure_delete", "ON")?;
Migrations::from_slice(MIGRATIONS).to_latest(&mut conn)?;
Ok(Store {
conn: Mutex::new(conn),
cipher: key.map(|k| XChaCha20Poly1305::new((&k).into())),
faults: Faults::default(),
})
}
#[doc(hidden)]
pub fn cap_size_for_test(&self, extra: u32) -> Result<()> {
let conn = self.lock();
let pages: i64 = conn.query_row("PRAGMA page_count", [], |r| r.get(0))?;
conn.pragma_update(None, "max_page_count", pages + extra as i64)?;
Ok(())
}
pub fn checkpoint(&self) -> Result<()> {
let conn = self.lock();
conn.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |_| Ok(()))?;
Ok(())
}
pub fn seal_legacy_outbox(&self, batch: usize) -> Result<u64> {
if self.cipher.is_none() {
return Ok(0);
}
let mut sealed = 0u64;
loop {
let mut conn = self.lock();
let tx = conn.transaction()?;
let rows: Vec<(String, String)> = {
let mut stmt = tx.prepare(
"SELECT msg_id, message FROM outbox
WHERE message != '' AND message NOT LIKE 'enc1:%' LIMIT ?1",
)?;
let rows = stmt.query_map(params![batch as i64], |r| Ok((r.get(0)?, r.get(1)?)))?;
rows.collect::<rusqlite::Result<Vec<_>>>()?
};
if rows.is_empty() {
tx.execute(
"INSERT INTO meta (key, value) VALUES ('outbox_sealed', ?1)
ON CONFLICT(key) DO UPDATE SET value=excluded.value",
params![crate::message::now_ts()],
)?;
tx.commit()?;
return Ok(sealed);
}
for (id, plain) in rows {
self.faults.check("migration.row")?;
tx.execute(
"UPDATE outbox SET message=?2 WHERE msg_id=?1",
params![id, self.seal(&plain)?],
)?;
sealed += 1;
}
tx.commit()?;
}
}
pub fn is_encrypted(&self) -> bool {
self.cipher.is_some()
}
fn seal(&self, body: &str) -> Result<String> {
match &self.cipher {
None => Ok(body.to_string()),
Some(c) => {
let nonce_bytes: [u8; 24] = rand::random();
let nonce = XNonce::from(nonce_bytes);
let ct = c
.encrypt(&nonce, body.as_bytes())
.map_err(|_| Error::Other("encrypt failed".into()))?;
let mut out = nonce_bytes.to_vec();
out.extend_from_slice(&ct);
Ok(format!(
"{ENC_PREFIX}{}",
data_encoding::BASE64.encode(&out)
))
}
}
}
fn unseal(&self, stored: &str) -> Option<String> {
let Some(b64) = stored.strip_prefix(ENC_PREFIX) else {
return Some(stored.to_string());
};
let c = self.cipher.as_ref()?;
let raw = data_encoding::BASE64.decode(b64.as_bytes()).ok()?;
if raw.len() < 24 {
return None;
}
let (n, ct) = raw.split_at(24);
let nonce = XNonce::from(<[u8; 24]>::try_from(n).ok()?);
let pt = c.decrypt(&nonce, ct).ok()?;
String::from_utf8(pt).ok()
}
fn lock(&self) -> std::sync::MutexGuard<'_, Connection> {
self.conn.lock().unwrap_or_else(|p| p.into_inner())
}
pub fn create_room(&self, room: &Room) -> Result<()> {
let conn = self.lock();
conn.execute(
"INSERT INTO rooms (id, name, about, owner, created, retention_days, class, paused, hold, home_node, home_hints)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
params![
room.id,
room.name,
room.about,
room.owner.to_string(),
room.created,
room.retention_days,
room.class.as_str(),
room.paused as i32,
room.hold as i32,
room.home_node,
serde_json::to_string(&room.home_hints)?,
],
)
.map_err(|e| match e {
rusqlite::Error::SqliteFailure(f, _) if f.code == rusqlite::ErrorCode::ConstraintViolation => {
Error::NameTaken(format!("a room named {} already exists here", room.name))
}
other => Error::Db(other),
})?;
Ok(())
}
pub fn update_room(&self, room: &Room) -> Result<()> {
let conn = self.lock();
conn.execute(
"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",
params![
room.id,
room.about,
room.retention_days,
room.class.as_str(),
room.paused as i32,
room.hold as i32,
room.home_node,
serde_json::to_string(&room.home_hints)?,
room.closed as i32,
],
)?;
Ok(())
}
fn row_to_room(row: &Row<'_>) -> rusqlite::Result<Room> {
let owner: String = row.get("owner")?;
let class: String = row.get("class")?;
let hints: String = row.get("home_hints")?;
Ok(Room {
id: row.get("id")?,
name: row.get("name")?,
about: row.get("about")?,
owner: owner.parse().map_err(|_| rusqlite::Error::InvalidQuery)?,
created: row.get("created")?,
retention_days: row.get("retention_days")?,
class: class.parse().unwrap_or(DataClass::Internal),
paused: row.get::<_, i32>("paused")? != 0,
hold: row.get::<_, i32>("hold")? != 0,
closed: row.get::<_, i32>("closed").unwrap_or(0) != 0,
home_node: row.get("home_node")?,
home_hints: serde_json::from_str(&hints).unwrap_or(serde_json::Value::Null),
})
}
pub fn room_by_name(&self, name: &str) -> Result<Option<Room>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM rooms WHERE name=?1",
params![name],
Self::row_to_room,
)
.optional()?)
}
pub fn room_by_id(&self, id: &str) -> Result<Option<Room>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM rooms WHERE id=?1",
params![id],
Self::row_to_room,
)
.optional()?)
}
pub fn room(&self, name_or_id: &str) -> Result<Room> {
if let Some(r) = self.room_by_name(name_or_id)? {
return Ok(r);
}
if let Some(r) = self.room_by_id(name_or_id)? {
return Ok(r);
}
Err(Error::NotInRoom(name_or_id.to_string()))
}
pub fn list_rooms(&self) -> Result<Vec<Room>> {
let conn = self.lock();
let mut stmt = conn.prepare("SELECT * FROM rooms ORDER BY created")?;
let rows = stmt.query_map([], Self::row_to_room)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn upsert_member(&self, m: &Member, identity: Option<&str>) -> Result<()> {
let conn = self.lock();
conn.execute(
"INSERT INTO members (room_id, name, key, kind, role, node, granted_by, expires_at, joined_at, last_seen, muted, revoked, profile, identity)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
ON CONFLICT(room_id, name) DO UPDATE SET
key=excluded.key, kind=excluded.kind, role=excluded.role, node=excluded.node,
granted_by=excluded.granted_by, expires_at=excluded.expires_at,
last_seen=COALESCE(excluded.last_seen, members.last_seen),
muted=excluded.muted, revoked=excluded.revoked, profile=excluded.profile,
identity=COALESCE(excluded.identity, members.identity)",
params![
m.room_id,
m.name,
m.key.to_string(),
m.kind.to_string(),
m.role.as_str(),
m.node,
m.granted_by.to_string(),
m.expires_at,
m.joined_at,
m.last_seen,
m.muted as i32,
m.revoked as i32,
serde_json::to_string(&m.profile)?,
identity,
],
)?;
Ok(())
}
fn row_to_member(row: &Row<'_>) -> rusqlite::Result<Member> {
let key: String = row.get("key")?;
let kind: String = row.get("kind")?;
let role: String = row.get("role")?;
let granted_by: String = row.get("granted_by")?;
let profile: String = row.get("profile")?;
Ok(Member {
room_id: row.get("room_id")?,
name: row.get("name")?,
key: key.parse().map_err(|_| rusqlite::Error::InvalidQuery)?,
kind: kind.parse().unwrap_or(Kind::Agent),
role: role.parse().unwrap_or(Role::Observer),
node: row.get("node")?,
granted_by: granted_by
.parse()
.map_err(|_| rusqlite::Error::InvalidQuery)?,
expires_at: row.get("expires_at")?,
joined_at: row.get("joined_at")?,
last_seen: row.get("last_seen")?,
muted: row.get::<_, i32>("muted")? != 0,
revoked: row.get::<_, i32>("revoked")? != 0,
profile: serde_json::from_str(&profile).unwrap_or(serde_json::Value::Null),
})
}
pub fn member_by_name(&self, room_id: &str, name: &str) -> Result<Option<Member>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM members WHERE room_id=?1 AND name=?2",
params![room_id, name],
Self::row_to_member,
)
.optional()?)
}
pub fn member_by_key(&self, room_id: &str, key: &PublicKey) -> Result<Option<Member>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM members WHERE room_id=?1 AND key=?2",
params![room_id, key.to_string()],
Self::row_to_member,
)
.optional()?)
}
pub fn members(&self, room_id: &str) -> Result<Vec<Member>> {
let conn = self.lock();
let mut stmt =
conn.prepare("SELECT * FROM members WHERE room_id=?1 ORDER BY joined_at, name")?;
let rows = stmt.query_map(params![room_id], Self::row_to_member)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn local_members(&self, room_id: &str) -> Result<Vec<LocalMember>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT * FROM members WHERE room_id=?1 AND identity IS NOT NULL ORDER BY joined_at",
)?;
let rows = stmt.query_map(params![room_id], |row| {
Ok(LocalMember {
member: Self::row_to_member(row)?,
identity: row.get("identity")?,
})
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn local_member(&self, room_id: &str, identity: &str) -> Result<Option<LocalMember>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM members WHERE room_id=?1 AND identity=?2",
params![room_id, identity],
|row| {
Ok(LocalMember {
member: Self::row_to_member(row)?,
identity: row.get("identity")?,
})
},
)
.optional()?)
}
pub fn set_member_role(
&self,
room_id: &str,
name: &str,
role: Role,
expires_at: Option<&str>,
) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE members SET role=?3, expires_at=?4 WHERE room_id=?1 AND name=?2",
params![room_id, name, role.as_str(), expires_at],
)?;
Ok(n == 1)
}
pub fn set_member_muted(&self, room_id: &str, name: &str, muted: bool) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE members SET muted=?3 WHERE room_id=?1 AND name=?2",
params![room_id, name, muted as i32],
)?;
Ok(n == 1)
}
pub fn set_member_revoked(&self, room_id: &str, name: &str) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE members SET revoked=1, node=NULL WHERE room_id=?1 AND name=?2",
params![room_id, name],
)?;
Ok(n == 1)
}
pub fn clear_member_node(&self, room_id: &str, name: &str) -> Result<()> {
let conn = self.lock();
conn.execute(
"UPDATE members SET node=NULL WHERE room_id=?1 AND name=?2",
params![room_id, name],
)?;
Ok(())
}
pub fn rename_room(&self, room_id: &str, new_name: &str) -> Result<()> {
let conn = self.lock();
conn.execute(
"UPDATE rooms SET name=?2 WHERE id=?1",
params![room_id, new_name],
)?;
Ok(())
}
pub fn touch_member(
&self,
room_id: &str,
name: &str,
node: Option<&str>,
now: &str,
) -> Result<()> {
let conn = self.lock();
conn.execute(
"UPDATE members SET last_seen=?3, node=COALESCE(?4, node) WHERE room_id=?1 AND name=?2",
params![room_id, name, now, node],
)?;
Ok(())
}
pub fn chain_head(&self, room_id: &str) -> Result<(u64, String)> {
let conn = self.lock();
Self::chain_head_in(&conn, room_id)
}
fn chain_head_in(conn: &Connection, room_id: &str) -> Result<(u64, String)> {
let head: Option<(i64, String)> = conn
.query_row(
"SELECT seq, chain_hash FROM messages WHERE room_id=?1 ORDER BY seq DESC LIMIT 1",
params![room_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
Ok(match head {
Some((seq, hash)) => (seq as u64, hash),
None => (0, GENESIS_PREV.to_string()),
})
}
pub fn append(&self, msg: &Message) -> Result<bool> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let written = self.append_in(&tx, msg)?;
tx.commit()?;
Ok(written)
}
fn append_in(&self, conn: &Connection, msg: &Message) -> Result<bool> {
if !msg.is_sequenced() {
return Err(Error::Invalid("message has no seq/prev yet".into()));
}
let exists: Option<i64> = conn
.query_row(
"SELECT seq FROM messages WHERE id=?1",
params![msg.id],
|r| r.get(0),
)
.optional()?;
if exists.is_some() {
Self::unqueue_in(conn, msg)?;
return Ok(false);
}
let (last_seq, last_hash) = Self::chain_head_in(conn, &msg.room)?;
if msg.seq != last_seq + 1 || msg.prev != last_hash {
return Err(Error::Invalid(format!(
"chain break in room {}: got seq {} prev {}, expected seq {} prev {}",
msg.room,
msg.seq,
&msg.prev[..msg.prev.len().min(16)],
last_seq + 1,
&last_hash[..16]
)));
}
conn.execute(
"INSERT INTO messages (room_id, seq, id, from_name, type, ts, prev, chain_hash, envelope, reply_to, received)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
params![
msg.room,
msg.seq as i64,
msg.id,
msg.from,
msg.kind.as_str(),
msg.ts,
msg.prev,
msg.chain_hash(),
serde_json::to_string(&msg.envelope())?,
msg.reply_to,
crate::message::now_ts(),
],
)?;
if msg.tombstone {
conn.execute(
"INSERT INTO contents (msg_id, body, deleted) VALUES (?1, '', 1)",
params![msg.id],
)?;
} else {
conn.execute(
"INSERT INTO contents (msg_id, body, deleted) VALUES (?1, ?2, 0)",
params![msg.id, self.seal(&serde_json::to_string(&msg.content())?)?],
)?;
for f in crate::files::refs(&msg.data).unwrap_or_default() {
conn.execute(
"INSERT OR IGNORE INTO file_refs (msg_id, room_id, hash, name, size, type, from_name, to_name)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![msg.id, msg.room, f.id, f.name, f.size as i64, f.mime, msg.from, msg.to],
)?;
}
}
Self::unqueue_in(conn, msg)?;
Ok(true)
}
fn unqueue_in(conn: &Connection, msg: &Message) -> Result<()> {
conn.execute(
"DELETE FROM outbox WHERE msg_id=?1 AND (sig IS NULL OR sig=?2)",
params![msg.id, msg.sig],
)?;
Ok(())
}
pub fn sequence_and_append(&self, msg: &mut Message) -> Result<bool> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let exists: Option<i64> = tx
.query_row(
"SELECT seq FROM messages WHERE id=?1",
params![msg.id],
|r| r.get(0),
)
.optional()?;
if exists.is_some() {
return Ok(false);
}
let (last_seq, last_hash) = Self::chain_head_in(&tx, &msg.room)?;
msg.sequence(last_seq + 1, &last_hash);
self.append_in(&tx, msg)?;
tx.commit()?;
Ok(true)
}
pub fn sequence_and_append_control(&self, msg: &mut Message, op: &ControlOp) -> Result<bool> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let exists: Option<i64> = tx
.query_row(
"SELECT seq FROM messages WHERE id=?1",
params![msg.id],
|r| r.get(0),
)
.optional()?;
if exists.is_some() {
return Ok(false);
}
let (last_seq, last_hash) = Self::chain_head_in(&tx, &msg.room)?;
msg.sequence(last_seq + 1, &last_hash);
self.append_in(&tx, msg)?;
Self::apply_control_in(&tx, &msg.room, op)?;
tx.commit()?;
Ok(true)
}
fn apply_control_in(conn: &Connection, room_id: &str, op: &ControlOp) -> Result<()> {
match op {
ControlOp::Grant { name, role, until } => {
conn.execute(
"UPDATE members SET role=?3, expires_at=?4 WHERE room_id=?1 AND name=?2",
params![room_id, name, role.as_str(), until],
)?;
}
ControlOp::Pause | ControlOp::Resume => {
conn.execute(
"UPDATE rooms SET paused=?2 WHERE id=?1",
params![room_id, matches!(op, ControlOp::Pause) as i32],
)?;
}
ControlOp::Mute { name } | ControlOp::Unmute { name } => {
conn.execute(
"UPDATE members SET muted=?3 WHERE room_id=?1 AND name=?2",
params![room_id, name, matches!(op, ControlOp::Mute { .. }) as i32],
)?;
}
ControlOp::Revoke { name } => {
conn.execute(
"UPDATE members SET revoked=1, node=NULL WHERE room_id=?1 AND name=?2",
params![room_id, name],
)?;
}
ControlOp::Hold { on } => {
conn.execute(
"UPDATE rooms SET hold=?2 WHERE id=?1",
params![room_id, *on as i32],
)?;
}
ControlOp::Rotated { .. } => {}
}
Ok(())
}
fn row_to_spend(row: &Row<'_>) -> rusqlite::Result<SpendRecord> {
Ok(SpendRecord {
approve_id: row.get("approve_id")?,
room_id: row.get("room_id")?,
action_hash: row.get("action_hash")?,
op_id: row.get("op_id")?,
spender: row.get("spender")?,
node: row.get("node")?,
at: row.get("at")?,
audit_seq: row.get::<_, i64>("audit_seq")? as u64,
})
}
pub fn spend_approve(&self, req: &SpendRequest, audit: &mut Message) -> Result<SpendOutcome> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let prior = tx
.query_row(
"SELECT * FROM spends WHERE approve_id=?1",
params![req.approve_id],
Self::row_to_spend,
)
.optional()?;
if let Some(rec) = prior {
if rec.op_id == req.op_id && rec.spender == req.spender && rec.room_id == req.room_id {
return Ok(SpendOutcome::AlreadyYours(rec));
}
return Ok(SpendOutcome::SpentByAnother);
}
let legacy: Option<String> = tx
.query_row(
"SELECT msg_id FROM approvals_used WHERE msg_id=?1",
params![req.approve_id],
|r| r.get(0),
)
.optional()?;
if legacy.is_some() {
return Ok(SpendOutcome::SpentByAnother);
}
let room = tx
.query_row(
"SELECT * FROM rooms WHERE id=?1",
params![req.room_id],
Self::row_to_room,
)
.optional()?
.ok_or_else(|| Error::NotInRoom(req.room_id.clone()))?;
if room.closed {
return Err(Error::Denied(format!("room {} was rotated", room.name)));
}
if room.paused {
return Err(Error::RoomPaused(room.name));
}
let member = |name: &str| -> Result<Option<Member>> {
Ok(tx
.query_row(
"SELECT * FROM members WHERE room_id=?1 AND name=?2",
params![req.room_id, name],
Self::row_to_member,
)
.optional()?)
};
let spender = member(&req.spender)?
.ok_or_else(|| Error::Denied(format!("{} is not a member", req.spender)))?;
if spender.revoked || spender.muted || spender.is_expired(&req.now) {
return Err(Error::Denied(format!(
"{} may not act on approves: revoked, muted or expired",
spender.name
)));
}
if spender.role == Role::Observer {
return Err(Error::Denied(format!(
"{} is an observer and may not act on approves",
spender.name
)));
}
let approve = tx
.query_row(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id WHERE m.id=?1 AND m.room_id=?2",
params![req.approve_id, req.room_id],
|r| self.row_to_message(r),
)
.optional()?
.ok_or_else(|| Error::Denied(format!("no approve {} in this room", req.approve_id)))?;
let invalid = |why: &str| Err(Error::Denied(format!("approve {}: {why}", approve.id)));
if approve.kind != MessageType::Approve {
return invalid("is not an approve");
}
if approve.action_hash.as_deref() != Some(req.action_hash.as_str()) {
return invalid("is for a different action");
}
match &approve.expires {
Some(exp) if exp.as_str() > req.now.as_str() => {}
_ => return invalid("has expired"),
}
if approve.ts.as_str() > req.now.as_str() {
return invalid("is dated in the future");
}
let approver = member(&approve.from)?
.ok_or_else(|| Error::Denied(format!("approver {} is not a member", approve.from)))?;
if approver.kind != Kind::Human {
return invalid("was not signed by a human key");
}
approver
.check_may_send(MessageType::Approve, approver.key == room.owner, &req.now)
.map_err(|e| Error::Denied(format!("approve {}: {e}", approve.id)))?;
approve.verify(&approver.key)?;
let (last_seq, last_hash) = Self::chain_head_in(&tx, &req.room_id)?;
audit.sequence(last_seq + 1, &last_hash);
self.append_in(&tx, audit)?;
tx.execute(
"INSERT INTO spends (approve_id, room_id, action_hash, op_id, spender, node, at, audit_seq)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
req.approve_id,
req.room_id,
req.action_hash,
req.op_id,
req.spender,
req.node,
req.now,
audit.seq as i64
],
)?;
tx.execute(
"INSERT OR IGNORE INTO approvals_used (msg_id, action_hash, used_at) VALUES (?1, ?2, ?3)",
params![req.approve_id, req.action_hash, req.now],
)?;
self.faults.check("spend.commit")?;
tx.commit()?;
Ok(SpendOutcome::Spent(SpendRecord {
approve_id: req.approve_id.clone(),
room_id: req.room_id.clone(),
action_hash: req.action_hash.clone(),
op_id: req.op_id.clone(),
spender: req.spender.clone(),
node: req.node.clone(),
at: req.now.clone(),
audit_seq: audit.seq,
}))
}
pub fn spend_of(&self, approve_id: &str) -> Result<Option<SpendRecord>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM spends WHERE approve_id=?1",
params![approve_id],
Self::row_to_spend,
)
.optional()?)
}
pub fn pending_spend(
&self,
room_id: &str,
action_hash: &str,
spender: &str,
fresh: &str,
now: &str,
) -> Result<String> {
let conn = self.lock();
conn.execute(
"INSERT OR IGNORE INTO pending_spends (room_id, action_hash, spender, op_id, created)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![room_id, action_hash, spender, fresh, now],
)?;
Ok(conn.query_row(
"SELECT op_id FROM pending_spends WHERE room_id=?1 AND action_hash=?2 AND spender=?3",
params![room_id, action_hash, spender],
|r| r.get(0),
)?)
}
pub fn pending_spend_done(
&self,
room_id: &str,
action_hash: &str,
spender: &str,
) -> Result<()> {
let conn = self.lock();
conn.execute(
"DELETE FROM pending_spends WHERE room_id=?1 AND action_hash=?2 AND spender=?3",
params![room_id, action_hash, spender],
)?;
Ok(())
}
fn row_to_message(&self, row: &Row<'_>) -> rusqlite::Result<Message> {
let envelope: String = row.get("envelope")?;
let body: Option<String> = row.get("body")?;
let deleted: Option<i32> = row.get("deleted")?;
let mut v: serde_json::Value =
serde_json::from_str(&envelope).map_err(|_| rusqlite::Error::InvalidQuery)?;
let (content, tombstone) = match (body, deleted) {
(Some(b), Some(0)) => match self.unseal(&b) {
Some(plain) => (serde_json::from_str(&plain).unwrap_or_default(), false),
None => (Content::default(), true),
},
_ => (Content::default(), true),
};
if let Some(obj) = v.as_object_mut() {
if tombstone {
obj.insert("tombstone".into(), serde_json::Value::Bool(true));
} else {
obj.remove("content_hash");
}
obj.insert("text".into(), serde_json::Value::String(content.text));
obj.insert(
"action".into(),
serde_json::to_value(content.action).unwrap_or(serde_json::Value::Null),
);
obj.insert("data".into(), content.data);
}
serde_json::from_value(v).map_err(|_| rusqlite::Error::InvalidQuery)
}
pub fn messages_after(&self, room_id: &str, after: u64, limit: u32) -> Result<Vec<Message>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
WHERE m.room_id=?1 AND m.seq>?2 ORDER BY m.seq LIMIT ?3",
)?;
let rows = stmt.query_map(params![room_id, after as i64, limit], |r| {
self.row_to_message(r)
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn messages_from_ts(&self, room_id: &str, since: &str, limit: u32) -> Result<Vec<Message>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
WHERE m.room_id=?1 AND m.ts>=?2 ORDER BY m.seq LIMIT ?3",
)?;
let rows = stmt.query_map(params![room_id, since, limit], |r| self.row_to_message(r))?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn message_by_id(&self, id: &str) -> Result<Option<Message>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id WHERE m.id=?1",
params![id],
|r| self.row_to_message(r),
)
.optional()?)
}
pub fn replies_to(&self, room_id: &str, msg_id: &str) -> Result<Vec<Message>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
WHERE m.room_id=?1 AND m.reply_to=?2 ORDER BY m.seq",
)?;
let rows = stmt.query_map(params![room_id, msg_id], |r| self.row_to_message(r))?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn messages_with_trace(&self, room_id: &str, trace: &str) -> Result<Vec<Message>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
WHERE m.room_id=?1 AND json_extract(m.envelope, '$.trace')=?2
ORDER BY m.seq LIMIT 5000",
)?;
let rows = stmt.query_map(params![room_id, trace], |r| self.row_to_message(r))?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn claim_holder(&self, room_id: &str, task_id: &str) -> Result<Option<String>> {
let mut holder: Option<String> = None;
for m in self.replies_to(room_id, task_id)? {
match m.kind {
MessageType::Claim if holder.is_none() => holder = Some(m.from.clone()),
MessageType::Release if holder.as_deref() == Some(m.from.as_str()) => holder = None,
_ => {}
}
}
Ok(holder)
}
pub fn approvals_for(&self, room_id: &str, action_hash: &str) -> Result<Vec<Message>> {
self.approvals_for_action(room_id, action_hash, false)
}
pub fn approvals_for_action(
&self,
room_id: &str,
action_hash: &str,
with_spent: bool,
) -> Result<Vec<Message>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted,
m.id IN (SELECT msg_id FROM approvals_used) AS spent
FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
WHERE m.room_id=?1 AND m.type='approve'
AND json_extract(m.envelope, '$.action_hash')=?2
AND (?3 OR m.id NOT IN (SELECT msg_id FROM approvals_used))
ORDER BY spent, m.seq",
)?;
let rows = stmt.query_map(params![room_id, action_hash, with_spent], |r| {
self.row_to_message(r)
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn approval_use(&self, msg_id: &str, action_hash: &str, now: &str) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"INSERT OR IGNORE INTO approvals_used (msg_id, action_hash, used_at) VALUES (?1, ?2, ?3)",
params![msg_id, action_hash, now],
)?;
Ok(n == 1)
}
pub fn tombstone_before(&self, room_id: &str, cutoff: &str) -> Result<u64> {
let conn = self.lock();
let n = conn.execute(
"UPDATE contents SET body='', deleted=1
WHERE deleted=0 AND msg_id IN (SELECT id FROM messages WHERE room_id=?1 AND ts<?2)",
params![room_id, cutoff],
)?;
conn.execute(
"DELETE FROM file_refs WHERE msg_id IN (SELECT id FROM messages WHERE room_id=?1 AND ts<?2)",
params![room_id, cutoff],
)?;
Ok(n as u64)
}
pub fn tombstone_message(&self, msg_id: &str) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE contents SET body='', deleted=1 WHERE msg_id=?1 AND deleted=0",
params![msg_id],
)?;
conn.execute("DELETE FROM file_refs WHERE msg_id=?1", params![msg_id])?;
Ok(n == 1)
}
pub fn message_count(&self, room_id: &str) -> Result<u64> {
let conn = self.lock();
let n: i64 = conn.query_row(
"SELECT COUNT(*) FROM messages WHERE room_id=?1",
params![room_id],
|r| r.get(0),
)?;
Ok(n as u64)
}
pub fn bookmark(&self, room_id: &str, reader: &str) -> Result<u64> {
let conn = self.lock();
let seq: Option<i64> = conn
.query_row(
"SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2",
params![room_id, reader],
|r| r.get(0),
)
.optional()?;
Ok(seq.unwrap_or(0) as u64)
}
pub fn set_bookmark(&self, room_id: &str, reader: &str, seq: u64) -> Result<()> {
let conn = self.lock();
conn.execute(
"INSERT INTO bookmarks (room_id, reader, seq) VALUES (?1, ?2, ?3)
ON CONFLICT(room_id, reader) DO UPDATE SET seq=MAX(bookmarks.seq, excluded.seq)",
params![room_id, reader, seq as i64],
)?;
Ok(())
}
fn row_to_delivery(row: &Row<'_>) -> rusqlite::Result<Delivery> {
let state: String = row.get("state")?;
Ok(Delivery {
room_id: row.get("room_id")?,
reader: row.get("reader")?,
seq: row.get::<_, i64>("seq")? as u64,
state: DeliveryState::parse(&state),
token: row.get("token")?,
lease_until: row.get("lease_until")?,
retry_at: row.get("retry_at")?,
attempt: row.get::<_, i64>("attempt")? as u32,
updated: row.get("updated")?,
})
}
fn bookmark_in(conn: &Connection, room_id: &str, reader: &str) -> Result<u64> {
let seq: Option<i64> = conn
.query_row(
"SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2",
params![room_id, reader],
|r| r.get(0),
)
.optional()?;
Ok(seq.unwrap_or(0) as u64)
}
fn advance_bookmark_in(conn: &Connection, room_id: &str, reader: &str) -> Result<u64> {
let start = Self::bookmark_in(conn, room_id, reader)?;
let mut bm = start;
'pages: loop {
let mut stmt = conn.prepare_cached(
"SELECT m.seq, m.from_name, m.type, d.state FROM messages m
LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
WHERE m.room_id=?1 AND m.seq>?3 ORDER BY m.seq LIMIT 500",
)?;
let rows: Vec<(i64, String, String, Option<String>)> = stmt
.query_map(params![room_id, reader, bm as i64], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))
})?
.collect::<rusqlite::Result<_>>()?;
let n = rows.len();
for (seq, from, kind, state) in rows {
let settled = !wanted(&from, &kind, reader)
|| matches!(state.as_deref(), Some("acked") | Some("quarantined"));
if !settled {
break 'pages;
}
bm = seq as u64;
}
if n < 500 {
break;
}
}
if bm > start {
conn.execute(
"INSERT INTO bookmarks (room_id, reader, seq) VALUES (?1, ?2, ?3)
ON CONFLICT(room_id, reader) DO UPDATE SET seq=MAX(bookmarks.seq, excluded.seq)",
params![room_id, reader, bm as i64],
)?;
}
Ok(bm)
}
pub fn delivery_lease(
&self,
room_id: &str,
reader: &str,
now: &str,
lease_until: &str,
max_attempts: u32,
) -> Result<Option<(Message, Delivery)>> {
let mut conn = self.lock();
let tx = conn.transaction()?;
tx.execute(
"UPDATE deliveries SET state='quarantined', token=NULL, updated=?3
WHERE room_id=?1 AND reader=?2 AND state='leased' AND lease_until<=?3 AND attempt>=?4",
params![room_id, reader, now, max_attempts],
)?;
let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
let mut seq: Option<i64> = tx
.query_row(
"SELECT seq FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq<=?3 AND (
state='replay'
OR (state='leased' AND lease_until<=?4)
OR (state='delayed' AND retry_at<=?4))
ORDER BY seq LIMIT 1",
params![room_id, reader, bm as i64, now],
|r| r.get(0),
)
.optional()?;
let mut after = bm as i64;
while seq.is_none() {
let rows: Vec<Candidate> = {
let mut stmt = tx.prepare_cached(
"SELECT m.seq, m.from_name, m.type, d.state, d.lease_until, d.retry_at
FROM messages m
LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
WHERE m.room_id=?1 AND m.seq>?3 ORDER BY m.seq LIMIT 500",
)?;
let rows = stmt.query_map(params![room_id, reader, after], |r| {
Ok((
r.get(0)?,
r.get(1)?,
r.get(2)?,
r.get(3)?,
r.get(4)?,
r.get(5)?,
))
})?;
rows.collect::<rusqlite::Result<_>>()?
};
if rows.is_empty() {
break;
}
for (s, from, kind, state, lease_until, retry_at) in &rows {
after = *s;
if !wanted(from, kind, reader) {
continue;
}
let free = match state.as_deref() {
None | Some("replay") => true,
Some("leased") => lease_until.as_deref().is_some_and(|t| t <= now),
Some("delayed") => retry_at.as_deref().is_some_and(|t| t <= now),
_ => false,
};
if free {
seq = Some(*s);
break;
}
}
}
let Some(seq) = seq else {
tx.commit()?;
return Ok(None);
};
let token = format!(
"d_{}",
data_encoding::HEXLOWER.encode(&rand::random::<[u8; 16]>())
);
tx.execute(
"INSERT INTO deliveries (room_id, reader, seq, state, token, lease_until, attempt, updated)
VALUES (?1, ?2, ?3, 'leased', ?4, ?5, 1, ?6)
ON CONFLICT(room_id, reader, seq) DO UPDATE SET
state='leased', token=excluded.token, lease_until=excluded.lease_until,
retry_at=NULL, attempt=deliveries.attempt+1, updated=excluded.updated",
params![room_id, reader, seq, token, lease_until, now],
)?;
self.faults.check("delivery.lease")?;
let msg = tx.query_row(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id WHERE m.room_id=?1 AND m.seq=?2",
params![room_id, seq],
|r| self.row_to_message(r),
)?;
let d = tx.query_row(
"SELECT * FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq=?3",
params![room_id, reader, seq],
Self::row_to_delivery,
)?;
tx.commit()?;
Ok(Some((msg, d)))
}
pub fn delivery_by_token(&self, token: &str) -> Result<Option<Delivery>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM deliveries WHERE token=?1",
params![token],
Self::row_to_delivery,
)
.optional()?)
}
pub fn delivery_settle(
&self,
token: &str,
how: &Settle,
now: &str,
max_attempts: u32,
) -> Result<Delivery> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let d = tx
.query_row(
"SELECT * FROM deliveries WHERE token=?1",
params![token],
Self::row_to_delivery,
)
.optional()?
.ok_or_else(|| {
Error::Denied(
"that delivery is no longer yours: it was settled, or its lease ran out and \
it was handed out again under a new token"
.into(),
)
})?;
let key = params![d.room_id, d.reader, d.seq as i64, now];
match (d.state, how) {
(DeliveryState::Acked, Settle::Ack) => {}
(DeliveryState::Leased, Settle::Ack) => {
tx.execute(
"UPDATE deliveries SET state='acked', lease_until=NULL, updated=?4
WHERE room_id=?1 AND reader=?2 AND seq=?3",
key,
)?;
}
(DeliveryState::Leased, Settle::Renew { lease_until }) => {
tx.execute(
"UPDATE deliveries SET lease_until=?5, updated=?4
WHERE room_id=?1 AND reader=?2 AND seq=?3",
params![d.room_id, d.reader, d.seq as i64, now, lease_until],
)?;
}
(DeliveryState::Leased, Settle::Nack { retry_at }) => {
let quarantine = d.attempt >= max_attempts;
tx.execute(
"UPDATE deliveries SET state=?5, token=NULL, lease_until=NULL, retry_at=?6, updated=?4
WHERE room_id=?1 AND reader=?2 AND seq=?3",
params![
d.room_id,
d.reader,
d.seq as i64,
now,
if quarantine { "quarantined" } else { "delayed" },
if quarantine { None } else { Some(retry_at) }
],
)?;
}
(state, _) => {
return Err(Error::Denied(format!(
"that delivery is {} and can only be acked again",
state.as_str()
)))
}
}
Self::advance_bookmark_in(&tx, &d.room_id, &d.reader)?;
let day_ago = chrono::DateTime::parse_from_rfc3339(now)
.map(|t| {
(t - chrono::Duration::days(1)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
})
.unwrap_or_default();
tx.execute(
"DELETE FROM deliveries WHERE room_id=?1 AND reader=?2 AND state='acked'
AND updated<?3 AND seq<=(SELECT seq FROM bookmarks WHERE room_id=?1 AND reader=?2)",
params![d.room_id, d.reader, day_ago],
)?;
let out = tx.query_row(
"SELECT * FROM deliveries WHERE room_id=?1 AND reader=?2 AND seq=?3",
params![d.room_id, d.reader, d.seq as i64],
Self::row_to_delivery,
)?;
tx.commit()?;
Ok(out)
}
pub fn delivery_ack_seqs(
&self,
room_id: &str,
reader: &str,
seqs: &[u64],
now: &str,
) -> Result<u64> {
let mut conn = self.lock();
let tx = conn.transaction()?;
for seq in seqs {
tx.execute(
"INSERT INTO deliveries (room_id, reader, seq, state, attempt, updated)
VALUES (?1, ?2, ?3, 'acked', 1, ?4)
ON CONFLICT(room_id, reader, seq) DO UPDATE SET
state='acked', token=NULL, lease_until=NULL, retry_at=NULL, updated=?4
WHERE deliveries.state != 'quarantined'",
params![room_id, reader, *seq as i64, now],
)?;
}
let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
tx.commit()?;
Ok(bm)
}
pub fn delivery_replay(
&self,
room_id: &str,
reader: &str,
seq: u64,
now: &str,
) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE deliveries SET state='replay', token=NULL, lease_until=NULL, retry_at=NULL,
attempt=0, updated=?4
WHERE room_id=?1 AND reader=?2 AND seq=?3 AND state='quarantined'",
params![room_id, reader, seq as i64, now],
)?;
Ok(n == 1)
}
pub fn delivery_peek(
&self,
room_id: &str,
reader: &str,
now: &str,
limit: u32,
) -> Result<Vec<Message>> {
let bm = {
let mut conn = self.lock();
let tx = conn.transaction()?;
let bm = Self::advance_bookmark_in(&tx, room_id, reader)?;
tx.commit()?;
bm
};
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT m.envelope, c.body, c.deleted FROM messages m
LEFT JOIN contents c ON c.msg_id = m.id
LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
WHERE m.room_id=?1 AND m.seq>?3
AND m.from_name != ?2 AND m.type NOT IN ('system', 'control')
AND (d.state IS NULL OR d.state IN ('leased', 'replay')
OR (d.state='delayed' AND d.retry_at<=?4))
ORDER BY m.seq LIMIT ?5",
)?;
let rows = stmt.query_map(params![room_id, reader, bm as i64, now, limit], |r| {
self.row_to_message(r)
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn deliveries(&self, room_id: &str, reader: Option<&str>) -> Result<Vec<Delivery>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT * FROM deliveries WHERE room_id=?1 AND (?2 IS NULL OR reader=?2)
AND state != 'acked' ORDER BY reader, seq",
)?;
let rows = stmt.query_map(params![room_id, reader], Self::row_to_delivery)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn delivery_next_due(
&self,
room_id: &str,
reader: &str,
now: &str,
) -> Result<Option<String>> {
let conn = self.lock();
Ok(conn.query_row(
"SELECT MIN(t) FROM (
SELECT lease_until AS t FROM deliveries
WHERE room_id=?1 AND reader=?2 AND state='leased' AND lease_until>?3
UNION ALL
SELECT retry_at AS t FROM deliveries
WHERE room_id=?1 AND reader=?2 AND state='delayed' AND retry_at>?3)",
params![room_id, reader, now],
|r| r.get(0),
)?)
}
pub fn wake_pending(
&self,
room_id: &str,
reader: &str,
now: &str,
trace_limit: u64,
) -> Result<WakePending> {
let conn = self.lock();
let bm = Self::bookmark_in(&conn, room_id, reader)?;
let mut stmt = conn.prepare_cached(
"SELECT m.seq, m.id, json_extract(m.envelope, '$.trace') FROM messages m
LEFT JOIN deliveries d ON d.room_id=m.room_id AND d.reader=?2 AND d.seq=m.seq
WHERE m.room_id=?1 AND m.from_name != ?2 AND m.type NOT IN ('system', 'control')
AND ((d.state IS NULL AND m.seq > ?3)
OR d.state='replay'
OR (d.state='leased' AND d.lease_until<=?4)
OR (d.state='delayed' AND d.retry_at<=?4))
ORDER BY m.seq LIMIT 10000",
)?;
let rows: Vec<(i64, String, Option<String>)> = stmt
.query_map(params![room_id, reader, bm as i64, now], |r| {
Ok((r.get(0)?, r.get(1)?, r.get(2)?))
})?
.collect::<rusqlite::Result<_>>()?;
let mut sent_on: HashMap<String, u64> = HashMap::new();
let mut out = WakePending::default();
for (seq, id, trace) in rows {
if let Some(t) = trace.filter(|t| !t.is_empty()) {
let n = match sent_on.get(&t) {
Some(n) => *n,
None => {
let n: i64 = conn.query_row(
"SELECT COUNT(*) FROM messages WHERE room_id=?1 AND from_name=?2
AND json_extract(envelope, '$.trace')=?3",
params![room_id, reader, t],
|r| r.get(0),
)?;
sent_on.insert(t.clone(), n as u64);
n as u64
}
};
if n > trace_limit {
out.looped += 1;
continue;
}
}
out.count += 1;
out.newest_seq = seq as u64;
out.newest_id = id;
}
Ok(out)
}
pub fn outbox_add(&self, msg: &Message) -> Result<()> {
let conn = self.lock();
let sealed = self.seal(&serde_json::to_string(msg)?)?;
self.faults.check("outbox.add")?;
conn.execute(
"INSERT OR IGNORE INTO outbox (msg_id, room_id, message, created, sender, state, updated, sig)
VALUES (?1, ?2, ?3, ?4, ?5, 'pending', ?6, ?7)",
params![
msg.id,
msg.room,
sealed,
msg.ts,
msg.from,
crate::message::now_ts(),
msg.sig
],
)?;
Ok(())
}
fn row_to_outbox(&self, row: &Row<'_>) -> rusqlite::Result<OutboxEntry> {
let state: String = row.get("state")?;
let stored: String = row.get("message")?;
let message = if stored.is_empty() {
None
} else {
self.unseal(&stored)
.and_then(|plain| serde_json::from_str::<Message>(&plain).ok())
};
Ok(OutboxEntry {
msg_id: row.get("msg_id")?,
room_id: row.get("room_id")?,
sender: row.get("sender")?,
state: state.parse().unwrap_or(OutboxState::Pending),
created: row.get("created")?,
attempts: row.get::<_, i64>("attempts")? as u32,
retry_at: row.get("retry_at")?,
reason: row.get("last_error")?,
reason_class: row.get("reason_class")?,
reason_code: row.get("reason_code")?,
updated: row.get("updated")?,
unreadable: !stored.is_empty() && message.is_none(),
message,
})
}
pub fn outbox_get(&self, msg_id: &str) -> Result<Option<OutboxEntry>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM outbox WHERE msg_id=?1",
params![msg_id],
|r| self.row_to_outbox(r),
)
.optional()?)
}
pub fn outbox_lanes(&self, room_id: &str) -> Result<Vec<(String, Vec<OutboxEntry>)>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT * FROM outbox WHERE room_id=?1 AND state IN ('pending', 'waiting')
ORDER BY sender, rowid",
)?;
let rows = stmt.query_map(params![room_id], |r| self.row_to_outbox(r))?;
let mut lanes: Vec<(String, Vec<OutboxEntry>)> = Vec::new();
for e in rows {
let e = e?;
match lanes.last_mut() {
Some((sender, lane)) if *sender == e.sender => lane.push(e),
_ => lanes.push((e.sender.clone(), vec![e])),
}
}
Ok(lanes)
}
pub fn outbox_entries(&self, room_id: Option<&str>) -> Result<Vec<OutboxEntry>> {
let conn = self.lock();
let mut stmt =
conn.prepare("SELECT * FROM outbox WHERE (?1 IS NULL OR room_id=?1) ORDER BY rowid")?;
let rows = stmt.query_map(params![room_id], |r| self.row_to_outbox(r))?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
fn outbox_set(&self, sql: &str, args: &[&dyn rusqlite::ToSql]) -> Result<bool> {
let conn = self.lock();
self.faults.check("outbox.state")?;
Ok(conn.execute(sql, args)? == 1)
}
pub fn outbox_wait(
&self,
msg_id: &str,
class: &str,
reason: &str,
code: Option<i32>,
retry_at: &str,
now: &str,
) -> Result<bool> {
self.outbox_set(
"UPDATE outbox SET state='waiting', attempts=attempts+1, retry_at=?3,
reason_class=?4, last_error=?5, reason_code=?6, updated=?2
WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
&[&msg_id, &now, &retry_at, &class, &reason, &code],
)
}
pub fn outbox_fail(
&self,
msg_id: &str,
reason: &str,
code: Option<i32>,
now: &str,
) -> Result<bool> {
self.outbox_set(
"UPDATE outbox SET state='failed', attempts=attempts+1, retry_at=NULL,
reason_class='refused', last_error=?3, reason_code=?4, updated=?2
WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
&[&msg_id, &now, &reason, &code],
)
}
#[allow(clippy::too_many_arguments)]
pub fn outbox_unknown(
&self,
msg_id: &str,
reason: &str,
code: Option<i32>,
retry_at: &str,
now: &str,
min_count: u32,
min_secs: i64,
) -> Result<OutboxState> {
let mut conn = self.lock();
let tx = conn.transaction()?;
let row: Option<(Option<i32>, Option<String>, i64)> = tx
.query_row(
"SELECT unknown_code, unknown_since, unknown_count FROM outbox
WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
params![msg_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let Some((last_code, since, count)) = row else {
return Err(Error::Invalid(format!(
"{msg_id} is not waiting to be sent"
)));
};
let (since, count) = match (last_code == code, since) {
(true, Some(since)) => (since, count as u32 + 1),
_ => (now.to_string(), 1),
};
let held_for = seconds_between(&since, now);
let state = if count >= min_count && held_for >= min_secs {
OutboxState::Quarantined
} else {
OutboxState::Waiting
};
self.faults.check("outbox.state")?;
tx.execute(
"UPDATE outbox SET state=?2, attempts=attempts+1, retry_at=?3,
reason_class='unknown', last_error=?4, reason_code=?5,
unknown_code=?5, unknown_since=?6, unknown_count=?7, updated=?8
WHERE msg_id=?1",
params![
msg_id,
state.as_str(),
if state == OutboxState::Waiting {
Some(retry_at)
} else {
None
},
reason,
code,
since,
count,
now
],
)?;
tx.commit()?;
Ok(state)
}
pub fn outbox_retry(&self, msg_id: &str, now: &str) -> Result<bool> {
self.outbox_set(
"UPDATE outbox SET state='pending', attempts=0, retry_at=NULL,
unknown_code=NULL, unknown_since=NULL, unknown_count=0, updated=?2
WHERE msg_id=?1 AND state IN ('failed', 'quarantined') AND message != ''",
&[&msg_id, &now],
)
}
pub fn outbox_drop(&self, msg_id: &str, now: &str) -> Result<bool> {
self.outbox_set(
"UPDATE outbox SET state='dropped', message='', retry_at=NULL, updated=?2
WHERE msg_id=?1 AND state IN ('failed', 'quarantined')",
&[&msg_id, &now],
)
}
pub fn outbox_wake(&self, room_id: &str, class: &str) -> Result<u64> {
let conn = self.lock();
let n = conn.execute(
"UPDATE outbox SET retry_at=NULL
WHERE room_id=?1 AND state='waiting' AND reason_class=?2",
params![room_id, class],
)?;
Ok(n as u64)
}
pub fn outbox_remove(&self, msg_id: &str) -> Result<()> {
let conn = self.lock();
conn.execute(
"DELETE FROM outbox WHERE msg_id=?1 AND state IN ('pending', 'waiting')",
params![msg_id],
)?;
Ok(())
}
pub fn outbox_count(&self, room_id: &str) -> Result<u64> {
Ok(self.outbox_counts(room_id)?.queued())
}
pub fn outbox_counts(&self, room_id: &str) -> Result<OutboxCounts> {
let conn = self.lock();
let mut stmt =
conn.prepare("SELECT state, COUNT(*) FROM outbox WHERE room_id=?1 GROUP BY state")?;
let rows = stmt.query_map(params![room_id], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?))
})?;
let mut c = OutboxCounts::default();
for row in rows {
let (state, n) = row?;
let n = n as u64;
match state.parse().unwrap_or(OutboxState::Pending) {
OutboxState::Pending => c.pending += n,
OutboxState::Waiting => c.waiting += n,
OutboxState::Failed => c.failed += n,
OutboxState::Quarantined => c.quarantined += n,
OutboxState::Dropped => c.dropped += n,
}
}
Ok(c)
}
pub fn invite_record(&self, inv: &Invite) -> Result<()> {
let conn = self.lock();
conn.execute(
"INSERT OR IGNORE INTO invites (nonce, room_id, name, created, expires) VALUES (?1, ?2, ?3, ?4, ?5)",
params![inv.nonce, inv.room_id, inv.name, inv.created, inv.expires],
)?;
Ok(())
}
pub fn invite_use(&self, nonce: &str, used_by: &str, now: &str) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"UPDATE invites SET used_at=?2, used_by=?3 WHERE nonce=?1 AND used_at IS NULL",
params![nonce, now, used_by],
)?;
Ok(n == 1)
}
fn count_since(
conn: &Connection,
room_id: &str,
from: Option<&str>,
since: &str,
) -> Result<u32> {
let n: i64 = match from {
Some(f) => conn.query_row(
"SELECT COUNT(*) FROM messages WHERE room_id=?1 AND from_name=?2 AND received>=?3",
params![room_id, f, since],
|r| r.get(0),
)?,
None => conn.query_row(
"SELECT COUNT(*) FROM messages WHERE room_id=?1 AND received>=?2",
params![room_id, since],
|r| r.get(0),
)?,
};
Ok(n as u32)
}
pub fn check_limits(&self, room_id: &str, from: &str, limits: &Limits) -> Result<bool> {
let conn = self.lock();
let now = chrono::Utc::now();
let minute_ago =
(now - chrono::Duration::minutes(1)).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let day_start = now
.date_naive()
.and_hms_opt(0, 0, 0)
.expect("midnight")
.and_utc()
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let per_minute = Self::count_since(&conn, room_id, Some(from), &minute_ago)?;
if per_minute >= limits.per_minute_per_sender {
let oldest: Option<String> = conn.query_row(
"SELECT MIN(received) FROM messages WHERE room_id=?1 AND from_name=?2 AND received>=?3",
params![room_id, from, minute_ago],
|r| r.get(0),
)?;
let wait = oldest
.map(|o| {
60 - seconds_between(
&o,
&now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
)
})
.unwrap_or(60)
.clamp(1, 60) as u64;
return Err(Error::OverBudget(
format!(
"{from} sent {per_minute} messages in the last minute; the limit is {}",
limits.per_minute_per_sender
),
Some(wait),
));
}
let today = Self::count_since(&conn, room_id, None, &day_start)?;
if today >= limits.daily_per_room {
let midnight = now
.date_naive()
.and_hms_opt(0, 0, 0)
.expect("midnight")
.and_utc()
+ chrono::Duration::days(1);
let wait = (midnight - now).num_seconds().max(1) as u64;
return Err(Error::OverBudget(
format!(
"room has used its daily budget of {} messages",
limits.daily_per_room
),
Some(wait),
));
}
let alert_at = limits.daily_per_room as u64 * limits.burst_alert_percent as u64 / 100;
Ok(today as u64 + 1 == alert_at)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoredFile {
pub room_id: String,
pub hash: String,
pub size: u64,
pub kind: String,
pub uploader: String,
pub created: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FileRefRow {
pub msg_id: String,
pub room_id: String,
pub file: crate::files::FileRef,
pub from: String,
pub to: Option<String>,
}
pub const ORPHAN_FILE_SECS: i64 = 24 * 3600;
impl Store {
fn row_to_ref(r: &Row) -> rusqlite::Result<FileRefRow> {
Ok(FileRefRow {
msg_id: r.get("msg_id")?,
room_id: r.get("room_id")?,
file: crate::files::FileRef {
id: r.get("hash")?,
name: r.get("name")?,
size: r.get::<_, i64>("size")? as u64,
mime: r.get("type")?,
},
from: r.get("from_name")?,
to: r.get("to_name")?,
})
}
pub fn file_refs(&self, room_id: &str, hash: &str) -> Result<Vec<FileRefRow>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT f.* FROM file_refs f JOIN messages m ON m.id = f.msg_id
WHERE f.room_id=?1 AND f.hash=?2 ORDER BY m.seq",
)?;
let rows = stmt.query_map(params![room_id, hash], Self::row_to_ref)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn file_refs_like(&self, room_id: &str, prefix: &str) -> Result<Vec<FileRefRow>> {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT f.* FROM file_refs f JOIN messages m ON m.id = f.msg_id
WHERE f.room_id=?1 AND substr(f.hash, 1, length(?2))=?2 ORDER BY f.hash, m.seq",
)?;
let rows = stmt.query_map(params![room_id, prefix], Self::row_to_ref)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn file_refs_of(&self, msg_id: &str) -> Result<Vec<FileRefRow>> {
let conn = self.lock();
let mut stmt = conn.prepare("SELECT * FROM file_refs WHERE msg_id=?1 ORDER BY rowid")?;
let rows = stmt.query_map(params![msg_id], Self::row_to_ref)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
fn row_to_file(r: &Row) -> rusqlite::Result<StoredFile> {
Ok(StoredFile {
room_id: r.get("room_id")?,
hash: r.get("hash")?,
size: r.get::<_, i64>("size")? as u64,
kind: r.get("kind")?,
uploader: r.get("uploader")?,
created: r.get("created")?,
})
}
pub fn file_stored(&self, room_id: &str, hash: &str) -> Result<Option<StoredFile>> {
let conn = self.lock();
Ok(conn
.query_row(
"SELECT * FROM files WHERE room_id=?1 AND hash=?2",
params![room_id, hash],
Self::row_to_file,
)
.optional()?)
}
pub fn files_stored(&self, room_id: &str) -> Result<Vec<StoredFile>> {
let conn = self.lock();
let mut stmt = conn.prepare("SELECT * FROM files WHERE room_id=?1 ORDER BY created")?;
let rows = stmt.query_map(params![room_id], Self::row_to_file)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
pub fn file_add(&self, f: &StoredFile) -> Result<bool> {
let conn = self.lock();
let n = conn.execute(
"INSERT OR IGNORE INTO files (room_id, hash, size, kind, uploader, created)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
params![
f.room_id,
f.hash,
f.size as i64,
f.kind,
f.uploader,
f.created
],
)?;
Ok(n == 1)
}
pub fn file_bytes(&self, room_id: &str, uploader: &str, since: &str) -> Result<(u64, u64)> {
let conn = self.lock();
let room: i64 = conn.query_row(
"SELECT COALESCE(SUM(size), 0) FROM files WHERE room_id=?1",
params![room_id],
|r| r.get(0),
)?;
let mine: i64 = conn.query_row(
"SELECT COALESCE(SUM(size), 0) FROM files WHERE room_id=?1 AND uploader=?2 AND created>=?3",
params![room_id, uploader, since],
|r| r.get(0),
)?;
Ok((room as u64, mine as u64))
}
pub fn file_fetched(&self, room_id: &str, hash: &str, member: &str, now: &str) -> Result<()> {
let conn = self.lock();
conn.execute(
"INSERT OR IGNORE INTO file_fetches (room_id, hash, member, at) VALUES (?1, ?2, ?3, ?4)",
params![room_id, hash, member, now],
)?;
Ok(())
}
pub fn file_sweep(
&self,
room_id: &str,
now: chrono::DateTime<chrono::Utc>,
keep_days: u64,
max_days: u64,
) -> Result<Vec<String>> {
let ts =
|t: chrono::DateTime<chrono::Utc>| t.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let days = |d: u64| chrono::Duration::days(d.min(100_000) as i64);
let too_old = ts(now - days(max_days));
let orphan = ts(now - chrono::Duration::seconds(ORPHAN_FILE_SECS));
let settled = ts(now - days(keep_days));
let members: Vec<String> = self
.members(room_id)?
.into_iter()
.filter(|m| !m.revoked)
.map(|m| m.name)
.collect();
let mut due = Vec::new();
for f in self.files_stored(room_id)? {
let refs = self.file_refs(room_id, &f.hash)?;
let go = if f.created < too_old {
true
} else if refs.is_empty() {
f.created < orphan
} else {
let mut audience: Vec<String> = Vec::new();
for r in &refs {
match &r.to {
Some(to) => audience.push(to.clone()),
None => audience.extend(members.iter().filter(|m| **m != r.from).cloned()),
}
}
audience.sort();
audience.dedup();
let fetched: HashMap<String, String> = {
let conn = self.lock();
let mut stmt = conn.prepare(
"SELECT member, at FROM file_fetches WHERE room_id=?1 AND hash=?2",
)?;
let rows = stmt.query_map(params![room_id, f.hash], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})?;
rows.collect::<rusqlite::Result<HashMap<_, _>>>()?
};
let mut last = f.created.clone();
let mut all = true;
for who in &audience {
match fetched.get(who) {
Some(at) => last = last.max(at.clone()),
None => all = false,
}
}
all && last < settled
};
if go {
due.push(f.hash);
}
}
let conn = self.lock();
for h in &due {
conn.execute(
"DELETE FROM files WHERE room_id=?1 AND hash=?2",
params![room_id, h],
)?;
conn.execute(
"DELETE FROM file_fetches WHERE room_id=?1 AND hash=?2",
params![room_id, h],
)?;
}
Ok(due)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::keys::{Identity, Kind};
use crate::message::{Draft, MessageType};
fn room(owner: &Identity) -> Room {
Room {
id: "r_test".into(),
name: "ops".into(),
about: "".into(),
owner: owner.public(),
created: "2026-01-01T00:00:00Z".into(),
retention_days: None,
class: DataClass::Internal,
paused: false,
hold: false,
closed: false,
home_node: "node".into(),
home_hints: serde_json::Value::Null,
}
}
fn msg(owner: &Identity, text: &str) -> Message {
Message::new(
Draft {
room: "r_test".into(),
from: "haris".into(),
text: text.into(),
kind: Some(MessageType::Chat),
..Default::default()
},
owner,
)
.unwrap()
}
#[test]
fn wake_pending_counts_what_waits_and_nothing_else() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let put = |from: &str, kind: MessageType, trace: Option<&str>| {
let mut m = Message::new(
Draft {
room: "r_test".into(),
from: from.into(),
text: "hi".into(),
kind: Some(kind),
trace: trace.map(String::from),
..Default::default()
},
&owner,
)
.unwrap();
s.sequence_and_append(&mut m).unwrap();
m
};
let now = "2026-09-27T12:00:00Z";
for _ in 0..3 {
put("haris", MessageType::Task, None);
}
put("bob", MessageType::Chat, None);
let newest = put("haris", MessageType::System, None);
let p = s.wake_pending("r_test", "bob", now, 5).unwrap();
assert_eq!((p.count, p.newest_seq), (3, 3));
assert_ne!(p.newest_id, newest.id);
let (leased, _) = s
.delivery_lease("r_test", "bob", now, "2026-09-27T12:10:00Z", 5)
.unwrap()
.unwrap();
assert_eq!(leased.seq, 1);
assert_eq!(s.wake_pending("r_test", "bob", now, 5).unwrap().count, 2);
assert_eq!(
s.wake_pending("r_test", "bob", "2026-09-27T12:11:00Z", 5)
.unwrap()
.count,
3
);
for _ in 0..6 {
put("bob", MessageType::Reply, Some("ping-pong"));
}
put("haris", MessageType::Reply, Some("ping-pong"));
put("haris", MessageType::Task, Some("fresh"));
let p = s.wake_pending("r_test", "bob", now, 5).unwrap();
assert_eq!((p.count, p.looped), (3, 1));
}
#[test]
fn rooms_and_members() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
assert!(matches!(
s.create_room(&room(&owner)),
Err(Error::NameTaken(_))
));
assert_eq!(s.room("ops").unwrap().id, "r_test");
assert!(matches!(s.room("nope"), Err(Error::NotInRoom(_))));
let m = Member {
room_id: "r_test".into(),
name: "haris".into(),
key: owner.public(),
kind: Kind::Human,
role: Role::Approver,
node: None,
granted_by: owner.public(),
expires_at: None,
joined_at: "2026-01-01T00:00:00Z".into(),
last_seen: None,
muted: false,
revoked: false,
profile: serde_json::Value::Null,
};
s.upsert_member(&m, Some("default")).unwrap();
assert_eq!(s.members("r_test").unwrap().len(), 1);
assert_eq!(s.local_members("r_test").unwrap()[0].identity, "default");
assert_eq!(
s.member_by_key("r_test", &owner.public())
.unwrap()
.unwrap()
.name,
"haris"
);
}
#[test]
fn chain_append_dedup_and_bookmarks() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
assert_eq!(
s.chain_head("r_test").unwrap(),
(0, GENESIS_PREV.to_string())
);
let mut a = msg(&owner, "one");
s.sequence_and_append(&mut a).unwrap();
assert_eq!(a.seq, 1);
assert_eq!(a.prev, GENESIS_PREV);
let mut b = msg(&owner, "two");
s.sequence_and_append(&mut b).unwrap();
assert_eq!(b.seq, 2);
assert_eq!(b.prev, a.chain_hash());
assert!(!s.append(&b).unwrap());
let mut c = msg(&owner, "three");
c.sequence(3, "sha256:bad");
assert!(s.append(&c).is_err());
c.sequence(3, &b.chain_hash());
assert!(s.append(&c).unwrap());
let all = s.messages_after("r_test", 0, 100).unwrap();
assert_eq!(all, vec![a.clone(), b.clone(), c.clone()]);
assert!(all[0].verify(&owner.public()).is_ok());
assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
s.set_bookmark("r_test", "haris", 2).unwrap();
s.set_bookmark("r_test", "haris", 1).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), 2);
assert_eq!(s.bookmark("r_test", "bob").unwrap(), 0);
assert_eq!(s.messages_after("r_test", 2, 100).unwrap(), vec![c]);
assert_eq!(s.message_count("r_test").unwrap(), 3);
}
#[test]
fn outbox_and_invites() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let m = msg(&owner, "queued");
s.outbox_add(&m).unwrap();
s.outbox_add(&m).unwrap();
assert_eq!(s.outbox_count("r_test").unwrap(), 1);
let lanes = s.outbox_lanes("r_test").unwrap();
assert_eq!(lanes.len(), 1);
assert_eq!(lanes[0].0, "haris");
assert_eq!(lanes[0].1[0].message.as_ref(), Some(&m));
s.outbox_remove(&m.id).unwrap();
assert_eq!(s.outbox_count("r_test").unwrap(), 0);
let inv = Invite::create(
crate::invite::InviteSpec {
room_id: "r_test".into(),
room_name: "ops".into(),
name: "bob".into(),
kind: Kind::Agent,
role: Role::TaskGiver,
home_node: "node".into(),
home_hints: serde_json::Value::Null,
for_node: None,
ttl_hours: None,
},
&owner,
)
.unwrap();
s.invite_record(&inv).unwrap();
assert!(s.invite_use(&inv.nonce, "key", "now").unwrap());
assert!(!s.invite_use(&inv.nonce, "key2", "now").unwrap());
assert!(!s.invite_use("i_unknown", "key", "now").unwrap());
}
#[test]
fn limits() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let limits = Limits {
per_minute_per_sender: 2,
daily_per_room: 3,
burst_alert_percent: 100,
};
s.check_limits("r_test", "haris", &limits).unwrap();
let mut a = msg(&owner, "1");
s.sequence_and_append(&mut a).unwrap();
let mut b = msg(&owner, "2");
s.sequence_and_append(&mut b).unwrap();
assert!(matches!(
s.check_limits("r_test", "haris", &limits),
Err(Error::OverBudget(..))
));
assert!(s.check_limits("r_test", "bob", &limits).unwrap());
}
#[test]
fn appending_the_same_message_twice_stores_it_once() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let mut a = msg(&owner, "1");
let mut again = a.clone();
assert!(s.sequence_and_append(&mut a).unwrap());
assert!(!s.sequence_and_append(&mut again).unwrap());
assert_eq!(again.seq, 0);
assert_eq!(s.message_count("r_test").unwrap(), 1);
}
#[test]
fn a_backdated_ts_still_counts_against_the_limits() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let limits = Limits {
per_minute_per_sender: 2,
daily_per_room: 100,
burst_alert_percent: 100,
};
for text in ["1", "2"] {
let mut m = msg(&owner, text);
m.ts = "2000-01-01T00:00:00Z".into();
s.sequence_and_append(&mut m).unwrap();
}
assert!(matches!(
s.check_limits("r_test", "haris", &limits),
Err(Error::OverBudget(..))
));
}
fn msg_from(who: &Identity, from: &str, text: &str) -> Message {
Message::new(
Draft {
room: "r_test".into(),
from: from.into(),
text: text.into(),
kind: Some(MessageType::Chat),
..Default::default()
},
who,
)
.unwrap()
}
fn temp_db(tag: &str) -> std::path::PathBuf {
static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let dir = std::env::temp_dir().join(format!(
"diavlos-store-{tag}-{}-{}-{}",
std::process::id(),
NEXT.fetch_add(1, std::sync::atomic::Ordering::SeqCst),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
dir.join("diavlos.db")
}
const T0: &str = "2026-01-01T00:00:00Z";
#[test]
fn the_outbox_accounts_for_every_message() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let msgs: Vec<Message> = (0..6)
.map(|i| {
msg_from(
&owner,
if i % 2 == 0 { "haris" } else { "bob" },
&format!("m{i}"),
)
})
.collect();
for m in &msgs {
s.outbox_add(m).unwrap();
}
let lanes = s.outbox_lanes("r_test").unwrap();
assert_eq!(lanes.len(), 2);
assert!(lanes.iter().all(|(_, l)| l.len() == 3));
let mut delivered = msgs[0].clone();
s.sequence_and_append(&mut delivered).unwrap();
s.outbox_wait(
&msgs[1].id,
"transport",
"offline",
Some(3),
"2026-01-01T00:01:00Z",
T0,
)
.unwrap();
s.outbox_fail(&msgs[2].id, "denied", Some(6), T0).unwrap();
s.outbox_fail(&msgs[3].id, "denied", Some(6), T0).unwrap();
s.outbox_drop(&msgs[3].id, T0).unwrap();
for i in 0..5 {
let now = format!("2026-01-01T0{i}:00:00Z");
s.outbox_unknown(&msgs[4].id, "huh", Some(1), &now, &now, 5, 3600)
.unwrap();
}
let c = s.outbox_counts("r_test").unwrap();
assert_eq!(
c,
OutboxCounts {
pending: 1,
waiting: 1,
failed: 1,
quarantined: 1,
dropped: 1
}
);
let in_chain = s.messages_after("r_test", 0, 100).unwrap().len() as u64;
assert_eq!(
in_chain + c.pending + c.waiting + c.failed + c.quarantined + c.dropped,
msgs.len() as u64
);
let dropped = s.outbox_get(&msgs[3].id).unwrap().unwrap();
assert_eq!(dropped.state, OutboxState::Dropped);
assert!(dropped.message.is_none());
assert!(s.outbox_retry(&msgs[4].id, T0).unwrap());
let back = s.outbox_get(&msgs[4].id).unwrap().unwrap();
assert_eq!(back.state, OutboxState::Pending);
assert_eq!(back.message.as_ref(), Some(&msgs[4]));
assert!(!s.outbox_retry(&msgs[5].id, T0).unwrap());
assert!(!s.outbox_drop(&msgs[5].id, T0).unwrap());
assert!(!s.outbox_retry(&msgs[3].id, T0).unwrap());
}
#[test]
fn an_unknown_answer_is_quarantined_only_after_both_count_and_time() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let m = msg(&owner, "x");
s.outbox_add(&m).unwrap();
let at = |mins: u32| format!("2026-01-01T{:02}:{:02}:00Z", mins / 60, mins % 60);
for i in 0..20 {
let st = s
.outbox_unknown(&m.id, "huh", Some(1), &at(i), &at(i), 5, 3600)
.unwrap();
assert_eq!(st, OutboxState::Waiting);
}
let st = s
.outbox_unknown(&m.id, "other", Some(9), &at(61), &at(61), 5, 3600)
.unwrap();
assert_eq!(st, OutboxState::Waiting);
for i in [90, 120, 180] {
let st = s
.outbox_unknown(&m.id, "other", Some(9), &at(i), &at(i), 5, 3600)
.unwrap();
assert_eq!(st, OutboxState::Waiting);
}
let st = s
.outbox_unknown(&m.id, "other", Some(9), &at(200), &at(200), 5, 3600)
.unwrap();
assert_eq!(st, OutboxState::Quarantined);
assert_eq!(s.outbox_count("r_test").unwrap(), 0);
assert_eq!(s.outbox_counts("r_test").unwrap().quarantined, 1);
}
#[test]
fn queued_messages_are_sealed_and_old_plain_ones_get_sealed_too() {
let owner = Identity::generate("haris", Kind::Human);
let key: [u8; 32] = rand::random();
let canary = "CANARY-4f1b9e-queued-text";
let path = temp_db("seal");
let ids: Vec<String> = {
let s = Store::open_with_key(&path, None).unwrap();
s.create_room(&room(&owner)).unwrap();
(0..3)
.map(|i| {
let m = msg(&owner, &format!("{canary} {i}"));
s.outbox_add(&m).unwrap();
m.id
})
.collect()
};
{
let s = Store::open_unsealed(&path, Some(key)).unwrap();
s.faults.arm_after("migration.row", 2);
assert!(s.seal_legacy_outbox(2).is_err());
let entries = s.outbox_entries(None).unwrap();
assert_eq!(entries.len(), 3);
assert!(
entries.iter().all(|e| e.message.is_some()),
"plain and sealed both read"
);
}
{
let s = Store::open_with_key(&path, Some(key)).unwrap();
let entries = s.outbox_entries(None).unwrap();
assert_eq!(entries.len(), 3);
for (e, id) in entries.iter().zip(&ids) {
assert_eq!(&e.msg_id, id);
assert!(e.message.as_ref().unwrap().text.starts_with(canary));
}
s.outbox_add(&msg(&owner, &format!("{canary} new")))
.unwrap();
s.outbox_fail(&ids[0], "no", Some(6), T0).unwrap();
s.outbox_fail(&ids[1], "no", Some(6), T0).unwrap();
s.outbox_drop(&ids[1], T0).unwrap();
s.checkpoint().unwrap();
}
let mut bytes = std::fs::read(&path).unwrap();
bytes.extend(std::fs::read(path.with_extension("db-wal")).unwrap_or_default());
let text = String::from_utf8_lossy(&bytes);
assert!(
!text.contains(canary),
"plain queued text left in the database files"
);
let s = Store::open_with_key(&path, None).unwrap();
let e = s.outbox_get(&ids[2]).unwrap().unwrap();
assert!(e.unreadable && e.message.is_none());
let _ = std::fs::remove_dir_all(path.parent().unwrap());
}
#[test]
fn a_full_disk_refuses_the_message_and_leaves_nothing_half_written() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory_with_key(rand::random()).unwrap();
s.create_room(&room(&owner)).unwrap();
s.cap_size_for_test(4).unwrap();
let big = "x".repeat(3000);
let mut kept = Vec::new();
let err = loop {
let m = msg(&owner, &big);
match s.outbox_add(&m) {
Ok(()) => kept.push(m.id),
Err(e) => break (e, m.id),
}
assert!(kept.len() < 1000, "never filled up");
};
match &err.0 {
Error::Db(rusqlite::Error::SqliteFailure(f, _)) => {
assert_eq!(f.code, rusqlite::ErrorCode::DiskFull)
}
other => panic!("expected a full disk, got {other}"),
}
assert!(s.outbox_get(&err.1).unwrap().is_none());
for id in &kept {
assert!(s.outbox_get(id).unwrap().unwrap().message.is_some());
}
let before = s.message_count("r_test").unwrap();
let mut m = msg(&owner, &big);
assert!(s.sequence_and_append(&mut m).is_err());
assert_eq!(s.message_count("r_test").unwrap(), before);
}
#[test]
fn a_queued_copy_of_a_message_already_in_the_chain_is_cleared() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let mut m = msg(&owner, "stored at the home, answer lost");
s.sequence_and_append(&mut m).unwrap();
s.outbox_add(&m).unwrap();
assert!(!s.append(&m).unwrap(), "not stored twice");
assert!(s.outbox_get(&m.id).unwrap().is_none());
let mut other = msg(&owner, "something else");
other.id = m.id.clone();
s.outbox_add(&other).unwrap();
let _ = s.append(&m);
assert!(s.outbox_get(&m.id).unwrap().is_some());
}
fn at(secs: i64) -> String {
(chrono::DateTime::parse_from_rfc3339(T0).unwrap() + chrono::Duration::seconds(secs))
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
}
fn owed(n: usize) -> (Store, Vec<u64>) {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let mut seqs = Vec::new();
for i in 0..n {
let mut m = msg_from(&owner, "bob", &format!("m{i}"));
s.sequence_and_append(&mut m).unwrap();
seqs.push(m.seq);
}
(s, seqs)
}
fn lease(s: &Store, now: i64) -> Option<Delivery> {
s.delivery_lease("r_test", "haris", &at(now), &at(now + 600), 5)
.unwrap()
.map(|(_, d)| d)
}
fn ack(s: &Store, d: &Delivery, now: i64) -> Result<Delivery> {
s.delivery_settle(d.token.as_deref().unwrap(), &Settle::Ack, &at(now), 5)
}
#[test]
fn the_bookmark_never_passes_a_message_still_owed() {
let (s, seqs) = owed(3);
let a = lease(&s, 0).unwrap();
let b = lease(&s, 0).unwrap();
let c = lease(&s, 0).unwrap();
assert_eq!((a.seq, b.seq, c.seq), (seqs[0], seqs[1], seqs[2]));
assert!(lease(&s, 0).is_none(), "all three are out");
ack(&s, &c, 1).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
ack(&s, &a, 1).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
ack(&s, &b, 1).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[2]);
ack(&s, &b, 2).unwrap();
}
#[test]
fn a_lease_that_runs_out_is_handed_out_again_and_the_old_token_is_refused() {
let (s, seqs) = owed(1);
let first = lease(&s, 0).unwrap();
assert_eq!(first.attempt, 1);
assert!(lease(&s, 599).is_none(), "still leased");
let second = lease(&s, 600).unwrap();
assert_eq!(second.seq, seqs[0]);
assert_eq!(second.attempt, 2);
assert_ne!(first.token, second.token);
assert!(matches!(ack(&s, &first, 700), Err(Error::Denied(_))));
let renew = Settle::Renew {
lease_until: at(5000),
};
assert!(s
.delivery_settle(first.token.as_deref().unwrap(), &renew, &at(700), 5)
.is_err());
assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
ack(&s, &second, 700).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
}
#[test]
fn a_late_ack_counts_if_nobody_was_handed_it_since() {
let (s, seqs) = owed(1);
let d = lease(&s, 0).unwrap();
ack(&s, &d, 5000).unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[0]);
}
#[test]
fn renew_holds_it_and_nack_hands_it_back_later() {
let (s, _) = owed(1);
let d = lease(&s, 0).unwrap();
let t = d.token.clone().unwrap();
s.delivery_settle(
&t,
&Settle::Renew {
lease_until: at(2000),
},
&at(500),
5,
)
.unwrap();
assert!(lease(&s, 1500).is_none(), "renewed past the first lease");
s.delivery_settle(&t, &Settle::Nack { retry_at: at(3000) }, &at(1600), 5)
.unwrap();
assert!(lease(&s, 2999).is_none(), "not before retry_at");
let again = lease(&s, 3000).unwrap();
assert_eq!(again.attempt, 2);
}
#[test]
fn five_strikes_quarantine_it_the_lane_moves_on_and_replay_brings_it_back() {
let (s, seqs) = owed(2);
let mut now = 0;
for i in 1..=5 {
let d = lease(&s, now).unwrap();
assert_eq!((d.seq, d.attempt), (seqs[0], i));
now += 600;
}
let next = lease(&s, now).unwrap();
assert_eq!(next.seq, seqs[1]);
let q = s.deliveries("r_test", Some("haris")).unwrap();
assert!(q
.iter()
.any(|d| d.seq == seqs[0] && d.state == DeliveryState::Quarantined));
ack(&s, &next, now).unwrap();
assert_eq!(
s.bookmark("r_test", "haris").unwrap(),
seqs[1],
"quarantine counts as settled"
);
assert!(s
.delivery_replay("r_test", "haris", seqs[0], &at(now))
.unwrap());
let back = lease(&s, now).unwrap();
assert_eq!((back.seq, back.attempt), (seqs[0], 1));
ack(&s, &back, now).unwrap();
assert!(lease(&s, now).is_none());
}
#[test]
fn own_and_housekeeping_messages_settle_by_themselves() {
let owner = Identity::generate("haris", Kind::Human);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let mut mine = msg(&owner, "from me");
s.sequence_and_append(&mut mine).unwrap();
let mut sys = msg(&owner, "joined");
sys.kind = MessageType::System;
s.sequence_and_append(&mut sys).unwrap();
let mut theirs = msg_from(&owner, "bob", "for you");
s.sequence_and_append(&mut theirs).unwrap();
let d = lease(&s, 0).unwrap();
assert_eq!(d.seq, theirs.seq);
assert_eq!(s.bookmark("r_test", "haris").unwrap(), sys.seq);
}
#[test]
fn peek_shows_what_is_owed_without_taking_it() {
let (s, seqs) = owed(3);
let d = lease(&s, 0).unwrap();
s.delivery_settle(
d.token.as_deref().unwrap(),
&Settle::Nack { retry_at: at(100) },
&at(1),
5,
)
.unwrap();
let peeked: Vec<u64> = s
.delivery_peek("r_test", "haris", &at(2), 10)
.unwrap()
.iter()
.map(|m| m.seq)
.collect();
assert_eq!(peeked, vec![seqs[1], seqs[2]], "the delayed one is not due");
let later: Vec<u64> = s
.delivery_peek("r_test", "haris", &at(100), 10)
.unwrap()
.iter()
.map(|m| m.seq)
.collect();
assert_eq!(later, seqs);
assert_eq!(lease(&s, 100).unwrap().seq, seqs[0]);
}
#[test]
fn read_and_ack_settles_exactly_what_was_read() {
let (s, seqs) = owed(3);
s.delivery_ack_seqs("r_test", "haris", &[seqs[1]], &at(0))
.unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), 0);
s.delivery_ack_seqs("r_test", "haris", &[seqs[0]], &at(0))
.unwrap();
assert_eq!(s.bookmark("r_test", "haris").unwrap(), seqs[1]);
assert_eq!(lease(&s, 0).unwrap().seq, seqs[2]);
}
#[test]
fn a_crash_after_the_lease_is_written_hands_it_out_again_later() {
let owner = Identity::generate("haris", Kind::Human);
let path = temp_db("lease");
let seq = {
let s = Store::open_with_key(&path, None).unwrap();
s.create_room(&room(&owner)).unwrap();
let mut m = msg_from(&owner, "bob", "work");
s.sequence_and_append(&mut m).unwrap();
lease(&s, 0).unwrap();
m.seq
};
let s = Store::open_with_key(&path, None).unwrap();
assert!(lease(&s, 10).is_none(), "the lease survived the restart");
let again = lease(&s, 600).unwrap();
assert_eq!((again.seq, again.attempt), (seq, 2));
s.faults.arm("delivery.lease");
assert!(s
.delivery_lease("r_test", "haris", &at(1300), &at(1900), 5)
.is_err());
let d = s.deliveries("r_test", Some("haris")).unwrap();
assert_eq!(d[0].attempt, 2);
let _ = std::fs::remove_dir_all(path.parent().unwrap());
}
#[test]
fn a_spend_and_its_audit_event_land_together_or_not_at_all() {
let owner = Identity::generate("haris", Kind::Human);
let bot = Identity::generate("bot", Kind::Agent);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
let member = |name: &str, id: &Identity, kind: Kind, role: Role| Member {
room_id: "r_test".into(),
name: name.into(),
key: id.public(),
kind,
role,
node: None,
granted_by: owner.public(),
expires_at: None,
joined_at: T0.into(),
last_seen: None,
muted: false,
revoked: false,
profile: serde_json::Value::Null,
};
s.upsert_member(&member("haris", &owner, Kind::Human, Role::Approver), None)
.unwrap();
s.upsert_member(&member("bot", &bot, Kind::Agent, Role::TaskGiver), None)
.unwrap();
let action: crate::message::Action =
serde_json::from_value(serde_json::json!({"verb":"deploy","target":"api","params":{}}))
.unwrap();
let mut q = Message::new(
Draft {
room: "r_test".into(),
from: "bot".into(),
text: "deploy?".into(),
kind: Some(MessageType::Question),
action: Some(action.clone()),
..Default::default()
},
&bot,
)
.unwrap();
s.sequence_and_append(&mut q).unwrap();
let mut a = Message::new(
Draft {
room: "r_test".into(),
from: "haris".into(),
text: "approved".into(),
kind: Some(MessageType::Approve),
reply_to: Some(q.id.clone()),
action_hash: Some(action.hash()),
expires: Some("2999-01-01T00:00:00Z".into()),
once: Some(true),
..Default::default()
},
&owner,
)
.unwrap();
s.sequence_and_append(&mut a).unwrap();
let req = SpendRequest {
room_id: "r_test".into(),
approve_id: a.id.clone(),
action_hash: action.hash(),
op_id: "op-1".into(),
spender: "bot".into(),
node: "n1".into(),
now: crate::message::now_ts(),
};
let audit = || msg(&owner, "spent");
let before = s.message_count("r_test").unwrap();
s.faults.arm("spend.commit");
assert!(s.spend_approve(&req, &mut audit()).is_err());
assert!(s.spend_of(&a.id).unwrap().is_none());
assert_eq!(s.message_count("r_test").unwrap(), before);
let SpendOutcome::Spent(rec) = s.spend_approve(&req, &mut audit()).unwrap() else {
panic!("not spent")
};
assert_eq!(s.message_count("r_test").unwrap(), before + 1);
assert_eq!(
s.spend_approve(&req, &mut audit()).unwrap(),
SpendOutcome::AlreadyYours(rec.clone())
);
let other = SpendRequest {
op_id: "op-2".into(),
..req.clone()
};
assert_eq!(
s.spend_approve(&other, &mut audit()).unwrap(),
SpendOutcome::SpentByAnother
);
assert_eq!(s.message_count("r_test").unwrap(), before + 1);
}
#[test]
fn files_are_indexed_kept_and_swept() {
let owner = Identity::generate("haris", Kind::Human);
let bot = Identity::generate("bot", Kind::Agent);
let s = Store::open_memory().unwrap();
s.create_room(&room(&owner)).unwrap();
for (name, id) in [("haris", &owner), ("bot", &bot)] {
s.upsert_member(
&Member {
room_id: "r_test".into(),
name: name.into(),
key: id.public(),
kind: Kind::Agent,
role: Role::TaskGiver,
node: None,
granted_by: owner.public(),
expires_at: None,
joined_at: T0.into(),
last_seen: None,
muted: false,
revoked: false,
profile: serde_json::Value::Null,
},
None,
)
.unwrap();
}
let hash = format!("sha256:{}", "a".repeat(64));
let orphan = format!("sha256:{}", "b".repeat(64));
let mut m = Message::new(
Draft {
room: "r_test".into(),
from: "haris".into(),
text: "log".into(),
kind: Some(MessageType::Chat),
data: serde_json::json!({"files": [{"id": hash, "name": "a.log", "size": 3, "type": "text/plain"}]}),
..Default::default()
},
&owner,
)
.unwrap();
assert!(s.sequence_and_append(&mut m).unwrap());
let refs = s.file_refs("r_test", &hash).unwrap();
assert_eq!(refs.len(), 1);
assert_eq!(refs[0].file.name, "a.log");
assert_eq!(s.file_refs_of(&m.id).unwrap()[0].from, "haris");
let day0 = chrono::DateTime::parse_from_rfc3339(T0)
.unwrap()
.with_timezone(&chrono::Utc);
let at = |d: i64| day0 + chrono::Duration::days(d);
let ts = |d: i64| at(d).to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
for h in [&hash, &orphan] {
s.file_add(&StoredFile {
room_id: "r_test".into(),
hash: h.clone(),
size: 3,
kind: "text".into(),
uploader: "haris".into(),
created: ts(0),
})
.unwrap();
}
assert_eq!(s.file_bytes("r_test", "haris", &ts(0)).unwrap(), (6, 6));
assert_eq!(s.file_bytes("r_test", "bot", &ts(0)).unwrap(), (6, 0));
assert_eq!(
s.file_sweep("r_test", at(2), 7, 30).unwrap(),
vec![orphan.clone()]
);
assert!(s.file_sweep("r_test", at(20), 7, 30).unwrap().is_empty());
s.file_fetched("r_test", &hash, "bot", &ts(20)).unwrap();
assert!(s.file_sweep("r_test", at(26), 7, 30).unwrap().is_empty());
assert_eq!(
s.file_sweep("r_test", at(28), 7, 30).unwrap(),
vec![hash.clone()]
);
assert!(s.file_stored("r_test", &hash).unwrap().is_none());
s.file_add(&StoredFile {
room_id: "r_test".into(),
hash: hash.clone(),
size: 3,
kind: "text".into(),
uploader: "haris".into(),
created: ts(0),
})
.unwrap();
assert_eq!(
s.file_sweep("r_test", at(31), 7, 30).unwrap(),
vec![hash.clone()]
);
assert!(s.tombstone_message(&m.id).unwrap());
assert!(s.file_refs("r_test", &hash).unwrap().is_empty());
}
}