"""Locked, atomic persistence for the local JSON VM inventory.""" from __future__ import annotations import fcntl import json import os import tempfile from collections.abc import Iterator from contextlib import contextmanager from pathlib import Path from .config import Settings from .domain import State, VmRecord from .errors import StateError, UvmError from .system import ensure_data_directories class StateStore: """Own the JSON state file so callers cannot race IP and VM allocation.""" def __init__(self, settings: Settings) -> None: self._settings = settings def initialize(self) -> None: ensure_data_directories(self._settings) with self._locked(): if not self._settings.state_path.exists(): self._save_unlocked(State()) else: try: os.chmod(self._settings.state_path, 0o600) except OSError as error: raise StateError( f"cannot update permissions on {self._settings.state_path}: {error}" ) from error def load(self) -> State: # Reads do not create /var/lib/uvm, so `uvm list` remains usable before install. # Atomic replacement makes an unlocked reader see either the old or new full document. return self._load_unlocked() @contextmanager def transaction(self) -> Iterator[State]: """Load and commit one state mutation while holding an exclusive lock.""" ensure_data_directories(self._settings) with self._locked(): state = self._load_unlocked() yield state self._save_unlocked(state) @contextmanager def operation_lock(self) -> Iterator[None]: """Serialize external VM lifecycle work across concurrent CLI processes.""" ensure_data_directories(self._settings) with self._locked_path(self._settings.operation_lock_path): yield def get(self, vm_id: str) -> VmRecord: state = self.load() try: return state.vms[vm_id] except KeyError as error: raise UvmError(f"VM not found: {vm_id}") from error def find_by_id_or_ip(self, identifier: str) -> VmRecord: state = self.load() if identifier in state.vms: return state.vms[identifier] for vm in state.vms.values(): if vm.guest_ip == identifier: return vm raise UvmError(f"VM not found: {identifier}") @contextmanager def _locked(self) -> Iterator[None]: with self._locked_path(self._settings.state_lock_path): yield @contextmanager def _locked_path(self, lock_path: Path) -> Iterator[None]: lock_path.parent.mkdir(parents=True, exist_ok=True) with lock_path.open("a+") as lock_file: fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX) try: yield finally: fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN) def _load_unlocked(self) -> State: path = self._settings.state_path if not path.exists(): return State() try: if path.stat().st_mode & 0o077: os.chmod(path, 0o600) with path.open(encoding="utf-8") as state_file: raw = json.load(state_file) except (OSError, json.JSONDecodeError) as error: raise StateError(f"cannot read {path}: {error}") from error return State.from_dict(raw) def _save_unlocked(self, state: State) -> None: path = self._settings.state_path path.parent.mkdir(parents=True, exist_ok=True) fd, temporary_path = tempfile.mkstemp( prefix=f".{path.name}.", suffix=".tmp", dir=path.parent, ) try: with os.fdopen(fd, "w", encoding="utf-8") as temporary_file: json.dump(state.to_dict(), temporary_file, indent=2, sort_keys=True) temporary_file.write("\n") temporary_file.flush() os.fsync(temporary_file.fileno()) os.replace(temporary_path, path) os.chmod(path, 0o600) self._fsync_parent(path) except OSError as error: raise StateError(f"cannot write {path}: {error}") from error finally: try: os.unlink(temporary_path) except FileNotFoundError: pass @staticmethod def _fsync_parent(path: Path) -> None: directory = os.open(str(path.parent), os.O_RDONLY) try: os.fsync(directory) finally: os.close(directory)