From 2e785fbdcfb7795feca0a702ea35bf281a61b16e Mon Sep 17 00:00:00 2001 From: Patrick Buckley Date: Tue, 2 Jun 2026 22:48:43 -0700 Subject: [PATCH] feat(attachments): per-node in-memory pending-upload buffer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The content-addressed model writes blob bytes to workstream_attachments only at send-commit (so every stored blob is born referenced). Pending uploads — between the upload request and the send that references them — live in this per-node, content-hash- keyed buffer, scoped to (ws_id, user_id) and bounded by a TTL + total-size ceiling (OOM-safety, not the removed per-user cap). ws->node affinity (HRW routing) keeps it process-local. Losing an unsent upload on crash/re-route is acceptable transient state. This replaces the persisted pending/reserved/consumed lifecycle + orphan-sweep. No consumers yet — the upload-endpoint rework and the send-commit drain wire onto it as the attachment cutover lands. --- tests/test_attachment_buffer.py | 102 ++++++++++++++++++ turnstone/core/attachment_buffer.py | 157 ++++++++++++++++++++++++++++ 2 files changed, 259 insertions(+) create mode 100644 tests/test_attachment_buffer.py create mode 100644 turnstone/core/attachment_buffer.py diff --git a/tests/test_attachment_buffer.py b/tests/test_attachment_buffer.py new file mode 100644 index 00000000..49f429d0 --- /dev/null +++ b/tests/test_attachment_buffer.py @@ -0,0 +1,102 @@ +"""Tests for the per-node pending-upload buffer (turnstone.core.attachment_buffer).""" + +from __future__ import annotations + +import hashlib + +from turnstone.core.attachment_buffer import ( + AttachmentBuffer, + StagedAttachment, + get_attachment_buffer, +) + + +def _stage( + buf: AttachmentBuffer, + *, + content: bytes = b"hi", + ws: str = "ws1", + user: str = "u1", + filename: str = "f.txt", + mime: str = "text/plain", + kind: str = "text", +) -> StagedAttachment: + return buf.stage( + ws_id=ws, user_id=user, filename=filename, mime_type=mime, kind=kind, content=content + ) + + +def test_stage_returns_content_hash_id_and_size() -> None: + buf = AttachmentBuffer() + entry = _stage(buf, content=b"hello") + assert entry.attachment_id == hashlib.sha256(b"hello").hexdigest() + assert entry.size_bytes == 5 + + +def test_stage_is_idempotent_for_identical_bytes() -> None: + buf = AttachmentBuffer() + a = _stage(buf, content=b"same") + b = _stage(buf, content=b"same") + assert a.attachment_id == b.attachment_id + assert len(buf.list_for(ws_id="ws1", user_id="u1")) == 1 # deduped by content hash + + +def test_get_enforces_scope() -> None: + buf = AttachmentBuffer() + entry = _stage(buf, ws="ws1", user="u1") + assert buf.get(entry.attachment_id, ws_id="ws1", user_id="u1") is not None + assert buf.get(entry.attachment_id, ws_id="ws2", user_id="u1") is None # wrong ws + assert buf.get(entry.attachment_id, ws_id="ws1", user_id="u2") is None # wrong user + + +def test_list_for_scopes_by_ws_and_user() -> None: + buf = AttachmentBuffer() + _stage(buf, content=b"a", ws="ws1", user="u1") + _stage(buf, content=b"b", ws="ws1", user="u1") + _stage(buf, content=b"c", ws="ws2", user="u1") + assert len(buf.list_for(ws_id="ws1", user_id="u1")) == 2 + assert len(buf.list_for(ws_id="ws2", user_id="u1")) == 1 + + +def test_discard_is_scope_checked() -> None: + buf = AttachmentBuffer() + entry = _stage(buf) + assert buf.discard(entry.attachment_id, ws_id="ws2", user_id="u1") is False + assert buf.discard(entry.attachment_id, ws_id="ws1", user_id="u1") is True + assert buf.get(entry.attachment_id, ws_id="ws1", user_id="u1") is None + + +def test_take_pops_in_scope_and_skips_others() -> None: + buf = AttachmentBuffer() + a = _stage(buf, content=b"a") + b = _stage(buf, content=b"b", ws="ws2") + taken = buf.take([a.attachment_id, b.attachment_id, "missing"], ws_id="ws1") + assert [t.attachment_id for t in taken] == [a.attachment_id] # only ws1's, missing skipped + assert buf.get(a.attachment_id, ws_id="ws1", user_id="u1") is None # popped + assert buf.get(b.attachment_id, ws_id="ws2", user_id="u1") is not None # other ws untouched + + +def test_ttl_eviction_on_access() -> None: + clock = [0.0] + buf = AttachmentBuffer(ttl_seconds=10.0, clock=lambda: clock[0]) + _stage(buf, content=b"x") + clock[0] = 11.0 # past the TTL + assert buf.list_for(ws_id="ws1", user_id="u1") == [] + + +def test_size_cap_evicts_oldest_first() -> None: + clock = [0.0] + buf = AttachmentBuffer(max_total_bytes=10, clock=lambda: clock[0]) + clock[0] = 1.0 + a = _stage(buf, content=b"aaaaa") # 5 bytes + clock[0] = 2.0 + b = _stage(buf, content=b"bbbbb") # +5 → 10, at the ceiling + clock[0] = 3.0 + c = _stage(buf, content=b"ccccc") # +5 → 15 > 10 → evict oldest (a) + ids = {e.attachment_id for e in buf.list_for(ws_id="ws1", user_id="u1")} + assert a.attachment_id not in ids + assert {b.attachment_id, c.attachment_id} <= ids + + +def test_singleton_getter_is_stable() -> None: + assert get_attachment_buffer() is get_attachment_buffer() diff --git a/turnstone/core/attachment_buffer.py b/turnstone/core/attachment_buffer.py new file mode 100644 index 00000000..e2a9a7c1 --- /dev/null +++ b/turnstone/core/attachment_buffer.py @@ -0,0 +1,157 @@ +"""Per-node in-memory buffer for pending (uploaded-but-unsent) attachments. + +In the content-addressed attachment model, blob bytes are written to +``workstream_attachments`` only at send-commit, so every *stored* blob is born +referenced (refcount >= 1). Between the upload request and the send that references it, +the bytes live here — keyed by their content hash, scoped to ``(ws_id, user_id)``, and +bounded by a TTL plus a total-size ceiling so a flood of unsent uploads can't exhaust a +node. ws->node affinity (HRW routing) guarantees the upload and the later send hit the +same node, so this buffer is process-local. + +Losing an unsent upload on a crash / restart / node re-route is acceptable: it is +pre-commit transient state, never committed data — the user re-uploads. This is the +deliberate simplification that replaces the persisted pending/reserved/consumed +lifecycle (and its orphan-sweep) the old upload flow carried. +""" + +from __future__ import annotations + +import hashlib +import threading +import time +from dataclasses import dataclass +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from collections.abc import Callable, Iterable + +# OOM-safety backstops (per node), not a product policy — the per-user upload cap was +# removed. Generous: the goal is only to bound a pathological flood of unsent uploads. +_DEFAULT_MAX_TOTAL_BYTES = 1024 * 1024 * 1024 # 1 GiB +_DEFAULT_TTL_SECONDS = 3600.0 # unsent uploads expire after an hour + + +@dataclass(frozen=True, slots=True) +class StagedAttachment: + """A pending upload held in memory until the send that references it commits.""" + + attachment_id: str # content hash (sha256 hex) — also the content-addressed blob id + ws_id: str + user_id: str + filename: str + mime_type: str + kind: str # 'image' | 'text' + content: bytes + staged_at: float + + @property + def size_bytes(self) -> int: + return len(self.content) + + +class AttachmentBuffer: + """Thread-safe, size- and TTL-bounded store of pending uploads, keyed by content hash.""" + + def __init__( + self, + *, + max_total_bytes: int = _DEFAULT_MAX_TOTAL_BYTES, + ttl_seconds: float = _DEFAULT_TTL_SECONDS, + clock: Callable[[], float] = time.monotonic, + ) -> None: + self._lock = threading.Lock() + self._entries: dict[str, StagedAttachment] = {} + self._max_total_bytes = max_total_bytes + self._ttl = ttl_seconds + self._clock = clock + + def stage( + self, *, ws_id: str, user_id: str, filename: str, mime_type: str, kind: str, content: bytes + ) -> StagedAttachment: + """Hold *content* until a send commits it; returns the staged entry. + + The id is the content hash, so re-uploading identical bytes is idempotent (and + dedupes against the eventual content-addressed blob). + """ + handle = hashlib.sha256(content).hexdigest() + entry = StagedAttachment( + attachment_id=handle, + ws_id=ws_id, + user_id=user_id, + filename=filename, + mime_type=mime_type, + kind=kind, + content=content, + staged_at=self._clock(), + ) + with self._lock: + self._evict_expired_locked() + self._entries[handle] = entry + self._evict_over_capacity_locked() + return entry + + def get(self, handle: str, *, ws_id: str, user_id: str) -> StagedAttachment | None: + """Fetch a pending upload, enforcing the uploader's scope (else ``None``).""" + with self._lock: + entry = self._entries.get(handle) + if entry is None or entry.ws_id != ws_id or entry.user_id != user_id: + return None + return entry + + def list_for(self, *, ws_id: str, user_id: str) -> list[StagedAttachment]: + """Pending uploads for ``(ws_id, user_id)``.""" + with self._lock: + self._evict_expired_locked() + return [e for e in self._entries.values() if e.ws_id == ws_id and e.user_id == user_id] + + def discard(self, handle: str, *, ws_id: str, user_id: str) -> bool: + """Drop a pending upload (scope-checked); ``True`` if it existed.""" + with self._lock: + entry = self._entries.get(handle) + if entry is None or entry.ws_id != ws_id or entry.user_id != user_id: + return False + del self._entries[handle] + return True + + def take(self, handles: Iterable[str], *, ws_id: str) -> list[StagedAttachment]: + """Pop the staged entries for *handles* at send-commit (ws-scoped). + + Missing / out-of-scope handles are skipped — the send proceeds with whatever + committed (a buffer eviction or crash between upload and send drops the upload). + """ + with self._lock: + taken: list[StagedAttachment] = [] + for handle in handles: + entry = self._entries.get(handle) + if entry is not None and entry.ws_id == ws_id: + taken.append(entry) + del self._entries[handle] + return taken + + # -- eviction (caller holds the lock) ------------------------------------ + def _evict_expired_locked(self) -> None: + if self._ttl <= 0: + return + cutoff = self._clock() - self._ttl + stale = [h for h, e in self._entries.items() if e.staged_at < cutoff] + for handle in stale: + del self._entries[handle] + + def _evict_over_capacity_locked(self) -> None: + total = sum(e.size_bytes for e in self._entries.values()) + if total <= self._max_total_bytes: + return + # Evict oldest-first until under the ceiling. + for handle, entry in sorted(self._entries.items(), key=lambda kv: kv[1].staged_at): + if total <= self._max_total_bytes: + break + total -= entry.size_bytes + del self._entries[handle] + + +_BUFFER = AttachmentBuffer() + + +def get_attachment_buffer() -> AttachmentBuffer: + """The per-node pending-upload buffer singleton.""" + return _BUFFER