use std::fs;
use std::path::PathBuf;
use std::process::{Command, Output};
use std::sync::atomic::{AtomicU64, Ordering};
static NEXT_FIXTURE: AtomicU64 = AtomicU64::new(0);
struct Fixture {
root: PathBuf,
}
impl Fixture {
fn new(label: &str) -> Self {
let ordinal = NEXT_FIXTURE.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"noxid-wo24-qa8-{label}-{}-{ordinal}",
std::process::id()
));
fs::create_dir_all(root.join("server/queues")).expect("create queue fixture");
fs::write(
root.join("Noxid.toml"),
"[app]\ntitle = \"WO-24 QA round 8\"\n",
)
.expect("write project manifest");
fs::write(root.join("package.json"), "{\"type\":\"module\"}\n").expect("write ESM marker");
fs::write(
root.join("server/host.js"),
"export const queues = Object.freeze({});\n",
)
.expect("write inert host");
fs::write(
root.join("server/queues/Probe.nox"),
"queue Probe { payload {} retry: 0 backoff: 1s }\n",
)
.expect("write probe queue");
Self { root }
}
fn write(&self, relative: &str, contents: &str) {
let path = self.root.join(relative);
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).expect("create fixture parent");
}
fs::write(path, contents).expect("write fixture file");
}
fn build(&self) -> Output {
Command::new(env!("CARGO_BIN_EXE_noxid"))
.args(["build", ".", "--out-dir", "dist"])
.current_dir(&self.root)
.output()
.expect("build queue fixture")
}
fn run_node_with_database(&self, name: &str, source: &str) -> Output {
self.write(&format!("dist/{name}.mjs"), source);
Command::new("node")
.arg(format!("{name}.mjs"))
.current_dir(self.root.join("dist"))
.env("DATABASE_URL", "postgres://qa.invalid/noxid")
.output()
.expect("execute generated queue runtime with fake database")
}
fn install_controlled_postgres(&self) {
self.write(
"dist/node_modules/postgres/package.json",
"{\"type\":\"module\",\"exports\":\"./index.js\"}\n",
);
self.write(
"dist/node_modules/postgres/index.js",
r#"export default function postgres() {
const sql = async (strings) => {
const query = strings.join("?");
if (query.includes("SELECT id, queue, payload")) return globalThis.__qaSelect();
return [];
};
sql.unsafe = async () => [];
sql.begin = async (callback) => callback(sql);
sql.end = async () => {};
sql.json = (value) => value;
return sql;
}
"#,
);
}
}
impl Drop for Fixture {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.root);
}
}
fn output_text(output: &Output) -> String {
format!(
"stdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
)
}
fn assert_success(output: &Output, context: &str) {
assert!(
output.status.success(),
"{context}:\n{}",
output_text(output)
);
}
#[test]
fn qa_round8_synchronous_duplicate_delivery_then_throw_fails_without_claim_or_clear() {
let fixture = Fixture::new("sync-delivery-then-throw");
assert_success(&fixture.build(), "build synchronous throw fixture");
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"sync-delivery-then-throw",
r#"let selectCalls = 0;
globalThis.__qaSelect = async () => { selectCalls += 1; return []; };
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
const unhandled = [];
const reported = [];
process.on("unhandledRejection", (reason) => { unhandled.push(reason); });
console.error = (error) => { reported.push(error); };
let setCalls = 0;
let clearCalls = 0;
let captured;
const worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 79,
setTimeout(callback) {
setCalls += 1;
captured = callback;
callback();
callback();
throw new Error("arm threw after delivery");
},
clearTimeout() { clearCalls += 1; },
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && setCalls === 0; turn += 1) await new Promise((resolve) => setImmediate(resolve));
await Promise.all([worker.stop(), worker.stop(), worker.stop()]);
captured();
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (selectCalls !== 1) throw new Error(`throwing arm launched ${selectCalls - 1} timer claim(s)`);
if (setCalls !== 1 || clearCalls !== 0) throw new Error(`Failed arm ownership changed: sets=${setCalls}, clears=${clearCalls}`);
if (reported.length !== 1 || reported[0]?.message !== "arm threw after delivery") throw new Error(`arm failure was not contained exactly once: ${reported.length}`);
if (unhandled.length !== 0) throw new Error(`arm failure leaked ${unhandled.length} rejection(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"contain a throwing arm after synchronous duplicate delivery",
);
}
#[test]
fn qa_round8_stop_after_sync_delivery_but_before_handle_publication_joins_the_attempt() {
let fixture = Fixture::new("stop-after-delivery-before-handle");
assert_success(&fixture.build(), "build pre-publication join fixture");
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"stop-after-delivery-before-handle",
r#"let selectCalls = 0;
let resolveSecond;
let secondStarted = false;
globalThis.__qaSelect = () => {
selectCalls += 1;
if (selectCalls === 1) return Promise.resolve([]);
secondStarted = true;
return new Promise((resolve) => { resolveSecond = resolve; });
};
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
const handle = Object.freeze({ id: "published-after-stop" });
let stopFromHook;
let stopSettled = false;
let setCalls = 0;
const cleared = [];
let worker;
worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 83,
setTimeout(callback) {
setCalls += 1;
callback();
callback();
stopFromHook = worker.stop().then(() => { stopSettled = true; });
return handle;
},
clearTimeout(observed) { cleared.push(observed); },
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && !secondStarted; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (!secondStarted || stopFromHook === undefined) throw new Error(`synchronous attempt was not launched: selects=${selectCalls}`);
await new Promise((resolve) => setImmediate(resolve));
if (stopSettled) throw new Error("stop settled before the pre-publication attempt");
resolveSecond([]);
await Promise.all([stopFromHook, worker.stop(), worker.stop()]);
if (selectCalls !== 2 || setCalls !== 1) throw new Error(`duplicate delivery changed work: selects=${selectCalls}, sets=${setCalls}`);
if (cleared.length !== 0) throw new Error(`already-delivered handle was cleared ${cleared.length} time(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"join synchronous work stopped before its handle is published",
);
}
#[test]
fn qa_round8_stop_before_handle_publication_clears_once_and_clear_delivery_is_inert() {
let fixture = Fixture::new("stop-before-handle-publication");
assert_success(
&fixture.build(),
"build pre-publication cancellation fixture",
);
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"stop-before-handle-publication",
r#"let selectCalls = 0;
globalThis.__qaSelect = async () => { selectCalls += 1; return []; };
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
const handle = Object.freeze({ id: "late-handle" });
let callback;
let stopFromSet;
let stopFromClear;
let setCalls = 0;
const cleared = [];
let worker;
worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 89,
setTimeout(next) {
setCalls += 1;
callback = next;
stopFromSet = worker.stop();
return handle;
},
clearTimeout(observed) {
cleared.push(observed);
callback();
callback();
stopFromClear = worker.stop();
},
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && stopFromClear === undefined; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (stopFromSet === undefined || stopFromClear === undefined) throw new Error("reentrant stop/clear sequence did not finish");
await Promise.all([stopFromSet, stopFromClear, worker.stop(), worker.stop()]);
callback();
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (selectCalls !== 1 || setCalls !== 1) throw new Error(`clear delivery restarted work: selects=${selectCalls}, sets=${setCalls}`);
if (cleared.length !== 1 || cleared[0] !== handle) throw new Error(`late handle was not cleared once by identity: ${cleared.length}`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"cancel before publication and ignore callback delivery from clearTimeout",
);
}
#[test]
fn qa_round8_stop_after_handle_publication_makes_synchronous_clear_delivery_inert() {
let fixture = Fixture::new("stop-after-handle-publication");
assert_success(
&fixture.build(),
"build post-publication cancellation fixture",
);
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"stop-after-handle-publication",
r#"let selectCalls = 0;
globalThis.__qaSelect = async () => { selectCalls += 1; return []; };
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
let callback;
let setCalls = 0;
let clearCalls = 0;
let reentrantStop;
const worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 97,
setTimeout(next) { setCalls += 1; callback = next; return null; },
clearTimeout(handle) {
if (handle !== null) throw new Error("falsy timer handle identity changed");
clearCalls += 1;
callback();
callback();
reentrantStop = worker.stop();
},
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && callback === undefined; turn += 1) await new Promise((resolve) => setImmediate(resolve));
await Promise.all([worker.stop(), worker.stop(), worker.stop()]);
await reentrantStop;
callback();
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (selectCalls !== 1 || setCalls !== 1) throw new Error(`post-publication clear restarted work: selects=${selectCalls}, sets=${setCalls}`);
if (clearCalls !== 1) throw new Error(`post-publication timer cleared ${clearCalls} time(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"ignore synchronous callback delivery while clearing a published handle",
);
}
#[test]
fn qa_round8_stale_and_duplicate_delivery_during_successor_arm_stays_one_shot() {
let fixture = Fixture::new("nested-successor-delivery");
assert_success(&fixture.build(), "build nested successor-delivery fixture");
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"nested-successor-delivery",
r#"const resolvers = new Map();
let selectCalls = 0;
globalThis.__qaSelect = () => {
selectCalls += 1;
const index = selectCalls;
return new Promise((resolve) => { resolvers.set(index, resolve); });
};
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
const callbacks = [];
const handles = [];
const cleared = [];
const worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 101,
setTimeout(callback) {
callbacks.push(callback);
const handle = Object.freeze({ id: callbacks.length });
handles.push(handle);
if (callbacks.length === 2) {
callbacks[0]();
callback();
callback();
callbacks[0]();
}
return handle;
},
clearTimeout(handle) { cleared.push(handle); },
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && !resolvers.has(1); turn += 1) await new Promise((resolve) => setImmediate(resolve));
resolvers.get(1)([]);
for (let turn = 0; turn < 20 && callbacks.length < 1; turn += 1) await new Promise((resolve) => setImmediate(resolve));
callbacks[0]();
callbacks[0]();
for (let turn = 0; turn < 20 && !resolvers.has(2); turn += 1) await new Promise((resolve) => setImmediate(resolve));
resolvers.get(2)([]);
for (let turn = 0; turn < 20 && !resolvers.has(3); turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (!resolvers.has(3) || callbacks.length !== 2) throw new Error(`nested successor did not launch exactly once: selects=${selectCalls}, callbacks=${callbacks.length}`);
let stopSettled = false;
const stopping = worker.stop().then(() => { stopSettled = true; });
await new Promise((resolve) => setImmediate(resolve));
if (stopSettled) throw new Error("stop did not join nested successor attempt");
resolvers.get(3)([]);
await stopping;
for (const callback of callbacks) { callback(); callback(); }
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (selectCalls !== 3) throw new Error(`nested stale/duplicate delivery started ${selectCalls} claims`);
if (cleared.length !== 0) throw new Error(`delivered handles were cleared ${cleared.length} time(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"keep nested stale and duplicate delivery inert during successor arming",
);
}
#[test]
fn qa_round8_thenable_handle_getter_can_deliver_and_stop_before_publication() {
let fixture = Fixture::new("thenable-getter-delivery");
assert_success(&fixture.build(), "build thenable-handle fixture");
fixture.install_controlled_postgres();
let node = fixture.run_node_with_database(
"thenable-getter-delivery",
r#"let selectCalls = 0;
let resolveSecond;
globalThis.__qaSelect = () => {
selectCalls += 1;
if (selectCalls === 1) return Promise.resolve([]);
return new Promise((resolve) => { resolveSecond = resolve; });
};
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
let thenReads = 0;
let stopFromGetter;
let stopSettled = false;
let clearCalls = 0;
let worker;
worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 103,
setTimeout(callback) {
const handle = Object.create(null);
Object.defineProperty(handle, "then", {
get() {
thenReads += 1;
callback();
callback();
stopFromGetter = worker.stop().then(() => { stopSettled = true; });
return (resolve) => resolve("observed");
},
});
return handle;
},
clearTimeout() { clearCalls += 1; },
onError(error) { throw error; },
});
for (let turn = 0; turn < 20 && resolveSecond === undefined; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (resolveSecond === undefined || stopFromGetter === undefined) throw new Error(`thenable getter did not launch its claim: selects=${selectCalls}`);
await new Promise((resolve) => setImmediate(resolve));
if (stopSettled) throw new Error("thenable getter stop abandoned the active claim");
resolveSecond([]);
await Promise.all([stopFromGetter, worker.stop(), worker.stop()]);
if (selectCalls !== 2 || thenReads !== 1) throw new Error(`thenable delivery repeated: selects=${selectCalls}, reads=${thenReads}`);
if (clearCalls !== 0) throw new Error(`already-delivered thenable handle was cleared ${clearCalls} time(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"handle reentrant thenable delivery before timer-handle publication",
);
}
#[test]
fn qa_round8_stop_in_running_joins_failure_settlement_but_not_notification() {
let fixture = Fixture::new("running-stop-notification");
fixture.write(
"server/host.js",
r#"export const queues = Object.freeze({
"queue:Probe": async (payload) => globalThis.__qaHandle(payload),
});
"#,
);
assert_success(&fixture.build(), "build Running-state stop fixture");
fixture.write(
"dist/node_modules/postgres/package.json",
"{\"type\":\"module\",\"exports\":\"./index.js\"}\n",
);
fixture.write(
"dist/node_modules/postgres/index.js",
r#"export default function postgres() {
const sql = async (strings) => {
const query = strings.join("?");
globalThis.__qaQueries.push(query);
if (query.includes("SELECT id, queue, payload")) {
if (globalThis.__qaClaimed) return [];
globalThis.__qaClaimed = true;
return [{ id: "qa-running", queue: "Probe", payload: {}, principal: null, attempts: 0, run_at: new Date("2028-02-29T12:00:00Z") }];
}
if (query.includes("state = 'dead-letter'")) throw new Error("failure update unavailable");
return [];
};
sql.unsafe = async () => [];
sql.begin = async (callback) => callback(sql);
sql.end = async () => {};
sql.json = (value) => value;
return sql;
}
"#,
);
let node = fixture.run_node_with_database(
"running-stop-notification",
r#"globalThis.__qaQueries = [];
globalThis.__qaClaimed = false;
let rejectHandler;
let handlerStarted = false;
globalThis.__qaHandle = () => {
handlerStarted = true;
return new Promise((_, reject) => { rejectHandler = reject; });
};
const { startQueueWorker, closeQueueDatabase } = await import("./server/handler.js");
let resolveNotification;
let notificationStarted = false;
const notification = new Promise((resolve) => { resolveNotification = resolve; });
let setCalls = 0;
const worker = startQueueWorker({
queue: "Probe",
worker: "qa",
now: "2028-02-29T12:00:00Z",
pollIntervalMs: 107,
setTimeout() { setCalls += 1; return 1; },
clearTimeout() {},
onError(error) {
if (error?.message !== "failure update unavailable") throw new Error(`unexpected notification ${error?.message}`);
notificationStarted = true;
return notification;
},
});
for (let turn = 0; turn < 20 && !handlerStarted; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (!handlerStarted) throw new Error("worker never entered Running(job)");
let firstStopped = false;
let secondStopped = false;
const firstStop = worker.stop().then(() => { firstStopped = true; });
const secondStop = worker.stop().then(() => { secondStopped = true; });
await new Promise((resolve) => setImmediate(resolve));
if (firstStopped || secondStopped) throw new Error("stop returned while the handler was still running");
rejectHandler(new Error("handler failed"));
for (let turn = 0; turn < 20 && !notificationStarted; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (!notificationStarted) throw new Error("failed Running attempt did not notify");
await Promise.race([
Promise.all([firstStop, secondStop]),
new Promise((_, reject) => setTimeout(() => reject(new Error("stop joined detached notification")), 100)),
]);
if (!firstStopped || !secondStopped) throw new Error("shared attempt join did not settle");
if (setCalls !== 0) throw new Error(`Stopping worker armed ${setCalls} successor timer(s)`);
resolveNotification();
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
await worker.stop();
if (setCalls !== 0) throw new Error(`late notification armed ${setCalls} timer(s)`);
await closeQueueDatabase();
"#,
);
assert_success(
&node,
"join Running work while keeping failure notification detached",
);
}