From 6018a036d670c248b7840338071e968c580e16ab Mon Sep 17 00:00:00 2001 From: Bretton Date: Thu, 13 Aug 2026 00:34:25 -0700 Subject: [PATCH] =?UTF-8?q?wip(task13):=20cycles=209-10=20=E2=80=94=20conf?= =?UTF-8?q?ig=20vars=20+=20validation,=20host-router=20wiring=20in=20main,?= =?UTF-8?q?=20user-origin=20inbox=20mount,=20federation=20opt-out=20lexico?= =?UTF-8?q?n?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 --- README.md | 46 +++++ cmd/tidepool/main.go | 41 +++- internal/config/config.go | 68 ++++++ internal/config/config_test.go | 207 +++++++++++++++---- internal/ingest/inbox.go | 9 + internal/ingest/inbox_handler_test.go | 74 +++++++ internal/personas/hostrouter.go | 13 +- internal/personas/hostrouter_test.go | 50 +++++ internal/personas/inbox_test.go | 149 +++++++++++++ internal/personas/personas.go | 11 + internal/personas/serving.go | 31 ++- lexicons/MANIFEST.sha256 | 1 + lexicons/federation_test.go | 89 ++++++++ lexicons/social/coves/bridge/federation.json | 25 +++ 14 files changed, 770 insertions(+), 44 deletions(-) create mode 100644 internal/ingest/inbox_handler_test.go create mode 100644 internal/personas/inbox_test.go create mode 100644 lexicons/federation_test.go create mode 100644 lexicons/social/coves/bridge/federation.json diff --git a/README.md b/README.md index fc68397..8d7877f 100644 --- a/README.md +++ b/README.md @@ -261,6 +261,8 @@ production**: | `RELAY_HOSTS` | *(optional)* | comma-separated relays to send `com.atproto.sync.requestCrawl` to on startup (each retried on a bounded budget — the relay calls back into `describeServer` before subscribing, which can race process start); in development the request is logged, never sent, unless `ALLOW_DEV_REQUEST_CRAWL` opts in | | `ALLOW_DEV_REQUEST_CRAWL` | off | dev-only: actually SEND `requestCrawl` to `RELAY_HOSTS` in development (exists for the e2e stack's local BigSky); refused in production, where sending is already the behavior | | `ADMIN_TOKEN` | `dev-admin-token` | bearer token protecting the `/admin` API | +| `AP_USER_ORIGIN` | `http://localhost:8091` | origin the Coves users' ActivityPub actors live under (e.g. `https://coves.social`). Baked into every `actor_id` minted under it, so serving derives URLs from the stored row and never from this value; its host must not be `BRIDGE_HOSTNAME` or a subdomain of it, which would shadow the bridged handle namespace (refused at startup) | +| `AP_HOST_FALLTHROUGH_DEV` | **on** in development | route unknown `Host`s to the bridge surface instead of refusing them with 421. The one dev flag that defaults ON — a laptop is reached by IP or tunnel hostname — and the only posture in production is off: setting it there is refused | | `BACKFILL_MAX_POSTS` | `100` | posts materialized per community backfill run | | `MINT_RATE_PER_MINUTE` / `MINT_BURST` | `60` / `120` | rate gate on inbound DID minting (PLC registrations are forever; unseen authors in delivered content trigger mints) | | `INGEST_WORKERS` | `4` | inbox queue worker-pool size | @@ -495,6 +497,19 @@ The contract is the lexicon at [`lexicons/social/coves/bridge/getVoteAggregates.json`](lexicons/social/coves/bridge/getVoteAggregates.json) and is versioned by nsid: breaking changes ship under a new name. +Its sibling under the same Tidepool-owned namespace is +[`lexicons/social/coves/bridge/federation.json`](lexicons/social/coves/bridge/federation.json) +(`key: literal:self`, one record per repo), the user-facing federation +preference. It is an **opt-OUT**: federation is on by default, so the record's +ABSENCE means enabled and it only ever exists to turn federation down. +`enabled: false` is a soft disable — the actor stops resolving via WebFinger +and stops delivering, while its actor document and already-federated +references stay intact. Adding `deleteRemote: true` escalates to the +destructive tier (ask peers to delete the user's federated content — +irreversible on their side). Deleting the record, or writing +`enabled: true`, restores the default under the SAME actor identity: the local +part is frozen at creation and never re-derived. + Counts reflect each distinct voter's **latest** state — flips (`Like` → `Dislike`) and `Undo`s are folded in, re-delivered activities are deduplicated by activity id. Votes on content the bridge never materialized @@ -515,6 +530,37 @@ to discard a stale update). Every `bridgedStats` write goes through the same lexicon validation and mapping bookkeeping as any other record commit, and a Lemmy edit that rebuilds a record carries an existing `bridgedStats` forward. +## Coves user origin (the `coves.social` AP surface) + +Coves users get ActivityPub identities of their own, served on +`AP_USER_ORIGIN` — a **second origin on the same listener**, distinct from the +bridge's `BRIDGE_HOSTNAME` surface. A `Host` router splits the two: the bridge +hostname and its bridged-handle subdomains (plus `localhost`, bare IPs, and an +absent `Host` — container healthchecks) reach the bridge; the user origin's own +`Host` reaches the user surface; anything else is refused with **421 Misdirected +Request** unless `AP_HOST_FALLTHROUGH_DEV` is on. When both names resolve to one +authority (the dev default, `localhost:8091`), the split falls back to the path: +the user surface answers first and its 404s fall through to the bridge. + +What the user origin serves: + +| Route | Purpose | +|---|---| +| `GET /.well-known/webfinger?resource=acct:alice@…` | discovery for a minted local part, scoped to the routed `Host` | +| `GET /ap/actor/{did}` | the user's `Person` document (`publicKey`, `inbox`, `endpoints.sharedInbox`, `outbox`, `published`) | +| `GET /ap/actor/{did}/outbox` | empty `OrderedCollection` — Lemmy requires the field, and a missing outbox rejects the whole actor | +| `POST /ap/inbox` | shared inbox; dispatched **verbatim** to the existing ingest pipeline (one verification, dedupe, and refusal taxonomy — never a second copy) | +| `GET /` | the origin's instance (`Application`) actor, republishing the bridge's key — Lemmy delivers `Delete{Person}` and other send-to-all-instances activities only to the inbox on that row | +| `GET /.well-known/nodeinfo`, `GET /nodeinfo/2.0` | software identification (`software.name: tidepool`) | + +An actor is minted lazily on a user's first federating interaction. Its local +part is derived once and then **frozen**: a native handle +(`alice.coves.social`) yields `alice`, anything else keeps its full handle +(`bretton.dev` stays `bretton.dev`), collisions take `-2`, `-3`, … , and a +later handle change refreshes only the cached display name, summary, and +avatar. The RSA private key is sealed with `BRIDGE_KEK` (AES-256-GCM, bound to +the DID) and never stored in the clear; only the public PEM is published. + ## Verifying with Jetstream **Automated:** the e2e harness (`make e2e`, above) runs a real Jetstream diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index f930a10..aad83e8 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -9,6 +9,7 @@ import ( "fmt" "log/slog" "net/http" + "net/url" "os" "os/signal" "syscall" @@ -23,6 +24,7 @@ import ( "tidepool/internal/identity" "tidepool/internal/ingest" "tidepool/internal/materialize" + "tidepool/internal/personas" "tidepool/internal/prune" "tidepool/internal/repo" "tidepool/internal/store" @@ -443,9 +445,46 @@ func run(logger *slog.Logger) error { } votesXRPC.Routes(router) + // The Coves user origin (task 13): AP Person actors for Coves users, + // served on AP_USER_ORIGIN's Host. It shares this listener with the + // bridge's own surface, and the Host router below decides which one a + // request belongs to. The inbox is handed the EXISTING ingest handler — + // the user origin publishes a shared inbox but never a second verify + // pipeline. + personasService, err := personas.New(personas.Options{ + DB: database, + Custodian: custodian, + UserOrigin: cfg.APUserOrigin, + ServiceActor: serviceActor, + InboxHandler: inbox.InboxHandler(), + }) + if err != nil { + return fmt.Errorf("user origin: %w", err) + } + + // Host routing wraps everything: the chi router keeps answering for the + // bridge hostname and its bridged-handle subdomains, the user origin + // answers for its own Host, and an unrecognized Host is refused with 421 + // unless AP_HOST_FALLTHROUGH_DEV is on. Both hosts naming one authority + // (the dev default) composes by path instead. + userHost, err := url.Parse(cfg.APUserOrigin) + if err != nil || userHost.Host == "" { + return fmt.Errorf("user origin: AP_USER_ORIGIN %q is not an absolute origin URL", cfg.APUserOrigin) + } + hostRouter, err := personas.NewHostRouter(personas.HostRouterOptions{ + ServiceHost: cfg.BridgeHostname, + ServiceHandler: router, + UserHost: userHost.Host, + UserHandler: personasService, + DevFallthrough: cfg.APHostFallthroughDev, + }) + if err != nil { + return fmt.Errorf("host router: %w", err) + } + server := &http.Server{ Addr: cfg.ListenAddr, - Handler: router, + Handler: hostRouter, ReadHeaderTimeout: readHeaderTimeout, WriteTimeout: writeTimeout, IdleTimeout: idleTimeout, diff --git a/internal/config/config.go b/internal/config/config.go index b1ba1e0..7ff693b 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -8,6 +8,7 @@ import ( "encoding/hex" "fmt" "log/slog" + "net/url" "os" "strconv" "strings" @@ -98,6 +99,20 @@ type Config struct { // mislead (the ALLOW_PRIVATE_FETCH pattern). Set // ALLOW_DEV_REQUEST_CRAWL=1 to enable. AllowDevRequestCrawl bool + // APUserOrigin is the origin Coves users' ActivityPub actors live under + // (AP_USER_ORIGIN, e.g. "https://coves.social"); dev defaults to + // http://localhost:8091. It seeds NEW actor rows only — serving derives + // every URL from the stored actor_id — and its host must not be + // BRIDGE_HOSTNAME or a subdomain of it, which would shadow the bridged + // handle namespace. + APUserOrigin string + // APHostFallthroughDev routes unknown Hosts to the service surface + // instead of refusing them (AP_HOST_FALLTHROUGH_DEV). Unlike the other + // dev flags this one defaults to TRUE in development — a laptop is + // reached by IP or tunnel hostname — and setting it in production is + // refused: an authenticated write surface must not answer under a Host + // an attacker chose. + APHostFallthroughDev bool // AdminToken is the bearer token protecting the /admin API (community // subscribe/unsubscribe/backfill). ADMIN_TOKEN; required in production, // dev default is a fixed, publicly known value. @@ -332,6 +347,33 @@ func Load(logger *slog.Logger) (*Config, error) { return nil, fmt.Errorf("config: ALLOW_DEV_REQUEST_CRAWL must not be set in production (production always sends requestCrawl)") } + // The Coves user origin. Required in production: it is baked into every + // actor_id this deployment mints, so a wrong or missing value is not a + // runtime inconvenience but a set of federated identities pointing at + // the wrong place, forever. + cfg.APUserOrigin, err = stringVar(logger, isDevelopment, "AP_USER_ORIGIN", "http://localhost:8091") + if err != nil { + return nil, err + } + if err := validateUserOrigin(cfg.APUserOrigin, cfg.BridgeHostname); err != nil { + return nil, err + } + + // Unlike every other dev flag this one defaults ON in development: a + // laptop is reached by IP, tunnel hostname, or whatever the tunnel + // minted this morning, and a default-off flag would 421 every local + // request. Production defaults it off and REFUSES it set — an + // authenticated write surface must not answer under a Host an attacker + // chose. + cfg.APHostFallthroughDev, err = boolVarDefault(logger, "AP_HOST_FALLTHROUGH_DEV", isDevelopment) + if err != nil { + return nil, err + } + if cfg.APHostFallthroughDev && !isDevelopment { + return nil, fmt.Errorf("config: AP_HOST_FALLTHROUGH_DEV must not be set in production " + + "(unknown Hosts are refused there)") + } + // Admin API auth: like the KEK, the dev default is fixed and public — // required in production. cfg.AdminToken, err = stringVar(logger, isDevelopment, "ADMIN_TOKEN", "dev-admin-token") @@ -544,6 +586,32 @@ func boolVarDefault(logger *slog.Logger, name string, fallback bool) (bool, erro return false, fmt.Errorf("config: %s must be a boolean (1/0, true/false, yes/no, on/off), got %q", name, raw) } +// validateUserOrigin refuses a user origin that would shadow the bridge's own +// handle namespace. Bridged handles are subdomains of BRIDGE_HOSTNAME resolved +// off r.Host, so a user origin AT that name or UNDER it would swallow them — +// and the Host router could not tell the two surfaces apart in the first +// place. The comparison is on host:port, because a different port is a +// different authority: the dev defaults are exactly that shape +// (BRIDGE_HOSTNAME localhost, user origin on :8091). Matching is on a label +// boundary, so "nottidepool.example" is not under "tidepool.example". +func validateUserOrigin(origin, bridgeHostname string) error { + parsed, err := url.Parse(origin) + if err != nil { + return fmt.Errorf("config: AP_USER_ORIGIN must be an absolute origin URL, got %q: %w", origin, err) + } + if parsed.Scheme == "" || parsed.Host == "" { + return fmt.Errorf("config: AP_USER_ORIGIN must be an absolute origin URL "+ + "(scheme and host), got %q", origin) + } + host := strings.ToLower(parsed.Host) + bridge := strings.ToLower(strings.TrimSpace(bridgeHostname)) + if host == bridge || strings.HasSuffix(host, "."+bridge) { + return fmt.Errorf("config: AP_USER_ORIGIN host %q must not be BRIDGE_HOSTNAME %q "+ + "or a subdomain of it: the bridged handle namespace lives there", host, bridge) + } + return nil +} + // stringVar returns the value of an environment variable. When unset it // falls back to the logged dev default in development and errors in // production. diff --git a/internal/config/config_test.go b/internal/config/config_test.go index dd2af7d..0305708 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -21,11 +21,29 @@ func clearConfigEnv(t *testing.T) { "ADMIN_TOKEN", "BACKFILL_MAX_POSTS", "MINT_RATE_PER_MINUTE", "MINT_BURST", "INGEST_WORKERS", "BRIDGE_SCHEME", "ALLOW_PRIVATE_FETCH", "ALLOW_DEV_REQUEST_CRAWL", "RELAY_HOSTS", + "AP_USER_ORIGIN", "AP_HOST_FALLTHROUGH_DEV", } { t.Setenv(name, "") } } +// setProductionEnv sets every variable production requires, so a test about +// ONE of them can blank exactly that one. Adding a new required variable +// here (AP_USER_ORIGIN was the last) then updates every production test at +// once instead of breaking them one at a time. +func setProductionEnv(t *testing.T) { + t.Helper() + clearConfigEnv(t) + t.Setenv("ENVIRONMENT", EnvironmentProduction) + t.Setenv("DATABASE_URL", "postgres://prod/db") + t.Setenv("LISTEN_ADDR", ":8080") + t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") + t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") + t.Setenv("BRIDGE_KEK", "sfDrM4bIeCJp01ZBTArLPJXNQlD7pcYFsod2An6UAF0=") // base64 form + t.Setenv("ADMIN_TOKEN", "prod-admin-token") + t.Setenv("AP_USER_ORIGIN", "https://coves.social") +} + func TestLoad_DevelopmentDefaults(t *testing.T) { clearConfigEnv(t) @@ -74,14 +92,7 @@ func TestLoad_ProductionRequiresValues(t *testing.T) { } func TestLoad_ProductionWithAllValues(t *testing.T) { - clearConfigEnv(t) - t.Setenv("ENVIRONMENT", EnvironmentProduction) - t.Setenv("DATABASE_URL", "postgres://prod/db") - t.Setenv("LISTEN_ADDR", ":8080") - t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") - t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") - t.Setenv("BRIDGE_KEK", "sfDrM4bIeCJp01ZBTArLPJXNQlD7pcYFsod2An6UAF0=") // base64 form - t.Setenv("ADMIN_TOKEN", "prod-admin-token") + setProductionEnv(t) cfg, err := Load(discardLogger()) require.NoError(t, err) @@ -97,13 +108,8 @@ func TestLoad_ProductionWithAllValues(t *testing.T) { } func TestLoad_ProductionRequiresAdminToken(t *testing.T) { - clearConfigEnv(t) - t.Setenv("ENVIRONMENT", EnvironmentProduction) - t.Setenv("DATABASE_URL", "postgres://prod/db") - t.Setenv("LISTEN_ADDR", ":8080") - t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") - t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") - t.Setenv("BRIDGE_KEK", "sfDrM4bIeCJp01ZBTArLPJXNQlD7pcYFsod2An6UAF0=") + setProductionEnv(t) + t.Setenv("ADMIN_TOKEN", "") _, err := Load(discardLogger()) require.Error(t, err, "production must never run on the public dev-default admin token") @@ -111,12 +117,8 @@ func TestLoad_ProductionRequiresAdminToken(t *testing.T) { } func TestLoad_ProductionRequiresKEK(t *testing.T) { - clearConfigEnv(t) - t.Setenv("ENVIRONMENT", EnvironmentProduction) - t.Setenv("DATABASE_URL", "postgres://prod/db") - t.Setenv("LISTEN_ADDR", ":8080") - t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") - t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") + setProductionEnv(t) + t.Setenv("BRIDGE_KEK", "") _, err := Load(discardLogger()) require.Error(t, err, "production must never run on the public dev-default KEK") @@ -171,14 +173,7 @@ func TestLoad_AllowDevRequestCrawl(t *testing.T) { } func TestLoad_AllowDevRequestCrawlRefusedInProduction(t *testing.T) { - clearConfigEnv(t) - t.Setenv("ENVIRONMENT", EnvironmentProduction) - t.Setenv("DATABASE_URL", "postgres://prod/db") - t.Setenv("LISTEN_ADDR", ":8080") - t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") - t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") - t.Setenv("BRIDGE_KEK", "sfDrM4bIeCJp01ZBTArLPJXNQlD7pcYFsod2An6UAF0=") - t.Setenv("ADMIN_TOKEN", "prod-admin-token") + setProductionEnv(t) t.Setenv("ALLOW_DEV_REQUEST_CRAWL", "1") _, err := Load(discardLogger()) @@ -187,14 +182,7 @@ func TestLoad_AllowDevRequestCrawlRefusedInProduction(t *testing.T) { } func TestLoad_BridgeSchemeHTTPRefusedInProduction(t *testing.T) { - clearConfigEnv(t) - t.Setenv("ENVIRONMENT", EnvironmentProduction) - t.Setenv("DATABASE_URL", "postgres://prod/db") - t.Setenv("LISTEN_ADDR", ":8080") - t.Setenv("BRIDGE_HOSTNAME", "tidepool.example") - t.Setenv("PLC_DIRECTORY_URL", "https://plc.directory") - t.Setenv("BRIDGE_KEK", "sfDrM4bIeCJp01ZBTArLPJXNQlD7pcYFsod2An6UAF0=") - t.Setenv("ADMIN_TOKEN", "prod-admin-token") + setProductionEnv(t) t.Setenv("BRIDGE_SCHEME", "http") _, err := Load(discardLogger()) @@ -203,3 +191,148 @@ func TestLoad_BridgeSchemeHTTPRefusedInProduction(t *testing.T) { assert.Contains(t, err.Error(), "production", "the refusal must come from the production branch, not generic scheme validation") } + +func TestLoad_APUserOrigin(t *testing.T) { + clearConfigEnv(t) + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + assert.Equal(t, "http://localhost:8091", cfg.APUserOrigin, + "the dev default serves user actors off the local listener") + + t.Setenv("AP_USER_ORIGIN", "https://coves.social") + cfg, err = Load(discardLogger()) + require.NoError(t, err) + assert.Equal(t, "https://coves.social", cfg.APUserOrigin) +} + +func TestLoad_APUserOriginRequiredInProduction(t *testing.T) { + setProductionEnv(t) + t.Setenv("AP_USER_ORIGIN", "") + + _, err := Load(discardLogger()) + require.Error(t, err, + "the origin is baked into every minted actor_id; production must state it") + assert.Contains(t, err.Error(), "AP_USER_ORIGIN") +} + +// TestLoad_APUserOriginMustNotShadowBridgeHostname: bridged handles are +// subdomains of BRIDGE_HOSTNAME resolved off r.Host, so a user origin at or +// under that name would swallow the handle namespace — and the Host router +// could not tell the two surfaces apart. +func TestLoad_APUserOriginMustNotShadowBridgeHostname(t *testing.T) { + for _, tc := range []struct { + name string + hostname string + userOrigin string + wantErr bool + }{ + { + name: "a distinct origin is fine", + hostname: "tidepool.example", + userOrigin: "https://coves.social", + }, + { + name: "the bridge hostname itself", + hostname: "tidepool.example", + userOrigin: "https://tidepool.example", + wantErr: true, + }, + { + name: "a subdomain of the bridge hostname", + hostname: "tidepool.example", + userOrigin: "https://users.tidepool.example", + wantErr: true, + }, + { + name: "a deep subdomain of the bridge hostname", + hostname: "tidepool.example", + userOrigin: "https://ap.users.tidepool.example", + wantErr: true, + }, + { + // Label-boundary check, both directions: nottidepool.example is + // not under tidepool.example. + name: "suffix-adjacent name is not a subdomain", + hostname: "tidepool.example", + userOrigin: "https://nottidepool.example", + }, + { + // The dev defaults are exactly this shape: BRIDGE_HOSTNAME + // localhost with the user origin on :8091. A different port is a + // different authority, so the comparison is on host:port. + name: "same name on another port", + hostname: "localhost", + userOrigin: "http://localhost:8091", + }, + { + name: "missing scheme", + hostname: "tidepool.example", + userOrigin: "coves.social", + wantErr: true, + }, + { + name: "not a URL at all", + hostname: "tidepool.example", + userOrigin: "://", + wantErr: true, + }, + } { + t.Run(tc.name, func(t *testing.T) { + clearConfigEnv(t) + t.Setenv("BRIDGE_HOSTNAME", tc.hostname) + t.Setenv("AP_USER_ORIGIN", tc.userOrigin) + + cfg, err := Load(discardLogger()) + if tc.wantErr { + require.Error(t, err, "AP_USER_ORIGIN %q with BRIDGE_HOSTNAME %q must be refused", + tc.userOrigin, tc.hostname) + assert.Contains(t, err.Error(), "AP_USER_ORIGIN") + return + } + require.NoError(t, err) + assert.Equal(t, tc.userOrigin, cfg.APUserOrigin) + }) + } +} + +// TestLoad_APHostFallthroughDev: unlike every other dev flag this one is ON +// by default in development — a laptop is reached by IP or tunnel hostname, +// and a default-off flag would make local runs 421 everything. +func TestLoad_APHostFallthroughDev(t *testing.T) { + clearConfigEnv(t) + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + assert.True(t, cfg.APHostFallthroughDev, "development falls through by default") + + t.Setenv("AP_HOST_FALLTHROUGH_DEV", "0") + cfg, err = Load(discardLogger()) + require.NoError(t, err) + assert.False(t, cfg.APHostFallthroughDev, + "a developer must be able to rehearse the production posture locally") + + t.Setenv("AP_HOST_FALLTHROUGH_DEV", "maybe") + _, err = Load(discardLogger()) + require.Error(t, err, "a default-ON flag disabled by a typo would be invisible") +} + +func TestLoad_APHostFallthroughRefusedInProduction(t *testing.T) { + setProductionEnv(t) + t.Setenv("AP_HOST_FALLTHROUGH_DEV", "1") + + _, err := Load(discardLogger()) + require.Error(t, err, + "an authenticated write surface must not answer under an attacker-chosen Host") + assert.Contains(t, err.Error(), "AP_HOST_FALLTHROUGH_DEV") +} + +// TestLoad_ProductionDefaultsFallthroughOff: leaving the flag unset in +// production must not inherit development's default-ON. +func TestLoad_ProductionDefaultsFallthroughOff(t *testing.T) { + setProductionEnv(t) + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + assert.False(t, cfg.APHostFallthroughDev, "production never falls through") +} diff --git a/internal/ingest/inbox.go b/internal/ingest/inbox.go index dce68b5..8ea4230 100644 --- a/internal/ingest/inbox.go +++ b/internal/ingest/inbox.go @@ -199,6 +199,15 @@ func (ib *Inbox) Routes(r chi.Router) { r.Get("/nodeinfo/2.0", ib.handleNodeInfo) } +// InboxHandler is the delivery handler mounted at /inbox, exported so a +// SECOND origin can serve the same pipeline without duplicating it. The Coves +// user origin's POST /ap/inbox dispatches straight here: signature +// verification, actor binding, dedupe, admission control, and the refusal +// taxonomy are one implementation, and a second copy would drift. +func (ib *Inbox) InboxHandler() http.Handler { + return http.HandlerFunc(ib.handleInbox) +} + // handleInbox receives one AP delivery: verify the HTTP signature, bind the // activity's actor to the signer, dedupe by activity id, enqueue for the // worker pool, 202. Everything heavier happens async — remote instances diff --git a/internal/ingest/inbox_handler_test.go b/internal/ingest/inbox_handler_test.go new file mode 100644 index 0000000..ab72cc4 --- /dev/null +++ b/internal/ingest/inbox_handler_test.go @@ -0,0 +1,74 @@ +package ingest + +import ( + "bytes" + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/ap" +) + +// The Coves user origin (task 13) serves its own /ap/inbox, but it must not +// grow a second delivery pipeline: it dispatches to the handler exported +// here. These tests pin that the exported handler IS the mounted one — same +// verification, same refusal taxonomy, same queueing — because a divergence +// would mean two inboxes with two security postures. + +// TestInboxHandler_MatchesMountedRoute: an unsigned delivery must be refused +// identically through both paths. +func TestInboxHandler_MatchesMountedRoute(t *testing.T) { + h := newHarness(t) + + body := []byte(`{"id":"https://lemmy.world/activities/like/unsigned","type":"Like",` + + `"actor":"https://lemmy.world/u/nobody","object":"` + pageID + `"}`) + newRequest := func(target string) *http.Request { + req := httptest.NewRequest(http.MethodPost, "https://"+bridgeHost+target, bytes.NewReader(body)) + req.Header.Set("Content-Type", ap.ContentTypeActivityJSON) + return req + } + + mounted := httptest.NewRecorder() + h.router.ServeHTTP(mounted, newRequest("/inbox")) + + exported := httptest.NewRecorder() + h.inbox.InboxHandler().ServeHTTP(exported, newRequest("/ap/inbox")) + + assert.Equal(t, http.StatusUnauthorized, mounted.Code, + "an unsigned delivery is a definitive refusal (4xx), not a retryable one") + assert.Equal(t, mounted.Code, exported.Code, + "the exported handler must refuse exactly as the mounted route does") + assert.Equal(t, mounted.Body.String(), exported.Body.String(), + "same refusal, same body: one implementation, two mount points") +} + +// TestInboxHandler_AcceptsSignedDelivery walks the whole pipeline through +// the exported handler — verification, dedupe, enqueue — so the test cannot +// pass on an error path alone. +func TestInboxHandler_AcceptsSignedDelivery(t *testing.T) { + h := newHarness(t) + alice := h.newRemoteActor("https://lemmy.world/u/exported", + person("https://lemmy.world/u/exported", "exported", nil)) + + const activityID = "https://lemmy.world/activities/like/exported" + body, err := json.Marshal(likeActivity(activityID, alice.id)) + require.NoError(t, err) + + // Addressed to the USER origin's inbox path: the handler must not care + // which mount point it was reached through. + req := httptest.NewRequest(http.MethodPost, "https://"+bridgeHost+"/ap/inbox", bytes.NewReader(body)) + req.Header.Set("Content-Type", ap.ContentTypeActivityJSON) + require.NoError(t, alice.signer().SignRequest(req, body)) + + rec := httptest.NewRecorder() + h.inbox.InboxHandler().ServeHTTP(rec, req) + require.Equal(t, http.StatusAccepted, rec.Code, "body=%s", rec.Body.String()) + + _, err = h.events.GetEvent(context.Background(), activityID) + assert.NoError(t, err, "an accepted delivery must be recorded for dedupe and queueing") +} diff --git a/internal/personas/hostrouter.go b/internal/personas/hostrouter.go index f33bfca..82b2053 100644 --- a/internal/personas/hostrouter.go +++ b/internal/personas/hostrouter.go @@ -70,10 +70,15 @@ type hostRouter struct { func (h *hostRouter) ServeHTTP(w http.ResponseWriter, r *http.Request) { host := normalizeHost(r.Host) switch { - case host == h.userHost && h.userHost == h.serviceHost: - // The default dev configuration points BRIDGE_HOSTNAME and - // AP_USER_ORIGIN at the same authority. One Host cannot pick a - // bucket, so the split moves to the path: see serveComposed. + case host == h.userHost && h.isServiceHost(host): + // BOTH buckets accept this authority, so one Host cannot pick one + // and the split moves to the path: see serveComposed. The test is + // "both buckets accept it", not "the two configured strings are + // equal" — the DEFAULT dev configuration sets BRIDGE_HOSTNAME to + // "localhost" and AP_USER_ORIGIN to "http://localhost:8091", two + // different strings naming one listener. Keyed on string equality, + // the user surface would swallow the whole listener and /healthz — + // the first thing a developer hits — would disappear. h.serveComposed(w, r) case host == h.userHost: h.userHandler.ServeHTTP(w, r) diff --git a/internal/personas/hostrouter_test.go b/internal/personas/hostrouter_test.go index 144803f..7783434 100644 --- a/internal/personas/hostrouter_test.go +++ b/internal/personas/hostrouter_test.go @@ -263,6 +263,56 @@ func TestHostRouter_CollidingHosts(t *testing.T) { }) } +// TestHostRouter_DevDefaultAuthorities is the wiring contract for the +// DEFAULT development configuration, where the two hosts are configured +// differently but name the same listener: BRIDGE_HOSTNAME is "localhost" +// while AP_USER_ORIGIN is "http://localhost:8091", so every local request +// arrives as Host "localhost:8091". +// +// That Host matches the user origin exactly AND satisfies the service +// bucket's loopback rule, so composition must key on "both buckets accept +// this authority", not on the two configured strings being equal. Keyed on +// string equality instead, the user surface swallows the whole listener and +// /healthz — the thing a developer hits first — disappears. +func TestHostRouter_DevDefaultAuthorities(t *testing.T) { + const devServiceHost = "localhost" // BRIDGE_HOSTNAME dev default + const devUserHost = "localhost:8091" + + newRouter := func(t *testing.T) (http.Handler, *marker, *marker) { + t.Helper() + service := newMarker("service") + user := newMarker("user", "/xrpc/_health", "/healthz") + router, err := NewHostRouter(HostRouterOptions{ + ServiceHost: devServiceHost, + ServiceHandler: service, + UserHost: devUserHost, + UserHandler: user, + DevFallthrough: true, + }) + require.NoError(t, err) + return router, service, user + } + + t.Run("the user surface answers its own routes", func(t *testing.T) { + router, _, user := newRouter(t) + rec := routeHost(t, router, "http", devUserHost, "/.well-known/webfinger") + assert.Equal(t, http.StatusOK, rec.Code, "body=%s", rec.Body.String()) + assert.Equal(t, "user", rec.Header().Get("X-Handler")) + assert.Equal(t, 1, user.calls) + }) + + for _, path := range []string{"/healthz", "/xrpc/_health"} { + t.Run("the service surface still answers "+path, func(t *testing.T) { + router, service, _ := newRouter(t) + rec := routeHost(t, router, "http", devUserHost, path) + require.Equal(t, http.StatusOK, rec.Code, + "a dev listener must keep serving %s; body=%s", path, rec.Body.String()) + assert.Equal(t, "service", rec.Header().Get("X-Handler")) + assert.Equal(t, 1, service.calls) + }) + } +} + // TestNewHostRouter_RequiresHandlers: a nil handler would nil-panic on the // first request of whichever bucket it was meant to serve. func TestNewHostRouter_RequiresHandlers(t *testing.T) { diff --git a/internal/personas/inbox_test.go b/internal/personas/inbox_test.go new file mode 100644 index 0000000..076f4da --- /dev/null +++ b/internal/personas/inbox_test.go @@ -0,0 +1,149 @@ +package personas + +import ( + "bytes" + "database/sql" + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/ap" + "tidepool/internal/identity" +) + +// recordingInbox stands in for ingest.Inbox's delivery handler and records +// exactly what reached it. +type recordingInbox struct { + calls int + method string + path string + host string + headers http.Header + body []byte + status int +} + +func (h *recordingInbox) ServeHTTP(w http.ResponseWriter, r *http.Request) { + h.calls++ + h.method = r.Method + h.path = r.URL.Path + h.host = r.Host + h.headers = r.Header.Clone() + var buf bytes.Buffer + _, _ = buf.ReadFrom(r.Body) + h.body = buf.Bytes() + status := h.status + if status == 0 { + status = http.StatusAccepted + } + w.Header().Set("X-Inbox", "reached") + w.WriteHeader(status) + _, _ = w.Write([]byte("delivered")) +} + +func newInboxService(t *testing.T, database *sql.DB, inbox http.Handler) *Service { + t.Helper() + custodian, err := identity.NewCustodian(testKEK) + require.NoError(t, err) + svc, err := New(Options{ + DB: database, + Custodian: custodian, + UserOrigin: userOrigin, + InboxHandler: inbox, + }) + require.NoError(t, err) + require.NotNil(t, svc) + return svc +} + +func postOnUserOrigin(h http.Handler, target string, body []byte, header http.Header) *httptest.ResponseRecorder { + req := httptest.NewRequest(http.MethodPost, "https://"+userHost+target, bytes.NewReader(body)) + req.Host = userHost + req.Header.Set("Content-Type", ap.ContentTypeActivityJSON) + for k, values := range header { + for _, v := range values { + req.Header.Add(k, v) + } + } + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + return rec +} + +// TestServeInbox_DispatchesVerbatim: the user origin publishes a shared +// inbox, but it must not grow a SECOND verify pipeline. The delivery is +// handed to the existing ingest inbox untouched — signature verification, +// authority binding, dedupe, admission control, and the refusal taxonomy +// all stay in one implementation. +func TestServeInbox_DispatchesVerbatim(t *testing.T) { + database := personasTestDB(t) + inbox := &recordingInbox{} + svc := newInboxService(t, database, inbox) + + body := []byte(`{"@context":"https://www.w3.org/ns/activitystreams",` + + `"id":"https://lemmy.world/activities/like/1","type":"Like",` + + `"actor":"https://lemmy.world/u/alice","object":"https://lemmy.world/post/1"}`) + rec := postOnUserOrigin(svc, "/ap/inbox", body, + http.Header{"Signature": []string{`keyId="https://lemmy.world/u/alice#main-key"`}}) + + require.Equal(t, 1, inbox.calls, "POST /ap/inbox must reach the ingest inbox") + assert.Equal(t, http.MethodPost, inbox.method) + assert.Equal(t, "/ap/inbox", inbox.path) + assert.Equal(t, userHost, inbox.host, "the routed Host must survive the handoff") + assert.Equal(t, body, inbox.body, "the body must arrive byte-identical: the digest covers it") + assert.Equal(t, `keyId="https://lemmy.world/u/alice#main-key"`, inbox.headers.Get("Signature"), + "the signature headers must not be rewritten on the way through") + assert.Equal(t, ap.ContentTypeActivityJSON, inbox.headers.Get("Content-Type")) + + // The inbox's own response is what the remote sees — including the + // refusal taxonomy, which Lemmy reads to decide retry vs drop. + assert.Equal(t, http.StatusAccepted, rec.Code) + assert.Equal(t, "reached", rec.Header().Get("X-Inbox")) + assert.Equal(t, "delivered", rec.Body.String()) +} + +// TestServeInbox_PropagatesRefusals: a 503 (retryable) must not be flattened +// into anything else on the way out, or Lemmy drops the delivery forever. +func TestServeInbox_PropagatesRefusals(t *testing.T) { + database := personasTestDB(t) + inbox := &recordingInbox{status: http.StatusServiceUnavailable} + svc := newInboxService(t, database, inbox) + + rec := postOnUserOrigin(svc, "/ap/inbox", []byte(`{"type":"Like"}`), nil) + assert.Equal(t, http.StatusServiceUnavailable, rec.Code, + "the ingest inbox owns the status; the user origin only routes") +} + +// TestServeInbox_Unconfigured: without an inbox handler the origin would be +// advertising an inbox it cannot serve. 404 is the honest answer. +func TestServeInbox_Unconfigured(t *testing.T) { + database := personasTestDB(t) + svc := newInboxService(t, database, nil) + + rec := postOnUserOrigin(svc, "/ap/inbox", []byte(`{"type":"Like"}`), nil) + assert.Equal(t, http.StatusNotFound, rec.Code) +} + +// TestServeInbox_RejectsGET: an inbox is write-only. The method check comes +// first, so the answer does not leak whether an inbox is wired up. +func TestServeInbox_RejectsGET(t *testing.T) { + database := personasTestDB(t) + + for _, tc := range []struct { + name string + inbox http.Handler + }{ + {"configured", &recordingInbox{}}, + {"unconfigured", nil}, + } { + t.Run(tc.name, func(t *testing.T) { + svc := newInboxService(t, database, tc.inbox) + rec := serveOnUserOrigin(svc, http.MethodGet, "/ap/inbox", nil) + assert.Equal(t, http.StatusMethodNotAllowed, rec.Code) + assert.Equal(t, http.MethodPost, rec.Header().Get("Allow")) + }) + } +} diff --git a/internal/personas/personas.go b/internal/personas/personas.go index 1158c44..895a7db 100644 --- a/internal/personas/personas.go +++ b/internal/personas/personas.go @@ -10,6 +10,7 @@ import ( "database/sql" stderrors "errors" "fmt" + "net/http" "net/url" "strconv" "strings" @@ -46,6 +47,12 @@ type Options struct { // from UserOrigin; only Key and CreatedAt are read. Optional: without // it the origin apex serves no instance actor (Lemmy tolerates that). ServiceActor *ap.ServiceActor + // InboxHandler receives POST /ap/inbox. The user origin does NOT run a + // second verify pipeline: deliveries go verbatim to the ingest inbox + // that already does signature verification, authority binding, dedupe, + // and queueing (ingest.Inbox.InboxHandler). Nil means the origin + // advertises an inbox it cannot serve, so the route 404s. + InboxHandler http.Handler } // Service mints and serves Coves user actors. @@ -61,6 +68,9 @@ type Service struct { // serviceActor is the bridge identity the origin apex republishes. Nil // means the apex publishes nothing. serviceActor *ap.ServiceActor + // inboxHandler is the ingest inbox this origin's shared inbox dispatches + // to. Nil means the route 404s. + inboxHandler http.Handler } // New builds a Service. UserOrigin is parsed once here: the host it yields @@ -78,6 +88,7 @@ func New(opts Options) (*Service, error) { userHost: host, serviceActor: opts.ServiceActor, + inboxHandler: opts.InboxHandler, }, nil } diff --git a/internal/personas/serving.go b/internal/personas/serving.go index 606f487..7eea081 100644 --- a/internal/personas/serving.go +++ b/internal/personas/serving.go @@ -63,6 +63,13 @@ func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) { return } s.handleNodeInfo(w, r) + case path == inboxPath: + // The method check comes FIRST, so a GET learns only that an inbox + // is write-only — never whether one is wired up here. + if !requireMethod(w, r, http.MethodPost) { + return + } + s.handleInbox(w, r) case strings.HasPrefix(path, actorPathPrefix): rest := strings.TrimPrefix(path, actorPathPrefix) if !isGET(w, r) { @@ -83,14 +90,34 @@ func (s *Service) ServeHTTP(w http.ResponseWriter, r *http.Request) { } func isGET(w http.ResponseWriter, r *http.Request) bool { - if r.Method != http.MethodGet { - w.Header().Set("Allow", http.MethodGet) + return requireMethod(w, r, http.MethodGet) +} + +func requireMethod(w http.ResponseWriter, r *http.Request, method string) bool { + if r.Method != method { + w.Header().Set("Allow", method) http.Error(w, "method not allowed", http.StatusMethodNotAllowed) return false } return true } +// handleInbox hands the delivery to the ingest inbox VERBATIM — same request, +// same body, same headers. The user origin publishes a shared inbox but must +// not grow a second verify pipeline: signature verification, actor binding, +// dedupe, admission control, and the refusal taxonomy Lemmy reads to decide +// retry-vs-drop all stay in one implementation, and the inbox's own response +// is what the remote sees. +func (s *Service) handleInbox(w http.ResponseWriter, r *http.Request) { + if s.inboxHandler == nil { + // Advertising an inbox this origin cannot serve; 404 is the honest + // answer. + http.NotFound(w, r) + return + } + s.inboxHandler.ServeHTTP(w, r) +} + // handleActorDocument serves the Person document. A DISABLED actor still // serves it: disabling removes an actor from discovery, not from the network, // and already-federated references to it must not become dangling. Paused diff --git a/lexicons/MANIFEST.sha256 b/lexicons/MANIFEST.sha256 index 1b37039..c1051cb 100644 --- a/lexicons/MANIFEST.sha256 +++ b/lexicons/MANIFEST.sha256 @@ -21,6 +21,7 @@ d26df0b986f1e21b2c13c1905c60257ae19e79e812a474b6d032c083939eff1c social/coves/a d3e077e34c9b9ccd8a8892148fe4ce7ab1de65873250a8906aec510d040bfa4c social/coves/aggregator/revokeApiKey.json 0ca7c339793fc8063a312bf39619983d98414fba498b05fee92c135c7ebea0fd social/coves/aggregator/service.json d688327a711491aeb58b75187f468f213d270adad0c0ca59b96fbef0229cd3f4 social/coves/aggregator/updateConfig.json +5a0ee9495bb2c0949f960a1d080af9f5f23cf696aea72187e02cdbeedd069151 social/coves/bridge/federation.json 020b4a33837455e17e1b0e258304f11241931d14929be099f33ca67a88fc2f49 social/coves/bridge/getVoteAggregates.json 94ef1e9a8c6c879697a787f53a48ebeed7edc3a05a90926991667df5a67008f5 social/coves/community/acceptance.json 88fb6259698d0097995200ed5d3a7887165f9cd29907c1c60880150c6568aebf social/coves/community/block.json diff --git a/lexicons/federation_test.go b/lexicons/federation_test.go new file mode 100644 index 0000000..bd78d79 --- /dev/null +++ b/lexicons/federation_test.go @@ -0,0 +1,89 @@ +package lexicons + +import ( + "testing" + + "github.com/bluesky-social/indigo/atproto/atdata" + "github.com/bluesky-social/indigo/atproto/lexicon" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// federationType is the opt-OUT record of decision 11: absence means +// federation is ON for native users, so the record only ever exists to turn +// it down. Coves' settings UI writes it; the bridge reads it. +const federationType = "social.coves.bridge.federation" + +// validate runs a record through the same path the materializer uses +// (atdata.UnmarshalJSON then lexicon.ValidateRecord), so a shape that +// validates here validates in production. +func validate(t *testing.T, raw string) error { + t.Helper() + catalog, err := Catalog() + require.NoError(t, err) + record, err := atdata.UnmarshalJSON([]byte(raw)) + require.NoError(t, err, "fixture must be valid JSON in the atproto data model") + return lexicon.ValidateRecord(catalog, record, federationType, lexicon.ValidateFlags(0)) +} + +// TestFederationLexiconLoads: the vendored file must be in the embedded +// catalog at all. A lexicon that ships but is not embedded fails silently — +// nothing validates against it. +func TestFederationLexiconLoads(t *testing.T) { + catalog, err := Catalog() + require.NoError(t, err) + + schema, err := catalog.Resolve(federationType) + require.NoError(t, err, "the embedded catalog must resolve %s", federationType) + require.NotNil(t, schema) + + def, ok := schema.Def.(lexicon.SchemaRecord) + require.True(t, ok, "main must be a record def, got %T", schema.Def) + // One federation preference per repo. The wire form of a fixed key is + // "literal:self" — indigo accepts only tid/nsid/any/literal:*, and the + // vendored lexicons spell it that way (community/profile.json, rules.json). + assert.Equal(t, "literal:self", def.Key) +} + +func TestFederationRecordShapes(t *testing.T) { + t.Run("soft disable", func(t *testing.T) { + assert.NoError(t, validate(t, + `{"$type":"`+federationType+`","enabled":false}`)) + }) + + t.Run("destructive tier", func(t *testing.T) { + assert.NoError(t, validate(t, + `{"$type":"`+federationType+`","enabled":false,"deleteRemote":true}`)) + }) + + t.Run("re-enabled", func(t *testing.T) { + assert.NoError(t, validate(t, + `{"$type":"`+federationType+`","enabled":true}`)) + }) + + t.Run("enabled is required", func(t *testing.T) { + err := validate(t, `{"$type":"`+federationType+`","deleteRemote":true}`) + require.Error(t, err, + "the opt-out must state its intent: a record with no enabled field is ambiguous") + assert.Contains(t, err.Error(), "enabled") + }) + + t.Run("enabled must be a boolean", func(t *testing.T) { + assert.Error(t, validate(t, `{"$type":"`+federationType+`","enabled":"false"}`), + `the string "false" is truthy in most languages; the lexicon must reject it`) + }) + + t.Run("deleteRemote must be a boolean", func(t *testing.T) { + assert.Error(t, validate(t, + `{"$type":"`+federationType+`","enabled":false,"deleteRemote":"yes"}`)) + }) + + t.Run("unknown fields are tolerated", func(t *testing.T) { + // Pinning ACTUAL behavior, verified against indigo: record objects + // are open unless a lexicon closes them, and none of the Coves + // lexicons do. A future field added upstream must not make today's + // bridge reject the record. + assert.NoError(t, validate(t, + `{"$type":"`+federationType+`","enabled":false,"futureField":"whatever"}`)) + }) +} diff --git a/lexicons/social/coves/bridge/federation.json b/lexicons/social/coves/bridge/federation.json new file mode 100644 index 0000000..48fe4db --- /dev/null +++ b/lexicons/social/coves/bridge/federation.json @@ -0,0 +1,25 @@ +{ + "lexicon": 1, + "id": "social.coves.bridge.federation", + "defs": { + "main": { + "type": "record", + "description": "A Coves user's ActivityPub federation preference, read by the Tidepool bridge. This record is an OPT-OUT: federation is on by default and its ABSENCE means enabled, so the record only ever exists to turn federation down. enabled=false is a soft disable — the user's AP actor stops resolving via WebFinger and stops delivering, while the actor document and its already-federated references stay intact. Adding deleteRemote=true escalates to the destructive tier, asking peers to delete the user's federated content. Re-enabling is either deleting this record or writing enabled=true; both restore the default-on state under the SAME actor identity, because the local part is frozen at creation and never re-derived.", + "key": "literal:self", + "record": { + "type": "object", + "required": ["enabled"], + "properties": { + "enabled": { + "type": "boolean", + "description": "Whether the user's content federates to the fediverse. Stated explicitly: a record with no enabled field is ambiguous about its own intent." + }, + "deleteRemote": { + "type": "boolean", + "description": "When disabling, also ask remote instances to delete the user's already-federated content. Destructive and irreversible on the remote side — peers that honor it cannot restore what they dropped. Meaningful only alongside enabled=false; omitted or false leaves remote copies in place." + } + } + } + } + } +} -- 2.51.2