diff --git a/README.md b/README.md index b15f091..56b9b51 100644 --- a/README.md +++ b/README.md @@ -44,3 +44,8 @@ a second production monitoring plane. The package records private installation and restore proof; delivery acceptance remains blocked in RTEL-WP-0002-T04. The [Prometheus mapping](docs/prometheus-mapping.md) supplies a tested report exporter and failure/absence rules for package integration. + +The [email acknowledgment component](docs/alert-acknowledgment.md) records an +explicit Railiance admin confirmation with a durable audit-core outbox. It is +not activated. Identity/policy adapters are implemented; native package admission, +SMTP password custody and recipient proof remain under T04. diff --git a/contracts/alert-acknowledgment.json b/contracts/alert-acknowledgment.json new file mode 100644 index 0000000..05b5b2b --- /dev/null +++ b/contracts/alert-acknowledgment.json @@ -0,0 +1,43 @@ +{ + "schema": "railiance-telemetry.alert-acknowledgment.v1", + "status": "implementation-candidate-not-activated", + "workplan_task": "RTEL-WP-0002-T04", + "recipient": "bernd.worsch@gmail.com", + "role": "railiance-admin", + "role_assignment": { + "requested_by": "Bernd Worsch", + "requested_on": "2026-09-28", + "identity_owner": "net-kingdom/key-cape", + "directory_username_candidate": "tegwick", + "verified_issuer_subject": null, + "applied": true, + "email_is_authority": false, + "directory_group": "railiance-admins", + "signed_role_claim_verified": false, + "issuer_mapping_deployed": true + }, + "tenant": "tenant:platform", + "actions": [ + "read", + "acknowledge" + ], + "resource_type": "telemetry-alert", + "policy_owner": "flex-auth", + "browser_origin": "https://telemetry.coulomb.social", + "browser_path": "/ack/alerts", + "callback_candidate": "https://telemetry.coulomb.social/ack/auth/callback", + "client_id_candidate": "railiance-telemetry-admin", + "sender": "Alertmanager native SMTP email integration", + "receiver_name": "railiance-admin-email", + "audit_source": "railiance-telemetry", + "audit_class": "telemetry.alert.acknowledged", + "audit_evidence_kind_proposed": "load-bearing", + "audit_write_only": true, + "audit_registration_applied": false, + "acknowledgment_resolves_alert": false, + "acknowledgment_silences_alert": false, + "from_address": "platform@coulomb.social", + "from_mailbox_status": "created-by-founder-password-custody-pending", + "smtp_kv_path": "platform/workloads/railiance-telemetry/smtp", + "smtp_password_field": "SMTP_PASSWORD" +} diff --git a/docs/alert-acknowledgment.md b/docs/alert-acknowledgment.md new file mode 100644 index 0000000..f87503b --- /dev/null +++ b/docs/alert-acknowledgment.md @@ -0,0 +1,163 @@ +# Email receipt acknowledgment + +The founder authorized email to `bernd.worsch@gmail.com`, a Railiance admin role +filled by that user, controlled failure/absence drills and acknowledgment audit +records on September 28. This work remains under RTEL-WP-0002-T04; no new task +or workplan is needed. The machine-readable record is +`contracts/alert-acknowledgment.json`. + +## User flow + +Alertmanager sends an email containing a link to one alert occurrence. The link +opens a review page at `https://telemetry.coulomb.social/ack/alerts`. Sign-in uses +the existing estate identity provider. A Railiance admin presses **Acknowledge +receipt**; the page confirms local recording and separately reports whether +audit-core has accepted its audit event. Clicking again returns the same record. +An email scanner or preview GET never acknowledges anything. Acknowledgment +does not resolve the fault, silence Alertmanager or suppress its repeat emails. + +The link contains the Alertmanager fingerprint and original start timestamp, +not a bearer credential. A fresh occurrence needs a fresh acknowledgment, even +if its labels/fingerprint are identical. The timestamp template preserves +RFC3339 nanoseconds. The corresponding webhook must arrive first; an early +click reports that the alert has not arrived and can be retried. + +## Implemented component + +`scripts/alert_ack.py` supplies a WSGI component with authenticated `/webhook` +and protected GET/POST `/ack/alerts` routes. It requires server-side identity +and authorization adapters at construction. `alert_service.py` supplies the +Waitress runtime, `alert_identity.py` supplies OIDC code/PKCE sessions, and +`alert_policy.py` validates native Flex Auth decisions. There is no anonymous +mode, trusted-email/header fallback or default allow decision. + +The private SQLite store commits the first human acknowledgment and exact +audit envelope in one transaction. Ordinary updates/deletes of acknowledgments, +alerts and audit payloads are refused. This does not protect against a database +administrator. Duplicate webhook calls and clicks preserve the original fact. +Only bounded alert name, fingerprint, start time and a digest of labels persist; +annotations and arbitrary alert details are not copied into email or audit. + +`scripts/alert_audit.py` supplies a bounded, redirect-refusing sender for +`POST /v1/events`, with stable `Idempotency-Key`. Credential provision remains +external. Exact `202 accepted` or `200 duplicate` and `audit:` are +required for delivery completion. Lost replies retain the event for replay; +schema/auth/conflict refusals retain it as blocked. The package must expose +pending/blocked debt and explicitly requeue a blocked row after repair. It must +schedule bounded drains and per-class reconciliation/heartbeat before claiming +live audit operation. The runtime drains every 30 seconds and reports readiness false while audit +debt remains. Per-class reconciliation/heartbeat and independent readback remain +activation gates. No retention deletion is installed. + +## Identity and authorization binding still required + +The email address is the notification destination. NetKingdom records associate +it with directory username `tegwick`; the explicit `railiance-admins` membership was applied using the native +identity-provisioner, preserving other memberships. KeyCape commit `1164f65` +was built, published and deployed successfully. A real login must still verify +the signed `(issuer, subject)` and role claim. Do not grant on an email match, domain match or client-supplied +header. Do not automatically equate `net-kingdom-admins` with `railiance-admin`. +The requested role is consumed here for `read` and `acknowledge` on platform +telemetry alerts; this integration grants no unrelated estate privileges. + +The package must bind a registered OIDC Authorization Code + PKCE session with +verified signature, issuer, audience, nonce, expiry, human provenance and MFA. +Keep tokens server-side; use Secure/HttpOnly/SameSite cookies and logout/expiry. +Populate `Actor` only from that verified session, with a random session-bound +CSRF value. The component separately enforces human/platform/role constraints, +session expiry, exact form Origin and CSRF. These are additional restrictions, +not a substitute for the policy decision. + +The Flex Auth adapter must request a fresh decision for every read or acknowledge, +with system `railiance-telemetry`, resource type `telemetry-alert`, resource +`alert:`, tenant `tenant:platform` and verified subject facts. +It must validate the native decision contract, exact actor/resource/action +binding, submitted request digest, admitted package/version/digest, lifetime +and supported obligations before returning the component's bounded receipt. +Both adapters are implemented. Signed RSA issuer fixtures exercise the browser +flow; the native Flex Auth evaluator validates exact request digests and rejects +wrong identities, missing roles/groups and stale MFA. These tests use synthetic +credentials and do not prove live authentication. The package must still admit +the OIDC client and enforced Kubernetes TokenReview caller. + +## Concrete activation packet for existing owners + +| Owner | Required binding | +| --- | --- | +| NetKingdom / key-cape | Register `railiance-telemetry-admin` and exact proposed callback `/ack/auth/callback`; verify Bernd's subject; apply and verify `railiance-admin` membership and MFA claims. | +| flex-auth | Admit workload caller, telemetry-alert resource and read/acknowledge package/assignments; return pinned package/version/digest and positive/negative fixtures. | +| railiance-platform | Provision dedicated SMTP, webhook and audit-sender custody through approved lanes; do not extract or copy email-connect's invitation credential into this app. | +| audit-core | Register source `railiance-telemetry`, exact tenant `tenant:platform`, write-only, `secret_policy=redact`, proposed load-bearing class; admit ingress and independent readback. | +| rapp-telemetry | Package session/PDP adapters, audit drain/debt monitoring and private persistent state; route `/ack/` separately from Grafana, keep `/webhook` private; integrate SMTP/template/webhook and verify restart/restore. | + +Use Alertmanager's native email integration with recipient +`bernd.worsch@gmail.com`, From `platform@coulomb.social`, template `railiance.alert.email` from +`templates/alert-email.tmpl` and receiver name `railiance-admin-email`. +The same receiver's webhook posts to the private acknowledgment service using +a dedicated credential file. Route only alerts with `owner=railiance-telemetry` +until other inventories are reviewed; the receiver rejects other scopes. +The package overlay has this shape (the private service and mounted credential +paths are candidates, not deployed resources): + +```yaml +templates: + - /etc/alertmanager/templates/alert-email.tmpl +receivers: + - name: railiance-admin-email + email_configs: + - to: bernd.worsch@gmail.com + from: platform@coulomb.social + require_tls: true + text: '{{ template "railiance.alert.email" . }}' + send_resolved: true + webhook_configs: + - url: http://telemetry-ack.telemetry.svc.cluster.local:8080/webhook + send_resolved: true + http_config: + authorization: + type: Bearer + credentials_file: /etc/alertmanager/secrets/telemetry-ack/webhook-token +``` + +The overlay must retain existing routes and supply approved global SMTP +smarthost/from/auth settings. Do not enable it with an unimplemented browser +adapter or assume the credential file/service already exists. +Use SMTP STARTTLS and mounted credential files through the package-owned secret +references. The existing email-connect lane proves IONOS is available but does +not admit a new SMTP consumer. No existing Secret was read during this work. + +The first live drill must retain email identifiers, alert occurrences, Bernd's +explicit acknowledgment events, exact audit references and independent readback. +Run a controlled failure and stopped-producer case, plus unauthorized account, +expired login, scanner GET, duplicate click and audit-unavailable checks. +Outside-node monitoring and recurring backups remain the existing T04 gates. + +## Verification + +```bash +RTEL_AUDIT_CORE_SOURCE=/home/worsch/audit-core python3 -m unittest discover -s tests -v +amtool template render --template.glob=templates/alert-email.tmpl \ + --template.text='{{ template "railiance.alert.email" . }}' +``` + +The optional receiver test uses actual audit-core ingestion and SQLite with +synthetic credentials: first accepted, lost reply, then duplicate after restart. +It is not production custody proof. The template follows the upstream +[notification data contract](https://prometheus.io/docs/alerting/latest/notifications/). +The founder created `platform@coulomb.social`. The reviewed platform helper +creates `platform/workloads/railiance-telemetry/smtp` without a password through +attended OpenBao login. The founder then adds `SMTP_PASSWORD` as a new version, +preserving existing public SMTP fields. +Live activation remains blocked on password provisioning, dedicated credential custody, +client/caller admission, runtime rollout and recipient/readback drills. +No email or production acknowledgment is claimed. + +Runtime validation (hash-pinned dependencies in requirements-runtime.lock): + +```bash +RTEL_FLEX_AUTH_BINARY=/tmp/rtel-flex-auth /tmp/rtel-ack-venv/bin/python -m unittest discover -s tests_runtime -v +``` + +The candidate workload, container recipe and runtime settings are owned by +`rapp-telemetry/acknowledgment`. The policy source and exact client registration +are in `integration/`; they are not active registrations. diff --git a/history/2026-09-28-alert-acknowledgment.md b/history/2026-09-28-alert-acknowledgment.md new file mode 100644 index 0000000..8f33be6 --- /dev/null +++ b/history/2026-09-28-alert-acknowledgment.md @@ -0,0 +1,45 @@ +# Email acknowledgment implementation — 2026-09-28 + +Existing task RTEL-WP-0002-T04 holds this work; no task/workplan added. +User decision f149e316-4ef4-4855-8453-bc9cdc938aad approves email to +bernd.worsch@gmail.com, explicit receipt confirmation, Railiance admin role, +audit-core evidence and controlled drills. Subsequent direction selected From +platform@coulomb.social and confirmed the mailbox needs setup. + +Implemented WSGI acknowledgment, immutable first receipt, atomic SQLite audit +outbox, bounded audit transport, OIDC code/PKCE sessions, native Flex Auth +binding/digest/lifetime checks, Waitress runtime and background audit draining. +GET is inert; POST requires verified human/platform/admin identity, a fresh PDP +allow, Origin and CSRF. Email links identify occurrences, not bearer credentials. + +Validation: 33 core tests with actual audit-core 3e42ca8 ingestion, including +accepted/lost-reply/reopened-store/duplicate; seven runtime tests with signed RSA +issuer fixtures, complete browser flow and the actual Flex Auth evaluator. +Synthetic identities/credentials do not prove production login or custody. +Native Alertmanager 0.28.1 template render passed. Hash-pinned runtime container +built; network-isolated read-only container returned HTTP 200 health and exited +cleanly on SIGTERM. Package candidates reside in rapp-telemetry/acknowledgment. +Policy and browser client registration candidates reside in integration/. + +Native changes: identity-provisioner verified tegwick's requested email and added +only railiance-admins, preserving other memberships. KeyCape 1164f65 maps that +explicit group to railiance-admin without granting platform-operator. Full Go +suite passed. Published immutable image: +sha256:6f79a2af1c695d39480173fad013facd84d21e7732860366af519302a2c496b8. +NetKingdom manifest server dry-run passed; diff changed only the image (plus +metadata). Deployment rolled out successfully. Signed role remains unverified +until a real login. No Kubernetes Secret values were read or printed. + +Remaining gates: mailbox setup; scoped SMTP/webhook/audit custody and receiver +registration; safe OIDC client registration and enforced policy caller admission; +application rollout; native login and actual failure/absence emails, Bernd's +acknowledgments and independent audit readback; outside-node watchdog and +recurring backups. Existing owner client helpers read complete Kubernetes +Secrets and cannot be used under the current environment orientation's rule. +No live email or acknowledgment is claimed. Workplan stays blocked, T04 wait. + +Subsequently the founder created platform@coulomb.social and requested an OpenBao +entry, with the password to be added by the founder as a new version. Platform +helper scripts/telemetry_smtp_entry.py is silent, CAS=0, never reads credential +values, and preserves any existing version. Four unit tests pass. Attended +founder OIDC/MFA execution is requested; no native KV creation is yet claimed. diff --git a/integration/fixtures.json b/integration/fixtures.json new file mode 100644 index 0000000..f66a665 --- /dev/null +++ b/integration/fixtures.json @@ -0,0 +1,12 @@ +[ + { + "id": "unknown-request-denied", + "request": { + "id": "unknown", + "subject": {"id": "unknown", "type": "human"}, + "action": "acknowledge", + "resource": {"id": "alert:unknown", "type": "telemetry-alert", "system": "railiance-telemetry"} + }, + "expect": {"effect": "deny"} + } +] diff --git a/integration/keycape-client.json b/integration/keycape-client.json new file mode 100644 index 0000000..1eeaf35 --- /dev/null +++ b/integration/keycape-client.json @@ -0,0 +1,20 @@ +{ + "clientId": "railiance-telemetry-admin", + "displayName": "Railiance Alert Acknowledgment", + "audience": "railiance-telemetry", + "redirectUris": [ + "https://telemetry.coulomb.social/ack/auth/callback" + ], + "allowedScopes": [ + "openid", + "profile", + "email", + "telemetry:read", + "telemetry:acknowledge" + ], + "grantTypes": [ + "authorization_code" + ], + "clientType": "public", + "mfaRequired": true +} diff --git a/integration/railiance-admin-directory.py b/integration/railiance-admin-directory.py new file mode 100644 index 0000000..bcb843f --- /dev/null +++ b/integration/railiance-admin-directory.py @@ -0,0 +1,42 @@ +"""Run inside the existing identity-provisioner pod; never print credentials. + +Default is read-only. --apply adds only the approved user's named group. +""" +import json +import os +import sys +from provisioner import LLDAPProvisioner + + +def main(): + apply = sys.argv[1:] == ['--apply'] + if sys.argv[1:] not in ([], ['--apply']): + raise ValueError() + if os.environ['LLDAP_URL'] != 'http://lldap.sso.svc.cluster.local:17170': + raise ValueError() + client = LLDAPProvisioner(base_url=os.environ['LLDAP_URL'], admin_password=os.environ['LLDAP_ADMIN_PASSWORD']) + token = client._login() + user = client._user(token, 'tegwick') + if not user or user['id'] != 'tegwick' or user['email'].lower() != 'bernd.worsch@gmail.com': + raise ValueError() + before = {g['displayName'] for g in user['groups']} + if apply and 'railiance-admins' not in before: + groups = client._gql(token, 'query { groups { id displayName } }', {})['groups'] + group = client._ensure_group(token, groups, 'railiance-admins') + client._add_group(token, 'tegwick', group) + after = {g['displayName'] for g in client._user(token, 'tegwick')['groups']} + if not before <= after or after - before - {'railiance-admins'}: + raise ValueError() + if apply and 'railiance-admins' not in after: + raise ValueError() + print(json.dumps({'mode': 'apply' if apply else 'inspect', 'directory_user': 'tegwick', + 'email_matches_requested_recipient': True, 'group': 'railiance-admins', + 'member': 'railiance-admins' in after, 'other_memberships_preserved': True, + 'signed_role_claim_verified': False})) + + +try: + main() +except Exception: + print('{"status":"directory-operation-refused"}') + sys.exit(1) diff --git a/integration/registry.json b/integration/registry.json new file mode 100644 index 0000000..cde5eaa --- /dev/null +++ b/integration/registry.json @@ -0,0 +1,11 @@ +{ + "subjects": [{ + "id": "uid=tegwick,ou=people,dc=netkingdom,dc=local", + "type": "Human", + "tenant": "tenant:platform", + "display_name": "Bernd Worsch", + "organization_relation": "ServiceProvider", + "roles": ["railiance-admin"] + }], + "resources": [] +} diff --git a/integration/telemetry-policy.md b/integration/telemetry-policy.md new file mode 100644 index 0000000..47de2df --- /dev/null +++ b/integration/telemetry-policy.md @@ -0,0 +1,71 @@ +--- +id: railiance-telemetry.alert-acknowledgment +name: Railiance admin alert receipt acknowledgment +namespace: railiance-telemetry:telemetry-alert +version: v1 +status: ready +package: flexauth.railiance_telemetry.alert_acknowledgment +allow_ttl: 30s +actions: [read, acknowledge] +owner: flex-auth +fixtures: [fixtures.json] +caring: + profile: caring-0.4.0-rc2 + enforce: false +activation: + mode: local +--- + +# Requested telemetry admin mandate + +Bernd Worsch authorized the named Railiance admin role and receipt actions on +September 28 under RTEL-WP-0002-T04. Native service caller admission and directory +membership remain required. This policy permits no alert silencing, resolution, +configuration change or unrelated estate operation. The caller must validate +the signed KeyCape session and supply its unchanged identity/assurance facts. + +```rego +import rego.v1 + +decision := {"effect": "allow", "reason": "railiance_admin_alert_receipt"} if { + input.tenant == "tenant:platform" + input.subject.type == "human" + input.subject.tenant == "tenant:platform" + is_string(input.subject.id) + input.subject.id != "" + input.subject.id == "uid=tegwick,ou=people,dc=netkingdom,dc=local" + "railiance-admin" in input.subject.attributes.roles + authentication := input.context.authentication + authentication.issuer == "https://kc.coulomb.social" + authentication.principal_type_source == "authentication-derived" + authentication.tenant_source == "directory-asserted" + "railiance-admin" in authentication.roles + "railiance-admins" in authentication.groups + assurance := authentication.assurance + assurance.level == "aal2" + assurance.mfa == true + assurance.source == "key-cape" + assurance.methods == ["pwd", "otp"] + is_number(assurance.at) + age := time.now_ns() / 1000000000 - assurance.at + age >= -30 + age <= 900 + input.resource.system == "railiance-telemetry" + input.resource.type == "telemetry-alert" + input.resource.tenant == "tenant:platform" + regex.match("^alert:[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$", input.resource.id) + input.action in {"read", "acknowledge"} +} else := {"effect": "deny", "reason": "telemetry_identity_or_scope_refused"} if { + true +} +``` + +```rego test +package flexauth.railiance_telemetry.alert_acknowledgment_test +import rego.v1 +import data.flexauth.railiance_telemetry.alert_acknowledgment + +test_unknown_request_denied if { + alert_acknowledgment.decision.effect == "deny" with input as {} +} +``` diff --git a/requirements-runtime.lock b/requirements-runtime.lock new file mode 100644 index 0000000..a30c230 --- /dev/null +++ b/requirements-runtime.lock @@ -0,0 +1,164 @@ +# This file was autogenerated by uv via the following command: +# uv pip compile requirements-runtime.txt --generate-hashes -o requirements-runtime.lock +cffi==2.1.1 \ + --hash=sha256:046bfc24911b37851ee1b51aab8bffe713d89c68c6a057b09484ce9fd5f69b4e \ + --hash=sha256:06c72bb76605a4b0cd0aad6930b69d4baf7dd5d806cfc409b824191099700e66 \ + --hash=sha256:0beceaabe56af686895136a2de78db54ecd8e4046b236b8fd6d6cb61389e9bf2 \ + --hash=sha256:154852545011f779917b11c78db2358d095da62a9a172b78ad0a583ee5adc0d0 \ + --hash=sha256:194cffa889098ced9976c3fc6340305e43f6303657d298da55366907c05c22d6 \ + --hash=sha256:19ee6127ee34de7d83ce3d371ebc5ed91addbdcc39f9ab15ce4eb35a4e534971 \ + --hash=sha256:1a18a57b58cfb21fc28d72e876acf10eaed67a1ed96226f92af4df681d571c4c \ + --hash=sha256:1aa5645c30469b09530c4ebca77ebf8f17618293c58f8549cb1a543a50236e7d \ + --hash=sha256:1dea0e4d7d4f11f619fe8c1d76caf49e24405b4b5743c0e3be16a500ecd930c9 \ + --hash=sha256:208f941bb9d18e768138677f0a6d2ce01f590df56043dda1df1535ac57c88517 \ + --hash=sha256:210019b6c7cf07f081b4c54635c8cf744377001350e29cc0f81c4377b4797735 \ + --hash=sha256:246fa40ce8645a614ff682e0b70f37134e460eaf93a775e0cbe3cca585a67a80 \ + --hash=sha256:25792eac27877609e7bb06d42ff88278a6624fff2ba9bbb523c09616b117e80f \ + --hash=sha256:27350daa11d4f10c540e6e89dada4c54feb7256ad03e9a4dc075ebad7ba360d1 \ + --hash=sha256:28907ab9bfb6aa13184cfc17c6b8e1023c5ab6fd7076d8c20a35e59fe04f8f29 \ + --hash=sha256:2ae64be792b8966f2c69538199728b290e34726562896df1e5dc8ffd8d8188e8 \ + --hash=sha256:31348097ff5bbe827ccc41795d4dd099d9f0625e7def00ee653c137a490c2a6c \ + --hash=sha256:3143d81e29e1e20a9ce10901ec369012947876596f75a222235965f2b7ae832e \ + --hash=sha256:3222ba5d678f80a030e6afbcc33dc1ae5cb45facabb61cee2c7016b8432fde48 \ + --hash=sha256:3311ed60d36f83378794e1009ac6258bafbf81f7888b4caa7b35a521e3f95813 \ + --hash=sha256:334644fbac4eff73d985a17a91226df55d0f394160c4cfb880e084c8f7161cac \ + --hash=sha256:34e261f78cb6ceaaa36f42f2613f4380d94d9c759a9c73c769ee6e0247364632 \ + --hash=sha256:363e05fa78e15116c3c32c210ee36884fd6b9afa6d440e47112c3bd511d64cb6 \ + --hash=sha256:398aff33cee2767e3e781d2554c54bd0dff386bb437581e0d8011fde1a942ec1 \ + --hash=sha256:3d22a20b1fb1632cc72c22f95f7b0d2961c3e1c235f245ba4c606c4771035659 \ + --hash=sha256:42a494cee34437f05546455144f2b5d9ac09b1face62bcfce597d2e521066688 \ + --hash=sha256:42e2f76b9455f5a9a844f770bf3e200ed3da0e15f5df3db9c31fe80b04b3d004 \ + --hash=sha256:42f6930c31dc7f50732c9ae793c2786c7b6b044195967bbdde40bb9be81c4cc0 \ + --hash=sha256:456a61fa52d579ebf9df2e9552ead5129855dbaff6c1e5a9b1bc408809bdc062 \ + --hash=sha256:471cee653ae88de62096552e6d24ccb4a5adb8c8c9f10b5054d0122c15bf2779 \ + --hash=sha256:49cbc70e6542d4ccccb936558d1064a8012541e78f821f955cff24e357776c94 \ + --hash=sha256:4a7c934f7360e8cd64fe9efadcbd10c7c6364f531e432b9a4bf5ccbc9e0e8b50 \ + --hash=sha256:4be96343e422f2dfcd12ab5c9f5aebe03f82f737c6bffeca6830b3875cb44aab \ + --hash=sha256:4f42141fc14250de6dde5ee7ea4432be017252d91f19c5ad043c084cea629cac \ + --hash=sha256:507a24c282e0f42f8ed737cf048572cbf580468da5555764a8331735e9c736b6 \ + --hash=sha256:51b31d1c98274844cfd7838ce00bfc27c7423a4dc00fc0772fc3331c2cc90676 \ + --hash=sha256:58acb8ab8e295e6c5ea12f888cbb13cf21511ef2a3303a23f4325c29d17fe5c1 \ + --hash=sha256:5a59cc1c4442bc3d5c703bf720b51138d0bfc173618807c9ee2490a7541dd3d9 \ + --hash=sha256:5bb4e7ea95dcd6a014a6fef62e62467d67d8e582326443f3d68e71d6320a9fcf \ + --hash=sha256:5c58fe613dc5e5336357eff555824a314d8e43282600435c8d1cb6a7a2fedd13 \ + --hash=sha256:5e7cecbaadb83884793e05828cee59b210b24583b9c7425d0ba6a754fe22eb4e \ + --hash=sha256:616f097f2fe415bc92a247f02e11f634e1f9e9a83d327e3c915c15089c87869e \ + --hash=sha256:63bbfd5ded17c4840ac07cd8f1c21ba9d9708141f840b324f422f41b207e3973 \ + --hash=sha256:64faea20f4e2613363a1a9b9c7dd73058f3ecd00133a511e72ad7c511658f527 \ + --hash=sha256:661c298b4821edebead0c91edd2b00374d67ad7c5a1f7a91d4442633b79d6a72 \ + --hash=sha256:68e62fe11f30d5ca8289242866f0a5291402d8529ca2178ab8afc5c9694ae890 \ + --hash=sha256:6a8dddef476fab96d066d578fc88526767b836ab5ab21754e1d5bf3879c31c7c \ + --hash=sha256:6e192623c49c94421616a5778fba35cf0d5a8d000650c1967ef4448ee5cdd990 \ + --hash=sha256:7225e4514edb64eb6740324353e0da0711954fd8d7da4576755b1c6e09b697cd \ + --hash=sha256:75f80557d1389eddbd0de2681f6a390a0c5338c31ddaa821381c203fc3fd50d9 \ + --hash=sha256:770de9db11e84213beec501cfcaa013b019820ca881e03344dea5844f7876d94 \ + --hash=sha256:7750c6449dff7864bb9bb27ddfb0267756189201a3afc911d82b3caacd70dfc3 \ + --hash=sha256:7bde5e4cc5c10140859842b9d383af292b22639a4dffb725314baf45968cef80 \ + --hash=sha256:7ce713ace7c0e4520535b42b77eaa742c16dab813978064913e5a3cf82973b41 \ + --hash=sha256:7da0c5eff80f0197f3b3d1232ec5a682a9325f4ae9016a78f5f5ca35f9ced1f5 \ + --hash=sha256:7dbb61fe3a7699468030f71bbe5f8a0e326a151daa91beb11a6fc1f980c55e1c \ + --hash=sha256:811bd1e21d32de12efca32393a0ab3f5133b54fce9bd44b8bd77ab07da14bf6a \ + --hash=sha256:8ef53b2de9bcb9197d31854256575d59dbac0cba72ac627bb291ef5eceb74be4 \ + --hash=sha256:937c0052c05a31ca1daf18de3158eed4dbfcb9cc107adbea227728d647be701e \ + --hash=sha256:9d2055050ea716bd38b7f7f1579c275386646b4894c155a3e2f3cd62ed41b7c6 \ + --hash=sha256:9f8d177621de5cb38ee3e731eda45d421db093ec0739f46a5594babda7987a98 \ + --hash=sha256:a2d7755bef5a12ed488f4ef1f1b69ee9191d7396083b755a5d2295f6edb4768b \ + --hash=sha256:a48d62ab9d6f4f98c983223a547af44be6ca3691074c31cecced6facd3ba2dc1 \ + --hash=sha256:a4f00aa42f75d6e4595e8866e748cc1705adc0cddfeb2ca86d0d03993d63ba03 \ + --hash=sha256:a6e721d4b0e45d5b65e87534470e67b18dcd092c83f68fba09f152b9cbc061af \ + --hash=sha256:a730a083190634c65cca36ba5f489531576ebd79bcd5c8e172130f6453127231 \ + --hash=sha256:a931079504ecc49efed7744c476a5c343a92fabf66dec2db95edb1b2fdc770e2 \ + --hash=sha256:aa9511c62d14da7aacc9b4bf51f3f697a621e83b2d6919008243c3aad168eea3 \ + --hash=sha256:ab36d55f9ed2d067327667c2fea18dda018eb628dd6347aa01dda6cf1f5d3836 \ + --hash=sha256:ad2c86c495b899d862ea0f4b42891b8713a3bd45dd4105c7fd51c2a72f39f3a5 \ + --hash=sha256:aeae0e330c9f6acd681f647d46cefd30c29f93e3392882e792e82080c9691399 \ + --hash=sha256:b0431303acaea1089ad4b3e9ce4e6518193def1118d4073ca848635ee4ea2e96 \ + --hash=sha256:b5bdfd1c873d4e093aabc0ca84c4ca6dbc4f752afb5c86f146d9742580c9da2e \ + --hash=sha256:baed1e86cc735622097354b9d1281406caf42ff42a886d29faa8e8d1630333be \ + --hash=sha256:c1453022f490d2459a11819d83ad1d586e9ff65a12ac3e705ffebd46d3685dcf \ + --hash=sha256:c26608d2222fb1e94487e4a387d85f13eb55d5ed725cb25a0c589ac4ee60e7bc \ + --hash=sha256:c7659f22557c5a0bc4855cd635f55edec690cc008a40768527762cb9fb263455 \ + --hash=sha256:c8c69575568085ba0b1b10c0249d779a214aea6f6522e949a0fc9fb0fcb449d0 \ + --hash=sha256:c8d2c9fd1f2d16f780d15127abb050d13d1a76c03a4bd87d7e4980e45e511e12 \ + --hash=sha256:ca82be1a1d406ecfe1d25dc16cb33488e5a16bf4438c9fb590484ea29d92478b \ + --hash=sha256:cc572dace3f60ef98d7b12ff411d20f5362feb31a0439eab0085bbfd349982d7 \ + --hash=sha256:d18e5ac0f2f03f4f518d3e23db0f0cad7faa1da8620e9c09461d443bbf6e6692 \ + --hash=sha256:d28630f5854ab07ab1fd4aba756de52326c82e6be15d414b12793f1975048b54 \ + --hash=sha256:d9c275eaacd24aa73f94ffd6de08fc3f932424d8b6c376f4bed7cde376fe7bc3 \ + --hash=sha256:da0e573f9f97159390c89d9f1a9e41908b66d408cc5b58d08cf3847d844c531b \ + --hash=sha256:dd31f52ea1086513bb9df30f8fcee9b8918323ae067a3d5b78bc826a000712be \ + --hash=sha256:dddad92b554513a31f272570678ba307fb9f618f05e3d4a5eacafff9eae03e1d \ + --hash=sha256:df423d40ee8654634421812bc3b196da3f9bd7d32929da813f8394c4348a5358 \ + --hash=sha256:df913725b79db7bcf03448f36b7bf8815363417d5b58deecf9305e3e30f0f21a \ + --hash=sha256:e0bcb7e0f677f543555d2adff3bf19c05f66cdb4796e5ff602442ab2fe3c4ef7 \ + --hash=sha256:e2d65b31f36619cda3999b78b2aa9632e76b78448e7a56fc4240824200e7c4fc \ + --hash=sha256:e6e8cff14d6fb0be70a09c0bdc58096f501952d04624ebf867e0e56da2df8960 \ + --hash=sha256:f16c709686a78c727bbbf059f92b0bf41c6fc60deec706d2dc19f529175a6125 \ + --hash=sha256:f24fb43132a4c6b4cb4eb029492919b2db645be6808d738f244fd146c03c32cb \ + --hash=sha256:f53e442b08449d42821fa4a4fba000095af9f62742a500f978a9f557ec44339a \ + --hash=sha256:f5cfbc5fe74540d335175b656c725d74d90e3730c626d92575eea35029d9afaa \ + --hash=sha256:f81b3b8f3d4e343550fa4baa0e479bba9f2d29ce9c2e9b51d1ce1718d7442fcf \ + --hash=sha256:f8ec5e643a9a937f64e1999eb9f75d072263751912dc5cd06d3c85f8f44be7c3 \ + --hash=sha256:fb92203a88b3d3053034db775110081c49d28be6551923805e039924093761e4 \ + --hash=sha256:fcd22650c908d7b7da162bbfaab594a1227a15d1643a98c68b122ac642fa2264 + # via cryptography +cryptography==50.0.1 \ + --hash=sha256:01f41478cf33fc605a6a089cd56d28b45c6c0b45a1928b61797f2621a04bac71 \ + --hash=sha256:05ba322c4da95b262a212c345af888ef2c37c88c0509756ea00a0e6d68850f23 \ + --hash=sha256:16c5ecd954b3330ebfb6605eca4fd952da8bef376551d5cc264534e3770a9ee6 \ + --hash=sha256:2a93d05e34d5f67fba6f891fe85d929999baa7195e853923ea6d7576c9e68c5e \ + --hash=sha256:2b34d76a652ea2b6faf777c35df230c5637842cd904e04f16230c3f9f03e4361 \ + --hash=sha256:2ebbfb0f1fed745e91796e3e1080a1440423fdae8ece1b995a1d80883a409054 \ + --hash=sha256:30a125032e5642a21ff816e021152bd4e7e94f03eff3f4b7fca41cd22bc3110f \ + --hash=sha256:330fbb252391c596f1ae42c5754449dc924e6ad012dca8efe0d703f9f2d12ec6 \ + --hash=sha256:359e62deae718bce96170e223fdcb6357e4fbd3bb7a3a75f4430763532560e49 \ + --hash=sha256:407fe2b6db00939c05c0e945e9914238f2f0a430974839429dafc82b1ee6bee5 \ + --hash=sha256:42be3bb70596b3abe4ac097b75be223e8b3ab614a0e5de068e3dcc54d71d6149 \ + --hash=sha256:4c4188f7c0cf655be5c06342b817ed0f9595b69ffa2b12026e5353eed29dea88 \ + --hash=sha256:51593d180cf6d179bde5c5d065bed81386b1f381656ae7d042b7ffc87a9895ad \ + --hash=sha256:51afcfceb15597cf2635068e4ac9a56b2abde622edde17f37d85fd7b5306497a \ + --hash=sha256:53e279950892dc102c6b4e52af03ae5ea92fac572a1ddab78ca73a997f62b69f \ + --hash=sha256:55d16b1ef3ee0958d893a977b19777887e546c9954ea81b200c3301a864013f2 \ + --hash=sha256:5dd9bda1c12b4162f6ff568eeb5e0ff956c28d14406e875cfe8a63a2d414ff20 \ + --hash=sha256:5fe002589592ed749ce77fe0695fcbd3500dd61d7d6db5858a7544c612fa8e45 \ + --hash=sha256:5fe939deeb161024a6be98229c953b6591fef1f41214497a78fe793a244c017f \ + --hash=sha256:693c99b49bd37d0d096e4334c10232c77248c415b98d35236094cdf96d57258b \ + --hash=sha256:76de83fbd91ac49c0feaaa983d0748fd7a53176afac5fb3bf7478d244f0eb527 \ + --hash=sha256:79bf008d1f9af6071c797ad133e39915dfee7614f18f18f4db9072eb715064a3 \ + --hash=sha256:804728ce710890870f3aaa344b2e161172d258d768ac139d02cfd9092d0d94e6 \ + --hash=sha256:8921d58f426793c5f1b47f0b59575780de9a095214958d0eb37d909593db8367 \ + --hash=sha256:8df2de9102026855887e4587084f6eabd80ed0f345b8ad8a7ac27ab9bf4723e0 \ + --hash=sha256:9cb3cb952cf5a8abd50c782a98a89d71699715e802fe349704b47f2425b42a94 \ + --hash=sha256:9dde0a357190eb3b1da1bb9ab750e9c85cba82ca5977aa0836cbb94e92611239 \ + --hash=sha256:9ebcdd5519be9b652a46f507817a74591774fc3d6923ac364e4dfa64e36b291b \ + --hash=sha256:a0b1a59e3a089064a0ec309e9428c8e3ae4e161419d20ac33600767e83fc658a \ + --hash=sha256:a255449073358275b64b67d3f595f268bbef70e72b6edb65e0c70c735bf739c9 \ + --hash=sha256:a8f40ea47330e71b594a7e246898f93177c259490c63183dbaf9e571d71ed9a5 \ + --hash=sha256:ac02b07824d4d1001bd4367599f839c19cb171924c796e52c23508ac14c2c0cc \ + --hash=sha256:aed8db4f6d71c51efb89530e12d9464e7bf2923d46c3205dc794a2a93f8c0648 \ + --hash=sha256:b8f852c65863251b9e3a1b8c150ce21e59b522dbb6a7d4bc80e680d38388e986 \ + --hash=sha256:be224a65493ec5b74a158ff22a5522ce4a5ca1e543c647a3a4730d4a09e5f959 \ + --hash=sha256:ca83d00d9e69cd5eb63f2e69c3a5a59e0cecae5ae14c6ae0b35830fe3b37bad0 \ + --hash=sha256:cbf74a81765ee67413503ca6e26dcc4f6f5a519822436cc0a1b97aab6c1b8a17 \ + --hash=sha256:d63ae8f6481fec907ac0f588eee8a90aefde112c633131fe540e5711ddbb5a4e \ + --hash=sha256:e22dfed744bd4002e909464cb23d2f0b05c6f3113a79ef2e9864a53db737c733 \ + --hash=sha256:e2ca8fd1b6b4b82a1c4cb02841d0837e3c12336c2e24b520ab8ab3b969733d8f \ + --hash=sha256:e74591e283fe6eb956416c929eb58262a719fe0311fd9054c62c3350ed8760d8 \ + --hash=sha256:f74455bb086a85d5e81246412602aaa97ed095e504cd40dd261ef50be42205bf \ + --hash=sha256:fb4b9672d389c738b175c4166e78310f8a70358886aacd9173ee03a85ffdc671 \ + --hash=sha256:fc3ed7ebd2a8c96f5b166de0ab9b624996bef3b07bbeb19364dfb78222c22c80 \ + --hash=sha256:fd3718b960d0b5dd213cdf03f3bcb7000e69dda0de8b956061947ff6bcff5558 \ + --hash=sha256:ff838d62ec1bfce4f9ba7fa16f4a7b554cd8d0c299e6be37502161a660c84eef + # via pyjwt +pycparser==3.0 \ + --hash=sha256:600f49d217304a5902ac3c37e1281c9fe94e4d0489de643a9504c5cdfdfc6b29 \ + --hash=sha256:b727414169a36b7d524c1c3e31839a521725078d7b2ff038656844266160a992 + # via cffi +pyjwt==2.15.0 \ + --hash=sha256:7a3742debf6b879e912dbb9819ceec1594be812452b78c5f2e2dfc56564954f8 \ + --hash=sha256:b11c5f9791d7bf51c2b39a81ed669f6b2dbbd669df2942f6c60167e9e3d1abe4 + # via -r requirements-runtime.txt +waitress==3.0.2 \ + --hash=sha256:682aaaf2af0c44ada4abfb70ded36393f0e307f4ab9456a215ce0020baefc31f \ + --hash=sha256:c56d67fd6e87c2ee598b76abdd4e96cfad1f24cacdea5078d382b1f9d7b5ed2e + # via -r requirements-runtime.txt diff --git a/requirements-runtime.txt b/requirements-runtime.txt new file mode 100644 index 0000000..b9415d1 --- /dev/null +++ b/requirements-runtime.txt @@ -0,0 +1,2 @@ +PyJWT[crypto]>=2.10,<3 +waitress>=3,<4 diff --git a/scripts/alert_ack.py b/scripts/alert_ack.py new file mode 100644 index 0000000..78b6b53 --- /dev/null +++ b/scripts/alert_ack.py @@ -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'

Receipt acknowledged

{title}

{status}

This does not resolve or silence the alert.

') + 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'

{title}

Started {html.escape(alert["starts_at"])}

' + '

Confirm that you received this alert. This does not resolve or silence it.

' + f'
' + '
') + 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 diff --git a/scripts/alert_audit.py b/scripts/alert_audit.py new file mode 100644 index 0000000..b9925d2 --- /dev/null +++ b/scripts/alert_audit.py @@ -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 diff --git a/scripts/alert_identity.py b/scripts/alert_identity.py new file mode 100644 index 0000000..a50b3cc --- /dev/null +++ b/scripts/alert_identity.py @@ -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) diff --git a/scripts/alert_policy.py b/scripts/alert_policy.py new file mode 100644 index 0000000..bb8858f --- /dev/null +++ b/scripts/alert_policy.py @@ -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 diff --git a/scripts/alert_service.py b/scripts/alert_service.py new file mode 100644 index 0000000..a8f83bd --- /dev/null +++ b/scripts/alert_service.py @@ -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() diff --git a/scripts/telemetry_http.py b/scripts/telemetry_http.py new file mode 100644 index 0000000..2106dbb --- /dev/null +++ b/scripts/telemetry_http.py @@ -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 diff --git a/templates/alert-email.tmpl b/templates/alert-email.tmpl new file mode 100644 index 0000000..762cc00 --- /dev/null +++ b/templates/alert-email.tmpl @@ -0,0 +1,15 @@ +{{ define "railiance.alert.email" }} +Railiance alert notification + +{{ range .Alerts }} +Alert: {{ .Labels.alertname }} +State: {{ .Status }} +Started: {{ .StartsAt.Format "2006-01-02T15:04:05.999999999Z07:00" }} + +Review and acknowledge receipt: +https://telemetry.coulomb.social/ack/alerts?fingerprint={{ .Fingerprint | urlquery }}&starts_at={{ .StartsAt.Format "2006-01-02T15:04:05.999999999Z07:00" | urlquery }} + +Sign in as a Railiance admin, then select Acknowledge receipt. +Opening this link does not acknowledge, resolve or silence the alert. +{{ end }} +{{ end }} diff --git a/tests/test_alert_ack.py b/tests/test_alert_ack.py new file mode 100644 index 0000000..ac13983 --- /dev/null +++ b/tests/test_alert_ack.py @@ -0,0 +1,158 @@ +from dataclasses import replace +import io +import json +import os +from pathlib import Path +import sqlite3 +import sys +import tempfile +import unittest +from urllib.parse import urlencode + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +from alert_ack import Actor, Application, Store, occurrence + +NOW = 1790553600 # 2026-09-28 UTC +START = '2026-09-28T00:00:00Z' + + +class AcknowledgmentTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.addCleanup(self.temp.cleanup) + self.path = Path(self.temp.name) / 'ack.db' + self.store = Store(self.path) + self.actor = Actor('https://issuer.example', 'test-human', 'tenant:platform', + ('railiance-admin',), NOW + 900, 'c' * 32) + self.decision = {'id': 'test-decision', 'effect': 'allow', 'expires_at': NOW + 30} + self.app = Application(self.store, 'https://telemetry.example', 'w' * 32, + lambda env: self.actor, lambda *args: self.decision, + clock=lambda: NOW) + self.alert = {'status': 'firing', 'fingerprint': '0123456789abcdef', + 'startsAt': START, 'labels': {'alertname': 'TestFailure', 'owner': 'railiance-telemetry'}} + self.payload = {'version': '4', 'receiver': 'railiance-admin-email', 'alerts': [self.alert]} + self.identity = occurrence(self.alert['fingerprint'], START) + self.query = urlencode({'fingerprint': self.alert['fingerprint'], 'starts_at': START}) + + def request(self, method='GET', path='/ack/alerts', body=b'', **extra): + env = {'REQUEST_METHOD': method, 'PATH_INFO': path, 'QUERY_STRING': self.query, + 'CONTENT_LENGTH': str(len(body)), 'wsgi.input': io.BytesIO(body), + 'HTTP_ORIGIN': 'https://telemetry.example'} + env.update(extra) + result = {} + raw = b''.join(self.app(env, lambda status, headers: result.update(status=status, headers=headers))) + return result['status'], raw.decode() + + def receive(self): + self.store.receive(self.payload, NOW) + + def post(self, **kw): + return self.request('POST', body=urlencode({'csrf': self.actor.csrf}).encode(), **kw) + + def test_scanner_get_does_not_acknowledge(self): + self.receive() + status, page = self.request() + self.assertEqual(status, '200 OK') + self.assertIn('Acknowledge receipt', page) + self.assertIsNone(self.store.get(self.identity)['occurred_at']) + with self.store.connect() as db: + self.assertEqual(db.execute('SELECT COUNT(*) FROM outbox').fetchone()[0], 0) + + def test_click_commits_ack_and_audit_once_across_restart(self): + self.receive() + self.assertIn('Audit delivery pending', self.post()[1]) + self.store = Store(self.path) + self.app.store = self.store + self.post() + with self.store.connect() as db: + rows = db.execute('SELECT body FROM outbox').fetchall() + self.assertEqual(len(rows), 1) + event = json.loads(rows[0]['body']) + self.assertEqual(event['data']['actor_subject'], 'test-human') + self.assertEqual(event['data']['role'], 'railiance-admin') + self.assertEqual(event['subject'], 'alert:' + self.identity) + self.assertEqual(event['type'], 'telemetry.alert.acknowledged') + + def test_failed_outbox_rolls_back_ack(self): + self.receive() + with self.store.connect() as db: + db.execute("CREATE TRIGGER simulate_full BEFORE INSERT ON outbox BEGIN SELECT RAISE(ABORT, 'full'); END") + self.assertEqual(self.post()[0], '503 Service Unavailable') + self.assertIsNone(self.store.get(self.identity)['occurred_at']) + + def test_wrong_actor_role_tenant_service_expiry_denied(self): + self.receive() + original = self.actor + for change in ({'roles': ()}, {'tenant': 'tenant:other'}, {'principal_type': 'service'}, {'expires_at': NOW}): + self.actor = replace(original, **change) + self.assertIn(self.post()[0], ('401 Unauthorized', '403 Forbidden')) + self.assertIsNone(self.store.get(self.identity)['occurred_at']) + + def test_email_or_headers_do_not_grant_role(self): + self.receive() + self.actor = replace(self.actor, roles=()) + self.assertEqual(self.post(HTTP_X_EMAIL='bernd.worsch@gmail.com', HTTP_X_ROLE='railiance-admin')[0], '403 Forbidden') + + def test_origin_csrf_and_expired_pdp_denied(self): + self.receive() + self.assertEqual(self.post(HTTP_ORIGIN='https://attacker.example')[0], '403 Forbidden') + self.assertEqual(self.request('POST', body=b'csrf=wrong')[0], '403 Forbidden') + self.decision['expires_at'] = NOW + self.assertEqual(self.post()[0], '403 Forbidden') + self.assertIsNone(self.store.get(self.identity)['occurred_at']) + + def test_webhook_auth_and_retry_and_batch_rollback(self): + raw = json.dumps(self.payload).encode() + self.assertEqual(self.request('POST', '/webhook', raw)[0], '401 Unauthorized') + for _ in range(2): + self.assertEqual(self.request('POST', '/webhook', raw, HTTP_AUTHORIZATION='Bearer ' + 'w' * 32)[0], '200 OK') + changed = json.loads(raw) + changed['alerts'][0]['labels']['alertname'] = 'Different' + with self.assertRaises(ValueError): self.store.receive(changed, NOW) + self.assertEqual(self.store.get(self.identity)['alertname'], 'TestFailure') + + def test_new_firing_occurrence_requires_new_ack(self): + self.receive() + self.post() + self.alert['startsAt'] = '2026-09-27T23:59:59Z' + other = self.store.receive(self.payload, NOW)[0] + self.assertNotEqual(other, self.identity) + self.assertIsNone(self.store.get(other)['occurred_at']) + + def test_lost_audit_receipt_replays_original_event(self): + self.receive() + self.post() + sent = [] + def lost(event): + sent.append(event) + raise TimeoutError() + self.store.drain(lost) + def duplicate(event): + self.assertEqual(event, sent[0]) + return 200, {'status': 'duplicate', 'reference': 'audit:' + event['id']} + self.store.drain(duplicate) + self.assertEqual(self.store.get(self.identity)['audit_status'], 'delivered') + self.assertIn('Audit record archived', self.request()[1]) + + def test_audit_refusal_retained_and_wrong_receipt_not_accepted(self): + self.receive() + self.post() + self.store.drain(lambda event: (202, {'status': 'accepted', 'reference': 'wrong'})) + self.assertEqual(self.store.get(self.identity)['audit_status'], 'pending') + self.store.drain(lambda event: (403, {})) + self.assertEqual(self.store.get(self.identity)['audit_status'], 'blocked') + self.assertIn('Audit delivery pending', self.request()[1]) + self.assertEqual(self.store.audit_debt(), {'blocked': 1}) + with self.store.connect() as db: + event = json.loads(db.execute('SELECT body FROM outbox').fetchone()[0]) + self.assertTrue(self.store.requeue(event['id'])) + self.store.drain(lambda value: (202, {'status': 'accepted', 'reference': 'audit:' + value['id']})) + self.assertEqual(self.store.audit_debt(), {}) + + def test_unsafe_database_refused(self): + os.chmod(self.path, 0o644) + with self.assertRaises(ValueError): Store(self.path) + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_alert_audit_native.py b/tests/test_alert_audit_native.py new file mode 100644 index 0000000..03e2424 --- /dev/null +++ b/tests/test_alert_audit_native.py @@ -0,0 +1,48 @@ +"""Opt in with RTEL_AUDIT_CORE_SOURCE; uses the real receiver, synthetic custody.""" +import io +import json +import os +from pathlib import Path +import sys +import tempfile +import unittest + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +from alert_ack import Actor, Store + + +@unittest.skipUnless(os.environ.get('RTEL_AUDIT_CORE_SOURCE'), 'set RTEL_AUDIT_CORE_SOURCE for real receiver check') +class NativeAuditTests(unittest.TestCase): + def test_real_receiver_accepts_and_deduplicates_after_lost_reply(self): + sys.path.insert(0, os.environ['RTEL_AUDIT_CORE_SOURCE']) + from audit_core.ingestion import IngestionApplication + from audit_core.senders import SenderIdentity, SenderRegistry + from audit_core.sqlite_backend import SQLiteAuditBackend + with tempfile.TemporaryDirectory() as tmp: + store = Store(Path(tmp) / 'ack.db') + identity = store.receive({'version': '4', 'receiver': 'railiance-admin-email', 'alerts': [ + {'status': 'firing', 'fingerprint': '0123456789abcdef', 'startsAt': '2026-09-28T00:00:00Z', + 'labels': {'alertname': 'ControlledFailure', 'owner': 'railiance-telemetry'}}]}, 1790553600)[0] + store.acknowledge(identity, Actor('https://issuer.example', 'fixture-human', 'tenant:platform', + ('railiance-admin',), 1790554500, 'c' * 32), 'fixture-decision', 1790553600) + sender = SenderIdentity(name='railiance-telemetry', tokens=('fixture-only',), + sources=frozenset({'railiance-telemetry'}), tenants=frozenset({'tenant:platform'}), + evidence_kind='load-bearing', may_read=False) + app = IngestionApplication(SQLiteAuditBackend(str(Path(tmp) / 'audit.db')), SenderRegistry([sender])) + calls = [] + def send(event): + raw = json.dumps(event).encode() + env = {'REQUEST_METHOD': 'POST', 'PATH_INFO': '/v1/events', + 'CONTENT_LENGTH': str(len(raw)), 'wsgi.input': io.BytesIO(raw), + 'HTTP_AUTHORIZATION': 'Bearer fixture-only', 'HTTP_IDEMPOTENCY_KEY': event['id']} + response = {} + body = b''.join(app(env, lambda status, headers: response.update(status=int(status[:3])))) + calls.append(response['status']) + if len(calls) == 1: + raise TimeoutError('simulated lost response') + return response['status'], json.loads(body) + store.drain(send) + store = Store(Path(tmp) / 'ack.db') + store.drain(send) + self.assertEqual(calls, [202, 200]) + self.assertEqual(store.get(identity)['audit_status'], 'delivered') diff --git a/tests/test_alert_audit_transport.py b/tests/test_alert_audit_transport.py new file mode 100644 index 0000000..65fc901 --- /dev/null +++ b/tests/test_alert_audit_transport.py @@ -0,0 +1,50 @@ +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +import json +from pathlib import Path +import sys +import threading +import unittest + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +from alert_audit import AuditTransport + + +class AuditTransportTests(unittest.TestCase): + def test_real_http_headers_and_redirect_refusal(self): + calls = [] + class Handler(BaseHTTPRequestHandler): + def log_message(self, *args): + pass + + def do_POST(self): + calls.append((self.path, self.headers['Authorization'], self.headers['Idempotency-Key'])) + body = json.loads(self.rfile.read(int(self.headers['Content-Length']))) + if body['id'] == 'redirect': + self.send_response(302) + self.send_header('Location', '/credential-leak') + self.end_headers() + return + self.send_response(202) + self.end_headers() + self.wfile.write(json.dumps({'status': 'accepted', 'reference': 'audit:' + body['id']}).encode()) + + server = ThreadingHTTPServer(('127.0.0.1', 0), Handler) + thread = threading.Thread(target=server.serve_forever) + thread.start() + try: + transport = AuditTransport(f'http://127.0.0.1:{server.server_port}', + lambda: 'fixture-credential', allow_internal_http=True) + self.assertEqual(transport({'id': 'test'}), (202, {'status': 'accepted', 'reference': 'audit:test'})) + self.assertEqual(transport({'id': 'redirect'}), (302, {})) + self.assertEqual(calls, [('/v1/events', 'Bearer fixture-credential', 'test'), + ('/v1/events', 'Bearer fixture-credential', 'redirect')]) + finally: + server.shutdown() + thread.join() + server.server_close() + + def test_public_cleartext_and_credentialed_origin_refused(self): + for origin in ('http://audit.example', 'https://user:password@audit.example', + 'https://audit.example/other', 'https://audit.example?next=elsewhere'): + with self.assertRaises(ValueError): + AuditTransport(origin, lambda: 'fixture', allow_internal_http=True) diff --git a/tests_runtime/test_identity_policy.py b/tests_runtime/test_identity_policy.py new file mode 100644 index 0000000..cfa75e6 --- /dev/null +++ b/tests_runtime/test_identity_policy.py @@ -0,0 +1,131 @@ +from dataclasses import replace +from datetime import datetime, timezone +import json +from pathlib import Path +import sys +import time +import unittest +from urllib.parse import parse_qs, urlsplit + +import jwt +from cryptography.hazmat.primitives.asymmetric import rsa + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +from alert_identity import Login, SCOPES +from alert_policy import Policy, request_digest +from alert_ack import Actor + + +class Issuer: + def __init__(self): + self.key = rsa.generate_private_key(public_exponent=65537, key_size=2048) + self.jwk = json.loads(jwt.algorithms.RSAAlgorithm.to_jwk(self.key.public_key())) + self.jwk.update(kid='fixture', use='sig', alg='RS256') + self.nonce = '' + self.change = lambda c: c + + def request(self, method, url, body=None, headers=None): + if url.endswith('/.well-known/openid-configuration'): + return 200, dict(issuer='https://issuer.example', authorization_endpoint='https://issuer.example/authorize', + token_endpoint='https://issuer.example/token', jwks_uri='https://issuer.example/jwks', + code_challenge_methods_supported=['S256']) + if url.endswith('/jwks'): + return 200, {'keys': [self.jwk]} + now = int(time.time()) + claims = dict(iss='https://issuer.example', sub='fixture-human', aud='railiance-telemetry-admin', + iat=now, exp=now + 600, nonce=self.nonce, tenant='tenant:platform', + tenant_source='registration', principal_type='human', groups=[], roles=['railiance-admin'], + assurance=dict(level='aal2', mfa=True, source='key-cape', methods=['pwd', 'otp'], at=now)) + claims = self.change(claims) + identity = jwt.encode(claims, self.key, algorithm='RS256', headers={'kid': 'fixture'}) + access = dict(claims, aud='railiance-telemetry', scope=' '.join(SCOPES)) + return 200, dict(id_token=identity, access_token=jwt.encode(access, self.key, algorithm='RS256', headers={'kid': 'fixture'})) + + +class IdentityTests(unittest.TestCase): + def setUp(self): + self.issuer = Issuer() + self.login = Login('https://issuer.example', 'https://telemetry.example', transport=self.issuer) + + def start(self): + url, browser = self.login.start('/ack/alerts?fingerprint=0123456789abcdef&starts_at=2026-09-28T00:00:00Z') + params = parse_qs(urlsplit(url).query) + self.issuer.nonce = params['nonce'][0] + self.assertEqual(params['code_challenge_method'], ['S256']) + return params['state'][0], browser + + def actor(self): + state, browser = self.start() + sid, target = self.login.finish(state, browser, 'fixture-code') + return self.login.session(sid) + + def test_verified_human_session_and_logout(self): + state, browser = self.start() + sid, target = self.login.finish(state, browser, 'fixture-code') + actor = self.login.session(sid) + self.assertEqual(actor.subject, 'fixture-human') + self.assertEqual(actor.roles, ('railiance-admin',)) + self.assertLessEqual(actor.expires_at, time.time() + 600) + self.login.logout(sid) + self.assertIsNone(self.login.session(sid)) + with self.assertRaises(ValueError): self.login.finish(state, browser, 'fixture-code') + + def test_wrong_browser_and_state_replay_refused(self): + state, browser = self.start() + with self.assertRaises(ValueError): self.login.finish(state, 'other', 'fixture-code') + with self.assertRaises(ValueError): self.login.finish(state, browser, 'fixture-code') + + def test_nonce_wrong_issuer_wrong_tenant_and_weak_mfa(self): + for changes in ({'nonce': 'wrong'}, {'iss': 'https://other.example'}, {'tenant': 'tenant:other'}, + {'principal_type': 'service'}, {'assurance': {'mfa': False}}, {'exp': 1}): + self.issuer.change = lambda claims, changes=changes: dict(claims, **changes) + state, browser = self.start() + with self.assertRaises(ValueError): self.login.finish(state, browser, 'fixture-code') + + def test_return_path_cannot_escape_surface(self): + for value in ('https://attacker.example/ack/alerts', '//attacker.example/ack/alerts', '/other'): + with self.assertRaises(ValueError): self.login.start(value) + + +class DecisionTransport: + def __init__(self): + self.change = lambda response: response + + def request(self, method, url, body=None, headers=None): + request = json.loads(body) + now = time.time() + stamp = lambda delta: datetime.fromtimestamp(now + delta, timezone.utc).isoformat() + response = dict(id='fixture-decision', request_id=request['id'], contract_version='flex-auth.decision-record.v1', + effect='allow', obligations=[], matched_policy_version='v1', subject=request['subject'], resource=request['resource'], + binding=dict(tenant=request['tenant'], action=request['action'], subject=request['subject'], resource=request['resource'], + context=request.get('context', {}), + submitted_request_digest=request_digest(request), request_digest='sha256:' + 'b' * 64), + provenance=dict(policy_package='telemetry.ack', policy_version='v1', policy_package_digest='sha256:' + 'a' * 64, + registry_snapshot_digest='sha256:' + 'c' * 64, evaluator='flex-auth/fixture', decision_time=stamp(0)), + lifetime=dict(kind='ttl', not_before=stamp(-1), expires_at=stamp(30))) + return 200, self.change(response) + + +class PolicyTests(unittest.TestCase): + def test_policy_pins_and_action_bindings_fail_closed(self): + actor = Actor('https://issuer.example', 'fixture-human', 'tenant:platform', ('railiance-admin',), + time.time() + 600, 'c' * 32, tenant_source='directory') + transport = DecisionTransport() + policy = Policy('https://policy.example', lambda: 'fixture-workload-token', 'telemetry.ack', 'v1', + 'sha256:' + 'a' * 64, transport=transport) + resource = 'alert:11111111-1111-1111-1111-111111111111' + self.assertEqual(policy(actor, 'acknowledge', resource)['effect'], 'allow') + mutations = [lambda r: dict(r, effect='deny'), lambda r: dict(r, obligations=['unknown']), + lambda r: dict(r, request_id='other'), + lambda r: dict(r, binding=dict(r['binding'], action='read')), + lambda r: dict(r, binding=dict(r['binding'], submitted_request_digest='sha256:' + '0' * 64)), + lambda r: dict(r, provenance=dict(r['provenance'], policy_package_digest='sha256:' + '0' * 64)), + lambda r: dict(r, lifetime=dict(r['lifetime'], expires_at='2000-01-01T00:00:00Z'))] + for change in mutations: + transport.change = change + self.assertEqual(policy(actor, 'acknowledge', resource)['effect'], 'deny') + self.assertEqual(policy(actor, 'delete', resource)['effect'], 'deny') + + +if __name__ == '__main__': + unittest.main() diff --git a/tests_runtime/test_policy_native.py b/tests_runtime/test_policy_native.py new file mode 100644 index 0000000..70adec3 --- /dev/null +++ b/tests_runtime/test_policy_native.py @@ -0,0 +1,42 @@ +from dataclasses import replace +import json +import os +from pathlib import Path +import subprocess +import sys +import tempfile +import time +import unittest + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / 'scripts')) +from alert_ack import Actor +from alert_policy import Policy + + +@unittest.skipUnless(os.environ.get('RTEL_FLEX_AUTH_BINARY'), 'set RTEL_FLEX_AUTH_BINARY') +class NativePolicyTests(unittest.TestCase): + def test_native_wire_digest_and_revoked_or_wrong_identity_denials(self): + with tempfile.TemporaryDirectory() as directory: + request_file = Path(directory) / 'request.json' + class Native: + def request(self, method, url, body=None, headers=None): + request_file.write_bytes(body) + result = subprocess.run([os.environ['RTEL_FLEX_AUTH_BINARY'], 'check', '--policy', + str(ROOT / 'integration/telemetry-policy.md'), '--registry', str(ROOT / 'integration/registry.json'), + '--request', str(request_file)], capture_output=True, check=True) + return 200, json.loads(result.stdout) + actor = Actor('https://kc.coulomb.social', 'uid=tegwick,ou=people,dc=netkingdom,dc=local', + 'tenant:platform', ('railiance-admin',), time.time() + 600, 'c' * 32, + assurance=dict(level='aal2', mfa=True, source='key-cape', methods=['pwd', 'otp'], at=int(time.time())), + groups=('railiance-admins',), tenant_source='directory') + # Native digest is pinned by the owner packet, not inferred from a response. + policy = Policy('https://policy.example', lambda: 'fixture-only', + 'railiance-telemetry.alert-acknowledgment', 'v1', + 'sha256:02938202cf75140d6ab638b8d0fbce6ebb22832354efc89cab040b13e8943f09', transport=Native()) + resource = 'alert:11111111-1111-1111-1111-111111111111' + self.assertEqual(policy(actor, 'acknowledge', resource)['effect'], 'allow') + for changed in (replace(actor, roles=()), replace(actor, groups=()), replace(actor, subject='other'), + replace(actor, tenant='tenant:other'), replace(actor, issuer='https://other.example'), + replace(actor, assurance=dict(actor.assurance, at=1))): + self.assertEqual(policy(changed, 'acknowledge', resource)['effect'], 'deny') diff --git a/tests_runtime/test_service.py b/tests_runtime/test_service.py new file mode 100644 index 0000000..7a62c64 --- /dev/null +++ b/tests_runtime/test_service.py @@ -0,0 +1,57 @@ +import io +from pathlib import Path +import sys +import tempfile +import time +import unittest +from urllib.parse import parse_qs, urlencode, urlsplit + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts')) +from alert_ack import Store +from alert_identity import Login +from alert_service import Router +from test_identity_policy import Issuer + + +class ServiceTests(unittest.TestCase): + def test_browser_login_confirmation_and_logout(self): + with tempfile.TemporaryDirectory() as tmp: + token = Path(tmp) / 'webhook-token' + token.write_text('w' * 32) + store = Store(Path(tmp) / 'state.db') + issuer = Issuer() + login = Login('https://issuer.example', 'https://telemetry.example', transport=issuer) + router = Router({'origin': 'https://telemetry.example', 'webhook_token_file': str(token)}, store, + login=login, policy=lambda *args: dict(id='fixture-decision', effect='allow', expires_at=time.time() + 30)) + query = urlencode(dict(fingerprint='0123456789abcdef', starts_at='2026-09-28T00:00:00Z')) + def call(path, method='GET', cookie='', query=query, body=b''): + env = dict(PATH_INFO=path, REQUEST_METHOD=method, HTTP_COOKIE=cookie, QUERY_STRING=query, + CONTENT_LENGTH=str(len(body)), HTTP_ORIGIN='https://telemetry.example') + env['wsgi.input'] = io.BytesIO(body) + result = {} + output = b''.join(router(env, lambda status, headers: result.update(status=status, headers=headers))) + return result, output.decode() + response, _ = call('/ack/alerts') + self.assertEqual(response['status'], '303 See Other') + headers = dict(response['headers']) + browser_cookie = headers['Set-Cookie'].split(';')[0] + params = parse_qs(urlsplit(headers['Location']).query) + issuer.nonce = params['nonce'][0] + response, _ = call('/ack/auth/callback', cookie=browser_cookie, + query=urlencode(dict(state=params['state'][0], code='fixture-code'))) + self.assertEqual(response['status'], '303 See Other') + cookies = [value for key, value in response['headers'] if key == 'Set-Cookie'] + session_cookie = next(c.split(';')[0] for c in cookies if c.startswith('__Host-rtel-session=')) + actor = login.session(session_cookie.split('=', 1)[1]) + identity = store.receive({'version': '4', 'receiver': 'railiance-admin-email', 'alerts': [ + {'status': 'firing', 'fingerprint': '0123456789abcdef', 'startsAt': '2026-09-28T00:00:00Z', + 'labels': {'alertname': 'TestFailure', 'owner': 'railiance-telemetry'}}]}, time.time())[0] + response, page = call('/ack/alerts', cookie=session_cookie) + self.assertIn('Acknowledge receipt', page) + self.assertIsNone(store.get(identity)['occurred_at']) + form = urlencode(dict(csrf=actor.csrf)).encode() + response, page = call('/ack/alerts', 'POST', session_cookie, body=form) + self.assertIn('Receipt acknowledged', page) + response, _ = call('/ack/logout', 'POST', session_cookie, query='', body=form) + self.assertEqual(response['status'], '200 OK') + self.assertIsNone(login.session(session_cookie.split('=', 1)[1])) diff --git a/workplans/RTEL-WP-0002-signal-contract.md b/workplans/RTEL-WP-0002-signal-contract.md index d22cd4b..e88d5b9 100644 --- a/workplans/RTEL-WP-0002-signal-contract.md +++ b/workplans/RTEL-WP-0002-signal-contract.md @@ -112,3 +112,28 @@ failure/absence acknowledgments, an outside-node watchdog and recurring backup ownership. The founder, Bernd Worsch, supplies recipient/admission decisions; rapp-telemetry and platform own runtime/custody integration. No additional task or workplan was opened, and no live schedule or notification was enabled. + +September 28 recipient decision: Bernd Worsch confirmed email to +`bernd.worsch@gmail.com`, requested the `railiance-admin` role for that user, +explicit link-based confirmation and acknowledgment records in audit-core, and +authorized controlled failure/absence drills. Recipient/channel choice is no +longer a blocker. The native directory group membership is now applied; the updated KeyCape +issuer is deployed. A real signed-role login remains unproved. + +Implemented the bounded acknowledgment component, email link template and audit +outbox/transport under this same T04. A GET never acknowledges; authenticated +human role, fresh policy decision, Origin and CSRF checks precede POST. The first +acknowledgment and audit envelope commit together; retries retain the same event. +Actual audit-core receiver code accepts then deduplicates after a lost reply and +process reopen. See `docs/alert-acknowledgment.md` for the concrete owner bindings. + +OIDC code/PKCE sessions, a native Flex Auth decision adapter, a background audit +worker and the package-owned container/manifests are now implemented. The native +policy evaluator and signed issuer fixtures pass. KeyCape 1164f65 is published +and deployed; the native directory role grant preserves existing memberships. + +T04 stays wait/blocked: the founder selected From `platform@coulomb.social` and +subsequently created the mailbox; password custody remains pending. Dedicated SMTP/webhook/audit custody, OIDC +client and enforced PDP caller admission, application rollout, signed identity, +real email/acknowledgment and independent audit readback remain. Outside-node +watchdog and recurring backup gates remain. No new task/workplan was created.