feat(automation): poisoned-run quarantine policy with durable failure signatures (SCHED-SOAK offline)

- poisoned_store.py: durable JSONL identical-failure counter (fcntl.flock + intra-process lock
  + thread-local reentrancy, append+fsync, torn-line tolerant, CAS quarantine epoch); no
  Alembic/ORM changes (046 migrations frozen).
- poisoned.py: N=3 identical infra failures -> typed POISONED_RUN_QUARANTINE, matching
  schedules disabled, exactly one blocked notification carrying the DLQ payload; success
  resets the pair; operator unquarantine (CAS) re-enables schedules and resets counters.
- REST: GET /scenario-automation/quarantine and POST /quarantine/{scenario_id}/release
  (scenario:automation MANAGE; 404 NOT_QUARANTINED / 409 STALE_QUARANTINE_VERSION).
- Tests: 13 policy + 2 API quarantine vectors; merged with the master ReadAcl/RetentionReceipts
  classes (conflict resolved keeping both). Focused 31 passed; 046 T023 offline slice.
This commit is contained in:
2026-09-22 12:35:39 +03:00
parent 93e97228e9
commit def2bf4038
5 changed files with 858 additions and 0 deletions

View File

@@ -46,6 +46,8 @@ from src.models.scenario_automation import AutomationPolicy, ScenarioNotificatio
from src.models.scenario_registry import ScenarioRegistryEntry
from src.services.dashboard_testing.automation.metrics import automation_metrics
from src.services.dashboard_testing.automation.eligibility import assert_automation_eligible
from src.services.dashboard_testing.automation.poisoned import unquarantine_scenario_automation
from src.services.dashboard_testing.automation.poisoned_store import get_poisoned_store
from src.services.dashboard_testing.automation.retention import tier_limits
from src.services.dashboard_testing.automation.trigger import dispatch_trigger_event
from src.services.dashboard_testing.execution.environment_policy import (
@@ -166,6 +168,14 @@ class PolicyRequest(BaseModel):
# #endregion Api.ScenarioAutomation.PolicyRequest
# #region Api.ScenarioAutomation.QuarantineReleaseRequest [C:1] [TYPE Class] [SEMANTICS scenario,automation,api,quarantine,recovery]
# @ingroup Api
class QuarantineReleaseRequest(BaseModel):
environment_id: str
expected_version: int = Field(..., ge=1)
# #endregion Api.ScenarioAutomation.QuarantineReleaseRequest
# #region Api.ScenarioAutomation.ApiTriggerRequest [C:1] [TYPE Class] [SEMANTICS scenario,automation,api,trigger,run]
# @ingroup Api
class ApiTriggerRequest(BaseModel):
@@ -394,6 +404,50 @@ def list_notifications(limit: int = 100, db=_DB, current_user=_READ):
# #endregion Api.ScenarioAutomation.Notifications
# #region Api.ScenarioAutomation.Quarantine [C:4] [TYPE Block] [SEMANTICS scenario,automation,api,quarantine,dlq,recovery,operator]
# @ingroup Api
# @BRIEF Operator-only read and CAS release of poisoned-run quarantines recorded by the durable failure counter.
# @RELATION CALLS -> [ScenarioAutomation.Poisoned.Unquarantine]
# @RELATION DEPENDS_ON -> [ScenarioAutomation.PoisonedStore]
# @POST A released quarantine re-enables matching schedules and resets the identical-failure counter.
# @INVARIANT Recovery requires scenario:automation MANAGE and a matching quarantine version; an unknown or stale
# release changes no schedule and no counter.
@router.get("/quarantine")
def list_quarantines(_user=_USER):
return get_poisoned_store().list_quarantines()
@router.post("/quarantine/{scenario_id}/release")
def release_quarantine(
scenario_id: str,
body: QuarantineReleaseRequest,
db=_DB,
_perm=_MANAGE,
current_user=_USER,
):
store = get_poisoned_store()
actor = str(getattr(current_user, "id", "") or getattr(current_user, "username", "") or "")
try:
result = unquarantine_scenario_automation(
db,
store,
scenario_id=scenario_id,
environment_id=body.environment_id,
expected_version=body.expected_version,
actor=actor,
)
db.commit()
except ValueError as exc:
db.rollback()
code = str(exc)
raise HTTPException(
status_code=404 if code == "NOT_QUARANTINED" else 409,
detail={"code": code},
) from exc
return result
# #endregion Api.ScenarioAutomation.Quarantine
# #region Api.ScenarioAutomation.EventDispatch [C:3] [TYPE Function] [SEMANTICS scenario,automation,api,event,trigger,persistence]
# @ingroup Api
# @BRIEF Persist policy-allowed event-triggered runs without request-time execution.

View File

@@ -0,0 +1,136 @@
# #region ScenarioAutomation.Poisoned [C:4] [TYPE Module] [SEMANTICS scenario,automation,poisoned,quarantine,dlq,notification,recovery]
# @defgroup ScenarioAutomation Poisoned-run policy: persisted quarantine, DLQ notification and operator recovery.
# @LAYER Service
# @BRIEF Turn an identical-failure decision from the durable counter into a scenario-automation quarantine:
# disable matching schedules, emit one typed blocked notification, and let an operator CAS-release it.
# @RELATION DEPENDS_ON -> [ScenarioAutomation.PoisonedStore]
# @RELATION DEPENDS_ON -> [ScenarioAutomation.Notify]
# @RELATION DEPENDS_ON -> [Models.ScenarioAutomation]
# @INVARIANT Quarantine disables every matching schedule and emits exactly one blocked notification per epoch.
# @INVARIANT Recovery is an explicit operator CAS on the quarantine version; no automatic retry clears quarantine.
# @RATIONALE The counter is file-durable but the automation lifecycle (schedule enabled flag, notification receipt)
# is ORM state, so the policy composes both under one caller-owned transaction.
# @REJECTED Auto-unquarantine after a later success was rejected — once poisoned, a human must confirm the fix;
# otherwise a flapping environment would silently re-arm the scheduler forever.
from __future__ import annotations
from typing import Any
from sqlalchemy.orm import Session
from src.core.logger import logger
from src.models.scenario_automation import ScenarioSchedule
from .notify import persist_notification
from .poisoned_store import POISONED_RUN_QUARANTINE, QUARANTINE_THRESHOLD, PoisonedRunStore
# #region ScenarioAutomation.Poisoned.Schedules [C:2] [TYPE Function] [SEMANTICS scenario,automation,poisoned,schedule,quarantine]
# @ingroup ScenarioAutomation
# @BRIEF Flip enabled on every schedule bound to the quarantined (scenario_id, environment_id) pair.
# @POST Returns the ids whose enabled flag actually changed; caller owns the transaction.
def _set_schedules_enabled(db: Session, scenario_id: str, environment_id: str, enabled: bool) -> list[str]:
rows = (
db.query(ScenarioSchedule)
.filter(ScenarioSchedule.scenario_id == scenario_id, ScenarioSchedule.environment_id == environment_id)
.all()
)
changed: list[str] = []
for row in rows:
if bool(row.enabled) != enabled:
row.enabled = enabled
changed.append(row.id)
db.flush()
return changed
# #endregion ScenarioAutomation.Poisoned.Schedules
# #region ScenarioAutomation.Poisoned.Evaluate [C:4] [TYPE Function] [SEMANTICS scenario,automation,poisoned,evaluate,notification,quarantine]
# @ingroup ScenarioAutomation
# @BRIEF Record an infra failure; on the N-th identical signature disable schedules and emit one DLQ notification.
# @PRE The store is durable and the DB session is caller-owned.
# @POST Retries never exceed N-1 identical signatures; a fresh quarantine writes one blocked notification carrying
# the typed POISONED_RUN_QUARANTINE reason and disables matching schedules.
# @SIDE_EFFECT Appends the file store, flushes schedule.enabled=False rows and one notification row; caller commits.
# @INVARIANT A notification is emitted exactly once per quarantine epoch (newly_quarantined only).
def evaluate_infra_failure(
db: Session,
store: PoisonedRunStore,
*,
scenario_id: str,
environment_id: str,
error_code: str,
run_id: str | None = None,
threshold: int = QUARANTINE_THRESHOLD,
) -> dict[str, Any]:
decision = store.record_failure(
scenario_id=scenario_id,
environment_id=environment_id,
error_code=error_code,
run_id=run_id,
threshold=threshold,
)
if decision["newly_quarantined"]:
decision["disabled_schedule_ids"] = _set_schedules_enabled(db, scenario_id, environment_id, False)
notification = persist_notification(
db,
event_type="blocked",
scenario_id=scenario_id,
run_id=run_id,
severity="warning",
payload={
"reason": POISONED_RUN_QUARANTINE,
"error_code": error_code,
"environment_id": environment_id,
"identical_count": decision["identical_count"],
"quarantine_version": decision["version"],
"dlq": "scenario_automation",
},
)
decision["notification_id"] = notification.id
logger.reflect("Scenario automation quarantined after identical infra failures", src="ScenarioAutomation.Poisoned.Evaluate", payload={
"scenario_id": scenario_id,
"environment_id": environment_id,
"error_code": error_code,
"identical_count": decision["identical_count"],
"quarantine_version": decision["version"],
"disabled_schedule_ids": decision["disabled_schedule_ids"],
})
return decision
# #endregion ScenarioAutomation.Poisoned.Evaluate
# #region ScenarioAutomation.Poisoned.Unquarantine [C:4] [TYPE Function] [SEMANTICS scenario,automation,poisoned,unquarantine,recovery,operator]
# @ingroup ScenarioAutomation
# @BRIEF Operator recovery: CAS-release the quarantine, reset the counter and re-enable matching schedules.
# @PRE The pair is quarantined; expected_version matches the current quarantine version; actor is the operator identity.
# @POST No active quarantine for the pair; matching schedules are enabled again; the store logged the release.
# @SIDE_EFFECT Appends an unquarantine store line, flushes schedule.enabled=True rows; caller commits.
# @RAISES ValueError("NOT_QUARANTINED") | ValueError("STALE_QUARANTINE_VERSION") — the REST layer maps these to 404/409.
def unquarantine_scenario_automation(
db: Session,
store: PoisonedRunStore,
*,
scenario_id: str,
environment_id: str,
expected_version: int,
actor: str,
) -> dict[str, Any]:
result = store.unquarantine(
scenario_id=scenario_id,
environment_id=environment_id,
expected_version=expected_version,
actor=actor,
)
result["enabled_schedule_ids"] = _set_schedules_enabled(db, scenario_id, environment_id, True)
logger.reflect("Scenario automation quarantine released", src="ScenarioAutomation.Poisoned.Unquarantine", payload={
"scenario_id": scenario_id,
"environment_id": environment_id,
"quarantine_version": result["version"],
"released_by": actor,
"enabled_schedule_ids": result["enabled_schedule_ids"],
})
return result
# #endregion ScenarioAutomation.Poisoned.Unquarantine
# #endregion ScenarioAutomation.Poisoned

View File

@@ -0,0 +1,337 @@
# #region ScenarioAutomation.PoisonedStore [C:4] [TYPE Module] [SEMANTICS scenario,automation,poisoned,quarantine,failure,jsonl,locking]
# @defgroup ScenarioAutomation Durable file-backed identical-failure counter for the poisoned-run policy.
# @LAYER Service
# @BRIEF Count identical infra-failure signatures per (scenario_id, environment_id, error_code) in an
# append-only JSONL file and expose the derived counters and quarantine epochs.
# @RELATION DEPENDS_ON -> [ScenarioAutomation.Poisoned]
# @INVARIANT The JSONL file is the single durable counter authority; it is read and appended only under an
# exclusive interprocess file lock, so concurrent workers cannot lose an identical-failure increment.
# @INVARIANT The quarantine threshold is inclusive: the Nth identical failure quarantines and no retry is permitted.
# @INVARIANT Different error_code signatures never sum — only consecutive identical signatures reach N.
# @INVARIANT A successful run resets that (scenario_id, environment_id) streak; quarantine is lifted only by the
# explicit operator unquarantine CAS on the quarantine version.
# @RATIONALE 046 closure freezes Alembic (D11), so the counter lives in a JSONL file next to the data root instead
# of a new ORM table or an aggregate over notification payloads (no uniqueness there). fcntl.flock plus
# an intra-process lock gives cross-process and cross-thread exclusion; a .lock file alongside the store
# keeps same-filesystem semantics and kernel-driven cleanup on process exit.
# @REJECTED A new scenario_run_failure_signatures table was rejected — 046 migrations are frozen.
# Summing notification-event payloads was rejected — ScenarioNotificationEvent has no uniqueness key,
# so a retried poison signal would double-count or lose increments.
# Fixed-name temp state files were rejected — two processes would collide; append + fsync under flock is
# atomic per line and preserves the full audit trail.
from __future__ import annotations
from collections.abc import Iterator
from contextlib import contextmanager
from datetime import UTC, datetime
import fcntl
import json
import os
from pathlib import Path
import threading
from typing import Any
QUARANTINE_THRESHOLD = 3
POISONED_RUN_QUARANTINE = "POISONED_RUN_QUARANTINE"
DEFAULT_STORE_RELATIVE_PATH = Path("data") / "dashboard_testing" / "poisoned_failures.jsonl"
# #region ScenarioAutomation.PoisonedStore.DefaultPath [C:1] [TYPE Function] [SEMANTICS scenario,automation,poisoned,path]
# @ingroup ScenarioAutomation
def default_poisoned_store_path() -> Path:
override = os.environ.get("SCENARIO_POISONED_STORE_PATH", "").strip()
return Path(override) if override else DEFAULT_STORE_RELATIVE_PATH
# #endregion ScenarioAutomation.PoisonedStore.DefaultPath
# #region ScenarioAutomation.PoisonedStore.Lock [C:3] [TYPE Function] [SEMANTICS scenario,automation,poisoned,locking,fcntl,concurrency]
# @ingroup ScenarioAutomation
# @BRIEF Acquire exclusive access to the store file: thread reentrancy + intra-process RLock + interprocess flock.
# @POST Returns only when this writer holds both the intra-process lock and the fcntl flock on the .lock file.
# @INVARIANT A concurrent process attempting the same path blocks on fcntl.flock; the lock releases on fd close.
_intra_locks: dict[Path, threading.RLock] = {}
_intra_locks_guard = threading.Lock()
_lock_state = threading.local()
@contextmanager
def _store_lock(store_path: Path) -> Iterator[None]:
resolved = store_path.resolve()
held = getattr(_lock_state, "held", None)
if held is None:
held = _lock_state.held = set()
if resolved in held:
yield
return
with _intra_locks_guard:
intra = _intra_locks.setdefault(resolved, threading.RLock())
with intra:
lock_path = resolved.parent / (resolved.name + ".lock")
lock_path.parent.mkdir(parents=True, exist_ok=True)
fd = os.open(str(lock_path), os.O_CREAT | os.O_RDWR, 0o644)
try:
fcntl.flock(fd, fcntl.LOCK_EX)
held.add(resolved)
try:
yield
finally:
held.discard(resolved)
finally:
fcntl.flock(fd, fcntl.LOCK_UN)
os.close(fd)
# #endregion ScenarioAutomation.PoisonedStore.Lock
# #region ScenarioAutomation.PoisonedStore.Record [C:1] [TYPE Function] [SEMANTICS scenario,automation,poisoned,record]
def _now_iso() -> str:
return datetime.now(UTC).isoformat()
def _record(
outcome: str,
scenario_id: str,
environment_id: str,
error_code: str,
*,
run_id: str | None = None,
identical_count: int | None = None,
version: int | None = None,
expected_version: int | None = None,
reason: str | None = None,
actor: str | None = None,
) -> dict[str, Any]:
entry: dict[str, Any] = {
"outcome": outcome,
"scenario_id": scenario_id,
"environment_id": environment_id,
"error_code": error_code,
"recorded_at": _now_iso(),
}
if run_id:
entry["run_id"] = run_id
if identical_count is not None:
entry["identical_count"] = identical_count
if version is not None:
entry["version"] = version
if expected_version is not None:
entry["expected_version"] = expected_version
if reason:
entry["reason"] = reason
if actor:
entry["actor"] = actor
return entry
# #endregion ScenarioAutomation.PoisonedStore.Record
# #region ScenarioAutomation.PoisonedStore.ReadAppend [C:2] [TYPE Function] [SEMANTICS scenario,automation,poisoned,jsonl,atomic]
# @ingroup ScenarioAutomation
# @BRIEF Read the JSONL audit log / append one fsynced line. A torn trailing line is skipped, not fatal.
def _read_records(path: Path) -> list[dict[str, Any]]:
if not path.exists():
return []
records: list[dict[str, Any]] = []
with path.open("r", encoding="utf-8") as handle:
for line in handle:
line = line.strip()
if not line:
continue
try:
record = json.loads(line)
except json.JSONDecodeError:
continue
if isinstance(record, dict):
records.append(record)
return records
def _append_record(path: Path, record: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
line = json.dumps(record, ensure_ascii=True, sort_keys=True, separators=(",", ":")) + "\n"
with path.open("a", encoding="utf-8") as handle:
handle.write(line)
handle.flush()
os.fsync(handle.fileno())
# #endregion ScenarioAutomation.PoisonedStore.ReadAppend
# #region ScenarioAutomation.PoisonedStore.Fold [C:3] [TYPE Function] [SEMANTICS scenario,automation,poisoned,fold,state]
# @ingroup ScenarioAutomation
# @BRIEF Fold the durable JSONL log into consecutive-failure counters, active quarantines and quarantine epochs.
# @POST Counters are per (scenario_id, environment_id, error_code); a success clears that pair; unquarantine clears
# both its quarantine and that pair's counters.
def _reset_pair(counters: dict[tuple[str, str, str], int], pair: tuple[str, str]) -> None:
for key in [key for key in counters if (key[0], key[1]) == pair]:
del counters[key]
def _fold(records: list[dict[str, Any]]) -> tuple[dict[tuple[str, str, str], int], dict[tuple[str, str], dict[str, Any]], dict[tuple[str, str], int]]:
counters: dict[tuple[str, str, str], int] = {}
quarantines: dict[tuple[str, str], dict[str, Any]] = {}
epochs: dict[tuple[str, str], int] = {}
for record in records:
scenario_id = str(record.get("scenario_id", ""))
environment_id = str(record.get("environment_id", ""))
error_code = str(record.get("error_code", ""))
pair = (scenario_id, environment_id)
outcome = record.get("outcome")
if outcome == "failure":
key = (scenario_id, environment_id, error_code)
counters[key] = counters.get(key, 0) + 1
elif outcome == "success":
_reset_pair(counters, pair)
elif outcome == "quarantine":
epochs[pair] = epochs.get(pair, 0) + 1
quarantines[pair] = {
"scenario_id": scenario_id,
"environment_id": environment_id,
"error_code": error_code,
"identical_count": int(record.get("identical_count", 0)),
"version": int(record.get("version", epochs[pair])),
"reason": str(record.get("reason", POISONED_RUN_QUARANTINE)),
"quarantined_at": record.get("recorded_at"),
}
elif outcome == "unquarantine":
quarantines.pop(pair, None)
_reset_pair(counters, pair)
return counters, quarantines, epochs
# #endregion ScenarioAutomation.PoisonedStore.Fold
# #region ScenarioAutomation.PoisonedStore.Store [C:4] [TYPE Class] [SEMANTICS scenario,automation,poisoned,store,quarantine,recovery]
# @ingroup ScenarioAutomation
# @BRIEF File-backed identical-failure counter with N=identical quarantine and operator CAS unquarantine.
class PoisonedRunStore:
def __init__(self, path: str | Path) -> None:
self.path = Path(path)
def _snapshot(self) -> tuple[dict[tuple[str, str, str], int], dict[tuple[str, str], dict[str, Any]], dict[tuple[str, str], int]]:
return _fold(_read_records(self.path))
# #region ScenarioAutomation.PoisonedStore.Store.RecordFailure [C:4] [TYPE Function] [SEMANTICS scenario,automation,poisoned,record,failure]
# @ingroup ScenarioAutomation
# @BRIEF Append one identical infra-failure and return retry/quarantine after N identical signatures.
# @PRE scenario_id, environment_id and error_code are non-empty.
# @POST The log has one new failure line; the Nth identical line is followed by a quarantine line.
# @SIDE_EFFECT Appends and fsyncs the JSONL store under an exclusive lock.
# @INVARIANT Distinct error_code counters never combine, and an already-quarantined pair never returns retry.
def record_failure(
self,
*,
scenario_id: str,
environment_id: str,
error_code: str,
run_id: str | None = None,
threshold: int = QUARANTINE_THRESHOLD,
) -> dict[str, Any]:
if not scenario_id or not environment_id or not error_code:
raise ValueError("poisoned signature requires scenario_id, environment_id and error_code")
with _store_lock(self.path):
counters, quarantines, epochs = self._snapshot()
pair = (scenario_id, environment_id)
key = (scenario_id, environment_id, error_code)
count = counters.get(key, 0) + 1
_append_record(self.path, _record("failure", scenario_id, environment_id, error_code, run_id=run_id))
existing = quarantines.get(pair)
if existing is not None:
return {
"action": "quarantine",
"newly_quarantined": False,
"reason": existing["reason"],
"identical_count": existing["identical_count"],
"version": existing["version"],
"quarantine": existing,
}
if count >= max(1, int(threshold)):
version = epochs.get(pair, 0) + 1
_append_record(self.path, _record("quarantine", scenario_id, environment_id, error_code, identical_count=count, version=version, reason=POISONED_RUN_QUARANTINE))
quarantine = {
"scenario_id": scenario_id,
"environment_id": environment_id,
"error_code": error_code,
"identical_count": count,
"version": version,
"reason": POISONED_RUN_QUARANTINE,
"quarantined_at": _now_iso(),
}
return {
"action": "quarantine",
"newly_quarantined": True,
"reason": POISONED_RUN_QUARANTINE,
"identical_count": count,
"version": version,
"quarantine": quarantine,
}
return {"action": "retry", "newly_quarantined": False, "reason": None, "identical_count": count, "version": None, "quarantine": None}
# #endregion ScenarioAutomation.PoisonedStore.Store.RecordFailure
# #region ScenarioAutomation.PoisonedStore.Store.RecordSuccess [C:3] [TYPE Function] [SEMANTICS scenario,automation,poisoned,record,success,reset]
# @ingroup ScenarioAutomation
# @BRIEF Append a success marker that resets that pair's identical-failure streak.
# @POST The pair's counters are cleared; an active quarantine is NOT cleared (operator CAS only).
def record_success(self, *, scenario_id: str, environment_id: str, run_id: str | None = None) -> dict[str, Any]:
if not scenario_id or not environment_id:
raise ValueError("success requires scenario_id and environment_id")
with _store_lock(self.path):
counters, _, _ = self._snapshot()
pair = (scenario_id, environment_id)
reset = sum(count for key, count in counters.items() if (key[0], key[1]) == pair)
_append_record(self.path, _record("success", scenario_id, environment_id, "", run_id=run_id))
return {"reset": reset}
# #endregion ScenarioAutomation.PoisonedStore.Store.RecordSuccess
# #region ScenarioAutomation.PoisonedStore.Store.ReadState [C:2] [TYPE Function] [SEMANTICS scenario,automation,poisoned,read,state]
# @ingroup ScenarioAutomation
# @BRIEF Read the derived counters and active quarantines without mutating the store.
def identical_count(self, *, scenario_id: str, environment_id: str, error_code: str) -> int:
with _store_lock(self.path):
counters, _, _ = self._snapshot()
return counters.get((scenario_id, environment_id, error_code), 0)
def quarantine_state(self, *, scenario_id: str, environment_id: str) -> dict[str, Any] | None:
with _store_lock(self.path):
_, quarantines, _ = self._snapshot()
return quarantines.get((scenario_id, environment_id))
def is_quarantined(self, *, scenario_id: str, environment_id: str) -> bool:
return self.quarantine_state(scenario_id=scenario_id, environment_id=environment_id) is not None
def list_quarantines(self) -> list[dict[str, Any]]:
with _store_lock(self.path):
_, quarantines, _ = self._snapshot()
return sorted(quarantines.values(), key=lambda item: (item["scenario_id"], item["environment_id"]))
# #endregion ScenarioAutomation.PoisonedStore.Store.ReadState
# #region ScenarioAutomation.PoisonedStore.Store.Unquarantine [C:4] [TYPE Function] [SEMANTICS scenario,automation,poisoned,unquarantine,cas,operator]
# @ingroup ScenarioAutomation
# @BRIEF Operator recovery: CAS on the quarantine version, append unquarantine, clear counters.
# @PRE The pair is quarantined and expected_version equals the current quarantine version; actor is non-empty.
# @POST The store has an unquarantine line; the pair is no longer quarantined and its counters are zero.
# @RAISES ValueError("NOT_QUARANTINED") | ValueError("STALE_QUARANTINE_VERSION").
def unquarantine(self, *, scenario_id: str, environment_id: str, expected_version: int, actor: str) -> dict[str, Any]:
if not str(actor or "").strip():
raise ValueError("operator identity required")
with _store_lock(self.path):
_, quarantines, _ = self._snapshot()
current = quarantines.get((scenario_id, environment_id))
if current is None:
raise ValueError("NOT_QUARANTINED")
if expected_version is None or int(expected_version) != int(current["version"]):
raise ValueError("STALE_QUARANTINE_VERSION")
_append_record(self.path, _record("unquarantine", scenario_id, environment_id, current.get("error_code", ""), expected_version=int(expected_version), actor=str(actor)))
return {**current, "released_by": str(actor)}
# #endregion ScenarioAutomation.PoisonedStore.Store.Unquarantine
# #endregion ScenarioAutomation.PoisonedStore.Store
# #region ScenarioAutomation.PoisonedStore.Retry [C:1] [TYPE Function] [SEMANTICS scenario,automation,poisoned,retry]
def should_retry(decision: dict[str, Any]) -> bool:
return decision.get("action") == "retry"
# #endregion ScenarioAutomation.PoisonedStore.Retry
# #region ScenarioAutomation.PoisonedStore.Factory [C:1] [TYPE Function] [SEMANTICS scenario,automation,poisoned,factory]
def get_poisoned_store() -> PoisonedRunStore:
return PoisonedRunStore(default_poisoned_store_path())
# #endregion ScenarioAutomation.PoisonedStore.Factory
# #endregion ScenarioAutomation.PoisonedStore

View File

@@ -45,6 +45,7 @@ from src.models.scenario_automation import (
)
from src.models.scenario_registry import ScenarioRegistryEntry, ScenarioRevision
from src.models.scenario_run import ScenarioRun
from src.services.dashboard_testing.automation.poisoned_store import PoisonedRunStore
from src.services.dashboard_testing.scenario.templates import (
ACTION_REGISTRY_VERSION,
action_registry_fingerprint,
@@ -416,6 +417,7 @@ class TestScenarioAutomationAnonymousReads:
"/api/scenario-automation/notifications",
"/api/scenario-automation/metrics",
"/api/scenario-automation/retention",
"/api/scenario-automation/quarantine",
)
def test_anonymous_reads_require_authenticated_principal(self, dashboard_testing_client):
@@ -798,4 +800,69 @@ class TestScenarioAutomationRetentionReceipts:
finally:
cleanup.close()
# #endregion Test.Api.ScenarioAutomation.RetentionReceipts
# #region Test.Api.ScenarioAutomation.QuarantineRelease [C:3] [TYPE Class] [SEMANTICS test,api,scenario,automation,quarantine,recovery,operator]
# @BRIEF Operator-only quarantine listing and CAS release; RBAC forbids a non-MANAGE principal.
# @RELATION VERIFIES -> [Api.ScenarioAutomation.Quarantine]
# @TEST_INVARIANT Api.ScenarioAutomation.Quarantine: release requires scenario:automation MANAGE and a matching
# quarantine version; a released pair re-enables schedules and a second release is 404.
# @TEST_EDGE missing_manage_scope -> 403 on quarantine release
class TestScenarioAutomationQuarantineRelease:
_SCENARIO = "77777777-7777-4777-8777-777777777777"
_ENV = "env-preprod-01"
def test_release_requires_manage_scope(self, monkeypatch, tmp_path):
monkeypatch.setenv("SCENARIO_POISONED_STORE_PATH", str(tmp_path / "poisoned.jsonl"))
user = _make_user_with_permissions([("scenario:automation", "TRIGGER")])
with _client_for(user) as client:
resp = client.post(
f"/api/scenario-automation/quarantine/{self._SCENARIO}/release",
json={"environment_id": self._ENV, "expected_version": 1},
)
assert resp.status_code == 403
def test_operator_release_roundtrip(self, dashboard_testing_client, monkeypatch, tmp_path):
client = dashboard_testing_client
store_path = tmp_path / "poisoned" / "failures.jsonl"
monkeypatch.setenv("SCENARIO_POISONED_STORE_PATH", str(store_path))
created = client.post(
"/api/scenario-automation/schedules",
json={"scenario_id": self._SCENARIO, "environment_id": self._ENV, "cron_expr": "0 7 * * *"},
)
assert created.status_code == 201, created.text
schedule_id = created.json()["id"]
disabled = client.patch(
f"/api/scenario-automation/schedules/{schedule_id}",
json={"scenario_id": self._SCENARIO, "environment_id": self._ENV, "cron_expr": "0 7 * * *", "enabled": False},
)
assert disabled.status_code == 200
store = PoisonedRunStore(store_path)
for _ in range(3):
store.record_failure(scenario_id=self._SCENARIO, environment_id=self._ENV, error_code="INFRA_TIMEOUT")
listed = client.get("/api/scenario-automation/quarantine")
assert listed.status_code == 200
assert any(
item["scenario_id"] == self._SCENARIO and item["version"] == 1
for item in listed.json()
)
released = client.post(
f"/api/scenario-automation/quarantine/{self._SCENARIO}/release",
json={"environment_id": self._ENV, "expected_version": 1},
)
assert released.status_code == 200, released.text
assert schedule_id in released.json()["enabled_schedule_ids"]
schedules = client.get("/api/scenario-automation/schedules")
assert next(item for item in schedules.json() if item["id"] == schedule_id)["enabled"] is True
assert client.get("/api/scenario-automation/quarantine").json() == []
again = client.post(
f"/api/scenario-automation/quarantine/{self._SCENARIO}/release",
json={"environment_id": self._ENV, "expected_version": 1},
)
assert again.status_code == 404
assert again.json()["detail"]["code"] == "NOT_QUARANTINED"
client.delete(f"/api/scenario-automation/schedules/{schedule_id}")
# #endregion Test.Api.ScenarioAutomation.QuarantineRelease
# #endregion Test.Api.ScenarioAutomation

View File

@@ -0,0 +1,264 @@
# #region Test.ScenarioAutomation.PoisonedRunPolicy [C:4] [TYPE Module] [SEMANTICS testing,scenario,automation,poisoned,quarantine,dlq,concurrency]
# @defgroup Tests for ScenarioAutomation.Poisoned — durable identical-failure counter, N=3 quarantine, recovery.
# @LAYER Test
# @RELATION BINDS_TO -> [ScenarioAutomation.Poisoned]
# @RELATION BINDS_TO -> [ScenarioAutomation.PoisonedStore]
# @TEST_INVARIANT ScenarioAutomation.PoisonedStore: The JSONL store is the durable counter authority; concurrent
# writers cannot lose an increment under the exclusive lock. -> VERIFIED_BY:
# test_concurrent_writers_never_lose_increments, test_count_survives_new_store_instance
# @TEST_INVARIANT ScenarioAutomation.PoisonedStore: The Nth identical failure quarantines and never returns retry;
# distinct error_codes never sum. -> VERIFIED_BY:
# test_third_identical_failure_quarantines_without_retry, test_distinct_error_codes_do_not_sum
# @TEST_INVARIANT ScenarioAutomation.Poisoned: A fresh quarantine emits exactly one blocked POISONED_RUN_QUARANTINE
# notification and disables matching schedules. -> VERIFIED_BY:
# test_quarantine_notifies_once_and_disables_schedule
# @TEST_INVARIANT ScenarioAutomation.Poisoned.Unquarantine: Recovery is an operator CAS on the quarantine version;
# a stale or absent quarantine changes nothing. -> VERIFIED_BY:
# test_stale_version_release_is_rejected, test_unquarantine_reenables_and_resets
# @TEST_EDGE missing_signature_fields -> record_failure raises ValueError, no store write
# @TEST_EDGE torn_trailing_line -> unparseable last JSONL line is ignored, prior counts survive
# @TEST_EDGE external_fail -> a success resets the streak but never clears an active quarantine
# @RATIONALE Hardcoded thresholds/versions (N=3, version 1) keep expectations independent of the implementation;
# concurrency proof uses real threads against the real file lock, not a mocked lock.
from __future__ import annotations
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
import pytest
from src.models.scenario_automation import ScenarioNotificationEvent, ScenarioSchedule
from src.services.dashboard_testing.automation.poisoned import (
evaluate_infra_failure,
unquarantine_scenario_automation,
)
from src.services.dashboard_testing.automation.poisoned_store import (
POISONED_RUN_QUARANTINE,
PoisonedRunStore,
should_retry,
)
_SCENARIO = "11111111-1111-4111-8111-111111111111"
_ENV = "env-preprod-01"
@pytest.fixture
def store(tmp_path: Path) -> PoisonedRunStore:
return PoisonedRunStore(tmp_path / "poisoned" / "failures.jsonl")
@pytest.fixture
def quarantined_schedule(db_session) -> ScenarioSchedule:
schedule = ScenarioSchedule(
scenario_id=_SCENARIO,
environment_id=_ENV,
revision_id="rev-poisoned-1",
cron_expr="0 7 * * 1-5",
enabled=True,
)
db_session.add(schedule)
db_session.commit()
return schedule
# #region Test.ScenarioAutomation.PoisonedRunPolicy.TwoIdenticalRetry [C:2] [TYPE Function] [SEMANTICS test,poisoned,retry,threshold]
def test_two_identical_failures_allow_retry(store: PoisonedRunStore):
"""Two identical signatures stay below N=3 and keep retrying."""
first = store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
second = store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert (first["action"], first["identical_count"]) == ("retry", 1)
assert (second["action"], second["identical_count"]) == ("retry", 2)
assert should_retry(second) is True
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.TwoIdenticalRetry
# #region Test.ScenarioAutomation.PoisonedRunPolicy.ThirdIdenticalQuarantine [C:2] [TYPE Function] [SEMANTICS test,poisoned,quarantine,threshold]
def test_third_identical_failure_quarantines_without_retry(store: PoisonedRunStore):
"""The third identical signature quarantines at version 1 and never permits retry."""
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
third = store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert (third["action"], third["newly_quarantined"]) == ("quarantine", True)
assert (third["identical_count"], third["version"]) == (3, 1)
assert third["reason"] == POISONED_RUN_QUARANTINE
assert should_retry(third) is False
fourth = store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert (fourth["action"], fourth["newly_quarantined"]) == ("quarantine", False)
assert should_retry(fourth) is False
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.ThirdIdenticalQuarantine
# #region Test.ScenarioAutomation.PoisonedRunPolicy.DistinctErrorCodes [C:3] [TYPE Function] [SEMANTICS test,poisoned,signature,distinct]
def test_distinct_error_codes_do_not_sum(store: PoisonedRunStore):
"""Different error_code signatures keep separate counters and never combine to N."""
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="PROVIDER_DOWN")
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 2
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="PROVIDER_DOWN") == 1
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="PROVIDER_DOWN")
final = store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="PROVIDER_DOWN")
assert final["action"] == "quarantine"
assert final["quarantine"]["error_code"] == "PROVIDER_DOWN"
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 2
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.DistinctErrorCodes
# #region Test.ScenarioAutomation.PoisonedRunPolicy.SuccessResets [C:3] [TYPE Function] [SEMANTICS test,poisoned,success,reset]
def test_success_resets_streak_but_not_quarantine(store: PoisonedRunStore):
"""A success clears the pair's counter; an active quarantine still needs an operator."""
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_success(scenario_id=_SCENARIO, environment_id=_ENV, run_id="run-ok-1")
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 0
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is True
store.record_success(scenario_id=_SCENARIO, environment_id=_ENV, run_id="run-ok-2")
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is True
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.SuccessResets
# #region Test.ScenarioAutomation.PoisonedRunPolicy.Notification [C:4] [TYPE Function] [SEMANTICS test,poisoned,notification,dlq,schedule]
def test_quarantine_notifies_once_and_disables_schedule(db_session, store: PoisonedRunStore, quarantined_schedule: ScenarioSchedule):
"""Nth failure disables matching schedules and writes exactly one typed blocked notification."""
for _ in range(2):
evaluate_infra_failure(db_session, store, scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", run_id="run-fail-1")
db_session.commit()
assert db_session.get(ScenarioSchedule, quarantined_schedule.id).enabled is True
third = evaluate_infra_failure(db_session, store, scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", run_id="run-fail-3")
db_session.commit()
assert third["newly_quarantined"] is True
assert quarantined_schedule.id in third["disabled_schedule_ids"]
assert db_session.get(ScenarioSchedule, quarantined_schedule.id).enabled is False
events = db_session.query(ScenarioNotificationEvent).all()
assert len(events) == 1
event = events[0]
assert (event.event_type, event.severity, event.scenario_id, event.run_id) == ("blocked", "warning", _SCENARIO, "run-fail-3")
assert event.payload["reason"] == POISONED_RUN_QUARANTINE
assert (event.payload["error_code"], event.payload["identical_count"], event.payload["quarantine_version"]) == ("INFRA_TIMEOUT", 3, 1)
fourth = evaluate_infra_failure(db_session, store, scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", run_id="run-fail-4")
db_session.commit()
assert fourth["newly_quarantined"] is False
assert db_session.query(ScenarioNotificationEvent).count() == 1
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.Notification
# #region Test.ScenarioAutomation.PoisonedRunPolicy.Recovery [C:4] [TYPE Function] [SEMANTICS test,poisoned,unquarantine,cas,operator]
def test_unquarantine_reenables_and_resets(db_session, store: PoisonedRunStore, quarantined_schedule: ScenarioSchedule):
"""Operator CAS release clears quarantine, resets the counter and re-enables schedules."""
for _ in range(3):
evaluate_infra_failure(db_session, store, scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", run_id="run-fail")
db_session.commit()
result = unquarantine_scenario_automation(
db_session, store, scenario_id=_SCENARIO, environment_id=_ENV, expected_version=1, actor="operator-7"
)
db_session.commit()
assert result["version"] == 1
assert result["released_by"] == "operator-7"
assert quarantined_schedule.id in result["enabled_schedule_ids"]
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 0
assert db_session.get(ScenarioSchedule, quarantined_schedule.id).enabled is True
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.Recovery
# #region Test.ScenarioAutomation.PoisonedRunPolicy.StaleRelease [C:3] [TYPE Function] [SEMANTICS test,poisoned,unquarantine,cas,stale]
def test_stale_version_release_is_rejected(store: PoisonedRunStore):
"""A stale CAS version or an absent quarantine raises and mutates nothing."""
for _ in range(3):
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
with pytest.raises(ValueError, match="STALE_QUARANTINE_VERSION"):
store.unquarantine(scenario_id=_SCENARIO, environment_id=_ENV, expected_version=99, actor="operator-7")
with pytest.raises(ValueError, match="NOT_QUARANTINED"):
store.unquarantine(scenario_id="22222222-2222-4222-8222-222222222222", environment_id=_ENV, expected_version=1, actor="operator-7")
with pytest.raises(ValueError, match="operator identity required"):
store.unquarantine(scenario_id=_SCENARIO, environment_id=_ENV, expected_version=1, actor="")
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is True
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.StaleRelease
# #region Test.ScenarioAutomation.PoisonedRunPolicy.Durability [C:4] [TYPE Function] [SEMANTICS test,poisoned,durable,concurrency,atomic]
def test_concurrent_writers_never_lose_increments(store: PoisonedRunStore):
"""Concurrent processes/threads appending the same signature produce an exact count."""
writers = 24
def write(_: int) -> None:
PoisonedRunStore(store.path).record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", threshold=10_000)
with ThreadPoolExecutor(max_workers=8) as pool:
list(pool.map(write, range(writers)))
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == writers
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.Durability
# #region Test.ScenarioAutomation.PoisonedRunPolicy.CountSurvivesReload [C:3] [TYPE Function] [SEMANTICS test,poisoned,durable,reload]
def test_count_survives_new_store_instance(store: PoisonedRunStore):
"""A fresh store instance reads the same derived state from the JSONL file."""
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
reloaded = PoisonedRunStore(store.path)
assert reloaded.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 2
reloaded.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
assert reloaded.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is True
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.CountSurvivesReload
# #region Test.ScenarioAutomation.PoisonedRunPolicy.TornLine [C:2] [TYPE Function] [SEMANTICS test,poisoned,jsonl,torn]
def test_torn_trailing_line_is_ignored(store: PoisonedRunStore):
"""An unparseable last line does not destroy prior identical-failure counts."""
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
with store.path.open("a", encoding="utf-8") as handle:
handle.write('{"outcome": "failure", "scenario_id": "torn"\n')
assert store.identical_count(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT") == 2
assert store.is_quarantined(scenario_id=_SCENARIO, environment_id=_ENV) is False
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.TornLine
# #region Test.ScenarioAutomation.PoisonedRunPolicy.MissingSignature [C:2] [TYPE Function] [SEMANTICS test,poisoned,validation,edge]
def test_missing_signature_fields_are_rejected(store: PoisonedRunStore):
"""An incomplete signature is refused before any store write."""
with pytest.raises(ValueError, match="poisoned signature requires"):
store.record_failure(scenario_id="", environment_id=_ENV, error_code="INFRA_TIMEOUT")
with pytest.raises(ValueError, match="poisoned signature requires"):
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="")
assert store.path.exists() is False
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.MissingSignature
# #region Test.ScenarioAutomation.PoisonedRunPolicy.NoInfiniteRetry [C:3] [TYPE Function] [SEMANTICS test,poisoned,threshold,infinite-retry]
def test_threshold_configurable_and_no_infinite_retry(store: PoisonedRunStore):
"""A custom threshold still blocks retry exactly at N and never loops past it."""
decisions = [
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT", threshold=2)
for _ in range(4)
]
assert [d["action"] for d in decisions] == ["retry", "quarantine", "quarantine", "quarantine"]
assert [should_retry(d) for d in decisions] == [True, False, False, False]
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.NoInfiniteRetry
# #region Test.ScenarioAutomation.PoisonedRunPolicy.QuarantineListEpoch [C:3] [TYPE Function] [SEMANTICS test,poisoned,list,epoch]
def test_requarantine_increments_epoch_version(store: PoisonedRunStore):
"""After a release, a new identical streak quarantines at version 2 (fresh CAS token)."""
for _ in range(3):
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
store.unquarantine(scenario_id=_SCENARIO, environment_id=_ENV, expected_version=1, actor="operator-7")
for _ in range(3):
store.record_failure(scenario_id=_SCENARIO, environment_id=_ENV, error_code="INFRA_TIMEOUT")
state = store.quarantine_state(scenario_id=_SCENARIO, environment_id=_ENV)
assert state is not None
assert state["version"] == 2
assert store.list_quarantines() == [state]
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy.QuarantineListEpoch
# #endregion Test.ScenarioAutomation.PoisonedRunPolicy