use crate::{agents::ForAgent, Db, Storelike};
use iroh::protocol::Router;
const STROKE_DATA: &str = "https://atomicdata.dev/ontology/canvas/strokeData";
const FOLDER_PROP: &str = "https://atomicdata.dev/ontology/canvas/folderId";
const CANVAS_CLASS: &str = "https://atomicdata.dev/ontology/canvas/Canvas";
const FOLDER_CLASS: &str = "https://atomicdata.dev/classes/Folder";
struct IrohPair {
db_a: Db,
db_b: Db,
drive: String,
node_id_a: String,
_router_a: Router,
ep_b: iroh::Endpoint,
}
async fn setup_pair(prefix: &str) -> IrohPair {
use crate::sync::peer;
let db_a = Db::init_temp(&format!("{prefix}_a")).await.unwrap();
let (agent_a, drive) = db_a.setup("Alice").await.unwrap();
let secret = agent_a.build_secret().unwrap();
let db_b = Db::init_temp(&format!("{prefix}_b")).await.unwrap();
db_b.load_agent_from_secret(&secret).await.unwrap();
let (node_id_a, router_a) = peer::start(db_a.clone()).await.unwrap();
let ep_b = iroh::Endpoint::builder()
.discovery_n0()
.discovery_local_network()
.bind()
.await
.unwrap();
let node_addr_a = router_a.endpoint().node_addr().await.unwrap();
ep_b.add_node_addr(node_addr_a).unwrap();
IrohPair {
db_a,
db_b,
drive,
node_id_a: node_id_a.to_string(),
_router_a: router_a,
ep_b,
}
}
async fn sync_b_from_a(pair: &IrohPair) -> usize {
use crate::sync::peer;
peer::sync_drive_with_peer_using(&pair.ep_b, &pair.node_id_a, &pair.drive, &pair.db_b, true)
.await
.expect("B→A sync should succeed")
}
async fn wait_until<F, Fut>(timeout: std::time::Duration, mut check: F) -> bool
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = bool>,
{
let deadline = tokio::time::Instant::now() + timeout;
while tokio::time::Instant::now() < deadline {
if check().await {
return true;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
false
}
async fn wait_for_live_peers(min: usize, timeout: std::time::Duration) {
let ok = wait_until(timeout, || async {
crate::sync::peer::live_peer_count() >= min
})
.await;
assert!(
ok,
"expected ≥{min} live peer(s), got {}",
crate::sync::peer::live_peer_count()
);
}
async fn stroke_count(db: &Db, canvas: &str) -> usize {
let r = db.get_resource(&canvas.into()).await.unwrap();
match r.get(STROKE_DATA) {
Ok(crate::Value::Json(serde_json::Value::Array(arr))) => arr.len(),
_ => 0,
}
}
async fn folder_id_on(db: &Db, canvas: &str) -> Option<String> {
let r = db.get_resource(&canvas.into()).await.ok()?;
r.get(FOLDER_PROP)
.ok()
.map(|v| v.to_string())
.filter(|s| !s.is_empty())
}
async fn assign_folder(db: &Db, canvas: &str, folder: &str) {
let mut r = db.get_resource(&canvas.into()).await.unwrap();
r.ensure_materialized().unwrap();
r.set_unsafe(FOLDER_PROP.into(), crate::Value::String(folder.into()))
.unwrap();
r.save_locally(db).await.unwrap();
}
#[tokio::test]
async fn e2e_hello_exchanges_device_names() {
use crate::sync::peer;
let pair = setup_pair("e2e_hello").await;
peer::set_device_name(&pair.db_a, "Alice's Laptop");
peer::set_device_name(&pair.db_b, "Bob's Phone");
let outcome = peer::sync_drive_with_peer_using_outcome(
&pair.ep_b,
&pair.node_id_a,
&pair.drive,
&pair.db_b,
true,
)
.await
.expect("B→A sync should succeed");
assert_eq!(
outcome.peer_name.as_deref(),
Some("Alice's Laptop"),
"initiator should see A's self-reported HELLO name"
);
let b_known = peer::get_known_peers(&pair.db_b);
let alice_on_b = b_known
.iter()
.find(|p| peer::normalize_node_id(&p.node_id) == peer::normalize_node_id(&pair.node_id_a));
assert_eq!(
alice_on_b.map(|p| p.name.as_str()),
Some("Alice's Laptop"),
"B (initiator) should have persisted A's HELLO name into known-peers"
);
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
let a_known = peer::get_known_peers(&pair.db_a);
let ep_b_node_id = pair.ep_b.node_id().to_string();
assert!(
!a_known
.iter()
.any(|p| peer::normalize_node_id(&p.node_id) == peer::normalize_node_id(&ep_b_node_id)),
"A (accept side) must NOT persist an unsolicited peer into \
known-peers — that hands them a permanent auto-reconnect slot \
with no pairing/consent (got {a_known:?})"
);
}
#[tokio::test]
async fn e2e_bidirectional_bulk_sync() {
let pair = setup_pair("e2e_bulk").await;
let canvas_a = pair
.db_a
.create_resource(
CANVAS_CLASS,
&pair.drive,
"Canvas A",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![
serde_json::json!({"color": 1}),
])),
)]),
)
.await
.unwrap();
let canvas_b = pair
.db_b
.create_resource(
CANVAS_CLASS,
&pair.drive,
"Canvas B",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![
serde_json::json!({"color": 2}),
])),
)]),
)
.await
.unwrap();
let imported = sync_b_from_a(&pair).await;
assert!(imported > 0, "B should import A's resources");
pair.db_b
.get_resource(&canvas_a.as_str().into())
.await
.expect("B should have A's canvas after bulk sync");
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
pair.db_a
.get_resource(&canvas_b.as_str().into())
.await
.expect("A should have B's canvas after bidirectional SYNC_PUSH");
}
#[tokio::test]
async fn e2e_stroke_append_after_sync() {
let pair = setup_pair("e2e_stroke").await;
let canvas = pair
.db_a
.create_resource(
CANVAS_CLASS,
&pair.drive,
"Stroke canvas",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![
serde_json::json!({"color": 1, "path": [[0.0, 0.0]]}),
])),
)]),
)
.await
.unwrap();
sync_b_from_a(&pair).await;
wait_for_live_peers(1, std::time::Duration::from_secs(3)).await;
assert_eq!(stroke_count(&pair.db_b, &canvas).await, 1);
let mut resource_a = pair
.db_a
.get_resource(&canvas.as_str().into())
.await
.unwrap();
resource_a.ensure_materialized().unwrap();
resource_a.init_undo();
resource_a
.push_list_item(
STROKE_DATA,
serde_json::json!({"color": 2, "width": 2.0, "path": [[1.0, 1.0]]}),
)
.unwrap();
resource_a.save_locally(&pair.db_a).await.unwrap();
let live_ok = wait_until(std::time::Duration::from_secs(3), || async {
stroke_count(&pair.db_b, &canvas).await == 2
})
.await;
assert!(live_ok, "B must see second stroke via live push");
}
#[tokio::test]
async fn e2e_canvas_folder_assignment_syncs() {
let pair = setup_pair("e2e_folder").await;
let folder = pair
.db_a
.create_resource(FOLDER_CLASS, &pair.drive, "Sketches", None)
.await
.unwrap();
let canvas = pair
.db_a
.create_resource(CANVAS_CLASS, &pair.drive, "Inbox", None)
.await
.unwrap();
sync_b_from_a(&pair).await;
wait_for_live_peers(1, std::time::Duration::from_secs(3)).await;
assign_folder(&pair.db_a, &canvas, &folder).await;
let live_ok = wait_until(std::time::Duration::from_secs(3), || async {
folder_id_on(&pair.db_b, &canvas).await.as_deref() == Some(folder.as_str())
})
.await;
if !live_ok {
let _ = sync_b_from_a(&pair).await;
}
assert_eq!(
folder_id_on(&pair.db_b, &canvas).await.as_deref(),
Some(folder.as_str()),
"B must see folderId after A assigns canvas to folder (live or bulk resync)"
);
}
#[tokio::test]
async fn e2e_new_resource_after_bulk_resync() {
let pair = setup_pair("e2e_new_res").await;
pair.db_a
.create_resource(CANVAS_CLASS, &pair.drive, "Seed", None)
.await
.unwrap();
sync_b_from_a(&pair).await;
let new_canvas = pair
.db_a
.create_resource(
CANVAS_CLASS,
&pair.drive,
"After sync",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![
serde_json::json!({"color": 99}),
])),
)]),
)
.await
.unwrap();
assert!(
pair.db_b
.get_resource(&new_canvas.as_str().into())
.await
.is_err(),
"B should not have the canvas before resync"
);
let imported = sync_b_from_a(&pair).await;
assert!(imported > 0, "second bulk sync should import new canvas");
let on_b = pair
.db_b
.get_resource(&new_canvas.as_str().into())
.await
.expect("B should have new canvas after bulk resync");
assert_eq!(
on_b.get(crate::urls::NAME).unwrap().to_string(),
"After sync"
);
}
#[tokio::test]
async fn e2e_engine_pull_after_iroh_bulk_sync() {
let pair = setup_pair("e2e_engine_pull").await;
let canvas = pair
.db_a
.create_resource(
CANVAS_CLASS,
&pair.drive,
"Pull test",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![serde_json::json!({"n": 1})])),
)]),
)
.await
.unwrap();
sync_b_from_a(&pair).await;
let mut resource_a = pair
.db_a
.get_resource(&canvas.as_str().into())
.await
.unwrap();
resource_a.ensure_materialized().unwrap();
resource_a.init_undo();
resource_a
.push_list_item(STROKE_DATA, serde_json::json!({"n": 2}))
.unwrap();
resource_a.save_locally(&pair.db_a).await.unwrap();
let drive_subject =
crate::Subject::from_raw(&pair.drive, pair.db_b.get_base_domain().as_deref());
let subjects = crate::sync::engine::collect_drive_subjects(&pair.db_b, &drive_subject).await;
let vvs = crate::sync::engine::build_drive_vvs(&pair.db_b, &subjects);
let hash = crate::sync::engine::compute_drive_hash(&vvs);
let frames = crate::sync::engine::handle_sync_vv(
&pair.drive,
&hash,
&[],
&std::collections::HashMap::new(),
&pair.db_a,
&ForAgent::Public,
)
.await;
let mut imported = 0;
for frame in frames {
if frame.first() == Some(&crate::sync::protocol::tag::SYNC_PUSH) {
if let Some(push) = crate::sync::protocol::decode_sync_push(&frame[1..]) {
let (count, _) = crate::sync::engine::import_sync_push(
&push,
&pair.db_b,
&ForAgent::Sudo,
false,
)
.await;
imported += count;
}
}
}
assert!(imported > 0, "engine pull should import A's edit");
assert_eq!(stroke_count(&pair.db_b, &canvas).await, 2);
}
#[tokio::test]
async fn e2e_managed_node_replicates_missing_drive() {
let pair = setup_pair("e2e_replicate").await;
let doc = pair
.db_a
.create_resource(
CANVAS_CLASS,
&pair.drive,
"Doc on A",
Some(vec![(
STROKE_DATA,
crate::Value::Json(serde_json::Value::Array(vec![serde_json::json!({"n": 1})])),
)]),
)
.await
.unwrap();
assert!(
!pair.db_b.has_resource_locally(&pair.drive),
"B should not host the drive before replicating"
);
let imported = sync_b_from_a(&pair).await;
assert!(imported > 0, "B should import the drive's resources");
pair.db_b
.get_resource(&doc.as_str().into())
.await
.expect("B should have A's resource after replicating");
let usage = pair
.db_b
.per_drive_usage(&[pair.drive.clone()])
.await
.unwrap();
let row = usage
.iter()
.find(|u| u.drive_subject == pair.drive)
.expect("usage row for the replicated drive");
assert!(
row.resource_count > 0,
"replicated drive should report resources, got {row:?}"
);
assert!(
pair.db_b.has_resource_locally(&pair.drive),
"B should host the drive after replicating"
);
}
#[tokio::test]
async fn pushing_a_workspace_to_an_empty_device_is_not_a_failure() {
use crate::sync::peer;
let db_server = Db::init_temp("push_server").await.unwrap();
db_server.setup("Server").await.unwrap();
let db_phone = Db::init_temp("push_phone").await.unwrap();
let (_agent, drive) = db_phone.setup("Alice").await.unwrap();
db_phone
.create_resource(crate::urls::FOLDER, &drive, "notes", None)
.await
.unwrap();
let (node_id, router) = peer::start(db_server.clone()).await.unwrap();
let ep_phone = iroh::Endpoint::builder()
.discovery_n0()
.discovery_local_network()
.bind()
.await
.unwrap();
ep_phone
.add_node_addr(router.endpoint().node_addr().await.unwrap())
.unwrap();
peer::sync_drive_with_peer_using(&ep_phone, &node_id.to_string(), &drive, &db_phone, true)
.await
.expect("pushing a workspace up must not be reported as a failure");
let landed = wait_until(std::time::Duration::from_secs(10), || async {
db_server.has_resource_locally(&drive)
})
.await;
assert!(landed, "the workspace should land on the always-on device");
}
#[tokio::test]
async fn peer_sync_says_why_when_a_different_agent_may_read_nothing() {
use crate::sync::peer;
let db_a = Db::init_temp("xagent_a").await.unwrap();
let (agent_a, drive) = db_a.setup("Alice").await.unwrap();
let mut drive_resource = db_a.get_resource(&drive.clone().into()).await.unwrap();
drive_resource
.set(
crate::urls::READ.into(),
vec![agent_a.subject.to_string()].into(),
&db_a,
)
.await
.unwrap();
drive_resource.save(&db_a).await.unwrap();
let db_b = Db::init_temp("xagent_b").await.unwrap();
db_b.setup("Bob").await.unwrap();
let (node_id_a, router_a) = peer::start(db_a.clone()).await.unwrap();
let ep_b = iroh::Endpoint::builder()
.discovery_n0()
.discovery_local_network()
.bind()
.await
.unwrap();
ep_b.add_node_addr(router_a.endpoint().node_addr().await.unwrap())
.unwrap();
let result =
peer::sync_drive_with_peer_using(&ep_b, &node_id_a.to_string(), &drive, &db_b, true).await;
let error = result.expect_err("an empty sync must not report success");
let message = error.to_string();
assert!(
message.contains("nothing synced"),
"the failure must say why, got: {message}"
);
assert!(
!db_b.has_resource_locally(&drive),
"a private drive may not cross to somebody else's device"
);
}
#[tokio::test]
async fn same_agent_peers_reconcile_the_agent_resource() {
let pair = setup_pair("agent_reconcile").await;
let agent_subject = pair.db_b.get_default_agent().unwrap().subject.to_string();
let before = pair
.db_b
.get_resource(&agent_subject.as_str().into())
.await
.unwrap();
assert!(
before.get(crate::urls::NAME).is_err(),
"B starts with a nameless stub agent"
);
sync_b_from_a(&pair).await;
let db_b = pair.db_b.clone();
let subject = agent_subject.clone();
let got_name = wait_until(std::time::Duration::from_secs(10), || {
let db_b = db_b.clone();
let subject = subject.clone();
async move {
db_b.get_resource(&subject.as_str().into())
.await
.ok()
.and_then(|r| r.get(crate::urls::NAME).ok().map(|v| v.to_string()))
== Some("Alice".to_string())
}
})
.await;
assert!(
got_name,
"B's agent resource should gain the name A holds, over the live link"
);
}
#[tokio::test]
async fn a_device_hands_its_agent_resource_to_a_different_account_server() {
use crate::sync::peer;
let db_server = Db::init_temp("agentx_server").await.unwrap();
db_server.setup("Server").await.unwrap();
let db_phone = Db::init_temp("agentx_phone").await.unwrap();
let (phone_agent, _drive) = db_phone.setup("Alice").await.unwrap();
let phone_agent_subject = phone_agent.subject.to_string();
let (node_id, router) = peer::start(db_server.clone()).await.unwrap();
let ep_phone = iroh::Endpoint::builder()
.discovery_n0()
.discovery_local_network()
.bind()
.await
.unwrap();
ep_phone
.add_node_addr(router.endpoint().node_addr().await.unwrap())
.unwrap();
peer::sync_drive_with_peer_using(&ep_phone, &node_id.to_string(), &_drive, &db_phone, true)
.await
.expect("sync should not fail");
let db_server_c = db_server.clone();
let subject = phone_agent_subject.clone();
let landed = wait_until(std::time::Duration::from_secs(10), || {
let db = db_server_c.clone();
let subject = subject.clone();
async move {
db.get_resource(&subject.as_str().into())
.await
.ok()
.and_then(|r| r.get(crate::urls::NAME).ok().map(|v| v.to_string()))
== Some("Alice".to_string())
}
})
.await;
assert!(
landed,
"the server should hold the phone's agent resource so a browser can find its drive"
);
}
#[tokio::test]
async fn a_peer_cannot_forge_a_third_agents_resource() {
use crate::agents::{Agent, ForAgent};
use crate::sync::peer;
let db_server = Db::init_temp("spoof_server").await.unwrap();
db_server.setup("Server").await.unwrap();
let eve = Agent::new(Some("Eve")).unwrap();
let victim = Agent::new(Some("Victim")).unwrap();
let mut cache = std::collections::HashMap::new();
let forged = peer::admitted_for_drive_for_test(
&db_server,
&ForAgent::AgentSubject(eve.subject.clone()),
&victim.subject.to_string(),
false,
&mut cache,
)
.await;
assert!(
!forged,
"a peer must not write an agent resource it holds no key for"
);
let own = peer::admitted_for_drive_for_test(
&db_server,
&ForAgent::AgentSubject(eve.subject.clone()),
&eve.subject.to_string(),
false,
&mut cache,
)
.await;
assert!(own, "a peer's own agent resource must be admitted");
}
#[tokio::test]
async fn a_completed_peer_sync_names_the_drive_to_reconnect_to() {
use crate::sync::peer;
let db_a = Db::init_temp("active_drive_a").await.unwrap();
let (agent_a, drive) = db_a.setup("Alice").await.unwrap();
let db_b = Db::init_temp("active_drive_b").await.unwrap();
let mut agent_b = agent_a.clone();
agent_b.initial_drive = None;
db_b.set_default_agent(agent_b);
assert_eq!(
db_b.get_active_drive(),
None,
"the second device starts out not knowing which drive is hers"
);
let (node_id_a, router_a) = peer::start(db_a.clone()).await.unwrap();
let ep_b = iroh::Endpoint::builder()
.discovery_n0()
.discovery_local_network()
.bind()
.await
.unwrap();
ep_b.add_node_addr(router_a.endpoint().node_addr().await.unwrap())
.unwrap();
peer::sync_drive_with_peer_using(&ep_b, &node_id_a.to_string(), &drive, &db_b, true)
.await
.expect("Alice's two devices should sync");
assert_eq!(
db_b.get_active_drive().as_deref(),
Some(drive.as_str()),
"after syncing a drive, the device must know to dial back for it"
);
}