refactor: split the library into fumi-core and make it embeddable
This commit is contained in:
parent
d3cd0afb4f
commit
8fb5661b53
28 changed files with 975 additions and 724 deletions
724
core/src/store.rs
Normal file
724
core/src/store.rs
Normal file
|
|
@ -0,0 +1,724 @@
|
|||
//! 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);
|
||||
CREATE TABLE IF NOT EXISTS sent (
|
||||
id BLOB PRIMARY KEY, recipient TEXT NOT NULL,
|
||||
envelope BLOB NOT NULL, sent_at INTEGER NOT NULL);
|
||||
";
|
||||
|
||||
/// 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>,
|
||||
}
|
||||
|
||||
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 = 1;
|
||||
|
||||
impl Store {
|
||||
pub fn open(path: &Path) -> Result<Store, Error> {
|
||||
let db = Connection::open(path)?;
|
||||
// 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).
|
||||
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 != 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()))
|
||||
}
|
||||
|
||||
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 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,
|
||||
) -> 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) VALUES (?1, ?2, ?3, ?4)",
|
||||
params![id, envelope, received_at, tier],
|
||||
)
|
||||
.map(|_| ())?;
|
||||
self.db()
|
||||
.execute(
|
||||
"INSERT OR IGNORE INTO seen (id, at) VALUES (?1, ?2)",
|
||||
params![id, received_at],
|
||||
)
|
||||
.map(|_| ())?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
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(())
|
||||
}
|
||||
|
||||
/// 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 FROM sent ORDER BY sent_at, id",
|
||||
None,
|
||||
),
|
||||
"all" => (
|
||||
"SELECT id, envelope, received_at, NULL FROM inbox ORDER BY received_at, id",
|
||||
None,
|
||||
),
|
||||
"requests" => (
|
||||
"SELECT id, envelope, received_at, NULL FROM inbox \
|
||||
WHERE tier = ?1 ORDER BY received_at, id",
|
||||
Some(TIER_REQUESTS as i8),
|
||||
),
|
||||
_ => (
|
||||
"SELECT id, envelope, received_at, NULL 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)?,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::account::{identity_seed, Account};
|
||||
|
||||
fn master() -> [u8; 32] {
|
||||
core::array::from_fn(|i| i as u8)
|
||||
}
|
||||
|
||||
fn temp_store(tag: &str) -> Store {
|
||||
let path = std::env::temp_dir().join(format!(
|
||||
"fumi-store-{tag}-{}-{}.db",
|
||||
std::process::id(),
|
||||
now()
|
||||
));
|
||||
Store::open(&path).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).unwrap();
|
||||
// A replayed id is not re-stored (sec 10).
|
||||
store.store_inbox(&id, b"envelope2", 100, true).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 = std::env::temp_dir().join(format!(
|
||||
"fumi-store-schema-{}-{}.db",
|
||||
std::process::id(),
|
||||
now()
|
||||
));
|
||||
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 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);
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue