From 491c2f42b459fee7ffcdd697969f6b3bd73cbe7c Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 30 Aug 2026 00:17:58 +0200 Subject: [PATCH] feat(porch-store): limiter delegates all counting to the key pool - delete all RateLimitCounter access from limiter.wo: query, increment, delete-then-insert, and the swallowing catch (e) nil -- the pool is now the only writer, so its serialization guarantee actually holds - Limiter gains pool/limit/trust_proxy fields; before() calls pool_count and acts on the Verdict; make_limiter takes a pool - key selection: req.principal first, else trust_proxy ? client_ip(req) : net.peer(req.conn); delete the dead req.ctx["verified_proxy"] branch - 429 on a spent window (Retry-After, X-RateLimit-*); 503 + Retry-After on a caught actor trap (saturated pool), request never let through - add Limiter.after(), registered alongside before() as both Mw and Aw (Cors's own shape) so the allowed path's X-RateLimit-* headers reach the response, not just req.ctx - scripts/web-app-accept.sh: three new gate legs -- threshold (N pass, N+1th 429), SIGTERM+restart (still limited from the WAL), and N genuinely-parallel curl clients on one key with an exact-count assertion Co-Authored-By: Claude Opus 5 (1M context) (cherry picked from commit a653dd0aa64711c43461126d32b4632cc8f64c7a) --- docs/examples/porch/middleware/limiter.wo | 124 ++++++++------- scripts/web-app-accept.sh | 182 ++++++++++++++++++++++ 2 files changed, 246 insertions(+), 60 deletions(-) diff --git a/docs/examples/porch/middleware/limiter.wo b/docs/examples/porch/middleware/limiter.wo index 9271965..e88acf5 100644 --- a/docs/examples/porch/middleware/limiter.wo +++ b/docs/examples/porch/middleware/limiter.wo @@ -1,96 +1,100 @@ --- porch/middleware/limiter.wo — rate limiter middleware backed by @table. --- Fixed-window counter with hybrid key (net.peer default, client_ip when verified proxy). --- Iteration 1 of the porch track. +-- porch/middleware/limiter.wo — rate limiter middleware. All counting is +-- delegated to the key pool (keypool.wo): this file never reads or writes +-- RateLimitCounter and holds no window arithmetic of its own. That is what +-- makes the pool's per-key serialization guarantee actually apply — a +-- store call here would be a second, uncoordinated writer. +-- Iteration 1 of the porch track, porch-store task 3. use time use http use net --- Limiter counts requests per key per window. --- On limit exceeded: returns 429 with Retry-After and X-RateLimit headers. --- Key selection: if req.ctx["verified_proxy"] == "true" → client_ip(req) (XFF left-most), --- else → net.peer(req.conn) (unforgeable "ip:port" or "unix"). --- Falls back to "principal:" when authenticated. +-- Limiter counts requests per key per window by calling into the shared +-- pool and acting on the Verdict it returns. +-- On limit exceeded: 429 with Retry-After and X-RateLimit-* headers. +-- On a saturated pool (the key's actor mailbox is full under load): 503 +-- with Retry-After. The request is refused, never let through — a limiter +-- that stops limiting under load is worse than no limiter, since +-- saturating the pool would otherwise be the bypass. +-- Key selection: req.principal wins when non-empty. Otherwise, trust_proxy +-- false (default) keys on net.peer(req.conn), which cannot be forged; +-- trust_proxy true keys on client_ip(req) (the left-most X-Forwarded-For +-- entry) — the app author's assertion that a proxy they control overwrites +-- that header. +-- The refused/saturated paths stamp their own headers directly on the Resp +-- they return. The allowed path has no Resp yet to stamp — before() stashes +-- the numbers on req.ctx, and `after` (same shape as Cors: register the one +-- value as both Mw and Aw) copies them onto whatever response the chain +-- eventually produces. pub class Limiter { - max: Int - window: Int -- window size in µs (e.g., 60_000_000 = 60s) + pool: Pool + limit: Int + window: Int -- window size in µs (e.g., 60_000_000 = 60s) + trust_proxy: Bool = false fn before(mut req: Req) -> ?Resp { - let key = limiter_key(req); - let now = time.ticks(); + let key = limiter_key(self, req); + let v = try pool_count(self.pool, key, self.limit, self.window) catch (e) nil; - -- Read existing counter - let hits = from c in RateLimitCounter where c.key == key take 1 select c; - let counter = RateLimitCounter { key: key, count: 0, window: now }; - let is_new_window = false; - - if len(hits) == 0 { - is_new_window = true; - } else { - counter = hits[0]; - -- Check if window has elapsed - if now - counter.window > self.window { - counter.count = 0; - counter.window = now; - is_new_window = true; - } + if v == nil { + let r = Resp { status: 503, headers: {}, body: "{\"error\":\"rate limiter saturated\"}" }; + set_header(r, "content-type", "application/json"); + set_header(r, "retry-after", "1"); + return r; } - -- Increment and check limit - counter.count = counter.count + 1; - - -- Headers for both allowed and limited responses - let limit_hdr = "${self.max}"; - let remaining = self.max - counter.count; + let limit_hdr = "${v.limit}"; + let remaining = v.limit - v.count; if remaining < 0 { remaining = 0; } - let remaining_hdr = "${remaining}"; - let reset_sec = (counter.window + self.window) / 1_000_000; - let reset_hdr = "${reset_sec}"; + let reset_hdr = "${v.reset_at / 1000}"; -- wall-clock ms -> Unix seconds - -- Store the counter (insert new or delete+insert for update) - if is_new_window { - try insert RateLimitCounter { key: counter.key, count: counter.count, window: counter.window } catch (e) nil; - } else { - -- Update existing: delete old, insert new - let to_delete = from c in RateLimitCounter where c.key == key take 1 select c; - if len(to_delete) > 0 { delete to_delete[0]; } - try insert RateLimitCounter { key: counter.key, count: counter.count, window: counter.window } catch (e) nil; - } - - -- Check limit - if counter.count > self.max { + if v.allowed == false { + let retry_sec = ((v.reset_at - time.now()) / 1000) + 1; let r = Resp { status: 429, headers: {}, body: "{\"error\":\"rate limit exceeded\"}" }; set_header(r, "content-type", "application/json"); set_header(r, "x-ratelimit-limit", limit_hdr); set_header(r, "x-ratelimit-remaining", "0"); set_header(r, "x-ratelimit-reset", reset_hdr); - set_header(r, "retry-after", "${((counter.window + self.window - now) / 1_000_000) + 1}"); + set_header(r, "retry-after", "${retry_sec}"); return r; } -- Allowed: attach headers to request for after-chain to stamp on response req.ctx["ratelimit_limit"] = limit_hdr; - req.ctx["ratelimit_remaining"] = remaining_hdr; + req.ctx["ratelimit_remaining"] = "${remaining}"; req.ctx["ratelimit_reset"] = reset_hdr; return nil; } + + -- Stamps the allowed-path numbers before() stashed. A refused/saturated + -- request never reaches here with anything to stamp (before() only + -- writes ctx on the allowed path), so this is a no-op for those. + fn after(req: Req, mut r: Resp) { + let limit_hdr = req.ctx["ratelimit_limit"]; + let remaining_hdr = req.ctx["ratelimit_remaining"]; + let reset_hdr = req.ctx["ratelimit_reset"]; + if limit_hdr == nil { return; } + if remaining_hdr == nil { return; } + if reset_hdr == nil { return; } + r.headers["x-ratelimit-limit"] = limit_hdr; + r.headers["x-ratelimit-remaining"] = remaining_hdr; + r.headers["x-ratelimit-reset"] = reset_hdr; + } } --- Default hybrid key function: net.peer unless verified proxy -pub fn limiter_key(req: Req) -> Text { - if req.ctx["verified_proxy"] == "true" { - let ip = client_ip(req); - if ip != "" { return "ip:${ip}"; } +-- Key selection: identity first, then the trust_proxy-gated peer address. +pub fn limiter_key(self: Limiter, req: Req) -> Text { + if req.principal != "" { return "principal:${req.principal}"; } + if self.trust_proxy { + return "ip:${client_ip(req)}"; } - -- net.Conn is a scalar (fd), so req.conn should work directly as the fd argument let peer = net.peer(req.conn); if peer != "" { return "ip:${peer}"; } - if req.principal != "" { return "principal:${req.principal}"; } return "unknown"; } -- Helper to build a Limiter with defaults -pub fn make_limiter(max: Int, window_sec: Int) -> Limiter { - return Limiter { max: max, window: window_sec * 1_000_000 }; -} \ No newline at end of file +pub fn make_limiter(pool: Pool, limit: Int, window_sec: Int) -> Limiter { + return Limiter { pool: pool, limit: limit, window: window_sec * 1_000_000 }; +} diff --git a/scripts/web-app-accept.sh b/scripts/web-app-accept.sh index 8e546a1..0d894ea 100755 --- a/scripts/web-app-accept.sh +++ b/scripts/web-app-accept.sh @@ -527,6 +527,188 @@ else bad "keypool" "compile: $(printf '%s' "$kp_out" | head -1)" fi +# ---- 17. porch-store task 3: the limiter delegates to the pool ---------- +# Same flattening trick as the keypool leg, but this one boots a REAL +# server: fiber-per-connection (mirrors web-app's ConnWorker — see App's +# own doc comment), Limiter registered as both Mw and Aw (Cors's own +# dual-role shape) so the allowed path's X-RateLimit-* headers actually +# reach the response, not just req.ctx. trust_proxy is on so every +# backgrounded client can land on ONE key by sending the same +# X-Forwarded-For value, regardless of its own ephemeral source port. +# +# ConnWorker's own state carries a bare actor handle, never a Pool: Pool +# is aliased by design (pool_count reads the same value on every request) +# and the compiler refuses a traced value in actor state or a message +# (WO-E222 — spawn placement makes every actor potentially remote). Each +# connection rebuilds a throwaway one-slot Pool from that handle instead. +LP="$W/limiter-check" +cp -r "$ROOT/docs/examples/porch" "$LP" +rm -f "$LP/wo.toml" +rm -rf "$LP/target" +cat >"$LP/limiter_check_main.wo" <<'WOEOF' +use net +use env +use http +use router +use middleware + +class Ping { + fn handle(req: Req) -> Resp { return ok_text("pong"); } +} + +fn build_app(slot: actor PoolMsg, limit: Int, window_us: Int) -> App { + let app = App { middleware: [], routes: [] }; + let p1 = Pool { actors: [PoolSlot { a: slot }] }; + let p2 = Pool { actors: [PoolSlot { a: slot }] }; + app.use_mw(Mw { m: Limiter { pool: p1, limit: limit, window: window_us, trust_proxy: true } }); + app.use_after(Aw { a: Limiter { pool: p2, limit: limit, window: window_us, trust_proxy: true } }); + app.get("/ping", Ping {}); + return app; +} + +class Conn { fd: net.Conn } + +class ConnWorker { + slot: actor PoolMsg + limit: Int + window: Int + fn receive(msg: Conn) { + let app = build_app(self.slot, self.limit, self.window); + app.handle_conn(msg.fd, 2000, 2000); + } +} + +fn main(args: multi Text) -> Int { + if len(args) < 3 { + print_err("usage: limiter_check "); + return 2; + } + let port = parse_int(args[0]); + if port == nil { print_err("bad port"); return 2; } + let limit = parse_int(args[1]); + if limit == nil { print_err("bad limit"); return 2; } + let window = parse_int(args[2]); + if window == nil { print_err("bad window"); 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, limit: limit, window: window }; + send(w, Conn { fd: c }); + } + } +} +WOEOF + +if lp_out="$("$WOC" --emit "$LP" -o "$LP/limiter_check.wob" 2>&1)"; then + ok "limiter: compiles against the pool (Mw+Aw, trust_proxy)" + + LPORT=$((PORT + 1)) + lhit() { # xff-value -> "STATUS\nHEADERS..." for one request keyed on it + curl -sD - -o /dev/null --max-time 5 -H "X-Forwarded-For: $1" -H "Host: a" "http://127.0.0.1:$LPORT/ping" | tr -d '\r' + } + lstatus() { printf '%s' "$1" | head -1 | awk '{print $2}'; } + lwait_listen() { # from-line + for _ in $(seq 1 40); do + tail -n "+$1" "$SRVLOG" 2>/dev/null | grep -q listening && return + sleep 0.1 + done + } + + # ---- 17a. threshold: of LIMIT+1 requests, first LIMIT pass, last 429 ---- + LDATA="$W/limiter-data"; mkdir -p "$LDATA" + LLIMIT=5 + LWINDOW_US=60000000 + printf '\n===== limiter check — port %s =====\n' "$LPORT" >>"$SRVLOG" + LEGFROM=$(( $(wc -l < "$SRVLOG") + 1 )) + WO_DATA="$LDATA" "$WOVM" "$LP/limiter_check.wob" "$LPORT" "$LLIMIT" "$LWINDOW_US" >>"$SRVLOG" 2>&1 & + SRV=$! + lwait_listen "$LEGFROM" + + codes="" + last="" + for i in $(seq 1 $((LLIMIT + 1))); do + last="$(lhit 6.6.6.6)" + codes="$codes$(lstatus "$last") " + done + want="" + for i in $(seq 1 $LLIMIT); do want="${want}200 "; done + want="${want}429 " + [[ "$codes" == "$want" ]] \ + && ok "limiter threshold: first $LLIMIT pass, request $((LLIMIT + 1)) is 429" \ + || bad "limiter-threshold" "codes=$codes want=$want" + printf '%s\n' "$last" | grep -qi '^retry-after:' \ + && ok "limiter 429 carries Retry-After" \ + || bad "limiter-429-retry-after" "$(printf '%s' "$last" | head -1)" + printf '%s\n' "$last" | grep -qi '^x-ratelimit-remaining: 0' \ + && ok "limiter 429 x-ratelimit-remaining is 0" \ + || bad "limiter-429-remaining" "$(printf '%s' "$last" | head -1)" + + allowed="$(lhit 6.6.6.7)" + printf '%s\n' "$allowed" | grep -qi "^x-ratelimit-limit: $LLIMIT\$" \ + && ok "limiter allowed path carries X-RateLimit-* headers (Mw+Aw)" \ + || bad "limiter-allowed-headers" "$(printf '%s' "$allowed" | head -1)" + + # ---- 17b. SIGTERM + restart: same WO_DATA, same key, still limited ---- + kill -TERM "$SRV" 2>/dev/null + stopped=1 + for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || { stopped=0; break; }; sleep 0.1; done + [[ $stopped -eq 0 ]] && ok "limiter: SIGTERM stops the server" || bad "limiter-stop" "still running" + SRV="" + + printf '\n===== limiter check — restart, port %s =====\n' "$LPORT" >>"$SRVLOG" + LEGFROM=$(( $(wc -l < "$SRVLOG") + 1 )) + WO_DATA="$LDATA" "$WOVM" "$LP/limiter_check.wob" "$LPORT" "$LLIMIT" "$LWINDOW_US" >>"$SRVLOG" 2>&1 & + SRV=$! + lwait_listen "$LEGFROM" + r="$(lhit 6.6.6.6)" + [[ "$(lstatus "$r")" == "429" ]] \ + && ok "limiter restart: counter replayed from the WAL, still limited" \ + || bad "limiter-restart" "got $(lstatus "$r")" + kill -TERM "$SRV" 2>/dev/null + for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || break; sleep 0.1; done + SRV="" + + # ---- 17c. concurrency: N genuinely-parallel requests, exact count ------- + # Backgrounded shell clients (curl `&` + explicit-PID `wait`), not asyncio + # in one process — a sequential version of this passes against the OLD + # read-modify-write limiter too and proves nothing. The exact count is + # read off ONE more sequential probe's X-RateLimit-Remaining afterward, + # never off the table directly: if the N parallel calls lost an + # increment, that number is wrong by exactly the lost count. + LCDATA="$W/limiter-cc-data"; mkdir -p "$LCDATA" + LCN=30 + LCLIMIT=1000 + printf '\n===== limiter check — concurrency, port %s =====\n' "$LPORT" >>"$SRVLOG" + LEGFROM=$(( $(wc -l < "$SRVLOG") + 1 )) + WO_DATA="$LCDATA" "$WOVM" "$LP/limiter_check.wob" "$LPORT" "$LCLIMIT" "$LWINDOW_US" >>"$SRVLOG" 2>&1 & + SRV=$! + lwait_listen "$LEGFROM" + + cc_pids=() + for i in $(seq 1 $LCN); do + ( curl -s -o /dev/null --max-time 10 -H "X-Forwarded-For: 6.6.6.8" -H "Host: a" "http://127.0.0.1:$LPORT/ping" ) & + cc_pids+=("$!") + done + for p in "${cc_pids[@]}"; do wait "$p"; done + + probe="$(lhit 6.6.6.8)" + remaining="$(printf '%s\n' "$probe" | grep -i '^x-ratelimit-remaining:' | awk '{print $2}')" + expected=$((LCLIMIT - (LCN + 1))) + [[ "$remaining" == "$expected" ]] \ + && ok "limiter concurrency: exact count after $LCN parallel requests (no lost increments)" \ + || bad "limiter-concurrency" "remaining=$remaining expected=$expected" + + 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 "limiter-compile" "$(printf '%s' "$lp_out" | head -1)" +fi + echo printf 'web-app-accept: %d checks, %d failures\n' "$((pass + fail))" "$fail" [[ $fail -eq 0 ]]