Implement authenticated alert receipt acknowledgments and audit delivery
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e6f1-443f-7783-9920-a16b2ffc467f
This commit is contained in:
parent
67283b66c2
commit
e7282e493d
25 changed files with 1881 additions and 0 deletions
307
scripts/alert_ack.py
Normal file
307
scripts/alert_ack.py
Normal file
|
|
@ -0,0 +1,307 @@
|
|||
"""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
|
||||
48
scripts/alert_audit.py
Normal file
48
scripts/alert_audit.py
Normal file
|
|
@ -0,0 +1,48 @@
|
|||
"""Bounded Audit Core transport for the acknowledgment outbox."""
|
||||
import json
|
||||
from urllib.error import HTTPError, URLError
|
||||
from urllib.parse import urlsplit
|
||||
from urllib.request import Request, HTTPRedirectHandler, ProxyHandler, build_opener
|
||||
|
||||
|
||||
class NoRedirect(HTTPRedirectHandler):
|
||||
def redirect_request(self, req, fp, code, msg, headers, newurl):
|
||||
return None
|
||||
|
||||
|
||||
class AuditTransport:
|
||||
def __init__(self, origin, token_provider, *, allow_internal_http=False):
|
||||
parsed = urlsplit(origin)
|
||||
internal = parsed.hostname in ('127.0.0.1', '::1') or (parsed.hostname or '').endswith(('.svc', '.svc.cluster.local'))
|
||||
if (parsed.scheme != 'https' and not (allow_internal_http and internal and parsed.scheme == 'http')
|
||||
or not parsed.hostname or parsed.path or parsed.query or parsed.fragment
|
||||
or parsed.username or parsed.password):
|
||||
raise ValueError('fixed audit origin required')
|
||||
self.url = origin + '/v1/events'
|
||||
self.token_provider = token_provider
|
||||
self.opener = build_opener(ProxyHandler({}), NoRedirect())
|
||||
|
||||
def __call__(self, event):
|
||||
token = self.token_provider()
|
||||
if (not isinstance(token, str) or not 1 <= len(token) <= 32768
|
||||
or any(ord(c) < 33 or ord(c) > 126 for c in token)):
|
||||
raise OSError('audit credential unavailable')
|
||||
request = Request(self.url, data=json.dumps(event, sort_keys=True, separators=(',', ':')).encode(),
|
||||
headers={'Content-Type': 'application/json', 'Authorization': 'Bearer ' + token,
|
||||
'Idempotency-Key': event['id']}, method='POST')
|
||||
try:
|
||||
with self.opener.open(request, timeout=5) as response:
|
||||
raw = response.read(65537)
|
||||
if len(raw) > 65536:
|
||||
raise OSError('audit response invalid')
|
||||
body = json.loads(raw)
|
||||
if not isinstance(body, dict):
|
||||
raise ValueError()
|
||||
return response.status, body
|
||||
except HTTPError as error:
|
||||
# Keep no upstream diagnostic body, URL, credential or response text.
|
||||
status = error.code
|
||||
error.close()
|
||||
return status, {}
|
||||
except (URLError, ValueError, UnicodeError):
|
||||
raise OSError('audit unavailable or invalid response') from None
|
||||
136
scripts/alert_identity.py
Normal file
136
scripts/alert_identity.py
Normal file
|
|
@ -0,0 +1,136 @@
|
|||
"""KeyCape OIDC code/PKCE sessions for the telemetry acknowledgment surface."""
|
||||
import base64
|
||||
import hashlib
|
||||
import secrets
|
||||
import threading
|
||||
import time
|
||||
from urllib.parse import urlencode, urlsplit
|
||||
|
||||
import jwt
|
||||
|
||||
from alert_ack import Actor, TENANT
|
||||
from telemetry_http import Transport, origin
|
||||
|
||||
SCOPES = ('openid', 'profile', 'email', 'telemetry:read', 'telemetry:acknowledge')
|
||||
|
||||
|
||||
class Login:
|
||||
def __init__(self, issuer, public_origin, *, transport=None, clock=time.time):
|
||||
self.issuer = origin(issuer)
|
||||
self.public_origin = origin(public_origin)
|
||||
self.callback = public_origin + '/ack/auth/callback'
|
||||
self.client = 'railiance-telemetry-admin'
|
||||
self.transport = transport or Transport()
|
||||
self.clock = clock
|
||||
self.lock = threading.RLock()
|
||||
self.pending, self.sessions = {}, {}
|
||||
|
||||
def prune(self):
|
||||
now = self.clock()
|
||||
self.pending = {k: v for k, v in self.pending.items() if v['expires'] > now}
|
||||
self.sessions = {k: v for k, v in self.sessions.items() if v.expires_at > now}
|
||||
|
||||
def metadata(self):
|
||||
status, data = self.transport.request('GET', self.issuer + '/.well-known/openid-configuration')
|
||||
if status != 200 or data.get('issuer') != self.issuer:
|
||||
raise ValueError('issuer unavailable')
|
||||
for key in ('authorization_endpoint', 'token_endpoint', 'jwks_uri'):
|
||||
p = urlsplit(data[key])
|
||||
if p.scheme + '://' + p.netloc != self.issuer or p.query or p.fragment or not p.path:
|
||||
raise ValueError('issuer endpoint refused')
|
||||
if 'S256' not in data.get('code_challenge_methods_supported', []):
|
||||
raise ValueError('PKCE unavailable')
|
||||
return data
|
||||
|
||||
def start(self, return_to):
|
||||
parsed = urlsplit(return_to)
|
||||
if parsed.scheme or parsed.netloc or parsed.fragment or parsed.path != '/ack/alerts' or len(return_to) > 1024:
|
||||
raise ValueError('invalid return path')
|
||||
metadata = self.metadata()
|
||||
state, browser, nonce, verifier = (secrets.token_urlsafe(32) for _ in range(4))
|
||||
with self.lock:
|
||||
self.prune()
|
||||
if len(self.pending) >= 1024:
|
||||
raise ValueError('login capacity')
|
||||
self.pending[state] = dict(browser=browser, nonce=nonce, verifier=verifier,
|
||||
expires=self.clock() + 300, return_to=return_to, metadata=metadata)
|
||||
challenge = base64.urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()).rstrip(b'=').decode()
|
||||
return metadata['authorization_endpoint'] + '?' + urlencode(dict(response_type='code',
|
||||
client_id=self.client, redirect_uri=self.callback, scope=' '.join(SCOPES),
|
||||
state=state, nonce=nonce, code_challenge=challenge, code_challenge_method='S256')), browser
|
||||
|
||||
def decode(self, token, audience, keys):
|
||||
if not isinstance(token, str) or not 1 <= len(token) <= 32768:
|
||||
raise ValueError('invalid token')
|
||||
header = jwt.get_unverified_header(token)
|
||||
if header.get('alg') != 'RS256' or not header.get('kid'):
|
||||
raise ValueError('invalid signing key')
|
||||
matches = [key for key in keys if key.get('kid') == header['kid'] and key.get('kty') == 'RSA'
|
||||
and key.get('use', 'sig') == 'sig' and key.get('alg', 'RS256') == 'RS256']
|
||||
if len(matches) != 1:
|
||||
raise ValueError('invalid signing key')
|
||||
claims = jwt.decode(token, jwt.PyJWK.from_dict(matches[0]).key, algorithms=['RS256'],
|
||||
issuer=self.issuer, audience=audience, options={
|
||||
'require': ['iss', 'aud', 'sub', 'iat', 'exp'], 'strict_aud': True})
|
||||
if (not isinstance(claims['sub'], str) or not claims['sub']
|
||||
or type(claims['exp']) is not int or type(claims['iat']) is not int
|
||||
or claims['exp'] <= self.clock() or claims['exp'] <= claims['iat']):
|
||||
raise ValueError('invalid claims')
|
||||
return claims
|
||||
|
||||
def finish(self, state, browser, code):
|
||||
with self.lock:
|
||||
self.prune()
|
||||
pending = self.pending.pop(state, None)
|
||||
if not pending or not browser or not secrets.compare_digest(browser, pending['browser']) or not 1 <= len(code) <= 4096:
|
||||
raise ValueError('invalid login state')
|
||||
try:
|
||||
status, tokens = self.transport.request('POST', pending['metadata']['token_endpoint'],
|
||||
urlencode(dict(grant_type='authorization_code', client_id=self.client,
|
||||
redirect_uri=self.callback, code=code, code_verifier=pending['verifier'])).encode(),
|
||||
{'Content-Type': 'application/x-www-form-urlencoded'})
|
||||
if status != 200:
|
||||
raise ValueError('code exchange failed')
|
||||
status, jwks = self.transport.request('GET', pending['metadata']['jwks_uri'])
|
||||
if status != 200 or not isinstance(jwks.get('keys'), list):
|
||||
raise ValueError('issuer keys unavailable')
|
||||
identity = self.decode(tokens['id_token'], self.client, jwks['keys'])
|
||||
access = self.decode(tokens['access_token'], 'railiance-telemetry', jwks['keys'])
|
||||
if identity.get('nonce') != pending['nonce']:
|
||||
raise ValueError('nonce mismatch')
|
||||
for key in ('sub', 'tenant', 'tenant_source', 'principal_type', 'roles', 'groups', 'assurance'):
|
||||
if key not in identity or identity[key] != access.get(key):
|
||||
raise ValueError('paired identity mismatch')
|
||||
if (access['tenant'] != TENANT or access['principal_type'] != 'human'
|
||||
or access['tenant_source'] not in ('directory', 'registration')
|
||||
or set(access.get('scope', '').split()) != set(SCOPES)):
|
||||
raise ValueError('human platform scope required')
|
||||
for key in ('roles', 'groups'):
|
||||
if not isinstance(access[key], list) or any(not isinstance(v, str) or not v for v in access[key]):
|
||||
raise ValueError('invalid identity claims')
|
||||
assurance = access['assurance']
|
||||
if (assurance.get('level') != 'aal2' or assurance.get('mfa') is not True
|
||||
or assurance.get('source') != 'key-cape' or assurance.get('methods') != ['pwd', 'otp']
|
||||
or type(assurance.get('at')) is not int or not -30 <= self.clock() - assurance['at'] <= 900):
|
||||
raise ValueError('fresh MFA required')
|
||||
actor = Actor(self.issuer, access['sub'], TENANT, tuple(access['roles']),
|
||||
min(access['exp'], identity['exp'], self.clock() + 900), secrets.token_urlsafe(32),
|
||||
assurance=dict(assurance), groups=tuple(access['groups']), tenant_source=access['tenant_source'])
|
||||
with self.lock:
|
||||
self.prune()
|
||||
if len(self.sessions) >= 1024:
|
||||
raise ValueError('session capacity')
|
||||
sid = secrets.token_urlsafe(32)
|
||||
self.sessions[sid] = actor
|
||||
return sid, pending['return_to']
|
||||
except (jwt.PyJWTError, KeyError, TypeError, AttributeError):
|
||||
raise ValueError('invalid issuer response') from None
|
||||
|
||||
def session(self, sid):
|
||||
with self.lock:
|
||||
self.prune()
|
||||
return self.sessions.get(sid)
|
||||
|
||||
def logout(self, sid):
|
||||
with self.lock:
|
||||
self.sessions.pop(sid, None)
|
||||
111
scripts/alert_policy.py
Normal file
111
scripts/alert_policy.py
Normal file
|
|
@ -0,0 +1,111 @@
|
|||
"""Consume native Flex Auth decisions with exact request and package binding.
|
||||
|
||||
Wire canonicalization follows the existing Informed Decision consumer contract:
|
||||
Go struct field order, sorted maps and Go JSON HTML escaping. No local allow rule.
|
||||
"""
|
||||
from datetime import datetime
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
import time
|
||||
import uuid
|
||||
|
||||
from telemetry_http import Transport, origin
|
||||
|
||||
DIGEST = re.compile(r'sha256:[0-9a-f]{64}')
|
||||
|
||||
|
||||
def request_digest(request):
|
||||
def ordered(value):
|
||||
if isinstance(value, dict):
|
||||
return {key: ordered(value[key]) for key in sorted(value)}
|
||||
if isinstance(value, list):
|
||||
return [ordered(v) for v in value]
|
||||
if value is None or type(value) in (str, bool) or type(value) is int and abs(value) <= 2**53:
|
||||
return value
|
||||
raise ValueError('unsupported policy input')
|
||||
def ref(value, fields):
|
||||
return {k: ordered(value[k]) for k in fields if value.get(k)}
|
||||
material = dict(tenant=request['tenant'],
|
||||
subject=ref(request['subject'], ('id', 'type', 'tenant', 'attributes')),
|
||||
action=request['action'],
|
||||
resource=ref(request['resource'], ('id', 'type', 'system', 'tenant', 'attributes')))
|
||||
if request.get('context'):
|
||||
material['context'] = ordered(request['context'])
|
||||
raw = json.dumps(material, ensure_ascii=False, separators=(',', ':'), allow_nan=False)
|
||||
for char, escaped in (('<', r'\u003c'), ('>', r'\u003e'), ('&', r'\u0026'), ('\u2028', r'\u2028'), ('\u2029', r'\u2029')):
|
||||
raw = raw.replace(char, escaped)
|
||||
return 'sha256:' + hashlib.sha256(raw.encode()).hexdigest()
|
||||
|
||||
|
||||
def timestamp(value):
|
||||
parsed = datetime.fromisoformat(value.replace('Z', '+00:00'))
|
||||
if parsed.tzinfo is None:
|
||||
raise ValueError('missing timezone')
|
||||
return parsed.timestamp()
|
||||
|
||||
|
||||
class Policy:
|
||||
def __init__(self, endpoint, token_provider, package, version, digest, *, transport=None, clock=time.time):
|
||||
self.endpoint = origin(endpoint, internal=True) + '/v1/check'
|
||||
if not package or not version or not DIGEST.fullmatch(digest):
|
||||
raise ValueError('pinned policy required')
|
||||
self.package, self.version, self.digest = package, version, digest
|
||||
self.token_provider, self.transport, self.clock = token_provider, transport or Transport(), clock
|
||||
|
||||
def __call__(self, actor, action, resource):
|
||||
deny = {'effect': 'deny', 'id': '', 'expires_at': 0}
|
||||
if action not in ('read', 'acknowledge') or not re.fullmatch(r'alert:[0-9a-f-]{36}', resource):
|
||||
return deny
|
||||
tenant_source = {'directory': 'directory-asserted', 'registration': 'registration-supplied'}.get(actor.tenant_source)
|
||||
if not tenant_source:
|
||||
return deny
|
||||
request = dict(id=str(uuid.uuid4()), tenant='tenant:platform',
|
||||
subject=dict(id=actor.subject, type=actor.principal_type, tenant=actor.tenant,
|
||||
attributes=dict(issuer=actor.issuer, roles=list(actor.roles), groups=list(actor.groups),
|
||||
assurance=actor.assurance, tenant_source=tenant_source,
|
||||
principal_type_source='authentication-derived')),
|
||||
action=action, resource=dict(id=resource, type='telemetry-alert', system='railiance-telemetry',
|
||||
tenant='tenant:platform'), policy_version=self.version)
|
||||
# Registry assignments and current authentication evidence are distinct.
|
||||
# Registry enrichment may replace subject attributes; it must not replace
|
||||
# the live signed MFA/role observations used for revocation/freshness.
|
||||
request['context'] = {'authentication': dict(request['subject']['attributes'])}
|
||||
started = self.clock()
|
||||
try:
|
||||
token = self.token_provider()
|
||||
if not isinstance(token, str) or not 1 <= len(token) <= 32768 or any(ord(c) < 33 or ord(c) > 126 for c in token):
|
||||
return deny
|
||||
status, data = self.transport.request('POST', self.endpoint,
|
||||
json.dumps(request, separators=(',', ':'), ensure_ascii=False).encode(),
|
||||
{'Content-Type': 'application/json', 'Authorization': 'Bearer ' + token})
|
||||
if (status != 200 or data.get('contract_version') != 'flex-auth.decision-record.v1'
|
||||
or data.get('request_id') != request['id'] or data.get('effect') != 'allow'
|
||||
or data.get('obligations', []) != [] or not isinstance(data.get('id'), str) or not data['id']):
|
||||
return deny
|
||||
binding, provenance = data['binding'], data['provenance']
|
||||
if (binding['submitted_request_digest'] != request_digest(request)
|
||||
or not DIGEST.fullmatch(binding['request_digest'])
|
||||
or binding['tenant'] != request['tenant'] or binding['action'] != action
|
||||
or binding.get('context', {}) != request.get('context', {})):
|
||||
return deny
|
||||
for key, fields in (('subject', ('id', 'type', 'tenant')),
|
||||
('resource', ('id', 'type', 'system', 'tenant'))):
|
||||
for name in fields:
|
||||
if binding[key].get(name) != request[key].get(name) or data[key].get(name) != binding[key].get(name):
|
||||
return deny
|
||||
if (provenance['policy_package'] != self.package or provenance['policy_version'] != self.version
|
||||
or provenance['policy_package_digest'] != self.digest
|
||||
or data['matched_policy_version'] != self.version
|
||||
or not DIGEST.fullmatch(provenance['registry_snapshot_digest'])
|
||||
or not provenance['evaluator'].startswith('flex-auth/')
|
||||
or not started - 30 <= timestamp(provenance['decision_time']) <= self.clock() + 30):
|
||||
return deny
|
||||
lifetime = data['lifetime']
|
||||
until = min(timestamp(lifetime['expires_at']), started + 30, actor.expires_at)
|
||||
if (lifetime['kind'] != 'ttl' or until <= self.clock()
|
||||
or timestamp(lifetime.get('not_before', provenance['decision_time'])) > self.clock()):
|
||||
return deny
|
||||
return {'effect': 'allow', 'id': data['id'], 'expires_at': until}
|
||||
except (OSError, KeyError, ValueError, TypeError, AttributeError):
|
||||
return deny
|
||||
129
scripts/alert_service.py
Normal file
129
scripts/alert_service.py
Normal file
|
|
@ -0,0 +1,129 @@
|
|||
#!/usr/bin/env python3
|
||||
"""Single-process private acknowledgment runtime, packaged by rapp-telemetry."""
|
||||
import argparse
|
||||
from http.cookies import SimpleCookie
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import secrets
|
||||
import signal
|
||||
import threading
|
||||
import time
|
||||
from urllib.parse import parse_qs
|
||||
|
||||
from alert_ack import Application, Store
|
||||
from alert_audit import AuditTransport
|
||||
from alert_identity import Login
|
||||
from alert_policy import Policy
|
||||
|
||||
|
||||
def credential(path):
|
||||
# Kubernetes projected credentials use symlinks; mount ownership is the
|
||||
# package trust boundary. Resolve once per call so rotations are picked up.
|
||||
with Path(path).open('rb') as source:
|
||||
raw = source.read(32769)
|
||||
value = raw.decode().strip()
|
||||
if not 32 <= len(value) <= 32768 or any(ord(c) < 33 or ord(c) > 126 for c in value):
|
||||
raise ValueError('invalid credential projection')
|
||||
return value
|
||||
|
||||
|
||||
class Router:
|
||||
def __init__(self, config, store, *, login=None, policy=None):
|
||||
self.config, self.store = config, store
|
||||
self.login = login or Login(config['issuer'], config['origin'])
|
||||
self.policy = policy or Policy(config['policy']['origin'],
|
||||
lambda: credential(config['policy']['token_file']), config['policy']['package'],
|
||||
config['policy']['version'], config['policy']['digest'])
|
||||
self.application = Application(store, config['origin'], credential(config['webhook_token_file']),
|
||||
self.actor, self.policy)
|
||||
self.last_drain = 0
|
||||
|
||||
@staticmethod
|
||||
def cookie(env, name):
|
||||
cookies = SimpleCookie()
|
||||
try:
|
||||
cookies.load(env.get('HTTP_COOKIE', ''))
|
||||
return cookies[name].value if name in cookies else ''
|
||||
except Exception:
|
||||
return ''
|
||||
|
||||
def actor(self, env):
|
||||
return self.login.session(self.cookie(env, '__Host-rtel-session'))
|
||||
|
||||
def __call__(self, env, start):
|
||||
def reply(status, text='', headers=()):
|
||||
start(status, [('Content-Type', 'text/plain; charset=utf-8'), ('Cache-Control', 'no-store'),
|
||||
('Referrer-Policy', 'no-referrer'), ('X-Content-Type-Options', 'nosniff'),
|
||||
('Content-Security-Policy', "default-src 'none'; frame-ancestors 'none'")] + list(headers))
|
||||
return [text.encode()]
|
||||
def cookie(name, value, age):
|
||||
return ('Set-Cookie', f'{name}={value}; Path=/; Secure; HttpOnly; SameSite=Lax; Max-Age={age}')
|
||||
path, method = env.get('PATH_INFO'), env.get('REQUEST_METHOD')
|
||||
try:
|
||||
if method == 'GET' and path in ('/healthz', '/readyz'):
|
||||
ready = path == '/healthz' or time.time() - self.last_drain < 90 and not self.store.audit_debt()
|
||||
return reply('200 OK' if ready else '503 Service Unavailable', 'ready' if ready else 'audit delivery pending')
|
||||
if path == '/ack/auth/callback' and method == 'GET':
|
||||
params = parse_qs(env.get('QUERY_STRING', ''), strict_parsing=True, max_num_fields=4)
|
||||
if any(len(v) != 1 for v in params.values()) or not {'state', 'code'} <= set(params):
|
||||
raise ValueError('invalid callback')
|
||||
sid, target = self.login.finish(params['state'][0], self.cookie(env, '__Host-rtel-login'), params['code'][0])
|
||||
return reply('303 See Other', headers=[('Location', target), cookie('__Host-rtel-session', sid, 900),
|
||||
cookie('__Host-rtel-login', '', 0)])
|
||||
if path == '/ack/logout' and method == 'POST':
|
||||
actor = self.actor(env)
|
||||
form = parse_qs(Application.body(env).decode(), max_num_fields=1, strict_parsing=True)
|
||||
if (not actor or env.get('HTTP_ORIGIN') != self.config['origin'] or set(form) != {'csrf'}
|
||||
or len(form['csrf']) != 1 or not secrets.compare_digest(actor.csrf, form['csrf'][0])):
|
||||
return reply('403 Forbidden', 'Invalid sign-out.')
|
||||
self.login.logout(self.cookie(env, '__Host-rtel-session'))
|
||||
return reply('200 OK', 'Signed out.', [cookie('__Host-rtel-session', '', 0)])
|
||||
if path == '/ack/alerts' and method == 'GET' and self.actor(env) is None:
|
||||
target = path + '?' + env.get('QUERY_STRING', '')
|
||||
url, browser = self.login.start(target)
|
||||
return reply('303 See Other', headers=[('Location', url), cookie('__Host-rtel-login', browser, 300)])
|
||||
# Refresh the webhook credential for projected rotation; never store it in SQLite.
|
||||
if path == '/webhook':
|
||||
self.application.webhook_token = credential(self.config['webhook_token_file'])
|
||||
return self.application(env, start)
|
||||
except (ValueError, OSError, KeyError, TypeError):
|
||||
return reply('503 Service Unavailable', 'Identity or service unavailable. Retry later.')
|
||||
|
||||
|
||||
def main():
|
||||
from waitress import serve
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument('--config', required=True, type=Path)
|
||||
args = parser.parse_args()
|
||||
os.umask(0o077)
|
||||
config = json.loads(args.config.read_text())
|
||||
Path(config['database']).parent.mkdir(mode=0o700, exist_ok=True)
|
||||
store = Store(config['database'])
|
||||
router = Router(config, store)
|
||||
transport = AuditTransport(config['audit']['origin'], lambda: credential(config['audit']['token_file']),
|
||||
allow_internal_http=True)
|
||||
stop = threading.Event()
|
||||
def drain():
|
||||
while not stop.is_set():
|
||||
try:
|
||||
store.drain(transport)
|
||||
router.last_drain = time.time()
|
||||
except (OSError, ValueError):
|
||||
router.last_drain = 0
|
||||
stop.wait(30)
|
||||
worker = threading.Thread(target=drain, daemon=True)
|
||||
worker.start()
|
||||
def shutdown(*args):
|
||||
raise SystemExit(0)
|
||||
signal.signal(signal.SIGTERM, shutdown)
|
||||
try:
|
||||
serve(router, host='0.0.0.0', port=8080, threads=4, max_request_body_size=32768,
|
||||
clear_untrusted_proxy_headers=True)
|
||||
finally:
|
||||
stop.set()
|
||||
worker.join(timeout=6)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
46
scripts/telemetry_http.py
Normal file
46
scripts/telemetry_http.py
Normal file
|
|
@ -0,0 +1,46 @@
|
|||
"""Bounded HTTP transport; endpoint configuration is never supplied by callers."""
|
||||
import json
|
||||
from http.client import HTTPException
|
||||
from urllib.error import HTTPError, URLError
|
||||
from urllib.parse import urlsplit
|
||||
from urllib.request import HTTPRedirectHandler, ProxyHandler, Request, build_opener
|
||||
|
||||
|
||||
def origin(value, internal=False):
|
||||
p = urlsplit(value)
|
||||
local = p.hostname in ('127.0.0.1', '::1') or (p.hostname or '').endswith(('.svc', '.svc.cluster.local'))
|
||||
if (p.scheme != 'https' and not (internal and local and p.scheme == 'http')
|
||||
or not p.hostname or p.username or p.password or p.path or p.query or p.fragment
|
||||
or any(c.isspace() for c in value)):
|
||||
raise ValueError('invalid fixed origin')
|
||||
return value
|
||||
|
||||
|
||||
class NoRedirect(HTTPRedirectHandler):
|
||||
def redirect_request(self, *args, **kwargs):
|
||||
return None
|
||||
|
||||
|
||||
class Transport:
|
||||
def __init__(self):
|
||||
self.opener = build_opener(ProxyHandler({}), NoRedirect())
|
||||
|
||||
def request(self, method, url, body=None, headers=None):
|
||||
try:
|
||||
try:
|
||||
response = self.opener.open(Request(url, data=body, headers=headers or {}, method=method), timeout=5)
|
||||
except HTTPError as error:
|
||||
response = error
|
||||
with response:
|
||||
status = response.code
|
||||
if status >= 300:
|
||||
return status, {}
|
||||
raw = response.read(262145)
|
||||
if len(raw) > 262144:
|
||||
raise ValueError()
|
||||
data = json.loads(raw)
|
||||
if not isinstance(data, dict):
|
||||
raise ValueError()
|
||||
return status, data
|
||||
except (OSError, URLError, HTTPException, ValueError, UnicodeError):
|
||||
raise OSError('upstream unavailable or invalid') from None
|
||||
Loading…
Add table
Add a link
Reference in a new issue