audit-core/tests/test_ingestion.py
tegwick 2f4e1adf66
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 1s
Add deployment manifests, custody-class guard and request counters
AUDIT-WP-0005-T03 (progress). Manifests validated --dry-run=server
--validate=strict against railiance01; not applied, since deployment is gated
on RAPP-POSTGRES-WP-0002 and T02 credentials. Nothing here mutates the cluster.

Conventions read off the deployed user-engine workload rather than invented:
digest-pinned image from forgejo.coulomb.social, runAsNonRoot with
RuntimeDefault seccomp, no privilege escalation, all capabilities dropped,
readOnlyRootFilesystem, probes on a named http port, same resource envelope.

The namespace carries railiance.io/postgres-client: platform-pg, which is what
platform-pg-consumer-ingress in rapp-postgres admits; without that label the
pod cannot reach the database at all.

NetworkPolicies default-deny both directions, then permit ingress from the
user-engine namespace only, a separately labelled operator read path, and
egress to PostgreSQL in databases plus DNS.

Three decisions worth naming. Liveness is /healthz while readiness is /readyz,
so a database outage drops the pod from the Service rather than restarting it
in a loop. readOnlyRootFilesystem enforces the empty-filesystem property rather
than trusting it, so the SQLite fallback physically cannot accumulate audit
records on ephemeral storage. AUDIT_CORE_REQUIRE_CUSTODY_CLASS=archive makes a
missing database URL a startup failure instead of a silent downgrade to the
development store.

Counters deferred from WP-0004-T06 are exposed as JSON at /v1/stats behind the
read privilege, not as Prometheus exposition format: the cluster runs no
Prometheus, no ServiceMonitor CRD and no other scrape target, so an exposition
endpoint would target a scrape path that does not exist. Usable with curl now
and a small step from /metrics later.

Tests 77 -> 80.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 17:42:43 +02:00

436 lines
16 KiB
Python

import io
import json
import pytest
from audit_core.ingestion import IngestionApplication
from audit_core.interface import BackendUnavailableError, RetentionPolicy
from audit_core.mock_file_backend import MockFileAuditBackend
from audit_core.sqlite_backend import SQLiteAuditBackend
def invoke(app, payload, *, token="opaque", key="evt-1", method="POST",
path="/v1/events", body=None, length=None):
raw = body if body is not None else json.dumps(payload).encode()
environ = {
"PATH_INFO": path,
"REQUEST_METHOD": method,
"CONTENT_LENGTH": str(len(raw)) if length is None else length,
"wsgi.input": io.BytesIO(raw),
"HTTP_AUTHORIZATION": f"Bearer {token}",
}
if key is not None:
environ["HTTP_IDEMPOTENCY_KEY"] = key
result = {}
out = b"".join(app(environ, lambda status, headers: result.update(status=status)))
return result["status"], (json.loads(out) if out else {})
def event(**overrides):
base = {
"id": "evt-1",
"type": "membership.added",
"source": "user-engine",
"subject": "membership-1",
"tenant": "tenant:friendly:binky",
"correlation_id": "corr-1",
"occurred_at": "2026-08-09T00:00:00+00:00",
"data": {"membership_id": "membership-1"},
}
base.update(overrides)
return base
@pytest.fixture
def app(tmp_path):
return IngestionApplication(SQLiteAuditBackend(str(tmp_path / "events.db")), "opaque")
# --- backend contract (T01) -------------------------------------------------
def test_refuses_a_non_durable_backend():
"""The development file backend must never become the production sink."""
with pytest.raises(ValueError, match="not durable"):
IngestionApplication(MockFileAuditBackend(base_dir="/tmp/unused"), "opaque")
def test_readiness_reports_custody_class(app):
status, body = invoke(app, None, path="/readyz", method="GET", body=b"")
assert status.startswith("200")
assert body["durable"] is True
assert body["custody_class"] == "development"
def test_accepted_events_are_durable_across_reopen(tmp_path):
path = str(tmp_path / "events.db")
first = IngestionApplication(SQLiteAuditBackend(path), "opaque")
assert invoke(first, event())[0].startswith("202")
reopened = IngestionApplication(SQLiteAuditBackend(path), "opaque")
status, body = invoke(reopened, event())
assert status.startswith("200") and body["status"] == "duplicate"
# --- idempotency and conflict (T01, T02) ------------------------------------
def test_accepts_once_and_replays_idempotently(app):
assert invoke(app, event())[0].startswith("202")
status, body = invoke(app, event())
assert status.startswith("200") and body["status"] == "duplicate"
def test_same_id_different_payload_is_a_conflict(app):
assert invoke(app, event())[0].startswith("202")
status, body = invoke(app, event(subject="membership-2"))
assert status.startswith("409")
assert body["error"] == "event_id_conflict"
# --- authentication (T02) ---------------------------------------------------
def test_rejects_wrong_credential(app):
assert invoke(app, event(), token="wrong")[0].startswith("401")
def test_non_ascii_credential_is_unauthorized_not_a_crash(app):
"""compare_digest raises TypeError on non-ASCII str; that is a 401."""
assert invoke(app, event(), token="wröng")[0].startswith("401")
# --- validation (T02, T07) --------------------------------------------------
@pytest.mark.parametrize("payload,expected", [
(event(source="somewhere-else"), "source_not_allowed"),
(event(occurred_at="2026-08-09T00:00:00"), "timestamp_missing_timezone"),
(event(occurred_at="not-a-date"), "invalid_timestamp"),
(event(tenant=""), "invalid_event"),
])
def test_rejects_bad_events(app, payload, expected):
status, body = invoke(app, payload)
assert status.startswith("400")
assert body["error"] == expected
def test_rejects_mismatched_idempotency_key(app):
status, body = invoke(app, event(), key="other")
assert status.startswith("400")
assert body["error"] == "idempotency_key_mismatch"
def test_rejects_malformed_json(app):
assert invoke(app, None, body=b"{not json")[0].startswith("400")
def test_rejects_empty_body(app):
status, body = invoke(app, None, body=b"")
assert status.startswith("400")
assert body["error"] == "empty_body"
def test_rejects_oversized_body(app):
status, body = invoke(app, None, body=b"x", length=str(512 * 1024))
assert status.startswith("400")
assert body["error"] == "payload_too_large"
def test_rejects_truncated_body(app):
status, body = invoke(app, None, body=b"{}", length="500")
assert status.startswith("400")
assert body["error"] == "truncated_body"
# --- routing (T07) ----------------------------------------------------------
@pytest.mark.parametrize("path,method", [
("/nope", "POST"),
("/nope", "GET"),
("/v1/events", "DELETE"),
])
def test_unknown_routes_are_not_found(app, path, method):
assert invoke(app, event(), path=path, method=method)[0].startswith("404")
def test_healthz_needs_no_credential(app):
status, _ = invoke(app, None, path="/healthz", method="GET", body=b"", token="wrong")
assert status.startswith("200")
# --- failure handling (T02) -------------------------------------------------
class _BrokenBackend:
@property
def retention_policy(self):
return RetentionPolicy("archive", 3650, True, True, durable=True)
def emit(self, event):
raise BackendUnavailableError("down")
def accept(self, event, payload_hash):
raise BackendUnavailableError("down")
class _ExplodingBackend(_BrokenBackend):
def accept(self, event, payload_hash):
raise RuntimeError("unexpected")
def test_backend_unavailable_is_retryable_503():
status, body = invoke(IngestionApplication(_BrokenBackend(), "opaque"), event())
assert status.startswith("503")
assert body["error"] == "backend_unavailable"
def test_unexpected_backend_error_still_returns_a_response():
"""No request path may terminate without calling start_response."""
status, body = invoke(IngestionApplication(_ExplodingBackend(), "opaque"), event())
assert status.startswith("500")
assert body["error"] == "internal_error"
# --- concurrency (T01) ------------------------------------------------------
def test_concurrent_duplicates_produce_exactly_one_record(tmp_path):
"""Racing submissions of one event: one acceptance, one custody record.
This is the assertion the whole service rests on, so it is exercised rather
than assumed. An earlier single-connection implementation passed every
serial test while letting two callers both be told they were first.
"""
import threading
backend = SQLiteAuditBackend(str(tmp_path / "race.db"))
app = IngestionApplication(backend, "opaque")
outcomes: list[str] = []
lock = threading.Lock()
barrier = threading.Barrier(16)
def submit():
barrier.wait()
status, body = invoke(app, event())
with lock:
outcomes.append(body.get("status", status))
threads = [threading.Thread(target=submit) for _ in range(16)]
for t in threads:
t.start()
for t in threads:
t.join()
assert outcomes.count("accepted") == 1, outcomes
assert outcomes.count("duplicate") == 15, outcomes
stored = backend.db.execute("SELECT COUNT(*) FROM events").fetchone()[0]
assert stored == 1
# --- sender identity binding (T03) ------------------------------------------
from audit_core.senders import SenderIdentity, SenderRegistry # noqa: E402
def bound_app(tmp_path, **kw):
identity = SenderIdentity(
name="user-engine",
tokens=kw.get("tokens", ("opaque",)),
sources=frozenset(kw.get("sources", {"user-engine"})),
tenants=frozenset(kw.get("tenants", {"tenant:friendly:binky"})),
may_read=kw.get("may_read", False),
secret_policy=kw.get("secret_policy", "redact"),
)
backend = SQLiteAuditBackend(str(tmp_path / "bound.db"))
return IngestionApplication(backend, SenderRegistry([identity])), backend
def test_credential_may_not_claim_another_tenant(tmp_path):
"""The property WP-0003 recorded as done but never implemented."""
app, _ = bound_app(tmp_path)
assert invoke(app, event())[0].startswith("202")
status, body = invoke(app, event(tenant="tenant:coulomb"))
assert status.startswith("400")
assert body["error"] == "tenant_not_allowed"
def test_credential_may_not_claim_another_source(tmp_path):
app, _ = bound_app(tmp_path)
status, body = invoke(app, event(source="issue-core"))
assert status.startswith("400")
assert body["error"] == "source_not_allowed"
def test_rotation_accepts_both_tokens(tmp_path):
"""Rotation must not need a delivery gap."""
app, _ = bound_app(tmp_path, tokens=("current", "next"))
assert invoke(app, event(), token="current")[0].startswith("202")
assert invoke(app, event(id="evt-2"), key="evt-2", token="next")[0].startswith("202")
assert invoke(app, event(id="evt-3"), key="evt-3", token="retired")[0].startswith("401")
def test_sender_credential_cannot_read_the_trail_back(tmp_path):
app, _ = bound_app(tmp_path, may_read=False)
assert invoke(app, event())[0].startswith("202")
status, body = invoke(app, None, path="/v1/events/evt-1", method="GET", body=b"")
assert status.startswith("403")
assert body["error"] == "read_forbidden"
# --- operator read surface (T05) --------------------------------------------
def test_lookup_by_event_id_and_correlation(app):
assert invoke(app, event())[0].startswith("202")
assert invoke(app, event(id="evt-2", correlation_id="corr-1"), key="evt-2")[0].startswith("202")
status, body = invoke(app, None, path="/v1/events/evt-1", method="GET", body=b"")
assert status.startswith("200")
assert body["event_id"] == "evt-1"
assert body["tenant"] == "tenant:friendly:binky"
status, body = invoke_query(app, "correlation_id=corr-1")
assert status.startswith("200")
assert {e["event_id"] for e in body["events"]} == {"evt-1", "evt-2"}
def invoke_query(app, query, token="opaque"):
environ = {
"PATH_INFO": "/v1/events",
"REQUEST_METHOD": "GET",
"QUERY_STRING": query,
"CONTENT_LENGTH": "0",
"wsgi.input": io.BytesIO(b""),
"HTTP_AUTHORIZATION": f"Bearer {token}",
}
result = {}
out = b"".join(app(environ, lambda status, headers: result.update(status=status)))
return result["status"], (json.loads(out) if out else {})
def test_unknown_event_id_is_not_found(app):
status, _ = invoke(app, None, path="/v1/events/nope", method="GET", body=b"")
assert status.startswith("404")
def test_correlation_lookup_requires_a_correlation_id(app):
status, body = invoke_query(app, "")
assert status.startswith("400")
assert body["error"] == "correlation_id_required"
def test_rejected_events_appear_as_dead_letters(app):
assert invoke(app, event(source="issue-core"))[0].startswith("400")
status, body = invoke(app, None, path="/v1/dead-letters", method="GET", body=b"")
assert status.startswith("200")
entry = body["dead_letters"][0]
assert entry["reason"] == "source_not_allowed"
assert entry["event_id"] == "evt-1"
assert entry["payload"] is not None
def test_secret_rejection_withholds_the_payload(tmp_path):
"""Storing the body of an event rejected for carrying secret-shaped
material would write that material into the audit store."""
app, _ = bound_app(tmp_path, secret_policy="reject", may_read=True)
assert invoke(app, event(data={"password": "hunter2"}))[0].startswith("400")
_, body = invoke(app, None, path="/v1/dead-letters", method="GET", body=b"")
entry = body["dead_letters"][0]
assert entry["reason"] == "secret_shaped_field"
assert entry["payload_withheld"] is True
assert entry["payload"] is None
assert entry["payload_hash"]
# --- redaction policy (T04) -------------------------------------------------
def test_default_policy_redacts_and_accepts(app):
"""Default is redact: losing the whole audit record over one field is
worse than storing it with that field masked."""
status, _ = invoke(app, event(data={"membership_id": "m-1", "auth_token": "s3cret"}))
assert status.startswith("202")
_, record = invoke(app, None, path="/v1/events/evt-1", method="GET", body=b"")
assert record["details"]["data"]["auth_token"] == "[redacted]"
assert record["details"]["data"]["membership_id"] == "m-1"
# The record must admit it was modified.
assert record["details"]["redaction"]["policy"] == "redact"
assert record["details"]["redaction"]["paths"] == ["data.auth_token"]
def test_reject_policy_is_available_per_sender(tmp_path):
app, _ = bound_app(tmp_path, secret_policy="reject")
status, body = invoke(app, event(data={"password": "x"}))
assert status.startswith("400")
assert body["error"] == "secret_shaped_field"
def test_nested_and_listed_secrets_are_redacted(app):
payload = event(data={"items": [{"api_secret": "a"}, {"ok": 1}], "n": {"private_key": "k"}})
assert invoke(app, payload)[0].startswith("202")
_, record = invoke(app, None, path="/v1/events/evt-1", method="GET", body=b"")
data = record["details"]["data"]
assert data["items"][0]["api_secret"] == "[redacted]"
assert data["items"][1]["ok"] == 1
assert data["n"]["private_key"] == "[redacted]"
assert set(record["details"]["redaction"]["paths"]) == {
"data.items[0].api_secret", "data.n.private_key",
}
def test_findings_are_counted_by_path_for_both_outcomes(tmp_path):
"""Counters name the field to fix, not just a total."""
app, backend = bound_app(tmp_path, may_read=True)
for i in range(3):
invoke(app, event(id=f"e{i}", data={"auth_token": "x"}), key=f"e{i}")
strict, _ = bound_app(tmp_path, secret_policy="reject")
invoke(strict, event(id="r1", data={"auth_token": "x"}), key="r1")
_, body = invoke(app, None, path="/v1/secret-findings", method="GET", body=b"")
rows = {(r["outcome"], r["field_path"]): r for r in body["secret_findings"]}
assert rows[("redacted", "data.auth_token")]["occurrences"] == 3
assert rows[("redacted", "data.auth_token")]["action"] == "membership.added"
assert rows[("redacted", "data.auth_token")]["source"] == "user-engine"
assert rows[("redacted", "data.auth_token")]["persisted"] is True
assert rows[("rejected", "data.auth_token")]["occurrences"] == 1
def test_counters_survive_restart(tmp_path):
"""The counters drive a fix in the sending service; that work outlives a
pod restart, so they are durable rather than in-memory."""
path = str(tmp_path / "counters.db")
first = IngestionApplication(SQLiteAuditBackend(path), "opaque")
invoke(first, event(data={"auth_token": "x"}))
reopened = IngestionApplication(SQLiteAuditBackend(path), "opaque")
_, body = invoke(reopened, None, path="/v1/secret-findings", method="GET", body=b"")
assert body["secret_findings"][0]["occurrences"] == 1
# --- deployment guards and counters (WP-0005-T03) ---------------------------
def test_required_custody_class_refuses_a_development_backend(tmp_path):
"""Losing AUDIT_CORE_DATABASE_URL must fail to start, not silently
downgrade custody to the development store."""
backend = SQLiteAuditBackend(str(tmp_path / "dev.db"))
with pytest.raises(ValueError, match="does not meet the required"):
IngestionApplication(backend, "opaque", require_custody_class="archive")
def test_counters_track_each_outcome(tmp_path):
app, _ = bound_app(tmp_path, may_read=True)
invoke(app, event()) # accepted
invoke(app, event()) # duplicate
invoke(app, event(subject="other")) # conflict
invoke(app, event(id="e2", tenant="tenant:coulomb"), key="e2") # rejected
invoke(app, event(), token="nope") # unauthorized
_, body = invoke(app, None, path="/v1/stats", method="GET", body=b"")
counts = body["counts"]
assert counts["accepted"] == 1
assert counts["duplicate"] == 1
assert counts["conflict"] == 1
assert counts["rejected"] == 1
assert counts["unauthorized"] == 1
assert body["since"]
def test_stats_require_the_read_privilege(tmp_path):
app, _ = bound_app(tmp_path, may_read=False)
status, _ = invoke(app, None, path="/v1/stats", method="GET", body=b"")
assert status.startswith("403")