208 lines
6.4 KiB
Python
208 lines
6.4 KiB
Python
from __future__ import annotations
|
|
|
|
import fcntl
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
from typing import Iterable
|
|
|
|
from .logging_setup import log_warning
|
|
|
|
|
|
DEFAULT_DATA_DIR = Path("./data")
|
|
_FILE_NAME = "posted.json"
|
|
_EMPTY_DOCUMENT: dict[str, list[str]] = {"posted": []}
|
|
|
|
|
|
class PostedStoreError(RuntimeError):
|
|
"""Raised when the posted-state store cannot be used."""
|
|
|
|
|
|
def posted_path(data_dir: Path | None = None) -> Path:
|
|
"""Resolve the path of ``posted.json``.
|
|
|
|
When ``data_dir`` is ``None``, the directory is read from the
|
|
``DATA_DIR`` environment variable, falling back to ``./data``
|
|
(the host bind-mount described in the README).
|
|
"""
|
|
if data_dir is None:
|
|
data_dir = Path(os.environ.get("DATA_DIR", str(DEFAULT_DATA_DIR)))
|
|
return Path(data_dir) / _FILE_NAME
|
|
|
|
|
|
def _coerce_list(value: object) -> list[str]:
|
|
"""Filter ``value`` down to a list of strings, dropping anything else."""
|
|
if not isinstance(value, list):
|
|
return []
|
|
return [item for item in value if isinstance(item, str)]
|
|
|
|
|
|
def _read_existing(path: Path) -> dict[str, list[str]] | None:
|
|
"""Return the parsed JSON document at ``path`` or ``None`` on errors.
|
|
|
|
A missing file is **not** an error: returns ``None`` so the caller
|
|
can treat it as the empty document. Malformed JSON or unreadable
|
|
bytes return ``None`` and a warning is emitted.
|
|
"""
|
|
try:
|
|
with path.open("r", encoding="utf-8") as fh:
|
|
data = json.load(fh)
|
|
except FileNotFoundError:
|
|
return None
|
|
except (json.JSONDecodeError, OSError) as exc:
|
|
log_warning("posted_store_corrupt", path=str(path), error=type(exc).__name__)
|
|
return None
|
|
if not isinstance(data, dict):
|
|
return None
|
|
return data
|
|
|
|
|
|
def _ensure_parent(path: Path) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
|
def _write_document(path: Path, document: dict[str, list[str]]) -> None:
|
|
"""Atomically write ``document`` to ``path`` via temp-file + rename."""
|
|
_ensure_parent(path)
|
|
pid = os.getpid()
|
|
tmp = path.with_name(f"{path.name}.tmp.{pid}")
|
|
try:
|
|
with tmp.open("w", encoding="utf-8") as fh:
|
|
json.dump(document, fh, sort_keys=False, indent=2)
|
|
fh.flush()
|
|
os.fsync(fh.fileno())
|
|
os.replace(tmp, path)
|
|
except Exception:
|
|
try:
|
|
tmp.unlink()
|
|
except FileNotFoundError:
|
|
pass
|
|
raise
|
|
|
|
|
|
def _lock_path(path: Path):
|
|
"""Return an open file descriptor on the lock file inside the data dir.
|
|
|
|
The lock file lives in the same directory as ``posted.json`` so the
|
|
flock is always on a single inode, regardless of whether the JSON
|
|
file currently exists. The descriptor is opened in append mode so
|
|
concurrent readers never truncate it.
|
|
"""
|
|
_ensure_parent(path)
|
|
lock = path.with_name(f".{path.name}.lock")
|
|
fd = os.open(str(lock), os.O_CREAT | os.O_RDWR, 0o644)
|
|
return lock, fd
|
|
|
|
|
|
class PostedStore:
|
|
"""Self-contained dedup store persisted as ``posted.json``.
|
|
|
|
The on-disk shape is::
|
|
|
|
{"posted": ["2014/2014-08-04-foo.md", ...]}
|
|
|
|
All mutating operations acquire an exclusive :mod:`fcntl` flock on a
|
|
sibling lock file so concurrent container runs cannot corrupt the
|
|
JSON document. The store never raises on a missing file, an
|
|
unreadable file, or an already-present path — it always reports
|
|
success in a way that the orchestrator can act on.
|
|
"""
|
|
|
|
def __init__(self, data_dir: Path | None = None) -> None:
|
|
self._path = posted_path(data_dir)
|
|
|
|
@property
|
|
def path(self) -> Path:
|
|
return self._path
|
|
|
|
def load(self) -> list[str]:
|
|
"""Read the posted list, creating an empty file when missing."""
|
|
_ensure_parent(self._path)
|
|
document = _read_existing(self._path)
|
|
if document is None:
|
|
if not self._path.exists():
|
|
_write_document(self._path, dict(_EMPTY_DOCUMENT))
|
|
return []
|
|
return _coerce_list(document.get("posted"))
|
|
|
|
def is_posted(self, relative_path: str) -> bool:
|
|
"""Return ``True`` when ``relative_path`` is already recorded."""
|
|
posted = self.load()
|
|
return relative_path in posted
|
|
|
|
def _with_lock(self, mutate):
|
|
lock_path, fd = _lock_path(self._path)
|
|
try:
|
|
fcntl.flock(fd, fcntl.LOCK_EX)
|
|
return mutate()
|
|
finally:
|
|
try:
|
|
fcntl.flock(fd, fcntl.LOCK_UN)
|
|
finally:
|
|
os.close(fd)
|
|
try:
|
|
lock_path.unlink()
|
|
except FileNotFoundError:
|
|
pass
|
|
|
|
def mark_posted(self, relative_path: str) -> bool:
|
|
"""Append ``relative_path`` to the stored list. Idempotent.
|
|
|
|
Returns ``True`` when the path was newly added, ``False`` when
|
|
it was already present.
|
|
"""
|
|
return self.mark_posted_many([relative_path]) != []
|
|
|
|
def mark_posted_many(self, relative_paths: Iterable[str]) -> list[str]:
|
|
"""Append any new entries from ``relative_paths`` in one write.
|
|
|
|
Returns the list of paths that were newly added (possibly
|
|
empty when the store already contained every supplied path).
|
|
"""
|
|
|
|
candidates = [p for p in relative_paths if isinstance(p, str) and p]
|
|
if not candidates:
|
|
return []
|
|
|
|
def _mutate() -> list[str]:
|
|
document = _read_existing(self._path) or dict(_EMPTY_DOCUMENT)
|
|
existing = _coerce_list(document.get("posted"))
|
|
added = [p for p in candidates if p not in existing]
|
|
if not added:
|
|
return []
|
|
existing.extend(added)
|
|
document["posted"] = existing
|
|
_write_document(self._path, document)
|
|
return added
|
|
|
|
return self._with_lock(_mutate)
|
|
|
|
|
|
def load_posted(data_dir: Path | None = None) -> list[str]:
|
|
return PostedStore(data_dir).load()
|
|
|
|
|
|
def is_posted(relative_path: str, data_dir: Path | None = None) -> bool:
|
|
return PostedStore(data_dir).is_posted(relative_path)
|
|
|
|
|
|
def mark_posted(relative_path: str, data_dir: Path | None = None) -> bool:
|
|
return PostedStore(data_dir).mark_posted(relative_path)
|
|
|
|
|
|
def mark_posted_many(
|
|
relative_paths: list[str], data_dir: Path | None = None
|
|
) -> list[str]:
|
|
return PostedStore(data_dir).mark_posted_many(relative_paths)
|
|
|
|
|
|
__all__ = [
|
|
"DEFAULT_DATA_DIR",
|
|
"PostedStore",
|
|
"PostedStoreError",
|
|
"is_posted",
|
|
"load_posted",
|
|
"mark_posted",
|
|
"mark_posted_many",
|
|
"posted_path",
|
|
] |