377 lines
12 KiB
Python
377 lines
12 KiB
Python
"""Process-shared reservation store for NemoTron dispatch side effects."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import time
|
|
import uuid
|
|
from pathlib import Path
|
|
from typing import Callable
|
|
|
|
try:
|
|
import fcntl
|
|
except ImportError: # pragma: no cover - exercised by compatibility monkeypatch
|
|
fcntl = None
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
STATE_VERSION = 1
|
|
PRODUCTION_STATE_PATH = Path(
|
|
"/app/data/ai_automation/nemotron_dispatch_dedupe_state.json"
|
|
)
|
|
|
|
|
|
def production_shared_reservation_required() -> bool:
|
|
return (
|
|
os.getenv("FLASK_ENV", "").strip().lower() in {"prod", "production"}
|
|
or os.getenv("USE_POSTGRESQL", "").strip().lower() == "true"
|
|
)
|
|
|
|
|
|
def _sku_key(sku: str) -> str:
|
|
value = str(sku or "").strip()
|
|
if not value:
|
|
return ""
|
|
return hashlib.sha256(f"nemotron-dispatch:{value}".encode("utf-8")).hexdigest()
|
|
|
|
|
|
class ReservationStateError(RuntimeError):
|
|
pass
|
|
|
|
|
|
class SharedFileReservationStore:
|
|
"""Atomic lease/commit state shared by processes using the same bind mount."""
|
|
|
|
def __init__(
|
|
self,
|
|
state_path: str | Path,
|
|
*,
|
|
committed_ttl_sec: int,
|
|
lease_ttl_sec: int,
|
|
clock: Callable[[], float] | None = None,
|
|
):
|
|
self.state_path = Path(state_path)
|
|
self.lock_path = self.state_path.with_suffix(self.state_path.suffix + ".lock")
|
|
self.committed_ttl_sec = max(1, int(committed_ttl_sec))
|
|
self.lease_ttl_sec = max(1, int(lease_ttl_sec))
|
|
self.clock = clock or time.time
|
|
|
|
@staticmethod
|
|
def _empty_state() -> dict:
|
|
return {"version": STATE_VERSION, "leases": {}, "committed": {}}
|
|
|
|
def _read_state(self) -> dict:
|
|
if not self.state_path.exists():
|
|
return self._empty_state()
|
|
try:
|
|
payload = json.loads(self.state_path.read_text(encoding="utf-8"))
|
|
except (OSError, json.JSONDecodeError) as exc:
|
|
raise ReservationStateError("reservation_state_unreadable") from exc
|
|
if not isinstance(payload, dict) or payload.get("version") != STATE_VERSION:
|
|
raise ReservationStateError("reservation_state_schema_invalid")
|
|
leases = payload.get("leases")
|
|
committed = payload.get("committed")
|
|
if not isinstance(leases, dict) or not isinstance(committed, dict):
|
|
raise ReservationStateError("reservation_state_schema_invalid")
|
|
for lease in leases.values():
|
|
if (
|
|
not isinstance(lease, dict)
|
|
or not isinstance(lease.get("token"), str)
|
|
or not isinstance(lease.get("expires_at"), (int, float))
|
|
or lease.get("phase", "reserved") not in {
|
|
"reserved",
|
|
"side_effect_started",
|
|
}
|
|
):
|
|
raise ReservationStateError("reservation_lease_schema_invalid")
|
|
for marker in committed.values():
|
|
if (
|
|
not isinstance(marker, dict)
|
|
or not isinstance(marker.get("until"), (int, float))
|
|
):
|
|
raise ReservationStateError("reservation_commit_schema_invalid")
|
|
return payload
|
|
|
|
def _write_state(self, payload: dict) -> None:
|
|
temporary = self.state_path.with_name(
|
|
f".{self.state_path.name}.{os.getpid()}.{uuid.uuid4().hex}.tmp"
|
|
)
|
|
data = json.dumps(payload, sort_keys=True, separators=(",", ":"))
|
|
try:
|
|
with temporary.open("w", encoding="utf-8") as handle:
|
|
handle.write(data)
|
|
handle.flush()
|
|
os.fsync(handle.fileno())
|
|
os.chmod(temporary, 0o640)
|
|
os.replace(temporary, self.state_path)
|
|
directory_fd = os.open(str(self.state_path.parent), os.O_RDONLY)
|
|
try:
|
|
os.fsync(directory_fd)
|
|
finally:
|
|
os.close(directory_fd)
|
|
finally:
|
|
try:
|
|
temporary.unlink(missing_ok=True)
|
|
except OSError:
|
|
pass
|
|
|
|
@staticmethod
|
|
def _cleanup(payload: dict, now: float) -> bool:
|
|
changed = False
|
|
for key in [
|
|
key
|
|
for key, lease in payload["leases"].items()
|
|
if float(lease["expires_at"]) <= now
|
|
]:
|
|
del payload["leases"][key]
|
|
changed = True
|
|
for key in [
|
|
key
|
|
for key, marker in payload["committed"].items()
|
|
if float(marker["until"]) <= now
|
|
]:
|
|
del payload["committed"][key]
|
|
changed = True
|
|
return changed
|
|
|
|
def _locked(self, operation):
|
|
if fcntl is None:
|
|
raise ReservationStateError("process_shared_lock_unavailable")
|
|
self.state_path.parent.mkdir(mode=0o750, parents=True, exist_ok=True)
|
|
with self.lock_path.open("a+", encoding="utf-8") as lock_handle:
|
|
os.chmod(self.lock_path, 0o640)
|
|
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_EX)
|
|
try:
|
|
payload = self._read_state()
|
|
now = float(self.clock())
|
|
cleanup_changed = self._cleanup(payload, now)
|
|
result, operation_changed = operation(payload, now)
|
|
if cleanup_changed or operation_changed:
|
|
self._write_state(payload)
|
|
return result
|
|
finally:
|
|
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_UN)
|
|
|
|
def reserve(self, sku: str) -> str | None:
|
|
key = _sku_key(sku)
|
|
if not key:
|
|
return None
|
|
|
|
def operation(payload, now):
|
|
if key in payload["leases"] or key in payload["committed"]:
|
|
return None, False
|
|
token = uuid.uuid4().hex
|
|
payload["leases"][key] = {
|
|
"token": token,
|
|
"expires_at": now + self.lease_ttl_sec,
|
|
"phase": "reserved",
|
|
}
|
|
return token, True
|
|
|
|
try:
|
|
return self._locked(operation)
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared reserve failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return None
|
|
|
|
def refresh(self, sku: str, token: str) -> bool:
|
|
key = _sku_key(sku)
|
|
|
|
def operation(payload, now):
|
|
lease = payload["leases"].get(key)
|
|
if not lease or lease.get("token") != str(token):
|
|
return False, False
|
|
ttl = (
|
|
self.committed_ttl_sec
|
|
if lease.get("phase") == "side_effect_started"
|
|
else self.lease_ttl_sec
|
|
)
|
|
lease["expires_at"] = now + ttl
|
|
return True, True
|
|
|
|
try:
|
|
return bool(key and self._locked(operation))
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared refresh failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return False
|
|
|
|
def begin_side_effect(self, sku: str, token: str) -> bool:
|
|
"""Persist the at-most-once boundary before invoking an external effect."""
|
|
key = _sku_key(sku)
|
|
|
|
def operation(payload, now):
|
|
lease = payload["leases"].get(key)
|
|
if not lease or lease.get("token") != str(token):
|
|
return False, False
|
|
lease["phase"] = "side_effect_started"
|
|
lease["expires_at"] = now + self.committed_ttl_sec
|
|
return True, True
|
|
|
|
try:
|
|
return bool(key and self._locked(operation))
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared side-effect boundary failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return False
|
|
|
|
def commit(self, sku: str, token: str) -> bool:
|
|
key = _sku_key(sku)
|
|
|
|
def operation(payload, now):
|
|
lease = payload["leases"].get(key)
|
|
if not lease or lease.get("token") != str(token):
|
|
return False, False
|
|
del payload["leases"][key]
|
|
payload["committed"][key] = {
|
|
"until": now + self.committed_ttl_sec,
|
|
}
|
|
return True, True
|
|
|
|
try:
|
|
return bool(key and self._locked(operation))
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared commit failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return False
|
|
|
|
def release(self, sku: str, token: str | None) -> bool:
|
|
key = _sku_key(sku)
|
|
|
|
def operation(payload, _now):
|
|
lease = payload["leases"].get(key)
|
|
if not lease or lease.get("token") != str(token or ""):
|
|
return False, False
|
|
del payload["leases"][key]
|
|
return True, True
|
|
|
|
try:
|
|
return bool(key and self._locked(operation))
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared release failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return False
|
|
|
|
def is_duplicate(self, sku: str) -> bool:
|
|
key = _sku_key(sku)
|
|
if not key:
|
|
return True
|
|
|
|
def operation(payload, _now):
|
|
return (
|
|
key in payload["leases"] or key in payload["committed"],
|
|
False,
|
|
)
|
|
|
|
try:
|
|
return bool(self._locked(operation))
|
|
except Exception as exc:
|
|
logger.error(
|
|
"[NemotronReservation] shared read failed closed: %s",
|
|
type(exc).__name__,
|
|
)
|
|
return True
|
|
|
|
|
|
def build_production_shared_reservation_store(
|
|
*,
|
|
committed_ttl_sec: int,
|
|
lease_ttl_sec: int,
|
|
) -> SharedFileReservationStore | None:
|
|
if not production_shared_reservation_required():
|
|
return None
|
|
return SharedFileReservationStore(
|
|
PRODUCTION_STATE_PATH,
|
|
committed_ttl_sec=committed_ttl_sec,
|
|
lease_ttl_sec=lease_ttl_sec,
|
|
)
|
|
|
|
|
|
def run_shared_reservation_canary(
|
|
*,
|
|
state_path: str | Path = PRODUCTION_STATE_PATH,
|
|
committed_ttl_sec: int = 4 * 3600,
|
|
lease_ttl_sec: int = 15 * 60,
|
|
) -> dict:
|
|
"""Exercise shared ownership semantics without business side effects."""
|
|
first = SharedFileReservationStore(
|
|
state_path,
|
|
committed_ttl_sec=committed_ttl_sec,
|
|
lease_ttl_sec=lease_ttl_sec,
|
|
)
|
|
second = SharedFileReservationStore(
|
|
state_path,
|
|
committed_ttl_sec=committed_ttl_sec,
|
|
lease_ttl_sec=lease_ttl_sec,
|
|
)
|
|
synthetic_sku = f"CANARY-NEMOTRON-RESERVATION-{uuid.uuid4().hex}"
|
|
token = first.reserve(synthetic_sku)
|
|
competing_token = second.reserve(synthetic_sku) if token else None
|
|
stale_release_rejected = (
|
|
second.release(synthetic_sku, "stale-canary-token") is False
|
|
if token
|
|
else False
|
|
)
|
|
side_effect_boundary_persisted = (
|
|
first.begin_side_effect(synthetic_sku, token) if token else False
|
|
)
|
|
owner_release_succeeded = (
|
|
first.release(synthetic_sku, token) if token else False
|
|
)
|
|
terminal_clear = (
|
|
first.is_duplicate(synthetic_sku) is False
|
|
if owner_release_succeeded
|
|
else False
|
|
)
|
|
checks = {
|
|
"first_owner_reserved": bool(token),
|
|
"competing_owner_rejected": competing_token is None,
|
|
"stale_token_release_rejected": stale_release_rejected,
|
|
"side_effect_boundary_persisted": side_effect_boundary_persisted,
|
|
"owner_release_succeeded": owner_release_succeeded,
|
|
"terminal_state_clear": terminal_clear,
|
|
}
|
|
success = all(checks.values())
|
|
return {
|
|
"success": success,
|
|
"status": "shared_reservation_canary_passed" if success else "shared_reservation_canary_failed",
|
|
"policy": "nemotron_process_shared_reservation_v1",
|
|
"state_path": str(Path(state_path)),
|
|
"checks": checks,
|
|
"check_count": len(checks),
|
|
"check_pass_count": sum(checks.values()),
|
|
"controlled_apply": {
|
|
"writes_shared_reservation_state": True,
|
|
"terminal_state_clear": terminal_clear,
|
|
"writes_database": False,
|
|
"sends_telegram": False,
|
|
"executes_business_tool": False,
|
|
"reads_secret": False,
|
|
},
|
|
}
|
|
|
|
|
|
__all__ = [
|
|
"PRODUCTION_STATE_PATH",
|
|
"ReservationStateError",
|
|
"SharedFileReservationStore",
|
|
"build_production_shared_reservation_store",
|
|
"production_shared_reservation_required",
|
|
"run_shared_reservation_canary",
|
|
]
|