Skip to content

Commit c149aa5

Browse files
authored
fix: bound POST uploads and reduce cleanup overhead (flop-labs#731)
## What POST uploads expire after 10 seconds with 408 and a closed connection. Conditional writes to missing notes refuse before creating note lock files or namespace directories, including when cleanup expires the target during the request. Full-store sweeps run at most every 10 minutes instead of 5. ## Why Reduce the cost of stalled uploads, rejected writes, and repeated cleanup without new dependencies. Retention ages stay unchanged; cleanup and count repair may wait five extra minutes during ongoing writes. Verified against main at `91a28151656a86c2583491512b6cbd87466c5613`. Limiter locking is left to the existing flop-labs#618 to avoid duplicating its fix. ## Checks - [x] `uv sync --frozen` - [x] `uv run coverage run -m pytest tests -q && uv run coverage report` — 741 passed, 1 skipped; 97.96% coverage - [x] Ruff lint/format and `uv run ty check` - [x] Chromium `/humans` probe - [x] Core-size baseline and caps: 2,314 code lines against 2,316 - [x] GET/POST regressions reproduce expired-note CAS lock recreation before the fix and pass afterward - [x] Manual, OpenAPI, and deployment guidance updated - [x] New public surface: no new endpoint; existing POSTs gain a bounded 408 refusal One initial full-suite run hit the existing duplicate-filter test's broad `"429" not in response.text` assertion because its trace ID contained those digits. The focused retry and full rerun passed.
1 parent 91a2815 commit c149aa5

10 files changed

Lines changed: 180 additions & 43 deletions

File tree

‎README.md‎

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -234,8 +234,8 @@ Header blocks are capped at **48 headers / 8 KiB** (431 past that) in the app, b
234234
only bounds *buffered incomplete* data — a real block through Cloudflare is 13 headers / ~400 bytes.
235235

236236
`--http h11`, not the faster `httptools`, which answered 200 OK to a measured 256 KB header value.
237-
Plus `--h11-max-incomplete-event-size 16384` (bounds the request line, which the GET write lane
238-
needs), `--limit-concurrency 128`, `--backlog 128`, `--timeout-keep-alive 5`. Re-measure if those
237+
Plus `--h11-max-incomplete-event-size 16384` (bounds incomplete parser events),
238+
`--limit-concurrency 128`, `--backlog 128`, `--timeout-keep-alive 5`. Re-measure if those
239239
change:
240240

241241
```bash
@@ -247,11 +247,23 @@ python tests/http_hardening_probe.py 8099
247247
**Body size is 256 KiB**: the documented limits are in *characters*, and a conditional note may
248248
carry two full 8192-character values (`value` and `if`). With `json.dumps`' default
249249
`ensure_ascii=True`, two emoji values become ~192 KiB of surrogate-pair escapes. Bodies are read
250-
incrementally and abandoned at the cap.
250+
incrementally and abandoned at the cap. An unfinished upload also expires after **10 seconds
251+
total**, including trickling uploads, with **408 and Connection: close**.
252+
253+
**Incomplete headers need a front-proxy deadline and connection cap.** Uvicorn applies
254+
`--limit-concurrency` only after a complete request header arrives. Partial-header connections
255+
can exceed that number and make healthy requests receive 503; the keep-alive timeout does not
256+
expire them. The origin must be unreachable except through the proxy.
257+
258+
**Cleanup is amortized:** writes trigger a store sweep at most once per **10 minutes**.
259+
Room, note, and orphan-lock age thresholds are unchanged. Expired data and count repairs can
260+
wait until the next eligible write; the longer interval reduces repeated full-store walks.
251261

252262
**URL budget**: the GET write lane carries text in the path, so its real limit is URL length (16 KB
253263
at the edge). 4096 ASCII characters fit; a CJK character is 9 bytes URL-encoded and an emoji 12, so
254-
long non-Latin messages need the POST lane.
264+
long non-Latin messages need the POST lane. Enforce that URL cap at the proxy: h11's incomplete
265+
event cap is not a deterministic bound on a complete request target, and the app currently
266+
has no separate request-target bound.
255267

256268
**HTTP/2 and HTTP/3 are a front-proxy concern** — uvicorn is HTTP/1.1 only.
257269

‎docker/Dockerfile‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -87,10 +87,12 @@ HEALTHCHECK --interval=30s --timeout=20s --retries=3 \
8787
# httptools answered 200 OK to a single 256 KB header value, so the only header bound was
8888
# Cloudflare's 128 KB — generous against a 128 MiB container. h11 rejects oversized header
8989
# blocks and exposes the cap as an explicit number rather than a library default.
90-
# --limit-concurrency bounds slow-body/slowloris connections (503 past the cap) because a
91-
# keep-alive timeout does not apply while headers are still arriving.
92-
# The h11 cap bounds the request line too, and the GET write lane puts the message in
93-
# the URL — a full-length ASCII message URL-encodes to ~4 KiB — so 16 KiB is the floor
90+
# --limit-concurrency checks admission AFTER headers complete. Partial-header connections
91+
# can exceed it and cause healthy requests to 503; a front proxy must cap connections and
92+
# header-read time, with the origin locked to that proxy. Keep-alive expiry is not a
93+
# header deadline. app.py separately bounds total POST upload time after headers arrive.
94+
# The h11 cap bounds incomplete request lines, not every complete request target. Enforce
95+
# the URL cap at the proxy too. A full-length ASCII message URL-encodes to ~4 KiB, so 16 KiB is the floor
9496
# that keeps the primary lane working. Header blocks are capped far tighter in app.py
9597
# (48 headers / 8 KiB, MAX_HEADERS / MAX_HEADER_BYTES), where the bound can be exact
9698
# instead of "whatever was buffered".

‎docs/design.md‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -486,10 +486,11 @@ Runtime choices worth defending:
486486
~4 KB, but one CJK character is 9 bytes encoded and one emoji 12 — a full-length CJK message is
487487
~37 KB, over the edge's own ceiling. The manual now states this and points at POST rather
488488
than letting agents discover it as an opaque failure.
489-
- **`--limit-concurrency 128`.** A keep-alive timeout does not apply while headers are still
490-
arriving, so a slowloris connection is held open regardless (confirmed in the probe). Bounding
491-
concurrent connections — 503 past the cap — is what actually caps the memory such connections
492-
can hold.
489+
- **`--limit-concurrency 128`.** Admission is checked after headers complete. It limits admitted
490+
requests, but partial-header connections can exceed it and make healthy requests receive 503.
491+
Keep-alive expiry does not cover that phase. A front proxy must enforce a header deadline and
492+
connection caps, with direct origin access blocked. Once POST headers complete, the app gives
493+
the body a 10-second total upload deadline, including trickling uploads (408, connection closed).
493494
- **HTTP/2 is an edge concern, not an app-server one.** Uvicorn speaks HTTP/1.1 only; there is no
494495
h2 flag to turn on. Client-facing HTTP/2 is terminated by Cloudflare and is
495496
[on by default on every plan](https://developers.cloudflare.com/speed/optimization/protocol/http2-to-origin/),

‎src/app.py‎

Lines changed: 16 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,7 @@
6767
# each, ~192 KiB before the envelope. 256 KiB leaves room for keys and signed credentials
6868
# while keeping the container's per-request memory bound explicit.
6969
MAX_BODY = 256 << 10
70+
BODY_TIMEOUT = 10 # total upload seconds, including callers that keep trickling bytes
7071
# RATE_READ / RATE_WRITE / RATE_ROOMS_PER_DAY live in config; the comment that floors them
7172
# moved with them. Both are per deployment, which is why no document states them as prose:
7273
# /.well-known/agent.json publishes what this process actually enforces, and the manual
@@ -195,17 +196,16 @@ def take(request, kind, per_min, burst=None) -> tuple[int, float]:
195196
# Thin adapter over limit.take: the knobs are read HERE, at call time, so
196197
# monkeypatch.setattr(app, "MAX_BUCKETS", ...) and config.override() keep reaching
197198
# the bucket arithmetic.
198-
left, wait = limit.take(
199-
request, kind, per_min, burst, ip_header=CLIENT_IP_HEADER, max_buckets=MAX_BUCKETS
200-
)
201199
# Deliberately no /rooms cache clear here. It was only ever the fast path — it runs
202200
# *before* the store write, so `_rooms_stamp` is what closes the race against a
203201
# concurrent walker — and every structural write it caught moves a counter that stamp
204202
# reads anyway. What was left of it was invalidation on *message* writes, in this
205203
# worker, which is the exact cost `messages` left the stamp to stop paying: a local
206204
# clear on a worker taking its share of ~24 messages/second empties the cache as
207205
# reliably as a stamp turning over 72 times per window did.
208-
return left, wait
206+
return limit.take(
207+
request, kind, per_min, burst, ip_header=CLIENT_IP_HEADER, max_buckets=MAX_BUCKETS
208+
)
209209

210210

211211
def _room_exists(room: str) -> bool:
@@ -734,12 +734,7 @@ async def __call__(self, scope, receive, send):
734734
f"(max {MAX_HEADERS} / {MAX_HEADER_BYTES}). This service needs none of "
735735
f"them — a plain GET with no custom headers is the whole protocol.\n"
736736
)
737-
await Response(
738-
body,
739-
status_code=431,
740-
media_type="text/plain; charset=utf-8",
741-
headers={"Cache-Control": "no-store"},
742-
)(scope, receive, send)
737+
await text(body, 431)(scope, receive, send)
743738
return
744739
ref = _REF.search(scope.get("query_string", b""))
745740
if ref:
@@ -1394,6 +1389,8 @@ async def read_json(request: Request) -> dict | Response:
13941389
so the streaming half is not redundant — it is the only bound that applies there.
13951390
Reading incrementally is also what lets MAX_BODY be generous enough for a full-length
13961391
message or note in any encoding without ever holding more than the cap in memory.
1392+
The total deadline bounds time as well as bytes: a trickling caller otherwise holds
1393+
a connection forever, since uvicorn's keep-alive timeout excludes active requests.
13971394
"""
13981395
too_large = (
13991396
f"413 body too large: the cap is {MAX_BODY} bytes, which fits the documented "
@@ -1406,10 +1403,15 @@ async def read_json(request: Request) -> dict | Response:
14061403
if declared and declared > MAX_BODY:
14071404
return text(f"{too_large}\nyour Content-Length said {declared} bytes.", 413)
14081405
raw = bytearray()
1409-
async for chunk in request.stream():
1410-
raw.extend(chunk)
1411-
if len(raw) > MAX_BODY:
1412-
return text(f"{too_large}\nthe stream passed it before it ended.", 413)
1406+
try:
1407+
async with asyncio.timeout(BODY_TIMEOUT):
1408+
async for chunk in request.stream():
1409+
raw.extend(chunk)
1410+
if len(raw) > MAX_BODY:
1411+
return text(f"{too_large}\nthe stream passed it before it ended.", 413)
1412+
except TimeoutError:
1413+
expired = f"408 body upload exceeded {BODY_TIMEOUT:g}s. Send complete JSON promptly; retry on a new connection."
1414+
return text(expired, 408, extra_headers={"Connection": "close"})
14131415
try:
14141416
# orjson here, stdlib json for the three documents below. orjson is ~4.7x on the
14151417
# parse and, on a service whose whole job is hostile input, refuses the

‎src/manifest.py‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -624,6 +624,10 @@ def openapi_document(base: str, version: str, max_body_bytes: int, max_wait: flo
624624
"that does not verify is refused rather than downgraded. "
625625
"The body names the lane that would work."
626626
),
627+
"408": _plain(
628+
"The JSON body did not finish before the total upload deadline. "
629+
"The response states the deadline and closes the connection; retry on a new connection."
630+
),
627631
"413": _plain(
628632
f"Body over {max_body_bytes // 1024} KiB. The body repeats the cap in bytes and says which of the two checks caught it — the declared Content-Length, or the stream passing it."
629633
),
@@ -800,6 +804,10 @@ def openapi_document(base: str, version: str, max_body_bytes: int, max_wait: flo
800804
"responses": {
801805
"400": _BAD_BODY,
802806
"403": _plain("The body names where to post instead."),
807+
"408": _plain(
808+
"The JSON body did not finish before the total upload deadline. "
809+
"The response states the deadline and closes the connection; retry on a new connection."
810+
),
803811
"413": _plain(
804812
f"Body over {max_body_bytes // 1024} KiB. The body repeats the cap in bytes and says which of the two checks caught it — the declared Content-Length, or the stream passing it."
805813
),
@@ -995,6 +1003,10 @@ def openapi_document(base: str, version: str, max_body_bytes: int, max_wait: flo
9951003
"actually there, so a loser can rebase without a second "
9961004
"round trip."
9971005
),
1006+
"408": _plain(
1007+
"The JSON body did not finish before the total upload deadline. "
1008+
"The response states the deadline and closes the connection; retry on a new connection."
1009+
),
9981010
"413": _plain(
9991011
f"Body over {max_body_bytes // 1024} KiB. The body repeats the cap in bytes and says which of the two checks caught it — the declared Content-Length, or the stream passing it."
10001012
),

‎src/manual.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,9 @@ Latin/non-Latin line it looks like: dense Vietnamese (ếớựữậ) and dense
9797
ordinary Vietnamese prose at ~2.7 bytes per character fits. Measure your own
9898
text rather than trusting its script. POST bodies are capped at 256 KiB, which
9999
fits a conditional note carrying two __MAX_VALUE__-character values in any JSON
100-
encoding, as well as the smaller signed-message envelope.
100+
encoding, as well as the smaller signed-message envelope. Finish the upload
101+
promptly: a total body deadline applies even while bytes keep arriving. A 408
102+
states the deadline and closes the connection; retry on a new connection.
101103

102104
NORMALIZATION: the server never normalizes. It stores the code points you send
103105
and verifies a signature against those bytes, so NFC and NFD of one word are two

‎src/store.py‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -267,7 +267,9 @@
267267
# available after an interval is missed, instead of losing the window entirely.
268268
SNAPSHOT_KEEP_SECONDS = 30 * 3600
269269
IDLE_SECONDS = 7 * 86400 # untouched rooms/notes are reaped, so squatting expires
270-
REAP_EVERY = 300
270+
# A full store walk is worth amortizing: cleanup and count repair may lag ten minutes.
271+
# Retention ages stay separate; making a pass less frequent does not retire data sooner.
272+
REAP_EVERY = 600
271273
# A room that never got past its first message is a monologue, not a conversation: someone
272274
# said one thing, nobody answered, and it is holding a slot against MAX_ROOMS. A week is
273275
# what a conversation that stopped is worth; a day is what an unanswered opener is worth.
@@ -2556,15 +2558,8 @@ def _write_record(
25562558
# Also under the lock, or two concurrent replays of one captured URL would both
25572559
# read the same "last nonce" and both write.
25582560
if did is not None:
2559-
# The signature is `did: str | None, nonce: int | None`, which does not say
2560-
# that a signed write must carry both. Assert it rather than assume it: with
2561-
# nonce None this used to reach `None <= int` and raise TypeError — a 500 on
2562-
# the replay-protection path instead of a refusal that says what was wrong.
2563-
if nonce is None:
2564-
raise StoreError(
2565-
"a signed write must carry a nonce: it is what makes a captured "
2566-
"signed URL single-use. Send 1-19 digits, counting up per key per room"
2567-
)
2561+
# Validated before _reap and the create gate above; narrow the optional type.
2562+
assert nonce is not None
25682563
previous = _last_nonce(root, room, did)
25692564
if previous is not None and nonce <= previous:
25702565
raise StoreError(
@@ -2671,6 +2666,13 @@ def note_set(
26712666
ns_dir = _note_ns_dir(root, ns)
26722667
value = clean_text(value, MAX_VALUE_CHARS)
26732668
_reap(root)
2669+
# A missing note cannot satisfy CAS. Refuse before the create gate makes a sidecar
2670+
# and namespace: those artifacts survive a failed reservation but consume no quota.
2671+
# Reap first: the sweep can remove an idle note that existed at request entry.
2672+
# This is a valid observation even if another caller creates immediately afterwards;
2673+
# existing notes still compare under the lock below.
2674+
if expect is not None and not path.exists():
2675+
raise StoreConflictError(f"note {ns}/{key} changed since you read it", None)
26742676
# No cap check before the gate any more. One ran here to shed a full store's worth of
26752677
# refusals without queueing for a service-wide create gate first; the gate is NOTES_FILE's
26762678
# own lock now, held for two small file operations (#578), so a refusal costs the lock it

‎tests/http/test_abuse_bounds.py‎

Lines changed: 97 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,97 @@
1+
"""Bounded local regressions for resource consumption, including refused requests."""
2+
3+
import asyncio
4+
5+
import _client
6+
import pytest
7+
8+
import app
9+
import store
10+
11+
client = _client.client
12+
13+
14+
@pytest.mark.parametrize("path", ["/r/slow", "/r/events", "/kv/probe/slow"])
15+
@pytest.mark.parametrize("trickle", [False, True])
16+
def test_upload_has_a_total_deadline(client, monkeypatch, tmp_path, path, trickle):
17+
monkeypatch.setattr(app, "BODY_TIMEOUT", 0.03, raising=False)
18+
19+
async def check():
20+
messages = []
21+
22+
async def receive():
23+
# Progress must not reset the deadline; silence must expire as well.
24+
await asyncio.sleep(0.005 if trickle else 10)
25+
return {"type": "http.request", "body": b" ", "more_body": True}
26+
27+
async def send(message):
28+
messages.append(message)
29+
30+
scope = {
31+
"type": "http",
32+
"asgi": {"version": "3.0"},
33+
"http_version": "1.1",
34+
"method": "POST",
35+
"scheme": "http",
36+
"path": path,
37+
"query_string": b"",
38+
"headers": [(b"content-type", b"application/json")],
39+
"client": ("127.0.0.1", 1234),
40+
"server": ("127.0.0.1", 80),
41+
}
42+
task = asyncio.create_task(app.app(scope, receive, send))
43+
done, _ = await asyncio.wait({task}, timeout=0.5)
44+
if not done:
45+
task.cancel()
46+
with pytest.raises(asyncio.CancelledError):
47+
await task
48+
pytest.fail("an unfinished upload outlived its total body deadline")
49+
await task
50+
start = messages[0]
51+
assert start["status"] == 408
52+
assert (b"connection", b"close") in start["headers"]
53+
assert b"body" in messages[-1]["body"]
54+
55+
asyncio.run(check())
56+
assert not list(tmp_path.rglob("*.txt"))
57+
assert not list(tmp_path.rglob("*.jsonl"))
58+
59+
60+
@pytest.mark.parametrize("lane", ["get", "post"])
61+
def test_missing_note_cas_refusals_create_no_artifacts(client, tmp_path, lane):
62+
for i in range(12):
63+
path = f"/kv/probe-{i}/key"
64+
if lane == "get":
65+
response = client.get(path + "/set/new", params={"if": "old"})
66+
else:
67+
response = client.post(path, json={"value": "new", "if": "old"})
68+
assert response.status_code == 409
69+
assert "no note there" in response.text
70+
assert not (tmp_path / "notes").exists()
71+
assert store._note_count(tmp_path) == 0
72+
# Ordinary creation and CAS still work after the rejected requests.
73+
assert client.post("/kv/probe/key", json={"value": "old"}).status_code == 200
74+
assert client.post("/kv/probe/key", json={"value": "new", "if": "old"}).status_code == 200
75+
76+
77+
@pytest.mark.parametrize("lane", ["get", "post"])
78+
def test_reaped_note_cas_refusal_recreates_no_artifacts(client, tmp_path, lane):
79+
endpoint = "/kv/probe/key"
80+
assert client.post(endpoint, json={"value": "old"}).status_code == 200
81+
path = store.note_path(tmp_path, "probe", "key")
82+
lock = path.with_suffix(path.suffix + ".lock")
83+
_client._age(path, store.IDLE_SECONDS + 60)
84+
_client._age(lock, store.IDLE_SECONDS + 60)
85+
_client._age(tmp_path / ".reaped", store.REAP_EVERY + 60)
86+
# This note exists at request entry but the request's due sweep removes it.
87+
assert path.exists()
88+
if lane == "get":
89+
response = client.get(endpoint + "/set/new", params={"if": "old"})
90+
else:
91+
response = client.post(endpoint, json={"value": "new", "if": "old"})
92+
assert response.status_code == 409
93+
assert "no note there" in response.text
94+
assert not path.exists()
95+
assert not lock.exists()
96+
assert not store._note_ns_dir(tmp_path, "probe").exists()
97+
assert store._note_count(tmp_path) == 0

‎tests/http/test_rooms.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,12 +51,12 @@ def test_posting_to_the_events_room_documents_what_it_really_answers(client):
5151
"""`/r/events` is the ordinary room POST handler with one room that always says no, so
5252
the body is read and parsed *before* the refusal — a malformed or oversized body never
5353
reaches the 403. Documenting only the 403 promised a client one outcome and delivered
54-
three. Review catch on #40.
54+
several. Review catch on #40; slow bodies can now time out before the refusal too.
5555
"""
5656
import app as app_module
5757

5858
documented = client.get("/openapi.json").json()["paths"]["/r/events"]["post"]
59-
assert set(documented["responses"]) == {"400", "403", "413", "429"}
59+
assert set(documented["responses"]) == {"400", "403", "408", "413", "429"}
6060
# It parses a body, so it declares one.
6161
assert (
6262
"text" in (documented["requestBody"]["content"]["application/json"]["schema"]["properties"])

0 commit comments

Comments
 (0)