mod common;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::time::{Duration, Instant};
use common::{Between, Device, PASSPHRASE, TestRelay, contains, data, drain, name, relay_bytes};
use efema::{Epoch, Error, Key, Limits, Secret, Transport, chain};
#[tokio::test(flavor = "multi_thread")]
async fn two_devices_sync_through_the_relay() {
let relay = TestRelay::start().await;
let (laptop, phone) = (Device::new(), Device::new());
let on_laptop = laptop.open(relay.transport(), "notes", 1).await;
let pushed = on_laptop.push([b"first".as_slice(), b"second"]).await.unwrap();
assert_eq!(pushed.positions, [2, 3]);
let on_phone = phone.open(relay.transport(), "notes", 1).await;
let received = drain(&on_phone).await;
assert_eq!(data(&received), [b"first".as_slice(), b"second"]);
assert!(received.iter().all(|r| !r.own && r.device == laptop.state.device()));
assert_eq!(on_phone.key().id(), on_laptop.key().id());
on_phone.push([b"from the phone"]).await.unwrap();
let received = drain(&on_laptop).await;
assert_eq!(data(&received), [b"first".as_slice(), b"second", b"from the phone"]);
assert_eq!(received.iter().map(|r| r.own).collect::<Vec<_>>(), [true, true, false]);
assert_eq!(received[2].device, phone.state.device());
}
#[tokio::test(flavor = "multi_thread")]
async fn the_relay_keeps_only_ciphertext() {
let relay = TestRelay::start().await;
let device = Device::new();
let client = device.open(relay.transport(), "notes", 1).await;
let marker = b"PLAINTEXT-MARKER-the-relay-must-never-see";
client.push([marker.as_slice(), marker.as_slice()]).await.unwrap();
let stored = relay_bytes(relay.dir.path());
assert!(stored.len() > marker.len(), "the relay's files were not found");
assert!(!contains(&stored, marker), "the relay stored an item in the clear");
assert!(!contains(&stored, PASSPHRASE), "the relay stored the passphrase");
assert!(!contains(&stored, client.key().to_bytes().as_slice()), "the relay stored the key");
assert!(!contains(&stored, device.state.device().as_bytes()), "the relay can see which device wrote");
let state = std::fs::read(device.state.path()).unwrap();
assert!(!contains(&state, client.key().to_bytes().as_slice()), "the state file holds the key unlocked");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_stream_opens_only_with_its_secret() {
let relay = TestRelay::start().await;
let first = Device::new();
let client = first.open(relay.transport(), "notes", 1).await;
let key = client.key().clone();
let stranger = Device::new();
let wrong = stranger.try_open(relay.transport(), "notes", 1, Secret::Passphrase(b"wrong")).await;
assert!(matches!(wrong, Err(Error::WrongPassphrase { .. })), "{wrong:?}");
let other_key = Key::generate().unwrap();
let mismatch = stranger.try_open(relay.transport(), "notes", 1, Secret::Key(other_key)).await;
assert!(matches!(mismatch, Err(Error::KeyMismatch { .. })), "{mismatch:?}");
let with_key = stranger.try_open(relay.transport(), "notes", 1, Secret::Key(key)).await.unwrap();
client.push([b"hello"]).await.unwrap();
assert_eq!(data(&drain(&with_key).await), [b"hello".as_slice()]);
let fresh = Device::new();
let refused = fresh.try_open(relay.transport(), "empty", 1, Secret::Key(Key::generate().unwrap())).await;
assert!(matches!(refused, Err(Error::NoStreamKey { .. })), "{refused:?}");
assert!(relay.transport().read(&name("empty"), None, 10).await.is_err(), "a refused open created a stream");
}
#[tokio::test(flavor = "multi_thread")]
async fn devices_creating_a_stream_at_once_agree_on_one_key() {
let relay = TestRelay::start().await;
let devices: Vec<Device> = (0..4).map(|_| Device::new()).collect();
let devices = Arc::new(devices);
let mut tasks = Vec::new();
for index in 0..devices.len() {
let devices = devices.clone();
let transport = relay.transport();
tasks.push(tokio::spawn(async move {
let client = devices[index].open(transport, "notes", 1).await;
client.push([format!("from {index}").into_bytes()]).await.unwrap();
client.key().id()
}));
}
let mut keys = Vec::new();
for task in tasks {
keys.push(task.await.unwrap());
}
assert!(keys.windows(2).all(|pair| pair[0] == pair[1]), "devices ended up with different keys: {keys:?}");
let reader = Device::new();
let client = reader.open(relay.transport(), "notes", 1).await;
let mut seen: Vec<String> =
drain(&client).await.iter().map(|r| String::from_utf8(r.data.clone()).unwrap()).collect();
seen.sort();
assert_eq!(seen, ["from 0", "from 1", "from 2", "from 3"]);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_device_that_joined_opens_without_the_network() {
let relay = TestRelay::start().await;
let device = Device::new();
device.open(relay.transport(), "notes", 1).await;
let gone = relay.url.clone();
relay.stop().await;
let device = device.restart();
let offline = efema::Relay::new(&gone).unwrap();
let client = device.try_open(offline, "notes", 1, Secret::Passphrase(PASSPHRASE)).await.unwrap();
let wrong = device.try_open(efema::Relay::new(&gone).unwrap(), "notes", 1, Secret::Passphrase(b"nope")).await;
assert!(matches!(wrong, Err(Error::WrongPassphrase { .. })), "{wrong:?}");
let unreachable = client.push([b"x"]).await.unwrap_err();
assert!(matches!(unreachable, Error::Transport(efema::TransportError::Unreachable { .. })), "{unreachable:?}");
}
#[tokio::test(flavor = "multi_thread")]
async fn the_cursor_survives_a_restart() {
let relay = TestRelay::start().await;
let writer = Device::new();
let on_writer = writer.open(relay.transport(), "notes", 1).await;
on_writer.push([b"one", b"two"]).await.unwrap();
let reader = Device::new();
let client = reader.open(relay.transport(), "notes", 1).await;
assert_eq!(drain(&client).await.len(), 2);
drop(client);
let reader = reader.restart();
on_writer.push([b"three"]).await.unwrap();
let client = reader.open(relay.transport(), "notes", 1).await;
assert_eq!(data(&drain(&client).await), [b"three".as_slice()]);
}
#[tokio::test(flavor = "multi_thread")]
async fn what_was_not_acknowledged_comes_again() {
let relay = TestRelay::start().await;
let device = Device::new();
let client = device.open(relay.transport(), "notes", 1).await;
client.push([b"one", b"two"]).await.unwrap();
let pulled = client.pull().await.unwrap();
assert_eq!(pulled.entries.len(), 2);
drop(client);
let device = device.restart();
let client = device.open(relay.transport(), "notes", 1).await;
let again = client.pull().await.unwrap();
assert_eq!(again.entries, pulled.entries);
client.ack(&again).await.unwrap();
assert!(client.pull().await.unwrap().entries.is_empty());
}
#[tokio::test(flavor = "multi_thread")]
async fn an_older_app_stops_before_a_newer_epoch_and_cannot_write_behind_it() {
let relay = TestRelay::start().await;
let (old, new) = (Device::new(), Device::new());
let on_old = old.open(relay.transport(), "notes", 1).await;
on_old.push([b"v1 a", b"v1 b"]).await.unwrap();
let on_new = new.open(relay.transport(), "notes", 2).await;
on_new.push([b"v2 a"]).await.unwrap();
on_old.push([b"v1 c"]).await.unwrap_err();
let pulled = on_old.pull().await.unwrap();
assert_eq!(data(&pulled.entries), [b"v1 a".as_slice(), b"v1 b"]);
assert!(pulled.more());
on_old.ack(&pulled).await.unwrap();
let stopped = on_old.pull().await.unwrap_err();
assert!(matches!(stopped, Error::NewerEpoch { seq: 4, ours: Epoch(1), entry: Epoch(2), .. }), "{stopped:?}");
assert_eq!(on_old.cursor().await.unwrap().seq, 3);
let refused = on_old.push([b"v1 d"]).await.unwrap_err();
assert!(matches!(refused, Error::EpochBehind { ours: Epoch(1), current: Epoch(2), .. }), "{refused:?}");
let received = drain(&on_new).await;
assert_eq!(data(&received), [b"v1 a".as_slice(), b"v1 b", b"v2 a"]);
assert_eq!(received.iter().map(|r| r.epoch.0).collect::<Vec<_>>(), [1, 1, 2]);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_push_is_cut_to_what_one_request_may_carry() {
let relay = TestRelay::start().await;
let device = Device::new();
let mut between = Between::new(relay.transport());
between.limits = Limits { max_body: 1000, ..Limits::V1 };
let appends = between.appends.clone();
let client = device.open(between, "notes", 1).await;
let opened = appends.load(Ordering::SeqCst);
let items: Vec<Vec<u8>> = (0..20u8).map(|i| vec![i; 150]).collect();
let pushed = client.push(&items).await.unwrap();
assert_eq!(pushed.positions, (2..22).collect::<Vec<u64>>());
let batches = appends.load(Ordering::SeqCst) - opened;
assert!(batches > 1, "twenty items of 150 bytes went in one request of at most 1000");
assert_eq!(data(&drain(&client).await), items.iter().map(Vec::as_slice).collect::<Vec<_>>());
let largest = client.max_item_len().await.unwrap();
assert!(largest < 1000 && largest > 800, "{largest}");
client.push([vec![0u8; largest]]).await.unwrap();
let sent = appends.load(Ordering::SeqCst);
let refused = client.push([vec![1u8; 10], vec![0u8; largest + 1]]).await.unwrap_err();
assert!(matches!(refused, Error::ItemTooLarge { .. }), "{refused:?}");
assert_eq!(appends.load(Ordering::SeqCst), sent, "part of a push that could not fit was sent");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_relay_that_lost_its_data_is_not_written_into() {
let relay = TestRelay::start().await;
let device = Device::new();
let client = device.open(relay.transport(), "notes", 1).await;
client.push([b"before"]).await.unwrap();
let dir = relay.stop().await;
std::fs::remove_dir_all(dir.path()).unwrap();
let relay = TestRelay::start_in(dir).await;
let client = device.open(relay.transport(), "notes", 1).await;
assert!(matches!(client.push([b"after"]).await, Err(Error::StreamGone { .. })));
assert!(matches!(client.pull().await, Err(Error::StreamGone { .. })));
assert!(relay.transport().read(&name("notes"), None, 10).await.is_err(), "the push created a stream");
let other = Device::new();
other.open(relay.transport(), "notes", 1).await.push([b"anew"]).await.unwrap();
assert!(matches!(client.push([b"after"]).await, Err(Error::StreamReplaced { .. })));
assert!(matches!(client.pull().await, Err(Error::StreamReplaced { .. })));
device.state.forget(&name("notes")).unwrap();
let client = device.open(relay.transport(), "notes", 1).await;
assert_eq!(data(&drain(&client).await), [b"anew".as_slice()]);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_relay_restored_from_an_old_backup_is_named() {
let relay = TestRelay::start().await;
let device = Device::new();
let client = device.open(relay.transport(), "notes", 1).await;
client.push([b"one"]).await.unwrap();
let dir = relay.stop().await;
let backup = tempfile::tempdir().unwrap();
for file in std::fs::read_dir(dir.path()).unwrap() {
let file = file.unwrap();
std::fs::copy(file.path(), backup.path().join(file.file_name())).unwrap();
}
let relay = TestRelay::start_in(dir).await;
let client = device.open(relay.transport(), "notes", 1).await;
client.push([b"two".as_slice(), b"three"]).await.unwrap();
drain(&client).await;
relay.stop().await;
let relay = TestRelay::start_in(backup).await;
let client = device.open(relay.transport(), "notes", 1).await;
let ahead = client.pull().await.unwrap_err();
assert!(matches!(ahead, Error::CursorAhead { cursor: 4, head: 2, .. }), "{ahead:?}");
let report = client.doctor().await;
assert!(report.problems().iter().any(|p| p.contains("State::forget")), "{report}");
let other = Device::new();
other.open(relay.transport(), "notes", 1).await.push([b"x", b"y"]).await.unwrap();
let diverged = client.pull().await.unwrap_err();
assert!(matches!(diverged, Error::CursorDiverged { seq: 4, .. }), "{diverged:?}");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_relay_cannot_alter_an_entry_unnoticed() {
let relay = TestRelay::start().await;
let writer = Device::new();
writer.open(relay.transport(), "notes", 1).await.push([b"good".as_slice(), b"to be altered"]).await.unwrap();
let mut careless = Between::new(relay.transport());
careless.tamper = Some(Arc::new(|_, page| {
if let Some(entry) = page.entries.iter_mut().find(|e| e.seq == 3) {
let last = entry.data.len() - 1;
entry.data[last] ^= 1;
}
}));
let reader = Device::new();
let client = reader.open(careless, "notes", 1).await;
let pulled = client.pull().await.unwrap();
assert_eq!(data(&pulled.entries), [b"good".as_slice()], "what came before the altered entry is delivered");
client.ack(&pulled).await.unwrap();
let broken = client.pull().await.unwrap_err();
assert!(matches!(broken, Error::BrokenChain { seq: 3, .. }), "{broken:?}");
let mut careful = Between::new(relay.transport());
careful.tamper = Some(Arc::new(|after, page| {
let mut previous = after.map_or_else(|| chain::genesis(&page.stream), |cursor| cursor.hash);
for entry in &mut page.entries {
if entry.seq == 3 {
let last = entry.data.len() - 1;
entry.data[last] ^= 1;
}
entry.hash = chain::link(&previous, entry.seq, entry.epoch, &entry.data);
previous = entry.hash;
}
}));
let reader = Device::new();
let client = reader.open(careful, "notes", 1).await;
let pulled = client.pull().await.unwrap();
client.ack(&pulled).await.unwrap();
let forged = client.pull().await.unwrap_err();
assert!(matches!(forged, Error::Inauthentic { seq: 3, .. }), "{forged:?}");
}
#[tokio::test(flavor = "multi_thread")]
async fn an_entry_no_client_wrote_is_refused_not_skipped() {
let relay = TestRelay::start().await;
let device = Device::new();
let client = device.open(relay.transport(), "notes", 1).await;
client.push([b"sealed"]).await.unwrap();
let plain = efema::wire::Batch { epoch: Epoch(1), entries: vec![b"plain text".to_vec()], stream: None };
relay.transport().append(&name("notes"), &plain).await.unwrap();
let pulled = client.pull().await.unwrap();
assert_eq!(data(&pulled.entries), [b"sealed".as_slice()]);
client.ack(&pulled).await.unwrap();
let foreign = client.pull().await.unwrap_err();
assert!(matches!(foreign, Error::ForeignEntry { seq: 3, .. }), "{foreign:?}");
relay.transport().append(&name("plain"), &plain).await.unwrap();
let refused = Device::new().try_open(relay.transport(), "plain", 1, Secret::Passphrase(PASSPHRASE)).await;
assert!(matches!(refused, Err(Error::ForeignEntry { seq: 1, .. })), "{refused:?}");
}
#[tokio::test(flavor = "multi_thread")]
async fn a_wait_wakes_on_a_new_entry_and_not_before() {
let relay = TestRelay::start().await;
let (a, b) = (Device::new(), Device::new());
let on_a = a.open(relay.transport(), "notes", 1).await;
let on_b = b.open(relay.transport(), "notes", 1).await;
drain(&on_a).await;
let started = Instant::now();
assert!(!on_a.wait(Duration::from_secs(1)).await.unwrap());
assert!(started.elapsed() >= Duration::from_millis(900));
let waiting = tokio::spawn(async move {
let woke = on_a.wait(Duration::from_secs(30)).await.unwrap();
(woke, on_a)
});
tokio::time::sleep(Duration::from_millis(200)).await;
let started = Instant::now();
on_b.push([b"wake up"]).await.unwrap();
let (woke, on_a) = waiting.await.unwrap();
assert!(woke && started.elapsed() < Duration::from_secs(10));
assert_eq!(data(&drain(&on_a).await), [b"wake up".as_slice()]);
}
#[tokio::test(flavor = "multi_thread")]
async fn the_doctor_reports_what_the_device_and_the_relay_know() {
let relay = TestRelay::start().await;
let (a, b) = (Device::new(), Device::new());
let on_a = a.open(relay.transport(), "notes", 1).await;
let on_b = b.open(relay.transport(), "notes", 1).await;
on_a.push([b"1", b"2", b"3"]).await.unwrap();
let report = on_b.doctor().await;
assert_eq!(report.behind(), Some(4), "{report}");
assert!(report.problems().is_empty(), "{report}");
assert_eq!(report.key.id, on_a.key().id());
let shown = report.to_string();
for line in ["relay ", "stream ", "cursor ", "key ", "device ", "epoch "] {
assert!(shown.contains(line), "the report has no `{line}` line:\n{shown}");
}
assert!(shown.contains("4 entries to read"), "{shown}");
assert!(shown.contains(&on_a.key().id().to_string()), "{shown}");
drain(&on_b).await;
assert_eq!(on_b.doctor().await.behind(), Some(0));
let url = relay.url.clone();
relay.stop().await;
let report = on_b.doctor().await;
assert!(report.relay.is_err() && !report.problems().is_empty(), "{report}");
assert!(report.to_string().contains(&url), "{report}");
}