diff --git a/CLAUDE.md b/CLAUDE.md index f1e896b..2441a7a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -97,6 +97,7 @@ Covered in `docs/runtime/database/02-wo-language.md § Concurrency Model` and `0 | `engine.rs` (Engine + Catalog) | `engine` | 2 | | `compile.rs` | `engine` | 2 | | `method.rs` (13b method executor) | `logic` | 6 | +| `pg.rs` + `mirror.rs` (16 Postgres backup mirror) | `db` | 3+ | | `server.rs` | `http` + `service` | 4 / 6 | | `bin/wo.rs` | stays in `rt` (the binary) | — | @@ -129,6 +130,7 @@ Stage-3 stubs (501) and policy-shaped 405/404 responses are **intentional and do - **`reference/crates/` is its own workspace.** Running `cargo build` at the root does not build v1. Running it in `reference/crates/` does. - **Parser identifiers vs. keywords.** `subscribe`, `receive`, `expect_abort`, `me`, `self`, and lowercase `insert`/`select` are NOT keywords in the lexer — they stay as plain idents (only SQL-layer `INSERT`/`SELECT` are keywords). The 13b statement parser matches `insert` as an ident; the select expression is recognised by the two-token shape `select {`. Adding these to the keyword map breaks `service rest "..." expose subscribe` and method bodies. - **Type-level annotations.** `@table(name: "...", index: [a, b])` before a `type`/`class` configures storage (it never toggles table-ness — every type IS a table). Unknown keys inside `@table(...)` are parse errors; unknown annotation *names* (`@foo`) skip silently. Engine secondary indexes are maintained ONLY via `Engine::row_insert`/`row_remove` — never touch `tables` directly or indexes drift. +- **The Postgres mirror is a backup, never a commit path.** `WO_PG=…` streams committed mutations to Postgres asynchronously (`mirror.rs`, one `wo-pg` thread, hand-rolled wire client in `pg.rs` — no crates). Reads and acks must never depend on it: taps sit AFTER `wal_log` accepts, use `try_send`, and drop loudly on overflow. Boot re-pushes all RAM state (`mirror_sync_all`), so RAM stays authoritative and Postgres is always reconstructible from a restart. See `docs/plan/16-postgres-mirror.md`. - **Parser skip-on-block.** Unknown triggers (`on update do ...`) are parsed-and-discarded by brace-depth-aware skipping. Object literals like `{ article_id: self.id }` inside trigger actions contain `}` that must not be mistaken for the type's outer close brace — the depth counter exists specifically because of this. Exception since 13b: `fn` inside a `class` parses into a real `MethodDecl` (body statements + expressions, executed by `method.rs`); `fn` inside a plain `type` still skips. - **Newline significance.** The lexer emits `Kind::Newline` tokens and the parser uses them to end policy/trigger lines. Do not filter newlines globally. - **Default-value parsing.** `= now()` is recognised explicitly as `DefaultExpr::Now`; anything else falls into an opaque-expression path that `engine::eval_default` then **omits from created rows** (computed fields display as empty, not as debug-printed tokens). diff --git a/README.md b/README.md index e7d307d..d12bb71 100644 --- a/README.md +++ b/README.md @@ -4,6 +4,11 @@ A declarative full-stack programming language. You write `.wo` files; the runtim Think **Go + Postgres + `net/http` + Phoenix LiveView, folded into one language and one binary.** +# persistant database + +- reads and writes database to RAM, persist data to postgres SQL. +- The entire database lives in RAM; every committed write is mirrored to PostgreSQL **as a backup** — asynchronously, behind the runtime's own WAL, never in the read or ack path. Set `WO_PG=postgres://user@host:5432/db` and every type's rows appear as a Postgres table (named by its `@table(name: ...)` annotation) that you can query with plain `psql`. Plan and phases: [`docs/plan/16-postgres-mirror.md`](docs/plan/16-postgres-mirror.md); try it: `just pricing-pg-demo`. + ## Quickstart ```bash @@ -17,28 +22,28 @@ See [`reference/rest/blog.rest`](reference/rest/blog.rest) for a preconfigured H ## What this repository contains -| Path | What it is | -| --- | --- | -| [`crates/rt/`](crates/rt/) | The new `.wo` language runtime — lexer, type-DSL parser, in-memory engine, axum REST server. Produces the `wo` binary. | -| [`crates/{ql,value,engine,txn,db,wal,sub,http,gen,policy,logic,service,ui,app}/`](crates/) | 14 empty placeholder crates scaffolded for Phases 2–6. Real code extracts from `rt/` as each phase activates. | -| [`docs/runtime/wo-language.md`](docs/runtime/wo-language.md) | **Start here.** The language overview: toolchain, hello-world, stdlib, client model. | -| [`docs/runtime/database.md`](docs/runtime/database.md) | The 7-phase engineering series that drives the runtime's design. | -| [`docs/examples/blog/`](docs/examples/blog/) | Sample `.wo` project: blog with articles, authors, tags, comments. ~200 lines. | -| [`docs/examples/ecommerce/`](docs/examples/ecommerce/) | Sample `.wo` project: storefront + live order-ops table + cross-paradigm checkout. ~300 lines. | -| [`prototypes/wo-db/`](prototypes/wo-db/) | C++ prototype of the query-layer engine (SQL + Cypher + document paths, `RETURNING` aliases, `LIVE` stub). ~2k lines, smoke tests pass. Reference implementation the Rust port follows. | -| [`reference/rest/`](reference/rest/) | `.rest` files (VS Code REST Client / JetBrains HTTP format) for manually testing the running prototype. | -| [`reference/crates/`](reference/crates/) | The v1 writeonce blog — 13 Rust crates implementing the original `.seg` + sidecar-index storage engine and `.htmlx` templating. Preserved as a nested workspace; see [`reference/README.md`](reference/README.md). | +| Path | What it is | +| ------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | +| [`crates/rt/`](crates/rt/) | The new `.wo` language runtime — lexer, type-DSL parser, in-memory engine, axum REST server. Produces the `wo` binary. | +| [`crates/{ql,value,engine,txn,db,wal,sub,http,gen,policy,logic,service,ui,app}/`](crates/) | 14 empty placeholder crates scaffolded for Phases 2–6. Real code extracts from `rt/` as each phase activates. | +| [`docs/runtime/wo-language.md`](docs/runtime/wo-language.md) | **Start here.** The language overview: toolchain, hello-world, stdlib, client model. | +| [`docs/runtime/database.md`](docs/runtime/database.md) | The 7-phase engineering series that drives the runtime's design. | +| [`docs/examples/blog/`](docs/examples/blog/) | Sample `.wo` project: blog with articles, authors, tags, comments. ~200 lines. | +| [`docs/examples/ecommerce/`](docs/examples/ecommerce/) | Sample `.wo` project: storefront + live order-ops table + cross-paradigm checkout. ~300 lines. | +| [`prototypes/wo-db/`](prototypes/wo-db/) | C++ prototype of the query-layer engine (SQL + Cypher + document paths, `RETURNING` aliases, `LIVE` stub). ~2k lines, smoke tests pass. Reference implementation the Rust port follows. | +| [`reference/rest/`](reference/rest/) | `.rest` files (VS Code REST Client / JetBrains HTTP format) for manually testing the running prototype. | +| [`reference/crates/`](reference/crates/) | The v1 writeonce blog — 13 Rust crates implementing the original `.seg` + sidecar-index storage engine and `.htmlx` templating. Preserved as a nested workspace; see [`reference/README.md`](reference/README.md). | ## Current stage The runtime is under active development. Each stage lands as an independently shippable cut: -| Stage | What works | Status | -| --- | --- | --- | -| **1** | `wo run ` discovers every `.wo` file under a directory | ✅ shipped | -| **2** | Type-DSL parser, in-memory engine, REST CRUD (`list` / `get` / `create` / `update` / `delete`) generated from `service rest` blocks, JSON bodies with auto-id, default-value seeding, partial-update PATCH | ✅ shipped — `cargo run -- run docs/examples/blog` | -| **3** | LIVE subscriptions over WebSocket, delta frames on commit, `me` / session layer | pending | -| **4+** | Transactional fns (`fn checkout in txn snapshot`), row-level policies, type-attached triggers, `##ui` SSR, WAL durability, codegen | see [docs/runtime/database.md](docs/runtime/database.md) | +| Stage | What works | Status | +| ------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------- | +| **1** | `wo run ` discovers every `.wo` file under a directory | ✅ shipped | +| **2** | Type-DSL parser, in-memory engine, REST CRUD (`list` / `get` / `create` / `update` / `delete`) generated from `service rest` blocks, JSON bodies with auto-id, default-value seeding, partial-update PATCH | ✅ shipped — `cargo run -- run docs/examples/blog` | +| **3** | LIVE subscriptions over WebSocket, delta frames on commit, `me` / session layer | pending | +| **4+** | Transactional fns (`fn checkout in txn snapshot`), row-level policies, type-attached triggers, `##ui` SSR, WAL durability, codegen | see [docs/runtime/database.md](docs/runtime/database.md) | `cargo test --lib` at the root runs 14 unit tests covering the lexer, parser, compiler, and engine. Stage-3 endpoints respond `501 Not Implemented` until they land. diff --git a/crates/rt/src/bin/wo.rs b/crates/rt/src/bin/wo.rs index 741f247..f0f06e5 100644 --- a/crates/rt/src/bin/wo.rs +++ b/crates/rt/src/bin/wo.rs @@ -25,6 +25,9 @@ USAGE: ENV: WO_LISTEN override the listen address (default: 127.0.0.1:8080) + WO_PG postgres://user[:pass]@host[:port]/db — mirror every + committed write to Postgres as a backup (reads stay in + RAM; see docs/plan/16-postgres-mirror.md) " ); } @@ -116,6 +119,25 @@ fn run(dir: PathBuf) -> anyhow::Result { } } + // Postgres backup mirror (plan 16b): RAM stays authoritative — the + // `wo-pg` thread receives every committed mutation on a bounded channel + // and upserts it as JSONB. Never in the ack path; off unless WO_PG set. + let mirror_tx = match std::env::var("WO_PG") { + Ok(url) => { + let cfg = rt::pg::PgConfig::from_url(&url) + .map_err(|e| anyhow::anyhow!("WO_PG: {e}"))?; + let tables: Vec<(String, String)> = catalog.order.iter() + .map(|name| (name.clone(), catalog.get(name).unwrap().storage_name.clone())) + .collect(); + let (tx, rx) = std::sync::mpsc::sync_channel(rt::mirror::QUEUE_CAP); + rt::mirror::spawn(cfg.clone(), rx, tables); + println!("[wo] postgres mirror: {}:{}/{} (backup only — reads stay in RAM)", + cfg.host, cfg.port, cfg.database); + Some(tx) + } + Err(_) => None, + }; + let bus = rt::shard::ShardBus::new(n)?; let catalog_for_workers = catalog.clone(); rt::runtime::scheduler::serve(&addr, move |id| { @@ -152,6 +174,12 @@ fn run(dir: PathBuf) -> anyhow::Result { Err(e) => eprintln!("[wo] shard {id}: WAL unavailable ({e}) — running non-durable"), } } + // Mirror attaches AFTER replay: the replayed state goes to Postgres + // once, as a boot-time bulk sync, then live mutations stream. + if let Some(tx) = &mirror_tx { + engine.attach_mirror(tx.clone()); + engine.mirror_sync_all(); + } let ctx = rt::shard::ShardCtx::new(id, n, engine, bus.clone()); let router = rt::server::router(ctx.clone(), &catalog_for_workers); let mail_fd = bus.mail_fd(id).as_raw_fd(); diff --git a/crates/rt/src/engine.rs b/crates/rt/src/engine.rs index 26faef8..7bb58ff 100644 --- a/crates/rt/src/engine.rs +++ b/crates/rt/src/engine.rs @@ -84,12 +84,22 @@ pub struct Engine { /// records here and journal undo entries; `commit_txn` emits one /// [`WalRec::Txn`] frame, `abort_txn` reverts RAM in reverse order. txn: Option, + /// Postgres backup mirror (plan 16b): committed mutations are cloned + /// onto this channel AFTER the WAL accepted them — the mirror never + /// gates an ack. `None` when `WO_PG` is unset. + mirror: Option, + /// Records dropped because the mirror channel was full/closed — + /// counted per shard, logged loudly but never blocking. + mirror_dropped: u64, } #[derive(Debug, Default)] struct TxnState { - wal: Vec, - undo: Vec, + wal: Vec, + undo: Vec, + /// Mirror records for this transaction — sent as ONE + /// [`crate::mirror::MirrorRec::Txn`] on commit, dropped on abort. + mirror: Vec, } /// Inverse of one applied mutation — enough to restore the pre-txn RAM state. @@ -129,7 +139,53 @@ impl Engine { } } Self { catalog, tables, indexes, next_id, id_step: n_shards.max(1) as i64, - wal: None, staged: false, txn: None } + wal: None, staged: false, txn: None, mirror: None, mirror_dropped: 0 } + } + + /// Attach the Postgres backup mirror (plan 16b). Like `attach_wal`, + /// this happens AFTER boot replay — replayed rows are pushed once via + /// [`mirror_sync_all`], not re-mirrored record by record. + pub fn attach_mirror(&mut self, tx: crate::mirror::MirrorSender) { + self.mirror = Some(tx); + } + + /// Enqueue this shard's ENTIRE current state as upserts — boot-time + /// initial sync so a fresh Postgres catches up with a replayed WAL. + pub fn mirror_sync_all(&mut self) { + if self.mirror.is_none() { return; } + let snapshot: Vec<(String, i64, Row)> = self.tables.iter() + .flat_map(|(ty, table)| table.iter() + .map(|(id, row)| (ty.clone(), *id, row.clone()))) + .collect(); + for (ty, id, row) in snapshot { + self.mirror_dispatch(crate::mirror::MirrorRec::Upsert { ty, id, row }); + } + } + + /// Route one committed mutation to the mirror: buffered while a method + /// transaction is open (sent atomically on commit), dispatched + /// immediately otherwise. No-op when no mirror is attached. + fn mirror_send(&mut self, rec: crate::mirror::MirrorRec) { + if self.mirror.is_none() { return; } + match self.txn.as_mut() { + Some(t) => t.mirror.push(rec), + None => self.mirror_dispatch(rec), + } + } + + /// Non-blocking send; a full or closed channel drops the record and + /// counts it — Postgres lags, clients never do (16b policy; the + /// dirty-flag resync is plan 16d). + fn mirror_dispatch(&mut self, rec: crate::mirror::MirrorRec) { + let Some(tx) = self.mirror.as_ref() else { return }; + if tx.try_send(rec).is_err() { + self.mirror_dropped += 1; + if self.mirror_dropped.is_power_of_two() { + eprintln!("[wo] pg mirror: queue full/closed — {} records dropped on this \ + shard (Postgres is behind RAM until resync, plan 16d)", + self.mirror_dropped); + } + } } /// Attach a per-commit WAL (fsync inside each mutation). Must happen @@ -245,6 +301,11 @@ impl Engine { self.apply_undo(t.undo); // never ack non-durable return Err(e); } + // Mirror the whole method as one atomic Postgres transaction — + // committed only; an aborted method never reaches this point. + if !t.mirror.is_empty() { + self.mirror_dispatch(crate::mirror::MirrorRec::Txn(t.mirror)); + } Ok(()) } @@ -336,6 +397,11 @@ impl Engine { self.row_remove(ty, id); // never ack non-durable return Err(e); } + if self.mirror.is_some() { + self.mirror_send(crate::mirror::MirrorRec::Upsert { + ty: ty.into(), id, row: row.clone(), + }); + } Ok(row) } @@ -361,6 +427,12 @@ impl Engine { self.row_insert(ty, id, prev); return Err(e); } + if self.mirror.is_some() { + // The mirror needs the FULL post-merge row, not the merge body. + self.mirror_send(crate::mirror::MirrorRec::Upsert { + ty: ty.into(), id, row: updated.clone(), + }); + } Ok(Some(updated)) } @@ -374,6 +446,7 @@ impl Engine { self.row_insert(ty, id, removed); // undo return Err(e); } + self.mirror_send(crate::mirror::MirrorRec::Delete { ty: ty.into(), id }); Ok(true) } diff --git a/crates/rt/src/lib.rs b/crates/rt/src/lib.rs index 32e8c5d..7e4b907 100644 --- a/crates/rt/src/lib.rs +++ b/crates/rt/src/lib.rs @@ -10,7 +10,9 @@ pub mod engine; pub mod http; pub mod lexer; pub mod method; +pub mod mirror; pub mod parser; +pub mod pg; pub mod runtime; pub mod server; pub mod shard; diff --git a/crates/rt/src/mirror.rs b/crates/rt/src/mirror.rs new file mode 100644 index 0000000..d341b50 --- /dev/null +++ b/crates/rt/src/mirror.rs @@ -0,0 +1,270 @@ +//! PostgreSQL backup mirror — plan 16b (`docs/plan/16-postgres-mirror.md`). +//! +//! RAM is authoritative; this module is the **backup mechanism**: every +//! committed mutation is cloned onto an mpsc channel by the shard engines +//! (after the WAL made it durable — never before, never gating the ack) and +//! a single dedicated `wo-pg` thread drains the channel into Postgres as +//! JSONB upserts. Reads never touch Postgres. +//! +//! Failure doctrine (16b): Postgres being down costs clients nothing — the +//! thread reconnects with capped backoff while the bounded channel absorbs +//! the burst; if the channel fills, records are dropped **loudly** (counted +//! and logged). Lossless catch-up (dirty-flag full resync) is plan 16d; +//! restore-from-Postgres at boot is 16e. +//! +//! Schema (16b): one table per type — `"" (id BIGINT PRIMARY +//! KEY, row JSONB NOT NULL)` — the `@table(name: "prices")` annotation +//! names the table. Typed-column projection is 16c. + +use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender}; +use std::time::Duration; + +use serde_json::Value; + +use crate::engine::Row; +use crate::pg::{escape_ident, escape_literal, Conn, PgConfig}; + +/// Mirror channel capacity. At ~200 bytes/record this bounds the buffered +/// backlog around a few tens of MB — enough to ride out a Postgres restart +/// under load without threatening the RAM budget. +pub const QUEUE_CAP: usize = 65_536; + +/// One committed mutation, as the mirror needs it. Unlike `WalRec::Update` +/// (which carries the merge body), `Upsert` always carries the FULL +/// post-merge row — the mirror's `ON CONFLICT ... DO UPDATE` replaces the +/// whole JSONB value. +#[derive(Debug, Clone)] +pub enum MirrorRec { + Upsert { ty: String, id: i64, row: Row }, + Delete { ty: String, id: i64 }, + /// One method transaction (plan 13b) — applied inside one Postgres + /// transaction, mirroring the WAL's atomic `WalRec::Txn` frame. + Txn(Vec), +} + +pub type MirrorSender = SyncSender; + +/// Spawn the `wo-pg` mirror thread. `tables` maps type name → storage +/// (table) name for every catalog type; DDL is bootstrapped on every +/// (re)connect so a fresh database works out of the box. +pub fn spawn( + cfg: PgConfig, + rx: Receiver, + tables: Vec<(String, String)>, +) -> std::thread::JoinHandle<()> { + std::thread::Builder::new() + .name("wo-pg".into()) + .spawn(move || run(cfg, rx, tables)) + .expect("spawn wo-pg mirror thread") +} + +fn run(cfg: PgConfig, rx: Receiver, tables: Vec<(String, String)>) { + let mut dropped: u64 = 0; + loop { + // (Re)connect with capped backoff, bootstrapping DDL each time. + let Some(mut c) = connect_with_backoff(&cfg, &tables, &rx, &mut dropped) else { + return; // channel closed — shutdown + }; + + eprintln!("[wo] pg mirror: connected to {}:{}/{} ({} tables)", + cfg.host, cfg.port, cfg.database, tables.len()); + if dropped > 0 { + eprintln!("[wo] pg mirror: WARNING — {dropped} records were dropped while \ + disconnected; Postgres is behind RAM until a resync (plan 16d)"); + } + + // Drain loop: batch what's queued, one round-trip per batch. + loop { + let batch = match next_batch(&rx) { + Some(b) => b, + None => return, // senders gone — shutdown + }; + if let Err(e) = apply_batch(&mut c, &tables, &batch) { + eprintln!("[wo] pg mirror: connection lost ({e}) — reconnecting"); + break; // outer loop reconnects + } + } + } +} + +/// Block for the next record, then opportunistically drain up to a batch. +/// `None` = all senders dropped (process shutting down). +fn next_batch(rx: &Receiver) -> Option> { + const BATCH: usize = 512; + let first = rx.recv().ok()?; + let mut batch = vec![first]; + while batch.len() < BATCH { + match rx.try_recv() { + Ok(rec) => batch.push(rec), + Err(_) => break, + } + } + Some(batch) +} + +/// `None` = every sender is gone (process shutting down). +fn connect_with_backoff( + cfg: &PgConfig, + tables: &[(String, String)], + rx: &Receiver, + dropped: &mut u64, +) -> Option { + let mut delay = Duration::from_millis(200); + loop { + match Conn::connect(cfg) { + Ok(mut c) => match bootstrap_ddl(&mut c, tables) { + Ok(()) => return Some(c), + Err(e) => eprintln!("[wo] pg mirror: DDL bootstrap failed ({e}) — retrying"), + }, + Err(e) => eprintln!("[wo] pg mirror: connect failed ({e}) — retrying in {delay:?}"), + } + // While waiting, keep the channel from silently backing up forever: + // absorb what we can into the void, counting the loss (16b policy — + // 16d replaces this with dirty-flag resync). + let wait_until = std::time::Instant::now() + delay; + loop { + let left = wait_until.saturating_duration_since(std::time::Instant::now()); + if left.is_zero() { break; } + match rx.recv_timeout(left.min(Duration::from_millis(100))) { + Ok(_) => { *dropped += 1; } + Err(RecvTimeoutError::Timeout) => {} + Err(RecvTimeoutError::Disconnected) => return None, + } + } + delay = (delay * 2).min(Duration::from_secs(5)); + } +} + +fn bootstrap_ddl(c: &mut Conn, tables: &[(String, String)]) -> Result<(), crate::pg::PgError> { + for (_, storage) in tables { + c.simple_query(&format!( + "CREATE TABLE IF NOT EXISTS {} (id BIGINT PRIMARY KEY, row JSONB NOT NULL)", + escape_ident(storage)))?; + } + Ok(()) +} + +/// Apply one batch. Statement/data errors are isolated per record and +/// logged (the batch continues); only I/O errors propagate (→ reconnect). +fn apply_batch( + c: &mut Conn, + tables: &[(String, String)], + batch: &[MirrorRec], +) -> Result<(), crate::pg::PgError> { + for rec in batch { + let sql = rec_sql(tables, rec); + match c.simple_query(&sql) { + Ok(_) => {} + Err(e) if e.severity == "CLIENT" => return Err(e), // socket-level: reconnect + Err(e) => eprintln!("[wo] pg mirror: statement rejected ({e}) — record skipped"), + } + } + Ok(()) +} + +/// Render one record as SQL. A `Txn` becomes BEGIN; …; COMMIT in a single +/// simple-query message — atomic on the Postgres side like its WAL frame. +fn rec_sql(tables: &[(String, String)], rec: &MirrorRec) -> String { + match rec { + MirrorRec::Upsert { ty, id, row } => { + let table = storage_for(tables, ty); + let json = serde_json::to_string(&Value::Object(row.clone())) + .unwrap_or_else(|_| "{}".into()); + format!( + "INSERT INTO {} (id, row) VALUES ({}, {}::jsonb) \ + ON CONFLICT (id) DO UPDATE SET row = EXCLUDED.row", + escape_ident(table), id, escape_literal(&json)) + } + MirrorRec::Delete { ty, id } => { + format!("DELETE FROM {} WHERE id = {}", escape_ident(storage_for(tables, ty)), id) + } + MirrorRec::Txn(recs) => { + let mut sql = String::from("BEGIN"); + for r in recs { + sql.push_str("; "); + sql.push_str(&rec_sql(tables, r)); + } + sql.push_str("; COMMIT"); + sql + } + } +} + +fn storage_for<'a>(tables: &'a [(String, String)], ty: &'a str) -> &'a str { + tables.iter() + .find(|(t, _)| t == ty) + .map(|(_, s)| s.as_str()) + .unwrap_or(ty) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn row(v: Value) -> Row { + match v { Value::Object(m) => m, _ => panic!() } + } + + #[test] + fn sql_rendering_upsert_delete_txn() { + let tables = vec![("Price".to_string(), "prices".to_string())]; + let up = MirrorRec::Upsert { + ty: "Price".into(), id: 2, + row: row(json!({"amount": 4999, "note": "it's"})), + }; + let sql = rec_sql(&tables, &up); + assert!(sql.starts_with(r#"INSERT INTO "prices" (id, row) VALUES (2, '{"#), "{sql}"); + assert!(sql.contains("''s"), "quote must be doubled: {sql}"); + assert!(sql.ends_with("ON CONFLICT (id) DO UPDATE SET row = EXCLUDED.row")); + + let del = MirrorRec::Delete { ty: "Price".into(), id: 7 }; + assert_eq!(rec_sql(&tables, &del), r#"DELETE FROM "prices" WHERE id = 7"#); + + // Unmapped type falls back to the type name. + let other = MirrorRec::Delete { ty: "Ghost".into(), id: 1 }; + assert_eq!(rec_sql(&tables, &other), r#"DELETE FROM "Ghost" WHERE id = 1"#); + + let txn = MirrorRec::Txn(vec![up, del]); + let sql = rec_sql(&tables, &txn); + assert!(sql.starts_with("BEGIN; ")); + assert!(sql.ends_with("; COMMIT")); + } + + /// Integration: full pipeline against a live server (WO_PG_TEST gated). + #[test] + fn mirror_pipeline_against_live_server() { + let Ok(url) = std::env::var("WO_PG_TEST") else { + eprintln!("mirror_pipeline: skipped (set WO_PG_TEST=postgres://... to run)"); + return; + }; + let cfg = PgConfig::from_url(&url).unwrap(); + { + let mut c = Conn::connect(&cfg).unwrap(); + c.simple_query(r#"DROP TABLE IF EXISTS "mirror_prices""#).unwrap(); + } + + let tables = vec![("Price".to_string(), "mirror_prices".to_string())]; + let (tx, rx) = std::sync::mpsc::sync_channel(QUEUE_CAP); + let handle = spawn(cfg.clone(), rx, tables); + + tx.send(MirrorRec::Upsert { + ty: "Price".into(), id: 1, row: row(json!({"amount": 100})), + }).unwrap(); + tx.send(MirrorRec::Txn(vec![ + MirrorRec::Upsert { ty: "Price".into(), id: 3, row: row(json!({"amount": 300})) }, + MirrorRec::Upsert { ty: "Price".into(), id: 1, row: row(json!({"amount": 150})) }, + ])).unwrap(); + tx.send(MirrorRec::Delete { ty: "Price".into(), id: 3 }).unwrap(); + drop(tx); // close channel → thread drains and exits + handle.join().unwrap(); + + let mut c = Conn::connect(&cfg).unwrap(); + let r = c.simple_query( + r#"SELECT id, row->>'amount' FROM "mirror_prices" ORDER BY id"#).unwrap(); + assert_eq!(r.rows.len(), 1, "id 3 deleted, id 1 remains: {:?}", r.rows); + assert_eq!(r.rows[0][0].as_deref(), Some("1")); + assert_eq!(r.rows[0][1].as_deref(), Some("150"), "txn upsert must have applied"); + c.simple_query(r#"DROP TABLE "mirror_prices""#).unwrap(); + } +} diff --git a/crates/rt/src/pg.rs b/crates/rt/src/pg.rs new file mode 100644 index 0000000..4b094ea --- /dev/null +++ b/crates/rt/src/pg.rs @@ -0,0 +1,460 @@ +//! Hand-rolled PostgreSQL wire-protocol client — plan 16a +//! (`docs/plan/16-postgres-mirror.md`). +//! +//! The mirror's outbound half: protocol v3 over a blocking +//! `std::net::TcpStream`, zero external crates — the same doctrine as the +//! hand-rolled HTTP layer and the CRC32 in `wal.rs`. Scope is exactly what +//! the backup mirror needs: +//! +//! * startup + auth: `trust`, `password` (cleartext), `md5` +//! (SCRAM-SHA-256 is plan 16f) +//! * the **simple query protocol** only (`Query` → `RowDescription` / +//! `DataRow` / `CommandComplete` / `ErrorResponse` / `ReadyForQuery`) — +//! no extended protocol, no prepared statements, no TLS +//! * literal/identifier escaping for SQL the mirror generates +//! +//! Protocol reference: PostgreSQL docs “Frontend/Backend Protocol” and +//! `reference/postgresql/src/include/libpq/` (research symlink). +//! +//! Blocking I/O is deliberate: the only caller is the dedicated `wo-pg` +//! mirror thread (plan 16b) — never a shard worker. + +use std::fmt; +use std::io::{self, Read, Write}; +use std::net::TcpStream; +use std::time::Duration; + +/// Parsed `postgres://user[:password]@host[:port]/database` URL. +/// (No percent-decoding — keep credentials URL-safe.) +#[derive(Debug, Clone)] +pub struct PgConfig { + pub user: String, + pub password: Option, + pub host: String, + pub port: u16, + pub database: String, +} + +impl PgConfig { + pub fn from_url(url: &str) -> Result { + let rest = url.strip_prefix("postgres://") + .or_else(|| url.strip_prefix("postgresql://")) + .ok_or_else(|| format!("WO_PG url must start with postgres:// — got {url}"))?; + let (userinfo, hostpart) = rest.split_once('@') + .ok_or_else(|| "WO_PG url needs user@host".to_string())?; + let (user, password) = match userinfo.split_once(':') { + Some((u, p)) => (u.to_string(), Some(p.to_string())), + None => (userinfo.to_string(), None), + }; + let (hostport, database) = hostpart.split_once('/') + .ok_or_else(|| "WO_PG url needs /database".to_string())?; + let (host, port) = match hostport.split_once(':') { + Some((h, p)) => (h.to_string(), + p.parse::().map_err(|_| format!("bad port `{p}`"))?), + None => (hostport.to_string(), 5432), + }; + if user.is_empty() || host.is_empty() || database.is_empty() { + return Err(format!("incomplete WO_PG url: {url}")); + } + Ok(PgConfig { user, password, host, port, database: database.to_string() }) + } +} + +/// A backend `ErrorResponse` (or client-side failure talking to it). +#[derive(Debug)] +pub struct PgError { + pub severity: String, + pub code: String, + pub message: String, +} + +impl fmt::Display for PgError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{} {}: {}", self.severity, self.code, self.message) + } +} + +impl PgError { + fn client(msg: impl Into) -> PgError { + PgError { severity: "CLIENT".into(), code: "XX000".into(), message: msg.into() } + } +} + +impl From for PgError { + fn from(e: io::Error) -> PgError { PgError::client(format!("io: {e}")) } +} + +/// Result of one simple query (possibly multi-statement). +#[derive(Debug, Default)] +pub struct QueryResult { + pub columns: Vec, + /// Text-format values, `None` = SQL NULL. Rows of the LAST result set. + pub rows: Vec>>, + /// One CommandComplete tag per statement, e.g. `INSERT 0 1`. + pub tags: Vec, +} + +pub struct Conn { + stream: TcpStream, +} + +impl Conn { + /// Connect and authenticate. Blocking, with a connect timeout. + pub fn connect(cfg: &PgConfig) -> Result { + let addr = format!("{}:{}", cfg.host, cfg.port); + let sockaddr = addr.parse() + .map_err(|_| { + // Not a literal ip:port — resolve via ToSocketAddrs. + PgError::client("resolve") + }); + let stream = match sockaddr { + Ok(sa) => TcpStream::connect_timeout(&sa, Duration::from_secs(5))?, + Err(_) => TcpStream::connect(&addr)?, // DNS path + }; + stream.set_nodelay(true).ok(); + stream.set_read_timeout(Some(Duration::from_secs(30)))?; + stream.set_write_timeout(Some(Duration::from_secs(30)))?; + let mut conn = Conn { stream }; + conn.startup(cfg)?; + Ok(conn) + } + + fn startup(&mut self, cfg: &PgConfig) -> Result<(), PgError> { + // StartupMessage: no type byte — i32 len | i32 196608 | k\0v\0 ... \0 + let mut body = Vec::new(); + body.extend_from_slice(&196_608i32.to_be_bytes()); // protocol 3.0 + for (k, v) in [("user", cfg.user.as_str()), + ("database", cfg.database.as_str()), + ("client_encoding", "UTF8"), + ("application_name", "wo-pg-mirror")] { + body.extend_from_slice(k.as_bytes()); body.push(0); + body.extend_from_slice(v.as_bytes()); body.push(0); + } + body.push(0); + let mut msg = Vec::with_capacity(body.len() + 4); + msg.extend_from_slice(&((body.len() as i32 + 4).to_be_bytes())); + msg.extend_from_slice(&body); + self.stream.write_all(&msg)?; + + // Authentication exchange, then drain to ReadyForQuery. + loop { + let (kind, payload) = self.read_message()?; + match kind { + b'R' => { + let auth = be_i32(&payload, 0)?; + match auth { + 0 => {} // AuthenticationOk + 3 => { // CleartextPassword + let pw = cfg.password.clone().ok_or_else(|| + PgError::client("server wants a password; none in WO_PG url"))?; + self.send_password(&pw)?; + } + 5 => { // MD5Password + 4B salt + let pw = cfg.password.clone().ok_or_else(|| + PgError::client("server wants md5 auth; no password in WO_PG url"))?; + let salt = payload.get(4..8).ok_or_else(|| + PgError::client("short md5 salt"))?; + // "md5" + md5hex(md5hex(password + user) + salt) + let inner = md5_hex(format!("{pw}{}", cfg.user).as_bytes()); + let mut outer_in = inner.into_bytes(); + outer_in.extend_from_slice(salt); + let digest = format!("md5{}", md5_hex(&outer_in)); + self.send_password(&digest)?; + } + 10 => return Err(PgError::client( + "server requires SCRAM-SHA-256 — not supported until plan 16f; \ + configure md5/password/trust auth for the mirror role")), + n => return Err(PgError::client(format!("unsupported auth type {n}"))), + } + } + b'S' | b'K' | b'N' => {} // ParameterStatus / BackendKeyData / Notice + b'Z' => return Ok(()), // ReadyForQuery + b'E' => return Err(parse_error(&payload)), + other => return Err(PgError::client(format!( + "unexpected message '{}' during startup", other as char))), + } + } + } + + fn send_password(&mut self, pw: &str) -> Result<(), PgError> { + let mut msg = Vec::with_capacity(pw.len() + 6); + msg.push(b'p'); + msg.extend_from_slice(&((pw.len() as i32 + 5).to_be_bytes())); + msg.extend_from_slice(pw.as_bytes()); + msg.push(0); + self.stream.write_all(&msg)?; + Ok(()) + } + + /// Run one simple query (may contain multiple `;`-separated statements — + /// the backend wraps them in an implicit transaction). Returns the last + /// result set + all command tags; a backend error is returned AFTER the + /// stream is drained to ReadyForQuery, so the connection stays usable. + pub fn simple_query(&mut self, sql: &str) -> Result { + let mut msg = Vec::with_capacity(sql.len() + 6); + msg.push(b'Q'); + msg.extend_from_slice(&((sql.len() as i32 + 5).to_be_bytes())); + msg.extend_from_slice(sql.as_bytes()); + msg.push(0); + self.stream.write_all(&msg)?; + + let mut out = QueryResult::default(); + let mut err: Option = None; + loop { + let (kind, payload) = self.read_message()?; + match kind { + b'T' => { // RowDescription + out.columns.clear(); + let n = be_i16(&payload, 0)? as usize; + let mut off = 2; + for _ in 0..n { + let name = read_cstr(&payload, off)?; + off += name.len() + 1 + 18; // 4+2+4+2+4+2 fixed fields + out.columns.push(name); + } + out.rows.clear(); // keep the last result set + } + b'D' => { // DataRow + let n = be_i16(&payload, 0)? as usize; + let mut off = 2; + let mut row = Vec::with_capacity(n); + for _ in 0..n { + let len = be_i32(&payload, off)?; + off += 4; + if len < 0 { row.push(None); continue; } + let len = len as usize; + let bytes = payload.get(off..off + len) + .ok_or_else(|| PgError::client("short DataRow"))?; + row.push(Some(String::from_utf8_lossy(bytes).into_owned())); + off += len; + } + out.rows.push(row); + } + b'C' => out.tags.push(read_cstr(&payload, 0)?), // CommandComplete + b'E' => { if err.is_none() { err = Some(parse_error(&payload)); } } + b'Z' => break, // ReadyForQuery + b'N' | b'S' | b'I' | b'G' | b'H' | b'W' => {} // notices etc. + other => return Err(PgError::client(format!( + "unexpected message '{}' in query response", other as char))), + } + } + match err { + Some(e) => Err(e), + None => Ok(out), + } + } + + /// Read one backend message: 1-byte type + i32 length (incl. itself). + fn read_message(&mut self) -> Result<(u8, Vec), PgError> { + let mut head = [0u8; 5]; + self.stream.read_exact(&mut head)?; + let len = i32::from_be_bytes([head[1], head[2], head[3], head[4]]); + if !(4..=64 * 1024 * 1024).contains(&len) { + return Err(PgError::client(format!("bad message length {len}"))); + } + let mut payload = vec![0u8; len as usize - 4]; + self.stream.read_exact(&mut payload)?; + Ok((head[0], payload)) + } +} + +// --- wire helpers --- + +fn be_i32(b: &[u8], off: usize) -> Result { + b.get(off..off + 4) + .map(|s| i32::from_be_bytes(s.try_into().unwrap())) + .ok_or_else(|| PgError::client("short message")) +} + +fn be_i16(b: &[u8], off: usize) -> Result { + b.get(off..off + 2) + .map(|s| i16::from_be_bytes(s.try_into().unwrap())) + .ok_or_else(|| PgError::client("short message")) +} + +fn read_cstr(b: &[u8], off: usize) -> Result { + let end = b[off..].iter().position(|&c| c == 0) + .ok_or_else(|| PgError::client("unterminated string"))?; + Ok(String::from_utf8_lossy(&b[off..off + end]).into_owned()) +} + +/// ErrorResponse / NoticeResponse: (field-code byte, cstring) pairs. +fn parse_error(payload: &[u8]) -> PgError { + let mut e = PgError { severity: "ERROR".into(), code: String::new(), message: String::new() }; + let mut off = 0; + while off < payload.len() && payload[off] != 0 { + let code = payload[off]; + let Ok(val) = read_cstr(payload, off + 1) else { break }; + off += 1 + val.len() + 1; + match code { + b'S' => e.severity = val, + b'C' => e.code = val, + b'M' => e.message = val, + _ => {} + } + } + e +} + +// --- SQL text helpers (the mirror builds statements as text) --- + +/// `'…'` literal with single quotes doubled. Standard-conforming strings +/// (the server default) treat backslashes literally, so quotes are the only +/// metacharacter. +pub fn escape_literal(s: &str) -> String { + let mut out = String::with_capacity(s.len() + 2); + out.push('\''); + for c in s.chars() { + if c == '\'' { out.push('\''); } + out.push(c); + } + out.push('\''); + out +} + +/// `"…"` identifier with double quotes doubled. +pub fn escape_ident(s: &str) -> String { + let mut out = String::with_capacity(s.len() + 2); + out.push('"'); + for c in s.chars() { + if c == '"' { out.push('"'); } + out.push(c); + } + out.push('"'); + out +} + +// --- hand-rolled MD5 (RFC 1321) — for the `md5` auth exchange only, the +// --- same no-crates spirit as the CRC32 in wal.rs. Not for new designs. + +pub fn md5_hex(data: &[u8]) -> String { + const S: [u32; 64] = [ + 7, 12, 17, 22, 7, 12, 17, 22, 7, 12, 17, 22, 7, 12, 17, 22, + 5, 9, 14, 20, 5, 9, 14, 20, 5, 9, 14, 20, 5, 9, 14, 20, + 4, 11, 16, 23, 4, 11, 16, 23, 4, 11, 16, 23, 4, 11, 16, 23, + 6, 10, 15, 21, 6, 10, 15, 21, 6, 10, 15, 21, 6, 10, 15, 21, + ]; + const K: [u32; 64] = [ + 0xd76aa478, 0xe8c7b756, 0x242070db, 0xc1bdceee, 0xf57c0faf, 0x4787c62a, + 0xa8304613, 0xfd469501, 0x698098d8, 0x8b44f7af, 0xffff5bb1, 0x895cd7be, + 0x6b901122, 0xfd987193, 0xa679438e, 0x49b40821, 0xf61e2562, 0xc040b340, + 0x265e5a51, 0xe9b6c7aa, 0xd62f105d, 0x02441453, 0xd8a1e681, 0xe7d3fbc8, + 0x21e1cde6, 0xc33707d6, 0xf4d50d87, 0x455a14ed, 0xa9e3e905, 0xfcefa3f8, + 0x676f02d9, 0x8d2a4c8a, 0xfffa3942, 0x8771f681, 0x6d9d6122, 0xfde5380c, + 0xa4beea44, 0x4bdecfa9, 0xf6bb4b60, 0xbebfbc70, 0x289b7ec6, 0xeaa127fa, + 0xd4ef3085, 0x04881d05, 0xd9d4d039, 0xe6db99e5, 0x1fa27cf8, 0xc4ac5665, + 0xf4292244, 0x432aff97, 0xab9423a7, 0xfc93a039, 0x655b59c3, 0x8f0ccc92, + 0xffeff47d, 0x85845dd1, 0x6fa87e4f, 0xfe2ce6e0, 0xa3014314, 0x4e0811a1, + 0xf7537e82, 0xbd3af235, 0x2ad7d2bb, 0xeb86d391, + ]; + + let mut msg = data.to_vec(); + let bit_len = (data.len() as u64).wrapping_mul(8); + msg.push(0x80); + while msg.len() % 64 != 56 { msg.push(0); } + msg.extend_from_slice(&bit_len.to_le_bytes()); + + let (mut a0, mut b0, mut c0, mut d0) = + (0x6745_2301u32, 0xefcd_ab89u32, 0x98ba_dcfeu32, 0x1032_5476u32); + + for chunk in msg.chunks_exact(64) { + let m: Vec = chunk.chunks_exact(4) + .map(|w| u32::from_le_bytes(w.try_into().unwrap())) + .collect(); + let (mut a, mut b, mut c, mut d) = (a0, b0, c0, d0); + for i in 0..64 { + let (f, g) = match i { + 0..=15 => ((b & c) | (!b & d), i), + 16..=31 => ((d & b) | (!d & c), (5 * i + 1) % 16), + 32..=47 => (b ^ c ^ d, (3 * i + 5) % 16), + _ => (c ^ (b | !d), (7 * i) % 16), + }; + let f2 = f.wrapping_add(a).wrapping_add(K[i]).wrapping_add(m[g]); + a = d; d = c; c = b; + b = b.wrapping_add(f2.rotate_left(S[i])); + } + a0 = a0.wrapping_add(a); + b0 = b0.wrapping_add(b); + c0 = c0.wrapping_add(c); + d0 = d0.wrapping_add(d); + } + + let mut out = String::with_capacity(32); + for word in [a0, b0, c0, d0] { + for byte in word.to_le_bytes() { + out.push_str(&format!("{byte:02x}")); + } + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn md5_matches_rfc_vectors() { + assert_eq!(md5_hex(b""), "d41d8cd98f00b204e9800998ecf8427e"); + assert_eq!(md5_hex(b"abc"), "900150983cd24fb0d6963f7d28e17f72"); + assert_eq!(md5_hex(b"message digest"), "f96b697d7cb7938d525a2f31aaf161d0"); + // > one block + assert_eq!( + md5_hex(b"12345678901234567890123456789012345678901234567890123456789012345678901234567890"), + "57edf4a22be3c955ac49da2e2107b67a"); + } + + #[test] + fn url_parse_covers_the_forms() { + let c = PgConfig::from_url("postgres://wo:secret@db.example:6432/prod").unwrap(); + assert_eq!((c.user.as_str(), c.password.as_deref(), c.host.as_str(), c.port, c.database.as_str()), + ("wo", Some("secret"), "db.example", 6432, "prod")); + let c = PgConfig::from_url("postgres://postgres@127.0.0.1/wo").unwrap(); + assert_eq!(c.port, 5432); + assert!(c.password.is_none()); + assert!(PgConfig::from_url("mysql://nope@x/y").is_err()); + assert!(PgConfig::from_url("postgres://user-only-no-host").is_err()); + } + + #[test] + fn escaping_doubles_quotes() { + assert_eq!(escape_literal("it's"), "'it''s'"); + assert_eq!(escape_literal(r#"back\slash"#), r#"'back\slash'"#); + assert_eq!(escape_ident(r#"we"ird"#), r#""we""ird""#); + } + + /// Integration: needs a reachable server — set WO_PG_TEST to run, e.g. + /// WO_PG_TEST=postgres://postgres@127.0.0.1:54329/wo cargo test pg_ + #[test] + fn pg_roundtrip_against_live_server() { + let Ok(url) = std::env::var("WO_PG_TEST") else { + eprintln!("pg_roundtrip: skipped (set WO_PG_TEST=postgres://... to run)"); + return; + }; + let cfg = PgConfig::from_url(&url).unwrap(); + let mut c = Conn::connect(&cfg).unwrap(); + + c.simple_query("DROP TABLE IF EXISTS wo_pg_smoke").unwrap(); + c.simple_query("CREATE TABLE wo_pg_smoke (id BIGINT PRIMARY KEY, row JSONB NOT NULL)").unwrap(); + c.simple_query(&format!( + "INSERT INTO wo_pg_smoke (id, row) VALUES (1, {}::jsonb) \ + ON CONFLICT (id) DO UPDATE SET row = EXCLUDED.row", + escape_literal(r#"{"amount":4999,"note":"it's fine"}"#))).unwrap(); + + let r = c.simple_query("SELECT row->>'amount', row->>'note' FROM wo_pg_smoke").unwrap(); + assert_eq!(r.rows.len(), 1); + assert_eq!(r.rows[0][0].as_deref(), Some("4999")); + assert_eq!(r.rows[0][1].as_deref(), Some("it's fine")); + + // A backend error must leave the connection usable. + assert!(c.simple_query("SELECT * FROM does_not_exist_xyz").is_err()); + let r = c.simple_query("SELECT count(*) FROM wo_pg_smoke").unwrap(); + assert_eq!(r.rows[0][0].as_deref(), Some("1")); + + // Multi-statement query = implicit transaction; both tags come back. + let r = c.simple_query( + "INSERT INTO wo_pg_smoke VALUES (2, '{}'::jsonb); DELETE FROM wo_pg_smoke WHERE id = 2" + ).unwrap(); + assert_eq!(r.tags.len(), 2); + c.simple_query("DROP TABLE wo_pg_smoke").unwrap(); + } +} diff --git a/docs/plan/00-kanban.md b/docs/plan/00-kanban.md index 3d6e4de..0d36b4f 100644 --- a/docs/plan/00-kanban.md +++ b/docs/plan/00-kanban.md @@ -34,13 +34,19 @@ Statuses: ✅ **done** · 🔄 **in progress** · ⬜ **not started** · ⏸ **p All numbers + find-and-fix stories: [09-concurrency-scaleout.md](09-concurrency-scaleout.md) shipped notes and the [benchmark table](../../prototypes/wo-rt-c/README.md). -### Track 3 — Storage & durability (plans 10–12) +### Track 3 — Storage & durability (plans 10–12, 16) — 🔄 in progress | Status | Phase | Doc | Notes | | --- | --- | --- | --- | | ⬜ | 10 storage foundations | [10](10-storage-foundations.md) | scope reduced: WAL framing/fallocate landed via 09c; `@table(name:, index:)` surface + RAM secondary indexes landed via 13 follow-up | | ⬜ | 11 WAL & recovery | [11](11-wal-and-recovery.md) | remaining: snapshots (`.data`), compaction, WAL rotation — replay core shipped in 09c | | ⬜ | 12 engine disk cutover | [12](12-engine-disk-cutover.md) | mmap arena engine (C phase B is the proving ground) | +| ✅ | 16a PG wire client | [16](16-postgres-mirror.md) | hand-rolled protocol v3 (`pg.rs`): trust/password/md5 auth, simple query — zero crates | +| ✅ | 16b PG backup mirror | [16](16-postgres-mirror.md) | `WO_PG=…`: async JSONB upserts behind the WAL ack; boot = full resync; RAM authoritative, reads never touch PG | +| ⬜ | 16c typed columns | [16](16-postgres-mirror.md) | catalog → typed columns; `@table` indexes → `CREATE INDEX` | +| ⬜ | 16d lossless resync | [16](16-postgres-mirror.md) | dirty-flag repair without restart; lag metrics on `/` | +| ⬜ | 16e restore from PG | [16](16-postgres-mirror.md) | `WO_PG_RESTORE=1` boot when the WAL is gone | +| ⬜ | 16f SCRAM auth | [16](16-postgres-mirror.md) | hand-rolled SHA-256/HMAC/PBKDF2 | ### Track 4 — Language & API (`.wo` on the wire) — 🔄 in progress diff --git a/docs/plan/16-postgres-mirror.md b/docs/plan/16-postgres-mirror.md new file mode 100644 index 0000000..f5e6694 --- /dev/null +++ b/docs/plan/16-postgres-mirror.md @@ -0,0 +1,92 @@ +# 16 — PostgreSQL mirror: RAM-authoritative database, Postgres as the backup + +> **Kanban: 🔄 in progress (Track 3 — Storage & durability)** — 16a ✅, 16b ✅ shipped; 16c–16f ⬜. Board: [00-kanban.md](00-kanban.md) + +**Context sources:** [`README.md` § persistent database](../../README.md) (the product goal this implements: *"reads and writes database to RAM, persist data to postgres SQL"*), [`../runtime/database/03-inmemory-engine.md`](../runtime/database/03-inmemory-engine.md) (RAM-resident doctrine: disk sits behind the read path, never in front), [`./09-concurrency-scaleout.md`](./09-concurrency-scaleout.md) (per-shard WAL + ack-after-fsync this rides behind), [`./13-class-model-live-pricing.md`](./13-class-model-live-pricing.md) (the Product/Price worked example; `@table(name: "prices")` names the mirrored table), [`../runtime/database/07-wo-seg-migration.md`](../runtime/database/07-wo-seg-migration.md) (the dual-write precedent), `reference/postgresql/` (research symlink — `src/include/libpq/` for the wire protocol), PostgreSQL docs *Frontend/Backend Protocol*. + +## Context + +The engine is RAM-resident by design and durable through its own per-shard WAL (09c). What's missing is an **external, queryable, operator-friendly copy** of the data — something a DBA can point `psql`, Grafana, or a nightly `pg_dump` at. That's what Postgres is here: **a backup mechanism**, not a storage engine. + +The contract, in one line each: + +- **The entire database lives in RAM.** Reads never touch Postgres — ever. +- **Writes go to RAM (+ WAL) and to Postgres** — but the Postgres write is asynchronous: the client's ack gates on the WAL fsync exactly as before; the mirror follows behind. +- **Postgres is disposable.** RAM is authoritative, so backup repair is always "re-push RAM state" — never a merge. + +Worked example throughout: `Product.set_price(amount)` from the [pricing demo](../examples/pricing/) → a `Price` row in RAM → the same row visible in `psql` as `SELECT * FROM prices`. + +## Design decisions (locked) + +1. **Async mirror, never a commit path.** The `wo-pg` thread is downstream of the commit: shard engines clone committed mutations onto a bounded channel (`try_send` — a full channel drops loudly, it never blocks a worker). Postgres being down costs clients nothing. +2. **Mirror what RAM holds.** The tap emits full post-merge rows (`MirrorRec::Upsert{ty, id, row}` — unlike `WalRec::Update`, which carries only the merge body), so a Postgres row is always byte-equivalent to its RAM row. Method transactions (13b) mirror as one `MirrorRec::Txn` → one `BEGIN…COMMIT` — an aborted method never reaches the channel at all. +3. **Hand-rolled wire client, zero new crates.** `crates/rt/src/pg.rs` speaks protocol v3 (startup, auth `trust`/`password`/`md5` with a hand-rolled MD5 — the CRC32 precedent; SCRAM is 16f) over a blocking `std::net::TcpStream`, **simple query protocol only**. Blocking is fine: the only caller is the dedicated mirror thread. +4. **One `wo-pg` thread per process.** N shard senders → one receiver; per-row ordering is preserved because a row's mutations always come from its owner shard (one FIFO sender). Batches drain opportunistically; statement errors are isolated per record (logged, skipped), socket errors reconnect with capped backoff. +5. **Boot = full resync.** Engines attach the mirror AFTER WAL replay and push their entire state as upserts (`mirror_sync_all`). Consequence: restarting `wo` against a fresh/empty/behind Postgres converges it — verified live (a record dropped during an outage reappeared after restart). +6. **Schema 16b: one JSONB table per type** — `"" (id BIGINT PRIMARY KEY, row JSONB NOT NULL)`, named by `@table(name: "prices")` (default: the type name). Upsert = `INSERT … ON CONFLICT (id) DO UPDATE`. Typed columns are 16c. +7. **Config: `WO_PG=postgres://user[:pass]@host[:port]/db`** env var (the `WO_DATA`/`WO_LISTEN` convention). Unset = mirror off, zero cost. +8. **The WAL stays the recovery source; Postgres is the backup of last resort.** Boot replay reads the WAL as today; restoring FROM Postgres (WAL lost) is 16e's explicit opt-in. + +## Sub-phase sequence + +### `16a` — hand-rolled wire client — ✅ shipped + +`crates/rt/src/pg.rs` (~450 lines, stdlib only): `PgConfig::from_url`, `Conn::connect` (startup + auth trust/cleartext/md5, RFC-1321 MD5 hand-rolled with test vectors), `simple_query` (RowDescription/DataRow/CommandComplete/ErrorResponse/ReadyForQuery; a backend error drains to ready so the connection stays usable), `escape_literal`/`escape_ident`. +**Exit (met):** unit tests (MD5 vectors, URL forms, escaping) green; gated integration test (`WO_PG_TEST=…`) round-trips DDL/upsert/select/error-recovery/multi-statement against `postgres:16`; md5-auth container connects with the right password and fails cleanly with the wrong one. + +### `16b` — the mirror pipeline — ✅ shipped + +`crates/rt/src/mirror.rs`: `MirrorRec{Upsert, Delete, Txn}`, `spawn(cfg, rx, tables)` → the `wo-pg` thread (DDL bootstrap per (re)connect, batched apply, per-record error isolation, capped-backoff reconnect, loud drop accounting). Engine tap (`crates/rt/src/engine.rs`): `attach_mirror`, `mirror_send` (txn-buffered like the WAL buffer — dispatched as one `Txn` on commit, dropped on abort), `mirror_sync_all` boot push; taps sit AFTER `wal_log` accepts in `create`/`update`/`delete`/`commit_txn`. Wiring in `bin/wo.rs` behind `WO_PG`. +**Exit (met, verified live on 2 shards):** `set_price` rows appear in `psql` under the `@table` name `prices` with full JSONB; the abort case (`amount: 0` → 409) mirrors **nothing**; `docker stop` mid-writes → writes keep acking 200, reads unaffected; fresh empty container → reconnect + DDL bootstrap + new writes flow; `wo` restart → WAL replay + bulk sync **heals the gap** (the row written during the outage appeared). 69 unit tests green; blog/ecommerce/hello unchanged without `WO_PG`. + +### `16c` — typed schema projection + +Catalog scalar fields become real columns (`Id`/`Int`/`Money`/`ref` → `BIGINT`, `Text`/`Timestamp`/unions → `TEXT`, `Bool` → `BOOLEAN`; arrays/structs stay in a residual `row JSONB`); `@table(index: [product, at])` → `CREATE INDEX IF NOT EXISTS`; schema evolution via `ADD COLUMN IF NOT EXISTS`. +**Exit:** `SELECT avg((row->>'amount')::bigint)` becomes `SELECT avg(amount) FROM prices WHERE product = 1` in psql, using the mirrored index. + +### `16d` — failure & lossless resync + +Replace drop-and-log with dirty-flag repair: channel overflow or reconnect marks shards dirty; workers re-enqueue their tables at tick (the 09b mail-eventfd mechanism), so convergence no longer waits for a process restart. Mirror lag + queue depth + drop counters surface on `GET /`. +**Exit:** kill Postgres under sustained write load, restart it → row counts converge with zero client errors and no `wo` restart. + +### `16e` — restore from Postgres + +Boot source of last resort when the WAL is gone: `WO_PG_RESTORE=1` makes each shard `SELECT id, row FROM …` its own partition (`(id-1) % n = shard`) before arming accept, seeding RAM and re-logging a fresh WAL. +**Exit:** `rm -rf wo-data` → boot with restore → `current_price` answers from the restored state; id high-water marks keep the interleave. + +### `16f` — SCRAM-SHA-256 auth + +Hand-rolled SHA-256 + HMAC + PBKDF2 (RFC 7677 exchange) so stock `postgres:16` works without `pg_hba` changes. +**Exit:** connects to an out-of-the-box scram-auth server; wrong password fails with the server's error. + +## Verification (16a/16b, reproducible) + +```bash +docker run -d --rm --name wo-pg -e POSTGRES_HOST_AUTH_METHOD=trust -e POSTGRES_DB=wo -p 54329:5432 postgres:16 +WO_PG_TEST=postgres://postgres@127.0.0.1:54329/wo cargo test --lib -- pg_ mirror_ # integration tests +just pricing-pg-demo # scripted end-to-end +``` + +| Check | Result | +| --- | --- | +| `set_price 4999/5999` → `psql: SELECT * FROM prices` | rows present, full JSONB, `@table` name honoured | +| Aborted method (`amount: 0` → 409) | nothing in Postgres — only committed txns mirror | +| Postgres stopped mid-writes | writes ack 200, reads unaffected, mirror retries with backoff | +| Fresh empty database on reconnect | DDL bootstrap recreates tables, stream resumes | +| `wo` restart against behind/empty Postgres | WAL replay + boot sync converges it (outage gap healed) | +| No `WO_PG` | zero behavioural change; 69 unit tests green | + +## Non-scope + +- **No reads from Postgres on any serving path** — doctrine; even 16e's restore happens before accept is armed. +- **No TLS** to Postgres (mirror a local/private endpoint; revisit with 16f). +- **No extended query protocol / prepared statements** — simple protocol is enough for a backup writer; revisit only if 16c profiling demands it. +- **No two-way sync / conflict resolution.** Postgres is write-only from writeonce's perspective (16e restore excepted); external writes to the mirrored tables are unsupported and will be overwritten. +- **No dependency creep.** `crates/rt` gains no crates for this plan — the wire client is part of the same hand-rolled surface as the HTTP layer. + +## Cross-references + +- [`./10-storage-foundations.md`](./10-storage-foundations.md) / [`11`](11-wal-and-recovery.md) / [`12`](12-engine-disk-cutover.md) — the native disk engine; the mirror is orthogonal (external queryable backup vs. native durability) and both sit behind the RAM read path. +- [`./13-class-model-live-pricing.md`](./13-class-model-live-pricing.md) — `@table(name:)` names the mirrored tables; 13b method txns map to Postgres txns; the pricing demo is the acceptance workload. +- [`../runtime/database/07-wo-seg-migration.md`](../runtime/database/07-wo-seg-migration.md) — the dual-write pattern precedent. +- [`./15-mcp-streamable-http.md`](./15-mcp-streamable-http.md) — the other "speak an established protocol, hand-rolled" track; 16a is to Postgres what 15a is to MCP. diff --git a/justfile b/justfile index e033670..2690246 100644 --- a/justfile +++ b/justfile @@ -73,6 +73,35 @@ pricing-demo port="8092": echo "--- live (13c pending, expect 501):"; curl -s -o /dev/null -w '%{http_code}\n' "$base/api/products/live" echo "--- delete 1 (expect 204):"; curl -s -X DELETE "$base/api/products/1" -o /dev/null -w '%{http_code}\n' +# Postgres backup mirror (plan 16): throwaway postgres:16 container, pricing +# demo with WO_PG, verify rows via psql, tear everything down. +# Needs docker + psql. RAM stays authoritative — psql is the backup view. +pricing-pg-demo port="8093" pgport="54331": + #!/usr/bin/env bash + set -euo pipefail + cargo build --bin wo + docker rm -f wo-pg-demo >/dev/null 2>&1 || true + docker run -d --rm --name wo-pg-demo -e POSTGRES_HOST_AUTH_METHOD=trust -e POSTGRES_DB=wo -p {{pgport}}:5432 postgres:16 >/dev/null + trap 'kill $server 2>/dev/null || true; docker rm -f wo-pg-demo >/dev/null 2>&1 || true' EXIT + for _ in $(seq 1 60); do psql -h 127.0.0.1 -p {{pgport}} -U postgres wo -c 'select 1' >/dev/null 2>&1 && break; sleep 0.5; done + WO_THREADS=1 WO_DATA=off WO_PG=postgres://postgres@127.0.0.1:{{pgport}}/wo WO_LISTEN=127.0.0.1:{{port}} \ + ./target/debug/wo run docs/examples/pricing & + server=$! + base=http://127.0.0.1:{{port}} + for _ in $(seq 1 40); do curl -s "$base/healthz" >/dev/null && break; sleep 0.25; done + echo + echo "--- create + set_price 4999, 5999 (RAM ack; mirror follows):" + curl -s -X POST "$base/api/products" -d '{"sku":"WO-001","name":"writeonce mug"}'; echo + curl -s -o /dev/null -X POST "$base/api/products/1/set_price" -d '{"amount":4999}' + curl -s -o /dev/null -X POST "$base/api/products/1/set_price" -d '{"amount":5999}' + echo "--- set_price 0 (aborts; must NOT reach Postgres):" + curl -s -X POST "$base/api/products/1/set_price" -d '{"amount":0}'; echo + sleep 1 + echo "--- psql: the backup view (table name from @table(name: \"prices\")):" + psql -h 127.0.0.1 -p {{pgport}} -U postgres wo -c "SELECT id, row->>'amount' AS amount, row->>'at' AS at FROM prices ORDER BY id" + echo "--- current_price from RAM (reads never touch Postgres):" + curl -s -X POST "$base/api/products/1/current_price"; echo + # main.wo in action: serve hello, run the full CRUD round-trip, shut down hello-demo port="8090": #!/usr/bin/env bash