986 lines
38 KiB
Rust
986 lines
38 KiB
Rust
//! Local state: one SQLite file, shared by both carriers (RNS.md sec 15).
|
|
//! Envelopes are stored sealed and opened on demand; no plaintext at rest.
|
|
//! The handle is `Send + Sync`: every operation locks an inner mutex, so a
|
|
//! GUI can hold one `Store` behind an executor and read it (for incremental
|
|
//! fetch progress) while a fetch runs on another thread.
|
|
|
|
use std::path::Path;
|
|
use std::sync::{Mutex, MutexGuard};
|
|
use std::time::{SystemTime, UNIX_EPOCH};
|
|
|
|
use rusqlite::{params, Connection};
|
|
|
|
use crate::account::Account;
|
|
use crate::address::Address;
|
|
use crate::crypto::ct_eq;
|
|
use crate::error::Error;
|
|
use crate::transport::{ID_LEN, KEY_LEN, TIER_MAIN, TIER_REQUESTS};
|
|
|
|
pub fn now() -> i64 {
|
|
SystemTime::now()
|
|
.duration_since(UNIX_EPOCH)
|
|
.expect("clock before 1970")
|
|
.as_secs() as i64
|
|
}
|
|
|
|
pub const SCHEMA: &str = "
|
|
-- One row, always present, so every update is a plain UPDATE.
|
|
CREATE TABLE IF NOT EXISTS state (
|
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
|
username TEXT, host TEXT, port INTEGER, scheme TEXT,
|
|
rotations INTEGER NOT NULL DEFAULT 0, -- sec 2 rotation index
|
|
after_time INTEGER NOT NULL DEFAULT 0, -- sec 6.1 FETCH cursor
|
|
after_id BLOB NOT NULL DEFAULT x'',
|
|
sync_ok INTEGER NOT NULL DEFAULT 1); -- may we replace the server's set?
|
|
INSERT OR IGNORE INTO state (id) VALUES (1);
|
|
CREATE TABLE IF NOT EXISTS servers (
|
|
host TEXT PRIMARY KEY, static BLOB NOT NULL, pinned_at INTEGER NOT NULL);
|
|
CREATE TABLE IF NOT EXISTS contacts (
|
|
address TEXT PRIMARY KEY, identity BLOB NOT NULL,
|
|
verified INTEGER NOT NULL, -- 1 when the key came from a self-certifying URI
|
|
seen_at INTEGER NOT NULL);
|
|
-- Correspondents admitted to this mailbox's main tier (sec 5.8). The
|
|
-- identity is frozen at acceptance because the token is derived from it: a
|
|
-- contact's later rotation must not change the token they already hold.
|
|
CREATE TABLE IF NOT EXISTS accepted (
|
|
address TEXT PRIMARY KEY, identity BLOB NOT NULL,
|
|
active INTEGER NOT NULL, added_at INTEGER NOT NULL);
|
|
-- Accept tokens received from correspondents, filed under the address that
|
|
-- issued them: an address outlives the keys behind it, so a token keeps
|
|
-- working across the issuer's rotations (sec 5.8).
|
|
CREATE TABLE IF NOT EXISTS tokens (
|
|
address TEXT PRIMARY KEY, token BLOB NOT NULL, seen_at INTEGER NOT NULL);
|
|
-- Every id ever fetched, so an envelope resent after we deleted it locally
|
|
-- is not stored again (sec 10).
|
|
CREATE TABLE IF NOT EXISTS seen (id BLOB PRIMARY KEY, at INTEGER NOT NULL);
|
|
CREATE TABLE IF NOT EXISTS inbox (
|
|
id BLOB PRIMARY KEY, envelope BLOB NOT NULL,
|
|
received_at INTEGER NOT NULL, tier INTEGER NOT NULL,
|
|
kept INTEGER NOT NULL DEFAULT 0); -- a server copy still exists (fetch --keep)
|
|
CREATE TABLE IF NOT EXISTS sent (
|
|
id BLOB PRIMARY KEY, recipient TEXT NOT NULL,
|
|
envelope BLOB NOT NULL, sent_at INTEGER NOT NULL);
|
|
-- Keys a contact used before its current one, with when each stopped being
|
|
-- current, in epoch milliseconds: the local record that a rotation
|
|
-- happened, so a later import can restore it (sec 8).
|
|
CREATE TABLE IF NOT EXISTS history (
|
|
address TEXT NOT NULL, identity BLOB NOT NULL, until INTEGER NOT NULL,
|
|
PRIMARY KEY (address, identity));
|
|
";
|
|
|
|
/// One sealed message in a folder, still encrypted.
|
|
pub struct Stored {
|
|
pub id: [u8; ID_LEN],
|
|
pub envelope: Vec<u8>,
|
|
pub at: i64,
|
|
/// Only sent copies carry a recipient.
|
|
pub recipient: Option<String>,
|
|
/// A server copy still exists, so a local delete must also remove it.
|
|
pub kept: bool,
|
|
}
|
|
|
|
pub struct Store {
|
|
db: Mutex<Connection>,
|
|
}
|
|
|
|
/// A contact row: address, identity, verified, and its acceptance state
|
|
/// (None when never accepted) for `fumi contacts`.
|
|
pub type ContactRow = (String, [u8; KEY_LEN], bool, Option<i64>);
|
|
|
|
/// The store's schema version, in SQLite's `user_version`. A store written
|
|
/// by a newer fumi is refused rather than misread; within a version the
|
|
/// schema is stable, and a bump comes with a documented migration.
|
|
pub const SCHEMA_VERSION: i32 = 2;
|
|
|
|
/// The migration from schema version 1: both changes are additive, so an
|
|
/// existing store is upgraded in place at open; only stores newer than this
|
|
/// build are refused.
|
|
const MIGRATE_FROM_1: &str = "
|
|
ALTER TABLE inbox ADD COLUMN kept INTEGER NOT NULL DEFAULT 0;
|
|
CREATE TABLE IF NOT EXISTS history (
|
|
address TEXT NOT NULL, identity BLOB NOT NULL, until INTEGER NOT NULL,
|
|
PRIMARY KEY (address, identity));
|
|
";
|
|
|
|
impl Store {
|
|
/// Opens the store at `path`. SQLite's own `:memory:` path works too;
|
|
/// `open_in_memory` spells that case without the magic string.
|
|
pub fn open(path: &Path) -> Result<Store, Error> {
|
|
Store::init(Connection::open(path)?)
|
|
}
|
|
|
|
/// An in-memory store: the same schema, version and contract, gone with
|
|
/// the handle. For tests, foreign importers assembling rows in code,
|
|
/// and FFI round-trips that never touch disk.
|
|
pub fn open_in_memory() -> Result<Store, Error> {
|
|
Store::init(Connection::open_in_memory()?)
|
|
}
|
|
|
|
fn init(db: Connection) -> Result<Store, Error> {
|
|
// WAL keeps concurrent readers cheap while a writer commits, and a
|
|
// short busy timeout absorbs the contention the mutex cannot (a
|
|
// second connection in the same process, or a CLI run alongside).
|
|
// An in-memory database journals in memory, so WAL applies to
|
|
// files only.
|
|
let journal: String = db.query_row("PRAGMA journal_mode", [], |row| row.get(0))?;
|
|
if journal != "memory" {
|
|
db.pragma_update(None, "journal_mode", "WAL")?;
|
|
}
|
|
db.busy_timeout(std::time::Duration::from_secs(5))?;
|
|
let version: i32 = db.query_row("PRAGMA user_version", [], |row| row.get(0))?;
|
|
if version == 0 {
|
|
db.execute_batch(SCHEMA)?;
|
|
db.pragma_update(None, "user_version", SCHEMA_VERSION)?;
|
|
} else if version == 1 {
|
|
db.execute_batch(MIGRATE_FROM_1)?;
|
|
db.pragma_update(None, "user_version", SCHEMA_VERSION)?;
|
|
} else if version != SCHEMA_VERSION {
|
|
return Err(Error::SchemaVersion {
|
|
found: version,
|
|
expected: SCHEMA_VERSION,
|
|
});
|
|
}
|
|
Ok(Store {
|
|
db: Mutex::new(db),
|
|
})
|
|
}
|
|
|
|
/// One lock acquisition; the guard's lifetime is the operation's.
|
|
fn db(&self) -> MutexGuard<'_, Connection> {
|
|
self.db.lock().expect("store mutex poisoned")
|
|
}
|
|
|
|
fn one<T>(&self, sql: &str, params: &[&dyn rusqlite::ToSql]) -> Result<Option<T>, Error>
|
|
where
|
|
T: rusqlite::types::FromSql,
|
|
{
|
|
Ok(
|
|
match self.db().query_row(sql, params, |row| row.get::<_, T>(0)) {
|
|
Ok(value) => Some(value),
|
|
Err(rusqlite::Error::QueryReturnedNoRows) => None,
|
|
Err(e) => return Err(e.into()),
|
|
},
|
|
)
|
|
}
|
|
|
|
/// The home address, when this identity has been registered or restored.
|
|
pub fn account(&self) -> Result<Option<Address>, Error> {
|
|
let row = self.db().query_row(
|
|
"SELECT username, host, port, scheme FROM state WHERE id = 1",
|
|
[],
|
|
|row| {
|
|
Ok((
|
|
row.get::<_, Option<String>>(0)?,
|
|
row.get::<_, Option<String>>(1)?,
|
|
row.get::<_, Option<i64>>(2)?,
|
|
row.get::<_, Option<String>>(3)?,
|
|
))
|
|
},
|
|
);
|
|
let (username, host, port, scheme) = match row {
|
|
Ok(tuple) => tuple,
|
|
Err(rusqlite::Error::QueryReturnedNoRows) => return Ok(None),
|
|
Err(e) => return Err(e.into()),
|
|
};
|
|
let (Some(username), Some(host)) = (username, host) else {
|
|
return Ok(None);
|
|
};
|
|
let scheme = scheme.unwrap_or_else(|| "tcp".to_string());
|
|
let port = port.unwrap_or(crate::address::DEFAULT_PORT as i64) as u16;
|
|
let text = if scheme == "rns" {
|
|
format!("smol+rns://{username}@{host}")
|
|
} else {
|
|
format!("{username}@{host}")
|
|
};
|
|
let mut addr = Address::parse(&text)
|
|
.map_err(|e| Error::Other(format!("stored account is not parseable: {e}")))?;
|
|
if addr.scheme == crate::address::Scheme::Tcp {
|
|
addr.port = port;
|
|
}
|
|
Ok(Some(addr))
|
|
}
|
|
|
|
pub fn set_account(&self, addr: &Address) -> Result<(), Error> {
|
|
let scheme = match addr.scheme {
|
|
crate::address::Scheme::Tcp => "tcp",
|
|
crate::address::Scheme::Rns => "rns",
|
|
};
|
|
self.db()
|
|
.execute(
|
|
"UPDATE state SET username = ?1, host = ?2, port = ?3, scheme = ?4 WHERE id = 1",
|
|
params![addr.user, addr.host, addr.port, scheme],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn rotations(&self) -> Result<u32, Error> {
|
|
Ok(self
|
|
.one::<i64>("SELECT rotations FROM state WHERE id = 1", &[])?
|
|
.unwrap_or(0) as u32)
|
|
}
|
|
|
|
pub fn set_rotations(&self, index: u32) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"UPDATE state SET rotations = ?1 WHERE id = 1",
|
|
params![index],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn cursor(&self) -> Result<(i64, [u8; ID_LEN]), Error> {
|
|
let (after_time, after_id): (i64, Vec<u8>) = self
|
|
.db()
|
|
.query_row(
|
|
"SELECT after_time, after_id FROM state WHERE id = 1",
|
|
[],
|
|
|row| Ok((row.get(0)?, row.get(1)?)),
|
|
)
|
|
.map_err(Error::from)?;
|
|
let after_id: [u8; ID_LEN] = if after_id.len() == ID_LEN {
|
|
after_id.try_into().unwrap()
|
|
} else {
|
|
[0u8; ID_LEN]
|
|
};
|
|
Ok((after_time, after_id))
|
|
}
|
|
|
|
pub fn set_cursor(&self, after_time: i64, after_id: &[u8; ID_LEN]) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"UPDATE state SET after_time = ?1, after_id = ?2 WHERE id = 1",
|
|
params![after_time, after_id],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Whether the local accept-token set may replace the server's: a client
|
|
/// restored from the master alone must not erase it (sec 4).
|
|
pub fn sync_ok(&self) -> Result<bool, Error> {
|
|
Ok(self
|
|
.one::<i64>("SELECT sync_ok FROM state WHERE id = 1", &[])?
|
|
.unwrap_or(1)
|
|
!= 0)
|
|
}
|
|
|
|
pub fn set_sync_ok(&self, ok: bool) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"UPDATE state SET sync_ok = ?1 WHERE id = 1",
|
|
params![ok as i64],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn server_pin(&self, host: &str) -> Result<Option<[u8; KEY_LEN]>, Error> {
|
|
Ok(self
|
|
.one::<Option<Vec<u8>>>("SELECT static FROM servers WHERE host = ?1", &[&host])?
|
|
.flatten()
|
|
.and_then(|k| k.try_into().ok()))
|
|
}
|
|
|
|
/// Removes a pin: the settings screen's unpin, the inverse of
|
|
/// `pin_server`. The next session against that host is trust on first
|
|
/// use again (sec 4, sec 8).
|
|
pub fn unpin_server(&self, host: &str) -> Result<(), Error> {
|
|
self.db()
|
|
.execute("DELETE FROM servers WHERE host = ?1", params![host])?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Every pin, host and key, ordered by host — the export's server map.
|
|
pub fn pins(&self) -> Result<Vec<(String, [u8; KEY_LEN])>, Error> {
|
|
let db = self.db();
|
|
let mut stmt = db.prepare("SELECT host, static FROM servers ORDER BY host")?;
|
|
let rows = stmt
|
|
.query_map([], |row| {
|
|
Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
|
|
})?
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
Ok(rows
|
|
.into_iter()
|
|
.filter_map(|(host, key)| Some((host, key.try_into().ok()?)))
|
|
.collect())
|
|
}
|
|
|
|
pub fn pin_server(&self, host: &str, key: &[u8; KEY_LEN]) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT INTO servers (host, static, pinned_at) VALUES (?1, ?2, ?3) \
|
|
ON CONFLICT (host) DO UPDATE SET static = ?2, pinned_at = ?3",
|
|
params![host, key, now()],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn contact(&self, address: &str) -> Result<Option<([u8; KEY_LEN], bool)>, Error> {
|
|
let row = self.db().query_row(
|
|
"SELECT identity, verified FROM contacts WHERE address = ?1",
|
|
[&address],
|
|
|row| Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, i64>(1)?)),
|
|
);
|
|
match row {
|
|
Ok((identity, verified)) => Ok(Some((
|
|
identity
|
|
.try_into()
|
|
.map_err(|_| Error::Other("stored contact identity is not 32 bytes".into()))?,
|
|
verified != 0,
|
|
))),
|
|
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
|
|
Err(e) => Err(e.into()),
|
|
}
|
|
}
|
|
|
|
pub fn save_contact(
|
|
&self,
|
|
address: &str,
|
|
identity: &[u8; KEY_LEN],
|
|
verified: bool,
|
|
) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT INTO contacts (address, identity, verified, seen_at) VALUES (?1, ?2, ?3, ?4) \
|
|
ON CONFLICT (address) DO UPDATE SET identity = ?2, verified = ?3, seen_at = ?4",
|
|
params![address, identity, verified as i64, now()],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Every contact with its acceptance state, for `fumi contacts`.
|
|
pub fn contact_rows(&self) -> Result<Vec<ContactRow>, Error> {
|
|
let db = self.db();
|
|
let mut stmt = db.prepare(
|
|
"SELECT c.address, c.identity, c.verified, a.active FROM contacts c \
|
|
LEFT JOIN accepted a ON a.address = c.address ORDER BY c.address",
|
|
)?;
|
|
let rows = stmt
|
|
.query_map([], |row| {
|
|
Ok((
|
|
row.get::<_, String>(0)?,
|
|
row.get::<_, Vec<u8>>(1)?,
|
|
row.get::<_, i64>(2)?,
|
|
row.get::<_, Option<i64>>(3)?,
|
|
))
|
|
})?
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
Ok(rows
|
|
.into_iter()
|
|
.filter_map(|(address, identity, verified, active)| {
|
|
Some((address, identity.try_into().ok()?, verified != 0, active))
|
|
})
|
|
.collect())
|
|
}
|
|
|
|
/// The accept tokens to push with AUTH, and whether to push at all (sec 4).
|
|
pub fn token_set(&self, account: &Account) -> Result<(u8, Vec<[u8; 32]>), Error> {
|
|
if !self.sync_ok()? {
|
|
return Ok((0, Vec::new()));
|
|
}
|
|
let db = self.db();
|
|
let mut stmt =
|
|
db.prepare("SELECT identity FROM accepted WHERE active = 1 ORDER BY added_at")?;
|
|
let tokens = stmt
|
|
.query_map([], |row| row.get::<_, Vec<u8>>(0))?
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
Ok((
|
|
1,
|
|
tokens
|
|
.iter()
|
|
.filter_map(|identity| {
|
|
let identity: [u8; KEY_LEN] = identity.as_slice().try_into().ok()?;
|
|
Some(account.token_for(&identity))
|
|
})
|
|
.collect(),
|
|
))
|
|
}
|
|
|
|
/// The address we know a signer by: a contact, or the Reply-To it signed
|
|
/// for itself. Naming a mailbox is not trusting a key, so nothing is
|
|
/// pinned here (sec 5.7, sec 8).
|
|
pub fn address_of(
|
|
&self,
|
|
sender: &[u8],
|
|
reply_to: Option<&str>,
|
|
) -> Result<Option<String>, Error> {
|
|
let db = self.db();
|
|
let mut stmt = db.prepare("SELECT address, identity FROM contacts")?;
|
|
let rows = stmt
|
|
.query_map([], |row| {
|
|
Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
|
|
})?
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
for (address, identity) in rows {
|
|
if ct_eq(&identity, sender) {
|
|
return Ok(Some(address));
|
|
}
|
|
}
|
|
if let Some(uri) = reply_to {
|
|
if let Ok(claimed) = Address::parse(uri) {
|
|
if claimed.identity.is_some_and(|key| ct_eq(&key, sender)) {
|
|
return Ok(Some(claimed.short()));
|
|
}
|
|
}
|
|
}
|
|
Ok(None)
|
|
}
|
|
|
|
/// Admits a correspondent to the main tier, freezing the identity the
|
|
/// token is derived from (sec 5.8): on conflict the identity is left
|
|
/// alone, so the token stays the one they already hold.
|
|
pub fn accept(&self, address: &str, identity: &[u8; KEY_LEN]) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT INTO accepted (address, identity, active, added_at) VALUES (?1, ?2, 1, ?3) \
|
|
ON CONFLICT (address) DO UPDATE SET active = 1",
|
|
params![address, identity, now()],
|
|
)
|
|
.map(|_| ())?;
|
|
self.set_sync_ok(true)?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn block(&self, address: &str) -> Result<bool, Error> {
|
|
let changed = self.db().execute(
|
|
"UPDATE accepted SET active = 0 WHERE address = ?1",
|
|
params![address],
|
|
)?;
|
|
Ok(changed > 0)
|
|
}
|
|
|
|
/// The acceptance state of an address, if it was ever accepted:
|
|
/// Some(true) main tier, Some(false) blocked, None never.
|
|
pub fn accepted_state(&self, address: &str) -> Result<Option<bool>, Error> {
|
|
Ok(self
|
|
.one::<i64>(
|
|
"SELECT active FROM accepted WHERE address = ?1",
|
|
&[&address],
|
|
)?
|
|
.map(|active| active != 0))
|
|
}
|
|
|
|
/// The identity frozen at acceptance, if this address is currently
|
|
/// accepted; the token for our own mailbox travels in the next message.
|
|
pub fn accepted_identity(&self, address: &str) -> Result<Option<[u8; KEY_LEN]>, Error> {
|
|
let identity: Option<Vec<u8>> = self.one(
|
|
"SELECT identity FROM accepted WHERE address = ?1 AND active = 1",
|
|
&[&address],
|
|
)?;
|
|
Ok(identity.and_then(|identity| identity.try_into().ok()))
|
|
}
|
|
|
|
/// A token a correspondent issued us, filed under their address (sec 5.8).
|
|
pub fn token_of(&self, address: &str) -> Result<Option<[u8; 32]>, Error> {
|
|
let token: Option<Vec<u8>> =
|
|
self.one("SELECT token FROM tokens WHERE address = ?1", &[&address])?;
|
|
Ok(token.and_then(|token| token.try_into().ok()))
|
|
}
|
|
|
|
pub fn save_token(&self, address: &str, token: &[u8; 32]) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT INTO tokens (address, token, seen_at) VALUES (?1, ?2, ?3) \
|
|
ON CONFLICT (address) DO UPDATE SET token = ?2, seen_at = ?3",
|
|
params![address, token, now()],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn seen(&self, id: &[u8; ID_LEN]) -> Result<bool, Error> {
|
|
Ok(self
|
|
.one::<i64>("SELECT 1 FROM seen WHERE id = ?1", &[id])?
|
|
.is_some())
|
|
}
|
|
|
|
pub fn store_inbox(
|
|
&self,
|
|
id: &[u8; ID_LEN],
|
|
envelope: &[u8],
|
|
received_at: i64,
|
|
requests_tier: bool,
|
|
kept: bool,
|
|
) -> Result<(), Error> {
|
|
let tier = if requests_tier {
|
|
TIER_REQUESTS
|
|
} else {
|
|
TIER_MAIN
|
|
};
|
|
self.db()
|
|
.execute(
|
|
"INSERT OR IGNORE INTO inbox (id, envelope, received_at, tier, kept) VALUES (?1, ?2, ?3, ?4, ?5)",
|
|
params![id, envelope, received_at, tier, kept as i64],
|
|
)
|
|
.map(|_| ())?;
|
|
self.db()
|
|
.execute(
|
|
"INSERT OR IGNORE INTO seen (id, at) VALUES (?1, ?2)",
|
|
params![id, received_at],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
/// A key a contact no longer uses, with when it stopped being current
|
|
/// (epoch ms). Written when a rotation displaces the key we held.
|
|
pub fn save_history(
|
|
&self,
|
|
address: &str,
|
|
identity: &[u8; KEY_LEN],
|
|
until: i64,
|
|
) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT OR IGNORE INTO history (address, identity, until) VALUES (?1, ?2, ?3)",
|
|
params![address, identity, until],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
/// The keys a contact used before its current one, oldest displacement
|
|
/// first: the local record that rotations happened (sec 8).
|
|
pub fn history(&self, address: &str) -> Result<Vec<([u8; KEY_LEN], i64)>, Error> {
|
|
let db = self.db();
|
|
let mut stmt =
|
|
db.prepare("SELECT identity, until FROM history WHERE address = ?1 ORDER BY until")?;
|
|
let rows = stmt
|
|
.query_map(params![address], |row| {
|
|
Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, i64>(1)?))
|
|
})?
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
Ok(rows
|
|
.into_iter()
|
|
.filter_map(|(identity, until)| Some((identity.try_into().ok()?, until)))
|
|
.collect())
|
|
}
|
|
|
|
pub fn store_sent(
|
|
&self,
|
|
id: &[u8; ID_LEN],
|
|
recipient: &str,
|
|
envelope: &[u8],
|
|
sent_at: i64,
|
|
) -> Result<(), Error> {
|
|
self.db()
|
|
.execute(
|
|
"INSERT OR IGNORE INTO sent (id, recipient, envelope, sent_at) VALUES (?1, ?2, ?3, ?4)",
|
|
params![id, recipient, envelope, sent_at],
|
|
)
|
|
.map(|_| ())?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Whether a message id is already stored in `folder` ("inbox" or "sent"),
|
|
/// so an import can skip rather than re-store.
|
|
pub fn stored(&self, folder: &str, id: &[u8; ID_LEN]) -> Result<bool, Error> {
|
|
let table = match folder {
|
|
"sent" => "sent",
|
|
_ => "inbox",
|
|
};
|
|
Ok(self
|
|
.one::<i64>(
|
|
&format!("SELECT 1 FROM {table} WHERE id = ?1"),
|
|
&[id],
|
|
)?
|
|
.is_some())
|
|
}
|
|
|
|
/// Removes one message from the local store — the reader's delete, which
|
|
/// the server-side DELETE (sec 6.1) knows nothing about. The id is marked
|
|
/// seen (sec 10), so a copy still sitting on the server is fetched again
|
|
/// without being re-stored: deleting locally must not resurrect the
|
|
/// message on the next fetch.
|
|
pub fn delete_local(&self, folder: &str, id: &[u8; ID_LEN]) -> Result<(), Error> {
|
|
let table = match folder {
|
|
"sent" => "sent",
|
|
_ => "inbox",
|
|
};
|
|
self.db()
|
|
.execute(&format!("DELETE FROM {table} WHERE id = ?1"), params![id])?;
|
|
self.db()
|
|
.execute("INSERT OR IGNORE INTO seen (id, at) VALUES (?1, ?2)", params![id, now()])?;
|
|
Ok(())
|
|
}
|
|
|
|
/// One folder, ordered by arrival or sending time. `folder` is "inbox",
|
|
/// "requests", "sent" or "all".
|
|
pub fn mail(&self, folder: &str) -> Result<Vec<Stored>, Error> {
|
|
let (sql, tier): (&str, Option<i8>) = match folder {
|
|
"sent" => (
|
|
"SELECT id, envelope, sent_at, recipient, 0 FROM sent ORDER BY sent_at, id",
|
|
None,
|
|
),
|
|
"all" => (
|
|
"SELECT id, envelope, received_at, NULL, kept FROM inbox ORDER BY received_at, id",
|
|
None,
|
|
),
|
|
"requests" => (
|
|
"SELECT id, envelope, received_at, NULL, kept FROM inbox \
|
|
WHERE tier = ?1 ORDER BY received_at, id",
|
|
Some(TIER_REQUESTS as i8),
|
|
),
|
|
_ => (
|
|
"SELECT id, envelope, received_at, NULL, kept FROM inbox \
|
|
WHERE tier = ?1 ORDER BY received_at, id",
|
|
Some(TIER_MAIN as i8),
|
|
),
|
|
};
|
|
let db = self.db();
|
|
let mut stmt = db.prepare(sql)?;
|
|
let rows = match tier {
|
|
Some(t) => stmt
|
|
.query_map(params![t], row_stored)?
|
|
.collect::<Result<Vec<_>, _>>()?,
|
|
None => stmt
|
|
.query_map([], row_stored)?
|
|
.collect::<Result<Vec<_>, _>>()?,
|
|
};
|
|
Ok(rows)
|
|
}
|
|
}
|
|
|
|
fn row_stored(row: &rusqlite::Row<'_>) -> rusqlite::Result<Stored> {
|
|
Ok(Stored {
|
|
id: row.get::<_, Vec<u8>>(0)?.try_into().unwrap(),
|
|
envelope: row.get(1)?,
|
|
at: row.get(2)?,
|
|
recipient: row.get(3)?,
|
|
kept: row.get::<_, i64>(4)? != 0,
|
|
})
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::account::{identity_seed, Account};
|
|
|
|
fn master() -> [u8; 32] {
|
|
core::array::from_fn(|i| i as u8)
|
|
}
|
|
|
|
/// A fresh scratch path per call: a monotonic counter, so two tests or
|
|
/// two runs cannot collide the way pid-and-second can.
|
|
fn scratch_path(tag: &str) -> std::path::PathBuf {
|
|
static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
|
|
let n = NEXT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
|
std::env::temp_dir().join(format!("fumi-store-{tag}-{n}-{}.db", std::process::id()))
|
|
}
|
|
|
|
fn temp_store(tag: &str) -> Store {
|
|
Store::open(&scratch_path(tag)).expect("open store")
|
|
}
|
|
|
|
#[test]
|
|
fn account_state_round_trip() {
|
|
let store = temp_store("account");
|
|
assert!(store.account().unwrap().is_none());
|
|
let addr = Address::parse("alice@example.org:1962").unwrap();
|
|
store.set_account(&addr).unwrap();
|
|
let back = store.account().unwrap().unwrap();
|
|
assert_eq!(back.short(), "alice@example.org:1962");
|
|
assert_eq!(back.port, 1962);
|
|
|
|
let rns = Address::parse("smol+rns://alice@8f2c1d0a7b6e5f4c3d2a1b0e9f8c7d6a").unwrap();
|
|
store.set_account(&rns).unwrap();
|
|
let back = store.account().unwrap().unwrap();
|
|
assert_eq!(back.scheme, crate::address::Scheme::Rns);
|
|
assert_eq!(back.short(), "alice@8f2c1d0a7b6e5f4c3d2a1b0e9f8c7d6a");
|
|
}
|
|
|
|
#[test]
|
|
fn rotation_index_and_cursor_persist() {
|
|
let store = temp_store("state");
|
|
assert_eq!(store.rotations().unwrap(), 0);
|
|
store.set_rotations(3).unwrap();
|
|
assert_eq!(store.rotations().unwrap(), 3);
|
|
assert_eq!(store.cursor().unwrap(), (0, [0u8; 32]));
|
|
store.set_cursor(17, &[9u8; 32]).unwrap();
|
|
assert_eq!(store.cursor().unwrap(), (17, [9u8; 32]));
|
|
// A short cursor blob (a fresh store's x'') reads as the start.
|
|
}
|
|
|
|
#[test]
|
|
fn sync_ok_gates_the_token_set() {
|
|
let store = temp_store("sync");
|
|
let account = Account::new(master(), 0).unwrap();
|
|
let key = Account::new(master(), 1).unwrap().me().pk();
|
|
store.accept("bob@example.org", &key).unwrap();
|
|
// Accepting turns sync on and pushes exactly one derived token.
|
|
assert!(store.sync_ok().unwrap());
|
|
let (sync, tokens) = store.token_set(&account).unwrap();
|
|
assert_eq!(sync, 1);
|
|
assert_eq!(tokens, vec![account.token_for(&key)]);
|
|
|
|
// A restored client must not erase the server's set (sec 4).
|
|
store.set_sync_ok(false).unwrap();
|
|
let (sync, tokens) = store.token_set(&account).unwrap();
|
|
assert_eq!(sync, 0);
|
|
assert!(tokens.is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn accept_freezes_identity_across_conflicts() {
|
|
let store = temp_store("accept");
|
|
let old = Account::new(master(), 1).unwrap().me().pk();
|
|
let new = Account::new(master(), 2).unwrap().me().pk();
|
|
store.accept("bob@example.org", &old).unwrap();
|
|
store.accept("bob@example.org", &new).unwrap();
|
|
assert_eq!(
|
|
store.accepted_identity("bob@example.org").unwrap(),
|
|
Some(old)
|
|
);
|
|
assert!(store.block("bob@example.org").unwrap());
|
|
assert_eq!(store.accepted_identity("bob@example.org").unwrap(), None);
|
|
// Blocking again is a harmless no-op; only a never-accepted address
|
|
// is an error, which the caller reports.
|
|
assert!(store.block("bob@example.org").unwrap());
|
|
assert!(!store.block("nobody@example.org").unwrap());
|
|
}
|
|
|
|
#[test]
|
|
fn address_of_matches_contacts_then_reply_to() {
|
|
let store = temp_store("address-of");
|
|
let bob = Account::new(master(), 1).unwrap().me().pk();
|
|
store.save_contact("bob@example.org", &bob, true).unwrap();
|
|
assert_eq!(
|
|
store.address_of(&bob, None).unwrap().as_deref(),
|
|
Some("bob@example.org")
|
|
);
|
|
// A Reply-To only names a signer whose key it carries (sec 5.7).
|
|
let carol = Account::new(master(), 2).unwrap().me().pk();
|
|
let uri = format!("smol://carol@example.org/{}", crate::crypto::b32(&carol));
|
|
assert_eq!(
|
|
store.address_of(&carol, Some(&uri)).unwrap().as_deref(),
|
|
Some("carol@example.org")
|
|
);
|
|
// Key mismatch: ordinary text, no address learned.
|
|
assert_eq!(
|
|
store
|
|
.address_of(
|
|
&carol,
|
|
Some(&format!(
|
|
"smol://mallory@example.org/{}",
|
|
crate::crypto::b32(&bob)
|
|
))
|
|
)
|
|
.unwrap(),
|
|
None
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn inbox_and_sent_store_sealed_and_deduplicate() {
|
|
let store = temp_store("mail");
|
|
let id = [3u8; 32];
|
|
store.store_inbox(&id, b"envelope", 100, false, false).unwrap();
|
|
// A replayed id is not re-stored (sec 10).
|
|
store.store_inbox(&id, b"envelope2", 100, true, false).unwrap();
|
|
assert!(store.seen(&id).unwrap());
|
|
let inbox = store.mail("inbox").unwrap();
|
|
assert_eq!(inbox.len(), 1);
|
|
assert_eq!(inbox[0].envelope, b"envelope");
|
|
assert_eq!(inbox[0].at, 100);
|
|
|
|
store
|
|
.store_sent(&id, "bob@example.org", b"sent-copy", 200)
|
|
.unwrap();
|
|
let sent = store.mail("sent").unwrap();
|
|
assert_eq!(sent.len(), 1);
|
|
assert_eq!(sent[0].recipient.as_deref(), Some("bob@example.org"));
|
|
}
|
|
|
|
#[test]
|
|
fn pins_and_contacts_round_trip() {
|
|
let store = temp_store("pins");
|
|
assert!(store.server_pin("example.org").unwrap().is_none());
|
|
store.pin_server("example.org", &[5u8; 32]).unwrap();
|
|
assert_eq!(store.server_pin("example.org").unwrap(), Some([5u8; 32]));
|
|
store.pin_server("example.org", &[6u8; 32]).unwrap();
|
|
assert_eq!(store.server_pin("example.org").unwrap(), Some([6u8; 32]));
|
|
|
|
let key = Account::new(master(), 1).unwrap().me().pk();
|
|
assert!(store.contact("bob@example.org").unwrap().is_none());
|
|
store.save_contact("bob@example.org", &key, false).unwrap();
|
|
assert_eq!(
|
|
store.contact("bob@example.org").unwrap(),
|
|
Some((key, false))
|
|
);
|
|
store.save_contact("bob@example.org", &key, true).unwrap();
|
|
assert_eq!(store.contact("bob@example.org").unwrap(), Some((key, true)));
|
|
}
|
|
|
|
#[test]
|
|
fn identity_seed_is_derived_not_stored() {
|
|
// A client must be able to re-derive every superseded key (sec 7).
|
|
let account = Account::new(master(), 4).unwrap();
|
|
assert_eq!(account.keys().len(), 5);
|
|
for (n, key) in account.keys().iter().enumerate() {
|
|
assert_eq!(
|
|
key.pk(),
|
|
Account::new(master(), n as u32).unwrap().me().pk()
|
|
);
|
|
let _ = identity_seed(&master(), n as u32);
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn schema_is_versioned_and_refuses_a_newer_store() {
|
|
let path = scratch_path("schema");
|
|
drop(Store::open(&path).expect("open store"));
|
|
let db = Connection::open(&path).unwrap();
|
|
let version: i32 = db.query_row("PRAGMA user_version", [], |row| row.get(0)).unwrap();
|
|
assert_eq!(version, SCHEMA_VERSION);
|
|
// A store from the future is refused, not misread.
|
|
db.pragma_update(None, "user_version", SCHEMA_VERSION + 1).unwrap();
|
|
drop(db);
|
|
match Store::open(&path) {
|
|
Err(Error::SchemaVersion { found, expected }) => {
|
|
assert_eq!((found, expected), (SCHEMA_VERSION + 1, SCHEMA_VERSION));
|
|
}
|
|
_ => panic!("a newer schema must be refused"),
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn v1_store_migrates_at_open() {
|
|
let path = scratch_path("v1");
|
|
// A schema-version-1 store: the inbox has no `kept` column and the
|
|
// history table does not exist yet.
|
|
let db = Connection::open(&path).unwrap();
|
|
db.execute_batch(
|
|
"CREATE TABLE state (
|
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
|
username TEXT, host TEXT, port INTEGER, scheme TEXT,
|
|
rotations INTEGER NOT NULL DEFAULT 0,
|
|
after_time INTEGER NOT NULL DEFAULT 0,
|
|
after_id BLOB NOT NULL DEFAULT x'',
|
|
sync_ok INTEGER NOT NULL DEFAULT 1);
|
|
INSERT INTO state (id) VALUES (1);
|
|
CREATE TABLE servers (
|
|
host TEXT PRIMARY KEY, static BLOB NOT NULL, pinned_at INTEGER NOT NULL);
|
|
CREATE TABLE contacts (
|
|
address TEXT PRIMARY KEY, identity BLOB NOT NULL,
|
|
verified INTEGER NOT NULL, seen_at INTEGER NOT NULL);
|
|
CREATE TABLE accepted (
|
|
address TEXT PRIMARY KEY, identity BLOB NOT NULL,
|
|
active INTEGER NOT NULL, added_at INTEGER NOT NULL);
|
|
CREATE TABLE tokens (
|
|
address TEXT PRIMARY KEY, token BLOB NOT NULL, seen_at INTEGER NOT NULL);
|
|
CREATE TABLE seen (id BLOB PRIMARY KEY, at INTEGER NOT NULL);
|
|
CREATE TABLE inbox (
|
|
id BLOB PRIMARY KEY, envelope BLOB NOT NULL,
|
|
received_at INTEGER NOT NULL, tier INTEGER NOT NULL);
|
|
CREATE TABLE sent (
|
|
id BLOB PRIMARY KEY, recipient TEXT NOT NULL,
|
|
envelope BLOB NOT NULL, sent_at INTEGER NOT NULL);
|
|
INSERT INTO inbox (id, envelope, received_at, tier)
|
|
VALUES (x'0101010101010101010101010101010101010101010101010101010101010101',
|
|
x'02', 7, 0);
|
|
PRAGMA user_version = 1;",
|
|
)
|
|
.unwrap();
|
|
drop(db);
|
|
|
|
let store = Store::open(&path).expect("v1 store migrates");
|
|
let db = Connection::open(&path).unwrap();
|
|
let version: i32 = db.query_row("PRAGMA user_version", [], |row| row.get(0)).unwrap();
|
|
drop(db);
|
|
assert_eq!(version, SCHEMA_VERSION);
|
|
// Existing rows survive with the column's default.
|
|
let mail = store.mail("all").unwrap();
|
|
assert_eq!(mail.len(), 1);
|
|
assert!(!mail[0].kept);
|
|
// The history table exists and round-trips.
|
|
store.save_history("alice@example.org", &[9u8; 32], 42).unwrap();
|
|
assert_eq!(store.history("alice@example.org").unwrap(), [([9u8; 32], 42)]);
|
|
}
|
|
|
|
#[test]
|
|
fn unpin_returns_to_trust_on_first_use() {
|
|
let store = temp_store("unpin");
|
|
store.pin_server("example.org", &[1u8; KEY_LEN]).unwrap();
|
|
store.unpin_server("example.org").unwrap();
|
|
assert_eq!(store.server_pin("example.org").unwrap(), None);
|
|
// Unpinning what was never pinned is not an error.
|
|
store.unpin_server("never.example.org").unwrap();
|
|
}
|
|
|
|
#[test]
|
|
fn kept_survives_the_folder_listing() {
|
|
let store = temp_store("kept");
|
|
store.store_inbox(&[1u8; 32], b"a", 1, false, true).unwrap();
|
|
store.store_inbox(&[2u8; 32], b"b", 2, false, false).unwrap();
|
|
let all = store.mail("all").unwrap();
|
|
assert!(all[0].kept && !all[1].kept);
|
|
}
|
|
|
|
#[test]
|
|
fn history_orders_by_displacement() {
|
|
let store = temp_store("history");
|
|
store.save_history("a@example.org", &[1u8; 32], 20).unwrap();
|
|
store.save_history("a@example.org", &[2u8; 32], 10).unwrap();
|
|
// A re-save of the same key is a no-op, like a non-rotation.
|
|
store.save_history("a@example.org", &[2u8; 32], 30).unwrap();
|
|
assert_eq!(
|
|
store.history("a@example.org").unwrap(),
|
|
[([2u8; 32], 10), ([1u8; 32], 20)]
|
|
);
|
|
assert!(store.history("b@example.org").unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn local_delete_marks_seen_so_a_refetch_does_not_restore() {
|
|
let store = temp_store("delete-local");
|
|
let id = [3u8; 32];
|
|
store.store_inbox(&id, b"envelope", 100, false, true).unwrap();
|
|
store.delete_local("inbox", &id).unwrap();
|
|
assert!(store.mail("all").unwrap().is_empty());
|
|
// Sec 10: the id is seen, so the fetch pipeline will not re-store
|
|
// the still-kept server copy. (The dedupe is the pipeline's seen
|
|
// check, not store_inbox's — a direct insert is a caller override.)
|
|
assert!(store.seen(&id).unwrap());
|
|
// Sent copies delete locally too, with the same protection.
|
|
let sent_id = [4u8; 32];
|
|
store.store_sent(&sent_id, "bob@example.org", b"sent", 50).unwrap();
|
|
store.delete_local("sent", &sent_id).unwrap();
|
|
assert!(store.mail("sent").unwrap().is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn in_memory_store_round_trips() {
|
|
let store = Store::open_in_memory().expect("memory store opens");
|
|
let addr = Address::parse("alice@example.org").unwrap();
|
|
store.set_account(&addr).unwrap();
|
|
assert_eq!(
|
|
store.account().unwrap().unwrap().short(),
|
|
"alice@example.org"
|
|
);
|
|
store.store_inbox(&[1u8; 32], b"envelope", 5, false, false).unwrap();
|
|
assert_eq!(store.mail("inbox").unwrap().len(), 1);
|
|
// The same handle through the path spelling.
|
|
let path = std::path::Path::new(":memory:");
|
|
assert!(Store::open(path).is_ok());
|
|
}
|
|
|
|
#[test]
|
|
fn store_is_shareable_across_threads() {
|
|
fn assert_send_sync<T: Send + Sync>() {}
|
|
assert_send_sync::<Store>();
|
|
let store = temp_store("threads");
|
|
std::thread::scope(|scope| {
|
|
for _ in 0..4 {
|
|
scope.spawn(|| {
|
|
store.set_rotations(7).unwrap();
|
|
store.rotations().unwrap();
|
|
});
|
|
}
|
|
});
|
|
assert_eq!(store.rotations().unwrap(), 7);
|
|
}
|
|
}
|