Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions config.example.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 8 additions & 1 deletion docs/files.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
10 changes: 10 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions src/gateway/api/routes/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@
"files_enabled",
"files_backend",
"files_local_dir",
"files_url",
"files_max_bytes",
"files_retention_hours",
"files_sweep_interval_sec",
Expand Down Expand Up @@ -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
Expand Down
24 changes: 23 additions & 1 deletion src/gateway/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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,
Expand Down
170 changes: 165 additions & 5 deletions src/gateway/services/file_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 <https://filesystem-spec.readthedocs.io>`_ 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
Expand Down Expand Up @@ -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 <https://filesystem-spec.readthedocs.io>`_ 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()
Expand All @@ -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)
Loading
Loading