railiance-telemetry/scripts/alert_ack.py
tegwick e7282e493d Implement authenticated alert receipt acknowledgments and audit delivery
Assistant: codex
Assistant-Model: gpt-6-astra
Assistant-Session: 01a0e6f1-443f-7783-9920-a16b2ffc467f
2026-09-28 11:11:33 +02:00

307 lines
16 KiB
Python

"""Durable alert acknowledgments behind owner-supplied identity and PDP adapters.
This module deliberately has no network entrypoint or trusted-header fallback.
The package must bind verified browser sessions and fresh authorization decisions.
"""
from dataclasses import dataclass, field
from contextlib import contextmanager
from datetime import datetime, timezone
import hashlib
import html
import hmac
import json
import os
from pathlib import Path
import re
import sqlite3
import stat
import time
from urllib.parse import parse_qs, urlencode
import uuid
from receiver import instant, strict_json
ROLE = 'railiance-admin'
TENANT = 'tenant:platform'
SOURCE = 'railiance-telemetry'
def encoded(value):
return json.dumps(value, sort_keys=True, separators=(',', ':'), allow_nan=False)
def occurrence(fingerprint, starts_at):
if not isinstance(fingerprint, str) or not re.fullmatch('[0-9a-f]{16}', fingerprint):
raise ValueError('invalid fingerprint')
# Preserve nanosecond precision in the wire timestamp; Python datetime only
# keeps microseconds. Alertmanager template and webhook use the same format.
if not isinstance(starts_at, str) or len(starts_at) > 40:
raise ValueError('invalid start time')
instant(starts_at)
return str(uuid.uuid5(uuid.NAMESPACE_URL, SOURCE + ':' + fingerprint + ':' + starts_at))
@dataclass(frozen=True)
class Actor:
"""Only constructed by the package's verified server-side session adapter."""
issuer: str
subject: str
tenant: str
roles: tuple
expires_at: float
csrf: str
principal_type: str = 'human'
assurance: dict = field(default_factory=dict)
groups: tuple = ()
tenant_source: str = ''
class Store:
def __init__(self, path):
path = Path(path)
parent = path.parent.lstat()
if (not stat.S_ISDIR(parent.st_mode) or parent.st_uid != os.getuid()
or parent.st_mode & 0o077):
raise ValueError('private owned state directory required')
fd = os.open(path, os.O_CREAT | os.O_RDWR | os.O_NOFOLLOW, 0o600)
try:
info = os.fstat(fd)
if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid()
or info.st_nlink != 1 or info.st_mode & 0o077):
raise ValueError('private owned database required')
finally:
os.close(fd)
self.path = str(path)
with self.connect() as db:
db.executescript('''
CREATE TABLE IF NOT EXISTS alerts (
id TEXT PRIMARY KEY, fingerprint TEXT NOT NULL,
starts_at TEXT NOT NULL, alertname TEXT NOT NULL,
labels_digest TEXT NOT NULL, received_at TEXT NOT NULL);
CREATE TABLE IF NOT EXISTS acknowledgments (
alert_id TEXT PRIMARY KEY REFERENCES alerts(id),
event_id TEXT UNIQUE NOT NULL, issuer TEXT NOT NULL,
subject TEXT NOT NULL, occurred_at TEXT NOT NULL,
decision_id TEXT NOT NULL);
CREATE TABLE IF NOT EXISTS outbox (
event_id TEXT PRIMARY KEY, body TEXT NOT NULL,
reference TEXT, status TEXT NOT NULL DEFAULT 'pending');
CREATE TRIGGER IF NOT EXISTS ack_immutable_update
BEFORE UPDATE ON acknowledgments BEGIN
SELECT RAISE(ABORT, 'immutable acknowledgment'); END;
CREATE TRIGGER IF NOT EXISTS ack_immutable_delete
BEFORE DELETE ON acknowledgments BEGIN
SELECT RAISE(ABORT, 'immutable acknowledgment'); END;
CREATE TRIGGER IF NOT EXISTS outbox_body_immutable
BEFORE UPDATE OF body, event_id ON outbox BEGIN
SELECT RAISE(ABORT, 'immutable audit event'); END;
CREATE TRIGGER IF NOT EXISTS outbox_immutable_delete
BEFORE DELETE ON outbox BEGIN
SELECT RAISE(ABORT, 'immutable audit event'); END;
CREATE TRIGGER IF NOT EXISTS alert_immutable_update
BEFORE UPDATE ON alerts BEGIN
SELECT RAISE(ABORT, 'immutable alert'); END;
CREATE TRIGGER IF NOT EXISTS alert_immutable_delete
BEFORE DELETE ON alerts BEGIN
SELECT RAISE(ABORT, 'immutable alert'); END;
''')
@contextmanager
def connect(self):
db = sqlite3.connect(self.path, timeout=10)
db.row_factory = sqlite3.Row
db.execute('PRAGMA foreign_keys=ON')
db.execute('PRAGMA synchronous=FULL')
try:
with db:
yield db
finally:
db.close()
def receive(self, payload, now):
if (payload.get('version') != '4' or payload.get('receiver') != 'railiance-admin-email'
or payload.get('truncatedAlerts', 0) != 0
or not isinstance(payload.get('alerts'), list)
or not 1 <= len(payload['alerts']) <= 100):
raise ValueError('invalid webhook')
rows = []
for alert in payload['alerts']:
labels = alert['labels']
if (alert['status'] not in ('firing', 'resolved') or not isinstance(labels, dict)
or not 1 <= len(labels) <= 32
or not all(isinstance(k, str) and isinstance(v, str)
and len(k) <= 100 and len(v) <= 256 for k, v in labels.items())
or not re.fullmatch('[A-Za-z_:][A-Za-z0-9_:]{0,99}', labels.get('alertname', ''))
or labels.get('owner') != SOURCE):
raise ValueError('invalid alert scope')
identity = occurrence(alert['fingerprint'], alert['startsAt'])
if instant(alert['startsAt']).timestamp() > now:
raise ValueError('future alert')
digest = hashlib.sha256(encoded(labels).encode()).hexdigest()
rows.append((identity, alert['fingerprint'], alert['startsAt'], labels['alertname'],
digest, datetime.fromtimestamp(now, timezone.utc).isoformat()))
with self.connect() as db:
db.execute('BEGIN IMMEDIATE')
for row in rows:
old = db.execute('SELECT labels_digest FROM alerts WHERE id=?', (row[0],)).fetchone()
if old and old['labels_digest'] != row[4]:
raise ValueError('alert identity collision')
db.execute('INSERT OR IGNORE INTO alerts VALUES (?,?,?,?,?,?)', row)
return [r[0] for r in rows]
def get(self, identity):
with self.connect() as db:
row = db.execute('''SELECT a.*, k.subject, k.occurred_at, o.reference,
o.status AS audit_status FROM alerts a
LEFT JOIN acknowledgments k ON a.id=k.alert_id
LEFT JOIN outbox o ON k.event_id=o.event_id WHERE a.id=?''', (identity,)).fetchone()
return dict(row) if row else None
def acknowledge(self, identity, actor, decision_id, now):
event_id = str(uuid.uuid4())
at = datetime.fromtimestamp(now, timezone.utc).isoformat()
with self.connect() as db:
db.execute('BEGIN IMMEDIATE')
alert = db.execute('SELECT * FROM alerts WHERE id=?', (identity,)).fetchone()
if alert is None:
raise ValueError('unknown alert')
old = db.execute('SELECT event_id FROM acknowledgments WHERE alert_id=?', (identity,)).fetchone()
if old:
return old['event_id']
event = dict(id=event_id, type='telemetry.alert.acknowledged', source=SOURCE,
subject='alert:' + identity, tenant=TENANT, correlation_id=identity,
occurred_at=at, data=dict(actor_issuer=actor.issuer,
actor_subject=actor.subject, role=ROLE, decision_id=decision_id,
fingerprint=alert['fingerprint'], starts_at=alert['starts_at'],
alertname=alert['alertname'], labels_sha256=alert['labels_digest'],
meaning='receipt acknowledged; not resolved or silenced'))
db.execute('INSERT INTO acknowledgments VALUES (?,?,?,?,?,?)',
(identity, event_id, actor.issuer, actor.subject, at, decision_id))
db.execute('INSERT INTO outbox(event_id,body) VALUES (?,?)', (event_id, encoded(event)))
return event_id
def drain(self, send):
"""send(event) returns (HTTP status, decoded body); transport owns custody.
Concurrent drains may replay identical events. The receiver's idempotency
contract makes that safe. No success without the exact receiver reference.
"""
with self.connect() as db:
rows = db.execute("SELECT event_id,body FROM outbox WHERE status='pending' LIMIT 20").fetchall()
for row in rows:
try:
status, body = send(json.loads(row['body']))
except (OSError, TimeoutError):
continue
if not isinstance(body, dict):
continue
accepted = (status, body.get('status')) in ((202, 'accepted'), (200, 'duplicate'))
reference = 'audit:' + row['event_id']
with self.connect() as db:
if accepted and body.get('reference') == reference:
db.execute("UPDATE outbox SET status='delivered',reference=? WHERE event_id=?", (reference, row['event_id']))
elif status in (400, 401, 403, 409, 422):
db.execute("UPDATE outbox SET status='blocked' WHERE event_id=? AND status='pending'", (row['event_id'],))
def requeue(self, event_id):
"""Explicit maintenance after fixing a receiver refusal; preserve bytes."""
with self.connect() as db:
return db.execute("UPDATE outbox SET status='pending' WHERE event_id=? AND status='blocked'",
(event_id,)).rowcount == 1
def audit_debt(self):
with self.connect() as db:
return {row['status']: row['n'] for row in db.execute(
"SELECT status,COUNT(*) AS n FROM outbox WHERE status!='delivered' GROUP BY status")}
class Application:
"""WSGI adapter. authenticate and authorize are required, never default-allow.
authenticate(environ) returns a verified Actor or None. authorize(actor,
action, resource) returns a fresh, binding-checked PDP decision receipt with
effect/id/expires_at; the package adapter must validate the native envelope.
"""
def __init__(self, store, origin, webhook_token, authenticate, authorize, clock=time.time):
from urllib.parse import urlsplit
parsed = urlsplit(origin)
if (parsed.scheme != 'https' or not parsed.hostname or parsed.path
or parsed.query or parsed.fragment or parsed.username or parsed.password):
raise ValueError('fixed HTTPS origin required')
if not isinstance(webhook_token, str) or len(webhook_token) < 32:
raise ValueError('dedicated webhook credential required')
if not callable(authenticate) or not callable(authorize):
raise ValueError('identity and policy adapters required')
self.store, self.origin, self.webhook_token = store, origin, webhook_token
self.authenticate, self.authorize, self.clock = authenticate, authorize, clock
def __call__(self, env, start):
def respond(code, body):
headers = [('Content-Type', 'text/html; charset=utf-8'), ('Cache-Control', 'no-store'),
('Content-Security-Policy', "default-src 'none'; form-action 'self'; frame-ancestors 'none'"),
('X-Content-Type-Options', 'nosniff'), ('Referrer-Policy', 'same-origin')]
start(code, headers)
return [body.encode()]
try:
method, path = env.get('REQUEST_METHOD'), env.get('PATH_INFO')
if path == '/webhook' and method == 'POST':
if not hmac.compare_digest(env.get('HTTP_AUTHORIZATION', ''), 'Bearer ' + self.webhook_token):
return respond('401 Unauthorized', 'Authentication required.')
raw = self.body(env)
self.store.receive(strict_json(raw), self.clock())
return respond('200 OK', 'Recorded.')
if path != '/ack/alerts' or method not in ('GET', 'POST'):
return respond('404 Not Found', 'Not found.')
actor = self.authenticate(env)
now = self.clock()
if (not isinstance(actor, Actor) or actor.expires_at <= now
or not actor.issuer or not actor.subject or len(actor.csrf) < 32):
return respond('401 Unauthorized', 'Sign in through the admitted identity provider.')
if actor.tenant != TENANT or actor.principal_type != 'human' or ROLE not in actor.roles:
return respond('403 Forbidden', 'Railiance admin role required.')
params = parse_qs(env.get('QUERY_STRING', ''), strict_parsing=True, max_num_fields=2)
if set(params) != {'fingerprint', 'starts_at'} or any(len(v) != 1 for v in params.values()):
raise ValueError('invalid link')
identity = occurrence(params['fingerprint'][0], params['starts_at'][0])
decision = self.authorize(actor, 'acknowledge' if method == 'POST' else 'read', 'alert:' + identity)
if (decision.get('effect') != 'allow' or not decision.get('id')
or decision.get('expires_at', 0) <= self.clock()):
return respond('403 Forbidden', 'Permission unavailable or refused.')
alert = self.store.get(identity)
if alert is None:
return respond('404 Not Found', 'Alert not received yet. Retry shortly.')
if method == 'POST':
if env.get('HTTP_ORIGIN') != self.origin:
return respond('403 Forbidden', 'Invalid origin.')
form = parse_qs(self.body(env).decode(), strict_parsing=True, max_num_fields=1)
if (set(form) != {'csrf'} or len(form['csrf']) != 1
or not hmac.compare_digest(form['csrf'][0], actor.csrf)):
return respond('403 Forbidden', 'Invalid confirmation.')
if actor.expires_at <= self.clock() or decision['expires_at'] <= self.clock():
return respond('403 Forbidden', 'Session or permission expired.')
self.store.acknowledge(identity, actor, decision['id'], self.clock())
alert = self.store.get(identity)
title = html.escape(alert['alertname'])
if alert['occurred_at']:
status = 'Audit record archived.' if alert['audit_status'] == 'delivered' else 'Audit delivery pending.'
return respond('200 OK', f'<h1>Receipt acknowledged</h1><p>{title}</p><p>{status}</p><p>This does not resolve or silence the alert.</p>')
query = html.escape(urlencode({k: v[0] for k, v in params.items()}), quote=True)
csrf = html.escape(actor.csrf, quote=True)
return respond('200 OK', f'<h1>{title}</h1><p>Started {html.escape(alert["starts_at"])}</p>'
'<p>Confirm that you received this alert. This does not resolve or silence it.</p>'
f'<form method="post" action="/ack/alerts?{query}"><input type="hidden" name="csrf" value="{csrf}">'
'<button type="submit">Acknowledge receipt</button></form>')
except (ValueError, KeyError, TypeError, AttributeError, UnicodeError):
return respond('400 Bad Request', 'Invalid request.')
except (OSError, sqlite3.Error):
return respond('503 Service Unavailable', 'Service unavailable. Retry later.')
@staticmethod
def body(env):
length = int(env.get('CONTENT_LENGTH', '0'))
if not 0 < length <= 32768:
raise ValueError('body size')
raw = env['wsgi.input'].read(length)
if len(raw) != length:
raise ValueError('incomplete body')
return raw