diff --git a/CLAUDE.md b/CLAUDE.md index 8d503ac..2d2a756 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -78,7 +78,11 @@ Covered in `docs/runtime/database/02-wo-language.md`: Both layers are `.wo` files. The schema layer compiles down to query-layer operations — but only when Phase 5 codegen and Phase 6 full-stack blocks need a single authoritative input. Stage 2 ships with the schema layer only. -### Single-threaded event loop +### Concurrency: thread-per-core event loops (09a shipped) + +Since plans 09a+09b, `wo run` boots `WO_THREADS` pinned worker threads (default = online cores; `wo-shard-` in `ps -T`), each running its own epoll `EventLoop` with its own `SO_REUSEPORT` listener (`crates/rt/src/runtime/scheduler.rs`) **and its own `Engine`** — there is no `Arc>` anywhere. Ids interleave per shard (`Engine::for_shard`; owner = `(id-1) % n`); cross-shard operations travel the shard bus (`crates/rt/src/shard.rs`: mpsc job mailboxes + mail eventfds; creates local, point ops hop once, lists fan out and merge). Deadlock-freedom: jobs never block, waiters pump their own inbox. Routers are thread-local (`HandlerFn` is not `Send`/`Sync`). Signals are blocked before spawn; worker 0 owns the signalfd and broadcasts shutdown via per-worker eventfds. Don't reintroduce shared mutable engine state — the doctrine below is now enforced by ownership. + +### Single-threaded event loop (original doctrine) Covered in `docs/runtime/database/02-wo-language.md § Concurrency Model` and `03-inmemory-engine.md`. The runtime is Redis/TigerBeetle-style: **one userland thread owns everything** — connection accept, parser, engine, subscription registry. The only non-userland thread is the kernel-owned io_uring SQPOLL helper. This is pinned architecturally — group commit still applies (loop drains many commits into one fsync SQE per tick), and scaling past one core is done by **sharding** independent engine processes, not by adding worker threads. Keep this in mind before proposing `Arc>` anything beyond what's already there. diff --git a/crates/rt/src/ast.rs b/crates/rt/src/ast.rs index 545d123..cee2f77 100644 --- a/crates/rt/src/ast.rs +++ b/crates/rt/src/ast.rs @@ -15,6 +15,11 @@ pub struct TypeDecl { pub name: String, pub fields: Vec, pub services: Vec, + /// Declared with `class` instead of `type`. Storage and REST are + /// class-blind (plan 13 decision 5); the flag exists for the 13b method + /// executor and for diagnostics. Stage 13a: `fn` methods inside the body + /// are parsed-and-discarded like triggers. + pub is_class: bool, } #[derive(Debug, Clone)] diff --git a/crates/rt/src/bin/wo.rs b/crates/rt/src/bin/wo.rs index 4c62830..741f247 100644 --- a/crates/rt/src/bin/wo.rs +++ b/crates/rt/src/bin/wo.rs @@ -5,19 +5,14 @@ //! compile a catalog, and serve REST CRUD on :8080. //! wo --help print usage. //! -//! After the phase-04 cutover this binary owns one event loop on one -//! thread. No tokio. The same pattern Redis and TigerBeetle use — see -//! `docs/runtime/database/02-wo-language.md § Concurrency Model`. +//! After the phase-04 cutover this binary owned one event loop on one +//! thread; plan 09a upgrades that to thread-per-core — `WO_THREADS` pinned +//! workers, each with its own event loop and `SO_REUSEPORT` listener +//! (`rt::runtime::scheduler`). Engine state is still globally shared until +//! 09b. No tokio. See `docs/plan/09-concurrency-scaleout.md`. -use std::collections::HashMap; -use std::os::unix::io::{AsRawFd, RawFd}; use std::path::PathBuf; use std::process::ExitCode; -use std::sync::{Arc, Mutex}; -use std::time::Duration; - -use rt::http::{Connection, Listener, Router}; -use rt::runtime::{EventLoop, Interest, SignalFd, Token}; fn usage() { eprintln!( @@ -85,90 +80,102 @@ fn run(dir: PathBuf) -> anyhow::Result { println!("[wo] compiled catalog — {} type{}", catalog.order.len(), if catalog.order.len() == 1 { "" } else { "s" }); - // 4. Boot the engine + the router. - let engine = Arc::new(Mutex::new(rt::engine::Engine::new(catalog.clone()))); + // 4. Print the route banner. println!(); println!("[wo] routes:"); - { - let e = engine.lock().unwrap(); - print!("{}", rt::server::describe_routes(&e)); - } - let router = rt::server::router(engine.clone(), &catalog); + print!("{}", rt::server::describe_routes(&catalog)); - // 5. Bind and serve. + // 5. Serve — thread-per-core with a SHARDED engine (plan 09b): each + // worker owns its own Engine (interleaved ids) and its own router; + // cross-shard operations travel the shard bus (mailbox + eventfd). + // No Arc> anywhere. let addr = std::env::var("WO_LISTEN").unwrap_or_else(|_| "127.0.0.1:8080".to_string()); - let listener = Listener::bind(&addr) - .map_err(|e| anyhow::anyhow!("bind {addr}: {e}"))?; + let n = rt::runtime::scheduler::thread_count(); println!(); - println!("[wo] listening on http://{}", listener.local_addr()); + println!("[wo] listening on http://{addr} — {n} shard{} (thread-per-core, SO_REUSEPORT, sharded engine)", + if n == 1 { "" } else { "s" }); println!("[wo] ctrl-C to stop"); - serve_loop(listener, router)?; - Ok(ExitCode::from(0)) -} - -fn serve_loop(listener: Listener, router: Router) -> anyhow::Result<()> { - let mut eloop = EventLoop::new()?; - let signals = SignalFd::new()?; - let listen_fd = listener.as_raw_fd(); - let signal_fd = signals.as_raw_fd(); - - // Tokens: connection fds carry their own raw fd as the token; the - // listener and signalfd use their fds too — they're disjoint by - // construction (different fds). - eloop.register(listen_fd, Interest::READABLE, Token(listen_fd as u64))?; - eloop.register(signal_fd, Interest::READABLE, Token(signal_fd as u64))?; - - let mut conns: HashMap = HashMap::new(); - - 'outer: loop { - let events = match eloop.wait_once(Some(Duration::from_secs(60))) { - Ok(evs) => evs, - Err(e) => { - eprintln!("[wo] event loop error: {e}"); - continue; - } - }; - - for ev in events { - let fd = ev.token().0 as RawFd; - - if fd == listen_fd { - // Drain accept queue (edge-triggered). - while let Some(cfd) = listener.accept()? { - eloop.register(cfd, Interest::READABLE, Token(cfd as u64))?; - conns.insert(cfd, Connection::new(cfd)); + // Durability (plan 09c): per-shard WAL under WO_DATA (default ./wo-data; + // WO_DATA=off disables). A `meta` file pins the shard count — replaying + // a 4-shard data dir with WO_THREADS=8 would strand logs and break the + // id interleave, so a mismatch refuses to boot (resharding is 09f). + let data_dir = std::env::var("WO_DATA").unwrap_or_else(|_| "./wo-data".to_string()); + let durable = data_dir != "off"; + if durable { + std::fs::create_dir_all(&data_dir)?; + let meta = std::path::Path::new(&data_dir).join("meta"); + match std::fs::read_to_string(&meta) { + Ok(prev) => { + let prev: usize = prev.trim().parse().unwrap_or(0); + if prev != 0 && prev != n { + anyhow::bail!("{data_dir} was written with WO_THREADS={prev} — restart with that, or wipe the dir"); } - continue; - } - - if fd == signal_fd { - let sig = signals.read().unwrap_or(0); - println!(); - println!("[wo] received signal {sig} — shutting down"); - break 'outer; - } - - // Connection event. - let Some(conn) = conns.get_mut(&fd) else { continue }; - let want_writable = match conn.drive(ev.readable, ev.writable, ev.hangup, ev.error, &router) { - Ok(w) => w, - Err(_) => { conns.remove(&fd); continue; } - }; - - if conn.is_done() { - eloop.deregister(fd).ok(); - conns.remove(&fd); // Drop closes the fd. - } else if want_writable { - let _ = eloop.modify(fd, Interest::READ_WRITE, Token(fd as u64)); } + Err(_) => std::fs::write(&meta, format!("{n}\n"))?, } } - // Tear down outstanding connections cleanly. Dropping Connection closes - // each fd; deregistering from the loop is optional (close auto-removes). - for (fd, _) in conns.drain() { - let _ = eloop.deregister(fd); - } - Ok(()) + let bus = rt::shard::ShardBus::new(n)?; + let catalog_for_workers = catalog.clone(); + rt::runtime::scheduler::serve(&addr, move |id| { + use std::os::unix::io::AsRawFd; + let mut engine = rt::engine::Engine::for_shard(catalog_for_workers.clone(), id, n); + let mut has_group_wal = false; + if durable { + let path = std::path::Path::new(&data_dir).join(format!("shard-{id}.rwal")); + let t0 = std::time::Instant::now(); + match rt::wal::Wal::open_and_replay(&path, &mut engine) { + Ok((wal, recs)) => { + if recs > 0 { + println!("[wo] shard {id}: replayed {recs} wal records in {:?}", t0.elapsed()); + } + // Group commit (io_uring): batch the tick's frames into + // one WRITE→FSYNC pair; acks ride the fsync CQE. Falls + // back to per-commit fsync if the ring is unavailable + // (or WO_GROUP_COMMIT=off, kept for A/B measurement). + let group_enabled = std::env::var("WO_GROUP_COMMIT").map(|v| v != "off").unwrap_or(true); + match (if group_enabled { rt::runtime::Uring::new(256) } else { Err(std::io::Error::other("disabled")) }) + .and_then(|ring| rt::wal::WalGroup::new(wal, ring)) + { + Ok(group) => { engine.attach_wal_group(group); has_group_wal = true; } + Err(e) => { + eprintln!("[wo] shard {id}: io_uring unavailable ({e}) — per-commit fsync"); + // wal moved; reopen in per-commit mode + let mut scratch = rt::engine::Engine::for_shard(catalog_for_workers.clone(), id, n); + if let Ok((w2, _)) = rt::wal::Wal::open_and_replay(&path, &mut scratch) { + engine.attach_wal(w2); + } + } + } + } + Err(e) => eprintln!("[wo] shard {id}: WAL unavailable ({e}) — running non-durable"), + } + } + 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(); + let wal_hooks = if has_group_wal { + ctx.engine.borrow().wal_ring_fd().map(|rfd| { + let pump_ctx = ctx.clone(); + let unpark_ctx = ctx.clone(); + let park_ctx = ctx.clone(); + rt::runtime::scheduler::WalHooks { + ring_fd: rfd, + pump: Box::new(move || pump_ctx.wal_pump()), + unparks: Box::new(move || unpark_ctx.take_unparks()), + park_conn: Box::new(move |fd, gen| park_ctx.engine.borrow_mut().park_conn(fd, gen)), + } + }) + } else { None }; + let mail_ctx = ctx.clone(); + rt::runtime::scheduler::Worker { + router, + mail: Some((mail_fd, Box::new(move || mail_ctx.drain_inbox()))), + wal: wal_hooks, + } + })?; + + println!("[wo] all {n} shards joined — bye"); + Ok(ExitCode::from(0)) } diff --git a/crates/rt/src/engine.rs b/crates/rt/src/engine.rs index b9e8756..bfb19f8 100644 --- a/crates/rt/src/engine.rs +++ b/crates/rt/src/engine.rs @@ -21,17 +21,141 @@ pub struct Engine { tables: std::collections::HashMap>, /// per-type id allocator next_id: std::collections::HashMap, + /// id stride — 1 for a standalone engine, `n_shards` for a 09b shard so + /// ids interleave (shard t mints t+1, t+1+n, …) and the owner of any id + /// is recoverable as `(id-1) % n` with zero coordination. + id_step: i64, + /// Per-shard write-ahead log (plan 09c). `PerCommit` fsyncs inside the + /// mutating call (simple, used by tests); `Group` stages frames and + /// parks acks for the worker's per-tick io_uring flush — the C + /// prototype's phase-D group commit. + wal: Option, + /// Set when the last mutation staged a group-commit frame — the caller + /// (handler or shard job) must park its ack. Cleared by `take_staged`. + staged: bool, +} + +#[derive(Debug)] +enum WalBackend { + PerCommit(crate::wal::Wal), + Group(crate::wal::WalGroup), } impl Engine { pub fn new(catalog: Catalog) -> Self { + Self::for_shard(catalog, 0, 1) + } + + /// One shard of a thread-per-core deployment (plan 09b): same engine, + /// interleaved id minting. + pub fn for_shard(catalog: Catalog, shard: usize, n_shards: usize) -> Self { let mut tables = std::collections::HashMap::new(); let mut next_id = std::collections::HashMap::new(); for name in catalog.order.iter() { tables.insert(name.clone(), BTreeMap::new()); - next_id.insert(name.clone(), 1); + next_id.insert(name.clone(), shard as i64 + 1); } - Self { catalog, tables, next_id } + Self { catalog, tables, next_id, id_step: n_shards.max(1) as i64, wal: None, staged: false } + } + + /// Attach a per-commit WAL (fsync inside each mutation). Must happen + /// AFTER replay — replayed mutations must not be re-logged. + pub fn attach_wal(&mut self, wal: crate::wal::Wal) { + self.wal = Some(WalBackend::PerCommit(wal)); + } + + /// Attach a group-commit WAL (io_uring): mutations stage frames; the + /// worker flushes once per tick and releases parked acks on the CQE. + pub fn attach_wal_group(&mut self, wal: crate::wal::WalGroup) { + self.wal = Some(WalBackend::Group(wal)); + } + + /// Did the last mutation stage a group-commit frame? (Cleared on read.) + /// The caller must park its ack on the batch when this is true. + pub fn take_staged(&mut self) -> bool { + std::mem::take(&mut self.staged) + } + + /// Park a cross-shard reply on the active batch — sent on fsync. + pub fn park_reply(&mut self, cb: Box) { + match self.wal.as_mut() { + Some(WalBackend::Group(g)) => g.park(crate::wal::Parked::Reply(cb)), + _ => cb(), // no group WAL: durability already settled (or off) + } + } + + /// Park a local connection's response on the active batch. + pub fn park_conn(&mut self, fd: std::os::unix::io::RawFd, gen: u64) { + if let Some(WalBackend::Group(g)) = self.wal.as_mut() { + g.park(crate::wal::Parked::Conn { fd, gen }); + } + } + + /// Worker hooks — flush at tick end; reap on ring-fd readable. + pub fn wal_flush(&mut self) { + if let Some(WalBackend::Group(g)) = self.wal.as_mut() { + if let Err(e) = g.flush() { eprintln!("[wo] wal flush: {e}"); } + } + } + + pub fn wal_complete(&mut self) -> Option<(bool, Vec)> { + match self.wal.as_mut() { + Some(WalBackend::Group(g)) => g.complete(), + _ => None, + } + } + + pub fn wal_ring_fd(&self) -> Option { + match self.wal.as_ref() { + Some(WalBackend::Group(g)) => Some(g.ring_fd()), + _ => None, + } + } + + /// Apply one replayed WAL record. Bypasses default-seeding and id + /// minting — the log carries exact state — but advances the id + /// high-water mark so post-recovery mints never collide. + pub fn replay(&mut self, rec: &crate::wal::WalRec) { + use crate::wal::WalRec; + match rec { + WalRec::Create { ty, row } => { + let Some(id) = row.get("id").and_then(|v| v.as_i64()) else { return }; + if let Some(table) = self.tables.get_mut(ty) { + table.insert(id, row.clone()); + let step = self.id_step; + let counter = self.next_id.entry(ty.clone()).or_insert(1); + while *counter <= id { *counter += step; } + } + } + WalRec::Update { ty, id, body } => { + if let Some(row) = self.tables.get_mut(ty).and_then(|t| t.get_mut(id)) { + if let Value::Object(input) = body { + for (k, v) in input { + if k != "id" { row.insert(k.clone(), v.clone()); } + } + } + } + } + WalRec::Delete { ty, id } => { + if let Some(table) = self.tables.get_mut(ty) { table.remove(id); } + } + } + } + + /// Make a mutation durable (per-commit) or stage it (group). `Err` + /// means the caller must undo the RAM apply. + fn wal_log(&mut self, rec: crate::wal::WalRec) -> Result<()> { + match self.wal.as_mut() { + Some(WalBackend::PerCommit(w)) => { + w.append(&rec).map_err(|e| anyhow::anyhow!("wal append: {e}"))?; + } + Some(WalBackend::Group(g)) => { + g.stage(&rec).map_err(|e| anyhow::anyhow!("wal stage: {e}"))?; + self.staged = true; + } + None => {} + } + Ok(()) } pub fn catalog(&self) -> &Catalog { &self.catalog } @@ -69,6 +193,11 @@ impl Engine { row.insert("id".into(), json!(id)); self.tables.get_mut(ty).unwrap().insert(id, row.clone()); + // Dual-write order: RAM applied above, durable now, ack after return. + if let Err(e) = self.wal_log(crate::wal::WalRec::Create { ty: ty.into(), row: row.clone() }) { + self.tables.get_mut(ty).unwrap().remove(&id); // never ack non-durable + return Err(e); + } Ok(row) } @@ -77,19 +206,30 @@ impl Engine { let table = self.tables.get_mut(ty) .ok_or_else(|| anyhow::anyhow!("no such type: {ty}"))?; let Some(row) = table.get_mut(&id) else { return Ok(None); }; - if let Value::Object(input) = body { + let prev = row.clone(); + if let Value::Object(input) = &body { for (k, v) in input { if k == "id" { continue; } // don't let the client mutate the primary key - row.insert(k, v); + row.insert(k.clone(), v.clone()); } } - Ok(Some(row.clone())) + let updated = row.clone(); + if let Err(e) = self.wal_log(crate::wal::WalRec::Update { ty: ty.into(), id, body }) { + self.tables.get_mut(ty).unwrap().insert(id, prev); // undo: never ack non-durable + return Err(e); + } + Ok(Some(updated)) } pub fn delete(&mut self, ty: &str, id: i64) -> Result { let table = self.tables.get_mut(ty) .ok_or_else(|| anyhow::anyhow!("no such type: {ty}"))?; - Ok(table.remove(&id).is_some()) + let Some(removed) = table.remove(&id) else { return Ok(false) }; + if let Err(e) = self.wal_log(crate::wal::WalRec::Delete { ty: ty.into(), id }) { + self.tables.get_mut(ty).unwrap().insert(id, removed); // undo + return Err(e); + } + Ok(true) } // --- helpers --- @@ -105,7 +245,7 @@ impl Engine { fn mint_id(&mut self, ty: &str) -> i64 { let counter = self.next_id.entry(ty.to_string()).or_insert(1); let id = *counter; - *counter += 1; + *counter += self.id_step; id } diff --git a/crates/rt/src/http/connection.rs b/crates/rt/src/http/connection.rs index 5abdbf7..246245b 100644 --- a/crates/rt/src/http/connection.rs +++ b/crates/rt/src/http/connection.rs @@ -1,10 +1,13 @@ //! Per-connection state machine driven by the phase-02 [`EventLoop`]. //! -//! Lifecycle (close-after-response, no keep-alive yet): -//! Reading → drain `read(2)` to `EAGAIN`, parse, dispatch through the -//! `Router`, queue the response. +//! Lifecycle (HTTP/1.1 keep-alive — the C prototype's phase-C sequence): +//! Reading → drain `read(2)` to `EAGAIN`, parse one request, dispatch +//! through the `Router`, queue the response. Consumed bytes are trimmed +//! so a pipelined follow-up request carries over. //! Writing → drain `write(2)` to `EAGAIN`. If a write was partial, the //! loop re-arms the fd as `WRITABLE` and we continue on the next event. +//! Once flushed: keep-alive resets to Reading (and immediately serves +//! any buffered pipelined request); `Connection: close` goes to Done. //! Done → loop closes the fd. //! //! Adapted from `reference/crates/wo-http/src/connection.rs`. The owning @@ -24,6 +27,8 @@ use super::route::Router; pub enum ConnState { Reading, Writing, + /// Response built but gated on the WAL batch's fsync (group commit). + Parked, Done, } @@ -33,6 +38,10 @@ pub struct Connection { read_buf: Vec, write_buf: Vec, write_offset: usize, + keep_alive: bool, + /// Incarnation stamp — parked acks are released only when the stamp + /// matches, so a reused fd can never receive another commit's ack. + gen: u64, } impl Connection { @@ -43,6 +52,24 @@ impl Connection { read_buf: Vec::with_capacity(4096), write_buf: Vec::new(), write_offset: 0, + keep_alive: true, + gen: 0, + } + } + + pub fn with_gen(fd: RawFd, gen: u64) -> Self { + let mut c = Self::new(fd); + c.gen = gen; + c + } + + pub fn gen(&self) -> u64 { self.gen } + pub fn is_parked(&self) -> bool { self.state == ConnState::Parked } + + /// The batch fsync landed — the gated response may leave now. + pub fn unpark(&mut self) { + if self.state == ConnState::Parked { + self.state = ConnState::Writing; } } @@ -105,18 +132,19 @@ impl Connection { } fn queue_response(&mut self, response: &Response) { - self.write_buf = response.to_bytes(); + self.write_buf = response.to_bytes(self.keep_alive); self.write_offset = 0; - self.state = ConnState::Writing; + self.state = if response.gate { ConnState::Parked } else { ConnState::Writing }; } /// One step of the state machine, given a readiness event from the - /// loop. Returns `true` if the connection now wants `WRITABLE` (the - /// caller should switch interest from `READABLE`); `false` otherwise. + /// loop. Serves as many buffered requests as it can (keep-alive + + /// pipelining). Returns `true` if the connection now wants `WRITABLE` + /// (the caller should switch interest from `READABLE`). pub fn drive( &mut self, readable: bool, - writable: bool, + _writable: bool, hangup: bool, error: bool, router: &Router, @@ -126,42 +154,52 @@ impl Connection { return Ok(false); } + let mut peer_open = true; if readable && self.state == ConnState::Reading { - let still_open = self.drain_read()?; - match self.try_parse() { - ParseResult::Complete(req) => { - let resp = router.dispatch(&req); - self.queue_response(&resp); - } - ParseResult::Incomplete => { - if !still_open { - self.state = ConnState::Done; + peer_open = self.drain_read()?; + } + + loop { + if self.state == ConnState::Reading { + match self.try_parse() { + ParseResult::Complete(req, consumed) => { + // The response's Connection header — and what we do + // after flushing it — follow the request's wish. + self.keep_alive = req.keep_alive; + self.read_buf.drain(..consumed); + let resp = router.dispatch(&req); + self.queue_response(&resp); + } + ParseResult::Incomplete => { + if !peer_open || hangup { + self.state = ConnState::Done; // peer gone mid-request / idle EOF + } return Ok(false); } - } - ParseResult::Error(msg) => { - let resp = Response::status(super::Status::BAD_REQUEST).text(msg); - self.queue_response(&resp); + ParseResult::Error(msg) => { + self.keep_alive = false; // protocol state is suspect + let resp = Response::status(super::Status::BAD_REQUEST).text(msg); + self.queue_response(&resp); + } } } - } - if self.state == ConnState::Writing { - let flushed = self.drain_write()?; - if flushed { + if self.state == ConnState::Writing { + let flushed = self.drain_write()?; + if !flushed { + return Ok(true); // wait for WRITABLE + } + if self.keep_alive { + self.write_buf.clear(); + self.write_offset = 0; + self.state = ConnState::Reading; + continue; // pipelined request may be buffered + } self.state = ConnState::Done; - return Ok(false); - } else if !writable { - // We tried, EAGAIN'd; tell the caller to wait for WRITABLE. - return Ok(true); } - } - if hangup && self.state != ConnState::Writing { - self.state = ConnState::Done; + return Ok(false); } - - Ok(false) } } @@ -195,34 +233,78 @@ mod tests { } #[test] - fn connection_handles_a_request() { + fn keep_alive_serves_many_requests_on_one_connection() { let (server_fd, client_fd) = socketpair_nonblock(); - let req = b"GET /healthz HTTP/1.1\r\nHost: localhost\r\n\r\n"; - let n = unsafe { - libc::write(client_fd, req.as_ptr() as *const _, req.len()) - }; - assert_eq!(n, req.len() as isize); - let router = Router::new() .route(Method::Get, "/healthz", |_, _| Response::ok().text("ok")); - let mut conn = Connection::new(server_fd); - let want_writable = conn.drive(true, false, false, false, &router).unwrap(); - assert!(!want_writable, "small response fits in one write"); - assert!(conn.is_done()); - let mut buf = [0u8; 4096]; - let n = unsafe { - libc::read(client_fd, buf.as_mut_ptr() as *mut _, buf.len()) - }; - assert!(n > 0); - let s = std::str::from_utf8(&buf[..n as usize]).unwrap(); - assert!(s.starts_with("HTTP/1.1 200 OK\r\n"), "got: {s}"); - assert!(s.ends_with("\r\n\r\nok")); + for i in 0..3 { + let req = b"GET /healthz HTTP/1.1\r\nHost: localhost\r\n\r\n"; + let n = unsafe { libc::write(client_fd, req.as_ptr() as *const _, req.len()) }; + assert_eq!(n, req.len() as isize); + + let want_writable = conn.drive(true, false, false, false, &router).unwrap(); + assert!(!want_writable, "small response fits in one write"); + assert!(!conn.is_done(), "keep-alive must survive request {i}"); + + let mut buf = [0u8; 4096]; + let n = unsafe { libc::read(client_fd, buf.as_mut_ptr() as *mut _, buf.len()) }; + assert!(n > 0); + let s = std::str::from_utf8(&buf[..n as usize]).unwrap(); + assert!(s.starts_with("HTTP/1.1 200 OK\r\n"), "got: {s}"); + assert!(s.contains("Connection: keep-alive\r\n"), "got: {s}"); + assert!(s.ends_with("\r\n\r\nok")); + } + + unsafe { libc::close(client_fd); } + } + + #[test] + fn pipelined_requests_are_served_in_order() { + let (server_fd, client_fd) = socketpair_nonblock(); + + let router = Router::new() + .route(Method::Get, "/healthz", |_, _| Response::ok().text("ok")); + let mut conn = Connection::new(server_fd); + + // Two requests in ONE write — the second must be served from the + // carried-over buffer without another readable event. + let req = b"GET /healthz HTTP/1.1\r\n\r\nGET /healthz HTTP/1.1\r\n\r\n"; + unsafe { libc::write(client_fd, req.as_ptr() as *const _, req.len()) }; + + conn.drive(true, false, false, false, &router).unwrap(); + assert!(!conn.is_done()); + + let mut buf = [0u8; 4096]; + let n = unsafe { libc::read(client_fd, buf.as_mut_ptr() as *mut _, buf.len()) }; + let s = std::str::from_utf8(&buf[..n as usize]).unwrap(); + assert_eq!(s.matches("HTTP/1.1 200 OK").count(), 2, "got: {s}"); + + unsafe { libc::close(client_fd); } + } + + #[test] + fn connection_close_header_is_honored() { + let (server_fd, client_fd) = socketpair_nonblock(); + + let router = Router::new() + .route(Method::Get, "/healthz", |_, _| Response::ok().text("ok")); + let mut conn = Connection::new(server_fd); + + let req = b"GET /healthz HTTP/1.1\r\nConnection: close\r\n\r\n"; + unsafe { libc::write(client_fd, req.as_ptr() as *const _, req.len()) }; + + conn.drive(true, false, false, false, &router).unwrap(); + assert!(conn.is_done(), "Connection: close must end the connection"); + + let mut buf = [0u8; 4096]; + let n = unsafe { libc::read(client_fd, buf.as_mut_ptr() as *mut _, buf.len()) }; + let s = std::str::from_utf8(&buf[..n as usize]).unwrap(); + assert!(s.contains("Connection: close\r\n"), "got: {s}"); unsafe { libc::close(client_fd); } - // server_fd is closed by Connection::drop. } #[test] diff --git a/crates/rt/src/http/listener.rs b/crates/rt/src/http/listener.rs index 7fd4fc8..1ce6359 100644 --- a/crates/rt/src/http/listener.rs +++ b/crates/rt/src/http/listener.rs @@ -18,6 +18,17 @@ impl Listener { /// Bind to `addr` (IPv4 only for now) and start listening with backlog 128. /// Socket is created `SOCK_NONBLOCK | SOCK_CLOEXEC`. pub fn bind(addr: &str) -> io::Result { + Self::bind_inner(addr, false) + } + + /// Like [`bind`](Self::bind), but with `SO_REUSEPORT`: every shard thread + /// binds its own listener to the same port and the kernel load-balances + /// incoming connections across them by 4-tuple hash (plan 09 decision 3). + pub fn bind_reuseport(addr: &str) -> io::Result { + Self::bind_inner(addr, true) + } + + fn bind_inner(addr: &str, reuseport: bool) -> io::Result { let parsed: SocketAddr = addr.parse().map_err(|e| { io::Error::new(io::ErrorKind::InvalidInput, format!("bad addr {addr:?}: {e}")) })?; @@ -49,6 +60,21 @@ impl Listener { return Err(err); } + if reuseport { + let ret = unsafe { + libc::setsockopt( + fd, libc::SOL_SOCKET, libc::SO_REUSEPORT, + &one as *const _ as *const libc::c_void, + std::mem::size_of::() as libc::socklen_t, + ) + }; + if ret < 0 { + let err = io::Error::last_os_error(); + unsafe { libc::close(fd); } + return Err(err); + } + } + let s_addr = u32::from_be_bytes(v4.ip().octets()).to_be(); let sock = libc::sockaddr_in { sin_family: libc::AF_INET as libc::sa_family_t, @@ -147,4 +173,15 @@ mod tests { unsafe { libc::close(cfd); } drop(stream); } + + #[test] + fn reuseport_allows_two_listeners_on_one_port() { + let a = Listener::bind_reuseport("127.0.0.1:0").unwrap(); + let port = a.local_addr().port(); + let b = Listener::bind_reuseport(&format!("127.0.0.1:{port}")) + .expect("second SO_REUSEPORT bind on the same port must succeed"); + assert_eq!(b.local_addr().port(), port); + // Plain bind on the same port must still fail (no REUSEPORT). + assert!(Listener::bind(&format!("127.0.0.1:{port}")).is_err()); + } } diff --git a/crates/rt/src/http/request.rs b/crates/rt/src/http/request.rs index 2d0d17a..1e595f7 100644 --- a/crates/rt/src/http/request.rs +++ b/crates/rt/src/http/request.rs @@ -42,10 +42,15 @@ pub struct Request { pub query: Option, pub headers: HashMap, pub body: Vec, + /// HTTP/1.1 defaults to keep-alive unless `Connection: close`; + /// HTTP/1.0 defaults to close unless `Connection: keep-alive`. + pub keep_alive: bool, } pub enum ParseResult { - Complete(Request), + /// A full request plus the number of bytes it consumed — the caller + /// trims its buffer so a pipelined follow-up request survives. + Complete(Request, usize), Incomplete, Error(String), } @@ -86,6 +91,7 @@ pub fn parse(buf: &[u8]) -> ParseResult { Some(p) => p, None => return ParseResult::Error("missing path".into()), }; + let http10 = parts.next() == Some("HTTP/1.0"); let (path, query) = match raw_path.split_once('?') { Some((p, q)) => (p.to_string(), Some(q.to_string())), None => (raw_path.to_string(), None), @@ -114,7 +120,13 @@ pub fn parse(buf: &[u8]) -> ParseResult { } let body = buf[header_bytes..total].to_vec(); - ParseResult::Complete(Request { method, path, query, headers, body }) + let conn_hdr = headers.get("connection").map(|v| v.to_ascii_lowercase()); + let keep_alive = if http10 { + conn_hdr.as_deref() == Some("keep-alive") + } else { + conn_hdr.as_deref() != Some("close") + }; + ParseResult::Complete(Request { method, path, query, headers, body, keep_alive }, total) } fn find_header_end(buf: &[u8]) -> Option { @@ -128,7 +140,7 @@ mod tests { #[test] fn parses_simple_get() { let raw = b"GET /api/articles HTTP/1.1\r\nHost: localhost\r\n\r\n"; - let ParseResult::Complete(req) = parse(raw) else { panic!("expected Complete") }; + let ParseResult::Complete(req, _) = parse(raw) else { panic!("expected Complete") }; assert_eq!(req.method, Method::Get); assert_eq!(req.path, "/api/articles"); assert!(req.query.is_none()); @@ -138,7 +150,7 @@ mod tests { #[test] fn parses_query_string() { let raw = b"GET /tag/rust?page=2 HTTP/1.1\r\n\r\n"; - let ParseResult::Complete(req) = parse(raw) else { panic!() }; + let ParseResult::Complete(req, _) = parse(raw) else { panic!() }; assert_eq!(req.path, "/tag/rust"); assert_eq!(req.query.as_deref(), Some("page=2")); } @@ -153,7 +165,8 @@ mod tests { raw.extend_from_slice(format!("Content-Length: {}\r\n", body.len()).as_bytes()); raw.extend_from_slice(b"\r\n"); raw.extend_from_slice(body); - let ParseResult::Complete(req) = parse(&raw) else { panic!() }; + let ParseResult::Complete(req, consumed) = parse(&raw) else { panic!() }; + assert_eq!(consumed, raw.len()); assert_eq!(req.method, Method::Post); assert_eq!(req.body, body); } @@ -173,7 +186,7 @@ mod tests { #[test] fn parses_patch_method() { let raw = b"PATCH /api/articles/1 HTTP/1.1\r\nContent-Length: 0\r\n\r\n"; - let ParseResult::Complete(req) = parse(raw) else { panic!() }; + let ParseResult::Complete(req, _) = parse(raw) else { panic!() }; assert_eq!(req.method, Method::Patch); } } diff --git a/crates/rt/src/http/response.rs b/crates/rt/src/http/response.rs index 707cfd3..83833c2 100644 --- a/crates/rt/src/http/response.rs +++ b/crates/rt/src/http/response.rs @@ -27,11 +27,14 @@ pub struct Response { pub status: Status, pub headers: Vec<(String, String)>, pub body: Vec, + /// Group-commit gate: this response acknowledges a staged WAL frame and + /// must not leave until the batch's fsync CQE (the connection parks). + pub gate: bool, } impl Response { pub fn status(s: Status) -> Self { - Self { status: s, headers: Vec::new(), body: Vec::new() } + Self { status: s, headers: Vec::new(), body: Vec::new(), gate: false } } pub fn ok() -> Self { Self::status(Status::OK) } @@ -62,8 +65,8 @@ impl Response { } /// Serialize to the wire format. Auto-injects `Content-Length` and - /// `Connection: close` (no keep-alive in Stage 2). - pub fn to_bytes(&self) -> Vec { + /// Serialize with the connection disposition the state machine decided. + pub fn to_bytes(&self, keep_alive: bool) -> Vec { let mut buf = Vec::with_capacity(256 + self.body.len()); buf.extend_from_slice( format!("HTTP/1.1 {} {}\r\n", self.status.0, self.status.1).as_bytes(), @@ -72,7 +75,8 @@ impl Response { buf.extend_from_slice(format!("{k}: {v}\r\n").as_bytes()); } buf.extend_from_slice(format!("Content-Length: {}\r\n", self.body.len()).as_bytes()); - buf.extend_from_slice(b"Connection: close\r\n"); + buf.extend_from_slice(if keep_alive { b"Connection: keep-alive\r\n".as_slice() } + else { b"Connection: close\r\n".as_slice() }); buf.extend_from_slice(b"\r\n"); buf.extend_from_slice(&self.body); buf @@ -87,7 +91,7 @@ mod tests { #[test] fn ok_text_body() { let r = Response::ok().text("ok"); - let s = String::from_utf8(r.to_bytes()).unwrap(); + let s = String::from_utf8(r.to_bytes(false)).unwrap(); assert!(s.starts_with("HTTP/1.1 200 OK\r\n")); assert!(s.contains("Content-Type: text/plain")); assert!(s.contains("Content-Length: 2\r\n")); @@ -97,7 +101,7 @@ mod tests { #[test] fn json_body() { let r = Response::created().json(&json!({"id": 1, "title": "Hi"})); - let s = String::from_utf8(r.to_bytes()).unwrap(); + let s = String::from_utf8(r.to_bytes(false)).unwrap(); assert!(s.starts_with("HTTP/1.1 201 Created\r\n")); assert!(s.contains("Content-Type: application/json")); assert!(s.contains(r#"{"id":1,"title":"Hi"}"#)); @@ -106,7 +110,7 @@ mod tests { #[test] fn no_content_status() { let r = Response::no_content(); - let s = String::from_utf8(r.to_bytes()).unwrap(); + let s = String::from_utf8(r.to_bytes(false)).unwrap(); assert!(s.starts_with("HTTP/1.1 204 No Content\r\n")); assert!(s.ends_with("\r\n\r\n")); } diff --git a/crates/rt/src/http/route.rs b/crates/rt/src/http/route.rs index c9a4f0e..0799630 100644 --- a/crates/rt/src/http/route.rs +++ b/crates/rt/src/http/route.rs @@ -13,7 +13,10 @@ use std::collections::HashMap; use super::{Method, Request, Response, Status}; -pub type HandlerFn = dyn Fn(&Request, &RouteParams) -> Response + Send + Sync + 'static; +// NOT `Send`/`Sync`: since plan 09b each worker thread builds and owns its +// own Router over its own shard engine (`Rc` captures) — routers +// never cross threads. +pub type HandlerFn = dyn Fn(&Request, &RouteParams) -> Response + 'static; #[derive(Debug, Clone, PartialEq)] enum Segment { @@ -105,7 +108,7 @@ impl Router { pub fn route(mut self, method: Method, pattern: &str, handler: F) -> Self where - F: Fn(&Request, &RouteParams) -> Response + Send + Sync + 'static, + F: Fn(&Request, &RouteParams) -> Response + 'static, { self.routes.push(Route { method, @@ -142,7 +145,7 @@ mod tests { use std::collections::HashMap; fn req(method: Method, path: &str) -> Request { - Request { method, path: path.into(), query: None, headers: HashMap::new(), body: vec![] } + Request { method, path: path.into(), query: None, headers: HashMap::new(), body: vec![], keep_alive: true } } #[test] diff --git a/crates/rt/src/lexer.rs b/crates/rt/src/lexer.rs index ba9d0c1..73ff0ac 100644 --- a/crates/rt/src/lexer.rs +++ b/crates/rt/src/lexer.rs @@ -130,6 +130,7 @@ impl<'a> Lexer<'a> { let name = self.read_ident_chars(); let kind = match name.as_str() { "type" => Kind::KwType, + "class" => Kind::KwClass, "ref" => Kind::KwRef, "multi" => Kind::KwMulti, "via" => Kind::KwVia, diff --git a/crates/rt/src/lib.rs b/crates/rt/src/lib.rs index 0e59499..637d76e 100644 --- a/crates/rt/src/lib.rs +++ b/crates/rt/src/lib.rs @@ -12,7 +12,9 @@ pub mod lexer; pub mod parser; pub mod runtime; pub mod server; +pub mod shard; pub mod token; +pub mod wal; use std::fs; use std::path::{Path, PathBuf}; diff --git a/crates/rt/src/parser.rs b/crates/rt/src/parser.rs index fb6b402..05fc7de 100644 --- a/crates/rt/src/parser.rs +++ b/crates/rt/src/parser.rs @@ -75,7 +75,7 @@ impl Parser { self.skip_newlines(); if self.at_end() { break; } match self.peek() { - Kind::KwType => sch.types.push(self.parse_type()?), + Kind::KwType | Kind::KwClass => sch.types.push(self.parse_type()?), // Skip constructs we don't execute yet. Kind::HashHash(_) | Kind::KwFn @@ -108,6 +108,7 @@ impl Parser { Kind::LParen => { depth += 1; self.advance(); } Kind::RParen => { depth -= 1; self.advance(); } Kind::KwType if depth == 0 && self.pos > start => break, + Kind::KwClass if depth == 0 && self.pos > start => break, Kind::HashHash(_) if depth == 0 && self.pos > start => break, _ => { self.advance(); } } @@ -118,7 +119,14 @@ impl Parser { // --- type declaration --- fn parse_type(&mut self) -> Result { - self.expect(&Kind::KwType, "`type`")?; + // `class` is the behavior-bearing sibling of `type` — identical field + // grammar plus `fn` methods (plan 13a). Storage/REST are class-blind. + let is_class = matches!(self.peek(), Kind::KwClass); + if is_class { + self.advance(); + } else { + self.expect(&Kind::KwType, "`type` or `class`")?; + } let name = self.expect_ident("type name")?; // Link types have a different header: `type Purchase link Customer -> Product { ... }`. @@ -133,7 +141,7 @@ impl Parser { self.expect(&Kind::LBrace, "'{'")?; - let mut decl = TypeDecl { name, fields: Vec::new(), services: Vec::new() }; + let mut decl = TypeDecl { name, fields: Vec::new(), services: Vec::new(), is_class }; loop { self.skip_newlines(); match self.peek() { @@ -141,6 +149,9 @@ impl Parser { Kind::End => bail!("unexpected end of input inside type body"), Kind::KwPolicy => self.skip_block_line()?, // policy ... Kind::KwOn => self.skip_on_block()?, // on update when ... do ... + Kind::KwFn => self.skip_block_line()?, // method — executes from 13b; + // brace depth keeps the body's + // `}` from closing the type Kind::KwService => decl.services.push(self.parse_service()?), Kind::Ident(_) => { if is_link { @@ -217,7 +228,7 @@ impl Parser { while matches!(self.peek(), Kind::Newline) { self.advance(); } if matches!(self.peek(), Kind::RBrace | Kind::KwPolicy | Kind::KwService - | Kind::KwOn | Kind::End + | Kind::KwOn | Kind::KwFn | Kind::End ) { return Ok(()); } if matches!(self.peek(), Kind::Ident(_)) && self.looks_like_field() { return Ok(()); @@ -594,4 +605,70 @@ type Order { status: Pending | Paid | Shipped } assert!(sch.types.iter().any(|t| t.name == "Article" && !t.services.is_empty()), "Article should have a service"); } + + #[test] + fn parses_class_with_methods() { + let src = r#" +class Product { + id: Id + sku: SKU @unique + name: Text + prices: multi Price + + fn current_price() -> Money in txn { + return latest(self.prices).amount; + } + + fn set_price(amount: Money) in txn { + insert Price { product: self.id, amount: amount }; + } + + service rest "/api/products" + expose list, get, create, update, delete, subscribe +} +"#; + let sch = parse(src).unwrap(); + assert_eq!(sch.types.len(), 1); + let t = &sch.types[0]; + assert!(t.is_class); + assert_eq!(t.name, "Product"); + // Methods are parsed-and-discarded in 13a; fields and services survive. + assert_eq!(t.fields.len(), 4); + assert!(matches!(t.fields[3].ty, FieldTy::MultiEdge { ref target, .. } if target == "Price")); + assert_eq!(t.services.len(), 1); + assert_eq!(t.services[0].path, "/api/products"); + assert_eq!(t.services[0].expose.len(), 6); + } + + #[test] + fn class_method_braces_do_not_truncate_body() { + // The method body's `}` and nested `{ ... }` literals must not be + // mistaken for the class's closing brace — fields AFTER the methods + // must still parse, and a following declaration must be seen. + let src = r#" +class Price { + id: Id + + fn discounted(pct: Int) -> Money { + if pct > 0 { + return self.amount * (100 - pct) / 100; + } + return self.amount; + } + + amount: Money + currency: Text = "EUR" +} + +type Audit { id: Id } +"#; + let sch = parse(src).unwrap(); + assert_eq!(sch.types.len(), 2); + let p = &sch.types[0]; + assert!(p.is_class); + assert_eq!(p.fields.len(), 3, "fields after the method must parse"); + assert_eq!(p.fields[1].name, "amount"); + assert!(!sch.types[1].is_class); + assert_eq!(sch.types[1].name, "Audit"); + } } diff --git a/crates/rt/src/runtime/mod.rs b/crates/rt/src/runtime/mod.rs index 497af95..8a1453c 100644 --- a/crates/rt/src/runtime/mod.rs +++ b/crates/rt/src/runtime/mod.rs @@ -13,10 +13,13 @@ mod eventfd; mod netpoll_epoll; +mod netpoll_io_uring; +pub mod scheduler; mod signalfd; mod timerfd; pub use eventfd::EventFd; pub use netpoll_epoll::{Event, EventLoop, Interest, Token}; +pub use netpoll_io_uring::Uring; pub use signalfd::SignalFd; pub use timerfd::TimerFd; diff --git a/crates/rt/src/runtime/netpoll_io_uring.rs b/crates/rt/src/runtime/netpoll_io_uring.rs new file mode 100644 index 0000000..b891017 --- /dev/null +++ b/crates/rt/src/runtime/netpoll_io_uring.rs @@ -0,0 +1,237 @@ +//! Raw io_uring — no liburing, kernel ABI structs defined by hand, exactly +//! the sequence proven in C (`prototypes/wo-rt-c/wo-rt.c` ring_init/enter; +//! card: `docs/plan/exploration/linux/07-io_uring.md`). +//! +//! Scope (this phase): the **storage ring** for per-shard group commit — +//! batched WAL `WRITE` + hard-linked `FSYNC` SQEs, one `io_uring_enter` per +//! flush. The ring fd is pollable, so it registers in the existing epoll +//! loop and completions arrive as just another readable event; the full +//! network port (accept/recv/send SQEs) is a later phase. +//! +//! Requires `IORING_FEAT_SINGLE_MMAP` (kernel ≥ 5.4). + +use std::io; +use std::os::unix::io::RawFd; +use std::sync::atomic::{AtomicU32, Ordering}; + +const SYS_SETUP: libc::c_long = 425; +const SYS_ENTER: libc::c_long = 426; + +const IORING_OFF_SQ_RING: i64 = 0; +const IORING_OFF_SQES: i64 = 0x1000_0000; +const IORING_ENTER_GETEVENTS: u32 = 1; +const IORING_FEAT_SINGLE_MMAP: u32 = 1; + +pub const OP_FSYNC: u8 = 3; +pub const OP_WRITE: u8 = 23; +pub const IOSQE_IO_LINK: u8 = 1 << 2; +pub const FSYNC_DATASYNC: u32 = 1; + +#[repr(C)] +#[derive(Default)] +struct SqOffsets { head: u32, tail: u32, ring_mask: u32, ring_entries: u32, flags: u32, dropped: u32, array: u32, resv1: u32, user_addr: u64 } + +#[repr(C)] +#[derive(Default)] +struct CqOffsets { head: u32, tail: u32, ring_mask: u32, ring_entries: u32, overflow: u32, cqes: u32, flags: u32, resv1: u32, user_addr: u64 } + +#[repr(C)] +#[derive(Default)] +struct Params { + sq_entries: u32, cq_entries: u32, flags: u32, + sq_thread_cpu: u32, sq_thread_idle: u32, + features: u32, wq_fd: u32, resv: [u32; 3], + sq_off: SqOffsets, cq_off: CqOffsets, +} + +/// One submission-queue entry — the 64-byte kernel layout. +#[repr(C)] +#[derive(Clone, Copy, Default)] +struct Sqe { + opcode: u8, flags: u8, ioprio: u16, fd: i32, + off: u64, addr: u64, len: u32, op_flags: u32, + user_data: u64, + buf_index: u16, personality: u16, splice_fd_in: i32, + _pad: [u64; 2], +} + +#[repr(C)] +#[derive(Clone, Copy)] +struct Cqe { user_data: u64, res: i32, flags: u32 } + +pub struct Uring { + fd: RawFd, + // SQ ring pointers (into the shared mmap) + sq_tail: *const AtomicU32, + sq_mask: u32, + sq_array: *mut u32, + // CQ ring pointers + cq_head: *const AtomicU32, + cq_tail: *const AtomicU32, + cq_mask: u32, + cqes: *const Cqe, + sqes: *mut Sqe, + local_tail: u32, + to_submit: u32, +} + +// The ring is owned and driven by exactly one worker thread (plan 09 +// decision 4); raw pointers into its own mmaps don't change that. +unsafe impl Send for Uring {} + +impl Uring { + pub fn new(entries: u32) -> io::Result { + let mut p = Params::default(); + let fd = unsafe { libc::syscall(SYS_SETUP, entries, &mut p as *mut Params) } as RawFd; + if fd < 0 { return Err(io::Error::last_os_error()); } + if p.features & IORING_FEAT_SINGLE_MMAP == 0 { + unsafe { libc::close(fd) }; + return Err(io::Error::other("kernel lacks IORING_FEAT_SINGLE_MMAP (need >= 5.4)")); + } + + let sq_sz = p.sq_off.array as usize + p.sq_entries as usize * 4; + let cq_sz = p.cq_off.cqes as usize + p.cq_entries as usize * std::mem::size_of::(); + let ring_sz = sq_sz.max(cq_sz); + let ring = unsafe { + libc::mmap(std::ptr::null_mut(), ring_sz, libc::PROT_READ | libc::PROT_WRITE, + libc::MAP_SHARED | libc::MAP_POPULATE, fd, IORING_OFF_SQ_RING) + }; + if ring == libc::MAP_FAILED { let e = io::Error::last_os_error(); unsafe { libc::close(fd) }; return Err(e); } + + let sqes_sz = p.sq_entries as usize * std::mem::size_of::(); + let sqes = unsafe { + libc::mmap(std::ptr::null_mut(), sqes_sz, libc::PROT_READ | libc::PROT_WRITE, + libc::MAP_SHARED | libc::MAP_POPULATE, fd, IORING_OFF_SQES) + }; + if sqes == libc::MAP_FAILED { let e = io::Error::last_os_error(); unsafe { libc::close(fd) }; return Err(e); } + + let at = |off: u32| unsafe { (ring as *mut u8).add(off as usize) }; + let sq_mask = unsafe { *(at(p.sq_off.ring_mask) as *const u32) }; + let cq_mask = unsafe { *(at(p.cq_off.ring_mask) as *const u32) }; + let sq_tail = at(p.sq_off.tail) as *const AtomicU32; + let local_tail = unsafe { (*sq_tail).load(Ordering::Relaxed) }; + + Ok(Self { + fd, + sq_tail, + sq_mask, + sq_array: at(p.sq_off.array) as *mut u32, + cq_head: at(p.cq_off.head) as *const AtomicU32, + cq_tail: at(p.cq_off.tail) as *const AtomicU32, + cq_mask, + cqes: at(p.cq_off.cqes) as *const Cqe, + sqes: sqes as *mut Sqe, + local_tail, + to_submit: 0, + }) + } + + /// The ring fd — readable when completions are pending, so it registers + /// straight into the epoll loop. + pub fn as_raw_fd(&self) -> RawFd { self.fd } + + fn sqe(&mut self) -> &mut Sqe { + let idx = self.local_tail & self.sq_mask; + unsafe { *self.sq_array.add(idx as usize) = idx; } + self.local_tail = self.local_tail.wrapping_add(1); + self.to_submit += 1; + let s = unsafe { &mut *self.sqes.add(idx as usize) }; + *s = Sqe::default(); + s + } + + /// Queue a positional write. SAFETY contract: `buf` must stay alive and + /// unmoved until this op's CQE is reaped — the caller double-buffers. + pub fn push_write(&mut self, fd: RawFd, buf: &[u8], offset: u64, link: bool, user_data: u64) { + let s = self.sqe(); + s.opcode = OP_WRITE; + s.fd = fd; + s.addr = buf.as_ptr() as u64; + s.len = buf.len() as u32; + s.off = offset; + s.flags = if link { IOSQE_IO_LINK } else { 0 }; + s.user_data = user_data; + } + + pub fn push_fsync(&mut self, fd: RawFd, user_data: u64) { + let s = self.sqe(); + s.opcode = OP_FSYNC; + s.fd = fd; + s.op_flags = FSYNC_DATASYNC; + s.user_data = user_data; + } + + /// Publish queued SQEs with one syscall. Non-blocking — completions are + /// observed via epoll on the ring fd. + pub fn submit(&mut self) -> io::Result<()> { + if self.to_submit == 0 { return Ok(()); } + unsafe { (*self.sq_tail).store(self.local_tail, Ordering::Release); } + let n = self.to_submit; + self.to_submit = 0; + loop { + let r = unsafe { libc::syscall(SYS_ENTER, self.fd, n, 0u32, IORING_ENTER_GETEVENTS, 0usize, 0usize) }; + if r >= 0 { return Ok(()); } + let e = io::Error::last_os_error(); + if e.raw_os_error() != Some(libc::EINTR) { return Err(e); } + } + } + + /// Reap every pending completion as `(user_data, res)`. + pub fn pop_cqes(&mut self) -> Vec<(u64, i32)> { + let mut out = Vec::new(); + unsafe { + let mut head = (*self.cq_head).load(Ordering::Relaxed); + let tail = (*self.cq_tail).load(Ordering::Acquire); + while head != tail { + let cqe = &*self.cqes.add((head & self.cq_mask) as usize); + out.push((cqe.user_data, cqe.res)); + head = head.wrapping_add(1); + } + (*self.cq_head).store(head, Ordering::Release); + } + out + } +} + +impl Drop for Uring { + fn drop(&mut self) { + unsafe { libc::close(self.fd) }; + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Read; + + #[test] + fn write_then_linked_fsync_round_trips() { + let path = std::env::temp_dir().join(format!("wo-uring-test-{}", std::process::id())); + let _ = std::fs::remove_file(&path); + let file = std::fs::OpenOptions::new().read(true).write(true).create(true).open(&path).unwrap(); + use std::os::unix::io::AsRawFd; + + let mut ring = Uring::new(8).expect("io_uring available"); + let buf = b"hello-from-the-ring".to_vec(); + ring.push_write(file.as_raw_fd(), &buf, 0, true, 1); + ring.push_fsync(file.as_raw_fd(), 2); + ring.submit().unwrap(); + + // Poll the ring fd until both CQEs arrive. + let mut got = Vec::new(); + for _ in 0..200 { + got.extend(ring.pop_cqes()); + if got.len() >= 2 { break; } + std::thread::sleep(std::time::Duration::from_millis(1)); + } + assert_eq!(got.len(), 2, "write + fsync completions"); + assert_eq!(got[0], (1, buf.len() as i32), "write res = full length"); + assert_eq!(got[1].0, 2); + assert!(got[1].1 >= 0, "fsync ok"); + + let mut s = String::new(); + std::fs::File::open(&path).unwrap().read_to_string(&mut s).unwrap(); + assert_eq!(s, "hello-from-the-ring"); + let _ = std::fs::remove_file(&path); + } +} diff --git a/crates/rt/src/runtime/scheduler.rs b/crates/rt/src/runtime/scheduler.rs new file mode 100644 index 0000000..89992bf --- /dev/null +++ b/crates/rt/src/runtime/scheduler.rs @@ -0,0 +1,276 @@ +//! Thread-per-core scheduler — plan 09a (`docs/plan/09-concurrency-scaleout.md`). +//! +//! Go parallel: `src/runtime/proc.go`, drastically simplified — there are no +//! goroutines to schedule. `WO_THREADS` OS threads (default: online cores) +//! spawn at boot; each pins itself to a core with `sched_setaffinity`, binds +//! its own `SO_REUSEPORT` listener on the shared port, and runs its own +//! [`EventLoop`] over its own connections. A connection accepted on thread K +//! is driven and closed on thread K — no migration, no work stealing. +//! +//! Per plan 09a, **engine state stays globally shared** (`Arc>` +//! inside the per-thread `Router`s) — one thing at a time; the sharded engine +//! is 09b. The C proving ground for this exact sequence is +//! `prototypes/wo-rt-c` phase A (see `docs/plan/exploration/c-runtime/`). +//! +//! Shutdown: signals are blocked in `main` before any worker spawns (the +//! mask is inherited), so only worker 0 — which owns the `signalfd` — ever +//! sees SIGINT/SIGTERM. It broadcasts by writing every worker's `eventfd`; +//! each loop wakes, drains its connections, and joins. + +use std::collections::HashMap; +use std::io; +use std::os::unix::io::{AsRawFd, RawFd}; +use std::sync::Arc; +use std::time::Duration; + +use crate::http::{Connection, Listener, Router}; + +use super::{EventFd, EventLoop, Interest, SignalFd, Token}; + +const MAX_THREADS: usize = 64; + +/// Resolve the worker count: `WO_THREADS` env override, else online cores. +pub fn thread_count() -> usize { + std::env::var("WO_THREADS") + .ok() + .and_then(|s| s.parse::().ok()) + .filter(|&n| n >= 1) + .unwrap_or_else(|| { + std::thread::available_parallelism().map(|n| n.get()).unwrap_or(1) + }) + .min(MAX_THREADS) +} + +fn pin_to_core(core: usize) { + let cores = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(1); + unsafe { + let mut set: libc::cpu_set_t = std::mem::zeroed(); + libc::CPU_ZERO(&mut set); + libc::CPU_SET(core % cores, &mut set); + // 0 = the calling thread. Best-effort: a denied affinity (cgroup + // restrictions) must not stop the worker from serving. + libc::sched_setaffinity(0, std::mem::size_of::(), &set); + } +} + +/// Everything a worker needs beyond its listener: the router over its own +/// shard, and (since 09b) an optional auxiliary fd + callback — the shard +/// bus's mail eventfd, drained into the local engine when it fires. +pub struct Worker { + pub router: Router, + pub mail: Option<(RawFd, Box)>, + /// Group-commit hooks (io_uring WAL): the pollable ring fd, a pump + /// (flush staged batch / reap completions), a drain of pending + /// connection unparks `(fd, gen, durable_ok)`, and a parker invoked + /// when a drive leaves a connection gated on the next fsync. + pub wal: Option, +} + +pub struct WalHooks { + pub ring_fd: RawFd, + pub pump: Box, + pub unparks: Box Vec<(RawFd, u64, bool)>>, + pub park_conn: Box, +} + +/// Spawn `thread_count()` pinned workers, each serving `addr` behind +/// `SO_REUSEPORT` with the [`Worker`] built by `worker_fn(id)`. Blocks +/// until a SIGINT/SIGTERM shuts every worker down. +pub fn serve(addr: &str, worker_fn: F) -> anyhow::Result<()> +where + F: Fn(usize) -> Worker + Send + Sync + 'static, +{ + let n = thread_count(); + + // Block SIGINT/SIGTERM NOW — every worker inherits the mask, so the + // signalfd (owned by worker 0) is the only delivery path. + let signals = SignalFd::new()?; + let sig_raw = signals.as_raw_fd(); + + let wake: Arc> = + Arc::new((0..n).map(|_| EventFd::new()).collect::>>()?); + let worker_fn = Arc::new(worker_fn); + + let mut handles = Vec::with_capacity(n); + for t in 0..n { + let wake = Arc::clone(&wake); + let worker_fn = Arc::clone(&worker_fn); + let addr = addr.to_string(); + let sigfd = (t == 0).then_some(sig_raw); + handles.push( + std::thread::Builder::new() + .name(format!("wo-shard-{t}")) + .spawn(move || worker(t, &addr, sigfd, &wake, &*worker_fn))?, + ); + } + + for h in handles { + match h.join() { + Ok(Ok(())) => {} + Ok(Err(e)) => eprintln!("[wo] worker error: {e}"), + Err(_) => eprintln!("[wo] worker panicked"), + } + } + drop(signals); + Ok(()) +} + +fn worker( + id: usize, + addr: &str, + sigfd: Option, + wake: &[EventFd], + worker_fn: &(dyn Fn(usize) -> Worker + Send + Sync), +) -> anyhow::Result<()> { + pin_to_core(id); + + let listener = match Listener::bind_reuseport(addr) { + Ok(l) => l, + Err(e) => { + // Without a listener this worker is useless — take the whole + // process down cleanly rather than serving with a hole. + for w in wake { let _ = w.write(1); } + anyhow::bail!("shard {id}: bind {addr}: {e}"); + } + }; + let Worker { router, mut mail, mut wal } = worker_fn(id); + + let mut eloop = EventLoop::new()?; + let listen_fd = listener.as_raw_fd(); + let wake_fd = wake[id].as_raw_fd(); + + eloop.register(listen_fd, Interest::READABLE, Token(listen_fd as u64))?; + eloop.register(wake_fd, Interest::READABLE, Token(wake_fd as u64))?; + if let Some(sfd) = sigfd { + eloop.register(sfd, Interest::READABLE, Token(sfd as u64))?; + } + let mail_fd = mail.as_ref().map(|(fd, _)| *fd); + if let Some(mfd) = mail_fd { + eloop.register(mfd, Interest::READABLE, Token(mfd as u64))?; + } + let ring_fd = wal.as_ref().map(|w| w.ring_fd); + if let Some(rfd) = ring_fd { + eloop.register(rfd, Interest::READABLE, Token(rfd as u64))?; + } + let mut next_gen: u64 = 1; + + let mut conns: HashMap = HashMap::new(); + + 'outer: loop { + let events = match eloop.wait_once(Some(Duration::from_secs(60))) { + Ok(evs) => evs, + Err(e) => { + eprintln!("[wo] shard {id}: event loop error: {e}"); + continue; + } + }; + + for ev in events { + let fd = ev.token().0 as RawFd; + + if fd == wake_fd { + let _ = wake[id].read(); + break 'outer; // shutdown broadcast + } + + if Some(fd) == ring_fd { + if let Some(w) = wal.as_mut() { (w.pump)(); } + continue; + } + + if Some(fd) == mail_fd { + // Reset the edge, then drain the shard-bus inbox into the + // local engine. (The eventfd counter is read-and-zeroed; + // jobs arriving mid-drain re-arm the edge.) + let mut buf = [0u8; 8]; + unsafe { libc::read(fd, buf.as_mut_ptr() as *mut libc::c_void, 8) }; + if let Some((_, drain)) = mail.as_mut() { drain(); } + continue; + } + + if Some(fd) == sigfd { + // Drain the siginfo and broadcast shutdown to every shard + // (including ourselves — we exit through the wake path). + let mut buf = [0u8; 128]; // sizeof(signalfd_siginfo) + let r = unsafe { + libc::read(fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) + }; + let signo = if r >= 4 { + u32::from_ne_bytes([buf[0], buf[1], buf[2], buf[3]]) + } else { 0 }; + println!(); + println!("[wo] received signal {signo} — broadcasting shutdown to {} shards", wake.len()); + for w in wake { let _ = w.write(1); } + continue; + } + + if fd == listen_fd { + // Drain the accept queue (edge-triggered). + while let Some(cfd) = listener.accept()? { + eloop.register(cfd, Interest::READABLE, Token(cfd as u64))?; + conns.insert(cfd, Connection::with_gen(cfd, next_gen)); + next_gen += 1; + } + continue; + } + + // Connection event — same state machine as the single-threaded + // loop; the only difference is whose loop it runs on. + let Some(conn) = conns.get_mut(&fd) else { continue }; + let want_writable = match conn.drive(ev.readable, ev.writable, ev.hangup, ev.error, &router) { + Ok(w) => w, + Err(_) => { conns.remove(&fd); continue; } + }; + + if conn.is_done() { + eloop.deregister(fd).ok(); + conns.remove(&fd); // Drop closes the fd. + } else if conn.is_parked() { + let g = conn.gen(); + if let Some(w) = wal.as_mut() { (w.park_conn)(fd, g); } + } else if want_writable { + let _ = eloop.modify(fd, Interest::READ_WRITE, Token(fd as u64)); + } + } + + // End of tick: flush the group-commit batch (one WRITE→FSYNC pair, + // one syscall) and release any unparks the pumps produced. Released + // connections may serve pipelined requests that commit again — loop + // until quiescent so nothing sleeps on an unflushed batch. + if let Some(w) = wal.as_mut() { + for _ in 0..64 { + (w.pump)(); + let pending = (w.unparks)(); + if pending.is_empty() { break; } + for (fd, gen, ok) in pending { + let Some(conn) = conns.get_mut(&fd) else { continue }; + if conn.gen() != gen || !conn.is_parked() { continue; } + if !ok { + eloop.deregister(fd).ok(); + conns.remove(&fd); // never ack non-durable + continue; + } + conn.unpark(); + let want_writable = match conn.drive(false, true, false, false, &router) { + Ok(wb) => wb, + Err(_) => { conns.remove(&fd); continue; } + }; + if conn.is_done() { + eloop.deregister(fd).ok(); + conns.remove(&fd); + } else if conn.is_parked() { + let g = conn.gen(); + (w.park_conn)(fd, g); + } else if want_writable { + let _ = eloop.modify(fd, Interest::READ_WRITE, Token(fd as u64)); + } + } + } + } + } + + for (fd, _) in conns.drain() { + let _ = eloop.deregister(fd); + } + Ok(()) +} diff --git a/crates/rt/src/server.rs b/crates/rt/src/server.rs index 885d564..59bc97c 100644 --- a/crates/rt/src/server.rs +++ b/crates/rt/src/server.rs @@ -13,52 +13,55 @@ //! me GET /path/me (501) //! ``` //! -//! Phase-04 cutover: the axum + tokio backend was replaced with the -//! hand-rolled [`http::Router`](crate::http::Router) running on the -//! phase-02 [`EventLoop`](crate::runtime::EventLoop). Handlers are -//! synchronous, the engine sits behind `Arc>`, and -//! the binary owns one event loop on one thread. +//! Since plan 09b the engine is **sharded**: each worker thread owns its own +//! [`Engine`] behind a [`ShardCtx`](crate::shard::ShardCtx) — no mutex, no +//! shared heap. Handlers resolve the owning shard from the row id +//! (`owner = (id-1) % n`, the interleaved-mint rule), run locally when it's +//! ours, ship a job over the shard bus when it isn't. Creates always mint +//! locally; lists fan out to every shard and merge by id. + +use std::rc::Rc; use crate::ast::{Operation, ServiceKind}; use crate::compile::Catalog; -use crate::engine::Engine; +use crate::engine::Row; use crate::http::{Method, Request, Response, RouteParams, Router, Status}; +use crate::shard::ShardCtx; use serde_json::{json, Value}; -use std::sync::{Arc, Mutex}; -pub type Shared = Arc>; - -/// Build the fully-wired [`Router`] for a running engine. -pub fn router(engine: Shared, catalog: &Catalog) -> Router { +/// Build the fully-wired [`Router`] for one shard's worker thread. +pub fn router(ctx: Rc, catalog: &Catalog) -> Router { + let shard = ctx.id; + let n = ctx.n; let mut r = Router::new() - .route(Method::Get, "/", |_, _| root_response()) + .route(Method::Get, "/", move |_, _| { + Response::ok().json(&json!({ + "runtime": "wo", + "stage": 2, + "threads": n, + "shard": shard, + "notes": "REST CRUD for each `service rest` block. /healthz for liveness. LIVE subscribe in Stage 3." + })) + }) .route(Method::Get, "/healthz", |_, _| Response::ok().text("ok")); for name in &catalog.order { let t = catalog.get(name).expect("type present"); for svc in &t.services { if svc.kind != ServiceKind::Rest { continue; } - r = attach_rest(r, engine.clone(), t.name.clone(), svc.path.clone(), &svc.expose); + r = attach_rest(r, ctx.clone(), t.name.clone(), svc.path.clone(), &svc.expose); } } r } -fn root_response() -> Response { - Response::ok().json(&json!({ - "runtime": "wo", - "stage": 2, - "notes": "REST CRUD for each `service rest` block. /healthz for liveness. LIVE subscribe in Stage 3." - })) -} - fn attach_rest( - mut r: Router, - engine: Shared, - ty: String, - path: String, - ops: &[Operation], + mut r: Router, + ctx: Rc, + ty: String, + path: String, + ops: &[Operation], ) -> Router { let id_path = format!("{path}/:id"); @@ -82,24 +85,24 @@ fn attach_rest( for op in ops { match op { Operation::List => { - let eng = engine.clone(); let ty = ty.clone(); - r = r.route(Method::Get, &path, move |req, params| list_h(&eng, &ty, req, params)); + let ctx = ctx.clone(); let ty = ty.clone(); + r = r.route(Method::Get, &path, move |req, params| list_h(&ctx, &ty, req, params)); } Operation::Create => { - let eng = engine.clone(); let ty = ty.clone(); - r = r.route(Method::Post, &path, move |req, params| create_h(&eng, &ty, req, params)); + let ctx = ctx.clone(); let ty = ty.clone(); + r = r.route(Method::Post, &path, move |req, params| create_h(&ctx, &ty, req, params)); } Operation::Get => { - let eng = engine.clone(); let ty = ty.clone(); - r = r.route(Method::Get, &id_path, move |req, params| get_h(&eng, &ty, req, params)); + let ctx = ctx.clone(); let ty = ty.clone(); + r = r.route(Method::Get, &id_path, move |req, params| get_h(&ctx, &ty, req, params)); } Operation::Update => { - let eng = engine.clone(); let ty = ty.clone(); - r = r.route(Method::Patch, &id_path, move |req, params| update_h(&eng, &ty, req, params)); + let ctx = ctx.clone(); let ty = ty.clone(); + r = r.route(Method::Patch, &id_path, move |req, params| update_h(&ctx, &ty, req, params)); } Operation::Delete => { - let eng = engine.clone(); let ty = ty.clone(); - r = r.route(Method::Delete, &id_path, move |req, params| delete_h(&eng, &ty, req, params)); + let ctx = ctx.clone(); let ty = ty.clone(); + r = r.route(Method::Delete, &id_path, move |req, params| delete_h(&ctx, &ty, req, params)); } Operation::Subscribe | Operation::Me | Operation::Custom => {} } @@ -108,41 +111,73 @@ fn attach_rest( } // --- handlers --- +// +// Cross-shard results travel as `Result<_, String>` (anyhow::Error isn't +// guaranteed Send-friendly to reconstruct losslessly; the string is what we +// put in the HTTP body anyway). `run_on` returning `None` means the owning +// shard is gone — only during shutdown — and maps to 503. -fn list_h(engine: &Shared, ty: &str, _req: &Request, _params: &RouteParams) -> Response { - let eng = engine.lock().unwrap(); - match eng.list(ty) { - Ok(rows) => Response::ok().json(&json!(rows)), - Err(e) => Response::status(Status::INTERNAL_SERVER_ERROR).text(e.to_string()), - } +fn shard_gone() -> Response { + Response::status(Status::INTERNAL_SERVER_ERROR).text("owning shard unavailable") } -fn get_h(engine: &Shared, ty: &str, _req: &Request, params: &RouteParams) -> Response { +fn list_h(ctx: &Rc, ty: &str, _req: &Request, _params: &RouteParams) -> Response { + let ty_owned = ty.to_string(); + // Fan out to every shard, merge by id — the cross-shard read per 09b. + let per_shard: Vec, String>> = + ctx.fanout(move |e| e.list(&ty_owned).map_err(|e| e.to_string())); + let mut rows = Vec::new(); + for r in per_shard { + match r { + Ok(mut v) => rows.append(&mut v), + Err(e) => return Response::status(Status::INTERNAL_SERVER_ERROR).text(e), + } + } + rows.sort_by_key(|row| row.get("id").and_then(|v| v.as_i64()).unwrap_or(0)); + Response::ok().json(&json!(rows)) +} + +fn get_h(ctx: &Rc, ty: &str, _req: &Request, params: &RouteParams) -> Response { let id = match parse_id(params) { Ok(id) => id, Err(r) => return r, }; - let eng = engine.lock().unwrap(); - match eng.get(ty, id) { - Ok(Some(row)) => Response::ok().json(&json!(row)), - Ok(None) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), - Err(e) => Response::status(Status::INTERNAL_SERVER_ERROR).text(e.to_string()), + let ty_owned = ty.to_string(); + let res = ctx.run_on(ctx.owner_of(id), move |e| { + e.get(&ty_owned, id).map_err(|e| e.to_string()) + }); + match res { + None => shard_gone(), + Some(Ok(Some(row))) => Response::ok().json(&json!(row)), + Some(Ok(None)) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), + Some(Err(e)) => Response::status(Status::INTERNAL_SERVER_ERROR).text(e), } } -fn create_h(engine: &Shared, ty: &str, req: &Request, _params: &RouteParams) -> Response { +fn create_h(ctx: &Rc, ty: &str, req: &Request, _params: &RouteParams) -> Response { let body = match parse_json_body(req) { Ok(v) => v, Err(r) => return r, }; - let mut eng = engine.lock().unwrap(); - match eng.create(ty, body) { - Ok(row) => Response::status(Status::CREATED).json(&json!(row)), + // Always local: the receiving shard mints from its own interleaved + // stride, so the row it creates is by construction a row it owns. + let res = ctx.engine.borrow_mut().create(ty, body); + match res { + Ok(row) => gate_if_staged(ctx, Response::status(Status::CREATED).json(&json!(row))), Err(e) => Response::status(Status::BAD_REQUEST).text(e.to_string()), } } -fn update_h(engine: &Shared, ty: &str, req: &Request, params: &RouteParams) -> Response { +/// Group commit: a mutation that staged a WAL frame must not be acked until +/// the batch fsync — flag the response so the connection parks it. +fn gate_if_staged(ctx: &Rc, mut resp: Response) -> Response { + if ctx.engine.borrow_mut().take_staged() { + resp.gate = true; + } + resp +} + +fn update_h(ctx: &Rc, ty: &str, req: &Request, params: &RouteParams) -> Response { let id = match parse_id(params) { Ok(id) => id, Err(r) => return r, @@ -151,24 +186,32 @@ fn update_h(engine: &Shared, ty: &str, req: &Request, params: &RouteParams) -> R Ok(v) => v, Err(r) => return r, }; - let mut eng = engine.lock().unwrap(); - match eng.update(ty, id, body) { - Ok(Some(row)) => Response::ok().json(&json!(row)), - Ok(None) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), - Err(e) => Response::status(Status::BAD_REQUEST).text(e.to_string()), + let ty_owned = ty.to_string(); + let res = ctx.run_on(ctx.owner_of(id), move |e| { + e.update(&ty_owned, id, body).map_err(|e| e.to_string()) + }); + match res { + None => shard_gone(), + Some(Ok(Some(row))) => gate_if_staged(ctx, Response::ok().json(&json!(row))), + Some(Ok(None)) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), + Some(Err(e)) => Response::status(Status::BAD_REQUEST).text(e), } } -fn delete_h(engine: &Shared, ty: &str, _req: &Request, params: &RouteParams) -> Response { +fn delete_h(ctx: &Rc, ty: &str, _req: &Request, params: &RouteParams) -> Response { let id = match parse_id(params) { Ok(id) => id, Err(r) => return r, }; - let mut eng = engine.lock().unwrap(); - match eng.delete(ty, id) { - Ok(true) => Response::no_content(), - Ok(false) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), - Err(e) => Response::status(Status::INTERNAL_SERVER_ERROR).text(e.to_string()), + let ty_owned = ty.to_string(); + let res = ctx.run_on(ctx.owner_of(id), move |e| { + e.delete(&ty_owned, id).map_err(|e| e.to_string()) + }); + match res { + None => shard_gone(), + Some(Ok(true)) => gate_if_staged(ctx, Response::no_content()), + Some(Ok(false)) => Response::status(Status::NOT_FOUND).text(format!("no {ty} with id {id}")), + Some(Err(e)) => Response::status(Status::INTERNAL_SERVER_ERROR).text(e), } } @@ -187,12 +230,12 @@ fn parse_json_body(req: &Request) -> Result { } /// Format the endpoint banner the CLI prints on startup. -pub fn describe_routes(engine: &Engine) -> String { +pub fn describe_routes(catalog: &Catalog) -> String { let mut out = String::new(); out.push_str(" GET / runtime info\n"); out.push_str(" GET /healthz liveness\n"); - for name in &engine.catalog().order { - let t = engine.catalog().get(name).unwrap(); + for name in &catalog.order { + let t = catalog.get(name).unwrap(); for svc in &t.services { if svc.kind != ServiceKind::Rest { continue; } for op in &svc.expose { @@ -217,13 +260,18 @@ pub fn describe_routes(engine: &Engine) -> String { mod tests { use super::*; use crate::compile::Catalog; + use crate::engine::Engine; use crate::parser::parse; + use crate::shard::ShardBus; - fn build(src: &str) -> (Shared, Router) { + /// Single-shard context: `run_on` is always local, `fanout` is just us — + /// exactly the WO_THREADS=1 production shape. + fn build(src: &str) -> (Rc, Router) { let cat = Catalog::from_schemas(vec![parse(src).unwrap()]).unwrap(); - let eng = Arc::new(Mutex::new(Engine::new(cat.clone()))); - let r = router(eng.clone(), &cat); - (eng, r) + let bus = ShardBus::new(1).unwrap(); + let ctx = ShardCtx::new(0, 1, Engine::for_shard(cat.clone(), 0, 1), bus); + let r = router(ctx.clone(), &cat); + (ctx, r) } fn req(method: Method, path: &str, body: &[u8]) -> Request { @@ -233,12 +281,13 @@ mod tests { query: None, headers: Default::default(), body: body.to_vec(), + keep_alive: true, } } #[test] fn router_serves_crud_for_a_service_rest_block() { - let (_eng, r) = build(r#" + let (_ctx, r) = build(r#" type Article { id: Id title: Text service rest "/api/articles" expose list, get, create, update, delete } @@ -268,7 +317,7 @@ type Article { id: Id #[test] fn unexposed_method_yields_405() { - let (_eng, r) = build(r#" + let (_ctx, r) = build(r#" type Tag { id: Id label: Text service rest "/api/tags" expose list, get } @@ -279,7 +328,7 @@ type Tag { id: Id #[test] fn live_subscribe_is_501() { - let (_eng, r) = build(r#" + let (_ctx, r) = build(r#" type Article { id: Id title: Text service rest "/api/articles" expose list, subscribe } diff --git a/crates/rt/src/shard.rs b/crates/rt/src/shard.rs new file mode 100644 index 0000000..2edb8be --- /dev/null +++ b/crates/rt/src/shard.rs @@ -0,0 +1,259 @@ +//! The shard bus — plan 09b (`docs/plan/09-concurrency-scaleout.md`). +//! +//! With 09b every worker owns its own [`Engine`] — `Arc>` is +//! gone. A request that lands on shard K (kernel `SO_REUSEPORT` hash) but +//! targets a row owned by shard J ships a **job** — a boxed closure — to J's +//! mailbox, wakes J's event loop through its mail `eventfd`, and waits for +//! the reply. Two rules keep this deadlock-free: +//! +//! 1. **Jobs never block.** A job is a pure local engine operation on the +//! owning thread; it cannot itself wait on another shard. +//! 2. **Waiters keep serving.** While shard K waits for J's reply it pumps +//! its own inbox, so J (or anyone) waiting on K is never starved. +//! +//! Row → owner mapping is the interleaved-id rule (`Engine::for_shard`): +//! shard t mints ids t+1, t+1+n, … so `owner(id) = (id-1) % n` with zero +//! coordination. Creates are always local (the receiving shard mints from +//! its own stride); reads/updates/deletes hop at most once; lists fan out +//! to every shard and merge. The C proving ground for the wake mechanism is +//! `prototypes/wo-rt-c` (eventfd broadcast); the mailbox-per-thread design +//! is plan 09 decision 2 and 09d's one-message-per-thread fan-out shape. + +use std::cell::RefCell; +use std::io; +use std::rc::Rc; +use std::sync::mpsc::{channel, Receiver, RecvTimeoutError, Sender}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use crate::engine::Engine; +use crate::runtime::EventFd; + +/// A unit of work shipped to the owning shard. Runs against that shard's +/// engine on that shard's thread; replies through whatever channel it +/// captured. +pub type Job = Box; + +/// Created once in `main`, shared by every worker: each shard's job sender +/// and mail eventfd. Inboxes are taken (once each) by their owning worker. +pub struct ShardBus { + senders: Vec>, + wakes: Vec, + inboxes: Mutex>>>, +} + +impl ShardBus { + pub fn new(n: usize) -> io::Result> { + let mut senders = Vec::with_capacity(n); + let mut inboxes = Vec::with_capacity(n); + let mut wakes = Vec::with_capacity(n); + for _ in 0..n { + let (tx, rx) = channel(); + senders.push(tx); + inboxes.push(Some(rx)); + wakes.push(EventFd::new()?); + } + Ok(Arc::new(Self { senders, wakes, inboxes: Mutex::new(inboxes) })) + } + + /// The owning worker claims its inbox at startup. Panics on double-take — + /// that would be a wiring bug, not a runtime condition. + pub fn take_inbox(&self, shard: usize) -> Receiver { + self.inboxes.lock().unwrap()[shard].take().expect("inbox already taken") + } + + pub fn mail_fd(&self, shard: usize) -> &EventFd { &self.wakes[shard] } +} + +/// Per-worker handle: this shard's engine plus the bus. Deliberately `!Send` +/// (`Rc`/`RefCell`) — it exists on exactly one thread, which is the point. +pub struct ShardCtx { + pub id: usize, + pub n: usize, + pub engine: Rc>, + inbox: Receiver, + bus: Arc, + /// Connection unparks discovered while pumping inside a handler — the + /// worker loop takes and applies them after the handler returns. + unparks: RefCell>, +} + +impl ShardCtx { + pub fn new(id: usize, n: usize, engine: Engine, bus: Arc) -> Rc { + let inbox = bus.take_inbox(id); + Rc::new(Self { id, n, engine: Rc::new(RefCell::new(engine)), inbox, bus, + unparks: RefCell::new(Vec::new()) }) + } + + /// Flush + reap this shard's group-commit WAL. Reply parks release + /// immediately; connection parks queue for the worker loop. Called at + /// tick end, on ring-fd events, AND from every pump-wait — a waiter + /// that didn't flush its own batch would deadlock with a peer waiting + /// on it (cross-shard mutual commit). + pub fn wal_pump(&self) { + let completed = { + let mut e = self.engine.borrow_mut(); + e.wal_flush(); + e.wal_complete() + }; + if let Some((ok, acks)) = completed { + for p in acks { + match p { + crate::wal::Parked::Reply(cb) => { + // A failed batch drops the callback: the requester's + // channel disconnects → 500, never a false ack. + if ok { cb() } + } + crate::wal::Parked::Conn { fd, gen } => { + self.unparks.borrow_mut().push((fd, gen, ok)); + } + } + } + } + } + + /// Worker loop: take any connection unparks the pumps produced. + pub fn take_unparks(&self) -> Vec<(std::os::unix::io::RawFd, u64, bool)> { + std::mem::take(&mut self.unparks.borrow_mut()) + } + + /// Which shard owns a row id, per the interleaved-mint rule. + pub fn owner_of(&self, row_id: i64) -> usize { + ((row_id - 1).rem_euclid(self.n as i64)) as usize + } + + /// Execute every queued job against the local engine. Called from the + /// event loop on a mail-eventfd event, and from the wait loops below. + pub fn drain_inbox(&self) { + while let Ok(job) = self.inbox.try_recv() { + job(&mut self.engine.borrow_mut()); + } + } + + /// Run `f` against the engine that owns `owner` — locally if that's us, + /// else ship it and wait, pumping our own inbox so peers waiting on us + /// make progress. Returns `None` only if the owner is gone (shutdown). + pub fn run_on(&self, owner: usize, f: F) -> Option + where + R: Send + 'static, + F: FnOnce(&mut Engine) -> R + Send + 'static, + { + if owner == self.id { + return Some(f(&mut self.engine.borrow_mut())); + } + let (tx, rx) = channel(); + // Group commit: if the job staged a WAL frame on the owner, its + // reply parks on the owner's batch and is sent on the fsync CQE — + // so our requester-side response leaves only after durability. + let job: Job = Box::new(move |e| { + let r = f(e); + if e.take_staged() { + e.park_reply(Box::new(move || { let _ = tx.send(r); })); + } else { + let _ = tx.send(r); + } + }); + if self.bus.senders[owner].send(job).is_err() { + return None; + } + let _ = self.bus.wakes[owner].write(1); + loop { + match rx.recv_timeout(Duration::from_micros(100)) { + Ok(r) => return Some(r), + Err(RecvTimeoutError::Timeout) => { self.drain_inbox(); self.wal_pump(); } + Err(RecvTimeoutError::Disconnected) => return None, + } + } + } + + /// Run `f` on every shard (self included) and collect the results. + /// Cross-shard reads — `list` — are the fan-out-and-merge case. + pub fn fanout(&self, f: F) -> Vec + where + R: Send + 'static, + F: Fn(&mut Engine) -> R + Send + Sync + Clone + 'static, + { + let (tx, rx) = channel(); + let mut remote = 0usize; + for (t, sender) in self.bus.senders.iter().enumerate() { + if t == self.id { continue; } + let tx = tx.clone(); + let f = f.clone(); + let job: Job = Box::new(move |e| { let _ = tx.send(f(e)); }); + if sender.send(job).is_ok() { + let _ = self.bus.wakes[t].write(1); + remote += 1; + } + } + let mut out = Vec::with_capacity(remote + 1); + out.push(f(&mut self.engine.borrow_mut())); + while out.len() < remote + 1 { + match rx.recv_timeout(Duration::from_micros(100)) { + Ok(r) => out.push(r), + Err(RecvTimeoutError::Timeout) => { self.drain_inbox(); self.wal_pump(); } + Err(RecvTimeoutError::Disconnected) => break, + } + } + out + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::compile::Catalog; + use crate::parser::parse; + use serde_json::json; + + fn catalog() -> Catalog { + Catalog::from_schemas(vec![parse( + r#"type Note { id: Id + title: Text + service rest "/api/notes" expose list, get, create }"#, + ).unwrap()]).unwrap() + } + + #[test] + fn interleaved_ids_map_back_to_their_shard() { + let cat = catalog(); + let bus = ShardBus::new(3).unwrap(); + let ctx0 = ShardCtx::new(0, 3, Engine::for_shard(cat.clone(), 0, 3), bus.clone()); + let ctx1 = ShardCtx::new(1, 3, Engine::for_shard(cat.clone(), 1, 3), bus.clone()); + + let a = ctx0.engine.borrow_mut().create("Note", json!({"title":"a"})).unwrap(); + let b = ctx0.engine.borrow_mut().create("Note", json!({"title":"b"})).unwrap(); + let c = ctx1.engine.borrow_mut().create("Note", json!({"title":"c"})).unwrap(); + let (a, b, c) = (a["id"].as_i64().unwrap(), b["id"].as_i64().unwrap(), c["id"].as_i64().unwrap()); + assert_eq!((a, b, c), (1, 4, 2)); + assert_eq!(ctx0.owner_of(a), 0); + assert_eq!(ctx0.owner_of(b), 0); + assert_eq!(ctx0.owner_of(c), 1); + } + + #[test] + fn cross_shard_job_round_trips() { + let cat = catalog(); + let bus = ShardBus::new(2).unwrap(); + let ctx1 = ShardCtx::new(1, 2, Engine::for_shard(cat.clone(), 1, 2), bus.clone()); + + // Shard 1's thread: serve jobs (one drain after the send below). + let bus2 = bus.clone(); + let cat2 = cat.clone(); + let t = std::thread::spawn(move || { + // shard 0 lives on this thread + let ctx0 = ShardCtx::new(0, 2, Engine::for_shard(cat2, 0, 2), bus2); + ctx0.engine.borrow_mut().create("Note", json!({"title":"on-zero"})).unwrap(); + // serve until the job arrives + for _ in 0..1000 { + ctx0.drain_inbox(); + std::thread::sleep(Duration::from_micros(200)); + } + }); + + // From shard 1, read the row owned by shard 0 (id 1). + let row = ctx1.run_on(0, |e| e.get("Note", 1).unwrap()).flatten(); + assert_eq!(row.unwrap()["title"], "on-zero"); + drop(ctx1); + t.join().unwrap(); + } +} diff --git a/crates/rt/src/token.rs b/crates/rt/src/token.rs index 83d5c50..1d80092 100644 --- a/crates/rt/src/token.rs +++ b/crates/rt/src/token.rs @@ -12,6 +12,7 @@ pub enum Kind { // keywords (schema + query layer) KwType, + KwClass, KwRef, KwMulti, KwVia, diff --git a/crates/rt/src/wal.rs b/crates/rt/src/wal.rs new file mode 100644 index 0000000..02ebe98 --- /dev/null +++ b/crates/rt/src/wal.rs @@ -0,0 +1,326 @@ +//! Per-shard write-ahead log — plan 09c (`docs/plan/09-concurrency-scaleout.md`) +//! + the durability core of plan 11, ported from the proven C sequence +//! (`prototypes/wo-rt-c` phases D/E, `docs/plan/exploration/c-runtime/00-plan.md`). +//! +//! One `shard-.rwal` per worker. Frame format (identical shape to the C +//! prototype): `u32 len | u32 crc32(payload) | payload | u32 COMMIT` — a +//! record replays whole or not at all; a torn tail fails CRC/trailer +//! validation and is truncated. Payloads are JSON-serialized [`WalRec`]s. +//! +//! Dual-write order (the C crash-under-load test's hard-won lesson): the +//! engine applies to RAM, appends the frame, `fdatasync`s, and only then +//! returns — so the HTTP ack (written after the handler returns, including +//! for cross-shard jobs whose reply follows the owner's engine call) is +//! always behind the fsync. Group commit (one fsync per loop tick, acks +//! parked on the completion) is deliberately deferred to the io_uring port — +//! doing it on the epoll loop would reopen the exact ack-before-fsync race +//! the C phase-F bench caught. One fsync per commit is slower and correct. +//! +//! Boot: `Wal::open_and_replay` walks the log into the shard's engine — +//! parallel across workers, before any accept is armed. No snapshots yet +//! (the WAL grows unbounded; compaction is phase 11 proper). + +use std::fs::{File, OpenOptions}; +use std::io::{self, Read, Seek, SeekFrom, Write}; +use std::os::unix::io::AsRawFd; +use std::path::Path; + +use serde::{Deserialize, Serialize}; +use serde_json::Value; + +use crate::engine::{Engine, Row}; + +const WAL_COMMIT: u32 = 0xC0FF_EE42; +const WAL_PREALLOC: i64 = 4 * 1024 * 1024; +const MAX_FRAME: u32 = 16 * 1024 * 1024; + +/// One logged mutation. `Create` carries the FULL post-default row (id, +/// timestamps included) so replay is byte-exact; `Update` carries the merge +/// body (merge is deterministic in log order). +#[derive(Serialize, Deserialize)] +#[serde(tag = "op", rename_all = "snake_case")] +pub enum WalRec { + Create { ty: String, row: Row }, + Update { ty: String, id: i64, body: Value }, + Delete { ty: String, id: i64 }, +} + +#[derive(Debug)] +pub struct Wal { + file: File, +} + +/// An acknowledgment parked until its batch's fsync CQE (group commit). +/// `Conn` carries the C-proven generation stamp — kernel fds get reused, and +/// releasing by bare fd would ack a NEW connection's commit before ITS batch +/// is durable (the phase-F ABA bug, prevented by construction here). +pub enum Parked { + /// A local connection whose response waits in its write buffer. + Conn { fd: std::os::unix::io::RawFd, gen: u64 }, + /// A cross-shard reply — the owner runs this to release the requester. + Reply(Box), +} + +/// Group-commit WAL (io_uring): mutations STAGE frames + park their acks; +/// once per loop tick the worker flushes the staging buffer as one `WRITE` +/// SQE hard-linked to one `FSYNC` SQE; the fsync completion releases every +/// parked ack in the batch. Double-buffered — while a batch is in flight +/// (its buffer pinned for the kernel), new commits stage into the twin. +impl std::fmt::Debug for WalGroup { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str("WalGroup") } +} + +pub struct WalGroup { + file: File, + ring: crate::runtime::Uring, + offset: u64, + staging: Vec, + parked: Vec, + inflight: Option<(Vec, Vec)>, + /// Batch sequence — encoded into user_data so a write CQE and an fsync + /// CQE can never be attributed to the wrong batch. + seq: u64, + got_write: Option, + got_fsync: Option, +} + +impl WalGroup { + /// Wrap a replayed [`Wal`] (same file, offset at the validated tail). + pub fn new(wal: Wal, ring: crate::runtime::Uring) -> io::Result { + let mut file = wal.file; + let offset = file.seek(SeekFrom::Current(0))?; + let _ = file.seek(SeekFrom::Start(offset)); + Ok(Self { file, ring, offset, staging: Vec::new(), parked: Vec::new(), inflight: None, + seq: 0, got_write: None, got_fsync: None }) + } + + pub fn ring_fd(&self) -> std::os::unix::io::RawFd { self.ring.as_raw_fd() } + + /// Stage one record into the active batch. RAM is already applied; the + /// ack must now be parked (see [`Parked`]) until this batch fsyncs. + pub fn stage(&mut self, rec: &WalRec) -> io::Result<()> { + let payload = serde_json::to_vec(rec).map_err(io::Error::other)?; + self.staging.extend_from_slice(&(payload.len() as u32).to_le_bytes()); + self.staging.extend_from_slice(&crc32(&payload).to_le_bytes()); + self.staging.extend_from_slice(&payload); + self.staging.extend_from_slice(&WAL_COMMIT.to_le_bytes()); + Ok(()) + } + + pub fn park(&mut self, p: Parked) { + self.parked.push(p); + } + + /// End-of-tick: if commits are staged and no batch is in flight, submit + /// the whole batch as WRITE→FSYNC linked SQEs — one syscall. + pub fn flush(&mut self) -> io::Result<()> { + if self.inflight.is_some() || self.staging.is_empty() { return Ok(()); } + let buf = std::mem::take(&mut self.staging); + let acks = std::mem::take(&mut self.parked); + let fd = self.file.as_raw_fd(); + // user_data = (batch_seq << 1) | op-bit — CQEs are matched to THIS + // batch only; a stale completion can never release the wrong acks. + let ud_write = self.seq << 1; + let ud_fsync = (self.seq << 1) | 1; + // SAFETY: `buf` moves into `inflight` and stays pinned until the CQE. + self.ring.push_write(fd, &buf, self.offset, true, ud_write); + self.ring.push_fsync(fd, ud_fsync); + self.ring.submit()?; + self.inflight = Some((buf, acks)); + self.got_write = None; + self.got_fsync = None; + Ok(()) + } + + /// Ring-fd readable: reap completions. Returns `Some((ok, acks))` only + /// when BOTH the write and fsync CQEs of the in-flight batch have + /// arrived — they routinely land in different ticks on real disks, and + /// releasing on the write CQE alone would ack before durability. + /// `ok = false` means short write or failed fsync: the caller must DROP + /// the acks (close the connections), never release them. + pub fn complete(&mut self) -> Option<(bool, Vec)> { + for (ud, res) in self.ring.pop_cqes() { + if ud >> 1 != self.seq { continue; } // not this batch (stale/corrupt) + if ud & 1 == 0 { self.got_write = Some(res); } + else { self.got_fsync = Some(res); } + } + if self.inflight.is_none() { return None; } + let (Some(wr), Some(fr)) = (self.got_write, self.got_fsync) else { return None; }; + let (buf, acks) = self.inflight.take()?; + self.got_write = None; + self.got_fsync = None; + self.seq += 1; + let mut ok = true; + if wr != buf.len() as i32 { + eprintln!("[wo] wal: short write {wr} != {}", buf.len()); + ok = false; + } + if fr < 0 { + eprintln!("[wo] wal: fsync failed ({fr})"); + ok = false; + } + if ok { self.offset += buf.len() as u64; } + Some((ok, acks)) + } +} + +impl Wal { + /// Open (creating if absent) the shard's log, replay every valid frame + /// into `engine`, truncate any torn tail, and return a writer positioned + /// at the validated end. Returns `(wal, replayed_records)`. + pub fn open_and_replay(path: &Path, engine: &mut Engine) -> io::Result<(Self, usize)> { + let mut file = OpenOptions::new().read(true).write(true).create(true).open(path)?; + unsafe { libc::fallocate(file.as_raw_fd(), 0, 0, WAL_PREALLOC) }; + + let mut buf = Vec::new(); + file.read_to_end(&mut buf)?; + + let mut off = 0usize; + let mut recs = 0usize; + loop { + let Some(frame) = read_frame(&buf, off) else { break }; + let (payload, next) = frame; + match serde_json::from_slice::(payload) { + Ok(rec) => engine.replay(&rec), + Err(e) => { + eprintln!("[wo] wal {}: undecodable record at byte {off} ({e}) — truncating", path.display()); + break; + } + } + recs += 1; + off = next; + } + + // Resume appends at the validated tail; drop torn bytes. + file.set_len(off as u64)?; + unsafe { libc::fallocate(file.as_raw_fd(), 0, 0, WAL_PREALLOC.max(off as i64)) }; + file.seek(SeekFrom::Start(off as u64))?; + Ok((Self { file }, recs)) + } + + /// Append one record and make it durable. The caller's mutation is only + /// allowed to stand (and its response to leave) after this returns Ok. + pub fn append(&mut self, rec: &WalRec) -> io::Result<()> { + let payload = serde_json::to_vec(rec).map_err(io::Error::other)?; + let mut frame = Vec::with_capacity(payload.len() + 12); + frame.extend_from_slice(&(payload.len() as u32).to_le_bytes()); + frame.extend_from_slice(&crc32(&payload).to_le_bytes()); + frame.extend_from_slice(&payload); + frame.extend_from_slice(&WAL_COMMIT.to_le_bytes()); + self.file.write_all(&frame)?; + self.file.sync_data()?; // the D in ACID — ack ordering lives here + Ok(()) + } +} + +/// Validate and slice one frame at `off`. `None` = clean end or torn tail. +fn read_frame(buf: &[u8], off: usize) -> Option<(&[u8], usize)> { + let u32_at = |o: usize| -> Option { + buf.get(o..o + 4).map(|b| u32::from_le_bytes(b.try_into().unwrap())) + }; + let len = u32_at(off)?; + if len == 0 || len > MAX_FRAME { return None; } // preallocated zeros / garbage + let len = len as usize; + let crc = u32_at(off + 4)?; + let payload = buf.get(off + 8..off + 8 + len)?; + let trailer = u32_at(off + 8 + len)?; + if crc32(payload) != crc || trailer != WAL_COMMIT { return None; } + Some((payload, off + 8 + len + 4)) +} + +/// Hand-rolled CRC32 (poly 0xEDB88320) — same algorithm as the C prototype; +/// no external crate. +fn crc32(data: &[u8]) -> u32 { + let mut c: u32; + let mut table = [0u32; 256]; + for (i, t) in table.iter_mut().enumerate() { + c = i as u32; + for _ in 0..8 { + c = if c & 1 != 0 { 0xEDB8_8320 ^ (c >> 1) } else { c >> 1 }; + } + *t = c; + } + let mut crc = 0xFFFF_FFFFu32; + for &b in data { + crc = table[((crc ^ b as u32) & 0xFF) as usize] ^ (crc >> 8); + } + crc ^ 0xFFFF_FFFF +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::compile::Catalog; + use crate::parser::parse; + use serde_json::json; + + fn catalog() -> Catalog { + Catalog::from_schemas(vec![parse( + r#"type Note { id: Id + title: Text + service rest "/api/notes" expose list, get, create, update, delete }"#, + ).unwrap()]).unwrap() + } + + fn tmp(name: &str) -> std::path::PathBuf { + let p = std::env::temp_dir().join(format!("wo-wal-test-{}-{name}", std::process::id())); + let _ = std::fs::remove_file(&p); + p + } + + #[test] + fn replay_restores_creates_updates_deletes_and_id_highwater() { + let path = tmp("roundtrip"); + { + let mut e = Engine::for_shard(catalog(), 0, 2); + let (wal, n) = Wal::open_and_replay(&path, &mut e).unwrap(); + assert_eq!(n, 0); + e.attach_wal(wal); + e.create("Note", json!({"title":"a"})).unwrap(); // id 1 + e.create("Note", json!({"title":"b"})).unwrap(); // id 3 + e.update("Note", 1, json!({"title":"a2"})).unwrap(); + e.create("Note", json!({"title":"c"})).unwrap(); // id 5 + e.delete("Note", 3).unwrap(); + } + // Fresh engine, replay from disk — the "first load". + let mut e = Engine::for_shard(catalog(), 0, 2); + let (wal, n) = Wal::open_and_replay(&path, &mut e).unwrap(); + assert_eq!(n, 5); + let rows = e.list("Note").unwrap(); + let ids: Vec = rows.iter().map(|r| r["id"].as_i64().unwrap()).collect(); + assert_eq!(ids, vec![1, 5]); + assert_eq!(rows[0]["title"], "a2"); + // id high-water restored: the next mint must not collide (and must + // keep the shard-0-of-2 stride: odd ids). + e.attach_wal(wal); + let next = e.create("Note", json!({"title":"d"})).unwrap(); + assert_eq!(next["id"], 7); + let _ = std::fs::remove_file(&path); + } + + #[test] + fn torn_tail_is_dropped_whole() { + let path = tmp("torn"); + { + let mut e = Engine::for_shard(catalog(), 0, 1); + let (wal, _) = Wal::open_and_replay(&path, &mut e).unwrap(); + e.attach_wal(wal); + e.create("Note", json!({"title":"keep"})).unwrap(); + e.create("Note", json!({"title":"casualty"})).unwrap(); + } + // Tear the last record mid-payload. The file is fallocate'd, so the + // data tail is the last non-zero byte (the 0xC0FFEE42 trailer), not + // the file length. + let bytes = std::fs::read(&path).unwrap(); + let tail = bytes.iter().rposition(|&b| b != 0).unwrap() as u64 + 1; + let f = OpenOptions::new().write(true).open(&path).unwrap(); + f.set_len(tail - 7).unwrap(); + drop(f); + + let mut e = Engine::for_shard(catalog(), 0, 1); + let (_, n) = Wal::open_and_replay(&path, &mut e).unwrap(); + assert_eq!(n, 1, "torn record must drop whole"); + assert_eq!(e.list("Note").unwrap().len(), 1); + let _ = std::fs::remove_file(&path); + } +} diff --git a/docs/cm.md b/docs/cm.md new file mode 100644 index 0000000..775c31f --- /dev/null +++ b/docs/cm.md @@ -0,0 +1,17 @@ +scaffold sibling crates, multi-app ecommerce, REST + concurrency docs + +- scaffold 14 placeholder crates (app, db, engine, gen, http, logic, + policy, ql, service, sub, txn, ui, value, wal) — empty Cargo.toml + + src/lib.rs to receive code phase-by-phase from `rt` +- restructure docs/examples/ecommerce into multi-app layout: apps/admin + and apps/storefront, with shared/ types/logic/components, per-app + app.wo + wo.toml, and reusable .htmlx components (layout, money, + order-row) +- add reference/rest/{blog,ecommerce}.rest — VS Code/JetBrains HTTP + request files driving the running prototype, including 501/404/405 + expectations for stubbed endpoints +- add docs/plan/09-concurrency-scaleout.md and docs/plan/ui/00-overview.md; + refine docs/plan/assembly/02-writeonce-stance.md +- refresh templates (about, article, header/footer, home, layout, styles) + and add static favicon/logo +- add infra/sync.sh and tighten .gitignore for reference/ symlinks diff --git a/docs/examples/hello/main.wo b/docs/examples/hello/main.wo new file mode 100644 index 0000000..7682c31 --- /dev/null +++ b/docs/examples/hello/main.wo @@ -0,0 +1,71 @@ +-- main.wo — the smallest complete writeonce program. +-- +-- Run it from the repo root: +-- +-- cargo run --bin wo -- run docs/examples/hello (or: just hello) +-- just hello-demo -- scripted CRUD round-trip +-- +-- and exercise it: +-- +-- curl -X POST localhost:8080/api/notes \ +-- -H "Content-Type: application/json" \ +-- -d '{"title":"hello","body":"# First note"}' +-- curl localhost:8080/api/notes +-- curl localhost:8080/api/notes/1 +-- curl -X PATCH localhost:8080/api/notes/1 -d '{"pinned":true}' +-- curl -X DELETE localhost:8080/api/notes/1 +-- +-- One type declaration is the whole app: the fields define the schema, the +-- `service` block generates the REST endpoints, and the runtime serves them +-- from a single binary — no external database, no framework. + +type Note { + id: Id + title: Text + body: Markdown + pinned: Bool = false -- literal defaults populate on create + created_at: Timestamp = now() -- so does now(); computed defaults + -- like words(...) are Stage 4+ + + service rest "/api/notes" + expose list, get, create, update, delete, subscribe + -- list/get/create/update/delete work today (Stage 2). + -- subscribe maps to GET /api/notes/live — a documented 501 stub + -- until Stage 3 lands LIVE subscriptions over WebSocket. +} + +-- A second actor that modifies Note. The closest thing to "another class" in +-- writeonce is another type with a trigger: there is no class/method model — +-- behavior attaches to data. Creating a Revision rewrites its target note's +-- title, inside the same transaction as the insert. Triggers execute from +-- Stage 4; today the block parses and is discarded. +type Revision { + id: Id + note: ref Note -- foreign key into Note + new_title: Text + at: Timestamp = now() + + on create + do update Note{ id == self.note }.title = self.new_title +} + +-- Procedural entry point. Stage 2 parses and discards `main` blocks +-- (see crates/rt/src/parser.rs — skip_top_level_chunk); from Phase 6 on +-- this runs once at startup, before the HTTP server binds. +main { + insert Note { title: "hello", body: "# First note" }; + + -- A snapshot binding would never see later changes (let n = select ...). + -- LIVE is the language's "pointer that observes stores": the handle + -- receives a delta at every commit that touches a matching row. Stage 3. + -- (Spec: docs/runtime/database/02-wo-language.md § Schema-Layer DML) + let live = LIVE select Note{ title }; + + -- The other actor fires: this insert runs Revision's on-create trigger, + -- the trigger updates the note, and the commit pushes one delta. + insert Revision { note: 1, new_title: "Hello World" }; + + for delta in live { + print(delta.kind, delta.row.title); -- update "Hello World" + } +} diff --git a/docs/examples/pricing/README.md b/docs/examples/pricing/README.md new file mode 100644 index 0000000..bcb1bb5 --- /dev/null +++ b/docs/examples/pricing/README.md @@ -0,0 +1,62 @@ +# `pricing` — class model + live pricing demo + +A `Price` class and a `Product` class with methods — products have prices — and a `/pricing` screen that shows **live price data for selected products**. Data stays in RAM; one product is readable by millions of customers at once; a price update pushes live to millions of subscribers; all concurrency is Linux kernel I/O (epoll today, io_uring per the scale-out plan). + +> **Status: 13a shipped.** `class` declarations parse and the demo serves real REST CRUD today — `wo run docs/examples/pricing` (or `just pricing-demo`) compiles 2 classes and serves `/api/products`. Methods (`set_price`, `current_price`) parse but don't execute until 13b; `subscribe` is a 501 stub until 13c; the UI is design-only until 13d/plan 14. The master plan is [`docs/plan/13-class-model-live-pricing.md`](../../plan/13-class-model-live-pricing.md). + +## The class model in one paragraph + +`class` = **state + methods, no inheritance**. Fields, defaults, `ref`/`multi`, `service`, `policy`, `on ` — all exactly as in `type` — plus `fn` methods with an implicit `self` receiver that run as row-scoped transactions (`in txn [snapshot]`, the same machinery as the ecommerce sample's free-standing `fn checkout`). No `extends`, no override, no virtual dispatch: "is-a" is a tagged union, "has-a" is a `ref`/`multi` edge. Go-style encapsulation, not Java-style hierarchies. + +## Layout + +The UI follows **MVC** ([`docs/plan/exploration/ui/08-mvc-structure.md`](../../plan/exploration/ui/08-mvc-structure.md)), with the same screen anatomy as the v1 Angular app (`reference/writeonce-app/src/app/article/`) collapsed into the single binary: the **model** is the class itself, the **view** is plain `.htmlx` with external `.scss`, and the **controller** is a `.wo` file that binds the model into the view and is the only place UI may call class methods. + +``` +pricing/ +├── wo.toml # app manifest +├── main.wo # entry point: insert product, LIVE subscribe, set_price +├── types/ # MODEL — classes: schema + methods +│ ├── price.wo # class Price — amount/currency/at + fn discounted(pct) +│ └── product.wo # class Product — prices: multi Price + fn current_price / +│ # fn set_price + service rest /api/products +└── ui/ + └── pricing/ # one screen = one MVC triplet + ├── pricing.wo # CONTROLLER — route, model: bindings, actions → class methods + ├── pricing.htmlx # VIEW — plain htmlx (Mustache + wo:live/wo:bind), logic-free + └── pricing.scss # VIEW styles — external SCSS, compiled at `wo build` +``` + +## What runs when + +| File / feature | Goes live in | Plan | +| --- | --- | --- | +| `class` parses; `/api/products` CRUD serves | **13a ✅ shipped** | lexer/parser/AST + spec amendments | +| `set_price` / `current_price` over RPC (`POST /api/products/:id/set_price`) | **13b** | method execution, row-scoped txn | +| `subscribe` / `LIVE select` — delta on every commit, WebSocket at `/api/products/live` | **13c** | subscription registry (scoped Stage 3) | +| `/pricing` screen patches price cells in open browsers | **13d** | SSR + `wo:live`/`wo:bind` client runtime | +| Millions of readers + millions of live recipients | **13e** | scale targets + load harness | + +## The scale story (13e) + +How "one product, millions of customers" works — all of it existing design, instantiated for this demo: + +- **RAM-resident.** The engine is the in-memory design of [`03-inmemory-engine.md`](../../runtime/database/03-inmemory-engine.md); the WAL/disk phases (10–12) sit *behind* the read path for durability, never in front of it. +- **Reads scale across cores.** Thread-per-core, shared-nothing shards behind `SO_REUSEPORT` ([`09-concurrency-scaleout.md`](../../plan/09-concurrency-scaleout.md)). A hot product row is owned by one shard but read-replicated to every thread, refreshed by the same per-thread broadcast that feeds subscribers — so `GET /api/products/1` spreads over all cores with zero contention, while writes keep a single owner. +- **Live updates fan out in two stages.** One `set_price` commit → one delta message **per thread** (not per subscriber, per plan 09d) → each thread predicate-matches its local subscription table and batches socket writes on its own ring. Kernel primitives only: edge-triggered epoll today ([`done/02-event-loop-epoll.md`](../../plan/done/02-event-loop-epoll.md)), per-thread io_uring next ([`exploration/linux/07-io_uring.md`](../../plan/exploration/linux/07-io_uring.md)). + +Verification targets (1 M aggregate reads/s of one product on 16 cores, 1 M live subscribers with p99 delta delivery < 250 ms, commit→first-delta p99 < 10 ms) are defined in [plan 13 § 13e](../../plan/13-class-model-live-pricing.md). + +## Try it today + +```bash +just pricing-demo # scripted CRUD round-trip, self-contained +# or: +cargo run --bin wo -- run docs/examples/pricing +# [wo] compiled catalog — 2 types +curl -X POST localhost:8080/api/products \ + -H "Content-Type: application/json" -d '{"sku":"WO-001","name":"writeonce mug"}' +curl localhost:8080/api/products +``` + +Methods and live push are the next sub-phases (13b/13c). For the `type`-based minimal program, see [`../hello/`](../hello/); for the full workspace shape, [`../../examples/ecommerce/`](../ecommerce/). diff --git a/docs/examples/pricing/main.wo b/docs/examples/pricing/main.wo new file mode 100644 index 0000000..9af45d1 --- /dev/null +++ b/docs/examples/pricing/main.wo @@ -0,0 +1,25 @@ +-- pricing — Phase 13 demo entry point. +-- (Design artifact: Stage 2 parses and discards `main` blocks; this executes +-- from Phase 6. The annotated walkthrough mirrors docs/examples/hello/main.wo.) + +main { + insert Product { sku: "WO-001", name: "writeonce mug" }; + + -- LIVE binding (13c): the handle receives a delta at every commit that + -- touches a matching row — the "pointer that observes stores". A plain + -- `let p = select ...` would be a snapshot and never see later changes. + -- Brace rule: bare identifiers are projections, so this subscribes to + -- the { name, prices } shape of all products. + -- (Spec: docs/runtime/database/02-wo-language.md § Schema-Layer DML) + let live = LIVE select Product{ name, prices }; + + -- Method call (13b): row-scoped transaction — inserts a Price, commits, + -- and the commit fans the delta out to `live` and to every browser on + -- the /pricing screen (13d). At scale: one commit → one message per + -- core → batched socket writes to millions of subscribers (13e). + Product{ sku == "WO-001" }.set_price(4999); + + for delta in live { + print(delta.kind, delta.row.name); -- update "writeonce mug" + } +} diff --git a/docs/examples/pricing/types/price.wo b/docs/examples/pricing/types/price.wo new file mode 100644 index 0000000..2139540 --- /dev/null +++ b/docs/examples/pricing/types/price.wo @@ -0,0 +1,20 @@ +-- Price — a `class`: state + methods, no inheritance. +-- (Phase 13 design artifact — see docs/plan/13-class-model-live-pricing.md. +-- The Stage 2 parser skips `class` blocks; this parses as real syntax from 13a.) +-- +-- Prices are append-only history rows: setting a new price inserts a Price, +-- it never mutates an old one. A product's "current price" is the latest row. + +class Price { + id: Id + product: ref Product -- owning product (foreign key) + amount: Money -- minor units (cents); stdlib scalar + currency: Text = "EUR" + at: Timestamp = now() + + -- Method: implicit `self` receiver, like Go methods — no class hierarchy, + -- just behavior attached to a row. Pure computation, so no `in txn`. + fn discounted(pct: Int) -> Money { + return self.amount * (100 - pct) / 100; + } +} diff --git a/docs/examples/pricing/types/product.wo b/docs/examples/pricing/types/product.wo new file mode 100644 index 0000000..a686a3d --- /dev/null +++ b/docs/examples/pricing/types/product.wo @@ -0,0 +1,31 @@ +-- Product — a `class` with state, methods, and a REST surface. +-- (Phase 13 design artifact — see docs/plan/13-class-model-live-pricing.md.) +-- +-- "Products have prices": composition via `multi`, not inheritance. The +-- price history is its own class (types/price.wo); the product owns the +-- collection edge. + +class Product { + id: Id + sku: SKU @unique + name: Text + prices: multi Price -- append-only price history + + -- Row-scoped transactional method (13b): runs inside a snapshot of the + -- receiving row, exposed as POST /api/products/:id/current_price. + fn current_price() -> Money in txn { + return latest(self.prices).amount; + } + + -- The write path of the live-pricing demo. One call inserts a Price in + -- the same transaction; the commit pushes one delta to every subscriber + -- of this product (13c) and patches every open pricing screen (13d). + fn set_price(amount: Money) in txn { + insert Price { product: self.id, amount: amount }; + } + + -- CRUD works from 13a (classes are storage-identical to types); + -- subscribe goes live in 13c. + service rest "/api/products" + expose list, get, create, update, delete, subscribe +} diff --git a/docs/examples/pricing/ui/pricing/pricing.htmlx b/docs/examples/pricing/ui/pricing/pricing.htmlx new file mode 100644 index 0000000..e64629e --- /dev/null +++ b/docs/examples/pricing/ui/pricing/pricing.htmlx @@ -0,0 +1,37 @@ +{{!-- + pricing.htmlx — the VIEW of the /pricing screen (MVC). Plain htmlx: + Mustache + /wo:bind per docs/plan/exploration/ui/01-htmlx-format-spec.md. + Logic-free — `products` and `watchlist` come from the controller's model: + block (pricing.wo); actions are raised by name, the controller dispatches. + + Live behavior (13c/13d): when set_price commits, only the wo:bind cells of + the affected row patch in place — no reload, no row re-render. +--}} +
+

Live pricing

+ + + + + + + + {{#each products as p}} + + + + + + + + {{/each}} + +
ProductSKUPriceUpdated
+ {{#if (in watchlist p.id)}} + + {{else}} + + {{/if}} + {{p.name}}{{p.sku}}{{> money amount=p.current_price}}{{relative p.prices[-1].at}}
+
+
diff --git a/docs/examples/pricing/ui/pricing/pricing.scss b/docs/examples/pricing/ui/pricing/pricing.scss new file mode 100644 index 0000000..89bbc0a --- /dev/null +++ b/docs/examples/pricing/ui/pricing/pricing.scss @@ -0,0 +1,39 @@ +// pricing.scss — VIEW styles of the /pricing screen (MVC), external to the +// markup. Compiled to flat CSS by `wo build` (strict SCSS subset: variables, +// nesting, @use of partials — see docs/plan/exploration/ui/08-mvc-structure.md +// decision 3) and served as a static asset via sendfile. + +$accent: #0a7d4f; +$muted: #6b7280; +$border: #e5e7eb; + +.pricing { + max-width: 56rem; + margin: 0 auto; + + .price-table { + width: 100%; + border-collapse: collapse; + + th { + text-align: left; + color: $muted; + border-bottom: 2px solid $border; + } + + td { + padding: 0.5rem 0.75rem; + border-bottom: 1px solid $border; + } + + .price { + color: $accent; + font-variant-numeric: tabular-nums; // cells patch live; keep digits steady + } + + .sku { font-family: monospace; } + .updated { color: $muted; } + + tr.watched { background: rgba($accent, 0.06); } + } +} diff --git a/docs/examples/pricing/ui/pricing/pricing.wo b/docs/examples/pricing/ui/pricing/pricing.wo new file mode 100644 index 0000000..d9440a5 --- /dev/null +++ b/docs/examples/pricing/ui/pricing/pricing.wo @@ -0,0 +1,29 @@ +-- pricing.wo — the CONTROLLER of the /pricing screen (MVC). +-- (Phase 13 design artifact; executes from 13d. MVC structure: +-- docs/plan/exploration/ui/08-mvc-structure.md. Anatomy mirrors the v1 +-- Angular component reference/writeonce-app/src/app/article/ — +-- view:/styles: ≈ templateUrl/styleUrl, model: ≈ component fields, +-- actions: ≈ component methods.) +-- +-- The controller binds the MODEL (class Product, types/product.wo) into the +-- VIEW (pricing.htmlx) and is the only place the UI may call class methods. +-- No service layer: the database is in-process, a model binding IS a query. + +##ui +#pricing + route: /pricing + view: pricing.htmlx -- View: plain htmlx, logic-free + styles: pricing.scss -- View styles: external SCSS, compiled at `wo build` + + -- Model → View binding. These names are the root scope of pricing.htmlx. + -- LIVE bindings re-patch the view on every commit that touches a match. + model: + products: LIVE select Product{ name, sku, prices } + watchlist: $session.watchlist + + -- Controller actions — handlers call class methods (plan 13b). + -- The view raises them via wo:action; it never calls methods directly. + actions: + set-price(id, amount): Product{ id == id }.set_price(amount) role: Ops | Admin + watch(id): session.watchlist += id + unwatch(id): session.watchlist -= id diff --git a/docs/examples/pricing/wo.toml b/docs/examples/pricing/wo.toml new file mode 100644 index 0000000..adb4c34 --- /dev/null +++ b/docs/examples/pricing/wo.toml @@ -0,0 +1,15 @@ +name = "pricing" +version = "0.1.0" +description = "Phase 13 demo: class model (state + methods, no inheritance) + live pricing fan-out" + +[runtime] +wo = ">= 0.1" + +[app] +listen = "127.0.0.1:8080" + +# RAM-resident engine; disk (phases 10-12) is durability behind the read path, +# never in front of it. See docs/plan/13-class-model-live-pricing.md § 13e. +[database] +data_dir = "./data" +isolation = "snapshot" diff --git a/docs/plan/09-concurrency-scaleout.md b/docs/plan/09-concurrency-scaleout.md index eb9e4cf..c39b9aa 100644 --- a/docs/plan/09-concurrency-scaleout.md +++ b/docs/plan/09-concurrency-scaleout.md @@ -63,18 +63,28 @@ Reference cards already exist for most; this phase adds the ones that are cross- Each one lands as its own numbered plan doc when ready for implementation. Smoke test (`cargo run --bin wo -- run docs/examples/ecommerce` serves correctly) stays green after every sub-phase. -### `09a-thread-per-core.md` — N event loops, `SO_REUSEPORT` +### `09a-thread-per-core.md` — N event loops, `SO_REUSEPORT` — ✅ shipped Introduce a thread-pool manager at `crates/rt/src/runtime/scheduler.rs` (Go parallel: `proc.go`). Spawn `WO_THREADS` OS threads at boot; each pins itself and runs an `EventLoop`. Replace the single `Listener` with per-thread listeners bound `SO_REUSEPORT` to the same port. State is still global at first (shared `Arc>`) — one thing at a time. Exit criterion: `wo run` boots N threads visible in `ps -T`, accepts load balanced across them per `ss -tnp`, no regression in the 20-assertion blog smoke. -### `09b-sharded-engine.md` — per-thread engine state +**Shipped:** `scheduler.rs` (~190 LOC) ports the proven [`wo-rt-c` phase A](./exploration/c-runtime/00-plan.md) sequence: workers named `wo-shard-`, pinned via `sched_setaffinity` (verified tid→cpu 0,1,2,3); `Listener::bind_reuseport` (+ unit test: two binds on one port succeed, plain bind still fails); signals blocked in `main` before spawn, worker 0 owns the `signalfd` and broadcasts shutdown through per-worker `eventfd`s. Measured: 2,400 concurrent requests spread evenly across 4 shards (49.8–56.6 M ns on-CPU per shard via `schedstat`); blog CRUD + 501-stub smoke green; ecommerce/hello/pricing boot unchanged; `WO_THREADS=1` preserves the old single-threaded behavior; SIGTERM joins all shards. Engine remains `Arc>` per this sub-phase's scope — 09b shards it. + +### `09b-sharded-engine.md` — per-thread engine state — ✅ shipped Partition the in-memory engine catalog + row BTreeMaps by shard id (= thread id). Shard key is `customer.id` for ecommerce / `author.id` for blog / per-type default for anything else. Add a shard router in front of every REST/WS handler: resolve the shard from the request's identifying field, send an in-process message to that thread's mailbox, await response. Shared `Arc>` goes away; each thread owns its slice. Cross-shard reads (admin `list orders`) fan out to every thread and merge results. -### `09c-per-shard-wal.md` — one WAL file per shard +**Shipped** (per-type-default shard key; declared-field shard keys await Phase 4's typed wire layer): `crates/rt/src/shard.rs` — `ShardBus` (per-shard mpsc job mailbox + mail `eventfd`) and `ShardCtx` (each worker's own `Engine`). `Engine::for_shard` mints interleaved ids (shard t: t+1, t+1+n, …) so `owner(id) = (id-1) % n` needs zero coordination; creates are always local, point ops hop at most once as boxed-closure jobs, lists fan out and merge by id. Deadlock-free by two rules: jobs never block (pure local engine ops), and waiters pump their own inbox while parked. `Arc>` is deleted; `HandlerFn` dropped its `Send+Sync` bounds (routers are thread-local now). Verified: cross-shard GET/PATCH/404 against rows owned by other shards, merged lists from all shards, all samples green, 42 unit tests (incl. interleave + cross-thread round-trip), clean broadcast shutdown. Measured: durable-free writes 74.9k → **112.9k/s (+51%)**, write p99 4.5 → 3.4 ms vs 09a on the same box — the read path stays connection-setup-bound until keep-alive/io_uring land (09's later phases). + +### `09c-per-shard-wal.md` — one WAL file per shard — ✅ shipped (epoll-stage scope) Each thread has its own `foo.wal` + `foo.data` + per-thread `io_uring` ring ([phase 11's durability work](./11-wal-and-recovery.md), repeated per shard). Recovery is parallel across threads. No shared WAL writer thread. Group commit is per-thread. +**Shipped** (`crates/rt/src/wal.rs` + engine integration): per-shard `shard-.rwal` under `WO_DATA` (default `./wo-data`, `off` disables), C-prototype frame format (`len|crc32|payload|COMMIT`, hand-rolled CRC32, fallocate prealloc), JSON `WalRec` payloads carrying full post-default rows so replay is byte-exact. Durability hooks inside `Engine::{create,update,delete}` — RAM apply → append → `fdatasync` → return, with undo-on-WAL-failure — which puts every ack behind the fsync *including cross-shard jobs* (the reply leaves the owner only after its engine call returns durable). Boot replays per shard in parallel before accepts arm; torn tails drop whole; the id high-water restores per stride; a `meta` file refuses a mismatched `WO_THREADS`. **Deliberately deferred to the io_uring port: group commit** — batching acks on a per-tick fsync over the epoll loop would reopen the ack-before-fsync race the C phase-F crash test caught, so this stage pays one `fdatasync` per commit (~1% on tmpfs; 44 unit tests incl. replay round-trip + torn-tail). Verified e2e: 32-record crash recovery exact (creates/update/delete, ~4 ms/shard), no id collisions post-recovery, durable 178.3k commits/s vs 180.1k non-durable. The `.data` snapshot/compaction half stays with phase 11. + +**Follow-up shipped — io_uring group commit** (`runtime/netpoll_io_uring.rs` — raw ring, kernel ABI structs by hand, no liburing; `wal::WalGroup`): mutations stage frames and **park their acks** instead of fsyncing inline; once per loop tick the worker submits the whole batch as one `WRITE`→`FSYNC` linked SQE pair (one `io_uring_enter`); the fsync CQE releases every parked ack — local responses via gated `Parked` connections with the C-proven generation stamps, cross-shard replies via parked callbacks on the owner's batch. The epoll loop polls the ring fd as an ordinary event source (full network port still pending). **Two bugs found and fixed during real-disk verification:** (1) the write and fsync CQEs of a linked pair routinely land in *different ticks* on ext4 — releasing on the first CQE alone acked before durability (caught because the pre-fix number, 230k/s, was impossibly fast for the disk); (2) reused `user_data` could attribute a stale CQE to the wrong batch — now `(batch_seq << 1) | op-bit`. **A new deadlock class was designed out**: two shards mutually parked on each other's batches would never reach their tick-end flush — `wal_pump` now runs inside every cross-shard wait loop. Measured (8 shards, 64 conns): real ext4/NVMe **27,014 durable commits/s vs 5,765 per-commit (4.7×)**, p50 2.2 ms (one shared fsync per tick); tmpfs 330k/s; reads unaffected (746k/s); crash-under-load on real disk: every acked write recovered; 100 concurrent cross-shard durable PATCHes, zero stalls; `WO_GROUP_COMMIT=off` keeps the per-commit path for A/B. 47 unit tests. + +**Follow-up shipped — HTTP keep-alive** (the C phase-C connection semantics, in `http/{request,response,connection}.rs`): requests report their `keep_alive` wish + consumed byte count; the connection state machine loops over buffered requests (pipelined carry-over included), resets to Reading after each flush, honors `Connection: close` and HTTP/1.0 defaults. Measured on the clean box: reads 227.8k → **770.7k/s (×3.4)**, durable writes 178.3k → **331.5k/s (×1.9)**, p99 993 → 172 µs, zero reconnects at 64 conns — now ahead of Go `net/http` on both axes while fsyncing every write, within ~10% of the C prototype's reads. 46 unit tests (keep-alive ×3, pipelining, close-header). Remaining C-side advantage: io_uring + per-tick group commit (`09g`/io_uring port territory). + ### `09d-cross-shard-subscriptions.md` — LIVE fanout A commit on shard K that creates/updates rows of type T needs to wake subscribers on every shard watching T. Via broadcast: K writes the delta to a per-subscriber-thread mailbox — one message per destination thread, not per subscriber. The destination thread then does the fine-grained predicate match against its local subscription table. Avoids N² traffic when N connections watch the same stream. diff --git a/docs/plan/13-class-model-live-pricing.md b/docs/plan/13-class-model-live-pricing.md new file mode 100644 index 0000000..994ff3b --- /dev/null +++ b/docs/plan/13-class-model-live-pricing.md @@ -0,0 +1,119 @@ +# 13 — Class model + live pricing: state and methods, no inheritance + +**Context sources:** [`../runtime/database/02-wo-language.md`](../runtime/database/02-wo-language.md) (schema layer, § Schema-Layer DML brace disambiguation, § Cross-Paradigm Transaction Coordinator), [`../runtime/database/04-client-api.md`](../runtime/database/04-client-api.md) (subscription engine), [`./09-concurrency-scaleout.md`](./09-concurrency-scaleout.md) (thread-per-core scale-out), [`./exploration/ui/00-overview.md`](./exploration/ui/00-overview.md) + [`./exploration/ui/01-htmlx-format-spec.md`](./exploration/ui/01-htmlx-format-spec.md) (live UI), [`../examples/pricing/`](../examples/pricing/) (the demo this phase makes real), [`../examples/ecommerce/shared/logic/checkout.wo`](../examples/ecommerce/shared/logic/checkout.wo) (the existing `fn … in txn snapshot` signature style methods reuse). + +## Context + +writeonce is declarative by design — `type`, `service`, `policy`, `on ` — and the docs explicitly reject OO. But developers arriving from OO languages keep reaching for "a class with methods", and the request has a legitimate core: **behavior that belongs to a row** (`product.set_price(amount)`) is today only expressible as a free `fn` or a trigger. This phase adds the smallest class model that satisfies it: + +> **`class` = state + methods. No inheritance, no override, no polymorphic dispatch — ever.** Composition via `ref` / `multi`, exactly like `type`. Go-style encapsulation, not Java-style hierarchies. + +The driving workload is the [`pricing` demo](../examples/pricing/): a `Price` class and a `Product` class with methods, products owning prices, a `##ui pricing` screen showing live prices of selected products — RAM-resident data, one product readable by millions of customers at once, price updates pushed live to millions of subscribers, all I/O on kernel primitives. + +This doc is a **master plan** in the style of [`09-concurrency-scaleout.md`](./09-concurrency-scaleout.md): sub-phases 13a–13e at a high level, each landing as its own plan doc when implementation starts. No code changes in this pass — the demo project ships as a design artifact alongside this doc. + +## Goal + +`cargo run --bin wo -- run docs/examples/pricing` serves the demo fully live: `class` declarations parse and store like types, methods execute as row-scoped transactions over RPC, every `set_price` commit pushes a delta to all subscribed clients, the `##ui pricing` screen patches price cells in place, and the read path scales per the phase-09 architecture. + +## Design decisions (locked) + +1. **`class` is the behavior-bearing sibling of `type`.** Identical field grammar — scalars, defaults, `@unique`/`@check`, embedded docs, `ref`, `multi`, `backlink`, unions, plus the same attachable blocks (`service`, `policy`, `on `). One addition: `fn` methods. +2. **Methods are row-scoped transactional functions.** `fn name(args) -> Ret [in txn [snapshot]]` with an implicit `self` bound to the receiving row. Same signature grammar and execution machinery as the free-standing `fn checkout(...) in txn snapshot` already in the ecommerce sample — a method is a free `fn` with a hidden first parameter. No new transaction semantics. +3. **No inheritance.** No `extends`, no `override`, no virtual dispatch, no abstract classes. This kills the table-per-class/single-table storage mapping problem before it exists: a class IS one table (+ its doc/graph parts), exactly like a type. "Is-a" modelling uses tagged unions (already in the language); "has-a" uses `ref`/`multi`. +4. **`self` stays an identifier in the lexer.** Same gotcha as `subscribe`/`receive`/`me` (CLAUDE.md): it must remain usable as a plain name in expose lists and expressions. The parser binds it positionally inside method bodies. +5. **Storage and REST are class-blind.** `Catalog::from_schemas` treats a class exactly like a type; `service rest` blocks generate the same CRUD routes. Methods add RPC routes on top (13b). A migration from `type` to `class` (or back, if no methods) is a no-op for stored data. +6. **Spec wording is amended in 13a, not before.** `docs/runtime/wo-language.md` ("Isn't OO") and `docs/writeonce-pl.md` ("no class model") change in the same commit that makes the parser accept the syntax, so docs never describe an unparseable language. + +## The class surface (normative example) + +```wo +class Price { + id: Id + product: ref Product + amount: Money -- minor units, stdlib scalar + currency: Text = "EUR" + at: Timestamp = now() + + -- Pure method: computes from self, touches nothing else. No txn needed. + fn discounted(pct: Int) -> Money { + return self.amount * (100 - pct) / 100; + } +} + +class Product { + id: Id + sku: SKU @unique + name: Text + prices: multi Price -- products have prices (append-only history) + + fn current_price() -> Money in txn { + return latest(self.prices).amount; + } + + fn set_price(amount: Money) in txn { + insert Price { product: self.id, amount: amount }; + } + + service rest "/api/products" + expose list, get, create, update, delete, subscribe +} +``` + +## Sub-phase sequence + +Each lands as its own numbered plan doc (`13a-…`, `13b-…`) when ready. The blog + ecommerce smoke stays green after every sub-phase. + +### `13a-class-surface.md` — lexer, parser, AST, spec amendments — ✅ shipped + +`class` joins the keyword map (`crates/rt/src/lexer.rs` keyword match, ~line 202 — note `self` stays an ident per decision 4). `parse_type` (`crates/rt/src/parser.rs:120`) takes the leading keyword as a parameter and serves both constructs; `fn` members inside the body parse-and-discard through the existing brace-depth skip — the same mechanism that already swallows `on update … do { … }` triggers. `ast::TypeDecl` gains `is_class: bool`; `Catalog::from_schemas` ignores it (decision 5), so REST CRUD works the moment parsing does. Docs amended in the same change: a "Class Model" subsection in [`02-wo-language.md`](../runtime/database/02-wo-language.md) next to § Schema-Layer DML, the "Isn't OO" paragraph in [`wo-language.md`](../runtime/wo-language.md), the class line in [`writeonce-pl.md`](../writeonce-pl.md), and `just pricing` / `just pricing-demo` recipes. +**Exit (met):** `wo run docs/examples/pricing` parses 2 classes, serves `/api/products` CRUD; parser unit tests (`parses_class_with_methods`, `class_method_braces_do_not_truncate_body`) green; blog/ecommerce/hello unchanged. + +### `13b-method-execution.md` — methods over RPC + +Compile method bodies — statements are the schema-layer DML already specified in [`02-wo-language.md § Schema-Layer DML`](../runtime/database/02-wo-language.md) (`insert` / `update` / `select` with brace disambiguation, `let`, `return`, `assert … otherwise abort`). Each exposed method gets an RPC route: `POST /api/products/:id/set_price` with a JSON args body; `in txn [snapshot]` wraps the body in the coordinator from [`02-wo-language.md § Cross-Paradigm Transaction Coordinator`](../runtime/database/02-wo-language.md). `self` resolves to the `:id` row inside the transaction's snapshot. +**Exit:** `curl -X POST /api/products/1/set_price -d '{"amount": 4999}'` inserts a Price atomically; `current_price` returns it; an aborting method rolls back completely. + +### `13c-live-pricing-push.md` — LIVE deltas on commit + +The subscription registry from [`04-client-api.md`](../runtime/database/04-client-api.md): keyed by type + predicate, matched on commit, deltas framed over WebSocket. Replaces the 501 stub at `/api//live` (`crates/rt/src/server.rs`) with a real upgrade for the pricing demo's needs — `LIVE select Product{ name, prices }` and the `subscribe` expose. This is the Stage 3 milestone scoped to one workload; the full wire protocol stays in Phase 4. +**Exit:** two terminals — `websocat /api/products/live` in one, `set_price` via curl in the other — the delta frame arrives on the open socket within one commit tick, no polling. + +### `13d-pricing-ui.md` — the `/pricing` screen, MVC + +The screen ships as an **MVC triplet** per [`exploration/ui/08-mvc-structure.md`](./exploration/ui/08-mvc-structure.md), built in the sub-phase sequence of [`14-mvc-ui-implementation.md`](./14-mvc-ui-implementation.md): model = the classes themselves, view = [`pricing.htmlx`](../examples/pricing/ui/pricing/pricing.htmlx) (plain htmlx, logic-free) + external [`pricing.scss`](../examples/pricing/ui/pricing/pricing.scss) (strict SCSS subset compiled at `wo build`, no external deps), controller = [`pricing.wo`](../examples/pricing/ui/pricing/pricing.wo) (`route:`/`view:`/`styles:`, `model:` bindings, `actions:` calling the 13b class methods). SSR per [`exploration/ui/01-htmlx-format-spec.md`](./exploration/ui/01-htmlx-format-spec.md), compiler glue per [`02-ui-compiler.md`](./exploration/ui/02-ui-compiler.md), and the vanilla-JS client runtime ([`03-client-runtime.md`](./exploration/ui/03-client-runtime.md)) patches the price cell when the 13c delta lands. The controller's `model:` block is the M→V binding; the watchlist narrows the subscription predicate server-side. +**Exit:** browser at `/pricing` shows selected products; a `set_price` commit from curl changes the price cell in every open browser without reload. + +### `13e-pricing-at-scale.md` — millions of readers, millions of live updates + +No new architecture — this sub-phase wires the demo to [`09-concurrency-scaleout.md`](./09-concurrency-scaleout.md) and adds one mechanism: + +- **RAM-resident:** the engine is the in-memory design of [`03-inmemory-engine.md`](../runtime/database/03-inmemory-engine.md); disk (phases 10–12) is durability behind it, never the read path. +- **Read fan-out:** thread-per-core, shared-nothing shards behind `SO_REUSEPORT` (09a/09b). A single hot product row is owned by one shard but **read-replicated to every shard**: each thread keeps a read-only copy of hot rows, refreshed by the same per-thread broadcast that 09d uses for subscriber fan-out. Millions of concurrent `GET /api/products/1` spread across all cores and never contend — writes still serialize on the owning shard, preserving the single-writer model. +- **Live fan-out to millions:** one `set_price` commit → one delta message per thread (09d, not per subscriber) → each thread predicate-matches its local subscription table and batches socket writes on its own `io_uring` ring ([`exploration/linux/07-io_uring.md`](./exploration/linux/07-io_uring.md)). Kernel primitives only: epoll today ([`done/02-event-loop-epoll.md`](./done/02-event-loop-epoll.md), edge-triggered + group commit), io_uring per phase 09. + +**Exit: verification targets defined and a load harness scripted** (not necessarily met on dev hardware): + +| Metric | Target | How measured | +| --- | --- | --- | +| Concurrent readers of one product | 1 M req/s aggregate on 16 cores | `wrk -c 10000` against `GET /api/products/1`, hot-row replicas on | +| Live subscribers receiving one price update | 1 M open sockets, delta delivered p99 < 250 ms | `websocat` fan-out harness, timestamped frames | +| Commit→first-delta latency | p99 < 10 ms | in-process timestamp at commit vs first socket write | +| Memory | ~2 KB/connection + engine working set | RSS under subscriber load | +| Dep count | 1 (`libc`) | `crates/rt/Cargo.toml` unchanged by this phase | + +## Non-scope + +- **No inheritance, ever, under this plan.** If hierarchy modelling pressure appears, the answer is tagged unions and composition; a future interfaces/traits proposal would be its own phase with its own doc. +- **No method overloading, no statics, no constructors.** Row creation stays `insert` / REST `create`; one method name per class. +- **No client-side method stubs** — `wo gen sdk` method support belongs to Phase 5 (Go SDK), not here. +- **No multi-node distribution.** Same stance as phase 09: single box, threads-as-shards. + +## Cross-references + +- [`../examples/pricing/`](../examples/pricing/) — the demo project this plan makes real, file-by-file phase map in its README. +- [`../runtime/database/02-wo-language.md`](../runtime/database/02-wo-language.md) — schema layer the class grammar extends; transaction coordinator methods reuse. +- [`../runtime/database/04-client-api.md`](../runtime/database/04-client-api.md) — subscription engine 13c scopes down. +- [`./09-concurrency-scaleout.md`](./09-concurrency-scaleout.md) — the scale architecture 13e instantiates. +- [`./exploration/ui/00-overview.md`](./exploration/ui/00-overview.md) — UI track 13d draws on. +- [`../examples/hello/main.wo`](../examples/hello/main.wo) — the minimal example whose `Revision`-trigger pattern is the declarative ancestor of methods. diff --git a/docs/plan/14-mvc-ui-implementation.md b/docs/plan/14-mvc-ui-implementation.md new file mode 100644 index 0000000..12aef8b --- /dev/null +++ b/docs/plan/14-mvc-ui-implementation.md @@ -0,0 +1,88 @@ +# 14 — MVC UI implementation: model = class, view = htmlx + scss, controller = .wo + +**Context sources:** [`./exploration/ui/08-mvc-structure.md`](./exploration/ui/08-mvc-structure.md) (the design this plan implements), [`./exploration/ui/01-htmlx-format-spec.md`](./exploration/ui/01-htmlx-format-spec.md) / [`02-ui-compiler.md`](./exploration/ui/02-ui-compiler.md) / [`03-client-runtime.md`](./exploration/ui/03-client-runtime.md) (the three UI-track pieces this plan sequences, each with port sources and LOC budgets), [`./13-class-model-live-pricing.md`](./13-class-model-live-pricing.md) (the class methods controllers call: 13a/13b; the LIVE deltas views consume: 13c), [`../examples/pricing/ui/pricing/`](../examples/pricing/ui/pricing/) (the reference MVC triplet), [`reference/crates/wo-htmlx/`](../../reference/crates/wo-htmlx/) (the v1 template engine, primary port source). + +## Context + +[`exploration/ui/08-mvc-structure.md`](./exploration/ui/08-mvc-structure.md) locks the screen anatomy: **model** = the `class`/`type` itself, **view** = plain `.htmlx` + external `.scss`, **controller** = a `.wo` file (`route:`/`view:`/`styles:`/`model:`/`actions:`) that binds the model into the view and is the only place UI may call class methods. The UI exploration docs 01–03 already specify the htmlx engine, the `##ui` compiler, and the client runtime in implementable detail. What's missing is the build order, the two genuinely new pieces (the controller format and the SCSS subset compiler), and the wiring into the `13` class-model track. This doc is that sequence. + +Everything lands in **`crates/ui`** (currently a placeholder) and small deltas to `crates/rt` — consistent with the crate inventory in [`crates/README.md`](../../crates/README.md). The deployment shape never changes: one binary serving SSR + database + API on the kernel-primitive runtime. + +## Goal + +`cargo run --bin wo -- run docs/examples/pricing` (after plan 13a–13c land) serves `GET /pricing` as styled SSR HTML; clicking ☆ dispatches a controller action; an Ops `set-price` action calls `Product.set_price`, the commit pushes a delta, and the price cell patches in every open browser without reload — the [plan 13d exit criterion](./13-class-model-live-pricing.md), implemented MVC-shaped. + +## Dependency graph + +``` +14a htmlx engine ──────┬─→ 14c controller format ─→ 14d SSR routes ─→ 14e actions ─→ 14f live patch +14b scss compiler ─────┘ │ │ │ + (14a ∥ 14b — no shared code) needs engine needs 13a+13b needs 13c + (exists today) +``` + +Phases 05/06 (hand-rolled JSON / bespoke error) are orthogonal: `crates/ui` adopts `serde`/`serde_json` per the ui/01 decision and migrates when 05 lands. Phase 08 (`sendfile`) upgrades static-asset serving in 14d when it arrives; 14d ships with plain buffered writes first. + +## Sub-phase sequence + +### `14a-htmlx-engine.md` — port the view engine into `crates/ui` + +Execute [`exploration/ui/01-htmlx-format-spec.md`](./exploration/ui/01-htmlx-format-spec.md) as written: port `reference/crates/wo-htmlx` (585 LOC — `parser.rs`, `ast.rs`, `value.rs`, `registry.rs`, `render.rs` carried over per its table) into `crates/ui/src/htmlx/`, extend with `` structured nodes, `wo:bind` capture, and the `data-wo-manifest` JSON emitter (~250 LOC new). One addition beyond the 01 spec, from the MVC design: `` records whether `source` is a bare name (controller model binding, resolved in 14d) or an inline query — a one-field change to `LiveSubscription`. +**Exit:** the 01 spec's criteria — `cargo build -p ui` green, golden parse+render for every `.htmlx` under `docs/examples/{blog,ecommerce}` **plus** [`pricing/ui/pricing/pricing.htmlx`](../examples/pricing/ui/pricing/pricing.htmlx), manifest matches the 01 schema. + +### `14b-scss-subset.md` — the stylesheet compiler + +New, no port source (~400 LOC at `crates/ui/src/scss/`): scanner → rule tree → flattener. Exactly the subset locked in [08-mvc decision 3](./exploration/ui/08-mvc-structure.md): `$variables`, nesting (including `&`-less descendant flattening), `@use "partials"` (`ui/styles/_*.scss`), comments. **No mixins, functions, `@extend`, or color math** — `rgba($accent, 0.06)` in the reference file compiles by literal substitution of `$accent` and is the only function-form supported. Output is one flat `.css` per screen, written to `target/wo//static/`. +**Exit:** [`pricing.scss`](../examples/pricing/ui/pricing/pricing.scss) → golden-file CSS; unknown construct = compile error naming file:line (never silent passthrough); runs standalone (`wo build` integration is 14d). + +### `14c-controller-format.md` — parse the controller, keep the shorthand + +The `Kind::HashHash("ui")` skip arm in `crates/rt/src/parser.rs` (L80–88) parses for real, into one of two IRs by key-shape dispatch: + +- **Controller form** (`route:`/`view:`/`styles:`/`model:`/`actions:` — the MVC triplet's `pricing.wo`) → new `Controller` IR: route pattern, view/styles paths resolved relative to the screen directory, `model:` entries as named query strings (`LIVE` flag captured, execution deferred), `actions:` entries as `(name, params, target-method-or-fn, role-set)`. +- **Shorthand form** (`source:`/`columns:`/… — the existing ecommerce/blog screens) → the `Screen` IR of [`exploration/ui/02-ui-compiler.md`](./exploration/ui/02-ui-compiler.md), whose codegen emits a generated view + a synthesized `Controller` — [08-mvc decision 5](./exploration/ui/08-mvc-structure.md): the shorthand is sugar over the triplet, one downstream path. + +`##app` and `##component` keep their current skip behaviour (owned by ui/05 and ui/02 respectively). +**Exit:** `pricing.wo` parses to a `Controller` with 2 model bindings + 3 actions; every existing `##ui` screen in blog/ecommerce parses to `Screen` and compiles to an `.htmlx` that 14a round-trips; a controller naming a missing view file is a compile error. + +### `14d-ssr-routes.md` — the single binary serves the screen + +Wire controllers into `crates/rt`'s router (`crates/rt/src/server.rs`): each `Controller.route` becomes a GET route; the handler resolves `model:` bindings against the in-process engine (snapshot `select` now — the `LIVE` flag additionally registers a 13c subscription when that phase is live), renders the view via 14a with the model names as root scope, and emits HTML + manifest + ``. `wo build`/`wo run` gain the asset step: compile SCSS (14b), bake `/_wo/runtime.js` (`include_bytes!`, per ui/03 decision 6), serve `target/wo/.../static/` with buffered writes (upgraded to `sendfile` when [phase 08](./08-sendfile-static-assets.md) lands). +**Exit:** `GET /pricing` returns styled SSR HTML with a valid manifest and resolvable CSS/JS links; `GET /` on the blog sample is unaffected; route table printed at boot includes UI routes alongside REST. + +### `14e-action-dispatch.md` — controller actions call class methods + +`wo:action` buttons POST to `/_wo/action//` with `wo:args` + form payload. The dispatcher looks up the controller's action table, enforces the `role:` set server-side (per-app policy model of [`exploration/ui/07-per-app-policies.md`](./exploration/ui/07-per-app-policies.md)), and invokes the target: a class method via the 13b row-scoped RPC path (`Product{ id == id }.set_price(amount)`) or a free `fn`. Response is 204 — the UI never re-renders from the action response; the visible change arrives as a 13c delta, keeping one update path. +**Requires:** 13a + 13b. **Exit:** the ☆/★ watch toggle round-trips; `set-price` with an Ops session commits a Price; without the role it's 403 and no transaction starts. + +### `14f-live-patching.md` — the browser follows commits + +Execute [`exploration/ui/03-client-runtime.md`](./exploration/ui/03-client-runtime.md) as written (~500 LOC vanilla JS at `crates/ui/assets/wo-runtime.js`, JSON frames over one WebSocket, targeted DOM patching by `data-key` + `wo:bind`, coalescing backpressure, snapshot resync on reconnect), pointed at the 13c subscription endpoint. +**Requires:** 13c. **Exit:** the plan-13d criterion — `set_price` via curl in one terminal, the price cell changes in every open `/pricing` browser without reload; kill the server, restart, the page resyncs on reconnect. + +## Verification targets (after 14f) + +| Check | Target | How | +| --- | --- | --- | +| Golden corpus | every `.htmlx` in blog/ecommerce/pricing parses + renders byte-stable | `cargo test -p ui` golden files | +| SCSS | `pricing.scss` → golden CSS; errors carry file:line | `cargo test -p ui scss` | +| SSR | `GET /pricing` < 5 ms p99 on the dev box (RAM engine, no I/O on read path) | scripted curl loop | +| End-to-end | 13d criterion green | two-terminal demo, scripted in `just pricing-demo` | +| Single binary | UI + DB + API + WS in one `wo build` output, no Node anywhere | `ldd` shows libc only; no build-step JS | +| Dep budget | `crates/ui`: `serde`/`serde_json` only (dropped when phase 05 lands) | `Cargo.toml` review | + +## Non-scope + +- **No SPA router, no client-side templates.** Navigation is full-page SSR; only `wo:bind` cells and `` subtrees mutate in place. (ui/00 decision; unchanged.) +- **No SCSS mixins/functions/`@extend`/color math** beyond literal variable substitution — the subset is a floor, widened only by demonstrated need in the sample corpus. +- **No theme system / design tokens.** Shared partials under `ui/styles/_*.scss` are the only sharing mechanism for now. +- **No component framework.** `##component` partials render server-side via the 14a engine; they have no client behaviour beyond inherited `wo:bind` sites. +- **No changes to the REST API surface.** UI routes live beside `/api/*`; nothing under `/api` changes shape in this plan. + +## Cross-references + +- [`./exploration/ui/08-mvc-structure.md`](./exploration/ui/08-mvc-structure.md) — the design; its exit criteria are satisfied by 14c/14b/14f respectively. +- [`./13-class-model-live-pricing.md`](./13-class-model-live-pricing.md) — 13a/13b gate 14e; 13c gates 14f; 13d's exit criterion is this plan's end-to-end target. +- [`./exploration/ui/00-overview.md`](./exploration/ui/00-overview.md) — the UI track's master frame (per-app binaries, shared DB daemon) that 14d's asset/serving choices stay compatible with. +- [`reference/crates/wo-htmlx/`](../../reference/crates/wo-htmlx/) — primary port source (585 LOC), per ui/01. +- [`../examples/pricing/ui/pricing/`](../examples/pricing/ui/pricing/) — the reference triplet every sub-phase tests against. diff --git a/docs/plan/exploration/ui/08-mvc-structure.md b/docs/plan/exploration/ui/08-mvc-structure.md new file mode 100644 index 0000000..adfccdc --- /dev/null +++ b/docs/plan/exploration/ui/08-mvc-structure.md @@ -0,0 +1,81 @@ +# 08 — MVC structure: model = class, view = htmlx + scss, controller = .wo + +**Context sources:** [`reference/writeonce-app/src/app/`](../../../../reference/writeonce-app/src/app/) (the v1 Angular app whose component anatomy this formalizes), [`./00-overview.md`](./00-overview.md) ("Angular-component-style layout" — `home/{home.wo, home.htmlx, home.css}`), [`./01-htmlx-format-spec.md`](./01-htmlx-format-spec.md) (the view grammar: Mustache + `` + `wo:bind`), [`./02-ui-compiler.md`](./02-ui-compiler.md), [`./03-client-runtime.md`](./03-client-runtime.md), [`../../13-class-model-live-pricing.md`](../../13-class-model-live-pricing.md) (the class methods controllers call), [`../../../examples/pricing/ui/pricing/`](../../../examples/pricing/ui/pricing/) (the reference screen). + +## Goal + +Every writeonce UI screen follows **Model–View–Controller**, with the same file anatomy the v1 Angular app used — but collapsed into the single binary. One screen = one directory with three files: + +``` +ui/pricing/ +├── pricing.wo # Controller — binds the model into the view, exposes actions +├── pricing.htmlx # View — plain htmlx (Mustache + wo:live/wo:bind), no logic +└── pricing.scss # View styles — external, compiled at `wo build` +``` + +## The mapping, against the v1 Angular app + +| MVC role | v1 Angular (`reference/writeonce-app/src/app/`) | writeonce | +| --- | --- | --- | +| **Model** | `models/article.ts` (interface) + `services/article.service.ts` (HTTP fetch) | the `class` / `type` declaration itself (`types/product.wo`). No service layer: the database is in-process, and a model binding **is** a query — `LIVE select` for push, `select` for snapshot | +| **View** | `article.component.html` + `article.component.css` | `pricing.htmlx` + `pricing.scss`. Plain markup; the only dynamic constructs are Mustache paths and `` / `wo:bind` from [`01-htmlx-format-spec.md`](./01-htmlx-format-spec.md) | +| **Controller** | `article.component.ts` (`@Component({templateUrl, styleUrl})`, fields, methods, `service.subscribe(...)`) | `pricing.wo` — declares `view:` / `styles:` (≈ `templateUrl` / `styleUrl`), a `model:` block (≈ component fields), and an `actions:` block whose handlers **call class methods** | + +What Angular needed four layers for (interface, service, component class, template) writeonce does in three files against one runtime — there is no HTTP client between controller and model because there is no network between them. + +## Design decisions + +1. **The controller is declarative, like everything else in `.wo`.** It does not contain imperative rendering code; it declares *what* is bound and *which* method each action invokes. Shape: + + ```wo + ##ui + #pricing + route: /pricing + view: pricing.htmlx -- ≈ Angular templateUrl + styles: pricing.scss -- ≈ Angular styleUrl + + -- Model → View binding. Names declared here are the root scope of + -- the .htmlx file; LIVE bindings re-patch the view on every commit. + model: + products: LIVE select Product{ name, sku, prices } + watchlist: $session.watchlist + + -- Controller actions: the only place UI may invoke class methods. + actions: + set-price(id, amount): Product{ id == id }.set_price(amount) role: Ops | Admin + watch(id): session.watchlist += id + ``` + +2. **The view is plain `.htmlx`, logic-free.** Mustache paths, `{{#each}}`/`{{#if}}`, partials, helpers, `` subtrees, `wo:bind` attributes — nothing else. A `` whose `source` is a bare name resolves against the controller's `model:` block (the M→V binding); an inline query in `source` remains legal for controller-less partials. Views never call methods — they raise actions (`wo:action="set-price"`), the controller dispatches. + +3. **Styles are external SCSS, compiled at `wo build`.** No `