Skip to content

Commit bbf3ec4

Browse files
authored
Create storage.py
1 parent c95995c commit bbf3ec4

1 file changed

Lines changed: 148 additions & 0 deletions

File tree

‎app/infra/storage.py‎

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
"""Durable blob storage: the system of record for uploads and artifacts.
2+
3+
Local disk dies with the pod, so a conversation's files must live
4+
somewhere durable (object storage in production). This module is that
5+
durability layer, behind a small key/bytes interface with two backends:
6+
7+
* ``LocalStorage`` -- a directory tree, the default (dev / single node).
8+
* ``S3Storage`` -- an S3-compatible bucket (prod / multi-node).
9+
10+
Two-tier model: this store is the *durable* tier (the source of
11+
truth). The sandbox harness still needs files as real paths in its
12+
working directory, so the per-conversation workspace remains a
13+
*sandbox-facing* tier -- a materialized working copy that is written
14+
through to this store on upload and can be rehydrated from it after a
15+
pod loss. So Storage never replaces the workspace; it backs it.
16+
17+
Keys are opaque, ``/``-joined strings (e.g. ``"<conversation>/<name>"``);
18+
a backend maps them to a path or an object key. This module is a pure
19+
infra primitive: no ``Session``, no ORM entity, no request scope -- it
20+
imports only stdlib + config (and, lazily, boto3 for S3).
21+
"""
22+
23+
from __future__ import annotations
24+
25+
from collections.abc import Iterator
26+
from pathlib import Path
27+
from typing import BinaryIO, Protocol, runtime_checkable
28+
29+
from .config import get_settings
30+
31+
32+
class StorageError(RuntimeError):
33+
"""A backend operation failed (missing key, transport error)."""
34+
35+
36+
@runtime_checkable
37+
class Storage(Protocol):
38+
"""Durable key/bytes store. Keys are opaque ``/``-joined strings."""
39+
40+
def put_bytes(self, key: str, data: bytes) -> None:
41+
"""Store ``data`` under ``key`` (overwrites)."""
42+
43+
def put_stream(self, key: str, stream: BinaryIO) -> int:
44+
"""Store a stream under ``key``; return the number of bytes written."""
45+
46+
def get_bytes(self, key: str) -> bytes:
47+
"""Return the bytes stored under ``key`` (raises if absent)."""
48+
49+
def open_stream(self, key: str) -> Iterator[bytes]:
50+
"""Yield the object's bytes in chunks (raises if absent)."""
51+
52+
def exists(self, key: str) -> bool: ...
53+
54+
def size(self, key: str) -> int | None:
55+
"""Byte length of the stored object, or None if absent."""
56+
57+
def delete(self, key: str) -> None:
58+
"""Remove ``key``; a missing key is not an error (idempotent)."""
59+
60+
61+
_CHUNK = 1024 * 1024
62+
63+
64+
def _safe_key(key: str) -> str:
65+
"""Reject keys that could escape the storage root.
66+
67+
Keys are ``/``-joined; ``..`` or an absolute component would let a
68+
crafted key walk out of a LocalStorage root or forge an S3 prefix.
69+
"""
70+
if not key or key.startswith("/") or ".." in Path(key).parts:
71+
raise StorageError(f"unsafe storage key: {key!r}")
72+
return key
73+
74+
75+
class LocalStorage:
76+
"""Filesystem-backed store rooted at a directory.
77+
78+
The default backend: durable enough for a single node, and the same
79+
on-disk layout the workspace already uses, so dev behavior is
80+
unchanged.
81+
"""
82+
83+
def __init__(self, root: str | Path) -> None:
84+
self._root = Path(root)
85+
86+
def _path(self, key: str) -> Path:
87+
return self._root / _safe_key(key)
88+
89+
def put_bytes(self, key: str, data: bytes) -> None:
90+
path = self._path(key)
91+
path.parent.mkdir(parents=True, exist_ok=True)
92+
path.write_bytes(data)
93+
94+
def put_stream(self, key: str, stream: BinaryIO) -> int:
95+
path = self._path(key)
96+
path.parent.mkdir(parents=True, exist_ok=True)
97+
written = 0
98+
with path.open("wb") as out:
99+
while chunk := stream.read(_CHUNK):
100+
written += len(chunk)
101+
out.write(chunk)
102+
return written
103+
104+
def get_bytes(self, key: str) -> bytes:
105+
try:
106+
return self._path(key).read_bytes()
107+
except OSError as exc:
108+
raise StorageError(f"cannot read {key!r}: {exc}") from exc
109+
110+
def open_stream(self, key: str) -> Iterator[bytes]:
111+
path = self._path(key)
112+
if not path.is_file():
113+
raise StorageError(f"no such key: {key!r}")
114+
115+
def _gen() -> Iterator[bytes]:
116+
with path.open("rb") as fh:
117+
while chunk := fh.read(_CHUNK):
118+
yield chunk
119+
120+
return _gen()
121+
122+
def exists(self, key: str) -> bool:
123+
return self._path(key).is_file()
124+
125+
def size(self, key: str) -> int | None:
126+
path = self._path(key)
127+
return path.stat().st_size if path.is_file() else None
128+
129+
def delete(self, key: str) -> None:
130+
self._path(key).unlink(missing_ok=True)
131+
132+
133+
def get_storage() -> Storage:
134+
"""The configured durable store (process-wide, chosen by settings)."""
135+
settings = get_settings()
136+
backend = settings.storage.backend
137+
if backend == "local":
138+
return LocalStorage(settings.storage.local_root)
139+
if backend == "s3":
140+
from .storage_s3 import S3Storage
141+
142+
return S3Storage(
143+
bucket=settings.storage.s3_bucket,
144+
prefix=settings.storage.s3_prefix,
145+
region=settings.storage.s3_region or None,
146+
endpoint_url=settings.storage.s3_endpoint_url or None,
147+
)
148+
raise ValueError(f"unknown storage backend: {backend!r}")

0 commit comments

Comments
 (0)