From 6ea4d551cb0b362e5c315f94ffebeb7bc678d5fd Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Mon, 28 Sep 2026 19:49:02 -0400 Subject: [PATCH 1/4] feat(eventing): GitHub sign-in and an approved-user list MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replaces the made-up identity from EB_AUTH_TOKENS with a real, externally verified one, and separates "I do not know you" from "I know you and you are not approved". - eventbridge/ghauth.py: OAuth device flow plus GET /user. GitHub does NOT issue a verifiable token for user login — it returns an opaque string with no signature and no claims — so EventBridge cannot verify locally and must ask GitHub who holds it. That makes sign-in depend on GitHub being reachable, and a failed lookup is a 401 rather than an allow. - The cache is load-bearing, not an optimisation: without it every request spends one of 5000 hourly API calls and adds GitHub's latency to the request path. Keyed by sha256(token), so it never holds a usable credential. Failures are not cached, so a revoked token stops working promptly. - No scopes are requested. GET /user returns the login for an unscoped token, so the app asks for the least access that answers "who is this". - auth.resolve() returns (identity, issuer, status, reason). 401 means unauthenticated; 403 means authenticated and not approved, and names the login that was refused because that is what makes it actionable. No WWW-Authenticate on a 403: retrying with another credential is not the fix. - ce_submitter_iss records who vouched — "github", or absent for a static token. Without it a reader cannot tell a verified identity from a name typed into an env var. - EB_AUTH_TOKENS still works, as a fallback that keeps tests off the network, keeps an offline demo possible, and gives an operator a break-glass credential when GitHub is unreachable. - An empty approved list denies everyone. The other reading — empty means everybody — would turn a missing variable into an open door. - Logins compare case-insensitively, because GitHub logins are. - CLI: login / logout / whoami. login prints the code and blocks rather than opening a browser, which fails silently over SSH and in a container. The token is stored 0600. 401/403 now print what to do instead of a traceback. - The client id is a committed default: the device flow has no client secret, so unlike an ntfy topic it is not a capability. Verified against a real GitHub account end to end: 401 unauthenticated, 401 on a garbage token, 403 for a real user off the list, and an approved user's login reaching the wire as ce_submitter:mrsabath / ce_submitter_iss:github. Tests: 450 passed, 5 skipped (was 410/5). No test touches the network. Refs: rossoctl/rossoctl#2606 Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/eventbridge/auth.py | 110 ++++- eventing/eventbridge/config.py | 24 +- eventing/eventbridge/config.toml | 8 + eventing/eventbridge/ghauth.py | 234 +++++++++++ eventing/eventbridge/group_service.py | 5 +- eventing/eventbridge/handlers.py | 32 +- eventing/eventbridge/kafka_out.py | 2 + eventing/shared/ce.py | 4 + .../skills/eventbridge/eventbridge-cli.py | 148 ++++++- eventing/tests/test_ghauth.py | 385 ++++++++++++++++++ 10 files changed, 920 insertions(+), 32 deletions(-) create mode 100644 eventing/eventbridge/ghauth.py create mode 100644 eventing/tests/test_ghauth.py diff --git a/eventing/eventbridge/auth.py b/eventing/eventbridge/auth.py index abec34e4..c535ae1d 100644 --- a/eventing/eventbridge/auth.py +++ b/eventing/eventbridge/auth.py @@ -3,6 +3,10 @@ Scope is deliberately narrow: this answers "who is asking me to run an agent?" on the two routes that CREATE work. It is not a general authorization layer. +Two mechanisms live here. `ghauth` verifies a GitHub sign-in and is the real +path; the static token map below is the fallback that keeps tests off the network +and an offline demo possible. `resolve()` is the entry point that picks. + Two properties worth stating, because both are easy to get wrong: * **Constant-time comparison.** Tokens are compared with @@ -51,35 +55,43 @@ def parse_tokens(raw: str) -> dict[str, str]: return out -def resolve_identity(environ: dict[str, Any], tokens: dict[str, str]) -> tuple[str | None, str | None]: - """Resolve the caller from a WSGI environ. - - Returns `(identity, error)`: - - * `(None, None)` — auth is disabled (no tokens configured). Callers treat - this as "allowed, anonymous", which keeps the default - demo path working unchanged. - * `(name, None)` — a valid credential for `name`. - * `(None, reason)` — reject with 401; `reason` is safe to return to the - client (it never echoes the presented token). +def _bearer(environ: dict[str, Any]) -> tuple[str | None, str | None]: + """Pull the bearer token out of the environ. `(token, error)`. Reads only headers. It must not touch `wsgi.input`: `handlers._read_json` reads the body from a non-seekable stream, so consuming it here would leave every downstream handler with an empty body. """ - if not tokens: - return None, None - header = environ.get("HTTP_AUTHORIZATION") or "" if not header: return None, "authentication required" - scheme, _, presented = header.partition(" ") if scheme.lower() != "bearer": return None, "unsupported authentication scheme; expected Bearer" presented = presented.strip() if not presented: return None, "empty bearer token" + return presented, None + + +def resolve_identity(environ: dict[str, Any], tokens: dict[str, str]) -> tuple[str | None, str | None]: + """Resolve the caller from a static token map. + + Returns `(identity, error)`: + + * `(None, None)` — auth is disabled (no tokens configured). Callers treat + this as "allowed, anonymous", which keeps the default + demo path working unchanged. + * `(name, None)` — a valid credential for `name`. + * `(None, reason)` — reject with 401; `reason` is safe to return to the + client (it never echoes the presented token). + """ + if not tokens: + return None, None + + presented, why = _bearer(environ) + if why: + return None, why # Compare against every configured token so the work done is independent of # which entry matches (and of whether any does). `compare_digest` on str @@ -95,3 +107,71 @@ def resolve_identity(environ: dict[str, Any], tokens: dict[str, str]) -> tuple[s if matched is None: return None, "invalid bearer token" return matched, None + + +# ---- the combined entry point ---------------------------------------------- + +def resolve(environ: dict[str, Any], cfg, *, cache=None, + fetch=None) -> tuple[str | None, str | None, int | None, str | None]: + """Resolve the caller. `(identity, issuer, status, reason)` — handlers call this. + + `issuer` records *who vouched* for the identity: `"github"` for a verified + sign-in, `None` for a static token. A reader of the event can then tell a + real identity from a name an operator typed into an environment variable. + + `status` is the HTTP status to refuse with, and it carries real information: + + * `401` — "I do not know you": no credential, a malformed one, or one GitHub + does not recognise. + * `403` — "I know exactly who you are, and you are not approved." A real, + authenticated person who is not on the list. + + Collapsing those into one answer would tell an operator less, and would tell + a user debugging their own access much less. + + Two mechanisms, checked in order: + + 1. **GitHub sign-in**, when a client id and an approved-user list are + configured. The real path. + 2. **Static tokens** (`EB_AUTH_TOKENS`), as a fallback. Tests must not reach + the network, and an offline demo has to stay possible, so this stays. + + With neither configured the result is all-`None` — allowed and + anonymous, which is what keeps the default demo working out of the box. + """ + from eventbridge import ghauth + + github_on = bool(getattr(cfg, "github_client_id", "")) or bool( + getattr(cfg, "allowed_users", frozenset())) + + if github_on: + presented, why = _bearer(environ) + if why: + # Fall back to a static token only when one could match; otherwise a + # missing header is simply unauthenticated. + if not cfg.auth_tokens: + return None, None, 401, why + name, why2 = resolve_identity(environ, cfg.auth_tokens) + return (name, None, None, None) if name else (None, None, 401, why2) + + login, err = ghauth.resolve( + presented, cache, **({"fetch": fetch} if fetch else {})) + if login is None: + # A static token is checked before giving up, so an operator can keep + # a break-glass credential alongside GitHub sign-in. + if cfg.auth_tokens: + name, _ = resolve_identity(environ, cfg.auth_tokens) + if name: + return name, None, None, None + return None, None, 401, err or "could not identify this token" + + if not ghauth.is_allowed(login, cfg.allowed_users): + # Name the login in the refusal: the user knows who they are, and + # being told which identity was refused is what makes it actionable. + return None, None, 403, f"{login} is not on the approved-user list" + return login, "github", None, None + + name, why = resolve_identity(environ, cfg.auth_tokens) + if why: + return None, None, 401, why + return name, None, None, None diff --git a/eventing/eventbridge/config.py b/eventing/eventbridge/config.py index 4828ef45..73e9d879 100644 --- a/eventing/eventbridge/config.py +++ b/eventing/eventbridge/config.py @@ -6,7 +6,7 @@ import tomllib from dataclasses import dataclass, field -from eventbridge import auth +from eventbridge import auth, ghauth @dataclass @@ -47,6 +47,17 @@ class Cfg: # only via EB_AUTH_TOKENS — deliberately never read from config.toml, which # is committed (tests/test_manifests.py pins the same rule for NTFY_TOKEN). auth_tokens: dict[str, str] = field(default_factory=dict) + # GitHub sign-in. The client id is PUBLIC — the OAuth device flow has no + # client secret, which is why it can be committed while an ntfy topic cannot. + github_client_id: str = "" + # Who may submit. Empty means nobody is approved: if GitHub sign-in is + # configured at all, an operator has to say who may use it, because the other + # reading ("empty means everybody") turns a missing variable into an open door. + allowed_users: frozenset[str] = frozenset() + # Token -> login lookups are cached for this long. Not an optimisation: a + # 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 ntfy: NtfyCfg = field(default_factory=NtfyCfg) @@ -78,6 +89,11 @@ def load() -> Cfg: phases=tuple(n.get("phases", cfg.ntfy.phases)), ) cfg.public_base_url = n.get("public_base_url", cfg.public_base_url) + g = d.get("github", {}) + cfg.github_client_id = g.get("client_id", cfg.github_client_id) + if g.get("allowed_users"): + cfg.allowed_users = frozenset( + str(u).strip().lower() for u in g["allowed_users"] if str(u).strip()) # env overrides e = os.environ.get @@ -98,6 +114,12 @@ def load() -> Cfg: auth_env = e("EB_AUTH_TOKENS") if auth_env: cfg.auth_tokens = auth.parse_tokens(auth_env) + cfg.github_client_id = e("EB_GITHUB_CLIENT_ID", cfg.github_client_id) + allowed_env = e("EB_ALLOWED_USERS") + if allowed_env: + 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))) 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/config.toml b/eventing/eventbridge/config.toml index 45328260..c1de2424 100644 --- a/eventing/eventbridge/config.toml +++ b/eventing/eventbridge/config.toml @@ -14,3 +14,11 @@ topic = "" token = "" phases = ["result", "error"] public_base_url = "http://127.0.0.1:8080" + +[github] +# The OAuth App client id is public: the device flow has no client secret, so +# unlike an ntfy topic this is safe to commit. Override with EB_GITHUB_CLIENT_ID. +client_id = "" +# Who may submit. Empty means nobody, so sign-in stays off until an operator +# names the approved logins. Override with EB_ALLOWED_USERS. +allowed_users = [] diff --git a/eventing/eventbridge/ghauth.py b/eventing/eventbridge/ghauth.py new file mode 100644 index 00000000..b69f2912 --- /dev/null +++ b/eventing/eventbridge/ghauth.py @@ -0,0 +1,234 @@ +"""GitHub sign-in for the submit path: device flow, then who the token belongs to. + +Stdlib only (`urllib`), as §1.1 requires. + +**GitHub does not issue a verifiable token for user login.** The device flow +returns an *opaque* access token — no signature, no claims, nothing to check +offline. (The JWKS at `token.actions.githubusercontent.com` is for Actions +workloads, not users.) So EventBridge cannot verify a token locally; it has to +ask GitHub who holds it, via `GET /user`. + +Three consequences, all deliberate rather than accidental: + +* **Sign-in depends on GitHub being reachable.** A failed lookup is a `401`, not + an allow. Failing closed is the only safe direction when the question is "who + is this". +* **The cache is load-bearing, not an optimisation.** Without it every request + costs an API call against a 5000/hour budget. It is keyed by a hash of the + token, never the token, so a memory dump or a log line cannot yield a working + credential. +* **No scopes are requested.** `GET /user` returns the login for an unscoped + token, so the app asks for the least access that answers the question. Nothing + here can read a repository. + +The approved-user list is separate from authentication on purpose: `401` means +"I do not know you", `403` means "I know you and you are not on the list". Those +are different answers and collapsing them tells an operator less. +""" +from __future__ import annotations + +import hashlib +import json +import time +import urllib.error +import urllib.parse +import urllib.request + +DEVICE_CODE_URL = "https://github.com/login/device/code" +ACCESS_TOKEN_URL = "https://github.com/login/oauth/access_token" +USER_URL = "https://api.github.com/user" +VERIFICATION_URL = "https://github.com/login/device" + +# GitHub's own default poll interval, used when a response omits one. +DEFAULT_INTERVAL_S = 5 +# Added to the interval when GitHub answers `slow_down`, as its docs specify. +SLOW_DOWN_PENALTY_S = 5 + +_UA = "rossoctl-eventing" + + +def _post_form(url: str, fields: dict[str, str], timeout: float) -> dict: + """POST a form, read JSON back. GitHub returns form-encoded unless asked.""" + body = urllib.parse.urlencode(fields).encode() + req = urllib.request.Request( + url, data=body, method="POST", + headers={"Accept": "application/json", "User-Agent": _UA, + "Content-Type": "application/x-www-form-urlencoded"}) + with urllib.request.urlopen(req, timeout=timeout) as r: + return json.loads(r.read() or b"{}") + + +# ---- device flow (the CLI side) --------------------------------------------- + +def start_device_flow(client_id: str, *, timeout: float = 10.0) -> dict: + """Ask GitHub for a device code and the code the user will type. + + Returns GitHub's response: `device_code`, `user_code`, `verification_uri`, + `expires_in`, `interval`. Raises `RuntimeError` on `device_flow_disabled`, + which means the app exists but device flow was never ticked on — the single + most common setup mistake, so it gets its own message. + """ + out = _post_form(DEVICE_CODE_URL, + # An empty scope is deliberate: see the module docstring. + {"client_id": client_id, "scope": ""}, timeout) + if "error" in out: + err = out.get("error") + if err == "device_flow_disabled": + raise RuntimeError( + "this OAuth App does not have Device Flow enabled — tick " + "'Enable Device Flow' on the app page at " + "https://github.com/settings/developers") + raise RuntimeError(f"GitHub refused the device code request: {err} " + f"({out.get('error_description', 'no detail')})") + if not out.get("device_code") or not out.get("user_code"): + raise RuntimeError(f"unexpected device code response: {out!r}") + return out + + +def poll_for_token(client_id: str, device_code: str, *, interval: float | None = None, + expires_in: float = 900.0, timeout: float = 10.0, + sleep=time.sleep, now=time.monotonic) -> str: + """Poll until the user authorises, then return the access token. + + Honours GitHub's pacing contract: `authorization_pending` means keep going, + `slow_down` means keep going but add five seconds. Polling faster than asked + gets the app rate-limited, which would break sign-in for everyone rather than + just this caller. + + `sleep` and `now` are injectable so tests exercise the pacing without + actually waiting. + """ + wait = float(interval or DEFAULT_INTERVAL_S) + deadline = now() + expires_in + while True: + if now() >= deadline: + raise TimeoutError( + "the device code expired before it was authorised — run login again") + sleep(wait) + out = _post_form(ACCESS_TOKEN_URL, { + "client_id": client_id, "device_code": device_code, + "grant_type": "urn:ietf:params:oauth:grant-type:device_code", + }, timeout) + token = out.get("access_token") + if token: + return str(token) + err = out.get("error") + if err == "authorization_pending": + continue + if err == "slow_down": + wait = float(out.get("interval", wait)) + SLOW_DOWN_PENALTY_S + continue + if err == "expired_token": + raise TimeoutError( + "the device code expired before it was authorised — run login again") + if err == "access_denied": + raise RuntimeError("sign-in was declined in the browser") + raise RuntimeError(f"GitHub refused the token request: {err} " + f"({out.get('error_description', 'no detail')})") + + +# ---- resolving a token to a login (the server side) ------------------------- + +def _token_key(token: str) -> str: + """Cache key. A hash, so the cache never holds a usable credential.""" + return hashlib.sha256(token.encode()).hexdigest() + + +def fetch_login(token: str, *, timeout: float = 5.0) -> tuple[str | None, str | None]: + """`GET /user` -> `(login, error)`. Exactly one is None. + + A 401 from GitHub means the token is unknown or revoked, which is the + caller's problem and safe to report. Anything else — a network failure, a 500, + a rate limit — is ours, and is reported as an unavailability rather than a + rejection, because telling a legitimate user "invalid credential" when + GitHub was merely unreachable sends them debugging the wrong thing. + """ + req = urllib.request.Request( + USER_URL, headers={"Authorization": f"Bearer {token}", + "Accept": "application/vnd.github+json", + "User-Agent": _UA}) + try: + with urllib.request.urlopen(req, timeout=timeout) as r: + body = json.loads(r.read() or b"{}") + except urllib.error.HTTPError as e: + if e.code == 401: + return None, "GitHub rejected this token" + if e.code == 403: + return None, "GitHub rate limit or access restriction (403)" + return None, f"GitHub returned HTTP {e.code}" + except Exception as e: # noqa: BLE001 — urllib raises a wide family here + return None, f"cannot reach GitHub: {e.__class__.__name__}" + login = body.get("login") + if not isinstance(login, str) or not login: + return None, "GitHub returned no login for this token" + return login, None + + +class LoginCache: + """Token hash -> login, with a TTL. + + Not an optimisation: without it every submitted request spends one of 5000 + hourly API calls, and adds GitHub's latency to the request path. Entries are + keyed by hash so the cache cannot be read back into a working credential. + + Negative results are not cached. A revoked token must stop working promptly, + and a GitHub outage must not pin a legitimate user to a failure for the whole + TTL. + """ + + def __init__(self, ttl_s: float = 300.0, now=time.monotonic) -> None: + self._ttl = float(ttl_s) + self._now = now + self._entries: dict[str, tuple[float, str]] = {} + + def get(self, token: str) -> str | None: + hit = self._entries.get(_token_key(token)) + if hit is None: + return None + expires, login = hit + if self._now() >= expires: + self._entries.pop(_token_key(token), None) + return None + return login + + def put(self, token: str, login: str) -> None: + self._entries[_token_key(token)] = (self._now() + self._ttl, login) + + def clear(self) -> None: + self._entries.clear() + + def __len__(self) -> int: + return len(self._entries) + + +def resolve(token: str, cache: LoginCache | None = None, *, + fetch=fetch_login) -> tuple[str | None, str | None]: + """Token -> `(login, error)`, consulting the cache first.""" + if cache is not None: + cached = cache.get(token) + if cached is not None: + return cached, None + login, err = fetch(token) + if login and cache is not None: + cache.put(token, login) + return login, err + + +def parse_allowed_users(raw: str) -> frozenset[str]: + """`"alice, bob"` -> `{"alice", "bob"}`, compared case-insensitively. + + GitHub logins are case-insensitive, so `Alice` and `alice` are one account. + Comparing case-sensitively would let a real approved user be refused because + of how they typed it, which reads as a broken deployment. + """ + return frozenset(u.strip().lower() for u in raw.split(",") if u.strip()) + + +def is_allowed(login: str, allowed: frozenset[str]) -> bool: + """Empty list means nobody is approved — deny. + + The alternative reading, "empty means everybody", would turn a missing + environment variable into an open door. If GitHub sign-in is configured at + all, the operator must say who may use it. + """ + return login.lower() in allowed diff --git a/eventing/eventbridge/group_service.py b/eventing/eventbridge/group_service.py index a9379bd4..9d4ed26b 100644 --- a/eventing/eventbridge/group_service.py +++ b/eventing/eventbridge/group_service.py @@ -55,7 +55,8 @@ def _publish_started(self, groupid, label, expected, min_success, deadline) -> N def submit_members(self, groupid: str, prompts: list[str], *, max_turns: int = 3, model: str | None = None, - submitter: str | None = None) -> list[str]: + submitter: str | None = None, + submitter_iss: str | None = None) -> list[str]: """Publish every member's request, recording membership first. Membership is recorded BEFORE the request is published so a fast agent's @@ -73,7 +74,7 @@ def submit_members(self, groupid: str, prompts: list[str], *, self.producer.publish_request( prompt=prompt, correlationid=corr, sessionuuid=sess, mode="start", model=model, max_turns=max_turns, subject="start", groupid=groupid, - submitter=submitter) + submitter=submitter, submitter_iss=submitter_iss) corrs.append(corr) return corrs diff --git a/eventing/eventbridge/handlers.py b/eventing/eventbridge/handlers.py index 9359ccf7..2bdc593e 100644 --- a/eventing/eventbridge/handlers.py +++ b/eventing/eventbridge/handlers.py @@ -7,7 +7,7 @@ import urllib.parse from typing import Any -from eventbridge import auth +from eventbridge import auth, ghauth from eventbridge.config import Cfg from eventbridge.correlation import REGEX as CORR_REGEX from eventbridge.correlation import Minter @@ -61,8 +61,15 @@ def _read_body(environ, limit: int) -> tuple[bytes | None, str | None]: return body, None -def _deny(start_response, reason: str) -> list[bytes]: - """401 with the Bearer challenge. `_json` already takes extra headers.""" +def _deny(start_response, reason: str, status: int = 401) -> list[bytes]: + """Refuse with 401 or 403. `_json` already takes extra headers. + + The challenge header goes on 401 only. On a 403 the credential was fine and + retrying with a different one is not the remedy, so advertising a scheme + would be misleading. + """ + if status == 403: + return _json(start_response, "403 Forbidden", {"error": reason}) return _json(start_response, "401 Unauthorized", {"error": reason}, [("WWW-Authenticate", auth.CHALLENGE)]) @@ -91,12 +98,15 @@ def __init__(self, cfg: Cfg, store: Store, producer: Producer, minter: Minter, self.producer = producer self.minter = minter self.groups = groups + # One cache for the process, so a burst of requests from the same user + # costs one GitHub call rather than one each. + self.logins = ghauth.LoginCache(cfg.github_cache_ttl_s) # ---- start ---- def start_agent(self, environ, start_response, **_): - submitter, why = auth.resolve_identity(environ, self.cfg.auth_tokens) - if why: - return _deny(start_response, why) + submitter, sub_iss, status, why = auth.resolve(environ, self.cfg, cache=self.logins) + if status: + return _deny(start_response, why, status) body = _read_json(environ) prompt = body.get("prompt") if not prompt: @@ -120,7 +130,7 @@ def start_agent(self, environ, start_response, **_): event_id = self.producer.publish_request( prompt=prompt, correlationid=corr, sessionuuid=sess, mode="start", model=model, max_turns=max_turns, subject="start", - groupid=groupid, submitter=submitter, + groupid=groupid, submitter=submitter, submitter_iss=sub_iss, ) out = { "correlationid": corr, "sessionuuid": sess, @@ -141,9 +151,9 @@ def create_group(self, environ, start_response, **_): the batch scale at all: every request is published before any is watched, so lag reaches N and KEDA scales past one pod. """ - submitter, why = auth.resolve_identity(environ, self.cfg.auth_tokens) - if why: - return _deny(start_response, why) + submitter, sub_iss, status, why = auth.resolve(environ, self.cfg, cache=self.logins) + if status: + return _deny(start_response, why, status) if self.groups is None: return _json(start_response, "503 Service Unavailable", {"error": "group support not enabled"}) @@ -179,7 +189,7 @@ def create_group(self, environ, start_response, **_): corrs = self.groups.submit_members( groupid, [str(x) for x in prompts], max_turns=int(body.get("max_turns", 3)), model=body.get("model"), - submitter=submitter) + submitter=submitter, submitter_iss=sub_iss) return _json(start_response, "202 Accepted", { "groupid": groupid, "created": True, "expected": expected or len(corrs), "members": corrs, diff --git a/eventing/eventbridge/kafka_out.py b/eventing/eventbridge/kafka_out.py index 8d8b3caa..5551740d 100644 --- a/eventing/eventbridge/kafka_out.py +++ b/eventing/eventbridge/kafka_out.py @@ -26,6 +26,7 @@ def publish_request( subject: str = "agent-request", groupid: str | None = None, submitter: str | None = None, + submitter_iss: str | None = None, ) -> str: event = ce.new_event( type=ce.TYPE_REQUEST, @@ -38,6 +39,7 @@ def publish_request( data={"prompt": prompt, "model": model, "max_turns": max_turns}, **({"groupid": groupid} if groupid else {}), **({ce.EXT_SUBMITTER: submitter} if submitter else {}), + **({ce.EXT_SUBMITTER_ISS: submitter_iss} if submitter_iss else {}), ) headers, value = ce.to_kafka_binary(event) future = self._prod.send(self._topic, key=correlationid.encode(), value=value, headers=headers) diff --git a/eventing/shared/ce.py b/eventing/shared/ce.py index 1119724d..e17d7a6c 100644 --- a/eventing/shared/ce.py +++ b/eventing/shared/ce.py @@ -46,6 +46,10 @@ # and forgeable by anyone with write access to the requests topic — it records # who EventBridge believes submitted, not cryptographic proof. 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 +# cannot tell a real identity from a name an operator typed into an env var. +EXT_SUBMITTER_ISS = "submitter_iss" CE_HEADER_PREFIX = "ce_" CORE_ATTRS = {"specversion", "type", "source", "id", "time", diff --git a/eventing/skills/eventbridge/eventbridge-cli.py b/eventing/skills/eventbridge/eventbridge-cli.py index 4a164268..e3d21a81 100644 --- a/eventing/skills/eventbridge/eventbridge-cli.py +++ b/eventing/skills/eventbridge/eventbridge-cli.py @@ -10,6 +10,7 @@ import argparse import json import os +import pathlib import shlex import sys import time @@ -73,16 +74,64 @@ def _add_target_flags(sp): "http://127.0.0.1:8080). A bare hostname is assumed https.") +def _token_path() -> pathlib.Path: + """Where `login` stores the GitHub token. + + XDG if set, else ~/.config — the same place the rest of a user's CLI state + lives, so it is findable and deletable without documentation. + """ + base = os.environ.get("XDG_CONFIG_HOME") or os.path.expanduser("~/.config") + return pathlib.Path(base) / "rossoctl-eventing" / "token" + + +def _read_token() -> str | None: + """The stored token, or $EVENTBRIDGE_TOKEN, or None. + + The environment wins: a CI job or a one-off shell should be able to act as a + different identity without disturbing an interactive login. + """ + env = os.environ.get("EVENTBRIDGE_TOKEN", "").strip() + if env: + return env + try: + tok = _token_path().read_text().strip() + return tok or None + except OSError: + return None + + def _req(method, path, body=None): data = json.dumps(body).encode() if body is not None else None + headers = {} + if body: + headers["Content-Type"] = "application/json" + token = _read_token() + if token: + headers["Authorization"] = f"Bearer {token}" req = urllib.request.Request( - BASE + path, data=data, method=method, - headers={"Content-Type": "application/json"} if body else {}, + BASE + path, data=data, method=method, headers=headers, ) try: with urllib.request.urlopen(req, timeout=15) as r: return json.loads(r.read()) - except urllib.error.HTTPError: + except urllib.error.HTTPError as e: + # 401 and 403 are the two an operator will actually hit, and they mean + # different things. A traceback for either sends someone reading urllib + # internals instead of fixing their access. + if e.code in (401, 403): + try: + detail = json.loads(e.read() or b"{}").get("error", "") + except Exception: # noqa: BLE001 + detail = "" + if e.code == 401: + print(f"✗ not signed in{': ' + detail if detail else ''}\n" + f" run: {sys.argv[0]} login", file=sys.stderr) + else: + print(f"✗ signed in, but not authorised" + f"{': ' + detail if detail else ''}\n" + f" ask an operator to add you to EB_ALLOWED_USERS", + file=sys.stderr) + raise SystemExit(1) from None raise except OSError as e: # Nothing listening, or the name does not resolve. Say what to change rather @@ -443,6 +492,87 @@ def cmd_selftest(args): print(f" · {note}") +# ---- sign-in ---------------------------------------------------------------- + +# The OAuth App client id is PUBLIC: the device flow has no client secret, which +# is why this can be a default in source. Override for a different app. +DEFAULT_GITHUB_CLIENT_ID = "Ov23liA2Z4jfbFRoJKZn" + + +def cmd_login(args): + """GitHub device flow: print a code, wait for the browser, store the token. + + Deliberately prints the code and the URL and then blocks. The alternative — + opening a browser automatically — fails silently over SSH and in a container, + which is where this is most often run. + """ + sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2])) + from eventbridge import ghauth + + client_id = args.client_id or os.environ.get( + "EB_GITHUB_CLIENT_ID") or DEFAULT_GITHUB_CLIENT_ID + try: + start = ghauth.start_device_flow(client_id) + except RuntimeError as e: + print(f"✗ {e}", file=sys.stderr) + raise SystemExit(1) from None + + print(f"\n Open {start.get('verification_uri', ghauth.VERIFICATION_URL)}") + print(f" Enter code: {start['user_code']}\n") + print(" Waiting for you to authorise in the browser (Ctrl-C to abort)...") + + try: + token = ghauth.poll_for_token( + client_id, start["device_code"], + interval=start.get("interval"), + expires_in=float(start.get("expires_in", 900))) + except KeyboardInterrupt: + print("\n✗ aborted", file=sys.stderr) + raise SystemExit(130) from None + except (RuntimeError, TimeoutError) as e: + print(f"✗ {e}", file=sys.stderr) + raise SystemExit(1) from None + + login, err = ghauth.fetch_login(token) + if login is None: + print(f"✗ signed in, but could not read the account: {err}", file=sys.stderr) + raise SystemExit(1) from None + + dest = _token_path() + dest.parent.mkdir(parents=True, exist_ok=True) + dest.write_text(token) + dest.chmod(0o600) # a token is a credential; not world-readable + print(f"✔ signed in as {login}") + print(f" token stored at {dest} (mode 600)") + + +def cmd_logout(args): + """Forget the stored token. Does not revoke it on GitHub.""" + dest = _token_path() + try: + dest.unlink() + print(f"✔ removed {dest}") + except FileNotFoundError: + print("• no stored token") + print(" to revoke access entirely: https://github.com/settings/applications") + + +def cmd_whoami(args): + """Who the stored token belongs to, asked of GitHub rather than assumed.""" + token = _read_token() + if not token: + print("• not signed in", file=sys.stderr) + raise SystemExit(1) + sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2])) + from eventbridge import ghauth + login, err = ghauth.fetch_login(token) + if login is None: + print(f"✗ {err}", file=sys.stderr) + raise SystemExit(1) + src = "$EVENTBRIDGE_TOKEN" if os.environ.get("EVENTBRIDGE_TOKEN") else str(_token_path()) + print(f"✔ {login} (token from {src})") + + def cmd_chat(args): r = _req("GET", f"/v0/agents/{args.correlationid}/turns") print(f"correlationid: {r['correlationid']} {len(r['turns'])} turns final={r['final']}") @@ -504,6 +634,18 @@ def main(): ) sub = p.add_subparsers(dest="op", required=True) + li = sub.add_parser("login", help="sign in with GitHub (device flow)") + li.add_argument("--client-id", default=None, metavar="ID", + help="OAuth App client id (default: $EB_GITHUB_CLIENT_ID, " + "else the built-in demo app)") + li.set_defaults(fn=cmd_login) + + lo = sub.add_parser("logout", help="forget the stored token") + lo.set_defaults(fn=cmd_logout) + + wai = sub.add_parser("whoami", help="show who the stored token belongs to") + wai.set_defaults(fn=cmd_whoami) + r = sub.add_parser("run", help="start an agent with a prompt and print its reply") r.add_argument("prompt") r.add_argument("--max-turns", type=int, default=3) diff --git a/eventing/tests/test_ghauth.py b/eventing/tests/test_ghauth.py new file mode 100644 index 00000000..c540b3cc --- /dev/null +++ b/eventing/tests/test_ghauth.py @@ -0,0 +1,385 @@ +"""GitHub sign-in: device flow, token -> login, and the approved-user list. + +No test here touches the network. The device flow and `GET /user` are exercised +through injected fakes, because a test suite that reaches GitHub would be slow, +flaky, rate-limited, and would fail entirely on a machine without credentials. +""" +from __future__ import annotations + +import io +import json + +import pytest + +from eventbridge import auth, ghauth +from eventbridge.config import Cfg as EbCfg + +# ---- the approved-user list ------------------------------------------------- + +def test_parse_allowed_users_splits_and_lowercases(): + assert ghauth.parse_allowed_users("Alice, BOB ,carol") == {"alice", "bob", "carol"} + + +def test_parse_allowed_users_ignores_blanks(): + assert ghauth.parse_allowed_users(" , ,alice, ") == {"alice"} + + +def test_parse_allowed_users_empty_is_empty(): + assert ghauth.parse_allowed_users("") == frozenset() + + +def test_is_allowed_is_case_insensitive(): + """GitHub logins are case-insensitive, so Alice and alice are one account. + Comparing case-sensitively would refuse a genuinely approved user.""" + allowed = ghauth.parse_allowed_users("mrsabath") + assert ghauth.is_allowed("mrsabath", allowed) + assert ghauth.is_allowed("MrSabath", allowed) + assert ghauth.is_allowed("MRSABATH", allowed) + + +def test_is_allowed_denies_when_the_list_is_empty(): + """The other reading — "empty means everybody" — would turn a missing + environment variable into an open door.""" + assert not ghauth.is_allowed("anyone", frozenset()) + + +def test_is_allowed_denies_a_non_member(): + assert not ghauth.is_allowed("stranger", ghauth.parse_allowed_users("alice,bob")) + + +# ---- the login cache -------------------------------------------------------- + +class FakeClock: + def __init__(self) -> None: + self.t = 1000.0 + + def __call__(self) -> float: + return self.t + + def advance(self, dt: float) -> None: + self.t += dt + + +def test_cache_returns_what_was_put(): + c = ghauth.LoginCache(300.0, now=FakeClock()) + c.put("tok", "alice") + assert c.get("tok") == "alice" + + +def test_cache_misses_an_unknown_token(): + assert ghauth.LoginCache(300.0).get("nope") is None + + +def test_cache_entry_expires(): + clock = FakeClock() + c = ghauth.LoginCache(300.0, now=clock) + c.put("tok", "alice") + clock.advance(299) + assert c.get("tok") == "alice" + clock.advance(2) + assert c.get("tok") is None + + +def test_cache_never_stores_the_token_itself(): + """A memory dump or a careless log of the cache must not yield a credential.""" + c = ghauth.LoginCache(300.0) + c.put("gho_supersecret", "alice") + assert "gho_supersecret" not in repr(c.__dict__) + assert all("gho_supersecret" not in k for k in c._entries) + + +def test_cache_separates_distinct_tokens(): + c = ghauth.LoginCache(300.0) + c.put("tok-a", "alice") + c.put("tok-b", "bob") + assert (c.get("tok-a"), c.get("tok-b")) == ("alice", "bob") + + +# ---- resolve (cache + fetch) ------------------------------------------------ + +def test_resolve_uses_the_cache_and_does_not_refetch(): + calls = [] + + def fetch(token): + calls.append(token) + return "alice", None + + c = ghauth.LoginCache(300.0) + assert ghauth.resolve("tok", c, fetch=fetch) == ("alice", None) + assert ghauth.resolve("tok", c, fetch=fetch) == ("alice", None) + assert len(calls) == 1, "second call should have been served from the cache" + + +def test_resolve_does_not_cache_a_failure(): + """A revoked token must stop working promptly, and a GitHub outage must not + pin a legitimate user to a failure for the whole TTL.""" + c = ghauth.LoginCache(300.0) + ghauth.resolve("tok", c, fetch=lambda t: (None, "GitHub rejected this token")) + assert len(c) == 0 + + +def test_resolve_works_without_a_cache(): + assert ghauth.resolve("tok", None, fetch=lambda t: ("bob", None)) == ("bob", None) + + +# ---- fetch_login error mapping --------------------------------------------- + +class FakeHTTPError(Exception): + def __init__(self, code): + self.code = code + + +def _patched_fetch(monkeypatch, *, raises=None, body=None): + import urllib.error + + class _Err(urllib.error.HTTPError): + def __init__(self, code): + self.code = code # bypass HTTPError's heavy __init__ + + def fake_urlopen(req, timeout=None): + if raises is not None: + raise _Err(raises) + + class _R: + def read(self_inner): + return json.dumps(body or {}).encode() + + def __enter__(self_inner): + return self_inner + + def __exit__(self_inner, *a): + return False + return _R() + + monkeypatch.setattr(ghauth.urllib.request, "urlopen", fake_urlopen) + + +def test_fetch_login_returns_the_login(monkeypatch): + _patched_fetch(monkeypatch, body={"login": "mrsabath"}) + assert ghauth.fetch_login("tok") == ("mrsabath", None) + + +def test_fetch_login_401_is_reported_as_a_rejection(monkeypatch): + _patched_fetch(monkeypatch, raises=401) + login, err = ghauth.fetch_login("tok") + assert login is None and "rejected" in err + + +def test_fetch_login_403_mentions_the_rate_limit(monkeypatch): + _patched_fetch(monkeypatch, raises=403) + login, err = ghauth.fetch_login("tok") + assert login is None and "rate limit" in err + + +def test_fetch_login_500_is_not_reported_as_a_bad_credential(monkeypatch): + """Telling a legitimate user "invalid credential" when GitHub was merely + broken sends them debugging the wrong thing.""" + _patched_fetch(monkeypatch, raises=500) + login, err = ghauth.fetch_login("tok") + assert login is None + assert "rejected" not in err and "500" in err + + +def test_fetch_login_network_failure_says_unreachable(monkeypatch): + def boom(req, timeout=None): + raise OSError("dns go boom") + monkeypatch.setattr(ghauth.urllib.request, "urlopen", boom) + login, err = ghauth.fetch_login("tok") + assert login is None and "cannot reach GitHub" in err + + +def test_fetch_login_missing_login_field(monkeypatch): + _patched_fetch(monkeypatch, body={"id": 1}) + login, err = ghauth.fetch_login("tok") + assert login is None and "no login" in err + + +# ---- device flow ------------------------------------------------------------ + +def _patch_post(monkeypatch, responses): + """Serve queued JSON responses to `_post_form`, recording the URLs called.""" + seen = [] + + def fake_post(url, fields, timeout): + seen.append((url, fields)) + return responses.pop(0) + + monkeypatch.setattr(ghauth, "_post_form", fake_post) + return seen + + +def test_start_device_flow_returns_the_codes(monkeypatch): + seen = _patch_post(monkeypatch, [{ + "device_code": "dc", "user_code": "WDJB-MJHT", + "verification_uri": ghauth.VERIFICATION_URL, "expires_in": 900, "interval": 5}]) + out = ghauth.start_device_flow("cid") + assert out["user_code"] == "WDJB-MJHT" + assert seen[0][0] == ghauth.DEVICE_CODE_URL + + +def test_start_device_flow_requests_no_scopes(monkeypatch): + """Least access that answers "who is this": GET /user needs no scope.""" + seen = _patch_post(monkeypatch, [{"device_code": "dc", "user_code": "U"}]) + ghauth.start_device_flow("cid") + assert seen[0][1]["scope"] == "" + + +def test_start_device_flow_explains_device_flow_disabled(monkeypatch): + """The single most common setup mistake, so it gets its own message.""" + _patch_post(monkeypatch, [{"error": "device_flow_disabled"}]) + with pytest.raises(RuntimeError, match="Enable Device Flow"): + ghauth.start_device_flow("cid") + + +def test_start_device_flow_reports_other_errors(monkeypatch): + _patch_post(monkeypatch, [{"error": "unauthorized_client", + "error_description": "bad client"}]) + with pytest.raises(RuntimeError, match="unauthorized_client"): + ghauth.start_device_flow("cid") + + +def test_poll_returns_the_token_once_authorised(monkeypatch): + _patch_post(monkeypatch, [ + {"error": "authorization_pending"}, + {"error": "authorization_pending"}, + {"access_token": "gho_abc"}, + ]) + slept = [] + tok = ghauth.poll_for_token("cid", "dc", interval=5, sleep=slept.append) + assert tok == "gho_abc" + assert slept == [5.0, 5.0, 5.0] + + +def test_poll_honours_slow_down_by_adding_five_seconds(monkeypatch): + """Polling faster than GitHub asks rate-limits the whole app, not just us.""" + _patch_post(monkeypatch, [ + {"error": "slow_down", "interval": 5}, + {"access_token": "gho_abc"}, + ]) + slept = [] + ghauth.poll_for_token("cid", "dc", interval=5, sleep=slept.append) + assert slept == [5.0, 10.0] + + +def test_poll_raises_on_expired_token(monkeypatch): + _patch_post(monkeypatch, [{"error": "expired_token"}]) + with pytest.raises(TimeoutError, match="expired"): + ghauth.poll_for_token("cid", "dc", interval=1, sleep=lambda s: None) + + +def test_poll_raises_when_the_user_declines(monkeypatch): + _patch_post(monkeypatch, [{"error": "access_denied"}]) + with pytest.raises(RuntimeError, match="declined"): + ghauth.poll_for_token("cid", "dc", interval=1, sleep=lambda s: None) + + +def test_poll_gives_up_at_the_deadline(monkeypatch): + _patch_post(monkeypatch, [{"error": "authorization_pending"}] * 50) + clock = FakeClock() + + def sleep(dt): + clock.advance(dt) + + with pytest.raises(TimeoutError): + ghauth.poll_for_token("cid", "dc", interval=5, expires_in=20, + sleep=sleep, now=clock) + + +# ---- auth.resolve: the 401 / 403 distinction ------------------------------- + +def _env(token: str | None = None) -> dict: + e = {"wsgi.input": io.BytesIO(b""), "REQUEST_METHOD": "POST"} + if token: + e["HTTP_AUTHORIZATION"] = f"Bearer {token}" + return e + + +def _cfg(**kw) -> EbCfg: + return EbCfg(tmpdir="/tmp/x", **kw) + + +def test_resolve_allows_anonymous_when_nothing_is_configured(): + ident, iss, status, why = auth.resolve(_env(), _cfg()) + assert (ident, iss, status, why) == (None, None, None, None) + + +def test_resolve_401_without_a_credential_when_github_is_on(): + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"})) + ident, iss, status, why = auth.resolve(_env(), cfg) + assert status == 401 and ident is None + + +def test_resolve_401_when_github_rejects_the_token(): + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"})) + ident, iss, status, why = auth.resolve( + _env("bad"), cfg, fetch=lambda t: (None, "GitHub rejected this token")) + assert status == 401 and "rejected" in why + + +def test_resolve_403_for_a_real_user_who_is_not_approved(): + """The demo-worthy case: authenticated, and still refused.""" + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"})) + ident, iss, status, why = auth.resolve( + _env("gho_x"), cfg, fetch=lambda t: ("mallory", None)) + assert status == 403 + assert "mallory" in why, "the refusal should name the identity it refused" + assert ident is None + + +def test_resolve_allows_an_approved_user_and_records_the_issuer(): + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"mrsabath"})) + ident, iss, status, why = auth.resolve( + _env("gho_x"), cfg, fetch=lambda t: ("mrsabath", None)) + assert (ident, iss, status) == ("mrsabath", "github", None) + + +def test_resolve_is_case_insensitive_about_the_approved_login(): + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"mrsabath"})) + ident, _, status, _ = auth.resolve( + _env("gho_x"), cfg, fetch=lambda t: ("MrSabath", None)) + assert status is None and ident == "MrSabath" + + +def test_a_static_token_still_works_as_a_break_glass_credential(): + """An operator keeping one alongside GitHub sign-in should not be locked out + when GitHub is unreachable.""" + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"}), + auth_tokens={"emergency": "operator"}) + ident, iss, status, _ = auth.resolve( + _env("emergency"), cfg, fetch=lambda t: (None, "cannot reach GitHub: OSError")) + assert (ident, iss, status) == ("operator", None, None) + + +def test_static_token_path_reports_no_issuer(): + """Absence of an issuer is the signal that this identity was not verified.""" + cfg = _cfg(auth_tokens={"tok": "alice"}) + ident, iss, status, _ = auth.resolve(_env("tok"), cfg) + assert (ident, iss, status) == ("alice", None, None) + + +def test_resolve_401_for_a_wrong_static_token(): + cfg = _cfg(auth_tokens={"tok": "alice"}) + _, _, status, _ = auth.resolve(_env("wrong"), cfg) + assert status == 401 + + +def test_github_on_with_only_an_allowed_list_still_enforces(): + """Setting the list but forgetting the client id must not silently allow.""" + cfg = _cfg(allowed_users=frozenset({"alice"})) + _, _, status, _ = auth.resolve(_env(), cfg) + assert status == 401 + + +def test_the_cache_means_one_github_call_for_repeated_requests(): + calls = [] + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"})) + cache = ghauth.LoginCache(300.0) + + def fetch(token): + calls.append(token) + return "alice", None + + for _ in range(5): + ident, _, status, _ = auth.resolve(_env("gho_x"), cfg, cache=cache, fetch=fetch) + assert status is None and ident == "alice" + assert len(calls) == 1 From bc883b409220e0559ec2e33c87fd39267f987e69 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 14:11:28 -0400 Subject: [PATCH 2/4] docs(eventing): add DESIGN_PHASE2 for identity on the event path The agentdocs record design decisions so their claims can be checked against the code. This change had none, so the reasoning behind it lived only in a PR description. Written as a delta over DESIGN_PHASE1, matching the existing phase documents, with an explicit "what does NOT change" section. What it records that is not obvious from the diff: - Why the design looks the way it does. GitHub does not issue a verifiable token for user login - the device flow returns an opaque string - so local validation is impossible and a GET /user lookup is forced. That single fact rules out the JWT-shaped design most readers will reach for, and makes the cache load-bearing rather than an optimisation. - Why 401 and 403 are kept distinct, and why the 403 names the login it refused: a user told only "forbidden" hunts for a broken token. - Why an empty approved list denies everyone rather than everybody. - That trust in EventBridge is load-bearing, since EventRunner never contacts GitHub. Documented as a property rather than left to be discovered in questions. - What ce_submitter is NOT worth while it stays unsigned, and exactly what would make it provable. - Why /continue and PUT /transcript are deliberately open, with the planned per-correlation HMAC fix for the first. - Why Keycloak and HMAC were both rejected for agent identity, including the ~170,000x speed advantage HMAC has and why it still fails the goal. - The blocker on the agent half: sign_event() has no production caller, so ER_REQUIRE_SIGNATURE=true is a kill switch. Verified again while writing this. - One environment trap found during testing: two EventBridge instances share a fixed responses consumer group, so they split partitions and the symptom looks like lost responses rather than a split group. Every file:line and section reference in the document was checked against the tree it describes. Refs: rossoctl/rossoctl#2606, rossoctl/rossoctl#2607 Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/DESIGN_PHASE2.md | 358 ++++++++++++++++++++++++++++ eventing/agentdocs/README.md | 5 + 2 files changed, 363 insertions(+) create mode 100644 eventing/agentdocs/DESIGN_PHASE2.md diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md new file mode 100644 index 00000000..6773ba9c --- /dev/null +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -0,0 +1,358 @@ +# DESIGN — Phase 2: identity on the event path + +Status: draft (revision 1) +Scope: **delta over `DESIGN_PHASE1.md`.** Read that first. This document records +only what changes when the demo stops trusting whoever can reach it. + +Phase 0 proved the wire shape. Phase 1 decided where the consumer runs and how +many of them there are. Phase 2 answers two questions neither phase asked: + +1. **Which person asked for this work?** — and is that person allowed to ask? +2. **Which agent produced this answer?** — and is that agent one we approved? + +```text + 👤 ──sign in──▶ 🌐 GitHub + │ │ GET /user -> login + │ Bearer ▼ + └──────────▶ ╔═══════════════════╗ + ║ 🔒 EventBridge ║ 401 no credential + ║ the PEP ║ 403 known, not approved + ╚═════════╤═════════╝ + │ ce_submitter, ce_submitter_iss + ▼ + Kafka:requests ─▶ EventRunner ─▶ Kafka:responses +``` + +The single new moving part in this half is GitHub: EventBridge verifies a +sign-in once at the edge, then records *who* on the request event. The second +half — proving which agent answered — has its groundwork in place (§4) but is +not yet wired, and this document says so plainly rather than describing it as +done. + +--- + +## 1. What does NOT change + +Stated explicitly, because the temptation in a security document is to redesign +things that already work: + +- **The CloudEvent contract** (Phase 0 §2). Two new extension attributes are + *added* (`submitter`, `submitter_iss`); nothing existing changes shape. + `ce.new_event(**attrs)` already accepts arbitrary attributes and + `to_kafka_binary` already emits every non-empty one as a `ce_*` header, so the + codec needed no change at all. +- **Kafka message key = `correlationid`**, the per-correlation FIFO router, the + `uuid5` session derivation, `--session-id` / `--resume` mechanics. +- **KEDA scaling on consumer lag** and scale-to-zero. Identity is checked at the + HTTP edge, so it never touches the scaling path. +- **Pure-Python discipline** (Phase 1 §1.1). No new runtime dependency: + `urllib`, `hmac`, `hashlib`, `json`. In particular **no JWT library**, because + there is no JWT to verify — see §2.1. +- **Auth is off by default.** With no client id, no approved-user list and no + static tokens configured, `resolve()` returns "allowed, anonymous" and the + Phase 0/1 demo behaves exactly as before. Every existing test passes + untouched; that is the check that this is additive. + +--- + +## 2. User identity + +### 2.1 The constraint that shapes everything + +**GitHub does not issue a verifiable token for user login.** The OAuth device +flow returns an *opaque* access token: no signature, no claims, nothing to check +offline. The JWKS at `token.actions.githubusercontent.com` is for Actions +workloads, not users, and there is no user-facing equivalent. + +This is the fact that rules out the obvious design. We cannot validate a token +locally the way AuthBridge validates a Keycloak JWT. EventBridge must ask GitHub +who holds the token, via `GET https://api.github.com/user`. + +Three consequences, all accepted deliberately: + +| Consequence | Why it is acceptable, and what it costs | +|---|---| +| Sign-in depends on GitHub being reachable | A failed lookup is a `401`, never an allow. Failing closed is the only safe direction for "who is this". A GitHub outage means nobody can submit — correct, and an operator can keep a static break-glass token (§2.5). | +| A lookup per request, against a 5000/hour budget | The cache (§2.3) is therefore **load-bearing, not an optimisation**. | +| GitHub's latency joins the request path | Same answer: the cache. A cache hit costs nothing. | + +### 2.2 Why the device flow + +The alternative — an OAuth web flow with a redirect — needs EventBridge to host +a callback endpoint at a URL GitHub can reach. The demo runs on a laptop behind +a VPN (Phase 1 §"Remote access caveat"), so that is precisely what it cannot do. + +The device flow inverts it: the CLI asks GitHub for a code, the user types the +code into a page GitHub already hosts, and the CLI polls. Nothing needs to reach +*us*. It also needs **no client secret**, which is why the client id is a +committed default rather than a capability like an ntfy topic. + +`login` prints the code and blocks rather than opening a browser. Auto-opening +fails silently over SSH and inside a container, which is where this is most +often run. + +**Scopes requested: none.** Verified against the live API — `GET /user` answers +with `x-accepted-oauth-scopes:` empty, so an unscoped token reads the login. The +demo asks for the least access that answers its question, and cannot read a +repository even if the token leaks. + +GitHub's pacing contract is honoured exactly: `authorization_pending` means keep +polling, `slow_down` means keep polling and add five seconds. Polling faster +than asked rate-limits the whole OAuth App, which would break sign-in for +everyone rather than just the impatient caller. + +### 2.3 The cache + +Token → login, TTL 300 s by default, keyed by **`sha256(token)`**. + +The hash is not decoration. A cache keyed by the raw token means a memory dump, +a careless `repr`, or a debug log yields a working credential. Keyed by hash, it +yields nothing. + +**Failures are not cached.** Two reasons, and they pull the same way: a revoked +token must stop working promptly, and a GitHub outage must not pin a legitimate +user to a failure for the whole TTL. + +### 2.4 `401` and `403` are different answers + +This is the design decision most worth defending, because collapsing them is +the easy thing to do. + +| Status | Meaning | Carries `WWW-Authenticate`? | +|---|---|---| +| `401` | "I do not know you." No credential, a malformed one, or one GitHub does not recognise. | Yes — retrying with a credential is the remedy. | +| `403` | "I know exactly who you are, and you are not approved." | **No** — retrying with another credential is *not* the remedy. | + +The `403` body names the login that was refused (`mrsabath is not on the +approved-user list`). A user who is told only "forbidden" goes looking for a +broken token; a user told *which identity* was refused knows to ask an operator. +That is the difference between an actionable error and a support ticket. + +An empty approved-user list **denies everyone**. The other reading — empty means +everybody — would turn a missing environment variable into an open door, which +is exactly the class of failure a security feature must not have. + +Logins compare case-insensitively, because GitHub logins are. Comparing exactly +would refuse a genuinely approved user over capitalisation, which reads as a +broken deployment rather than a policy. + +### 2.5 Static tokens remain, on purpose + +`EB_AUTH_TOKENS` is not deprecated. It serves three things GitHub sign-in cannot: + +- **Tests must not reach the network.** A suite that calls GitHub is slow, + flaky, rate-limited, and fails on a machine without credentials. +- **An offline demo has to stay possible.** Conference wifi is not a dependency + worth accepting. +- **A break-glass credential.** When GitHub is unreachable, an operator with a + static token can still drive the system. `resolve()` checks it before giving + up for exactly this reason. + +### 2.6 Identity on the event, and what it is worth + +Two attributes ride the request: + +```text +ce_submitter: mrsabath +ce_submitter_iss: github +``` + +`submitter_iss` exists because without it a reader cannot tell a verified GitHub +login from a name an operator typed into an environment variable. Absent issuer +means "static token" — the weaker claim, visible as such. + +**Both are unsigned.** Anything with write access to the `requests` topic can +forge them, and the broker is plaintext. The honest claim after this phase is: + +> A real GitHub user, on an approved list, authorised this request — as recorded +> by EventBridge. + +**Not** "the event proves who submitted it." Making it provable needs `submitter` +inside `signing.SIGNED_ATTRS` *and* a producer that signs. See §4. + +### 2.7 Where the check lives, and the trust that follows + +EventBridge is the **policy enforcement point**. It verifies once, then attests +by recording. EventRunner never contacts GitHub. + +That means **compromising EventBridge means being able to claim any user.** This +is standard PEP design, and the alternative is worse: passing the user's GitHub +token through to every runner would spread a live credential across every +workload and make each one a lookup client. Stated here so it is a documented +property rather than a discovery during questions. + +--- + +## 3. What is deliberately left open + +### 3.1 `/continue` is unauthenticated + +`ntfy.py` emits an ntfy `http` action so a notification has a **Continue…** +button. That action is a recipe serialised into the notification, so any +credential it carries comes to rest in four places outside our control: the +payload sent to `ntfy.sh`, ntfy's message store, the phone's notification +history, and every other subscriber of the topic. + +A long-lived bearer token must not go there. So `/continue` stays open, and the +reasoning is not a shrug: continuing requires already knowing an unguessable +correlationid, which is a **capability URL** — the same model the HTML transcript +already relies on. Creating *new* work is the privileged act. + +**Planned fix**: derive `key = HMAC(server_secret, correlationid + exp)`, bake +`?k=` into the action URL, and accept either a bearer token or a valid key. +A leaked notification then grants one conversation, with an expiry, instead of +the API. Stdlib `hmac`, nothing stored. + +### 3.2 `PUT /transcript` is unauthenticated + +EventRunner uses it for session checkpointing and has no credential concept +anywhere in `eventrunner/config.py`. Gating it breaks `/continue` after a cold +pod — the Phase 1 §16 Gap B path. Fixing it properly means giving EventRunner an +identity, which is §4's territory. + +### 3.3 Not addressed at all + +- **Kafka is plaintext**, with no transport authentication. SASL and ACLs were + evaluated and rejected for this demo: the Kafka authorizer is global rather + than per-listener, so enabling it either denies every existing PLAINTEXT client + or, with `allow.everyone.if.no.acl.found=true`, makes the ACL demo vacuous. + More importantly they authenticate the *connection*, not the payload — they + cannot tell a real EventRunner from anything else holding valid credentials. +- **No rate limiting.** Refusing an invalid request is cheap but not free. +- **The approved lists are files.** They record what an operator approved, not + what a platform attested. + +--- + +## 4. Agent identity: groundwork, not yet 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**. + +### 4.1 What exists + +- `signing.sign_event(event, seed, kid=None)` writes a key id into the JWS + protected header. Because the header is part of the signed input, **the `kid` + cannot be swapped** to relabel an event as coming from another agent. + Omitting it is byte-identical to before, so this was additive — no + canonicalisation break and no signature migration. +- `signing.token_kid(token)` reads the `kid` *before* verification, to choose a + key. It is a hint until verification succeeds with the key it named. +- `shared/keyset.py` maps `kid` → Ed25519 public key from a JSON file. **That + file is the authorization list**: an unknown or absent `kid` has no key and the + event is refused. + +`select(None)` returns a key only when exactly one is approved. With several, +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 + +**`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. + +### 4.3 Why not Keycloak, and why not HMAC + +**Keycloak client credentials** prove an agent *holds a secret* — which a +compromised pod also holds. It adds a token-issuing dependency while proving the +least of the available options. + +**HMAC** is symmetric: EventBridge would hold the key it verifies with, so it +could forge any agent's response and any agent could forge another's. That fails +the goal as stated. It is ~170,000× faster than the hand-rolled Ed25519 +(0.0013 ms/op against ~150 ms), which is tempting and still wrong here. + +**Ed25519** proves possession of a private key that never leaves the runner, and +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`, `submitter_iss` 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. + +--- + +## 5. Configuration + +| Variable | Default | Effect | +|---|---|---| +| `EB_GITHUB_CLIENT_ID` | from `config.toml` | OAuth App client id. **Public** — the device flow has no client secret. | +| `EB_ALLOWED_USERS` | empty | Comma-separated approved logins. Empty denies everyone. | +| `EB_GITHUB_CACHE_TTL_S` | `300` | Token → login cache lifetime. | +| `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. + +--- + +## 6. Verification + +**Tests: 490 passed, 5 skipped.** Baseline on the same tree is 450/5, so this +adds 40 and regresses nothing. No test reaches the network — the device flow and +`GET /user` are exercised through injected fakes. + +Verified end to end against a real GitHub account, not only in unit tests: + +| Check | Result | +|---|---| +| `login` browser flow | `✔ signed in as mrsabath`, token stored `0600` | +| No credential | `401` + `WWW-Authenticate` | +| Unknown token | `401` | +| Real user, not on the list | `403` — `mrsabath is not on the approved-user list` | +| Approved user | `202`, agent ran, `final=True` | +| On the wire | `ce_submitter:mrsabath`, `ce_submitter_iss:github` | +| Group of 5 | every member carried both attributes | + +### 6.1 One environment finding worth recording + +During the group test the completion counter read 1/5 while all five members had +run. Cause: several EventBridge instances were running against one Kafka, and the +**response consumer uses a fixed shared group** (`eventbridge-responses`) while +the requests mirror uses a per-PID one. The instances therefore split the twelve +response partitions between them and no single one saw every response. + +Not a defect in this change, and not a bug in Phase 1 either — a single deployed +EventBridge is the intended topology, and Phase 1 lists multi-replica EventBridge +as deferred. But it is a real trap for anyone running two copies on a laptop: +**symptoms look like lost responses, not like a split consumer group.** Worth +knowing before a demo. + +--- + +## 7. Reading order for whoever picks this up + +- **Running it?** `README_PHASE1.md`, then `EB_GITHUB_CLIENT_ID` and + `EB_ALLOWED_USERS` from §5. +- **Changing the identity model?** §2.1 first — the opaque-token constraint is + what rules out the design most people reach for. +- **Finishing agent identity?** §4.2 for the blocker, then §4.4 in order. +- **Presenting it?** §2.6 and §3 — what the controls do *not* prove. A security + demo that oversells its guarantee is worse than one that does not exist. diff --git a/eventing/agentdocs/README.md b/eventing/agentdocs/README.md index 3eec783a..1cab8b21 100644 --- a/eventing/agentdocs/README.md +++ b/eventing/agentdocs/README.md @@ -18,6 +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. | ## Reading order @@ -27,8 +28,12 @@ and what is *not* verified. - **Changing it?** [`DESIGN_PHASE0.md`](DESIGN_PHASE0.md) for the wire contract that must not change, then [`DESIGN_PHASE1.md`](DESIGN_PHASE1.md) for the deployment model. +- **Working on auth or identity?** [`DESIGN_PHASE2.md`](DESIGN_PHASE2.md) §2.1 first: + the opaque-token constraint is what rules out the design most people reach for. - **Debugging something odd?** [`IMPLEMENTATION_REPORT1.md`](IMPLEMENTATION_REPORT1.md) §4 — most surprises are already recorded there with their cause. +- **Presenting it?** [`DESIGN_PHASE2.md`](DESIGN_PHASE2.md) §2.6 and §3, for what the + identity controls do *not* prove. The `PHASE0` documents describe a laptop demo and remain accurate for that; where Phase 1 supersedes them it says so explicitly rather than editing them in place. From 2c94c9c5bedaedf02a129d7d0769cca225fa265d Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 14:30:38 -0400 Subject: [PATCH 3/4] fix(eventing): CloudEvents-compliant attribute name, safe token write Review fixes. 1. submitter_iss -> submitteriss (must-fix). CloudEvents v1.0 requires attribute names to be lower-case [a-z0-9] only: no underscore. It was the only one of eleven EXT_* constants to break the rule, and the codec here could not catch it - to_kafka_binary/from_kafka_binary only add and strip the ce_ prefix and validate nothing, so a bad name round-trips locally and is rejected or silently dropped by a spec-compliant SDK, an HTTP-binding gateway or a Knative broker a hop later. Renamed while nothing is persisted; later it would be a migration. Added test_roundtrip_binary.py assertions over every EXT_* constant so the next extension cannot repeat it, plus the 20-character SHOULD limit. Verified the test fails when the old name is restored. 2. The token write no longer has a permissive window. write_text creates at the process umask, so there was a moment where a credential was group- and world-readable. Now os.open(..., 0o600) at creation, and mode=0o700 on the directory. Kept the trailing chmod: os.open applies its mode only when it creates the file, so a re-login over a file left loose by an earlier version would otherwise keep 0644 - confirmed by experiment. 3. Removed a dead fallback in auth.resolve. That branch ran only when _bearer had already failed, and resolve_identity re-reads the same header through _bearer, so name was unconditionally None and why2 == why. The break-glass path is the branch below, where a token WAS presented and GitHub could not vouch for it; verified still working after the deletion. DESIGN_PHASE2 updated for the rename, with the naming rule and how it escaped local testing recorded in 2.6 - the point of agentdocs being that the next person does not rediscover it. Not changed, deliberately: LoginCache still evicts only on lookup of the same key, so it grows with distinct valid tokens seen. Bounded by the approved-user count, failures are not cached, and adding a sweep would be unexercised code at demo scale. Noted here rather than silently accepted. Tests: 492 passed, 5 skipped. ruff clean under the pinned 0.11.4. Refs: rossoctl/rossoctl#2606 Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/DESIGN_PHASE2.md | 24 ++++++++++----- eventing/eventbridge/auth.py | 12 ++++---- eventing/shared/ce.py | 8 ++++- .../skills/eventbridge/eventbridge-cli.py | 13 ++++++-- eventing/tests/test_roundtrip_binary.py | 30 +++++++++++++++++++ 5 files changed, 70 insertions(+), 17 deletions(-) diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md index 6773ba9c..2444ba54 100644 --- a/eventing/agentdocs/DESIGN_PHASE2.md +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -18,7 +18,7 @@ many of them there are. Phase 2 answers two questions neither phase asked: ║ 🔒 EventBridge ║ 401 no credential ║ the PEP ║ 403 known, not approved ╚═════════╤═════════╝ - │ ce_submitter, ce_submitter_iss + │ ce_submitter, ce_submitteriss ▼ Kafka:requests ─▶ EventRunner ─▶ Kafka:responses ``` @@ -37,7 +37,7 @@ Stated explicitly, because the temptation in a security document is to redesign things that already work: - **The CloudEvent contract** (Phase 0 §2). Two new extension attributes are - *added* (`submitter`, `submitter_iss`); nothing existing changes shape. + *added* (`submitter`, `submitteriss`); nothing existing changes shape. `ce.new_event(**attrs)` already accepts arbitrary attributes and `to_kafka_binary` already emits every non-empty one as a `ce_*` header, so the codec needed no change at all. @@ -153,14 +153,24 @@ broken deployment rather than a policy. Two attributes ride the request: ```text -ce_submitter: mrsabath -ce_submitter_iss: github +ce_submitter: mrsabath +ce_submitteriss: github ``` -`submitter_iss` exists because without it a reader cannot tell a verified GitHub +`submitteriss` exists because without it a reader cannot tell a verified GitHub login from a name an operator typed into an environment variable. Absent issuer means "static token" — the weaker claim, visible as such. +**On the spelling.** CloudEvents v1.0 requires attribute names to be lower-case +`[a-z0-9]` only — no underscore, hyphen or upper case — because an event crosses +several hops and protocols disagree about metadata case-sensitivity. This first +shipped as `submitter_iss` and was caught in review, not by the code: the codec +here only adds and strips the `ce_` prefix, so a non-compliant name round-trips +locally and is rejected or silently dropped by a spec-compliant SDK, an +HTTP-binding gateway or a Knative broker further along. `test_roundtrip_binary.py` +now asserts the rule over every `EXT_*` constant, so the next extension cannot +repeat it. + **Both are unsigned.** Anything with write access to the `requests` topic can forge them, and the broker is plaintext. The honest claim after this phase is: @@ -286,7 +296,7 @@ from. 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`, `submitter_iss` and the missing `groupid` to +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. @@ -328,7 +338,7 @@ Verified end to end against a real GitHub account, not only in unit tests: | Unknown token | `401` | | Real user, not on the list | `403` — `mrsabath is not on the approved-user list` | | Approved user | `202`, agent ran, `final=True` | -| On the wire | `ce_submitter:mrsabath`, `ce_submitter_iss:github` | +| On the wire | `ce_submitter:mrsabath`, `ce_submitteriss:github` | | Group of 5 | every member carried both attributes | ### 6.1 One environment finding worth recording diff --git a/eventing/eventbridge/auth.py b/eventing/eventbridge/auth.py index c535ae1d..addd3f2c 100644 --- a/eventing/eventbridge/auth.py +++ b/eventing/eventbridge/auth.py @@ -147,12 +147,12 @@ def resolve(environ: dict[str, Any], cfg, *, cache=None, if github_on: presented, why = _bearer(environ) if why: - # Fall back to a static token only when one could match; otherwise a - # missing header is simply unauthenticated. - if not cfg.auth_tokens: - return None, None, 401, why - name, why2 = resolve_identity(environ, cfg.auth_tokens) - return (name, None, None, None) if name else (None, None, 401, why2) + # No usable credential at all. Falling back to a static token here + # would be dead code: `resolve_identity` re-reads the same header via + # `_bearer` and fails for the same reason. The break-glass path is the + # branch below, which is the case that matters — a token WAS presented + # and GitHub could not vouch for it. + return None, None, 401, why login, err = ghauth.resolve( presented, cache, **({"fetch": fetch} if fetch else {})) diff --git a/eventing/shared/ce.py b/eventing/shared/ce.py index e17d7a6c..6bc77427 100644 --- a/eventing/shared/ce.py +++ b/eventing/shared/ce.py @@ -49,7 +49,13 @@ # 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 # cannot tell a real identity from a name an operator typed into an env var. -EXT_SUBMITTER_ISS = "submitter_iss" +# +# Spelled without a separator because CloudEvents v1.0 requires attribute names +# to be lower-case [a-z0-9] only — no underscore. This codec would not have +# caught `submitter_iss`: to_kafka_binary/from_kafka_binary just add and strip the +# `ce_` prefix and validate nothing, so it round-trips here and is rejected or +# silently dropped by a spec-compliant consumer a hop later. +EXT_SUBMITTER_ISS = "submitteriss" CE_HEADER_PREFIX = "ce_" CORE_ATTRS = {"specversion", "type", "source", "id", "time", diff --git a/eventing/skills/eventbridge/eventbridge-cli.py b/eventing/skills/eventbridge/eventbridge-cli.py index e3d21a81..433bb794 100644 --- a/eventing/skills/eventbridge/eventbridge-cli.py +++ b/eventing/skills/eventbridge/eventbridge-cli.py @@ -539,9 +539,16 @@ def cmd_login(args): raise SystemExit(1) from None dest = _token_path() - dest.parent.mkdir(parents=True, exist_ok=True) - dest.write_text(token) - dest.chmod(0o600) # a token is a credential; not world-readable + # 0o700 on the directory and 0o600 at creation, rather than writing first and + # tightening after: `write_text` creates at the process umask, which leaves a + # window — however short — where a credential is group- and world-readable. + dest.parent.mkdir(parents=True, exist_ok=True, mode=0o700) + fd = os.open(dest, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + with os.fdopen(fd, "w") as f: + f.write(token) + # os.open applies its mode only when it CREATES the file, so a re-login over + # a file left loose by an earlier version would keep the old permissions. + dest.chmod(0o600) print(f"✔ signed in as {login}") print(f" token stored at {dest} (mode 600)") diff --git a/eventing/tests/test_roundtrip_binary.py b/eventing/tests/test_roundtrip_binary.py index 8e7429f2..abeca5fd 100644 --- a/eventing/tests/test_roundtrip_binary.py +++ b/eventing/tests/test_roundtrip_binary.py @@ -1,5 +1,6 @@ """Wire contract: CloudEvent binary-mode round-trip through headers+value.""" import json +import re import uuid from shared import ce @@ -78,3 +79,32 @@ def test_empty_value_decodes_to_none_data(): assert value == b"" evt2 = ce.from_kafka_binary(headers, value) assert evt2.data is None + + +def test_every_extension_attribute_name_is_cloudevents_compliant(): + """CloudEvents v1.0: attribute names MUST be lower-case [a-z0-9] only. + + No underscore, hyphen, dot or upper case. The spec restricts the set because + an event traverses several hops and some protocols treat metadata as + case-sensitive while others do not. + + This codec cannot catch a violation on its own: `to_kafka_binary` and + `from_kafka_binary` only add and strip the `ce_` prefix, so a bad name + round-trips happily here and is rejected — or silently dropped — by a + spec-compliant SDK, an HTTP-binding gateway or a Knative broker further + along. `submitter_iss` shipped in review for exactly that reason. Hence a + test rather than a comment. + """ + ext_names = [v for k, v in vars(ce).items() + if k.startswith("EXT_") and isinstance(v, str)] + assert ext_names, "no EXT_* constants found — has ce.py been restructured?" + bad = [n for n in ext_names if not re.fullmatch(r"[a-z][a-z0-9]*", n)] + assert not bad, f"not CloudEvents-compliant attribute names: {bad}" + + +def test_extension_names_are_terse_enough_to_survive_a_gateway(): + """The spec SHOULD-limit is 20 characters. Not a hard failure upstream, but + a name over it is a smell worth catching while renaming is still free.""" + long = [v for k, v in vars(ce).items() + if k.startswith("EXT_") and isinstance(v, str) and len(v) > 20] + assert not long, f"extension names over the 20-char SHOULD limit: {long}" From 304925d124a3ab6fe667f0ef93c9d5cd6c199731 Mon Sep 17 00:00:00 2001 From: Mariusz Sabath Date: Tue, 29 Sep 2026 17:35:59 -0400 Subject: [PATCH 4/4] fix(eventing): correct the revocation claim, check static tokens first Addresses @aslom's three non-blocking review findings on #878. 1. The revocation claim was wrong, in three places. Not caching failures does nothing for revocation: a cache hit short-circuits resolve() before fetch_login runs, so a revoked token keeps authenticating until its POSITIVE entry expires. Reproduced the measurement - accepted at t+299s, refused at t+300s. The window is the full TTL. Corrected ghauth.py, test_ghauth.py and DESIGN_PHASE2 2.3, and added test_a_revoked_token_keeps_working_until_its_positive_entry_expires plus a shorter-TTL case so the claim cannot drift again. The second half of the old reasoning - that an outage must not pin a legitimate user to a refusal - is correct and is what not-caching-failures actually buys; it now says only that. 2. Static tokens are now checked BEFORE GitHub. The old order made break-glass slowest exactly when it was needed: during an outage every static-token request paid a full fetch_login timeout on a call that could never succeed, against the same budget 2.3 says the cache protects. Static is local and constant-time, so it goes first. Nothing is shadowed - a static secret would have to deliberately collide with a live gho_-shaped token. Two tests: one asserting no GitHub call happens for a static token, one asserting real sign-in still resolves. 3. resolve()'s docstring said "and" where github_on is an or. Fixed, and documented what each half-configured state actually does, since they differ: client-id-only 403s everyone, list-only accepts any PAT from an approved login with no OAuth App involved. Tests: 496 passed, 5 skipped (was 492/5). ruff check and format both clean. Refs: rossoctl/rossoctl#2606 Assisted-By: Claude (Anthropic AI) Signed-off-by: Mariusz Sabath --- eventing/agentdocs/DESIGN_PHASE2.md | 27 ++++++++--- eventing/eventbridge/auth.py | 43 +++++++++++++----- eventing/eventbridge/ghauth.py | 12 +++-- eventing/tests/test_ghauth.py | 69 ++++++++++++++++++++++++++++- 4 files changed, 129 insertions(+), 22 deletions(-) diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md index 2444ba54..28379122 100644 --- a/eventing/agentdocs/DESIGN_PHASE2.md +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -1,6 +1,6 @@ # DESIGN — Phase 2: identity on the event path -Status: draft (revision 1) +Status: draft (revision 2) Scope: **delta over `DESIGN_PHASE1.md`.** Read that first. This document records only what changes when the demo stops trusting whoever can reach it. @@ -109,9 +109,20 @@ The hash is not decoration. A cache keyed by the raw token means a memory dump, a careless `repr`, or a debug log yields a working credential. Keyed by hash, it yields nothing. -**Failures are not cached.** Two reasons, and they pull the same way: a revoked -token must stop working promptly, and a GitHub outage must not pin a legitimate -user to a failure for the whole TTL. +**Failures are not cached**, so a GitHub outage cannot pin a legitimate user to a +refusal for the whole TTL — each request retries. + +**It does not speed up revocation**, and an earlier revision of this document +claimed it did. A cache hit short-circuits `resolve()` before `fetch_login` runs, +so a token revoked on GitHub keeps authenticating until its *positive* entry +expires. Measured with an injected clock: still accepted at t+299 s, refused at +t+300 s. The window is the full TTL. + +That is bounded, configurable via `EB_GITHUB_CACHE_TTL_S`, and 300 s is a +defensible trade against a 5000/hour budget — but it is a real window, and +`test_a_revoked_token_keeps_working_until_its_positive_entry_expires` now pins it +so the claim cannot drift again. Lower the TTL if prompt revocation matters more +than API calls. ### 2.4 `401` and `403` are different answers @@ -145,8 +156,12 @@ broken deployment rather than a policy. - **An offline demo has to stay possible.** Conference wifi is not a dependency worth accepting. - **A break-glass credential.** When GitHub is unreachable, an operator with a - static token can still drive the system. `resolve()` checks it before giving - up for exactly this reason. + static token can still drive the system. `resolve()` checks the static map + **before** GitHub, precisely so the fallback is fastest when it is needed: the + reverse order made every break-glass request pay a full `fetch_login` timeout + on a call that could never succeed, against the budget §2.3 says the cache + exists to protect. Nothing is shadowed — a static secret would have to + deliberately collide with a live `gho_`-shaped token. ### 2.6 Identity on the event, and what it is worth diff --git a/eventing/eventbridge/auth.py b/eventing/eventbridge/auth.py index addd3f2c..3d13e6c2 100644 --- a/eventing/eventbridge/auth.py +++ b/eventing/eventbridge/auth.py @@ -129,12 +129,25 @@ def resolve(environ: dict[str, Any], cfg, *, cache=None, Collapsing those into one answer would tell an operator less, and would tell a user debugging their own access much less. - Two mechanisms, checked in order: + GitHub sign-in is active when a client id **or** an approved-user list is set + — deliberately an `or`, so a half-configured deployment fails closed rather + than silently falling back to no authentication. The two halves behave + differently, and it is worth knowing which mistake you have made: - 1. **GitHub sign-in**, when a client id and an approved-user list are - configured. The real path. - 2. **Static tokens** (`EB_AUTH_TOKENS`), as a fallback. Tests must not reach - the network, and an offline demo has to stay possible, so this stays. + * **client id only** — every request gets `403`, because the approved list is + empty and an empty list approves nobody. + * **approved list only** — any GitHub token from an approved login is + accepted, with no OAuth App involved. Useful for a quick test with a + personal access token; not what you want in a deployment. + + Within that, checked in order: + + 1. **Static tokens** (`EB_AUTH_TOKENS`) — a local constant-time compare, and + the break-glass path when GitHub is unreachable, so it goes first. + 2. **GitHub sign-in** — the real path. + + Static tokens also keep tests off the network and an offline demo possible, + which is why they are not removed now that sign-in exists. With neither configured the result is all-`None` — allowed and anonymous, which is what keeps the default demo working out of the box. @@ -154,15 +167,23 @@ def resolve(environ: dict[str, Any], cfg, *, cache=None, # and GitHub could not vouch for it. return None, None, 401, why + # Static tokens first, because they are the break-glass path and this is a + # local constant-time compare. Checking GitHub first made the fallback + # slowest exactly when it is needed: during an outage every break-glass + # request paid a full `fetch_login` timeout on a call that could never + # succeed, against the same API budget the cache exists to protect. + # + # Nothing is shadowed by the order. A GitHub token is `gho_`/`ghp_`-shaped + # and an operator choosing a static secret that collides with a live + # GitHub token would have to do so deliberately. + if cfg.auth_tokens: + name, _ = resolve_identity(environ, cfg.auth_tokens) + if name: + return name, None, None, None + login, err = ghauth.resolve( presented, cache, **({"fetch": fetch} if fetch else {})) if login is None: - # A static token is checked before giving up, so an operator can keep - # a break-glass credential alongside GitHub sign-in. - if cfg.auth_tokens: - name, _ = resolve_identity(environ, cfg.auth_tokens) - if name: - return name, None, None, None return None, None, 401, err or "could not identify this token" if not ghauth.is_allowed(login, cfg.allowed_users): diff --git a/eventing/eventbridge/ghauth.py b/eventing/eventbridge/ghauth.py index b69f2912..01ccee40 100644 --- a/eventing/eventbridge/ghauth.py +++ b/eventing/eventbridge/ghauth.py @@ -171,9 +171,15 @@ class LoginCache: hourly API calls, and adds GitHub's latency to the request path. Entries are keyed by hash so the cache cannot be read back into a working credential. - Negative results are not cached. A revoked token must stop working promptly, - and a GitHub outage must not pin a legitimate user to a failure for the whole - TTL. + Negative results are not cached, so a GitHub outage cannot pin a legitimate + user to a failure for the whole TTL — each request retries. + + That does **not** speed up revocation. A revoked token keeps working until its + positive entry expires, because a cache hit short-circuits `resolve()` before + `fetch_login` is ever called: the window is the full TTL, 300 s by default. + Bounded and configurable via `EB_GITHUB_CACHE_TTL_S`, and 300 s is a + defensible trade against spending the 5000/hour API budget — but it is a real + window, not an absence of one. Lower the TTL if that matters more than calls. """ def __init__(self, ttl_s: float = 300.0, now=time.monotonic) -> None: diff --git a/eventing/tests/test_ghauth.py b/eventing/tests/test_ghauth.py index c540b3cc..bf9f7c9b 100644 --- a/eventing/tests/test_ghauth.py +++ b/eventing/tests/test_ghauth.py @@ -111,13 +111,51 @@ def fetch(token): def test_resolve_does_not_cache_a_failure(): - """A revoked token must stop working promptly, and a GitHub outage must not - pin a legitimate user to a failure for the whole TTL.""" + """So a GitHub outage cannot pin a legitimate user to a failure for the whole + TTL — each request retries rather than replaying a cached refusal. + + Note this does NOT speed up revocation: see + `test_a_revoked_token_keeps_working_until_its_positive_entry_expires`.""" c = ghauth.LoginCache(300.0) ghauth.resolve("tok", c, fetch=lambda t: (None, "GitHub rejected this token")) assert len(c) == 0 +def test_a_revoked_token_keeps_working_until_its_positive_entry_expires(): + """The revocation window is the POSITIVE TTL, and it is worth pinning. + + A cache hit short-circuits `resolve()` before `fetch_login` runs, so a token + revoked on GitHub keeps authenticating until its entry expires. Not caching + failures does nothing for this — that only stops an outage pinning a + legitimate user to a refusal. + + Bounded and configurable, and 300 s is a defensible trade against the + 5000/hour budget. But it is a real window, and a design doc claiming + otherwise is worse than one that states it. + """ + clock = FakeClock() + cache = ghauth.LoginCache(300.0, now=clock) + assert ghauth.resolve("tok", cache, fetch=lambda t: ("alice", None)) == ("alice", None) + + revoked = lambda t: (None, "GitHub rejected this token") # noqa: E731 + assert ghauth.resolve("tok", cache, fetch=revoked)[0] == "alice", "cache hit" + clock.advance(299) + assert ghauth.resolve("tok", cache, fetch=revoked)[0] == "alice", "still inside the TTL" + clock.advance(2) + assert ghauth.resolve("tok", cache, fetch=revoked)[0] is None, "TTL expired -> refused" + + +def test_a_shorter_ttl_shortens_the_revocation_window(): + """The knob that actually controls it, so `EB_GITHUB_CACHE_TTL_S` is not + mistaken for a pure performance setting.""" + clock = FakeClock() + cache = ghauth.LoginCache(5.0, now=clock) + ghauth.resolve("tok", cache, fetch=lambda t: ("alice", None)) + clock.advance(6) + assert ghauth.resolve("tok", cache, + fetch=lambda t: (None, "GitHub rejected this token"))[0] is None + + def test_resolve_works_without_a_cache(): assert ghauth.resolve("tok", None, fetch=lambda t: ("bob", None)) == ("bob", None) @@ -340,6 +378,33 @@ def test_resolve_is_case_insensitive_about_the_approved_login(): assert status is None and ident == "MrSabath" +def test_break_glass_does_not_pay_a_doomed_github_call(): + """Static tokens are checked BEFORE GitHub, so the fallback is fastest exactly + when it is needed. The previous order meant every break-glass request during + an outage paid a full fetch_login timeout on a call that could never succeed, + against the same API budget the cache exists to protect.""" + calls = [] + + def fetch(token): + calls.append(token) + return None, "cannot reach GitHub: OSError" + + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"alice"}), + auth_tokens={"emergency": "operator"}) + ident, iss, status, _ = auth.resolve(_env("emergency"), cfg, fetch=fetch) + assert (ident, iss, status) == ("operator", None, None) + assert calls == [], "a static token must not trigger a GitHub lookup" + + +def test_a_github_token_is_unaffected_by_the_static_first_order(): + """The reorder must not shadow real sign-in.""" + cfg = _cfg(github_client_id="cid", allowed_users=frozenset({"mrsabath"}), + auth_tokens={"emergency": "operator"}) + ident, iss, status, _ = auth.resolve( + _env("gho_real"), cfg, fetch=lambda t: ("mrsabath", None)) + assert (ident, iss, status) == ("mrsabath", "github", None) + + def test_a_static_token_still_works_as_a_break_glass_credential(): """An operator keeping one alongside GitHub sign-in should not be locked out when GitHub is unreachable."""