use std::time::Duration;
use anyhow::{anyhow, bail, Context, Result};
use chacha20poly1305::aead::{Aead, KeyInit, Payload as AeadPayload};
use chacha20poly1305::{Key, XChaCha20Poly1305, XNonce};
use ed25519_dalek::{Signature, Signer, SigningKey, VerifyingKey};
use rusqlite::{params, Connection, OptionalExtension};
use serde::{Deserialize, Serialize};
pub const RELAY: &str = "https://relay.snyvi.com";
pub fn relay() -> String {
std::env::var("SNYVI_RELAY")
.ok()
.filter(|s| !s.trim().is_empty())
.map(|s| s.trim().trim_end_matches('/').to_string())
.unwrap_or_else(|| RELAY.to_string())
}
pub const KEY_NAME: &str = "SNYVI_PEER_KEY";
pub const CODE_TTL: i64 = 10 * 60;
const VERSION: u8 = 1;
const SPAKE_ID: &[u8] = b"snyvi peer v1";
const HELLO_AAD: &[u8] = b"snyvi hello v1";
pub const NOTE_CHARS: usize = 200;
const HTTP_TIMEOUT: Duration = Duration::from_secs(40);
pub const SEND_MAX: usize = 8 * 1024 * 1024;
pub const FRAME_MAX: usize = SEND_MAX + 4096;
pub const INLINE_MAX: usize = 128 * 1024;
pub const PING_EVERY: Duration = Duration::from_secs(45);
pub const BACKOFF_MAX: Duration = Duration::from_secs(300);
pub const TAKEN_KEPT: i64 = 8 * 24 * 60 * 60;
#[derive(Clone)]
pub struct Identity {
sign: SigningKey,
boxk: x25519_dalek::StaticSecret,
}
impl std::fmt::Debug for Identity {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Identity")
.field("address", &self.address())
.finish()
}
}
impl Identity {
pub fn load_or_mint(secrets: &crate::secrets::Secrets) -> Result<Identity> {
if let Some(v) = secrets.value(0, KEY_NAME) {
if let Some(id) = Identity::from_secret(&v) {
return Ok(id);
}
bail!("the peer key kept as {KEY_NAME} is not one this snyvi can read; remove it to pair afresh");
}
let mut seed = [0u8; 64];
getrandom::fill(&mut seed)
.map_err(|e| anyhow!("reading random bytes for the peer key: {e}"))?;
let id = Identity::from_seed(&seed);
secrets
.keep(0, KEY_NAME, &b64(&seed))
.context("keeping the peer key")?;
Ok(id)
}
pub fn load(secrets: &crate::secrets::Secrets) -> Option<Identity> {
secrets
.value(0, KEY_NAME)
.and_then(|v| Identity::from_secret(&v))
}
fn from_secret(v: &str) -> Option<Identity> {
let bytes = unb64(v.trim())?;
(bytes.len() == 64).then(|| Identity::from_seed(bytes.as_slice().try_into().unwrap()))
}
pub fn from_seed(seed: &[u8; 64]) -> Identity {
let mut s = [0u8; 32];
s.copy_from_slice(&seed[..32]);
let mut b = [0u8; 32];
b.copy_from_slice(&seed[32..]);
Identity {
sign: SigningKey::from_bytes(&s),
boxk: x25519_dalek::StaticSecret::from(b),
}
}
pub fn address(&self) -> String {
b64(self.sign.verifying_key().as_bytes())
}
pub fn box_public(&self) -> [u8; 32] {
x25519_dalek::PublicKey::from(&self.boxk).to_bytes()
}
pub fn relay_auth(&self, method: &str, path: &str) -> String {
let secs = crate::store::now();
let msg = format!("snyvi-relay-v1\n{method}\n{path}\n{secs}");
let sig = self.sign.sign(msg.as_bytes());
format!("{secs}.{}", b64(&sig.to_bytes()))
}
fn box_key(&self, peer_box: &[u8; 32], from_me: bool) -> [u8; 32] {
let shared = self
.boxk
.diffie_hellman(&x25519_dalek::PublicKey::from(*peer_box));
let mine = self.box_public();
let (a, b) = if from_me {
(&mine, peer_box)
} else {
(peer_box, &mine)
};
let mut material = Vec::with_capacity(96);
material.extend_from_slice(shared.as_bytes());
material.extend_from_slice(a);
material.extend_from_slice(b);
blake3::derive_key("snyvi peer v1 box", &material)
}
}
const WORDS: &str = include_str!("peer_words.txt");
fn words() -> Vec<&'static str> {
WORDS
.lines()
.map(str::trim)
.filter(|w| !w.is_empty())
.collect()
}
const CHECK: &[u8] = b"bcdfghjkmnpqrstvwxz";
fn check(words: &str) -> String {
let h = blake3::hash(format!("snyvi code v1 {words}").as_bytes());
h.as_bytes()[..3]
.iter()
.map(|b| CHECK[*b as usize % CHECK.len()] as char)
.collect()
}
pub fn mint_code() -> Result<String> {
let list = words();
let mut r = [0u8; 8];
getrandom::fill(&mut r).map_err(|e| anyhow!("reading random bytes for the code: {e}"))?;
let n = u64::from_le_bytes(r);
let pick = |i: u32| list[((n >> (i * 16)) & 0xffff) as usize % list.len()];
let w = format!("{}-{}-{}", pick(0), pick(1), pick(2));
let c = check(&w);
Ok(format!("{w}-{c}"))
}
pub fn normalize_code(typed: &str) -> std::result::Result<String, &'static str> {
let parts: Vec<String> = typed
.trim()
.to_lowercase()
.split(|c: char| !c.is_ascii_alphanumeric())
.filter(|p| !p.is_empty())
.map(str::to_string)
.collect();
if parts.len() != 4 {
return Err("a code is three words and three letters, like ocean-ladder-fish-kpm");
}
let list = words();
if !parts[..3].iter().all(|p| list.contains(&p.as_str())) {
return Err("one of the words is not one a code is made of; ask for it again");
}
let w = parts[..3].join("-");
if check(&w) != parts[3] {
return Err("the last three letters do not match the words; one was heard wrong");
}
Ok(format!("{w}-{}", parts[3]))
}
pub fn room_of(code: &str) -> String {
blake3::hash(format!("snyvi room v1 {code}").as_bytes())
.to_hex()
.to_string()
}
const EMOJI: [&str; 64] = [
"🍎", "🍌", "🍒", "🍇", "🍋", "🍓", "🥝", "🍑", "🥕", "🌽", "🍄", "🌶️", "🥑", "🍞", "🧀", "🍪",
"🐶", "🐱", "🐭", "🐰", "🦊", "🐻", "🐼", "🐨", "🐸", "🐧", "🦉", "🐝", "🦋", "🐢", "🐙", "🦀",
"🌸", "🌻", "🌵", "🍀", "🌙", "⭐", "☀️", "🌈", "⚡", "❄️", "🔥", "💧", "🌊", "🍁", "🪐", "🌍",
"⚽", "🎸", "🎹", "🎲", "🎈", "🎁", "🔑", "🔔", "⏰", "📚", "✏️", "🧭", "🔭", "🚲", "⛵", "🚀",
];
pub fn emoji(a: &[u8; 32], b: &[u8; 32]) -> String {
let (x, y) = if a <= b { (a, b) } else { (b, a) };
let mut m = Vec::with_capacity(64);
m.extend_from_slice(x);
m.extend_from_slice(y);
let h = blake3::derive_key("snyvi peer v1 emoji", &m);
h[..4]
.iter()
.map(|i| EMOJI[(*i as usize) % 64])
.collect::<Vec<_>>()
.join(" ")
}
pub const CONTENT_V: u32 = 1;
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase", tag = "kind")]
pub enum Content {
Document {
title: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
lang: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
file: Option<String>,
name: String,
id: String,
#[serde(flatten)]
at: Folder,
},
Note {
text: String,
name: String,
#[serde(flatten)]
at: Folder,
},
#[serde(other)]
Other,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct Folder {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repo: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub remote: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub branch: Option<String>,
#[serde(default, skip_serializing_if = "v_unset")]
pub v: u32,
}
fn v_unset(n: &u32) -> bool {
*n == 0
}
impl Folder {
pub fn named(&self) -> bool {
self.repo.is_some() || self.remote.is_some()
}
}
pub fn safe_path(p: &str) -> Option<String> {
let p = p.replace('\\', "/");
if p.is_empty() || p.len() > 400 || p.starts_with('/') || p.chars().any(char::is_control) {
return None;
}
let parts: Vec<&str> = p.split('/').collect();
if parts.first().is_some_and(|f| f.contains(':')) {
return None;
}
if parts
.iter()
.any(|c| c.is_empty() || *c == "." || *c == "..")
{
return None;
}
Some(p)
}
fn pack(content: &Content, body: &[u8]) -> Vec<u8> {
let meta = serde_json::to_vec(content).expect("content is json");
let mut out = Vec::with_capacity(4 + meta.len() + body.len());
out.extend_from_slice(&(meta.len() as u32).to_le_bytes());
out.extend_from_slice(&meta);
out.extend_from_slice(body);
out
}
fn unpack(payload: &[u8]) -> Result<(Content, Vec<u8>)> {
if payload.len() < 4 {
bail!("a payload too short to carry anything");
}
let n = u32::from_le_bytes(payload[..4].try_into().unwrap()) as usize;
if payload.len() < 4 + n {
bail!("a payload whose header outruns it");
}
let content: Content =
serde_json::from_slice(&payload[4..4 + n]).context("the frame's content")?;
Ok((content, payload[4 + n..].to_vec()))
}
pub fn seal(me: &Identity, peer: &Peer, content: &Content, body: &[u8]) -> Result<Vec<u8>> {
let peer_box = peer.box_bytes()?;
let key = me.box_key(&peer_box, true);
let mut nonce = [0u8; 24];
getrandom::fill(&mut nonce).map_err(|e| anyhow!("reading random bytes for a nonce: {e}"))?;
let mut out = Vec::with_capacity(1 + 32 + 24 + body.len() + 256);
out.push(VERSION);
out.extend_from_slice(me.sign.verifying_key().as_bytes());
out.extend_from_slice(&nonce);
let cipher = XChaCha20Poly1305::new(Key::from_slice(&key));
let plain = pack(content, body);
let ct = cipher
.encrypt(
XNonce::from_slice(&nonce),
AeadPayload {
msg: &plain,
aad: &out[..33],
},
)
.map_err(|_| anyhow!("sealing the frame"))?;
out.extend_from_slice(&ct);
let sig = me.sign.sign(&out);
out.extend_from_slice(&sig.to_bytes());
Ok(out)
}
#[cfg(test)]
pub fn sender_of(frame: &[u8]) -> Option<String> {
(frame.len() >= 33 && frame[0] == VERSION).then(|| b64(&frame[1..33]))
}
pub fn open(me: &Identity, peer: &Peer, frame: &[u8]) -> Result<(Content, Vec<u8>)> {
if frame.len() < 1 + 32 + 24 + 16 + 64 {
bail!("a frame too short to be one");
}
if frame[0] != VERSION {
bail!("a frame of a version this snyvi does not read");
}
if b64(&frame[1..33]) != peer.sign_key {
bail!("a frame that names another sender");
}
let (signed, sig) = frame.split_at(frame.len() - 64);
let pinned = peer.verifying_key()?;
let sig: [u8; 64] = sig.try_into().unwrap();
let sig = Signature::from_bytes(&sig);
pinned
.verify_strict(signed, &sig)
.map_err(|_| anyhow!("a frame whose signature is not {}'s", peer.name))?;
let peer_box = peer.box_bytes()?;
let key = me.box_key(&peer_box, false);
let cipher = XChaCha20Poly1305::new(Key::from_slice(&key));
let plain = cipher
.decrypt(
XNonce::from_slice(&signed[33..57]),
AeadPayload {
msg: &signed[57..],
aad: &signed[..33],
},
)
.map_err(|_| anyhow!("a frame that does not open with {}'s key", peer.name))?;
unpack(&plain)
}
pub fn frame_id(doc_id: &str, peer_sign_key: &str) -> String {
blake3::hash(format!("snyvi frame v1 {doc_id} {peer_sign_key}").as_bytes())
.to_hex()
.to_string()
}
#[derive(Serialize, Deserialize)]
struct Hello {
sign: String,
#[serde(rename = "box")]
boxk: String,
name: String,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase", tag = "state")]
pub enum PairState {
Waiting,
Done {
emoji: String,
name: String,
peer: i64,
},
Failed {
why: String,
},
}
fn http() -> ureq::Agent {
ureq::config::Config::builder()
.timeout_global(Some(HTTP_TIMEOUT))
.http_status_as_error(false)
.build()
.new_agent()
}
pub fn pair(me: &Identity, code: &str, my_name: &str, until: i64) -> Result<(Peer, String)> {
use spake2::{Ed25519Group, Identity as SpakeId, Password, Spake2};
let room = room_of(code);
let mut side = [0u8; 16];
getrandom::fill(&mut side).map_err(|e| anyhow!("reading random bytes: {e}"))?;
let side = hex(&side);
let agent = http();
let base = relay();
let (state, msg) = Spake2::<Ed25519Group>::start_symmetric(
&Password::new(code.as_bytes()),
&SpakeId::new(SPAKE_ID),
);
let theirs = exchange(&agent, &base, &room, "spake", &side, &msg, until)?;
let key = state
.finish(&theirs)
.map_err(|_| anyhow!("the other side had a different code"))?;
let key: [u8; 32] = blake3::derive_key("snyvi peer v1 hello", &key);
let hello = Hello {
sign: me.address(),
boxk: b64(&me.box_public()),
name: my_name.trim().chars().take(60).collect(),
};
let mut nonce = [0u8; 24];
getrandom::fill(&mut nonce).map_err(|e| anyhow!("reading random bytes: {e}"))?;
let cipher = XChaCha20Poly1305::new(Key::from_slice(&key));
let mut sealed = nonce.to_vec();
sealed.extend(
cipher
.encrypt(
XNonce::from_slice(&nonce),
AeadPayload {
msg: &serde_json::to_vec(&hello)?,
aad: HELLO_AAD,
},
)
.map_err(|_| anyhow!("sealing the hello"))?,
);
let theirs = exchange(&agent, &base, &room, "hello", &side, &sealed, until)?;
if theirs.len() < 24 + 16 {
bail!("the other side's hello is not one");
}
let plain = cipher
.decrypt(
XNonce::from_slice(&theirs[..24]),
AeadPayload {
msg: &theirs[24..],
aad: HELLO_AAD,
},
)
.map_err(|_| anyhow!("the other side's hello does not open: a different code"))?;
let h: Hello = serde_json::from_slice(&plain).context("the other side's hello")?;
let peer = Peer {
id: 0,
sign_key: h.sign,
box_key: h.boxk,
name: if h.name.trim().is_empty() {
"a friend".to_string()
} else {
h.name.trim().to_string()
},
paired_at: 0,
muted: false,
removed_at: 0,
last_from: 0,
last_to: 0,
desk_id: 0,
};
let their_sign = peer.verifying_key()?;
peer.box_bytes()?;
if peer.sign_key == me.address() {
bail!("that is this snyvi's own code");
}
let glyphs = emoji(me.sign.verifying_key().as_bytes(), their_sign.as_bytes());
let _ = inbox(me);
Ok((peer, glyphs))
}
fn exchange(
agent: &ureq::Agent,
base: &str,
room: &str,
stage: &str,
side: &str,
mine: &[u8],
until: i64,
) -> Result<Vec<u8>> {
let url = format!("{base}/room/{room}/{stage}");
let mut resp = agent
.put(&url)
.header("x-snyvi-side", side)
.send(mine)
.context("reaching the relay")?;
match resp.status().as_u16() {
200 => return read_all(&mut resp),
202 => {}
s => return Err(pair_refused(s)),
}
let mut bell = true;
loop {
let left = until - crate::store::now();
if left <= 0 {
bail!("the code ran out before the other side typed it");
}
let wait = if bell {
match ring_wait(base, room, stage, side, left.min(25)) {
Bell::Rang | Bell::Quiet => 0,
Bell::Absent => {
bell = false;
left.min(25)
}
}
} else {
left.min(25)
};
let asked = std::time::Instant::now();
let mut resp = agent
.get(&format!("{url}?wait={wait}"))
.header("x-snyvi-side", side)
.call()
.context("reaching the relay")?;
match resp.status().as_u16() {
200 => return read_all(&mut resp),
204 if wait > 0 && asked.elapsed() < Duration::from_secs(1) => {
std::thread::sleep(Duration::from_secs(1))
}
204 => {}
s => return Err(pair_refused(s)),
}
}
}
fn pair_refused(status: u16) -> anyhow::Error {
match status {
409 => anyhow!("someone else already used this code"),
410 => anyhow!("this code was already used"),
429 => anyhow!("{BUSY}; try again in a minute"),
s => anyhow!("the relay answered {s} to the pairing"),
}
}
#[derive(Debug, PartialEq, Eq)]
enum Bell {
Rang,
Quiet,
Absent,
}
fn ring_wait(base: &str, room: &str, stage: &str, side: &str, secs: i64) -> Bell {
let Ok(rt) = tokio::runtime::Handle::try_current() else {
return Bell::Absent;
};
let url = format!("{}/room/{room}/ws?side={side}", ws_base(base));
rt.block_on(async move {
use futures_util::StreamExt;
use tokio_tungstenite::tungstenite::Message;
let connect = tokio::time::timeout(
Duration::from_secs(10),
tokio_tungstenite::connect_async(url.as_str()),
);
let Ok(Ok((mut socket, _))) = connect.await else {
return Bell::Absent;
};
let ring = async {
while let Some(Ok(msg)) = socket.next().await {
if matches!(&msg, Message::Text(t) if rings_for(t.as_str()) == Some(stage)) {
return Bell::Rang;
}
}
Bell::Quiet
};
let bell = tokio::time::timeout(Duration::from_secs(secs.max(1) as u64), ring)
.await
.unwrap_or(Bell::Quiet);
let _ = socket.close(None).await;
bell
})
}
fn rings_for(text: &str) -> Option<&str> {
let v: serde_json::Value = serde_json::from_str(text).ok()?;
match v.get("ready")?.as_str()? {
"spake" => Some("spake"),
"hello" => Some("hello"),
_ => None,
}
}
fn read_all(resp: &mut ureq::http::Response<ureq::Body>) -> Result<Vec<u8>> {
Ok(resp
.body_mut()
.with_config()
.limit(FRAME_MAX as u64)
.read_to_vec()?)
}
#[derive(Clone, Debug, Deserialize)]
pub struct Waiting {
pub id: String,
pub sender: String,
pub size: u64,
}
#[derive(Debug, PartialEq, Eq)]
pub enum Deposit {
Sent,
Later(&'static str),
}
pub const BUSY: &str = "the relay is busy";
pub fn deposit(me: &Identity, to: &str, id: &str, frame: &[u8]) -> Result<Deposit> {
let path = format!("/to/{to}");
let mut resp = http()
.post(&format!("{}{path}", relay()))
.header("x-snyvi-id", id)
.header("x-snyvi-from", &me.address())
.header("x-snyvi-auth", &me.relay_auth("POST", &path))
.send(frame)
.context("reaching the relay")?;
let status = resp.status().as_u16();
if status == 200 || status == 201 {
return Ok(Deposit::Sent);
}
let body = String::from_utf8_lossy(&read_all(&mut resp)?)
.trim()
.to_string();
match later(status, &body) {
Some(why) => Ok(Deposit::Later(why)),
None => bail!("the relay answered {status}: {body}"),
}
}
fn later(status: u16, body: &str) -> Option<&'static str> {
match status {
429 if body.contains("busy") => Some(BUSY),
429 => Some("their mailbox is full"),
404 => Some("their snyvi has not been online for a while"),
_ => None,
}
}
pub fn inbox(me: &Identity) -> Result<Vec<Waiting>> {
let path = format!("/inbox/{}", me.address());
let mut resp = http()
.get(&format!("{}{path}", relay()))
.header("x-snyvi-auth", &me.relay_auth("GET", &path))
.call()
.context("reaching the relay")?;
if resp.status().as_u16() != 200 {
bail!("the relay answered {} to the inbox", resp.status());
}
#[derive(Deserialize)]
struct List {
frames: Vec<Waiting>,
}
let l: List = resp.body_mut().read_json().context("the relay's listing")?;
Ok(l.frames)
}
pub fn fetch(me: &Identity, id: &str) -> Result<Vec<u8>> {
let path = format!("/inbox/{}/{id}", me.address());
let mut resp = http()
.get(&format!("{}{path}", relay()))
.header("x-snyvi-auth", &me.relay_auth("GET", &path))
.call()
.context("reaching the relay")?;
if resp.status().as_u16() != 200 {
bail!("the relay answered {} to a frame", resp.status());
}
read_all(&mut resp)
}
pub fn ack(me: &Identity, id: &str) -> Result<()> {
let path = format!("/inbox/{}/{id}", me.address());
let resp = http()
.delete(&format!("{}{path}", relay()))
.header("x-snyvi-auth", &me.relay_auth("DELETE", &path))
.call()
.context("reaching the relay")?;
match resp.status().as_u16() {
204 | 404 => Ok(()),
s => bail!("the relay answered {s} to an ack"),
}
}
pub fn relay_ws(address: &str) -> String {
ws_of(&relay(), address)
}
fn ws_of(relay: &str, address: &str) -> String {
format!("{}/inbox/{address}", ws_base(relay))
}
fn ws_base(relay: &str) -> String {
if let Some(rest) = relay.strip_prefix("https://") {
format!("wss://{rest}")
} else if let Some(rest) = relay.strip_prefix("http://") {
format!("ws://{rest}")
} else {
relay.to_string()
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Pushed {
Frame(Waiting),
}
pub fn ack_message(id: &str) -> String {
serde_json::json!({ "ack": id }).to_string()
}
pub fn backoff(attempt: u32) -> Duration {
let base = Duration::from_secs(1u64 << attempt.clamp(0, 20)).min(BACKOFF_MAX);
let mut b = [0u8; 1];
let _ = getrandom::fill(&mut b);
base + base.mul_f64(f64::from(b[0]) / 255.0 * 0.3)
}
#[derive(Clone, Debug, Default, Serialize, PartialEq, Eq)]
pub struct Peer {
pub id: i64,
pub sign_key: String,
pub box_key: String,
pub name: String,
pub paired_at: i64,
pub muted: bool,
#[serde(skip_serializing_if = "is_zero")]
pub removed_at: i64,
pub last_from: i64,
pub last_to: i64,
pub desk_id: i64,
}
fn is_zero(n: &i64) -> bool {
*n == 0
}
impl Peer {
fn verifying_key(&self) -> Result<VerifyingKey> {
let b = unb64(&self.sign_key)
.filter(|b| b.len() == 32)
.ok_or_else(|| anyhow!("a pinned key that is not one"))?;
VerifyingKey::from_bytes(b.as_slice().try_into().unwrap())
.map_err(|_| anyhow!("a pinned key that is not one"))
}
fn box_bytes(&self) -> Result<[u8; 32]> {
let b = unb64(&self.box_key)
.filter(|b| b.len() == 32)
.ok_or_else(|| anyhow!("a pinned box key that is not one"))?;
Ok(b.as_slice().try_into().unwrap())
}
pub fn project_root(&self) -> String {
format!("peer:{}", self.sign_key)
}
pub fn project_name(&self) -> String {
format!("From {}", self.name)
}
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
pub struct PeerNote {
pub id: i64,
pub peer_id: i64,
pub from: String,
pub text: String,
pub arrived_at: i64,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
pub struct Offer {
pub id: i64,
pub peer_id: i64,
pub to: String,
pub doc_id: String,
pub title: String,
pub by: String,
pub pane: String,
pub offered_at: i64,
}
pub const SCHEMA: &str = r#"
CREATE TABLE IF NOT EXISTS peers (
id INTEGER PRIMARY KEY,
sign_key TEXT NOT NULL UNIQUE,
box_key TEXT NOT NULL,
name TEXT NOT NULL,
paired_at INTEGER NOT NULL,
muted INTEGER NOT NULL DEFAULT 0,
removed_at INTEGER NOT NULL DEFAULT 0,
last_from INTEGER NOT NULL DEFAULT 0,
last_to INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS peer_outbox (
id TEXT PRIMARY KEY,
peer_id INTEGER NOT NULL REFERENCES peers(id),
doc_id TEXT NOT NULL,
queued_at INTEGER NOT NULL,
sent_at INTEGER NOT NULL DEFAULT 0,
tries INTEGER NOT NULL DEFAULT 0,
error TEXT NOT NULL DEFAULT ''
);
CREATE TABLE IF NOT EXISTS peer_notes (
id INTEGER PRIMARY KEY,
peer_id INTEGER NOT NULL REFERENCES peers(id),
text TEXT NOT NULL,
arrived_at INTEGER NOT NULL,
taken_at INTEGER NOT NULL DEFAULT 0,
removed_at INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS peer_offers (
id INTEGER PRIMARY KEY,
peer_id INTEGER NOT NULL REFERENCES peers(id),
doc_id TEXT NOT NULL,
pane TEXT NOT NULL DEFAULT '',
by TEXT NOT NULL DEFAULT '',
offered_at INTEGER NOT NULL,
answered_at INTEGER NOT NULL DEFAULT 0,
sent INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS peer_taken (
id TEXT PRIMARY KEY,
at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS peer_held (
id TEXT PRIMARY KEY,
peer_id INTEGER NOT NULL REFERENCES peers(id),
bytes BLOB NOT NULL,
held_at INTEGER NOT NULL
);
"#;
pub const COLUMNS_1_22: [&str; 4] = [
"ALTER TABLE projects ADD COLUMN repo TEXT NOT NULL DEFAULT ''",
"ALTER TABLE projects ADD COLUMN remote TEXT NOT NULL DEFAULT ''",
"ALTER TABLE projects ADD COLUMN printed_at INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE docs ADD COLUMN peer_key TEXT NOT NULL DEFAULT ''",
];
pub const HELD_KEPT: i64 = 30 * 86_400;
pub const COLUMNS_1_19: [&str; 2] = [
"ALTER TABLE peers ADD COLUMN desk_id INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE peer_outbox ADD COLUMN text TEXT NOT NULL DEFAULT ''",
];
const PEER_COLS: &str =
"id, sign_key, box_key, name, paired_at, muted, removed_at, last_from, last_to, desk_id";
fn row_peer(r: &rusqlite::Row) -> rusqlite::Result<Peer> {
Ok(Peer {
id: r.get(0)?,
sign_key: r.get(1)?,
box_key: r.get(2)?,
name: r.get(3)?,
paired_at: r.get(4)?,
muted: r.get::<_, i64>(5)? != 0,
removed_at: r.get(6)?,
last_from: r.get(7)?,
last_to: r.get(8)?,
desk_id: r.get(9)?,
})
}
pub fn list(conn: &Connection) -> Result<Vec<Peer>> {
let rows = conn
.prepare(&format!(
"SELECT {PEER_COLS} FROM peers ORDER BY removed_at != 0, LOWER(name), id"
))?
.query_map([], row_peer)?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn get(conn: &Connection, id: i64) -> Result<Option<Peer>> {
Ok(conn
.query_row(
&format!("SELECT {PEER_COLS} FROM peers WHERE id = ?1"),
params![id],
row_peer,
)
.optional()?)
}
pub fn by_sign_key(conn: &Connection, key: &str) -> Result<Option<Peer>> {
Ok(conn
.query_row(
&format!("SELECT {PEER_COLS} FROM peers WHERE sign_key = ?1"),
params![key],
row_peer,
)
.optional()?)
}
pub fn by_name(conn: &Connection, name: &str) -> Result<Option<Peer>> {
Ok(conn
.query_row(
&format!("SELECT {PEER_COLS} FROM peers WHERE removed_at = 0 AND LOWER(name) = LOWER(?1) ORDER BY id LIMIT 1"),
params![name.trim()],
row_peer,
)
.optional()?)
}
pub fn pin(conn: &Connection, p: &Peer, now: i64) -> Result<Peer> {
conn.execute(
"INSERT INTO peers(sign_key, box_key, name, paired_at) VALUES(?1, ?2, ?3, ?4)
ON CONFLICT(sign_key) DO UPDATE SET box_key = excluded.box_key, paired_at = excluded.paired_at, removed_at = 0",
params![p.sign_key, p.box_key, p.name, now],
)?;
Ok(by_sign_key(conn, &p.sign_key)?.expect("just pinned"))
}
pub fn rename(conn: &Connection, id: i64, name: &str) -> Result<bool> {
let name: String = name.trim().chars().take(60).collect();
if name.is_empty() {
return Ok(false);
}
Ok(conn.execute(
"UPDATE peers SET name = ?2 WHERE id = ?1",
params![id, name],
)? > 0)
}
pub fn mute(conn: &Connection, id: i64, muted: bool) -> Result<bool> {
Ok(conn.execute(
"UPDATE peers SET muted = ?2 WHERE id = ?1",
params![id, muted as i64],
)? > 0)
}
pub fn remove(conn: &Connection, id: i64, now: i64) -> Result<bool> {
Ok(conn.execute(
"UPDATE peers SET removed_at = ?2 WHERE id = ?1 AND removed_at = 0",
params![id, now],
)? > 0)
}
pub fn restore(conn: &Connection, id: i64) -> Result<bool> {
Ok(conn.execute(
"UPDATE peers SET removed_at = 0 WHERE id = ?1 AND removed_at != 0",
params![id],
)? > 0)
}
pub fn set_desk(conn: &Connection, id: i64, desk_id: i64) -> Result<bool> {
Ok(conn.execute(
"UPDATE peers SET desk_id = ?2 WHERE id = ?1",
params![id, desk_id.max(0)],
)? > 0)
}
pub fn touch(conn: &Connection, id: i64, from: bool, now: i64) -> Result<()> {
let col = if from { "last_from" } else { "last_to" };
conn.execute(
&format!("UPDATE peers SET {col} = ?2 WHERE id = ?1"),
params![id, now],
)?;
Ok(())
}
pub fn queue(conn: &Connection, peer: &Peer, doc_id: &str, now: i64) -> Result<String> {
let id = frame_id(doc_id, &peer.sign_key);
conn.execute(
"INSERT INTO peer_outbox(id, peer_id, doc_id, queued_at) VALUES(?1, ?2, ?3, ?4)
ON CONFLICT(id) DO UPDATE SET queued_at = excluded.queued_at, sent_at = 0, tries = 0, error = ''",
params![id, peer.id, doc_id, now],
)?;
Ok(id)
}
pub fn queue_note(conn: &Connection, peer: &Peer, text: &str, now: i64) -> Result<String> {
let text: String = text.trim().chars().take(NOTE_CHARS).collect();
let mut nonce = [0u8; 16];
getrandom::fill(&mut nonce)
.map_err(|e| anyhow!("reading random bytes for a line's id: {e}"))?;
let id = blake3::hash(&nonce).to_hex().to_string();
conn.execute(
"INSERT INTO peer_outbox(id, peer_id, doc_id, text, queued_at) VALUES(?1, ?2, '', ?3, ?4)",
params![id, peer.id, text, now],
)?;
Ok(id)
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Unsent {
pub id: String,
pub peer_id: i64,
pub doc_id: String,
pub text: String,
pub tries: i64,
}
pub fn unsent(conn: &Connection) -> Result<Vec<Unsent>> {
let rows = conn
.prepare(
"SELECT o.id, o.peer_id, o.doc_id, o.text, o.tries FROM peer_outbox o JOIN peers p ON p.id = o.peer_id
WHERE o.sent_at = 0 AND p.removed_at = 0 ORDER BY o.queued_at, o.id",
)?
.query_map([], |r| {
Ok(Unsent {
id: r.get(0)?,
peer_id: r.get(1)?,
doc_id: r.get(2)?,
text: r.get(3)?,
tries: r.get(4)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn sent(conn: &Connection, id: &str, now: i64) -> Result<()> {
conn.execute(
"UPDATE peer_outbox SET sent_at = ?2, error = '' WHERE id = ?1",
params![id, now],
)?;
Ok(())
}
pub fn failed(conn: &Connection, id: &str, why: &str) -> Result<()> {
conn.execute(
"UPDATE peer_outbox SET tries = tries + 1, error = ?2 WHERE id = ?1",
params![id, why.chars().take(200).collect::<String>()],
)?;
Ok(())
}
pub fn waiting(conn: &Connection, id: &str, why: &str) -> Result<()> {
conn.execute(
"UPDATE peer_outbox SET error = ?2 WHERE id = ?1",
params![id, why.chars().take(200).collect::<String>()],
)?;
Ok(())
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
pub struct Outgoing {
pub id: String,
pub peer_id: i64,
pub what: String,
pub queued_at: i64,
pub tries: i64,
pub error: String,
}
pub fn outgoing(conn: &Connection) -> Result<Vec<Outgoing>> {
let rows = conn
.prepare(
"SELECT o.id, o.peer_id, COALESCE(d.title, o.text), o.queued_at, o.tries, o.error
FROM peer_outbox o JOIN peers p ON p.id = o.peer_id LEFT JOIN docs d ON d.id = o.doc_id AND o.doc_id != ''
WHERE o.sent_at = 0 AND p.removed_at = 0 ORDER BY o.queued_at, o.id",
)?
.query_map([], |r| {
Ok(Outgoing {
id: r.get(0)?,
peer_id: r.get(1)?,
what: r.get(2)?,
queued_at: r.get(3)?,
tries: r.get(4)?,
error: r.get(5)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn retry(conn: &Connection, id: &str) -> Result<bool> {
Ok(conn.execute(
"UPDATE peer_outbox SET tries = 0, error = '' WHERE id = ?1 AND sent_at = 0",
params![id],
)? > 0)
}
pub fn hold(conn: &Connection, id: &str, peer_id: i64, bytes: &[u8], now: i64) -> Result<()> {
conn.execute(
"INSERT OR IGNORE INTO peer_held(id, peer_id, bytes, held_at) VALUES(?1, ?2, ?3, ?4)",
params![id, peer_id, bytes, now],
)?;
Ok(())
}
pub fn held(conn: &Connection) -> Result<Vec<(String, i64, Vec<u8>)>> {
let rows = conn
.prepare("SELECT id, peer_id, bytes FROM peer_held ORDER BY held_at, id")?
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn unhold(conn: &Connection, id: &str) -> Result<()> {
conn.execute("DELETE FROM peer_held WHERE id = ?1", params![id])?;
Ok(())
}
pub fn prune_held(conn: &Connection, before: i64) -> Result<usize> {
Ok(conn.execute("DELETE FROM peer_held WHERE held_at < ?1", params![before])?)
}
pub fn note_arrived(conn: &Connection, peer_id: i64, text: &str, now: i64) -> Result<i64> {
let text: String = text.trim().chars().take(NOTE_CHARS).collect();
conn.execute(
"INSERT INTO peer_notes(peer_id, text, arrived_at) VALUES(?1, ?2, ?3)",
params![peer_id, text, now],
)?;
Ok(conn.last_insert_rowid())
}
pub fn notes_waiting(conn: &Connection) -> Result<Vec<PeerNote>> {
let rows = conn
.prepare(
"SELECT n.id, n.peer_id, p.name, n.text, n.arrived_at FROM peer_notes n JOIN peers p ON p.id = n.peer_id
WHERE n.taken_at = 0 AND n.removed_at = 0 AND p.removed_at = 0 ORDER BY n.arrived_at, n.id",
)?
.query_map([], |r| {
Ok(PeerNote {
id: r.get(0)?,
peer_id: r.get(1)?,
from: r.get(2)?,
text: r.get(3)?,
arrived_at: r.get(4)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn settle_note(conn: &Connection, id: i64, what: &str, now: i64) -> Result<bool> {
let n = match what {
"taken" => conn.execute(
"UPDATE peer_notes SET taken_at = ?2 WHERE id = ?1 AND taken_at = 0",
params![id, now],
)?,
"remove" => conn.execute(
"UPDATE peer_notes SET removed_at = ?2 WHERE id = ?1 AND removed_at = 0",
params![id, now],
)?,
"restore" => conn.execute(
"UPDATE peer_notes SET removed_at = 0, taken_at = 0 WHERE id = ?1 AND (removed_at != 0 OR taken_at != 0)",
params![id],
)?,
_ => return Ok(false),
};
Ok(n > 0)
}
pub fn offer(
conn: &Connection,
peer_id: i64,
doc_id: &str,
pane: &str,
by: &str,
now: i64,
) -> Result<i64> {
conn.execute(
"INSERT INTO peer_offers(peer_id, doc_id, pane, by, offered_at) VALUES(?1, ?2, ?3, ?4, ?5)",
params![
peer_id,
doc_id,
pane,
by.chars().take(60).collect::<String>(),
now
],
)?;
Ok(conn.last_insert_rowid())
}
pub fn offers_open(conn: &Connection) -> Result<Vec<Offer>> {
let rows = conn
.prepare(
"SELECT o.id, o.peer_id, p.name, o.doc_id, COALESCE(d.title, ''), o.by, o.pane, o.offered_at
FROM peer_offers o JOIN peers p ON p.id = o.peer_id LEFT JOIN docs d ON d.id = o.doc_id
WHERE o.answered_at = 0 AND p.removed_at = 0 ORDER BY o.offered_at, o.id",
)?
.query_map([], |r| {
Ok(Offer {
id: r.get(0)?,
peer_id: r.get(1)?,
to: r.get(2)?,
doc_id: r.get(3)?,
title: r.get(4)?,
by: r.get(5)?,
pane: r.get(6)?,
offered_at: r.get(7)?,
})
})?
.collect::<std::result::Result<_, _>>()?;
Ok(rows)
}
pub fn offer_get(conn: &Connection, id: i64) -> Result<Option<Offer>> {
Ok(offers_open(conn)?.into_iter().find(|o| o.id == id))
}
pub fn answer_offer(conn: &Connection, id: i64, sent: bool, now: i64) -> Result<bool> {
Ok(conn.execute(
"UPDATE peer_offers SET answered_at = ?2, sent = ?3 WHERE id = ?1 AND answered_at = 0",
params![id, now, sent as i64],
)? > 0)
}
pub fn reopen_offer(conn: &Connection, id: i64) -> Result<bool> {
Ok(conn.execute(
"UPDATE peer_offers SET answered_at = 0 WHERE id = ?1 AND answered_at != 0 AND sent = 0",
params![id],
)? > 0)
}
pub fn drop_offers_of(conn: &Connection, pane: &str, now: i64) -> Result<usize> {
Ok(conn.execute(
"UPDATE peer_offers SET answered_at = ?2, sent = 0 WHERE pane = ?1 AND answered_at = 0",
params![pane, now],
)?)
}
pub fn taken(conn: &Connection, id: &str) -> Result<bool> {
Ok(conn
.query_row(
"SELECT 1 FROM peer_taken WHERE id = ?1",
params![id],
|_| Ok(()),
)
.optional()?
.is_some())
}
pub fn take(conn: &Connection, id: &str, now: i64) -> Result<()> {
conn.execute(
"INSERT OR REPLACE INTO peer_taken (id, at) VALUES (?1, ?2)",
params![id, now],
)?;
Ok(())
}
pub fn prune_taken(conn: &Connection, before: i64) -> Result<usize> {
Ok(conn.execute("DELETE FROM peer_taken WHERE at < ?1", params![before])?)
}
pub fn clear(conn: &Connection) -> Result<()> {
conn.execute_batch("DELETE FROM peer_held; DELETE FROM peer_taken; DELETE FROM peer_offers; DELETE FROM peer_notes; DELETE FROM peer_outbox; DELETE FROM peers;")?;
Ok(())
}
pub fn b64(bytes: &[u8]) -> String {
const A: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_";
let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
for chunk in bytes.chunks(3) {
let n = chunk.iter().fold(0u32, |acc, b| (acc << 8) | *b as u32) << (8 * (3 - chunk.len()));
for i in 0..chunk.len() + 1 {
out.push(A[((n >> (18 - 6 * i)) & 63) as usize] as char);
}
}
out
}
pub fn unb64(s: &str) -> Option<Vec<u8>> {
let val = |c: u8| -> Option<u32> {
Some(match c {
b'A'..=b'Z' => c - b'A',
b'a'..=b'z' => c - b'a' + 26,
b'0'..=b'9' => c - b'0' + 52,
b'-' => 62,
b'_' => 63,
_ => return None,
} as u32)
};
let s = s.trim_end_matches('=').as_bytes();
if s.len() % 4 == 1 {
return None;
}
let mut out = Vec::with_capacity(s.len() * 3 / 4);
for chunk in s.chunks(4) {
let mut n = 0u32;
for &c in chunk {
n = (n << 6) | val(c)?;
}
n <<= 6 * (4 - chunk.len());
for i in 0..chunk.len() - 1 {
out.push(((n >> (16 - 8 * i)) & 255) as u8);
}
}
Some(out)
}
fn hex(bytes: &[u8]) -> String {
bytes.iter().map(|b| format!("{b:02x}")).collect()
}
#[cfg(test)]
mod tests {
use super::*;
fn two() -> (Identity, Identity) {
let mut a = [7u8; 64];
a[0] = 1;
let mut b = [9u8; 64];
b[0] = 2;
(Identity::from_seed(&a), Identity::from_seed(&b))
}
fn as_peer(id: &Identity, name: &str) -> Peer {
Peer {
id: 1,
sign_key: id.address(),
box_key: b64(&id.box_public()),
name: name.into(),
paired_at: 0,
muted: false,
removed_at: 0,
last_from: 0,
last_to: 0,
desk_id: 0,
}
}
#[test]
fn base64url_round_trips_and_is_what_the_relay_spells() {
for n in 0..70 {
let v: Vec<u8> = (0..n).map(|i| (i * 37 % 256) as u8).collect();
let s = b64(&v);
assert!(!s.contains('='), "{s}");
assert_eq!(unb64(&s).unwrap(), v, "{n}");
}
assert_eq!(b64(b"hello"), "aGVsbG8");
assert_eq!(unb64("aGVsbG8=").unwrap(), b"hello");
assert!(unb64("a").is_none());
assert!(unb64("a+b/").is_none());
assert_eq!(b64(&[0u8; 32]).len(), 43, "an address is 43 characters");
}
#[test]
fn a_frame_opens_for_its_recipient_and_for_nobody_else() {
let (sunny, trapti) = two();
let content = Content::Document {
title: "Plan for the garden".into(),
lang: Some("md".into()),
file: None,
name: "Sunny".into(),
id: "abc123".into(),
at: Folder::default(),
};
let frame = seal(
&sunny,
&as_peer(&trapti, "Trapti"),
&content,
b"# the plan\n",
)
.unwrap();
assert_eq!(frame[0], VERSION);
assert_eq!(sender_of(&frame).as_deref(), Some(sunny.address().as_str()));
let (got, body) = open(&trapti, &as_peer(&sunny, "Sunny"), &frame).unwrap();
assert_eq!(got, content);
assert_eq!(body, b"# the plan\n");
let mut bad = frame.clone();
bad[60] ^= 1;
let e = open(&trapti, &as_peer(&sunny, "Sunny"), &bad)
.unwrap_err()
.to_string();
assert!(e.contains("signature"), "{e}");
let mut bad = frame.clone();
let n = bad.len() - 1;
bad[n] ^= 1;
assert!(open(&trapti, &as_peer(&sunny, "Sunny"), &bad).is_err());
let stranger = Identity::from_seed(&[3u8; 64]);
let e = open(&trapti, &as_peer(&stranger, "X"), &frame)
.unwrap_err()
.to_string();
assert!(e.contains("another sender"), "{e}");
assert!(open(&stranger, &as_peer(&sunny, "Sunny"), &frame).is_err());
assert!(open(&trapti, &as_peer(&sunny, "Sunny"), &frame[..100]).is_err());
let mut v = frame.clone();
v[0] = 2;
assert!(sender_of(&v).is_none());
assert!(open(&trapti, &as_peer(&sunny, "Sunny"), &v).is_err());
}
#[test]
fn a_note_travels_too_and_the_id_is_one_per_document_and_friend() {
let (sunny, trapti) = two();
let frame = seal(
&sunny,
&as_peer(&trapti, "Trapti"),
&Content::Note {
text: "water the beans".into(),
name: "Sunny".into(),
at: Folder::default(),
},
b"",
)
.unwrap();
let (got, body) = open(&trapti, &as_peer(&sunny, "Sunny"), &frame).unwrap();
assert!(matches!(got, Content::Note { ref text, .. } if text == "water the beans"));
assert!(body.is_empty());
assert_eq!(
frame_id("d1", &trapti.address()),
frame_id("d1", &trapti.address())
);
assert_ne!(
frame_id("d1", &trapti.address()),
frame_id("d1", &sunny.address())
);
assert_eq!(frame_id("d1", "k").len(), 64);
}
#[test]
fn a_code_is_three_words_a_check_and_ten_minutes_of_room() {
let code = mint_code().unwrap();
let parts: Vec<&str> = code.split('-').collect();
assert_eq!(parts.len(), 4, "{code}");
assert_eq!(parts[3].len(), 3);
assert_eq!(normalize_code(&code).unwrap(), code);
assert_eq!(
normalize_code(&format!(" {} ", code.to_uppercase().replace('-', " "))).unwrap(),
code
);
assert!(normalize_code("ocean ladder").is_err());
assert!(
normalize_code("ocean-ladder-xyzzy-abc").is_err(),
"not a word"
);
let mut wrong = parts[..3].join("-");
wrong.push_str(if parts[3] == "bbb" { "-ccc" } else { "-bbb" });
assert!(normalize_code(&wrong).unwrap_err().contains("heard wrong"));
assert_eq!(room_of(&code).len(), 64);
assert_ne!(room_of(&code), room_of("acid-acorn-acre-bbb"));
assert_eq!(words().len(), 1296);
}
#[test]
fn the_emoji_are_four_and_the_same_from_either_side() {
let (a, b) = two();
let x = emoji(
a.sign.verifying_key().as_bytes(),
b.sign.verifying_key().as_bytes(),
);
let y = emoji(
b.sign.verifying_key().as_bytes(),
a.sign.verifying_key().as_bytes(),
);
assert_eq!(x, y);
assert_eq!(x.split(' ').count(), 4, "{x}");
let c = Identity::from_seed(&[5u8; 64]);
assert_ne!(
x,
emoji(
a.sign.verifying_key().as_bytes(),
c.sign.verifying_key().as_bytes()
)
);
}
#[test]
fn the_identity_is_kept_once_and_read_back_the_same() {
let dir = crate::store::tempdir::Dir::new("snyvi-peer-id");
let s = crate::secrets::Secrets::file_only(dir.path.join("keys.json"));
assert!(
Identity::load(&s).is_none(),
"nothing minted before anyone pairs"
);
let a = Identity::load_or_mint(&s).unwrap();
let b = Identity::load_or_mint(&s).unwrap();
assert_eq!(a.address(), b.address());
assert_eq!(a.box_public(), b.box_public());
assert_eq!(Identity::load(&s).unwrap().address(), a.address());
let auth = a.relay_auth("GET", "/inbox/x");
let (secs, sig) = auth.split_once('.').unwrap();
assert!(secs.parse::<i64>().is_ok());
assert_eq!(unb64(sig).unwrap().len(), 64);
}
#[test]
fn friends_are_pinned_listed_renamed_muted_removed_and_restored() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch("CREATE TABLE docs (id TEXT PRIMARY KEY, title TEXT NOT NULL);")
.unwrap();
conn.execute_batch(SCHEMA).unwrap();
for c in COLUMNS_1_19 {
conn.execute_batch(c).unwrap();
}
let (sunny, trapti) = two();
let t = pin(&conn, &as_peer(&trapti, "Trapti"), 100).unwrap();
assert_eq!(t.name, "Trapti");
assert_eq!(t.paired_at, 100);
assert_eq!(list(&conn).unwrap().len(), 1);
assert!(by_name(&conn, " trapti ").unwrap().is_some());
assert!(rename(&conn, t.id, "T").unwrap());
assert!(
!rename(&conn, t.id, " ").unwrap(),
"a name is not nothing"
);
assert!(mute(&conn, t.id, true).unwrap());
assert!(set_desk(&conn, t.id, 4).unwrap());
assert_eq!(get(&conn, t.id).unwrap().unwrap().desk_id, 4);
assert!(set_desk(&conn, t.id, 0).unwrap());
assert!(get(&conn, t.id).unwrap().unwrap().muted);
assert!(remove(&conn, t.id, 200).unwrap());
assert!(
by_name(&conn, "T").unwrap().is_none(),
"removed is off the list"
);
assert!(restore(&conn, t.id).unwrap());
assert!(by_name(&conn, "T").unwrap().is_some());
let mut again = as_peer(&trapti, "Trapti");
again.box_key = b64(&sunny.box_public());
let t2 = pin(&conn, &again, 300).unwrap();
assert_eq!(t2.id, t.id);
assert_eq!(t2.name, "T", "the reader's name for them stands");
assert_eq!(t2.box_key, b64(&sunny.box_public()));
let fid = queue(&conn, &t2, "doc1", 400).unwrap();
assert_eq!(queue(&conn, &t2, "doc1", 401).unwrap(), fid);
assert_eq!(unsent(&conn).unwrap().len(), 1);
failed(&conn, &fid, "offline").unwrap();
assert_eq!(unsent(&conn).unwrap()[0].tries, 1);
waiting(&conn, &fid, "their mailbox is full").unwrap();
waiting(&conn, &fid, "their mailbox is full").unwrap();
assert_eq!(unsent(&conn).unwrap()[0].tries, 1);
sent(&conn, &fid, 402).unwrap();
assert!(unsent(&conn).unwrap().is_empty());
let l1 = queue_note(&conn, &t2, " water the beans ", 410).unwrap();
let l2 = queue_note(&conn, &t2, "water the beans", 411).unwrap();
assert_ne!(l1, l2);
let waiting = unsent(&conn).unwrap();
assert_eq!(waiting.len(), 2);
assert_eq!(waiting[0].text, "water the beans");
assert_eq!(waiting[0].doc_id, "");
sent(&conn, &l1, 412).unwrap();
sent(&conn, &l2, 412).unwrap();
let n = note_arrived(&conn, t2.id, " water the beans ", 500).unwrap();
assert_eq!(notes_waiting(&conn).unwrap()[0].text, "water the beans");
assert!(settle_note(&conn, n, "taken", 501).unwrap());
assert!(notes_waiting(&conn).unwrap().is_empty());
assert!(settle_note(&conn, n, "restore", 502).unwrap());
assert!(settle_note(&conn, n, "remove", 503).unwrap());
assert!(!settle_note(&conn, n, "eat", 503).unwrap());
conn.execute("INSERT INTO docs VALUES('doc1', 'Plan')", [])
.unwrap();
let o = offer(&conn, t2.id, "doc1", "pane-a", "Claude", 600).unwrap();
let open = offers_open(&conn).unwrap();
assert_eq!(open[0].title, "Plan");
assert_eq!(open[0].to, "T");
assert!(answer_offer(&conn, o, true, 601).unwrap());
assert!(!answer_offer(&conn, o, true, 601).unwrap(), "answered once");
assert!(!reopen_offer(&conn, o).unwrap(), "a sent offer stays sent");
let no = offer(&conn, t2.id, "doc1", "pane-b", "Claude", 601).unwrap();
assert!(answer_offer(&conn, no, false, 601).unwrap());
assert!(reopen_offer(&conn, no).unwrap(), "Not now has an Undo");
assert_eq!(offers_open(&conn).unwrap().len(), 1);
assert!(answer_offer(&conn, no, false, 601).unwrap());
offer(&conn, t2.id, "doc1", "pane-a", "Claude", 602).unwrap();
assert_eq!(drop_offers_of(&conn, "pane-a", 603).unwrap(), 1);
assert!(offers_open(&conn).unwrap().is_empty());
assert!(!taken(&conn, "f1").unwrap());
take(&conn, "f1", 700).unwrap();
take(&conn, "f1", 701).unwrap();
take(&conn, "f2", 800).unwrap();
assert!(taken(&conn, "f1").unwrap());
assert_eq!(prune_taken(&conn, 750).unwrap(), 1);
assert!(!taken(&conn, "f1").unwrap());
assert!(taken(&conn, "f2").unwrap());
clear(&conn).unwrap();
assert!(list(&conn).unwrap().is_empty());
assert!(!taken(&conn, "f2").unwrap());
}
#[test]
fn the_link_knows_its_address_and_its_messages() {
let (a, _) = two();
assert_eq!(
ws_of("http://127.0.0.1:8799", &a.address()),
format!("ws://127.0.0.1:8799/inbox/{}", a.address())
);
assert_eq!(
ws_of(RELAY, &a.address()),
format!("wss://relay.snyvi.com/inbox/{}", a.address())
);
assert!(relay_ws(&a.address()).starts_with("ws"));
let pushed: Pushed = serde_json::from_str(
r#"{"frame":{"id":"ab","sender":"k","size":12,"at":1700000000000}}"#,
)
.unwrap();
let Pushed::Frame(w) = pushed;
assert_eq!((w.id.as_str(), w.sender.as_str(), w.size), ("ab", "k", 12));
assert!(
serde_json::from_str::<Pushed>(r#"{"hello":{}}"#).is_err(),
"a message this daemon does not know is not a frame"
);
assert!(serde_json::from_str::<Pushed>("pong").is_err());
assert_eq!(ack_message("ab"), r#"{"ack":"ab"}"#);
}
#[test]
fn a_deposit_that_cannot_go_now_waits_and_says_why() {
assert_eq!(
later(429, "the relay is busy; try again in a minute"),
Some(BUSY)
);
assert_eq!(
later(429, "the mailbox is full; try later"),
Some("their mailbox is full")
);
assert!(
later(404, "no such mailbox").is_some(),
"an idle friend's cleared mailbox"
);
assert_eq!(
later(401, "not the key's signature"),
None,
"a real failure"
);
assert_eq!(later(413, "a frame is at most 8 MB"), None);
}
#[test]
fn the_doorbell_and_the_pairing_speak_plainly() {
assert_eq!(ws_base("https://relay.snyvi.com"), "wss://relay.snyvi.com");
assert_eq!(ws_base("http://127.0.0.1:8799"), "ws://127.0.0.1:8799");
assert_eq!(rings_for(r#"{"ready":"spake"}"#), Some("spake"));
assert_eq!(rings_for(r#"{"ready":"hello"}"#), Some("hello"));
assert_eq!(rings_for(r#"{"ready":"other"}"#), None);
assert_eq!(rings_for("pong"), None);
assert_eq!(
pair_refused(429).to_string(),
"the relay is busy; try again in a minute"
);
assert_eq!(pair_refused(410).to_string(), "this code was already used");
assert_eq!(
ring_wait("http://127.0.0.1:9", &"a".repeat(64), "spake", "ab", 1),
Bell::Absent
);
}
#[test]
fn the_backoff_climbs_and_stops() {
let mut last = Duration::ZERO;
for attempt in 0..12 {
let b = backoff(attempt);
let base = Duration::from_secs(1 << attempt).min(BACKOFF_MAX);
assert!(
b >= base && b <= base.mul_f64(1.3),
"attempt {attempt}: {b:?}"
);
assert!(base >= last);
last = base;
}
assert!(
backoff(40) <= BACKOFF_MAX.mul_f64(1.3),
"capped, not overflowed"
);
}
#[test]
fn a_frame_says_its_folder_and_older_and_newer_frames_still_open() {
let doc = Content::Document {
title: "Plan".into(),
lang: None,
file: Some("PLAN.md".into()),
name: "Trapti".into(),
id: "d1".into(),
at: Folder {
repo: Some("r".repeat(32)),
remote: None,
path: Some("docs/PLAN.md".into()),
key: None,
branch: Some("main".into()),
v: CONTENT_V,
},
};
let (got, body) = unpack(&pack(&doc, b"x")).unwrap();
assert_eq!(got, doc);
assert_eq!(body, b"x");
let old = br#"{"kind":"document","title":"T","name":"S","id":"i","file":"a.md"}"#;
let got: Content = serde_json::from_slice(old).unwrap();
assert!(matches!(got, Content::Document { ref at, .. } if *at == Folder::default()));
#[derive(Deserialize)]
#[serde(rename_all = "lowercase", tag = "kind")]
enum Was {
Document {
title: String,
file: Option<String>,
name: String,
id: String,
},
Note {
text: String,
name: String,
},
}
let json = serde_json::to_vec(&doc).unwrap();
let was: Was = serde_json::from_slice(&json).unwrap();
assert!(
matches!(was, Was::Document { ref title, ref file, .. } if title == "Plan" && file.as_deref() == Some("PLAN.md"))
);
let note = Content::Note {
text: "hi".into(),
name: "S".into(),
at: Folder {
repo: Some("r".into()),
v: CONTENT_V,
..Default::default()
},
};
let was: Was = serde_json::from_slice(&serde_json::to_vec(¬e).unwrap()).unwrap();
assert!(matches!(was, Was::Note { ref text, .. } if text == "hi"));
let newer = br#"{"kind":"receipt","of":"d1","v":2}"#;
assert_eq!(
serde_json::from_slice::<Content>(newer).unwrap(),
Content::Other
);
}
#[test]
fn a_path_from_a_friend_stays_inside_the_folder() {
assert_eq!(safe_path("docs/PLAN.md").as_deref(), Some("docs/PLAN.md"));
assert_eq!(safe_path("docs\\PLAN.md").as_deref(), Some("docs/PLAN.md"));
assert_eq!(safe_path("PLAN.md").as_deref(), Some("PLAN.md"));
for bad in [
"",
"/etc/passwd",
"../x.md",
"docs/../../x",
"C:/x.md",
"c:x.md",
"a//b",
"./a",
"a/\u{1b}[2J",
] {
assert_eq!(safe_path(bad), None, "{bad:?}");
}
assert_eq!(safe_path(&"a/".repeat(201)), None);
}
}