"""PostgreSQL storage: schema and connection helpers (psycopg 3).""" from __future__ import annotations import json import psycopg from psycopg.rows import dict_row from .util import now_iso SCHEMA = [ """CREATE TABLE IF NOT EXISTS users ( id BIGSERIAL PRIMARY KEY, username TEXT NOT NULL UNIQUE, role TEXT NOT NULL DEFAULT 'member', created_at TEXT NOT NULL )""", """CREATE TABLE IF NOT EXISTS tokens ( token_hash TEXT PRIMARY KEY, user_id BIGINT NOT NULL REFERENCES users(id), name TEXT NOT NULL DEFAULT '', created_at TEXT NOT NULL, last_used_at TEXT, revoked BOOLEAN NOT NULL DEFAULT FALSE )""", """CREATE TABLE IF NOT EXISTS settings ( scope TEXT NOT NULL, scope_id BIGINT NOT NULL, key TEXT NOT NULL, value_json TEXT NOT NULL, version INTEGER NOT NULL DEFAULT 1, updated_by BIGINT NOT NULL, updated_at TEXT NOT NULL, PRIMARY KEY (scope, scope_id, key) )""", """CREATE TABLE IF NOT EXISTS plugins ( id TEXT PRIMARY KEY, name TEXT NOT NULL DEFAULT '', description TEXT NOT NULL DEFAULT '', latest_version TEXT, updated_at TEXT NOT NULL )""", """CREATE TABLE IF NOT EXISTS plugin_versions ( plugin_id TEXT NOT NULL REFERENCES plugins(id), version TEXT NOT NULL, channel TEXT NOT NULL DEFAULT 'stable', sha256 TEXT NOT NULL, size_bytes BIGINT NOT NULL, manifest_json TEXT NOT NULL, min_harness_version TEXT NOT NULL DEFAULT '0.0.0', artifact_path TEXT NOT NULL, published_by BIGINT NOT NULL, published_at TEXT NOT NULL, PRIMARY KEY (plugin_id, version) )""", """CREATE TABLE IF NOT EXISTS audit_log ( id BIGSERIAL PRIMARY KEY, actor BIGINT NOT NULL, action TEXT NOT NULL, detail_json TEXT NOT NULL DEFAULT '{}', created_at TEXT NOT NULL )""", "CREATE INDEX IF NOT EXISTS idx_audit_created ON audit_log (created_at)", ] def connect(cfg) -> psycopg.Connection: dsn = cfg.database_url if "connect_timeout" not in dsn: # Bound every connection attempt: never hang a request on the DB. dsn += ("&" if "?" in dsn else "?") + "connect_timeout=5" conn = psycopg.connect(dsn, row_factory=dict_row) return conn def init_db(cfg) -> None: conn = connect(cfg) try: with conn.cursor() as cur: for stmt in SCHEMA: cur.execute(stmt) conn.commit() finally: conn.close() def record_audit(conn, actor: int, action: str, detail: dict) -> None: """Append an audit row and commit it together with the change it describes.""" conn.execute( "INSERT INTO audit_log (actor, action, detail_json, created_at) VALUES (%s, %s, %s, %s)", (actor, action, json.dumps(detail), now_iso()), ) conn.commit() def reset_schema(cfg) -> None: """Test helper: wipe everything and recreate the schema.""" conn = connect(cfg) try: conn.autocommit = True with conn.cursor() as cur: cur.execute("DROP SCHEMA IF EXISTS public CASCADE") cur.execute("CREATE SCHEMA public") for stmt in SCHEMA: cur.execute(stmt) finally: conn.close()