diff --git a/landing/src/Docs.svelte b/landing/src/Docs.svelte index ce50789..939747b 100644 --- a/landing/src/Docs.svelte +++ b/landing/src/Docs.svelte @@ -13,6 +13,7 @@ { id: 'introduction', title: 'Introduction' }, { id: 'philosophy', title: 'Philosophy & the bet' }, { id: 'honesty', title: 'Is it production-ready?' }, + { id: 'crates', title: 'The workspace map' }, ], }, { @@ -23,6 +24,7 @@ { id: 'features', title: 'Feature flags' }, { id: 'configuration', title: 'Configuration' }, { id: 'layout', title: 'Directory & conventions' }, + { id: 'cli', title: 'The CLI' }, ], }, { @@ -31,16 +33,24 @@ { id: 'routing', title: 'Routing' }, { id: 'requests', title: 'Requests & the Ctx' }, { id: 'responses', title: 'Responses & errors' }, - { id: 'middleware', title: 'Middleware & groups' }, + { id: 'middleware', title: 'Middleware, CORS & guards' }, { id: 'validation', title: 'Validation' }, + { id: 'static', title: 'Static files' }, ], }, { group: 'Database', items: [ { id: 'models', title: 'Models' }, + { id: 'relations', title: 'Relations' }, { id: 'queries', title: 'The query builder' }, - { id: 'backend', title: 'Backends: SQLite & Postgres' }, + { id: 'backend', title: 'Backends & capabilities' }, + { id: 'concurrency', title: 'Locks, isolation, bulk' }, + { id: 'advisory', title: 'Advisory locks' }, + { id: 'json', title: 'JSON path queries' }, + { id: 'search', title: 'Full-text & hybrid search' }, + { id: 'vectors', title: 'Embeddings & vectors' }, + { id: 'reactive', title: 'Reactive queries' }, { id: 'migrations', title: 'Migrations' }, { id: 'kv', title: 'The key/value store' }, ], @@ -51,25 +61,35 @@ { id: 'agents', title: 'The agent surface' }, { id: 'tools', title: 'Defining tools' }, { id: 'streaming', title: 'Streaming & SSE' }, + { id: 'websockets', title: 'WebSockets' }, + { id: 'pubsub', title: 'PubSub' }, + { id: 'channels', title: 'Channels & presence' }, { id: 'queues', title: 'Queues' }, + { id: 'actors', title: 'Actors & supervision' }, ], }, { group: 'Framework services', items: [ { id: 'auth', title: 'Authentication' }, - { id: 'sessions', title: 'Sessions' }, + { id: 'sessions', title: 'Sessions & CSRF' }, { id: 'mail', title: 'Mail' }, { id: 'storage', title: 'File storage' }, { id: 'events', title: 'Event sourcing' }, { id: 'templates', title: 'Templates' }, + { id: 'collections', title: 'Collections' }, + { id: 'crypto', title: 'Crypto primitives' }, ], }, { - group: 'Architecture & deployment', + group: 'Architecture & operations', items: [ { id: 'hexagonal', title: 'Hexagonal architecture' }, { id: 'testing', title: 'Testing' }, + { id: 'repl', title: 'The REPL' }, + { id: 'internals', title: 'Inside the server' }, + { id: 'listeners', title: 'Listeners' }, + { id: 'options', title: 'Tuning & limits' }, { id: 'ops', title: 'Operational endpoints' }, { id: 'deploying', title: 'Deploying' }, { id: 'security', title: 'Security posture' }, @@ -112,19 +132,21 @@ } // --- code samples (strings so braces stay literal) --- - const cInstall = `# The default features suit most apps: derive + orm + validate + ai. + const cInstall = `# The defaults — derive + orm + validate — suit a typical app. cargo add sutegi -# In Cargo.toml, pick the pillars you want. The HTTP core is always present. +# In Cargo.toml, pick the pillars you want. The HTTP core and the agent +# tool surface are always present; there is no "ai" feature to enable. [dependencies] -# Single-node app with the SQLite backend + graceful shutdown: -sutegi = { version = "*", features = ["sqlite", "graceful"] } +# Single-node app: bundled SQLite + graceful shutdown. +sutegi = { version = "0.10", features = ["sqlite", "graceful"] } -# Multi-pod app on Postgres, with the durable queue: -# sutegi = { version = "*", features = ["postgres", "queue", "graceful"] } +# Multi-pod app: Postgres, the durable queue, cross-pod realtime. +# sutegi = { version = "0.10", features = [ +# "postgres", "queue", "channels", "pubsub-postgres", "graceful" ] } -# Nothing but the HTTP core — ~394 KB, no ORM, no agent layer: -# sutegi = { version = "*", default-features = false }`; +# Nothing but the HTTP core — ~394 KB, no ORM, no derives: +# sutegi = { version = "0.10", default-features = false }`; const cFirstApp = `use sutegi::prelude::*; @@ -138,7 +160,8 @@ fn main() -> std::io::Result<()> { }`; const cFirstRun = `cargo run -# [sutegi] hello on http://0.0.0.0:8080 +# sutegi · hello on http://0.0.0.0:8080 +# ops: /__health /__ready /__metrics /__introspect curl localhost:8080/hello/ada # -> hi, ada curl localhost:8080/__introspect # the whole app surface, as JSON`; @@ -152,6 +175,31 @@ let hosts = cfg.list("ALLOWED_HOSTS"); // comma-separated -> Vec cfg.require_all(&["DATABASE_URL", "API_KEY"])?; // fail fast, listing every missing key let db_cfg = cfg.prefixed("DB_"); // DB_HOST/DB_PORT -> HOST/PORT`; + const cEnv = `# read by .serve() itself +HOST=0.0.0.0 # bind host (default 0.0.0.0; argv[1] overrides both) +PORT=8080 # bind port (default 8080) +WORKERS=8 # HTTP worker threads (default 8) + +# read by the Postgres backend (Pg::from_env / Pool::from_env) +DATABASE_URL=postgres://user:pw@host:5432/db +# …or the standard discrete vars: +PGHOST=localhost PGPORT=5432 PGUSER=postgres PGPASSWORD=… PGDATABASE=postgres + +# read by the SQLite backend helper Db::open_or_memory("DATABASE_PATH") +DATABASE_PATH=app.db # unset -> an in-memory database + +# read by Mailer::from_env() +MAIL_DRIVER=log|memory|smtp|sendmail +MAIL_FROM=noreply@example.com +MAIL_HOST=… MAIL_PORT=… MAIL_USERNAME=… MAIL_PASSWORD=…`; + + const cCli = `sutegi new blog # scaffold an app in the conventional layout +sutegi make:model Post # src/models/post.rs (table: posts) +sutegi make:route health # src/routes/health.rs with register(app) +sutegi introspect [addr] # pretty-print a running app's /__introspect +sutegi repl 127.0.0.1:8080 # drive a running app over the agent contract +sutegi version | help`; + const cRouting = `App::new("api") .get("/", "Health check", |_| "ok") .get("/todos/:id", "Show a todo", |c| { @@ -163,6 +211,10 @@ let db_cfg = cfg.prefixed("DB_"); // DB_HOST/DB_PORT -> HOST/PORT`; }) .put("/todos/:id", "Replace a todo", |_| status(204)) .delete("/todos/:id", "Delete a todo", |_| status(204)) + // Any other verb goes through the generic .route(...): + .route(Method::Patch, "/todos/:id", "Patch a todo", |_| status(204)) + // A rest pattern captures the remainder of the path: + .get("/files/*path", "Serve a file", |c| c.param("path").unwrap_or("").to_string()) // Group a prefix + shared middleware, then register routes inside it: .group("/admin", vec![mw(require_key)], |g| { g.get("/stats", "Admin stats", |_| "…") @@ -175,16 +227,19 @@ let db_cfg = cfg.prefixed("DB_"); // DB_HOST/DB_PORT -> HOST/PORT`; let query = c.query(); // BTreeMap of ?a=b let page = c.query().get("page").cloned(); - // Headers, cookies, the raw body - let auth = c.header("authorization"); // Option<&str> - let bytes = &c.req.body; // &[u8] — the raw request body + // Headers, the raw body, the peer + let auth = c.header("authorization"); // Option<&str>, case-insensitive + let bytes = &c.req.body; // &[u8] — the raw request body + let peer = c.req.peer.as_deref(); // Option<&str> — the socket address // Bodies, parsed let body: Json = c.json()?; // application/json let form = c.form(); // application/x-www-form-urlencoded // Shared application state, registered once with .state(...) - let db = c.db::(); // the pooled DB handle + let db = c.db::(); // the pooled DB handle + let cfg = c.state::(); // panics if not registered + let opt = c.try_state::(); // Option<&Mailer> Ok(body) })`; @@ -193,15 +248,20 @@ let db_cfg = cfg.prefixed("DB_"); // DB_HOST/DB_PORT -> HOST/PORT`; |_| Json::obj(vec![("ok", Json::Bool(true))]) // 200 application/json |_| (201, some_json) // an explicit (status, body) |_| status(204) // a bare status code +|_| no_content() // 204, spelled out +|_| html(200, "

hi

") // 200 text/html |_| redirect("/login") // 302 |c| c.model::("id").map(|t| t.to_json()) // Result — ? just works // Errors carry a status, a message, and optional per-field detail. -// Returning Err(...) from a handler maps to the right HTTP response. Err(Error::not_found("no such todo")) // 404 Err(Error::unauthorized("log in first")) // 401 Err(Error::unprocessable("bad shape") // 422 with structured fields - .with_fields(errors.to_json()))`; + .with_fields(errors.to_json())) + +// 4xx messages ARE the API and render verbatim. A 5xx does not: +// its message goes to stderr for the operator and the client gets +// {"error":"internal error"} — a SQL string or a disk path never leaks.`; const cMiddleware = `// Before-middleware returns Some(Response) to short-circuit, None to continue. fn require_key(req: &Request) -> Option { @@ -212,11 +272,21 @@ fn require_key(req: &Request) -> Option { } App::new("api") - // App-wide middleware: + // App-wide middleware, in registration order: .middleware(logger()) // log every request .middleware(rate_limit(100, Duration::from_secs(60))) - .after(secure_headers()) // after-middleware rewrites the response + // After-middleware rewrites the outgoing response: + .after(secure_headers()) .after(cors("https://app.example.com")) + // A browser frontend on another origin that must send its session cookie + // needs the credentialed pair — plain cors() is not enough: + .after(cors_credentialed("https://app.example.com")) + .middleware(cors_preflight_credentialed("https://app.example.com", "GET, POST", "content-type")) + // Gate the whole /__ agent + ops surface (probes stay open): + .ops_guard(|req| match req.header("authorization") { + Some(t) if t == "Bearer ops-token" => None, + _ => Some(status(401)), + }) // Or scope middleware to a group: .group("/admin", vec![mw(require_key)], |g| { g.get("/users", "List users", |_| "…") @@ -234,11 +304,16 @@ struct Todo { // (a) Model-driven: the ruleset was generated by #[derive(Validate)]. .post("/todos", "Create", |c| { - let todo: Todo = c.validated()?; // parse + validate + hydrate, or 422 + let todo: Todo = c.validated()?; // JSON body: parse + validate + hydrate, or 422 let id = todo.save(c.db::())?; Ok::<_, Error>((201, Todo { id, ..todo }.to_json())) }) +// The same, over the other three input surfaces: +let todo: Todo = c.validated_form()?; // application/x-www-form-urlencoded +let filter: Filter = c.validated_query()?; // the query string +let key: Key = c.validated_path()?; // the path parameters + // (b) Ad-hoc: a ruleset for a shape that has no model. let rules = Ruleset::new() .field("email", &[Rule::Required, Rule::Email]) @@ -247,81 +322,305 @@ let rules = Ruleset::new() .field("password_confirmation", &[Rule::Same("password".into())]); let body = c.validate(&rules)?; // Err -> { "email": ["must be a valid email"] }`; + const cStatic = `App::new("site") + .get("/api/health", "Health check", |_| "ok") // API routes first… + .static_dir("/assets", "public/assets") // …then the file trees + .static_dir("/", "dist") // dist/index.html is the site root + .serve() + +// Routes match in registration order, so register static_dir last. +// A directory (or the bare prefix) serves its index.html; traversal, +// dotfiles and backslashes are 404s, never a read outside the root.`; + const cModels = `#[derive(Model, Validate)] #[model(table = "todos")] // omit to infer snake_case + plural struct Todo { #[model(primary)] id: i64, // the DB assigns this on insert + #[model(unique)] + slug: String, // a UNIQUE constraint in the schema + #[model(index)] + owner_id: i64, // a secondary index title: String, - done: bool, // round-trips as a real bool + #[model(default = false)] + done: bool, // DEFAULT false, so an ADD COLUMN is safe note: Option, // Option -> a nullable column + #[model(column = "created_ms")] + created: i64, // a column name that differs from the field + #[model(vector(dim = 384))] + embedding: Vec, // vector(384) on PG, TEXT on SQLite #[model(skip)] cached: bool, // not persisted; default-initialised } let db = Db::open_or_memory("DATABASE_PATH"); // pooled, Send + Sync + Clone -Todo::migrate(&db).unwrap(); +Todo::migrate(&db).unwrap(); // dev-mode additive sync let all: Vec = Todo::all_typed(&db)?; // typed reads let one: Option = Todo::find_typed(&db, 1.into())?; let count = Todo::count(&db)?; -let mut todo = Todo { id: 0, title: "ship".into(), done: false, note: None, cached: false }; let id = todo.save(&db)?; // insert; returns the new pk +Todo::update(&db, 1.into(), &[("done", true.into())])?; // by primary key +Todo::delete(&db, 1.into())?; // In a handler, route-model binding parses the param, loads the row, or 404s: .get("/todos/:id", "show", |c| c.model::("id").map(|t| t.to_json()))`; + const cRelations = `#[derive(Model)] +struct User { + #[model(primary)] id: i64, + name: String, + #[model(has_many(Post, foreign_key = "author_id"))] + posts: Vec, // not a column +} + +#[derive(Model)] +struct Post { + #[model(primary)] id: i64, + title: String, + #[model(index)] author_id: i64, + #[model(belongs_to(User, foreign_key = "author_id", on_delete = "cascade"))] + author: Option, // the FK flows into the schema IR +} + +// The derive generates one batch loader per relation: two queries, never N+1. +let users = User::all_typed(&db)?; +let users = User::with_posts(&db, users)?; // WHERE author_id IN (…) +for u in &users { println!("{} wrote {}", u.name, u.posts.len()); } + +let posts = Post::with_author(&db, Post::all_typed(&db)?)?;`; + const cQueries = `// The parameterized query builder — never string-concatenated, and identifiers // are guarded against injection (important when an AI tool arg reaches a column). let overdue = Todo::query() .filter("done", "=", false.into()) .filter("due", "<", now.into()) + .filter_in("owner_id", vec![1.into(), 2.into()]) .or_group(&[("priority", "=", "high".into()), ("pinned", "=", true.into())]) .where_not_null("assignee") .like("title", "%urgent%") - .order_by("due", false) // ascending + .join("users", "users.id", "todos.owner_id") // JOIN / LEFT JOIN + .group_by(&["users.name"]).distinct() + .where_raw("created_at > ?", vec![0.into()]) // the explicit escape hatch + .order_by("due", false) // false = ASCENDING; true = DESC .limit(20) .offset(0); -let rows: Vec = db.fetch(&overdue)?; // typed +let rows: Vec = db.fetch(&overdue)?; // typed +let one: Option = db.fetch_one(&overdue)?; let page: Page = db.paginate_typed(&overdue, 1, 20)?; // .items / .total / .has_next() +let n = db.count(&overdue)?; +let any = db.exists(&overdue)?; + +// Writes have builders too — and RETURNING, so no re-SELECT race: +UpdateBuilder::table("todos").set("done", true.into()) + .filter("id", "=", 5.into()).returning(&["id", "title"]); +DeleteBuilder::table("todos").filter("id", "=", 5.into()); -// Raw SQL is always the escape hatch: +// Raw SQL is always available: let rows = db.query("SELECT count(*) AS n FROM todos WHERE done = ?", &[true.into()])?;`; const cBackend = `// Single node — an embedded file, nothing to run. -let db = Db::open_or_memory("DATABASE_PATH"); +let db = Db::open_or_memory("DATABASE_PATH"); // or Db::memory() / Db::open_pool("app.db", 8) // Multi-pod — the SAME model code, Postgres underneath. -let pg = Pg::from_env()?; // reads the PG_* / DATABASE_URL environment +let pg = Pg::from_env(8)?; // DATABASE_URL, else PGHOST/PGPORT/…; 8 pooled conns Todo::migrate(&pg).unwrap(); let all = Todo::all_typed(&pg)?; // an identical call site // Write your domain against the trait and choose the store at boot: fn active_count(store: &impl Backend) -> Result { - Todo::query().filter("done", "=", false.into()).count_on(store) -}`; + store.count(&Todo::query().filter("done", "=", false.into())) +} + +// Transactions work on either backend; the closure gets a Backend, so the +// query builder and Model helpers run inside the transaction: +db.transact(|tx| { tx.insert("todos", &[("title", "x".into())], "id")?; Ok(()) })?;`; + + const cCaps = `// Ask the store what it can do — never find out from a dialect SQL error. +let caps = db.capabilities(); +if caps.skip_locked { /* the claim shape is safe to use */ } +if caps.json_contains { /* @> containment is available */ } +match caps.advisory_locks { + CapScope::Cluster => elect_a_leader(&db)?, // Postgres + CapScope::Process => run_locally(&db)?, // SQLite + CapScope::None => run_unlocked(&db)?, +} + +// Publish the block so an agent can read it too: +App::new("api") + .register_capabilities(db.capabilities()) // -> "capabilities" in /__introspect + .state(db) + .serve() + +// Gated features return one uniform error instead of leaking dialect SQL: +// Err("unsupported: json_contains is not available on sqlite")`; + + const cConcurrency = `// Row locks — the work-queue claim shape (Postgres; see the matrix above). +db.transact(|tx| { + let claimed = tx.select( + &QueryBuilder::table("jobs") + .filter("state", "=", "ready".into()) + .order_by("id", false) + .limit(1) + .for_update() // or .for_share() + .skip_locked(), // or .nowait() + )?; + Ok(claimed) +})?; + +// Isolation levels — a lost-update race surfaces as PG error 40001, so retry. +db.transact_with(Isolation::Serializable, |tx| { + let n = read_counter(tx)?; + tx.execute("UPDATE counters SET n = ? WHERE id = 1", &[(n + 1).into()]) +})?; + +// RETURNING on DML — the affected rows in one round-trip. +let rows = db.update_returning( + &UpdateBuilder::table("todos") + .set("done", true.into()) + .filter("id", "=", 7.into()) + .returning(&["id", "title", "done"]), +)?; + +// Bulk insert — multi-row VALUES anywhere, wire-native COPY on Postgres +// (measured 30.8x faster than row-at-a-time at 5k rows). +let n = db.insert_many("events", &["id", "kind", "payload"], &rows)?;`; + + const cAdvisory = `use std::time::Duration; + +// Try once; None = someone else holds it. +if let Some(_guard) = db.try_lock("nightly-report")? { + run_report(&db)?; +} // dropping the guard releases + +// Wait up to 5s, then give up. +let guard = db.lock("reindex", Duration::from_secs(5))?; + +// The singleton-job shape: at most one pod runs f. Ok(None) = another pod did. +db.with_lock("janitor", Duration::ZERO, || sweep(&db))?; + +// Leader election — hold it for the process lifetime; followers retry. +std::thread::spawn(move || loop { + if let Ok(Some(_leader)) = db.try_lock("scheduler-leader") { + run_scheduler(&db); // returns only if the scheduler stops + } + std::thread::sleep(Duration::from_secs(5)); +}); + +// An operator can take the same lock from psql: +// SELECT pg_try_advisory_lock();`; + + const cJson = `// WHERE inside the document, typed comparison: +let hot = db.select( + &QueryBuilder::table("docs") + .where_json("meta", "$.stats.views", ">", 50.into()), +)?; + +// Project a path out as a column: +let rows = db.select( + &QueryBuilder::table("docs") + .select(&["id"]) + .select_json("meta", "$.author.name", "author"), +)?; + +// Containment — Postgres only (capabilities().json_contains): +let posts = db.select( + &QueryBuilder::table("docs") + .where_json_contains("meta", Json::obj(vec![("kind", Json::str("post"))])), +)?; + +// The grammar is $.key.nested[0].deeper — identifier keys and [n] indexes. +// A malformed path is a builder error, and the compiled path is always +// BOUND as a parameter: SQLite json_extract(col, ?), Postgres col #>> ?.`; + + const cSearch = `use sutegi::orm::search; + +search::setup(&db, "docs", "id", &["title", "body"])?; // idempotent DDL + +// Lexical, ranked best-first, with "_rank" attached to each row: +let hits = search::search(&db, "docs", "id", &["title", "body"], + "rust \\"job queue\\" -django", 20)?; + +// Hybrid: the lexical leg + a vector leg, fused with reciprocal-rank fusion +// (sum of 1/(60+rank)); "_score" attached. One code path on both engines. +let hits = search::hybrid_search(&db, "docs", "id", &["title", "body"], + "rust queue", "embedding", &query_vec, 10)?; + +// Tell agents what is searchable, without source access: +App::new("api").register_search("docs", &["title", "body"])`; + + const cVectors = `use sutegi::orm::embedding::{self, Metric}; + +#[derive(Model)] +struct Doc { + #[model(primary)] id: i64, + body: String, + #[model(vector(dim = 384))] embedding: Vec, +} + +// Portable brute force — correct on every backend, ideal for SQLite: +let hits: Vec<(Doc, f32)> = embedding::nearest_typed::( + &db, &Doc::query(), "embedding", &query_vec, 10, Metric::Cosine, +)?; + +// Pushdown — ORDER BY col <=> ? LIMIT k, straight to pgvector's ANN index. +// Requires capabilities().vector (Postgres + the pgvector extension). +let hits = embedding::nearest_pushdown_typed::( + &pg, &Doc::query(), "embedding", &query_vec, 10, Metric::Cosine, +)?; + +// Values travel in pgvector's canonical [1,2,3] text form, so the same row +// round-trips identically on either backend. Lower distance = closer.`; + + const cReactive = `use sutegi::orm::watch::Watcher; + +let watcher = Watcher::postgres(&pg)?; // one per process (Watcher::sqlite(&db)) +let sub = watcher.watch( + Todo::query().filter("done", "=", false.into()), + "id", // the pk the diff keys on +)?; + +for row in sub.rows() { /* the result at watch time */ } +while let Some(change) = sub.recv_timeout(Duration::from_secs(30)) { + // Change { table, added, updated, removed } — only when the watched + // result actually moved. to_json() is broadcast-ready: + hub.broadcast("todos:lobby", "changed", &change.to_json()); +} + +// Postgres: watch() idempotently installs a statement-level _sutegi_watch_ +// trigger that pg_notify's a shared channel, received on a dedicated LISTEN +// session — so ANY pod's committed write (psql included) wakes every pod's +// watchers, and a rolled-back write never fires. SQLite: update_hook on every +// pooled connection, this process only.`; const cMigrations = `use sutegi::prelude::*; fn migrations() -> Migrator { - Migrator::new().add(Migration::reversible( - "20260701_000001", // ordered id - "create_todos", - |db| db.migrate_schema(&Todo::schema()), // up - |db| db.execute("DROP TABLE todos", &[]).map(|_| ()), // down - )) + Migrator::new().load_dir(sutegi::migrate::MIGRATIONS_DIR).expect("load migrations") } fn main() -> std::io::Result<()> { let db = Db::open_or_memory("DATABASE_PATH"); - // \`myapp migrate | migrate:status | migrate:rollback\` runs and exits: - if sutegi::migrate::dispatch(&migrations(), &db) { return Ok(()); } - migrations().run(&db).expect("migrate"); // otherwise apply pending, then serve + let models = sutegi::schemas![Todo, User]; // the desired state + // migrate | :rollback | :status | :gen | :plan | :drift | :fresh + if sutegi::migrate::dispatch_full(&migrations(), &db, &models, + sutegi::migrate::MIGRATIONS_DIR) { + return Ok(()); // a subcommand ran; exit + } + migrations().run(&db).expect("migrate"); // else apply pending, then serve App::new("todo").state(db).serve() }`; + const cMigrateCli = `myapp migrate:gen create_todos # diff models <-> shadow schema; write the file +myapp migrate:plan # …show what :gen would write, without writing +myapp migrate # apply all pending migrations +myapp migrate:status # the ledger: check applied, ? orphan +myapp migrate:rollback 1 # roll back the last n batches (default 1) +myapp migrate:drift # three-way report: DB vs migrations vs models +myapp migrate:fresh # roll everything back and re-run (dev only)`; + const cKv = `use sutegi::orm::kv::Kv; let kv = Kv::new(db); // over SQLite *or* Postgres — same API @@ -330,7 +629,10 @@ kv.migrate()?; kv.set("config", "theme", &Json::str("dark"))?; // namespace, key, value let theme = kv.get("config", "theme")?; // Option let all = kv.scan("flags")?; // Vec<(String, Json)> +let some = kv.scan_prefix("flags", "beta_")?; +let n = kv.count("flags")?; kv.delete("config", "theme")?; +kv.clear("flags")?; // It is ordinary app state — no Arc>: App::new("settings") @@ -344,10 +646,19 @@ App::new("settings") .serve()`; const cAgents = `curl localhost:8080/__introspect -# { "framework": "sutegi", "name": "...", "routes": [...], "models": [...], "tools": [...] } +# { +# "framework": "sutegi", "version": "0.10.0", "name": "todo", +# "routes": [ { "method": "GET", "pattern": "/todos/:id", "doc": "…" } ], +# "models": [ … ], "tools": [ … ], +# "capabilities": { "backend": "postgres", "skip_locked": true, … }, +# "search": [ { "table": "docs", "columns": ["title","body"] } ], +# "listeners": [ { "name": "statsd", "doc": "Ingests statsd on udp/8125." } ], +# "endpoints": { "introspect": "/__introspect", "health": "/__health", +# "ready": "/__ready", "metrics": "/__metrics" } +# } curl localhost:8080/__tools -# [ { "name": "create_todo", "description": "...", "input_schema": {...}, "streaming": false } ] +# [ { "name": "create_todo", "description": "…", "input_schema": {…}, "streaming": false } ] curl -X POST localhost:8080/__tools/create_todo -d '{"title":"ship sutegi"}' # args are validated against the tool's schema -> 422 on a bad shape`; @@ -370,31 +681,186 @@ curl -X POST localhost:8080/__tools/create_todo -d '{"title":"ship sutegi"}' for token in prompt.split(' ') { sink.data(token)?; } sink.event("done", "{}") }) - .serve()`; + .serve() + +// POST /__tools/create_todo -> invoke (422 on a bad shape) +// POST /__tools/stream_answer/stream -> the SSE variant`; const cStreaming = `.get("/stream", "SSE demo", |_| sse(|sink| { for token in answer().split(' ') { sink.data(token)?; // each frame is flushed immediately } + sink.comment("keep-alive")?; // a ':'-prefixed no-op frame sink.event("done", "{}") // a named event })) -// Raw byte streams work the same way via stream(status, content_type, producer).`; - - const cQueues = `use sutegi::queue::{Queue, Workers}; +// Raw byte streams — NDJSON, large exports — take the same shape: +.get("/export", "Stream rows", |_| stream(200, "application/x-ndjson", |sink| { + for row in rows() { sink.write_str(&format!("{}\\n", row.to_json()))?; } + Ok(()) +})) -let mut queue = Queue::new(Pg::from_env()?.pool().clone()); -queue.migrate()?; // creates the sutegi_jobs table -queue.register("notify", |payload| { - let to = payload.get("to").and_then(Json::as_str).unwrap_or(""); - /* send the email … */ Ok(()) // Err -> retried with backoff +// Regular responses use keep-alive; streams are close-framed by design.`; + + const cWs = `use sutegi::prelude::*; + +App::new("live") + // Tune the engine BEFORE the first .ws(...) — that call starts the reactor. + .ws_config(WsConfig { shards: 4, ping_interval: Duration::from_secs(20), + ..WsConfig::default() }) + .ws("/socket", "A raw WebSocket endpoint.", Ws::new() + // A cookie-authenticated socket MUST pin its origins (CSWSH): + .check_origin(["https://app.example.com"]) + .authorize(|req| req.header("authorization").is_some()) + .on_open(|conn, _req| conn.send_text("welcome")) + .on_message(|conn, msg| match msg { + Msg::Text(t) => conn.send_text(&t.to_uppercase()), + Msg::Binary(b) => conn.send_binary(&b), + }) + .on_close(|conn, code| println!("{} left ({code})", conn.id()))) + .serve() + +// Broadcast: encode the frame ONCE and share the Arc across every queue. +let frame = text_frame("maintenance at noon"); +for conn in &conns { conn.send_shared(&frame); }`; + + const cPubsub = `use sutegi::prelude::*; // Broker, BrokerExt, PubSub (+ PgPubSub) + +// Single pod: an in-process broker. Clone-cheap — it shares one registry. +let broker = PubSub::new(); + +// Many pods: the same Broker trait over PG LISTEN/NOTIFY. Nothing else changes. +// let broker = PgPubSub::connect(&pg_config)?; // sutegi_pg::Config + +// Subscribe with a closure (BrokerExt::on); the id is for unsubscribe. +let id = broker.on("orders", |msg: &str| { handle(msg); }); +broker.publish("orders", &Json::obj(vec![("id", Json::Int(42))]).to_string()); +broker.unsubscribe("orders", id); + +// Local delivery is synchronous and in subscription order. Cross-pod, +// PgPubSub echoes back its own instance id and skips it — so a publish is +// delivered exactly once locally and once per other pod. +// NOTIFY payloads cap at ~8 KB: try_publish returns the error instead of +// dropping it silently. Ship ids, not blobs; delivery is at-most-once.`; + + const cChannels = `use sutegi::prelude::*; + +let hub = Channels::new() + .channel( + Channel::new("room:*") + .doc("A chat room. Join with a nick; messages fan out to the room.") + .join_schema("A display name.", schema::object( + vec![("nick", schema::string("Display name"))], &["nick"])) + .on_join(|socket, payload| { + let nick = payload.pointer("/nick").and_then(Json::as_str) + .ok_or_else(|| Json::str("nick required"))?; + socket.assign("nick", Json::str(nick)); + Presence::track(socket, nick, Json::obj(vec![])); // feature "presence" + Ok(Json::Null) // rides the ok reply + }) + .on("new_msg", |socket, payload| { + socket.broadcast("new_msg", payload); // all members, all pods + Reply::None + }) + .on_leave(|socket, _reason| { + socket.broadcast_from("left", &Json::obj(vec![])); + }), + ) + // .broker(PgPubSub::connect(&pg_cfg)?) // <- the only cross-pod change + .check_origin(["https://app.example.com"]) // MUST set if cookies auth the socket + .build(); + +App::new("chat") + .channels("/channels", "The chat socket.", hub.clone()) // + the /__channels manifest + .serve()?; + +// From anywhere — an HTTP handler, a background thread, the REPL: +hub.broadcast("room:1", "announcement", &Json::str("maintenance at noon"));`; + + const cChannelsWire = `// One JSON object per text frame — an object, not a positional array, so +// /__channels alone teaches an agent the protocol: +{"topic":"room:1","event":"new_msg","ref":"3","join_ref":"1","payload":{"body":"hi"}} + +// Control events are stg:-prefixed and reserved: +// stg:join stg:leave stg:reply {status:"ok"|"error", response} stg:error stg:close +// Heartbeats are a ref'd push on topic "stg", event "heartbeat".`; + + const cChannelsJs = `// sutegi_channels::JS_CLIENT is a bundled ~4 KB dependency-free client: +// auto-reconnect with capped backoff, automatic rejoin, heartbeat liveness. +const socket = new SutegiSocket("/channels"); +socket.connect(); +const room = socket.channel("room:1", {nick: "ada"}); +room.on("new_msg", p => render(p)); +room.join().receive("ok", () => {}); +room.push("new_msg", {body: "hello"});`; + + const cQueues = `use std::sync::Arc; +use sutegi::queue::Queue; + +let mut queue = Queue::new(db.clone()); // ANY Backend: Db or Pg +queue.migrate()?; // creates sutegi_jobs +queue.register("notify", |job| { + let to = job.payload().get("to").and_then(Json::as_str).unwrap_or(""); + if job.is_last_attempt() { /* write a user-visible failure */ } + job.heartbeat()?; // push the visibility window forward + Ok(()) // Err -> retried with backoff }); -// Enqueue from a handler; return immediately. +// Enqueue from a handler and return immediately: queue.dispatch("notify", Json::obj(vec![("to", Json::str("a@b.com"))]))?; -// Start N worker threads. Any pod can claim the next job (FOR UPDATE SKIP LOCKED). -let workers: Workers = std::sync::Arc::new(queue).start(4);`; +// …or shape the dispatch: its own pool, one in flight per key, 3 tries. +queue.job("video.ingest", Json::obj(vec![("id", Json::str("abc"))])) + .queue("video") + .unique("yt:abc") // hands back the live row's id instead of a duplicate + .priority(10) // runs ahead of older work in the same queue + .max_attempts(3) + .delay(Duration::from_secs(30)) + .dispatch()?; + +let queue = Arc::new(queue); +let fast = Arc::clone(&queue).start(4); // 4 workers on "default" +let slow = Arc::clone(&queue).start_on("video", 1); // 1 on the slow queue +// … later: fast.stop(); slow.stop(); + +// Ops: +queue.failed(20)?; // the dead letters +queue.retry(job_id, 3)?; // revive one +queue.purge_failed(Duration::ZERO)?; +queue.stats()?; queue.stats_for("video")?; +queue.cross_pod(); // true on Postgres — the honest answer`; + + const cActors = `use sutegi::prelude::*; + +struct Counter { n: u64 } +enum Msg { Bump, Get(ReplyTo) } + +impl Actor for Counter { + type Msg = Msg; + fn handle(&mut self, msg: Msg) { + match msg { + Msg::Bump => self.n += 1, + Msg::Get(reply) => reply.reply(self.n), + } + } +} + +let counter = spawn(Counter { n: 0 }); +counter.tell(Msg::Bump)?; // cast; TellError::Full = backpressure +let n = counter.ask(Msg::Get, Duration::from_secs(1))?; // call + +// OTP-style supervision: the factory re-runs on restart, so a crashed child +// comes back as a FRESH value — never with the state it died holding. +let sup = Supervisor::new("pipeline") + .strategy(Strategy::OneForOne) // RestForOne | OneForAll + .intensity(3, Duration::from_secs(5)) // exceeded -> stop all, Failed + .child(ChildSpec::new("worker-1", || Worker::new())) + .child(ChildSpec::new("flaky-api", || ApiClient::new()) + .restart(Restart::Transient) // Permanent | Transient | Temporary + .backoff(Duration::from_millis(250))) + .start(); + +App::new("myapp").actors(sup.registry()).serve()?; // GET /__actors`; const cAuth = `use std::sync::Arc; @@ -403,26 +869,36 @@ users.migrate()?; let tokens = Arc::new(Tokens::new(db.clone())); tokens.migrate()?; -let auth = Arc::new(Auth::new( - users, - Sessions::new(secret.as_bytes()), // .insecure() drops Secure for local http:// -)); +let auth = Arc::new( + Auth::new(users, Sessions::new(secret.as_bytes())) // .insecure() for local http:// + .remember(Remember::new(db.clone())), // selector/validator cookies +); +let throttle = Throttle::new(db.clone()); // 5 attempts / 60s by default App::new("app") .state(auth.clone()) - .post("/register", "Sign up", move |c| { - let body = c.json()?; - let user = auth_reg.users.register(email, password, name).map_err(Error::unprocessable)?; - Ok::<_, Error>(auth_reg.login(c.req, &user, json(201, &user.to_json()))) + .post("/login", "Log in", move |c| { + let key = format!("login:{}|{}", email, c.req.peer.as_deref().unwrap_or("")); + if let Some(retry_after) = throttle.too_many(&key)? { + return Err(Error::new(429, "too many attempts").with_fields( + Json::obj(vec![("retry_after", Json::Int(retry_after))]))); + } + match auth.users.authenticate(email, password)? { + Some(u) => { throttle.clear(&key)?; + Ok::<_, Error>(auth.login_remembered(c.req, &u, remember, resp)) } + None => { throttle.hit(&key)?; Err(Error::unauthorized("bad credentials")) } + } }) - .post("/login", "Log in", move |c| { /* authenticate -> auth.login(...) */ }) - .get("/me", "Current user", move |c| match auth_me.current(c.req)? { - Some(u) => Ok::<_, Error>(json(200, &u.to_json())), + // identify() = session-or-remember revival; attach() sets the fresh cookies. + .get("/me", "Current user", move |c| match auth_me.identify(c.req)? { + Some(id) => Ok::<_, Error>(id.attach(json(200, &id.user.to_json()))), None => Err(Error::unauthorized("unauthenticated")), }) // Guards are just middleware: - .group("/admin", vec![mw(require_role(auth.clone(), "admin"))], |g| { /* … */ }) - .group("/api", vec![mw(require_token(tokens.clone()))], |g| { /* stg_ bearer tokens */ }) + .group("/admin", vec![mw(require_role(auth.clone(), "admin")), + mw(require_verified(auth.clone())), + mw(require_csrf(auth.clone()))], |g| { /* … */ }) + .group("/api", vec![mw(require_token(tokens.clone()))], |g| { /* stg_ bearer tokens */ }) .serve()`; const cSessions = `// Signed-cookie sessions (HMAC-SHA256). No server-side store needed. @@ -432,37 +908,67 @@ let sessions = Sessions::new(secret.as_bytes()); let mut s = sessions.load(c.req); s.set("last_item", Json::str("sku-42")); Ok::<_, Error>(sessions.save(&s, json(200, &Json::obj(vec![("ok", Json::Bool(true))])))) -})`; +}) + +// CSRF lives inside the signed session — get-or-mint, then verify in constant time. +let mut s = sessions.load(c.req); +let token = sessions.csrf(&mut s)?; // 32 random bytes, stable per session +let ok = sessions.verify_csrf(&s, presented); +// Auth::csrf(req, resp) is the handler shape; the require_csrf guard enforces +// X-CSRF-Token on mutating methods (419 on mismatch) and passes reads and +// Authorization-header callers — a bearer client carries no ambient credential. + +// For callers that collect cookies before touching a response: +let set_cookie = sessions.cookie_for(&s);`; const cMail = `// Configure once from the environment (MAIL_* vars pick the driver: -// Log for dev, SMTP/Sendmail for real delivery). Drivers implement one Transport method. +// log for dev, memory for tests, smtp/sendmail for real delivery). let mailer = Mailer::from_env()?; let email = Email::new() .to("ada@example.com") .subject("Welcome") .text("Thanks for signing up.") - .html("

Thanks for signing up.

"); + .html("

Thanks for signing up.

"); // both set -> multipart/alternative mailer.send(email)?; -// Or a themed, Laravel-notification-style message — HTML card + text from the same blocks: -let msg = MailMessage::new() +// Or a themed, notification-style message — HTML card + text from the same blocks: +let theme = Theme::new("Acme") + .brand_color("#7c3aed") + .logo_url("https://cdn.example.com/logo.png") + .footer("Acme Inc · Bilbao"); + +mailer.send(theme.message() + .subject("Welcome!") .greeting("Hi Ada,") .line("Your account is ready.") - .action("Verify email", "https://app.example.com/verify?token=…") - .line("If you didn't sign up, ignore this email."); -mailer.send(msg.build("Acme"))?;`; + .action("Verify email", &url) // brand-colored button + fallback link + .note("This link is valid for 24 hours.") + .email()? + .to("ada@example.com"))?; + +// A hosted provider is one method — Transport hands you the structured Email +// AND its rendered RFC 2822 form; post whichever your API wants: +struct Resend { key: String } +impl Transport for Resend { + fn send(&self, email: &Email, _raw: &str, id: &str) -> Result { + my_http_post("https://api.resend.com/emails", &self.key, &to_json(email))?; + Ok(id.to_string()) + } +}`; const cStorage = `// One trait, three backends. Swap the type you construct, not the call sites. -let store = FsStorage::new("files")?; // single-node, one directory -let store = S3Store::r2(&acct, "media", &ak, &sk) // …or a Cloudflare R2 bucket +let store = FsStorage::new("data/files")?; // single-node, one directory +// let store = DbStorage::new(pg); // blobs over the Backend seam +let store = S3Store::r2(&account, "media", &ak, &sk) // …or a Cloudflare R2 bucket .storage(SystemCurl::new()); -.put("/files/:name", "Upload", |c| { - let ct = c.header("content-type").unwrap_or(""); - c.state::().put(c.param("name").unwrap(), &c.req.body, ct)?; - Ok::<_, Error>(status(201)) -}) +store.put("reports/q2.pdf", &bytes, "application/pdf")?; +let meta: Option = store.stat("reports/q2.pdf")?; +let objects = store.list("reports/")?; +let reader = store.get_reader("reports/q2.pdf")?; // real streaming reads +store.delete("reports/q2.pdf")?; + .get("/files/:name", "Download", |c| -> Result { let store = c.state::(); match store.stat(c.param("name").unwrap())? { @@ -473,8 +979,8 @@ let store = S3Store::r2(&acct, "media", &ak, &sk) // …or a Cloudfla } }) -// Agent-native S3: mint a time-limited URL and let the agent move the bytes itself. -let s3 = S3Store::new(&bucket, ®ion, &access, &secret); // .with_endpoint(...) for R2/MinIO +// Agent-native S3: mint a time-limited URL and let the agent move the bytes. +let s3 = S3Store::new("bucket", "eu-central-1", &ak, &sk); // .with_endpoint(…) for R2/MinIO let url = s3.presign_put("reports/q2.pdf", 900)?; // seconds`; const cEvents = `use sutegi::events::{event, Aggregate, EventStore, Expected, Projections, StoredEvent}; @@ -496,9 +1002,10 @@ impl Aggregate for Account { let store = EventStore::new(db.clone()); store.migrate()?; -// Append with optimistic concurrency; fold state back on demand. +// Append with optimistic concurrency (Any | NoStream | Version(n)): store.append("account-42", Expected::Any, &[event("deposited", amount_payload(100))])?; let (account, version) = store.load::("account-42")?; // balance = 100 +// append_tx composes with a transaction you already own. // A checkpointed projection maintains a read model, exactly once, rebuildable: let mut projections = Projections::new(db.clone()); @@ -507,21 +1014,61 @@ projections.register("account_balances", |e, tx| { }); let _workers = std::sync::Arc::new(projections).start();`; - const cTemplates = `// A Blade-lite engine over Json contexts. {{ }} escapes; {!! !!} is raw. -let src = "\\ -

Hi {{ name }}

-@if(admin) -

You are an admin.

-@endif -
    @foreach(items as item)
  • {{ item }}
  • @endforeach
"; - -let mut views = Templates::new(); -views.register("home", src); -let html = views.render("home", &Json::obj(vec![ - ("name", Json::str("Ada")), - ("admin", Json::Bool(true)), - ("items", Json::arr(vec![Json::str("a"), Json::str("b")])), -]))?;`; + const cTemplates = `let mut views = Templates::new(); +views.add("row", "
  • {{ item.name }}@if(item.admin) *@endif
  • ")?; +views.add("list", "
      @foreach(users as item)@include(row)@endforeach
    ")?; + +let html = views.render("list", &Json::obj(vec![ + ("users", Json::arr(vec![ + Json::obj(vec![("name", Json::str("Ada")), ("admin", Json::Bool(true))]), + ])), +]))?; + +// {{ }} escapes, {!! !!} does not; dot paths reach into nested objects. +// @if / @else, @foreach … as … (with loop.index / loop.first / loop.last), +// @include for partials. Templates compile once to an AST and report +// line-numbered errors. The mail layer's themed HTML renders through this.`; + + const cCollections = `use sutegi::collect; + +let report = collect(orders) + .filter(|o| o.paid) + .group_by(|o| o.country.clone()) // HashMap> + .into_iter() + .map(|(country, os)| format!("{country}: {}", os.sum_by(|o| o.total))) + .collect::>(); + +// Numeric chains read left-to-right: +let total: i64 = collect(vec![1, 2, 3, 4]).filter(|n| n % 2 == 0).map(|n| n * 10).sum(); + +// filter/reject, map/filter_map/flat_map, pluck, partition, chunk, unique, +// sort/sort_by/sort_by_key, take/skip, reduce, sum_by, implode, each, +// tap/pipe. It Derefs to [T] and round-trips through Vec, so it costs +// nothing over doing the work by hand.`; + + const cCrypto = `use sutegi::crypto; + +// Hashing & MACs — the primitives the rest of the framework is built on. +let digest = crypto::sha256(b"payload"); +let mut h = crypto::Sha256::new(); // incremental: constant memory +h.update(&chunk); let digest = h.finalize(); // S3 bodies never live twice +let mac = crypto::hmac_sha256(secret, message); +let dk = crypto::pbkdf2_hmac_sha256(pw, &salt, 600_000); + +// Per-purpose subkeys, so a signing key and an encryption key never alias: +let enc_key = crypto::hkdf_sha256(master, b"", b"session-enc", 32); + +// Two-way encryption: ChaCha20-Poly1305 (RFC 8439). seal() prepends a fresh +// random nonce (nonce || ct || tag, +28 bytes), so nonce reuse is impossible +// by construction. ChaCha over AES deliberately: add-rotate-xor is +// constant-time in plain software, where AES table lookups leak via cache. +let sealed = crypto::seal(&key, b"secret")?; +let plain = crypto::open(&key, &sealed); // Option> + +// Utilities: constant_time_eq (piped through black_box), random_bytes (OS +// entropy), hex/from_hex, base64_encode/decode, now_secs/now_millis. +// Known-answer tested against the RFC vectors, with deterministic fuzz +// coverage for round-trip, bit-flip, AAD-mismatch and truncation.`; const cHex = `// Domain use case, written against a port trait — no HTTP, no SQL in sight. impl UseCase for CreateTodo { @@ -534,7 +1081,7 @@ impl UseCase for CreateTodo { } } -// Inbound HTTP adapter — respond_created maps AppResult to the right HTTP response: +// Inbound HTTP adapter — respond_created maps AppResult to the right response: .post("/todos", "Create", move |c| { let title = c.json()?.get("title").and_then(Json::as_str).unwrap_or("").to_string(); respond_created(create.execute(title)) @@ -556,21 +1103,128 @@ fn creates_a_todo() { let resp = handle(Request::post("/todos", br#"{"title":"x"}"#)); assert_eq!(resp.status, 201); } -// See crates/sutegi/tests/server.rs for the full end-to-end suite that boots a real server.`; +// The tool surface is reachable the same way — POST /__tools/:name through +// service(), no server, no agent. For end-to-end coverage the framework's own +// suite boots a real server over a loopback socket: crates/sutegi/tests/server.rs.`; + + const cRepl = `// In-process: consume the built App; data commands light up via .db(...). +Repl::new(app).db(db).run()?; + +// Or against a RUNNING app with no source access — the agent contract, +// driven by a human: +// sutegi repl 127.0.0.1:8080 + +sutegi> routes +GET /api/todos List todos +sutegi> tools +create_todo [unary ] Create a todo +sutegi> call create_todo {"title":"ship"} +{ "id": 1, "title": "ship", "done": false } +sutegi> q todos where done = false order id desc limit 5 +sutegi> sql SELECT count(*) FROM todos +sutegi> kv scan flags +sutegi> events account-42 10 +sutegi> jobs + +// Line editing is plain stdin (zero deps); wrap with rlwrap for history.`; + + const cInternals = `accept() -> a fixed thread pool (WORKERS, default 8) + | + | one connection per worker thread, blocking I/O + v + parse request line + headers (bounded: max_header_bytes, header_timeout) + | + v + ops guard -> global middleware -> route match -> group middleware + | | + | v + | handler(&Ctx) + | | + v v + after-middleware <--------------------------- IntoResponse + | + v + write response -> keep-alive (keep_alive_idle / keep_alive_max) or close + -> or detach: Body::Upgrade hands the socket to the + ws reactor and frees the worker immediately + +// A handler panic is caught per request and becomes a 500 — one bad request +// never takes a worker down. (In release the workspace builds with +// panic = "abort", so the pod supervisor owns that failure instead.) +// /__metrics counts requests total, in-flight, and by status class.`; + + const cListeners = `use sutegi::prelude::*; +use std::net::UdpSocket; +use std::time::Duration; + +App::new("metrics-demo") + .state(db) + .listener("statsd", "Ingests statsd counters on udp/8125.", |ctx| { + let sock = UdpSocket::bind("0.0.0.0:8125").unwrap(); + // Bounded read + a should_stop() check each lap IS the shutdown + // contract: serve() joins listener threads before returning. + sock.set_read_timeout(Some(Duration::from_millis(250))).unwrap(); + let mut buf = [0u8; 1500]; + while !ctx.should_stop() { + if let Ok((n, _from)) = sock.recv_from(&mut buf) { + ingest(ctx.db::(), &buf[..n]); // the same state handlers see + } + } + }) + .serve() + +// The closure runs once, on a thread named sutegi-listener-statsd, started by +// run / run_until / run_graceful / serve. App::service() never spawns them, so +// the in-process request closure stays socket-free for tests and benches. +// GET /__introspect gains: "listeners": [ { "name": …, "doc": … } ]`; + + const cOptions = `use std::time::Duration; + +App::new("api") + .workers(16) // HTTP threads (WORKERS env wins) + .max_body(8 * 1024 * 1024) // 413 above this (default 2 MiB) + .request_timeout(Some(Duration::from_secs(60))) // per-socket read/write + .limits(Limits { // …or replace the whole set + max_body: 2 * 1024 * 1024, // default 2 MiB + max_header_bytes: 64 * 1024, // default 64 KiB + timeout: Some(Duration::from_secs(30)), + header_timeout: Duration::from_secs(15), // whole-request deadline + keep_alive_idle: Duration::from_secs(5), // idle keep-alive pins a thread + keep_alive_max: 100, // requests per connection + }) + .ws_config(WsConfig { // the realtime engine, if enabled + shards: 0, // 0 = one per core + max_frame: 1 << 20, + max_message: 1 << 20, + ping_interval: Duration::from_secs(30), + idle_timeout: Duration::from_secs(75), // must exceed ping_interval + max_buffered: 1 << 20, // slow consumers get dropped + max_connections: 1 << 20, + max_connections_per_ip: 1024, + raise_nofile: true, // lift RLIMIT_NOFILE at start + }) + .serve()`; const cOps = `let ready = db.clone(); // the probe keeps its own pooled handle App::new("api") .state(db) .readiness(move || ready.query("SELECT 1", &[]).is_ok()) - .serve()?; // reads HOST/PORT/WORKERS; graceful SIGTERM drain + .register_capabilities(caps) + .actors(sup.registry()) // GET /__actors (feature "actors") + .get("/__migrations", "Migration status + drift.", + move |_| json(200, &report)) // your own /__-mounted route + .ops_guard(|req| gate(req)) // gates every /__ EXCEPT the probes + .serve()?; // HOST/PORT/WORKERS; SIGTERM drain // Always on, no feature required: // GET /__health liveness (200 while up) GET /__ready readiness (200/503) -// GET /__metrics Prometheus text GET /__introspect the full app surface`; +// GET /__metrics Prometheus text GET /__introspect the full app surface +// GET /__tools the LLM manifest GET /__channels the channel manifest`; const cDeploy = `./ontzi up 3 # 3 replicas behind an nginx LB on http://localhost:8080 ./ontzi curl /api/todos ./ontzi logs +./ontzi down ./ontzi k8s apply # promote deploy/k8s/ — probes, drain, Prometheus annotations wired`; @@ -588,6 +1242,7 @@ App::new("api")
    + @@ -636,13 +1291,15 @@ App::new("api") framework for Rust with a single unusual constraint: it has zero third-party runtime dependencies. No tokio, no serde, no hyper, not even a Postgres driver crate. The HTTP/1.1 parser, the JSON codec, the router, the ORM, the Postgres - wire driver, and the agent tool layer are all original code, built on the standard library. + wire protocol, the WebSocket reactor, the crypto primitives and the agent tool layer are all original + code, built on the standard library.

    - If you have used Laravel, the shape will feel familiar: an expressive router, an Eloquent-style ORM, - first-class validation, a queue, mail, sessions, storage, and a CLI that scaffolds the conventional - pieces. The difference is what sits underneath — nothing you did not choose to compile in — and one - addition Laravel never had to think about: + If you have used Laravel, the shape will feel familiar: an expressive router, an Eloquent-style ORM with + relations and migrations, first-class validation, a durable queue, mail, sessions, storage, and a CLI that + scaffolds the conventional pieces. If you have used Phoenix, the realtime half will: channels, presence, + and pushed live queries. The difference is what sits underneath — nothing you did not choose to compile in + — and one addition neither had to think about: an AI agent is a first-class user of your app, able to discover and drive every route, model, and tool over plain JSON with no SDK.

    @@ -652,12 +1309,12 @@ App::new("api") They are written to be read in order the first time and grepped after. The path is deliberate:

      -
    • Getting started — install, boot a first app, and learn how features and configuration work.
    • +
    • Getting started — install, boot a first app, and learn how features, configuration and the CLI work.
    • The basics — routing, the request context, responses, middleware, and validation: the everyday request loop.
    • -
    • Database — models, the query builder, the one Backend trait behind SQLite and Postgres, migrations, and the KV store.
    • -
    • Agents & realtime — the introspection surface, tools, streaming, and the durable queue.
    • -
    • Framework services — auth, sessions, mail, storage, event sourcing, and templates: opt-in pillars you reach for as needed.
    • -
    • Architecture & deployment — hexagonal structure, testing, the operational endpoints, deploying, and an honest security posture.
    • +
    • Database — models and relations, the query builder, the one Backend trait behind SQLite and Postgres, then the deeper seams: locks, JSON paths, search, embeddings, live queries, migrations, KV.
    • +
    • Agents & realtime — the introspection surface, tools, streaming, WebSockets, pubsub, channels, the durable queue, and actors.
    • +
    • Framework services — auth, sessions, mail, storage, event sourcing, templates, collections, crypto: opt-in pillars you reach for as needed.
    • +
    • Architecture & operations — hexagonal structure, testing, the REPL, what the server actually does per request, every tuning knob, the operational endpoints, deploying, and an honest security posture.
    Start here
    @@ -676,20 +1333,28 @@ App::new("api") Frameworks usually force a trade: batteries-included but heavy, or tiny but bare. sutegi bets you can refuse the trade if you build every layer on std and make each one an opt-in compile-time feature. Compile in only the HTTP core and you get a ~394 KB binary with - no async runtime; switch on sqlite, ai and graceful and you get an - ergonomic, agent-native service — with nothing else along for the ride. + no async runtime; switch on sqlite and graceful and you get an ergonomic, + agent-native service — with nothing else along for the ride. The full todo example, every + pillar plus bundled SQLite, is ~1.31 MB.

    The second bet is agent-native design. Because every route, model, and tool registers its own metadata, the framework assembles a complete, machine-readable description of your app for free. An LLM points at - /__introspect, reads the surface, and calls /__tools — the same application you - built for humans is drivable by a model without a line of glue. + /__introspect, reads the surface — including what the store can do — and calls + /__tools. The same application you built for humans is drivable by a model without a line of + glue.

    The third principle is that the type you hold is the only thing that changes - when you scale. Handlers, models, and validation are written once against traits; moving from a - single SQLite file to a fleet of Postgres-backed pods swaps a constructor, not your code. The cost of - these bets is a large hand-rolled surface and a young ecosystem — which the next page addresses head-on. + when you scale. Handlers, models, validation, the queue, the event store, file storage and + channel broadcasts are written once against traits; moving from a single SQLite file to a fleet of + Postgres-backed pods swaps a constructor, not your code. Where the engines genuinely differ, a + capability bit says so out loud and the gated call fails with a + named error instead of a dialect SQL surprise. +

    +

    + The cost of these bets is a large hand-rolled surface and a young ecosystem — which the next page + addresses head-on.

    @@ -704,12 +1369,20 @@ App::new("api")

    The core is small enough to read in an afternoon and is exercised hard. Beyond the unit suite, a deterministic, pure-std fuzz and differential harness - hammers every hand-rolled surface — JSON, HTTP, the crypto primitives, the Postgres wire protocol, - templating, SigV4 — and runs in CI as a required gate. Building that harness caught and fixed several real - bugs (a JSON stack-overflow DoS, an HTTP unbounded line-read, an unchecked PG frame panic, a SCRAM - iteration-count DoS, and more), and the JSON parser is checked round-trip against serde_json - on hundreds of thousands of cases. The query builder guards identifiers against SQL injection — important - precisely because an AI tool argument can reach a column or sort slot. + hammers every hand-rolled surface — JSON, HTTP, the crypto primitives, the Postgres wire protocol, the + RFC 6455 WebSocket codec, templating, SigV4 — and runs in CI as a required gate. Building that + harness caught and fixed several real bugs (a JSON stack-overflow DoS, an HTTP unbounded line-read, an + unchecked PG frame panic, a SCRAM iteration-count DoS, and more), and the JSON parser is checked + round-trip against serde_json on hundreds of thousands of cases. The query builder guards + identifiers against SQL injection, and the search grammar sanitizes input before it can reach engine query + syntax — both matter precisely because an AI tool argument can reach a column, a sort slot or a match + expression. +

    +

    + The concurrency claims are tested against real servers, not asserted: queue claim exclusivity under six + concurrent workers, crash recovery through an expired lease, a lost-update race surfacing as + 40001 under Serializable, a rolled-back write never waking a watcher, and two OS + processes chatting through one PostgreSQL over channels.

    @@ -728,17 +1401,18 @@ App::new("api") sutegi does not ship TLS. The intended posture is to terminate HTTPS at a load balancer or service mesh — standard, and fine for the front door. The genuine limitation is in-cluster Postgres and SMTP, often expected encrypted; today those - connections must stay inside a trusted network boundary. TLS is the one primitive we will not - hand-roll; the plan is a single curated, audited dependency (rustls) behind an opt-in - tls feature, added when a real consumer needs it. The mail and storage layers already have an - adapter seam so it drops in without a rewrite. + connections must stay inside a trusted network boundary. (Object storage is already solved: the + storage transport seam borrows the system curl for + https, so S3/R2 works today without a TLS stack in the tree.) TLS is the one primitive we + will not hand-roll; the plan is a single curated, audited dependency (rustls) behind + an opt-in tls feature, added when a real consumer needs it.

    It is a coherent solo bet, not a mature ecosystem

    - Laravel took years and a large community. sutegi is a broad, coherent framework built quickly by one - maintainer — its breadth currently runs ahead of its battle-tested depth, and no real production workload - has yet run on it under load. That is not a defect to paper over; it is the honest framing. + Laravel and Phoenix took years and large communities. sutegi is a broad, coherent framework built quickly + by one maintainer — its breadth currently runs ahead of its battle-tested depth, and no real production + workload has yet run on it under load. That is not a defect to paper over; it is the honest framing.

    Three tiers you can actually claim

    @@ -758,6 +1432,44 @@ App::new("api")
    +
    +

    The workspace map

    +

    + sutegi is one facade crate over twenty-odd small ones. You depend on sutegi and switch + features on; the table is here so you know what you are reading when you open the source, and so an + error message from sutegi_orm or sutegi_pg tells you where to look. +

    +
    + + + + + + + + + + + + + + + + + + + + + + + + + + +
    CrateFeatureResponsibility
    sutegi-jsoncoreJSON value, parser, serializer (deterministic key order).
    sutegi-httpcoreHTTP/1.1 parsing + the blocking thread-pool server on std::net.
    sutegi-webcoreRouter, App, middleware, groups, streaming, static files, /__introspect, and the whole agent tool surface.
    sutegi-cryptocoreSHA-256/1, MD5, HMAC, PBKDF2, HKDF, ChaCha20-Poly1305, base64, CSPRNG.
    sutegi-ormormSchema IR, query builder, the Backend trait, migrations/diff, KV, JSON paths, search, embeddings, watchers, and the two runnable backends.
    sutegi-pgpostgresPure-std PostgreSQL driver: wire protocol v3, SCRAM-SHA-256, COPY, LISTEN, pool.
    sutegi-macrosderive#[derive(Model)] / #[derive(Validate)]. Build-time only — syn/quote never reach your binary.
    sutegi-validatevalidateRulesets and a JSON Schema subset validator, one structured error shape.
    sutegi-queuequeueDurable job queue over the Backend seam.
    sutegi-eventseventsAppend-only event store, aggregates, checkpointed projections.
    sutegi-wswsRFC 6455 codec + the sharded kqueue/epoll reactor.
    sutegi-pubsubpubsubThe Broker seam: in-process, or PgPubSub over LISTEN/NOTIFY.
    sutegi-channelschannelsPhoenix-style channels, presence, the /__channels manifest, the JS client.
    sutegi-actorsactorsActor processes, OTP-style supervision trees, /__actors.
    sutegi-sessionsessionSigned-cookie sessions + CSRF tokens.
    sutegi-authauthUsers, passwords, guards, API tokens, remember-me, throttling, verification/reset.
    sutegi-mailmailEmail builder, RFC 2822/MIME rendering, themed messages, the Transport seam.
    sutegi-templatetemplateBlade-lite engine over Json contexts.
    sutegi-storagestorageThe Storage trait: local fs, DB blobs, S3/R2 + SigV4, over an injected HTTP transport.
    sutegi-hexagonhexagonPorts & adapters primitives: UseCase, AppError, respond.
    sutegi-replreplThe tinker-style shell, in-process or over the wire.
    sutegi-clibinaryThe sutegi command: scaffold, introspect, repl.
    +
    +
    +
    Getting started
    @@ -766,7 +1478,7 @@ App::new("api") Every sutegi app is an ordinary Rust binary — cargo new, add the crate, write a main. There is no separate runtime to install and nothing trailing behind the binary at run time. Choose the feature pillars you want; the HTTP core (json + http + - web) is always present. + web + crypto) is always present, and so is the agent tool surface.

    {@render code(cInstall, 'install')}
    @@ -801,22 +1513,50 @@ App::new("api")

    Feature flags

    - sutegi’s pillars are Cargo features on the facade crate. Only json, http - and web are compiled unconditionally; everything else is opt-in, so your binary contains - exactly the surface you use. The defaults — derive, orm, validate, - ai — suit a typical app. + sutegi’s pillars are Cargo features on the facade crate. Only json, http, + web and crypto are compiled unconditionally; everything else is opt-in, so your + binary contains exactly the surface you use. The defaults are + derive, orm, validate.

    -
      -
    • sqlite / postgres — the two runnable ORM backends (single-node / multi-pod).
    • -
    • graceful — SIGTERM/SIGINT draining for rolling deploys.
    • -
    • queue — the durable, cross-pod job queue (Postgres-backed).
    • -
    • session / auth / auth-mail — cookies, the user system, and email verification/reset.
    • -
    • mail / template / storage / storage-db / events / hex — the remaining services and the hexagonal toolkit.
    • -
    +
    + + + + + + + + + + + + + + + + + + + + + +
    FeatureDefaultGives you
    ormyesSchema, query builder, Backend trait, migrations, KV, JSON paths, search, embeddings, watchers.
    deriveyes#[derive(Model)] / #[derive(Validate)] (build-time only).
    validateyesRequest + tool validation, Ctx::validate/validated*.
    sqliteThe bundled single-node backend, Db.
    postgresThe pure-std multi-pod backend, Pg.
    gracefulSIGTERM/SIGINT draining for rolling deploys.
    queueThe durable job queue (SQLite or Postgres).
    eventsEvent store, aggregates, projections.
    wsApp::ws on the kqueue/epoll reactor.
    pubsub / pubsub-postgresThe in-process broker / cross-pod PgPubSub.
    channels / presencePhoenix-style channels + /__channels / who’s-online tracking.
    actorsActors, supervision trees, /__actors.
    session / auth / auth-mailCookies + CSRF / the user system / verification & reset flows.
    mail / templateThe mailer (pulls template) / the Blade-lite engine.
    storage / storage-dbLocal fs + S3/R2 objects / blobs over the Backend seam.
    hexagonPorts & adapters primitives.
    replThe tinker-style shell (data commands light up with orm).
    +

    - Turn everything off with default-features = false for a minimal HTTP service, then add back - precisely what a given deployment needs. + Features compose the obvious way: auth implies session + orm, + channels implies ws + pubsub, mail implies + template. Turn everything off with default-features = false for a minimal HTTP + service, then add back precisely what a given deployment needs.

    +
    +
    There is no ai feature any more
    +

    + It was removed in 0.6.0. The tool surface — App::tool/stream_tool, the + schema helpers, ToolCtx, /__tools — lives in + sutegi-web and is always compiled. Drop ai from your feature list; no code + change is needed. +

    +
    @@ -829,10 +1569,10 @@ App::new("api")

    {@render code(cConfig, 'config')}

    - .serve() reads HOST (default 0.0.0.0), PORT (default - 8080) and WORKERS from the environment on its own, so a bare app is already - 12-factor without touching Config. + You do not have to use it: the framework reads a handful of conventional variables itself, so a bare app + is already 12-factor.

    + {@render code(cEnv, 'env')}
    @@ -849,35 +1589,63 @@ App::new("api")

    The repo’s examples/ directory is the fastest way in: hello (minimal), todo (every pillar in ~60 lines), auth, events, kv, - storage, hexagonal, and redactor (an agent tool service). + storage, hexagonal, redactor (an agent tool service), + chat (channels + presence, single-pod or cross-pod), ws-chat, + ws-load and http-load (the stress harnesses behind the numbers quoted here).

    +
    +

    The CLI

    +

    + The sutegi binary is a scaffolder and a remote control. Scaffolding follows rigid conventions + on purpose — one right shape per artifact — so an LLM can extend a sutegi app correctly with minimal + context. introspect and repl need no source access at all: they speak the same + agent contract an LLM would. +

    + {@render code(cCli, 'cli')} +

    + Migration verbs are not here — they live in your own binary, because they need your model set. + See Migrations. +

    +
    +
    The basics

    Routing

    - Register routes with .get, .post, .put, .delete (or - the generic .route(method, …)). Each takes a URL pattern, a short doc string — which - shows up in /__introspect, so keep it meaningful — and a handler closure. Path parameters use - :name and are read with c.param("name"). Related routes can share a prefix and + Register routes with .get, .post, .put, .delete — or + .route(method, …) for anything else, including PATCH. Each takes a URL + pattern, a short doc string — which shows up in /__introspect, so keep it meaningful — and a + handler closure. Path parameters use :name and are read with c.param("name"); a + trailing *name captures the rest of the path. Related routes can share a prefix and middleware via .group.

    {@render code(cRouting, 'routing')} +

    + The router splits on segments after trimming leading and trailing slashes, so /todos/1 and + /todos/1/ are the same route. Patterns are matched in registration order; a path that exists + under another method answers 405 rather than 404. +

    Requests & the Ctx

    The &Ctx is your window into a request and into shared application state. Path parameters, - the query string, headers, cookies, and the raw body are all reachable from it, and the body can be parsed - as JSON (c.json()) or a form (c.form()). State registered once with - .state(value) is retrieved by type with c.state::<T>() — and, with the ORM - feature, the database handle with c.db::<Db>(). + the query string, headers, the peer address and the raw body are all reachable from it, and the body can be + parsed as JSON (c.json()) or a form (c.form()). State registered once with + .state(value) is retrieved by type with c.state::<T>() (or + try_state if it may be absent) — and, with the ORM feature, the database handle with + c.db::<Db>().

    {@render code(cRequests, 'requests')} +

    + State is keyed by type, one value per type: registering the same type twice replaces it. Wrap a handle in + a newtype when you need two of the same kind. +

    @@ -892,14 +1660,24 @@ App::new("api")

    Error carries a status, a message, and optional per-field detail. The constructors (Error::bad_request, unauthorized, forbidden, - not_found, unprocessable, internal) cover the common statuses, and + not_found, unprocessable, internal, or + Error::new(status, msg) for anything else) cover the common statuses, and .with_fields(json) attaches structured validation detail.

    {@render code(cResponses, 'responses')} +
    +
    5xx bodies are redacted
    +

    + A 4xx message is part of your API and renders verbatim. A 5xx message is an + internal detail — a SQL error carrying your schema, a subprocess’s stderr, a path on disk — so it + is logged to stderr and the client gets {'{'}"error":"internal error"{'}'}. Anything the + caller is meant to act on belongs in a 4xx. +

    +
    -

    Middleware & groups

    +

    Middleware, CORS & guards

    Middleware comes in two flavours. Before-middleware (Fn(&Request) -> Option<Response>) runs ahead of the handler and can @@ -911,8 +1689,37 @@ App::new("api") {@render code(cMiddleware, 'middleware')}

    The batteries are built in: logger(), rate_limit(max, per), bearer, - basic, cors, cors_preflight, and secure_headers. + basic, cors, cors_preflight, cors_credentialed, + cors_preflight_credentialed, and secure_headers.

    +
    +
    A cookie-authenticated frontend needs the credentialed pair
    +

    + Plain cors stamps Access-Control-Allow-Origin and nothing else — correct for a + public API, quietly useless for a browser app on another origin that must send a session cookie. Three + things have to be right and each fails silently alone: Allow-Credentials must be on the + real response and not only the preflight; Allow-Headers: * stops being a wildcard + once credentials are in play, so the first JSON POST is refused by a server that believes + it allows everything; and a preflight without Max-Age is paid on every mutating call. The + credentialed pair handles all three and is idempotent, so composing both halves cannot emit a duplicate + Allow-Origin — which a browser rejects outright while curl -i shows a + perfectly good 204. +

    +
    +
    +
    ops_guard is the one that protects /__
    +

    + Introspection and tool invocation are open by default — that is the + agent-native contract. In any deployment where the agent surface must not be public, set an + ops_guard: it runs ahead of the global middleware chain for every /__-mounted + route (/__introspect, /__metrics, /__tools*, + /__channels, /__actors, your own), while /__health and + /__ready stay open so orchestrator probes need no credential. It is gated on the segments + the router actually matches, not on the raw path string — see + Security posture for why that distinction was a CVE-shaped + bug. +

    +
    @@ -921,8 +1728,10 @@ App::new("api") Never trust the request body. Deriving Validate alongside Model reads the #[validate(…)] field attributes and generates the type’s ruleset at build time. In a handler, c.validated::<Todo>() parses the body, validates it, and hydrates a typed - value — or returns a 422 with structured, per-field messages. For shapes without a model, - build an ad-hoc Ruleset and call c.validate. + value — or returns a 422 with structured, per-field messages. The same ruleset drives the + other three input surfaces (validated_form, validated_query, + validated_path). For shapes without a model, build an ad-hoc Ruleset and call + c.validate.

    {@render code(cValidation, 'validation')}

    @@ -930,86 +1739,374 @@ App::new("api") Integer, Number, Bool), formats (Email, Url, Alpha, AlphaNum), bounds (Min/Max, Between, MinLen/MaxLen), and relational rules (In, - Same). + Same). Every path produces the same error shape: + {'{'} field: [messages] {'}'}.

    Agents get this for free
    -

    AI tool arguments are validated against each tool’s JSON schema automatically, returning the same structured errors.

    +

    + AI tool arguments are validated against each tool’s JSON schema by a second validator in the same + crate (types, required, enum, bounds), returning the same structured errors — + and, for streaming tools, before the stream opens, so a malformed call still gets a normal JSON + 422. +

    +
    +

    Static files

    +

    + App::static_dir(prefix, dir) mounts a directory as a rest route. It is enough to serve a + built SPA next to your API: point / at dist and register it last, since routes + match in registration order. +

    + {@render code(cStatic, 'static')} +

    + Traversal attempts, dotfiles and backslashes are 404s rather than reads outside the root; a + directory (or the bare prefix) serves its index.html. +

    +
    +
    Database

    Models

    Deriving Model on a plain struct makes it the single source of truth for a table: schema, - migrations, typed reads, JSON serialization, and save() all come from it. The - #[model(…)] attributes tune the mapping — primary marks the key, - table = "…" overrides the inferred snake-case plural name, and skip keeps - a field out of the database. Bools round-trip as real bools and Option<T> becomes a - nullable column. + migrations, typed reads, JSON serialization, and save() all come from it. Bools round-trip as + real bools and Option<T> becomes a nullable column.

    {@render code(cModels, 'models')} +

    + The #[model(…)] attributes are the mapping controls, and everything except + skip flows into the schema IR the migration diff reads: +

    +
      +
    • table = "…" (struct level) — override the inferred snake_case plural name.
    • +
    • primary — the primary key. column = "…" — a differing column name.
    • +
    • unique / index — a UNIQUE constraint / a secondary index.
    • +
    • default = <lit> — a column default (int, float, bool or string), which is what makes adding a NOT NULL column a safe migration.
    • +
    • vector / vector(dim = N) — an embedding column on a Vec<f32>.
    • +
    • has_many / has_one / belongs_to — a relation, not a column.
    • +
    • skip — not persisted, not serialized, default-initialised on load.
    • +
    +

    + The derive also generates from_input() — a lenient hydrate from a partial client payload, + which is what tool closures use on already-validated args — and to_json(). Both macros run at + build time, so syn/quote never reach your binary; turn them off with + default-features = false and hand-write the Model impl if you prefer. +

    +
    + +
    +

    Relations

    +

    + A relation field is declared with an attribute, is not a column, and gets its own generated batch loader. + User::with_posts(&db, users) runs one WHERE author_id IN (…) query for + the whole batch and attaches the typed children — the classic two-query strategy, so N+1 never arises. +

    + {@render code(cRelations, 'relations')} +

    + belongs_to also carries the foreign key into the schema, including + on_delete = "cascade", so the migration diff sees it. Join keys are integers (the usual + primary/foreign-key case), and loading is explicit rather than lazy: there is no hidden query behind a + field access. +

    The query builder

    For anything beyond all/find, Model::query() returns a fluent, - fully parameterized query builder — filter, or_group, where_null, - like, join/left_join, group_by, order_by, - limit/offset, and where_raw as the explicit escape hatch. Run it - through a backend to get JSON rows, typed values (fetch), or a Page<T> - (paginate_typed). Values are always bound as parameters, and identifiers are validated against - an allowlist — an AI tool argument cannot smuggle SQL into a column or sort slot. + fully parameterized query builder. Run it through a backend to get JSON rows + (select), typed values (fetch/fetch_one), a count, an existence + check, or a Page<T> (paginate_typed). Values are always bound as + parameters, and identifiers are validated against an allowlist — an AI tool argument cannot smuggle SQL + into a column or sort slot.

    {@render code(cQueries, 'queries')} +
    +
    order_by’s bool is descending
    +

    + .order_by("due", false) is ascending; + true is descending. Chain several for multi-column ordering. +

    +
    +

    + QueryBuilder, UpdateBuilder and DeleteBuilder can also be used + directly against a table name when there is no model, and build() hands you the + (sql, params) pair — which is why the core ships no driver at all and the default binary stays + tiny. +

    -

    Backends: SQLite & Postgres

    +

    Backends & capabilities

    This is the key architectural story. The ORM is written against a Backend trait, not a concrete engine. Db (SQLite, the sqlite feature) is the single-node store — one embedded file, zero operations. Pg (Postgres, via a pure-std wire driver, the postgres feature) is the multi-pod store. Both implement - Backend, so Model is written once and every call site is identical: swap - Db for Pg::from_env()? and your handlers do not change. + Backend, so Model is written once and every call site is identical.

    {@render code(cBackend, 'backend')}

    The trait is small — five required primitives (query, execute, insert, upsert, migrate) — and everything else - (select, count, paginate, transactions…) is a default method - implemented once on top. Write your domain against &impl Backend and choose the store at - boot: SQLite for local dev, edge, and single-node; Postgres when you scale to many pods. + (select, count, paginate, transactions, locks, bulk insert…) + is a default method implemented once on top. It is object-safe, so &dyn Backend works; + the typed helpers are Self: Sized-gated. +

    +

    Where the engines differ, a bit says so

    +

    + Some things genuinely cannot be papered over. Backend::capabilities() returns a + BackendCaps descriptor — the honest default is everything off, so a backend never advertises + what it has not implemented — and a gated call returns + unsupported: <cap> is not available on <backend> before any SQL is sent. + App::register_capabilities(db.capabilities()) publishes the block in + /__introspect, so an agent reads what the store can do instead of finding out from a dialect + error. +

    + {@render code(cCaps, 'caps')} +
    + + + + + + + + + + + + + + + +
    CapabilitySQLitePostgres
    advisory_locksprocesscluster
    live_queriesprocesscluster
    isolation_levelsyesyes
    returning_dmlyesyes
    json_pathyesyes
    ftsyesyes
    row_locks / skip_lockednoyes
    bulk_copynoyes
    json_containsnoyes
    listen_notifynoyes
    vectornoyes
    +
    +

    + A capability describes the framework surface, not the underlying C library: bundled SQLite ships + JSON1, FTS5 and RETURNING, but a bit stays off until sutegi exposes the feature through the + builder. Pick the backend for the deployment: SQLite for local dev, edge, and single-node; Postgres when + you scale to many pods or need cluster-wide coordination. +

    +
    + +
    +

    Locks, isolation, bulk

    +

    + Four throughput and correctness primitives, all through the same Backend seam and all gated + on the capability bits above. +

    + {@render code(cConcurrency, 'concurrency')} +
      +
    • Row locks. The builder stores the request; the executing backend emits it. Postgres emits FOR UPDATE [SKIP LOCKED|NOWAIT]. SQLite treats plain for_update/for_share as a no-op — its write transaction already holds the whole database, strictly coarser — but skip_locked/nowait error: altered contention semantics are the point, and SQLite cannot express them.
    • +
    • Isolation. Postgres gets BEGIN ISOLATION LEVEL …. SQLite is always serializable, so levels map to when the write lock is taken (Serializable → BEGIN EXCLUSIVE, RepeatableRead → BEGIN IMMEDIATE, ReadCommitted → BEGIN) — stronger than asked is honest, weaker would not be.
    • +
    • RETURNING on DML. Same syntax on both engines (SQLite ≥ 3.35; the bundled build is newer), routed through query so the rows come back.
    • +
    • Bulk insert. The default batches multi-row INSERT … VALUES under the placeholder budget on any backend; Postgres overrides it with wire-native COPY FROM STDIN, with tabs, newlines, backslashes and NULLs escaped correctly (hostile content round-trips in the tests).
    • +
    +
    + +
    +

    Advisory locks

    +

    + Named locks are the coordination primitive for “exactly one of us”: singleton jobs, janitors, + leader election, a migration mutex. try_lock returns an RAII LockGuard; + with_lock is the singleton-job shape, where Ok(None) means another pod ran it. +

    + {@render code(cAdvisory, 'advisory')} +
      +
    • Postgres is cluster-scoped via pg_try_advisory_lock on a dedicated session — never a pooled connection, so a leader lock held for the process lifetime cannot starve request traffic. Release is closing that session, which is exactly what the server does when a holder crashes, so crash-release needs no cleanup path.
    • +
    • SQLite is process-scoped: a named-mutex registry keyed per database file. Two OS processes on the same file do not contend — believe the capability.
    • +
    • Inside a Postgres transaction, try_lock is transaction-scoped (pg_try_advisory_xact_lock) and releases at COMMIT, not at guard drop.
    • +
    • with_lock reconnects per call on Postgres — fine for a janitor loop, wrong in a request hot path. These are for coordination, not thousands of concurrent holds.
    • +
    +
    + +
    +

    JSON path queries

    +

    + Query inside JSON columns through the builder — document-store mode for data whose schema is not + known up front, which is exactly the shape agent-authored payloads arrive in. +

    + {@render code(cJson, 'json')} +

    + The path grammar is a deliberate subset — identifier keys and [n] indexes — because exotic + keys need per-engine quoting rules while identifiers compile everywhere. This is the first builder feature + whose SQL shape differs per engine, so the builder stores parsed segments and the executing + backend compiles them: json_extract(col, ?) with typed values on SQLite, + col #>> ? with a bound {'{'}a,b,0{'}'} text array and value-driven casts on + Postgres (::numeric so 9 < 10 compares as numbers). Containment is Postgres + only; SQLite errors honestly rather than emulating subset semantics with json_each walks. +

    +
    +
    Indexing
    +

    + A jsonb GIN index for containment-heavy tables is not emitted by migrations yet — the + schema IR learns non-btree index kinds alongside the search work. Add one by hand if @> + gets hot: CREATE INDEX ON docs USING GIN (meta). +

    +
    +
    + + + +
    +

    Embeddings & vectors

    +

    + Embeddings are a first-class column type, not a blob you decode by hand. A + #[model(vector(dim = N))] field becomes vector(N) on Postgres (pgvector) and + TEXT on SQLite, and travels in pgvector’s canonical [1,2,3] text form so + the same value round-trips identically on either backend. +

    + {@render code(cVectors, 'vectors')} +

    + Two search paths share the same distance semantics (lower is closer, matching pgvector’s operators): + portable brute force, which loads candidates and ranks them in Rust, and pushdown, which sends + ORDER BY col <=> ? LIMIT k to the database and uses its ANN index. Metric + covers cosine, L2 and inner product. sutegi does not compute embeddings — call whatever model you use and + store the vector. +

    +
    +
    Dimension changes are invisible to drift
    +

    + A vector(dim) column reflects back dimensionless from Postgres introspection (the dimension + lives in a catalog the reflection does not read), so changing it is not detected by + migrate:drift. Declare that change in a migration explicitly. +

    +
    +
    + +
    +

    Reactive queries

    +

    + watch(query) gives you the current rows plus pushed diffs + whenever the watched result actually moves. It is the primitive that unifies LISTEN/NOTIFY, the + event-store wakeups and cross-pod pubsub into one shape — and it turns a channel topic into a live query + in one loop. +

    + {@render code(cReactive, 'reactive')} +

    + The semantics are table-coarse requery-diff (v1): on a change to a watched table — debounced 25 ms, + bursts coalesced — each watcher re-runs its query and diffs by primary key. A new pk is + added, a vanished pk is removed, a same-pk-different-row is + updated, and a write that does not move the watched result emits + nothing (the requery runs, the diff is empty, nothing is sent). Twenty + writes in one window typically become one Change carrying twenty rows. +

    +

    + Guardrails: 1024 live subscriptions per watcher, dropping a subscription unregisters it, and dropping the + watcher shuts the worker down — interrupting the blocked LISTEN session on Postgres. One watcher per + process per backend handle. On SQLite, attach the watcher before serving traffic, and remember that writes + from other processes on the same file are invisible.

    Migrations

    - A Migrator holds an ordered list of Migrations, each with an id, a name, and up - (and optionally down) closures. Migration::reversible takes both directions; - migrate_schema(&T::schema()) creates a table straight from a derived model. - sutegi::migrate::dispatch wires up the migrate, migrate:status, and - migrate:rollback subcommands — call it early in main and it runs the requested - command then returns true so you can exit before serving. + sutegi turns a change to a #[derive(Model)] struct into a versioned, reversible, + backend-portable migration — the TypeORM generate workflow without the live-database + nondeterminism. There are two kinds: declarative migrations, a list of + schema ops stored as JSON, generated by diffing and reversible for free; and + closure migrations, hand-written up/down + closures for data backfills and DDL the diff engine does not model.

    {@render code(cMigrations, 'migrations')} + {@render code(cMigrateCli, 'migratecli')} +

    Diff against history, not the live database

    +

    + migrate:gen computes changes by diffing your models against the shadow schema — the + schema you get by folding every existing migration’s ops in memory. It never consults your live + database to decide what to generate, so the same repository state always produces the same migration no + matter what anyone’s local database happens to contain. The live database is consulted for exactly + one thing: drift detection, which reports both + DB vs migrations (a hand-edit, or migrations not fully applied) and + models vs migrations (you changed a model and have not run :gen). + sutegi::migrate::report_json is the same report as JSON, ready to mount at + /__migrations for agents and dashboards. +

    +

    Safety classes

    +
      +
    • Safe — create a table, add a nullable or defaulted column, create an index, add a foreign key, widen a column.
    • +
    • NeedsData — add a NOT NULL column with no default, or tighten a column to NOT NULL. Valid only against an empty table; give it a default or write a backfill.
    • +
    • Destructive — drop a table or column, a lossy type change. :gen writes these but flags them loudly; the committed, reviewed file is your approval.
    • +
    +

    + Renames are never guessed. A dropped column plus an added column of + the same type generate a drop + add pair with a “possible rename?” warning — edit the file to a + rename op if that is what you meant, which preserves the data. Applied migrations are checksummed in + _sutegi_migrations, so editing a file after the fact fails the next migrate + rather than silently diverging (Migrator::repair re-stamps after a deliberate edit). On + Postgres, run/rollback take a session advisory lock, so many pods booting at + once serialize instead of racing. +

    +
    +
    Dev-mode sync
    +

    + Model::migrate (and sutegi::orm::migrate::sync) apply + additive, non-destructive changes directly — create missing tables, add columns, indexes and + foreign keys — with no migration file. They never drop anything and refuse, pointing you at + migrate:gen, any change that needs a real migration. Use this locally; use files in + production. +

    +
    +
    +
    Honest caveats
    +

    + down restores schema shape, not data — a dropped column comes back empty. The + schema IR covers tables, columns, secondary indexes and single-column foreign keys; views, triggers, + sequences, CHECK constraints and multi-column FKs are out of scope (use a closure migration). On SQLite, + changing a column’s type or nullability is a full table rebuild — rows survive, but it is heavier + than the in-place ALTER Postgres does. +

    +

    The key/value store

    - Kv<B> is a namespaced JSON key/value store over either backend — handy for - config, feature flags, cached values, or small shared state that does not deserve a table. Values are - arbitrary Json; keys are grouped by namespace and support set/get/ - delete/keys/scan/scan_prefix/count/ - clear. It is ordinary application state. + Kv<B> is a namespaced JSON key/value store over either backend — one table, + single-statement reads and writes — handy for config, feature flags, cached values, or small shared state + that does not deserve a schema. Values are arbitrary Json; keys are grouped by namespace.

    {@render code(cKv, 'kv')} +

    + On a single SQLite node it is the natural home for config and flags; on Postgres it works for small + shared state. It is ordinary application state — no Arc<Mutex<…>>. +

    @@ -1017,13 +2114,20 @@ App::new("api")
    Agents & realtime

    The agent surface

    - Because every route, model, and tool registers its own metadata, sutegi exposes your whole app as - machine-readable JSON with no extra work. /__introspect returns the full surface — routes - with their docs, models with their schemas, and tools with their input schemas. An agent reads that once, - then calls tools through /__tools. There is no SDK and no source access required: the app you - built for humans is the app the model drives. + Because every route, model, tool, capability and searchable table registers its own metadata, sutegi + exposes your whole app as machine-readable JSON with no extra work. /__introspect returns the + full surface; an agent reads it once, then calls tools through /__tools. There is no SDK and + no source access required: the app you built for humans is the app the model drives.

    {@render code(cAgents, 'agents')} +

    + The contract is three steps — discover (/__introspect), read the manifest + (/__tools), invoke (POST /__tools/:name, or + /__tools/:name/stream for SSE) — and it is extended, not replaced, by the realtime surfaces: + /__channels teaches an agent the channel protocol, /__actors shows live process + status. Everything under /__ except the two probes sits behind + ops_guard. +

    The full contract

    @@ -1041,9 +2145,14 @@ App::new("api") that shares app state (c.db::<Db>(), c.state::<T>()). Build the argument schema with the schema:: helpers (object, string, integer, boolean, array). .stream_tool(…) is the - same, but the closure also gets an SseSink to stream results as Server-Sent Events. + same, but the closure also gets an SseSink to stream results as Server-Sent Events, and the + manifest marks it "streaming": true so an agent knows to hit the SSE endpoint.

    {@render code(cTools, 'tools')} +

    + A tool and an HTTP route can — and usually should — be two adapters over one + use case. Write the logic once; expose it to both audiences. +

    @@ -1052,23 +2161,183 @@ App::new("api") Because the server is blocking and thread-per-connection, streaming is trivial and naturally backpressured — there is no executor to fight. sse(producer) gives the producer an SseSink with data, event, and comment; each frame is - flushed immediately. It is the same transport that carries live LLM tokens back to a UI and the one - .stream_tool rides on. Regular responses use keep-alive; streams are close-framed by design. + flushed immediately. stream(status, content_type, producer) is the raw-bytes equivalent for + NDJSON and large exports.

    {@render code(cStreaming, 'streaming')} +

    + It is the same transport that carries live LLM tokens back to a UI and the one + .stream_tool rides on. Regular responses use keep-alive; streams are close-framed by design, + which is valid HTTP/1.1 and needs no chunked encoding. Behind nginx, set + proxy_buffering off — the ontzi config already does. +

    +
    + +
    +

    WebSockets

    +

    + The ws feature adds App::ws. The HTTP side stays blocking thread-per-connection, + but an upgraded socket detaches into a sharded kqueue/epoll reactor — + no async runtime, no futures, just poller syscalls — so the worker thread is freed immediately and an idle + connection costs ~340 bytes of user-space RSS and zero threads. +

    + {@render code(cWs, 'ws')} +

    + Measured on a dev laptop: 80,000 live sockets at 0.0% idle CPU, a broadcast enqueue of 80k shared-Arc + frames in ~1.5 ms, and 5k-fleet delivery at p50 15 ms / max 30 ms end-to-end. The codec is + strict RFC 6455 — masking required, minimal length encodings, control-frame rules, close-code + validation, UTF-8 enforcement — with a deterministic fuzz suite behind it. +

    +
      +
    • Callbacks run inline on the shard, which is what guarantees per-connection ordering. Keep them CPU-quick; push blocking work to your own threads and answer later through the cloneable Conn.
    • +
    • Broadcast by cloning one encoded frame. text_frame/binary_frame encode once into an Arc; send_shared puts it on each queue.
    • +
    • Slow consumers are dropped at max_buffered; ping and idle sweeps close dead sockets; RLIMIT_NOFILE is raised at startup.
    • +
    • Every knob lives in WsConfig, and ws_config must come before the first .ws(…) — that call starts the reactor.
    • +
    +
    +
    Cookie-authenticated sockets must pin their origins
    +

    + A browser sends cookies on a cross-origin WebSocket handshake and the same-origin policy does not stop + it. If your socket is authenticated by a cookie, call check_origin([…]) — otherwise + any page on the internet can open an authenticated socket as your user (CSWSH). Bearer-token clients are + unaffected: they carry no ambient credential. +

    +
    +
    + +
    +

    PubSub

    +

    + Topic fan-out behind one Broker trait, so the layer above it — channels, presence, your own + code — does not know whether it is talking to one pod or a fleet. PubSub is the in-process + broker; PgPubSub (pubsub-postgres) is the same trait over Postgres + LISTEN/NOTIFY. +

    + {@render code(cPubsub, 'pubsub')} +

    + PgPubSub uses one shared PG channel with topics inside a JSON envelope — immune to the 63-byte + identifier truncation trap — plus a lazy publisher with one transparent retry and a listener that + reconnects with capped backoff. Fan-out cost on the database is per message, not per subscriber: + one connection per pod for listening, one for publishing. +

    +
    +
    At-most-once, by design
    +

    + Delivery is fire-and-forget (Phoenix’s contract too): a pod that is reconnecting misses messages + sent meanwhile, and NOTIFY payloads cap at ~8 KB. If a message must survive, put it in + a table (or the event store) and broadcast that there is + news. +

    +
    +
    + +
    +

    Channels & presence

    +

    + Channels (channels) are the realtime identity feature: many topics multiplexed over one + socket, with joins, replies, broadcasts, per-membership state — and an agent manifest. Phoenix is the + design reference; the deltas are deliberate. +

    + {@render code(cChannels, 'channels')} +

    + Broadcasts ride the pubsub Broker seam, so + the same channel code is single-pod on the in-process broker and cross-pod on + PgPubSub with zero changes — verified by two OS processes chatting through a real + PostgreSQL. A broadcast to a large room encodes the frame once and takes each reactor shard’s lock + once, so the 80k-socket transport numbers apply unchanged. +

    +

    The wire protocol

    + {@render code(cChannelsWire, 'chwire')} +

    + GET /__channels returns the full manifest — envelope shape, control events, every + channel’s pattern, docs and join/event schemas — which is enough for an agent to join and speak over + a raw WebSocket with no client library. That manifest is part of the agent contract, next to + /__introspect and /__tools. +

    + {@render code(cChannelsJs, 'chjs')} +

    Join lifecycle notes

    +
      +
    • Inside on_join the member is not admitted yet: a broadcast there reaches the room but not the joiner, and a push lands before the join reply. Use socket.after_join(…) for welcome pushes — it runs after the ok reply is on the wire, and never runs if the join is refused.
    • +
    • A join on an already-joined topic replaces the membership: the old one gets the leave callback with LeaveReason::Rejoin, and assigns start fresh.
    • +
    • assigns are per-membership JSON state (socket.assign / assign_get), gone on leave or disconnect.
    • +
    +

    Presence

    +

    + The presence feature adds who’s-online tracking: Presence::track / + untrack / list. The tracked member receives presence_state (the full + view, after its join reply); the room receives presence_diff {'{'}joins, leaves{'}'} on every + change. Untrack is automatic on leave, rejoin and disconnect, and several memberships may track the same + key — one user, many tabs, each contributing a meta. +

    +
    +
    Presence is heartbeat-based, not a CRDT
    +

    + Each pod re-publishes its local state every presence_heartbeat (default 30 s) and + expires pods silent for ~2.5×, reporting their members as leaves. So a crashed pod’s users can + linger up to ~75 s, and partition conflicts resolve by expiry. That is the right trade for a + “who’s online” sidebar; keep anything stronger in a table. +

    +

    Queues

    - Some work should not block the response. The durable queue (the queue feature) is - Postgres-backed and cross-pod: it claims jobs with - FOR UPDATE SKIP LOCKED, so any pod can pull the next one, and a visibility-timeout retry - recovers work from a crashed worker. Register a handler by name, dispatch a JSON payload - (optionally delayed), and start N worker threads. A job survives a pod restart and - dead-letters once its retries are spent. + Some work should not block the response. The durable queue (queue) runs over the + Backend seam, so one jobs table and one set of SQL work on bundled SQLite (a single box) and + on Postgres (many pods). Jobs survive a crash: the claim stamps a lease instead of deleting the row, so a + dead worker’s job becomes visible again after the visibility timeout — at-least-once. Retries, + delays, priorities and dedupe keys are columns, so a restart forgets nothing.

    {@render code(cQueues, 'queues')} +

    + Claims are exclusive on both backends: FOR UPDATE SKIP LOCKED wherever + capabilities().skip_locked says it exists, and on SQLite the serialized writer already + provides it, since the second UPDATE … RETURNING no longer sees the claimed row. Postgres + additionally makes that exclusivity cross-pod; queue.cross_pod() tells you which + guarantee you actually have instead of letting the docs imply the stronger one. +

    +
      +
    • Named queues with their own pools — .queue("video") plus start_on("video", 1), so a slow job class cannot starve a fast one.
    • +
    • Dedupe keys — .unique("yt:abc") hands back the live row’s id instead of enqueueing a second copy, backed by a partial unique index that deliberately excludes dead letters, so a failure never owns a key forever.
    • +
    • A panicking handler is a failed job, not a lost worker — the panic is caught, the row retries or dead-letters like any other, and the pool keeps going.
    • +
    • Dispatch wakes an idle worker over a condvar instead of making it wait out the poll interval; the interval stays as the safety net for delayed jobs and work enqueued by another pod.
    • +
    • JobCtx gives a handler what it needs to behave well: is_last_attempt() to tell a retryable blip from a terminal failure before writing a user-visible error, heartbeat() to outlive the visibility timeout, and should_stop() for loops.
    • +
    +

    + Timestamps are epoch milliseconds supplied by the caller rather than SQL now() — which is + what lets one statement work in both dialects, and makes schedules testable without sleeping. Tuning: + visibility_timeout, poll_interval, retry_backoff. +

    +
    + +
    +

    Actors & supervision

    +

    + The actors feature is the OTP half of the stack: isolated processes with typed mailboxes, and + supervision trees that restart them when they crash. An Actor owns its state outright on its + own thread — no locks, no Sync bound on the state — and other threads communicate by message + through an ActorRef. +

    + {@render code(cActors, 'actors')} +
      +
    • Bounded mailbox (default 1024): a full mailbox fails tell fast with TellError::Full, so backpressure is explicit rather than a hidden unbounded queue.
    • +
    • Let it crash. A panic in handle kills only that actor and is reported as ExitReason::Crashed. Lifecycle hooks started() / stopped(&reason) run on every (re)start and exit. stop() is queued behind pending messages, so the mailbox drains first.
    • +
    • Restart intensity (default 3 per 5 s) fails the supervisor and stops all children when exceeded, so a crash-looping child cannot burn a thread forever. RestForOne/OneForAll stop dependents in reverse order and restart them in start order.
    • +
    • SupervisorHandle::child_ref(name) is the whereis analog — re-fetch after a crash, because an old ref points at the dead generation’s queue. start() returns only once every child’s first generation is up.
    • +
    • App::actors(registry) mounts GET /__actors: state, mailbox depth, restart counts and last crash message per actor — the Observer-lite half of the agent contract, ops_guard-gated like the rest.
    • +
    +
    +
    One OS thread per actor
    +

    + This is supervision and fault isolation, not BEAM-scale lightweight processes: thousands of actors, not + millions. There are no links or monitors beyond the supervisor relationship, no distribution (Postgres + is the bus), and no hot code reload. In release the workspace builds with + panic = "abort", so catch_unwind never runs there — dev and test builds get the + full crash-restart semantics, and the production posture is “the pod supervisor restarts the + process”. +

    +
    @@ -1076,36 +2345,55 @@ App::new("api")
    Framework services

    Authentication

    - The auth feature is a complete user system over either backend. Users handles - registration and authentication with PBKDF2-HMAC-SHA256 passwords (600k iterations by default, per OWASP; - the hash never leaves the store). Auth issues signed-cookie logins with a server-enforced - expiry baked into the signed payload. Tokens mints stg_ bearer tokens for agents - — shown in plaintext once, stored hashed, and surviving logout. The guards - (require_auth, require_role, require_token) are just middleware you - attach to a group. + The auth feature is a complete user system over either backend — everything Laravel’s + auth scaffolding does, with no third-party dependency. Users handles registration and + authentication with PBKDF2-HMAC-SHA256 passwords (600k iterations by default, per OWASP; the hash never + leaves the store). Auth issues signed-cookie logins. Tokens mints + stg_ bearer tokens for agents. The guards are just middleware you attach to a group.

    {@render code(cAuth, 'auth')} +

    What the production pieces do

    +
      +
    • Remember me — selector/validator cookies where only the validator’s SHA-256 is stored, the validator rotates on every use (a stolen-then-replayed copy revokes the row, surfacing the theft), tokens are bound to the password hash at mint, and server-side expiry (default 30 days) ignores whatever the client claims.
    • +
    • Sessions bound to the password hash — login stamps a fingerprint of the current PHC string into the signed session, and a stale binding reads as anonymous. Changing a password logs out every other device on its next store-checking request.
    • +
    • Login throttling — a DB-backed fixed window (default 5 attempts / 60 s, keyed however you like; the convention is login:<email>|<ip>), so every pod counts the same attempts.
    • +
    • CSRF — a token minted inside the signed session and compared in constant time; the require_csrf guard enforces X-CSRF-Token on mutating methods (419, Laravel’s “Page Expired”) while passing reads and Authorization-header callers.
    • +
    • Auto-rehash at login — a verified password whose stored iteration count is below the store’s is transparently re-hashed (best-effort; a rehash failure never fails a valid login). Raise the work factor once and the fleet upgrades itself credential-by-credential.
    • +
    • Guards — require_auth, require_role, require_verified, require_csrf, require_token. Auth::identify is the revival point (session-or-remember) and its Identified::attach sets the fresh cookies, because sutegi middleware cannot set cookies on pass-through.
    • +
    • API tokens — the agent door. Plaintext shown once, only the SHA-256 stored; issue_expiring adds a deadline, every successful verify stamps last_used_at, and tokens survive logout.
    • +
    • Profile operations — change_password (verifies the current one first), set_name, and set_email (normalizes, checks uniqueness, and resets verified_at so the new address must re-verify).
    • +
    +

    + Unknown emails burn the same PBKDF2 time as wrong passwords, so the endpoint does not become a user + enumeration oracle. +

    Verification & reset

    Add the auth-mail feature and AuthMail for email verification and - password-reset flows — reset tokens are enumeration-safe and bound to the current password hash, so they - are stateless and single-use. + password-reset flows: built-in text+HTML templates, signed expiring links (24 h / 1 h), no + state tables. Reset tokens are enumeration-safe (an unknown email is a silent Ok) and bound + to the current password hash — the moment the password changes, every outstanding link dies, which makes + them single-use without storing anything.

    -

    Sessions

    +

    Sessions & CSRF

    The session feature provides signed-cookie sessions (HMAC-SHA256) with no server-side store: Sessions::new(secret), then load a session off the request, set/ - get/remove values, and save it onto the response. Call - .insecure() in local http:// development to drop the cookie’s - Secure flag. This is the machinery Auth is built on; use it directly for - lightweight per-visitor state like a cart or a wizard step. + get/remove values, and save it onto the response. The expiry is + stamped inside the signed payload, so a stolen cookie dies on schedule no matter what the client + claims. Call .insecure() in local http:// development to drop the + cookie’s Secure flag.

    {@render code(cSessions, 'sessions')} +

    + This is the machinery Auth is built on; use it directly + for lightweight per-visitor state like a cart or a wizard step. +

    @@ -1113,17 +2401,26 @@ App::new("api")

    The mail feature gives you an Email builder with real RFC 2822/MIME rendering (multipart/alternative, encoded-words, header-injection folding) and a one-method - Transport seam. Built-in drivers cover Log (dev default), Memory - (tests), a pure-std Smtp client, and Sendmail; Mailer::from_env() - picks one from MAIL_* variables. For notification-style mail, MailMessage builds - a themed HTML card and matching plain text from the same fluent blocks. + Transport seam. Built-in drivers cover log (dev default — messages print, + nothing escapes), memory (test assertions), a pure-std smtp client (EHLO, AUTH + PLAIN/LOGIN, dot-stuffing) and sendmail (pipe to the local Postfix — the VPS shape); + Mailer::from_env() picks one from MAIL_*.

    {@render code(cMail, 'mail')} +

    + Theme + MailMessage produce a clean, email-client-safe HTML card + and a matching plain-text part from the same blocks, so every message + is multipart/alternative for free. The outer chrome is a + template source — swap it wholesale with + Theme::layout(…) while block rendering keeps working. +

    -
    A hosted provider is ~10 lines
    +
    The SMTP client has no TLS

    - A Resend/SendGrid transport is a small adapter over any HTTP client — the same seam that lets sutegi dodge - the TLS wall until it lands (see the posture page). + Same stance as the Postgres driver: point it at an in-cluster relay or Mailpit on + localhost:1025. For hosted delivery over the public internet, a Resend/SendGrid/Postmark/SES + adapter is ~10 lines over any HTTP client — the same seam that lets sutegi dodge the TLS wall until it + lands (see the posture page).

    @@ -1134,17 +2431,38 @@ App::new("api") The storage feature abstracts object storage behind one Storage trait (put/get/stat/delete/list/ get_reader, with traversal-validated keys). FsStorage writes to a local directory - (atomic temp-and-rename); DbStorage (the storage-db feature) stores blobs in - SQLite or Postgres for multi-pod files with no new infrastructure; and S3Storage puts the same - trait on a real bucket — AWS S3, Cloudflare R2, MinIO, Garage. It moves the bytes over an injected - HttpTransport, which is how an S3 client exists here with no TLS stack and no dependency: - SystemCurl borrows the system curl for https, PlainHttp is pure std - for a store on a trusted path. Requests are SigV4-signed over the real payload hash, and downloads are - ETag-verified. The same credentials also presign, which is the agent-native trick: mint a time-limited URL - and let the client (or agent) move the bytes straight to the object store — they never pass through your - server. + (atomic temp-and-rename, real streaming reads); DbStorage (storage-db) stores + blobs in SQLite or Postgres — multi-pod files with no new infrastructure, honest ceiling a few MB per + object; and S3Storage puts the same trait on a real bucket: AWS S3, Cloudflare R2, MinIO, + Garage, Ceph RGW.

    {@render code(cStorage, 'storage')} +

    The transport seam (how S3 works with no TLS in the tree)

    +

    + S3Storage never opens a socket itself. It hands a signed request to an + HttpTransport — one method — and two implementations ship with it. SystemCurl + delegates the https handshake and certificate verification to the system curl, so + the crypto that must not be hand-rolled isn’t and the dependency count stays at zero. + PlainHttp is pure std and refuses + https rather than pretending — for a store on a trusted path: in-cluster MinIO, a + sidecar Garage, a dev container. Your own client is impl HttpTransport for MyClient. +

    +

    + Every request is signed with the real payload hash + (x-amz-content-sha256, never UNSIGNED-PAYLOAD), so a body altered in flight is + refused by the store — integrity that holds even over PlainHttp. Uploads and downloads are + ETag-verified when the store reports a plain MD5 (multipart and encrypted ETags are skipped, + not faked; verify_etag(false) opts out). list follows continuation tokens and is + bounded by max_list_keys (default 100,000), so a ten-million-object bucket + errors instead of silently truncating. Signing is verified against + AWS’s published known-answer vectors, for presigned URLs and signed headers alike. +

    +

    + S3Store deliberately does not implement Storage: minting a URL is a different + contract than moving bytes. S3Store::storage(transport) is the crossing point, and on its own + S3Store is a pure-std SigV4 presigner — which is the agent-native trick: mint a time-limited + URL and let the client or agent move the bytes straight to the object store, never through your server. +

    @@ -1153,11 +2471,15 @@ App::new("api") The events feature is an append-only event store over the Backend seam. You append events to a per-entity stream with an Expected version for optimistic concurrency (Any, NoStream, or Version(n)), and fold current state - back on demand by implementing Aggregate::apply. Projections are checkpointed - consumers: the handler’s writes and its checkpoint bump commit in one transaction, giving - exactly-once read models that you can reset and rebuild from the log. + back on demand by implementing Aggregate::apply. Global log positions are gap-free.

    {@render code(cEvents, 'events')} +

    + Projections are checkpointed consumers: the handler’s writes and its checkpoint bump + commit in one transaction, giving exactly-once read models that you can reset and rebuild + from the log. append_tx composes with a transaction you already own, and position-race + retries back off quadratically so a thundering herd cannot exhaust one writer’s retries. +

    @@ -1167,31 +2489,79 @@ App::new("api") {'{{ escaped }}'} and {'{!! raw !!}'} dot-path interpolation, @if/@else, @foreach … as … (with loop.index/first/last), and @include for partials. - Templates compile once to an AST and report line-numbered errors. It also powers the themed HTML in the - mail layer.

    {@render code(cTemplates, 'templates')} +

    + Templates compile once to an AST and report line-numbered errors. It also powers the themed HTML in the + mail layer. +

    - +
    +

    Collections

    +

    + collect(..) wraps any iterable in a Collection<T> — a fluent, chainable API + for the everyday shaping that raw Iterator makes verbose. It is part of the facade, no feature + needed. +

    + {@render code(cCollections, 'collections')} +

    + It is a thin layer over Vec<T>: it Derefs to [T] and + round-trips through Vec and iterators, so it adds no allocation over doing the work by hand — + and you can drop back to plain iterator code at any point in a chain. +

    +
    + +
    +

    Crypto primitives

    +

    + sutegi::crypto is always compiled — it is what the sessions, auth, SigV4 signing, SCRAM + handshake and WebSocket handshake are built from — and it is directly usable. +

    + {@render code(cCrypto, 'crypto')} +

    + The choices are deliberate: ChaCha20-Poly1305 over AES because add-rotate-xor is constant-time in plain + software where AES table lookups leak through cache timing; a fresh random nonce prepended by + seal so nonce reuse is impossible by construction; HKDF so a signing key and an encryption + key derived from one master secret never alias; and constant_time_eq piped through + black_box so the optimizer cannot rewrite it into an early exit. The raw stream cipher and + one-shot Poly1305 stay private, because those are the pieces that are dangerous to hold directly. +

    +
    +
    Read the posture page first
    +

    + This is hand-rolled cryptography, known-answer tested and fuzzed but + not independently audited. Use it inside a trusted boundary; + do not make it the only thing between hostile traffic and your secrets. +

    +
    +
    + +
    -
    Architecture & deployment
    +
    Architecture & operations

    Hexagonal architecture

    - As an app grows, the hex toolkit keeps it honest. Your domain stays pure; the application + As an app grows, the hexagon toolkit keeps it honest. Your domain stays pure; the application layer depends on port traits (an outbound TodoRepository, say); and adapters — an - HTTP route, an AI tool, a repo over either Backend — plug in at the edges. UseCase - is the inbound-port trait, AppError/AppResult are transport-agnostic with a - canonical HTTP mapping, and respond/respond_created are the glue that turns an - AppResult into a Response. One use case can back both a route and a tool, over - whichever store the composition root injects — and it is fully testable without starting a server. + HTTP route, an AI tool, a repo over either Backend — plug in at the edges. + UseCase is the inbound-port trait, AppError/AppResult are + transport-agnostic with a canonical HTTP mapping, and respond/respond_created are + the glue that turns an AppResult into a Response.

    {@render code(cHex, 'hex')} +

    + One use case can back both a route and a tool, over whichever store the composition root injects — and it + is fully testable without starting a server. Command, Query, Event + and EventBus are there when you want the CQRS vocabulary too. +

    Full guide

    docs/HEXAGONAL.md - covers the dependency rule, layer responsibilities, layout, and testing strategy in depth. + covers the dependency rule, layer responsibilities, layout, and testing strategy in depth, and + examples/hexagonal is a worked reference with two interchangeable repositories + (in-memory ↔ SQLite) selected at the composition root.

    @@ -1201,51 +2571,216 @@ App::new("api")

    App::service() returns the app as a plain Fn(Request) -> Response, so you can exercise the whole routing, state, validation, and tool surface in process — no socket, no port, - no async harness. Back it with Db::memory() for a fresh database per test. For full - end-to-end coverage, the framework’s own suite boots a real server over a loopback socket; see - crates/sutegi/tests/server.rs for the pattern. + no async harness. Back it with Db::memory() for a fresh database per test.

    {@render code(cTesting, 'testing')} +

    + The pieces that need a server or a real engine are testable too, and the framework’s own suite shows + how: Mailer’s memory driver for mail assertions, the queue’s + caller-supplied millisecond timestamps so schedules need no sleeping, an in-process S3 stub over a real + socket for the storage wire suite, and crates/sutegi/tests/server.rs for the full loopback + end-to-end pattern. +

    + + +
    +

    The REPL

    +

    + The repl feature is a tinker-style interactive shell. It is deliberately + not a Rust evaluator — that would need a compiler toolchain and + third-party machinery — but a command shell over the surfaces your app already exposes: the agent contract + (introspection, the tool manifest, tool invocation with SSE frames printed live, raw HTTP through your own + routes) and, with a Backend attached, the data layer (raw SQL, a where-clause + query DSL, KV, the event store, the job queue). +

    + {@render code(cRepl, 'repl')} +

    + Two transports, one command set. In-process it consumes the built App; remote mode drives a + running app over plain HTTP with no source access — exactly the way an LLM does — which is why + sutegi repl <addr> works against any sutegi app. Repl::eval(line) is the + programmatic seam the loop is built on, so the same commands are scriptable and testable. +

    +
    + +
    +

    Inside the server

    +

    + There is no executor and no hidden machinery: a fixed thread pool accepts connections and each worker + handles one connection with blocking I/O. That is the whole model, and it is why streaming is trivial, + backpressure is free, and a stack trace points at your handler rather than at a runtime. +

    + {@render code(cInternals, 'internals')} +

    + The consequences are worth internalising, because they are what you tune: +

    +
      +
    • A connection costs a thread while it is open. Keep-alive is therefore capped twice — keep_alive_idle (default 5 s, deliberately much shorter than the socket timeout) and keep_alive_max (default 100 requests) — so an idle client cannot pin a worker indefinitely.
    • +
    • Reading a request is bounded twice too. timeout bounds a single stalled recv; header_timeout bounds the total time a peer may take to deliver request line, headers and body, so a byte-per-interval dripper cannot hold a worker (slowloris, CWE-400).
    • +
    • Panics are isolated per request. A handler panic becomes a 500 rather than a downed worker — in dev and test builds, where the workspace unwinds. Release builds use panic = "abort" to drop the unwind tables, which moves that responsibility to your pod supervisor.
    • +
    • WebSockets escape the model entirely. An upgrade returns Body::Upgrade, the reactor adopts the socket, and the worker is released — which is how one process holds 80k connections on eight threads.
    • +
    • The release profile is tuned for size: opt-level = "z", LTO, one codegen unit, no unwind tables, stripped symbols. That is where ~394 KB comes from.
    • +
    +

    + Ordering is fixed and worth knowing when you place a guard: ops guard → global middleware → + /__metrics and /__introspect → route match → group middleware → handler → + after-middleware. The two probes are matched before everything, so an orchestrator never needs a + credential. Metrics are recorded on every path, including short-circuits. +

    +
    + +
    +

    Listeners

    +

    + App::listener registers a long-running loop — a UDP ingest port, a raw TCP protocol, a + discovery beacon — that runs on its own thread for the life of the server and shuts down with it. It is a + lifecycle seam, not a protocol: sutegi does not frame your packets. + std::net already gives you UdpSocket and TcpListener; what a + hand-spawned thread cannot do is see shared state, drain gracefully, or show up where an agent can + discover it. This closes those three gaps and nothing more. +

    + {@render code(cListeners, 'listeners')} +

    + ListenerCtx is deliberately tiny: should_stop() (true once shutdown has begun), + state::<T>() / try_state::<T>() — the same typed state handlers and + tools see — db::<B>() as sugar for a backend, and name(). +

    +
    +
    The shutdown contract is cooperative
    +

    + On shutdown the server stops accepting, drains in-flight HTTP requests, then + joins every listener thread before returning. That join waits for + your loop to notice, so never park in an unbounded blocking read: set a read timeout on the socket and + poll should_stop() each lap. A listener that never polls holds the process until the + orchestrator’s grace period expires and SIGKILL lands — honest, but avoidable. +

    +
    +

    + A panicking listener does not restart and does not take the server down: + the panic is caught, one line lands on stderr, and the thread ends. If the loop must survive crashes, own + that policy — wrap the body in a retry loop, or run it under an + actor supervisor and keep the listener as the thin socket shim. + Discovery is name and doc only; which port and protocol the loop speaks belongs in the doc string, where + agents read it. +

    +
    + +
    +

    Tuning & limits

    +

    + Every knob is a builder method with a documented default; nothing here is required to run. The HTTP + defaults are conservative on purpose — raise max_body before you accept large uploads, and + remember that in a thread-per-connection server the keep-alive settings are a capacity decision, not a + cosmetic one. +

    + {@render code(cOptions, 'options')} +
    + + + + + + + + + + + + + + + + + + +
    SettingDefaultWhat it bounds
    workers8HTTP threads; WORKERS in the environment overrides it.
    max_body2 MiBRequest body size; 413 above it.
    max_header_bytes64 KiBTotal header bytes; 413 above it.
    timeout30 sOne stalled socket read or write.
    header_timeout15 sThe whole request delivery, start to finish.
    keep_alive_idle5 sIdle time between requests on one connection.
    keep_alive_max100Requests served per connection.
    ws.shards0 (per core)Reactor threads.
    ws.max_frame / max_message1 MiBOne frame / an assembled message (close 1009 above).
    ws.ping_interval / idle_timeout30 s / 75 sServer ping cadence / drop a silent socket.
    ws.max_buffered1 MiBPer-connection outbound queue; a slower consumer is dropped.
    ws.max_connections / per_ip1,048,576 / 1024Process-wide cap / single-source exhaustion.
    queue.visibility_timeout—How long a claimed job may run before it becomes visible again.
    actor mailbox1024Queued messages before tell fails with Full.
    +
    +

    + Pool sizes are constructor arguments rather than builder methods: Db::open_pool(path, n) and + Pg::from_env(n). Size the Postgres pool with the advisory-lock note in mind — a held lock uses + its own dedicated connection, outside the pool. +

    Operational endpoints

    - Four operational endpoints are always on, no feature required. /__health is liveness (200 - while the process is up), /__ready runs the probe you register with .readiness(…) + Four endpoints are always on, no feature required. /__health is liveness (200 while the + process is up), /__ready runs the probe you register with .readiness(…) and returns 200 or 503, /__metrics exposes Prometheus text (requests total, in-flight, by - status class), and /__introspect is the full surface. Because Db is - Clone, clone a handle for the readiness probe before you hand ownership to .state(). + status class), and /__introspect is the full surface — routes, models, tools, capabilities, + searchable tables and listeners. Features add more endpoints — + /__tools, /__channels, /__actors — and you can mount your own, like + a read-only /__migrations.

    {@render code(cOps, 'ops')} +

    + Because Db is Clone, clone a handle for the readiness probe before you hand + ownership to .state(). The probes are intentionally credential-free and disclose nothing; + everything else under /__ should sit behind an + ops_guard in any deployment where the agent + surface is not meant to be public. +

    Deploying

    - .serve() already does the right thing for a rolling update: it drains in-flight requests on - SIGTERM before exiting. ontzi (Basque: vessel) wraps Docker Compose to - run the horizontally-scaled shape locally — N replicas behind an nginx load balancer configured with - proxy_buffering off so SSE streams pass straight through — and promotes the same shape to - Kubernetes with manifests that already wire probes, graceful drain, and Prometheus annotations. For a - single box, a provisioning script installs the binary as a hardened systemd unit behind nginx instead. + .serve() already does the right thing for a rolling update: it traps SIGTERM/SIGINT, stops + accepting new connections, and drains in-flight requests before exiting. (run(addr) serves + forever and run_until(addr, flag) gives manual control without the signal feature.) + ontzi (Basque: vessel) wraps Docker Compose to run the horizontally-scaled shape + locally — N replicas behind an nginx load balancer configured with proxy_buffering off so SSE + streams pass straight through — and promotes the same shape to Kubernetes with manifests that already wire + probes, graceful drain and Prometheus annotations. For a single box, a provisioning script installs the + binary as a hardened systemd unit behind nginx instead.

    {@render code(cDeploy, 'deploy')}

    Pick the backend for the deployment, not the code: one instance runs on SQLite (embedded, zero-ops); many - pods run on Postgres plus the durable queue. The request surface is stateless and scales horizontally - either way. + pods run on Postgres, which is also what turns the queue’s exclusivity, the advisory locks, the + watchers and the channel broker from process-scoped into cluster-scoped. The request surface is stateless + and scales horizontally either way, and the binary is small enough that + requests: 32Mi is a reasonable ask.

    Security posture

    - sutegi ships panic isolation (a handler panic becomes a 500, not a downed worker), - configurable body/header size limits, slowloris socket timeouts, per-IP rate limiting, secure-header and - CORS middleware, bearer/basic guards, and signed-cookie sessions. Passwords are PBKDF2-HMAC-SHA256 PHC - strings; agent tokens are stored hashed. The query builder guards identifiers against injection, and the - fuzz and differential harness runs as a required CI gate. + sutegi ships panic isolation, bounded bodies and headers, two slowloris deadlines, capped keep-alive, + per-IP rate limiting, secure-header and CORS middleware, bearer/basic guards, signed-cookie sessions with + server-side expiry, CSRF tokens, login throttling and hash-bound sessions. Passwords are + PBKDF2-HMAC-SHA256 PHC strings; agent tokens are stored hashed. 5xx messages never reach the client. The + query builder guards identifiers, the search grammar sanitizes input before it can reach engine syntax, + JSON paths and search queries are bound as parameters, and the fuzz and differential harness runs as a + required CI gate. +

    +

    Two traps worth reading even if you never touch the code

    +
      +
    • + Gate on router segments, not on the path string. The ops guard used + to test req.path.starts_with("/__") — but the router trims and splits, so + //__tools/x reached the same route while failing that test, leaving tool invocation and + introspection reachable without the configured credential (CWE-288). The check and the dispatch now + agree by construction. Audit your reverse proxy too: + Apache’s idiomatic <LocationMatch "^/__"> has the identical blind spot, so a + deployment can look doubly protected and be neither. (MergeSlashes, on by default since + 2.4.39, collapses it — which makes the anchored form correct by accident and one directive away from + not being. ^/+__ is free.) +
    • +
    • + A credentialed WebSocket needs check_origin. The + same-origin policy does not stop a cross-origin handshake from carrying cookies. See + WebSockets. +
    • +
    +

    + The S3 path is hardened in the places subprocess-based HTTP usually is not: credentials go to + curl on stdin so Authorization is invisible to ps, the protocol is + pinned with --proto =https, redirect following is off so a 3xx cannot replay a signature at + an attacker-chosen host, TLS floors at 1.2, bodies are capped, certificate verification cannot be + disabled, PUT bodies stage through a 0600 O_EXCL temp file, and + S3Store’s Debug redacts the secret key.

    Read this before deploying to hostile traffic
    @@ -1330,6 +2865,38 @@ App::new("api") margin: 0.85rem 0; } .prose-doc :global(.list li) { margin: 0.5rem 0; } + .prose-doc :global(.tbl-wrap) { + overflow-x: auto; + margin: 1.1rem 0; + border: 1px solid rgba(255, 255, 255, 0.08); + border-radius: 0.6rem; + } + .prose-doc :global(.tbl) { + width: 100%; + border-collapse: collapse; + font-size: 14px; + color: #b4b4c2; + } + .prose-doc :global(.tbl th) { + text-align: left; + font-weight: 600; + color: #cfcfe0; + font-size: 11px; + text-transform: uppercase; + letter-spacing: 0.05em; + padding: 0.7rem 0.9rem; + background: rgba(255, 255, 255, 0.03); + border-bottom: 1px solid rgba(255, 255, 255, 0.1); + white-space: nowrap; + } + .prose-doc :global(.tbl td) { + padding: 0.6rem 0.9rem; + border-bottom: 1px solid rgba(255, 255, 255, 0.05); + vertical-align: top; + line-height: 1.55; + } + .prose-doc :global(.tbl tr:last-child td) { border-bottom: none; } + .prose-doc :global(.tbl td:first-child) { white-space: nowrap; } .prose-doc :global(.callout) { border-radius: 0.6rem; padding: 1rem 1.15rem; diff --git a/landing/src/Landing.svelte b/landing/src/Landing.svelte index 6d1e24a..f569de7 100644 --- a/landing/src/Landing.svelte +++ b/landing/src/Landing.svelte @@ -120,10 +120,11 @@ // --- code snippets (kept as strings so braces are literal text) --- const codeCargo = `[dependencies] -# default = ["derive", "orm", "validate", "ai"] +# default = ["derive", "orm", "validate"] — the agent tool surface is core sutegi = { version = "*", features = ["sqlite", "graceful"] } # single-node -# multi-pod: swap in Postgres + the durable queue instead -# sutegi = { version = "*", features = ["postgres", "queue", "graceful"] } +# multi-pod: Postgres, the durable queue, cross-pod realtime +# sutegi = { version = "*", features = [ +# "postgres", "queue", "channels", "pubsub-postgres", "graceful" ] } # minimal HTTP service, nothing else compiled in: # sutegi = { version = "*", default-features = false }`; @@ -290,21 +291,52 @@ let body = c.validate(&rules)?; // Err -> { "email": ["… valid email …"] } sink.event("done", "{}") }))`, }, + { + id: 'realtime', icon: Server, kicker: 'Step 7b', title: 'Push to the browser', + lead: 'For live UIs, upgrade instead of polling. An upgraded socket detaches from the HTTP worker into a sharded kqueue/epoll reactor — no async runtime — so an idle connection costs ~340 bytes and zero threads (measured: 80k live sockets at 0.0% idle CPU). On top of that, channels gives you Phoenix-style topics, joins, replies and presence, and broadcasts ride the pubsub Broker seam: adding .broker(PgPubSub::connect(&cfg)?) is the only change between one pod and a fleet.', + code: `let hub = Channels::new() + .channel(Channel::new("room:*") + .doc("A chat room. Join with a nick.") + .on_join(|socket, payload| { + socket.assign("nick", payload.pointer("/nick").cloned().unwrap()); + Ok(Json::Null) + }) + .on("new_msg", |socket, payload| { + socket.broadcast("new_msg", payload); // all members, all pods + Reply::None + })) + // .broker(PgPubSub::connect(&pg_cfg)?) // <- cross-pod, that's it + .check_origin(["https://app.example.com"]) // MUST set if cookies auth it + .build(); + +App::new("chat") + .channels("/channels", "The chat socket.", hub.clone()) + .serve() // + GET /__channels: the manifest an agent joins from`, + tip: 'A bundled ~4 KB dependency-free JS client (sutegi_channels::JS_CLIENT) handles reconnect, rejoin and heartbeats. And watch(query) turns a topic into a live query: any pod’s committed write pushes a Change {added, updated, removed} diff you can broadcast straight through.', + }, { id: 'jobs', icon: Cpu, kicker: 'Step 8', title: 'Defer the slow work', - lead: 'Some work shouldn’t block the response — sending mail, calling a webhook. The durable queue (queue feature) is Postgres-backed and cross-pod: it claims jobs with FOR UPDATE SKIP LOCKED, so any pod can pull the next one, and visibility-timeout retries recover from a crashed worker. Register a handler by name, dispatch a JSON payload, and start N worker threads.', - code: `use sutegi::queue::{Queue, Workers}; + lead: 'Some work shouldn’t block the response — sending mail, calling a webhook. The durable queue (queue feature) runs over the same Backend seam, so one jobs table and one set of SQL work on bundled SQLite and on Postgres. The claim stamps a lease instead of deleting the row, so a crashed worker’s job becomes visible again after the visibility timeout. Register a handler by name, dispatch a payload, and start N worker threads.', + code: `use std::sync::Arc; +use sutegi::queue::Queue; -let mut queue = Queue::new(Pg::from_env(8)?.pool().clone()); +let mut queue = Queue::new(db.clone()); // any Backend: Db or Pg queue.migrate()?; // creates sutegi_jobs -queue.register("notify", |payload| { - let to = payload.get("to").and_then(Json::as_str).unwrap_or(""); +queue.register("notify", |job| { + let to = job.payload().get("to").and_then(Json::as_str).unwrap_or(""); /* send … */ Ok(()) // Err -> retried w/ backoff }); queue.dispatch("notify", Json::obj(vec![("to", Json::str("a@b.com"))]))?; -let workers: Workers = std::sync::Arc::new(queue).start(4); // cross-pod`, - note: 'The old in-process Job trait and Queue::new(4) worker-pool are gone — the queue is durable now, so a job survives a pod restart and dead-letters after its retries are spent.', + +// …or shape it: its own pool, one in flight per key, 3 tries. +queue.job("video.ingest", payload).queue("video").unique("yt:abc") + .priority(10).max_attempts(3).dispatch()?; + +let queue = Arc::new(queue); +let fast = Arc::clone(&queue).start(4); // 4 on "default" +let slow = Arc::clone(&queue).start_on("video", 1); // 1 on the slow queue`, + note: 'Claims are exclusive on both backends — FOR UPDATE SKIP LOCKED on Postgres, and on SQLite the serialized writer already gives it. Postgres additionally makes that exclusivity cross-pod; queue.cross_pod() tells you which guarantee you actually have.', }, ], }, @@ -449,7 +481,7 @@ App::new("api")
    0
    runtime deps
    -
    362 KB
    core binary
    +
    394 KB
    core binary
    std
    only
    @@ -553,8 +585,12 @@ App::new("api") { icon: FileCode, t: '#[derive(Model, Validate)]', d: 'Schema, migrations, JSON and the validation ruleset — all from one struct, at build time.' }, { icon: Plug, t: 'Kv store', d: 'Namespaced JSON key/value over either backend: set / get / scan / delete.' }, { icon: Radio, t: 'Streaming & SSE', d: 'Stream bytes or Server-Sent Events with natural backpressure; same transport as stream tools.' }, - { icon: Cpu, t: 'Durable queue', d: 'Postgres-backed, cross-pod: SKIP LOCKED claim, visibility-timeout retries, dead-letter.' }, + { icon: Cpu, t: 'Durable queue', d: 'Over the Backend seam — SQLite or Postgres: leased claims, retries, priorities, dedupe, dead-letter.' }, { icon: Zap, t: 'Agent-native', d: '.tool() / .stream_tool() closures auto-mount /__tools; /__introspect exposes the whole surface.' }, + { icon: Server, t: 'Realtime', d: 'RFC 6455 sockets on a kqueue/epoll reactor (80k idle conns, zero threads), Phoenix-style channels + presence, cross-pod over PG.' }, + { icon: Layers, t: 'Actors & supervision', d: 'Typed mailboxes, OTP restart policies and strategies, live status at /__actors.' }, + { icon: Search, t: 'Search & embeddings', d: 'One grammar over tsvector and FTS5, vector columns, and RRF hybrid retrieval in one call.' }, + { icon: ShieldCheck, t: 'Production auth', d: 'PBKDF2 passwords, remember-me, login throttling, CSRF, hash-bound sessions, hashed agent tokens.' }, ] as f} {@const Icon = f.icon}
    @@ -578,7 +614,7 @@ App::new("api") {#each [ { icon: Zap, t: 'Agent tool servers', d: 'Expose capabilities to an LLM with a built-in manifest and validated invocation — no glue layer.' }, { icon: Boxes, t: 'Internal microservices', d: 'Start single-node on SQLite; scale to many pods on Postgres + the durable queue without rewriting handlers.' }, - { icon: Plug, t: 'Edge & embedded', d: 'A ~362 KB binary with no async runtime and one embedded SQLite file fits where a full stack will not.' }, + { icon: Plug, t: 'Edge & embedded', d: 'A ~394 KB binary with no async runtime and one embedded SQLite file fits where a full stack will not.' }, { icon: FileCode, t: 'LLM-generated apps', d: 'Rigid scaffolding conventions mean a model can extend the codebase correctly with minimal context.' }, { icon: Server, t: 'JSON APIs & CRUD', d: 'Routing + typed models + validation cover the everyday backend without pulling a framework zoo.' }, { icon: Radio, t: 'Streaming endpoints', d: 'SSE token streams for chat/AI UIs, backpressured by the thread-per-connection model.' }, @@ -704,6 +740,14 @@ App::new("api")
    Hexagonal guide →
    The dependency rule, layer responsibilities, layout, and testing strategy.
    + +
    Inside the server →
    +
    The request lifecycle, panic isolation, keep-alive, and every tuning knob.
    +
    + +
    Backends & capabilities →
    +
    Locks, isolation, JSON paths, search, live queries — and which ones your store actually has.
    +