feat(kanban): stop mixed queue on rejection
req: sync/009 req: sync/010
This commit is contained in:
@@ -19,10 +19,12 @@ let maxObservedInFlight = 0;
|
||||
let stopped = false;
|
||||
|
||||
class UploadError extends Error {
|
||||
constructor(status, retryable) {
|
||||
constructor(status, retryable, reason) {
|
||||
super(`sync upload failed with ${status}`);
|
||||
this.name = "UploadError";
|
||||
this.status = status;
|
||||
this.retryable = retryable;
|
||||
this.reason = reason;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -132,10 +134,17 @@ function setOnline(online) {
|
||||
root.setAttribute("data-sync-connection", online ? "online" : "offline");
|
||||
}
|
||||
|
||||
function setManualRetryAvailable(available) {
|
||||
const retry = root.querySelector("[data-sync-retry]");
|
||||
retry.disabled = !available;
|
||||
if (available) root.setAttribute("data-sync-manual-retry", "available");
|
||||
else root.removeAttribute("data-sync-manual-retry");
|
||||
}
|
||||
|
||||
function scheduleManualRetry(command, error) {
|
||||
clearTimeout(retryTimer);
|
||||
root.setAttribute("data-sync-error", error instanceof Error ? error.message : String(error));
|
||||
root.setAttribute("data-sync-manual-retry", "available");
|
||||
setManualRetryAvailable(true);
|
||||
setPhase("offline", "Sync is offline after bounded retries; the durable command remains queued. Retry now when ready.");
|
||||
root.dispatchEvent(new CustomEvent("kanban:sync-exhausted", { detail: { commandId: command.id, attempts: MAX_ATTEMPTS } }));
|
||||
}
|
||||
@@ -157,7 +166,11 @@ async function upload(command) {
|
||||
await new Promise((resolve) => setTimeout(resolve, delay));
|
||||
continue;
|
||||
}
|
||||
if (!response.ok) throw new UploadError(response.status, response.status >= 500);
|
||||
if (!response.ok) {
|
||||
const problem = await response.json().catch(() => ({}));
|
||||
const reason = typeof problem.error === "string" ? problem.error : "unclassified rejection";
|
||||
throw new UploadError(response.status, response.status >= 500, reason);
|
||||
}
|
||||
return response.json();
|
||||
} catch (error) {
|
||||
if (error instanceof UploadError && !error.retryable) throw error;
|
||||
@@ -190,7 +203,7 @@ async function continuePendingWork() {
|
||||
root.setAttribute("data-sync-pending-count", String(commands.length));
|
||||
if (commands.length === 0) return;
|
||||
if (uploadsThisRun >= uploadLimit) {
|
||||
root.setAttribute("data-sync-manual-retry", "available");
|
||||
setManualRetryAvailable(true);
|
||||
setPhase("backpressured", `Upload limit ${uploadLimit} reached; ${commands.length} durable command${commands.length === 1 ? " remains" : "s remain"} queued. Retry now to continue.`);
|
||||
return;
|
||||
}
|
||||
@@ -285,8 +298,23 @@ async function synchronize(command) {
|
||||
});
|
||||
} catch (error) {
|
||||
synchronizing = false;
|
||||
if (error instanceof UploadError && !error.retryable) {
|
||||
setOnline(true);
|
||||
clearTimeout(leaseTimer);
|
||||
root.setAttribute("data-sync-error", error.message);
|
||||
root.setAttribute("data-sync-error-kind", "permanent-rejection");
|
||||
root.setAttribute("data-sync-error-status", String(error.status));
|
||||
root.setAttribute("data-sync-error-reason", error.reason);
|
||||
root.setAttribute("data-sync-rejected-command-id", command.id);
|
||||
const remaining = await pendingCommands(database);
|
||||
root.setAttribute("data-sync-pending-count", String(remaining.length));
|
||||
setManualRetryAvailable(false);
|
||||
setPhase("rejected", `Command ${command.id} was permanently rejected (${error.status}: ${error.reason}); ${remaining.length} durable command${remaining.length === 1 ? " remains" : "s remain"} queued for review.`);
|
||||
await releaseUploaderLease(database);
|
||||
root.setAttribute("data-sync-leader", "false");
|
||||
return;
|
||||
}
|
||||
setOnline(false);
|
||||
if (error instanceof UploadError && !error.retryable) throw error;
|
||||
scheduleManualRetry(command, error);
|
||||
}
|
||||
}
|
||||
@@ -320,6 +348,7 @@ async function start() {
|
||||
root.setAttribute("data-sync-in-flight", "0");
|
||||
root.setAttribute("data-sync-max-observed-in-flight", "0");
|
||||
root.setAttribute("data-sync-pending-count", String(commands.length));
|
||||
setManualRetryAvailable(false);
|
||||
if (commands.length === 0) {
|
||||
setPhase("idle", "No pending commands.");
|
||||
return;
|
||||
@@ -329,12 +358,13 @@ async function start() {
|
||||
if (!event.target.closest("[data-sync-retry]")) return;
|
||||
uploadsThisRun = 0;
|
||||
root.setAttribute("data-sync-uploaded-this-run", "0");
|
||||
root.removeAttribute("data-sync-manual-retry");
|
||||
setManualRetryAvailable(false);
|
||||
const [next] = await pendingCommands(database);
|
||||
if (next) synchronize(validatePending(next)).catch(failPermanently);
|
||||
});
|
||||
window.addEventListener("online", async () => {
|
||||
if (root.getAttribute("data-sync-phase") !== "offline") return;
|
||||
setManualRetryAvailable(false);
|
||||
const [next] = await pendingCommands(database);
|
||||
if (next) synchronize(validatePending(next)).catch(failPermanently);
|
||||
});
|
||||
|
||||
@@ -335,7 +335,7 @@ async fn pending_local_command_uploads_with_bounded_retry_and_is_removed_on_ack(
|
||||
driver.goto(&format!("http://{app_addr}/sync-demo")).await?;
|
||||
wait_until(
|
||||
&driver,
|
||||
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'failed'",
|
||||
"return document.querySelector('[data-kanban-sync]')?.getAttribute('data-sync-phase') === 'rejected'",
|
||||
)
|
||||
.await?;
|
||||
let rejected = driver
|
||||
@@ -349,7 +349,10 @@ async fn pending_local_command_uploads_with_bounded_retry_and_is_removed_on_ack(
|
||||
assert_eq!(rejected["pending"], "1");
|
||||
assert_eq!(rejected["attempts"], "1");
|
||||
assert_eq!(rejected["error"], "sync upload failed with 409");
|
||||
assert_eq!(rejected["status"], "Sync failed; the durable command remains queued.");
|
||||
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
|
||||
@@ -369,6 +372,134 @@ async fn pending_local_command_uploads_with_bounded_retry_and_is_removed_on_ack(
|
||||
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', 1);
|
||||
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: 1, actor: 'mixed', session: 'mixed-session',
|
||||
causal, kind: 'reorder_card', cardId, 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', 1);
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user