diff --git a/backend/.env.example b/backend/.env.example index 7e1391d3..e32fe8c2 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -163,4 +163,8 @@ GROQ_MODELS=["llama-3.1-8b-instant","llama-3.1-70b-versatile","mixtral-8x7b-3276 # Location of the tamper-evident audit database written by the Flask ML API. # Point this at a durable, backed-up volume in production. Defaults to # audit_log.db beside api.py when unset. -# AUDIT_DB_PATH=/var/lib/spam-detection/audit_log.db \ No newline at end of file +# AUDIT_DB_PATH=/var/lib/spam-detection/audit_log.db +# Retention window (days) for the admin prune routine. Records older than this +# are deleted and the surviving chain is re-sealed. A non-positive value +# disables pruning. Defaults to 90. +# AUDIT_RETENTION_DAYS=90 diff --git a/backend/api.py b/backend/api.py index 1c5922d2..2528d057 100644 --- a/backend/api.py +++ b/backend/api.py @@ -1676,6 +1676,54 @@ def feedback_stats(): ) +# ============================================ +# AUDIT TRAIL QUERY (issue #1023) +# ============================================ + + +@app.route("/audit", methods=["GET"]) +@validate_request +@validate_internal_request +def get_audit_records(): + """Admin-only view over the tamper-evident audit trail (issue #1023). + + Gated by the same service-to-service internal secret as every other + privileged route; on top of that it requires the trusted backend to forward + the caller's authenticated identity (X-User-Username) and admin role + (X-User-Role), mirroring the Node admin gate. Supports exact-match filters + (actor, action, resource), an inclusive ISO-8601 time window (since, until) + and limit/offset pagination. The response also reports whether the stored + chain still verifies, so an operator sees integrity status alongside data. + """ + username = _require_username() + if not username: + raise ApiError( + ErrorCode.MISSING_USERNAME, "Missing X-User-Username header", 401 + ) + if not _is_admin_request(): + raise ApiError( + ErrorCode.FORBIDDEN, "Audit trail access requires the admin role", 403 + ) + + records = audit_store.query( + actor=request.args.get("actor"), + action=request.args.get("action"), + resource=request.args.get("resource"), + since=request.args.get("since"), + until=request.args.get("until"), + limit=request.args.get("limit", default=100, type=int), + offset=request.args.get("offset", default=0, type=int), + ) + return jsonify( + { + "success": True, + "count": len(records), + "records": records, + "chain_intact": audit_store.verify_chain(), + } + ) + + # ============================================ # EMAIL HEADER ANALYSIS # ============================================ @@ -2124,6 +2172,16 @@ def _require_username(): return request.headers.get("X-User-Username") +def _is_admin_request(): + """Whether the trusted backend forwarded an admin role for this caller. + + Reuses the existing X-User-* header convention and the ``admin`` role from + the published /api/roles vocabulary; the internal-secret gate upstream + guarantees the header can only originate from the trusted backend. + """ + return request.headers.get("X-User-Role", "").strip().lower() == "admin" + + @app.route("/imap/connect", methods=["POST"]) @validate_request @validate_internal_request diff --git a/backend/audit_store.py b/backend/audit_store.py index 06c9e329..617d48de 100644 --- a/backend/audit_store.py +++ b/backend/audit_store.py @@ -9,10 +9,17 @@ chain end to end and reports whether it is intact, giving operators a cheap integrity check over the whole trail. +Records are otherwise immutable and append-only. The one sanctioned mutation is +:func:`prune`, which enforces a retention window by deleting expired records and +then re-links the survivors from the genesis anchor so the retained trail keeps +verifying. :func:`query` exposes the trail for the admin-only ``GET /audit`` +endpoint with field filters, a time window and pagination. + The store is deliberately dependency-free (stdlib ``sqlite3`` only) and is meant -to be called fail-soft: a write failure must never break the request being -audited (see ``api.audit_log``). The database location is configurable via -``AUDIT_DB_PATH`` so deployments can point it at durable storage. +to be called fail-soft on the write path: a write failure must never break the +request being audited (see ``api.audit_log``). The database location is +configurable via ``AUDIT_DB_PATH`` and the retention window via +``AUDIT_RETENTION_DAYS``. >>> import os, tempfile >>> db = os.path.join(tempfile.mkdtemp(), "audit.db") @@ -22,7 +29,7 @@ True """ -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone import hashlib import json import os @@ -32,10 +39,14 @@ __all__ = [ "GENESIS_HASH", "DB_PATH", + "DEFAULT_RETENTION_DAYS", + "MAX_QUERY_LIMIT", "get_db_connection", "init_db", "append", "verify_chain", + "query", + "prune", ] # Default location for the audit database. Overridable so deployments can point @@ -45,6 +56,14 @@ str(Path(__file__).resolve().parent / "audit_log.db"), ) +# Age (in days) beyond which records are eligible for pruning when a caller does +# not pass an explicit window. A non-positive value disables pruning. +DEFAULT_RETENTION_DAYS = int(os.getenv("AUDIT_RETENTION_DAYS", "90")) + +# Upper bound on how many records a single query may return, so a caller cannot +# ask for an unbounded page. +MAX_QUERY_LIMIT = 1000 + # prev_hash of the first record. A fixed, all-zero digest anchors the chain so # the genesis record is verified with the same rule as every later one. GENESIS_HASH = "0" * 64 @@ -54,6 +73,9 @@ # payload; ordering/deletion is instead caught by the prev_hash linkage. _HASHED_FIELDS = ("actor", "action", "resource", "request_id", "status", "timestamp") +# Columns a caller may filter ``query`` on by exact match. +_FILTER_COLUMNS = ("actor", "action", "resource") + def get_db_connection(db_path=None): path = db_path if db_path is not None else DB_PATH @@ -158,6 +180,111 @@ def verify_chain(db_path=None): return True +def query( + actor=None, + action=None, + resource=None, + since=None, + until=None, + limit=100, + offset=0, + db_path=None, +): + """Return stored records, newest first, matching the given filters. + + ``actor``/``action``/``resource`` are exact-match filters. ``since`` and + ``until`` bound the ``timestamp`` column (inclusive) and are compared as + ISO-8601 UTC strings, whose lexical order matches chronological order. + ``limit`` is clamped to :data:`MAX_QUERY_LIMIT` and ``offset`` to a + non-negative value. Each result is a plain ``dict`` of all columns. + """ + init_db(db_path) + + clauses = [] + params = [] + for column, value in ( + ("actor", actor), + ("action", action), + ("resource", resource), + ): + if value is not None: + clauses.append(f"{column} = ?") + params.append(str(value)) + if since is not None: + clauses.append("timestamp >= ?") + params.append(str(since)) + if until is not None: + clauses.append("timestamp <= ?") + params.append(str(until)) + + where = f"WHERE {' AND '.join(clauses)}" if clauses else "" + safe_limit = max(1, min(int(limit), MAX_QUERY_LIMIT)) + safe_offset = max(0, int(offset)) + params.extend([safe_limit, safe_offset]) + + with get_db_connection(db_path) as conn: + rows = conn.execute( + f""" + SELECT * FROM audit_records + {where} + ORDER BY id DESC + LIMIT ? OFFSET ? + """, + params, + ).fetchall() + return [dict(row) for row in rows] + + +def prune(retention_days=None, now=None, db_path=None): + """Delete records older than the retention window and re-seal the chain. + + ``retention_days`` defaults to :data:`DEFAULT_RETENTION_DAYS`; a + non-positive window disables pruning and returns ``0``. Records whose + ``timestamp`` predates ``now - retention_days`` are removed, then the + surviving records are re-linked from the genesis anchor so + :func:`verify_chain` continues to pass over the retained trail. Returns the + number of records deleted. + """ + days = DEFAULT_RETENTION_DAYS if retention_days is None else int(retention_days) + if days <= 0: + return 0 + + init_db(db_path) + reference = now or datetime.now(timezone.utc) + cutoff = (reference - timedelta(days=days)).isoformat() + + with get_db_connection(db_path) as conn: + cursor = conn.execute( + "DELETE FROM audit_records WHERE timestamp < ?", (cutoff,) + ) + deleted = cursor.rowcount + if deleted: + _reseal(conn) + conn.commit() + return deleted + + +def _reseal(conn): + """Recompute prev_hash/record_hash for every surviving record in order. + + Pruning removes the oldest links, which would otherwise strand the earliest + survivor's ``prev_hash``. Re-sealing walks the remaining rows from the + genesis anchor and rewrites their linkage so the retained trail is once + again a valid chain. This is the only place stored hashes are rewritten and + it runs only as part of sanctioned retention maintenance. + """ + rows = conn.execute("SELECT * FROM audit_records ORDER BY id ASC").fetchall() + running_prev = GENESIS_HASH + for row in rows: + record = {k: row[k] for k in _HASHED_FIELDS} + record_hash = _compute_hash(running_prev, record) + conn.execute( + "UPDATE audit_records SET prev_hash = ?, record_hash = ? WHERE id = ?", + (running_prev, record_hash, row["id"]), + ) + running_prev = record_hash + + def _canonical(record): """Deterministic serialization of the signed fields. diff --git a/backend/openapi_spec.py b/backend/openapi_spec.py index b3d5ecb2..7eee1334 100644 --- a/backend/openapi_spec.py +++ b/backend/openapi_spec.py @@ -636,6 +636,52 @@ def _extended_paths(): }, } }, + "/audit": { + "get": { + "summary": "Query the tamper-evident audit trail (admin only)", + "operationId": "getAuditRecords", + "tags": ["System"], + "description": ( + "Returns persisted audit records, newest first. Requires " + "the internal secret plus an admin caller identified by the " + "X-User-Username and X-User-Role headers the trusted backend " + "forwards. Supports exact-match filters, an inclusive " + "ISO-8601 time window and limit/offset pagination." + ), + "parameters": [ + _query_param("actor", "Filter by the acting username."), + _query_param("action", "Filter by audited action name."), + _query_param("resource", "Filter by resource type."), + _query_param( + "since", + "Inclusive lower bound on the record timestamp " + "(ISO-8601 UTC).", + ), + _query_param( + "until", + "Inclusive upper bound on the record timestamp " + "(ISO-8601 UTC).", + ), + _query_param( + "limit", + "Maximum records to return (clamped to 1000).", + schema={"type": "integer", "default": 100}, + ), + _query_param( + "offset", + "Number of records to skip for pagination.", + schema={"type": "integer", "default": 0}, + ), + ], + "responses": { + "200": _json_response( + "Matching audit records plus chain-integrity status.", + {"$ref": "#/components/schemas/AuditQueryResponse"}, + ), + "401": _error_response("Missing X-User-Username header."), + "403": _error_response( + "Caller is not admin or lacks the internal secret." + ), "/cache-stats": { "get": { "summary": "/predict response cache statistics", @@ -1108,6 +1154,42 @@ def _extended_schemas(): }, }, }, + "AuditRecord": { + "type": "object", + "description": ( + "One persisted audit entry. `prev_hash`/`record_hash` form the " + "SHA-256 chain that makes the trail tamper-evident." + ), + "properties": { + "id": {"type": "integer"}, + "actor": {"type": "string"}, + "action": {"type": "string"}, + "resource": {"type": "string"}, + "request_id": {"type": "string"}, + "status": {"type": "integer"}, + "timestamp": {"type": "string", "description": "UTC ISO-8601."}, + "prev_hash": {"type": "string"}, + "record_hash": {"type": "string"}, + }, + }, + "AuditQueryResponse": { + "type": "object", + "properties": { + "success": {"type": "boolean"}, + "count": { + "type": "integer", + "description": "Number of records in this page.", + }, + "records": { + "type": "array", + "items": {"$ref": "#/components/schemas/AuditRecord"}, + }, + "chain_intact": { + "type": "boolean", + "description": "Whether the stored hash chain still verifies.", + }, + }, + }, "AuthUrlResponse": { "type": "object", "properties": {"auth_url": {"type": "string"}}, diff --git a/backend/tests/test_audit_query.py b/backend/tests/test_audit_query.py new file mode 100644 index 00000000..a4c1685b --- /dev/null +++ b/backend/tests/test_audit_query.py @@ -0,0 +1,115 @@ +"""Tests for audit trail querying and retention pruning (issue #1023).""" + +from datetime import datetime, timedelta, timezone +from pathlib import Path +import sys + +import pytest + +BASE_DIR = Path(__file__).resolve().parents[2] +BACKEND_DIR = BASE_DIR / "backend" + +sys.path.insert(0, str(BACKEND_DIR)) + +import audit_store + +NOW = datetime(2026, 7, 29, 12, 0, 0, tzinfo=timezone.utc) + + +def _ts(days_ago): + return (NOW - timedelta(days=days_ago)).isoformat() + + +@pytest.fixture +def store(tmp_path, monkeypatch): + db_path = tmp_path / "audit_query_test.db" + monkeypatch.setattr(audit_store, "DB_PATH", str(db_path)) + audit_store.init_db() + return audit_store + + +def test_query_returns_newest_first(store): + store.append("alice", "predict", "message", "req-1", 200) + store.append("bob", "predict", "message", "req-2", 200) + store.append("carol", "predict", "message", "req-3", 200) + + records = store.query() + + assert [r["actor"] for r in records] == ["carol", "bob", "alice"] + + +def test_filter_by_actor_action_and_resource(store): + store.append("alice", "predict", "message", "req-1", 200) + store.append("alice", "reload_model", "model", "req-2", 200) + store.append("bob", "predict", "message", "req-3", 200) + + assert {r["request_id"] for r in store.query(actor="alice")} == {"req-1", "req-2"} + assert {r["request_id"] for r in store.query(action="predict")} == { + "req-1", + "req-3", + } + assert {r["request_id"] for r in store.query(resource="model")} == {"req-2"} + # Filters compose (AND semantics). + assert {r["request_id"] for r in store.query(actor="alice", action="predict")} == { + "req-1" + } + + +def test_time_window_filters_are_inclusive(store): + store.append("alice", "predict", "message", "old", 200, timestamp=_ts(10)) + store.append("alice", "predict", "message", "mid", 200, timestamp=_ts(5)) + store.append("alice", "predict", "message", "new", 200, timestamp=_ts(1)) + + assert {r["request_id"] for r in store.query(since=_ts(5))} == {"mid", "new"} + assert {r["request_id"] for r in store.query(until=_ts(5))} == {"old", "mid"} + assert {r["request_id"] for r in store.query(since=_ts(5), until=_ts(5))} == {"mid"} + + +def test_pagination_limit_and_offset(store): + for i in range(5): + store.append("alice", "predict", "message", f"req-{i}", 200) + + first_page = store.query(limit=2) + second_page = store.query(limit=2, offset=2) + + assert [r["request_id"] for r in first_page] == ["req-4", "req-3"] + assert [r["request_id"] for r in second_page] == ["req-2", "req-1"] + + +def test_limit_is_clamped_to_sane_bounds(store): + for i in range(3): + store.append("alice", "predict", "message", f"req-{i}", 200) + + # Oversized limit returns everything without error; a zero limit still + # yields at least one row rather than an empty/invalid query. + assert len(store.query(limit=10_000)) == 3 + assert len(store.query(limit=0)) == 1 + + +def test_prune_deletes_records_older_than_window(store): + store.append("alice", "predict", "message", "ancient", 200, timestamp=_ts(120)) + store.append("alice", "predict", "message", "old", 200, timestamp=_ts(45)) + store.append("alice", "predict", "message", "recent", 200, timestamp=_ts(5)) + + deleted = store.prune(retention_days=30, now=NOW) + + assert deleted == 2 + remaining = {r["request_id"] for r in store.query()} + assert remaining == {"recent"} + + +def test_prune_reseals_so_chain_still_verifies(store): + store.append("alice", "predict", "message", "ancient", 200, timestamp=_ts(120)) + store.append("bob", "predict", "message", "old", 200, timestamp=_ts(45)) + store.append("carol", "predict", "message", "recent", 200, timestamp=_ts(5)) + + store.prune(retention_days=30, now=NOW) + + assert store.verify_chain() is True + + +def test_prune_is_a_noop_for_nonpositive_window(store): + store.append("alice", "predict", "message", "old", 200, timestamp=_ts(400)) + + assert store.prune(retention_days=0, now=NOW) == 0 + assert len(store.query()) == 1