stackchain-dashboard/src/idempotency.py
timmy 965a73ab49
All checks were successful
CI / lint (pull_request) Successful in 3m24s
CI / build-release (pull_request) Successful in 6s
CI / browser-journey (pull_request) Successful in 5m2s
CI / release-candidate (pull_request) Has been skipped
feat: rotate shared private-state keys (Closes #1237)
2026-08-21 21:12:54 +00:00

185 lines
7.6 KiB
Python

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