diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md index 28379122..f8fc0b0c 100644 --- a/eventing/agentdocs/DESIGN_PHASE2.md +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -249,12 +249,13 @@ identity, which is §4's territory. --- -## 4. Agent identity: groundwork, not yet wired +## 4. Agent identity: wired The second question — *which agent produced this answer?* — matters more than it -first appears. Today anything with write access to the `responses` topic gets its -output stored, rendered in the HTML transcript, and pushed to the operator's -phone **as a legitimate agent answer**. +first appears. Before this, anything with write access to the `responses` topic +got its output stored, rendered in the HTML transcript, and pushed to the +operator's phone **as a legitimate agent answer**. With verification enabled it +lands as a rejection instead. ### 4.1 What exists @@ -273,14 +274,21 @@ phone **as a legitimate agent answer**. an unnamed token is ambiguous, and guessing would mean accepting a signature from *any* approved agent for an event that named none of them. -### 4.2 The blocker, stated plainly +### 4.2 The blocker, now cleared -**`sign_event()` has no production caller.** `ER_REQUIRE_SIGNATURE=true` today -rejects 100% of traffic — it is a kill switch, not a feature. -`IMPLEMENTATION_REPORT1.md` §811-814 already says "EventBridge does not sign." -Nothing about signed attributes means anything until that is fixed, which is why -`submitter` is **not** in `SIGNED_ATTRS` yet: adding it before the signing side -exists would only invalidate canonicalisation twice. +This section used to read "**`sign_event()` has no production caller**" — +`ER_REQUIRE_SIGNATURE=true` rejected 100% of traffic, a kill switch rather than a +feature. That is fixed. EventBridge signs the requests and group events it +publishes, EventRunner signs terminal responses, and each verifies the other's +output against the approved-key set. + +`submitter`, `submitteriss` and `groupid` joined `SIGNED_ATTRS` **after** the +signing side existed, in that order deliberately: covering them earlier would +have invalidated canonicalisation twice for no benefit. Nothing had ever +published a signed event, so the change needed no compatibility flag — but if +signing is already enabled somewhere, both services must be upgraded together, +because every grouped request now carries a signed `groupid` that an older +verifier omits when it recomputes. ### 4.3 Why not Keycloak, and why not HMAC @@ -298,27 +306,61 @@ upgrades cleanly: SPIRE later distributes the same keys rooted in workload attestation, and the verification code does not change — only where keys come from. -### 4.4 The remaining work - -1. Sign requests in `kafka_out.publish_request`, after `ce.new_event` (which - fills `id`/`time`, both signed) and before `to_kafka_binary`. -2. Sign responses in `emit.py` between event construction and serialisation. - Hold the seed **on the Emitter** — its constructor runs once, so none of the - seven `emit()` call sites change. Sign **terminal events only**: `emit()` is - on the hot path for every `stdout` frame and Ed25519 costs ~150 ms here. -3. Verify responses in `kafka_in.py` while the `CloudEvent` is still in hand. - Two hazards: the decode and store write are **not** inside a `try`, so a raise - kills the consumer thread — the verification path must degrade, never raise. - And group events are published by EventBridge itself, carry `groupid` but no - `correlationid`, so they need either their own `kid` or a skip. -4. Then add `submitter`, `submitteriss` and the missing `groupid` to - `SIGNED_ATTRS`. `DESIGN_PHASE1.md` §21.9.9 requires `groupid`; its absence - means signatures currently say nothing about batch membership. - -On failure, set `phase="error"`. That reuses machinery already wired: ntfy -**priority 5** with an error tag, a visually distinct bubble in the live SSE -transcript, and `raw_json` persistence for audit. The demo artifact costs -nothing extra. +### 4.4 How it is wired + +1. **Requests** are signed in `kafka_out.publish_request`, after `ce.new_event` + fills `id`/`time` (both signed) and before `to_kafka_binary`, so the signature + covers exactly what goes on the wire. +2. **Terminal responses only** are signed in `emit.py`. The seed lives on the + Emitter, whose constructor runs once, so none of the seven `emit()` call sites + changed. `emit()` runs per `stdout` frame and Ed25519 costs ~150-200 ms here, + so signing every frame would add minutes to a chatty run. + + **The verifier has to match that policy, and initially did not.** Verifying + all-or-nothing rewrote every streamed frame of every genuine run to + `phase="error"` — found by running it against a live broker, not by any unit + test, because every test until then used terminal events. So an event carrying + **no** signature and **not** terminal is passed through; an unsigned *terminal* + event is still refused, since that is the one the transcript presents as the + answer. A frame that presents a bad signature is still checked — the exemption + is for absence, not for failure. + + The honest limit that remains: this proves *who finished a run*, not *what it + said along the way*. Forged `final=false` frames still render. Closing that needs + either cheap signatures or a signed digest chain across frames. +3. **Requests are verified** in `consume.py`, keyed on the token's `kid` when a + keyset is configured and falling back to the single-key path otherwise. A + rejection commits the offset without running, because a bad signature is still + bad on redelivery, and increments `rejected_unsigned` so it is distinguishable + from the stale-request drop. +4. **Responses are verified** in `kafka_in.py` while the `CloudEvent` is still in + hand — the check needs `.attrs`/`.data`, which `envelope_dict` has flattened. + The decision itself is a pure function in `shared/signing.py`, which is what + makes this path testable at all: the consumer had no behavioural test before. +5. **Group events are signed** under EventBridge's own `kid`, and the verifier + accepts *only* that kid for them. Skipping them instead would have left the one + hole worth closing — a forged `group.completed` ends a batch early and fires a + "finished" notification for work that never ran. A flat keyset is not enough + here: any approved runner would otherwise do. +6. `submitter`, `submitteriss` and `groupid` joined `SIGNED_ATTRS` last (§4.2). + +**Both hazards in step 4 were real.** `from_kafka_binary` and `insert_response` +are not inside a `try`, so anything raising in the verification path would end the +consume loop for the life of the pod — silently, with the process still healthy. +The guard around the decision call fails closed only where enforcement is on: if +the verifier itself is broken, an unverifiable event is not evidence of anything. + +On failure the event is stored with `phase="error"` rather than dropped — a drop is +indistinguishable from an agent that never answered. That reuses machinery already +wired: ntfy **priority 5** with an error tag (the error body comes from +`data["text"]`, so the key name is load-bearing), a red card in the live SSE +transcript, and `raw_json` persistence for audit. + +**Rollout is two flags.** A keyset alone verifies and logs while storing events +unchanged; `EB_REQUIRE_RESPONSE_SIGNATURE=true` is what rewrites them. Enforcement +mutates persisted rows and pages a phone, so there is a step where the reject rate +is observable first. EventBridge refuses to start with a keyset but no +`EB_SIGNING_KID`, since there would be nothing to attribute a group event to. --- @@ -332,9 +374,28 @@ nothing extra. | `EB_AUTH_TOKENS` | empty | `name:token,...` fallback. Empty plus no GitHub config means auth is off. | | `EVENTBRIDGE_TOKEN` | unset | CLI: overrides the stored token, so a shell can act as another identity. | -The client id lives in `config.toml` because it is not a capability. `test_manifests.py` -pins the opposite rule for `NTFY_TOPIC`/`NTFY_TOKEN`, and that distinction is the -point: one is public by construction, the others grant access. +Agent identity (§4). Every one of these is off by default, so the e2e path is +unaffected until an operator opts in: + +| Variable | Default | Effect | +|---|---|---| +| `EB_SIGNING_KEY_PATH` | empty | Ed25519 seed EventBridge signs requests and group events with. Empty = no signing. | +| `EB_SIGNING_KID` | empty | Names EventBridge's key. Also the **only** kid accepted on group lifecycle events. | +| `EB_VERIFY_KEYSET_PATH` | empty | Approved-key set for responses. Empty = no verification. | +| `EB_REQUIRE_RESPONSE_SIGNATURE` | `false` | `false` = verify and log (audit); `true` = rewrite failures to `phase=error`. | +| `ER_SIGNING_KEY_PATH` | empty | Seed EventRunner signs terminal responses with. | +| `ER_SIGNING_KID` | empty | Names this runner's key. Needed once more than one runner is approved. | +| `ER_REQUIRE_SIGNATURE` | `false` | Refuse unsigned or badly-signed requests. | +| `ER_VERIFY_KEYSET_PATH` | empty | Approved-key set for requests. Empty falls back to `ER_VERIFY_KEY_PATH`'s single key. | + +Seed paths name **Secret** mounts and are env-only, never `config.toml`. Keyset +paths name **ConfigMaps** — only public keys belong in them, and `keyset.load()` +cannot tell a seed from a public key by length, so the file's location is what keeps +the distinction reviewable. + +The GitHub client id lives in `config.toml` because it is not a capability. +`test_manifests.py` pins the opposite rule for `NTFY_TOPIC`/`NTFY_TOKEN`, and that +distinction is the point: one is public by construction, the others grant access. --- diff --git a/eventing/agentdocs/IMPLEMENTATION_REPORT1.md b/eventing/agentdocs/IMPLEMENTATION_REPORT1.md index edd03a1a..8c9b768a 100644 --- a/eventing/agentdocs/IMPLEMENTATION_REPORT1.md +++ b/eventing/agentdocs/IMPLEMENTATION_REPORT1.md @@ -812,6 +812,9 @@ That is RQ-1 behaving exactly as designed, observed by accident. tested against the RFC vectors and the flag is wired through `consume.py`, but no publisher signs requests yet, so the verify-and-reject path has not been exercised on a cluster. EventBridge does not sign. + *(Resolved in Phase 2: EventBridge signs requests and group events, EventRunner + signs terminal responses, and both verify against the approved-key set. The + reject paths now have tests. See `DESIGN_PHASE2.md` §4.)* - **Resource requests and limits are still placeholders** (§18). A real run has not produced numbers. - **Single-broker, ephemeral storage.** A broker restart loses all topic data, and @@ -824,7 +827,7 @@ That is RQ-1 behaving exactly as designed, observed by accident. ## 9. Signing: the cost of the pure-Python rule §1.1 bans C extensions, which rules out `cryptography`, so Ed25519 is implemented -from RFC 8032 in `eventrunner/signing.py` (~120 lines) and checked against the +from RFC 8032 in `shared/signing.py` (~120 lines) and checked against the RFC's own test vectors — all three pass for key derivation, signing and verification, plus tamper, wrong-key, malformed-input and `alg` confusion cases. diff --git a/eventing/agentdocs/README.md b/eventing/agentdocs/README.md index 1cab8b21..1c9c7ecf 100644 --- a/eventing/agentdocs/README.md +++ b/eventing/agentdocs/README.md @@ -18,7 +18,7 @@ and what is *not* verified. | [`DESIGN_PHASE1.md`](DESIGN_PHASE1.md) | The Phase 1 design — a delta over Phase 0, not a replacement: KEDA scaling on consumer lag, scale-to-zero, the three cluster findings that shaped it, the §16 gaps (A: rebalance floor, B: ephemeral session state, C: idle replay), §3.2 on supporting a local Kind cluster, and §21 designing agent **groups** — batch fan-out with a tracked fan-in, its own page and exactly two notifications. | | [`IMPLEMENTATION_REPORT1.md`](IMPLEMENTATION_REPORT1.md) | What Phase 1 built, the measured results on both clusters, 28 findings including two real bugs only a scale-to-zero deployment could expose, and an honest account of what is still blocked and why. | | [`README_PHASE1.md`](README_PHASE1.md) | The Phase 1 runbook in full detail, covering all four ways to run it: locally without containers, under Docker, on a Kind cluster, and on OpenShift. | -| [`DESIGN_PHASE2.md`](DESIGN_PHASE2.md) | The Phase 2 design — identity on the event path. Why GitHub's opaque user token rules out local verification and forces a `GET /user` lookup plus a load-bearing cache; why `401` and `403` are kept distinct; what `ce_submitter` is and is not worth while it remains unsigned; the `kid` and approved-key-set groundwork for proving *which agent* answered, and the blocker that nothing signs yet. | +| [`DESIGN_PHASE2.md`](DESIGN_PHASE2.md) | The Phase 2 design — identity on the event path. Why GitHub's opaque user token rules out local verification and forces a `GET /user` lookup plus a load-bearing cache; why `401` and `403` are kept distinct; what `ce_submitter` is and is not worth, and what a signature over it does and does not prove; how the `kid` and the approved-key set prove *which agent* answered, why group events are pinned to EventBridge's own key, and why a rejected response is stored as an error rather than dropped. | ## Reading order diff --git a/eventing/eventbridge/__main__.py b/eventing/eventbridge/__main__.py index 44546917..146f3ecf 100644 --- a/eventing/eventbridge/__main__.py +++ b/eventing/eventbridge/__main__.py @@ -21,6 +21,7 @@ from eventbridge.openapi import spec from eventbridge.router import Dispatcher from eventbridge.store import Store +from shared import keyset, signing from shared.pidfile import PidFile @@ -46,10 +47,40 @@ def main() -> int: for corr in store.all_correlations(limit=10000): minter.remember(corr) + # §11 — key material is loaded ONCE, here, and deliberately not caught. A bridge + # that believes it is signing but is not fails silently; one that will not start + # says so in `kubectl logs` before it accepts a request. + seed = None + if cfg.signing_key_path: + seed = signing.load_seed(cfg.signing_key_path) + print(f"[eventbridge] request signing ON as kid={cfg.signing_kid or '(unnamed)'} " + f"(seed {cfg.signing_key_path})") + else: + print("[eventbridge] request signing OFF (set EB_SIGNING_KEY_PATH to enable)") + + # Loaded once, not per record: live reload is deliberately absent (keyset.py), so a + # mid-run edit must not silently widen the set of agents this bridge trusts. + ks = keyset.load_if_set(cfg.verify_keyset_path) + if ks is not None: + mode = "ENFORCING" if cfg.require_response_signature else "audit only" + print(f"[eventbridge] response verification ON ({mode}) — {len(ks)} kid(s): " + f"{', '.join(ks.kids)}") + if not cfg.signing_kid: + # Without a bridge kid there is nothing to compare a group event against, + # so any approved runner could forge one. Refuse rather than accept-any. + raise SystemExit( + "[eventbridge] EB_VERIFY_KEYSET_PATH is set but EB_SIGNING_KID is not: " + "group lifecycle events could not be attributed to this bridge. " + "Set EB_SIGNING_KID to the kid naming EventBridge's key.") + else: + print("[eventbridge] response verification OFF " + "(set EB_VERIFY_KEYSET_PATH to enable)") + # The producer needs the RESPONSES topic too: group lifecycle events go there, # not on requests, because EventRunner would try to execute anything on requests. producer = Producer(cfg.kafka_bootstrap, cfg.request_topic, cfg.source_uri, - response_topic=cfg.response_topic) + response_topic=cfg.response_topic, + seed=seed, kid=cfg.signing_kid or None) groups = GroupService(cfg, store, producer, minter) ntfy = NtfyPublisher(cfg.ntfy, cfg.public_base_url, store=store) @@ -59,7 +90,10 @@ def main() -> int: consumer = Consumer(cfg.kafka_bootstrap, cfg.response_topic, store, on_event=ntfy.submit, on_group_event=groups.on_group_event, - on_member_event=groups.on_member_event) + on_member_event=groups.on_member_event, + keyset=ks, + require_signature=cfg.require_response_signature, + bridge_kid=cfg.signing_kid or None) consumer.start() # Back-fill prompts from the requests topic — also gives us prompt visibility diff --git a/eventing/eventbridge/auth.py b/eventing/eventbridge/auth.py index 3d13e6c2..da05d5ed 100644 --- a/eventing/eventbridge/auth.py +++ b/eventing/eventbridge/auth.py @@ -18,12 +18,21 @@ onto the request event as `ce_submitter`, which is the whole point — a `401` tells you nothing after the fact, a recorded submitter does. -What this is NOT: the submitter attribute is **unsigned**. Anyone who can write -to the Kafka `requests` topic can forge it, and the broker is plaintext. The -honest claim is "EventBridge refuses unauthenticated submissions and records who -it believes submitted this" — not "this event proves who submitted it." Proving -it needs `submitter` inside `signing.SIGNED_ATTRS` and a producer that actually -signs (today nothing does; see agentdocs/IMPLEMENTATION_REPORT1.md §811-814). +What this is NOT, and the two limits that still apply. `submitter` is now inside +`signing.SIGNED_ATTRS` and EventBridge signs the requests it publishes, so when +signing is configured the attribute cannot be altered in flight without invalidating +the signature. But: + +* **Signing is opt-in.** With no `EB_SIGNING_KEY_PATH` the events are unsigned, the + broker is plaintext, and anyone who can write to the `requests` topic can forge a + submitter. The claim only holds where verification is actually enabled. +* **A signature proves the assertion, not the identity.** It shows EventBridge said + this, not that the name is real. A static `EB_AUTH_TOKENS` entry is a name an + operator typed into an env var; `ce_submitteriss` is what distinguishes it from a + login GitHub verified. + +So the honest claim is "EventBridge refuses unauthenticated submissions, records who +it believes submitted, and — when signing is on — makes that record tamper-evident." """ from __future__ import annotations diff --git a/eventing/eventbridge/config.py b/eventing/eventbridge/config.py index 73e9d879..e9dcad90 100644 --- a/eventing/eventbridge/config.py +++ b/eventing/eventbridge/config.py @@ -58,6 +58,25 @@ class Cfg: # GitHub API call per request would spend a 5000/hour budget and add GitHub's # latency to the request path. github_cache_ttl_s: float = 300.0 + # §11 — signing the requests and group events EventBridge publishes. Empty + # disables it, which is the default: a key is a capability, and the demo has to + # work with none. Env only, never config.toml — the file it points at is a + # Secret mount (same rule as auth_tokens and NTFY_TOKEN above). + signing_key_path: str = "" + # The kid written into the protected header, naming EventBridge's own key in the + # approved set. Also the kid group lifecycle events must be signed by: the set is + # otherwise flat, so without this any approved runner could forge a + # group.completed and end a batch early. + signing_kid: str = "" + # §11 — the approved-key set used to verify responses coming back off the + # responses topic. Empty disables verification entirely. A ConfigMap path, not a + # Secret: only public keys belong in it. + verify_keyset_path: str = "" + # Whether a failed verification is ENFORCED. With a keyset but this false, + # EventBridge verifies and logs but stores the event unchanged — audit mode. + # Enforcement rewrites persisted rows and raises a priority-5 notification, so + # there has to be a way to watch the reject rate before turning it on. + require_response_signature: bool = False ntfy: NtfyCfg = field(default_factory=NtfyCfg) @@ -120,6 +139,15 @@ def load() -> Cfg: cfg.allowed_users = ghauth.parse_allowed_users(allowed_env) cfg.github_cache_ttl_s = float(e("EB_GITHUB_CACHE_TTL_S", str(cfg.github_cache_ttl_s))) + # §11. Env only: the seed path names a Secret mount, and keeping the whole + # signing block in one layer means it cannot be half-configured from a + # committed file. + cfg.signing_key_path = e("EB_SIGNING_KEY_PATH", cfg.signing_key_path) + cfg.signing_kid = e("EB_SIGNING_KID", cfg.signing_kid) + cfg.verify_keyset_path = e("EB_VERIFY_KEYSET_PATH", cfg.verify_keyset_path) + cfg.require_response_signature = ( + e("EB_REQUIRE_RESPONSE_SIGNATURE", + "true" if cfg.require_response_signature else "false").lower() == "true") cfg.ntfy.enabled = (e("NTFY_ENABLED", "true" if cfg.ntfy.enabled else "false").lower() == "true") cfg.ntfy.base_url = e("NTFY_BASE_URL", cfg.ntfy.base_url) diff --git a/eventing/eventbridge/kafka_in.py b/eventing/eventbridge/kafka_in.py index a0dd3ff6..73ffe880 100644 --- a/eventing/eventbridge/kafka_in.py +++ b/eventing/eventbridge/kafka_in.py @@ -1,4 +1,22 @@ -"""KafkaConsumer thread — polls responses, writes SQLite, fans out to ntfy.""" +"""KafkaConsumer thread — polls responses, writes SQLite, fans out to ntfy. + +§11: this is where a response is checked against the approved-key set. Anything with +write access to the responses topic otherwise gets its output stored, rendered in the +transcript and pushed to the operator's phone **as a legitimate agent answer**. + +A rejected event is stored as `phase="error"` rather than dropped. That is deliberate: +dropping it silently is indistinguishable from an agent that never answered, while +`phase="error"` reuses machinery already wired — a red card in the HTML transcript and +an ntfy priority-5 alert — so the forgery attempt is visible instead of invisible. + +**What "reviewable" requires.** `Store.insert_response` derives both `data_json` and +`raw_json` from the single dict it is handed, so rewriting the envelope in place would +overwrite the evidence with the notice about it: the refused payload, the phase it +claimed, and the attributes the signature covered would all be gone, leaving a +signature that can no longer be checked against anything. The original is therefore +preserved under `data["rejected"]` — attrs, payload, claimed source and signature — so +an incident can be verified offline rather than merely logged. +""" from __future__ import annotations import threading @@ -7,7 +25,7 @@ from kafka import KafkaConsumer from eventbridge.store import Store -from shared import ce +from shared import ce, signing class Consumer(threading.Thread): @@ -20,6 +38,9 @@ def __init__( group_id: str = "eventbridge-responses", on_group_event: Callable[[dict], None] | None = None, on_member_event: Callable[[dict], None] | None = None, + keyset=None, + require_signature: bool = False, + bridge_kid: str | None = None, ) -> None: super().__init__(daemon=True, name="kafka-responses-consumer") self._bootstrap_servers = bootstrap @@ -29,8 +50,20 @@ def __init__( self._on_group_event = on_group_event self._on_member_event = on_member_event self._group = group_id + # §11. `keyset=None` means verification is off, which is the default and + # exactly today's behaviour. `bridge_kid` pins group lifecycle events to + # EventBridge's own key, so an approved runner cannot forge a group.completed. + self._keyset = keyset + self._require_sig = require_signature + self._bridge_kid = bridge_kid + self._rejected = 0 self._stopping = threading.Event() + @property + def rejected(self) -> int: + """Responses stored as phase=error because they did not verify.""" + return self._rejected + def stop(self) -> None: self._stopping.set() @@ -49,7 +82,54 @@ def run(self) -> None: if self._stopping.is_set(): break evt = ce.from_kafka_binary(rec.headers or [], rec.value) + # §11: verify while the CloudEvent is still in hand — the check + # needs `.attrs`/`.data`, which envelope_dict has already flattened. + # + # The try is not belt-and-braces. `from_kafka_binary` above and + # `insert_response` below are NOT inside one, so a raise anywhere in + # here ends the for, ends the while, and the thread is gone — while + # the process stays up and the pod still reports healthy. The + # verification path must degrade, never raise. + ok, why = True, "not checked" + if self._keyset is not None: + try: + ok, why = signing.response_decision( + evt, self._keyset, require=self._require_sig, + bridge_kid=self._bridge_kid) + except Exception as e: # noqa: BLE001 + # Fail closed only where enforcement is on: if the verifier + # itself is broken, an unverifiable event is not evidence of + # anything, and silently accepting it defeats the control. + ok, why = (not self._require_sig), f"verifier raised: {e!r}" + print(f"[kafka_in] verification error: {e!r}") d = ce.envelope_dict(evt) + if not ok: + self._rejected += 1 + print(f"[kafka_in] unverified response on " + f"{d.get('correlationid') or d.get('groupid')}: {why}") + # `text` is load-bearing: ntfy reads data["text"] for the error + # body, so anything else shows up on the phone as + # "(error, see raw)". str() because insert_response json.dumps + # this dict outside any try — a non-serialisable reason would + # kill the thread by a second route. + # + # `rejected` carries the event as it actually arrived. + # `insert_response` derives BOTH data_json and raw_json from + # this one dict, so overwriting `phase`/`data` in place would + # destroy the forensic record while the docstring above still + # promised it — leaving a signature whose covered attributes no + # longer exist, and nothing for an operator to review. + d = dict(d, phase="error", data={ + "text": f"unverified response rejected: {why}", + "signature_rejected": True, + "reason": str(why), + "rejected": {"phase": evt.get("phase"), + "final": evt.get("final"), + "source": evt.get("source"), + "signature": evt.get("signature"), + "attrs": dict(evt.attrs), + "data": evt.data}, + }) # §21.2: route on type. A group lifecycle event carries `groupid` # but no `correlationid`, so handing it to insert_response would # violate that table's (correlationid, sequence) primary key. diff --git a/eventing/eventbridge/kafka_out.py b/eventing/eventbridge/kafka_out.py index 5551740d..2e0aa649 100644 --- a/eventing/eventbridge/kafka_out.py +++ b/eventing/eventbridge/kafka_out.py @@ -3,16 +3,42 @@ from kafka import KafkaProducer -from shared import ce +from shared import ce, signing class Producer: + """Publishes requests, and group lifecycle events onto the responses topic. + + §11: when a seed is configured, every event published here is signed under `kid` + before serialisation, so EventRunner can refuse a request from anything that is + not an approved submitter. `kid` rides in the JWS protected header, which is + signed input — it cannot be swapped to impersonate another key. + + **Cost.** The Ed25519 here is pure Python (`cryptography` is a C extension and is + banned by §1.1) and takes ~150-200 ms per signature. `publish_request` is called + once per group member by `GroupService.submit_members`, in a loop, inside one HTTP + request — so a 100-member batch spends ~20 s signing while the caller waits. That + is accepted: a batch launch is an operator action, not a hot path. It is recorded + here rather than discovered later, and a thread pool would not fix it (the GIL + serialises pure-Python signing anyway). If signing ever becomes mandatory at + scale, the fix is the `cryptography` dependency conversation, not concurrency. + + The number an operator actually hits: ~200 members is ~40 s, past a common 30 s + client timeout. `POST /v0/groups` honours an `Idempotency-Key` header, so the + retry after such a timeout returns the original batch instead of launching a + second one — the failure is survivable, but only if the caller sends the header. + """ + def __init__(self, bootstrap: str, request_topic: str, source_uri: str, - response_topic: str | None = None) -> None: + response_topic: str | None = None, + seed: bytes | None = None, kid: str | None = None) -> None: self._prod = KafkaProducer(bootstrap_servers=bootstrap, acks="all", linger_ms=5) self._topic = request_topic self._response_topic = response_topic self._source = source_uri + # Held on the instance so no call site has to know about signing. + self._seed = seed + self._kid = kid def publish_request( self, @@ -41,6 +67,9 @@ def publish_request( **({ce.EXT_SUBMITTER: submitter} if submitter else {}), **({ce.EXT_SUBMITTER_ISS: submitter_iss} if submitter_iss else {}), ) + # After new_event (which fills `id` and `time`, both signed) and before + # serialisation, so the signature covers exactly what goes on the wire. + signing.sign_into(event, self._seed, self._kid) headers, value = ce.to_kafka_binary(event) future = self._prod.send(self._topic, key=correlationid.encode(), value=value, headers=headers) future.get(timeout=5) @@ -54,6 +83,11 @@ def publish_group_event(self, *, type_: str, groupid: str, an agent run to execute. Responses is also where EventBridge's own consumer and ntfy publisher already listen, so the event gets stored, notified and audited with no new plumbing. + + Signed under the same `kid` as requests, and the responses-side verifier + accepts *only* that kid here. Skipping these instead would leave the one hole + worth closing: a forged `group.completed` ends a batch early and fires a + "finished" notification for work that never ran. """ if not self._response_topic: raise RuntimeError("Producer has no response_topic; cannot publish group events") @@ -65,6 +99,7 @@ def publish_group_event(self, *, type_: str, groupid: str, groupid=groupid, data=data, ) + signing.sign_into(event, self._seed, self._kid) headers, value = ce.to_kafka_binary(event) self._prod.send(self._response_topic, key=groupid.encode(), value=value, headers=headers).get(timeout=5) diff --git a/eventing/eventrunner/__main__.py b/eventing/eventrunner/__main__.py index 76bb81ba..3b70d159 100644 --- a/eventing/eventrunner/__main__.py +++ b/eventing/eventrunner/__main__.py @@ -11,6 +11,7 @@ from eventrunner.router import Router from eventrunner.runner import log_forwarded_env, run_agent from eventrunner.transcript import TranscriptStore +from shared import keyset, signing from shared.heartbeat import Heartbeat from shared.pidfile import PidFile @@ -55,13 +56,38 @@ def main() -> int: heartbeat = Heartbeat(cfg.heartbeat_path) heartbeat.touch() # before Kafka, so an absent file always means "never started" - emitter = Emitter(cfg.kafka_bootstrap, cfg.response_topic, cfg.source_uri) + # §11 — key material is loaded ONCE, here, and deliberately not caught. A runner + # that believes it is signing but is not, or that cannot read the set it verifies + # against, is worse than one that refuses to start: the first failure mode is + # silent and the second is in `kubectl logs` before any request is accepted. + seed = None + if cfg.signing_key_path: + seed = signing.load_seed(cfg.signing_key_path) + print(f"[eventrunner] response signing ON as kid={cfg.signing_kid or '(unnamed)'} " + f"(seed {cfg.signing_key_path}) — terminal events only") + else: + print("[eventrunner] response signing OFF (set ER_SIGNING_KEY_PATH to enable)") + + # Load once, never per record: keyset.py states that live reload is deliberately + # absent, and re-reading per event would let a mid-run edit silently widen the + # set of agents this runner trusts. + ks = keyset.load_if_set(cfg.verify_keyset_path) + if ks is not None: + print(f"[eventrunner] request verification uses the approved set " + f"{cfg.verify_keyset_path} — {len(ks)} kid(s): {', '.join(ks.kids)}") + elif cfg.require_signature: + print("[eventrunner] request verification uses a single key " + f"({cfg.verify_key_path or cfg.signing_key_path}); " + "set ER_VERIFY_KEYSET_PATH for an approved-agent list") + + emitter = Emitter(cfg.kafka_bootstrap, cfg.response_topic, cfg.source_uri, + seed=seed, kid=cfg.signing_kid or None) def _run(event): run_agent(cfg, emitter, event, transcripts=transcripts) router = Router(_run, cfg.max_concurrent) - consumer = Consumer(cfg, router, heartbeat=heartbeat) + consumer = Consumer(cfg, router, heartbeat=heartbeat, keyset=ks) consumer.start() stop_evt = threading.Event() diff --git a/eventing/eventrunner/config.py b/eventing/eventrunner/config.py index 499446b7..4c8e2704 100644 --- a/eventing/eventrunner/config.py +++ b/eventing/eventrunner/config.py @@ -96,6 +96,14 @@ class Cfg: require_signature: bool = False signing_key_path: str = "" verify_key_path: str = "" + # The approved-key set: kid -> public key. When set, request verification + # selects a key by the token's kid instead of assuming one key, which is what + # makes it an allowlist of several agents rather than a single-key check. Empty + # falls back to verify_key_path, so the older one-key deployment is unchanged. + verify_keyset_path: str = "" + # The kid this runner names in the responses it signs. Empty still signs, and a + # verifier holding exactly one key accepts it; name it as soon as there are two. + signing_kid: str = "" def load() -> Cfg: @@ -140,9 +148,11 @@ def load() -> Cfg: cfg.eventbridge_url = e("ER_EVENTBRIDGE_URL", cfg.eventbridge_url) cfg.transcript_max_bytes = int(e("ER_TRANSCRIPT_MAX_BYTES", str(cfg.transcript_max_bytes))) - cfg.require_signature = _bool(e("ER_REQUIRE_SIGNATURE"), False) - cfg.signing_key_path = e("ER_SIGNING_KEY_PATH", cfg.signing_key_path) - cfg.verify_key_path = e("ER_VERIFY_KEY_PATH", cfg.verify_key_path) + cfg.require_signature = _bool(e("ER_REQUIRE_SIGNATURE"), False) + cfg.signing_key_path = e("ER_SIGNING_KEY_PATH", cfg.signing_key_path) + cfg.verify_key_path = e("ER_VERIFY_KEY_PATH", cfg.verify_key_path) + cfg.verify_keyset_path = e("ER_VERIFY_KEYSET_PATH", cfg.verify_keyset_path) + cfg.signing_kid = e("ER_SIGNING_KID", cfg.signing_kid) base = e("TMPDIR", "/tmp").rstrip("/") cfg.tmpdir = f"{base}/rossoctl-keda1" diff --git a/eventing/eventrunner/consume.py b/eventing/eventrunner/consume.py index 73683db7..8f267b62 100644 --- a/eventing/eventrunner/consume.py +++ b/eventing/eventrunner/consume.py @@ -66,7 +66,8 @@ class Consumer(threading.Thread): def __init__(self, cfg: Cfg, router: Router, *, heartbeat: Heartbeat | None = None, ledger: OffsetLedger | None = None, - consumer_factory=None) -> None: + consumer_factory=None, + keyset=None) -> None: super().__init__(daemon=True, name="kafka-requests-consumer") self._cfg = cfg self._router = router @@ -90,6 +91,11 @@ def __init__(self, cfg: Cfg, router: Router, *, self._seeded: set = set() self._connect_attempts = 0 self._skipped_stale = 0 + # §11: the approved-agent set, or None to fall back to the single-key path. + # Loaded once by __main__ — never re-read here, so the trust set cannot widen + # mid-run. + self._keyset = keyset + self._rejected_unsigned = 0 self._paused = False # ---- lifecycle ---- @@ -108,6 +114,13 @@ def ledger(self) -> OffsetLedger: def skipped_stale(self) -> int: return self._skipped_stale + @property + def rejected_unsigned(self) -> int: + """Requests refused for a missing or bad signature. Distinct from + `skipped_stale`: both commit the offset without running, and an operator + needs to know which one is happening.""" + return self._rejected_unsigned + def _build_consumer(self) -> KafkaConsumer: return KafkaConsumer( bootstrap_servers=self._cfg.kafka_bootstrap, @@ -271,9 +284,12 @@ def _handle(self, rec) -> None: return if self._cfg.require_signature: - from eventrunner import signing - ok, why = signing.verify_event(evt, self._cfg) + from shared import signing + # With a keyset this is an allowlist keyed on the token's `kid`; without + # one it stays the single-key check it has always been. + ok, why = signing.verify_request(evt, self._cfg, self._keyset) if not ok: + self._rejected_unsigned += 1 _elog(f"rejecting unsigned/badly-signed request corr={corr}: {why} " f"(ER_REQUIRE_SIGNATURE=true) — offset committed, not retried") self._skip(topic, partition, offset) diff --git a/eventing/eventrunner/emit.py b/eventing/eventrunner/emit.py index 714dc839..f127981a 100644 --- a/eventing/eventrunner/emit.py +++ b/eventing/eventrunner/emit.py @@ -6,14 +6,35 @@ from kafka import KafkaProducer -from shared import ce +from shared import ce, signing class Emitter: - def __init__(self, bootstrap: str, response_topic: str, source_uri: str) -> None: + """§11: signs **terminal events only**, and what that does and does not prove. + + `emit()` runs for every `stdout` frame an agent produces, and the pure-Python + Ed25519 here costs ~150-200 ms per signature, so signing every frame would add + minutes to a chatty run. Terminal events are one per run, where the cost is + invisible against an agent that already took seconds. + + The honest consequence: a signature on the terminal event proves **who finished a + run**, not **what the run said along the way**. Anything with write access to the + responses topic can still forge `final=false` frames for a live correlation, and + they will render in the transcript. Closing that needs either cheap signatures or + a signed digest chain across frames; neither is in scope here, and claiming + otherwise while demoing would be wrong. + + The seed lives on the instance, so none of the seven `emit()` call sites in + `runner.py` know signing exists. + """ + + def __init__(self, bootstrap: str, response_topic: str, source_uri: str, + seed: bytes | None = None, kid: str | None = None) -> None: self._prod = KafkaProducer(bootstrap_servers=bootstrap, acks="all", linger_ms=5) self._topic = response_topic self._source = source_uri + self._seed = seed + self._kid = kid self._seq_lock = threading.Lock() self._seq_by_corr: dict[str, int] = {} @@ -56,6 +77,10 @@ def emit(self, *, correlationid: str, sessionuuid: str, sequence: int, data=data, **attrs, ) + # Terminal events only — see the class docstring for the cost and the limit + # that buys. `final` is already the parameter, so no call site changes. + if final: + signing.sign_into(event, self._seed, self._kid) headers, value = ce.to_kafka_binary(event) self._prod.send(self._topic, key=correlationid.encode(), diff --git a/eventing/k8s/base/configmap.yaml b/eventing/k8s/base/configmap.yaml index 1f81c568..e9a8aef9 100644 --- a/eventing/k8s/base/configmap.yaml +++ b/eventing/k8s/base/configmap.yaml @@ -37,3 +37,20 @@ data: ER_HEARTBEAT_MAX_AGE_S: "90" # §11, feature-flagged off: the e2e path is unaffected. ER_REQUIRE_SIGNATURE: "false" + # §11 — the approved-key set, kid -> Ed25519 public key. This file IS the + # authorization list: an event whose kid is not in it is refused. A ConfigMap + # path, never a Secret; only public keys belong in it. Empty disables the + # kid-aware path and falls back to ER_VERIFY_KEY_PATH's single key. + ER_VERIFY_KEYSET_PATH: "" + EB_VERIFY_KEYSET_PATH: "" + # The kid each service names in what it signs. EB_SIGNING_KID is also the only + # kid accepted on group lifecycle events, so an approved runner cannot forge a + # group.completed and end a batch early. + ER_SIGNING_KID: "" + EB_SIGNING_KID: "" + # Audit before enforce. With a keyset configured but this false, EventBridge + # verifies responses and logs what failed while storing them unchanged; true + # rewrites a failure to phase=error, which is persisted and raises a + # priority-5 notification. Seed paths (ER_/EB_SIGNING_KEY_PATH) are NOT here — + # they name Secret mounts and belong in the Deployment beside the volume. + EB_REQUIRE_RESPONSE_SIGNATURE: "false" diff --git a/eventing/shared/ce.py b/eventing/shared/ce.py index 6bc77427..13a9432f 100644 --- a/eventing/shared/ce.py +++ b/eventing/shared/ce.py @@ -41,10 +41,15 @@ EXT_SIGNATURE = "signature" # Phase 1 §21: the batch a correlation belongs to. At most one per correlation. EXT_GROUPID = "groupid" -# The authenticated caller that submitted this request, from EB_AUTH_TOKENS. -# Absent when auth is disabled. NOT in signing.SIGNED_ATTRS, so it is unsigned -# and forgeable by anyone with write access to the requests topic — it records -# who EventBridge believes submitted, not cryptographic proof. +# The authenticated caller that submitted this request, from EB_AUTH_TOKENS or a +# GitHub sign-in. Absent when auth is disabled. +# +# In signing.SIGNED_ATTRS, so when EventBridge signs, this attribute cannot be +# changed in flight without invalidating the signature. What that proves is bounded: +# *EventBridge asserted this submitter*, not that the submitter is who they claim — +# that is EXT_SUBMITTER_ISS's job. On an UNSIGNED event it remains forgeable by +# anyone with write access to the requests topic, which is why verification has to be +# enabled for it to mean anything at all. EXT_SUBMITTER = "submitter" # Who vouched for EXT_SUBMITTER: "github" when the login came from a verified # GitHub sign-in, absent when it came from a static token. Without this a reader diff --git a/eventing/shared/keyset.py b/eventing/shared/keyset.py index b9cd7ba5..2f081bb0 100644 --- a/eventing/shared/keyset.py +++ b/eventing/shared/keyset.py @@ -1,16 +1,10 @@ """The approved-key set: `kid` -> Ed25519 public key. Stdlib only. -**Not yet enforced.** Nothing consults this module today — `signing.verify_event` -resolves a single key from `ER_VERIFY_KEY_PATH`/`ER_SIGNING_KEY_PATH` and ignores -`kid` entirely. The primitives and their tests are in place so that wiring it in is -a small change; it also needs a config field for the path, which does not exist yet. -Until then `kid` is signature-covered but unused: a token naming an unknown key id -still verifies, because nothing selects a key by it. - -Once wired, this file **is** the authorization list: a signed event whose `kid` is -not in the set is refused, and so is one carrying no signature at all. Approving an -agent means adding its public key here; revoking it means removing the entry. That -is a deliberate design choice rather than a placeholder for something richer: +`signing.verify_with_keyset` consults this module, so this file **is** the +authorization list: a signed event whose `kid` is not in the set is refused, and so is +one carrying no signature at all. Approving an agent means adding its public key here; +revoking it means removing the entry. That is a deliberate design choice rather than a +placeholder for something richer: * It needs no issuer, no network call, and no clock, so it works identically on a laptop and in a cluster. diff --git a/eventing/eventrunner/signing.py b/eventing/shared/signing.py similarity index 58% rename from eventing/eventrunner/signing.py rename to eventing/shared/signing.py index 2a2893cb..bf939e82 100644 --- a/eventing/eventrunner/signing.py +++ b/eventing/shared/signing.py @@ -1,7 +1,9 @@ """Detached JWS over the CloudEvent envelope. DESIGN_PHASE1.md §11. -Feature-flagged and **disabled by default**, so the e2e path is unaffected and -signing can be enabled independently of causation binding. +Feature-flagged and **off by default**, so the e2e path is unaffected. Both services +use this module: EventBridge signs the requests and group events it publishes, +EventRunner signs terminal responses, and each verifies what the other produced +against `shared.keyset` — which is the authorization list, not merely a key lookup. Two design constraints shape this: @@ -31,8 +33,11 @@ import hashlib import json import pathlib +import sys from typing import Any +from shared import ce + # ---- Ed25519 (RFC 8032), pure Python --------------------------------------- _P = 2 ** 255 - 19 @@ -76,11 +81,18 @@ def _scalar_mult(p: tuple[int, int], e: int) -> tuple[int, int]: `if e & 1` branches on secret bits when `e` is the secret scalar from `_secret_scalar`, and `_edwards_add` does a modular inversion per addition, so - wall-clock time varies with the scalar's Hamming weight. That is acceptable - *here* only because of where this runs: signing is feature-flagged off - (`ER_REQUIRE_SIGNATURE=false`), the seed never leaves the pod, and nothing - exposes a remote timing oracle over `sign()`. Verification uses only public - inputs, so it is not the sensitive direction. + wall-clock time varies with the scalar's Hamming weight. Verification uses only + public inputs, so it is not the sensitive direction. + + **This justification weakened when signing got production callers.** It used to + rest on "nothing exposes a remote timing oracle over `sign()`", which is no longer + strictly true: with `EB_SIGNING_KEY_PATH` set, EventBridge signs on the HTTP + request path, so a caller who can time `POST /v0/agents` observes something + correlated with the scalar. What still makes it acceptable is that the signal is + buried under a Kafka round trip and a ~150 ms pure-Python operation whose variance + dwarfs the leak, the seed never leaves the pod, and signing remains opt-in. It is + a real if impractical weakness rather than a non-issue — do not promote this to a + trust boundary that assumes constant time. If the pure-Python constraint (§1.1) is ever relaxed, `cryptography`'s Ed25519 is the better trade than hardening this by hand. @@ -172,7 +184,15 @@ def verify(message: bytes, signature: bytes, pub: bytes) -> bool: # verifier must not accept a signature that covered less than it thinks. SIGNED_ATTRS = ("specversion", "type", "source", "id", "time", "subject", "datacontenttype", "correlationid", "sessionuuid", "sequence", - "phase", "final", "mode", "causationid") + "phase", "final", "mode", "causationid", + # Added once signing had real callers. `submitter`/`submitteriss` + # were held back deliberately: covering them before anything signed + # would have invalidated canonicalisation twice for no benefit. + # `groupid` is required by DESIGN_PHASE1.md §21.9.9 — without it a + # signature says nothing about which batch an event belongs to, so a + # forged `groupid` could move a response into another batch and + # corrupt its fan-in counts. + "submitter", "submitteriss", "groupid") def data_bytes(data: Any) -> bytes: @@ -256,6 +276,37 @@ def sign_event(event, seed: bytes, kid: str | None = None) -> str: return f"{protected}..{_b64u(sig)}" +def sign_into(event, seed: bytes | None, kid: str | None = None) -> bool: + """Assign `event.attrs["signature"]` when a seed is configured. Returns whether + it signed. + + `sign_event` does not mutate, so every producer would otherwise repeat the same + three lines; this is the one place that decides what "signing is enabled" means. + + **Never raises.** A signing failure here would turn a key-configuration mistake + into a total publish outage — every request rejected because one seed file has a + stray byte. Degrading to unsigned is the lesser harm, and it is safe only because + the verifying side is what enforces: an unsigned event is refused there when + enforcement is on. The loud failure belongs at startup instead, where the seed is + loaded once and a bad path stops the process. + """ + if seed is None: + return False + try: + # `sign()` does not check the seed length — only `public_key()` does — so a + # truncated key file would produce a well-formed signature that no verifier + # can ever match, reported to the operator as success. Check it here, where + # production signing enters, so the failure names the key instead. + if len(seed) != 32: + raise ValueError(f"an Ed25519 seed is 32 bytes, got {len(seed)}") + event.attrs["signature"] = sign_event(event, seed, kid) + return True + except Exception as e: # noqa: BLE001 - see above: publishing unsigned beats not publishing + print(f"[signing] could not sign {event.get('id')!r}, publishing unsigned: {e!r}", + file=sys.stderr, flush=True) + return False + + def token_kid(token: str) -> str | None: """The `kid` from a detached-JWS token's protected header, or None. @@ -302,6 +353,111 @@ def verify_signature(event, pub: bytes) -> tuple[bool, str]: return True, "ok" +# ---- authorization: the keyset as an allowlist ------------------------------- + +def verify_with_keyset(event, ks, *, expect_kid: str | None = None) -> tuple[bool, str]: + """(ok, reason) for one event against an approved-key set. + + The composition `tests/test_keyset.py` previously had to hand-wire: read the + `kid`, select the key it names, then let the signature decide. The `kid` is only + a hint until that last step succeeds, because the header is signed input — so a + token naming an unapproved key is refused for *naming* it, not because the `kid` + itself was disbelieved. + + `expect_kid` pins the signer to one identity. Group lifecycle events use it: the + keyset is otherwise flat, so any approved runner could forge a `group.completed` + and end a batch early. Passing it restricts a class of event to one key while + leaving the rest of the set alone. + + **Never raises**, so a caller inside a consumer loop needs no guard of its own to + stay alive. Every exit returns a distinct reason, because a rejection nobody can + explain gets diagnosed as "signing is broken" and switched off. + """ + token = event.get("signature") + if not token: + return False, "no ce_signature attribute" + kid = token_kid(token) + if expect_kid and kid != expect_kid: + return False, (f"expected a signature from kid {expect_kid!r}, " + f"got {kid!r}") + pub = ks.select(kid) + if pub is None: + if kid is None: + return False, (f"the token names no kid and the approved set holds " + f"{len(ks)} keys, so it is ambiguous") + return False, f"kid {kid!r} is not in the approved key set" + return verify_signature(event, pub) + + +def verify_request(event, cfg, ks=None) -> tuple[bool, str]: + """(ok, reason) for a request event, using whichever key source is configured. + + Prefers the keyset, which is an allowlist of many approved agents. Falls back to + `verify_event`'s single-key path so `ER_REQUIRE_SIGNATURE=true` with only + `ER_VERIFY_KEY_PATH` set keeps behaving as it did — that combination predates the + keyset and is still the simplest useful deployment. + """ + if ks is not None: + return verify_with_keyset(event, ks) + return verify_event(event, cfg) + + +def _is_terminal(event) -> bool: + """Whether a signature is expected on this event. + + Terminal responses are signed (`emit()` signs on `final`), and so is every group + lifecycle event. A group event carries no `final` attribute at all, so testing + `final` alone would classify it as an unsigned intermediate frame and wave it + through — which is exactly the forged `group.completed` this is meant to catch. + """ + if ce.is_group_event(event): + return True + return str(event.get("final", "")).lower() == "true" + + +def response_decision(event, ks, *, require: bool, + bridge_kid: str | None = None) -> tuple[bool, str]: + """(accept_as_is, reason) for one event off the responses topic. + + Pure: no I/O, no logging, and **it does not touch `event`** — the caller owns the + rewrite, which is what makes this testable without a Kafka consumer. + + Collapsing three questions into one answer is deliberate; it leaves the caller no + policy to get wrong: + + * `ks is None` — verification is not configured. Accept, as today. + * verified — accept. + * not verified and `require` false — **audit mode**: accept, but hand back the + reason so the caller can log it. Enforcement rewrites persisted rows and pages + a phone, so there has to be a way to watch the reject rate first. + * not verified and `require` true — reject; the caller stores it as `phase=error`. + + `bridge_kid` pins group lifecycle events to EventBridge's own key. It is only + applied when set, so a single-key deployment — where `KeySet.select(None)` returns + the sole key and nothing needs to name a kid — keeps working untouched. + + **Unsigned non-terminal frames are accepted, and that is not a loophole being + left open — it is the signing policy on the other side.** `emit()` signs terminal + events only, because it runs for every `stdout` frame and a signature costs + ~150-200 ms; verifying all-or-nothing would rewrite every streamed frame of every + genuine run to `phase=error`. So an event that carries no signature AND is not + terminal is passed through, while an unsigned **terminal** event is still refused — + that is the one the transcript presents as the answer, and refusing it is the whole + control. A forged intermediate frame therefore still renders (see `emit.py`: this + proves who *finished* a run, not what it said along the way), but it can no longer + masquerade as the result. + """ + if ks is None: + return True, "verification not enabled" + expect = bridge_kid if (bridge_kid and ce.is_group_event(event)) else None + ok, why = verify_with_keyset(event, ks, expect_kid=expect) + if ok: + return True, why + if not event.get("signature") and not _is_terminal(event): + return True, "unsigned non-terminal frame (signing covers terminal events)" + return (not require), why + + # ---- key loading ------------------------------------------------------------ def load_seed(path: str | pathlib.Path) -> bytes: @@ -343,18 +499,21 @@ def verify_event(event, cfg) -> tuple[bool, str]: # with a whitespace edge byte is not truncated out of the 32-byte path. data = pathlib.Path(key_path).read_bytes() seed_or_pub = data if len(data) == 32 else data.strip() - pub = (bytes(seed_or_pub) if len(seed_or_pub) == 32 - else public_key(load_seed(key_path))) + # 32 bytes on disk is ambiguous between a seed and a public key, and so is a + # hex/base64 file that decodes to 32 — `load_seed` decodes either and cannot + # tell them apart. So collect both readings and try each: the encoded-public-key + # case used to be missed entirely, because only `public_key(load_seed(...))` was + # tried and that derives the wrong key from a public one. + if len(seed_or_pub) == 32: + candidates = [bytes(seed_or_pub), public_key(bytes(seed_or_pub))] + else: + decoded = load_seed(key_path) + candidates = [decoded, public_key(decoded)] except Exception as e: # noqa: BLE001 return False, f"cannot load verification key from {key_path}: {e}" - # A 32-byte file is ambiguous between seed and public key; try both. - ok, why = verify_signature(event, pub) - if ok: - return True, why - try: - ok2, why2 = verify_signature(event, public_key(load_seed(key_path))) - if ok2: - return True, why2 - except Exception: # noqa: BLE001 - pass + why = "no candidate key verified the signature" + for pub in candidates: + ok, why = verify_signature(event, pub) + if ok: + return True, why return False, why diff --git a/eventing/tests/test_consume_phase1.py b/eventing/tests/test_consume_phase1.py index d744eb48..99c1a8d2 100644 --- a/eventing/tests/test_consume_phase1.py +++ b/eventing/tests/test_consume_phase1.py @@ -7,7 +7,9 @@ consuming nothing. So the central assertion in several of these tests is not "it worked" but "the thread is still alive". """ +import binascii import datetime as dt +import json import threading import time @@ -17,9 +19,24 @@ from eventrunner.config import Cfg, load from eventrunner.consume import Consumer, event_age_s from eventrunner.offsets import OffsetLedger -from shared import ce +from shared import ce, keyset +from shared import signing as S from shared.heartbeat import Heartbeat +# RFC 8032 vectors 1 and 2: an approved runner and a rogue. Fixed, never generated, +# so a failure is reproducible. +SEED = binascii.unhexlify( + "9d61b19deffd5a60ba844af492ec2cc44449c5697b326919703bac031cae7f60") +PUB = S.public_key(SEED) +SEED2 = binascii.unhexlify( + "4ccd089b28ff96da9db6c346ec114e0f5b8a319f35aba624da8cf6ed4fb8a6fb") + + +def approved_keyset(tmp_path, keys=None): + p = tmp_path / "agents.json" + p.write_text(json.dumps(keys or {"runner-01": PUB.hex()})) + return keyset.load(str(p)) + # ---- fakes ------------------------------------------------------------------ class Rec: @@ -29,7 +46,14 @@ def __init__(self, topic, partition, offset, headers, value): def request_record(corr="brave-otter-4718", *, offset=0, partition=0, - topic="kev1-requests", mode="start", when=None, prompt="hi"): + topic="kev1-requests", mode="start", when=None, prompt="hi", + seed=None, kid=None): + """One Kafka record carrying a request event. + + `seed` signs it before serialisation, the same order `kafka_out.publish_request` + uses — so a signed record here travels the wire exactly as a real one does, + rather than being signed over attributes the codec would have changed. + """ attrs = {} if when is not None: attrs["time"] = when @@ -37,6 +61,7 @@ def request_record(corr="brave-otter-4718", *, offset=0, partition=0, datacontenttype="application/json", correlationid=corr, sessionuuid=ce.session_uuid(corr), mode=mode, data={"prompt": prompt}, **attrs) + S.sign_into(evt, seed, kid) headers, value = ce.to_kafka_binary(evt) return Rec(topic, partition, offset, headers, value) @@ -486,3 +511,138 @@ def sequence(): assert {("kev1-requests", 0): 1} in fake.commits, \ "a commit must still land while draining" assert fake.paused_tps, "intake stops by pausing the assignment" + + +# ---- §11: signature verification, the reject branches ------------------------ +# +# `consume.py`'s verification branch had no tests at all: until this change nothing +# signed, so there was no way to reach the accept side and the reject side was +# indistinguishable from "always rejects". Each test below asserts the same three +# things the stale-request guard does, because a rejection that loses the offset is a +# poison pill that blocks the partition forever: +# 1. the request did not reach the router, +# 2. a named counter says why it was refused, +# 3. the offset still advanced. + +def _signed_setup(tmp_path, *, seed, kid, ks=None, **cfg_over): + tp = TopicPartition("kev1-requests", 0) + rec = request_record(offset=0, seed=seed, kid=kid) + fake = FakeConsumer([{tp: [rec]}], committed={tp: OffsetAndMetadata(0, "", -1)}) + router = RecordingRouter() + cfg = cfg_for(tmp_path, require_signature=True, **cfg_over) + return fake, router, Consumer(cfg, router, keyset=ks) + + +def test_a_request_from_an_approved_agent_runs(tmp_path): + """The accept side, which was unreachable before anything signed.""" + fake, router, c = _signed_setup(tmp_path, seed=SEED, kid="runner-01", + ks=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: router.submitted) + assert len(router.submitted) == 1, "an approved, correctly signed request must run" + assert c.rejected_unsigned == 0 + + +def test_an_unsigned_request_is_rejected_when_signatures_are_required(tmp_path): + fake, router, c = _signed_setup(tmp_path, seed=None, kid=None, + ks=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: fake.commits) + assert router.submitted == [], "an unsigned request must not reach the router" + assert c.rejected_unsigned == 1 + assert {("kev1-requests", 0): 1} in fake.commits, "but its offset must advance" + + +def test_a_request_signed_by_an_unapproved_key_is_rejected(tmp_path): + """The forgery the demo turns on: a real signature from a key nobody approved. + + The attacker holds a valid Ed25519 key and can produce a structurally perfect + signature. Being unforgeable is not the point — being *unapproved* is. + """ + fake, router, c = _signed_setup(tmp_path, seed=SEED2, kid="runner-99", + ks=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: fake.commits) + assert router.submitted == [], "an unapproved signer must not reach the router" + assert c.rejected_unsigned == 1 + assert {("kev1-requests", 0): 1} in fake.commits + + +def test_a_rogue_key_claiming_an_approved_kid_is_rejected(tmp_path): + """The nastier case: the attacker knows an approved kid but not its key. The kid + resolves, so only the signature check stands between them and a run.""" + fake, router, c = _signed_setup(tmp_path, seed=SEED2, kid="runner-01", + ks=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: fake.commits) + assert router.submitted == [] + assert c.rejected_unsigned == 1 + + +def test_an_unnamed_token_is_rejected_when_the_approved_set_is_ambiguous(tmp_path): + """Two approved keys and a token naming neither: accepting it would mean taking a + signature from ANY approved agent for an event that claimed none of them.""" + ks = approved_keyset(tmp_path, {"runner-01": PUB.hex(), + "runner-02": S.public_key(SEED2).hex()}) + fake, router, c = _signed_setup(tmp_path, seed=SEED, kid=None, ks=ks) + drive(c, fake, until=lambda: fake.commits) + assert router.submitted == [] + assert c.rejected_unsigned == 1 + + +def test_an_unnamed_token_runs_against_a_single_key_set(tmp_path): + """The friendly deployment: one approved key, so nothing has to name it.""" + fake, router, c = _signed_setup(tmp_path, seed=SEED, kid=None, + ks=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: router.submitted) + assert len(router.submitted) == 1 + + +def test_the_single_key_path_still_works_without_a_keyset(tmp_path): + """Back-compat: ER_REQUIRE_SIGNATURE with only ER_VERIFY_KEY_PATH predates the + keyset and is still the simplest useful deployment.""" + key = tmp_path / "verify.hex" + key.write_text(PUB.hex()) + fake, router, c = _signed_setup(tmp_path, seed=SEED, kid=None, ks=None, + verify_key_path=str(key)) + drive(c, fake, until=lambda: router.submitted) + assert len(router.submitted) == 1, "a single configured key must still verify" + + +def test_verification_is_skipped_entirely_when_not_required(tmp_path): + """The default path: an unsigned request runs, and no signing work happens.""" + tp = TopicPartition("kev1-requests", 0) + fake = FakeConsumer([{tp: [request_record(offset=0)]}], + committed={tp: OffsetAndMetadata(0, "", -1)}) + router = RecordingRouter() + c = Consumer(cfg_for(tmp_path), router) # require_signature defaults False + drive(c, fake, until=lambda: router.submitted) + assert len(router.submitted) == 1 + assert c.rejected_unsigned == 0 + + +def test_a_stale_request_is_dropped_before_its_signature_is_checked(tmp_path): + """Ordering, pinned: a replayed day-old event is history, not an attack, and must + not cost a ~150 ms verification each to discard.""" + tp = TopicPartition("kev1-requests", 0) + old = request_record(offset=0, when=_iso(86400)) # stale AND unsigned + fake = FakeConsumer([{tp: [old]}], committed={tp: OffsetAndMetadata(0, "", -1)}) + router = RecordingRouter() + c = Consumer(cfg_for(tmp_path, require_signature=True, max_request_age_s=3600), + router, keyset=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: fake.commits) + assert router.submitted == [] + assert c.skipped_stale == 1, "the age guard must fire..." + assert c.rejected_unsigned == 0, "...instead of the signature check" + + +def test_the_consumer_thread_survives_a_rejected_request(tmp_path): + """The failure mode this whole file exists for: a reject path that kills the + thread leaves a pod that reports healthy and consumes nothing.""" + tp = TopicPartition("kev1-requests", 0) + bad = request_record(corr="brave-otter-4718", offset=0, seed=SEED2, kid="runner-99") + good = request_record(corr="calm-badger-1234", offset=1, seed=SEED, kid="runner-01") + fake = FakeConsumer([{tp: [bad, good]}], committed={tp: OffsetAndMetadata(0, "", -1)}) + router = RecordingRouter() + c = Consumer(cfg_for(tmp_path, require_signature=True), router, + keyset=approved_keyset(tmp_path)) + drive(c, fake, until=lambda: router.submitted) + assert c.rejected_unsigned == 1, "the first record was refused" + assert len(router.submitted) == 1, "and the loop went on to process the second" + assert router.submitted[0]["correlationid"] == "calm-badger-1234" diff --git a/eventing/tests/test_kafka_in.py b/eventing/tests/test_kafka_in.py new file mode 100644 index 00000000..f6a36e1a --- /dev/null +++ b/eventing/tests/test_kafka_in.py @@ -0,0 +1,355 @@ +"""EventBridge's responses consumer: §11 verification, and the thread that must not die. + +This module had no behavioural test at all before signing was wired — it was only +grepped as source text by `test_groups.py`. That mattered, because `kafka_in.run()` +decodes and writes to SQLite *outside* any `try`: anything that raises between the poll +and the store write ends the `for`, ends the `while`, and the consumer thread is gone +while the process stays up and the pod still reports healthy. A consumer that silently +stopped consuming is the worst failure this system has, so the verification path added +here has to degrade rather than raise — and that property needs a test, not a comment. + +Most of what matters is decided by `signing.response_decision`, which is pure; those +tests live in `test_keyset.py`. What is left here is the part only the loop can show: +that a rejected event is stored rather than dropped, and that the loop survives. +""" +from __future__ import annotations + +import binascii +import json +import pathlib + +from eventbridge import kafka_in +from eventbridge.store import Store +from shared import ce, keyset +from shared import signing as S + +# RFC 8032 vectors: the bridge's own key, an approved runner, and a rogue. +SEED_EB = binascii.unhexlify( + "9d61b19deffd5a60ba844af492ec2cc44449c5697b326919703bac031cae7f60") +SEED_R1 = binascii.unhexlify( + "4ccd089b28ff96da9db6c346ec114e0f5b8a319f35aba624da8cf6ed4fb8a6fb") +SEED_ROGUE = binascii.unhexlify( + "c5aa8df43f9f837bedb7442f31dcb7b166d38535076f094b85ce3a2e0b4458f7") + + +def _keyset(tmp_path): + p = tmp_path / "agents.json" + p.write_text(json.dumps({"eb-01": S.public_key(SEED_EB).hex(), + "runner-01": S.public_key(SEED_R1).hex()})) + return keyset.load(str(p)) + + +class _Rec: + def __init__(self, headers, value): + self.headers, self.value = headers, value + + +def _response(corr="brave-otter-4718", *, seed=None, kid=None, seq=1, + phase="result", text="the real answer", **over): + attrs = dict(type=ce.TYPE_RESPONSE, source="rossoctl://eventrunner/test", + datacontenttype="application/json", correlationid=corr, + sessionuuid=ce.session_uuid(corr), sequence=seq, phase=phase, + final="true") + attrs.update(over) + evt = ce.new_event(data={"text": text}, **attrs) + S.sign_into(evt, seed, kid) + return _Rec(*ce.to_kafka_binary(evt)) + + +def _group_event(gid="g-123", *, seed=None, kid=None, + type_=ce.TYPE_GROUP_COMPLETED): + evt = ce.new_event(type=type_, source="rossoctl://eventbridge/test", + datacontenttype="application/json", groupid=gid, + data={"reason": "all done"}) + S.sign_into(evt, seed, kid) + return _Rec(*ce.to_kafka_binary(evt)) + + +class _FakeKafka: + """Iterable KafkaConsumer stand-in: `kafka_in.run()` does `for rec in c`. + + Stopping is driven from inside iteration and has to happen at exactly the right + moment. `run()` checks `stopping` before the `while`, so stopping beforehand skips + the loop entirely; it also checks at the top of each record, so stopping *before* + yielding the last one makes `run()` break without processing it. The flag is + therefore set when iteration is resumed after the final record — by then its body + has already run, and the outer `while` exits instead of spinning on an exhausted + iterator forever. + """ + + def __init__(self, records, consumer): + self._records = list(records) + self._consumer = consumer + self.closed = False + + def __iter__(self): + while self._records: + yield self._records.pop(0) + # Every record has been processed by now: end the outer while loop. + self._consumer.stop() + + def close(self): + self.closed = True + + +def _drain(tmp_path, records, monkeypatch, *, store=None, **kw): + """Run the consumer over `records` to exhaustion, synchronously. + + `run()` is called directly rather than through `start()`: the record list is + finite, so there is nothing to wait for and a real thread would only add a race. + """ + store = store or Store(pathlib.Path(tmp_path) / "eb") + seen: list[dict] = [] + c = kafka_in.Consumer("broker:9092", "responses", store, + on_event=seen.append, **kw) + fake = _FakeKafka(records, c) + monkeypatch.setattr(kafka_in, "KafkaConsumer", lambda *a, **k: fake) + c.run() + return c, store, seen, fake + + +# ---- the default path: nothing configured, nothing changes ------------------- + +def test_an_unsigned_response_is_stored_normally_when_verification_is_off( + tmp_path, monkeypatch): + """Today's behaviour, which must be exactly preserved as the default.""" + c, store, seen, _ = _drain(tmp_path, [_response()], monkeypatch) + rows = store.events_for("brave-otter-4718") + assert len(rows) == 1 + assert rows[0]["phase"] == "result", "an unverified event must not be rewritten" + assert rows[0]["data"]["text"] == "the real answer" + assert c.rejected == 0 + + +# ---- verification on --------------------------------------------------------- + +def test_a_response_from_an_approved_agent_is_stored_unchanged(tmp_path, monkeypatch): + c, store, _, _ = _drain( + tmp_path, [_response(seed=SEED_R1, kid="runner-01")], monkeypatch, + keyset=_keyset(tmp_path), require_signature=True, bridge_kid="eb-01") + rows = store.events_for("brave-otter-4718") + assert rows[0]["phase"] == "result" and c.rejected == 0 + assert rows[0]["data"]["text"] == "the real answer" + + +def test_a_forged_response_is_stored_as_an_error_not_dropped(tmp_path, monkeypatch): + """The demo, and the reason rejection is not a silent drop. + + A dropped event is indistinguishable from an agent that never answered. Stored as + `phase=error` it becomes a red card in the transcript, a priority-5 notification, + and a retained `raw_json` row — the forgery attempt is evidence rather than absence. + """ + forged = _response(seed=None, text="Transfer approved. Ship the goods.") + c, store, seen, _ = _drain(tmp_path, [forged], monkeypatch, + keyset=_keyset(tmp_path), require_signature=True, + bridge_kid="eb-01") + rows = store.events_for("brave-otter-4718") + assert len(rows) == 1, "the event must still be persisted" + assert rows[0]["phase"] == "error", "and rewritten so it renders as a rejection" + assert rows[0]["data"]["signature_rejected"] is True + # ntfy reads data["text"] for the error body; any other key shows up on the + # phone as "(error, see raw)". + assert "unverified response rejected" in rows[0]["data"]["text"] + assert "no ce_signature" in rows[0]["data"]["reason"] + assert c.rejected == 1 + # The forged text must not be presented AS the answer... + assert "Ship the goods" not in rows[0]["data"]["text"] + assert seen and seen[0]["phase"] == "error", "downstream sees the rejection too" + + +def test_a_rejected_event_is_retained_for_audit(tmp_path, monkeypatch): + """...but it must still be retained, which is the stated reason for storing a + rejection rather than dropping it. + + `insert_response` derives BOTH `data_json` and `raw_json` from one dict, so + mutating the envelope in place destroyed the forensic record while the module + docstring still promised it — leaving a signature whose covered attributes no + longer existed and nothing for an operator to review. Raised by @aslom on #879. + """ + forged = _response(seed=None, seq=4, phase="result", + text="Transfer approved. Ship the goods.") + c, store, _, _ = _drain(tmp_path, [forged], monkeypatch, + keyset=_keyset(tmp_path), require_signature=True, + bridge_kid="eb-01") + raw = store.raw_events_for("brave-otter-4718")[0] + kept = raw["data"]["rejected"] + assert kept["data"]["text"] == "Transfer approved. Ship the goods.", \ + "the payload that was refused must be reviewable" + assert kept["phase"] == "result", "including the phase it claimed to be" + assert kept["attrs"]["sequence"] == "4" + assert kept["source"] == "rossoctl://eventrunner/test" + assert raw["phase"] == "error", "while the stored phase still drives the red card" + assert c.rejected == 1 + + +def test_the_retained_signature_can_still_be_re_checked(tmp_path, monkeypatch): + """The sharper half of the same point: a retained signature is only worth keeping + if the attributes it covered are kept with it. Here a validly-signed event is + rejected for naming an unapproved kid, and the record is complete enough to verify + offline — which is what makes an incident reviewable rather than just logged.""" + rec = _response(seq=5, seed=SEED_ROGUE, kid="runner-99") + c, store, _, _ = _drain(tmp_path, [rec], monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01") + kept = store.raw_events_for("brave-otter-4718")[0]["data"]["rejected"] + assert c.rejected == 1 + # Rebuild the event exactly as it arrived and verify it against the rogue key. + replayed = ce.CloudEvent(attrs=dict(kept["attrs"]), data=kept["data"]) + ok, why = S.verify_signature(replayed, S.public_key(SEED_ROGUE)) + assert ok, f"the retained record must still verify against the key that signed it: {why}" + assert S.token_kid(kept["signature"]) == "runner-99", \ + "so an operator can see which key id the forgery claimed" + + +def test_a_response_signed_by_an_unapproved_key_is_rejected(tmp_path, monkeypatch): + """A structurally perfect signature from a key nobody approved.""" + c, store, _, _ = _drain( + tmp_path, [_response(seed=SEED_ROGUE, kid="runner-99")], monkeypatch, + keyset=_keyset(tmp_path), require_signature=True, bridge_kid="eb-01") + assert store.events_for("brave-otter-4718")[0]["phase"] == "error" + assert c.rejected == 1 + + +def test_a_tampered_payload_is_rejected(tmp_path, monkeypatch): + """Signed by an approved agent, then the body was changed in flight.""" + rec = _response(seed=SEED_R1, kid="runner-01") + rec.value = json.dumps({"text": "Transfer approved. Ship the goods."}).encode() + c, store, _, _ = _drain(tmp_path, [rec], monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01") + assert store.events_for("brave-otter-4718")[0]["phase"] == "error" + assert c.rejected == 1 + + +def test_audit_mode_reports_without_rewriting(tmp_path, monkeypatch): + """A keyset alone verifies and logs; only enforcement rewrites what users see. + + That ordering exists so an operator can watch the reject rate on real traffic + before it starts painting bubbles red and paging a phone. + """ + c, store, _, _ = _drain(tmp_path, [_response(seed=None)], monkeypatch, + keyset=_keyset(tmp_path), require_signature=False, + bridge_kid="eb-01") + rows = store.events_for("brave-otter-4718") + assert rows[0]["phase"] == "result", "audit mode must not rewrite" + assert rows[0]["data"]["text"] == "the real answer" + assert c.rejected == 0 + + +def test_a_real_runs_streamed_frames_survive_enforcement(tmp_path, monkeypatch): + """What a genuine run actually looks like on the topic, under enforcement. + + `emit()` signs terminal events only, so a run publishes N unsigned `stdout` + frames and one signed terminal. Verifying all-or-nothing turned every frame of + every real run red — caught by running this against a live broker, not by a unit + test, which is why this one exists. + """ + recs = [ + _response(seq=1, phase="stdout", text="thinking...", final="false"), + _response(seq=2, phase="stdout", text="still working", final="false"), + _response(seq=3, phase="result", text="2+2 is 4", + seed=SEED_R1, kid="runner-01"), + ] + c, store, _, _ = _drain(tmp_path, recs, monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01") + rows = store.events_for("brave-otter-4718") + assert [r["phase"] for r in rows] == ["stdout", "stdout", "result"], \ + "a genuine run must render unchanged" + assert c.rejected == 0 + assert rows[-1]["data"]["text"] == "2+2 is 4" + + +def test_a_forged_terminal_is_rejected_among_genuine_frames(tmp_path, monkeypatch): + """The demo in its realistic shape: the forgery arrives alongside real streaming + output, and only it is rewritten.""" + recs = [ + _response(seq=1, phase="stdout", text="thinking...", final="false"), + _response(seq=2, phase="result", text="2+2 is 4", + seed=SEED_R1, kid="runner-01"), + _response(seq=3, phase="result", text="Transfer approved. Ship the goods."), + ] + c, store, _, _ = _drain(tmp_path, recs, monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01") + rows = store.events_for("brave-otter-4718") + assert [r["phase"] for r in rows] == ["stdout", "result", "error"] + assert c.rejected == 1 + assert "Ship the goods" not in rows[2]["data"]["text"], \ + "the forgery must not be presented as the answer" + assert rows[1]["data"]["text"] == "2+2 is 4", "the genuine answer is untouched" + + +# ---- group lifecycle events -------------------------------------------------- + +def test_a_group_event_signed_by_the_bridge_is_accepted(tmp_path, monkeypatch): + """Group events carry `groupid` and no `correlationid`, so they take the routing + branch that never reaches insert_response — they must still verify.""" + groups: list[dict] = [] + c, _, _, _ = _drain(tmp_path, [_group_event(seed=SEED_EB, kid="eb-01")], + monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01", + on_group_event=groups.append) + assert c.rejected == 0 + assert groups and groups[0]["groupid"] == "g-123" + assert groups[0].get("phase") != "error" + + +def test_an_approved_runner_cannot_forge_a_group_event(tmp_path, monkeypatch): + """The hole a flat keyset would leave: a forged `group.completed` ends a batch + early and fires a "finished" notification for work that never ran. runner-01 is + genuinely approved — it is just not the bridge.""" + groups: list[dict] = [] + c, _, _, _ = _drain(tmp_path, [_group_event(seed=SEED_R1, kid="runner-01")], + monkeypatch, keyset=_keyset(tmp_path), + require_signature=True, bridge_kid="eb-01", + on_group_event=groups.append) + assert c.rejected == 1 + assert groups and groups[0]["phase"] == "error", \ + "the group event is marked rejected rather than settling the group" + + +# ---- the property this file exists for -------------------------------------- + +def test_the_consumer_survives_a_verifier_that_raises(tmp_path, monkeypatch): + """The most important test here. + + `from_kafka_binary` and `insert_response` are not inside a try, so an exception + escaping the verification path would end the consume loop for the life of the pod — + silently, with the process still healthy. Both records must be processed. + """ + def boom(*a, **k): + raise RuntimeError("keyset backend exploded") + monkeypatch.setattr(S, "response_decision", boom) + + recs = [_response(corr="brave-otter-4718", seed=SEED_R1, kid="runner-01"), + _response(corr="calm-badger-1234", seed=SEED_R1, kid="runner-01")] + c, store, seen, fake = _drain(tmp_path, recs, monkeypatch, + keyset=_keyset(tmp_path), require_signature=True, + bridge_kid="eb-01") + assert len(seen) == 2, "the loop must process both records, not die on the first" + assert c.rejected == 2, "and fail closed while enforcement is on" + for corr in ("brave-otter-4718", "calm-badger-1234"): + assert store.events_for(corr)[0]["phase"] == "error" + assert fake.closed, "the consumer is still closed cleanly on the way out" + + +def test_a_verifier_that_raises_fails_open_when_not_enforcing(tmp_path, monkeypatch): + """Audit mode must not start rejecting because the verifier broke: nothing is + being enforced, so a broken check has no opinion to act on.""" + def boom(*a, **k): + raise RuntimeError("keyset backend exploded") + monkeypatch.setattr(S, "response_decision", boom) + + c, store, seen, _ = _drain(tmp_path, [_response(seed=SEED_R1, kid="runner-01")], + monkeypatch, keyset=_keyset(tmp_path), + require_signature=False, bridge_kid="eb-01") + assert len(seen) == 1 + assert c.rejected == 0 + assert store.events_for("brave-otter-4718")[0]["phase"] == "result" + + +def test_an_undecodable_record_still_ends_the_loop_as_before(tmp_path, monkeypatch): + """Not a regression this change introduces, but worth pinning what it does NOT fix: + the decode at the top of the loop is still outside any try. Verification was made + safe; the pre-existing decode hazard is unchanged and out of scope here.""" + src = pathlib.Path(kafka_in.__file__).read_text() + assert "evt = ce.from_kafka_binary(rec.headers or [], rec.value)" in src + assert "enable_auto_commit=True" in src, ( + "the live consumer must keep committing — test_groups.py pins this too") diff --git a/eventing/tests/test_keyset.py b/eventing/tests/test_keyset.py index cc3fa6bf..95a4ff1f 100644 --- a/eventing/tests/test_keyset.py +++ b/eventing/tests/test_keyset.py @@ -14,8 +14,8 @@ import pytest -from eventrunner import signing as S from shared import ce, keyset +from shared import signing as S # RFC 8032 test vector 1 — the same seed the signing tests use. SEED = binascii.unhexlify( @@ -220,6 +220,53 @@ def test_payload_tampering_on_the_wire_is_detected(): assert not ok +def _tampered_on_the_wire(attr, value, **signed_attrs): + """Sign an event, then rewrite one `ce_` header in flight. Returns (ok, why).""" + e = _event(**signed_attrs) + e.attrs["signature"] = S.sign_event(e, SEED, kid="runner-01") + headers, body = ce.to_kafka_binary(e) + headers = [(k, value if k == f"ce_{attr}" else v) for k, v in headers] + return S.verify_signature(ce.from_kafka_binary(headers, body), PUB) + + +def test_the_submitter_is_covered_by_the_signature(): + """What makes the claim in eventbridge/auth.py true rather than aspirational. + + `submitter` records who EventBridge believes submitted a request. Until it was in + SIGNED_ATTRS, anything with topic write access could rewrite that name — so the + audit trail named whoever the last writer chose. + """ + ok, why = _tampered_on_the_wire("submitter", b"attacker", + submitter="mrsabath", submitteriss="github") + assert not ok and "does not verify" in why + + +def test_the_submitter_issuer_is_covered_by_the_signature(): + """Without this, a forged `submitteriss=github` upgrades a name an operator typed + into an env var to a login GitHub appears to have verified.""" + ok, _ = _tampered_on_the_wire("submitteriss", b"github", + submitter="mrsabath", submitteriss="static") + assert not ok + + +def test_the_groupid_is_covered_by_the_signature(): + """DESIGN_PHASE1.md §21.9.9 requires it: a forged `groupid` moves a response into + another batch and corrupts that batch's fan-in counts.""" + ok, _ = _tampered_on_the_wire("groupid", b"g-victim", groupid="g-mine") + assert not ok + + +def test_adding_a_groupid_after_signing_is_detected(): + """The absence of an attribute is signed too, not just its value — otherwise an + ungrouped response could be adopted into a batch it was never part of.""" + e = _event() # no groupid + e.attrs["signature"] = S.sign_event(e, SEED, kid="runner-01") + headers, body = ce.to_kafka_binary(e) + headers.append(("ce_groupid", b"g-victim")) + ok, _ = S.verify_signature(ce.from_kafka_binary(headers, body), PUB) + assert not ok + + def test_text_plain_payload_roundtrips(): e = _event(datacontenttype="text/plain") e.data = "hello" @@ -234,9 +281,10 @@ def test_text_plain_payload_roundtrips(): def test_approved_agent_verifies_and_unapproved_does_not(tmp_path): """The intended mechanism, composed by hand: the keyset as authorization list. - NB this wires `token_kid` -> `select` -> `verify_signature` itself, because no - production path does yet — see the note in shared/keyset.py. It proves the - primitives compose, not that the runner enforces anything. + Kept hand-wired even though `signing.verify_with_keyset` now does exactly this in + production, because a test that reaches through the primitives one at a time says + where a failure is — `select` returned nothing, or the signature did not verify — + which a single boolean from the composed function cannot. """ ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) @@ -263,3 +311,188 @@ def test_a_rogue_key_claiming_an_approved_kid_is_refused(tmp_path): pub = ks.select(S.token_kid(back.get("signature"))) assert pub == PUB # kid resolves... assert not S.verify_signature(back, pub)[0] # ...but the signature does not + + +# ---- sign_into: the one place that decides "signing is enabled" -------------- + +def test_sign_into_is_a_no_op_without_a_seed(): + """The default path, asserted rather than assumed: no seed, no attribute.""" + e = _event() + assert S.sign_into(e, None) is False + assert "signature" not in e.attrs + + +def test_sign_into_signs_and_the_result_verifies(): + e = _event() + assert S.sign_into(e, SEED, "runner-01") is True + assert S.verify_signature(e, PUB) == (True, "ok") + assert S.token_kid(e.attrs["signature"]) == "runner-01" + + +def test_sign_into_degrades_rather_than_raising_on_a_bad_seed(): + """A key-config mistake must not take out publishing entirely. + + Unsigned is safe because the verifying side refuses unsigned events when + enforcement is on; a raise here would reject 100% of traffic at the producer. + """ + e = _event() + assert S.sign_into(e, b"\x01" * 31) is False # 31 bytes is not an Ed25519 seed + assert "signature" not in e.attrs + + +# ---- verify_with_keyset: the composed authorization check -------------------- + +def test_verify_with_keyset_accepts_an_approved_signer(tmp_path): + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + e = _event() + S.sign_into(e, SEED, "runner-01") + back = ce.from_kafka_binary(*ce.to_kafka_binary(e)) + ok, why = S.verify_with_keyset(back, ks) + assert ok, why + + +def test_verify_with_keyset_reports_a_missing_signature(tmp_path): + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + ok, why = S.verify_with_keyset(_event(), ks) + assert not ok and "no ce_signature" in why + + +def test_verify_with_keyset_reports_an_unapproved_kid(tmp_path): + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + e = _event() + S.sign_into(e, SEED2, "runner-99") + ok, why = S.verify_with_keyset(e, ks) + assert not ok and "'runner-99' is not in the approved key set" in why + + +def test_verify_with_keyset_reports_an_ambiguous_unnamed_token(tmp_path): + """Two approved keys and a token naming neither: guessing would accept a + signature from *any* approved agent for an event that claimed none of them.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex(), + "runner-02": PUB2.hex()})) + e = _event() + S.sign_into(e, SEED) # no kid + ok, why = S.verify_with_keyset(e, ks) + assert not ok and "ambiguous" in why + + +def test_verify_with_keyset_accepts_an_unnamed_token_against_a_single_key(tmp_path): + """The friendly path: one approved key means nothing has to name it.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + e = _event() + S.sign_into(e, SEED) # no kid + ok, why = S.verify_with_keyset(e, ks) + assert ok, why + + +def test_expect_kid_pins_the_signer_to_one_identity(tmp_path): + """What stops an approved runner forging an EventBridge group event.""" + ks = keyset.load(_write(tmp_path, {"eventbridge": PUB.hex(), + "runner-01": PUB2.hex()})) + e = _event(type=ce.TYPE_GROUP_STARTED) + S.sign_into(e, SEED2, "runner-01") # a genuinely approved key... + ok, why = S.verify_with_keyset(e, ks, expect_kid="eventbridge") + assert not ok and "expected a signature from kid 'eventbridge'" in why + + bridge = _event(type=ce.TYPE_GROUP_STARTED) + S.sign_into(bridge, SEED, "eventbridge") + assert S.verify_with_keyset(bridge, ks, expect_kid="eventbridge")[0] + + +# ---- response_decision: enabled? verified? enforced? ------------------------- + +def test_response_decision_accepts_everything_when_no_keyset_is_configured(): + """Today's behaviour, which must survive untouched as the default.""" + ok, why = S.response_decision(_event(), None, require=True) + assert ok and "not enabled" in why + + +def test_response_decision_accepts_a_verified_response(tmp_path): + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + e = _event() + S.sign_into(e, SEED, "runner-01") + assert S.response_decision(e, ks, require=True) == (True, "ok") + + +def test_response_decision_rejects_an_unsigned_response_when_enforcing(tmp_path): + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + ok, why = S.response_decision(_event(), ks, require=True) + assert not ok and "no ce_signature" in why + + +def test_response_decision_reports_but_accepts_in_audit_mode(tmp_path): + """The two-flag rollout in one assertion: a keyset alone verifies and explains, + without yet rewriting anything a user sees.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + ok, why = S.response_decision(_event(), ks, require=False) + assert ok, "audit mode must not reject" + assert "no ce_signature" in why, "but it must still say what was wrong" + + +def test_response_decision_requires_the_bridge_kid_on_group_events(tmp_path): + ks = keyset.load(_write(tmp_path, {"eventbridge": PUB.hex(), + "runner-01": PUB2.hex()})) + forged = _event(type=ce.TYPE_GROUP_COMPLETED) + S.sign_into(forged, SEED2, "runner-01") + ok, why = S.response_decision(forged, ks, require=True, bridge_kid="eventbridge") + assert not ok and "expected a signature from kid 'eventbridge'" in why + + # A plain response from that same runner is still fine — the pin is per-class. + resp = _event() + S.sign_into(resp, SEED2, "runner-01") + assert S.response_decision(resp, ks, require=True, bridge_kid="eventbridge")[0] + + +def test_an_unsigned_non_terminal_frame_is_accepted(tmp_path): + """Matching the signing policy, not a loophole. + + `emit()` signs terminal events only, because it runs per stdout frame at + ~150-200 ms a signature. Verifying all-or-nothing would rewrite every streamed + frame of every genuine run to phase=error — which is what happened the first time + this was run against a real broker. + """ + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + frame = _event(sequence=1, phase="stdout", final="false") + ok, why = S.response_decision(frame, ks, require=True) + assert ok, why + assert "non-terminal" in why + + +def test_an_unsigned_terminal_response_is_still_rejected(tmp_path): + """The line the exemption must not cross: the terminal event is the one the + transcript presents as the answer, so refusing it unsigned is the whole control.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + ok, why = S.response_decision(_event(phase="result", final="true"), ks, require=True) + assert not ok and "no ce_signature" in why + + +def test_a_non_terminal_frame_that_claims_a_signature_is_still_verified(tmp_path): + """The exemption is for events carrying NO signature. One that presents a bad + signature is lying about something, and gets checked.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + frame = _event(sequence=1, phase="stdout", final="false") + S.sign_into(frame, SEED2, "runner-01") # approved kid, wrong key + ok, why = S.response_decision(frame, ks, require=True) + assert not ok and "does not verify" in why + + +def test_an_unsigned_group_event_is_rejected_despite_having_no_final_attribute(tmp_path): + """Group events carry no `final` at all, so testing that alone would classify one + as an unsigned intermediate frame and wave it through — which is precisely the + forged `group.completed` the bridge kid exists to catch.""" + ks = keyset.load(_write(tmp_path, {"eventbridge": PUB.hex()})) + grp = _event(type=ce.TYPE_GROUP_COMPLETED, groupid="g-1") + grp.attrs.pop("final", None) + assert "final" not in grp.attrs + ok, why = S.response_decision(grp, ks, require=True, bridge_kid="eventbridge") + assert not ok and "no ce_signature" in why + + +def test_response_decision_does_not_mutate_the_event(tmp_path): + """What "pure" buys: the caller owns the rewrite, so this is safe to call on the + hot path of a consumer loop without copying first.""" + ks = keyset.load(_write(tmp_path, {"runner-01": PUB.hex()})) + e = _event() + before_attrs, before_data = dict(e.attrs), dict(e.data) + S.response_decision(e, ks, require=True) + assert e.attrs == before_attrs and e.data == before_data diff --git a/eventing/tests/test_signing.py b/eventing/tests/test_signing.py index 42828d79..24e41bb6 100644 --- a/eventing/tests/test_signing.py +++ b/eventing/tests/test_signing.py @@ -7,13 +7,14 @@ The canonicalization tests matter just as much: §11 says outright that signer and verifier agreeing byte-for-byte "is the part that will bite". """ +import base64 import binascii import json import pytest -from eventrunner import signing as S from shared import ce +from shared import signing as S # ---- RFC 8032 §7.1 test vectors -------------------------------------------- @@ -117,8 +118,19 @@ def test_canonical_binds_the_payload_by_digest(): def test_canonical_changes_when_any_signed_attribute_changes(): + """Derived from SIGNED_ATTRS rather than listed by hand. + + A hardcoded list silently stops covering whatever is added to the tuple next, + which is exactly what happened when `submitter`/`submitteriss`/`groupid` were + added. `specversion` is excluded because a different value is not a different + event but a different envelope format, and `datacontenttype` because changing it + changes how `data` is encoded rather than only the attribute. + """ + skip = {"specversion", "datacontenttype"} + covered = [a for a in S.SIGNED_ATTRS if a not in skip] + assert len(covered) >= 15, "SIGNED_ATTRS shrank unexpectedly" base = S.canonical(_event().attrs, None) - for attr in ("id", "correlationid", "sessionuuid", "mode", "phase", "causationid"): + for attr in covered: other = S.canonical(_event(**{attr: "different"}).attrs, None) assert other != base, f"{attr} is in SIGNED_ATTRS but did not affect the digest" @@ -279,3 +291,41 @@ def test_verify_event_uses_the_configured_key(tmp_path): e.attrs["signature"] = S.sign_event(e, seed) ok, why = S.verify_event(e, cfg) assert ok, why + + +@pytest.mark.parametrize("form", ["raw-pub", "hex-pub", "b64-pub", + "raw-seed", "hex-seed", "b64-seed"]) +def test_verify_event_accepts_a_key_in_any_documented_form(tmp_path, form): + """`ER_VERIFY_KEY_PATH` may hold a seed or a public key, in raw, hex or base64. + + A 32-byte file cannot be told apart from an encoded one by content, so all six + combinations have to work. A hex-encoded PUBLIC key used to fail: only + `public_key(load_seed(path))` was tried for encoded files, which derives a + different key when the file already holds a public one. + """ + from eventrunner.config import Cfg + seed = binascii.unhexlify(RFC_VECTORS[1][0]) + pub = S.public_key(seed) + material = seed if form.endswith("seed") else pub + encode = {"raw": lambda b: b, + "hex": lambda b: b.hex().encode(), + "b64": base64.b64encode}[form.split("-")[0]] + key = tmp_path / "key.bin" + key.write_bytes(encode(material)) + e = _event() + e.attrs["signature"] = S.sign_event(e, seed) + ok, why = S.verify_event(e, Cfg(verify_key_path=str(key))) + assert ok, f"{form}: {why}" + + +def test_verify_event_still_rejects_an_unrelated_key(tmp_path): + """Accepting more *encodings* must not mean accepting more *keys*.""" + from eventrunner.config import Cfg + seed, other = (binascii.unhexlify(RFC_VECTORS[1][0]), + binascii.unhexlify(RFC_VECTORS[2][0])) + key = tmp_path / "other.hex" + key.write_text(S.public_key(other).hex()) + e = _event() + e.attrs["signature"] = S.sign_event(e, seed) + ok, why = S.verify_event(e, Cfg(verify_key_path=str(key))) + assert not ok and "does not verify" in why