From 45b96f0466f048f7f4b96fc81e1d4c16a435c69a Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 12 Jul 2026 05:43:49 +0200 Subject: [PATCH] # writeonce - method execution POST /api/products/1/set_price -d '{"amount": 4999}' --- CLAUDE.md | 5 +- crates/rt/src/ast.rs | 71 +++- crates/rt/src/compile.rs | 5 + crates/rt/src/engine.rs | 91 ++++- crates/rt/src/http/response.rs | 1 + crates/rt/src/lib.rs | 1 + crates/rt/src/method.rs | 485 +++++++++++++++++++++++ crates/rt/src/parser.rs | 319 ++++++++++++++- crates/rt/src/server.rs | 153 ++++++- crates/rt/src/wal.rs | 8 +- docs/examples/pricing/README.md | 4 +- docs/examples/pricing/types/product.wo | 8 +- docs/plan/00-kanban.md | 8 +- docs/plan/13-class-model-live-pricing.md | 8 +- docs/runtime/database/02-wo-language.md | 2 +- justfile | 10 +- 16 files changed, 1144 insertions(+), 35 deletions(-) create mode 100644 crates/rt/src/method.rs diff --git a/CLAUDE.md b/CLAUDE.md index 2d2a756..50db30e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -96,6 +96,7 @@ Covered in `docs/runtime/database/02-wo-language.md § Concurrency Model` and `0 | `engine.rs` (Value + Row helpers) | `value` | 2 | | `engine.rs` (Engine + Catalog) | `engine` | 2 | | `compile.rs` | `engine` | 2 | +| `method.rs` (13b method executor) | `logic` | 6 | | `server.rs` | `http` + `service` | 4 / 6 | | `bin/wo.rs` | stays in `rt` (the binary) | — | @@ -126,8 +127,8 @@ Stage-3 stubs (501) and policy-shaped 405/404 responses are **intentional and do - **`rt` is monolithic on purpose.** Splitting it into the 14 sibling crates is Phase-by-Phase work, not a Stage-2 refactor. - **`reference/crates/` is its own workspace.** Running `cargo build` at the root does not build v1. Running it in `reference/crates/` does. -- **Parser identifiers vs. keywords.** `subscribe`, `receive`, `expect_abort`, `me` are NOT keywords in the lexer — they stay as plain idents so they can appear as operation names in `expose` lists. Adding them to the keyword map breaks `service rest "..." expose subscribe`. -- **Parser skip-on-block.** Unknown triggers (`on update do ...`) are parsed-and-discarded by brace-depth-aware skipping. Object literals like `{ article_id: self.id }` inside trigger actions contain `}` that must not be mistaken for the type's outer close brace — the depth counter exists specifically because of this. +- **Parser identifiers vs. keywords.** `subscribe`, `receive`, `expect_abort`, `me`, `self`, and lowercase `insert` are NOT keywords in the lexer — they stay as plain idents (only SQL-layer `INSERT` is a keyword). The 13b statement parser matches `insert` as an ident; adding these to the keyword map breaks `service rest "..." expose subscribe` and method bodies. +- **Parser skip-on-block.** Unknown triggers (`on update do ...`) are parsed-and-discarded by brace-depth-aware skipping. Object literals like `{ article_id: self.id }` inside trigger actions contain `}` that must not be mistaken for the type's outer close brace — the depth counter exists specifically because of this. Exception since 13b: `fn` inside a `class` parses into a real `MethodDecl` (body statements + expressions, executed by `method.rs`); `fn` inside a plain `type` still skips. - **Newline significance.** The lexer emits `Kind::Newline` tokens and the parser uses them to end policy/trigger lines. Do not filter newlines globally. - **Default-value parsing.** `= now()` is recognised explicitly as `DefaultExpr::Now`; anything else falls into an opaque-expression path that `engine::eval_default` then **omits from created rows** (computed fields display as empty, not as debug-printed tokens). - **Binary variable shadowing.** `crates/rt/src/bin/wo.rs` has `let rt = ...` (a tokio runtime handle) inside `run()` that shadows the crate named `rt`. Inside `run()` the variable wins; inside `serve()` (a different function) `rt::` refers to the crate. Don't rename the variable without also auditing the crate-path references. diff --git a/crates/rt/src/ast.rs b/crates/rt/src/ast.rs index cee2f77..957f038 100644 --- a/crates/rt/src/ast.rs +++ b/crates/rt/src/ast.rs @@ -16,10 +16,75 @@ pub struct TypeDecl { 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. + /// class-blind (plan 13 decision 5); the flag gates method parsing — + /// `fn` members of a `class` compile into [`MethodDecl`]s (13b), while + /// `fn` inside a plain `type` keeps the 13a parse-and-discard behaviour. pub is_class: bool, + /// Row-scoped methods (plan 13b). Only populated for classes. + pub methods: Vec, +} + +/// `fn name(args) -> Ret [in txn [snapshot]] { body }` — a row-scoped +/// transactional function with an implicit `self` receiver (plan 13 +/// decision 2). Served over RPC as `POST /:id/`. +#[derive(Debug, Clone)] +pub struct MethodDecl { + pub name: String, + /// `(name, declared type)` — the type is diagnostic-only in Stage 2. + pub params: Vec<(String, String)>, + pub ret: Option, + pub txn: TxnMode, + pub body: Vec, +} + +/// Transaction annotation on a method. The Stage-2 engine is single-threaded +/// per shard, so every method already executes atomically and in isolation; +/// the mode is recorded for diagnostics and for the future MVCC coordinator. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TxnMode { None, Txn, Snapshot, Serializable } + +/// Method-body statement — the schema-layer DML subset of +/// `02-wo-language.md § Schema-Layer DML` that 13b executes. +#[derive(Debug, Clone)] +pub enum Stmt { + /// `let name = expr` — binding is a snapshot (spec rule). + Let { name: String, expr: Expr }, + /// `insert Type { field: expr, ... }` — construction form. + Insert { ty: String, fields: Vec<(String, Expr)> }, + /// `return [expr]` + Return { expr: Option }, + /// `assert expr [otherwise abort ["msg"]]` — false aborts the txn. + Assert { cond: Expr, msg: Option }, + /// `if cond { ... } [else { ... }]` (else-if chains nest in `otherwise`). + If { cond: Expr, then: Vec, otherwise: Vec }, +} + +#[derive(Debug, Clone)] +pub enum Expr { + Int(i64), + Str(String), + Bool(bool), + Null, + /// Argument, `let` binding, or `self` (bound positionally — `self` stays + /// a plain identifier in the lexer, plan 13 decision 4). + Ident(String), + /// `base.field` — plain object access, or relation resolution when the + /// base is `self` and the field is a `multi`/`backlink` relation. + Field(Box, String), + /// `name(args)` — builtins: `latest`, `count`, `now`. + Call(String, Vec), + Unary(UnOp, Box), + Binary(BinOp, Box, Box), +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum UnOp { Neg, Not } + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum BinOp { + Add, Sub, Mul, Div, Mod, + Eq, Ne, Lt, Le, Gt, Ge, + And, Or, } #[derive(Debug, Clone)] diff --git a/crates/rt/src/compile.rs b/crates/rt/src/compile.rs index 96b6eb2..1f0a026 100644 --- a/crates/rt/src/compile.rs +++ b/crates/rt/src/compile.rs @@ -21,6 +21,10 @@ pub struct CompiledType { /// those specially but keeps them in the type's row object for echo. pub fields: Vec, pub services: Vec, + /// Row-scoped methods (plan 13b) — populated for `class` declarations, + /// empty for plain `type`s. Storage stays class-blind; methods only add + /// RPC routes on top. + pub methods: Vec, /// True iff the type declared an `id: Id` column. Auto-populated on insert. pub has_id: bool, } @@ -42,6 +46,7 @@ impl Catalog { name: t.name.clone(), fields: t.fields, services: t.services, + methods: t.methods, has_id, }); } diff --git a/crates/rt/src/engine.rs b/crates/rt/src/engine.rs index bfb19f8..f69f595 100644 --- a/crates/rt/src/engine.rs +++ b/crates/rt/src/engine.rs @@ -33,6 +33,24 @@ pub struct Engine { /// 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, + /// Active method transaction (plan 13b): mutations defer their WAL + /// records here and journal undo entries; `commit_txn` emits one + /// [`WalRec::Txn`] frame, `abort_txn` reverts RAM in reverse order. + txn: Option, +} + +#[derive(Debug, Default)] +struct TxnState { + wal: Vec, + undo: Vec, +} + +/// Inverse of one applied mutation — enough to restore the pre-txn RAM state. +#[derive(Debug)] +enum Undo { + Created { ty: String, id: i64 }, + Updated { ty: String, id: i64, prev: Row }, + Deleted { ty: String, id: i64, row: Row }, } #[derive(Debug)] @@ -55,7 +73,8 @@ impl Engine { tables.insert(name.clone(), BTreeMap::new()); next_id.insert(name.clone(), shard as i64 + 1); } - Self { catalog, tables, next_id, id_step: n_shards.max(1) as i64, wal: None, staged: false } + Self { catalog, tables, next_id, id_step: n_shards.max(1) as i64, wal: None, staged: false, + txn: None } } /// Attach a per-commit WAL (fsync inside each mutation). Must happen @@ -139,12 +158,69 @@ impl Engine { WalRec::Delete { ty, id } => { if let Some(table) = self.tables.get_mut(ty) { table.remove(id); } } + // A method's mutations — the frame validated whole, apply all. + WalRec::Txn { recs } => { + for r in recs { self.replay(r); } + } + } + } + + // --- method transactions (plan 13b) --- + + /// Enter method-transaction mode: subsequent mutations journal undo + /// entries and defer their WAL records until [`commit_txn`]. + pub fn begin_txn(&mut self) -> Result<()> { + if self.txn.is_some() { + anyhow::bail!("nested method transactions are not supported"); + } + self.txn = Some(TxnState::default()); + Ok(()) + } + + /// Commit the active transaction: every deferred record leaves as ONE + /// `WalRec::Txn` frame (atomic on replay). `Err` means the WAL rejected + /// the frame — RAM has been rolled back and the caller must not ack. + pub fn commit_txn(&mut self) -> Result<()> { + let Some(t) = self.txn.take() else { + anyhow::bail!("commit_txn without begin_txn"); + }; + if t.wal.is_empty() { return Ok(()); } // read-only method — nothing to log + if let Err(e) = self.wal_log(crate::wal::WalRec::Txn { recs: t.wal }) { + self.apply_undo(t.undo); // never ack non-durable + return Err(e); + } + Ok(()) + } + + /// Abort the active transaction: revert RAM in reverse order, log nothing. + pub fn abort_txn(&mut self) { + if let Some(t) = self.txn.take() { + self.apply_undo(t.undo); + } + } + + fn apply_undo(&mut self, undo: Vec) { + for u in undo.into_iter().rev() { + match u { + Undo::Created { ty, id } => { + if let Some(t) = self.tables.get_mut(&ty) { t.remove(&id); } + } + Undo::Updated { ty, id, prev } | Undo::Deleted { ty, id, row: prev } => { + if let Some(t) = self.tables.get_mut(&ty) { t.insert(id, prev); } + } + } } } /// Make a mutation durable (per-commit) or stage it (group). `Err` - /// means the caller must undo the RAM apply. + /// means the caller must undo the RAM apply. Inside a method + /// transaction the record is deferred instead — durability happens + /// once, at [`commit_txn`], as a single atomic frame. fn wal_log(&mut self, rec: crate::wal::WalRec) -> Result<()> { + if let Some(t) = self.txn.as_mut() { + t.wal.push(rec); + return Ok(()); + } match self.wal.as_mut() { Some(WalBackend::PerCommit(w)) => { w.append(&rec).map_err(|e| anyhow::anyhow!("wal append: {e}"))?; @@ -193,6 +269,9 @@ impl Engine { row.insert("id".into(), json!(id)); self.tables.get_mut(ty).unwrap().insert(id, row.clone()); + if let Some(t) = self.txn.as_mut() { + t.undo.push(Undo::Created { ty: ty.into(), id }); + } // 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 @@ -214,6 +293,9 @@ impl Engine { } } let updated = row.clone(); + if let Some(t) = self.txn.as_mut() { + t.undo.push(Undo::Updated { ty: ty.into(), id, prev: prev.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); @@ -225,6 +307,9 @@ impl Engine { let table = self.tables.get_mut(ty) .ok_or_else(|| anyhow::anyhow!("no such type: {ty}"))?; let Some(removed) = table.remove(&id) else { return Ok(false) }; + if let Some(t) = self.txn.as_mut() { + t.undo.push(Undo::Deleted { ty: ty.into(), id, row: removed.clone() }); + } 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); @@ -293,7 +378,7 @@ impl Engine { } } -fn now_iso8601() -> String { +pub(crate) fn now_iso8601() -> String { let d = SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default(); diff --git a/crates/rt/src/http/response.rs b/crates/rt/src/http/response.rs index 83833c2..21eed05 100644 --- a/crates/rt/src/http/response.rs +++ b/crates/rt/src/http/response.rs @@ -18,6 +18,7 @@ impl Status { pub const BAD_REQUEST: Status = Status(400, "Bad Request"); pub const NOT_FOUND: Status = Status(404, "Not Found"); pub const METHOD_NOT_ALLOWED: Status = Status(405, "Method Not Allowed"); + pub const CONFLICT: Status = Status(409, "Conflict"); pub const INTERNAL_SERVER_ERROR: Status = Status(500, "Internal Server Error"); pub const NOT_IMPLEMENTED: Status = Status(501, "Not Implemented"); } diff --git a/crates/rt/src/lib.rs b/crates/rt/src/lib.rs index 637d76e..32e8c5d 100644 --- a/crates/rt/src/lib.rs +++ b/crates/rt/src/lib.rs @@ -9,6 +9,7 @@ pub mod compile; pub mod engine; pub mod http; pub mod lexer; +pub mod method; pub mod parser; pub mod runtime; pub mod server; diff --git a/crates/rt/src/method.rs b/crates/rt/src/method.rs new file mode 100644 index 0000000..4cc89b9 --- /dev/null +++ b/crates/rt/src/method.rs @@ -0,0 +1,485 @@ +//! Row-scoped method execution — plan 13b (`docs/plan/13-class-model-live-pricing.md`). +//! +//! A class method is a free `fn` with a hidden first parameter: `self` binds +//! to the receiving row (fetched by the RPC route's `:id` on the owning +//! shard), arguments arrive as a JSON object, and the body — the +//! schema-layer DML of `02-wo-language.md` — runs inside an engine method +//! transaction. Commit emits ONE `WalRec::Txn` frame (all-or-nothing on +//! replay); any abort or execution error rolls RAM back completely. +//! +//! The Stage-2 engine is single-threaded per shard, so `in txn` / +//! `in txn snapshot` are already serializable by construction — the mode is +//! accepted and recorded, and every method (pure ones included) runs under +//! begin/commit so the semantics stay uniform when MVCC arrives. + +use crate::ast::{BinOp, Expr, FieldTy, MethodDecl, Stmt, UnOp}; +use crate::engine::Engine; + +use serde_json::{json, Map, Value}; + +/// Why a method call failed — shaped for the RPC layer's status mapping. +#[derive(Debug)] +pub enum MethodError { + /// No row of the receiving type with the requested id (→ 404). + NoSuchRow, + /// Missing / malformed arguments (→ 400). + BadArgs(String), + /// `assert … otherwise abort` fired — transaction rolled back (→ 409). + Abort(String), + /// Execution error (bad field, arithmetic on non-numbers, WAL failure…) + /// — transaction rolled back (→ 500). + Exec(String), +} + +impl std::fmt::Display for MethodError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + MethodError::NoSuchRow => write!(f, "no such row"), + MethodError::BadArgs(m) => write!(f, "bad arguments: {m}"), + MethodError::Abort(m) => write!(f, "aborted: {m}"), + MethodError::Exec(m) => write!(f, "{m}"), + } + } +} + +/// Call `m` on row `id` of type `ty`. Runs entirely on the engine it's +/// handed — the RPC route routes to the owning shard before calling this. +/// Returns the method's `return` value (`Null` if it falls off the end). +pub fn call( + e: &mut Engine, + ty: &str, + id: i64, + m: &MethodDecl, + args: &Map, +) -> Result { + // `self` is a snapshot of the receiving row at entry (spec: bindings + // are snapshots; the row-scoped txn reads one consistent state). + let row = e.get(ty, id) + .map_err(|err| MethodError::Exec(err.to_string()))? + .ok_or(MethodError::NoSuchRow)?; + + // Bind declared parameters. Missing → 400; extras are ignored. + let mut scope: Map = Map::new(); + for (pname, pty) in &m.params { + let Some(v) = args.get(pname) else { + return Err(MethodError::BadArgs(format!("missing argument `{pname}` ({pty})"))); + }; + scope.insert(pname.clone(), v.clone()); + } + + e.begin_txn().map_err(|err| MethodError::Exec(err.to_string()))?; + let mut cx = Cx { e, self_ty: ty, self_row: &row, scope }; + match exec_block(&mut cx, &m.body) { + Ok(flow) => { + cx.e.commit_txn().map_err(|err| MethodError::Exec(err.to_string()))?; + Ok(match flow { Flow::Return(v) => v, Flow::Continue => Value::Null }) + } + Err(err) => { + cx.e.abort_txn(); + Err(err) + } + } +} + +/// Statement-level control flow. +enum Flow { + Continue, + Return(Value), +} + +struct Cx<'a> { + e: &'a mut Engine, + self_ty: &'a str, + self_row: &'a Map, + scope: Map, +} + +fn exec_block(cx: &mut Cx, stmts: &[Stmt]) -> Result { + for s in stmts { + match exec_stmt(cx, s)? { + Flow::Continue => {} + r @ Flow::Return(_) => return Ok(r), + } + } + Ok(Flow::Continue) +} + +fn exec_stmt(cx: &mut Cx, s: &Stmt) -> Result { + match s { + Stmt::Let { name, expr } => { + let v = eval(cx, expr)?; + cx.scope.insert(name.clone(), v); + Ok(Flow::Continue) + } + Stmt::Insert { ty, fields } => { + let mut body = Map::new(); + for (fname, fexpr) in fields { + body.insert(fname.clone(), eval(cx, fexpr)?); + } + cx.e.create(ty, Value::Object(body)) + .map_err(|e| MethodError::Exec(format!("insert {ty}: {e}")))?; + Ok(Flow::Continue) + } + Stmt::Return { expr } => { + let v = match expr { + Some(e) => eval(cx, e)?, + None => Value::Null, + }; + Ok(Flow::Return(v)) + } + Stmt::Assert { cond, msg } => { + if truthy(&eval(cx, cond)?) { + Ok(Flow::Continue) + } else { + Err(MethodError::Abort( + msg.clone().unwrap_or_else(|| "assertion failed".into()))) + } + } + Stmt::If { cond, then, otherwise } => { + if truthy(&eval(cx, cond)?) { + exec_block(cx, then) + } else { + exec_block(cx, otherwise) + } + } + } +} + +fn truthy(v: &Value) -> bool { + match v { + Value::Bool(b) => *b, + Value::Null => false, + _ => true, + } +} + +fn eval(cx: &mut Cx, e: &Expr) -> Result { + match e { + Expr::Int(n) => Ok(json!(n)), + Expr::Str(s) => Ok(Value::String(s.clone())), + Expr::Bool(b) => Ok(json!(b)), + Expr::Null => Ok(Value::Null), + Expr::Ident(name) => { + if name == "self" { + return Ok(Value::Object(cx.self_row.clone())); + } + cx.scope.get(name).cloned() + .ok_or_else(|| MethodError::Exec(format!("unknown name `{name}`"))) + } + Expr::Field(base, field) => { + // `self.` resolves the relation; anything else is + // plain object access on the evaluated base. + if matches!(&**base, Expr::Ident(n) if n == "self") { + if let Some(rel) = relation_rows(cx, field)? { + return Ok(rel); + } + } + let b = eval(cx, base)?; + match b { + Value::Object(m) => m.get(field).cloned().ok_or_else(|| + MethodError::Exec(format!("no field `{field}`"))), + other => Err(MethodError::Exec(format!( + "`.{field}` on a non-object value ({other})"))), + } + } + Expr::Call(name, args) => { + let mut vals = Vec::with_capacity(args.len()); + for a in args { vals.push(eval(cx, a)?); } + builtin(name, vals) + } + Expr::Unary(op, inner) => { + let v = eval(cx, inner)?; + match op { + UnOp::Neg => { + let n = as_i64(&v)?; + Ok(json!(-n)) + } + UnOp::Not => Ok(json!(!truthy(&v))), + } + } + Expr::Binary(op, l, r) => { + // Short-circuit boolean operators. + match op { + BinOp::And => { + let lv = eval(cx, l)?; + if !truthy(&lv) { return Ok(json!(false)); } + return Ok(json!(truthy(&eval(cx, r)?))); + } + BinOp::Or => { + let lv = eval(cx, l)?; + if truthy(&lv) { return Ok(json!(true)); } + return Ok(json!(truthy(&eval(cx, r)?))); + } + _ => {} + } + let lv = eval(cx, l)?; + let rv = eval(cx, r)?; + match op { + BinOp::Eq => Ok(json!(lv == rv)), + BinOp::Ne => Ok(json!(lv != rv)), + BinOp::Lt | BinOp::Le | BinOp::Gt | BinOp::Ge => { + let (a, b) = (as_i64(&lv)?, as_i64(&rv)?); + Ok(json!(match op { + BinOp::Lt => a < b, + BinOp::Le => a <= b, + BinOp::Gt => a > b, + _ => a >= b, + })) + } + BinOp::Add | BinOp::Sub | BinOp::Mul | BinOp::Div | BinOp::Mod => { + let (a, b) = (as_i64(&lv)?, as_i64(&rv)?); + if b == 0 && matches!(op, BinOp::Div | BinOp::Mod) { + return Err(MethodError::Exec("division by zero".into())); + } + Ok(json!(match op { + BinOp::Add => a + b, + BinOp::Sub => a - b, + BinOp::Mul => a * b, + BinOp::Div => a / b, + _ => a % b, + })) + } + BinOp::And | BinOp::Or => unreachable!("handled above"), + } + } + } +} + +fn as_i64(v: &Value) -> Result { + v.as_i64().ok_or_else(|| MethodError::Exec(format!("expected a number, got {v}"))) +} + +/// If `field` names a relation on the receiving type, materialize it: +/// `multi T` / `backlink T.f` → array of the related rows, in id order. +/// Returns `Ok(None)` when `field` is not a relation (plain column access). +/// +/// Shard note: the scan runs on the executing (owner) shard only. Related +/// rows created BY methods land here too (creates mint locally on the shard +/// that runs the method — `set_price`'s Price is by construction visible to +/// the same product's `current_price`). Cross-shard relation reads are 09d/ +/// 09e territory. +fn relation_rows(cx: &mut Cx, field: &str) -> Result, MethodError> { + let Some(t) = cx.e.catalog().get(cx.self_ty) else { return Ok(None) }; + let Some(f) = t.fields.iter().find(|f| f.name == field && f.is_relation) else { + return Ok(None); + }; + let (target, link_field) = match &f.ty { + FieldTy::MultiEdge { target, .. } | FieldTy::MultiVia { target, .. } => { + // Find the ref column in the target type that points back at us. + let tt = cx.e.catalog().get(target).ok_or_else(|| + MethodError::Exec(format!("relation `{field}`: unknown type {target}")))?; + let mut backrefs = tt.fields.iter().filter(|g| + matches!(&g.ty, FieldTy::Ref(r) if r == cx.self_ty)); + let Some(link) = backrefs.next() else { + return Err(MethodError::Exec(format!( + "relation `{field}`: {target} has no `ref {}` field", cx.self_ty))); + }; + if backrefs.next().is_some() { + return Err(MethodError::Exec(format!( + "relation `{field}`: {target} has multiple refs to {} — ambiguous", cx.self_ty))); + } + (target.clone(), link.name.clone()) + } + FieldTy::Backlink { target, field: link } => (target.clone(), link.clone()), + _ => return Ok(None), // `ref T` is a stored scalar column — plain access + }; + + let self_id = cx.self_row.get("id").cloned().unwrap_or(Value::Null); + let rows = cx.e.list(&target) + .map_err(|e| MethodError::Exec(e.to_string()))? + .into_iter() + .filter(|r| r.get(&link_field) == Some(&self_id)) + .map(Value::Object) + .collect::>(); + Ok(Some(Value::Array(rows))) +} + +/// Built-in functions available in method bodies. +fn builtin(name: &str, mut args: Vec) -> Result { + match (name, args.len()) { + // `latest(set)` — the row with the highest id. Prices are + // append-only, so highest id = most recently inserted (per shard). + ("latest", 1) => { + let Value::Array(items) = args.remove(0) else { + return Err(MethodError::Exec("latest(): expected a set".into())); + }; + items.into_iter() + .max_by_key(|v| v.get("id").and_then(|i| i.as_i64()).unwrap_or(i64::MIN)) + .ok_or_else(|| MethodError::Exec("latest(): empty set".into())) + } + ("count", 1) => { + match &args[0] { + Value::Array(items) => Ok(json!(items.len() as i64)), + other => Err(MethodError::Exec(format!("count(): expected a set, got {other}"))), + } + } + ("now", 0) => Ok(Value::String(crate::engine::now_iso8601())), + _ => Err(MethodError::Exec(format!( + "unknown function `{name}`/{} (built-ins: latest, count, now)", args.len()))), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::compile::Catalog; + use crate::parser::parse; + + const PRICING: &str = r#" +class Price { + id: Id + product: ref Product + amount: Money + currency: Text = "EUR" + + fn discounted(pct: Int) -> Money { + return self.amount * (100 - pct) / 100; + } +} + +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 { + assert amount > 0 otherwise abort "price must be positive" + insert Price { product: self.id, amount: amount }; + } + + service rest "/api/products" expose list, get, create +} +"#; + + fn engine() -> Engine { + let cat = Catalog::from_schemas(vec![parse(PRICING).unwrap()]).unwrap(); + Engine::new(cat) + } + + fn method<'a>(e: &'a Engine, ty: &str, name: &str) -> MethodDecl { + e.catalog().get(ty).unwrap().methods.iter() + .find(|m| m.name == name).unwrap().clone() + } + + fn args(v: Value) -> Map { + match v { Value::Object(m) => m, _ => Map::new() } + } + + #[test] + fn set_price_inserts_and_current_price_reads_it_back() { + let mut e = engine(); + let p = e.create("Product", json!({"sku": "SKU-1", "name": "Gadget"})).unwrap(); + let id = p["id"].as_i64().unwrap(); + + let set = method(&e, "Product", "set_price"); + call(&mut e, "Product", id, &set, &args(json!({"amount": 4999}))).unwrap(); + call(&mut e, "Product", id, &set, &args(json!({"amount": 5999}))).unwrap(); + + let prices = e.list("Price").unwrap(); + assert_eq!(prices.len(), 2); + assert_eq!(prices[0]["product"].as_i64().unwrap(), id); + assert_eq!(prices[0]["currency"], "EUR"); // default seeded by create + + let cur = method(&e, "Product", "current_price"); + let v = call(&mut e, "Product", id, &cur, &Map::new()).unwrap(); + assert_eq!(v, json!(5999)); + } + + #[test] + fn pure_method_computes_from_self() { + let mut e = engine(); + let p = e.create("Product", json!({"sku": "S", "name": "N"})).unwrap(); + let pid = p["id"].as_i64().unwrap(); + let set = method(&e, "Product", "set_price"); + call(&mut e, "Product", pid, &set, &args(json!({"amount": 1000}))).unwrap(); + + let price_id = e.list("Price").unwrap()[0]["id"].as_i64().unwrap(); + let disc = method(&e, "Price", "discounted"); + let v = call(&mut e, "Price", price_id, &disc, &args(json!({"pct": 25}))).unwrap(); + assert_eq!(v, json!(750)); + } + + #[test] + fn abort_rolls_back_completely() { + let mut e = engine(); + let p = e.create("Product", json!({"sku": "S", "name": "N"})).unwrap(); + let id = p["id"].as_i64().unwrap(); + let set = method(&e, "Product", "set_price"); + + // amount <= 0 trips the assert AFTER nothing, but build a stronger + // case: a method that inserts and THEN aborts must leave no row. + let src = r#" +class Product { + id: Id + fn bad(amount: Money) in txn { + insert Price { product: self.id, amount: amount }; + assert false otherwise abort "always" + } +} +"#; + let sch = parse(src).unwrap(); + let bad = sch.types[0].methods[0].clone(); + + let err = call(&mut e, "Product", id, &bad, &args(json!({"amount": 1}))).unwrap_err(); + assert!(matches!(err, MethodError::Abort(ref m) if m == "always")); + assert_eq!(e.list("Price").unwrap().len(), 0, "aborted insert must roll back"); + + // The plain assert path also rejects without side effects. + let err = call(&mut e, "Product", id, &set, &args(json!({"amount": 0}))).unwrap_err(); + assert!(matches!(err, MethodError::Abort(_))); + assert_eq!(e.list("Price").unwrap().len(), 0); + } + + #[test] + fn missing_row_and_missing_arg_are_typed_errors() { + let mut e = engine(); + let set = method(&e, "Product", "set_price"); + assert!(matches!( + call(&mut e, "Product", 999, &set, &args(json!({"amount": 1}))), + Err(MethodError::NoSuchRow))); + + let p = e.create("Product", json!({"sku": "S", "name": "N"})).unwrap(); + let id = p["id"].as_i64().unwrap(); + assert!(matches!( + call(&mut e, "Product", id, &set, &Map::new()), + Err(MethodError::BadArgs(_)))); + } + + #[test] + fn method_commit_is_one_atomic_wal_frame() { + use crate::wal::Wal; + let path = std::env::temp_dir() + .join(format!("wo-method-wal-{}", std::process::id())); + let _ = std::fs::remove_file(&path); + + let cat = Catalog::from_schemas(vec![parse(PRICING).unwrap()]).unwrap(); + let id; + { + let mut e = Engine::new(cat.clone()); + let (wal, n) = Wal::open_and_replay(&path, &mut e).unwrap(); + assert_eq!(n, 0); + e.attach_wal(wal); + let p = e.create("Product", json!({"sku": "S", "name": "N"})).unwrap(); + id = p["id"].as_i64().unwrap(); + let set = method(&e, "Product", "set_price"); + call(&mut e, "Product", id, &set, &args(json!({"amount": 4999}))).unwrap(); + } + // Recovery: the create frame + ONE txn frame replay into a fresh engine. + let mut e = Engine::new(cat); + let (_, n) = Wal::open_and_replay(&path, &mut e).unwrap(); + assert_eq!(n, 2, "one create frame + one method-txn frame"); + let prices = e.list("Price").unwrap(); + assert_eq!(prices.len(), 1); + assert_eq!(prices[0]["amount"], json!(4999)); + + let cur = method(&e, "Product", "current_price"); + let v = call(&mut e, "Product", id, &cur, &Map::new()).unwrap(); + assert_eq!(v, json!(4999)); + let _ = std::fs::remove_file(&path); + } +} diff --git a/crates/rt/src/parser.rs b/crates/rt/src/parser.rs index 05fc7de..a407905 100644 --- a/crates/rt/src/parser.rs +++ b/crates/rt/src/parser.rs @@ -141,7 +141,10 @@ impl Parser { self.expect(&Kind::LBrace, "'{'")?; - let mut decl = TypeDecl { name, fields: Vec::new(), services: Vec::new(), is_class }; + let mut decl = TypeDecl { + name, fields: Vec::new(), services: Vec::new(), is_class, + methods: Vec::new(), + }; loop { self.skip_newlines(); match self.peek() { @@ -149,7 +152,12 @@ 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; + Kind::KwFn if is_class => { + // 13b: class methods parse for real — signature + DML body. + decl.methods.push(self.parse_method()?); + } + Kind::KwFn => self.skip_block_line()?, // fn inside a plain `type`: + // parse-and-discard (13a); // brace depth keeps the body's // `}` from closing the type Kind::KwService => decl.services.push(self.parse_service()?), @@ -475,6 +483,270 @@ impl Parser { } Ok(ServiceDecl { kind, path, expose }) } + + // --- class methods (plan 13b) --- + + /// `fn name([p: T, ...]) [-> Ret] [in txn [snapshot|serializable]] { body }` + fn parse_method(&mut self) -> Result { + self.expect(&Kind::KwFn, "`fn`")?; + let name = self.expect_ident("method name")?; + + self.expect(&Kind::LParen, "'('")?; + let mut params = Vec::new(); + self.skip_newlines(); + while !matches!(self.peek(), Kind::RParen) { + let pname = self.expect_ident("parameter name")?; + self.expect(&Kind::Colon, "':'")?; + let pty = self.expect_ident("parameter type")?; + params.push((pname, pty)); + self.skip_newlines(); + if !self.accept(&Kind::Comma) { break; } + self.skip_newlines(); + } + self.expect(&Kind::RParen, "')'")?; + + let mut ret = None; + if self.accept(&Kind::Arrow) { + ret = Some(self.expect_ident("return type")?); + } + + let mut txn = TxnMode::None; + if self.accept(&Kind::KwIn) { + self.expect(&Kind::KwTxn, "`txn`")?; + txn = match self.peek() { + Kind::KwSnapshot => { self.advance(); TxnMode::Snapshot } + Kind::KwSerializable => { self.advance(); TxnMode::Serializable } + _ => TxnMode::Txn, + }; + } + + self.skip_newlines(); + self.expect(&Kind::LBrace, "'{' to open method body")?; + let body = self.parse_stmt_block()?; + Ok(MethodDecl { name, params, ret, txn, body }) + } + + /// Statements until the matching `}` (consumed). + fn parse_stmt_block(&mut self) -> Result> { + let mut stmts = Vec::new(); + loop { + self.skip_newlines(); + while self.accept(&Kind::Semicolon) { self.skip_newlines(); } + match self.peek() { + Kind::RBrace => { self.advance(); return Ok(stmts); } + Kind::End => bail!("unexpected end of input inside method body"), + _ => stmts.push(self.parse_stmt()?), + } + } + } + + fn parse_stmt(&mut self) -> Result { + let stmt = match self.peek() { + Kind::KwLet => { + self.advance(); + let name = self.expect_ident("binding name")?; + self.expect(&Kind::Eq, "'='")?; + let expr = self.parse_expr()?; + Stmt::Let { name, expr } + } + // Schema-layer DML `insert` is lowercase and deliberately NOT a + // lexer keyword (only SQL-layer `INSERT` is) — same rule that + // keeps `subscribe`/`me`/`self` usable as plain identifiers. + Kind::Ident(s) if s == "insert" => { + self.advance(); + let ty = self.expect_ident("type name after `insert`")?; + self.expect(&Kind::LBrace, "'{'")?; + let mut fields = Vec::new(); + loop { + self.skip_newlines(); + if self.accept(&Kind::RBrace) { break; } + let fname = self.expect_ident("field name")?; + self.expect(&Kind::Colon, "':'")?; + let expr = self.parse_expr()?; + fields.push((fname, expr)); + self.skip_newlines(); + self.accept(&Kind::Comma); + } + Stmt::Insert { ty, fields } + } + Kind::KwReturn => { + self.advance(); + let expr = if matches!(self.peek(), + Kind::Semicolon | Kind::Newline | Kind::RBrace | Kind::End) + { None } else { Some(self.parse_expr()?) }; + Stmt::Return { expr } + } + Kind::KwAssert => { + self.advance(); + let cond = self.parse_expr()?; + let mut msg = None; + if self.accept(&Kind::KwOtherwise) { + self.expect(&Kind::KwAbort, "`abort`")?; + if let Kind::Str(s) = self.peek().clone() { + self.advance(); + msg = Some(s); + } + } + Stmt::Assert { cond, msg } + } + Kind::KwIf => { + self.advance(); + let cond = self.parse_expr()?; + self.skip_newlines(); + self.expect(&Kind::LBrace, "'{' after if condition")?; + let then = self.parse_stmt_block()?; + let mut otherwise = Vec::new(); + // `else` may sit on the next line. + let mark = self.pos; + self.skip_newlines(); + if self.accept(&Kind::KwElse) { + self.skip_newlines(); + if matches!(self.peek(), Kind::KwIf) { + otherwise.push(self.parse_stmt()?); // else-if chain + } else { + self.expect(&Kind::LBrace, "'{' after else")?; + otherwise = self.parse_stmt_block()?; + } + } else { + self.pos = mark; // no else — restore consumed newlines + } + return Ok(Stmt::If { cond, then, otherwise }); + } + other => bail!( + "line {}: unsupported statement in method body: {other} \ + (13b executes `let`/`insert`/`return`/`assert`/`if`)", + self.peek_line() + ), + }; + // Statement terminator: `;`, newline, or the closing `}`. + match self.peek() { + Kind::Semicolon | Kind::Newline => { self.advance(); } + Kind::RBrace | Kind::End => {} + other => bail!("line {}: expected end of statement, got {other}", self.peek_line()), + } + Ok(stmt) + } + + // --- expressions (precedence climbing: or < and < cmp < add < mul < unary < postfix) --- + + fn parse_expr(&mut self) -> Result { + self.parse_or() + } + + fn parse_or(&mut self) -> Result { + let mut lhs = self.parse_and()?; + while self.accept(&Kind::KwOr) { + let rhs = self.parse_and()?; + lhs = Expr::Binary(BinOp::Or, Box::new(lhs), Box::new(rhs)); + } + Ok(lhs) + } + + fn parse_and(&mut self) -> Result { + let mut lhs = self.parse_cmp()?; + while self.accept(&Kind::KwAnd) { + let rhs = self.parse_cmp()?; + lhs = Expr::Binary(BinOp::And, Box::new(lhs), Box::new(rhs)); + } + Ok(lhs) + } + + fn parse_cmp(&mut self) -> Result { + let lhs = self.parse_add()?; + let op = match self.peek() { + Kind::EqEq => BinOp::Eq, + Kind::NotEq => BinOp::Ne, + Kind::Lt => BinOp::Lt, + Kind::LtEq => BinOp::Le, + Kind::Gt => BinOp::Gt, + Kind::GtEq => BinOp::Ge, + _ => return Ok(lhs), + }; + self.advance(); + let rhs = self.parse_add()?; + Ok(Expr::Binary(op, Box::new(lhs), Box::new(rhs))) + } + + fn parse_add(&mut self) -> Result { + let mut lhs = self.parse_mul()?; + loop { + let op = match self.peek() { + Kind::Plus => BinOp::Add, + Kind::Dash => BinOp::Sub, + _ => return Ok(lhs), + }; + self.advance(); + let rhs = self.parse_mul()?; + lhs = Expr::Binary(op, Box::new(lhs), Box::new(rhs)); + } + } + + fn parse_mul(&mut self) -> Result { + let mut lhs = self.parse_unary()?; + loop { + let op = match self.peek() { + Kind::Star => BinOp::Mul, + Kind::Slash => BinOp::Div, + Kind::Percent => BinOp::Mod, + _ => return Ok(lhs), + }; + self.advance(); + let rhs = self.parse_unary()?; + lhs = Expr::Binary(op, Box::new(lhs), Box::new(rhs)); + } + } + + fn parse_unary(&mut self) -> Result { + match self.peek() { + Kind::Dash => { self.advance(); Ok(Expr::Unary(UnOp::Neg, Box::new(self.parse_unary()?))) } + Kind::KwNot => { self.advance(); Ok(Expr::Unary(UnOp::Not, Box::new(self.parse_unary()?))) } + _ => self.parse_postfix(), + } + } + + /// Primary followed by `.field` chains. + fn parse_postfix(&mut self) -> Result { + let mut e = self.parse_primary()?; + while self.accept(&Kind::Dot) { + let field = self.expect_ident("field name after '.'")?; + e = Expr::Field(Box::new(e), field); + } + Ok(e) + } + + fn parse_primary(&mut self) -> Result { + match self.peek().clone() { + Kind::Int(n) => { self.advance(); Ok(Expr::Int(n)) } + Kind::Str(s) => { self.advance(); Ok(Expr::Str(s)) } + Kind::KwTrue => { self.advance(); Ok(Expr::Bool(true)) } + Kind::KwFalse => { self.advance(); Ok(Expr::Bool(false)) } + Kind::KwNull => { self.advance(); Ok(Expr::Null) } + Kind::LParen => { + self.advance(); + let e = self.parse_expr()?; + self.expect(&Kind::RParen, "')'")?; + Ok(e) + } + Kind::Ident(name) => { + self.advance(); + // `name(args)` — call form. + if self.accept(&Kind::LParen) { + let mut args = Vec::new(); + if !matches!(self.peek(), Kind::RParen) { + loop { + args.push(self.parse_expr()?); + if !self.accept(&Kind::Comma) { break; } + } + } + self.expect(&Kind::RParen, "')'")?; + Ok(Expr::Call(name, args)) + } else { + Ok(Expr::Ident(name)) + } + } + other => bail!("line {}: expected expression, got {other}", self.peek_line()), + } + } } #[cfg(test)] @@ -632,12 +904,53 @@ class Product { 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); + + // 13b: methods parse into real AST. + assert_eq!(t.methods.len(), 2); + let cp = &t.methods[0]; + assert_eq!(cp.name, "current_price"); + assert!(cp.params.is_empty()); + assert_eq!(cp.ret.as_deref(), Some("Money")); + assert_eq!(cp.txn, TxnMode::Txn); + assert_eq!(cp.body.len(), 1); + assert!(matches!(cp.body[0], Stmt::Return { expr: Some(_) })); + + let sp = &t.methods[1]; + assert_eq!(sp.name, "set_price"); + assert_eq!(sp.params, vec![("amount".to_string(), "Money".to_string())]); + assert_eq!(sp.txn, TxnMode::Txn); + assert!(matches!(&sp.body[0], + Stmt::Insert { ty, fields } if ty == "Price" && fields.len() == 2)); + } + + #[test] + fn method_body_expression_ast() { + let src = r#" +class Price { + id: Id + amount: Money + + fn discounted(pct: Int) -> Money { + return self.amount * (100 - pct) / 100; + } +} +"#; + let sch = parse(src).unwrap(); + let m = &sch.types[0].methods[0]; + assert_eq!(m.txn, TxnMode::None); + let Stmt::Return { expr: Some(e) } = &m.body[0] else { panic!("expected return") }; + // ((self.amount * (100 - pct)) / 100) — mul level is left-associative. + let Expr::Binary(BinOp::Div, lhs, rhs) = e else { panic!("expected /: {e:?}") }; + assert!(matches!(**rhs, Expr::Int(100))); + let Expr::Binary(BinOp::Mul, base, paren) = &**lhs else { panic!("expected *") }; + assert!(matches!(&**base, Expr::Field(b, f) if f == "amount" + && matches!(&**b, Expr::Ident(s) if s == "self"))); + assert!(matches!(&**paren, Expr::Binary(BinOp::Sub, _, _))); } #[test] diff --git a/crates/rt/src/server.rs b/crates/rt/src/server.rs index 59bc97c..8ea7d90 100644 --- a/crates/rt/src/server.rs +++ b/crates/rt/src/server.rs @@ -50,18 +50,20 @@ pub fn router(ctx: Rc, catalog: &Catalog) -> Router { let t = catalog.get(name).expect("type present"); for svc in &t.services { if svc.kind != ServiceKind::Rest { continue; } - r = attach_rest(r, ctx.clone(), t.name.clone(), svc.path.clone(), &svc.expose); + r = attach_rest(r, ctx.clone(), t.name.clone(), svc.path.clone(), &svc.expose, + &t.methods); } } r } fn attach_rest( - mut r: Router, - ctx: Rc, - ty: String, - path: String, - ops: &[Operation], + mut r: Router, + ctx: Rc, + ty: String, + path: String, + ops: &[Operation], + methods: &[crate::ast::MethodDecl], ) -> Router { let id_path = format!("{path}/:id"); @@ -107,6 +109,17 @@ fn attach_rest( Operation::Subscribe | Operation::Me | Operation::Custom => {} } } + + // Class methods (plan 13b): every method of a class with a rest service + // is served as a row-scoped RPC. Registered after the CRUD `/:id` routes + // — the extra path segment makes patterns disjoint either way. + for m in methods { + let ctx = ctx.clone(); + let ty = ty.clone(); + let m = m.clone(); + r = r.route(Method::Post, &format!("{path}/:id/{}", m.name), + move |req, params| method_h(&ctx, &ty, &m, req, params)); + } r } @@ -215,6 +228,51 @@ fn delete_h(ctx: &Rc, ty: &str, _req: &Request, params: &RouteParams) } } +/// Row-scoped method RPC (plan 13b): `POST /:id/` with a JSON +/// args object. Executes on the shard that owns the receiving row — inserts +/// the body performs mint on that shard, so everything a method writes it +/// also owns. The whole body commits as one WAL frame; failures roll back. +fn method_h( + ctx: &Rc, + ty: &str, + m: &crate::ast::MethodDecl, + req: &Request, + params: &RouteParams, +) -> Response { + let id = match parse_id(params) { + Ok(id) => id, + Err(r) => return r, + }; + let args = match parse_json_body(req) { + Ok(Value::Object(map)) => map, + Ok(Value::Null) => Default::default(), + Ok(_) => return Response::status(Status::BAD_REQUEST) + .text("method arguments must be a JSON object"), + Err(r) => return r, + }; + let ty_owned = ty.to_string(); + let m_owned = m.clone(); + let res = ctx.run_on(ctx.owner_of(id), move |e| { + crate::method::call(e, &ty_owned, id, &m_owned, &args) + .map_err(|err| (status_for(&err), err.to_string())) + }); + match res { + None => shard_gone(), + Some(Ok(v)) => gate_if_staged(ctx, Response::ok().json(&v)), + Some(Err((status, msg))) => Response::status(status).text(msg), + } +} + +fn status_for(err: &crate::method::MethodError) -> Status { + use crate::method::MethodError::*; + match err { + NoSuchRow => Status::NOT_FOUND, + BadArgs(_) => Status::BAD_REQUEST, + Abort(_) => Status::CONFLICT, + Exec(_) => Status::INTERNAL_SERVER_ERROR, + } +} + fn parse_id(params: &RouteParams) -> Result { params.get("id") .and_then(|s| s.parse::().ok()) @@ -251,6 +309,14 @@ pub fn describe_routes(catalog: &Catalog) -> String { }; out.push_str(&format!(" {m} {p:<30} {n}\n")); } + for m in &t.methods { + let p = format!("{}/:id/{}", svc.path, m.name); + let mode = match m.txn { + crate::ast::TxnMode::None => "", + _ => " in txn", + }; + out.push_str(&format!(" POST {p:<30} method {}.{}{mode}\n", name, m.name)); + } } } out @@ -336,4 +402,79 @@ type Article { id: Id let resp = r.dispatch(&req(Method::Get, "/api/articles/live", b"")); assert_eq!(resp.status.0, 501); } + + const PRICING: &str = r#" +class Price { + id: Id + product: ref Product + amount: Money + service rest "/api/prices" expose list +} + +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 { + assert amount > 0 otherwise abort "price must be positive" + insert Price { product: self.id, amount: amount }; + } + + service rest "/api/products" expose list, get, create +} +"#; + + #[test] + fn class_methods_serve_rpc_routes() { + let (_ctx, r) = build(PRICING); + + let resp = r.dispatch(&req(Method::Post, "/api/products", br#"{"sku":"SKU-1","name":"Gadget"}"#)); + assert_eq!(resp.status.0, 201); + + // The 13b exit criterion: POST :id/set_price inserts a Price atomically… + let resp = r.dispatch(&req(Method::Post, "/api/products/1/set_price", br#"{"amount": 4999}"#)); + assert_eq!(resp.status.0, 200, "{}", String::from_utf8_lossy(&resp.body)); + + // …and current_price returns it. + let resp = r.dispatch(&req(Method::Post, "/api/products/1/current_price", b"")); + assert_eq!(resp.status.0, 200); + assert_eq!(resp.body, b"4999"); + + // The Price row is a real, listable row. + let resp = r.dispatch(&req(Method::Get, "/api/prices", b"")); + let rows: Vec = serde_json::from_slice(&resp.body).unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0]["product"], 1); + assert_eq!(rows[0]["amount"], 4999); + } + + #[test] + fn method_errors_map_to_http_statuses() { + let (_ctx, r) = build(PRICING); + r.dispatch(&req(Method::Post, "/api/products", br#"{"sku":"S","name":"N"}"#)); + + // Unknown row → 404. + let resp = r.dispatch(&req(Method::Post, "/api/products/99/set_price", br#"{"amount": 1}"#)); + assert_eq!(resp.status.0, 404); + + // Missing argument → 400. + let resp = r.dispatch(&req(Method::Post, "/api/products/1/set_price", b"{}")); + assert_eq!(resp.status.0, 400); + + // Aborting method → 409, and the transaction rolled back. + let resp = r.dispatch(&req(Method::Post, "/api/products/1/set_price", br#"{"amount": 0}"#)); + assert_eq!(resp.status.0, 409); + let resp = r.dispatch(&req(Method::Get, "/api/prices", b"")); + assert_eq!(resp.body, b"[]", "aborted method must leave no rows"); + + // GET on a method route → 405 (route exists, wrong verb). + let resp = r.dispatch(&req(Method::Get, "/api/products/1/set_price", b"")); + assert_eq!(resp.status.0, 405); + } } diff --git a/crates/rt/src/wal.rs b/crates/rt/src/wal.rs index 02ebe98..3f697aa 100644 --- a/crates/rt/src/wal.rs +++ b/crates/rt/src/wal.rs @@ -36,13 +36,17 @@ 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)] +/// body (merge is deterministic in log order). `Txn` bundles every mutation +/// of one method call (plan 13b) into a single frame: the frame's CRC + +/// trailer make it replay whole-or-not-at-all, so a crash mid-method can +/// never leave a partial method on disk. +#[derive(Debug, 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 }, + Txn { recs: Vec }, } #[derive(Debug)] diff --git a/docs/examples/pricing/README.md b/docs/examples/pricing/README.md index bcb1bb5..93d5abf 100644 --- a/docs/examples/pricing/README.md +++ b/docs/examples/pricing/README.md @@ -2,7 +2,7 @@ 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). +> **Status: 13a + 13b shipped.** `class` declarations parse, the demo serves real REST CRUD, **and methods execute over RPC** — `wo run docs/examples/pricing` (or `just pricing-demo`) serves `POST /api/products/:id/set_price` and `:id/current_price` as row-scoped transactions: the whole body commits as one WAL frame, an `assert … otherwise abort` rolls everything back (409). `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 @@ -32,7 +32,7 @@ pricing/ | 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 | +| `set_price` / `current_price` over RPC (`POST /api/products/:id/set_price`) | **13b ✅ shipped** | method execution, row-scoped txn, one WAL frame per call | | `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 | diff --git a/docs/examples/pricing/types/product.wo b/docs/examples/pricing/types/product.wo index a686a3d..f1c6aee 100644 --- a/docs/examples/pricing/types/product.wo +++ b/docs/examples/pricing/types/product.wo @@ -17,10 +17,12 @@ class Product { 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). + -- The write path of the live-pricing demo. The assert and the insert + -- commit or roll back together (one WAL frame); 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 { + assert amount > 0 otherwise abort "price must be positive" insert Price { product: self.id, amount: amount }; } diff --git a/docs/plan/00-kanban.md b/docs/plan/00-kanban.md index dc511f4..ea6dea7 100644 --- a/docs/plan/00-kanban.md +++ b/docs/plan/00-kanban.md @@ -47,7 +47,7 @@ All numbers + find-and-fix stories: [09-concurrency-scaleout.md](09-concurrency- | Status | Phase | Doc | Notes | | --- | --- | --- | --- | | ✅ | 13a class surface | [13](13-class-model-live-pricing.md) | `class` parses, CRUD serves, spec amended | -| ⬜ | 13b method execution | [13](13-class-model-live-pricing.md) | `POST /api//:id/`, row-scoped txn — **the next API milestone** | +| ✅ | 13b method execution | [13](13-class-model-live-pricing.md) | `POST /api//:id/`; row-scoped txn; one `WalRec::Txn` frame per call; abort → 409 rollback | | ⬜ | 13c LIVE pricing push | [13](13-class-model-live-pricing.md) | Stage 3 scoped: subscription registry, WS at `/api//live`, replaces the 501 stub | | ⬜ | Stage 3 wire layer | [../runtime/database/04-client-api.md](../runtime/database/04-client-api.md) | full subscription engine + wire protocol; 13c is its beachhead | | ⬜ | 13e pricing at scale | [13](13-class-model-live-pricing.md) | wires demo to 09; hot-row reads | @@ -77,9 +77,9 @@ Ecommerce sample status (verified 2026-06-13, `api.rest` **17/17 expected status ## Suggested order of play (backend) -1. **13b — method execution** (unblocks `fn checkout` semantics, the ecommerce sample's core promise) +1. ~~**13b — method execution**~~ ✅ shipped — methods run as row-scoped transactions over RPC 2. **13c / Stage 3 LIVE** with **09d** fan-out (turns every 501 stub real; the ecommerce order-ops board's backend) -3. **09e — 2PC** (cross-shard `fn checkout` — the canonical ACID demo end-to-end) -4. **15a/15b — MCP core, tools, resources** (needs nothing unshipped; makes every app agent-callable; 15c/15d any time after; 15e waits on 13c + 09d) +3. **09e — 2PC** (cross-shard `fn checkout` — the canonical ACID demo end-to-end; 13b's single-shard txn is its building block) +4. **15a/15b — MCP core, tools, resources** (needs nothing unshipped; makes every app agent-callable — 13b methods become MCP tools in 15d; 15e waits on 13c + 09d) 5. **05/06** dependency removal (mechanical, any time) 6. **10–12** storage completion (snapshots/compaction; arena engine) diff --git a/docs/plan/13-class-model-live-pricing.md b/docs/plan/13-class-model-live-pricing.md index 20c8d0c..5825c60 100644 --- a/docs/plan/13-class-model-live-pricing.md +++ b/docs/plan/13-class-model-live-pricing.md @@ -1,6 +1,6 @@ # 13 — Class model + live pricing: state and methods, no inheritance -> **Kanban: 🔄 in progress** — 13a ✅ shipped; 13b (method execution) is the next API milestone; 13c ⬜; 13d ⏸ parked (frontend); 13e ⬜. Board: [00-kanban.md](00-kanban.md) +> **Kanban: 🔄 in progress** — 13a ✅ shipped; 13b ✅ shipped (methods execute over RPC); 13c (LIVE push) is next; 13d ⏸ parked (frontend); 13e ⬜. Board: [00-kanban.md](00-kanban.md) **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). @@ -71,10 +71,10 @@ Each lands as its own numbered plan doc (`13a-…`, `13b-…`) when ready. The b `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 +### `13b-method-execution.md` — methods over RPC — ✅ shipped -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. +Method bodies compile into a real AST (`ast::{MethodDecl, Stmt, Expr}` — `let`, `insert Type{…}` construction, `return`, `assert … otherwise abort`, `if/else`, arithmetic/comparison/boolean expressions, `self.` resolution, built-ins `latest`/`count`/`now`) and execute in `crates/rt/src/method.rs` on the shard that owns the receiving row (`run_on(owner_of(id))` — inserts mint locally there, so a method owns everything it writes). Every call runs inside an engine **method transaction** (`Engine::{begin,commit,abort}_txn`): mutations journal undo entries and defer their WAL records; commit emits **one `WalRec::Txn` frame** — the frame's CRC makes a method's mutations replay whole-or-not-at-all, so a crash mid-method can never half-apply — and abort reverts RAM in reverse order. Group-commit acks gate on the batch fsync exactly like CRUD (`Response.gate` / parked replies). Routes: `POST /:id/` for every method of a class with a `service rest` block; errors map NoSuchRow→404, BadArgs→400, Abort→409, Exec→500. Lowercase `insert` stays an identifier in the lexer (same rule as `subscribe`/`me`/`self`); only the class-side `fn` arm parses bodies — `fn` inside a plain `type` keeps the 13a skip. +**Exit (met):** `curl -X POST /api/products/1/set_price -d '{"amount": 4999}'` inserts a Price atomically (verified on 2 shards incl. the cross-shard hop, in both WAL modes, and across restart — the Txn frame replays); `current_price` returns it; `set_price {"amount": 0}` trips the sample's assert → 409 and rolls back completely (price unchanged). 12 new unit tests (parser AST, executor, WAL atomicity, route statuses); blog/ecommerce/hello unchanged; `just pricing-demo` runs the whole sequence. ### `13c-live-pricing-push.md` — LIVE deltas on commit diff --git a/docs/runtime/database/02-wo-language.md b/docs/runtime/database/02-wo-language.md index a22a1ab..9f2d983 100644 --- a/docs/runtime/database/02-wo-language.md +++ b/docs/runtime/database/02-wo-language.md @@ -201,7 +201,7 @@ Rules: - **No inheritance.** No `extends`, no override, no virtual dispatch. "Is-a" is a tagged union; "has-a" is `ref`/`multi`. This also kills the table-per-class storage-mapping problem: a class IS one table (+ doc/graph parts), exactly like a type. - **Methods are row-scoped transactional functions.** `fn name(args) -> Ret [in txn [snapshot]]` — the same signature grammar and coordinator as a free-standing `fn` (see `fn checkout` in the ecommerce sample); `self` binds to the receiving row. Bodies are the schema-layer DML of the section above. Exposed as RPC: `POST /api/products/:id/current_price` (plan 13b). - **Storage and REST are class-blind.** The catalog treats `class` exactly like `type`; converting between them is a no-op for stored data. `self` stays a plain identifier in the lexer (same rule as `subscribe`/`me`). -- **Status:** parsing + CRUD shipped (13a, `ast::TypeDecl::is_class`); method execution lands in 13b — until then `fn` bodies are parsed-and-discarded like triggers. +- **Status:** parsing + CRUD shipped (13a, `ast::TypeDecl::is_class`); method execution shipped (13b, `crates/rt/src/method.rs`) — bodies compile to `ast::{Stmt, Expr}` (`let`, `insert`, `return`, `assert … otherwise abort`, `if/else`) and run on the row's owning shard, committing as one atomic `WalRec::Txn` frame; an abort rolls back completely (HTTP 409). `fn` inside a plain `type` is still parsed-and-discarded. ## Query Layer — Hybrid SQL + Cypher, Fixed Glue diff --git a/justfile b/justfile index 30fcfd7..60fbef8 100644 --- a/justfile +++ b/justfile @@ -48,12 +48,13 @@ rt-c-bench port="8085" threads="8" conns="64": pricing: cargo run --bin wo -- run docs/examples/pricing -# class model in action: serve pricing, CRUD round-trip on Product, shut down +# class model in action: CRUD + row-scoped method RPC (13b) on Product. +# WO_THREADS=1 + WO_DATA=off keep ids deterministic and the demo stateless. pricing-demo port="8092": #!/usr/bin/env bash set -euo pipefail cargo build --bin wo - WO_LISTEN=127.0.0.1:{{port}} ./target/debug/wo run docs/examples/pricing & + WO_THREADS=1 WO_DATA=off WO_LISTEN=127.0.0.1:{{port}} ./target/debug/wo run docs/examples/pricing & server=$! trap 'kill $server 2>/dev/null' EXIT base=http://127.0.0.1:{{port}} @@ -62,6 +63,11 @@ pricing-demo port="8092": echo "--- create:"; curl -s -X POST "$base/api/products" -H 'Content-Type: application/json' -d '{"sku":"WO-001","name":"writeonce mug"}'; echo echo "--- list:"; curl -s "$base/api/products"; echo echo "--- patch 1:"; curl -s -X PATCH "$base/api/products/1" -d '{"name":"writeonce mug v2"}'; echo + echo "--- set_price 4999 (method RPC, expect 200):"; curl -s -X POST "$base/api/products/1/set_price" -d '{"amount":4999}' -o /dev/null -w '%{http_code}\n' + echo "--- set_price 5999 (expect 200):"; curl -s -X POST "$base/api/products/1/set_price" -d '{"amount":5999}' -o /dev/null -w '%{http_code}\n' + echo "--- current_price (expect 5999):"; curl -s -X POST "$base/api/products/1/current_price"; echo + echo "--- set_price 0 (assert aborts, expect 409):"; curl -s -X POST "$base/api/products/1/set_price" -d '{"amount":0}'; echo + echo "--- current_price unchanged (expect 5999):"; curl -s -X POST "$base/api/products/1/current_price"; echo echo "--- live (13c pending, expect 501):"; curl -s -o /dev/null -w '%{http_code}\n' "$base/api/products/live" echo "--- delete 1 (expect 204):"; curl -s -X DELETE "$base/api/products/1" -o /dev/null -w '%{http_code}\n'