"""Authorized HIST backtest writers — never INSERT OR REPLACE on immutable rows. Score/manifest idempotency (Codex 2886717e / e4b06c5a): NOT inherent to INSERT ... ON CONFLICT DO NOTHING. BEFORE INSERT identity-lock triggers abort before SQLite conflict handling when a lock already exists. Writers pre-check existing rows, and on IntegrityError re-read to return idempotent_same_payload ONLY when the stored payload matches. Changed payloads fail closed. Concurrent same-payload retries succeed without weakening DB locks. Authorized corrections append new events — they do NOT UPDATE immutable rows. v0.3.2 also guards every PRIMARY/UNIQUE conflict target. """ from __future__ import annotations import json import sqlite3 from typing import Any, Optional from srp.backtest.guards import ( BacktestAuthError, BacktestPolicyError, resolve_authorization, validate_score_row, ) class BacktestWriteError(ValueError): pass def _row_payload(row: sqlite3.Row | tuple) -> dict[str, Any]: if isinstance(row, sqlite3.Row): return { "score": row["score"], "score_label": row["score_label"], "status_policy_version": row["status_policy_version"], "split_role": row["split_role"], "input_obs_ids_json": row["input_obs_ids_json"], "audit_json": row["audit_json"], } return { "score": row[0], "score_label": row[1], "status_policy_version": row[2], "split_role": row[3], "input_obs_ids_json": row[4], "audit_json": row[5], } def _payload_equal(existing: dict[str, Any], *, score, score_label, status_policy_version, split_role, input_obs_ids_json, audit_json) -> bool: def norm(v): return None if v is None else v return ( norm(existing["score"]) == norm(score) and existing["score_label"] == score_label and norm(existing["status_policy_version"]) == norm(status_policy_version) and existing["split_role"] == split_role and norm(existing["input_obs_ids_json"]) == norm(input_obs_ids_json) and existing["audit_json"] == audit_json ) def insert_score_event( conn: sqlite3.Connection, *, run_id: str, construct_id: str, score_date: str, score: Optional[float], score_label: str, split_role: str, audit_json: str, status_policy_version: Optional[str] = None, input_obs_ids_json: Optional[str] = None, require_auth: bool = True, ) -> str: """Insert a score event or succeed idempotently on same payload. Returns: 'inserted' | 'idempotent_same_payload' Raises BacktestWriteError on replacement attempt / construct mismatch / auth failure. Never uses INSERT OR REPLACE. Codex 2886717e: raw INSERT ... ON CONFLICT DO NOTHING is NOT inherently idempotent — the BEFORE INSERT identity-lock trigger aborts before conflict handling when a lock already exists. This writer pre-checks, then recovers from IntegrityError by re-reading and returning idempotent_same_payload ONLY for an identical stored payload. Changed payloads fail closed. Concurrent same-payload retries succeed without weakening locks. """ validate_score_row( score=score, score_label=score_label, status_policy_version=status_policy_version, ) pin = conn.execute( """ SELECT p.construct_id, p.rule_id, p.authorized, p.authorization_decision_id FROM backtest_runs r JOIN backtest_formula_pins p ON p.pin_id = r.pin_id WHERE r.run_id = ? """, (run_id,), ).fetchone() if pin is None: raise BacktestWriteError(f"unknown run_id: {run_id}") pin_construct = pin["construct_id"] if isinstance(pin, sqlite3.Row) else pin[0] rule_id = pin["rule_id"] if isinstance(pin, sqlite3.Row) else pin[1] pin_auth = int(pin["authorized"] if isinstance(pin, sqlite3.Row) else pin[2]) decision_id = pin["authorization_decision_id"] if isinstance(pin, sqlite3.Row) else pin[3] if construct_id != pin_construct: raise BacktestWriteError( f"construct mismatch: score={construct_id} pin={pin_construct} (run→pin binding)" ) if require_auth: resolve_authorization( conn, rule_id=rule_id, authorization_decision_id=decision_id, pin_authorized_cache=pin_auth, expected_construct_id=construct_id, ) def _payload_outcome_for_existing(): row = conn.execute( """ SELECT score, score_label, status_policy_version, split_role, input_obs_ids_json, audit_json FROM backtest_score_events WHERE run_id=? AND construct_id=? AND score_date=? """, (run_id, construct_id, score_date), ).fetchone() if row is None: return None if _payload_equal( _row_payload(row), score=score, score_label=score_label, status_policy_version=status_policy_version, split_role=split_role, input_obs_ids_json=input_obs_ids_json, audit_json=audit_json, ): return "idempotent_same_payload" raise BacktestWriteError( "refusing replacement of immutable score event; " "use append-only correction events if a correction is authorized" ) existing_outcome = _payload_outcome_for_existing() if existing_outcome == "idempotent_same_payload": return existing_outcome # existing_outcome is None here (mismatch raises) try: cur = conn.execute( """ INSERT INTO backtest_score_events (run_id, construct_id, score_date, score, score_label, status_policy_version, split_role, input_obs_ids_json, audit_json) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(run_id, construct_id, score_date) DO NOTHING """, ( run_id, construct_id, score_date, score, score_label, status_policy_version, split_role, input_obs_ids_json, audit_json, ), ) except sqlite3.IntegrityError as e: # Lock-before-conflict or concurrent race: same payload => success; else fail closed. try: outcome = _payload_outcome_for_existing() except BacktestWriteError as be: if "refusing replacement" in str(be): raise raise BacktestWriteError(f"score insert rejected by DB guards: {e}") from e if outcome == "idempotent_same_payload": return outcome raise BacktestWriteError(f"score insert rejected by DB guards: {e}") from e if cur.rowcount == 0: outcome = _payload_outcome_for_existing() if outcome == "idempotent_same_payload": return outcome raise BacktestWriteError("score insert did not persist") row = conn.execute( """ SELECT score, score_label, status_policy_version, split_role, input_obs_ids_json, audit_json FROM backtest_score_events WHERE run_id=? AND construct_id=? AND score_date=? """, (run_id, construct_id, score_date), ).fetchone() if row is None: raise BacktestWriteError("score insert did not persist") if not _payload_equal( _row_payload(row), score=score, score_label=score_label, status_policy_version=status_policy_version, split_role=split_role, input_obs_ids_json=input_obs_ids_json, audit_json=audit_json, ): raise BacktestWriteError("immutable score identity occupied by different payload") conn.commit() return "inserted" def authorized_update_projection( conn: sqlite3.Connection, *, run_id: str, status: Optional[str] = None, summary_json: Optional[str] = None, metrics_json: Optional[str] = None, exposure_variant_ledger_json: Optional[str] = None, ) -> None: """Authorized mutable projection update (NOT score events). exposure_variant_ledger_json is a cache only — append exposure/variant events first. """ cols = [] vals: list[Any] = [] if status is not None: cols.append("status=?") vals.append(status) if summary_json is not None: cols.append("summary_json=?") vals.append(summary_json) if metrics_json is not None: cols.append("metrics_json=?") vals.append(metrics_json) if exposure_variant_ledger_json is not None: cols.append("exposure_variant_ledger_json=?") vals.append(exposure_variant_ledger_json) if not cols: return cols.append("updated_at=datetime('now')") vals.append(run_id) cur = conn.execute( f"UPDATE backtest_run_projections SET {', '.join(cols)} WHERE run_id=?", vals, ) if cur.rowcount == 0: raise BacktestWriteError(f"no projection row for run_id={run_id}") conn.commit() def append_exposure_variant_event( conn: sqlite3.Connection, *, run_id: str, event_kind: str, payload: dict[str, Any] | str, ) -> int: payload_json = payload if isinstance(payload, str) else json.dumps(payload, sort_keys=True) cur = conn.execute( """ INSERT INTO backtest_exposure_variant_events (run_id, event_kind, payload_json) VALUES (?, ?, ?) """, (run_id, event_kind, payload_json), ) conn.commit() return int(cur.lastrowid) def append_runtime_provenance_event( conn: sqlite3.Connection, *, run_id: str, new_git_or_worker_rev: str, reason: str, ) -> int: prior = conn.execute( "SELECT git_or_worker_rev FROM backtest_runs WHERE run_id=?", (run_id,), ).fetchone() if prior is None: raise BacktestWriteError(f"unknown run_id: {run_id}") prior_rev = prior[0] if not isinstance(prior, sqlite3.Row) else prior["git_or_worker_rev"] cur = conn.execute( """ INSERT INTO backtest_runtime_provenance_events (run_id, prior_git_or_worker_rev, new_git_or_worker_rev, reason) VALUES (?, ?, ?, ?) """, (run_id, prior_rev, new_git_or_worker_rev, reason), ) conn.commit() return int(cur.lastrowid) def insert_input_manifest( conn: sqlite3.Connection, *, manifest_id: str, pin_id: str, artifact_kind: str, artifact_path: str, artifact_sha256: str, raw_observation_ids_json: Optional[str] = None, notes: Optional[str] = None, ) -> str: """Insert an input manifest or succeed idempotently on same payload. Guards both conflict targets (manifest_id PK and alternate UNIQUE tuple). Returns: 'inserted' | 'idempotent_same_payload' Raises BacktestWriteError on replacement / changed payload. Never uses INSERT OR REPLACE. """ if len(artifact_sha256) != 64: raise BacktestWriteError("artifact_sha256 must be 64 hex chars") def _existing(): by_id = conn.execute( """ SELECT manifest_id, pin_id, artifact_kind, artifact_path, artifact_sha256, raw_observation_ids_json, notes FROM backtest_input_manifests WHERE manifest_id=? """, (manifest_id,), ).fetchone() by_nk = conn.execute( """ SELECT manifest_id, pin_id, artifact_kind, artifact_path, artifact_sha256, raw_observation_ids_json, notes FROM backtest_input_manifests WHERE pin_id=? AND artifact_kind=? AND artifact_path=? AND artifact_sha256=? """, (pin_id, artifact_kind, artifact_path, artifact_sha256), ).fetchone() return by_id, by_nk def _match(row) -> bool: if row is None: return False get = (lambda k, i: row[k] if isinstance(row, sqlite3.Row) else row[i]) return ( get("manifest_id", 0) == manifest_id and get("pin_id", 1) == pin_id and get("artifact_kind", 2) == artifact_kind and get("artifact_path", 3) == artifact_path and get("artifact_sha256", 4) == artifact_sha256 and (get("raw_observation_ids_json", 5) == raw_observation_ids_json) and (get("notes", 6) == notes) ) by_id, by_nk = _existing() if by_id is not None or by_nk is not None: if _match(by_id) or _match(by_nk): # Both conflict targets must refer to the same immutable row when present. if by_id is not None and by_nk is not None: id0 = by_id["manifest_id"] if isinstance(by_id, sqlite3.Row) else by_id[0] id1 = by_nk["manifest_id"] if isinstance(by_nk, sqlite3.Row) else by_nk[0] if id0 != id1: raise BacktestWriteError( "refusing replacement of immutable input manifest; " "conflict targets point at different rows" ) if _match(by_id if by_id is not None else by_nk): return "idempotent_same_payload" raise BacktestWriteError( "refusing replacement of immutable input manifest; " "use a new manifest_id and distinct natural key if authorized" ) try: cur = conn.execute( """ INSERT INTO backtest_input_manifests (manifest_id, pin_id, artifact_kind, artifact_path, artifact_sha256, raw_observation_ids_json, notes) VALUES (?, ?, ?, ?, ?, ?, ?) """, ( manifest_id, pin_id, artifact_kind, artifact_path, artifact_sha256, raw_observation_ids_json, notes, ), ) except sqlite3.IntegrityError as e: by_id, by_nk = _existing() if _match(by_id) or _match(by_nk): return "idempotent_same_payload" raise BacktestWriteError(f"manifest insert rejected by DB guards: {e}") from e if cur.rowcount == 0: raise BacktestWriteError("manifest insert did not persist") conn.commit() return "inserted" def register_authorization_decision( conn: sqlite3.Connection, *, decision_id: str, rule_id: str, construct_id: str, decided_by: str, decision: str = "authorize_backtest", evidence_json: Optional[str] = None, ) -> None: conn.execute( """ INSERT INTO backtest_authorization_decisions (decision_id, rule_id, construct_id, decided_by, decision, evidence_json) VALUES (?, ?, ?, ?, ?, ?) """, (decision_id, rule_id, construct_id, decided_by, decision, evidence_json), ) conn.commit()