From c97de237ef625bb7cfc4f0a524c9e17d09dc0323 Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 30 Aug 2026 01:13:54 +0200 Subject: [PATCH] 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) (cherry picked from commit eae1b06cdc38e4766d4e27a66b02264f47164e99) --- docs/examples/porch/middleware/idempotent.wo | 153 +++++++--------- docs/examples/porch/middleware/keypool.wo | 99 ++++++++++- scripts/web-app-accept.sh | 177 +++++++++++++++++++ 3 files changed, 333 insertions(+), 96 deletions(-) diff --git a/docs/examples/porch/middleware/idempotent.wo b/docs/examples/porch/middleware/idempotent.wo index c58e7ea..dbe8bd1 100644 --- a/docs/examples/porch/middleware/idempotent.wo +++ b/docs/examples/porch/middleware/idempotent.wo @@ -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 = {}; - 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, 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}"; -} \ No newline at end of file diff --git a/docs/examples/porch/middleware/keypool.wo b/docs/examples/porch/middleware/keypool.wo index 925aea4..4100739 100644 --- a/docs/examples/porch/middleware/keypool.wo +++ b/docs/examples/porch/middleware/keypool.wo @@ -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 = {}; + 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: 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 = {}; + 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 + }); +} diff --git a/scripts/web-app-accept.sh b/scripts/web-app-accept.sh index 6c71e0f..1ddfdea 100755 --- a/scripts/web-app-accept.sh +++ b/scripts/web-app-accept.sh @@ -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 "); + 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 ]]