#![cfg(feature = "offline-sync")]
use autumn_web::sync::{
Change, ChangeOutcome, LwwResolver, MemorySyncBackend, Op, PullResponse, PushRequest,
SyncBackend, SyncScope,
};
use chrono::{Duration, Utc};
use serde_json::json;
const SCOPE: &str = SyncScope::GLOBAL;
fn change(change_id: &str, pk: &str, op: Op, payload: Option<serde_json::Value>) -> Change {
Change {
change_id: change_id.to_owned(),
collection: "conformance".to_owned(),
pk: pk.to_owned(),
op,
payload,
base_version: 0,
updated_at: Utc::now(),
}
}
fn push(device: &str, changes: Vec<Change>) -> PushRequest {
PushRequest {
device_id: device.to_owned(),
changes,
}
}
#[allow(clippy::too_many_lines)] pub fn run_backend_conformance(backend: &dyn SyncBackend) {
let resolver = LwwResolver;
assert_eq!(backend.latest_version().expect("latest"), 0);
assert_eq!(backend.tombstone_horizon(SCOPE).expect("horizon"), 0);
let seed = push(
"device-a",
vec![
change(
"00000000-0000-4000-8000-000000000001",
"n1",
Op::Upsert,
Some(json!({"title": "one"})),
),
change(
"00000000-0000-4000-8000-000000000002",
"n2",
Op::Upsert,
Some(json!({"title": "two"})),
),
],
);
let response = backend
.apply_push(SCOPE, &seed, &resolver)
.expect("seed push");
let versions: Vec<i64> = response
.outcomes
.iter()
.map(|o| match o {
ChangeOutcome::Applied { version } => *version,
other => panic!("expected Applied, got {other:?}"),
})
.collect();
assert!(versions[0] > 0);
assert!(versions[1] > versions[0], "versions strictly increase");
assert_eq!(backend.latest_version().expect("latest"), versions[1]);
let retry = backend
.apply_push(SCOPE, &seed, &resolver)
.expect("retry push");
let retry_versions: Vec<i64> = retry
.outcomes
.iter()
.map(|o| match o {
ChangeOutcome::AlreadyApplied { version } => *version,
other => panic!("retry must dedup, got {other:?}"),
})
.collect();
assert_eq!(
retry_versions, versions,
"already_applied must carry the versions of the first application"
);
assert_eq!(backend.latest_version().expect("latest"), versions[1]);
let PullResponse::Ok {
rows,
next_cursor,
tombstone_horizon,
} = backend.pull_since(SCOPE, 0, 100, 0).expect("pull all")
else {
panic!("cursor 0 never requires a resync");
};
assert_eq!(rows.len(), 2);
assert!(rows.windows(2).all(|w| w[0].version < w[1].version));
assert_eq!(next_cursor, versions[1]);
assert_eq!(tombstone_horizon, 0);
let PullResponse::Ok {
rows, next_cursor, ..
} = backend
.pull_since(SCOPE, versions[1], 100, versions[1])
.expect("pull caught-up")
else {
panic!("caught-up cursor never requires a resync");
};
assert!(rows.is_empty());
assert_eq!(next_cursor, versions[1], "empty page keeps the cursor");
let PullResponse::Ok { rows, .. } = backend.pull_since(SCOPE, 0, 1, 0).expect("pull limited")
else {
panic!("cursor 0 never requires a resync");
};
assert_eq!(rows.len(), 1, "limit caps the page");
let mut winner = change(
"00000000-0000-4000-8000-000000000003",
"n1",
Op::Upsert,
Some(json!({"title": "one-b"})),
);
winner.base_version = 0; winner.updated_at = Utc::now() + Duration::seconds(30);
let winner_request = push("device-b", vec![winner]);
let response = backend
.apply_push(SCOPE, &winner_request, &resolver)
.expect("conflict push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"stale base_version must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(row.version > versions[1], "resolved rows get a NEW version");
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("one-b"),
"newer client write wins LWW"
);
let winner_version = row.version;
let winner_row = row.clone();
let latest_before_replay = backend.latest_version().expect("latest");
let replay = backend
.apply_push(SCOPE, &winner_request, &resolver)
.expect("replay push");
let ChangeOutcome::Resolved { row } = &replay.outcomes[0] else {
panic!(
"a retry of a Resolved change must replay Resolved, got {:?}",
replay.outcomes[0]
);
};
assert_eq!(
row, &winner_row,
"the replay must carry the originally resolved row"
);
assert_eq!(
backend.latest_version().expect("latest"),
latest_before_replay,
"a replay must not assign versions"
);
let mut loser = change(
"00000000-0000-4000-8000-000000000004",
"n1",
Op::Upsert,
Some(json!({"title": "stale"})),
);
loser.base_version = versions[0]; loser.updated_at = Utc::now() - Duration::seconds(3600);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![loser]), &resolver)
.expect("losing conflict push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"stale base_version must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
row.version > winner_version,
"even KeepServer re-versions the row"
);
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("one-b"),
"server content survives a losing push"
);
let mut delete = change(
"00000000-0000-4000-8000-000000000005",
"n2",
Op::Delete,
None,
);
delete.base_version = versions[1];
let response = backend
.apply_push(SCOPE, &push("device-a", vec![delete]), &resolver)
.expect("delete push");
assert!(matches!(
response.outcomes[0],
ChangeOutcome::Applied { .. }
));
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, 0, 100, 0)
.expect("pull with tombstone")
else {
panic!("cursor 0 never requires a resync");
};
let n2 = rows.iter().find(|r| r.pk == "n2").expect("n2 row");
assert!(n2.deleted, "deletes replicate as tombstones");
assert_eq!(n2.payload, None);
let latest = backend.latest_version().expect("latest");
let removed = backend.gc_tombstones(latest).expect("gc");
assert_eq!(removed, 1);
assert_eq!(backend.tombstone_horizon(SCOPE).expect("horizon"), latest);
assert_eq!(backend.gc_tombstones(latest).expect("re-gc"), 0);
let PullResponse::Ok { rows, .. } = backend.pull_since(SCOPE, 0, 100, 0).expect("pull post-gc")
else {
panic!("cursor 0 never requires a resync");
};
assert!(rows.iter().all(|r| !r.deleted), "GC'd tombstones are gone");
let live_version = rows.first().expect("a live row survives GC").version;
let stale = backend
.pull_since(SCOPE, versions[0], 100, versions[0])
.expect("stale pull");
assert!(
matches!(stale, PullResponse::FullResyncRequired { tombstone_horizon } if tombstone_horizon == latest),
"a session starting behind the horizon must be told to resync, got {stale:?}"
);
let mid_page = backend
.pull_since(SCOPE, live_version, 100, 0)
.expect("mid-pagination pull");
assert!(
matches!(mid_page, PullResponse::Ok { .. }),
"a from-0 session paging past a sub-horizon cursor must get rows, got {mid_page:?}"
);
let latest_before_gc = backend.latest_version().expect("latest");
backend.gc_tombstones(i64::MAX).expect("gc with huge up_to");
assert_eq!(
backend.tombstone_horizon(SCOPE).expect("horizon"),
latest_before_gc,
"the horizon must never exceed the newest committed version"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000006",
"n3",
Op::Upsert,
Some(json!({"title": "post-gc"})),
)],
),
&resolver,
)
.expect("post-gc push");
let ChangeOutcome::Applied { version } = response.outcomes[0] else {
panic!("expected Applied, got {:?}", response.outcomes[0]);
};
assert!(
version > latest_before_gc,
"new versions must land above the clamped horizon"
);
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, latest_before_gc, 100, latest_before_gc)
.expect("post-gc pull")
else {
panic!("a cursor at the horizon never requires a resync");
};
assert!(
rows.iter().any(|r| r.pk == "n3"),
"a row created after a huge-up_to GC must be delivered, got {rows:?}"
);
let latest_before_reject = backend.latest_version().expect("latest");
let bad_batch = push(
"device-a",
vec![
change(
"00000000-0000-4000-8000-00000000000a",
"n4",
Op::Upsert,
Some(json!({"title": "valid sibling"})),
),
change(
"00000000-0000-4000-8000-00000000000b",
"n5",
Op::Upsert,
None,
),
],
);
let err = backend
.apply_push(SCOPE, &bad_batch, &resolver)
.expect_err("an upsert without a payload must be rejected");
assert!(
matches!(err, autumn_web::sync::SyncError::Protocol(_)),
"payload-less upserts are protocol errors, got {err:?}"
);
assert_eq!(
backend.latest_version().expect("latest"),
latest_before_reject,
"a rejected batch must not assign versions"
);
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, latest_before_reject, 100, latest_before_reject)
.expect("pull after rejected push")
else {
panic!("a caught-up cursor never requires a resync");
};
assert!(
rows.is_empty(),
"nothing from a rejected batch may be applied, got {rows:?}"
);
let response = backend
.apply_push(
SCOPE,
&push("device-a", vec![bad_batch.changes[0].clone()]),
&resolver,
)
.expect("re-push of the valid sibling");
assert!(
matches!(response.outcomes[0], ChangeOutcome::Applied { .. }),
"the valid sibling of a rejected batch must apply on retry, got {:?}",
response.outcomes[0]
);
let latest_before_dup = backend.latest_version().expect("latest");
let dup_batch = push(
"device-a",
vec![
change(
"00000000-0000-4000-8000-000000000030",
"dup-a",
Op::Upsert,
Some(json!({"title": "first"})),
),
change(
"00000000-0000-4000-8000-000000000030",
"dup-b",
Op::Upsert,
Some(json!({"title": "second"})),
),
],
);
let err = backend
.apply_push(SCOPE, &dup_batch, &resolver)
.expect_err("duplicate change_ids within one batch must be rejected");
assert!(
matches!(err, autumn_web::sync::SyncError::Protocol(_)),
"duplicate change_ids are protocol errors, got {err:?}"
);
assert_eq!(
backend.latest_version().expect("latest"),
latest_before_dup,
"a rejected batch must not assign versions"
);
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, latest_before_dup, 100, latest_before_dup)
.expect("pull after rejected duplicate batch")
else {
panic!("a caught-up cursor never requires a resync");
};
assert!(
rows.is_empty(),
"nothing from a rejected batch may be applied, got {rows:?}"
);
let response = backend
.apply_push(
SCOPE,
&push("device-a", vec![dup_batch.changes[0].clone()]),
&resolver,
)
.expect("re-push after the rejected duplicate batch");
assert!(
matches!(response.outcomes[0], ChangeOutcome::Applied { .. }),
"a change from a rejected batch must apply cleanly on retry, got {:?}",
response.outcomes[0]
);
let latest_before_stale = backend.latest_version().expect("latest");
let mut stale_edit = change(
"00000000-0000-4000-8000-00000000000c",
"n2",
Op::Upsert,
Some(json!({"title": "resurrected?"})),
);
stale_edit.base_version = versions[1]; stale_edit.updated_at = Utc::now() - Duration::seconds(3600);
let stale_request = push("device-offline", vec![stale_edit]);
let response = backend
.apply_push(SCOPE, &stale_request, &resolver)
.expect("stale push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"an offline-past-GC edit must resolve, not clean-apply, got {:?}",
response.outcomes[0]
);
};
assert!(
row.deleted && row.payload.is_none(),
"the GC'd deletion must win under LWW, got {row:?}"
);
assert!(
row.version > latest_before_stale,
"the resolution gets a fresh version so every device converges"
);
let stale_tombstone = row.clone();
let replay = backend
.apply_push(SCOPE, &stale_request, &resolver)
.expect("stale replay");
let ChangeOutcome::Resolved { row } = &replay.outcomes[0] else {
panic!(
"a retry of the stale push must replay Resolved, got {:?}",
replay.outcomes[0]
);
};
assert_eq!(
row, &stale_tombstone,
"the replay must carry the original tombstone"
);
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, latest_before_stale, 100, latest_before_stale)
.expect("pull after stale push")
else {
panic!("an at-horizon cursor never requires a resync");
};
assert!(
rows.iter().all(|r| r.pk != "n2" || r.deleted),
"n2 must not come back as a live row: {rows:?}"
);
let removed = backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("re-gc the fresh tombstone");
assert_eq!(removed, 1, "the resolution tombstone is GC'd again");
let mut skewed_edit = change(
"00000000-0000-4000-8000-00000000000e",
"n2",
Op::Upsert,
Some(json!({"title": "clock cheat"})),
);
skewed_edit.base_version = versions[1];
skewed_edit.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-skewed", vec![skewed_edit]), &resolver)
.expect("skewed push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a future-clocked stale edit must still resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
row.deleted && row.payload.is_none(),
"a clock ahead of the server must not resurrect the GC'd row, got {row:?}"
);
for (attempt, change_id) in [
"00000000-0000-4000-8000-000000000011",
"00000000-0000-4000-8000-000000000012",
]
.iter()
.enumerate()
{
let mut repeat_edit = change(
change_id,
"n2",
Op::Upsert,
Some(json!({"title": "resurrected via repeat?"})),
);
repeat_edit.base_version = versions[1];
repeat_edit.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-skewed", vec![repeat_edit]), &resolver)
.expect("repeat stale push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"repeat stale push #{attempt} must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
row.deleted && row.payload.is_none(),
"repeat stale push #{attempt} must stay server-winning even with \
a fast clock, got {row:?}"
);
}
let response = backend
.apply_push(
SCOPE,
&push(
"device-offline",
vec![change(
"00000000-0000-4000-8000-00000000000d",
"brand-new",
Op::Upsert,
Some(json!({"title": "fresh insert"})),
)],
),
&resolver,
)
.expect("fresh insert push");
assert!(
matches!(response.outcomes[0], ChangeOutcome::Applied { .. }),
"a base-0 insert must stay a clean apply, got {:?}",
response.outcomes[0]
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000040",
"reborn",
Op::Upsert,
Some(json!({"title": "first life"})),
)],
),
&resolver,
)
.expect("first-life push");
let ChangeOutcome::Applied {
version: first_life,
} = response.outcomes[0]
else {
panic!(
"first life must clean-apply, got {:?}",
response.outcomes[0]
);
};
let mut kill = change(
"00000000-0000-4000-8000-000000000041",
"reborn",
Op::Delete,
None,
);
kill.base_version = first_life;
backend
.apply_push(SCOPE, &push("device-a", vec![kill]), &resolver)
.expect("first-life delete");
let removed = backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("gc before rebirth");
assert!(removed >= 1, "the first-life tombstone must be GC'd");
let response = backend
.apply_push(
SCOPE,
&push(
"device-b",
vec![change(
"00000000-0000-4000-8000-000000000042",
"reborn",
Op::Upsert,
Some(json!({"title": "second life"})),
)],
),
&resolver,
)
.expect("rebirth push");
let ChangeOutcome::Applied {
version: second_life,
} = response.outcomes[0]
else {
panic!(
"a base-0 recreate must stay a clean apply, got {:?}",
response.outcomes[0]
);
};
let mut ghost_edit = change(
"00000000-0000-4000-8000-000000000043",
"reborn",
Op::Upsert,
Some(json!({"title": "ghost edit"})),
);
ghost_edit.base_version = first_life; ghost_edit.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-a", vec![ghost_edit]), &resolver)
.expect("ghost edit push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a pre-horizon stale edit must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row.version > second_life
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("second life"),
"the recreated row must win over a pre-horizon stale edit, got {row:?}"
);
let after_ghost_edit = row.version;
let mut ghost_delete = change(
"00000000-0000-4000-8000-000000000044",
"reborn",
Op::Delete,
None,
);
ghost_delete.base_version = first_life;
ghost_delete.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-a", vec![ghost_delete]), &resolver)
.expect("ghost delete push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a pre-horizon stale delete must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted && row.version > after_ghost_edit,
"a pre-horizon stale delete must not kill the new incarnation, got {row:?}"
);
let after_ghosts = row.version;
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, after_ghosts - 1, 100, after_ghosts - 1)
.expect("pull the reborn row")
else {
panic!("an at-feed cursor never requires a resync");
};
assert!(
rows.iter().any(|r| r.pk == "reborn"
&& !r.deleted
&& r.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("second life")),
"the second life must still be pulled intact, got {rows:?}"
);
let mut third_life = change(
"00000000-0000-4000-8000-000000000045",
"reborn",
Op::Upsert,
Some(json!({"title": "third life"})),
);
third_life.base_version = second_life; third_life.updated_at = Utc::now() + Duration::seconds(120);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![third_life]), &resolver)
.expect("post-horizon conflict push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a post-horizon stale edit must resolve, got {:?}",
response.outcomes[0]
);
};
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("third life"),
"post-horizon conflicts must still reach the (LWW) resolver"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000050",
"bystander",
Op::Upsert,
Some(json!({"title": "one"})),
)],
),
&resolver,
)
.expect("bystander create");
let ChangeOutcome::Applied {
version: bystander_v1,
} = response.outcomes[0]
else {
panic!("bystander must clean-apply, got {:?}", response.outcomes[0]);
};
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000051",
"victim-b",
Op::Upsert,
Some(json!({"title": "doomed"})),
)],
),
&resolver,
)
.expect("victim create");
let ChangeOutcome::Applied { version: victim_v } = response.outcomes[0] else {
panic!("victim must clean-apply, got {:?}", response.outcomes[0]);
};
let victim_kill = Change {
base_version: victim_v,
..change(
"00000000-0000-4000-8000-000000000052",
"victim-b",
Op::Delete,
None,
)
};
backend
.apply_push(SCOPE, &push("device-a", vec![victim_kill]), &resolver)
.expect("victim delete");
let removed = backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("unrelated gc");
assert!(removed >= 1, "victim-b's tombstone must be GC'd");
assert!(
backend.tombstone_horizon(SCOPE).expect("horizon") > bystander_v1,
"the unrelated GC must have raised the horizon above bystander's base"
);
let update = Change {
base_version: bystander_v1,
..change(
"00000000-0000-4000-8000-000000000053",
"bystander",
Op::Upsert,
Some(json!({"title": "two"})),
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![update]), &resolver)
.expect("bystander update");
assert!(matches!(
response.outcomes[0],
ChangeOutcome::Applied { .. }
));
let mut offline_edit = change(
"00000000-0000-4000-8000-000000000054",
"bystander",
Op::Upsert,
Some(json!({"title": "three"})),
);
offline_edit.base_version = bystander_v1;
offline_edit.updated_at = Utc::now() + Duration::seconds(60);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![offline_edit]), &resolver)
.expect("offline edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a same-incarnation stale edit must resolve, got {:?}",
response.outcomes[0]
);
};
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("three"),
"the resolver must run despite the unrelated GC — LWW takes the \
newer client edit"
);
let mut losing_edit = change(
"00000000-0000-4000-8000-000000000055",
"bystander",
Op::Upsert,
Some(json!({"title": "ancient"})),
);
losing_edit.base_version = bystander_v1;
losing_edit.updated_at = Utc::now() - Duration::seconds(3600);
let response = backend
.apply_push(SCOPE, &push("device-d", vec![losing_edit]), &resolver)
.expect("losing offline edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"the losing direction must also resolve, got {:?}",
response.outcomes[0]
);
};
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("three"),
"LWW keeps the server content for an older client edit"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000056",
"phoenix",
Op::Upsert,
Some(json!({"title": "first flight"})),
)],
),
&resolver,
)
.expect("phoenix create");
let ChangeOutcome::Applied {
version: phoenix_v1,
} = response.outcomes[0]
else {
panic!("phoenix must clean-apply, got {:?}", response.outcomes[0]);
};
let phoenix_kill = Change {
base_version: phoenix_v1,
..change(
"00000000-0000-4000-8000-000000000057",
"phoenix",
Op::Delete,
None,
)
};
let response = backend
.apply_push(SCOPE, &push("device-a", vec![phoenix_kill]), &resolver)
.expect("phoenix delete");
let ChangeOutcome::Applied {
version: phoenix_tomb,
} = response.outcomes[0]
else {
panic!(
"the delete must clean-apply, got {:?}",
response.outcomes[0]
);
};
let rebirth = Change {
base_version: phoenix_tomb,
..change(
"00000000-0000-4000-8000-000000000058",
"phoenix",
Op::Upsert,
Some(json!({"title": "second flight"})),
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![rebirth]), &resolver)
.expect("phoenix recreate");
assert!(matches!(
response.outcomes[0],
ChangeOutcome::Applied { .. }
));
let mut phoenix_ghost = change(
"00000000-0000-4000-8000-000000000059",
"phoenix",
Op::Upsert,
Some(json!({"title": "ghost flight"})),
);
phoenix_ghost.base_version = phoenix_v1;
phoenix_ghost.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![phoenix_ghost]), &resolver)
.expect("phoenix ghost edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"an old-incarnation base must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("second flight"),
"the recreate must win over the previous incarnation without any \
GC involved, got {row:?}"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000060",
"lazarus",
Op::Upsert,
Some(json!({"title": "first life"})),
)],
),
&resolver,
)
.expect("lazarus create");
let ChangeOutcome::Applied {
version: lazarus_v1,
} = response.outcomes[0]
else {
panic!("lazarus must clean-apply, got {:?}", response.outcomes[0]);
};
let lazarus_kill = Change {
base_version: lazarus_v1,
..change(
"00000000-0000-4000-8000-000000000061",
"lazarus",
Op::Delete,
None,
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![lazarus_kill]), &resolver)
.expect("lazarus delete");
let ChangeOutcome::Applied {
version: lazarus_tomb,
} = response.outcomes[0]
else {
panic!(
"the delete must clean-apply, got {:?}",
response.outcomes[0]
);
};
let removed = backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("gc lazarus tombstone");
assert_eq!(removed, 1, "exactly the lazarus tombstone is dropped");
assert_eq!(
backend.tombstone_horizon(SCOPE).expect("horizon"),
lazarus_tomb,
"the horizon is the newest dropped tombstone version"
);
let pull = backend
.pull_since(SCOPE, lazarus_tomb, 100, lazarus_tomb)
.expect("at-horizon pull");
assert!(
matches!(pull, PullResponse::Ok { .. }),
"a session at the horizon must not be told to resync, got {pull:?}"
);
let risen = Change {
base_version: lazarus_tomb,
..change(
"00000000-0000-4000-8000-000000000062",
"lazarus",
Op::Upsert,
Some(json!({"title": "risen"})),
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![risen]), &resolver)
.expect("at-horizon recreate");
let ChangeOutcome::Applied {
version: lazarus_v2,
} = response.outcomes[0]
else {
panic!(
"a recreate based on the GC'd delete it pulled must apply, got {:?}",
response.outcomes[0]
);
};
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, lazarus_v2 - 1, 100, lazarus_v2 - 1)
.expect("pull risen lazarus")
else {
panic!("an at-feed cursor never requires a resync");
};
assert!(
rows.iter().any(|r| r.pk == "lazarus" && !r.deleted),
"the recreate must be live, got {rows:?}"
);
let mut lazarus_ghost = change(
"00000000-0000-4000-8000-000000000063",
"lazarus",
Op::Upsert,
Some(json!({"title": "ghost"})),
);
lazarus_ghost.base_version = lazarus_v1;
lazarus_ghost.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![lazarus_ghost]), &resolver)
.expect("lazarus ghost edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a first-life base must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("risen"),
"the at-horizon recreate must be a REAL new incarnation, got {row:?}"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000064",
"mummy",
Op::Upsert,
Some(json!({"title": "wrapped"})),
)],
),
&resolver,
)
.expect("mummy create");
let ChangeOutcome::Applied { version: mummy_v1 } = response.outcomes[0] else {
panic!("mummy must clean-apply, got {:?}", response.outcomes[0]);
};
let mummy_kill = Change {
base_version: mummy_v1,
..change(
"00000000-0000-4000-8000-000000000065",
"mummy",
Op::Delete,
None,
)
};
backend
.apply_push(SCOPE, &push("device-a", vec![mummy_kill]), &resolver)
.expect("mummy delete");
backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("gc mummy tombstone");
let mut mummy_stale = change(
"00000000-0000-4000-8000-000000000066",
"mummy",
Op::Upsert,
Some(json!({"title": "unwrapped?"})),
);
mummy_stale.base_version = mummy_v1; mummy_stale.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-b", vec![mummy_stale]), &resolver)
.expect("mummy stale push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a strictly-pre-horizon base must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
row.deleted && row.payload.is_none(),
"a base that never saw the delete must stay server-winning, got {row:?}"
);
let mummy_tombstone = row.version;
let mummy_recreate = Change {
base_version: mummy_tombstone,
..change(
"00000000-0000-4000-8000-000000000067",
"mummy",
Op::Upsert,
Some(json!({"title": "second wrapping"})),
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![mummy_recreate]), &resolver)
.expect("mummy recreate over the materialized tombstone");
assert!(
matches!(response.outcomes[0], ChangeOutcome::Applied { .. }),
"a base equal to the tombstone's version saw the delete, got {:?}",
response.outcomes[0]
);
let mut mummy_ghost = change(
"00000000-0000-4000-8000-000000000068",
"mummy",
Op::Upsert,
Some(json!({"title": "ancient curse"})),
);
mummy_ghost.base_version = mummy_v1;
mummy_ghost.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-c", vec![mummy_ghost]), &resolver)
.expect("mummy ghost edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a first-incarnation base must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("second wrapping"),
"the recreate must survive the ghost, got {row:?}"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-000000000070",
"graveyard",
Op::Upsert,
Some(json!({"title": "first tenant"})),
)],
),
&resolver,
)
.expect("graveyard create");
let ChangeOutcome::Applied { version: grave_v1 } = response.outcomes[0] else {
panic!("create must clean-apply, got {:?}", response.outcomes[0]);
};
let grave_kill1 = Change {
base_version: grave_v1,
..change(
"00000000-0000-4000-8000-000000000071",
"graveyard",
Op::Delete,
None,
)
};
let response = backend
.apply_push(SCOPE, &push("device-a", vec![grave_kill1]), &resolver)
.expect("first delete");
let ChangeOutcome::Applied {
version: first_grave_tomb,
} = response.outcomes[0]
else {
panic!("delete must clean-apply, got {:?}", response.outcomes[0]);
};
let grave_rebirth = Change {
base_version: first_grave_tomb,
..change(
"00000000-0000-4000-8000-000000000072",
"graveyard",
Op::Upsert,
Some(json!({"title": "second tenant"})),
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![grave_rebirth]), &resolver)
.expect("recreate");
let ChangeOutcome::Applied { version: grave_v2 } = response.outcomes[0] else {
panic!("recreate must clean-apply, got {:?}", response.outcomes[0]);
};
let grave_kill2 = Change {
base_version: grave_v2,
..change(
"00000000-0000-4000-8000-000000000073",
"graveyard",
Op::Delete,
None,
)
};
let response = backend
.apply_push(SCOPE, &push("device-b", vec![grave_kill2]), &resolver)
.expect("second delete");
let ChangeOutcome::Applied {
version: second_grave_tomb,
} = response.outcomes[0]
else {
panic!("delete must clean-apply, got {:?}", response.outcomes[0]);
};
let redundant_delete = Change {
base_version: second_grave_tomb,
..change(
"00000000-0000-4000-8000-000000000074",
"graveyard",
Op::Delete,
None,
)
};
let response = backend
.apply_push(SCOPE, &push("device-c", vec![redundant_delete]), &resolver)
.expect("redundant delete");
assert!(
matches!(response.outcomes[0], ChangeOutcome::Applied { .. }),
"a no-op delete over a tombstone is a valid ack, got {:?}",
response.outcomes[0]
);
let mut third_tenant = change(
"00000000-0000-4000-8000-000000000075",
"graveyard",
Op::Upsert,
Some(json!({"title": "third tenant"})),
);
third_tenant.base_version = second_grave_tomb;
third_tenant.updated_at = Utc::now() + Duration::seconds(60);
let response = backend
.apply_push(SCOPE, &push("device-d", vec![third_tenant]), &resolver)
.expect("recreate from the original tombstone");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a same-incarnation recreate must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("third tenant"),
"the recreate from the pre-redundant-delete tombstone must win, \
got {row:?}"
);
let mut grave_ghost = change(
"00000000-0000-4000-8000-000000000076",
"graveyard",
Op::Upsert,
Some(json!({"title": "first tenant returns"})),
);
grave_ghost.base_version = grave_v1;
grave_ghost.updated_at = Utc::now() + Duration::days(365);
let response = backend
.apply_push(SCOPE, &push("device-a", vec![grave_ghost]), &resolver)
.expect("ghost edit");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"a first-incarnation base must resolve, got {:?}",
response.outcomes[0]
);
};
assert!(
!row.deleted
&& row
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("third tenant"),
"old-incarnation ghosts must still lose to the revived row, got {row:?}"
);
let response = backend
.apply_push(
SCOPE,
&push(
"device-a",
vec![change(
"00000000-0000-4000-8000-00000000000f",
"nullable",
Op::Upsert,
Some(serde_json::Value::Null),
)],
),
&resolver,
)
.expect("null-payload push");
let ChangeOutcome::Applied {
version: null_version,
} = response.outcomes[0]
else {
panic!(
"a null-payload upsert is a valid document, got {:?}",
response.outcomes[0]
);
};
let PullResponse::Ok { rows, .. } = backend
.pull_since(SCOPE, null_version - 1, 100, null_version - 1)
.expect("pull null row")
else {
panic!("an at-feed cursor never requires a resync");
};
let null_row = rows
.iter()
.find(|r| r.pk == "nullable")
.expect("the null-payload row must be pulled");
assert!(!null_row.deleted, "a null payload is not a tombstone");
assert_eq!(
null_row.payload,
Some(serde_json::Value::Null),
"the null document must survive the wire intact"
);
let mut null_conflict = change(
"00000000-0000-4000-8000-000000000010",
"nullable",
Op::Upsert,
Some(serde_json::Value::Null),
);
null_conflict.base_version = 0; null_conflict.updated_at = Utc::now() + Duration::seconds(60);
let null_request = push("device-b", vec![null_conflict]);
let response = backend
.apply_push(SCOPE, &null_request, &resolver)
.expect("null conflict push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!("stale base must resolve, got {:?}", response.outcomes[0]);
};
assert_eq!(row.payload, Some(serde_json::Value::Null));
let resolved_null_row = row.clone();
let replay = backend
.apply_push(SCOPE, &null_request, &resolver)
.expect("null replay");
let ChangeOutcome::Resolved { row } = &replay.outcomes[0] else {
panic!(
"the retry must replay Resolved, got {:?}",
replay.outcomes[0]
);
};
assert_eq!(
row, &resolved_null_row,
"the dedup snapshot must round-trip a null payload"
);
let tenant_change = |change_id: &str, title: &str| Change {
collection: "scoped".to_owned(),
..change(change_id, "s1", Op::Upsert, Some(json!({"title": title})))
};
let shared_change_id = "00000000-0000-4000-8000-000000000020";
let response = backend
.apply_push(
"tenant-a",
&push("device-a", vec![tenant_change(shared_change_id, "alpha")]),
&resolver,
)
.expect("tenant-a push");
let ChangeOutcome::Applied { version: version_a } = response.outcomes[0] else {
panic!(
"tenant-a push must clean-apply, got {:?}",
response.outcomes[0]
);
};
let response = backend
.apply_push(
"tenant-b",
&push("device-a", vec![tenant_change(shared_change_id, "beta")]),
&resolver,
)
.expect("tenant-b push");
let ChangeOutcome::Applied { version: version_b } = response.outcomes[0] else {
panic!(
"the same device/change/pk in another scope is a distinct fresh \
row and change, got {:?}",
response.outcomes[0]
);
};
assert!(version_b > version_a, "one global sequence spans scopes");
for (scope, version, title) in [
("tenant-a", version_a, "alpha"),
("tenant-b", version_b, "beta"),
] {
let PullResponse::Ok {
rows, next_cursor, ..
} = backend.pull_since(scope, 0, 100, 0).expect("tenant pull")
else {
panic!("cursor 0 never requires a resync");
};
assert_eq!(
rows.len(),
1,
"{scope} must see exactly its own row, got {rows:?}"
);
assert_eq!(rows[0].version, version);
assert_eq!(
rows[0]
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some(title),
"{scope} must pull its own payload"
);
assert_eq!(next_cursor, version);
}
let replay = backend
.apply_push(
"tenant-a",
&push("device-a", vec![tenant_change(shared_change_id, "alpha")]),
&resolver,
)
.expect("tenant-a replay");
assert!(
matches!(
replay.outcomes[0],
ChangeOutcome::AlreadyApplied { version } if version == version_a
),
"the dedup record must be scope-local, got {:?}",
replay.outcomes[0]
);
let mut cross = tenant_change("00000000-0000-4000-8000-000000000021", "cross-scope ghost");
cross.base_version = version_a; cross.updated_at = Utc::now() + Duration::seconds(60);
let response = backend
.apply_push("tenant-b", &push("device-x", vec![cross]), &resolver)
.expect("cross-scope-shaped push");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"an unverifiable base must resolve within its own scope, got {:?}",
response.outcomes[0]
);
};
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("beta"),
"a pre-incarnation base claim must not clobber tenant-b's row"
);
let tenant_b_version = row.version;
let del = Change {
collection: "scoped".to_owned(),
base_version: tenant_b_version,
..change(
"00000000-0000-4000-8000-000000000022",
"s1",
Op::Delete,
None,
)
};
let response = backend
.apply_push("tenant-b", &push("device-a", vec![del]), &resolver)
.expect("tenant-b delete");
assert!(matches!(
response.outcomes[0],
ChangeOutcome::Applied { .. }
));
let PullResponse::Ok { rows, .. } = backend
.pull_since("tenant-b", 0, 100, 0)
.expect("tenant-b post-delete pull")
else {
panic!("cursor 0 never requires a resync");
};
assert!(
rows.len() == 1 && rows[0].deleted,
"tenant-b's row is a tombstone, got {rows:?}"
);
let PullResponse::Ok { rows, .. } = backend
.pull_since("tenant-a", 0, 100, 0)
.expect("tenant-a post-delete pull")
else {
panic!("cursor 0 never requires a resync");
};
assert!(
rows.len() == 1
&& !rows[0].deleted
&& rows[0].version == version_a
&& rows[0]
.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str())
== Some("alpha"),
"tenant-a's row must be untouched by tenant-b's writes, got {rows:?}"
);
let horizon_a_before = backend
.tombstone_horizon("tenant-a")
.expect("tenant-a horizon");
let horizon_global_before = backend.tombstone_horizon(SCOPE).expect("global horizon");
let removed = backend
.gc_tombstones(backend.latest_version().expect("latest"))
.expect("tenant-b gc");
assert!(removed >= 1, "tenant-b's tombstone must be dropped");
let horizon_b = backend
.tombstone_horizon("tenant-b")
.expect("tenant-b horizon");
assert!(
horizon_b > version_a,
"tenant-b's horizon must cover its dropped tombstone"
);
assert_eq!(
backend
.tombstone_horizon("tenant-a")
.expect("tenant-a horizon"),
horizon_a_before,
"another scope's GC must not move tenant-a's horizon"
);
assert_eq!(
backend.tombstone_horizon(SCOPE).expect("global horizon"),
horizon_global_before,
"another scope's GC must not move the global scope's horizon"
);
let mut update_a = tenant_change("00000000-0000-4000-8000-000000000023", "alpha 2");
update_a.base_version = version_a;
let response = backend
.apply_push("tenant-a", &push("device-a", vec![update_a]), &resolver)
.expect("tenant-a update");
assert!(matches!(
response.outcomes[0],
ChangeOutcome::Applied { .. }
));
let mut conflict_a = tenant_change("00000000-0000-4000-8000-000000000024", "alpha 3");
conflict_a.base_version = version_a;
conflict_a.updated_at = Utc::now() + Duration::seconds(60);
let response = backend
.apply_push("tenant-a", &push("device-y", vec![conflict_a]), &resolver)
.expect("tenant-a conflict");
let ChangeOutcome::Resolved { row } = &response.outcomes[0] else {
panic!(
"tenant-a's stale base must resolve, got {:?}",
response.outcomes[0]
);
};
assert_eq!(
row.payload
.as_ref()
.and_then(|p| p.get("title"))
.and_then(|v| v.as_str()),
Some("alpha 3"),
"tenant-a's resolver must still run (and LWW take the client) — \
tenant-b's horizon must not force KeepServer here"
);
let pull = backend
.pull_since("tenant-a", version_a, 100, version_a)
.expect("tenant-a stale-session pull");
assert!(
matches!(pull, PullResponse::Ok { .. }),
"tenant-b's GC must not force tenant-a resyncs, got {pull:?}"
);
let pull = backend
.pull_since("tenant-b", version_b, 100, version_b)
.expect("tenant-b stale-session pull");
assert!(
matches!(pull, PullResponse::FullResyncRequired { .. }),
"tenant-b's own stale session must still resync, got {pull:?}"
);
let removed = backend
.gc_applied(Utc::now() + Duration::seconds(3600))
.expect("gc applied");
assert!(
removed >= 5,
"all dedup records older than the cutoff must go, removed {removed}"
);
assert_eq!(
backend
.gc_applied(Utc::now() + Duration::seconds(3600))
.expect("re-gc applied"),
0,
"dedup GC is idempotent"
);
}
#[test]
fn memory_backend_passes_conformance() {
run_backend_conformance(&MemorySyncBackend::new());
}