fix(porch-store): close the ephemeral-row race, not just shrink it
- Root cause of the residual: a 4xx/5xx row lived under the bare key, so a second message could delete-and-replace it before the FIRST caller's own middleware-side read (necessarily outside receive, WO-E226) ever ran -- the owner itself could read back a LATER message's answer, not just a duplicate reading a stale one - Fix: a 4xx/5xx miss is never stored under the bare key at all. Each such attempt gets its own row, keyed by a nonce carried back in the scalar reply's low digits, so no other message for the same bare key ever touches it -- the decision AND the row's identity are both fixed inside the one serialized receive call - Disclosed trade-off: that row is never revisited by a bare-key lookup, so it is never TTL-pruned either -- permanent per failed attempt, the same no-sweeper trade-off this codebase already makes elsewhere, not a new one - Gate leg 18e: reran 20x in isolation against the fix with zero 500-500 or 200-200 outcomes (was reproducible before) - §18's SIGTERM-stop check now force-kills on timeout before clearing $SRV, instead of matching §14/§17b's own gap where a still-running process escapes the exit trap too -- an orphan no longer survives past this leg regardless of the assertion's own outcome Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> (cherry picked from commit 464147a9ddf3a53cb1946637c3f517ba615ee358)
This commit is contained in:
parent
d98ff82027
commit
e296d0541d
3 changed files with 135 additions and 45 deletions
|
|
@ -60,13 +60,18 @@ pub class Idempotent {
|
||||||
return r;
|
return r;
|
||||||
}
|
}
|
||||||
|
|
||||||
-- outcome 1 (2xx/3xx, a durable replay target) or 3 (4xx/5xx, a
|
-- Outcome 1: the bare key names a durable (2xx/3xx) replay target.
|
||||||
-- one-read relay only): either way the actor already committed this
|
-- Outcome 3: this call's own 4xx/5xx answer lives under a row keyed
|
||||||
-- row before returning, so it is there to read right now.
|
-- by a nonce (the reply's low digits) that nobody else's message
|
||||||
let hits = from k in IdempotencyKey where k.key == key take 1 select k;
|
-- for this same bare key ever writes to -- rebuild that exact key
|
||||||
|
-- rather than reading the bare one, so a race with a LATER message
|
||||||
|
-- for this key (which never touches this row) can't hand back the
|
||||||
|
-- wrong response. This file never deletes either kind of row.
|
||||||
|
let lookup_key = key;
|
||||||
|
if outcome == 3 { lookup_key = "${key}#eph:${raw % 1_000_000_000}"; }
|
||||||
|
let hits = from k in IdempotencyKey where k.key == lookup_key take 1 select k;
|
||||||
if len(hits) == 0 { return server_error(); }
|
if len(hits) == 0 { return server_error(); }
|
||||||
let row = hits[0];
|
let stored = json.decode(hits[0].response) as IdempotentStoredResp;
|
||||||
let stored = json.decode(row.response) as IdempotentStoredResp;
|
|
||||||
if stored == nil { return server_error(); }
|
if stored == nil { return server_error(); }
|
||||||
-- `.. ""` forces a fresh, independently-owned Text for every key/value
|
-- `.. ""` forces a fresh, independently-owned Text for every key/value
|
||||||
-- copied out of the decoded record: json.decode's Text values do not
|
-- copied out of the decoded record: json.decode's Text values do not
|
||||||
|
|
@ -75,15 +80,7 @@ pub class Idempotent {
|
||||||
-- otherwise) — concat is documented to always allocate new owned text.
|
-- otherwise) — concat is documented to always allocate new owned text.
|
||||||
let hdrs: map<Text, Text> = {};
|
let hdrs: map<Text, Text> = {};
|
||||||
for k, v in stored.headers { hdrs[k .. ""] = v .. ""; }
|
for k, v in stored.headers { hdrs[k .. ""] = v .. ""; }
|
||||||
let result = Resp { status: stored.status, headers: hdrs, body: stored.body .. "" };
|
return Resp { status: stored.status, headers: hdrs, body: stored.body .. "" };
|
||||||
if outcome == 3 {
|
|
||||||
-- Never a replay target: delete now so the next attempt with this
|
|
||||||
-- key (an immediate retry, most likely) is a genuine miss and
|
|
||||||
-- re-executes the handler, instead of caching a transient failure
|
|
||||||
-- for the TTL.
|
|
||||||
delete row;
|
|
||||||
}
|
|
||||||
return result;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -108,8 +108,9 @@ class KeyActor {
|
||||||
-- Expired: lazy delete (no sweeper exists), fall through to miss.
|
-- Expired: lazy delete (no sweeper exists), fall through to miss.
|
||||||
delete stored;
|
delete stored;
|
||||||
} else if stored.digest == msg.digest {
|
} else if stored.digest == msg.digest {
|
||||||
-- Same request (the owner's own retry, or a duplicate that
|
-- A row lives under the BARE key only when it is durable (see
|
||||||
-- waited in the mailbox): the row already holds the response.
|
-- below) -- an ephemeral one never does -- so any hit here is
|
||||||
|
-- already a stable, valid replay target.
|
||||||
return pool_pack(1, 0);
|
return pool_pack(1, 0);
|
||||||
} else {
|
} else {
|
||||||
-- Same key, a different request: refuse rather than serve the
|
-- Same key, a different request: refuse rather than serve the
|
||||||
|
|
@ -118,6 +119,9 @@ class KeyActor {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
-- Miss: fresh key, an expired row just deleted, or the prior
|
||||||
|
-- attempt (if any) was ephemeral and so is invisible to the
|
||||||
|
-- bare-key lookup above -- all three run the handler fresh.
|
||||||
let resp = msg.handler.handle(msg.req);
|
let resp = msg.handler.handle(msg.req);
|
||||||
let hdrs: map<Text, Text> = {};
|
let hdrs: map<Text, Text> = {};
|
||||||
let ct = resp.headers["content-type"];
|
let ct = resp.headers["content-type"];
|
||||||
|
|
@ -125,27 +129,39 @@ class KeyActor {
|
||||||
let stored_json = json.encode(IdempotentStoredResp {
|
let stored_json = json.encode(IdempotentStoredResp {
|
||||||
status: resp.status, headers: hdrs, body: resp.body
|
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 {
|
if resp.status >= 200 and resp.status < 400 {
|
||||||
-- 2xx/3xx: a real answer worth replaying for the TTL.
|
-- 2xx/3xx: a real answer worth replaying for the TTL, stored
|
||||||
|
-- under the bare key -- one stable row per key, unguarded same
|
||||||
|
-- as kind 1's own insert (this actor is the only writer for
|
||||||
|
-- this key; a @unique violation here would mean something is
|
||||||
|
-- genuinely wrong, not a race to paper over -- a swallowed
|
||||||
|
-- failure would answer "stored" for a response never written).
|
||||||
|
insert IdempotencyKey {
|
||||||
|
key: msg.key, response: stored_json, created_at: now, digest: msg.digest
|
||||||
|
};
|
||||||
return pool_pack(1, 0);
|
return pool_pack(1, 0);
|
||||||
}
|
}
|
||||||
-- 4xx/5xx: the row above exists only so the scalar-only reply can
|
-- 4xx/5xx: never a replay target, so it does NOT go under the bare
|
||||||
-- still hand the caller its exact response (WO-E226 -- a Resp
|
-- key -- a bare-key row is memoryless (this file's own doc above),
|
||||||
-- cannot ride the mailbox). It must NOT survive to answer a later
|
-- but a SHARED, mutable row is not: a second message racing the
|
||||||
-- retry: caching a transient 500 for the TTL (default 24h) would
|
-- first could delete-and-replace it before the first caller's own
|
||||||
-- make every retry fail until it expires, worse than no idempotency
|
-- middleware-side read (necessarily outside receive -- WO-E226,
|
||||||
-- at all. idempotent.wo reads this row once and deletes it.
|
-- a Resp cannot ride the mailbox) ever runs, so the FIRST caller
|
||||||
return pool_pack(3, 0);
|
-- could read back the SECOND caller's answer. Every miss instead
|
||||||
|
-- gets its own row, keyed by a nonce carried back in the scalar's
|
||||||
|
-- low digits (the same slot pool_pack's remaining_ms uses for
|
||||||
|
-- kind 1) so idempotent.wo can reconstruct the exact same key and
|
||||||
|
-- read only ever what THIS call produced -- immune to any other
|
||||||
|
-- message touching this bare key, ever. (Trade-off, disclosed:
|
||||||
|
-- unlike the bare-key row, this one is never revisited by a bare-
|
||||||
|
-- key lookup, so it is never lazily pruned by TTL either -- it is
|
||||||
|
-- a permanent row per failed attempt. There is no sweeper in this
|
||||||
|
-- codebase by design; this is that same trade-off, not a new one.)
|
||||||
|
let nonce = now % 1_000_000_000;
|
||||||
|
insert IdempotencyKey {
|
||||||
|
key: "${msg.key}#eph:${nonce}", response: stored_json, created_at: now, digest: msg.digest
|
||||||
|
};
|
||||||
|
return pool_pack(3, nonce);
|
||||||
}
|
}
|
||||||
|
|
||||||
-- kind 1: count.
|
-- kind 1: count.
|
||||||
|
|
@ -245,15 +261,18 @@ pub fn pool_count(pool: Pool, key: Text, limit: Int, window: Int) -> Verdict {
|
||||||
}
|
}
|
||||||
|
|
||||||
-- The begin accessor idempotent.wo calls: hides pool_select/call the same
|
-- 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
|
-- way pool_count does. Returns the raw packed outcome — 1 means the
|
||||||
-- response is durably in IdempotencyKey (fresh store or matched replay);
|
-- BARE key names a durable (2xx/3xx) replay target; 2 means a digest
|
||||||
-- 2 means a digest mismatch (422, nothing to read); 3 means a 4xx/5xx
|
-- mismatch (422, nothing to read); 3 means this call's own 4xx/5xx
|
||||||
-- miss whose row is a one-read-then-delete relay only, never a replay
|
-- answer lives under "${key}#eph:${raw % 1_000_000_000}" instead --
|
||||||
-- target (see the kind-2 arm above). Never 0: idempotent.wo's own
|
-- the low digits of the reply are that row's own nonce, not a window
|
||||||
-- `try ... catch (e) nil` cannot tell a literal 0 reply apart from a
|
-- size (kind 1's use of the same slot), so idempotent.wo can rebuild
|
||||||
-- trapped call, so the encoding avoids it on purpose. A trapped call (a
|
-- the exact key and read only ever what THIS call produced. Never 0:
|
||||||
-- saturated mailbox) propagates to the caller uncaught, same as
|
-- idempotent.wo's own `try ... catch (e) nil` cannot tell a literal 0
|
||||||
-- pool_count -- the middleware's own try/catch answers 503.
|
-- 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 {
|
pub fn pool_begin(pool: Pool, key: Text, digest: Text, ttl: Int, req: Req, handler: Handler) -> Int {
|
||||||
let a = pool_select(pool, key);
|
let a = pool_select(pool, key);
|
||||||
return call(a, PoolMsg {
|
return call(a, PoolMsg {
|
||||||
|
|
|
||||||
|
|
@ -794,12 +794,40 @@ class FlakyHandler {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
-- second reviewer follow-up: the SAME fails-once handler and table,
|
||||||
|
-- separate from FlakyMark/FlakyHandler above so the sequential leg
|
||||||
|
-- (18d) can't consume the one-time failure this leg (18e) needs -- but
|
||||||
|
-- this time hit by two GENUINELY concurrent duplicates, to pin that
|
||||||
|
-- neither ever receives a replayed 5xx from the other's ephemeral row.
|
||||||
|
@table(name: "flaky_marks2")
|
||||||
|
class FlakyMark2 {
|
||||||
|
n: Int
|
||||||
|
}
|
||||||
|
|
||||||
|
class FlakyHandler2 {
|
||||||
|
fn handle(req: Req) -> Resp {
|
||||||
|
let n = len(from f in FlakyMark2 select f);
|
||||||
|
insert FlakyMark2 { n: 1 };
|
||||||
|
if n == 0 { return server_error(); }
|
||||||
|
return ok_json("{\"ok\":true}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
class FlakyCount2 {
|
||||||
|
fn handle(req: Req) -> Resp {
|
||||||
|
let n = len(from f in FlakyMark2 select f);
|
||||||
|
return ok_json("{\"count\":${n}}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn build_app(slot: actor PoolMsg) -> App {
|
fn build_app(slot: actor PoolMsg) -> App {
|
||||||
let app = App { middleware: [], routes: [] };
|
let app = App { middleware: [], routes: [] };
|
||||||
let p = Pool { actors: [PoolSlot { a: slot }] };
|
let p = Pool { actors: [PoolSlot { a: slot }] };
|
||||||
app.post("/create", Idempotent { key_header: "idempotency-key", pool: p, inner: SlowHandler {} });
|
app.post("/create", Idempotent { key_header: "idempotency-key", pool: p, inner: SlowHandler {} });
|
||||||
app.post("/flaky", Idempotent { key_header: "idempotency-key", pool: p, inner: FlakyHandler {} });
|
app.post("/flaky", Idempotent { key_header: "idempotency-key", pool: p, inner: FlakyHandler {} });
|
||||||
|
app.post("/flaky2", Idempotent { key_header: "idempotency-key", pool: p, inner: FlakyHandler2 {} });
|
||||||
app.get("/execs", ExecCount {});
|
app.get("/execs", ExecCount {});
|
||||||
|
app.get("/flaky2count", FlakyCount2 {});
|
||||||
return app;
|
return app;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -938,8 +966,54 @@ if ip_out="$("$WOC" --emit "$IP" -o "$IP/idempotent_check.wob" 2>&1)"; then
|
||||||
&& ok "idempotent: a transient 5xx is not replayed -- retry re-executes (500 then 200)" \
|
&& ok "idempotent: a transient 5xx is not replayed -- retry re-executes (500 then 200)" \
|
||||||
|| bad "idempotent-5xx-not-cached" "first=$f1 second=$f2 want 500 then 200"
|
|| bad "idempotent-5xx-not-cached" "first=$f1 second=$f2 want 500 then 200"
|
||||||
|
|
||||||
|
# ---- 18e. gate leg: a 5xx is never replayed to a CONCURRENT duplicate ---
|
||||||
|
# (coordinator follow-up on 18d's residual). Two backgrounded clients,
|
||||||
|
# launched together, same key, against a handler that fails only its
|
||||||
|
# first-ever invocation. Whichever message the actor's mailbox happens
|
||||||
|
# to process first gets that real failure; the second message must find
|
||||||
|
# the row ephemeral and re-run the handler itself -- never read back a
|
||||||
|
# replayed 500. Which of the two clients goes first is a race this test
|
||||||
|
# cannot pin, so it asserts the UNORDERED outcome instead: the statuses
|
||||||
|
# are exactly one 500 and one 200 (both requests genuinely executed --
|
||||||
|
# FlakyMark2 count = 2). The pre-fix behavior (middleware-side delete,
|
||||||
|
# racing the actor) would show 500 and 500 with count = 1 whenever the
|
||||||
|
# duplicate is dequeued before the owner's delete lands -- deterministic
|
||||||
|
# either way, no ordering assumption needed.
|
||||||
|
( s="$(curl -s -o "$W/i12a.body" -w '%{http_code}' --max-time 5 -X POST \
|
||||||
|
-H "Host: a" -H "Idempotency-Key: leg12-key" -H "Content-Type: text/plain" \
|
||||||
|
--data-binary "x" "http://127.0.0.1:$IPORT/flaky2")"; echo "$s" >"$W/i12a.status" ) &
|
||||||
|
cf1=$!
|
||||||
|
( s="$(curl -s -o "$W/i12b.body" -w '%{http_code}' --max-time 5 -X POST \
|
||||||
|
-H "Host: a" -H "Idempotency-Key: leg12-key" -H "Content-Type: text/plain" \
|
||||||
|
--data-binary "x" "http://127.0.0.1:$IPORT/flaky2")"; echo "$s" >"$W/i12b.status" ) &
|
||||||
|
cf2=$!
|
||||||
|
wait "$cf1" "$cf2"
|
||||||
|
g1="$(cat "$W/i12a.status")"; g2="$(cat "$W/i12b.status")"
|
||||||
|
gsorted="$(printf '%s\n%s\n' "$g1" "$g2" | sort | tr '\n' ' ')"
|
||||||
|
[[ "$gsorted" == "200 500 " ]] \
|
||||||
|
&& ok "idempotent: concurrent duplicates never replay a 5xx (one 500, one 200)" \
|
||||||
|
|| bad "idempotent-5xx-concurrent" "g1=$g1 g2=$g2 want one 500 and one 200"
|
||||||
|
fc2="$(curl -s --max-time 5 -H "Host: a" "http://127.0.0.1:$IPORT/flaky2count" \
|
||||||
|
| grep -o '"count":[0-9]*' | cut -d: -f2)"
|
||||||
|
[[ "$fc2" == "2" ]] \
|
||||||
|
&& ok "idempotent: both concurrent attempts genuinely executed (FlakyMark2 count = 2)" \
|
||||||
|
|| bad "idempotent-5xx-concurrent-execs" "FlakyMark2 count=$fc2 want 2"
|
||||||
|
|
||||||
kill -TERM "$SRV" 2>/dev/null
|
kill -TERM "$SRV" 2>/dev/null
|
||||||
for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || break; sleep 0.1; done
|
istopped=1
|
||||||
|
for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || { istopped=0; break; }; sleep 0.1; done
|
||||||
|
[[ $istopped -eq 0 ]] && ok "idempotent: SIGTERM stops the server" || bad "idempotent-stop" "still running"
|
||||||
|
if [[ $istopped -eq 1 ]]; then
|
||||||
|
# §14/§17b's own pattern clears SRV here unconditionally, which is
|
||||||
|
# exactly how an orphan survives past this leg: the EXIT trap only
|
||||||
|
# kills a non-empty $SRV, so a still-running process that this loop
|
||||||
|
# gave up on would otherwise keep the port bound for the NEXT run of
|
||||||
|
# this whole script. Force it dead right here instead of trusting the
|
||||||
|
# trap -- the bad-verdict line above already told the reader SIGTERM
|
||||||
|
# alone did not work.
|
||||||
|
kill -9 "$SRV" 2>/dev/null
|
||||||
|
for _ in $(seq 1 20); do kill -0 "$SRV" 2>/dev/null || break; sleep 0.1; done
|
||||||
|
fi
|
||||||
SRV=""
|
SRV=""
|
||||||
else
|
else
|
||||||
bad "idempotent-compile" "$(printf '%s' "$ip_out" | head -1)"
|
bad "idempotent-compile" "$(printf '%s' "$ip_out" | head -1)"
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue