diff --git a/eventing/agentdocs/DESIGN_PHASE2.md b/eventing/agentdocs/DESIGN_PHASE2.md new file mode 100644 index 00000000..28379122 --- /dev/null +++ b/eventing/agentdocs/DESIGN_PHASE2.md @@ -0,0 +1,383 @@ +# DESIGN β€” Phase 2: identity on the event path + +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. + +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_submitteriss + β–Ό + 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`, `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. +- **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**, 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 + +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 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 + +Two attributes ride the request: + +```text +ce_submitter: mrsabath +ce_submitteriss: 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: + +> 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`, `submitteriss` and the missing `groupid` to + `SIGNED_ATTRS`. `DESIGN_PHASE1.md` Β§21.9.9 requires `groupid`; its absence + means signatures currently say nothing about batch membership. + +On failure, set `phase="error"`. That reuses machinery already wired: ntfy +**priority 5** with an error tag, a visually distinct bubble in the live SSE +transcript, and `raw_json` persistence for audit. The demo artifact costs +nothing extra. + +--- + +## 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_submitteriss: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. diff --git a/eventing/eventbridge/auth.py b/eventing/eventbridge/auth.py index abec34e4..3d13e6c2 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,92 @@ 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. + + 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: + + * **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. + """ + 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: + # 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 + + # 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: + 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..01ccee40 --- /dev/null +++ b/eventing/eventbridge/ghauth.py @@ -0,0 +1,240 @@ +"""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, 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: + 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..6bc77427 100644 --- a/eventing/shared/ce.py +++ b/eventing/shared/ce.py @@ -46,6 +46,16 @@ # 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. +# +# 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 4a164268..433bb794 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,94 @@ 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() + # 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)") + + +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 +641,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..bf9f7c9b --- /dev/null +++ b/eventing/tests/test_ghauth.py @@ -0,0 +1,450 @@ +"""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(): + """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) + + +# ---- 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_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.""" + 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 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}"