AUDIT-WP-0004-T04, closing the workplan. Decision (Bernd): default to redaction, allow rejection per sender. Losing an audit record over one field is worse than storing it masked, but a higher-assurance channel must be able to refuse rather than mask. secret_policy is set per sender identity in AUDIT_CORE_SENDERS and defaults to redact. Detection now covers the whole payload at any depth, including lists, rather than only the top level of data. Under redaction the value is masked and the key is preserved: dropping the key would hide that the sender transmitted the field at all, which is exactly what an operator needs in order to stop it. The stored record carries details.redaction with policy and affected paths, so a reader never has to infer whether what they see is what was sent. Idempotency is unaffected - the payload hash is taken over the original request body, so redaction is deterministic and a resubmission still reconciles as a duplicate. Both outcomes are counted durably by sender, source, action and field path, exposed at GET /v1/secret-findings. Per-path aggregation is the point: the actionable unit is "stop emitting data.auth.token on membership.added", not "there were 47 redactions". Counters survive restart because the fix they drive lives in another service. Contract doc updated to match. Tests 46 -> 50. WP-0004 is finished. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
433 lines
17 KiB
Python
433 lines
17 KiB
Python
"""Authenticated, idempotent HTTP ingestion for user-engine outbox events.
|
|
|
|
Response contract (AUDIT-WP-0004-T02). Senders key their retry behaviour off
|
|
these, so they are part of the interface, not an implementation detail:
|
|
|
|
=== ========== =========================================================
|
|
202 accepted Event is durably in custody. Do not retry.
|
|
200 duplicate Exact resubmission of an event already in custody. Do not
|
|
retry; delivery already succeeded.
|
|
400 rejected Malformed or disallowed. Retrying will not help — dead
|
|
letter it.
|
|
401 unauthorized Credential missing or invalid. Do not retry without a
|
|
new credential.
|
|
409 conflict The event id is held with a different payload. Retrying
|
|
will not help; this indicates a sender bug or id reuse.
|
|
503 unavailable Not accepted, but retryable. Retry with backoff.
|
|
500 error Unexpected fault. Not accepted. Retryable with backoff.
|
|
=== ========== =========================================================
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import signal
|
|
import sys
|
|
from datetime import datetime, timezone
|
|
from http import HTTPStatus
|
|
from typing import Any
|
|
from urllib.parse import parse_qs
|
|
|
|
from audit_core.interface import (
|
|
AuditEvent,
|
|
BackendUnavailableError,
|
|
EventConflictError,
|
|
EventValidationError,
|
|
IdempotentAuditBackend,
|
|
)
|
|
from audit_core.redaction import (
|
|
POLICY_REDACT,
|
|
SecretFieldRejection,
|
|
apply_policy,
|
|
finding_from_path,
|
|
)
|
|
from audit_core.senders import SenderIdentity, SenderRegistry, development_registry
|
|
from audit_core.sqlite_backend import SQLiteAuditBackend
|
|
|
|
MAX_BODY_BYTES = 256 * 1024
|
|
|
|
log = logging.getLogger("audit_core.ingestion")
|
|
|
|
|
|
class IngestionApplication:
|
|
"""WSGI application accepting user-engine outbox events.
|
|
|
|
Writes through the audit backend contract rather than to storage directly,
|
|
so a 202 means a backend with a declared retention policy acknowledged the
|
|
event.
|
|
"""
|
|
|
|
def __init__(
|
|
self, backend: IdempotentAuditBackend, senders: SenderRegistry | str
|
|
) -> None:
|
|
policy = backend.retention_policy
|
|
if not policy.durable:
|
|
# The mock file backend declares durable=False. Refusing it here is
|
|
# what stops a development sink from silently becoming the
|
|
# production one (AUDIT-WP-0004-T01).
|
|
raise ValueError(
|
|
f"backend custody_class={policy.custody_class!r} is not durable; "
|
|
"refusing to accept audit events against it"
|
|
)
|
|
if isinstance(senders, str):
|
|
if not senders:
|
|
raise ValueError("bearer token is required")
|
|
senders = development_registry(senders)
|
|
self.backend = backend
|
|
self.senders = senders
|
|
|
|
def __call__(self, environ, start_response):
|
|
try:
|
|
return self._handle(environ, start_response)
|
|
except Exception:
|
|
# Nothing may escape: an unhandled exception here means
|
|
# start_response is never called and the sender sees a dropped
|
|
# connection it cannot classify.
|
|
log.exception("unhandled error in ingestion request")
|
|
return self._json(
|
|
start_response, HTTPStatus.INTERNAL_SERVER_ERROR, {"error": "internal_error"}
|
|
)
|
|
|
|
def _handle(self, environ, start_response):
|
|
path = environ.get("PATH_INFO", "")
|
|
if path == "/healthz":
|
|
return self._json(start_response, HTTPStatus.OK, {"status": "ok"})
|
|
if path == "/readyz":
|
|
return self._readiness(start_response)
|
|
method = environ.get("REQUEST_METHOD")
|
|
identity = self.senders.authenticate(environ.get("HTTP_AUTHORIZATION"))
|
|
if identity is None:
|
|
return self._json(
|
|
start_response, HTTPStatus.UNAUTHORIZED, {"error": "unauthorized"}
|
|
)
|
|
|
|
if method == "GET" and (
|
|
path.startswith("/v1/events")
|
|
or path in ("/v1/dead-letters", "/v1/secret-findings")
|
|
):
|
|
return self._read(start_response, environ, path, identity)
|
|
|
|
if path != "/v1/events" or method != "POST":
|
|
return self._json(start_response, HTTPStatus.NOT_FOUND, {"error": "not_found"})
|
|
|
|
if not identity.may_write:
|
|
return self._json(start_response, HTTPStatus.FORBIDDEN, {"error": "write_forbidden"})
|
|
|
|
raw = b""
|
|
payload: Any = {}
|
|
try:
|
|
raw = self._read_body(environ)
|
|
payload = json.loads(raw)
|
|
event = normalize(payload, environ.get("HTTP_IDEMPOTENCY_KEY"), identity)
|
|
except SecretFieldRejection as exc:
|
|
self._count_secrets(payload, identity, "rejected", exc.findings)
|
|
self._dead_letter(raw, str(exc), identity)
|
|
return self._json(start_response, HTTPStatus.BAD_REQUEST, {"error": str(exc)})
|
|
except (ValueError, TypeError, KeyError, json.JSONDecodeError) as exc:
|
|
self._dead_letter(raw, str(exc), identity)
|
|
return self._json(start_response, HTTPStatus.BAD_REQUEST, {"error": str(exc)})
|
|
|
|
redaction = event.details.get("redaction")
|
|
if redaction:
|
|
self._count_secrets(
|
|
payload, identity, "redacted",
|
|
[finding_from_path(p) for p in redaction["paths"]],
|
|
)
|
|
|
|
try:
|
|
result = self.backend.accept(event, hashlib.sha256(raw).hexdigest())
|
|
except EventConflictError as exc:
|
|
log.warning("event conflict: %s", exc)
|
|
return self._json(start_response, HTTPStatus.CONFLICT, {"error": "event_id_conflict"})
|
|
except EventValidationError as exc:
|
|
return self._json(start_response, HTTPStatus.BAD_REQUEST, {"error": str(exc)})
|
|
except BackendUnavailableError as exc:
|
|
log.error("backend unavailable: %s", exc)
|
|
return self._json(
|
|
start_response, HTTPStatus.SERVICE_UNAVAILABLE, {"error": "backend_unavailable"}
|
|
)
|
|
|
|
return self._json(
|
|
start_response,
|
|
HTTPStatus.OK if result.duplicate else HTTPStatus.ACCEPTED,
|
|
{
|
|
"status": "duplicate" if result.duplicate else "accepted",
|
|
"reference": result.reference,
|
|
},
|
|
)
|
|
|
|
def _read(self, start_response, environ, path: str, identity):
|
|
"""Operator read surface (AUDIT-WP-0004-T05).
|
|
|
|
Read is a distinct privilege from write: a sender credential must not
|
|
be able to read the audit trail back.
|
|
"""
|
|
if not identity.may_read:
|
|
return self._json(start_response, HTTPStatus.FORBIDDEN, {"error": "read_forbidden"})
|
|
|
|
query = parse_qs(environ.get("QUERY_STRING", ""))
|
|
try:
|
|
if path == "/v1/dead-letters":
|
|
return self._json(
|
|
start_response, HTTPStatus.OK,
|
|
{"dead_letters": self.backend.dead_letters(_limit(query))},
|
|
)
|
|
if path == "/v1/secret-findings":
|
|
return self._json(
|
|
start_response, HTTPStatus.OK,
|
|
{"secret_findings": self.backend.secret_findings(_limit(query))},
|
|
)
|
|
if path == "/v1/events":
|
|
correlation = (query.get("correlation_id") or [""])[0]
|
|
if not correlation:
|
|
return self._json(
|
|
start_response, HTTPStatus.BAD_REQUEST,
|
|
{"error": "correlation_id_required"},
|
|
)
|
|
return self._json(
|
|
start_response, HTTPStatus.OK,
|
|
{"events": self.backend.by_correlation(correlation, _limit(query))},
|
|
)
|
|
event_id = path[len("/v1/events/"):]
|
|
record = self.backend.get(event_id) if event_id else None
|
|
if record is None:
|
|
return self._json(start_response, HTTPStatus.NOT_FOUND, {"error": "not_found"})
|
|
return self._json(start_response, HTTPStatus.OK, record)
|
|
except BackendUnavailableError as exc:
|
|
log.error("read failed: %s", exc)
|
|
return self._json(
|
|
start_response, HTTPStatus.SERVICE_UNAVAILABLE, {"error": "backend_unavailable"}
|
|
)
|
|
|
|
def _count_secrets(self, payload: Any, identity, outcome: str, findings) -> None:
|
|
"""Count secret-shaped fields by path, sender, source and action.
|
|
|
|
Counted so the sending service can be fixed. Aggregation is by field
|
|
path rather than by event, because the actionable unit is "stop
|
|
emitting ``data.auth.token`` on ``membership.added``", not "there were
|
|
47 redactions".
|
|
"""
|
|
counter = getattr(self.backend, "count_secret_findings", None)
|
|
if not callable(counter) or not findings:
|
|
return
|
|
source = str(payload.get("source") or "") if isinstance(payload, dict) else ""
|
|
action = str(payload.get("type") or "") if isinstance(payload, dict) else ""
|
|
try:
|
|
counter(
|
|
sender=identity.name, source=source, action=action,
|
|
outcome=outcome, findings=findings,
|
|
)
|
|
except BackendUnavailableError as exc:
|
|
# A missed counter must never change the event's outcome.
|
|
log.error("could not count secret findings: %s", exc)
|
|
|
|
def _dead_letter(self, raw: bytes, reason: str, identity) -> None:
|
|
"""Record a rejection so an operator can see what the sender dropped."""
|
|
recorder = getattr(self.backend, "record_rejection", None)
|
|
if not callable(recorder):
|
|
return
|
|
event_id = None
|
|
try:
|
|
parsed = json.loads(raw)
|
|
if isinstance(parsed, dict):
|
|
event_id = str(parsed.get("id") or "") or None
|
|
except (ValueError, TypeError):
|
|
pass
|
|
try:
|
|
recorder(
|
|
event_id=event_id,
|
|
reason=reason,
|
|
payload_hash=hashlib.sha256(raw).hexdigest(),
|
|
sender=identity.name,
|
|
payload=raw.decode("utf-8", "replace"),
|
|
)
|
|
except BackendUnavailableError as exc:
|
|
# A rejection we could not record is worth a log line, but it must
|
|
# not turn a 400 into a 503 — the event is still rejected.
|
|
log.error("could not record dead letter: %s", exc)
|
|
|
|
def _read_body(self, environ) -> bytes:
|
|
try:
|
|
length = int(environ.get("CONTENT_LENGTH") or 0)
|
|
except (TypeError, ValueError):
|
|
raise ValueError("invalid_content_length") from None
|
|
if length <= 0:
|
|
raise ValueError("empty_body")
|
|
if length > MAX_BODY_BYTES:
|
|
raise ValueError("payload_too_large")
|
|
raw = environ["wsgi.input"].read(length)
|
|
if len(raw) != length:
|
|
raise ValueError("truncated_body")
|
|
return raw
|
|
|
|
def _readiness(self, start_response):
|
|
try:
|
|
health = getattr(self.backend, "health", None)
|
|
if callable(health):
|
|
health()
|
|
except BackendUnavailableError as exc:
|
|
log.error("readiness failed: %s", exc)
|
|
return self._json(
|
|
start_response, HTTPStatus.SERVICE_UNAVAILABLE, {"status": "unavailable"}
|
|
)
|
|
policy = self.backend.retention_policy
|
|
return self._json(
|
|
start_response,
|
|
HTTPStatus.OK,
|
|
{"status": "ok", "custody_class": policy.custody_class, "durable": policy.durable},
|
|
)
|
|
|
|
@staticmethod
|
|
def _json(start_response, status: HTTPStatus, payload: dict[str, Any]):
|
|
body = json.dumps(payload).encode()
|
|
start_response(
|
|
f"{status.value} {status.phrase}",
|
|
[("Content-Type", "application/json"), ("Content-Length", str(len(body)))],
|
|
)
|
|
return [body]
|
|
|
|
|
|
def normalize(
|
|
payload: dict[str, Any],
|
|
idempotency_key: str | None,
|
|
identity: SenderIdentity | None = None,
|
|
) -> AuditEvent:
|
|
required = (
|
|
"id", "type", "source", "subject", "tenant", "correlation_id", "occurred_at", "data",
|
|
)
|
|
if not isinstance(payload, dict) or any(not payload.get(key) for key in required):
|
|
raise ValueError("invalid_event")
|
|
if idempotency_key != payload["id"]:
|
|
raise ValueError("idempotency_key_mismatch")
|
|
source = str(payload["source"])
|
|
tenant = str(payload["tenant"])
|
|
# The claimed source and tenant are checked against what this credential is
|
|
# permitted to assert, not against a literal (AUDIT-WP-0004-T03).
|
|
if identity is not None:
|
|
if not identity.permits_source(source):
|
|
raise ValueError("source_not_allowed")
|
|
if not identity.permits_tenant(tenant):
|
|
raise ValueError("tenant_not_allowed")
|
|
elif source != "user-engine":
|
|
raise ValueError("source_not_allowed")
|
|
observed_at = _normalize_timestamp(payload["occurred_at"])
|
|
policy = identity.secret_policy if identity is not None else POLICY_REDACT
|
|
# Raises SecretFieldRejection under the reject policy; that exception
|
|
# carries the findings so the caller can count them.
|
|
data, findings = apply_policy(payload, policy)
|
|
details: dict[str, Any] = {
|
|
"correlation_id": str(payload["correlation_id"]),
|
|
"data": data,
|
|
}
|
|
if findings:
|
|
# The stored record differs from what the sender transmitted. Say so in
|
|
# the record itself rather than leaving it to be inferred.
|
|
details["redaction"] = {
|
|
"policy": policy,
|
|
"paths": [f.path for f in findings],
|
|
}
|
|
return AuditEvent(
|
|
event_id=str(payload["id"]),
|
|
observed_at=observed_at,
|
|
tenant=tenant,
|
|
scope="tenant",
|
|
source=source,
|
|
action=str(payload["type"]),
|
|
resource=str(payload["subject"]),
|
|
outcome="recorded",
|
|
actor=None,
|
|
details=details,
|
|
)
|
|
|
|
|
|
def _normalize_timestamp(value: Any) -> str:
|
|
"""Parse an event timestamp, requiring an explicit offset.
|
|
|
|
A naive timestamp is ambiguous by up to a day, which is not good enough for
|
|
an audit trail — the sender must say which offset it meant.
|
|
"""
|
|
try:
|
|
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
|
|
except (TypeError, ValueError):
|
|
raise ValueError("invalid_timestamp") from None
|
|
if parsed.tzinfo is None:
|
|
raise ValueError("timestamp_missing_timezone")
|
|
return parsed.astimezone(timezone.utc).isoformat()
|
|
|
|
|
|
def _limit(query: dict[str, list[str]], default: int = 100, ceiling: int = 1000) -> int:
|
|
try:
|
|
return max(1, min(int((query.get("limit") or [default])[0]), ceiling))
|
|
except (TypeError, ValueError):
|
|
return default
|
|
|
|
|
|
def serve(app, host: str, port: int, threads: int, timeout: int) -> None:
|
|
"""Serve ``app``, preferring a production WSGI server (T06).
|
|
|
|
waitress is the intended production server and is installed in the image
|
|
via the ``serve`` extra. The fallback is a threaded wsgiref server with a
|
|
socket timeout — bounded rather than good, and loud about which one is in
|
|
use so a deployment cannot quietly end up on the fallback.
|
|
"""
|
|
try:
|
|
from waitress import serve as waitress_serve
|
|
except ImportError:
|
|
log.warning(
|
|
"waitress not installed — falling back to a threaded wsgiref server. "
|
|
"Install the 'serve' extra for production (AUDIT-WP-0004-T06)."
|
|
)
|
|
_serve_fallback(app, host, port, timeout)
|
|
return
|
|
|
|
log.info("serving on waitress host=%s port=%s threads=%s", host, port, threads)
|
|
waitress_serve(
|
|
app, host=host, port=port, threads=threads,
|
|
channel_timeout=timeout, ident="audit-core",
|
|
)
|
|
|
|
|
|
def _serve_fallback(app, host: str, port: int, timeout: int) -> None:
|
|
from socketserver import ThreadingMixIn
|
|
from wsgiref.simple_server import WSGIServer, make_server
|
|
|
|
class ThreadedWSGIServer(ThreadingMixIn, WSGIServer):
|
|
daemon_threads = True
|
|
# Without this a slow or idle client holds a worker indefinitely; the
|
|
# original single-threaded server let one such client block every
|
|
# sender.
|
|
timeout = timeout
|
|
|
|
with make_server(host, port, app, server_class=ThreadedWSGIServer) as server:
|
|
server.socket.settimeout(timeout)
|
|
|
|
def shutdown(signum, _frame):
|
|
log.info("received signal %s, shutting down", signum)
|
|
server.shutdown()
|
|
|
|
for sig in (signal.SIGTERM, signal.SIGINT):
|
|
signal.signal(sig, shutdown)
|
|
log.info("serving on threaded wsgiref host=%s port=%s", host, port)
|
|
server.serve_forever()
|
|
|
|
|
|
def main() -> None:
|
|
logging.basicConfig(
|
|
level=os.environ.get("AUDIT_CORE_LOG_LEVEL", "INFO"),
|
|
format='{"ts":"%(asctime)s","level":"%(levelname)s","logger":"%(name)s","msg":"%(message)s"}',
|
|
stream=sys.stdout,
|
|
)
|
|
backend = SQLiteAuditBackend(
|
|
os.environ.get("AUDIT_CORE_DATABASE_PATH", "/data/audit-core.db")
|
|
)
|
|
app = IngestionApplication(backend, SenderRegistry.from_env())
|
|
serve(
|
|
app,
|
|
host=os.environ.get("AUDIT_CORE_HOST", "0.0.0.0"),
|
|
port=int(os.environ.get("AUDIT_CORE_HTTP_PORT", "8080")),
|
|
threads=int(os.environ.get("AUDIT_CORE_THREADS", "8")),
|
|
timeout=int(os.environ.get("AUDIT_CORE_REQUEST_TIMEOUT", "30")),
|
|
)
|