use std::fs;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::thread;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
struct ComposePostgres {
file: PathBuf,
project: String,
}
impl ComposePostgres {
fn start() -> Option<Self> {
if !Command::new("docker")
.arg("info")
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.is_ok_and(|status| status.success())
{
eprintln!(
"WO-26 Postgres queue drain: SKIP (Docker daemon unavailable; capability, budget, and job-count behavior were not exercised against PostgreSQL)"
);
return None;
}
if !Command::new("docker")
.args(["compose", "version"])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.is_ok_and(|status| status.success())
{
eprintln!(
"WO-26 Postgres queue drain: SKIP (Docker Compose unavailable; capability, budget, and job-count behavior were not exercised against PostgreSQL)"
);
return None;
}
let nonce = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock after epoch")
.as_nanos();
let postgres = Self {
file: Path::new(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/wo19-postgres/compose.yaml"),
project: format!("noxidwo26{}{}", std::process::id(), nonce),
};
let output = postgres
.command()
.args(["up", "--detach"])
.output()
.expect("start queue-drain Postgres");
assert!(
output.status.success(),
"Docker was available but queue-drain Postgres failed to start: {}",
String::from_utf8_lossy(&output.stderr)
);
for _ in 0..480 {
if postgres
.command()
.args([
"exec",
"--no-TTY",
"postgres",
"pg_isready",
"-U",
"noxid_test",
"-d",
"noxid_test",
])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.expect("probe queue-drain Postgres")
.success()
{
return Some(postgres);
}
thread::sleep(Duration::from_millis(250));
}
panic!("queue-drain Postgres did not become ready within 120 seconds");
}
fn command(&self) -> Command {
let mut command = Command::new("docker");
command
.args(["compose", "-f"])
.arg(&self.file)
.args(["--project-name", &self.project]);
command
}
fn database_url(&self) -> String {
let output = self
.command()
.args(["port", "postgres", "5432"])
.output()
.expect("resolve queue-drain Postgres port");
assert!(output.status.success());
let mapping = String::from_utf8(output.stdout).expect("UTF-8 port mapping");
let port = mapping
.trim()
.rsplit_once(':')
.map(|(_, port)| port)
.expect("mapped port");
format!("postgres://noxid_test:noxid_test@127.0.0.1:{port}/noxid_test")
}
}
impl Drop for ComposePostgres {
fn drop(&mut self) {
let _ = self
.command()
.args(["down", "--volumes", "--remove-orphans"])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status();
}
}
#[test]
fn postgres_queue_drain_honors_capability_budget_and_counts() {
let Some(postgres) = ComposePostgres::start() else {
return;
};
let root = std::env::temp_dir().join(format!(
"noxid-wo26-postgres-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock after epoch")
.as_nanos()
));
fs::create_dir_all(root.join("server/queues")).expect("create queue project");
fs::write(root.join("Noxid.toml"), "[app]\ntitle = \"Queue drain\"\n").expect("write config");
fs::write(root.join("package.json"), "{\"type\":\"module\"}\n").expect("write package");
fs::write(
root.join("server/queues/Deliver.nox"),
"queue Deliver { payload { label: String } retry: 0 backoff: 1s }\n",
)
.expect("write queue");
fs::write(
root.join("server/host.js"),
r#"export const queues = Object.freeze({
"queue:Deliver": async ({ label }) => {
await new Promise((resolve) => setTimeout(resolve, 20));
return label;
},
});
export async function authorize({ capability, environment }) {
environment.seenCapability = capability;
return environment.allow === true;
}
"#,
)
.expect("write host");
let build = Command::new(env!("CARGO_BIN_EXE_noxid"))
.args(["build", ".", "--out-dir", "dist"])
.current_dir(&root)
.output()
.expect("build queue project");
assert!(
build.status.success(),
"{}",
String::from_utf8_lossy(&build.stderr)
);
#[cfg(unix)]
{
let repository_modules = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../node_modules");
if !repository_modules.join("postgres").exists() {
eprintln!(
"WO-26 Postgres queue drain: SKIP (workspace postgres Node driver unavailable)"
);
let _ = fs::remove_dir_all(&root);
return;
}
std::os::unix::fs::symlink(&repository_modules, root.join("node_modules"))
.expect("link admitted Postgres driver");
}
fs::write(
root.join("dist/exercise.mjs"),
r#"import { enqueue, fetchEndpoint, queueStatus, closeQueueDatabase } from "./server/handler.js";
for (const label of ["one", "two", "three"]) await enqueue("Deliver", { label });
const url = "http://noxid.test/_noxid/queue/drain";
let environment = { allow: false };
let response = await fetchEndpoint(new Request(url, { method: "POST" }), environment, { noxidQueueDrain: true, queueDrainBudgetMs: 5 });
let body = await response.json();
if (response.status !== 403 || body.error?.code !== "QUEUE_DRAIN_CAPABILITY_DENIED" || environment.seenCapability !== "queue.drain") throw new Error(`capability guard failed: ${response.status} ${JSON.stringify(body)}`);
let status = await queueStatus("Deliver");
if (!status.some((entry) => entry.state === "pending" && entry.count === 3)) throw new Error(`denial claimed work: ${JSON.stringify(status)}`);
environment = { allow: true };
response = await fetchEndpoint(new Request(url, { method: "POST" }), environment, { noxidQueueDrain: true, queueDrainBudgetMs: 5 });
body = await response.json();
if (response.status !== 200 || body.counts?.claimed !== 1 || body.counts?.completed !== 1 || body.counts?.retried !== 0 || body.counts?.deadLettered !== 0) throw new Error(`budget/count response changed: ${response.status} ${JSON.stringify(body)}`);
status = await queueStatus("Deliver");
if (!status.some((entry) => entry.state === "pending" && entry.count === 2)) throw new Error(`budget claimed too many jobs: ${JSON.stringify(status)}`);
response = await fetchEndpoint(new Request(url, { method: "POST" }), environment, { noxidQueueDrain: true, queueDrainBudgetMs: 1000 });
body = await response.json();
if (response.status !== 200 || body.counts?.claimed !== 2 || body.counts?.completed !== 2) throw new Error(`final drain counts changed: ${response.status} ${JSON.stringify(body)}`);
await closeQueueDatabase();
"#,
)
.expect("write drain exercise");
let node = Command::new("node")
.arg("exercise.mjs")
.current_dir(root.join("dist"))
.env("DATABASE_URL", postgres.database_url())
.output()
.expect("run PostgreSQL drain contract");
assert!(
node.status.success(),
"stdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&node.stdout),
String::from_utf8_lossy(&node.stderr)
);
let _ = fs::remove_dir_all(root);
}