高温氦气物性补全;三通四通能量计算bug修正
This commit is contained in:
1 parent
6fc9afe41d
commit
2b07d996cf
133 files changed
+5230
-253310
No files matched your search
@@ -0,0 +1,327 @@
|
||||
"""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.")
|
||||
Reference in new issue
Block a user