writeonce/docs/examples/porch/middleware/keypool.wo
shoney.arickathil 359d21a57f fix(porch-store): widen idempotent replay headers, real per-conn sharding
- stored/replayed headers widen from content-type only to an allowlist
  (content-type, location, etag, cache-control), matched case-insensitively
  -- a redirect() lost its Location on its own first response, not just replay
- add pool_slots(Pool) -> multi PoolSlot and pool_of(multi PoolSlot) -> Pool
- Pool is demand-promoted to traced (WO-E222) and can't live in actor
  state; PoolSlot/multi PoolSlot never is, the same shape chat/main.wo's
  Room already holds directly -- this is what lets an app actually shard
  across N actors per connection instead of a forced one-slot pool
- log a genuine pool_select trap instead of silently folding it into 503
- fix stale comments: the prune below IS a delete-then-insert (of a
  fresh row, not the same one) contradicting the doc comment above it;
  the catch shape referenced in two comments had changed

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
(cherry picked from commit b738269314f01a95dee1341437c7661ce9e28730)
2026-09-15 01:15:30 +02:00

344 lines
15 KiB
Text

-- porch/middleware/keypool.wo — the key pool: an actor per shard, picked by
-- hash of the key, that serializes rate-limit counting (this file, kind 1)
-- and idempotency begin (Task 4, kind 2) against the @table rows in
-- store.wo. This is the only file that knows a pool exists — the
-- middlewares call through it and never touch RateLimitCounter
-- themselves. IdempotencyKey is the one exception: the response has to
-- travel through that table (a Resp cannot ride the mailbox — see
-- below), so idempotent.wo reads the row a begin call already committed.
--
-- `call`'s reply crosses the actor boundary as a single copyable scalar
-- (WO-E226 — no class, no Text can ride it). The exact count is decided
-- atomically inside `receive`; `pool_count` packs it with the window's
-- remaining time into one Int and unpacks that into the `Verdict` callers
-- actually read, so the packing never leaks outside this file. kind 2
-- (Task 4) reuses the exact same pool_pack scheme for its outcome code —
-- WO-E226 forces every `receive` in the program to agree on one return
-- type, so a second encoding is not an option.
use time
use http
use json
-- To a pool actor. kind 1 = count (this file); kind 2 = begin (Task 4
-- fills in the arm — the fields below are already shaped for it: the
-- bare idempotency key travels in `key`, the body digest in `digest`,
-- and the actor runs `handler` against `req` itself so a duplicate waits
-- in the mailbox rather than needing a held reply).
class PoolMsg {
kind: Int
key: Text
limit: Int -- count: max requests per window
window: Int -- count: window size, µs
digest: Text -- begin: sha256(method|path|body), Task 4
req: Req -- begin: the request, Task 4
handler: Handler -- begin: the route's handler, invoked inside receive, Task 4
}
-- What the limiter reads back from a count. `allowed` and `limit` are
-- filled in by `pool_count` — the caller already knows `limit`, it is the
-- one it sent. `count` and `reset_at` come from the actor.
class Verdict {
allowed: Bool
count: Int
limit: Int
reset_at: Int -- wall-clock ms (time.now()) when this key's window resets
}
-- PoolMsg requires `req`/`handler` on every construction (an actor
-- message's fields are all required, like RoomMsg's `writer` in
-- docs/examples/chat/main.wo). A count message has no request to run, so
-- it fills those two with an inert placeholder — same shape as chat's
-- dummy_writer() for RoomMsg's shutdown message.
class NullHandler {
fn handle(req: Req) -> Resp {
return Resp { status: 500, headers: {}, body: "" };
}
}
fn dummy_req() -> Req {
return Req {
method: "", path: "", params: {}, query: {}, headers: {},
body: "", principal: "", ctx: {}, conn: 0 - 1
};
}
-- A live Req arriving at a Handler is a borrow (Handler.handle's signature
-- fixes that, in router.wo — not this file's to change): its map fields
-- are references that cannot outlive the caller's scope, so forwarding
-- them as-is into an actor message is refused (WO-E222 — the same
-- aliasing rule Pool's own doc comment above describes). Copying each map
-- field into a brand-new map, then building a brand-new Req from that plus
-- the plain scalars, produces a value with no other referrer — the same
-- shape dummy_req() already sends, just carrying the real request.
fn fresh_req(r: Req) -> Req {
let params: map<Text, Text> = {};
for k, v in r.params { params[k] = v; }
let query: map<Text, Text> = {};
for k, v in r.query { query[k] = v; }
let headers: map<Text, Text> = {};
for k, v in r.headers { headers[k] = v; }
let ctx: map<Text, Text> = {};
for k, v in r.ctx { ctx[k] = v; }
return Req {
method: r.method, path: r.path, params: params, query: query,
headers: headers, body: r.body, principal: r.principal, ctx: ctx,
conn: r.conn
};
}
-- One actor per shard. Reads the row for the key, decides, and writes the
-- new count by assigning to the row's field — that writes through and
-- maintains indexes; never delete-then-insert as an update of the SAME
-- row (the window prune below IS a delete-then-insert, but of a fresh
-- row for the new window — the stale row is retired, not mutated).
class KeyActor {
fn receive(msg: PoolMsg) -> Int {
if msg.kind == 2 {
-- Task 4: idempotency begin. msg.window carries the TTL here (both
-- are µs durations; kind 1 has no use for a TTL and kind 2 has no
-- use for a window, so the one field serves both). A miss runs
-- msg.handler right here, inside receive, so a duplicate already
-- queued behind this message dequeues to a settled row instead of
-- a race.
let now = time.ticks();
let hits = from k in IdempotencyKey where k.key == msg.key take 1 select k;
if len(hits) > 0 {
let stored = hits[0];
if now - stored.created_at > msg.window {
-- Expired: lazy delete (no sweeper exists), fall through to miss.
delete stored;
} else if stored.digest == msg.digest {
-- A row lives under the BARE key only when it is durable (see
-- below) -- an ephemeral one never does -- so any hit here is
-- already a stable, valid replay target.
return pool_pack(1, 0);
} else {
-- Same key, a different request: refuse rather than serve the
-- other request's response.
return pool_pack(2, 0);
}
}
-- Miss: fresh key, an expired row just deleted, or the prior
-- attempt (if any) was ephemeral and so is invisible to the
-- bare-key lookup above -- all three run the handler fresh.
let resp = msg.handler.handle(msg.req);
-- Allowlist, not denylist: an allowlist fails safe when Resp grows a
-- new header later (excluded from replay until reviewed, never
-- replayed by accident from day one). content-type alone dropped
-- Location off every redirect() and any Etag/Cache-Control a handler
-- set; matched case-insensitively since set_header writes the name
-- verbatim (a handler using "Content-Type" was silently losing it).
let hdrs: map<Text, Text> = {};
for hk, hv in resp.headers {
let lhk = to_lower(hk);
if lhk == "content-type" or lhk == "location" or lhk == "etag" or lhk == "cache-control" {
hdrs[lhk] = hv;
}
}
let stored_json = json.encode(IdempotentStoredResp {
status: resp.status, headers: hdrs, body: resp.body
});
if resp.status >= 200 and resp.status < 400 {
-- 2xx/3xx: a real answer worth replaying for the TTL, stored
-- under the bare key -- one stable row per key, unguarded same
-- as kind 1's own insert (this actor is the only writer for
-- this key; a @unique violation here would mean something is
-- genuinely wrong, not a race to paper over -- a swallowed
-- failure would answer "stored" for a response never written).
insert IdempotencyKey {
key: msg.key, response: stored_json, created_at: now, digest: msg.digest
};
return pool_pack(1, 0);
}
-- 4xx/5xx: never a replay target, so it does NOT go under the bare
-- key -- a bare-key row is memoryless (this file's own doc above),
-- but a SHARED, mutable row is not: a second message racing the
-- first could delete-and-replace it before the first caller's own
-- middleware-side read (necessarily outside receive -- WO-E226,
-- a Resp cannot ride the mailbox) ever runs, so the FIRST caller
-- could read back the SECOND caller's answer. Every miss instead
-- gets its own row, keyed by a nonce carried back in the scalar's
-- low digits (the same slot pool_pack's remaining_ms uses for
-- kind 1) so idempotent.wo can reconstruct the exact same key and
-- read only ever what THIS call produced -- immune to any other
-- message touching this bare key, ever. That same unguessability
-- (the nonce is never handed to anyone but this one call() reply)
-- is also why idempotent.wo deletes this row right after reading
-- it: nothing else can ever construct this exact key, so nothing
-- else is deleted out from under. Without that delete, the row
-- would linger forever (no sweeper exists) AND time.ticks() % 1e9
-- wraps every ~1000s, so a later failed attempt for the SAME
-- bare key landing on the same nonce would collide with it --
-- reintroducing a stale-replay risk on wraparound, or poisoning
-- this actor's own unguarded insert above. Deleting it removes
-- both, not just the storage growth.
let nonce = now % 1_000_000_000;
insert IdempotencyKey {
key: "${msg.key}#eph:${nonce}", response: stored_json, created_at: now, digest: msg.digest
};
return pool_pack(3, nonce);
}
-- kind 1: count.
let now = time.ticks();
let hits = from c in RateLimitCounter where c.key == msg.key take 1 select c;
if len(hits) == 0 {
insert RateLimitCounter { key: msg.key, count: 1, window: now };
return pool_pack(1, msg.window / 1000);
}
let row = hits[0];
if now - row.window > msg.window {
-- the window fully elapsed: prune the stale row rather than reset it
-- in place — resetting keeps one row forever for every key ever
-- seen, an unbounded leak for IP-keyed limiting. There is no
-- sweeper; this lazy expiry on access is it.
delete row;
insert RateLimitCounter { key: msg.key, count: 1, window: now };
return pool_pack(1, msg.window / 1000);
}
row.count = row.count + 1;
let remaining_us = row.window + msg.window - now;
if remaining_us < 0 { remaining_us = 0; }
return pool_pack(row.count, remaining_us / 1000);
}
}
-- Packs (count, remaining-ms-in-window) into one Int: count * 1e9 +
-- remaining_ms, remaining_ms clamped to stay under 1e9 (~11.5 days —
-- far past any realistic rate-limit window). That clamp only blurs the
-- advisory reset header on an absurdly long window; it never touches the
-- count, which is the correctness-critical half.
fn pool_pack(count: Int, remaining_ms: Int) -> Int {
let r = remaining_ms;
if r < 0 { r = 0; }
if r >= 1_000_000_000 { r = 999_999_999; }
return count * 1_000_000_000 + r;
}
-- One actor address per slot. `multi actor PoolMsg` does not parse (a
-- `multi`'s element type is one token) — chat/main.wo's RoomRef wraps an
-- actor handle in a one-field class for exactly this reason, mirrored
-- here as PoolSlot.
class PoolSlot {
a: actor PoolMsg
}
class Pool {
actors: multi PoolSlot
}
-- Spawns n identical actors and returns the pool. n is a capacity knob:
-- too small and a hot key's mailbox saturates under load (a `call` trap,
-- answered 503 by the middleware — never a silent bypass). n < 1 is a
-- caller misconfiguration, not a capacity choice, and guarding it HERE
-- (not in pool_select's division) is what matters: every pool_select call
-- runs inside the middleware's own `try ... catch (e) { print_err(...);
-- nil }`, so a mod-by-zero trap there would be misreported as ordinary
-- 503 saturation forever (though now at least logged, not silently
-- swallowed), never surfacing the real bug on its own.
pub fn make_pool(n: Int) -> Pool {
let count = n;
if count < 1 { count = 1; }
let actors: multi PoolSlot = [];
let i = 0;
while i < count {
push(actors, PoolSlot { a: spawn KeyActor {} });
i = i + 1;
}
return Pool { actors: actors };
}
-- Pool itself is demand-promoted to traced (WO-E222) the moment an app
-- aliases it — e.g. Limiter/Idempotent's own `pool: Pool` field, read on
-- every request without being consumed — so it can never live in an
-- actor's state or a message. PoolSlot is not: WO-E222's contains_traced
-- check only recurses into a field typed as a class name (or a `multi`/
-- `map` of one); `a: actor PoolMsg` is an actor handle, a different case
-- entirely, so it never pulls PoolSlot (or `multi PoolSlot`) into the
-- traced set the way wrapping it in Pool does. An actor CAN hold `multi
-- PoolSlot` directly in its own state — the exact shape chat/main.wo's
-- `Room { members: multi Mem }` already uses for a multi of actor
-- handles — which is what makes real per-connection sharding possible:
-- call make_pool(n) ONCE at process start, hand pool_slots(pool) to every
-- connection actor's spawn, and each one rebuilds a transient Pool via
-- pool_of(self.slots) wherever Limiter/Idempotent needs one. Calling
-- make_pool per connection instead (the natural misreading of this pair
-- sitting right after a capacity-sizing knob) gives every connection its
-- own actors and silently restores the lost-increment race this whole
-- design exists to prevent.
--
-- Both functions copy field-by-field, the same trick fresh_req uses above:
-- an actor handle is a plain, freely-copyable scalar (not traced), so
-- rebuilding each PoolSlot by value produces a list with no lingering
-- alias into the traced Pool (pool_slots) or the caller's own copy
-- (pool_of) — never a value some other reader could still be holding.
pub fn pool_slots(p: Pool) -> multi PoolSlot {
let out: multi PoolSlot = [];
for s in p.actors { push(out, PoolSlot { a: s.a }); }
return out;
}
pub fn pool_of(s: multi PoolSlot) -> Pool {
let out: multi PoolSlot = [];
for x in s { push(out, PoolSlot { a: x.a }); }
return Pool { actors: out };
}
-- Hashes a key to one of the pool's actors — sum of bytes modulo n, a
-- shard selector, not a security hash. The same key always selects the
-- same actor, which is the entire per-key serialization mechanism.
pub fn pool_select(pool: Pool, key: Text) -> actor PoolMsg {
let sum = 0;
let i = 0;
while i < len(key) {
sum = sum + byte_at(key, i);
i = i + 1;
}
let idx = sum % len(pool.actors);
return pool.actors[idx].a;
}
-- The count accessor every later task's limiter calls. Unpacks the
-- actor's scalar reply into the Verdict the limiter reads.
pub fn pool_count(pool: Pool, key: Text, limit: Int, window: Int) -> Verdict {
let a = pool_select(pool, key);
let raw = call(a, PoolMsg {
kind: 1, key: key, limit: limit, window: window,
digest: "", req: dummy_req(), handler: NullHandler {}
});
let count = raw / 1_000_000_000;
let remaining_ms = raw % 1_000_000_000;
return Verdict {
allowed: count <= limit,
count: count,
limit: limit,
reset_at: time.now() + remaining_ms
};
}
-- The begin accessor idempotent.wo calls: hides pool_select/call the same
-- way pool_count does. Returns the raw packed outcome — 1 means the
-- BARE key names a durable (2xx/3xx) replay target; 2 means a digest
-- mismatch (422, nothing to read); 3 means this call's own 4xx/5xx
-- answer lives under "${key}#eph:${raw % 1_000_000_000}" instead --
-- the low digits of the reply are that row's own nonce, not a window
-- size (kind 1's use of the same slot), so idempotent.wo can rebuild
-- the exact key and read only ever what THIS call produced. Never 0:
-- idempotent.wo's own catch-and-log-nil handler cannot tell a literal 0
-- reply apart from a trapped call, so the encoding avoids it on
-- purpose. A trapped call (a saturated mailbox) propagates to the
-- caller uncaught, same as pool_count -- the middleware's own
-- try/catch answers 503.
pub fn pool_begin(pool: Pool, key: Text, digest: Text, ttl: Int, req: Req, handler: Handler) -> Int {
let a = pool_select(pool, key);
return call(a, PoolMsg {
kind: 2, key: key, limit: 0, window: ttl,
digest: digest, req: fresh_req(req), handler: handler
});
}