# writeonce - method execution

POST /api/products/1/set_price -d '{"amount": 4999}'
This commit is contained in:
shoney.arickathil 2026-07-12 05:43:49 +02:00
parent e5203b289d
commit 45b96f0466
16 changed files with 1144 additions and 35 deletions

View file

@ -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.

View file

@ -16,10 +16,75 @@ pub struct TypeDecl {
pub fields: Vec<Field>,
pub services: Vec<ServiceDecl>,
/// 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<MethodDecl>,
}
/// `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 <service-path>/:id/<name>`.
#[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<String>,
pub txn: TxnMode,
pub body: Vec<Stmt>,
}
/// 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<Expr> },
/// `assert expr [otherwise abort ["msg"]]` — false aborts the txn.
Assert { cond: Expr, msg: Option<String> },
/// `if cond { ... } [else { ... }]` (else-if chains nest in `otherwise`).
If { cond: Expr, then: Vec<Stmt>, otherwise: Vec<Stmt> },
}
#[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<Expr>, String),
/// `name(args)` — builtins: `latest`, `count`, `now`.
Call(String, Vec<Expr>),
Unary(UnOp, Box<Expr>),
Binary(BinOp, Box<Expr>, Box<Expr>),
}
#[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)]

View file

@ -21,6 +21,10 @@ pub struct CompiledType {
/// those specially but keeps them in the type's row object for echo.
pub fields: Vec<Field>,
pub services: Vec<ServiceDecl>,
/// 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<MethodDecl>,
/// 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,
});
}

View file

@ -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<TxnState>,
}
#[derive(Debug, Default)]
struct TxnState {
wal: Vec<crate::wal::WalRec>,
undo: Vec<Undo>,
}
/// 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<Undo>) {
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();

View file

@ -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");
}

View file

@ -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;

485
crates/rt/src/method.rs Normal file
View file

@ -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<String, Value>,
) -> Result<Value, MethodError> {
// `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<String, Value> = 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<String, Value>,
scope: Map<String, Value>,
}
fn exec_block(cx: &mut Cx, stmts: &[Stmt]) -> Result<Flow, MethodError> {
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<Flow, MethodError> {
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<Value, MethodError> {
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.<relation>` 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<i64, MethodError> {
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<Option<Value>, 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::<Vec<_>>();
Ok(Some(Value::Array(rows)))
}
/// Built-in functions available in method bodies.
fn builtin(name: &str, mut args: Vec<Value>) -> Result<Value, MethodError> {
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<String, Value> {
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);
}
}

View file

@ -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 <read|write|...> ...
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<MethodDecl> {
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<Vec<Stmt>> {
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<Stmt> {
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<Expr> {
self.parse_or()
}
fn parse_or(&mut self) -> Result<Expr> {
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<Expr> {
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<Expr> {
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<Expr> {
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<Expr> {
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<Expr> {
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<Expr> {
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<Expr> {
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]

View file

@ -50,18 +50,20 @@ pub fn router(ctx: Rc<ShardCtx>, 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<ShardCtx>,
ty: String,
path: String,
ops: &[Operation],
mut r: Router,
ctx: Rc<ShardCtx>,
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<ShardCtx>, ty: &str, _req: &Request, params: &RouteParams)
}
}
/// Row-scoped method RPC (plan 13b): `POST <path>/:id/<method>` 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<ShardCtx>,
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<i64, Response> {
params.get("id")
.and_then(|s| s.parse::<i64>().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::Value> = 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);
}
}

View file

@ -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<WalRec> },
}
#[derive(Debug)]

View file

@ -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 |

View file

@ -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 };
}

View file

@ -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/<t>/:id/<method>`, row-scoped txn — **the next API milestone** |
| ✅ | 13b method execution | [13](13-class-model-live-pricing.md) | `POST /api/<t>/:id/<method>`; 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/<t>/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)

View file

@ -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.<relation>` 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 <service-path>/:id/<method>` 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

View file

@ -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

View file

@ -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'