-- 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 = {}; for k, v in r.params { params[k] = v; } let query: map = {}; for k, v in r.query { query[k] = v; } let headers: map = {}; for k, v in r.headers { headers[k] = v; } let ctx: map = {}; 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. 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); let hdrs: map = {}; let ct = resp.headers["content-type"]; if ct != nil { hdrs["content-type"] = ct; } 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) nil`, so a -- mod-by-zero trap there would be swallowed and misreported as ordinary -- 503 saturation forever, never surfacing the real bug. 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 }; } -- 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 `try ... catch (e) nil` 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 }); }