"""aria-atomic-write — 3-Layer atomic writes mit File-Lock + Retry.

Adopted aus Hermes `cron/jobs.py:430` + `cron/scheduler.py:1676` (siehe
`brain/02-Wissen/openclaw-hermes-code-spelunking-2026-05-13.md` §2).

Layer:
1. `tempfile.mkstemp` im selben Verzeichnis (kein cross-device move)
2. `f.flush()` + `os.fsync()` vor rename
3. `os.replace(tmp, target)` (POSIX-atomic)
4. optional `chmod` nach replace
5. `fcntl.flock(LOCK_EX | LOCK_NB)` für coordinated writes
6. Exponential-Backoff-Retry mit stale-timeout bei Lock-Contention

CLI:
    python3 aria-atomic-write.py write <path> < stdin
    python3 aria-atomic-write.py demo  # parallel-write smoke-test
"""
from __future__ import annotations

import contextlib
import fcntl
import json
import os
import random
import sys
import tempfile
import time
from pathlib import Path
from typing import Any, Iterator

DEFAULT_MODE = 0o600
DEFAULT_RETRIES = 8
DEFAULT_BASE_DELAY = 0.05  # 50 ms
DEFAULT_FACTOR = 2.0
DEFAULT_RANDOMIZE = True
DEFAULT_STALE_SECONDS = 15.0


class LockContention(RuntimeError):
    """Raised when a file-lock could not be acquired within retries."""


def atomic_write_bytes(path: Path | str, data: bytes, mode: int = DEFAULT_MODE) -> None:
    """Atomar Bytes nach `path` schreiben. mkstemp + fsync + os.replace + chmod."""
    target = Path(path)
    target.parent.mkdir(parents=True, exist_ok=True)
    fd, tmp_str = tempfile.mkstemp(
        dir=str(target.parent),
        prefix=f".{target.name}.",
        suffix=".tmp",
    )
    tmp_path = Path(tmp_str)
    try:
        with os.fdopen(fd, "wb") as f:
            f.write(data)
            f.flush()
            os.fsync(f.fileno())
        os.chmod(tmp_path, mode)
        os.replace(tmp_path, target)
    except Exception:
        with contextlib.suppress(FileNotFoundError):
            tmp_path.unlink()
        raise


def atomic_write_text(
    path: Path | str,
    content: str,
    *,
    encoding: str = "utf-8",
    mode: int = DEFAULT_MODE,
) -> None:
    """Atomar Text nach `path` schreiben."""
    atomic_write_bytes(path, content.encode(encoding), mode=mode)


def atomic_write_json(
    path: Path | str,
    data: Any,
    *,
    indent: int = 2,
    ensure_ascii: bool = False,
    mode: int = DEFAULT_MODE,
) -> None:
    """Atomar JSON nach `path` schreiben."""
    payload = json.dumps(data, indent=indent, ensure_ascii=ensure_ascii)
    atomic_write_text(path, payload, mode=mode)


@contextlib.contextmanager
def file_lock(
    lock_path: Path | str,
    *,
    exclusive: bool = True,
    retries: int = DEFAULT_RETRIES,
    base_delay: float = DEFAULT_BASE_DELAY,
    factor: float = DEFAULT_FACTOR,
    randomize: bool = DEFAULT_RANDOMIZE,
    stale_seconds: float = DEFAULT_STALE_SECONDS,
) -> Iterator[None]:
    """File-Lock mit Exponential-Backoff. Stale-Lock-Recovery: wenn Lock-Datei
    älter als `stale_seconds` ist, wird sie ignoriert (Crashed-Holder-Recovery).
    """
    lp = Path(lock_path)
    lp.parent.mkdir(parents=True, exist_ok=True)
    flag = fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH
    flag |= fcntl.LOCK_NB
    delay = base_delay
    last_err: BaseException | None = None
    for attempt in range(retries + 1):
        # Stale-Recovery: alte Lock-Files entfernen (best-effort)
        if lp.exists():
            try:
                age = time.time() - lp.stat().st_mtime
                if age > stale_seconds:
                    with contextlib.suppress(FileNotFoundError):
                        lp.unlink()
            except OSError:
                pass
        try:
            fd = os.open(str(lp), os.O_CREAT | os.O_WRONLY, 0o600)
        except OSError as exc:
            last_err = exc
            time.sleep(delay)
            delay *= factor
            if randomize:
                delay *= 0.5 + random.random()
            continue
        try:
            fcntl.flock(fd, flag)
        except OSError as exc:
            os.close(fd)
            last_err = exc
            if attempt == retries:
                raise LockContention(
                    f"lock not acquired on {lp} after {retries + 1} attempts: {exc}"
                ) from exc
            time.sleep(delay)
            delay *= factor
            if randomize:
                delay *= 0.5 + random.random()
            continue
        try:
            os.write(fd, f"{os.getpid()}@{time.time()}\n".encode())
            os.fsync(fd)
            yield
        finally:
            try:
                fcntl.flock(fd, fcntl.LOCK_UN)
            finally:
                os.close(fd)
                with contextlib.suppress(FileNotFoundError):
                    lp.unlink()
        return
    raise LockContention(f"lock not acquired on {lp}: {last_err}")


# ---------------------------------------------------------------- CLI / demo
def _cli_write(target: str) -> int:
    content = sys.stdin.read()
    atomic_write_text(target, content)
    print(f"wrote {target} ({len(content)} chars)", file=sys.stderr)
    return 0


def _cli_demo() -> int:
    import multiprocessing as mp

    target = Path("/tmp/aria-atomic-demo/out.txt")
    lock = Path("/tmp/aria-atomic-demo/.out.lock")
    target.parent.mkdir(parents=True, exist_ok=True)
    with contextlib.suppress(FileNotFoundError):
        target.unlink()

    def worker(tag: str) -> None:
        with file_lock(lock):
            existing = target.read_text() if target.exists() else ""
            time.sleep(0.05)
            atomic_write_text(target, existing + f"{tag}-{os.getpid()}\n")

    procs = [mp.Process(target=worker, args=(f"P{i}",)) for i in range(5)]
    for p in procs:
        p.start()
    for p in procs:
        p.join()
    print("--- demo output ---")
    print(target.read_text())
    print(f"line count: {len(target.read_text().splitlines())}")
    return 0


def main() -> int:
    if len(sys.argv) < 2:
        print(__doc__)
        return 2
    cmd = sys.argv[1]
    if cmd == "write":
        if len(sys.argv) < 3:
            print("usage: aria-atomic-write.py write <path>", file=sys.stderr)
            return 2
        return _cli_write(sys.argv[2])
    if cmd == "demo":
        return _cli_demo()
    print(f"unknown command: {cmd}", file=sys.stderr)
    return 2


if __name__ == "__main__":
    sys.exit(main())
