Files
dsh/server/dsh_sync/db.py
T

108 lines
3.5 KiB
Python

"""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()