#[path = "support/storage_contract.rs"]
mod storage_contract;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::{Command, Output, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::thread;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use storage_contract::{BUILTIN_STORAGE_DRIVERS, STORAGE_API_CONTRACT, StorageContractDriver};
static NEXT_FIXTURE: AtomicU64 = AtomicU64::new(0);
struct Fixture {
root: PathBuf,
}
impl Fixture {
fn new(label: &str, driver: StorageContractDriver) -> Self {
let ordinal = NEXT_FIXTURE.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"noxid-wo27-{label}-{}-{ordinal}",
std::process::id()
));
fs::create_dir_all(root.join("server/api")).expect("create WO-27 fixture");
fs::write(root.join("package.json"), "{\"type\":\"module\"}\n")
.expect("write package marker");
fs::write(
root.join("Noxid.toml"),
format!(
"[app]\ntitle = \"WO-27 storage\"\n\n[server]\nstorage = \"{}\"\ndb_pool = 3\n",
driver.manifest_value
),
)
.expect("write manifest");
fs::write(
root.join("server/host.js"),
"import { storage } from \"noxid:server\";\nexport const stores = { direct: storage(\"direct\") };\n",
)
.expect("write storage host");
fs::write(
root.join("server/api/contract.get.nox"),
"endpoint Contract { result: Int handler { return 1 } }\n",
)
.expect("write storage fixture endpoint");
Self { root }
}
fn build(&self) -> Output {
Command::new(env!("CARGO_BIN_EXE_noxid"))
.args(["build", ".", "--out-dir", "dist"])
.current_dir(&self.root)
.output()
.expect("build WO-27 storage fixture")
}
fn run_node(&self, name: &str, source: &str, database_url: Option<&str>) -> Output {
let script = self.root.join("dist").join(format!("{name}.mjs"));
fs::write(&script, source).expect("write storage exercise");
let mut command = Command::new("node");
command
.arg(script.file_name().expect("script filename"))
.current_dir(self.root.join("dist"))
.env("NOXID_STORAGE_DIR", self.root.join("storage"));
if let Some(database_url) = database_url {
command.env("DATABASE_URL", database_url);
}
command.output().expect("run storage exercise")
}
}
impl Drop for Fixture {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.root);
}
}
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-27 Postgres storage integration: SKIP (Docker daemon unavailable; shared rate limit, idempotency, schema, and TTL sweep were not exercised)"
);
return None;
}
if !Command::new("docker")
.args(["compose", "version"])
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.is_ok_and(|status| status.success())
{
eprintln!(
"WO-27 Postgres storage integration: SKIP (Docker Compose unavailable; shared rate limit, idempotency, schema, and TTL sweep were not exercised)"
);
return None;
}
let nonce = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock after epoch")
.as_nanos();
let ordinal = NEXT_FIXTURE.fetch_add(1, Ordering::Relaxed);
let postgres = Self {
file: Path::new(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/wo19-postgres/compose.yaml"),
project: format!("noxidwo27{}{}{}", std::process::id(), nonce, ordinal),
};
let output = postgres
.command()
.args(["up", "--detach"])
.output()
.expect("start storage Postgres");
assert!(
output.status.success(),
"Docker was available but storage 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 storage Postgres")
.success()
{
return Some(postgres);
}
thread::sleep(Duration::from_millis(250));
}
panic!("storage 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 storage 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();
}
}
fn assert_success(output: &Output, phase: &str) {
assert!(
output.status.success(),
"{phase} failed\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
fn exercise_concurrent_postgres_cold_start(fixture: &Fixture, database_url: &str) {
let ready = fixture.root.join("cold-ready.txt");
let start = fixture.root.join("cold-start.txt");
let script = fixture.root.join("dist/cold-start.mjs");
fs::write(
&script,
r#"import fs from "node:fs";
import { storage } from "./server/noxid-server.js";
fs.appendFileSync(process.env.READY_FILE, `${process.pid}\n`);
while (!fs.existsSync(process.env.START_FILE)) await new Promise((resolve) => setTimeout(resolve, 2));
await storage("cold-start").set(String(process.pid), true);
process.exit(0);
"#,
)
.expect("write concurrent cold-start exercise");
let replicas = 6;
let mut children = Vec::new();
for _ in 0..replicas {
children.push(
Command::new("node")
.arg(script.file_name().expect("cold-start filename"))
.current_dir(fixture.root.join("dist"))
.env("DATABASE_URL", database_url)
.env("READY_FILE", &ready)
.env("START_FILE", &start)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn cold Postgres replica"),
);
}
let deadline = Instant::now() + Duration::from_secs(15);
while fs::read_to_string(&ready)
.unwrap_or_default()
.lines()
.count()
< replicas
{
assert!(
Instant::now() < deadline,
"not all cold Postgres replicas reached the barrier"
);
thread::sleep(Duration::from_millis(10));
}
fs::write(&start, "go\n").expect("release cold Postgres replicas");
for child in children {
let output = child
.wait_with_output()
.expect("wait for cold Postgres replica");
assert_success(&output, "initialize Postgres storage concurrently");
}
}
#[cfg(unix)]
fn link_workspace_node_modules(fixture: &Fixture) -> bool {
let modules = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../node_modules");
if !modules.join("postgres").exists() {
eprintln!(
"WO-27 Postgres storage integration: SKIP (workspace postgres Node driver unavailable; shared storage paths were not exercised)"
);
return false;
}
std::os::unix::fs::symlink(modules, fixture.root.join("node_modules"))
.expect("link admitted postgres driver");
true
}
#[test]
fn storage_contract_is_driver_parameterized_and_postgres_sweeps_ttl_rows() {
let postgres = ComposePostgres::start();
let database_url = postgres.as_ref().map(ComposePostgres::database_url);
for driver in BUILTIN_STORAGE_DRIVERS {
if driver.needs_redis {
continue;
}
if driver.needs_database && postgres.is_none() {
continue;
}
let fixture = Fixture::new(driver.name, *driver);
let build = fixture.build();
assert_success(&build, &format!("build {} storage", driver.name));
if driver.needs_database {
let runtime = fs::read_to_string(fixture.root.join("dist/server/noxid-server.js"))
.expect("read Postgres storage runtime");
assert!(runtime.contains("CREATE TABLE IF NOT EXISTS _noxid_storage"));
assert!(runtime.contains("postgres(url, { max: 3 })"));
}
if driver.needs_database && !link_workspace_node_modules(&fixture) {
continue;
}
if driver.needs_database {
exercise_concurrent_postgres_cold_start(
&fixture,
database_url.as_deref().expect("Postgres database URL"),
);
}
let first = fixture.run_node(
"contract-first",
STORAGE_API_CONTRACT,
database_url.as_deref().filter(|_| driver.needs_database),
);
assert_success(&first, &format!("execute {} storage contract", driver.name));
let second = fixture.run_node(
"contract-second",
&storage_contract::persistence_check(driver.persists_across_processes),
database_url.as_deref().filter(|_| driver.needs_database),
);
assert_success(
&second,
&format!("verify {} storage persistence", driver.name),
);
if driver.needs_database {
let inspect = fixture.run_node(
"inspect-postgres",
r#"import postgres from "postgres";
const sql = postgres(process.env.DATABASE_URL, { max: 1 });
const columns = await sql`SELECT column_name, data_type FROM information_schema.columns WHERE table_name = '_noxid_storage' ORDER BY ordinal_position`;
const expected = [["namespace","text"],["key","text"],["value","jsonb"],["expires_at","timestamp with time zone"]];
if (JSON.stringify(columns.map((row) => [row.column_name, row.data_type])) !== JSON.stringify(expected)) throw new Error(`schema changed: ${JSON.stringify(columns)}`);
const expired = await sql`SELECT count(*)::int AS count FROM _noxid_storage WHERE expires_at IS NOT NULL AND expires_at <= now()`;
if (expired[0]?.count !== 0) throw new Error(`expired rows were not swept: ${JSON.stringify(expired)}`);
await sql.end({ timeout: 1 });
"#,
database_url.as_deref(),
);
assert_success(&inspect, "inspect postgres storage schema and sweep");
}
}
}
#[test]
fn two_processes_share_postgres_rate_limit_and_idempotency_replay() {
let Some(postgres) = ComposePostgres::start() else {
return;
};
let driver = BUILTIN_STORAGE_DRIVERS
.iter()
.find(|driver| driver.name == "postgres")
.copied()
.expect("postgres driver descriptor");
let fixture = Fixture::new("shared-endpoints", driver);
fs::create_dir_all(fixture.root.join("server/api")).expect("create endpoint directory");
fs::write(
fixture.root.join("server/api/rate.post.nox"),
"endpoint SharedRate { result: Int limit: 1 per minute per ip handler { return 1 } }\n",
)
.expect("write rate endpoint");
fs::write(
fixture.root.join("server/api/save.post.nox"),
"endpoint SharedReplay { body { value: Int } result: Int idempotent handler { return value } }\n",
)
.expect("write idempotent endpoint");
let build = fixture.build();
assert_success(&build, "build shared Postgres endpoint fixture");
if !link_workspace_node_modules(&fixture) {
return;
}
let database_url = postgres.database_url();
let first = fixture.run_node(
"instance-one",
r#"import { fetch as handle } from "./server/handler.js";
let response = await handle(new Request("http://noxid.test/api/rate", { method: "POST" }), { ip: "203.0.113.27" });
if (response.status !== 200) throw new Error(`first rate failed: ${response.status} ${await response.text()}`);
response = await handle(new Request("http://noxid.test/api/save", { method: "POST", headers: { "content-type": "application/json", "idempotency-key": "shared" }, body: JSON.stringify({ value: 7 }) }), { ip: "203.0.113.28" });
const body = await response.json();
if (response.status !== 200 || body.value !== 7) throw new Error(`first replay failed: ${response.status} ${JSON.stringify(body)}`);
process.exit(0);
"#,
Some(&database_url),
);
assert_success(&first, "seed Postgres state from instance one");
let second = fixture.run_node(
"instance-two",
r#"import { fetch as handle } from "./server/handler.js";
let response = await handle(new Request("http://noxid.test/api/rate", { method: "POST" }), { ip: "203.0.113.27" });
let body = await response.json();
if (response.status !== 429 || body.error?.code !== "ENDPOINT_RATE_LIMITED") throw new Error(`combined rate limit failed: ${response.status} ${JSON.stringify(body)}`);
response = await handle(new Request("http://noxid.test/api/save", { method: "POST", headers: { "content-type": "application/json", "idempotency-key": "shared" }, body: JSON.stringify({ value: 99 }) }), { ip: "203.0.113.28" });
body = await response.json();
if (response.status !== 200 || body.value !== 7) throw new Error(`shared replay failed: ${response.status} ${JSON.stringify(body)}`);
process.exit(0);
"#,
Some(&database_url),
);
assert_success(&second, "reuse Postgres state from instance two");
}