From 17fdc212c2872c676432cbac3540d22b9608ba4f Mon Sep 17 00:00:00 2001 From: John Coffey Date: Fri, 14 Aug 2026 15:21:55 -0700 Subject: [PATCH] Give ingest a real tenant identity (write-routing deferred, disclosed) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ingest tenant-awareness was named "undesigned, not just unbuilt" across CLAUDE.md/threat-model.md/the runbook since early Phase 4 -- the last major standing gap. Scoping was agreed via AskUserQuestion: a config-supplied tenant_id + shared-secret token ingest validates (smaller real implementation, no new PKI), over per-tenant mTLS certs. This change builds that identity mechanism end to end and attaches it to every record at the point it enters the system; it deliberately does NOT build per-tenant write-routing for ClickHouse or Tantivy -- that's real, separately-scoped follow-up work, disclosed explicitly everywhere this was previously called undesigned, not silently left half-done. New pieces: - metadata/migrations/0034 + enterprise/internal/rbacstore/ ingest_credentials.go: a per-tenant bearer credential, only its SHA-256 hash ever persisted (same reasoning a password gets hashed, not stored raw) -- CreateIngestCredential returns the plaintext exactly once, ValidateIngestCredential/RevokeIngestCredential/ ListIngestCredentialsForTenant round it out. - enterprise-auth gains -create-ingest-credential-tenant/ -list-ingest-credentials-tenant/-revoke-ingest-credential (same offline-operator-flag shape as every other credential-minting flag in this binary) and a new POST /internal/authorize-ingest endpoint (internal/authhandler) validating a presented token and resolving its tenant -- a genuinely different credential type from session-backed /internal/authorize, so it doesn't touch session.Manager at all. - ingest (AGPL core) gains an optional TenantResolver (internal/grpcserver, nil by default) and its HTTP client implementation (internal/tenantresolver.HTTPResolver) -- a plain HTTP call to enterprise-auth's new endpoint, never an enterprise/ import, same "network boundary, not import boundary" shape api/authz.HTTPAuthorizer already uses for the query path. PushBatch now requires an `authorization: Bearer ` gRPC metadata entry once a resolver is configured, fails the whole batch closed on a missing/invalid credential (never falls back to "no tenant"), and attaches the resolved tenant ID to every record as a `tenant_id` Kafka message header before producing it. Verified with real round trips at every layer, no Docker needed: rbacstore's credential CRUD (skip-gated on live Postgres, same as every other rbacstore integration test this phase), authhandler's new endpoint (real HTTP via httptest, including the regression test that a session token must not validate as an ingest credential), tenantresolver (real HTTP client against httptest, same pattern as authz.HTTPAuthorizer's own tests), and grpcserver's PushBatch (fake resolver/producer -- no resolver leaves messages unchanged, a configured resolver attaches the right header or fails closed on a bad/missing token). Helm: ingest.requireTenantCredential (default false) is a deliberate, separate opt-in from enterprise.enabled -- turning ENTERPRISE_AUTH_URL on for ingest requires every agent to already hold a credential or be refused outright, so it must not default on just because enterprise.enabled does (same reasoning api.yaml's ENTERPRISE_AUTH_URL isn't tied to enterprise.enabled directly either). docker-compose.yml leaves it unset, same as ever. Docs updated everywhere this was called "undesigned": CLAUDE.md, docs/architecture.md, docs/security/threat-model.md (including its summary table, now split into "identity: built" vs "write-routing: not yet"), docs/phase-4-runbook.md (new §13), enterprise/README.md. --- CLAUDE.md | 32 +++-- deploy/helm/sentry/templates/ingest.yaml | 18 +++ deploy/helm/sentry/values.yaml | 7 + docker-compose.yml | 10 ++ docs/architecture.md | 29 ++-- docs/phase-4-runbook.md | 73 +++++++++- docs/security/threat-model.md | 46 ++++-- enterprise/README.md | 73 +++++++--- enterprise/cmd/enterprise-auth/main.go | 68 ++++++++- .../internal/authhandler/authhandler.go | 58 +++++++- .../internal/authhandler/authhandler_test.go | 104 +++++++++++++- .../internal/rbacstore/ingest_credentials.go | 110 +++++++++++++++ .../internal/rbacstore/rbacstore_test.go | 119 ++++++++++++++++ ingest/cmd/ingest/main.go | 14 +- ingest/internal/config/config.go | 7 + ingest/internal/config/config_test.go | 7 + ingest/internal/grpcserver/server.go | 95 ++++++++++++- ingest/internal/grpcserver/server_test.go | 131 +++++++++++++++++- .../internal/tenantresolver/tenantresolver.go | 65 +++++++++ .../tenantresolver/tenantresolver_test.go | 56 ++++++++ .../0034_create_ingest_credentials.sql | 17 +++ 21 files changed, 1071 insertions(+), 68 deletions(-) create mode 100644 enterprise/internal/rbacstore/ingest_credentials.go create mode 100644 ingest/internal/tenantresolver/tenantresolver.go create mode 100644 ingest/internal/tenantresolver/tenantresolver_test.go create mode 100644 metadata/migrations/0034_create_ingest_credentials.sql diff --git a/CLAUDE.md b/CLAUDE.md index 9a5e06c..3585e5c 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -228,14 +228,30 @@ than one `tenant_memberships` row now gets a real `GET pending-login token, distinct from a real session by both Go type and JWT claim name — a real token-confusion bug this design's own tests caught before it shipped) instead of the flat refusal Phase 4 shipped -with earlier. What still keeps this phase from being done: the actual -tenant-picker *page* doesn't exist (`web` has no session/cookie-handling -code at all yet, and `enterprise-auth` has no CORS middleware for a -cross-origin `fetch` with credentials — both real, separately-scoped -frontend gaps), and ingest itself has no tenant concept for either -storage engine (every record lands in the one shared ClickHouse database -and Tantivy index no matter what — undesigned, not just unbuilt). Full -accounting: +with earlier. Ingest tenant-awareness — the gap this section used to +call "undesigned" — now has a real, if intentionally partial, design: +`ingest` (AGPL core) gained an optional `TenantResolver` +(`ingest/internal/grpcserver`), a per-tenant bearer credential an agent +presents (minted via `enterprise-auth +-create-ingest-credential-tenant=`, validated over the network via a +new `POST /internal/authorize-ingest` endpoint — never an `enterprise/` +import, same boundary shape as `api/authz.Authorizer`), and the +resolved tenant ID is attached to every record as a `tenant_id` Kafka +message header before it's produced. **What's still deferred, clearly**: +nothing downstream reads that header yet — neither `ingest`'s own +ClickHouse writer nor `search`'s independent Redpanda consumer route a +record's write into a per-tenant destination, so every record still +lands in the one shared ClickHouse database/Tantivy index regardless of +which tenant it's now correctly tagged with. That write-routing split +(likely another "second binary," mirroring `enterprise-api`) is real, +scoped, remaining work — attaching a verified tenant identity as early +as possible was deliberately built as a self-contained first step, not +the whole feature. What still keeps this phase from being done: the +actual tenant-picker *page* doesn't exist (`web` has no session/cookie- +handling code at all yet, and `enterprise-auth` has no CORS middleware +for a cross-origin `fetch` with credentials — both real, separately- +scoped frontend gaps), and per-tenant write-routing for ingest per the +above. Full accounting: `/docs/security/threat-model.md`; step-by-step verification procedure (not yet run against a live cluster in this environment): `/docs/phase-4-runbook.md`. The rest of this section describes the exit diff --git a/deploy/helm/sentry/templates/ingest.yaml b/deploy/helm/sentry/templates/ingest.yaml index f0b8cf3..a567683 100644 --- a/deploy/helm/sentry/templates/ingest.yaml +++ b/deploy/helm/sentry/templates/ingest.yaml @@ -32,6 +32,24 @@ spec: secretKeyRef: name: {{ .Release.Name }}-clickhouse key: password + {{- if and .Values.enterprise.enabled .Values.ingest.requireTenantCredential }} + # Enables ingest/internal/grpcserver.TenantResolver. + # Deliberately its OWN opt-in, not folded into + # enterprise.enabled directly (same reasoning + # api.yaml/enterprise-api.yaml's ENTERPRISE_AUTH_URL isn't + # set just because enterprise.enabled is true -- see that + # env var's own comment there): turning this on requires + # every agent to already present a valid `Authorization: + # Bearer ` (minted via `enterprise-auth + # -create-ingest-credential-tenant=`) or be refused + # outright, which would silently break ingest for any + # not-yet-reconfigured agent if it defaulted on alongside + # enterprise.enabled. Off (the default) leaves every record + # without a tenant_id header, same as every Phase 0-3 + # deployment. + - name: ENTERPRISE_AUTH_URL + value: "http://{{ .Release.Name }}-enterprise-auth:8082" + {{- end }} ports: - name: grpc containerPort: 4317 diff --git a/deploy/helm/sentry/values.yaml b/deploy/helm/sentry/values.yaml index 19212d2..dbf5146 100644 --- a/deploy/helm/sentry/values.yaml +++ b/deploy/helm/sentry/values.yaml @@ -71,6 +71,13 @@ ingest: # "boring, well-understood" preference as everywhere else in this # repo -- use cert-manager or an equivalent, don't hand-roll it here). tlsSecretName: "" + # Only meaningful when enterprise.enabled is also true -- see + # templates/ingest.yaml's ENTERPRISE_AUTH_URL comment for why this is + # its own deliberate opt-in, not folded into enterprise.enabled + # directly: turning it on requires every agent to already present a + # valid ingest credential (`enterprise-auth + # -create-ingest-credential-tenant=`) or be refused outright. + requireTenantCredential: false search: image: diff --git a/docker-compose.yml b/docker-compose.yml index 7d4b469..828deb1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -156,6 +156,16 @@ services: # TLS_*_FILE env vars are left at their defaults # (/etc/sentry-ingest/{server,server-key,ca}.pem) — matches where # the volume below mounts the generated dev certs. + # + # ENTERPRISE_AUTH_URL is deliberately NOT set here (see + # ingest/internal/grpcserver's TenantResolver): with it unset, + # PushBatch attaches no tenant_id header to any record, matching + # every Phase 0-3 deployment's behavior. Setting it to + # "http://enterprise-auth:8082" would require every agent to + # present a valid `Authorization: Bearer ` (minted + # via `enterprise-auth -create-ingest-credential-tenant=`) or + # be refused outright -- not turned on here since nothing in this + # compose file provisions one. volumes: - ./hack/dev-certs/out:/etc/sentry-ingest:ro diff --git a/docs/architecture.md b/docs/architecture.md index 75f0667..c61d583 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -139,13 +139,23 @@ escape hatch is opaque to any compiler-injected filter. ran in the environment it was built in — Tantivy is an embedded library, so the cross-tenant isolation probe needed no live database or Docker to execute for real, and it passed. -- Neither storage engine's isolation extends to *ingest*: every record - `ingest` produces lands in the one shared ClickHouse database and the - one shared (default) Tantivy index regardless of tenant. A - newly-provisioned tenant's database/index are real and isolated at - query time — and permanently empty until something upstream of - `chrunner`/`searchclient` becomes tenant-aware on the write side, - which is undesigned, not merely unbuilt. +- **Ingest identity is now built, though write-routing isn't.** `ingest` + (AGPL core) gained an optional `TenantResolver` + (`ingest/internal/grpcserver`): an agent presents a per-tenant bearer + credential (`enterprise-auth -create-ingest-credential-tenant=` + mints one, only its hash stored), validated over the network via a new + `POST /internal/authorize-ingest` endpoint (never an `enterprise/` + import — same "network boundary, not import boundary" shape + `api/authz.Authorizer` already uses), and the resolved tenant ID rides + as a `tenant_id` Kafka message header on every record produced. What + isn't built yet: neither `ingest`'s own ClickHouse writer nor + `search`'s independent Redpanda consumer reads that header back to + route the write anywhere per-tenant — every record still lands in the + one shared ClickHouse database and Tantivy index regardless of tenant, + correctly tagged but not yet isolated at write time. That per-tenant + write-routing split is real, scoped remaining work (likely another + "second binary," mirroring `enterprise-api`), not something this + change claims to have closed. - `deploy/operator`'s `Tenant` CRD and `enterprise-api -provision-tenant` are now unified, deliberately lightweight: `-provision-tenant` stays the sole real actor (ClickHouse + `rbacstore`), and now also syncs its @@ -167,8 +177,9 @@ plain `api`), sharing a host-port/network-alias trick so `alerting`/ `web` need no conditional config either way. With both storage engines' connection/index-layer mechanisms built, deployment topology enforced at both the Helm and docker-compose layers, and the two provisioning -mechanisms unified, the largest remaining gap is ingest's lack of -tenant-awareness, which is undesigned, not merely unbuilt. +mechanisms unified, the largest remaining gap is ingest's per-tenant +*write-routing* (identity is now attached at ingest time; nothing +downstream of Redpanda consumes it yet to isolate the write, see above). ## Licensing boundary diff --git a/docs/phase-4-runbook.md b/docs/phase-4-runbook.md index fe7d2b2..9a62f8b 100644 --- a/docs/phase-4-runbook.md +++ b/docs/phase-4-runbook.md @@ -554,6 +554,63 @@ all -- a cross-origin `fetch` with credentials from `web`'s origin to actual picker UI is real, separately-scoped frontend work; this section only closes the backend half. +## 13. Ingest tenant identity (no per-tenant write-routing yet) + +The identity mechanism was chosen deliberately (config-supplied +tenant_id + a shared-secret token ingest validates, not per-tenant +mTLS certs -- smaller real implementation, no new PKI). Verified in +this environment without Docker or a live enterprise-auth, using the +same fake-client-at-every-layer discipline as everything else in this +runbook that doesn't need a live stack: + +```sh +cd enterprise +go test ./internal/rbacstore/... -run IngestCredential -v +# skip-gated (RBACSTORE_TEST_POSTGRES_ADDR) -- CreateIngestCredential/ +# ValidateIngestCredential/RevokeIngestCredential round trip, and the +# regression test that only a SHA-256 hash is ever persisted, never the +# plaintext token. + +go test ./internal/authhandler/... -run AuthorizeIngest -v +# real HTTP round trip against POST /internal/authorize-ingest with a +# fake credential validator -- proves a session token (service or +# human) does NOT work as an ingest credential, since this endpoint +# never calls session.Manager.Validate at all. + +cd ../ingest +go test ./internal/tenantresolver/... -v +# real HTTP round trip (httptest), same shape as api/authz. +# HTTPAuthorizer's own tests -- forwards the bearer token, parses +# tenant_id, treats a non-2xx or an empty tenant_id as an error. + +go test ./internal/grpcserver/... -run 'Resolver|TenantHeader' -v +# PushBatch with a fake TenantResolver: no resolver configured -> +# unchanged behavior, no tenant_id header at all; resolver configured -> +# every produced Kafka message carries a tenant_id header matching the +# resolved tenant; missing or invalid bearer token -> the whole batch is +# refused (codes.Unauthenticated), fail-closed, never falls back to "no +# tenant." +``` + +**Not built, and explicitly scoped out for now**: per-tenant write +routing. Neither `ingest/internal/consumer` (the ClickHouse writer) nor +`search/src/consumer.rs` (a completely independent Redpanda consumer, +not called through `ingest` at all -- see that file) reads the +`tenant_id` Kafka header back to route a record's write into a +per-tenant ClickHouse database or Tantivy index. Every record still +lands in the one shared destination regardless of tenant, correctly +tagged but not yet isolated at write time -- see CLAUDE.md and +`/docs/security/threat-model.md`'s "Read this first" for the full +disclosure. Also not built: any Helm/`docker-compose.yml` wiring that +issues an agent a real ingest credential automatically (`enterprise- +auth -create-ingest-credential-tenant=` is, like every other +credential-minting flag in this codebase, a manual operator action) -- +`deploy/helm/sentry/values.yaml`'s `ingest.requireTenantCredential` +(default `false`) only turns on *validation*, deliberately not folded +into `enterprise.enabled` directly, since flipping that flag with no +agents holding a credential yet would refuse all ingest traffic outright +rather than degrading gracefully. + ## Known gaps (do not treat this phase as done without reading these) Full accounting: `/docs/security/threat-model.md`. Headline items: @@ -578,11 +635,17 @@ Full accounting: `/docs/security/threat-model.md`. Headline items: that split (declarative request vs. imperative provisioning action) is intentional, not the "two disconnected sources of truth" gap this bullet used to describe. -- **Ingest has no tenant concept for either storage engine.** Every - record `ingest` produces lands in the one shared ClickHouse database - and the one shared Tantivy index no matter what. A newly-provisioned - tenant's storage is real, isolated at query time, and permanently - empty until this changes — undesigned, not just unbuilt. +- **Ingest now has a real tenant identity (§13), but no per-tenant + write-routing yet.** An agent presents a bearer credential + (`enterprise-auth -create-ingest-credential-tenant=`), + `ingest/internal/grpcserver.TenantResolver` validates it (fail-closed) + and attaches the resolved tenant ID to every record as a `tenant_id` + Kafka message header. Nothing downstream reads that header back yet -- + every record still lands in the one shared ClickHouse database and the + one shared Tantivy index no matter what. A newly-provisioned tenant's + storage is real, isolated at query time, and permanently empty until + the write-routing split is built (a real, scoped follow-up, no longer + an undesigned one). - **Human SSO login now works for both OIDC (§3a) and SAML (§3b)** -- each verified with a real fake IdP (genuine cryptographic signing and verification), not yet a real external IdP or a running diff --git a/docs/security/threat-model.md b/docs/security/threat-model.md index c811143..7d5c848 100644 --- a/docs/security/threat-model.md +++ b/docs/security/threat-model.md @@ -68,15 +68,30 @@ provisioned, pointing at the same ClickHouse/Postgres. The Helm chart makes the *default*, chart-managed path correct; it isn't a runtime guard against misconfiguration. -**Ingest is not tenant-aware for either storage engine**, and this is -more load-bearing than it sounds: `chrunner`/`searchclient` prove *read* -isolation given tenant-scoped data exists, but nothing writes -tenant-scoped data yet. Every record `ingest` produces lands in the one -shared ClickHouse database and the one shared (default) Tantivy index, -regardless of tenant. A newly-provisioned tenant's ClickHouse database -and Tantivy index are real, isolated, and queryable through -`enterprise-api` — and permanently empty, until ingest itself becomes -tenant-aware, which is undesigned, not just unbuilt. +**Ingest now has a real tenant identity, but no per-tenant write +routing yet** — a narrower, more precise gap than "not tenant-aware at +all." `chrunner`/`searchclient` prove *read* isolation given tenant- +scoped data exists; a new optional `ingest/internal/grpcserver. +TenantResolver` closes the "does a record know which tenant it belongs +to" half by validating a per-tenant bearer credential an agent presents +(`enterprise-auth -create-ingest-credential-tenant=` mints one; only +its SHA-256 hash is ever stored) against a new `POST +/internal/authorize-ingest` endpoint, and attaching the resolved tenant +ID to every record as a `tenant_id` Kafka message header before +producing it — fail-closed: once a resolver is configured, a missing or +invalid credential refuses the whole batch, never falls back to "no +tenant." What's still missing is the "does that identity actually +change where the record is written" half: neither `ingest`'s own +ClickHouse writer nor `search`'s independent Redpanda consumer reads +that header back to route the write anywhere per-tenant yet. Every +record still lands in the one shared ClickHouse database and the one +shared (default) Tantivy index, regardless of tenant — correctly tagged, +not yet isolated at write time. A newly-provisioned tenant's ClickHouse +database and Tantivy index remain real, isolated, and queryable through +`enterprise-api` — and permanently empty, until that write-routing split +is built (likely another "second binary," mirroring `enterprise-api` +itself), which is now scoped, disclosed remaining work, not an +undesigned gap. ## System overview @@ -114,10 +129,12 @@ sentryctl ──▶ api, alerting (Bearer token when SENTRYCTL_TOKEN is set) ``` Ingest path (agent → Redpanda → ingest → ClickHouse, and Redpanda → -search → Tantivy) carries no tenant concept at all yet either — every -ingested log record lands in the one shared `logs` table/index. Tenant -isolation for *ingest*, not just query, is out of scope for what's built -so far and is not separately designed in +search → Tantivy): `ingest` now resolves and tags each record with a +real tenant ID (see "Read this first" above), but nothing downstream +routes on it yet — every ingested log record still lands in the one +shared `logs` table/index. Tenant isolation for the *write* path is +still out of scope for what's built so far and is not separately +designed in `/docs/phase-4-isolation-design.md`; named here as a gap that design doc doesn't yet cover, not just an implementation gap. @@ -420,7 +437,8 @@ terms: | `system.*` ClickHouse metadata isolation | **Built, not live-verified** — same caveat as above | | Tantivy per-tenant index routing (`search/src/registry.rs`) | **Enforced, verified live** — real Tantivy indices, real cross-tenant probe, all passing | | Tantivy tenant_id resolution (`enterprise/internal/searchclient`) | **Enforced, verified live** — real gRPC wire-level test | -| Ingest tenant-awareness (ClickHouse and Tantivy both) | **Not implemented, undesigned** — every ingested record lands in the single shared database/index regardless of tenant | +| Ingest tenant *identity* (credential validation, tagging) | **Built and tested** — fail-closed `TenantResolver`, `tenant_id` Kafka header attached per record | +| Ingest tenant *write-routing* (ClickHouse and Tantivy both) | **Not implemented, now scoped** — every record still lands in the single shared database/index regardless of tenant; consuming the tenant_id header to route the write is real, disclosed remaining work | | Deployment actually routing traffic to `enterprise-api` (Helm) | **Enforced** — `api`/`enterprise-api` are mutually exclusive, same flag as RBAC/audit/SSO | | Deployment actually routing traffic to `enterprise-api` (docker-compose) | **Enforced** — `api`/`enterprise-api` are mutually exclusive via `COMPOSE_PROFILES`, same flag choice as Helm's `enterprise.enabled`; verified via `docker compose config`, not an actual `docker compose up` in this environment | | Human SSO login — OIDC | **Built, verified with a real fake IdP** (not yet tried against a real external IdP) | diff --git a/enterprise/README.md b/enterprise/README.md index 48d6fb2..1b5190d 100644 --- a/enterprise/README.md +++ b/enterprise/README.md @@ -182,24 +182,60 @@ silently left out: a cross-origin `fetch` with credentials from `web`'s origin needs it), neither of which is verifiable in this environment without a live backend and a browser session to exercise. -- **Ingest tenant-awareness, for either storage engine** -- `chrunner`/ - `searchclient` prove read isolation given tenant-scoped data exists, - but nothing writes it: every record `ingest` produces still lands in - the single shared ClickHouse database and the single shared Tantivy - index. A newly-provisioned tenant's storage is real and isolated, and - permanently empty. Undesigned, not just unbuilt -- see - `/docs/security/threat-model.md`. -- Any deployment-topology mechanism that actually routes traffic to - `enterprise-api` instead of `api` -- both binaries exist, - `docker-compose.yml` includes `enterprise-api` available but not - wired into `web`'s default base URL, and the Helm chart has no - service for it at all yet. **This is now the single largest gap** -- - both storage engines' isolation mechanisms themselves are built. +- **Ingest write-routing, for either storage engine** -- identity is now + real (see "Ingest tenant identity" below), but nothing consumes it + yet: `chrunner`/`searchclient` prove read isolation given tenant- + scoped data exists, and every record `ingest` produces is now tagged + with a real tenant ID, but neither `ingest/internal/consumer` (the + ClickHouse writer) nor `search/src/consumer.rs` (a completely + independent Redpanda consumer) reads that tag back to route the write + anywhere per-tenant. Every record still lands in the single shared + ClickHouse database and the single shared Tantivy index regardless of + tenant. A newly-provisioned tenant's storage is real and isolated, and + permanently empty. Now scoped, disclosed remaining work, not an + undesigned gap -- see `/docs/security/threat-model.md`. + +Deployment-topology routing (does traffic actually reach `enterprise-api` +instead of `api`) is no longer deferred -- both `deploy/helm/sentry` and +`docker-compose.yml` make it a single-flag choice now (`enterprise. +enabled` / `COMPOSE_PROFILES`), see CLAUDE.md. + +## Ingest tenant identity + +`ingest` (AGPL core) gained an optional `TenantResolver` +(`ingest/internal/grpcserver`) -- nil by default, the same "off unless +configured" shape as every other optional integration point in this +codebase. When `ENTERPRISE_AUTH_URL` is set, `PushBatch` requires an +`authorization: Bearer ` gRPC metadata entry on every call, +resolves it via a new `POST /internal/authorize-ingest` endpoint on +*this* service (`internal/authhandler`, backed by a new +`ingest_credentials` table in `internal/rbacstore` -- only a SHA-256 +hash of the token is ever stored), and attaches the resolved tenant ID +to every record as a `tenant_id` Kafka message header before producing +it. Fail-closed: once a resolver is configured, a missing or invalid +credential refuses the whole batch, never falls back to "no tenant." + +Mint a credential with `-create-ingest-credential-tenant=` (prints +the plaintext token exactly once -- see `ingest_credentials`' migration +comment for why it can't be recovered again, only reissued); +`-list-ingest-credentials-tenant=`/`-revoke-ingest-credential=` +manage existing ones. `ingest`'s own HTTP client +(`ingest/internal/tenantresolver.HTTPResolver`) is the piece that +actually calls `/internal/authorize-ingest` -- never an `enterprise/` +import (`ingest` is AGPL core), same "network boundary, not import +boundary" shape `api/authz.HTTPAuthorizer` already uses for the query +path. + +**What this does not do**: change where a record is actually written. +See "Deliberately deferred" above -- attaching a verified tenant +identity as early as possible (right where the credential is presented) +was built as a self-contained first step; per-tenant write-routing for +both storage engines is separate, scoped follow-up work. ## Package layout ``` -cmd/enterprise-auth/ config loading, OIDC discovery at startup, health/authorize/features endpoints, -mint-service-token, -create-tenant, -grant-membership-*, -revoke-membership-*, -list-memberships-tenant +cmd/enterprise-auth/ config loading, OIDC discovery at startup, health/authorize/features/authorize-ingest endpoints, -mint-service-token, -create-tenant, -grant-membership-*, -revoke-membership-*, -list-memberships-tenant, -create-ingest-credential-tenant, -list-ingest-credentials-tenant, -revoke-ingest-credential cmd/enterprise-api/ multi-tenant-aware alternative to api/cmd/api -- see its own doc comment internal/tenant/ the ID type -- see its package doc comment before touching it internal/oidc/ coreos/go-oidc wiring: discovery, login redirect, code exchange + ID token verification @@ -218,8 +254,13 @@ internal/apiconfig/ enterprise-api's own env-var config internal/config/ enterprise-auth's env-var config ``` -Future additions: ingest tenant-awareness (undesigned), and real -deployment-topology wiring for `enterprise-api` -- see "Status" above. +`ingest/internal/tenantresolver` (AGPL core, not enterprise/, since +ingest must never import enterprise/) is the client side of `internal/ +authhandler`'s new `POST /internal/authorize-ingest` -- see "Ingest +tenant identity" above. + +Future additions: per-tenant write-routing for ingest (ClickHouse and +Tantivy both) -- see "Ingest tenant identity" above. ## Why OIDC and SAML aren't hand-rolled diff --git a/enterprise/cmd/enterprise-auth/main.go b/enterprise/cmd/enterprise-auth/main.go index 52a4dbb..6d728f3 100644 --- a/enterprise/cmd/enterprise-auth/main.go +++ b/enterprise/cmd/enterprise-auth/main.go @@ -79,6 +79,9 @@ func main() { revokeTenant := flag.String("revoke-membership-tenant", "", "tenant id to revoke a membership from -- both -revoke-membership-* flags are required together") revokeUserEmail := flag.String("revoke-membership-user-email", "", "email of the user whose tenant_memberships row to delete") listMembershipsTenant := flag.String("list-memberships-tenant", "", "print every user with a membership in this tenant (id, email, display name, role) and exit") + createIngestCredentialTenant := flag.String("create-ingest-credential-tenant", "", "mint a new ingest bearer token for this tenant, print it once, and exit -- see ingest/internal/grpcserver.TenantResolver") + listIngestCredentialsTenant := flag.String("list-ingest-credentials-tenant", "", "print every ingest credential's id/created_at for this tenant (never the token itself -- only its hash is stored) and exit") + revokeIngestCredential := flag.String("revoke-ingest-credential", "", "delete an ingest credential by id (see -list-ingest-credentials-tenant) and exit") // -healthcheck: same self-check mode as api/-healthcheck (see that // binary's doc comment) -- enterprise-auth's image is distroless too. healthcheck := flag.Bool("healthcheck", false, "self-check mode for Docker's HEALTHCHECK") @@ -132,6 +135,15 @@ func main() { if *listMembershipsTenant != "" { os.Exit(runListMemberships(ctx, logger, rbac, *listMembershipsTenant)) } + if *createIngestCredentialTenant != "" { + os.Exit(runCreateIngestCredential(ctx, logger, rbac, *createIngestCredentialTenant)) + } + if *listIngestCredentialsTenant != "" { + os.Exit(runListIngestCredentials(ctx, logger, rbac, *listIngestCredentialsTenant)) + } + if *revokeIngestCredential != "" { + os.Exit(runRevokeIngestCredential(ctx, logger, rbac, *revokeIngestCredential)) + } // oidcProvider stays nil (loginhandler.RegisterRoutes then registers // nothing) unless OIDC is actually configured -- matches every other @@ -187,7 +199,7 @@ func main() { OIDCEnabled: cfg.OIDC.IssuerURL != "", SAMLEnabled: cfg.SAML.IDPMetadataURL != "", } - authhandler.New(logger, sessionManager, features).RegisterRoutes(mux) + authhandler.New(logger, sessionManager, features, rbac).RegisterRoutes(mux) loginhandler.New(logger, oidcProvider, samlProvider, sessionManager, rbac, cfg.PostLoginRedirectURL, cfg.SelectTenantRedirectURL).RegisterRoutes(mux) srv := &http.Server{Addr: cfg.HTTPListenAddr, Handler: mux} @@ -355,6 +367,60 @@ func runListMemberships(ctx context.Context, logger *slog.Logger, rbac *rbacstor return 0 } +// runCreateIngestCredential mints a new ingest bearer token for a +// tenant and prints it to stdout exactly once -- rbacstore only ever +// stores its hash (see ingest_credentials's doc comment), so this +// output is the only chance to capture the plaintext. An agent presents +// it as an `Authorization: Bearer ` gRPC metadata entry on every +// PushBatch call; ingest resolves it to a tenant via +// POST /internal/authorize-ingest. +func runCreateIngestCredential(ctx context.Context, logger *slog.Logger, rbac *rbacstore.Store, tenantID string) int { + if _, err := rbac.GetTenant(ctx, tenantID); err != nil { + logger.Error("looking up tenant", "tenant_id", tenantID, "error", err) + return 1 + } + token, err := rbac.CreateIngestCredential(ctx, tenantID) + if err != nil { + logger.Error("creating ingest credential", "error", err) + return 1 + } + fmt.Println(token) + return 0 +} + +func runListIngestCredentials(ctx context.Context, logger *slog.Logger, rbac *rbacstore.Store, tenantID string) int { + if _, err := rbac.GetTenant(ctx, tenantID); err != nil { + logger.Error("looking up tenant", "tenant_id", tenantID, "error", err) + return 1 + } + creds, err := rbac.ListIngestCredentialsForTenant(ctx, tenantID) + if err != nil { + logger.Error("listing ingest credentials", "error", err) + return 1 + } + if len(creds) == 0 { + fmt.Println("(no ingest credentials)") + return 0 + } + for _, c := range creds { + fmt.Printf("%s\t%s\n", c.ID, c.CreatedAt.Format(time.RFC3339)) + } + return 0 +} + +func runRevokeIngestCredential(ctx context.Context, logger *slog.Logger, rbac *rbacstore.Store, id string) int { + if err := rbac.RevokeIngestCredential(ctx, id); err != nil { + if err == rbacstore.ErrNotFound { + logger.Error("no ingest credential with this id", "id", id) + } else { + logger.Error("revoking ingest credential", "error", err) + } + return 1 + } + logger.Info("revoked ingest credential", "id", id) + return 0 +} + // runHealthcheck mirrors api/cmd/api/main.go's runHealthcheck exactly -- // see that function's doc comment for why this execs the binary against // itself rather than using an external tool. diff --git a/enterprise/internal/authhandler/authhandler.go b/enterprise/internal/authhandler/authhandler.go index 6498fc9..ea2d7d0 100644 --- a/enterprise/internal/authhandler/authhandler.go +++ b/enterprise/internal/authhandler/authhandler.go @@ -7,9 +7,16 @@ // signed tokens with a different Role claim, so one validation path // handles both, and the Role claim (not which header carried it) is what // determines whether the result looks like a human or a service identity. +// +// POST /internal/authorize-ingest is a sibling endpoint, same network- +// boundary shape but for a different caller (`ingest`, core/AGPL, not +// api/authz) and a different credential type (an ingest bearer token +// checked against rbacstore's ingest_credentials table, not a +// session.Manager JWT) -- see ingestCredentialValidator's doc comment. package authhandler import ( + "context" "encoding/json" "log/slog" "net/http" @@ -32,19 +39,32 @@ type Features struct { SAMLEnabled bool } -type Handler struct { - logger *slog.Logger - manager *session.Manager - features Features +// ingestCredentialValidator is the narrow interface POST +// /internal/authorize-ingest needs -- *rbacstore.Store is the production +// implementation. Unlike session-backed /internal/authorize, this +// endpoint validates a completely different credential type (an ingest +// bearer token, checked against enterprise/internal/rbacstore's +// ingest_credentials table, never a session.Manager-signed JWT), so it +// needs a dependency session.Manager alone can't supply. +type ingestCredentialValidator interface { + ValidateIngestCredential(ctx context.Context, token string) (tenantID string, err error) } -func New(logger *slog.Logger, manager *session.Manager, features Features) *Handler { - return &Handler{logger: logger, manager: manager, features: features} +type Handler struct { + logger *slog.Logger + manager *session.Manager + features Features + ingestCredentials ingestCredentialValidator +} + +func New(logger *slog.Logger, manager *session.Manager, features Features, ingestCredentials ingestCredentialValidator) *Handler { + return &Handler{logger: logger, manager: manager, features: features, ingestCredentials: ingestCredentials} } func (h *Handler) RegisterRoutes(mux *http.ServeMux) { mux.HandleFunc("POST /internal/authorize", h.handleAuthorize) mux.HandleFunc("GET /auth/features", h.handleFeatures) + mux.HandleFunc("POST /internal/authorize-ingest", h.handleAuthorizeIngest) } type featuresResponse struct { @@ -100,6 +120,32 @@ func (h *Handler) handleAuthorize(w http.ResponseWriter, r *http.Request) { }) } +type authorizeIngestResponse struct { + TenantID string `json:"tenant_id"` +} + +// handleAuthorizeIngest is ingest/internal/grpcserver.HTTPTenantResolver's +// server side -- ingest calls this once per PushBatch (with the bearer +// token the agent presented) to resolve which tenant the batch belongs +// to, the network-boundary equivalent of api/authz.HTTPAuthorizer +// calling /internal/authorize, for a different credential type. +func (h *Handler) handleAuthorizeIngest(w http.ResponseWriter, r *http.Request) { + token := bearerToken(r.Header.Get("Authorization")) + if token == "" { + http.Error(w, "no credentials presented", http.StatusUnauthorized) + return + } + + tenantID, err := h.ingestCredentials.ValidateIngestCredential(r.Context(), token) + if err != nil { + http.Error(w, "invalid ingest credential", http.StatusUnauthorized) + return + } + + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(authorizeIngestResponse{TenantID: tenantID}) +} + func bearerToken(header string) string { const prefix = "Bearer " if !strings.HasPrefix(header, prefix) { diff --git a/enterprise/internal/authhandler/authhandler_test.go b/enterprise/internal/authhandler/authhandler_test.go index 140bde0..2503de3 100644 --- a/enterprise/internal/authhandler/authhandler_test.go +++ b/enterprise/internal/authhandler/authhandler_test.go @@ -1,6 +1,7 @@ package authhandler import ( + "context" "encoding/json" "io" "log/slog" @@ -11,13 +12,37 @@ import ( "github.com/sentry/sentry/enterprise/internal/session" ) +// fakeIngestCredentialValidator is an in-memory stand-in for +// *rbacstore.Store's ValidateIngestCredential, keyed by token. +type fakeIngestCredentialValidator struct { + tenantByToken map[string]string +} + +func newFakeIngestCredentialValidator() *fakeIngestCredentialValidator { + return &fakeIngestCredentialValidator{tenantByToken: map[string]string{}} +} + +func (f *fakeIngestCredentialValidator) ValidateIngestCredential(_ context.Context, token string) (string, error) { + tenantID, ok := f.tenantByToken[token] + if !ok { + return "", errNotFound + } + return tenantID, nil +} + +var errNotFound = &fakeNotFoundError{} + +type fakeNotFoundError struct{} + +func (*fakeNotFoundError) Error() string { return "not found" } + func testHandler(t *testing.T) (*Handler, *session.Manager) { t.Helper() m, err := session.NewManager([]byte("this-is-a-32-byte-test-signing-key!")) if err != nil { t.Fatalf("session.NewManager: %v", err) } - return New(slog.New(slog.NewTextHandler(io.Discard, nil)), m, Features{}), m + return New(slog.New(slog.NewTextHandler(io.Discard, nil)), m, Features{}, newFakeIngestCredentialValidator()), m } func doAuthorize(t *testing.T, h *Handler, mutate func(*http.Request)) *httptest.ResponseRecorder { @@ -121,7 +146,7 @@ func TestFeaturesReflectsConfiguredMechanisms(t *testing.T) { if err != nil { t.Fatalf("session.NewManager: %v", err) } - h := New(slog.New(slog.NewTextHandler(io.Discard, nil)), m, Features{OIDCEnabled: true, SAMLEnabled: false}) + h := New(slog.New(slog.NewTextHandler(io.Discard, nil)), m, Features{OIDCEnabled: true, SAMLEnabled: false}, newFakeIngestCredentialValidator()) mux := http.NewServeMux() h.RegisterRoutes(mux) @@ -173,3 +198,78 @@ func TestAuthorizeTokenFromWrongManagerIsUnauthorized(t *testing.T) { t.Fatalf("status = %d, want 401", rec.Code) } } + +func doAuthorizeIngest(t *testing.T, h *Handler, mutate func(*http.Request)) *httptest.ResponseRecorder { + t.Helper() + mux := http.NewServeMux() + h.RegisterRoutes(mux) + req := httptest.NewRequest(http.MethodPost, "/internal/authorize-ingest", nil) + if mutate != nil { + mutate(req) + } + rec := httptest.NewRecorder() + mux.ServeHTTP(rec, req) + return rec +} + +func TestAuthorizeIngestResolvesTenant(t *testing.T) { + m, err := session.NewManager([]byte("this-is-a-32-byte-test-signing-key!")) + if err != nil { + t.Fatalf("session.NewManager: %v", err) + } + validator := newFakeIngestCredentialValidator() + validator.tenantByToken["real-token"] = "acme" + h := New(slog.New(slog.NewTextHandler(io.Discard, nil)), m, Features{}, validator) + + rec := doAuthorizeIngest(t, h, func(r *http.Request) { + r.Header.Set("Authorization", "Bearer real-token") + }) + if rec.Code != http.StatusOK { + t.Fatalf("status = %d, body = %s", rec.Code, rec.Body.String()) + } + var body authorizeIngestResponse + if err := json.NewDecoder(rec.Body).Decode(&body); err != nil { + t.Fatalf("decoding response: %v", err) + } + if body.TenantID != "acme" { + t.Fatalf("TenantID = %q, want acme", body.TenantID) + } +} + +func TestAuthorizeIngestNoCredentialsIsUnauthorized(t *testing.T) { + h, _ := testHandler(t) + rec := doAuthorizeIngest(t, h, nil) + if rec.Code != http.StatusUnauthorized { + t.Fatalf("status = %d, want 401", rec.Code) + } +} + +func TestAuthorizeIngestUnknownTokenIsUnauthorized(t *testing.T) { + h, _ := testHandler(t) + rec := doAuthorizeIngest(t, h, func(r *http.Request) { + r.Header.Set("Authorization", "Bearer not-a-real-token") + }) + if rec.Code != http.StatusUnauthorized { + t.Fatalf("status = %d, want 401", rec.Code) + } +} + +// TestAuthorizeIngestRejectsSessionToken is the regression test for the +// two /internal/authorize* endpoints validating genuinely different +// credential types: a real session.Manager-signed token (a service +// token or human session) must not work as an ingest credential, since +// it was never checked against rbacstore.ValidateIngestCredential -- +// this endpoint doesn't call session.Manager.Validate at all. +func TestAuthorizeIngestRejectsSessionToken(t *testing.T) { + h, m := testHandler(t) + sessionToken, err := m.IssueServiceToken("alerting") + if err != nil { + t.Fatalf("IssueServiceToken: %v", err) + } + rec := doAuthorizeIngest(t, h, func(r *http.Request) { + r.Header.Set("Authorization", "Bearer "+sessionToken) + }) + if rec.Code != http.StatusUnauthorized { + t.Fatalf("status = %d, want 401 (a session token must not validate as an ingest credential)", rec.Code) + } +} diff --git a/enterprise/internal/rbacstore/ingest_credentials.go b/enterprise/internal/rbacstore/ingest_credentials.go new file mode 100644 index 0000000..8939d0a --- /dev/null +++ b/enterprise/internal/rbacstore/ingest_credentials.go @@ -0,0 +1,110 @@ +package rbacstore + +import ( + "context" + "crypto/rand" + "crypto/sha256" + "encoding/base64" + "encoding/hex" + "errors" + "fmt" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5" +) + +// IngestCredential is one ingest_credentials row -- see +// metadata/migrations/0034_create_ingest_credentials.sql's doc comment +// for why only a hash is stored. Consumed by +// ingest/internal/grpcserver.TenantResolver (an HTTP call to +// enterprise-auth's POST /internal/authorize-ingest, which calls +// ValidateIngestCredential below) so an agent's records can be +// attributed to a tenant at the point they enter the system. +type IngestCredential struct { + ID string + TenantID string + CreatedAt time.Time +} + +func hashIngestToken(token string) string { + sum := sha256.Sum256([]byte(token)) + return hex.EncodeToString(sum[:]) +} + +// CreateIngestCredential generates a new bearer token for tenantID and +// returns the plaintext exactly once -- only its hash is ever persisted +// (see this file's package doc comment). There is no way to retrieve a +// lost token again; the only recovery is issuing a new one +// (RevokeIngestCredential + CreateIngestCredential), the same "can't +// recover, can only reissue" UX every real API-key system uses. +func (s *Store) CreateIngestCredential(ctx context.Context, tenantID string) (token string, err error) { + raw := make([]byte, 32) + if _, err := rand.Read(raw); err != nil { + return "", fmt.Errorf("rbacstore: generating ingest credential: %w", err) + } + token = base64.RawURLEncoding.EncodeToString(raw) + + _, err = s.pool.Exec(ctx, ` + INSERT INTO ingest_credentials (id, tenant_id, token_hash) + VALUES ($1, $2, $3)`, + uuid.NewString(), tenantID, hashIngestToken(token)) + if err != nil { + return "", fmt.Errorf("rbacstore: creating ingest credential: %w", err) + } + return token, nil +} + +// ValidateIngestCredential hashes the presented token and looks up which +// tenant it belongs to via an indexed exact-match on the UNIQUE +// token_hash column -- the only production call site is +// enterprise-auth's POST /internal/authorize-ingest handler, per an +// agent's PushBatch request. +func (s *Store) ValidateIngestCredential(ctx context.Context, token string) (tenantID string, err error) { + row := s.pool.QueryRow(ctx, `SELECT tenant_id FROM ingest_credentials WHERE token_hash = $1`, hashIngestToken(token)) + if err := row.Scan(&tenantID); err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return "", ErrNotFound + } + return "", fmt.Errorf("rbacstore: validating ingest credential: %w", err) + } + return tenantID, nil +} + +// RevokeIngestCredential deletes a credential by ID (not by token -- +// the plaintext is never stored, so revocation has to name the row some +// other way; ListIngestCredentialsForTenant is what an operator uses to +// find the ID). +func (s *Store) RevokeIngestCredential(ctx context.Context, id string) error { + tag, err := s.pool.Exec(ctx, `DELETE FROM ingest_credentials WHERE id = $1`, id) + if err != nil { + return fmt.Errorf("rbacstore: revoking ingest credential: %w", err) + } + if tag.RowsAffected() == 0 { + return ErrNotFound + } + return nil +} + +// ListIngestCredentialsForTenant never returns the plaintext token (it +// isn't stored) -- just enough (ID, creation time) for an operator to +// decide which one to revoke. +func (s *Store) ListIngestCredentialsForTenant(ctx context.Context, tenantID string) ([]IngestCredential, error) { + rows, err := s.pool.Query(ctx, ` + SELECT id, tenant_id, created_at FROM ingest_credentials + WHERE tenant_id = $1 ORDER BY created_at`, tenantID) + if err != nil { + return nil, fmt.Errorf("rbacstore: listing ingest credentials: %w", err) + } + defer rows.Close() + + var out []IngestCredential + for rows.Next() { + var c IngestCredential + if err := rows.Scan(&c.ID, &c.TenantID, &c.CreatedAt); err != nil { + return nil, fmt.Errorf("rbacstore: scanning ingest credential: %w", err) + } + out = append(out, c) + } + return out, rows.Err() +} diff --git a/enterprise/internal/rbacstore/rbacstore_test.go b/enterprise/internal/rbacstore/rbacstore_test.go index e166d4b..538f075 100644 --- a/enterprise/internal/rbacstore/rbacstore_test.go +++ b/enterprise/internal/rbacstore/rbacstore_test.go @@ -771,3 +771,122 @@ func TestListMembershipsForTenant(t *testing.T) { t.Fatalf("unexpected editor entry: %+v", got) } } + +func TestCreateAndValidateIngestCredential(t *testing.T) { + s := testStore(t) + ctx := context.Background() + tenantID := "test-tenant-" + uniqueSuffix() + if _, err := s.CreateTenant(ctx, tenantID, "Test Tenant"); err != nil { + t.Fatalf("CreateTenant: %v", err) + } + + token, err := s.CreateIngestCredential(ctx, tenantID) + if err != nil { + t.Fatalf("CreateIngestCredential: %v", err) + } + if token == "" { + t.Fatal("expected a non-empty token") + } + + got, err := s.ValidateIngestCredential(ctx, token) + if err != nil { + t.Fatalf("ValidateIngestCredential: %v", err) + } + if got != tenantID { + t.Fatalf("ValidateIngestCredential tenant = %q, want %q", got, tenantID) + } +} + +func TestValidateIngestCredentialRejectsUnknownToken(t *testing.T) { + s := testStore(t) + if _, err := s.ValidateIngestCredential(context.Background(), "not-a-real-token"); err != ErrNotFound { + t.Fatalf("ValidateIngestCredential error = %v, want ErrNotFound", err) + } +} + +// TestIngestCredentialTokenNeverStoredAsPlaintext is the regression test +// for this table's whole reason for hashing: the raw token string must +// not appear anywhere in the persisted row (only its hash), so a +// database leak doesn't hand out usable credentials. +func TestIngestCredentialTokenNeverStoredAsPlaintext(t *testing.T) { + s := testStore(t) + ctx := context.Background() + tenantID := "test-tenant-" + uniqueSuffix() + if _, err := s.CreateTenant(ctx, tenantID, "Test Tenant"); err != nil { + t.Fatalf("CreateTenant: %v", err) + } + token, err := s.CreateIngestCredential(ctx, tenantID) + if err != nil { + t.Fatalf("CreateIngestCredential: %v", err) + } + + var stored string + row := s.pool.QueryRow(ctx, `SELECT token_hash FROM ingest_credentials WHERE tenant_id = $1`, tenantID) + if err := row.Scan(&stored); err != nil { + t.Fatalf("reading stored token_hash: %v", err) + } + if stored == token { + t.Fatal("the plaintext token must never be stored directly in token_hash") + } + if stored != hashIngestToken(token) { + t.Fatalf("stored hash = %q, want sha256(token) = %q", stored, hashIngestToken(token)) + } +} + +func TestRevokeIngestCredential(t *testing.T) { + s := testStore(t) + ctx := context.Background() + tenantID := "test-tenant-" + uniqueSuffix() + if _, err := s.CreateTenant(ctx, tenantID, "Test Tenant"); err != nil { + t.Fatalf("CreateTenant: %v", err) + } + token, err := s.CreateIngestCredential(ctx, tenantID) + if err != nil { + t.Fatalf("CreateIngestCredential: %v", err) + } + creds, err := s.ListIngestCredentialsForTenant(ctx, tenantID) + if err != nil || len(creds) != 1 { + t.Fatalf("ListIngestCredentialsForTenant = (%+v, %v), want exactly one", creds, err) + } + + if err := s.RevokeIngestCredential(ctx, creds[0].ID); err != nil { + t.Fatalf("RevokeIngestCredential: %v", err) + } + if _, err := s.ValidateIngestCredential(ctx, token); err != ErrNotFound { + t.Fatalf("ValidateIngestCredential after revoke = %v, want ErrNotFound", err) + } +} + +func TestRevokeIngestCredentialNotFound(t *testing.T) { + s := testStore(t) + if err := s.RevokeIngestCredential(context.Background(), uuid.NewString()); err != ErrNotFound { + t.Fatalf("RevokeIngestCredential error = %v, want ErrNotFound", err) + } +} + +func TestListIngestCredentialsForTenantExcludesOtherTenants(t *testing.T) { + s := testStore(t) + ctx := context.Background() + tenantID := "test-tenant-" + uniqueSuffix() + otherTenantID := "test-tenant-" + uniqueSuffix() + if _, err := s.CreateTenant(ctx, tenantID, "Test Tenant"); err != nil { + t.Fatalf("CreateTenant: %v", err) + } + if _, err := s.CreateTenant(ctx, otherTenantID, "Other Tenant"); err != nil { + t.Fatalf("CreateTenant other: %v", err) + } + if _, err := s.CreateIngestCredential(ctx, tenantID); err != nil { + t.Fatalf("CreateIngestCredential: %v", err) + } + if _, err := s.CreateIngestCredential(ctx, otherTenantID); err != nil { + t.Fatalf("CreateIngestCredential other: %v", err) + } + + creds, err := s.ListIngestCredentialsForTenant(ctx, tenantID) + if err != nil { + t.Fatalf("ListIngestCredentialsForTenant: %v", err) + } + if len(creds) != 1 || creds[0].TenantID != tenantID { + t.Fatalf("unexpected credentials: %+v", creds) + } +} diff --git a/ingest/cmd/ingest/main.go b/ingest/cmd/ingest/main.go index f6b504a..e3e9c1f 100644 --- a/ingest/cmd/ingest/main.go +++ b/ingest/cmd/ingest/main.go @@ -26,6 +26,7 @@ import ( "github.com/sentry/sentry/ingest/internal/consumer" "github.com/sentry/sentry/ingest/internal/grpcserver" "github.com/sentry/sentry/ingest/internal/producer" + "github.com/sentry/sentry/ingest/internal/tenantresolver" ) func main() { @@ -53,7 +54,18 @@ func main() { if *mode == "server" || *mode == "all" { p := producer.New(cfg.Redpanda) defer p.Close() - srv := grpcserver.New(logger, cfg.GRPC, cfg.TLS, p) + // resolver stays nil (every batch's tenant_id header is simply + // never set) unless ENTERPRISE_AUTH_URL is configured -- matches + // every other "off unless configured" optional dependency in + // this codebase. + var resolver grpcserver.TenantResolver + if cfg.EnterpriseAuthURL != "" { + resolver = tenantresolver.New(cfg.EnterpriseAuthURL) + logger.Info("ingest tenant resolution configured", "enterprise_auth_url", cfg.EnterpriseAuthURL) + } else { + logger.Info("ENTERPRISE_AUTH_URL not set -- ingest records carry no tenant_id, single-tenant behavior") + } + srv := grpcserver.New(logger, cfg.GRPC, cfg.TLS, p, resolver) g.Go(func() error { return srv.Run(ctx) }) } diff --git a/ingest/internal/config/config.go b/ingest/internal/config/config.go index 40955a8..42e5326 100644 --- a/ingest/internal/config/config.go +++ b/ingest/internal/config/config.go @@ -17,6 +17,12 @@ type Config struct { Redpanda RedpandaConfig ClickHouse ClickHouseConfig Batch BatchConfig + // EnterpriseAuthURL enables per-tenant ingest credential validation + // (internal/grpcserver.TenantResolver) when set -- empty (the + // default) is a documented no-op, same "off unless configured" shape + // as every other optional enterprise integration point in this + // codebase (e.g. api's own ENTERPRISE_AUTH_URL). + EnterpriseAuthURL string } type GRPCConfig struct { @@ -70,6 +76,7 @@ func Load() (Config, error) { Username: getenv("CLICKHOUSE_USERNAME", "default"), Password: getenv("CLICKHOUSE_PASSWORD", ""), }, + EnterpriseAuthURL: getenv("ENTERPRISE_AUTH_URL", ""), } maxSize, err := strconv.Atoi(getenv("CONSUMER_BATCH_MAX_SIZE", "500")) diff --git a/ingest/internal/config/config_test.go b/ingest/internal/config/config_test.go index 6ec7f8a..63a6f2c 100644 --- a/ingest/internal/config/config_test.go +++ b/ingest/internal/config/config_test.go @@ -19,12 +19,16 @@ func TestLoadDefaults(t *testing.T) { if cfg.Batch.FlushIntervalMS != 2000 { t.Errorf("Batch.FlushIntervalMS = %d, want 2000", cfg.Batch.FlushIntervalMS) } + if cfg.EnterpriseAuthURL != "" { + t.Errorf("EnterpriseAuthURL = %q, want empty (tenant resolution off by default)", cfg.EnterpriseAuthURL) + } } func TestLoadOverridesFromEnv(t *testing.T) { t.Setenv("GRPC_LISTEN_ADDR", ":9999") t.Setenv("REDPANDA_BROKERS", "a:9092,b:9092") t.Setenv("CONSUMER_BATCH_MAX_SIZE", "10") + t.Setenv("ENTERPRISE_AUTH_URL", "http://enterprise-auth:8082") cfg, err := Load() if err != nil { @@ -39,6 +43,9 @@ func TestLoadOverridesFromEnv(t *testing.T) { if cfg.Batch.MaxSize != 10 { t.Errorf("Batch.MaxSize = %d, want 10", cfg.Batch.MaxSize) } + if cfg.EnterpriseAuthURL != "http://enterprise-auth:8082" { + t.Errorf("EnterpriseAuthURL = %q, want http://enterprise-auth:8082", cfg.EnterpriseAuthURL) + } } func TestLoadInvalidBatchSizeErrors(t *testing.T) { diff --git a/ingest/internal/grpcserver/server.go b/ingest/internal/grpcserver/server.go index 85e04f9..398691a 100644 --- a/ingest/internal/grpcserver/server.go +++ b/ingest/internal/grpcserver/server.go @@ -4,6 +4,22 @@ // happen exactly once, here, rather than in either downstream consumer) // and otherwise forwards records unchanged onto Redpanda — normalization // into the ClickHouse row shape happens later, on the consumer side. +// +// If a TenantResolver is configured, PushBatch also resolves which +// tenant the call's bearer credential belongs to and attaches it as a +// "tenant_id" Kafka message header on every record produced -- the first +// step of Phase 4's ingest tenant-awareness (see +// /docs/phase-4-runbook.md and CLAUDE.md's "ingest itself has no tenant +// concept" gap). Deliberately scoped no further than that for now: +// nothing downstream (this package's own consumer, or `search`'s +// separate Redpanda consumer) reads that header yet to route a record's +// write into a per-tenant ClickHouse database/Tantivy index -- every +// record still lands in the one shared destination either way, tenant_id +// header or not. That's real, disclosed, deferred follow-up work, not +// silently incomplete: attaching a verifiable tenant identity as early +// as possible (right where the credential is actually presented) is a +// self-contained, independently valuable step on its own, and it's what +// any later per-tenant write-routing work will consume. package grpcserver import ( @@ -11,12 +27,14 @@ import ( "fmt" "log/slog" "net" + "strings" "github.com/google/uuid" "github.com/segmentio/kafka-go" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials" + "google.golang.org/grpc/metadata" "google.golang.org/grpc/status" "google.golang.org/protobuf/proto" @@ -24,6 +42,12 @@ import ( logsv1 "github.com/sentry/sentry/proto/sentry/logs/v1" ) +// TenantIDHeaderKey is the Kafka message header a resolved tenant ID is +// attached under -- exported so internal/consumer (or a future per- +// tenant write-routing consumer) can read it back by the same name +// without duplicating the literal. +const TenantIDHeaderKey = "tenant_id" + type Server struct { logsv1.UnimplementedLogIngestServer @@ -31,6 +55,7 @@ type Server struct { grpcCfg config.GRPCConfig tlsCfg config.TLSConfig producer batchProducer + resolver TenantResolver } // batchProducer is the subset of *producer.Producer this package depends @@ -39,8 +64,22 @@ type batchProducer interface { WriteBatch(ctx context.Context, msgs []kafka.Message) error } -func New(logger *slog.Logger, grpcCfg config.GRPCConfig, tlsCfg config.TLSConfig, p batchProducer) *Server { - return &Server{logger: logger, grpcCfg: grpcCfg, tlsCfg: tlsCfg, producer: p} +// TenantResolver validates an ingest credential (a bearer token +// presented via gRPC metadata, `authorization: Bearer `) and +// resolves which tenant it belongs to. nil is a deliberate no-op: every +// record's Kafka message gets no tenant_id header at all, matching every +// ingest deployment's behavior before per-tenant ingest credentials +// existed. The real implementation +// (ingest/internal/tenantresolver.HTTPResolver) is a plain HTTP client +// calling enterprise-auth's /internal/authorize-ingest -- never an +// enterprise/ import, since this package is AGPL core (same "network +// boundary, not import boundary" shape api/authz.Authorizer uses). +type TenantResolver interface { + ResolveTenant(ctx context.Context, token string) (tenantID string, err error) +} + +func New(logger *slog.Logger, grpcCfg config.GRPCConfig, tlsCfg config.TLSConfig, p batchProducer, resolver TenantResolver) *Server { + return &Server{logger: logger, grpcCfg: grpcCfg, tlsCfg: tlsCfg, producer: p, resolver: resolver} } // Run blocks serving gRPC until ctx is canceled, then gracefully stops. @@ -77,6 +116,28 @@ func (s *Server) PushBatch(ctx context.Context, req *logsv1.PushBatchRequest) (* return &logsv1.PushBatchResponse{Accepted: 0}, nil } + // tenantID stays empty (no header attached below) unless a resolver + // is actually configured -- single-tenant deployments never present + // a bearer credential and never need to. Once a resolver IS + // configured, a missing/invalid credential fails the whole batch + // closed rather than falling back to "no tenant" -- exactly the + // same fail-closed shape enterprise/internal/chrunner.Registry.RunSQL + // uses on the read side, applied here at the point data enters the + // system. + var tenantID string + if s.resolver != nil { + token, ok := bearerTokenFromContext(ctx) + if !ok { + return nil, status.Error(codes.Unauthenticated, "missing bearer credential") + } + resolved, err := s.resolver.ResolveTenant(ctx, token) + if err != nil { + s.logger.Error("resolving ingest tenant", "batch_id", req.GetBatchId(), "error", err) + return nil, status.Error(codes.Unauthenticated, "invalid ingest credential") + } + tenantID = resolved + } + msgs := make([]kafka.Message, 0, len(req.GetRecords())) for _, rec := range req.GetRecords() { // Assigned here, once, before this record is produced to @@ -92,10 +153,14 @@ func (s *Server) PushBatch(ctx context.Context, req *logsv1.PushBatchRequest) (* if err != nil { return nil, status.Errorf(codes.InvalidArgument, "marshaling record: %v", err) } - msgs = append(msgs, kafka.Message{ + msg := kafka.Message{ Key: []byte(rec.GetHost()), Value: val, - }) + } + if tenantID != "" { + msg.Headers = []kafka.Header{{Key: TenantIDHeaderKey, Value: []byte(tenantID)}} + } + msgs = append(msgs, msg) } if err := s.producer.WriteBatch(ctx, msgs); err != nil { @@ -103,6 +168,26 @@ func (s *Server) PushBatch(ctx context.Context, req *logsv1.PushBatchRequest) (* return nil, status.Errorf(codes.Unavailable, "writing to transport: %v", err) } - s.logger.Debug("batch produced to redpanda", "batch_id", req.GetBatchId(), "records", len(req.GetRecords())) + s.logger.Debug("batch produced to redpanda", "batch_id", req.GetBatchId(), "records", len(req.GetRecords()), "tenant_id", tenantID) return &logsv1.PushBatchResponse{Accepted: uint32(len(req.GetRecords()))}, nil } + +// bearerTokenFromContext reads the same "authorization: Bearer " +// gRPC metadata shape HTTP's Authorization header uses -- an agent sets +// this once per PushBatch call (see the agent's grpc.rs), not per +// record. +func bearerTokenFromContext(ctx context.Context) (string, bool) { + md, ok := metadata.FromIncomingContext(ctx) + if !ok { + return "", false + } + values := md.Get("authorization") + if len(values) == 0 { + return "", false + } + const prefix = "Bearer " + if !strings.HasPrefix(values[0], prefix) { + return "", false + } + return strings.TrimPrefix(values[0], prefix), true +} diff --git a/ingest/internal/grpcserver/server_test.go b/ingest/internal/grpcserver/server_test.go index 538fe0a..fec1f4b 100644 --- a/ingest/internal/grpcserver/server_test.go +++ b/ingest/internal/grpcserver/server_test.go @@ -2,12 +2,16 @@ package grpcserver import ( "context" + "errors" "io" "log/slog" "sync" "testing" "github.com/segmentio/kafka-go" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/status" "google.golang.org/protobuf/proto" "github.com/sentry/sentry/ingest/internal/config" @@ -32,8 +36,35 @@ func (f *fakeProducer) WriteBatch(_ context.Context, msgs []kafka.Message) error return nil } +// fakeResolver is an in-memory stand-in for +// ingest/internal/tenantresolver.HTTPResolver, keyed by token. +type fakeResolver struct { + tenantByToken map[string]string +} + +func (f *fakeResolver) ResolveTenant(_ context.Context, token string) (string, error) { + tenantID, ok := f.tenantByToken[token] + if !ok { + return "", errors.New("fakeResolver: unknown token") + } + return tenantID, nil +} + func newTestServer(p batchProducer) *Server { - return New(slog.New(slog.NewTextHandler(io.Discard, nil)), config.GRPCConfig{}, config.TLSConfig{}, p) + return New(slog.New(slog.NewTextHandler(io.Discard, nil)), config.GRPCConfig{}, config.TLSConfig{}, p, nil) +} + +func newTestServerWithResolver(p batchProducer, resolver TenantResolver) *Server { + return New(slog.New(slog.NewTextHandler(io.Discard, nil)), config.GRPCConfig{}, config.TLSConfig{}, p, resolver) +} + +// contextWithBearerToken builds an incoming gRPC context carrying an +// "authorization: Bearer " metadata entry -- the shape a real +// grpc-go server hands PushBatch once TLS/framing is stripped away, so +// this exercises the same metadata.FromIncomingContext path production +// traffic does, not a shortcut around it. +func contextWithBearerToken(token string) context.Context { + return metadata.NewIncomingContext(context.Background(), metadata.Pairs("authorization", "Bearer "+token)) } func TestPushBatchAssignsRecordID(t *testing.T) { @@ -124,3 +155,101 @@ func TestPushBatchEmptyRecordsIsANoOp(t *testing.T) { t.Fatalf("expected no batches written for an empty request, got %d", len(fp.written)) } } + +// TestPushBatchNoResolverAttachesNoTenantHeader is the regression test +// for single-tenant deployments' behavior staying unchanged: with no +// TenantResolver configured, records are produced exactly as before -- +// no tenant_id header at all -- even with a bearer token present (it's +// simply never inspected). +func TestPushBatchNoResolverAttachesNoTenantHeader(t *testing.T) { + fp := &fakeProducer{} + s := newTestServer(fp) + + req := &logsv1.PushBatchRequest{Records: []*logsv1.LogRecord{{Host: "h1", Message: "one"}}} + if _, err := s.PushBatch(contextWithBearerToken("irrelevant"), req); err != nil { + t.Fatalf("PushBatch() error = %v", err) + } + + fp.mu.Lock() + defer fp.mu.Unlock() + for _, h := range fp.written[0][0].Headers { + if h.Key == TenantIDHeaderKey { + t.Fatalf("expected no %s header with no resolver configured, got %q", TenantIDHeaderKey, h.Value) + } + } +} + +func TestPushBatchWithResolverAttachesTenantHeader(t *testing.T) { + fp := &fakeProducer{} + resolver := &fakeResolver{tenantByToken: map[string]string{"real-token": "acme"}} + s := newTestServerWithResolver(fp, resolver) + + req := &logsv1.PushBatchRequest{Records: []*logsv1.LogRecord{ + {Host: "h1", Message: "one"}, + {Host: "h1", Message: "two"}, + }} + if _, err := s.PushBatch(contextWithBearerToken("real-token"), req); err != nil { + t.Fatalf("PushBatch() error = %v", err) + } + + fp.mu.Lock() + defer fp.mu.Unlock() + if len(fp.written[0]) != 2 { + t.Fatalf("expected 2 messages written, got %d", len(fp.written[0])) + } + for _, msg := range fp.written[0] { + found := false + for _, h := range msg.Headers { + if h.Key == TenantIDHeaderKey { + found = true + if string(h.Value) != "acme" { + t.Fatalf("%s header = %q, want acme", TenantIDHeaderKey, h.Value) + } + } + } + if !found { + t.Fatalf("expected every record to carry a %s header", TenantIDHeaderKey) + } + } +} + +func TestPushBatchWithResolverRejectsMissingToken(t *testing.T) { + fp := &fakeProducer{} + resolver := &fakeResolver{tenantByToken: map[string]string{"real-token": "acme"}} + s := newTestServerWithResolver(fp, resolver) + + req := &logsv1.PushBatchRequest{Records: []*logsv1.LogRecord{{Host: "h1", Message: "one"}}} + _, err := s.PushBatch(context.Background(), req) // no bearer token in context at all + if status.Code(err) != codes.Unauthenticated { + t.Fatalf("PushBatch() error = %v, want Unauthenticated", err) + } + + fp.mu.Lock() + defer fp.mu.Unlock() + if len(fp.written) != 0 { + t.Fatal("a batch with no bearer token must never reach the producer once a resolver is configured") + } +} + +// TestPushBatchWithResolverRejectsInvalidToken is the fail-closed +// regression test: a resolver configured but a token it doesn't +// recognize must refuse the whole batch, never fall back to "no tenant" +// (which would silently defeat the point of requiring a credential at +// all). +func TestPushBatchWithResolverRejectsInvalidToken(t *testing.T) { + fp := &fakeProducer{} + resolver := &fakeResolver{tenantByToken: map[string]string{"real-token": "acme"}} + s := newTestServerWithResolver(fp, resolver) + + req := &logsv1.PushBatchRequest{Records: []*logsv1.LogRecord{{Host: "h1", Message: "one"}}} + _, err := s.PushBatch(contextWithBearerToken("wrong-token"), req) + if status.Code(err) != codes.Unauthenticated { + t.Fatalf("PushBatch() error = %v, want Unauthenticated", err) + } + + fp.mu.Lock() + defer fp.mu.Unlock() + if len(fp.written) != 0 { + t.Fatal("a batch with an invalid token must never reach the producer once a resolver is configured") + } +} diff --git a/ingest/internal/tenantresolver/tenantresolver.go b/ingest/internal/tenantresolver/tenantresolver.go new file mode 100644 index 0000000..d75fede --- /dev/null +++ b/ingest/internal/tenantresolver/tenantresolver.go @@ -0,0 +1,65 @@ +// Package tenantresolver is ingest's HTTP client for resolving an +// agent-presented ingest credential to a tenant -- calls enterprise- +// auth's POST /internal/authorize-ingest over the network, never +// importing enterprise/ (ingest is AGPL core; enterprise/ is +// commercial-licensed and must never be imported by core code -- same +// "network boundary, not import boundary" shape api/authz.HTTPAuthorizer +// already uses for the query path, and enterprise-auth's own doc +// comment on POST /internal/authorize-ingest). nil (no resolver +// configured) is grpcserver.Server's documented no-op default -- +// single-tenant deployments never construct one, and every record's +// TenantID stays empty, exactly like before per-tenant ingest +// credentials existed. +package tenantresolver + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "time" +) + +type HTTPResolver struct { + baseURL string + http *http.Client +} + +func New(baseURL string) *HTTPResolver { + return &HTTPResolver{baseURL: baseURL, http: &http.Client{Timeout: 3 * time.Second}} +} + +type authorizeIngestResponse struct { + TenantID string `json:"tenant_id"` +} + +// ResolveTenant implements grpcserver.TenantResolver. Forwards only the +// bearer token itself, nothing else about the caller's request -- same +// "forward exactly the credential, never the rest of the request" +// discipline api/authz.HTTPAuthorizer already follows. +func (r *HTTPResolver) ResolveTenant(ctx context.Context, token string) (string, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.baseURL+"/internal/authorize-ingest", nil) + if err != nil { + return "", fmt.Errorf("tenantresolver: building request: %w", err) + } + req.Header.Set("Authorization", "Bearer "+token) + + resp, err := r.http.Do(req) + if err != nil { + return "", fmt.Errorf("tenantresolver: calling enterprise-auth: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("tenantresolver: enterprise-auth returned status %d", resp.StatusCode) + } + + var body authorizeIngestResponse + if err := json.NewDecoder(resp.Body).Decode(&body); err != nil { + return "", fmt.Errorf("tenantresolver: decoding response: %w", err) + } + if body.TenantID == "" { + return "", fmt.Errorf("tenantresolver: enterprise-auth returned an empty tenant_id") + } + return body.TenantID, nil +} diff --git a/ingest/internal/tenantresolver/tenantresolver_test.go b/ingest/internal/tenantresolver/tenantresolver_test.go new file mode 100644 index 0000000..e5af47a --- /dev/null +++ b/ingest/internal/tenantresolver/tenantresolver_test.go @@ -0,0 +1,56 @@ +package tenantresolver + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" +) + +func TestResolveTenantForwardsTokenAndParsesTenantID(t *testing.T) { + var gotAuth string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotAuth = r.Header.Get("Authorization") + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(authorizeIngestResponse{TenantID: "acme"}) + })) + defer srv.Close() + + res := New(srv.URL) + tenantID, err := res.ResolveTenant(context.Background(), "real-token") + if err != nil { + t.Fatalf("ResolveTenant: %v", err) + } + if tenantID != "acme" { + t.Fatalf("tenantID = %q, want acme", tenantID) + } + if gotAuth != "Bearer real-token" { + t.Fatalf("Authorization header = %q, want Bearer real-token", gotAuth) + } +} + +func TestResolveTenantNon2xxIsAnError(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusUnauthorized) + })) + defer srv.Close() + + res := New(srv.URL) + if _, err := res.ResolveTenant(context.Background(), "bad-token"); err == nil { + t.Fatal("expected an error for a 401 response from enterprise-auth") + } +} + +func TestResolveTenantRejectsEmptyTenantID(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(authorizeIngestResponse{}) + })) + defer srv.Close() + + res := New(srv.URL) + if _, err := res.ResolveTenant(context.Background(), "some-token"); err == nil { + t.Fatal("expected an error when enterprise-auth returns an empty tenant_id despite a 200") + } +} diff --git a/metadata/migrations/0034_create_ingest_credentials.sql b/metadata/migrations/0034_create_ingest_credentials.sql new file mode 100644 index 0000000..64f39d1 --- /dev/null +++ b/metadata/migrations/0034_create_ingest_credentials.sql @@ -0,0 +1,17 @@ +-- Per-tenant bearer credentials an agent presents to `ingest` (see +-- ingest/internal/grpcserver's TenantResolver) so a record can be +-- attributed to a tenant at the point it enters the system, rather than +-- landing in the one shared ClickHouse database/Tantivy index every +-- record lands in today. Only the SHA-256 hash of the token is stored -- +-- same reasoning a password gets hashed, not stored raw: enterprise-auth +-- only ever needs to check "does the presented token match," never to +-- recover the plaintext, so there's no reason to keep it recoverable. +-- Losing the plaintext means issuing a new credential, not resetting +-- this one -- the plaintext is returned exactly once, at creation. +CREATE TABLE IF NOT EXISTS ingest_credentials +( + id UUID PRIMARY KEY, + tenant_id TEXT NOT NULL REFERENCES tenants(id) ON DELETE CASCADE, + token_hash TEXT NOT NULL UNIQUE, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +)