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, queue_worker: bool) -> Self {
let ordinal = NEXT_FIXTURE.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"noxid-wo24-queue-{label}-{}-{ordinal}",
std::process::id()
));
fs::create_dir_all(root.join("server/queues")).expect("create queue fixture");
fs::write(
root.join("Noxid.toml"),
format!(
"[app]\ntitle = \"Queue runtime\"\n\n[server]\nqueue_worker = {queue_worker}\n\n[deploy]\nadapter = \"node\"\n"
),
)
.expect("write queue config");
fs::write(root.join("package.json"), "{\"type\":\"module\"}\n")
.expect("write module marker");
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 source");
}
fn noxid(&self, arguments: &[&str]) -> Output {
Command::new(env!("CARGO_BIN_EXE_noxid"))
.args(arguments)
.current_dir(&self.root)
.output()
.expect("run noxid")
}
fn build(&self) -> Output {
self.noxid(&["build", ".", "--out-dir", "dist"])
}
}
impl Drop for Fixture {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.root);
}
}
fn queue_fixture(label: &str, queue_worker: bool) -> Fixture {
let fixture = Fixture::new(label, queue_worker);
fixture.write(
"server/queues/SendReceipt.nox",
r#"queue SendReceipt {
payload { email: String orderId: Int note: Optional<String> }
retry: 2
backoff: 30s
}
"#,
);
fixture.write(
"server/host.js",
r#"export const queues = Object.freeze({
"queue:SendReceipt": async (payload, context) => ({ payload, context }),
});
"#,
);
fixture
}
fn assert_success(output: &Output, context: &str) {
assert!(
output.status.success(),
"{context}:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[test]
fn queue_only_project_emits_exact_manifests_worker_runtime_and_node_embedding() {
let fixture = queue_fixture("artifacts", true);
assert_success(&fixture.build(), "build queue-only project");
let handler = fs::read_to_string(fixture.root.join("dist/server/handler.js"))
.expect("read queue handler");
for contract in [
"CREATE TABLE IF NOT EXISTS _noxid_jobs",
"FOR UPDATE SKIP LOCKED",
"queue:SendReceipt",
"export async function enqueue",
"export function startQueueWorker",
"export async function queueStatus",
"QUEUE_PAYLOAD_DRIFT",
] {
assert!(handler.contains(contract), "handler omitted {contract}");
}
let manifest = fs::read_to_string(fixture.root.join("dist/server/queues.manifest.json"))
.expect("read queue manifest");
assert!(manifest.contains("\"retry\":2"));
assert!(manifest.contains("\"backoffMs\":30000"));
assert!(manifest.contains("\"hostKey\":\"queue:SendReceipt\""));
let app =
fs::read_to_string(fixture.root.join("dist/app.manifest.json")).expect("read app manifest");
assert!(app.contains("\"queues\": 1") && app.contains("\"queueWorker\": true"));
let adapted = fixture.noxid(&["adapt", ".", "--out-dir", "deploy", "--adapter", "node"]);
assert_success(&adapted, "adapt queue project for Node");
let node = fs::read_to_string(fixture.root.join("deploy/server.mjs")).expect("read node app");
assert!(
node.contains("const startQueueWorker = noxidServerHandler.startQueueWorker;")
&& node.contains("startQueueWorker();")
);
let package =
fs::read_to_string(fixture.root.join("deploy/package.json")).expect("read node package");
assert!(package.contains("\"postgres\": \"3.4.9\""));
}
#[test]
fn composed_ssr_entry_reexports_enqueue() {
let fixture = queue_fixture("composed-ssr-enqueue", false);
fixture.write(
"src/routes/+page.nox",
r#"component QueuePage {
route { title: "Queue" render: ssr }
render { mode: server hydrate: never }
view { <main><h1>Queue</h1></main> }
}
"#,
);
assert_success(&fixture.build(), "build composed SSR queue project");
let actions = fs::read_to_string(fixture.root.join("dist/server/actions.js"))
.expect("read composed actions module");
let handler = fs::read_to_string(fixture.root.join("dist/server/handler.js"))
.expect("read composed server entry");
assert!(actions.contains("export async function enqueue"));
assert!(handler.contains("export { enqueue,"));
fixture.write(
"dist/check-composed-enqueue.mjs",
r#"import { enqueue } from "./server/handler.js";
let code = null;
try { await enqueue("SendReceipt", { email: "missing-order-id" }); }
catch (error) { code = error?.code; }
if (code !== "QUEUE_PAYLOAD_TYPE") throw new Error(`composed enqueue surface returned ${code}`);
"#,
);
let output = Command::new("node")
.arg("check-composed-enqueue.mjs")
.current_dir(fixture.root.join("dist"))
.output()
.expect("run composed enqueue check");
assert_success(&output, "reach enqueue through composed SSR entry");
}
#[test]
fn enqueue_payload_and_options_validate_before_database_or_host_execution() {
let fixture = queue_fixture("validation", false);
assert_success(&fixture.build(), "build validation fixture");
let script = fixture.root.join("dist/check.mjs");
fs::write(
&script,
r#"import { enqueue } from "./server/handler.js";
const cases = [
["bad scalar", () => enqueue("SendReceipt", 1), "QUEUE_PAYLOAD_TYPE"],
["missing", () => enqueue("SendReceipt", { email: "a" }), "QUEUE_PAYLOAD_TYPE"],
["wrong type", () => enqueue("SendReceipt", { email: "a", orderId: "1" }), "QUEUE_PAYLOAD_TYPE"],
["extra", () => enqueue("SendReceipt", { email: "a", orderId: 1, extra: true }), "QUEUE_PAYLOAD_TYPE"],
["unknown", () => enqueue("Missing", {}), "QUEUE_NOT_FOUND"],
["options", () => enqueue("SendReceipt", { email: "a", orderId: 1 }, { retry: true }), "QUEUE_OPTIONS_INVALID"],
["runAt", () => enqueue("SendReceipt", { email: "a", orderId: 1 }, { runAt: "never" }), "QUEUE_RUN_AT_INVALID"],
];
for (const [label, invoke, code] of cases) {
let observed = null;
try { await invoke(); } catch (error) { observed = error.code; }
if (observed !== code) throw new Error(`${label}: expected ${code}, saw ${observed}`);
}
"#,
)
.expect("write queue validation script");
let output = Command::new("node")
.arg("check.mjs")
.current_dir(fixture.root.join("dist"))
.output()
.expect("run queue validation script");
assert_success(&output, "execute queue boundary validation");
let missing = fixture.noxid(&["queue", "status", "--queue", "Missing"]);
assert!(!missing.status.success());
assert!(String::from_utf8_lossy(&missing.stderr).contains("error[QUEUE_NOT_FOUND]"));
}
#[test]
fn queue_layout_and_worker_configuration_fail_closed() {
let no_queue = Fixture::new("worker-without-queue", true);
let rejected = no_queue.build();
assert!(!rejected.status.success());
assert!(String::from_utf8_lossy(&rejected.stderr).contains("QUEUE_WORKER_REQUIRES_QUEUE"));
let fixture = queue_fixture("layout", false);
fixture.write(
"server/queues/Nested/Bad.nox",
"queue Bad { payload { value: Int } retry: 1 backoff: 1s }\n",
);
let rejected = fixture.build();
assert!(!rejected.status.success());
assert!(String::from_utf8_lossy(&rejected.stderr).contains("QUEUE_LAYOUT_INVALID"));
}
#[test]
fn queue_worker_state_machine_has_a_total_table_and_hostile_timer_delivery_is_inert() {
let fixture = queue_fixture("worker-state-machine", false);
assert_success(&fixture.build(), "build worker state-machine fixture");
let handler_path = fixture.root.join("dist/server/handler.js");
let mut handler = fs::read_to_string(&handler_path).expect("read generated queue handler");
handler.push_str(
"\nexport { QUEUE_WORKER_STATE_NAMES, QUEUE_WORKER_EVENT_NAMES, QUEUE_WORKER_TRANSITIONS, queueWorkerState, queueWorkerEvent, queueWorkerTransition };\n",
);
fs::write(&handler_path, handler).expect("expose state machine to its generated-runtime test");
fixture.write(
"dist/worker-state-machine.mjs",
r#"import {
QUEUE_WORKER_STATE_NAMES as states,
QUEUE_WORKER_EVENT_NAMES as events,
QUEUE_WORKER_TRANSITIONS as table,
queueWorkerState as makeState,
queueWorkerEvent as makeEvent,
queueWorkerTransition as transition,
startQueueWorker,
} from "./server/handler.js";
const expected = Object.freeze({
Idle: Object.freeze({ Start: "Claiming", Arm: "Scheduled", Deliver: "Idle", Claimed: "Idle", Settle: "Idle", Stop: "Stopped", ArmFailed: "Idle", Notify: "Idle" }),
Scheduled: Object.freeze({ Start: "Scheduled", Arm: "Scheduled", Deliver: "Claiming", Claimed: "Scheduled", Settle: "Scheduled", Stop: "Stopped", ArmFailed: "Failed", Notify: "Scheduled" }),
Claiming: Object.freeze({ Start: "Claiming", Arm: "Claiming", Deliver: "Claiming", Claimed: "Running", Settle: "Idle", Stop: "Stopping", ArmFailed: "Failed", Notify: "Claiming" }),
Running: Object.freeze({ Start: "Running", Arm: "Running", Deliver: "Running", Claimed: "Running", Settle: "Idle", Stop: "Stopping", ArmFailed: "Running", Notify: "Running" }),
Stopping: Object.freeze({ Start: "Stopping", Arm: "Stopping", Deliver: "Stopping", Claimed: "Stopping", Settle: "Stopped", Stop: "Stopping", ArmFailed: "Stopping", Notify: "Stopping" }),
Stopped: Object.freeze({ Start: "Stopped", Arm: "Stopped", Deliver: "Stopped", Claimed: "Stopped", Settle: "Stopped", Stop: "Stopped", ArmFailed: "Stopped", Notify: "Stopped" }),
Failed: Object.freeze({ Start: "Failed", Arm: "Failed", Deliver: "Failed", Claimed: "Failed", Settle: "Failed", Stop: "Failed", ArmFailed: "Failed", Notify: "Failed" }),
});
if (JSON.stringify(table) !== JSON.stringify(expected)) throw new Error("generated transition table differs from the worker contract");
if (states.length !== 7 || events.length !== 8) throw new Error(`incomplete state/event axes: ${states.length} x ${events.length}`);
const token = Object.freeze({ id: 1 });
const staleToken = Object.freeze({ id: 0 });
const job = Object.freeze({ id: "job" });
const failure = new Error("arm failed");
function stateFor(name) {
if (name === "Scheduled") return makeState(name, token);
if (name === "Running") return makeState(name, job);
if (name === "Failed") return makeState(name, failure);
return makeState(name);
}
function eventFor(name) {
if (name === "Arm" || name === "Deliver") return makeEvent(name, { token });
if (name === "Claimed") return makeEvent(name, { job });
if (name === "ArmFailed") return makeEvent(name, { token, error: failure });
return makeEvent(name);
}
for (const stateName of states) {
for (const eventName of events) {
const before = stateFor(stateName);
const after = transition(before, eventFor(eventName));
const target = expected[stateName][eventName];
if (after.name !== target) throw new Error(`${stateName} x ${eventName}: expected ${target}, saw ${after.name}`);
if (target === stateName && after !== before) throw new Error(`${stateName} x ${eventName}: no-op replaced the state`);
}
}
const armed = transition(makeState("Idle"), makeEvent("Arm", { token }));
if (armed.detail !== token) throw new Error("Scheduled did not retain its exact token");
const running = transition(makeState("Claiming"), makeEvent("Claimed", { job }));
if (running.detail !== job) throw new Error("Running did not retain its exact job");
for (const eventName of ["Deliver", "ArmFailed"]) {
const scheduled = makeState("Scheduled", token);
const after = transition(scheduled, makeEvent(eventName, { token: staleToken, error: failure }));
if (after !== scheduled) throw new Error(`stale ${eventName} changed Scheduled`);
}
const failed = transition(makeState("Scheduled", token), makeEvent("ArmFailed", { token, error: failure }));
if (failed.name !== "Failed" || failed.detail !== failure) throw new Error("arm failure was not contained in Failed(error)");
const unhandled = [];
process.on("unhandledRejection", (reason) => { unhandled.push(reason); });
let syncSets = 0;
let syncErrors = 0;
let firstCallback;
let lateCallback;
const cleared = [];
const synchronous = startQueueWorker({
pollIntervalMs: 71,
setTimeout(callback) {
syncSets += 1;
if (syncSets === 1) {
firstCallback = callback;
callback();
callback();
} else {
lateCallback = callback;
}
return syncSets;
},
clearTimeout(handle) { cleared.push(handle); },
onError() { syncErrors += 1; },
});
for (let turn = 0; turn < 24 && (syncSets < 2 || syncErrors < 2); turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (syncSets !== 2 || syncErrors !== 2) throw new Error(`synchronous duplicate delivery escaped one-shot ownership: sets=${syncSets}, errors=${syncErrors}`);
await Promise.all([synchronous.stop(), synchronous.stop(), synchronous.stop()]);
firstCallback();
lateCallback();
firstCallback();
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (syncSets !== 2 || syncErrors !== 2 || cleared.length !== 1 || cleared[0] !== 2) {
throw new Error(`late delivery changed the stopped worker: sets=${syncSets}, errors=${syncErrors}, clears=${cleared}`);
}
let throwingSets = 0;
let throwingErrors = 0;
const throwing = startQueueWorker({
pollIntervalMs: 73,
setTimeout() { throwingSets += 1; throw new Error("timer arm threw"); },
clearTimeout() { throw new Error("Failed worker attempted timer cleanup"); },
onError() { throwingErrors += 1; },
});
for (let turn = 0; turn < 24 && throwingSets === 0; turn += 1) await new Promise((resolve) => setImmediate(resolve));
await Promise.all([throwing.stop(), throwing.stop(), throwing.stop()]);
for (let turn = 0; turn < 4; turn += 1) await new Promise((resolve) => setImmediate(resolve));
if (throwingSets !== 1 || throwingErrors !== 1) throw new Error(`throwing arm was not terminal: sets=${throwingSets}, errors=${throwingErrors}`);
if (unhandled.length !== 0) throw new Error(`worker hooks leaked ${unhandled.length} rejection(s)`);
"#,
);
let output = Command::new("node")
.arg("worker-state-machine.mjs")
.current_dir(fixture.root.join("dist"))
.env_remove("DATABASE_URL")
.output()
.expect("execute worker state-machine contract");
assert_success(&output, "assert the complete worker transition table");
}