import json import sqlite3 import time from dataclasses import dataclass from pathlib import Path from typing import Any, Callable from src.private_state import connect_private_sqlite from src.state_encryption import PrivateStateCipher, private_state_encryption_config @dataclass(frozen=True) class Reservation: state: str response: Any = None class IdempotencyLedgerBusy(RuntimeError): """Raised when a durable completion cannot acquire the SQLite write lock.""" class IdempotencyLedger: """Small durable ledger for replaying successful Gitea mutations.""" def __init__( self, path: str | Path, *, ttl_seconds: float, max_entries: int, lock_timeout_seconds: float = 0.1, clock: Callable[[], float] = time.time, encryption_key: bytes | None = None, ) -> None: self.path = Path(path) self.ttl_seconds = ttl_seconds self.max_entries = max_entries self.lock_timeout_seconds = lock_timeout_seconds self.clock = clock self._cipher = PrivateStateCipher( encryption_key if encryption_key is not None else private_state_encryption_config(), store="idempotency-ledger", ) self._initialize() def _connect(self) -> sqlite3.Connection: connection = connect_private_sqlite(self.path, timeout=self.lock_timeout_seconds) connection.execute( f"PRAGMA busy_timeout = {max(1, int(self.lock_timeout_seconds * 1000))}" ) return connection def _initialize(self) -> None: with self._connect() as connection: connection.execute( """ CREATE TABLE IF NOT EXISTS idempotency_operations ( key TEXT PRIMARY KEY, fingerprint TEXT NOT NULL, status TEXT NOT NULL CHECK(status IN ('pending', 'completed')), response_json TEXT, created_at REAL NOT NULL, completed_at REAL ) """ ) connection.execute( "CREATE INDEX IF NOT EXISTS idempotency_completed_at_idx " "ON idempotency_operations(completed_at)" ) @staticmethod def _fingerprint(value: tuple[Any, ...]) -> str: return json.dumps(value, separators=(",", ":"), sort_keys=True) def _seal(self, value: str, *, key: str, field: str) -> str: return self._cipher.seal(value, binding=f"{key}:{field}") def _open(self, payload: str, *, key: str, field: str) -> tuple[str, bool]: value, legacy = self._cipher.open(payload, binding=f"{key}:{field}") if legacy: return json.dumps(value, separators=(",", ":"), sort_keys=True), True if not isinstance(value, str): raise RuntimeError("idempotency ledger payload is invalid") return value, legacy def reserve(self, key: str, fingerprint: tuple[Any, ...]) -> Reservation: encoded = self._fingerprint(fingerprint) now = self.clock() try: with self._connect() as connection: connection.execute("BEGIN IMMEDIATE") connection.execute( "DELETE FROM idempotency_operations " "WHERE status = 'completed' AND completed_at <= ?", (now - self.ttl_seconds,), ) row = connection.execute( "SELECT fingerprint, status, response_json, created_at FROM idempotency_operations " "WHERE key = ?", (key,), ).fetchone() if row is not None: stored_fingerprint, legacy_fingerprint = self._open( row[0], key=key, field="fingerprint" ) if legacy_fingerprint: connection.execute( "UPDATE idempotency_operations SET fingerprint = ? " "WHERE key = ? AND fingerprint = ?", ( self._seal( stored_fingerprint, key=key, field="fingerprint" ), key, row[0], ), ) if stored_fingerprint != encoded: return Reservation("conflict") if row[1] == "completed": stored_response, legacy_response = self._open( row[2], key=key, field="response" ) if legacy_response: connection.execute( "UPDATE idempotency_operations SET response_json = ? " "WHERE key = ? AND response_json = ?", ( self._seal( stored_response, key=key, field="response" ), key, row[2], ), ) return Reservation("completed", json.loads(stored_response)) if row[3] <= now - self.ttl_seconds: return Reservation("uncertain") return Reservation("pending") count = connection.execute( "SELECT COUNT(*) FROM idempotency_operations " "WHERE status = 'completed' OR created_at > ?", (now - self.ttl_seconds,), ).fetchone()[0] if count >= self.max_entries: completed = connection.execute( "SELECT key FROM idempotency_operations " "WHERE status = 'completed' ORDER BY completed_at, rowid LIMIT 1" ).fetchone() if completed is None: return Reservation("busy") connection.execute( "DELETE FROM idempotency_operations WHERE key = ?", (completed[0],), ) connection.execute( "INSERT INTO idempotency_operations " "(key, fingerprint, status, created_at) VALUES (?, ?, 'pending', ?)", (key, self._seal(encoded, key=key, field="fingerprint"), now), ) except sqlite3.OperationalError as exc: if "locked" in str(exc).lower() or "busy" in str(exc).lower(): return Reservation("busy") raise return Reservation("reserved") def clear(self) -> None: with self._connect() as connection: connection.execute("DELETE FROM idempotency_operations") def complete(self, key: str, response: Any) -> None: encoded = json.dumps(response, separators=(",", ":"), sort_keys=True) try: with self._connect() as connection: connection.execute( "UPDATE idempotency_operations SET status = 'completed', " "response_json = ?, completed_at = ? WHERE key = ?", (self._seal(encoded, key=key, field="response"), self.clock(), key), ) except sqlite3.OperationalError as exc: if "locked" in str(exc).lower() or "busy" in str(exc).lower(): raise IdempotencyLedgerBusy from exc raise