From fc5d0cf12762d145b9086d420bf1bfebe1421d1e Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 17:49:07 -0400 Subject: [PATCH 1/8] refactor(eventing): move signing.py to shared/ so EventBridge can use it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wiring #2607 makes EventBridge sign the requests it publishes and verify the responses it consumes, so it needs the signing primitives. They lived in `eventrunner/`, which EventBridge has never imported from. This is not tidiness. `Dockerfile-eventbridge` COPYs only `shared/` and `eventbridge/` — never `eventrunner/` — so an `eventbridge -> eventrunner` import would pass every local test and then ImportError inside the container. `shared/` already holds `ce.py` and `keyset.py`, the two modules both services use, so signing belongs beside them. Pure move plus four reference updates. Verified behaviour-neutral by diffing `pytest --collect-only` before and after: the two test sets are identical. 496 passed, 5 skipped; ruff clean under the pinned 0.11.4. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/IMPLEMENTATION_REPORT1.md | 2 +- eventing/eventrunner/consume.py | 2 +- eventing/{eventrunner => shared}/signing.py | 0 eventing/tests/test_keyset.py | 2 +- eventing/tests/test_signing.py | 2 +- 5 files changed, 4 insertions(+), 4 deletions(-) rename eventing/{eventrunner => shared}/signing.py (100%) diff --git a/eventing/agentdocs/IMPLEMENTATION_REPORT1.md b/eventing/agentdocs/IMPLEMENTATION_REPORT1.md index edd03a1a..6f8d940d 100644 --- a/eventing/agentdocs/IMPLEMENTATION_REPORT1.md +++ b/eventing/agentdocs/IMPLEMENTATION_REPORT1.md @@ -824,7 +824,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/eventrunner/consume.py b/eventing/eventrunner/consume.py index 73683db7..18c252d8 100644 --- a/eventing/eventrunner/consume.py +++ b/eventing/eventrunner/consume.py @@ -271,7 +271,7 @@ def _handle(self, rec) -> None: return if self._cfg.require_signature: - from eventrunner import signing + from shared import signing ok, why = signing.verify_event(evt, self._cfg) if not ok: _elog(f"rejecting unsigned/badly-signed request corr={corr}: {why} " diff --git a/eventing/eventrunner/signing.py b/eventing/shared/signing.py similarity index 100% rename from eventing/eventrunner/signing.py rename to eventing/shared/signing.py diff --git a/eventing/tests/test_keyset.py b/eventing/tests/test_keyset.py index cc3fa6bf..5885b875 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( diff --git a/eventing/tests/test_signing.py b/eventing/tests/test_signing.py index 42828d79..1a47302f 100644 --- a/eventing/tests/test_signing.py +++ b/eventing/tests/test_signing.py @@ -12,8 +12,8 @@ import pytest -from eventrunner import signing as S from shared import ce +from shared import signing as S # ---- RFC 8032 §7.1 test vectors -------------------------------------------- From d7628ef739955ce3bdc0741cf7abe3cdb1b83e03 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 17:54:09 -0400 Subject: [PATCH 2/8] feat(eventing): add the signing/verification helpers #2607 wires up MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Four pure functions in `shared/signing.py`, plus their tests. Nothing is wired to a producer or consumer yet — that follows in separate commits, so a review of the policy is separable from a review of the plumbing. * `sign_into(event, seed, kid)` — the one place that decides what "signing is enabled" means. `sign_event` does not mutate, so every producer would otherwise repeat the same three lines. Never raises: a bad seed degrades to publishing unsigned rather than rejecting 100% of traffic, which is safe only because the verifying side refuses unsigned events when enforcement is on. * `verify_with_keyset(event, ks, expect_kid=)` — the composition `tests/test_keyset.py` previously had to hand-wire: read the `kid`, select the key it names, let the signature decide. `expect_kid` pins a class of event to one identity; group lifecycle events need it because the keyset is otherwise flat and any approved runner could forge a `group.completed`. * `verify_request(event, cfg, ks)` — keyset when configured, else the existing single-key `verify_event`, so `ER_REQUIRE_SIGNATURE=true` with only `ER_VERIFY_KEY_PATH` set keeps behaving as it does today. * `response_decision(event, ks, require=, bridge_kid=)` — collapses "is verification configured?", "did it verify?" and "do we enforce?" into one answer, leaving the caller no policy to get wrong. Pure and non-mutating, which is what makes the responses path testable without a Kafka consumer at all. Every rejection returns a distinct reason. A rejection nobody can explain gets diagnosed as "signing is broken" and switched off. Found while testing: `sign()` never validated seed length — only `public_key()` did — so a truncated key file produced a well-formed signature that no verifier could ever match, reported to the operator as success. `sign_into` now checks, where production signing enters. 15 new tests. 511 passed, 5 skipped; ruff clean under the pinned 0.11.4. (The baseline is 496, not the 492 an older handoff recorded; verified by collecting the pristine tree.) Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/shared/keyset.py | 16 ++-- eventing/shared/signing.py | 113 ++++++++++++++++++++++++++ eventing/tests/test_keyset.py | 147 +++++++++++++++++++++++++++++++++- 3 files changed, 262 insertions(+), 14 deletions(-) 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/shared/signing.py b/eventing/shared/signing.py index 2a2893cb..e39bd75b 100644 --- a/eventing/shared/signing.py +++ b/eventing/shared/signing.py @@ -31,8 +31,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 @@ -256,6 +259,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 +336,85 @@ 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 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. + """ + 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 + return (not require), why + + # ---- key loading ------------------------------------------------------------ def load_seed(path: str | pathlib.Path) -> bytes: diff --git a/eventing/tests/test_keyset.py b/eventing/tests/test_keyset.py index 5885b875..812243e7 100644 --- a/eventing/tests/test_keyset.py +++ b/eventing/tests/test_keyset.py @@ -234,9 +234,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 +264,143 @@ 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_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 From f8595299c0bab5f65f93d54651235bf5931d1696 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 20:49:46 -0400 Subject: [PATCH 3/8] feat(eventing): sign requests and responses, verify both (#2607) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `sign_event()` finally has production callers. Until now `ER_REQUIRE_SIGNATURE=true` rejected 100% of traffic — a kill switch, not a feature — because nothing on either side ever signed. Signing (opt-in, off by default): * EventBridge signs requests in `publish_request`, after `new_event` fills `id` and `time` (both signed) and before serialisation. * EventBridge signs group lifecycle events under the same kid. Skipping them would have left the hole worth closing: a forged `group.completed` ends a batch early and fires a "finished" notification for work that never ran. * EventRunner signs **terminal responses only**. `emit()` runs per stdout frame and a signature costs ~150-200 ms, so signing every frame would add minutes to a chatty run. The seed lives on the Emitter, so none of the seven `emit()` call sites changed. Verification: * EventRunner resolves the key by the token's `kid` when a keyset is configured, making it an allowlist of several agents; falls back to the existing single-key path so `ER_REQUIRE_SIGNATURE` + `ER_VERIFY_KEY_PATH` is unchanged. Rejections increment `rejected_unsigned` and keep the existing commit-and-do-not-retry semantics. * EventBridge verifies responses and stores a failure as `phase="error"` rather than dropping it — dropping is indistinguishable from an agent that never answered, while `phase=error` reuses a red transcript card, ntfy priority 5 and `raw_json` audit retention. Group events are pinned to EventBridge's own kid. Two hazards in `kafka_in.py` handled explicitly. `from_kafka_binary` and `insert_response` are not inside a try, so a raise in the verification path would end the consume loop permanently while the pod still reported healthy — hence the guard, which fails closed only where enforcement is on. And the error payload uses `data["text"]` because that is what ntfy renders; anything else shows up on the phone as "(error, see raw)". Rollout is two flags so enforcement is never the first step: a keyset alone verifies and logs while storing events unchanged (audit mode); `EB_REQUIRE_RESPONSE_SIGNATURE=true` is what rewrites them. EventBridge refuses to start with a keyset but no `EB_SIGNING_KID`, since there would be nothing to attribute a group event to. Key material loads once at startup and is deliberately uncaught: a service that believes it is signing but is not fails silently, whereas one that will not start says so in `kubectl logs`. Approved kids are printed, because a rejection caused by a stale ConfigMap is otherwise invisible. Verified with real key files outside the test suite: an EB-signed request verifies at the runner, a runner-signed terminal response verifies at the bridge, an unsigned forgery is refused and then accepted once verification is switched off, and an approved runner cannot forge a group event. 511 passed, 5 skipped — every pre-existing test unchanged; ruff clean. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/eventbridge/__main__.py | 38 ++++++++++++++++++- eventing/eventbridge/config.py | 28 ++++++++++++++ eventing/eventbridge/kafka_in.py | 62 ++++++++++++++++++++++++++++++- eventing/eventbridge/kafka_out.py | 34 ++++++++++++++++- eventing/eventrunner/__main__.py | 30 ++++++++++++++- eventing/eventrunner/config.py | 16 ++++++-- eventing/eventrunner/consume.py | 20 +++++++++- eventing/eventrunner/emit.py | 29 ++++++++++++++- eventing/k8s/base/configmap.yaml | 17 +++++++++ 9 files changed, 259 insertions(+), 15 deletions(-) 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/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..c755df93 100644 --- a/eventing/eventbridge/kafka_in.py +++ b/eventing/eventbridge/kafka_in.py @@ -1,4 +1,15 @@ -"""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, an +ntfy priority-5 alert, and `raw_json` retained for audit — so the forgery attempt is +visible and reviewable instead of invisible. +""" from __future__ import annotations import threading @@ -7,7 +18,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 +31,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 +43,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 +75,39 @@ 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. + d["phase"] = "error" + d["data"] = {"text": f"unverified response rejected: {why}", + "signature_rejected": True, "reason": str(why)} # §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..30d21e1b 100644 --- a/eventing/eventbridge/kafka_out.py +++ b/eventing/eventbridge/kafka_out.py @@ -3,16 +3,37 @@ 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. + """ + 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 +62,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 +78,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 +94,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 18c252d8..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, @@ -272,8 +285,11 @@ def _handle(self, rec) -> None: if self._cfg.require_signature: from shared import signing - ok, why = signing.verify_event(evt, self._cfg) + # 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" From 669b932b04fb66a73a9aa7d68f82daeb70b4ad33 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 20:55:14 -0400 Subject: [PATCH 4/8] test(eventing): cover the request reject branches, and fix verify_event's key forms MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `consume.py`'s verification branch had no tests at all. Until signing had a producer there was no way to reach the accept side, which left the reject side indistinguishable from "always rejects" — the same property a kill switch has. Ten tests, each asserting the three things the stale-request guard asserts, because a rejection that loses the offset is a poison pill that blocks the partition forever: the request did not reach the router, a named counter says why, and the offset still advanced. Covered: an approved agent runs; unsigned is refused; a real signature from an unapproved key is refused (the forgery the demo turns on — being unforgeable was never the point, being unapproved is); a rogue key claiming an approved kid is refused; an unnamed token is refused when the set is ambiguous but accepted when there is exactly one key; the single-key path still works with no keyset; the default path skips verification entirely; a stale request is dropped by the age guard *before* costing a ~150 ms verification; and the consume loop survives a rejection and goes on to process the next record. Writing the back-compat test surfaced a pre-existing bug in `verify_event`: a hex- or base64-encoded **public key** never verified. A 32-byte file is ambiguous between a seed and a public key and the code already tried both readings for the raw case, but for an encoded file it only tried `public_key(load_seed(path))`, which derives the wrong key when the file already holds a public one. Since `ER_VERIFY_KEY_PATH` is documented to accept either, all six combinations now do, pinned by a parametrized test. Widening the accepted *encodings* does not widen the accepted *keys* — a separate test asserts an unrelated key is still refused. 528 passed, 5 skipped; ruff clean. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/shared/signing.py | 27 +++-- eventing/tests/test_consume_phase1.py | 164 +++++++++++++++++++++++++- eventing/tests/test_signing.py | 39 ++++++ 3 files changed, 216 insertions(+), 14 deletions(-) diff --git a/eventing/shared/signing.py b/eventing/shared/signing.py index e39bd75b..235f3fa8 100644 --- a/eventing/shared/signing.py +++ b/eventing/shared/signing.py @@ -456,18 +456,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_signing.py b/eventing/tests/test_signing.py index 1a47302f..25dcd00d 100644 --- a/eventing/tests/test_signing.py +++ b/eventing/tests/test_signing.py @@ -7,6 +7,7 @@ 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 @@ -279,3 +280,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 From 13b49d32e00d2ea2969d7a6b3d77632f87d9a883 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 21:01:57 -0400 Subject: [PATCH 5/8] test(eventing): first behavioural tests for the responses consumer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `eventbridge/kafka_in.py` had no behavioural test at all — it was only grepped as source text by `test_groups.py`. That was the gap worth closing before adding a verification path to it, because the module decodes and writes to SQLite outside any `try`: anything raising between the poll and the store write ends the `for`, ends the `while`, and the thread is gone while the process stays up and the pod reports healthy. Eleven tests. The two that matter most assert the loop survives a verifier that raises — processing both records rather than dying on the first — and that it fails closed while enforcing but open in audit mode, because a broken check has no opinion to act on when nothing is being enforced. Both were confirmed load-bearing by removing the guard and watching exactly those two fail with the exception escaping. The rest: the default path stores an unsigned event unchanged; an approved agent's response is untouched; a forged response is stored as `phase=error` rather than dropped, with the forged text absent from what a reader sees as the answer; an unapproved key and a tampered payload are both rejected; audit mode reports without rewriting; a bridge-signed group event is accepted while an approved runner forging one is not. Driving the loop needed care worth recording: `run()` checks `stopping` before the `while` and again at the top of each record, so the fake consumer sets the flag when iteration resumes after the last record — stopping any earlier skips the record, any later spins forever on an exhausted iterator. 539 passed, 5 skipped; ruff clean. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/tests/test_kafka_in.py | 270 ++++++++++++++++++++++++++++++++ 1 file changed, 270 insertions(+) create mode 100644 eventing/tests/test_kafka_in.py diff --git a/eventing/tests/test_kafka_in.py b/eventing/tests/test_kafka_in.py new file mode 100644 index 00000000..a42ef18e --- /dev/null +++ b/eventing/tests/test_kafka_in.py @@ -0,0 +1,270 @@ +"""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 payload must not survive into what a reader sees as the answer. + assert "Ship the goods" not in json.dumps(rows[0]["data"]) + assert seen and seen[0]["phase"] == "error", "downstream sees the rejection too" + + +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 + + +# ---- 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") From 889ca529aa53d69af0c7e3ef892d0d0b5be6b3e6 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Wed, 30 Sep 2026 10:10:44 -0400 Subject: [PATCH 6/8] feat(eventing): cover submitter, submitteriss and groupid by the signature MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Last step of #2607, deliberately last: adding these before a producer signed would have invalidated canonicalisation twice for no benefit. `groupid` is required by DESIGN_PHASE1.md §21.9.9 and its absence was the sharpest gap — a signature said nothing about batch membership, so a forged `groupid` could move a response into another batch and corrupt its fan-in counts. Tests now pin all three on the wire, including that *adding* a `groupid` to a signed ungrouped event is detected: the absence of an attribute is signed too, or an ungrouped response could be adopted into a batch it was never part of. No compatibility flag. Nothing had ever published a signed event, and both services ship from this repo against one ConfigMap, so no supported configuration has a signing producer meeting a verifying consumer at different versions. The release note that matters: if signing is already enabled, upgrade both together, because every grouped request now carries a signed `groupid` an older verifier omits when it recomputes. `test_canonical_changes_when_any_signed_attribute_changes` now derives its list from `SIGNED_ATTRS` instead of hardcoding it — the hardcoded version silently stopped covering whatever was added next, which is exactly what happened here. Documentation corrected rather than left contradicting the code. Four places claimed the old state: `shared/ce.py` said submitter was unsigned and forgeable, `eventbridge/auth.py` said proving it "needs a producer that actually signs (today nothing does)", `shared/keyset.py` said nothing consulted it, and `agentdocs/README.md` summarised Phase 2 as blocked on nothing signing. Each now states the bounded claim instead: a signature proves *EventBridge asserted this submitter*, not that the submitter is who they say, and on an unsigned event the attribute stays forgeable — so the property only holds where verification is on. Also revised `_scalar_mult`'s side-channel note, which justified non-constant-time scalar multiplication partly on "nothing exposes a remote timing oracle over sign()". With EB_SIGNING_KEY_PATH set, EventBridge signs on the HTTP request path, so that is no longer strictly true; the note now says what still makes it acceptable rather than resting on a premise this change removed. 543 passed, 5 skipped; ruff clean. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/DESIGN_PHASE2.md | 120 +++++++++++++------ eventing/agentdocs/IMPLEMENTATION_REPORT1.md | 3 + eventing/agentdocs/README.md | 2 +- eventing/eventbridge/auth.py | 21 +++- eventing/shared/ce.py | 13 +- eventing/shared/signing.py | 33 +++-- eventing/tests/test_keyset.py | 47 ++++++++ eventing/tests/test_signing.py | 13 +- 8 files changed, 197 insertions(+), 55 deletions(-) diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md index 28379122..50f6593d 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,50 @@ 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 honest limit: + this proves *who finished a run*, not *what it said along the way* — forged + `final=false` frames are still possible and still render. +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 +363,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 6f8d940d..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 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/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/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/signing.py b/eventing/shared/signing.py index 235f3fa8..a6cbc0b6 100644 --- a/eventing/shared/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: @@ -79,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. @@ -175,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: diff --git a/eventing/tests/test_keyset.py b/eventing/tests/test_keyset.py index 812243e7..2921daa1 100644 --- a/eventing/tests/test_keyset.py +++ b/eventing/tests/test_keyset.py @@ -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" diff --git a/eventing/tests/test_signing.py b/eventing/tests/test_signing.py index 25dcd00d..24e41bb6 100644 --- a/eventing/tests/test_signing.py +++ b/eventing/tests/test_signing.py @@ -118,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" From b8fac88b3b1cdcc1e05e1db1fa04a04c4a69099f Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Wed, 30 Sep 2026 10:17:20 -0400 Subject: [PATCH 7/8] fix(eventing): match response verification to the terminal-only signing policy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Caught by running the demo against the live broker, which is the only reason it was caught: every test until now used terminal events, so all of them passed. `emit()` signs terminal events only — a signature costs ~150-200 ms and `emit()` runs for every stdout frame. But `response_decision` verified all-or-nothing, so a genuine run's streamed frames carry no signature and were all rewritten to `phase="error"`. With enforcement on, every real run would have rendered as a column of red cards with a single valid answer at the end. An event that carries **no** signature and is **not** terminal is now passed through. The security property is unchanged where it matters: * an unsigned **terminal** event is still refused — that is the one the transcript presents as the answer, and refusing it is the entire control; * a non-terminal frame that *presents* a signature is still verified, because the exemption is for absence, not for failure; * group lifecycle events count as terminal. They carry no `final` attribute at all, so testing `final` alone would have classified one as an unsigned intermediate frame and waved it through — exactly the forged `group.completed` the bridge kid exists to catch. Verified on the real broker end to end: a genuine run (one streamed frame plus a signed terminal) renders unchanged while a forged terminal on the same correlation is rejected. Also verified the signed-request path with all four signed attributes travelling as `ce_` headers, and that tampering with `submitter` in flight fails verification. Six new tests, including the realistic shape — a forgery arriving alongside real streaming output, where only it is rewritten. 549 passed, 5 skipped; ruff clean. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/DESIGN_PHASE2.md | 17 +++++++++-- eventing/shared/signing.py | 26 +++++++++++++++++ eventing/tests/test_kafka_in.py | 40 +++++++++++++++++++++++++ eventing/tests/test_keyset.py | 45 +++++++++++++++++++++++++++++ 4 files changed, 125 insertions(+), 3 deletions(-) diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md index 50f6593d..f8fc0b0c 100644 --- a/eventing/agentdocs/DESIGN_PHASE2.md +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -314,9 +314,20 @@ from. 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 honest limit: - this proves *who finished a run*, not *what it said along the way* — forged - `final=false` frames are still possible and still render. + 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 diff --git a/eventing/shared/signing.py b/eventing/shared/signing.py index a6cbc0b6..bf939e82 100644 --- a/eventing/shared/signing.py +++ b/eventing/shared/signing.py @@ -402,6 +402,19 @@ def verify_request(event, cfg, ks=None) -> tuple[bool, str]: 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. @@ -422,6 +435,17 @@ def response_decision(event, ks, *, require: bool, `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" @@ -429,6 +453,8 @@ def response_decision(event, ks, *, require: bool, 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 diff --git a/eventing/tests/test_kafka_in.py b/eventing/tests/test_kafka_in.py index a42ef18e..ac1c57d8 100644 --- a/eventing/tests/test_kafka_in.py +++ b/eventing/tests/test_kafka_in.py @@ -191,6 +191,46 @@ def test_audit_mode_reports_without_rewriting(tmp_path, monkeypatch): 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 json.dumps(rows[2]["data"]) + + # ---- group lifecycle events -------------------------------------------------- def test_a_group_event_signed_by_the_bridge_is_accepted(tmp_path, monkeypatch): diff --git a/eventing/tests/test_keyset.py b/eventing/tests/test_keyset.py index 2921daa1..95a4ff1f 100644 --- a/eventing/tests/test_keyset.py +++ b/eventing/tests/test_keyset.py @@ -443,6 +443,51 @@ def test_response_decision_requires_the_bridge_kid_on_group_events(tmp_path): 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.""" From 3264e8dd9700b57c65facb0da9ed2d88fe36d331 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Wed, 30 Sep 2026 15:33:43 -0400 Subject: [PATCH 8/8] fix(eventing): keep the forensic record when a response is rejected MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit @aslom's must-fix on #879, and he was exactly right: the audit trail the reject path promises was not kept. `Store.insert_response` derives BOTH `data_json` and `raw_json` from the single dict it is handed, so mutating the envelope in place overwrote the evidence with the notice about it. The refused payload, the phase it claimed, and the attributes the signature covered were all gone — leaving a signature that could no longer be checked against anything, and nothing for an operator following the docstring to review. Reproduced before fixing, matching his output exactly: forged payload was: {"text":"TRANSFER THE FUNDS — forged answer"} is it anywhere in raw_json? -> False original phase 'result' in raw_json? -> False signature retained? -> True The envelope is now copied rather than mutated, and the original is preserved under `data["rejected"]` — attrs, payload, claimed source and signature. What a reader sees is unchanged: `phase="error"`, the rejection notice in `data["text"]`, and the forged text never presented as the answer. Two tests, the second being the sharper one: a retained signature is only worth keeping if the attributes it covered are kept with it, so a validly-signed event rejected for naming an unapproved kid is rebuilt from the stored record and verified offline against the key that signed it. That is what makes an incident reviewable rather than merely logged. Two existing assertions were checking `"Ship the goods" not in json.dumps(data)`, which asserted the bug. They now assert the real property — absent from `data["text"]`, so not presented as the answer — while retention is asserted separately. Also added the number an operator actually hits, per his second note: ~200 batch members is ~40 s of signing, past a common 30 s client timeout. `POST /v0/groups` honours `Idempotency-Key` (verified: `HTTP_IDEMPOTENCY_KEY` in handlers.py, and a key hit returns the original group), so the retry is survivable — but only if the caller sends the header, which is the part worth saying. 551 passed, 5 skipped; `ruff check` and `ruff format --check` clean. Verified against the live broker: a forgery published with raw Kafka access renders as a rejection while its payload, claimed source and claimed phase survive in `raw_json`. Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/eventbridge/kafka_in.py | 34 +++++++++++++++++---- eventing/eventbridge/kafka_out.py | 5 +++ eventing/tests/test_kafka_in.py | 51 +++++++++++++++++++++++++++++-- 3 files changed, 81 insertions(+), 9 deletions(-) diff --git a/eventing/eventbridge/kafka_in.py b/eventing/eventbridge/kafka_in.py index c755df93..73ffe880 100644 --- a/eventing/eventbridge/kafka_in.py +++ b/eventing/eventbridge/kafka_in.py @@ -6,9 +6,16 @@ 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, an -ntfy priority-5 alert, and `raw_json` retained for audit — so the forgery attempt is -visible and reviewable instead of invisible. +`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 @@ -105,9 +112,24 @@ def run(self) -> None: # "(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. - d["phase"] = "error" - d["data"] = {"text": f"unverified response rejected: {why}", - "signature_rejected": True, "reason": str(why)} + # + # `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 30d21e1b..2e0aa649 100644 --- a/eventing/eventbridge/kafka_out.py +++ b/eventing/eventbridge/kafka_out.py @@ -22,6 +22,11 @@ class Producer: 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, diff --git a/eventing/tests/test_kafka_in.py b/eventing/tests/test_kafka_in.py index ac1c57d8..f6a36e1a 100644 --- a/eventing/tests/test_kafka_in.py +++ b/eventing/tests/test_kafka_in.py @@ -152,11 +152,54 @@ def test_a_forged_response_is_stored_as_an_error_not_dropped(tmp_path, monkeypat assert "unverified response rejected" in rows[0]["data"]["text"] assert "no ce_signature" in rows[0]["data"]["reason"] assert c.rejected == 1 - # The forged payload must not survive into what a reader sees as the answer. - assert "Ship the goods" not in json.dumps(rows[0]["data"]) + # 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( @@ -228,7 +271,9 @@ def test_a_forged_terminal_is_rejected_among_genuine_frames(tmp_path, monkeypatc 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 json.dumps(rows[2]["data"]) + 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 --------------------------------------------------