Files
hemx/examples/kanban/tests/browser_e2e.rs
T
slhx agent fbc0c9db30 feat(kanban): retry explicit keep-local choice
req: sync/009

req: sync/010

req: sync/011

req: sync/016
2026-07-13 19:47:39 +02:00

2123 lines
100 KiB
Rust

use hemx_test::TestProcess;
use std::fs;
use std::net::TcpListener;
use std::process::Command;
use std::time::Duration;
use thirtyfour::prelude::*;
const STARTUP_TIMEOUT: Duration = Duration::from_secs(12);
#[tokio::test]
async fn server_first_route_does_not_load_optional_client_assets() -> WebDriverResult<()> {
// test req: performance/006
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app.env("HEMX_KANBAN_ADDR", &app_addr);
let _app = TestProcess::start(app, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
driver.find(By::Css("section[data-hemx-root='kanban']")).await?;
let loaded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
navigator.serviceWorker.getRegistrations().then(async (registrations) => {
const databases = indexedDB.databases ? await indexedDB.databases() : [];
done({
resources: performance.getEntriesByType('resource').map((entry) => new URL(entry.name).pathname),
scripts: [...document.scripts].map((script) => ({ src: new URL(script.src).pathname, type: script.type })),
clientRoots: document.querySelectorAll('[data-hemx-client-module]').length,
serviceWorkers: registrations.length,
databases: databases.map((database) => database.name),
});
}).catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert!(loaded["error"].is_null(), "browser inspection failed: {loaded}");
assert_eq!(loaded["clientRoots"], 0);
assert_eq!(loaded["serviceWorkers"], 0);
assert_eq!(loaded["databases"].as_array().map(Vec::len), Some(0));
let scripts = loaded["scripts"].as_array().expect("document scripts");
assert_eq!(scripts.len(), 1);
let runtime_path = scripts[0]["src"].as_str().expect("runtime script path");
assert!(
runtime_path.starts_with("/hemx.") && runtime_path.ends_with(".js"),
"unexpected server runtime asset: {loaded}"
);
assert_eq!(scripts[0]["type"], "");
let resources = loaded["resources"]
.as_array()
.expect("resource timing entries")
.iter()
.filter_map(|value| value.as_str())
.collect::<Vec<_>>();
assert!(resources.contains(&runtime_path), "runtime was not loaded: {loaded}");
for optional in [
"/hemx.client.js",
"/kanban_client.js",
"/kanban_client_bg.wasm",
"/app.js",
"/offline.js",
] {
assert!(
!resources.contains(&optional),
"server-first route loaded optional client asset {optional}: {loaded}"
);
}
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn idempotent_server_command_is_acknowledged_after_reconnect() -> WebDriverResult<()> {
// test req: sync/001 req: sync/005 req: sync/006 req: sync/007
// test req: sync/008 req: sync/012 req: sync/013
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app.env("HEMX_KANBAN_ADDR", &app_addr);
let _app = TestProcess::start(app, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
driver.find(By::Css("section[data-hemx-root='kanban']")).await?;
driver
.execute(
r#"
window.__syncProof = { ready: false, opens: 0, events: [], error: null };
(async () => {
const endpoint = '/sync/commands?command_id=actor-1%3A1&card_id=1';
const firstResponse = await fetch(endpoint, { method: 'POST' });
const first = await firstResponse.json();
const duplicateResponse = await fetch(endpoint, { method: 'POST' });
const duplicate = await duplicateResponse.json();
const conflictResponse = await fetch('/sync/commands?command_id=actor-1%3A1&card_id=2', { method: 'POST' });
const conflict = await conflictResponse.json();
window.__syncProof.command = {
firstStatus: firstResponse.status,
duplicateStatus: duplicateResponse.status,
first,
duplicate,
conflictStatus: conflictResponse.status,
conflict,
};
const source = new EventSource('/sync/acknowledgements?after=0&reconnect=browser-proof');
source.onopen = () => { window.__syncProof.opens += 1; };
source.addEventListener('acknowledgement', (event) => {
window.__syncProof.events.push({ id: event.lastEventId, acknowledgement: JSON.parse(event.data) });
window.__syncProof.ready = true;
source.close();
});
source.onerror = () => {
if (source.readyState === EventSource.CLOSED && !window.__syncProof.ready) {
window.__syncProof.error = 'acknowledgement stream closed';
}
};
})().catch((error) => { window.__syncProof.error = String(error); });
return true;
"#,
Vec::new(),
)
.await?;
wait_until(
&driver,
"return window.__syncProof.ready === true || window.__syncProof.error !== null",
)
.await?;
let proof = driver
.execute("return window.__syncProof", Vec::new())
.await?
.json()
.clone();
assert!(proof["error"].is_null(), "sync failed: {proof}");
assert!(
proof["opens"].as_u64().is_some_and(|opens| opens >= 2),
"transport did not reconnect: {proof}"
);
assert_eq!(proof["command"]["firstStatus"], 200);
assert_eq!(proof["command"]["duplicateStatus"], 200);
assert_eq!(proof["command"]["first"], proof["command"]["duplicate"]);
assert_eq!(proof["command"]["first"]["commandId"], "actor-1:1");
assert_eq!(proof["command"]["first"]["serverSequence"], 1);
assert_eq!(proof["command"]["first"]["cardId"], 1);
assert_eq!(proof["command"]["first"]["canonicalColumn"], "done");
assert_eq!(proof["command"]["first"]["status"], "accepted");
assert_eq!(proof["command"]["conflictStatus"], 409);
assert_eq!(
proof["command"]["conflict"]["error"],
"command_id was already used for a different payload"
);
assert_eq!(proof["events"].as_array().map(Vec::len), Some(1));
assert_eq!(proof["events"][0]["id"], "1");
assert_eq!(
proof["events"][0]["acknowledgement"],
proof["command"]["first"]
);
driver.refresh().await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "3"]));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn pending_local_command_uploads_with_bounded_retry_and_is_removed_on_ack(
) -> WebDriverResult<()> {
// test req: sync/004 req: sync/009 req: sync/010 req: sync/016 req: sync/017
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_FAIL_FIRST_SYNC", "1");
let _app = TestProcess::start(app, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onerror = () => done({ error: open.error && open.error.name });
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'sync-actor:1', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'sync-actor', session: 'sync-session',
causal: 1, kind: 'reorder_card', cardId: '1', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true, "failed to seed durable command: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'acknowledged' || document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'failed'",
)
.await?;
let proof = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), pending: root.getAttribute('data-sync-pending-count'), pendingBeforeAck: root.getAttribute('data-sync-pending-before-ack'), attempts: root.getAttribute('data-sync-attempts'), maxAttempts: root.getAttribute('data-sync-max-attempts'), backoffBase: root.getAttribute('data-sync-last-backoff-base-ms'), backoff: root.getAttribute('data-sync-last-backoff-ms'), opens: root.getAttribute('data-sync-transport-opens'), uploadSequence: root.getAttribute('data-sync-upload-sequence'), ackSequence: root.getAttribute('data-sync-ack-sequence'), canonicalColumn: root.getAttribute('data-sync-canonical-column'), status: root.querySelector('[role=status]').textContent, error: root.getAttribute('data-sync-error') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(proof["phase"], "acknowledged", "sync failed: {proof}");
assert_eq!(proof["pending"], "0");
assert_eq!(proof["pendingBeforeAck"], "1");
assert_eq!(proof["attempts"], "2", "unexpected retry state: {proof}");
assert_eq!(proof["maxAttempts"], "3");
assert_eq!(proof["backoffBase"], "25");
assert!(
proof["backoff"]
.as_str()
.and_then(|value| value.parse::<u64>().ok())
.is_some_and(|delay| (25..50).contains(&delay)),
"retry jitter left its bounded interval: {proof}"
);
assert!(
proof["opens"].as_str().and_then(|value| value.parse::<u64>().ok()).is_some_and(|opens| opens >= 2),
"transport did not reconnect: {proof}"
);
assert_eq!(proof["uploadSequence"], "1");
assert_eq!(proof["ackSequence"], "1");
assert_eq!(proof["canonicalColumn"], "done");
assert_eq!(proof["status"], "Command sync-actor:1 acknowledged in done.");
assert!(proof["error"].is_null());
let queue_count = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const request = open.result.transaction('commands', 'readonly').objectStore('commands').count();
request.onsuccess = () => done(request.result);
request.onerror = () => done({ error: request.error && request.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(queue_count, 0);
let rejected_seed = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'sync-actor:1', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'sync-actor', session: 'sync-session',
causal: 2, kind: 'reorder_card', cardId: '2', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(rejected_seed["seeded"], true, "failed to seed rejected command");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'rejected'",
)
.await?;
let rejected = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { pending: root.getAttribute('data-sync-pending-count'), attempts: root.getAttribute('data-sync-attempts'), error: root.getAttribute('data-sync-error'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(rejected["pending"], "1");
assert_eq!(rejected["attempts"], "1");
assert_eq!(rejected["error"], "sync upload failed with 409");
assert_eq!(
rejected["status"],
"Command sync-actor:1 was permanently rejected (409: command_id was already used for a different payload); 1 durable command remains queued for review."
);
driver.goto(&format!("http://{app_addr}/")).await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "3"]));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn account_partition_hides_replay_and_export_until_owner_returns() -> WebDriverResult<()> {
// test req: sync/020 req: security/004 req: auth/005 req: operations/002
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env(
"HEMX_KANBAN_SESSION_ALICE_ALPHA_EDITOR",
"test-token-alice-alpha-editor",
)
.env(
"HEMX_KANBAN_SESSION_BOB_ALPHA_VIEWER",
"test-token-bob-alpha-viewer",
)
.env(
"HEMX_KANBAN_SESSION_CAROL_BETA_EDITOR",
"test-token-carol-beta-editor",
)
.env("HEMX_KANBAN_SYNC_FAILURES", "3");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
document.cookie = 'hemx_kanban_session=test-token-carol-beta-editor; Path=/; SameSite=Strict';
const open = indexedDB.open('hemx-kanban-v1', 2);
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'auth:1', schemaVersion: 2, accountPartition: 'alpha:alice', actor: 'alice-device', session: 'enqueue-session',
causal: 1, kind: 'reorder_card', cardId: '1', targetColumn: 'done',
eventKind: 'click', key: null, enqueuedPrincipal: 'alice', enqueuedTenant: 'alpha',
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true, "failed to seed auth queue: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'idle'",
)
.await?;
let cross_tenant = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), account: root.getAttribute('data-sync-account-partition'), pending: root.getAttribute('data-sync-pending-count'), attempts: root.getAttribute('data-sync-attempts'), leakedId: document.body.textContent.includes('auth:1'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(cross_tenant["phase"], "idle");
assert_eq!(cross_tenant["account"], "beta:carol");
assert_eq!(cross_tenant["pending"], "0");
assert!(cross_tenant["attempts"].is_null());
assert_eq!(cross_tenant["leakedId"], false);
assert_eq!(cross_tenant["status"], "No pending commands.");
assert_eq!(command_count(&driver).await?, 1);
driver
.execute(
"document.cookie = 'hemx_kanban_session=; Path=/; Max-Age=0; SameSite=Strict'; return true;",
Vec::new(),
)
.await?;
driver.refresh().await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'failed'",
)
.await?;
let signed_out = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { error: root.getAttribute('data-sync-error'), pending: root.getAttribute('data-sync-pending-count'), account: root.getAttribute('data-sync-account-partition') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(signed_out["error"], "account context failed with 401");
assert!(signed_out["pending"].is_null());
assert!(signed_out["account"].is_null());
assert_eq!(command_count(&driver).await?, 1);
let before_authorized = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
fetch('/').then((response) => response.text()).then((html) => {
const page = new DOMParser().parseFromString(html, 'text/html');
done(page.querySelector('[data-key="1"]').closest('section').querySelector('h2').textContent);
}).catch((error) => done(`error:${error}`));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(before_authorized, "Backlog");
driver
.execute(
"document.cookie = 'hemx_kanban_session=test-token-bob-alpha-viewer; Path=/; SameSite=Strict'; return true;",
Vec::new(),
)
.await?;
driver.refresh().await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'idle'",
)
.await?;
let switched_user = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { account: root.getAttribute('data-sync-account-partition'), pending: root.getAttribute('data-sync-pending-count'), attempts: root.getAttribute('data-sync-attempts'), leakedId: document.body.textContent.includes('auth:1') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(switched_user["account"], "alpha:bob");
assert_eq!(switched_user["pending"], "0");
assert!(switched_user["attempts"].is_null());
assert_eq!(switched_user["leakedId"], false);
assert_eq!(command_count(&driver).await?, 1);
let export_boundary = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); const button = root.querySelector('[data-sync-export]'); button.click(); return { exportDisabled: button.disabled, exported: root.getAttribute('data-sync-exported-count'), leakedId: document.body.textContent.includes('auth:1') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(export_boundary["exportDisabled"], true);
assert!(export_boundary["exported"].is_null());
assert_eq!(export_boundary["leakedId"], false);
driver
.execute(
"document.cookie = 'hemx_kanban_session=test-token-alice-alpha-editor; Path=/; SameSite=Strict'; return true;",
Vec::new(),
)
.await?;
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'offline' && root?.getAttribute('data-sync-pending-count') === '1'",
)
.await?;
driver
.execute(
"window.__exportPayload = null; const create = URL.createObjectURL; URL.createObjectURL = (blob) => { blob.text().then((text) => { window.__exportPayload = JSON.parse(text); }); return create(blob); }; return true;",
Vec::new(),
)
.await?;
driver.find(By::Css("[data-sync-export]")).await?.click().await?;
wait_until(&driver, "return window.__exportPayload !== null").await?;
let owner_export = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { payload: window.__exportPayload, exported: root.getAttribute('data-sync-exported-count') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(owner_export["exported"], "1");
assert_eq!(owner_export["payload"]["accountPartition"], "alpha:alice");
assert_eq!(owner_export["payload"]["commands"].as_array().unwrap().len(), 1);
assert_eq!(owner_export["payload"]["commands"][0]["id"], "auth:1");
assert_eq!(owner_export["payload"]["commands"][0]["accountPartition"], "alpha:alice");
driver.find(By::Css("[data-sync-retry]")).await?.click().await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'acknowledged' && root?.getAttribute('data-sync-pending-count') === '0'",
)
.await?;
let authorized = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { sequence: root.getAttribute('data-sync-ack-sequence'), command: root.getAttribute('data-sync-ack-command-id'), column: root.getAttribute('data-sync-canonical-column'), pending: root.getAttribute('data-sync-pending-count') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(authorized["sequence"], "1");
assert_eq!(authorized["command"], "auth:1");
assert_eq!(authorized["column"], "done");
assert_eq!(authorized["pending"], "0");
assert_eq!(command_count(&driver).await?, 0);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn canonical_snapshot_and_history_are_tenant_scoped() -> WebDriverResult<()> {
// test req: sync/007 req: security/004 req: auth/005
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env(
"HEMX_KANBAN_SESSION_ALICE_ALPHA_EDITOR",
"test-token-alice-alpha-editor",
)
.env(
"HEMX_KANBAN_SESSION_BOB_ALPHA_VIEWER",
"test-token-bob-alpha-viewer",
)
.env(
"HEMX_KANBAN_SESSION_CAROL_BETA_EDITOR",
"test-token-carol-beta-editor",
);
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let proof = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const session = (token) => { document.cookie = `hemx_kanban_session=${token}; Path=/; SameSite=Strict`; };
const command = (id, card) => fetch(`/sync/commands?command_id=${encodeURIComponent(id)}&card_id=${card}&column=done`, { method: 'POST' });
(async () => {
session('test-token-alice-alpha-editor');
const aliceWrite = await command('alice-history:1', 1);
session('test-token-carol-beta-editor');
const carolWrite = await command('carol-history:1', 2);
const betaSnapshotResponse = await fetch('/sync/snapshot', { cache: 'no-store' });
const betaSnapshot = await betaSnapshotResponse.json();
const betaHistoryResponse = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
const betaHistory = await betaHistoryResponse.text();
session('test-token-bob-alpha-viewer');
const alphaSnapshotResponse = await fetch('/sync/snapshot', { cache: 'no-store' });
const alphaSnapshot = await alphaSnapshotResponse.json();
const alphaHistoryResponse = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
const alphaHistory = await alphaHistoryResponse.text();
document.cookie = 'hemx_kanban_session=; Path=/; Max-Age=0; SameSite=Strict';
const signedOutSnapshot = await fetch('/sync/snapshot', { cache: 'no-store' });
const signedOutHistory = await fetch('/sync/acknowledgements?after=0', { headers: { Accept: 'text/event-stream' }, cache: 'no-store' });
done({
aliceWrite: aliceWrite.status,
carolWrite: carolWrite.status,
betaSnapshotStatus: betaSnapshotResponse.status,
betaSnapshot,
betaHistoryStatus: betaHistoryResponse.status,
betaHistory,
alphaSnapshotStatus: alphaSnapshotResponse.status,
alphaSnapshot,
alphaHistoryStatus: alphaHistoryResponse.status,
alphaHistory,
signedOutSnapshotStatus: signedOutSnapshot.status,
signedOutHistoryStatus: signedOutHistory.status,
});
})().catch((error) => done({ error: String(error), stack: error.stack }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert!(proof.get("error").is_none(), "browser proof failed: {proof}");
assert_eq!(proof["aliceWrite"], 200);
assert_eq!(proof["carolWrite"], 200);
assert_eq!(proof["betaSnapshotStatus"], 200);
assert_eq!(
proof["betaSnapshot"],
serde_json::json!({
"schemaVersion": 1,
"serverSequence": 1,
"cards": [{ "id": 2, "column": "done" }]
})
);
assert_eq!(proof["betaHistoryStatus"], 200);
let beta_history = proof["betaHistory"].as_str().unwrap();
assert!(beta_history.contains("carol-history:1"));
assert!(!beta_history.contains("alice-history:1"));
assert!(!beta_history.contains("\"cardId\":1"));
assert_eq!(proof["alphaSnapshotStatus"], 200);
assert_eq!(
proof["alphaSnapshot"],
serde_json::json!({
"schemaVersion": 1,
"serverSequence": 1,
"cards": [
{ "id": 1, "column": "done" },
{ "id": 3, "column": "done" }
]
})
);
assert_eq!(proof["alphaHistoryStatus"], 200);
let alpha_history = proof["alphaHistory"].as_str().unwrap();
assert!(alpha_history.contains("alice-history:1"));
assert!(!alpha_history.contains("carol-history:1"));
assert!(!alpha_history.contains("\"cardId\":2"));
assert_eq!(proof["signedOutSnapshotStatus"], 401);
assert_eq!(proof["signedOutHistoryStatus"], 401);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn schema_upgrade_preserves_queued_order_and_local_intent() -> WebDriverResult<()> {
// test req: sync/004 req: sync/010 req: sync/014
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_SYNC_FAILURES", "3");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1', 1);
open.onupgradeneeded = () => {
const database = open.result;
database.createObjectStore('commands', { keyPath: 'id' });
database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction(['commands', 'meta'], 'readwrite');
const commands = tx.objectStore('commands');
for (const [causal, cardId, eventKind, key] of [[1, '1', 'click', null], [2, '2', 'keydown', 'Enter'], [3, '3', 'click', null]]) {
commands.add({
id: `upgrade:${causal}`, schemaVersion: 1, actor: 'upgrade', session: 'legacy-session',
causal, kind: 'reorder_card', cardId, eventKind, key,
});
}
tx.objectStore('meta').put(3, 'causal');
tx.oncomplete = () => done({ version: open.result.version });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["version"], 1, "failed to seed legacy queue: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'offline'",
)
.await?;
let migrated = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const root = document.querySelector('[data-kanban-sync]');
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const database = open.result;
const tx = database.transaction(['commands', 'meta'], 'readonly');
const commands = tx.objectStore('commands').getAll();
const receipt = tx.objectStore('meta').get('commandSchemaMigration');
tx.oncomplete = () => done({
databaseVersion: database.version,
commandSchema: root.getAttribute('data-sync-command-schema'),
migrationFrom: root.getAttribute('data-sync-migration-from'),
migrationTo: root.getAttribute('data-sync-migration-to'),
migratedCount: root.getAttribute('data-sync-migrated-count'),
pending: root.getAttribute('data-sync-pending-count'),
receipt: receipt.result,
commands: commands.result.sort((a, b) => a.causal - b.causal).map(({ id, schemaVersion, cardId, targetColumn, eventKind, key }) => ({ id, schemaVersion, cardId, targetColumn, eventKind, key })),
});
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(migrated["databaseVersion"], 3);
assert_eq!(migrated["commandSchema"], "2");
assert_eq!(migrated["migrationFrom"], "1");
assert_eq!(migrated["migrationTo"], "3");
assert_eq!(migrated["migratedCount"], "3");
assert_eq!(migrated["pending"], "3");
assert_eq!(
migrated["receipt"],
serde_json::json!({ "from": 1, "to": 3, "migrated": 3 })
);
assert_eq!(
migrated["commands"],
serde_json::json!([
{ "id": "upgrade:1", "schemaVersion": 2, "cardId": "1", "targetColumn": "done", "eventKind": "click", "key": null },
{ "id": "upgrade:2", "schemaVersion": 2, "cardId": "2", "targetColumn": "done", "eventKind": "keydown", "key": "Enter" },
{ "id": "upgrade:3", "schemaVersion": 2, "cardId": "3", "targetColumn": "done", "eventKind": "click", "key": null }
])
);
driver.find(By::Css("[data-sync-retry]")).await?.click().await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'backpressured' && root?.getAttribute('data-sync-pending-count') === '1'",
)
.await?;
driver.find(By::Css("[data-sync-retry]")).await?.click().await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'acknowledged' && root?.getAttribute('data-sync-pending-count') === '0'",
)
.await?;
driver
.execute(
r#"
window.__upgradeReplay = [];
const source = new EventSource('/sync/acknowledgements?after=0');
source.addEventListener('acknowledgement', (event) => {
window.__upgradeReplay.push({ id: event.lastEventId, body: JSON.parse(event.data) });
if (window.__upgradeReplay.length === 3) source.close();
});
return true;
"#,
Vec::new(),
)
.await?;
wait_until(&driver, "return window.__upgradeReplay.length === 3").await?;
let replay = driver
.execute("return window.__upgradeReplay", Vec::new())
.await?
.json()
.clone();
assert_eq!(
replay,
serde_json::json!([
{ "id": "1", "body": { "commandId": "upgrade:1", "cardId": 1, "canonicalColumn": "done", "serverSequence": 1, "status": "accepted" } },
{ "id": "2", "body": { "commandId": "upgrade:2", "cardId": 2, "canonicalColumn": "done", "serverSequence": 2, "status": "accepted" } },
{ "id": "3", "body": { "commandId": "upgrade:3", "cardId": 3, "canonicalColumn": "done", "serverSequence": 3, "status": "accepted" } }
])
);
assert_eq!(command_count(&driver).await?, 0);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn mixed_queue_removes_accepted_prefix_and_retains_rejected_tail() -> WebDriverResult<()> {
// test req: sync/009 req: sync/010
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command.env("HEMX_KANBAN_ADDR", &app_addr);
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
const commands = tx.objectStore('commands');
for (const [causal, cardId] of [[1, '1'], [2, '999'], [3, '2']]) {
commands.add({
id: `mixed:${causal}`, schemaVersion: 2, accountPartition: 'demo:demo', actor: 'mixed', session: 'mixed-session',
causal, kind: 'reorder_card', cardId, targetColumn: 'done', eventKind: 'click', key: null,
});
}
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true, "failed to seed mixed queue: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'rejected'",
)
.await?;
tokio::time::sleep(Duration::from_millis(150)).await;
let rejected = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); const retry = root.querySelector('[data-sync-retry]'); return { phase: root.getAttribute('data-sync-phase'), pending: root.getAttribute('data-sync-pending-count'), uploaded: root.getAttribute('data-sync-uploaded-total'), ackSequence: root.getAttribute('data-sync-ack-sequence'), attempts: root.getAttribute('data-sync-attempts'), inFlight: root.getAttribute('data-sync-in-flight'), maxInFlight: root.getAttribute('data-sync-max-observed-in-flight'), kind: root.getAttribute('data-sync-error-kind'), errorStatus: root.getAttribute('data-sync-error-status'), reason: root.getAttribute('data-sync-error-reason'), rejectedId: root.getAttribute('data-sync-rejected-command-id'), retryDisabled: retry.disabled, status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(rejected["phase"], "rejected");
assert_eq!(rejected["pending"], "2");
assert_eq!(rejected["uploaded"], "1");
assert_eq!(rejected["ackSequence"], "1");
assert_eq!(rejected["attempts"], "1");
assert_eq!(rejected["inFlight"], "0");
assert_eq!(rejected["maxInFlight"], "1");
assert_eq!(rejected["kind"], "permanent-rejection");
assert_eq!(rejected["errorStatus"], "400");
assert_eq!(rejected["reason"], "unknown card_id");
assert_eq!(rejected["rejectedId"], "mixed:2");
assert_eq!(rejected["retryDisabled"], true);
assert_eq!(
rejected["status"],
"Command mixed:2 was permanently rejected (400: unknown card_id); 2 durable commands remain queued for review."
);
let queued = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const request = open.result.transaction('commands', 'readonly').objectStore('commands').getAll();
request.onsuccess = () => done(request.result.sort((a, b) => a.causal - b.causal).map(({ id, cardId }) => ({ id, cardId })));
request.onerror = () => done({ error: request.error && request.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(
queued,
serde_json::json!([
{ "id": "mixed:2", "cardId": "999" },
{ "id": "mixed:3", "cardId": "2" }
])
);
driver.goto(&format!("http://{app_addr}/")).await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[1]["title"], "Doing");
assert_eq!(canonical[1]["cards"], serde_json::json!(["2"]));
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "3"]));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn upload_backpressure_keeps_pending_work_visible_and_recoverable() -> WebDriverResult<()> {
// test req: sync/017
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command.env("HEMX_KANBAN_ADDR", &app_addr);
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
const commands = tx.objectStore('commands');
for (let causal = 1; causal <= 3; causal += 1) {
commands.add({
id: `pressure:${causal}`, schemaVersion: 2, accountPartition: 'demo:demo', actor: 'pressure', session: 'pressure-session',
causal, kind: 'reorder_card', cardId: String(causal), targetColumn: 'done', eventKind: 'click', key: null,
});
}
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true, "failed to seed pending work: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'backpressured'",
)
.await?;
let bounded = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), pending: root.getAttribute('data-sync-pending-count'), limit: root.getAttribute('data-sync-upload-limit'), uploaded: root.getAttribute('data-sync-uploaded-this-run'), total: root.getAttribute('data-sync-uploaded-total'), maxInFlight: root.getAttribute('data-sync-max-observed-in-flight'), sequence: root.getAttribute('data-sync-ack-sequence'), manual: root.getAttribute('data-sync-manual-retry'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(bounded["phase"], "backpressured");
assert_eq!(bounded["pending"], "1");
assert_eq!(bounded["limit"], "2");
assert_eq!(bounded["uploaded"], "2");
assert_eq!(bounded["total"], "2");
assert_eq!(bounded["maxInFlight"], "1");
assert_eq!(bounded["sequence"], "2");
assert_eq!(bounded["manual"], "available");
assert_eq!(
bounded["status"],
"Upload limit 2 reached; 1 durable command remains queued. Retry now to continue."
);
assert_eq!(command_count(&driver).await?, 1);
driver.find(By::Css("[data-sync-retry]")).await?.click().await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'acknowledged' && root?.getAttribute('data-sync-pending-count') === '0'",
)
.await?;
let recovered = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { pending: root.getAttribute('data-sync-pending-count'), uploaded: root.getAttribute('data-sync-uploaded-this-run'), total: root.getAttribute('data-sync-uploaded-total'), maxInFlight: root.getAttribute('data-sync-max-observed-in-flight'), sequence: root.getAttribute('data-sync-ack-sequence'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(recovered["pending"], "0");
assert_eq!(recovered["uploaded"], "1");
assert_eq!(recovered["total"], "3");
assert_eq!(recovered["maxInFlight"], "1");
assert_eq!(recovered["sequence"], "3");
assert_eq!(recovered["status"], "Command pressure:3 acknowledged in done.");
assert_eq!(command_count(&driver).await?, 0);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn two_tabs_coordinate_single_uploader_and_takeover_without_duplicate_application(
) -> WebDriverResult<()> {
// test req: sync/018
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_SYNC_FAILURES", "3");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'tabs:1', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'tabs', session: 'tabs-session',
causal: 1, kind: 'reorder_card', cardId: '1', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true);
driver
.execute(
"sessionStorage.setItem('hemx-kanban-sync-tab-id', 'leader-seed'); return true;",
Vec::new(),
)
.await?;
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'offline'",
)
.await?;
let leader = driver.window().await?;
let follower = driver.new_tab().await?;
driver.switch_to_window(follower.clone()).await?;
driver.goto(&format!("http://{app_addr}/")).await?;
driver
.execute(
"sessionStorage.setItem('hemx-kanban-sync-tab-id', 'follower-seed'); return true;",
Vec::new(),
)
.await?;
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'standby'",
)
.await?;
let standby = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), leader: root.getAttribute('data-sync-leader'), attempts: root.getAttribute('data-sync-attempts'), owner: root.getAttribute('data-sync-lease-owner'), tab: root.getAttribute('data-sync-tab-id'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(standby["phase"], "standby");
assert_eq!(standby["leader"], "false");
assert!(standby["attempts"].is_null());
assert_ne!(standby["owner"], standby["tab"]);
assert_eq!(standby["status"], "Another tab owns sync; waiting for lease takeover.");
driver.switch_to_window(leader).await?;
let first = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), leader: root.getAttribute('data-sync-leader'), attempts: root.getAttribute('data-sync-attempts'), pending: root.getAttribute('data-sync-pending-count'), owner: root.getAttribute('data-sync-lease-owner'), tab: root.getAttribute('data-sync-tab-id') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(first["phase"], "offline", "leader lost ownership: {first}");
assert_eq!(first["leader"], "true");
assert_eq!(first["attempts"], "3");
assert_eq!(first["pending"], "1");
driver.close_window().await?;
driver.switch_to_window(follower).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'acknowledged'",
)
.await?;
let takeover = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), attempts: root.getAttribute('data-sync-attempts'), pending: root.getAttribute('data-sync-pending-count'), sequence: root.getAttribute('data-sync-ack-sequence'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(takeover["phase"], "acknowledged");
assert_eq!(takeover["attempts"], "1");
assert_eq!(takeover["pending"], "0");
assert_eq!(takeover["sequence"], "1");
assert_eq!(takeover["status"], "Command tabs:1 acknowledged in done.");
assert_eq!(command_count(&driver).await?, 0);
driver
.execute(
r#"
window.__tabReplay = [];
const source = new EventSource('/sync/acknowledgements?after=0');
source.addEventListener('acknowledgement', (event) => {
window.__tabReplay.push({ id: event.lastEventId, body: JSON.parse(event.data) });
setTimeout(() => source.close(), 25);
});
return true;
"#,
Vec::new(),
)
.await?;
wait_until(&driver, "return window.__tabReplay.length === 1").await?;
tokio::time::sleep(Duration::from_millis(50)).await;
let replay = driver
.execute("return window.__tabReplay", Vec::new())
.await?
.json()
.clone();
assert_eq!(replay.as_array().map(Vec::len), Some(1));
assert_eq!(replay[0]["id"], "1");
assert_eq!(replay[0]["body"]["commandId"], "tabs:1");
driver.goto(&format!("http://{app_addr}/")).await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "3"]));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn exhausted_offline_retries_keep_command_until_later_reconnect() -> WebDriverResult<()> {
// test req: sync/004 req: sync/010 req: sync/011 req: sync/014 req: sync/016
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_SYNC_FAILURES", "3");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'offline-actor:1', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'offline-actor', session: 'offline-session',
causal: 1, kind: 'reorder_card', cardId: '2', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded["seeded"], true, "failed to seed offline command: {seeded}");
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'offline'",
)
.await?;
let exhausted = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), connection: root.getAttribute('data-sync-connection'), pending: root.getAttribute('data-sync-pending-count'), attempts: root.getAttribute('data-sync-attempts'), maxAttempts: root.getAttribute('data-sync-max-attempts'), manual: root.getAttribute('data-sync-manual-retry'), error: root.getAttribute('data-sync-error'), status: root.querySelector('[role=status]').textContent, retry: root.querySelector('[data-sync-retry]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(exhausted["phase"], "offline");
assert_eq!(exhausted["connection"], "offline");
assert_eq!(exhausted["pending"], "1");
assert_eq!(exhausted["attempts"], "3");
assert_eq!(exhausted["maxAttempts"], "3");
assert_eq!(exhausted["manual"], "available");
assert_eq!(exhausted["error"], "sync upload failed with 503");
assert_eq!(
exhausted["status"],
"Sync is offline after bounded retries; the durable command remains queued. Retry now when ready."
);
assert_eq!(exhausted["retry"], "Retry sync now");
let queued = command_count(&driver).await?;
assert_eq!(queued, 1);
driver
.find(By::Css("[data-sync-retry]"))
.await?
.click()
.await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'acknowledged'",
)
.await?;
let converged = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { connection: root.getAttribute('data-sync-connection'), pending: root.getAttribute('data-sync-pending-count'), attempts: root.getAttribute('data-sync-attempts'), sequence: root.getAttribute('data-sync-ack-sequence'), column: root.getAttribute('data-sync-canonical-column'), status: root.querySelector('[role=status]').textContent, error: root.getAttribute('data-sync-error') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(converged["connection"], "online");
assert_eq!(converged["pending"], "0");
assert_eq!(converged["attempts"], "1");
assert_eq!(converged["sequence"], "1");
assert_eq!(converged["column"], "done");
assert_eq!(converged["status"], "Command offline-actor:1 acknowledged in done.");
assert!(converged["error"].is_null());
assert_eq!(command_count(&driver).await?, 0);
driver.goto(&format!("http://{app_addr}/")).await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["2", "3"]));
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn missing_history_rebase_and_user_conflict_resolution_preserve_suffix() -> WebDriverResult<()>
{
// test req: sync/007 req: sync/010 req: sync/011 req: sync/020
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_RETAINED_AFTER", "1");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded_server = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
fetch('/sync/commands?command_id=history%3A1&card_id=1', { method: 'POST' })
.then(async (response) => done({ status: response.status, body: await response.json() }))
.catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded_server["status"], 200);
assert_eq!(seeded_server["body"]["serverSequence"], 1);
let seeded_local = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onupgradeneeded = () => {
const database = open.result;
if (!database.objectStoreNames.contains('commands')) database.createObjectStore('commands', { keyPath: 'id' });
if (!database.objectStoreNames.contains('meta')) database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
tx.objectStore('commands').add({
id: 'history:2', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'history', session: 'history-session',
causal: 2, kind: 'reorder_card', cardId: '2', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded_local["seeded"], true);
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'rebased'",
)
.await?;
let fallback = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), uploadSequence: root.getAttribute('data-sync-upload-sequence'), ackSequence: root.getAttribute('data-sync-ack-sequence'), snapshotSequence: root.getAttribute('data-sync-snapshot-sequence'), snapshotSchema: root.getAttribute('data-sync-snapshot-schema'), snapshotCards: root.getAttribute('data-sync-snapshot-card-count'), pending: root.getAttribute('data-sync-pending-count'), rebasePending: root.getAttribute('data-sync-rebase-pending-count'), decision: root.getAttribute('data-sync-rebase-decision'), reason: root.getAttribute('data-sync-rebase-reason'), canonicalColumn: root.getAttribute('data-sync-canonical-column'), status: root.querySelector('[role=status]').textContent, error: root.getAttribute('data-sync-error') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(fallback["phase"], "rebased");
assert_eq!(fallback["uploadSequence"], "2");
assert_eq!(fallback["ackSequence"], "2");
assert_eq!(fallback["snapshotSequence"], "2");
assert_eq!(fallback["snapshotSchema"], "1");
assert_eq!(fallback["snapshotCards"], "3");
assert_eq!(fallback["pending"], "0");
assert_eq!(fallback["rebasePending"], "1", "unexpected rebase state: {fallback}");
assert_eq!(fallback["decision"], "converged");
assert_eq!(fallback["reason"], "intent-already-canonical");
assert_eq!(fallback["canonicalColumn"], "done");
assert_eq!(
fallback["status"],
"Canonical snapshot 2 already satisfies history:2; committed and removed the pending command."
);
assert!(fallback["error"].is_null());
assert_eq!(command_count(&driver).await?, 0);
let committed = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const tx = open.result.transaction('meta', 'readonly');
const meta = tx.objectStore('meta');
const cursor = meta.get('acknowledgementCursor');
const snapshot = meta.get('canonicalSnapshot');
tx.oncomplete = () => done({ cursor: cursor.result, snapshot: snapshot.result });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(committed["cursor"], 2);
assert_eq!(committed["snapshot"]["schemaVersion"], 1);
assert_eq!(committed["snapshot"]["serverSequence"], 2);
driver.goto(&format!("http://{app_addr}/")).await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "2", "3"]));
let divergent_server = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
(async () => {
const local = await fetch('/sync/commands?command_id=history%3A3&card_id=2&column=done', { method: 'POST' });
const remote = await fetch('/sync/commands?command_id=remote%3A4&card_id=2&column=doing', { method: 'POST' });
done({
local: { status: local.status, body: await local.json() },
remote: { status: remote.status, body: await remote.json() },
});
})().catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(divergent_server["local"]["status"], 200);
assert_eq!(divergent_server["local"]["body"]["serverSequence"], 3);
assert_eq!(divergent_server["remote"]["status"], 200);
assert_eq!(divergent_server["remote"]["body"]["serverSequence"], 4);
assert_eq!(divergent_server["remote"]["body"]["canonicalColumn"], "doing");
let seeded_conflict = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
const commands = tx.objectStore('commands');
commands.add({
id: 'history:3', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'history', session: 'history-session',
causal: 3, kind: 'reorder_card', cardId: '2', targetColumn: 'done', eventKind: 'click', key: null,
});
commands.add({
id: 'history:4', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'history', session: 'history-session',
causal: 4, kind: 'reorder_card', cardId: '1', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded_conflict["seeded"], true);
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'conflicted'",
)
.await?;
let conflicted = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { phase: root.getAttribute('data-sync-phase'), uploadSequence: root.getAttribute('data-sync-upload-sequence'), snapshotSequence: root.getAttribute('data-sync-snapshot-sequence'), pending: root.getAttribute('data-sync-pending-count'), rebasePending: root.getAttribute('data-sync-rebase-pending-count'), decision: root.getAttribute('data-sync-rebase-decision'), reason: root.getAttribute('data-sync-rebase-reason'), canonicalColumn: root.getAttribute('data-sync-canonical-column'), resolutionDisabled: root.querySelector('[data-sync-use-canonical]').disabled, status: root.querySelector('[role=status]').textContent, error: root.getAttribute('data-sync-error') }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(conflicted["phase"], "conflicted", "unexpected conflict state: {conflicted}");
assert_eq!(conflicted["uploadSequence"], "3");
assert_eq!(conflicted["snapshotSequence"], "4");
assert_eq!(conflicted["pending"], "2");
assert_eq!(conflicted["rebasePending"], "2");
assert_eq!(conflicted["decision"], "conflicted");
assert_eq!(conflicted["reason"], "canonical-state-diverged");
assert_eq!(conflicted["canonicalColumn"], "doing");
assert_eq!(conflicted["resolutionDisabled"], false);
assert_eq!(
conflicted["status"],
"Canonical snapshot 4 conflicts with history:3 (canonical-state-diverged); the pending command remains queued."
);
assert!(conflicted["error"].is_null());
assert_eq!(command_count(&driver).await?, 2);
let preserved = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const tx = open.result.transaction(['commands', 'meta'], 'readonly');
const commands = tx.objectStore('commands').getAll();
const cursor = tx.objectStore('meta').get('acknowledgementCursor');
const snapshot = tx.objectStore('meta').get('canonicalSnapshot');
tx.oncomplete = () => done({ commands: commands.result.sort((left, right) => left.causal - right.causal), cursor: cursor.result, snapshot: snapshot.result });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(preserved["commands"][0]["id"], "history:3");
assert_eq!(preserved["commands"][1]["id"], "history:4");
assert_eq!(preserved["cursor"], 2);
assert_eq!(preserved["snapshot"]["serverSequence"], 2);
driver.goto(&format!("http://{app_addr}/")).await?;
let divergent_board = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(divergent_board[1]["title"], "Doing");
assert_eq!(divergent_board[1]["cards"], serde_json::json!(["2"]));
assert_eq!(divergent_board[2]["cards"], serde_json::json!(["1", "3"]));
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'conflicted'",
)
.await?;
driver
.find(By::Css("[data-sync-use-canonical]"))
.await?
.click()
.await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'rebased' && root?.getAttribute('data-sync-pending-count') === '0'",
)
.await?;
let resolved = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { resolution: root.getAttribute('data-sync-conflict-resolution'), resolvedCommand: root.getAttribute('data-sync-resolved-command-id'), pending: root.getAttribute('data-sync-pending-count'), ackSequence: root.getAttribute('data-sync-ack-sequence'), canonicalColumn: root.getAttribute('data-sync-canonical-column'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(resolved["resolution"], "used-canonical-state");
assert_eq!(resolved["resolvedCommand"], "history:3");
assert_eq!(resolved["pending"], "0");
assert_eq!(resolved["ackSequence"], "5");
assert_eq!(resolved["canonicalColumn"], "done");
assert_eq!(
resolved["status"],
"Canonical snapshot 5 already satisfies history:4; committed and removed the pending command."
);
assert_eq!(command_count(&driver).await?, 0);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn keep_local_retry_preserves_conflicted_command_and_suffix_order() -> WebDriverResult<()> {
// test req: sync/009 req: sync/010 req: sync/011 req: sync/016
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let mut app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_RETAINED_AFTER", "1");
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
.expect("start ready hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let seeded_server = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const command = (id, card, column = 'done') => fetch(`/sync/commands?command_id=${encodeURIComponent(id)}&card_id=${card}&column=${column}`, { method: 'POST' });
(async () => {
const statuses = [];
statuses.push((await command('keep:1', 1)).status);
statuses.push((await command('keep:2', 2)).status);
statuses.push((await command('keep:3', 2)).status);
statuses.push((await command('keep:external', 2, 'doing')).status);
done({ statuses });
})().catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded_server["statuses"], serde_json::json!([200, 200, 200, 200]));
let seeded_local = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1', 3);
open.onupgradeneeded = () => {
const database = open.result;
const commands = database.createObjectStore('commands', { keyPath: 'id' });
commands.createIndex('byAccountPartition', 'accountPartition');
database.createObjectStore('meta');
};
open.onsuccess = () => {
const tx = open.result.transaction('commands', 'readwrite');
const commands = tx.objectStore('commands');
commands.add({
id: 'keep:3', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'keep', session: 'keep-session',
causal: 3, kind: 'reorder_card', cardId: '2', targetColumn: 'done', eventKind: 'click', key: null,
});
commands.add({
id: 'keep:4', schemaVersion: 2, accountPartition: 'demo:demo', actor: 'keep', session: 'keep-session',
causal: 4, kind: 'reorder_card', cardId: '1', targetColumn: 'done', eventKind: 'click', key: null,
});
tx.oncomplete = () => done({ seeded: true });
tx.onabort = () => done({ error: tx.error && tx.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(seeded_local["seeded"], true);
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'conflicted'",
)
.await?;
driver
.execute(
r#"
window.__keepRejections = 0;
window.__keepFailures = 0;
const originalFetch = window.fetch.bind(window);
window.fetch = (input, init = {}) => {
const url = new URL(typeof input === 'string' ? input : input.url, location.href);
if (url.pathname === '/sync/commands' && url.searchParams.get('command_id') === 'keep:3:keep:4') {
if (window.__keepRejections < 1) {
window.__keepRejections += 1;
return Promise.resolve(new Response(JSON.stringify({ kind: 'command-conflict', error: 'injected stale keep-local decision' }), {
status: 409,
headers: { 'content-type': 'application/json' },
}));
}
if (window.__keepFailures < 3) {
window.__keepFailures += 1;
return Promise.resolve(new Response(JSON.stringify({ kind: 'transport-failure', error: 'injected keep-local retry' }), {
status: 503,
headers: { 'content-type': 'application/json' },
}));
}
}
return originalFetch(input, init);
};
return true;
"#,
Vec::new(),
)
.await?;
driver
.find(By::Css("[data-sync-keep-local]"))
.await?
.click()
.await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'resolution-rejected'",
)
.await?;
let rejected = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const root = document.querySelector('[data-kanban-sync]');
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const request = open.result.transaction('commands', 'readonly').objectStore('commands').getAll();
request.onsuccess = () => done({
resolution: root.getAttribute('data-sync-conflict-resolution'),
resolutionDisabled: root.querySelector('[data-sync-keep-local]').disabled,
pending: root.getAttribute('data-sync-pending-count'),
rejections: window.__keepRejections,
commands: request.result.sort((left, right) => left.causal - right.causal).map(({ id, causal, cardId }) => ({ id, causal, cardId })),
});
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(rejected["resolution"], "keep-local-rejected");
assert_eq!(rejected["resolutionDisabled"], false);
assert_eq!(rejected["pending"], "2");
assert_eq!(rejected["rejections"], 1);
assert_eq!(
rejected["commands"],
serde_json::json!([
{ "id": "keep:3", "causal": 3, "cardId": "2" },
{ "id": "keep:4", "causal": 4, "cardId": "1" }
])
);
driver
.find(By::Css("[data-sync-keep-local]"))
.await?
.click()
.await?;
wait_until(
&driver,
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'offline'",
)
.await?;
let retrying = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const root = document.querySelector('[data-kanban-sync]');
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const request = open.result.transaction('commands', 'readonly').objectStore('commands').getAll();
request.onsuccess = () => done({
phase: root.getAttribute('data-sync-phase'),
resolution: root.getAttribute('data-sync-conflict-resolution'),
resolutionCommand: root.getAttribute('data-sync-resolution-command-id'),
resolvedCommand: root.getAttribute('data-sync-resolved-command-id'),
attempts: root.getAttribute('data-sync-attempts'),
pending: root.getAttribute('data-sync-pending-count'),
manualRetry: root.getAttribute('data-sync-manual-retry'),
failures: window.__keepFailures,
commands: request.result.sort((left, right) => left.causal - right.causal).map(({ id, causal, cardId }) => ({ id, causal, cardId })),
});
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(retrying["phase"], "offline");
assert_eq!(retrying["resolution"], "keep-local-pending");
assert_eq!(retrying["resolutionCommand"], "keep:3:keep:4");
assert_eq!(retrying["resolvedCommand"], "keep:3");
assert_eq!(retrying["attempts"], "3");
assert_eq!(retrying["pending"], "2");
assert_eq!(retrying["manualRetry"], "available");
assert_eq!(retrying["failures"], 3);
assert_eq!(
retrying["commands"],
serde_json::json!([
{ "id": "keep:3", "causal": 3, "cardId": "2" },
{ "id": "keep:4", "causal": 4, "cardId": "1" }
])
);
driver.find(By::Css("[data-sync-retry]")).await?.click().await?;
wait_until(
&driver,
"const root = document.querySelector('[data-kanban-sync]'); return root?.getAttribute('data-sync-phase') === 'rebased' && root?.getAttribute('data-sync-pending-count') === '0'",
)
.await?;
let resolved = driver
.execute(
"const root = document.querySelector('[data-kanban-sync]'); return { resolution: root.getAttribute('data-sync-conflict-resolution'), resolvedCommand: root.getAttribute('data-sync-resolved-command-id'), pending: root.getAttribute('data-sync-pending-count'), ackSequence: root.getAttribute('data-sync-ack-sequence'), canonicalColumn: root.getAttribute('data-sync-canonical-column'), status: root.querySelector('[role=status]').textContent }",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(resolved["resolution"], "kept-local-change");
assert_eq!(resolved["resolvedCommand"], "keep:3");
assert_eq!(resolved["pending"], "0");
assert_eq!(resolved["ackSequence"], "6");
assert_eq!(resolved["canonicalColumn"], "done");
assert_eq!(
resolved["status"],
"Canonical snapshot 6 already satisfies keep:4; committed and removed the pending command."
);
assert_eq!(command_count(&driver).await?, 0);
Ok(())
}
.await;
let quit = driver.quit().await;
result.and(quit)
}
#[tokio::test]
async fn canonical_acknowledgement_survives_server_restart() -> WebDriverResult<()> {
// test req: sync/001 req: sync/005 req: sync/007 req: sync/008 req: sync/013
let app_port = available_port();
let app_addr = format!("127.0.0.1:{app_port}");
let store = std::env::temp_dir().join(format!(
"hemx-kanban-sync-{}-{app_port}.json",
std::process::id()
));
let _ = fs::remove_file(&store);
let mut first_app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
first_app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_SYNC_STORE", &store);
let first_app = TestProcess::start(
first_app_command,
"hemx-kanban-first",
&app_addr,
STARTUP_TIMEOUT,
)
.expect("start first hemx-kanban");
let webdriver_port = available_port();
let webdriver_addr = format!("127.0.0.1:{webdriver_port}");
let mut webdriver = Command::new("geckodriver");
webdriver.arg("--port").arg(webdriver_port.to_string());
let _webdriver = TestProcess::start(webdriver, "geckodriver", &webdriver_addr, STARTUP_TIMEOUT)
.expect("start ready geckodriver");
let mut caps = DesiredCapabilities::firefox();
caps.set_headless()?;
let driver = WebDriver::new(&format!("http://{webdriver_addr}"), caps).await?;
let result = async {
driver.goto(&format!("http://{app_addr}/")).await?;
let accepted = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
fetch('/sync/commands?command_id=restart-proof%3A1&card_id=1', { method: 'POST' })
.then(async (response) => done({ status: response.status, body: await response.json() }))
.catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(accepted["status"], 200);
assert_eq!(accepted["body"]["serverSequence"], 1);
assert_eq!(accepted["body"]["canonicalColumn"], "done");
assert!(store.is_file(), "server did not materialize sync store");
let persisted = fs::read_to_string(&store).expect("read sync store");
assert!(persisted.contains("restart-proof:1"));
assert!(persisted.contains("\"schemaVersion\": 2"));
assert!(persisted.contains("\"tenant\": \"demo\""));
drop(first_app);
let mut second_app_command = Command::new(env!("CARGO_BIN_EXE_hemx-kanban-example"));
second_app_command
.env("HEMX_KANBAN_ADDR", &app_addr)
.env("HEMX_KANBAN_SYNC_STORE", &store);
let second_app = TestProcess::start(
second_app_command,
"hemx-kanban-second",
&app_addr,
STARTUP_TIMEOUT,
)
.expect("restart hemx-kanban from durable sync store");
let after_restart = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
fetch('/sync/commands?command_id=restart-proof%3A1&card_id=1', { method: 'POST' })
.then(async (response) => done({ status: response.status, body: await response.json() }))
.catch((error) => done({ error: String(error) }));
"#,
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(after_restart["status"], 200);
assert_eq!(after_restart["body"], accepted["body"]);
driver
.execute(
r#"
window.__restartReplay = null;
const source = new EventSource('/sync/acknowledgements?after=0');
source.addEventListener('acknowledgement', (event) => {
window.__restartReplay = { id: event.lastEventId, body: JSON.parse(event.data) };
source.close();
});
return true;
"#,
Vec::new(),
)
.await?;
wait_until(&driver, "return window.__restartReplay !== null").await?;
let replay = driver
.execute("return window.__restartReplay", Vec::new())
.await?
.json()
.clone();
assert_eq!(replay["id"], "1");
assert_eq!(replay["body"], accepted["body"]);
driver.refresh().await?;
let canonical = driver
.execute(
"return [...document.querySelectorAll('section.column')].map((column) => ({ title: column.querySelector('h2').textContent, cards: [...column.querySelectorAll('[data-key]')].map((card) => card.dataset.key) }))",
Vec::new(),
)
.await?
.json()
.clone();
assert_eq!(canonical[2]["title"], "Done");
assert_eq!(canonical[2]["cards"], serde_json::json!(["1", "3"]));
drop(second_app);
Ok(())
}
.await;
let quit = driver.quit().await;
let _ = fs::remove_file(&store);
result.and(quit)
}
async fn command_count(driver: &WebDriver) -> WebDriverResult<u64> {
let count = driver
.execute_async(
r#"
const done = arguments[arguments.length - 1];
const open = indexedDB.open('hemx-kanban-v1');
open.onsuccess = () => {
const request = open.result.transaction('commands', 'readonly').objectStore('commands').count();
request.onsuccess = () => done(request.result);
request.onerror = () => done({ error: request.error && request.error.name });
};
"#,
Vec::new(),
)
.await?
.json()
.clone();
Ok(count.as_u64().expect("durable command count"))
}
async fn wait_until(driver: &WebDriver, script: &str) -> WebDriverResult<()> {
for _ in 0..200 {
if driver.execute(script, Vec::new()).await?.json().as_bool() == Some(true) {
return Ok(());
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
let snapshot = driver
.execute(
"const sync = document.querySelector('[data-kanban-sync]'); return { url: location.href, body: document.body.textContent, sync: sync ? Object.fromEntries([...sync.attributes].map((attribute) => [attribute.name, attribute.value])) : null }",
Vec::new(),
)
.await?
.json()
.clone();
panic!("browser condition timed out: {script}; snapshot: {snapshot}");
}
fn available_port() -> u16 {
TcpListener::bind("127.0.0.1:0")
.expect("reserve browser test port")
.local_addr()
.expect("browser test address")
.port()
}