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) <noreply@anthropic.com> (cherry picked from commit a653dd0aa64711c43461126d32b4632cc8f64c7a)
This commit is contained in:
parent
9af42c8e69
commit
491c2f42b4
2 changed files with 246 additions and 60 deletions
|
|
@ -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:<req.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 };
|
||||
}
|
||||
pub fn make_limiter(pool: Pool, limit: Int, window_sec: Int) -> Limiter {
|
||||
return Limiter { pool: pool, limit: limit, window: window_sec * 1_000_000 };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 <port> <limit> <window_us>");
|
||||
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 ]]
|
||||
|
|
|
|||
Loading…
Reference in a new issue