328 lines
14 KiB
Python
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.")
|