- Miss path only marks a response a durable replay target (outcome 1) when status is 2xx/3xx; a 4xx/5xx gets outcome 3 instead - Outcome 3's row is a one-shot relay: the scalar reply still can't carry a Resp (WO-E226), so the row exists only to hand the exact response back once, then idempotent.wo deletes it -- a retry with the same key is a genuine miss and re-executes, instead of caching a 500 for the 24h default TTL - Reviewer finding: caching any status meant a transient failure was replayed verbatim until TTL expiry, worse than no idempotency at all - Gate leg 18d: FlakyHandler fails once then succeeds; same key twice must answer 500 then 200 -- confirmed failing (500, 500) before the fix, passing after Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> (cherry picked from commit e61015f2065a7c6aec6f5c2439e78b53f796bab3)
263 lines
11 KiB
Text
263 lines
11 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.
|
|
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 {
|
|
-- Same request (the owner's own retry, or a duplicate that
|
|
-- waited in the mailbox): the row already holds the response.
|
|
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);
|
|
}
|
|
}
|
|
|
|
let resp = msg.handler.handle(msg.req);
|
|
let hdrs: map<Text, Text> = {};
|
|
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
|
|
});
|
|
-- unguarded, same as kind 1's own insert: this actor is the only
|
|
-- writer for this key (messages are processed one at a time), so a
|
|
-- @unique violation here would mean something is genuinely wrong,
|
|
-- not a race to paper over. A swallowed failure would answer a
|
|
-- stored/ephemeral outcome for a response that was never actually
|
|
-- written -- let it trap instead, so a saturated-looking 503 is
|
|
-- what the middleware answers, never a false success.
|
|
insert IdempotencyKey {
|
|
key: msg.key, response: stored_json, created_at: now, digest: msg.digest
|
|
};
|
|
if resp.status >= 200 and resp.status < 400 {
|
|
-- 2xx/3xx: a real answer worth replaying for the TTL.
|
|
return pool_pack(1, 0);
|
|
}
|
|
-- 4xx/5xx: the row above exists only so the scalar-only reply can
|
|
-- still hand the caller its exact response (WO-E226 -- a Resp
|
|
-- cannot ride the mailbox). It must NOT survive to answer a later
|
|
-- retry: caching a transient 500 for the TTL (default 24h) would
|
|
-- make every retry fail until it expires, worse than no idempotency
|
|
-- at all. idempotent.wo reads this row once and deletes it.
|
|
return pool_pack(3, 0);
|
|
}
|
|
|
|
-- 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).
|
|
pub fn make_pool(n: Int) -> Pool {
|
|
let actors: multi PoolSlot = [];
|
|
let i = 0;
|
|
while i < n {
|
|
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 a 2xx/3xx
|
|
-- response is durably in IdempotencyKey (fresh store or matched replay);
|
|
-- 2 means a digest mismatch (422, nothing to read); 3 means a 4xx/5xx
|
|
-- miss whose row is a one-read-then-delete relay only, never a replay
|
|
-- target (see the kind-2 arm above). 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
|
|
});
|
|
}
|