Skip to content

Commit 804eef1

Browse files
authored
Merge pull request #5 from TrustPager/fix/no-repeat-write-on-5xx
Never repeat a write after a server error that may have landed
2 parents 78b6fa1 + ba3e553 commit 804eef1

5 files changed

Lines changed: 217 additions & 10 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,10 @@ Format follows [Keep a Changelog](https://keepachangelog.com/en/1.1.0/). Version
5454
- **`make-thumbnail` resolves its studio and brand explicitly** when a working directory holds more than one thumbnail studio or brand kit, instead of rendering into the first one it finds.
5555
- **Onboarding (`/start-here`) redirected to consultation-first (founder-ruled 2026-07-03).** The Day-1 win is now the collaborative consultative conversation (reflect understanding, draw out the goal and the owner's own theory of the blocker, then think alongside them with the reasoning shown), decided by an engagement gauge, rather than a built artifact handed over on the spot. Any build is deferred to a recommendation-with-alternatives at the end; a terse owner still gets a fast tangible win. The assistant now mirrors the owner's register. See `docs/architecture/2026-07-03-collaborative-consultation-design.md`.
5656

57+
### Fixed
58+
59+
- **A write is no longer repeated after a server error that may have landed.** The API layer retried every server error, writes included, so a create or a send that went through and then errored could happen twice (a second opportunity, a second email). Reads still retry; a write retries only when the driver says the server refused it before running it, which for TrustPager is its edge's momentary "too busy" 503. Rate-limit retries are unchanged. A write that is not retried now says it may already have gone through, so an agent checks before trying again instead of reading "try again" literally. The new optional `DriverConfig.refused_before_running` carries that judgement, so the kernel stays vendor neutral.
60+
5761
---
5862

5963
## [1.0.0] - 2026-06-29

‎drivers/trustpager/__init__.py‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,21 @@
7575
# read crossed the sea (measured 2026-09-28).
7676
HOME_REGION_HEADERS = {"x-region": "ap-southeast-2"}
7777

78+
79+
def _refused_before_running(code: int, headers: dict[str, str], body: bytes) -> bool:
80+
"""TrustPager's edge briefly refuses a request when it is too busy to run it.
81+
82+
That refusal is a 503 carrying SERVICE_DEGRADED in its sb-error-code header
83+
(or its body). The request never ran, so the kernel may repeat it even for
84+
a write. Any other server error may have landed first and is not repeated.
85+
Most of them arrive in the first seconds after a quarter hour (measured
86+
2026-10-07).
87+
"""
88+
if code != 503:
89+
return False
90+
return "SERVICE_DEGRADED" in headers.get("sb-error-code", "") or b"SERVICE_DEGRADED" in body
91+
92+
7893
# =============================================================================
7994
# The single TrustPager DriverConfig. Constructing it registers the tp_ secret
8095
# pattern with the redaction registry (DriverConfig.__post_init__).
@@ -86,6 +101,7 @@
86101
error_map=TP_ERROR_MAP,
87102
approval_url=APPROVAL_URL,
88103
extra_headers=HOME_REGION_HEADERS,
104+
refused_before_running=_refused_before_running,
89105
)
90106

91107

‎kernel/runtime/transport.py‎

Lines changed: 59 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,12 @@
1313
- The offline guard (is_offline()) runs BEFORE cfg.key_resolver() so a real
1414
key is never read in tests/CI. tests/test_safety.py depends on this order.
1515
- The actual network call goes through the module-level `_http` indirection
16-
so tests can monkeypatch kernel.runtime.transport._http without a socket.
17-
- Retry/backoff for 429 (honouring Retry-After) and 5xx is preserved here.
16+
so tests can monkeypatch kernel.runtime.transport._http without a socket
17+
(and backoff waits through `_sleep`, for the same reason).
18+
- Retry/backoff lives here: 429 (honouring Retry-After) for any method; a
19+
5xx only for a read, or for a write the driver says never ran
20+
(DriverConfig.refused_before_running). Any other server error may have
21+
landed after the write did, so repeating it could do it twice.
1822
- Per-code messages come from cfg.error_map; absent codes fall back to a
1923
GENERIC kernel message. No vendor literal appears in this file.
2024
"""
@@ -36,6 +40,8 @@
3640
DEFAULT_TIMEOUT_SECONDS = 30
3741
DEFAULT_RETRIES_ON_429 = 3
3842
DEFAULT_RETRIES_ON_5XX = 2
43+
# Methods that change nothing, so a server error is always safe to repeat.
44+
SAFE_TO_REPEAT_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
3945

4046

4147
@dataclass
@@ -59,6 +65,13 @@ class DriverConfig:
5965
(missing keys are tolerated).
6066
approval_url: where the operator approves a queued (202) write.
6167
extra_headers: optional headers merged into every request.
68+
refused_before_running:
69+
optional predicate (http_code, headers, body) -> bool,
70+
headers as a lowercase-keyed dict and body as raw
71+
bytes. True means the server refused the request
72+
BEFORE running it (e.g. a momentarily overloaded edge),
73+
so it is safe to repeat even a write. Without it, a
74+
5xx on a write is never retried.
6275
"""
6376

6477
base_url: str
@@ -69,6 +82,7 @@ class DriverConfig:
6982
error_map: dict[int | str, str]
7083
approval_url: str
7184
extra_headers: dict[str, str] | None = None
85+
refused_before_running: Callable[[int, dict[str, str], bytes], bool] | None = None
7286

7387
def __post_init__(self) -> None:
7488
# Register this driver's key shape so any key that ever surfaces in a
@@ -130,6 +144,26 @@ def _http(req: urllib.request.Request, timeout: int):
130144
return urllib.request.urlopen(req, timeout=timeout)
131145

132146

147+
def _sleep(seconds: float) -> None:
148+
"""Backoff wait. Indirection so tests can skip real waiting."""
149+
time.sleep(seconds)
150+
151+
152+
def _never_ran(cfg: DriverConfig, code: int, headers: Any, body: bytes) -> bool:
153+
"""Ask the driver whether the server refused this request before running it.
154+
155+
A predicate that raises is treated as "it may have run", so a bug in a
156+
driver can only ever make the kernel retry less, never more.
157+
"""
158+
if cfg.refused_before_running is None:
159+
return False
160+
try:
161+
hdrs = {str(k).lower(): str(v) for k, v in (headers.items() if headers else [])}
162+
return bool(cfg.refused_before_running(code, hdrs, body or b""))
163+
except Exception:
164+
return False
165+
166+
133167
def request(cfg: DriverConfig, method: str, path: str,
134168
params: dict[str, Any] | None = None,
135169
body: dict[str, Any] | None = None,
@@ -145,8 +179,12 @@ def request(cfg: DriverConfig, method: str, path: str,
145179
146180
Retry behaviour:
147181
- 429: retries up to DEFAULT_RETRIES_ON_429 times, honouring Retry-After
148-
- 5xx: retries up to DEFAULT_RETRIES_ON_5XX times with exponential backoff
149-
- Network errors: no retry — bubbles up immediately
182+
- 5xx: retries up to DEFAULT_RETRIES_ON_5XX times with exponential
183+
backoff, but only for a read (SAFE_TO_REPEAT_METHODS) or when
184+
cfg.refused_before_running says the request never ran. A write that
185+
got any other server error is NOT repeated: it may already have
186+
landed, and a retry would create a second record or send twice.
187+
- Network errors: no retry, bubbles up immediately
150188
151189
The offline guard fires BEFORE cfg.key_resolver() so a real key is never
152190
read in tests/CI.
@@ -203,20 +241,31 @@ def request(cfg: DriverConfig, method: str, path: str,
203241
if e.code == 429 and _attempt < DEFAULT_RETRIES_ON_429:
204242
retry_after = e.headers.get("Retry-After") if hasattr(e, "headers") else None
205243
wait = int(retry_after) if retry_after and retry_after.isdigit() else 2 ** _attempt
206-
time.sleep(min(wait, 30))
244+
_sleep(min(wait, 30))
207245
return request(cfg, method, path, params=params, body=body,
208246
timeout=timeout, extra_headers=extra_headers,
209247
_attempt=_attempt + 1)
210248

211-
# 5xx — server error. Retry a couple of times then give up.
212-
if 500 <= e.code < 600 and _attempt < DEFAULT_RETRIES_ON_5XX:
213-
time.sleep(2 ** _attempt)
249+
# 5xx: retry a read, or a write the server refused before running.
250+
# Any other server error on a write may have landed already.
251+
is_5xx = 500 <= e.code < 600
252+
is_read = method.upper() in SAFE_TO_REPEAT_METHODS
253+
never_ran = is_5xx and not is_read and _never_ran(
254+
cfg, e.code, getattr(e, "headers", None), detail_raw)
255+
if is_5xx and _attempt < DEFAULT_RETRIES_ON_5XX and (is_read or never_ran):
256+
_sleep(2 ** _attempt)
214257
return request(cfg, method, path, params=params, body=body,
215258
timeout=timeout, extra_headers=extra_headers,
216259
_attempt=_attempt + 1)
217260

218-
raise BOSError(_format_http_error(cfg, e.code, path, url, detail_str,
219-
detail_parsed)) from None
261+
message = _format_http_error(cfg, e.code, path, url, detail_str, detail_parsed)
262+
if is_5xx and not is_read and not never_ran:
263+
# The caller is often an agent that reads "try again" literally.
264+
message += (
265+
"\nThis was a write and it may have gone through before the error. "
266+
"Check whether it did before trying again, or it could happen twice."
267+
)
268+
raise BOSError(message) from None
220269
except urllib.error.URLError as e:
221270
raise BOSError(
222271
f"Could not reach the API.\n"

‎tests/test_transport_offline.py‎

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -307,5 +307,116 @@ def test_exact_429_template_renders(self):
307307
self.assertIn("Slow down: too many", str(ctx.exception))
308308

309309

310+
class TestRetryPolicy(unittest.TestCase):
311+
"""Which server errors are repeated.
312+
313+
A read is always safe to repeat. A write is repeated only when the driver
314+
says the server refused it before running it; any other server error may
315+
have landed already, and repeating it would create a second record or send
316+
twice. Driven by a synthetic cfg whose refusal marker is a made-up header.
317+
"""
318+
319+
def setUp(self):
320+
self._saved_patterns = redaction._snapshot_patterns()
321+
self._prev_offline = os.environ.get("BOS_OFFLINE")
322+
os.environ.pop("BOS_OFFLINE", None)
323+
self._saved_http = transport._http
324+
self._saved_sleep = transport._sleep
325+
transport._sleep = lambda seconds: None
326+
self.calls = []
327+
328+
def tearDown(self):
329+
redaction._restore_patterns(self._saved_patterns)
330+
transport._http = self._saved_http
331+
transport._sleep = self._saved_sleep
332+
if self._prev_offline is None:
333+
os.environ.pop("BOS_OFFLINE", None)
334+
else:
335+
os.environ["BOS_OFFLINE"] = self._prev_offline
336+
337+
def _cfg(self, predicate=None):
338+
return transport.DriverConfig(
339+
base_url="https://example.invalid/api",
340+
key_resolver=lambda: "fake_key_value",
341+
secret_pattern=r"fakesecret_[A-Za-z0-9]{6,}",
342+
error_map={},
343+
approval_url="https://example.invalid/approvals",
344+
refused_before_running=predicate,
345+
)
346+
347+
@staticmethod
348+
def _refusal_marker(code, headers, body):
349+
return code == 503 and headers.get("x-fake-refusal") == "never-ran"
350+
351+
def _script(self, *outcomes):
352+
"""Fake _http that plays outcomes in order: (code, headers) or 'ok'."""
353+
def fake_http(req, timeout):
354+
self.calls.append(req.get_method())
355+
outcome = outcomes[min(len(self.calls) - 1, len(outcomes) - 1)]
356+
if outcome == "ok":
357+
return _FakeResp(200, b'{"data": {"id": "x1"}}')
358+
code, hdrs = outcome
359+
raise urllib.error.HTTPError(
360+
url="https://example.invalid/api/things", code=code,
361+
msg="Synthetic", hdrs=hdrs, fp=io.BytesIO(b'{}'),
362+
)
363+
transport._http = fake_http
364+
365+
def test_read_is_repeated_on_a_server_error(self):
366+
self._script((500, None), "ok")
367+
out = transport.request(self._cfg(), "GET", "things")
368+
self.assertEqual(out, {"data": {"id": "x1"}})
369+
self.assertEqual(self.calls, ["GET", "GET"])
370+
371+
def test_write_is_not_repeated_on_a_server_error(self):
372+
for method in ("POST", "PATCH", "PUT", "DELETE"):
373+
self.calls = []
374+
self._script((500, None), "ok")
375+
with self.assertRaises(BOSError) as ctx:
376+
transport.request(self._cfg(self._refusal_marker), method, "things", body={"k": "v"})
377+
self.assertEqual(len(self.calls), 1, f"{method} must not be repeated")
378+
self.assertIn("may have gone through", str(ctx.exception))
379+
380+
def test_write_is_repeated_when_the_server_refused_before_running(self):
381+
self._script((503, {"X-Fake-Refusal": "never-ran"}), "ok")
382+
out = transport.request(self._cfg(self._refusal_marker), "POST", "things", body={"k": "v"})
383+
self.assertEqual(out, {"data": {"id": "x1"}})
384+
self.assertEqual(self.calls, ["POST", "POST"])
385+
386+
def test_refused_every_time_gives_up_without_the_maybe_landed_warning(self):
387+
self._script((503, {"X-Fake-Refusal": "never-ran"}))
388+
with self.assertRaises(BOSError) as ctx:
389+
transport.request(self._cfg(self._refusal_marker), "POST", "things", body={"k": "v"})
390+
self.assertEqual(len(self.calls), 1 + transport.DEFAULT_RETRIES_ON_5XX)
391+
self.assertNotIn("may have gone through", str(ctx.exception))
392+
393+
def test_plain_503_on_a_write_is_not_repeated(self):
394+
self._script((503, None), "ok")
395+
with self.assertRaises(BOSError):
396+
transport.request(self._cfg(self._refusal_marker), "POST", "things", body={"k": "v"})
397+
self.assertEqual(len(self.calls), 1)
398+
399+
def test_no_predicate_means_writes_are_never_repeated(self):
400+
self._script((503, {"X-Fake-Refusal": "never-ran"}), "ok")
401+
with self.assertRaises(BOSError):
402+
transport.request(self._cfg(None), "POST", "things", body={"k": "v"})
403+
self.assertEqual(len(self.calls), 1)
404+
405+
def test_a_predicate_that_raises_is_treated_as_may_have_run(self):
406+
def broken(code, headers, body):
407+
raise RuntimeError("driver bug")
408+
self._script((503, {"X-Fake-Refusal": "never-ran"}), "ok")
409+
with self.assertRaises(BOSError):
410+
transport.request(self._cfg(broken), "POST", "things", body={"k": "v"})
411+
self.assertEqual(len(self.calls), 1)
412+
413+
def test_429_still_repeats_a_write(self):
414+
# Rate limiting refuses before running, so it stays safe for writes.
415+
self._script((429, {}), "ok")
416+
out = transport.request(self._cfg(), "POST", "things", body={"k": "v"})
417+
self.assertEqual(out, {"data": {"id": "x1"}})
418+
self.assertEqual(self.calls, ["POST", "POST"])
419+
420+
310421
if __name__ == "__main__":
311422
unittest.main()

‎tests/test_trustpager_driver.py‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,5 +106,32 @@ def test_every_call_asks_for_the_database_region(self):
106106
self.assertEqual((TP_CFG.extra_headers or {}).get("x-region"), "ap-southeast-2")
107107

108108

109+
class TestRefusedBeforeRunning(unittest.TestCase):
110+
"""Only the edge's busy refusal lets the kernel repeat a write."""
111+
112+
CODE = "SUPABASE_EDGE_RUNTIME_SERVICE_DEGRADED"
113+
114+
def setUp(self):
115+
from drivers.trustpager import TP_CFG, _refused_before_running
116+
self.cfg = TP_CFG
117+
self.check = _refused_before_running
118+
119+
def test_wired_into_the_driver_config(self):
120+
self.assertIs(self.cfg.refused_before_running, self.check)
121+
122+
def test_busy_refusal_in_the_header(self):
123+
self.assertTrue(self.check(503, {"sb-error-code": self.CODE}, b""))
124+
125+
def test_busy_refusal_named_only_in_the_body(self):
126+
self.assertTrue(self.check(503, {}, ('{"code":"%s"}' % self.CODE).encode()))
127+
128+
def test_any_other_503_may_have_run(self):
129+
self.assertFalse(self.check(503, {}, b'{"error":{"message":"store not connected"}}'))
130+
131+
def test_other_server_errors_may_have_run(self):
132+
for code in (500, 502, 504):
133+
self.assertFalse(self.check(code, {"sb-error-code": self.CODE}, b""), code)
134+
135+
109136
if __name__ == "__main__":
110137
unittest.main()

0 commit comments

Comments
 (0)