From 42e39e41de27488f660cf935c9ee150008f0bd21 Mon Sep 17 00:00:00 2001 From: daavoo Date: Wed, 9 Sep 2026 11:11:41 +0200 Subject: [PATCH] feat(files): fsspec backend so any filesystem can hold uploaded files files_backend: fsspec plus a files_url (gcs://, abfs://, s3://, sftp://, file://, ...) and files_storage_options reach whatever filesystem fsspec has an implementation installed for, through the same FileStore protocol the local and boto3 S3 backends implement. fsspec was already in the tree through any-llm and is now a declared dependency. Co-Authored-By: Claude Fable 5.1 --- config.example.yml | 4 + docs/files.md | 9 +- pyproject.toml | 10 ++ src/gateway/api/routes/settings.py | 2 + src/gateway/core/config.py | 24 +++- src/gateway/services/file_store.py | 170 ++++++++++++++++++++++++++- tests/unit/test_fsspec_file_store.py | 148 +++++++++++++++++++++++ uv.lock | 2 + 8 files changed, 362 insertions(+), 7 deletions(-) create mode 100644 tests/unit/test_fsspec_file_store.py diff --git a/config.example.yml b/config.example.yml index fcf3e90c4b..c33e2373a1 100644 --- a/config.example.yml +++ b/config.example.yml @@ -82,6 +82,10 @@ providers: # File and image normalization. See docs/files.md. # files_enabled: true # files_local_dir: "./otari-files" +# Any fsspec filesystem instead of a local directory or boto3 S3: +# files_backend: fsspec +# files_url: "gcs://my-bucket/otari-files" +# files_storage_options: { project: "my-project" } # files_max_bytes: 536870912 # files_retention_hours: 168 # files_sweep_interval_sec: 3600 diff --git a/docs/files.md b/docs/files.md index 961cd0212b..253020c997 100644 --- a/docs/files.md +++ b/docs/files.md @@ -157,7 +157,14 @@ in order: See [config.example.yml](../config.example.yml) for the full list. Key knobs: - `files_enabled`, `files_backend`, `files_local_dir`, `files_max_bytes`, -`files_retention_hours`: upload storage. An expired file answers 404 at once, +`files_retention_hours`: upload storage. `files_backend` is `local` (a +directory), `s3` (boto3, `files_s3_*`), or `fsspec`: any filesystem +[fsspec](https://filesystem-spec.readthedocs.io) has an implementation for, +named by `files_url` (`gcs://bucket/prefix`, `abfs://container/prefix`, +`s3://bucket/prefix`, `sftp://host/path`, `file:///path`, ...) with the +implementation's own keyword arguments in `files_storage_options`. Install the +implementation package for the protocol (`gcsfs`, `adlfs`, `s3fs`, `paramiko`); +most read their standard credential environment variables on their own. An expired file answers 404 at once, and the background sweep (`files_sweep_interval_sec`, hourly by default, `0` to disable) then reclaims its bytes and row along with those of deleted files. - `file_understanding_enabled`: master switch for content normalization. diff --git a/pyproject.toml b/pyproject.toml index 57024172dc..5587735a3d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -45,6 +45,11 @@ dependencies = [ # let a proof taken on one domain admit addresses at another. "idna>=3.7", "fastapi>=0.115.0", + # The generic files backend (`FsspecFileStore`): one URL reaches whichever + # filesystem the operator has an fsspec implementation installed for. + # Already in the tree through any-llm, declared because the gateway imports + # it directly. + "fsspec>=2024.6.0", "genai-prices>=0.1.0", "markitdown[docx,pptx,xlsx,pdf]>=0.1.0", "mcp>=1.28.1,<2.0.0", @@ -196,6 +201,11 @@ ignore_missing_imports = true module = ["trafilatura.*"] ignore_missing_imports = true +[[tool.mypy.overrides]] +# fsspec ships no stubs; the file store drives it through a handful of calls. +module = ["fsspec", "fsspec.*"] +ignore_missing_imports = true + [[tool.mypy.overrides]] # Optional/untyped extraction deps imported at the root (e.g. `import pypdfium2`, # `from markitdown import ...`). Bare root names are required — a `foo.*` pattern diff --git a/src/gateway/api/routes/settings.py b/src/gateway/api/routes/settings.py index 67c5b41e87..52372c2cb3 100644 --- a/src/gateway/api/routes/settings.py +++ b/src/gateway/api/routes/settings.py @@ -112,6 +112,7 @@ "files_enabled", "files_backend", "files_local_dir", + "files_url", "files_max_bytes", "files_retention_hours", "files_sweep_interval_sec", @@ -204,6 +205,7 @@ "oauth_github_client_secret", "web_search_provider_api_key", "web_search_backend_token", + "files_storage_options", # Structured blocks. ``ConfigField.value`` is bool/int/float/str/list[str], # so a dict or a nested model has no representation here at all. Each of # these has its own surface where it can be rendered as what it is diff --git a/src/gateway/core/config.py b/src/gateway/core/config.py index f7e1d60322..9382864f92 100644 --- a/src/gateway/core/config.py +++ b/src/gateway/core/config.py @@ -947,7 +947,10 @@ class GatewayConfig(BaseSettings): ) files_backend: str = Field( default="local", - description="Blob backend for uploaded file bytes: 'local' (filesystem) or 's3'. Future: 'gcs'.", + description=( + "Blob backend for uploaded file bytes: 'local' (a directory), 's3' (boto3), or 'fsspec' " + "(any filesystem fsspec has an implementation for, named by files_url)." + ), ) files_local_dir: str = Field( default="./otari-files", @@ -972,6 +975,25 @@ class GatewayConfig(BaseSettings): "'us-east-1' when unset." ), ) + files_url: str | None = Field( + default=None, + description=( + "Root URL for the 'fsspec' files backend, e.g. 'gcs://bucket/otari-files', " + "'abfs://container/prefix', 's3://bucket/prefix', 'sftp://host/path' or " + "'file:///var/lib/otari/files'. The protocol picks the fsspec implementation, which " + "must be installed (gcsfs, adlfs, s3fs, paramiko, ...). Required when files_backend " + "is 'fsspec'." + ), + ) + files_storage_options: dict[str, Any] = Field( + default_factory=dict, + description=( + "Keyword arguments for the fsspec implementation behind files_url: credentials, " + "endpoint URLs, regions, project ids. Passed through untouched and never logged; " + "most implementations also read their standard environment variables, so this " + "can usually stay empty." + ), + ) files_max_bytes: int = Field( default=512 * 1024 * 1024, ge=1, diff --git a/src/gateway/services/file_store.py b/src/gateway/services/file_store.py index b8af3bdfd8..1c2a9fdefd 100644 --- a/src/gateway/services/file_store.py +++ b/src/gateway/services/file_store.py @@ -11,18 +11,20 @@ instead of buffering an entire file, which is what actually bounds memory use for concurrent large uploads (see issue #156). -Only a local-filesystem backend ships today; ``S3FileStore`` / ``GCSFileStore`` -can implement the same :class:`FileStore` protocol without touching callers. +Three backends implement the :class:`FileStore` protocol: a local directory, +S3 through boto3, and :class:`FsspecFileStore`, which reaches any filesystem +`fsspec `_ has an implementation for +(GCS, Azure, SFTP, HDFS, WebDAV, and S3 again) from one ``files_url``. """ from __future__ import annotations import asyncio import tempfile -from collections.abc import AsyncGenerator, AsyncIterator, Iterator +from collections.abc import AsyncGenerator, AsyncIterator, Iterator, Mapping from contextlib import asynccontextmanager, contextmanager from pathlib import Path -from typing import IO, TYPE_CHECKING, Protocol, runtime_checkable +from typing import IO, TYPE_CHECKING, Any, Protocol, runtime_checkable from gateway.core.config import GatewayConfig from gateway.log_config import logger @@ -367,6 +369,159 @@ async def delete(self, storage_ref: str) -> None: await asyncio.to_thread(self._client.delete_object, Bucket=self._bucket, Key=storage_ref) +@contextmanager +def _translate_fsspec_errors(storage_ref: str) -> Iterator[None]: + """Re-raise whatever an fsspec implementation threw as the ``OSError`` family. + + fsspec's own filesystems raise ``FileNotFoundError`` and ``PermissionError`` + for the common cases, but a third-party implementation may surface its + client's exception class instead (a botocore or google-api error), and the + route and sweep callers only know ``OSError``, exactly as they do for the S3 + backend. A missing object stays ``FileNotFoundError`` so callers can tell + "already gone" from "broken". + """ + try: + yield + except FileNotFoundError: + raise + except OSError as exc: + msg = f"fsspec operation failed for {storage_ref!r}: {exc}" + raise OSError(msg) from exc + except Exception as exc: # noqa: BLE001 — a backend's own client error + msg = f"fsspec operation failed for {storage_ref!r}: {exc}" + raise OSError(msg) from exc + + +class FsspecFileStore: + """A :class:`FileStore` over any `fsspec `_ filesystem. + + ``url`` names the root the store writes under, ``s3://bucket/otari-files``, + ``gcs://bucket/prefix``, ``abfs://container/prefix``, ``file:///var/otari``, + ``memory://`` and so on; whatever protocol fsspec can resolve with the + implementation packages installed (``s3fs``, ``gcsfs``, ``adlfs``, ...). + ``storage_options`` go to that implementation as its constructor keyword + arguments, which is where credentials, endpoints and regions live, so they + are never logged here. + + Every call goes through fsspec's synchronous API on a worker thread, the way + the S3 backend drives boto3: the async implementations exist only for a few + protocols, and the sync API is the one every implementation has. + """ + + def __init__(self, url: str, storage_options: Mapping[str, Any] | None = None) -> None: + try: + from fsspec.core import url_to_fs + except ImportError as exc: # pragma: no cover - fsspec is a declared dependency + msg = "FsspecFileStore requires fsspec" + raise ImportError(msg) from exc + + fs, root = url_to_fs(url, **dict(storage_options or {})) + self._fs = fs + self._root = root.rstrip("/") + + def _resolve(self, storage_ref: str) -> str: + """Join ``storage_ref`` under the root, rejecting anything that could leave it. + + A server-generated ref has no ``..`` in it; this is defense-in-depth for + the day one comes from elsewhere, matching the local backend. + """ + parts = storage_ref.split("/") + if not storage_ref or storage_ref.startswith("/") or any(part in ("", ".", "..") for part in parts): + msg = f"Invalid storage_ref escapes the file store root: {storage_ref!r}" + raise ValueError(msg) + return f"{self._root}/{storage_ref}" if self._root else storage_ref + + def _mkparent(self, path: str) -> None: + # Object stores have no directories and treat this as a no-op; a + # filesystem-like backend needs it before the first write into a shard. + self._fs.makedirs(path.rsplit("/", 1)[0], exist_ok=True) + + async def put(self, file_id: str, data: bytes) -> str: + ref = _shard_key(file_id) + path = self._resolve(ref) + + def _write() -> None: + self._mkparent(path) + self._fs.pipe_file(path, data) + + with _translate_fsspec_errors(ref): + await asyncio.to_thread(_write) + return ref + + async def get(self, storage_ref: str) -> bytes: + path = self._resolve(storage_ref) + with _translate_fsspec_errors(storage_ref): + data: bytes = await asyncio.to_thread(self._fs.cat_file, path) + return data + + async def put_stream(self, file_id: str, chunks: AsyncIterator[bytes]) -> tuple[str, int]: + ref = _shard_key(file_id) + path = self._resolve(ref) + total = 0 + + def _open() -> IO[bytes]: + self._mkparent(path) + handle: IO[bytes] = self._fs.open(path, "wb") + return handle + + def _discard_partial() -> None: + try: + self._fs.rm(path) + except FileNotFoundError: + pass + + with _translate_fsspec_errors(ref): + handle = await asyncio.to_thread(_open) + try: + try: + async for chunk in chunks: + total += len(chunk) + with _translate_fsspec_errors(ref): + await asyncio.to_thread(handle.write, chunk) + finally: + # Object-store handles upload on close, so the close is part of + # the write and its failure is a write failure. Shielded like the + # local backend's: this also runs while a cancellation unwinds. + with _translate_fsspec_errors(ref): + await asyncio.shield(asyncio.to_thread(handle.close)) + except BaseException: + try: + await asyncio.shield(asyncio.to_thread(_discard_partial)) + except Exception as cleanup_exc: # noqa: BLE001 + logger.warning("put_stream: failed to remove partial blob %s: %s", ref, cleanup_exc) + raise + return ref, total + + async def get_stream(self, storage_ref: str) -> AsyncGenerator[bytes, None]: + path = self._resolve(storage_ref) + with _translate_fsspec_errors(storage_ref): + handle: IO[bytes] = await asyncio.to_thread(self._fs.open, path, "rb") + try: + while True: + with _translate_fsspec_errors(storage_ref): + chunk = await asyncio.to_thread(handle.read, _STREAM_CHUNK_BYTES) + if not chunk: + break + yield chunk + finally: + try: + await asyncio.shield(asyncio.to_thread(handle.close)) + except Exception as close_exc: # noqa: BLE001 + logger.warning("get_stream: failed to close handle for %s: %s", storage_ref, close_exc) + + async def delete(self, storage_ref: str) -> None: + path = self._resolve(storage_ref) + + def _rm() -> None: + try: + self._fs.rm(path) + except FileNotFoundError: + logger.debug("file_store delete: %s already absent", storage_ref) + + with _translate_fsspec_errors(storage_ref): + await asyncio.to_thread(_rm) + + def build_file_store(config: GatewayConfig) -> FileStore: """Construct the configured :class:`FileStore` backend.""" backend = config.files_backend.strip().lower() @@ -377,5 +532,10 @@ def build_file_store(config: GatewayConfig) -> FileStore: msg = "files_s3_bucket is required when files_backend is 's3'" raise ValueError(msg) return S3FileStore(config.files_s3_bucket, config.files_s3_endpoint_url, config.files_s3_region) - msg = f"Unsupported files_backend: {config.files_backend!r} (supported: 'local', 's3')" + if backend == "fsspec": + if not config.files_url: + msg = "files_url is required when files_backend is 'fsspec'" + raise ValueError(msg) + return FsspecFileStore(config.files_url, config.files_storage_options) + msg = f"Unsupported files_backend: {config.files_backend!r} (supported: 'local', 's3', 'fsspec')" raise ValueError(msg) diff --git a/tests/unit/test_fsspec_file_store.py b/tests/unit/test_fsspec_file_store.py new file mode 100644 index 0000000000..2a6dce7f29 --- /dev/null +++ b/tests/unit/test_fsspec_file_store.py @@ -0,0 +1,148 @@ +"""Unit tests for the fsspec-backed file store. + +Runs on fsspec's built-in ``memory://`` and ``file://`` filesystems, so the +suite needs no cloud implementation package and no network. What it proves is +the adapter's own contract (refs, streaming, cleanup, error translation); the +cloud implementations are fsspec's to keep working. +""" + +from __future__ import annotations + +import asyncio +from collections.abc import AsyncIterator +from pathlib import Path + +import fsspec +import pytest + +from gateway.core.config import GatewayConfig +from gateway.services.file_store import FsspecFileStore, build_file_store + + +async def _iter(chunks: list[bytes]) -> AsyncIterator[bytes]: + for chunk in chunks: + yield chunk + + +@pytest.fixture +def memory_root() -> str: + # The memory filesystem is process-global; give each test its own prefix + # and clear it afterwards so one test's blobs never show up in another. + root = "memory://otari-test" + fs = fsspec.filesystem("memory") + if fs.exists("otari-test"): + fs.rm("otari-test", recursive=True) + return root + + +@pytest.mark.asyncio +async def test_put_get_roundtrip(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + ref = await store.put("file-abcdef0123", b"hello bytes") + assert ref == "ab/file-abcdef0123" + assert await store.get(ref) == b"hello bytes" + + +@pytest.mark.asyncio +async def test_put_stream_and_get_stream_roundtrip(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + payload = b"x" * (2 * 1024 * 1024 + 5) + ref, size = await store.put_stream("file-streamtest01", _iter([payload[:1000], payload[1000:]])) + assert size == len(payload) + collected = bytearray() + async for chunk in store.get_stream(ref): + collected.extend(chunk) + assert bytes(collected) == payload + + +@pytest.mark.asyncio +async def test_put_stream_removes_partial_blob_on_failure(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + + async def _failing() -> AsyncIterator[bytes]: + yield b"partial" + raise RuntimeError("client went away") + + with pytest.raises(RuntimeError): + await store.put_stream("file-partial00001", _failing()) + assert not fsspec.filesystem("memory").exists("otari-test/pa/file-partial00001") + + +@pytest.mark.asyncio +async def test_put_stream_removes_partial_blob_on_cancellation(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + started = asyncio.Event() + + async def _slow() -> AsyncIterator[bytes]: + yield b"first" + started.set() + await asyncio.sleep(30) + yield b"never" + + task = asyncio.create_task(store.put_stream("file-cancel000001", _slow())) + await started.wait() + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + assert not fsspec.filesystem("memory").exists("otari-test/ca/file-cancel000001") + + +@pytest.mark.asyncio +async def test_missing_blob_is_file_not_found(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + with pytest.raises(FileNotFoundError): + await store.get("no/file-nope") + with pytest.raises(FileNotFoundError): + async for _ in store.get_stream("no/file-nope"): + pass + + +@pytest.mark.asyncio +async def test_delete_is_idempotent(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + ref = await store.put("file-deleteme0001", b"x") + await store.delete(ref) + await store.delete(ref) + with pytest.raises(FileNotFoundError): + await store.get(ref) + + +@pytest.mark.asyncio +async def test_rejects_refs_that_could_leave_the_root(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + for bad in ("../escape", "/absolute", "a//b", "a/./b", ""): + with pytest.raises(ValueError): + await store.get(bad) + + +@pytest.mark.asyncio +async def test_backend_client_errors_become_oserror(memory_root: str) -> None: + store = FsspecFileStore(memory_root) + + def _boom(*_args: object, **_kwargs: object) -> None: + raise RuntimeError("some client's own exception class") + + store._fs.cat_file = _boom + with pytest.raises(OSError, match="fsspec operation failed"): + await store.get("ab/file-abcdef0123") + + +@pytest.mark.asyncio +async def test_local_file_protocol_writes_under_the_root(tmp_path: Path) -> None: + store = FsspecFileStore(f"file://{tmp_path}") + ref = await store.put("file-abcdef0123", b"on disk") + assert (tmp_path / "ab" / "file-abcdef0123").read_bytes() == b"on disk" + assert await store.get(ref) == b"on disk" + + +def test_build_file_store_fsspec_requires_url() -> None: + cfg = GatewayConfig(files_backend="fsspec") + with pytest.raises(ValueError, match="files_url"): + build_file_store(cfg) + + +def test_build_file_store_fsspec(tmp_path: Path) -> None: + cfg = GatewayConfig( + files_backend="fsspec", files_url=f"file://{tmp_path}", files_storage_options={"auto_mkdir": True} + ) + assert isinstance(build_file_store(cfg), FsspecFileStore) diff --git a/uv.lock b/uv.lock index a48a1ed207..b926a9f766 100644 --- a/uv.lock +++ b/uv.lock @@ -1025,6 +1025,7 @@ dependencies = [ { name = "cryptography" }, { name = "dnspython" }, { name = "fastapi" }, + { name = "fsspec" }, { name = "genai-prices" }, { name = "idna" }, { name = "markitdown", extra = ["docx", "pdf", "pptx", "xlsx"] }, @@ -1086,6 +1087,7 @@ requires-dist = [ { name = "cryptography", specifier = ">=50.0.0" }, { name = "dnspython", specifier = ">=2.7.0" }, { name = "fastapi", specifier = ">=0.115.0" }, + { name = "fsspec", specifier = ">=2024.6.0" }, { name = "genai-prices", specifier = ">=0.1.0" }, { name = "idna", specifier = ">=3.7" }, { name = "markitdown", extras = ["docx", "pptx", "xlsx", "pdf"], specifier = ">=0.1.0" },