noxid-cli 0.2.0

The Noxid compiler command line: check, build, test, adapt, and the agent surface
#[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");
}