"""Private filesystem backend used for tests and controlled local vertical slices.""" from __future__ import annotations import json import os import stat import tempfile from pathlib import Path from typing import BinaryIO, Mapping from arr_storage.contracts import BackendError, BackendObject, valid_object_key _METADATA_SUFFIX = ".arr-metadata.json" class FilesystemObjectBackend: """An immutable local backend with exclusive creates and private permissions.""" def __init__(self, root: Path, *, create: bool = False) -> None: if create: root.mkdir(parents=True, exist_ok=True, mode=0o700) self._root = root.resolve() if not self._root.is_dir(): raise ValueError("object backend root must be an existing directory") @property def root(self) -> Path: return self._root def put_file( self, object_key: str, source: str, mime_type: str, metadata: Mapping[str, str], *, if_absent: bool, ) -> BackendObject: if not if_absent: raise BackendError("invalid_request") destination = self._object_path(object_key) metadata_path = self._metadata_path(object_key) source_path = Path(source) try: source_stat = source_path.lstat() except OSError: raise BackendError("unavailable") from None if source_path.is_symlink() or not stat.S_ISREG(source_stat.st_mode): raise BackendError("invalid_request") destination.parent.mkdir(parents=True, exist_ok=True, mode=0o700) if destination.exists() or metadata_path.exists(): raise BackendError("conflict") descriptor = -1 destination_created = False try: descriptor = os.open(destination, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) destination_created = True with source_path.open("rb") as source_handle, os.fdopen(descriptor, "wb") as target: descriptor = -1 while True: chunk = source_handle.read(1024 * 1024) if not chunk: break target.write(chunk) target.flush() os.fsync(target.fileno()) self._write_metadata(metadata_path, mime_type, metadata) except FileExistsError: if destination_created: destination.unlink(missing_ok=True) raise BackendError("conflict") from None except BackendError: if destination_created: destination.unlink(missing_ok=True) raise except OSError: if destination_created: destination.unlink(missing_ok=True) raise BackendError("unavailable") from None finally: if descriptor >= 0: os.close(descriptor) return self.head(object_key) def head(self, object_key: str) -> BackendObject: object_path = self._object_path(object_key) metadata_path = self._metadata_path(object_key) try: if object_path.is_symlink() or metadata_path.is_symlink(): raise BackendError("unavailable") object_stat = object_path.stat() if not stat.S_ISREG(object_stat.st_mode): raise BackendError("not_found") raw = metadata_path.read_text(encoding="utf-8") metadata_payload = json.loads(raw) except FileNotFoundError: raise BackendError("not_found") from None except BackendError: raise except (OSError, UnicodeError, json.JSONDecodeError): raise BackendError("unavailable") from None if ( not isinstance(metadata_payload, dict) or set(metadata_payload) != {"content_type", "metadata"} or not isinstance(metadata_payload.get("content_type"), str) or not isinstance(metadata_payload.get("metadata"), dict) or any( not isinstance(key, str) or not isinstance(value, str) for key, value in metadata_payload["metadata"].items() ) ): raise BackendError("unavailable") metadata = dict(metadata_payload["metadata"]) metadata.setdefault("arr-backend-content-type", metadata_payload["content_type"]) return BackendObject( object_key=object_key, byte_size=object_stat.st_size, metadata=metadata, ) def open_reader(self, object_key: str) -> BinaryIO: path = self._object_path(object_key) try: if path.is_symlink(): raise BackendError("unavailable") descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)) opened_stat = os.fstat(descriptor) if not stat.S_ISREG(opened_stat.st_mode): os.close(descriptor) raise BackendError("not_found") return os.fdopen(descriptor, "rb") except FileNotFoundError: raise BackendError("not_found") from None except BackendError: raise except OSError: raise BackendError("unavailable") from None def copy_object( self, source_key: str, destination_key: str, metadata: Mapping[str, str], *, if_absent: bool, ) -> BackendObject: source = self.head(source_key) mime_type = source.metadata.get("arr-mime-type", "application/octet-stream") try: with tempfile.NamedTemporaryFile(prefix="arr-object-copy-", delete=False) as temporary: temporary_path = Path(temporary.name) with self.open_reader(source_key) as reader: while True: chunk = reader.read(1024 * 1024) if not chunk: break temporary.write(chunk) temporary.flush() os.fsync(temporary.fileno()) return self.put_file( destination_key, str(temporary_path), mime_type, metadata, if_absent=if_absent, ) finally: if "temporary_path" in locals(): temporary_path.unlink(missing_ok=True) def delete_object(self, object_key: str) -> None: object_path = self._object_path(object_key) metadata_path = self._metadata_path(object_key) try: object_path.unlink(missing_ok=True) metadata_path.unlink(missing_ok=True) except OSError: raise BackendError("unavailable") from None def _object_path(self, object_key: str) -> Path: if not valid_object_key(object_key): raise BackendError("invalid_request") candidate = self._root.joinpath(*object_key.split("/")) try: candidate.parent.resolve().relative_to(self._root) except (OSError, ValueError): raise BackendError("invalid_request") from None return candidate def _metadata_path(self, object_key: str) -> Path: object_path = self._object_path(object_key) return object_path.with_name(object_path.name + _METADATA_SUFFIX) @staticmethod def _write_metadata( destination: Path, mime_type: str, metadata: Mapping[str, str], ) -> None: if any(not isinstance(key, str) or not isinstance(value, str) for key, value in metadata.items()): raise BackendError("invalid_request") payload = json.dumps( {"content_type": mime_type, "metadata": dict(metadata)}, ensure_ascii=True, sort_keys=True, separators=(",", ":"), ).encode("utf-8") descriptor = -1 created = False try: descriptor = os.open(destination, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) created = True with os.fdopen(descriptor, "wb") as target: descriptor = -1 target.write(payload) target.flush() os.fsync(target.fileno()) except FileExistsError: raise BackendError("conflict") from None except OSError: if created: destination.unlink(missing_ok=True) raise BackendError("unavailable") from None finally: if descriptor >= 0: os.close(descriptor)