Files

328 lines
14 KiB
Python

"""Quota-controlled, append-only result archives in PROJECT/simresults.
Reservations cover state AND future output blocks before the worker may write.
A global process lock serializes grants; per-run leases protect active writers.
Grants are persisted in bounded windows; individual blocks consume local credit
without scanning history or synchronously rewriting the manifest.
Only recognized archives are evicted, oldest first. Unknown user files count
toward the quota but are never deleted. File sizes are logical bytes, not the
filesystem's allocation-unit overhead.
"""
from __future__ import annotations
from contextlib import contextmanager
import ctypes
import json
import os
from pathlib import Path
import re
import stat
import struct
import threading
import time
from uuid import uuid4
import zlib
from .cache_storage import acquire_cache_lease, _acquire_lock_file, _reparse, _require_directory
RESULT_ROOT = Path(__file__).resolve().parents[3] / "simresults"
DEFAULT_LIMIT_BYTES = 1024**3
_RESERVATION_WINDOW_BYTES = 16 * 1024**2
_RUN = re.compile(r"run-([0-9a-f]{64})\Z")
_FILES = frozenset(("manifest.json", "manifest.tmp", "states.bin", "outputs.bin"))
_HEADER = struct.Struct("<8sQQQQ")
class ResultQuotaError(RuntimeError):
pass
def result_archive_path(result_id: str, *, root: Path | None = None) -> Path:
"""Locate an archive by ID under this installation's result root.
Never use a saved absolute directory as a locator: older manifests may
contain one from a different machine. Moving a stopped project and starting
it again establishes its new root through __file__ above.
"""
if not isinstance(result_id, str) or _RUN.fullmatch(result_id) is None:
raise ValueError("Invalid result archive ID.")
base = Path(root or RESULT_ROOT).absolute()
_require_directory(base)
path = base / result_id
_require_directory(path)
return path
def result_limit_bytes() -> int:
raw = os.environ.get("SIMULATION_RESULT_STORAGE_MB", str(DEFAULT_LIMIT_BYTES // 1024**2))
if not re.fullmatch(r"[0-9]+", raw) or int(raw) < 1:
raise ValueError("SIMULATION_RESULT_STORAGE_MB must be a positive integer in MiB.")
return int(raw) * 1024**2
def _alive(pid: int | None) -> bool:
if not pid:
return False
if os.name == "nt":
kernel = ctypes.WinDLL("kernel32", use_last_error=True)
kernel.OpenProcess.argtypes = [ctypes.c_ulong, ctypes.c_int, ctypes.c_ulong]
kernel.OpenProcess.restype = ctypes.c_void_p
kernel.WaitForSingleObject.argtypes = [ctypes.c_void_p, ctypes.c_ulong]
kernel.CloseHandle.argtypes = [ctypes.c_void_p]
handle = kernel.OpenProcess(0x100000, False, pid) # SYNCHRONIZE
if not handle:
return ctypes.get_last_error() == 5 # Access denied: conservatively pin.
try:
return kernel.WaitForSingleObject(handle, 0) != 0
finally:
kernel.CloseHandle(handle)
try:
os.kill(pid, 0)
return True
except ProcessLookupError:
return False
except PermissionError:
return True
def _size(path: Path) -> int:
info = path.lstat()
if _reparse(info) or not stat.S_ISDIR(info.st_mode):
return info.st_size
return sum(_size(child) for child in path.iterdir())
def _record(path: Path) -> dict | None:
try:
info = (path / "manifest.json").lstat()
if _reparse(info) or not stat.S_ISREG(info.st_mode):
return None
data = json.loads((path / "manifest.json").read_bytes())
if (data.get("format") != "simulation-blocks-v1" or
data.get("id") != path.name or type(data.get("reservedBytes")) is not int or
data["reservedBytes"] < 0 or type(data.get("createdAt")) is not int or
type(data.get("finishedAt", data["createdAt"])) is not int or
data.get("status") not in ("running", "completed", "cancelled", "failed", "interrupted") or
(data.get("workerPid") is not None and
(type(data["workerPid"]) is not int or data["workerPid"] <= 0))):
return None
return data
except (OSError, ValueError, AttributeError):
return None
def scan_blocks(path: Path, *, expected_columns: int | None = None) -> dict:
"""Validate a committed prefix using bounded reads; ignore an incomplete tail.
A CRC mismatch is reported as corruption, never a valid sample. This reader
works after process termination without loading the trajectory into RAM.
"""
result = {"blocks": 0, "samples": 0, "validBytes": 0, "storedUntil": None,
"incompleteTail": False, "corrupt": False}
if not path.exists():
return result
with path.open("rb") as stream:
size = os.fstat(stream.fileno()).st_size
while stream.tell() < size:
raw = stream.read(_HEADER.size)
if len(raw) != _HEADER.size:
result["incompleteTail"] = True
break
magic, sequence, rows, columns, expected_crc = _HEADER.unpack(raw)
if (magic != b"SIMBLK01" or sequence != result["blocks"] or not 0 < rows <= 1024 or
not columns or (expected_columns is not None and columns != expected_columns)):
result["corrupt"] = True
break
length = rows * columns * 8
if length + 8 > size - stream.tell():
result["incompleteTail"] = True
break
# The first column is always time; at most 1024 timestamps.
times = stream.read(rows * 8)
crc = zlib.crc32(times)
remaining = length - len(times)
while remaining:
block = stream.read(min(65536, remaining))
if not block:
raise OSError("Result block changed during inspection.")
crc = zlib.crc32(block, crc)
remaining -= len(block)
if stream.read(8) != b"COMMIT01" or crc != expected_crc:
result["corrupt"] = True
break
result["blocks"] += 1
result["samples"] += rows
result["validBytes"] = stream.tell()
result["storedUntil"] = struct.unpack_from("<d", times, len(times) - 8)[0]
return result
class ResultArchive:
def __init__(self, metadata: dict, *, root: Path | None = None, limit_bytes: int | None = None):
self.root = Path(root or RESULT_ROOT).absolute()
_require_directory(self.root, create=True)
self.limit = result_limit_bytes() if limit_bytes is None else limit_bytes
if self.limit <= 0:
raise ValueError("Result quota must be positive.")
self.key = uuid4().hex + uuid4().hex
self.path = self.root / ("run-" + self.key)
self.lease = acquire_cache_lease(self.root, "models", self.key)
self.closed = False
self._mutex = threading.RLock()
self._credit = 0
# This is accounting credit, not an allocated memory/disk buffer. Bound
# speculative grants for small quotas and concurrent simulations.
self._reservation_window = min(_RESERVATION_WINDOW_BYTES, self.limit // 16)
self.data = {"format": "simulation-blocks-v1", "id": self.path.name,
"createdAt": time.time_ns(), "status": "running", "workerPid": None,
"metadata": metadata, "reservedBytes": 0}
# Covers final scalars, JSON framing and both copies during atomic update.
self.metadata_budget = max(65536, len(json.dumps(self.data).encode()) * 8)
self.data["reservedBytes"] = self.metadata_budget
try:
with self._locked():
self._make_room(self.metadata_budget)
self.path.mkdir()
self._publish()
except BaseException:
self.lease.close()
raise
@contextmanager
def _locked(self):
with _acquire_lock_file(self.root, "models", "0" * 64, "results-quota.lock",
exclusive=True, blocking=True):
yield
def _publish(self):
encoded = json.dumps(self.data, ensure_ascii=False, allow_nan=False,
separators=(",", ":")).encode("utf-8")
if len(encoded) * 2 > self.metadata_budget:
raise ResultQuotaError("Result metadata exceeded its reserved space.")
temporary = self.path / "manifest.tmp"
with temporary.open("wb") as stream:
stream.write(encoded)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, self.path / "manifest.json")
def _make_room(self, extra: int) -> int:
"""Make room for required bytes and return free bytes BEFORE this grant.
Caller holds the global quota lock through publishing the reservation.
Optional ahead-of-write credit must never trigger extra eviction.
"""
usage = 0
candidates = []
for path in self.root.iterdir():
info = path.lstat()
actual = _size(path)
match = _RUN.fullmatch(path.name)
record = (_record(path) if match and stat.S_ISDIR(info.st_mode) and not _reparse(info) else None)
if record is None:
usage += actual
continue
lease = acquire_cache_lease(self.root, "models", match[1], exclusive=True, blocking=False)
active = lease is None or (record.get("status") == "running" and _alive(record.get("workerPid")))
if lease is not None:
lease.close()
usage += max(actual, record["reservedBytes"]) if active else actual
if not active and path != self.path:
candidates.append((record.get("finishedAt", record["createdAt"]), path, match[1], actual))
if usage + extra <= self.limit:
return self.limit - usage
for _, path, key, actual in sorted(candidates):
lease = acquire_cache_lease(self.root, "models", key, exclusive=True, blocking=False)
if lease is None:
continue
with lease:
# No recursive deletion, links, unknown files, or paths outside
# the configured root. User-owned additions prevent eviction.
if path.resolve().parent != self.root.resolve():
continue
files = list(path.iterdir())
if any(p.name not in _FILES or _reparse(p.lstat()) or
not stat.S_ISREG(p.lstat().st_mode) for p in files):
continue
for item in files:
item.unlink()
path.rmdir()
usage -= actual
if usage + extra <= self.limit:
return self.limit - usage
raise ResultQuotaError("Result storage quota exhausted; active runs and unrelated files were retained.")
def worker_started(self, pid: int):
with self._mutex, self._locked():
self._ensure_open()
self.data["workerPid"] = pid
self._publish()
def _ensure_open(self):
if self.closed:
raise OSError("Result archive is closed.")
def reserve(self, amount: int):
if type(amount) is not int or not 0 < amount <= self.limit:
raise ResultQuotaError("Invalid or oversized result block reservation.")
with self._mutex:
self._ensure_open()
if amount <= self._credit:
self._credit -= amount
return
needed = amount - self._credit
with self._locked():
available = self._make_room(needed)
grant = min(available, max(needed, self._reservation_window))
previous = self.data["reservedBytes"]
self.data["reservedBytes"] += grant
try:
# Publish BEFORE granting permission to write, so another
# process (or an orphan worker) cannot reuse our quota.
self._publish()
except BaseException:
self.data["reservedBytes"] = previous
raise
self._credit += grant - amount
def finish(self, payload: dict | None = None, error: str | None = None) -> dict:
with self._mutex:
self._ensure_open()
return self._finish(payload, error)
def _finish(self, payload: dict | None, error: str | None) -> dict:
try:
metadata = self.data["metadata"]
states = scan_blocks(self.path / "states.bin",
expected_columns=len(metadata["stateColumns"]) if "stateColumns" in metadata
else max(1, len(metadata["stateKeys"])) + 1 if "stateKeys" in metadata else None)
outputs = scan_blocks(self.path / "outputs.bin",
expected_columns=len(metadata["outputColumns"]) if "outputColumns" in metadata
else max(1, len(metadata["variables"])) + 1 if "variables" in metadata else None)
status = payload.get("status", "failed") if payload is not None else "interrupted"
if states["corrupt"] or outputs["corrupt"]:
status = "failed"
summary = {"id": self.path.name, "directory": f"simresults/{self.path.name}",
"directoryBase": "project", "limitBytes": self.limit,
"status": status, "states": states, "outputs": outputs}
with self._locked():
# Release unused ahead-of-write credit even if a result reader
# subsequently pins this completed archive with a shared lease.
self.data["reservedBytes"] -= self._credit
self._credit = 0
self.data.update(status=status, finishedAt=time.time_ns(), blocks=summary,
error=error, workerPid=None)
if payload is not None:
self.data["result"] = {k: v for k, v in payload.items()
if k not in ("series", "buildDetails", "resultStorage")}
self._publish()
return summary
finally:
self.closed = True
self.lease.close()
def close(self):
with self._mutex:
if not self.closed:
self.finish(error="Execution ended before a result was returned.")