feat(porch-store): idempotency rebuilt so the actor runs the handler

- Delete before/after flow: it stored on a miss, so two duplicates both
  missed and both ran; its 10s "in flight" check fired on fast legit
  replays and never on a real collision
- Idempotent now wraps the route's Handler and hands request + handler to
  the pool; a duplicate waits in the actor's mailbox, not a held reply
- keypool.wo kind-2 arm: digest match replays, mismatch refuses (422), a
  miss runs the handler inside receive and stores status/body/
  content-type, all via the same pool_pack(count, remaining_ms) scalar
  kind-1 uses (WO-E226 forces one return type)
- Outcome codes start at 1, never 0: idempotent.wo's try/catch cannot
  tell a literal 0 reply apart from a trapped call
- fresh_req() copies a borrowed Req's map fields into a new Req before it
  crosses the actor boundary (WO-E222: aliased graphs can't cross heaps)
- Reading a stored row back forces fresh Text via `.. ""` on every field
  copied out of json.decode's result -- decoded Text does not survive
  being handed onward once the decoded record goes out of scope
- insert is unguarded (kind 1's own convention): a swallowed failure
  would answer "stored" for a response never written
- web-app-accept.sh: leg 18a/b/c -- byte-identical replay off an ExecMark
  row count, digest mismatch is 422, genuinely parallel duplicates run
  the handler exactly once

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
(cherry picked from commit eae1b06cdc38e4766d4e27a66b02264f47164e99)
This commit is contained in:
shoney.arickathil 2026-08-30 01:13:54 +02:00
parent 37a63192c6
commit c97de237ef
3 changed files with 333 additions and 96 deletions

View file

@ -1,110 +1,85 @@
-- porch/middleware/idempotent.wo — idempotency middleware backed by @table.
-- Stores successful responses keyed by Idempotency-Key header (+ optional body digest).
-- Replays stored response on subsequent requests with same key.
-- Iteration 1 of the porch track.
use time
-- porch/middleware/idempotent.wo — idempotency, actor-run.
-- Iteration 1 of the porch track, porch-store task 4.
--
-- Rebuilt, not patched: the old before/after shape ran the handler and
-- stored the response afterward, so two simultaneous duplicates both
-- missed and both executed — before() has nothing to find until after()
-- runs, which is too late for the request that is racing it. The fix is
-- that the ACTOR runs the handler: Idempotent wraps the route's own
-- Handler and hands both the request and that handler to the pool. A
-- duplicate for the same key then simply waits in the actor's mailbox
-- (one message at a time) and is dequeued once the owner's receive has
-- already committed the row — no in-flight heuristic needed, because
-- there is no window where a duplicate can see "nothing yet".
--
-- `call`'s reply is a copyable scalar only (WO-E226): keypool.wo's kind-2
-- arm answers with the SAME pool_pack(count, remaining_ms) encoding kind-1
-- uses, never the response itself. The response travels through the
-- @table instead — the actor stores it, this file reads the same row
-- back and builds the Resp from it, so owner and duplicate answer with
-- byte-identical bytes structurally, not by careful bookkeeping.
use http
use json
-- Idempotent middleware: before (check/replay) + after (store on miss).
-- Key composition: "idem:${header}" or "idem:${header}:${sha256(method|path|body)}"
-- Only stores status, body, content-type (allowlist). Never replays Set-Cookie, Date, etc.
-- In-flight collision: returns 409 if key exists but response not yet stored.
-- Lazy expiry: deletes expired keys on access (TTL default 24h).
-- Idempotent wraps a route's Handler. Registration: Idempotent { key_header:
-- "Idempotency-Key", pool: p, inner: CreateOrder {} } in place of the bare
-- handler in app.post(...). Only the digest allowlists content-type for
-- replay (never Set-Cookie, never Date) — cookies arrive in porch 2.
pub class Idempotent {
key_header: Text -- e.g., "Idempotency-Key"
include_body: Bool = true -- digest method+path+body into key
ttl: Int = 86_400_000_000 -- 24h in µs
pool: Pool
inner: Handler -- the real route handler; the actor runs this
include_body: Bool = true -- digest method+path+body, refuse a key reused with a different one
ttl: Int = 86_400_000_000 -- 24h in µs, lazy-expired on access
fn before(mut req: Req) -> ?Resp {
fn handle(req: Req) -> Resp {
let header_val = req.headers[self.key_header];
if header_val == nil { return nil; }
if header_val == nil { return self.inner.handle(req); }
let key = idempotent_key(self, header_val, req);
let now = time.ticks();
-- Look up existing key
let hits = from k in IdempotencyKey where k.key == key take 1 select k;
if len(hits) > 0 {
let stored = hits[0];
-- Check expiry
if now - stored.created_at > self.ttl {
-- Expired: delete and treat as miss
delete stored;
} else {
-- Check if response is stored (created_at within last 10s = in-flight)
if now - stored.created_at < 10_000_000 {
-- In-flight collision: another request with same key is being processed
let r = Resp { status: 409, headers: {}, body: "{\"error\":\"idempotency key in flight\"}" };
set_header(r, "content-type", "application/json");
return r;
}
-- Valid stored response: decode and replay
let resp_json = stored.response;
-- Parse JSON response (status, headers, body)
let resp = json.decode(resp_json) as IdempotentStoredResp;
if resp != nil {
let r = Resp { status: resp.status, headers: resp.headers, body: resp.body };
return r;
}
-- Corrupted stored response: delete and fall through to miss
delete stored;
}
let key = "idem:${self.key_header}:${header_val}";
let digest = "";
if self.include_body {
let digest_input = "${req.method}|${req.path}|${req.body}";
digest = base64_encode(bytes_slice(sha256(bytes_of_text(digest_input)), 0, 16));
}
-- Miss: mark request so after() knows to store the response
req.ctx["idem_miss"] = "true";
req.ctx["idem_key"] = key;
return nil;
}
let raw = try pool_begin(self.pool, key, digest, self.ttl, req, self.inner) catch (e) nil;
if raw == nil {
let r = Resp { status: 503, headers: {}, body: "{\"error\":\"idempotency store saturated\"}" };
set_header(r, "content-type", "application/json");
set_header(r, "retry-after", "1");
return r;
}
fn after(req: Req, mut resp: Resp) {
-- Only store on successful responses (2xx/3xx) and only if before() was a miss
if req.ctx["idem_miss"] != "true" { return; }
if resp.status < 200 { return; }
if resp.status >= 400 { return; }
if raw / 1_000_000_000 == 2 {
-- Same key, a different request: refuse rather than serve the
-- other request's response.
let r = Resp { status: 422, headers: {}, body: "{\"error\":\"idempotency key reused with a different request\"}" };
set_header(r, "content-type", "application/json");
return r;
}
let key = req.ctx["idem_key"];
if key == nil { return; }
-- Allowlist headers for replay: only content-type
-- Stored (fresh miss or matched replay): the actor already committed
-- this row before returning, so it is there to read.
let hits = from k in IdempotencyKey where k.key == key take 1 select k;
if len(hits) == 0 { return server_error(); }
let stored = json.decode(hits[0].response) as IdempotentStoredResp;
if stored == nil { return server_error(); }
-- `.. ""` forces a fresh, independently-owned Text for every key/value
-- copied out of the decoded record: json.decode's Text values do not
-- survive being handed onward as-is once the decoded record itself
-- goes out of scope (a stale row read back corrupted mid-response
-- otherwise) — concat is documented to always allocate new owned text.
let hdrs: map<Text, Text> = {};
let ct = resp.headers["content-type"];
if ct != nil { hdrs["content-type"] = ct; }
let stored = IdempotentStoredResp {
status: resp.status,
headers: hdrs,
body: resp.body
};
let resp_json = json.encode(stored);
let now = time.ticks();
-- Compute digest of method|path|body
let digest_input = "${req.method}|${req.path}|${req.body}";
let digest = sha256(bytes_of_text(digest_input));
let digest_hex = base64_encode(bytes_slice(digest, 0, 16));
try insert IdempotencyKey { key: key, response: resp_json, created_at: now, digest: digest_hex } catch (e) nil;
for k, v in stored.headers { hdrs[k .. ""] = v .. ""; }
return Resp { status: stored.status, headers: hdrs, body: stored.body .. "" };
}
}
-- Internal typedef for JSON decode of stored response
-- JSON shape of IdempotencyKey.response. Allowlisted headers only:
-- content-type, never Set-Cookie or Date.
typedef IdempotentStoredResp = {
status: Int,
headers: map<Text, Text>,
body: Text
}
-- Key composition function
pub fn idempotent_key(self: Idempotent, header_val: Text, req: Req) -> Text {
if self.key_header == "" { return "idem:${header_val}"; }
if self.include_body == false { return "idem:${self.key_header}:${header_val}"; }
-- Include body digest: sha256(method|path|body)
let digest_input = "${req.method}|${req.path}|${req.body}";
let digest = sha256(bytes_of_text(digest_input));
-- Take first 16 chars of hex digest for brevity (base64_encode of bytes)
let short_digest = base64_encode(bytes_slice(digest, 0, 16));
return "idem:${self.key_header}:${header_val}:${short_digest}";
}

View file

@ -2,17 +2,23 @@
-- 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 or
-- IdempotencyKey themselves.
-- 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.
-- 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
@ -57,16 +63,79 @@ fn dummy_req() -> Req {
};
}
-- 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: look up the idempotency key, replay on a digest match,
-- refuse on a mismatch, or run msg.handler and store the result.
-- Placeholder until then.
return 0 - 1;
-- 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
-- outcome 1 (stored) 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
};
return pool_pack(1, 0);
}
-- kind 1: count.
@ -164,3 +233,19 @@ pub fn pool_count(pool: Pool, key: Text, limit: Int, window: Int) -> Verdict {
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
-- response is now in IdempotencyKey (fresh store or matched replay), 2
-- means a digest mismatch (422, nothing to read). 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
});
}

View file

@ -735,6 +735,183 @@ else
bad "limiter-compile" "$(printf '%s' "$lp_out" | head -1)"
fi
# ---- 18. porch-store task 4: idempotency runs the handler inside the actor --
# Same flattening trick as the keypool/limiter legs. Idempotent wraps a
# SLOW route handler (500ms) so two genuinely-parallel duplicates actually
# overlap in the pool's mailbox. The handler's side effect (ExecMark) is
# a real @table row count read back over GET /execs -- never a log line,
# per the brief. One server serves all three legs with distinct keys, so
# the exec count accumulates 1 -> 2 -> 3 across them.
IP="$W/idempotent-check"
cp -r "$ROOT/docs/examples/porch" "$IP"
rm -f "$IP/wo.toml"
rm -rf "$IP/target"
cat >"$IP/idempotent_check_main.wo" <<'WOEOF'
use net
use env
use http
use router
use middleware
use time
@table(name: "exec_marks")
class ExecMark {
n: Int
}
-- a deliberately slow handler: the concurrency leg's workload. Records
-- one row per REAL execution so a duplicate that wrongly ran it too
-- shows up as a row-count of 2, never as a log line.
class SlowHandler {
fn handle(req: Req) -> Resp {
insert ExecMark { n: 1 };
let n = len(from e in ExecMark select e);
time.sleep(500);
return ok_json("{\"echo\":\"${req.body}\",\"exec\":${n}}");
}
}
class ExecCount {
fn handle(req: Req) -> Resp {
let n = len(from e in ExecMark select e);
return ok_json("{\"count\":${n}}");
}
}
fn build_app(slot: actor PoolMsg) -> App {
let app = App { middleware: [], routes: [] };
let p = Pool { actors: [PoolSlot { a: slot }] };
app.post("/create", Idempotent { key_header: "idempotency-key", pool: p, inner: SlowHandler {} });
app.get("/execs", ExecCount {});
return app;
}
class Conn { fd: net.Conn }
class ConnWorker {
slot: actor PoolMsg
fn receive(msg: Conn) {
let app = build_app(self.slot);
app.handle_conn(msg.fd, 5000, 5000);
}
}
fn main(args: multi Text) -> Int {
if len(args) < 1 {
print_err("usage: idempotent_check <port>");
return 2;
}
let port = parse_int(args[0]);
if port == nil { print_err("bad port"); return 2; }
let ka: actor PoolMsg = spawn KeyActor {};
let srv = net.listen("127.0.0.1", port);
print("listening on 127.0.0.1:${port}");
while true {
if env.stopping() { net.close(srv); return 0; }
let c = net.accept_dl(srv, 250);
if c != nil {
let w: actor Conn = spawn ConnWorker { slot: ka };
send(w, Conn { fd: c });
}
}
}
WOEOF
if ip_out="$("$WOC" --emit "$IP" -o "$IP/idempotent_check.wob" 2>&1)"; then
ok "idempotent: compiles (actor-run handler, digest, pool_begin)"
IPORT=$((PORT + 2))
IDATA="$W/idempotent-data"; mkdir -p "$IDATA"
printf '\n===== idempotent check — port %s =====\n' "$IPORT" >>"$SRVLOG"
LEGFROM=$(( $(wc -l < "$SRVLOG") + 1 ))
WO_DATA="$IDATA" "$WOVM" "$IP/idempotent_check.wob" "$IPORT" >>"$SRVLOG" 2>&1 &
SRV=$!
iwait_listen() {
for _ in $(seq 1 40); do
tail -n "+$LEGFROM" "$SRVLOG" 2>/dev/null | grep -q listening && return
sleep 0.1
done
}
iwait_listen
ipost() { # key body outfile -> prints STATUS, leaves the body in outfile
curl -s -o "$3" -w '%{http_code}' --max-time 5 -X POST \
-H "Host: a" -H "Idempotency-Key: $1" -H "Content-Type: text/plain" \
--data-binary "$2" "http://127.0.0.1:$IPORT/create"
}
iexecs() { # -> the ExecMark row count
curl -s --max-time 5 -H "Host: a" "http://127.0.0.1:$IPORT/execs" \
| grep -o '"count":[0-9]*' | cut -d: -f2
}
# ---- 18a. gate leg: replay is exact (brief step 8) ----------------------
s1="$(ipost leg8-key hello "$W/i8a.body")"
s2="$(ipost leg8-key hello "$W/i8b.body")"
[[ "$s1" == "200" && "$s2" == "200" ]] \
&& ok "idempotent: same key + same body both answer 200" \
|| bad "idempotent-replay-status" "s1=$s1 s2=$s2"
if cmp -s "$W/i8a.body" "$W/i8b.body"; then
ok "idempotent: replay is byte-identical"
else
bad "idempotent-replay-bytes" "$(cat "$W/i8a.body") != $(cat "$W/i8b.body")"
fi
ec="$(iexecs)"
[[ "$ec" == "1" ]] \
&& ok "idempotent: handler ran exactly once (ExecMark row count = 1)" \
|| bad "idempotent-replay-execs" "ExecMark count=$ec want 1"
# ---- 18b. gate leg: digest mismatch is 422 (brief step 9) ---------------
m1="$(ipost leg9-key bodyA "$W/i9a.body")"
m2="$(ipost leg9-key bodyB "$W/i9b.body")"
[[ "$m1" == "200" && "$m2" == "422" ]] \
&& ok "idempotent: same key + different body is 422, not 200" \
|| bad "idempotent-mismatch-status" "m1=$m1 m2=$m2"
if ! cmp -s "$W/i9a.body" "$W/i9b.body"; then
ok "idempotent: 422 body is the refusal, not the other request's response"
else
bad "idempotent-mismatch-bytes" "422 body equals the first request's stored response"
fi
ec="$(iexecs)"
[[ "$ec" == "2" ]] \
&& ok "idempotent: the refused request never ran the handler (ExecMark row count = 2)" \
|| bad "idempotent-mismatch-execs" "ExecMark count=$ec want 2"
# ---- 18c. gate leg: concurrent duplicates (brief step 10) ---------------
# Two backgrounded curl clients, launched together, hitting the SAME
# slow (500ms) handler through the SAME key -- a sequential version of
# this passes against the old before/after flow too and proves nothing.
t0=$(date +%s%3N)
( s="$(ipost leg10-key samebody "$W/i10a.body")"; echo "$s" >"$W/i10a.status" ) &
cc1=$!
( s="$(ipost leg10-key samebody "$W/i10b.body")"; echo "$s" >"$W/i10b.status" ) &
cc2=$!
wait "$cc1" "$cc2"
t1=$(date +%s%3N)
elapsed=$((t1 - t0))
cs1="$(cat "$W/i10a.status")"; cs2="$(cat "$W/i10b.status")"
[[ "$cs1" == "200" && "$cs2" == "200" ]] \
&& ok "idempotent concurrency: both parallel duplicates answer 200" \
|| bad "idempotent-cc-status" "cs1=$cs1 cs2=$cs2"
if cmp -s "$W/i10a.body" "$W/i10b.body"; then
ok "idempotent concurrency: both clients got the same response body"
else
bad "idempotent-cc-bytes" "$(cat "$W/i10a.body") != $(cat "$W/i10b.body")"
fi
[[ "$elapsed" -lt 900 ]] \
&& ok "idempotent concurrency: genuinely overlapped (${elapsed}ms, serial would be ~1000ms+)" \
|| bad "idempotent-cc-elapsed" "${elapsed}ms"
ec="$(iexecs)"
[[ "$ec" == "3" ]] \
&& ok "idempotent concurrency: exactly one execution despite 2 parallel duplicates (ExecMark row count = 3)" \
|| bad "idempotent-cc-execs" "ExecMark count=$ec want 3"
kill -TERM "$SRV" 2>/dev/null
for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || break; sleep 0.1; done
SRV=""
else
bad "idempotent-compile" "$(printf '%s' "$ip_out" | head -1)"
fi
echo
printf 'web-app-accept: %d checks, %d failures\n' "$((pass + fail))" "$fail"
[[ $fail -eq 0 ]]