fix(sync): bound acknowledgement stream
req: operations/004
This commit is contained in:
@@ -25,6 +25,13 @@ use std::sync::{Arc, Mutex};
|
|||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
const COLUMNS: [(&str, &str); 3] = [("backlog", "Backlog"), ("doing", "Doing"), ("done", "Done")];
|
const COLUMNS: [(&str, &str); 3] = [("backlog", "Backlog"), ("doing", "Doing"), ("done", "Done")];
|
||||||
|
const ACKNOWLEDGEMENT_STREAM_BUFFER_LIMIT: usize = 64;
|
||||||
|
const ACKNOWLEDGEMENT_HEARTBEAT_INTERVAL: Duration = Duration::from_secs(15);
|
||||||
|
const ACKNOWLEDGEMENT_RECONNECT_BACKOFF: [Duration; 3] = [
|
||||||
|
Duration::from_millis(100),
|
||||||
|
Duration::from_millis(250),
|
||||||
|
Duration::from_millis(500),
|
||||||
|
];
|
||||||
|
|
||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
struct AppState {
|
struct AppState {
|
||||||
@@ -32,6 +39,7 @@ struct AppState {
|
|||||||
sync: Mutex<SyncState>,
|
sync: Mutex<SyncState>,
|
||||||
sync_store: Option<SyncStore>,
|
sync_store: Option<SyncStore>,
|
||||||
sync_sessions: SyncSessionTokens,
|
sync_sessions: SyncSessionTokens,
|
||||||
|
acknowledgement_heartbeat_interval: Duration,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Default)]
|
#[derive(Default)]
|
||||||
@@ -437,6 +445,12 @@ async fn main() {
|
|||||||
sync: Mutex::new(sync),
|
sync: Mutex::new(sync),
|
||||||
sync_store,
|
sync_store,
|
||||||
sync_sessions: SyncSessionTokens::from_env(),
|
sync_sessions: SyncSessionTokens::from_env(),
|
||||||
|
acknowledgement_heartbeat_interval: std::env::var("HEMX_KANBAN_ACK_HEARTBEAT_MS")
|
||||||
|
.ok()
|
||||||
|
.and_then(|value| value.parse::<u64>().ok())
|
||||||
|
.filter(|milliseconds| *milliseconds > 0)
|
||||||
|
.map(Duration::from_millis)
|
||||||
|
.unwrap_or(ACKNOWLEDGEMENT_HEARTBEAT_INTERVAL),
|
||||||
});
|
});
|
||||||
|
|
||||||
let app = Router::new()
|
let app = Router::new()
|
||||||
@@ -789,14 +803,17 @@ async fn sync_acknowledgements(
|
|||||||
if let Some(key) = reconnect_key {
|
if let Some(key) = reconnect_key {
|
||||||
let attempts = sync.reconnects.entry(key).or_default();
|
let attempts = sync.reconnects.entry(key).or_default();
|
||||||
*attempts += 1;
|
*attempts += 1;
|
||||||
if *attempts == 1 {
|
if let Some(backoff) = usize::try_from(*attempts - 1)
|
||||||
|
.ok()
|
||||||
|
.and_then(|index| ACKNOWLEDGEMENT_RECONNECT_BACKOFF.get(index))
|
||||||
|
{
|
||||||
return Ok(Sse::new(
|
return Ok(Sse::new(
|
||||||
stream::iter([Ok(Event::default()
|
stream::iter([Ok(Event::default()
|
||||||
.comment("reconnect")
|
.comment(format!("reconnect-attempt-{attempts}"))
|
||||||
.retry(Duration::from_millis(25)))])
|
.retry(*backoff))])
|
||||||
.boxed(),
|
.boxed(),
|
||||||
)
|
)
|
||||||
.keep_alive(KeepAlive::default()));
|
.keep_alive(KeepAlive::new().interval(state.acknowledgement_heartbeat_interval)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let first_available = sync
|
let first_available = sync
|
||||||
@@ -817,7 +834,17 @@ async fn sync_acknowledgements(
|
|||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
let history_missing = after < latest
|
let history_missing = after < latest
|
||||||
&& first_available.is_none_or(|first_sequence| after.saturating_add(1) < first_sequence);
|
&& first_available.is_none_or(|first_sequence| after.saturating_add(1) < first_sequence);
|
||||||
let events = if history_missing {
|
let pending_count = sync
|
||||||
|
.acknowledgements
|
||||||
|
.values()
|
||||||
|
.filter(|acknowledgement| {
|
||||||
|
visible_acknowledgement(principal, acknowledgement)
|
||||||
|
&& acknowledgement.server_sequence > after
|
||||||
|
&& acknowledgement.server_sequence > sync.retained_after
|
||||||
|
})
|
||||||
|
.count();
|
||||||
|
let slow_consumer = pending_count > ACKNOWLEDGEMENT_STREAM_BUFFER_LIMIT;
|
||||||
|
let events = if history_missing || slow_consumer {
|
||||||
vec![Ok(Event::default()
|
vec![Ok(Event::default()
|
||||||
.id(latest.to_string())
|
.id(latest.to_string())
|
||||||
.event("snapshot-required")
|
.event("snapshot-required")
|
||||||
@@ -826,6 +853,9 @@ async fn sync_acknowledgements(
|
|||||||
"firstAvailable": first_available,
|
"firstAvailable": first_available,
|
||||||
"latest": latest,
|
"latest": latest,
|
||||||
"snapshotUrl": "/sync/snapshot",
|
"snapshotUrl": "/sync/snapshot",
|
||||||
|
"reason": if slow_consumer { "slow-consumer" } else { "missing-history" },
|
||||||
|
"pendingCount": pending_count,
|
||||||
|
"bufferLimit": ACKNOWLEDGEMENT_STREAM_BUFFER_LIMIT,
|
||||||
}))
|
}))
|
||||||
.expect("serializable missing history event"))]
|
.expect("serializable missing history event"))]
|
||||||
} else {
|
} else {
|
||||||
@@ -845,7 +875,24 @@ async fn sync_acknowledgements(
|
|||||||
})
|
})
|
||||||
.collect::<Vec<Result<Event, Infallible>>>()
|
.collect::<Vec<Result<Event, Infallible>>>()
|
||||||
};
|
};
|
||||||
Ok(Sse::new(stream::iter(events).boxed()).keep_alive(KeepAlive::default()))
|
drop(sync);
|
||||||
|
let heartbeat_interval = state.acknowledgement_heartbeat_interval;
|
||||||
|
let heartbeat = stream::unfold(heartbeat_interval, |interval| async move {
|
||||||
|
tokio::time::sleep(interval).await;
|
||||||
|
Some((
|
||||||
|
Ok(Event::default()
|
||||||
|
.event("heartbeat")
|
||||||
|
.data("{\"status\":\"ok\"}")),
|
||||||
|
interval,
|
||||||
|
))
|
||||||
|
});
|
||||||
|
Ok(
|
||||||
|
Sse::new(stream::iter(events).chain(heartbeat).boxed()).keep_alive(
|
||||||
|
KeepAlive::new()
|
||||||
|
.interval(heartbeat_interval)
|
||||||
|
.text("heartbeat"),
|
||||||
|
),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn registry(state: Arc<AppState>) -> impl DispatchRegistry {
|
fn registry(state: Arc<AppState>) -> impl DispatchRegistry {
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ const MIGRATION_KEY = "commandSchemaMigration";
|
|||||||
const MAX_ATTEMPTS = 3;
|
const MAX_ATTEMPTS = 3;
|
||||||
const BACKOFF_MS = [25, 50];
|
const BACKOFF_MS = [25, 50];
|
||||||
const REQUEST_TIMEOUT_MS = 1_000;
|
const REQUEST_TIMEOUT_MS = 1_000;
|
||||||
|
const ACKNOWLEDGEMENT_STREAM_BUFFER_LIMIT = 64;
|
||||||
const root = document.querySelector("[data-kanban-sync]");
|
const root = document.querySelector("[data-kanban-sync]");
|
||||||
const TAB_ID = sessionStorage.getItem("hemx-kanban-sync-tab-id") || crypto.randomUUID();
|
const TAB_ID = sessionStorage.getItem("hemx-kanban-sync-tab-id") || crypto.randomUUID();
|
||||||
const LEASE_MS = 5000;
|
const LEASE_MS = 5000;
|
||||||
@@ -481,6 +482,17 @@ async function synchronize(command) {
|
|||||||
source.addEventListener("open", () => {
|
source.addEventListener("open", () => {
|
||||||
opens += 1;
|
opens += 1;
|
||||||
root.setAttribute("data-sync-transport-opens", String(opens));
|
root.setAttribute("data-sync-transport-opens", String(opens));
|
||||||
|
root.setAttribute("data-sync-stream-state", "open");
|
||||||
|
});
|
||||||
|
source.addEventListener("heartbeat", () => {
|
||||||
|
const heartbeats = Number(root.getAttribute("data-sync-heartbeats") || "0") + 1;
|
||||||
|
root.setAttribute("data-sync-heartbeats", String(heartbeats));
|
||||||
|
root.setAttribute("data-sync-stream-state", "healthy");
|
||||||
|
});
|
||||||
|
source.addEventListener("error", () => {
|
||||||
|
const reconnects = Number(root.getAttribute("data-sync-reconnects") || "0") + 1;
|
||||||
|
root.setAttribute("data-sync-reconnects", String(reconnects));
|
||||||
|
root.setAttribute("data-sync-stream-state", "reconnecting");
|
||||||
});
|
});
|
||||||
source.addEventListener("acknowledgement", async (event) => {
|
source.addEventListener("acknowledgement", async (event) => {
|
||||||
const canonical = JSON.parse(event.data);
|
const canonical = JSON.parse(event.data);
|
||||||
@@ -631,6 +643,7 @@ async function runLeaseLoop(command) {
|
|||||||
async function start() {
|
async function start() {
|
||||||
if (!root) return;
|
if (!root) return;
|
||||||
root.setAttribute("data-sync-request-timeout-ms", String(REQUEST_TIMEOUT_MS));
|
root.setAttribute("data-sync-request-timeout-ms", String(REQUEST_TIMEOUT_MS));
|
||||||
|
root.setAttribute("data-sync-stream-buffer-limit", String(ACKNOWLEDGEMENT_STREAM_BUFFER_LIMIT));
|
||||||
const contextResponse = await fetchWithTimeout("/sync/context", { credentials: "same-origin", cache: "no-store" });
|
const contextResponse = await fetchWithTimeout("/sync/context", { credentials: "same-origin", cache: "no-store" });
|
||||||
if (!contextResponse.ok) throw new Error(`account context failed with ${contextResponse.status}`);
|
if (!contextResponse.ok) throw new Error(`account context failed with ${contextResponse.status}`);
|
||||||
const context = await contextResponse.json();
|
const context = await contextResponse.json();
|
||||||
@@ -701,7 +714,10 @@ window.addEventListener("pagehide", () => {
|
|||||||
stopped = true;
|
stopped = true;
|
||||||
clearTimeout(leaseTimer);
|
clearTimeout(leaseTimer);
|
||||||
clearTimeout(retryTimer);
|
clearTimeout(retryTimer);
|
||||||
acknowledgementSource?.close();
|
if (acknowledgementSource) {
|
||||||
|
acknowledgementSource.close();
|
||||||
|
root.setAttribute("data-sync-stream-state", "cancelled");
|
||||||
|
}
|
||||||
acknowledgementSource = undefined;
|
acknowledgementSource = undefined;
|
||||||
for (const controller of activeRequests) {
|
for (const controller of activeRequests) {
|
||||||
controller.abort(new DOMException("sync cancelled because page is hidden", "AbortError"));
|
controller.abort(new DOMException("sync cancelled because page is hidden", "AbortError"));
|
||||||
|
|||||||
@@ -2196,6 +2196,113 @@ async fn adversarial_wire_inputs_are_rejected_before_partial_application() -> We
|
|||||||
result.and(quit)
|
result.and(quit)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn acknowledgement_stream_bounds_reconnect_buffering_heartbeat_and_cancellation(
|
||||||
|
) -> WebDriverResult<()> {
|
||||||
|
// test req: operations/004
|
||||||
|
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_ACK_HEARTBEAT_MS", "25");
|
||||||
|
let _app = TestProcess::start(app_command, "hemx-kanban", &app_addr, STARTUP_TIMEOUT)
|
||||||
|
.expect("start 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}/sync-demo")).await?;
|
||||||
|
wait_until(
|
||||||
|
&driver,
|
||||||
|
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'idle'",
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
let proof = driver
|
||||||
|
.execute_async(
|
||||||
|
r#"
|
||||||
|
const done = arguments[arguments.length - 1];
|
||||||
|
(async () => {
|
||||||
|
for (let index = 1; index <= 65; index += 1) {
|
||||||
|
const query = new URLSearchParams({
|
||||||
|
command_id: `buffer:${index}`,
|
||||||
|
card_id: '1',
|
||||||
|
column: 'done',
|
||||||
|
});
|
||||||
|
const response = await fetch(`/sync/commands?${query}`, { method: 'POST' });
|
||||||
|
if (!response.ok) throw new Error(`command ${index} failed with ${response.status}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
const observe = (url, eventName, timeoutMs = 5000) => new Promise((resolve, reject) => {
|
||||||
|
const started = performance.now();
|
||||||
|
const source = new EventSource(url);
|
||||||
|
let opens = 0;
|
||||||
|
let errors = 0;
|
||||||
|
const timeout = setTimeout(() => {
|
||||||
|
source.close();
|
||||||
|
reject(new Error(`${eventName} timed out after ${timeoutMs} ms`));
|
||||||
|
}, timeoutMs);
|
||||||
|
source.addEventListener('open', () => { opens += 1; });
|
||||||
|
source.addEventListener('error', () => { errors += 1; });
|
||||||
|
source.addEventListener(eventName, (event) => {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
const data = JSON.parse(event.data);
|
||||||
|
source.close();
|
||||||
|
resolve({
|
||||||
|
data,
|
||||||
|
opens,
|
||||||
|
errors,
|
||||||
|
elapsedMs: performance.now() - started,
|
||||||
|
cancelled: source.readyState === EventSource.CLOSED,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
const slowConsumer = await observe(
|
||||||
|
'/sync/acknowledgements?after=0&reconnect=slow-consumer-proof',
|
||||||
|
'snapshot-required',
|
||||||
|
);
|
||||||
|
const heartbeat = await observe(
|
||||||
|
'/sync/acknowledgements?after=65&reconnect=heartbeat-proof',
|
||||||
|
'heartbeat',
|
||||||
|
);
|
||||||
|
done({ slowConsumer, heartbeat });
|
||||||
|
})().catch((error) => done({ error: String(error), stack: error?.stack }));
|
||||||
|
"#,
|
||||||
|
Vec::new(),
|
||||||
|
)
|
||||||
|
.await?
|
||||||
|
.json()
|
||||||
|
.clone();
|
||||||
|
assert!(proof["error"].is_null(), "stream bounds failed: {proof}");
|
||||||
|
assert_eq!(proof["slowConsumer"]["data"]["reason"], "slow-consumer");
|
||||||
|
assert_eq!(proof["slowConsumer"]["data"]["pendingCount"], 65);
|
||||||
|
assert_eq!(proof["slowConsumer"]["data"]["bufferLimit"], 64);
|
||||||
|
assert_eq!(proof["slowConsumer"]["opens"], 4);
|
||||||
|
assert!(proof["slowConsumer"]["errors"].as_u64().is_some_and(|errors| errors >= 3));
|
||||||
|
assert!(proof["slowConsumer"]["elapsedMs"]
|
||||||
|
.as_f64()
|
||||||
|
.is_some_and(|elapsed| (700.0..5_000.0).contains(&elapsed)));
|
||||||
|
assert_eq!(proof["slowConsumer"]["cancelled"], true);
|
||||||
|
assert_eq!(proof["heartbeat"]["data"]["status"], "ok");
|
||||||
|
assert_eq!(proof["heartbeat"]["opens"], 4);
|
||||||
|
assert!(proof["heartbeat"]["errors"].as_u64().is_some_and(|errors| errors >= 3));
|
||||||
|
assert_eq!(proof["heartbeat"]["cancelled"], true);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
.await;
|
||||||
|
let quit = driver.quit().await;
|
||||||
|
result.and(quit)
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn sync_requests_timeout_and_cancel_on_pagehide() -> WebDriverResult<()> {
|
async fn sync_requests_timeout_and_cancel_on_pagehide() -> WebDriverResult<()> {
|
||||||
// test req: operations/003
|
// test req: operations/003
|
||||||
|
|||||||
Reference in New Issue
Block a user