postgress as mirror database
This commit is contained in:
parent
7ae3a20af1
commit
10cf6e4737
10 changed files with 988 additions and 21 deletions
|
|
@ -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 <Ident> {`. 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).
|
||||
|
|
|
|||
39
README.md
39
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 <dir>` 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 <dir>` 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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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<ExitCode> {
|
|||
}
|
||||
}
|
||||
|
||||
// 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<ExitCode> {
|
|||
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();
|
||||
|
|
|
|||
|
|
@ -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<TxnState>,
|
||||
/// 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<crate::mirror::MirrorSender>,
|
||||
/// 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<crate::wal::WalRec>,
|
||||
undo: Vec<Undo>,
|
||||
wal: Vec<crate::wal::WalRec>,
|
||||
undo: Vec<Undo>,
|
||||
/// Mirror records for this transaction — sent as ONE
|
||||
/// [`crate::mirror::MirrorRec::Txn`] on commit, dropped on abort.
|
||||
mirror: Vec<crate::mirror::MirrorRec>,
|
||||
}
|
||||
|
||||
/// 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)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
270
crates/rt/src/mirror.rs
Normal file
270
crates/rt/src/mirror.rs
Normal file
|
|
@ -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 — `"<storage_name>" (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<MirrorRec>),
|
||||
}
|
||||
|
||||
pub type MirrorSender = SyncSender<MirrorRec>;
|
||||
|
||||
/// 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<MirrorRec>,
|
||||
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<MirrorRec>, 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<MirrorRec>) -> Option<Vec<MirrorRec>> {
|
||||
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<MirrorRec>,
|
||||
dropped: &mut u64,
|
||||
) -> Option<Conn> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
460
crates/rt/src/pg.rs
Normal file
460
crates/rt/src/pg.rs
Normal file
|
|
@ -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<String>,
|
||||
pub host: String,
|
||||
pub port: u16,
|
||||
pub database: String,
|
||||
}
|
||||
|
||||
impl PgConfig {
|
||||
pub fn from_url(url: &str) -> Result<PgConfig, String> {
|
||||
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::<u16>().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<String>) -> PgError {
|
||||
PgError { severity: "CLIENT".into(), code: "XX000".into(), message: msg.into() }
|
||||
}
|
||||
}
|
||||
|
||||
impl From<io::Error> 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<String>,
|
||||
/// Text-format values, `None` = SQL NULL. Rows of the LAST result set.
|
||||
pub rows: Vec<Vec<Option<String>>>,
|
||||
/// One CommandComplete tag per statement, e.g. `INSERT 0 1`.
|
||||
pub tags: Vec<String>,
|
||||
}
|
||||
|
||||
pub struct Conn {
|
||||
stream: TcpStream,
|
||||
}
|
||||
|
||||
impl Conn {
|
||||
/// Connect and authenticate. Blocking, with a connect timeout.
|
||||
pub fn connect(cfg: &PgConfig) -> Result<Conn, PgError> {
|
||||
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<QueryResult, PgError> {
|
||||
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<PgError> = 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<u8>), 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<i32, PgError> {
|
||||
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<i16, PgError> {
|
||||
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<String, PgError> {
|
||||
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<u32> = 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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
92
docs/plan/16-postgres-mirror.md
Normal file
92
docs/plan/16-postgres-mirror.md
Normal file
|
|
@ -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** — `"<storage_name>" (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.
|
||||
29
justfile
29
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
|
||||
|
|
|
|||
Loading…
Reference in a new issue