diff --git a/docs/PRD_ADMIN_MODERATION.md b/docs/PRD_ADMIN_MODERATION.md index b2c45c5..f020fae 100644 --- a/docs/PRD_ADMIN_MODERATION.md +++ b/docs/PRD_ADMIN_MODERATION.md @@ -34,7 +34,7 @@ This is label-integrated moderation, rather than a local deletion feature with f - Another instance must be able to opt into Coves.social's moderation decisions. - **A removed post's thread returns `NotFound`** (2026-09-20). Its comments stay in their authors' repositories but are not reachable through the thread endpoint or profile activity. Preserved discussion behind a redacted header is later scope and would need a versioned thread endpoint. - **Removed comments reuse the existing deleted-comment placeholder** (2026-09-20): `record` absent, `isDeleted`, `deletionReason: moderator`, plus optional moderation attribution. No versioned comment endpoint. The published `commentView` is corrected in place under an approved evolution-rule exception (section 14.4). -- **Removal covers Coves-served media bytes** (2026-09-20), on cache hits as well as misses, with the cached bytes purged. A removal blocks the record's (owner DID, blob CID) pair; a removal with reason `illegal-content` also blocks that CID for every owner, so a re-upload from another account is not served. +- **Removal covers Coves-served media bytes** (2026-09-20), on cache hits as well as misses, with the cached bytes purged. A removal blocks the record's (owner DID, blob CID) pair; a removal with reason `illegal-content` also blocks that CID for every owner, so a re-upload from another account is not served. Decided 2026-09-28 (Q-I6): an image the author adds to the subject after an `illegal-content` removal (by an edit, a recreate, or a create the consumer does not index) gets an owner-scoped block only; every-owner blocks cover only the CIDs indexed when the admin acted. Restoring the decision lifts both. - **Simultaneous author deletion and removal shows the author-deleted placeholder** (2026-09-20); moderation history is retained. - **The public modlog names the acting admin** (2026-09-20). The API carries the actor DID (`#actorRef`); clients resolve and display the handle. - **No admin inspection or retained-evidence storage in this release** (2026-09-20). No `getSubject` endpoint; an admin reviewing a restore reads the current record from its PDS. @@ -348,7 +348,7 @@ Proposed API family under `social.coves.moderation.*`: remove content, restore a Removing a record reference alone does not prevent someone replaying its cached image URL. The image proxy currently addresses content by DID/blob CID, so implementation must inventory every Coves-controlled media/cache route before claiming complete suppression. -Decided 2026-09-20: removal covers Coves-served media bytes, enforced on cache hits as well as misses, with cached bytes purged. Media served directly from a third-party CDN (Bluesky embeds) is outside Coves' control; removal stops Coves from handing out those URLs. Shared-blob semantics are decided too: block the (owner DID, blob CID) pair for every removal, and block the CID for every owner when the reason is `illegal-content`. "Owner DID" means the repository that holds the blob, which is not always the human author: legacy bridged posts live in a Tidepool community repository, so several authors' legacy posts can share one community-DID/CID pair and a pair block then suppresses that image on all of them. That consequence is accepted. The media inventory and acceptance cases must cover both legacy community-owned blobs and author-owned `postv2` blobs. Restoring the decision lifts its blocks. Test both rules, including warm-cache requests, before rollout. +Decided 2026-09-20: removal covers Coves-served media bytes, enforced on cache hits as well as misses, with cached bytes purged. Media served directly from a third-party CDN (Bluesky embeds) is outside Coves' control; removal stops Coves from handing out those URLs. Shared-blob semantics are decided too: block the (owner DID, blob CID) pair for every removal, and block the CID for every owner when the reason is `illegal-content`. Decided 2026-09-28 (Q-I6): the every-owner block covers only the CIDs indexed when the admin acted; an image the author adds after the removal (by an edit, a recreate, or a create the consumer does not index) gets an owner-scoped block only, so an author cannot get a copy of someone else's image blocked for every owner. "Owner DID" means the repository that holds the blob, which is not always the human author: legacy bridged posts live in a Tidepool community repository, so several authors' legacy posts can share one community-DID/CID pair and a pair block then suppresses that image on all of them. That consequence is accepted. The media inventory and acceptance cases must cover both legacy community-owned blobs and author-owned `postv2` blobs. Restoring the decision lifts its blocks. Test both rules, including warm-cache requests, before rollout. Raw evidence retention and admin inspection are separate from public serving. Avoid copying content into the action ledger; choose a retention policy before implementing retained-evidence storage. A moderation restore cannot reverse an independent legal/storage purge. diff --git a/docs/PRD_CSAM_SCANNING.md b/docs/PRD_CSAM_SCANNING.md index ffc832f..a445272 100644 --- a/docs/PRD_CSAM_SCANNING.md +++ b/docs/PRD_CSAM_SCANNING.md @@ -1,7 +1,7 @@ # PRD: CSAM Scanning via Cloudflare + Media Choke Point -**Status:** Workstreams 1 and 2 implemented (AppView + Caddy). Workstream 3 (Cloudflare zone config) is manual dashboard/DNS work, not yet done. Workstream 4 (takedown runbook) not started. -**Last updated:** 2026-07-27 +**Status:** Workstreams 1 and 2 implemented (AppView + Caddy). Workstream 3 (Cloudflare zone config) is manual dashboard/DNS work: as of 2026-09-28 Cloudflare proxies `img.coves.social`, and no Cloudflare API token with the Cache Purge permission exists yet. The other WS3 steps (cache rule, CSAM Scanning Tool, SSL mode) are not confirmed here. Workstream 4 (takedown runbook) not started. +**Last updated:** 2026-09-29 ## Problem @@ -21,7 +21,7 @@ Cloudflare can only scan what is **proxied through Cloudflare and cached at its - `tdpl.io` **cannot** be CDN-proxied at all: on-demand TLS for bridged-handle certs requires DNS pointing directly at the origin (`Caddyfile` catch-all block), and Cloudflare wildcard proxying doesn't cover `*.*.tdpl.io` anyway. - `pds.coves.me` serves the atproto sync surface (firehose WebSockets, relay traffic) — proxying it through Cloudflare is possible but risky and unnecessary. -However, we already have the right choke point built: the **image proxy** (`internal/core/imageproxy/`, route `GET /img/{preset}/plain/{did}/{cid}`). It resolves *any* DID to its PDS (including the bridge PDS), fetches the blob, transforms it, and serves it with `Cache-Control: public, max-age=31536000, immutable` + ETag — ideal for edge caching. The presets registry already includes `content_preview`, `content_full`, and `embed_thumbnail`, not just avatars/banners. +However, we already have the right choke point built: the **image proxy** (`internal/core/imageproxy/`, route `GET /img/{preset}/plain/{did}/{cid}`). It resolves *any* DID to its PDS (including the bridge PDS), fetches the blob, transforms it, and serves it with `Cache-Control: public, max-age=86400` + ETag — cacheable at the edge for a day, which bounds how long an unpurged edge or browser keeps an image after a moderation removal. The presets registry already includes `content_preview`, `content_full`, and `embed_thumbnail`, not just avatars/banners. **Decision:** Do NOT put the whole site (or the PDS, or tdpl.io) behind Cloudflare. Instead: @@ -59,7 +59,7 @@ The `img.coves.social` site block is in the production `Caddyfile`: `/img/*` rev Also done: - CSP: the `coves.social` `img-src` is now `'self' data: https://img.coves.social`. -- Error responses from the proxy carry `Cache-Control: no-store` (`writeErrorResponse` in `internal/api/handlers/imageproxy/handler.go`), so a transient PDS timeout or an unpropagated DID can't be pinned at the edge for the year the success path advertises. +- Error responses from the proxy carry `Cache-Control: no-store` (`writeErrorResponse` in `internal/api/handlers/imageproxy/handler.go`), so a transient PDS timeout, an unpropagated DID or a moderation block can't be pinned at the edge or in a browser. Remaining at deploy time: **the bind-mount trap** — Caddyfile changes require `docker compose up -d --force-recreate caddy`, not just a `git pull` + reload. @@ -68,7 +68,7 @@ Remaining at deploy time: **the bind-mount trap** — Caddyfile changes require On the `coves.social` zone (we already own it — DNS-01 tokens exist): 1. **DNS**: `img.coves.social` A/AAAA → OVH origin IP, **Proxied** (orange cloud). All other records stay DNS-only (grey) — especially anything under `tdpl.io` and `coves.me`. -2. **Cache**: add a Cache Rule for `img.coves.social/*`: *Eligible for cache*, respect origin `Cache-Control`. Blobs are content-addressed (CID in URL) so immutable caching is correct. Optionally enable Tiered Cache. +2. **Cache**: add a Cache Rule for `img.coves.social/*`: *Eligible for cache*, with Edge TTL set to respect origin `Cache-Control` headers, and set the zone's Browser Cache TTL to *Respect Existing Headers*. Blobs are content-addressed (CID in URL), but caching is no longer immutable: a moderation removal must stop serving an image, so the origin advertises `public, max-age=86400` and an overriding edge or browser TTL would outlive that bound. Optionally enable Tiered Cache. 3. **Enable CSAM Scanning Tool**: Dashboard → Caching → Configuration → CSAM Scanning Tool → Configure. Provide a monitored role address (e.g. `abuse@coves.social`, forwarded to admins) and verify it. Agree to the service-specific terms. 4. **SSL mode**: Full (strict) for the zone (origin has valid certs via Caddy). 5. Do **not** enable Cloudflare features that interfere with API semantics on other hostnames — only `img` is proxied, so blast radius is zero. @@ -104,6 +104,8 @@ Phase 1 can be a documented manual runbook using existing tools (psql, PDS admin | Direct PDS `getBlob` remains publicly fetchable | Required by atproto sync (relays, other AppViews) | API no longer emits these URLs; optionally rate-limit `getBlob` at Caddy for non-relay UAs | | Only known-hash CSAM is detected | Fuzzy hash lists can't catch novel content | Community reporting (`internal/core/adminreports/`) + moderator review remain the backstop | | Video blobs unscanned | Image proxy is stills-only | Track as separate workstream | +| Browser caches, and any shared cache nobody purges, keep a removed image for up to one day | A moderation removal purges only the proxy's own disk cache; a copy already served with `Cache-Control: public, max-age=86400` stays valid wherever it was stored | The one-day `max-age` bounds the exposure without any configuration, including for self-hosters with no CDN; the optional CDN purge on removal, once built, will shorten it for an edge that is configured for it | +| Responses cached under the pre-deploy `public, max-age=31536000, immutable` header | Copies stored before the one-day header shipped keep their original lifetime | Browsers keep them for up to a year and nothing server-side can reach them; the edge keeps them until a one-time Purge Everything of the `coves.social` zone after deploy (Rollout order step 4; WS3 proxies only `img.coves.social`, so this drops only cached images) | | Bridge PDS stores blobs regardless of scanning | Blobs land before any serve-time scan | Phase 2 ingest scanning; instance allow/blocklist at the bridge is the coarse control | | `record.embed` still carries blob references | Post and comment responses include the verbatim atproto record, whose embed is unprojected by design (the lexicon calls it verbatim). A client *could* build a `getBlob` URL from it | Neither client reads `record.embed` today. Coves image URLs use the proxy, with the explicit foreign Bluesky CDN exception below; this does not mean no blob reference reaches a client. The same record bytes are public on the PDS regardless. Revisit if a client starts reading it | | Resolved Bluesky post images, avatars, and link-preview thumbnails | Approved direct-CDN exception, 2026-09-19: `social.coves.embed.post` resolution preserves validated HTTPS `cdn.bsky.app` image URLs from Bluesky's resolved views, including media on the one-level quoted post. These foreign images bypass the Coves image proxy and Coves scanning edge; URL validation is not content scanning. No PDS blob fetch or CDN-to-blob-reference conversion is used. Videos are outside this change | Fetch and serving boundaries reject unsafe media URLs while preserving post and external-card metadata; see `internal/core/blueskypost/cdn_url.go` and `projection.go`. Backend implementation only: this records the accepted scanning exception, not production deployment or security verification | @@ -113,11 +115,12 @@ Phase 1 can be a documented manual runbook using existing tools (psql, PDS admin WS1 and WS2 are both in the tree, so they ship together. The ordering constraint that remains is **DNS before deploy**: the AppView will start emitting `https://img.coves.social/...` URLs the moment it boots with the new config, so that hostname has to resolve and serve first or every image 404s. 1. **DNS + Cloudflare (WS3)** — create the `img.coves.social` A/AAAA record pointing at the OVH origin, **Proxied** (orange cloud). Every other record stays DNS-only, especially `tdpl.io` and `coves.me`. Set the zone to Full (strict). Add the cache rule for `img.coves.social/*`. Enable the CSAM Scanning Tool with a verified role address. -2. **Deploy Caddy** — `docker compose up -d --force-recreate caddy` (bind-mount trap). Verify `curl -sD- -o /dev/null https://img.coves.social/img/avatar/plain//` returns 200 with `Cache-Control: public, max-age=31536000, immutable`, and that `https://img.coves.social/` 404s. (Use `-sD- -o /dev/null`, not `-I`: the image route is registered GET-only and chi does not map HEAD to it, so `-I` returns 405 and shows none of the cache headers.) +2. **Deploy Caddy** — `docker compose up -d --force-recreate caddy` (bind-mount trap). Verify `curl -sD- -o /dev/null https://img.coves.social/img/avatar/plain//` returns 200 with `Cache-Control: public, max-age=86400`, and that `https://img.coves.social/` 404s. (Use `-sD- -o /dev/null`, not `-I`: the image route is registered GET-only and chi does not map HEAD to it, so `-I` returns 405 and shows none of the cache headers.) 3. **Deploy the AppView** with `IMAGE_PROXY_BASE_URL=https://img.coves.social`. Startup now fails loudly on a misconfigured proxy rather than silently falling back. Verify feeds render, then watch proxy cache hit rate and origin bandwidth. -4. **Client follow-ups** — ship the two `coves-frontend` type/`extractEmbedUrl` changes noted in WS1. -5. WS4 runbook written and dry-run before announcing Lemmy federation more broadly. -6. Phase 2 (ingest-time hash matching) scheduled after federation traffic is real. +4. **Purge the edge once** — after the first AppView deploy that sends `Cache-Control: public, max-age=86400`, run Purge Everything on the `coves.social` zone in the Cloudflare dashboard. Edge copies stored under the old `public, max-age=31536000, immutable` header otherwise keep serving for up to a year, including images removed since. Browser copies cannot be purged (WS5). +5. **Client follow-ups** — ship the two `coves-frontend` type/`extractEmbedUrl` changes noted in WS1. +6. WS4 runbook written and dry-run before announcing Lemmy federation more broadly. +7. Phase 2 (ingest-time hash matching) scheduled after federation traffic is real. ## Open questions diff --git a/internal/api/handlers/imageproxy/avatar_serving_test.go b/internal/api/handlers/imageproxy/avatar_serving_test.go index b5e2725..e528328 100644 --- a/internal/api/handlers/imageproxy/avatar_serving_test.go +++ b/internal/api/handlers/imageproxy/avatar_serving_test.go @@ -139,14 +139,12 @@ func TestImageProxy_ServesRealPDSAvatar(t *testing.T) { assertImageSize(t, body, 360, 360) }) - t.Run("the response is immutably cacheable", func(t *testing.T) { + t.Run("the response is cacheable for one day", func(t *testing.T) { resp, _ := fetch(t, url, nil) - // A preset plus a content-addressed CID names bytes that can never - // change, so the response is safe to cache forever — that is the whole - // economic argument for the proxy, and a weakened header here would - // quietly send every view back to the PDS. - assert.Equal(t, "public, max-age=31536000, immutable", resp.Header.Get("Cache-Control")) + // The CID identifies stable bytes, but moderation can later block their + // serving. Without a purge, a browser or CDN must revalidate within a day. + assert.Equal(t, "public, max-age=86400", resp.Header.Get("Cache-Control")) assert.Equal(t, fmt.Sprintf(`"avatar_small-%s"`, avatar.avatarCID), resp.Header.Get("ETag")) }) diff --git a/internal/api/handlers/imageproxy/blocked_status_test.go b/internal/api/handlers/imageproxy/blocked_status_test.go index fc9cb1f..8faf01d 100644 --- a/internal/api/handlers/imageproxy/blocked_status_test.go +++ b/internal/api/handlers/imageproxy/blocked_status_test.go @@ -84,5 +84,5 @@ func TestHandler_HandleImage_BlockedIsIndistinguishableFromAFailedFetch(t *testi assert.Equalf(t, "no-store", blocked.Header().Get("Cache-Control"), "a refusal must stay uncacheable like every other error on this route: it sits behind a CDN "+ - "whose success responses advertise a one-year immutable lifetime") + "whose success responses can be cached for one day; a cached refusal could outlive its cause") } diff --git a/internal/api/handlers/imageproxy/cache_headers_test.go b/internal/api/handlers/imageproxy/cache_headers_test.go new file mode 100644 index 0000000..09a1131 --- /dev/null +++ b/internal/api/handlers/imageproxy/cache_headers_test.go @@ -0,0 +1,132 @@ +package imageproxy + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + "time" + + "Coves/internal/core/imageproxy" + + "github.com/go-chi/chi/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestHandler_RoutedImageCacheHeaders(t *testing.T) { + const path = "/img/avatar/plain/" + validTestDID + "/" + validTestCID + const cacheControl = "public, max-age=86400" + const etag = `"avatar-` + validTestCID + `"` + + route := func(service Service) http.Handler { + router := chi.NewRouter() + router.Get("/img/{preset}/plain/{did}/{cid}", NewHandler(service, resolverForPDS("https://pds.example.com")).HandleImage) + return router + } + request := func(router http.Handler, ifNoneMatch string) *httptest.ResponseRecorder { + req := httptest.NewRequest(http.MethodGet, path, nil) + if ifNoneMatch != "" { + req.Header.Set("If-None-Match", ifNoneMatch) + } + response := httptest.NewRecorder() + router.ServeHTTP(response, req) + return response + } + + t.Run("200 and matching 304 advertise the same one-day policy and ETag", func(t *testing.T) { + service := &mockService{getImageFunc: func(context.Context, string, string, string, string) ([]byte, error) { + return []byte("image bytes"), nil + }} + router := route(service) + served := request(router, "") + require.Equal(t, http.StatusOK, served.Code) + assert.Equal(t, cacheControl, served.Header().Get("Cache-Control")) + assert.NotContains(t, served.Header().Get("Cache-Control"), "s-maxage") + assert.NotContains(t, served.Header().Get("Cache-Control"), "immutable") + assert.Equal(t, etag, served.Header().Get("ETag")) + + conditional := request(router, served.Header().Get("ETag")) + require.Equal(t, http.StatusNotModified, conditional.Code) + assert.Equal(t, cacheControl, conditional.Header().Get("Cache-Control")) + assert.Equal(t, served.Header().Get("Cache-Control"), conditional.Header().Get("Cache-Control")) + assert.Equal(t, served.Header().Get("ETag"), conditional.Header().Get("ETag")) + }) + + for _, test := range []struct { + name string + serveFirst bool + conditional bool + }{ + // The mock service returns ErrBlobBlocked itself and never calls the + // isBlobBlockedFunc stub: this case covers only the handler's mapping. + {name: "service ErrBlobBlocked maps to an uncacheable 404"}, + // The handler calls IsBlobBlocked itself before it may answer 304. + {name: "block check refuses a matching conditional after a 200", serveFirst: true, conditional: true}, + } { + t.Run(test.name, func(t *testing.T) { + blocked := false + service := &mockService{ + isBlobBlockedFunc: func(context.Context, string, string) (bool, error) { return blocked, nil }, + getImageFunc: func(context.Context, string, string, string, string) ([]byte, error) { + if blocked { + return nil, imageproxy.ErrBlobBlocked + } + return []byte("previously served image"), nil + }, + } + router := route(service) + if test.serveFirst { + require.Equal(t, http.StatusOK, request(router, "").Code) + } + blocked = true + ifNoneMatch := "" + if test.conditional { + ifNoneMatch = etag + } + response := request(router, ifNoneMatch) + require.Equal(t, http.StatusNotFound, response.Code) + assert.Equal(t, "no-store", response.Header().Get("Cache-Control")) + assert.Empty(t, response.Header().Get("ETag")) + assert.NotContains(t, response.Body.String(), "previously served image") + }) + } + + t.Run("blocked warm disk cache cannot serve stale bytes", func(t *testing.T) { + cache, err := imageproxy.NewDiskCache(t.TempDir(), 1, 0) + require.NoError(t, err) + require.NoError(t, cache.Set("avatar", validTestDID, validTestCID, []byte("cached image bytes"))) + cached, found, err := cache.Get("avatar", validTestDID, validTestCID) + require.NoError(t, err) + require.True(t, found) + require.Equal(t, []byte("cached image bytes"), cached) + + processor, err := imageproxy.NewProcessor(imageproxy.DefaultMaxSourceMegapixels) + require.NoError(t, err) + service, err := imageproxy.NewService(cache, processor, imageproxy.NewPDSFetcher(time.Second, 1), + blockedBlobChecker{}, imageproxy.Config{ + MaxConcurrentProcesses: 1, ProcessQueueWait: time.Second, MaxInFlightRequests: 1, + }) + require.NoError(t, err) + response := request(route(service), "") + require.Equal(t, http.StatusNotFound, response.Code) + assert.Equal(t, "no-store", response.Header().Get("Cache-Control")) + assert.Empty(t, response.Header().Get("ETag")) + assert.NotContains(t, response.Body.String(), "cached image bytes") + }) + + t.Run("failed fetch is not cacheable", func(t *testing.T) { + service := &mockService{getImageFunc: func(context.Context, string, string, string, string) ([]byte, error) { + return nil, imageproxy.ErrPDSFetchFailed + }} + response := request(route(service), "") + require.Equal(t, http.StatusBadGateway, response.Code) + assert.Equal(t, "no-store", response.Header().Get("Cache-Control")) + }) +} + +type blockedBlobChecker struct{} + +func (blockedBlobChecker) IsBlocked(context.Context, string, string) (bool, error) { + return true, nil +} diff --git a/internal/api/handlers/imageproxy/handler.go b/internal/api/handlers/imageproxy/handler.go index 7265d11..5f21fc7 100644 --- a/internal/api/handlers/imageproxy/handler.go +++ b/internal/api/handlers/imageproxy/handler.go @@ -24,6 +24,11 @@ import ( // the pressure being shed, so the two must never drift apart. var processorBusyRetryAfterSeconds = strconv.Itoa(int(imageproxy.DefaultProcessQueueWait / time.Second)) +// successCacheControl is the cache policy of a served image and of its 304. +// One day bounds how long a browser or an unpurged shared cache keeps an image +// after a moderation removal, with no configuration and no CDN required. +const successCacheControl = "public, max-age=86400" + // Service defines the interface for the image proxy service. // This interface is implemented by the imageproxy package's service layer. type Service interface { @@ -77,8 +82,10 @@ func (h *Handler) HandleImage(w http.ResponseWriter, r *http.Request) { return } - // Validate DID format (must be did:plc: or did:web:) - if err := imageproxy.ValidateDID(did); err != nil { + // Accept only canonical did:plc or did:web owner spellings. Anything else + // is refused here, before the block check, the cache lookup, DID + // resolution and the fetch. + if err := imageproxy.ValidateOwnerDID(did); err != nil { writeErrorResponse(w, http.StatusBadRequest, "invalid DID format") return } @@ -108,6 +115,8 @@ func (h *Handler) HandleImage(w http.ResponseWriter, r *http.Request) { writeErrorResponse(w, http.StatusNotFound, "blob not found") return } + w.Header().Set("Cache-Control", successCacheControl) + w.Header().Set("ETag", etag) w.WriteHeader(http.StatusNotModified) return } @@ -146,7 +155,7 @@ func (h *Handler) HandleImage(w http.ResponseWriter, r *http.Request) { // Set response headers w.Header().Set("Content-Type", "image/jpeg") - w.Header().Set("Cache-Control", "public, max-age=31536000, immutable") + w.Header().Set("Cache-Control", successCacheControl) w.Header().Set("ETag", etag) // Write image data @@ -279,11 +288,9 @@ func logBlockCheckFailure(err error, did, cid string) { // For the image proxy, we use simple text responses rather than JSON // since the expected response is binary image data. // -// Errors are explicitly uncacheable. Success responses advertise a one-year -// immutable lifetime, which is correct for content-addressed blobs but -// catastrophic for a failure: this route sits behind a CDN, and a transient -// PDS timeout or a DID that had not yet propagated would otherwise be pinned -// at the edge for a year, long after the image became fetchable. +// Errors and blocked responses are no-store: a cached error or block would +// outlive its cause, still refusing the image after the PDS recovers or the +// block is lifted. func writeErrorResponse(w http.ResponseWriter, status int, message string) { w.Header().Set("Content-Type", "text/plain; charset=utf-8") w.Header().Set("Cache-Control", "no-store") diff --git a/internal/api/handlers/imageproxy/handler_test.go b/internal/api/handlers/imageproxy/handler_test.go index 399c6f1..a2930f0 100644 --- a/internal/api/handlers/imageproxy/handler_test.go +++ b/internal/api/handlers/imageproxy/handler_test.go @@ -159,7 +159,7 @@ func TestHandler_HandleImage_Success(t *testing.T) { // Verify Cache-Control cacheControl := w.Header().Get("Cache-Control") - expectedCacheControl := "public, max-age=31536000, immutable" + expectedCacheControl := "public, max-age=86400" if cacheControl != expectedCacheControl { t.Errorf("Expected Cache-Control %q, got %q", expectedCacheControl, cacheControl) } @@ -689,11 +689,11 @@ func TestHandler_HandleImage_InvalidCID(t *testing.T) { } } -// This route sits behind a CDN and advertises a one-year immutable lifetime on -// success, which is correct for content-addressed blobs. Inheriting anything -// cacheable on an error would pin a transient failure — a PDS timeout, a DID -// that had not propagated yet — at the edge long after the image became -// fetchable. Every error path must therefore say no-store. +// This route sits behind a CDN, and a browser or CDN may cache successful +// images for one day. +// Caching an error would let a transient failure — a PDS timeout or a DID +// that had not propagated yet — outlive its cause. Every error path must +// therefore say no-store. func TestHandler_HandleImage_ErrorsAreNeverCacheable(t *testing.T) { tests := []struct { name string diff --git a/internal/api/handlers/imageproxy/owner_did_spelling_test.go b/internal/api/handlers/imageproxy/owner_did_spelling_test.go new file mode 100644 index 0000000..2029bce --- /dev/null +++ b/internal/api/handlers/imageproxy/owner_did_spelling_test.go @@ -0,0 +1,124 @@ +package imageproxy + +import ( + "bytes" + "context" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "Coves/internal/atproto/identity" + + "github.com/go-chi/chi/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type countingOwnerDIDService struct { + *mockService + getCalls int +} + +func (service *countingOwnerDIDService) GetImageResolvingPDS(ctx context.Context, preset, did, cid string, resolvePDS func(context.Context) (string, error)) ([]byte, error) { + service.getCalls++ + return service.mockService.GetImageResolvingPDS(ctx, preset, did, cid, resolvePDS) +} + +func TestHandlerRefusesNoncanonicalOwnerDIDBeforeService(t *testing.T) { + const preset = "content_preview" + imageBytes := []byte{0xff, 0xd8, 0xff, 0xe0} + for _, test := range []struct { + name, did string + }{ + {name: "uppercase plc identifier", did: "did:plc:Testauthor1"}, + {name: "uppercase web host", did: "did:web:Example.test"}, + {name: "percent-encoded web host letter", did: "did:web:%65xample.test"}, + {name: "path-based web DID", did: "did:web:example.test:a:b"}, + } { + for _, conditional := range []bool{false, true} { + name := "ordinary" + if conditional { + name = "matching If-None-Match" + } + t.Run(test.name+"/"+name, func(t *testing.T) { + blockChecks, resolutions := 0, 0 + service := &countingOwnerDIDService{mockService: &mockService{ + isBlobBlockedFunc: func(context.Context, string, string) (bool, error) { + blockChecks++ + return false, nil + }, + getImageFunc: func(context.Context, string, string, string, string) ([]byte, error) { + return imageBytes, nil + }, + }} + resolver := &mockIdentityResolver{resolveDIDFunc: func(_ context.Context, did string) (*identity.DIDDocument, error) { + resolutions++ + return &identity.DIDDocument{DID: did, Service: []identity.Service{{ + Type: "AtprotoPersonalDataServer", ServiceEndpoint: "https://pds.example.test", + }}}, nil + }} + handler := NewHandler(service, resolver) + router := chi.NewRouter() + router.Get("/img/{preset}/plain/{did}/{cid}", func(w http.ResponseWriter, r *http.Request) { + assert.Equal(t, test.did, chi.URLParam(r, "did"), "chi must pass the exact raw owner spelling to the handler") + handler.HandleImage(w, r) + }) + // Construct the path directly so %65 stays in RawPath rather than + // being escaped a second time before chi routes it. + request := httptest.NewRequest(http.MethodGet, "/img/"+preset+"/plain/"+test.did+"/"+validTestCID, nil) + if strings.Contains(test.did, "%") { + require.Contains(t, request.URL.RawPath, "%65") + } + if conditional { + request.Header.Set("If-None-Match", `"`+preset+`-`+validTestCID+`"`) + } + response := httptest.NewRecorder() + router.ServeHTTP(response, request) + assert.Equal(t, http.StatusBadRequest, response.Code) + assert.Equal(t, "invalid DID format", response.Body.String()) + assert.Equal(t, "no-store", response.Header().Get("Cache-Control")) + assert.NotContains(t, response.Header().Get("Content-Type"), "image/") + assert.False(t, bytes.Equal(imageBytes, response.Body.Bytes()), "no image bytes can escape") + assert.Zero(t, blockChecks, "no moderation block check before DID refusal") + assert.Zero(t, service.getCalls, "no cache lookup or image fetch before DID refusal") + assert.Zero(t, resolutions, "no PDS resolution before DID refusal") + }) + } + } +} + +func TestHandlerServesCanonicalWebOwnerDIDThroughRouter(t *testing.T) { + const preset = "content_preview" + const owner = "did:web:example.test" + const pdsURL = "https://pds.example.test" + imageBytes := []byte{0xff, 0xd8, 0xff, 0xe0} + var resolvedDIDs, fetchedDIDs, fetchedPDSURLs []string + service := &countingOwnerDIDService{mockService: &mockService{ + getImageFunc: func(_ context.Context, _, did, _, resolvedPDSURL string) ([]byte, error) { + fetchedDIDs = append(fetchedDIDs, did) + fetchedPDSURLs = append(fetchedPDSURLs, resolvedPDSURL) + return imageBytes, nil + }, + }} + resolver := &mockIdentityResolver{resolveDIDFunc: func(_ context.Context, did string) (*identity.DIDDocument, error) { + resolvedDIDs = append(resolvedDIDs, did) + return &identity.DIDDocument{DID: did, Service: []identity.Service{{ + Type: "AtprotoPersonalDataServer", ServiceEndpoint: pdsURL, + }}}, nil + }} + router := chi.NewRouter() + router.Get("/img/{preset}/plain/{did}/{cid}", NewHandler(service, resolver).HandleImage) + + response := httptest.NewRecorder() + router.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/img/"+preset+"/plain/"+owner+"/"+validTestCID, nil)) + + require.Equal(t, http.StatusOK, response.Code, "body: %s", response.Body.String()) + assert.Equal(t, imageBytes, response.Body.Bytes()) + assert.Equal(t, "image/jpeg", response.Header().Get("Content-Type")) + assert.Equal(t, successCacheControl, response.Header().Get("Cache-Control")) + assert.Equal(t, 1, service.getCalls, "the canonical did:web owner must reach the service") + assert.Equal(t, []string{owner}, resolvedDIDs, "the canonical did:web owner must be resolved as routed") + assert.Equal(t, []string{owner}, fetchedDIDs) + assert.Equal(t, []string{pdsURL}, fetchedPDSURLs) +} diff --git a/internal/api/handlers/imageproxy/roundtrip_serving_test.go b/internal/api/handlers/imageproxy/roundtrip_serving_test.go index ae8191d..1431e98 100644 --- a/internal/api/handlers/imageproxy/roundtrip_serving_test.go +++ b/internal/api/handlers/imageproxy/roundtrip_serving_test.go @@ -194,9 +194,8 @@ func TestImageProxy_EmittedURLsAreFetchable(t *testing.T) { assertImageSize(t, body, 1000, 1000) }) - // Errors must be uncacheable: this route advertises a one-year immutable - // lifetime on success and sits behind a CDN, so a cacheable failure would - // outlive the condition that caused it by a year. + // Errors must be uncacheable: success can be cached for one day behind a + // CDN, but a cached failure could outlive the condition that caused it. t.Run("an unresolvable blob returns an uncacheable error", func(t *testing.T) { // A decodable CID the PDS does not hold, so the failure is the fetch; // an undecodable one would be refused before it. diff --git a/internal/api/routes/moderation_did_spelling_integration_test.go b/internal/api/routes/moderation_did_spelling_integration_test.go new file mode 100644 index 0000000..1bc2f24 --- /dev/null +++ b/internal/api/routes/moderation_did_spelling_integration_test.go @@ -0,0 +1,110 @@ +//go:build integration + +package routes_test + +import ( + "bytes" + "io" + "net/http" + "strings" + "testing" + + "Coves/internal/db/postgres" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestModerationMediaBlockHoldsForEveryDIDSpelling(t *testing.T) { + const preset = "content_preview" + h, _ := newModerationMediaHarness(t, false) + imageX, imageY := mediaImageCID("blocked plc spelling"), mediaImageCID("blocked web spelling") + webOwner := "did:web:mediaowner.test" + fixtures.User(t, h.db, testkit.UniqueIDWithPrefix(t, "webowner")+".test", webOwner) + h.remove(t, h.comment(t, h.ownerA, imageX), "social.coves.moderation.defs#reasonSpam") + h.remove(t, h.comment(t, webOwner, imageY), "social.coves.moderation.defs#reasonSpam") + store := postgres.NewModerationRepository(h.db) + for _, blockedPair := range []mediaBlobKey{{h.ownerA, imageX}, {webOwner, imageY}} { + blocked, err := store.IsBlocked(t.Context(), blockedPair.did, blockedPair.cid) + require.NoError(t, err) + require.True(t, blocked, "removal must have installed an owner-scoped block") + } + otherBlocked, err := store.IsBlocked(t.Context(), h.ownerB, imageX) + require.NoError(t, err) + require.False(t, otherBlocked, "X must remain fetchable from another owner") + h.pds.mu.Lock() + h.pds.known[mediaBlobKey{h.ownerB, imageX}] = true + h.pds.mu.Unlock() + + requestImage := func(t *testing.T, did, blobCID string, conditional bool) (int, http.Header, []byte) { + t.Helper() + request, err := http.NewRequestWithContext(t.Context(), http.MethodGet, + h.proxy.URL+"/img/"+preset+"/plain/"+did+"/"+blobCID, nil) + require.NoError(t, err) + if strings.Contains(did, "%") { + require.NotEmpty(t, request.URL.RawPath, "escaped owner must reach chi as a raw path") + } + if conditional { + request.Header.Set("If-None-Match", `"`+preset+`-`+blobCID+`"`) + } + response, err := h.proxy.Client().Do(request) + require.NoError(t, err) + defer response.Body.Close() + body, err := io.ReadAll(response.Body) + require.NoError(t, err) + return response.StatusCode, response.Header, body + } + assertRefused := func(t *testing.T, did, blobCID, stage string, seeded []byte, conditional bool) { + t.Helper() + blobRequestsBefore := h.pds.totalBlobRequests() + status, headers, body := requestImage(t, did, blobCID, conditional) + assert.Equal(t, http.StatusBadRequest, status, "%s: %s must be refused as an invalid DID", stage, did) + assert.Equal(t, blobRequestsBefore, h.pds.totalBlobRequests(), "%s: %s must not reach the PDS", stage, did) + assert.Equal(t, "no-store", headers.Get("Cache-Control"), "%s: rejected DID must not be cached", stage) + assert.NotContains(t, headers.Get("Content-Type"), "image/", "%s: rejected DID cannot serve image bytes", stage) + assert.False(t, bytes.Equal(seeded, body), "%s: rejected DID must not serve cached image bytes", stage) + } + + seededImage := testkit.TestPNG(32, 32) + for _, spelling := range []struct { + name, did, blockedDID, blobCID string + }{ + {name: "uppercase plc identifier", did: strings.Replace(h.ownerA, ":test", ":Test", 1), blockedDID: h.ownerA, blobCID: imageX}, + {name: "uppercase web host", did: "did:web:Mediaowner.test", blockedDID: webOwner, blobCID: imageY}, + {name: "percent-encoded web host letter", did: "did:web:%6dediaowner.test", blockedDID: webOwner, blobCID: imageY}, + } { + t.Run(spelling.name, func(t *testing.T) { + require.NotEqual(t, spelling.blockedDID, spelling.did, "the case must request a different spelling, not the blocked DID itself") + assertRefused(t, spelling.did, spelling.blobCID, "cold", seededImage, false) + require.NoError(t, h.cache.Set(preset, spelling.did, spelling.blobCID, seededImage)) + cached, found, err := h.cache.Get(preset, spelling.did, spelling.blobCID) + require.NoError(t, err) + require.True(t, found, "noncanonical spelling must actually be warm on disk") + require.Equal(t, seededImage, cached) + assertRefused(t, spelling.did, spelling.blobCID, "warm", seededImage, false) + assertRefused(t, spelling.did, spelling.blobCID, "matching If-None-Match", seededImage, true) + }) + } + + t.Run("distinct web spellings colliding in disk cache", func(t *testing.T) { + first := "did:web:example.test:a_b" + second := "did:web:example.test:a:b" + require.Equal(t, h.cachePath(preset, first, imageY), h.cachePath(preset, second, imageY), + "the two spellings must actually share a disk-cache directory") + require.NoError(t, h.cache.Set(preset, first, imageY, seededImage)) + cached, found, err := h.cache.Get(preset, second, imageY) + require.NoError(t, err) + require.True(t, found, "the second spelling must be able to see the seeded entry without DID validation") + require.Equal(t, seededImage, cached) + assertRefused(t, second, imageY, "cache-key collision", seededImage, false) + }) + + assert.Zero(t, h.pds.totalBlobRequests(), "refused owner spellings must never cause a blob request, even on a cold cache miss") + status, headers, body := requestImage(t, h.ownerB, imageX, false) + assert.Equal(t, http.StatusOK, status, "a different owner must still serve the same CID") + assert.Contains(t, headers.Get("Content-Type"), "image/") + assert.NotEmpty(t, body, "the positive control must serve image bytes") + assert.Equal(t, 1, h.pds.totalBlobRequests(), "only the positive control may reach the fake PDS") +} diff --git a/internal/api/routes/moderation_media_integration_test.go b/internal/api/routes/moderation_media_integration_test.go index 022da38..634cce8 100644 --- a/internal/api/routes/moderation_media_integration_test.go +++ b/internal/api/routes/moderation_media_integration_test.go @@ -60,9 +60,10 @@ func base58MediaCID(t *testing.T, canonical string) string { type mediaBlobKey struct{ did, cid string } type mediaPDS struct { - mu sync.Mutex - counts map[mediaBlobKey]int - known map[mediaBlobKey]bool + mu sync.Mutex + counts map[mediaBlobKey]int + known map[mediaBlobKey]bool + blobRequests int } func (p *mediaPDS) serve(w http.ResponseWriter, r *http.Request) { @@ -70,6 +71,9 @@ func (p *mediaPDS) serve(w http.ResponseWriter, r *http.Request) { http.NotFound(w, r) return } + p.mu.Lock() + p.blobRequests++ + p.mu.Unlock() // Like the reference PDS, look the blob up by its parsed CID, so every // multibase encoding of one CID names the same blob. parsed, err := cid.Decode(r.URL.Query().Get("cid")) @@ -98,6 +102,12 @@ func (p *mediaPDS) count(did, cid string) int { return p.counts[mediaBlobKey{did, cid}] } +func (p *mediaPDS) totalBlobRequests() int { + p.mu.Lock() + defer p.mu.Unlock() + return p.blobRequests +} + type mediaPDSResolver struct{ url string } func (r mediaPDSResolver) ResolveDID(_ context.Context, did string) (*identity.DIDDocument, error) { diff --git a/internal/atproto/jetstream/moderation_edit_media_acceptance_test.go b/internal/atproto/jetstream/moderation_edit_media_acceptance_test.go new file mode 100644 index 0000000..a8af66b --- /dev/null +++ b/internal/atproto/jetstream/moderation_edit_media_acceptance_test.go @@ -0,0 +1,324 @@ +//go:build integration + +package jetstream + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "io" + "net/http" + "net/http/httptest" + "sync" + "testing" + "time" + + imagehandler "Coves/internal/api/handlers/imageproxy" + "Coves/internal/api/routes" + "Coves/internal/atproto/identity" + "Coves/internal/core/imageproxy" + "Coves/internal/core/moderation" + "Coves/internal/core/posts" + "Coves/internal/crypto/credentialcipher/credentialciphertest" + "Coves/internal/db/postgres" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + "github.com/go-chi/chi/v5" + "github.com/ipfs/go-cid" + "github.com/multiformats/go-multihash" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const editMediaIllegalReason = "social.coves.moderation.defs#reasonIllegalContent" + +func editMediaCID(label string) string { + digest, err := multihash.Sum([]byte(label), multihash.SHA2_256, -1) + if err != nil { + panic(err) + } + return cid.NewCidV1(cid.Raw, digest).String() +} + +type editMediaPDS struct { + mu sync.Mutex + known map[string]map[string]bool +} + +func (p *editMediaPDS) know(owner string, blobCIDs ...string) { + p.mu.Lock() + defer p.mu.Unlock() + if p.known[owner] == nil { + p.known[owner] = make(map[string]bool) + } + for _, blobCID := range blobCIDs { + p.known[owner][blobCID] = true + } +} + +func (p *editMediaPDS) serve(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet || r.URL.Path != "/xrpc/com.atproto.sync.getBlob" { + http.NotFound(w, r) + return + } + parsed, err := cid.Decode(r.URL.Query().Get("cid")) + if err != nil { + http.NotFound(w, r) + return + } + p.mu.Lock() + known := p.known[r.URL.Query().Get("did")][parsed.String()] + p.mu.Unlock() + if !known { + http.NotFound(w, r) + return + } + w.Header().Set("Content-Type", "image/png") + _, _ = w.Write(testkit.TestPNG(32, 32)) +} + +type editMediaResolver struct{ pdsURL string } + +func (r editMediaResolver) ResolveDID(_ context.Context, did string) (*identity.DIDDocument, error) { + return &identity.DIDDocument{DID: did, Service: []identity.Service{{ + Type: "AtprotoPersonalDataServer", ServiceEndpoint: r.pdsURL, + }}}, nil +} + +func (editMediaResolver) Resolve(context.Context, string) (*identity.Identity, error) { + return nil, errors.New("blob owner must be resolved by DID") +} + +func (editMediaResolver) ResolveHandle(context.Context, string) (string, string, error) { + return "", "", errors.New("blob owner must be resolved by DID") +} +func (editMediaResolver) Purge(context.Context, string) error { return nil } + +type editMediaHarness struct { + db *sql.DB + proxy *httptest.Server + pds *editMediaPDS + moderator moderation.Service + posts *PostEventConsumer + comments *CommentEventConsumer + otherOwner string +} + +func newEditMediaHarness(t *testing.T) *editMediaHarness { + t.Helper() + db := testkit.DB(t) + fixture := newPV2Fixture(t, db) + otherOwner := fixtures.DID(testkit.UniqueIDWithPrefix(t, "othermedia")) + pds := &editMediaPDS{known: make(map[string]map[string]bool)} + pdsServer := httptest.NewServer(http.HandlerFunc(pds.serve)) + t.Cleanup(pdsServer.Close) + cache, err := imageproxy.NewDiskCache(t.TempDir(), 1, 0) + require.NoError(t, err) + processor, err := imageproxy.NewProcessor(imageproxy.DefaultMaxSourceMegapixels) + require.NoError(t, err) + store := postgres.NewModerationRepository(db) + proxyService, err := imageproxy.NewService(cache, processor, + imageproxy.NewPDSFetcher(30*time.Second, 10, imageproxy.WithPrivateHostsAllowed()), + store, imageproxy.DefaultConfig()) + require.NoError(t, err) + router := chi.NewRouter() + routes.RegisterImageProxyRoutes(router, imagehandler.NewHandler(proxyService, editMediaResolver{pdsURL: pdsServer.URL})) + proxy := httptest.NewServer(router) + t.Cleanup(proxy.Close) + reconciler := moderation.NewMediaReconciler(store, fixtures.InstanceDID(), proxyService) + moderator := moderation.NewService( + moderation.NewRepositorySubjectReader(postgres.NewPostRepository(db), postgres.NewCommentRepository(db)), + store, moderation.Config{ + InstanceDID: fixtures.InstanceDID(), IdempotencyRetention: 24 * time.Hour, + MaxLiveIdempotencyKeys: 1000, Purger: proxyService, + }, + ) + return &editMediaHarness{ + db: db, proxy: proxy, pds: pds, moderator: moderator, otherOwner: otherOwner, + posts: NewPostEventConsumer(postgres.NewPostRepository(db), + postgres.NewCommunityRepository(db, credentialciphertest.Fixed()), fixture.users, db, + WithAdmissions(fixture.admissions), WithDeletedAccounts(postgres.NewDeletedAccountRepository(db)), + WithPostMediaReconciler(reconciler)), + comments: NewCommentEventConsumer(postgres.NewCommentRepository(db), db, WithCommentMediaReconciler(reconciler)), + } +} + +func (h *editMediaHarness) imageStatus(t *testing.T, owner, blobCID string) int { + t.Helper() + request, err := http.NewRequestWithContext(t.Context(), http.MethodGet, + h.proxy.URL+"/img/content_preview/plain/"+owner+"/"+blobCID, nil) + require.NoError(t, err) + response, err := h.proxy.Client().Do(request) + require.NoError(t, err) + defer response.Body.Close() + body, err := io.ReadAll(response.Body) + require.NoError(t, err) + if response.StatusCode == http.StatusOK { + require.NotEmpty(t, body, "the image proxy must serve image bytes") + require.Equal(t, "image/jpeg", response.Header.Get("Content-Type")) + } + return response.StatusCode +} + +func (h *editMediaHarness) remove(t *testing.T, subject moderation.StrongRef, version string) *moderation.MutationResult { + t.Helper() + result, err := h.moderator.RemoveContent(t.Context(), fixtures.DID(testkit.UniqueIDWithPrefix(t, "mediaadmin")), moderation.RemoveContentRequest{ + Subject: subject, ExpectedVersion: version, IdempotencyKey: "remove-" + testkit.UniqueID(t), Reason: editMediaIllegalReason, + }) + require.NoError(t, err) + require.Equal(t, moderation.OutcomeApplied, result.Outcome) + require.NotNil(t, result.Action) + return result +} + +func (h *editMediaHarness) restore(t *testing.T, subject moderation.StrongRef, removed *moderation.MutationResult) string { + t.Helper() + result, err := h.moderator.RestoreContent(t.Context(), fixtures.DID(testkit.UniqueIDWithPrefix(t, "mediaadmin")), moderation.RestoreContentRequest{ + ActionID: removed.Action.ID, ReviewedSubject: &subject, ExpectedVersion: removed.State.Version, + IdempotencyKey: "restore-" + testkit.UniqueID(t), Reason: "social.coves.moderation.defs#reasonModeratorDiscretion", + }) + require.NoError(t, err) + require.Equal(t, moderation.OutcomeApplied, result.Outcome) + return result.State.Version +} + +func editMediaCommentEvent(root moderation.StrongRef, operation, rkey, rev, recordCID string, eventTime time.Time, blobCIDs ...string) *JetstreamEvent { + var record map[string]interface{} + if operation != "delete" { + record = map[string]interface{}{ + "$type": moderation.CommentCollection, "content": "image comment", "createdAt": eventTime.Format(time.RFC3339), + "reply": map[string]interface{}{ + "root": map[string]interface{}{"uri": root.URI, "cid": root.CID}, + "parent": map[string]interface{}{"uri": root.URI, "cid": root.CID}, + }, + "embed": postModerationRecord(blobCIDs...)["embed"], + } + } + return revCommitEvent(pv2Author, moderation.CommentCollection, operation, rkey, rev, recordCID, eventTime.UnixMicro(), record) +} + +// A consumer may block an edit-introduced blob only for its author. The original +// admin removal retains its ownerless block, while a later removal covers the +// blobs indexed at that later instant. +func TestModerationEditIntroducedMediaIsOwnerScoped(t *testing.T) { + for _, variant := range []struct { + name string + kind string + // indexesEdit is true when the consumer indexes B, so the stored row + // names it at the second removal. + indexesEdit bool + }{ + {name: "comment edit", kind: "comment", indexesEdit: true}, + {name: "postv2 edit", kind: "post edit", indexesEdit: true}, + {name: "postv2 delete and recreate", kind: "post recreate"}, + {name: "postv2 already-indexed create", kind: "post conflict"}, + } { + t.Run(variant.name, func(t *testing.T) { + h := newEditMediaHarness(t) + imageA := editMediaCID(variant.name + " image A") + imageB := editMediaCID(variant.name + " image B") + for _, owner := range []string{pv2Author, h.otherOwner} { + h.pds.know(owner, imageA, imageB) + } + revs := increasingTIDs(t, 4) + started := time.Now().Add(-time.Minute) + postRkey := testkit.TID() + postSubject := moderation.StrongRef{URI: pv2URI(pv2Author, postRkey), CID: editMediaCID(variant.name + " post record")} + require.NoError(t, h.posts.HandleEvent(t.Context(), pv2Event(pv2Author, "create", postRkey, revs[0], + postSubject.CID, started.UnixMicro(), postModerationRecord(imageA)))) + subject := postSubject + var commentRkey string + if variant.kind == "comment" { + commentRkey = testkit.TID() + subject = moderation.StrongRef{ + URI: "at://" + pv2Author + "/" + moderation.CommentCollection + "/" + commentRkey, + CID: editMediaCID(variant.name + " comment record"), + } + require.NoError(t, h.comments.HandleEvent(t.Context(), editMediaCommentEvent(postSubject, + "create", commentRkey, revs[0], subject.CID, started, imageA))) + } + // Warm the other owner's copy of B before reconciliation: a block must + // win over disk cache, and a scoped purge must leave this copy intact. + require.Equal(t, http.StatusOK, h.imageStatus(t, h.otherOwner, imageB)) + removed := h.remove(t, subject, "v0") + switch variant.kind { + case "comment": + require.NoError(t, h.comments.HandleEvent(t.Context(), editMediaCommentEvent(postSubject, + "update", commentRkey, revs[1], editMediaCID(variant.name+" edited record"), started.Add(time.Second), imageA, imageB))) + case "post edit": + require.NoError(t, h.posts.HandleEvent(t.Context(), pv2Event(pv2Author, "update", postRkey, revs[1], + editMediaCID(variant.name+" edited record"), started.Add(time.Second).UnixMicro(), postModerationRecord(imageA, imageB)))) + case "post recreate": + require.NoError(t, h.posts.HandleEvent(t.Context(), pv2Event(pv2Author, "delete", postRkey, revs[1], + "", started.Add(time.Second).UnixMicro(), nil))) + require.NoError(t, h.posts.HandleEvent(t.Context(), pv2Event(pv2Author, "create", postRkey, revs[2], + editMediaCID(variant.name+" recreated record"), started.Add(2*time.Second).UnixMicro(), postModerationRecord(imageB)))) + case "post conflict": + // Exercise the real create/ON CONFLICT path, whose incoming embed + // never replaces the already-indexed post row. + embed, err := json.Marshal(postModerationRecord(imageB)["embed"]) + require.NoError(t, err) + embedJSON := string(embed) + applied, err := h.posts.indexPostIfRevWins(t.Context(), &posts.Post{ + URI: subject.URI, CID: editMediaCID(variant.name + " incoming record"), RKey: postRkey, + AuthorDID: pv2Author, CommunityDID: pv2Community, Embed: &embedJSON, + CreatedAt: started, IndexedAt: started.Add(time.Second), + }, revs[1]) + require.NoError(t, err) + require.False(t, applied, "an existing post must take the ON CONFLICT path") + default: + t.Fatalf("unknown variant kind %q", variant.kind) + } + var indexedEmbed string + if variant.kind == "comment" { + require.NoError(t, h.db.QueryRowContext(t.Context(), `SELECT embed::text FROM comments WHERE uri = $1`, subject.URI).Scan(&indexedEmbed)) + } else { + require.NoError(t, h.db.QueryRowContext(t.Context(), `SELECT embed::text FROM posts WHERE uri = $1`, subject.URI).Scan(&indexedEmbed)) + } + if variant.indexesEdit { + require.Contains(t, indexedEmbed, imageB, "the consumer must index the edited image") + } else { + require.Contains(t, indexedEmbed, imageA, "the original post row must remain indexed") + } + assert.Equal(t, 1, countRows(t, h.db, `SELECT count(*) FROM moderation_media_blocks + WHERE owner_did = $1 AND blob_cid = $2 AND active`, pv2Author, imageB), + "the new image must have exactly one active author-owned block") + assert.Equal(t, 1, countRows(t, h.db, `SELECT count(*) FROM moderation_media_blocks + WHERE action_id = $1 AND owner_did = $2 AND blob_cid = $3 AND active`, removed.Action.ID, pv2Author, imageB), + "the author-owned block belongs to the removal being restored") + assert.Zero(t, countRows(t, h.db, `SELECT count(*) FROM moderation_media_blocks + WHERE blob_cid = $1 AND owner_did IS NULL AND active`, imageB), + "edit-introduced B must have no ownerless block") + assert.Equal(t, http.StatusOK, h.imageStatus(t, h.otherOwner, imageB), "another owner's warm B stays available") + assert.Equal(t, http.StatusNotFound, h.imageStatus(t, pv2Author, imageB), "the author's B is blocked") + assert.Equal(t, http.StatusNotFound, h.imageStatus(t, h.otherOwner, imageA), "the initial removal blocks A for every owner") + + currentCID := subject.CID + if variant.indexesEdit { + currentCID = editMediaCID(variant.name + " edited record") + } + currentSubject := moderation.StrongRef{URI: subject.URI, CID: currentCID} + version := h.restore(t, currentSubject, removed) + for _, owner := range []string{pv2Author, h.otherOwner} { + for _, blobCID := range []string{imageA, imageB} { + assert.Equal(t, http.StatusOK, h.imageStatus(t, owner, blobCID), "restore must release both blobs for %s", owner) + } + } + h.remove(t, currentSubject, version) + indexedImages := []string{imageA} + if variant.indexesEdit { + indexedImages = append(indexedImages, imageB) + } + for _, indexedImage := range indexedImages { + assert.Equal(t, http.StatusNotFound, h.imageStatus(t, h.otherOwner, indexedImage), + "a new illegal-content removal blocks each currently indexed CID for every owner, even after a warm cache hit") + } + if !variant.indexesEdit { + assert.Equal(t, http.StatusOK, h.imageStatus(t, h.otherOwner, imageB), + "a new removal must not block another owner's B, which the stored row does not name") + } + }) + } +} diff --git a/internal/atproto/jetstream/post_moderation_consumer_test.go b/internal/atproto/jetstream/post_moderation_consumer_test.go index 7e05c02..9e642da 100644 --- a/internal/atproto/jetstream/post_moderation_consumer_test.go +++ b/internal/atproto/jetstream/post_moderation_consumer_test.go @@ -191,11 +191,11 @@ func TestModerationPostConsumerReplayAndDuplicateKeepRemoval(t *testing.T) { func TestModerationPostConsumerEditReconcilesNewImages(t *testing.T) { for _, scenario := range []struct { - name, reason string - ownerless bool + name, reason string + illegalContent bool }{ {name: "spam owner block", reason: "social.coves.moderation.defs#reasonSpam"}, - {name: "illegal content ownerless block", reason: "social.coves.moderation.defs#reasonIllegalContent", ownerless: true}, + {name: "illegal content owner-scoped block", reason: "social.coves.moderation.defs#reasonIllegalContent", illegalContent: true}, } { t.Run(scenario.name, func(t *testing.T) { f := newPostModerationConsumerFixture(t) @@ -227,15 +227,18 @@ func TestModerationPostConsumerEditReconcilesNewImages(t *testing.T) { assert.Equal(t, 1, ownedBlocks, "the edited image needs an active author-owned block on the original action") assert.Contains(t, f.purger.calls, postModerationPurgeCall{ownerDID: pv2Author, blobCID: postModerationCIDTwo, committed: true}, "the consumer must purge the author-owned image after the block commits") - if scenario.ownerless { + if scenario.illegalContent { var ownerlessBlocks int require.NoError(t, f.db.QueryRowContext(t.Context(), ` SELECT count(*) FROM moderation_media_blocks WHERE action_id = $1 AND owner_did IS NULL AND blob_cid = $2 AND active `, removed.Action.ID, postModerationCIDTwo).Scan(&ownerlessBlocks)) - assert.Equal(t, 1, ownerlessBlocks, "illegal content must also block the edited image for every owner") - assert.Contains(t, f.purger.calls, postModerationPurgeCall{blobCID: postModerationCIDTwo, committed: true}, - "the ownerless image cache must be purged after commit") + assert.Zero(t, ownerlessBlocks, "an edit must not add an ownerless image block") + assert.Equal(t, []postModerationPurgeCall{{ownerDID: pv2Author, blobCID: postModerationCIDTwo, committed: true}}, f.purger.calls, + "only the author's image cache is purged after commit") + blocked, checkErr := postgres.NewModerationRepository(f.db).IsBlocked(t.Context(), pv2Other, postModerationCIDTwo) + require.NoError(t, checkErr) + assert.False(t, blocked, "another owner's new image must remain available") } }) } @@ -243,11 +246,11 @@ func TestModerationPostConsumerEditReconcilesNewImages(t *testing.T) { func TestModerationPostConsumerDeleteThenCreateRemainsNotFound(t *testing.T) { for _, scenario := range []struct { - name, reason string - ownerless bool + name, reason string + illegalContent bool }{ {name: "spam owner block", reason: "social.coves.moderation.defs#reasonSpam"}, - {name: "illegal content ownerless block", reason: "social.coves.moderation.defs#reasonIllegalContent", ownerless: true}, + {name: "illegal content owner-scoped block", reason: "social.coves.moderation.defs#reasonIllegalContent", illegalContent: true}, } { t.Run(scenario.name, func(t *testing.T) { f := newPostModerationConsumerFixture(t) @@ -271,13 +274,16 @@ func TestModerationPostConsumerDeleteThenCreateRemainsNotFound(t *testing.T) { `, removed.Action.ID, pv2Author, postModerationCIDTwo), "the recreated image needs an active author-owned block on the original action") assert.Contains(t, f.purger.calls, postModerationPurgeCall{ownerDID: pv2Author, blobCID: postModerationCIDTwo, committed: true}, "the consumer must purge the recreated image after the block commits") - if scenario.ownerless { - assert.Equal(t, 1, countRows(t, f.db, ` + if scenario.illegalContent { + assert.Zero(t, countRows(t, f.db, ` SELECT count(*) FROM moderation_media_blocks WHERE action_id = $1 AND owner_did IS NULL AND blob_cid = $2 AND active - `, removed.Action.ID, postModerationCIDTwo), "illegal content must also block the recreated image for every owner") - assert.Contains(t, f.purger.calls, postModerationPurgeCall{blobCID: postModerationCIDTwo, committed: true}, - "the ownerless image cache must be purged after commit") + `, removed.Action.ID, postModerationCIDTwo), "a recreate must not add an ownerless image block") + assert.Equal(t, []postModerationPurgeCall{{ownerDID: pv2Author, blobCID: postModerationCIDTwo, committed: true}}, f.purger.calls, + "only the author's image cache is purged after commit") + blocked, checkErr := postgres.NewModerationRepository(f.db).IsBlocked(t.Context(), pv2Other, postModerationCIDTwo) + require.NoError(t, checkErr) + assert.False(t, blocked, "another owner's recreated image must remain available") } }) } @@ -306,14 +312,16 @@ func TestModerationPostConsumerAlreadyIndexedCreateBlocksIncomingImages(t *testi SELECT count(*) FROM moderation_media_blocks WHERE action_id = $1 AND owner_did = $2 AND blob_cid = $3 AND active `, removed.Action.ID, pv2Author, postModerationCIDTwo), "the discarded create's image needs an author-owned block") - assert.Equal(t, 1, countRows(t, f.db, ` + assert.Zero(t, countRows(t, f.db, ` SELECT count(*) FROM moderation_media_blocks WHERE action_id = $1 AND owner_did IS NULL AND blob_cid = $2 AND active - `, removed.Action.ID, postModerationCIDTwo), "illegal content must block the discarded create's image for every owner") + `, removed.Action.ID, postModerationCIDTwo), "a discarded create must not add an ownerless block") assert.ElementsMatch(t, []postModerationPurgeCall{ {ownerDID: pv2Author, blobCID: postModerationCIDTwo, committed: true}, - {blobCID: postModerationCIDTwo, committed: true}, }, f.purger.calls) + blocked, checkErr := postgres.NewModerationRepository(f.db).IsBlocked(t.Context(), pv2Other, postModerationCIDTwo) + require.NoError(t, checkErr) + assert.False(t, blocked, "another owner's incoming image must remain available") } func TestModerationPostConsumerRemovalIsPerURI(t *testing.T) { diff --git a/internal/core/comments/comment_moderation_consumer_integration_test.go b/internal/core/comments/comment_moderation_consumer_integration_test.go index 75a0efb..0045b7d 100644 --- a/internal/core/comments/comment_moderation_consumer_integration_test.go +++ b/internal/core/comments/comment_moderation_consumer_integration_test.go @@ -92,7 +92,7 @@ func TestModerationCommentConsumerReconcilesRemovedImages(t *testing.T) { remove bool }{ {name: "edit adds image under spam removal", reason: moderationTestReason, operation: "update", initialImage: true, remove: true}, - {name: "illegal content edit adds ownerless block", reason: "social.coves.moderation.defs#reasonIllegalContent", operation: "update", remove: true}, + {name: "illegal content edit adds owner-scoped block", reason: "social.coves.moderation.defs#reasonIllegalContent", operation: "update", remove: true}, {name: "duplicate create after removal does not purge", reason: moderationTestReason, operation: "duplicate", initialImage: true, remove: true}, {name: "author delete then recreate same URI", reason: moderationTestReason, operation: "recreate", remove: true}, {name: "newer-rev re-create of the active row", reason: moderationTestReason, operation: "recreate-active", initialImage: true, remove: true}, @@ -243,12 +243,15 @@ func TestModerationCommentConsumerReconcilesRemovedImages(t *testing.T) { require.NotEmpty(t, purger.calls, "new pair must be purged") assert.Equal(t, consumerPurgeCall{ownerDID: authorDID, blobCID: newImageCID, visible: true}, purger.calls[0], "pair purge must happen after commit") if scenario.reason == "social.coves.moderation.defs#reasonIllegalContent" { - assert.Equal(t, 2, blocks, "illegal content blocks both the owner pair and every owner") - require.Len(t, purger.calls, 2) - assert.Contains(t, purger.calls, consumerPurgeCall{blobCID: newImageCID, visible: true}, "ownerless purge must happen after commit") + assert.Equal(t, 1, blocks, "an edit adds only the author's image block") + var ownerlessBlocks int + require.NoError(t, db.QueryRowContext(ctx, `SELECT count(*) FROM moderation_media_blocks WHERE blob_cid = $1 AND owner_did IS NULL AND active`, newImageCID).Scan(&ownerlessBlocks)) + assert.Zero(t, ownerlessBlocks, "an edit must not add an ownerless image block") + assert.Equal(t, []consumerPurgeCall{{ownerDID: authorDID, blobCID: newImageCID, visible: true}}, purger.calls, + "only the author's image cache is purged after commit") otherOwnerBlocked, checkErr := moderationRepo.IsBlocked(ctx, fixtures.DID("otherimageowner"), newImageCID) require.NoError(t, checkErr) - assert.True(t, otherOwnerBlocked) + assert.False(t, otherOwnerBlocked, "another owner's new image must remain available") } else { assert.Equal(t, 1, blocks, "spam blocks only the image owner's pair") require.Len(t, purger.calls, 1) diff --git a/internal/core/imageproxy/moderation_test.go b/internal/core/imageproxy/moderation_test.go index 7f7b759..2567cad 100644 --- a/internal/core/imageproxy/moderation_test.go +++ b/internal/core/imageproxy/moderation_test.go @@ -45,6 +45,43 @@ func TestImageProxyService_BlocksBeforeReadingCache(t *testing.T) { } } +// The handler refuses these first; the service refuses them too, so a second +// caller cannot miss an owner-scoped block or read another spelling's cache. +func TestImageProxyService_RefusesNoncanonicalOwnerDID(t *testing.T) { + for _, test := range []struct{ name, owner string }{ + {name: "uppercase plc identifier", owner: "did:plc:Moderationimageowner"}, + {name: "path-based web DID sharing a cache directory", owner: "did:web:example.test:a:b"}, + } { + t.Run(test.name, func(t *testing.T) { + cache := NewMockCache() + cache.SetCacheData("avatar", test.owner, moderationTestCID, []byte("cached secret")) + fetcher := NewMockFetcher([]byte("fetched secret"), nil) + var checks, resolutions atomic.Int32 + checker := blockCheckFunc(func(context.Context, string, string) (bool, error) { + checks.Add(1) + return false, nil + }) + service, err := NewService(cache, NewMockProcessor([]byte("processed"), nil), fetcher, checker, DefaultConfig()) + require.NoError(t, err) + + data, err := service.GetImageResolvingPDS(t.Context(), "avatar", test.owner, moderationTestCID, func(context.Context) (string, error) { + resolutions.Add(1) + return "https://pds.example.com", nil + }) + assert.ErrorIs(t, err, ErrInvalidDID) + assert.Empty(t, data) + blocked, err := service.IsBlobBlocked(t.Context(), test.owner, moderationTestCID) + assert.ErrorIs(t, err, ErrInvalidDID) + assert.False(t, blocked) + + assert.Zero(t, checks.Load(), "no block lookup for a noncanonical owner") + assert.Zero(t, cache.GetCalls(), "a warm entry under a noncanonical owner must not be read") + assert.Zero(t, resolutions.Load(), "no PDS resolution for a noncanonical owner") + assert.Zero(t, fetcher.Calls()) + }) + } +} + func TestImageProxyService_BlockCheckFailureFailsClosed(t *testing.T) { cache := NewMockCache() cache.SetCacheData("avatar", moderationTestOwner, moderationTestCID, []byte("cached secret")) diff --git a/internal/core/imageproxy/service.go b/internal/core/imageproxy/service.go index c377296..487b4fc 100644 --- a/internal/core/imageproxy/service.go +++ b/internal/core/imageproxy/service.go @@ -148,7 +148,7 @@ func NewService(cache Cache, processor Processor, fetcher Fetcher, blocks BlockC const defaultPublicationTimeout = 10 * time.Second // GetImageResolvingPDS implements Service. The service flow is: -// 1. Validate preset exists +// 1. Validate preset exists and the owner DID is canonical // 2. Check moderation, then cache for (preset, did, cid) - return if hit // 3. Resolve the DID's PDS with resolvePDS (misses only) // 4. Acquire an admission slot, waiting at most ProcessQueueWait @@ -171,6 +171,11 @@ func (s *ImageProxyService) GetImageResolvingPDS( if err != nil { return nil, err } + // A second spelling of the owner would miss its block and could read + // another spelling's cache directory; see ValidateOwnerDID. + if err := ValidateOwnerDID(did); err != nil { + return nil, err + } // Step 2: Check moderation before reading even a warm cache entry. blocked, err := s.IsBlobBlocked(ctx, did, cid) @@ -365,8 +370,12 @@ type BlockChecker interface { IsBlocked(ctx context.Context, ownerDID, cid string) (bool, error) } -// IsBlobBlocked reports whether moderation blocks serving the blob. +// IsBlobBlocked reports whether moderation blocks serving the blob. A +// noncanonical owner DID returns ErrInvalidDID: blocks match one spelling. func (s *ImageProxyService) IsBlobBlocked(ctx context.Context, did, cid string) (bool, error) { + if err := ValidateOwnerDID(did); err != nil { + return false, err + } blocked, err := s.blocks.IsBlocked(ctx, did, cid) if err != nil { return false, fmt.Errorf("%w: %w", ErrBlockCheckFailed, err) diff --git a/internal/core/imageproxy/validation.go b/internal/core/imageproxy/validation.go index 72d4671..e4547ae 100644 --- a/internal/core/imageproxy/validation.go +++ b/internal/core/imageproxy/validation.go @@ -1,6 +1,7 @@ package imageproxy import ( + "regexp" "strings" "github.com/bluesky-social/indigo/atproto/syntax" @@ -25,6 +26,35 @@ func ValidateDID(did string) error { return nil } +// ownerDIDPattern is the canonical owner spelling: a lowercase did:plc, or a +// lowercase did:web host with an optional %3A port that has no leading zero. +var ownerDIDPattern = regexp.MustCompile(`^(?:did:plc:[a-z0-9._-]+|did:web:[a-z0-9-]+(?:\.[a-z0-9-]+)*(?:%3A[1-9][0-9]{0,4})?)$`) + +// ValidateOwnerDID accepts only the canonical spelling of a did:plc or did:web +// blob owner, so that each owner has exactly one string. Two lookups key on it: +// - the media block check compares owner_did = $1 exactly, so an uppercase or +// percent-escaped spelling of a blocked owner would not match its block; +// - the disk cache directory is the DID with ":" rewritten to "_" and ".." +// stripped, so did:web:x:a:b and did:web:x:a_b would share cached bytes. +// +// The pattern therefore refuses uppercase, every percent escape except one +// %3A port, and any colon after the method. +// +// A real PLC identifier is 24 base32 characters. The pattern also allows ".", +// "_" and "-" so that test fixtures such as did:plc:apikey_aggregator pass; +// the cache path rewrites none of those characters, so they cannot alias. +// ValidateDID must run first: it refuses "..", which the pattern accepts and +// the cache strips (did:plc:ab..cd would share did:plc:abcd's directory). +func ValidateOwnerDID(did string) error { + if err := ValidateDID(did); err != nil { + return err + } + if !ownerDIDPattern.MatchString(did) { + return ErrInvalidDID + } + return nil +} + // ValidateCID validates that a CID string is a valid content identifier. // It uses the Indigo library's syntax.ParseCID for consistent validation across the codebase. // Returns ErrInvalidCID if the CID is invalid. diff --git a/internal/core/imageproxy/validation_test.go b/internal/core/imageproxy/validation_test.go index 26a52d6..fc27739 100644 --- a/internal/core/imageproxy/validation_test.go +++ b/internal/core/imageproxy/validation_test.go @@ -6,8 +6,53 @@ import ( "github.com/ipfs/go-cid" "github.com/multiformats/go-multibase" + "github.com/stretchr/testify/require" ) +func TestValidateOwnerDID(t *testing.T) { + for _, test := range []struct { + name, did string + accepted bool + }{ + {name: "canonical plc", did: "did:plc:z72i7hdynmk6r22z27h6tvur", accepted: true}, + {name: "fixture plc", did: "did:plc:testauthor1", accepted: true}, + {name: "canonical web", did: "did:web:example.test", accepted: true}, + {name: "subdomain and hyphen", did: "did:web:sub.example-host.test", accepted: true}, + {name: "web port", did: "did:web:localhost%3A8080", accepted: true}, + {name: "uppercase plc identifier", did: "did:plc:testAuthor1"}, + {name: "uppercase web host", did: "did:web:Example.test"}, + {name: "escaped web host letter", did: "did:web:%65xample.test"}, + {name: "escaped plc identifier", did: "did:plc:%61bc"}, + {name: "lowercase port escape", did: "did:web:localhost%3a8080"}, + {name: "port escape without digits", did: "did:web:localhost%3A"}, + {name: "web port leading zero", did: "did:web:localhost%3A08080"}, + {name: "web port zero", did: "did:web:localhost%3A0"}, + {name: "web path segments", did: "did:web:example.test:a:b"}, + {name: "web path underscore", did: "did:web:example.test:a_b"}, + {name: "web host underscore", did: "did:web:exa_mple.test"}, + {name: "empty web host", did: "did:web:"}, + {name: "empty plc identifier", did: "did:plc:"}, + {name: "key method", did: "did:key:z6MkhaXgBZDvotDkL5257faiztiGiC2QtKLGpbnnEGta2doK"}, + {name: "uppercase method", did: "did:PLC:abc"}, + {name: "plc traversal", did: "did:plc:../../../etc/passwd"}, + {name: "plc path", did: "did:plc:abc/def"}, + // The pattern alone accepts these; the disk cache strips ".." from the + // directory name, so each would share did:plc:abcd's (or did:plc:ab.cd's) + // cache entries. Only ValidateDID's ".." guard refuses them. + {name: "plc double dot", did: "did:plc:ab..cd"}, + {name: "plc triple dot", did: "did:plc:ab...cd"}, + } { + t.Run(test.name, func(t *testing.T) { + err := ValidateOwnerDID(test.did) + if test.accepted { + require.NoError(t, err) + } else { + require.ErrorIs(t, err, ErrInvalidDID) + } + }) + } +} + func TestValidateDID(t *testing.T) { tests := []struct { name string diff --git a/internal/core/moderation/media.go b/internal/core/moderation/media.go index 30d0ccf..d683282 100644 --- a/internal/core/moderation/media.go +++ b/internal/core/moderation/media.go @@ -3,6 +3,7 @@ package moderation import ( "context" "database/sql" + "fmt" "log/slog" ) @@ -21,8 +22,11 @@ type TransactionBinder interface { BindTransaction(tx *sql.Tx) MediaTransaction } -// MediaReconciler keeps media blocks in step with a removed subject's -// indexed images when a consumer rewrites the subject. +// MediaReconciler adds owner-scoped blocks for images introduced after a +// subject is removed. Ownerless (every-owner) blocks exist only for an +// illegal-content removal and cover only the images indexed when the admin +// removed the subject; reconciliation never adds one. Both reconcile methods +// therefore refuse an empty owner DID, which the store records as ownerless. type MediaReconciler struct { binder TransactionBinder instanceDID string @@ -35,7 +39,8 @@ func NewMediaReconciler(binder TransactionBinder, instanceDID string, purger Med } // ReconcileTx blocks images newly present on a subject with an active -// removal, inside the caller's transaction, and returns the new blocks. +// removal for that subject's owner, inside the caller's transaction, and +// returns the new blocks. func (r *MediaReconciler) ReconcileTx(ctx context.Context, tx *sql.Tx, subjectURI string) ([]MediaBlock, error) { bound := r.binder.BindTransaction(tx) action, err := bound.ActiveRemoval(ctx, r.instanceDID, subjectURI) @@ -49,13 +54,17 @@ func (r *MediaReconciler) ReconcileTx(ctx context.Context, tx *sql.Tx, subjectUR if subject == nil { return nil, ErrSubjectNotIndexed } - return bound.InsertNewMediaBlocks(ctx, imageMediaBlocks(subject, action)) + if subject.OwnerDID == "" { + return nil, fmt.Errorf("%w: indexed subject has no owner DID", ErrInvalidSubject) + } + return bound.InsertNewMediaBlocks(ctx, ownerImageMediaBlocks(subject, action)) } // ReconcileIncomingTx blocks blobs of incoming content the consumer did not // index, such as a recreate of a deleted post whose tombstone is kept. The // owner's repository still serves those blobs, and the stored row does not // name them, so the caller passes the owner and CIDs from the incoming record. +// These new blocks apply only to that owner, regardless of removal reason. func (r *MediaReconciler) ReconcileIncomingTx(ctx context.Context, tx *sql.Tx, subjectURI, ownerDID string, blobCIDs []string) ([]MediaBlock, error) { if len(blobCIDs) == 0 { return nil, nil @@ -65,7 +74,10 @@ func (r *MediaReconciler) ReconcileIncomingTx(ctx context.Context, tx *sql.Tx, s if err != nil || action == nil { return nil, err } - return bound.InsertNewMediaBlocks(ctx, imageMediaBlocks(&indexedSubject{OwnerDID: ownerDID, BlobCIDs: blobCIDs}, action)) + if ownerDID == "" { + return nil, fmt.Errorf("%w: incoming content has no owner DID", ErrInvalidSubject) + } + return bound.InsertNewMediaBlocks(ctx, ownerImageMediaBlocks(&indexedSubject{OwnerDID: ownerDID, BlobCIDs: blobCIDs}, action)) } // Purge removes cached bytes of newly blocked blobs after commit. @@ -96,7 +108,9 @@ func purgeMediaBlocks(purger MediaPurger, blocks []MediaBlock) { } } -func imageMediaBlocks(subject *indexedSubject, action *Action) []MediaBlock { +// ownerImageMediaBlocks returns one owner-scoped block per distinct CID. +// Reconciliation uses it directly, and imageMediaBlocks builds on it. +func ownerImageMediaBlocks(subject *indexedSubject, action *Action) []MediaBlock { var blocks []MediaBlock seen := make(map[string]bool) for _, cid := range subject.BlobCIDs { @@ -105,9 +119,20 @@ func imageMediaBlocks(subject *indexedSubject, action *Action) []MediaBlock { } seen[cid] = true blocks = append(blocks, MediaBlock{OwnerDID: subject.OwnerDID, BlobCID: cid, ActionID: action.ID}) - if action.Reason == illegalContentReason { - blocks = append(blocks, MediaBlock{BlobCID: cid, ActionID: action.ID}) - } + } + return blocks +} + +// imageMediaBlocks returns the blocks an admin removal installs: the owner +// blocks plus, for illegal content, an every-owner block per CID. +func imageMediaBlocks(subject *indexedSubject, action *Action) []MediaBlock { + ownerBlocks := ownerImageMediaBlocks(subject, action) + if action.Reason != illegalContentReason { + return ownerBlocks + } + var blocks []MediaBlock + for _, ownerBlock := range ownerBlocks { + blocks = append(blocks, ownerBlock, MediaBlock{BlobCID: ownerBlock.BlobCID, ActionID: action.ID}) } return blocks } diff --git a/internal/core/moderation/post_rules_test.go b/internal/core/moderation/post_rules_test.go index 1cd4567..2bd3f68 100644 --- a/internal/core/moderation/post_rules_test.go +++ b/internal/core/moderation/post_rules_test.go @@ -231,10 +231,11 @@ func TestRestorePostContentReviewsCurrentPostAndDeactivatesBlocks(t *testing.T) } type postRulesMediaTransaction struct { - post moderation.IndexedPost - action moderation.Action - calls []string - blocks []moderation.MediaBlock + post moderation.IndexedPost + comment *moderation.IndexedComment + action moderation.Action + calls []string + blocks []moderation.MediaBlock } func (transaction *postRulesMediaTransaction) ActiveRemoval(context.Context, string, string) (*moderation.Action, error) { @@ -244,6 +245,9 @@ func (transaction *postRulesMediaTransaction) ActiveRemoval(context.Context, str func (transaction *postRulesMediaTransaction) ReadIndexedComment(context.Context, string) (*moderation.IndexedComment, error) { transaction.calls = append(transaction.calls, "ReadIndexedComment") + if transaction.comment != nil { + return transaction.comment, nil + } return &moderation.IndexedComment{}, nil } @@ -284,4 +288,96 @@ func TestReconcileRemovedPostMediaReadsIndexedPost(t *testing.T) { assert.Equal(t, want, blocks) } +func TestReconcileIllegalContentMediaIsOwnerScoped(t *testing.T) { + for _, test := range []struct { + name, uri, ownerDID string + blobCIDs []string + incoming bool + readCall string + }{ + { + name: "indexed post", uri: postRulesURI, ownerDID: postRulesAuthorDID, + blobCIDs: []string{removeRulesFirstImage, removeRulesSecondImage}, readCall: "ReadIndexedPost", + }, + { + name: "indexed comment", uri: removeRulesURI, ownerDID: removeRulesAuthorDID, + blobCIDs: []string{removeRulesFirstImage}, readCall: "ReadIndexedComment", + }, + { + name: "incoming post blobs not indexed", uri: postRulesURI, ownerDID: postRulesAuthorDID, + blobCIDs: []string{removeRulesSecondImage}, incoming: true, + }, + } { + t.Run(test.name, func(t *testing.T) { + bound := &postRulesMediaTransaction{ + post: moderation.IndexedPost{ + URI: postRulesURI, CID: postRulesCID, + OwnerDID: postRulesAuthorDID, BlobCIDs: test.blobCIDs, + }, + comment: &moderation.IndexedComment{ + URI: removeRulesURI, CID: removeRulesCID, + OwnerDID: removeRulesAuthorDID, ImageCIDs: test.blobCIDs, + }, + action: moderation.Action{ID: "illegal-removal", Reason: removeRulesIllegal}, + } + purger := &removeRulesPurger{} + reconciler := moderation.NewMediaReconciler(postRulesMediaBinder{bound}, removeRulesInstanceDID, purger) + var blocks []moderation.MediaBlock + var err error + if test.incoming { + blocks, err = reconciler.ReconcileIncomingTx(t.Context(), nil, test.uri, test.ownerDID, test.blobCIDs) + assert.Equal(t, []string{"ActiveRemoval", "InsertNewMediaBlocks"}, bound.calls) + } else { + blocks, err = reconciler.ReconcileTx(t.Context(), nil, test.uri) + assert.Equal(t, []string{"ActiveRemoval", test.readCall, "InsertNewMediaBlocks"}, bound.calls) + } + require.NoError(t, err) + wantBlocks := make([]moderation.MediaBlock, 0, len(test.blobCIDs)) + wantPurges := make([]removeRulesOwnerPurge, 0, len(test.blobCIDs)) + for _, blobCID := range test.blobCIDs { + wantBlocks = append(wantBlocks, moderation.MediaBlock{OwnerDID: test.ownerDID, BlobCID: blobCID, ActionID: bound.action.ID}) + wantPurges = append(wantPurges, removeRulesOwnerPurge{test.ownerDID, blobCID}) + } + assert.Equal(t, wantBlocks, bound.blocks, "only owner-scoped blocks may be inserted") + assert.Equal(t, wantBlocks, blocks, "only owner-scoped blocks may be returned") + reconciler.Purge(blocks) + assert.Equal(t, wantPurges, purger.ownerPurges, "new blobs must use PurgeOwnerBlob") + assert.Empty(t, purger.blobPurges, "reconciliation must never call PurgeBlob") + }) + } +} + +// The store writes an empty owner as an every-owner block, so reconciliation +// must refuse an empty owner rather than insert anything. +func TestReconcileRefusesEmptyOwnerDID(t *testing.T) { + for _, test := range []struct { + name, uri string + incoming bool + }{ + {name: "indexed post", uri: postRulesURI}, + {name: "indexed comment", uri: removeRulesURI}, + {name: "incoming post blobs not indexed", uri: postRulesURI, incoming: true}, + } { + t.Run(test.name, func(t *testing.T) { + blobCIDs := []string{removeRulesFirstImage} + bound := &postRulesMediaTransaction{ + post: moderation.IndexedPost{URI: postRulesURI, CID: postRulesCID, BlobCIDs: blobCIDs}, + comment: &moderation.IndexedComment{URI: removeRulesURI, CID: removeRulesCID, ImageCIDs: blobCIDs}, + action: moderation.Action{ID: "spam-removal", Reason: removeRulesSpam}, + } + reconciler := moderation.NewMediaReconciler(postRulesMediaBinder{bound}, removeRulesInstanceDID, nil) + var blocks []moderation.MediaBlock + var err error + if test.incoming { + blocks, err = reconciler.ReconcileIncomingTx(t.Context(), nil, test.uri, "", blobCIDs) + } else { + blocks, err = reconciler.ReconcileTx(t.Context(), nil, test.uri) + } + require.ErrorIs(t, err, moderation.ErrInvalidSubject) + assert.Empty(t, blocks) + assert.NotContains(t, bound.calls, "InsertNewMediaBlocks", "an empty owner must never reach the store") + }) + } +} + var _ moderation.MediaTransaction = (*postRulesMediaTransaction)(nil) diff --git a/tests/e2e/moderation_contract_test.go b/tests/e2e/moderation_contract_test.go index f6e7e0b..316cfbe 100644 --- a/tests/e2e/moderation_contract_test.go +++ b/tests/e2e/moderation_contract_test.go @@ -4,6 +4,8 @@ package e2e import ( "context" + "encoding/json" + "errors" "fmt" "net/http" "net/url" @@ -244,6 +246,8 @@ func TestModerationCommentRemovalContract(t *testing.T) { func TestModerationPostRemovalContract(t *testing.T) { p := newPipeline(t) author := p.IndexedAccount(t, "mpa") + authorToken := p.AppView.SignIn(t, author) + authorView := p.AppView.As(authorToken) community := indexedCommunity(t, p, "mpa", author.DID) admin := testkit.ModerationAdmin(t, 1) rkey := testkit.TID() @@ -282,6 +286,64 @@ func TestModerationPostRemovalContract(t *testing.T) { } return response.Posts[0], nil } + // GetBinary preserves the actual response bytes so content leakage is checked + // across the whole response, not just in the decoded post union member. + readAuthorPost := func() (map[string]any, []byte, error) { + path := "/xrpc/social.coves.community.post.get?" + url.Values{"uris": {uri}}.Encode() + response, err := authorView.GetBinary(t.Context(), path) + if err != nil { + return nil, nil, err + } + var out struct { + Posts []map[string]any `json:"posts"` + } + if err := json.Unmarshal(response.Body, &out); err != nil { + return nil, nil, fmt.Errorf("decoding author post.get: %w", err) + } + if len(out.Posts) != 1 { + return nil, nil, fmt.Errorf("author post.get returned %d union members for one URI", len(out.Posts)) + } + return out.Posts[0], response.Body, nil + } + requireModeratedAuthorPost := func(forbiddenContent ...string) { + t.Helper() + post, body, err := readAuthorPost() + require.NoError(t, err) + require.Equal(t, "social.coves.community.post.defs#moderatedPost", post["$type"]) + require.Equal(t, uri, post["uri"]) + moderation, ok := post["moderation"].(map[string]any) + require.True(t, ok, "author's moderated post must carry moderation details") + require.Equal(t, "removed", moderation["state"]) + sources, ok := moderation["sources"].([]any) + require.True(t, ok, "moderation sources must be an array") + require.Len(t, sources, 1) + source, ok := sources[0].(map[string]any) + require.True(t, ok, "moderation source must be an object") + require.Equal(t, communityInstanceDID, source["authorityDid"]) + scope, ok := source["scope"].(map[string]any) + require.True(t, ok, "moderation scope must be an object") + require.Equal(t, "instance", scope["kind"]) + for _, key := range []string{"record", "title", "embed"} { + require.NotContains(t, post, key, "moderated post must not expose content") + } + for _, text := range forbiddenContent { + require.NotContains(t, string(body), text, "moderated response leaked post content") + } + } + readAuthorFeedURIs := func(method string, params url.Values) []string { + t.Helper() + var feed struct { + Feed []feedItemView `json:"feed"` + } + require.NoError(t, authorView.Query(t.Context(), method, params, &feed)) + uris := make([]string, 0, len(feed.Feed)) + for _, item := range feed.Feed { + uris = append(uris, item.Post.URI) + } + return uris + } + authorFeedParams := url.Values{"actor": {author.DID}, "limit": {"25"}} + communityFeedParams := url.Values{"community": {community.DID}, "sort": {"new"}, "limit": {"50"}} var imageURL string p.Await(t, "the accepted image post to serve through post.get", func() (bool, error) { post, err := readPost() @@ -308,12 +370,16 @@ func TestModerationPostRemovalContract(t *testing.T) { return ok && imageURL != "", nil }) require.Contains(t, imageURL, image.CID()) - requireServesImage(t, p, "accepted post image", imageURL) + servedImage := requireServesImage(t, p, "accepted post image", imageURL) + require.Equal(t, "public, max-age=86400", servedImage.Header.Get("Cache-Control")) + require.NotContains(t, servedImage.Header.Get("Cache-Control"), "s-maxage") parsedImage, err := url.Parse(imageURL) require.NoError(t, err) require.True(t, strings.HasPrefix(parsedImage.Path, "/img/"), "the image must be served through the AppView proxy") require.Contains(t, communityFeedURIs(t, p, community.DID), uri, "the accepted post must appear in the feed before its removal can prove exclusion") + require.Contains(t, readAuthorFeedURIs("social.coves.actor.getPosts", authorFeedParams), uri) + require.Contains(t, readAuthorFeedURIs("social.coves.communityFeed.getCommunity", communityFeedParams), uri) p.Await(t, "search to find the accepted post before removal", func() (bool, error) { search, err := queryPostSearch(p, needle, community.DID) if err != nil { @@ -335,6 +401,11 @@ func TestModerationPostRemovalContract(t *testing.T) { node, found := thread.find(commentURI) return found && node.Comment.Record["content"] == commentText, nil }, withReadCadence()) + // The wait above can spend 19 of getComments' 20 reads per minute, and the + // removal phase reads the thread three more times. + p.FreshReadQuota(t, "removed-post-thread") + // As copies the client IP, and FreshReadQuota just replaced it. + authorView = p.AppView.As(authorToken) stateToken := admin.ServiceAuth(t, communityInstanceDID, subjectStateMethod) readState := func() (moderationSubjectStateResponse, error) { @@ -401,19 +472,35 @@ func TestModerationPostRemovalContract(t *testing.T) { removed, err := moderated() require.NoError(t, err) require.True(t, removed, "post.get must serve a content-free moderatedPost immediately after removal") + requireModeratedAuthorPost(title, content) missingRootURI := authorPostURI(author.DID, testkit.TID()) _, missingRootErr := p.Thread(context.Background(), missingRootURI, nil) missingRoot := requireXRPCRefusal(t, missingRootErr, http.StatusNotFound, "RootNotFound", "a never-indexed post thread") _, removedRootErr := p.Thread(context.Background(), uri, nil) removedRoot := requireXRPCRefusal(t, removedRootErr, http.StatusNotFound, "RootNotFound", "an instance-removed post thread") require.Equal(t, missingRoot.XRPCError, removedRoot.XRPCError) + authorRootErr := authorView.Query(t.Context(), "social.coves.community.comment.getComments", + url.Values{"post": {uri}}, nil) + authorRoot := requireXRPCRefusal(t, authorRootErr, http.StatusNotFound, "RootNotFound", "the author's instance-removed post thread") + require.Equal(t, removedRoot.XRPCError, authorRoot.XRPCError) + require.NotContains(t, readAuthorFeedURIs("social.coves.actor.getPosts", authorFeedParams), uri, + "the author must not find an instance-removed post in their own feed") + require.NotContains(t, readAuthorFeedURIs("social.coves.communityFeed.getCommunity", communityFeedParams), uri, + "the author must not find an instance-removed post in the community feed") require.NotContains(t, communityFeedURIs(t, p, community.DID), uri) search, err := queryPostSearch(p, needle, community.DID) require.NoError(t, err) require.Empty(t, search.Feed, "the removed post must not appear in search for its unique title") + var blockedImageHeaders http.Header pathBlocked := func(path string) (bool, error) { _, err := p.AppView.GetBinary(context.Background(), path) if testkit.IsStatus(err, http.StatusNotFound) { + if path == parsedImage.Path { + var statusError *testkit.StatusError + if errors.As(err, &statusError) { + blockedImageHeaders = statusError.Header + } + } return true, nil } if err != nil { @@ -424,13 +511,16 @@ func TestModerationPostRemovalContract(t *testing.T) { p.Await(t, "the removed post's cached image to return 404", func() (bool, error) { return pathBlocked(parsedImage.Path) }) + require.NotNil(t, blockedImageHeaders, "the blocked image's 404 response headers were not captured") + require.Equal(t, "no-store", blockedImageHeaders.Get("Cache-Control")) editImage := author.UploadBlob(t, testkit.TestPNG(96, 96), "image/png") require.NotEqual(t, image.CID(), editImage.CID()) editImagePath := strings.Replace(parsedImage.Path, image.CID(), editImage.CID(), 1) require.NotEqual(t, parsedImage.Path, editImagePath) editedTitle := title + " edited" - editedCID := writePost(editedTitle, "edited while removed", image, editImage) + editedContent := "edited while removed" + editedCID := writePost(editedTitle, editedContent, image, editImage) require.NotEqual(t, createdCID, editedCID) // Observe the edit in the indexed post before testing the overlay or media // reconciliation. A still-hidden pre-edit view proves neither behavior. @@ -442,10 +532,12 @@ func TestModerationPostRemovalContract(t *testing.T) { } return indexed.State.CurrentSubject.CID == editedCID, nil }) + awaitStatus(t, p, uri, community.DID, "pending_reacceptance", + "the author's edit to invalidate the old community acceptance") // The edited CID is not admitted yet, so an anonymous viewer gets notFound: - // removal never widens access to an unadmitted CID. The author's #moderatedPost - // view in this state is proven at T1 by TestModerationPostConsumerEditReconcilesNewImages, - // because this tier cannot hold an AppView-sealed OAuth session. + // removal never widens access to an unadmitted CID. The author could read + // their own pending post before the removal, so they keep a view, but only a + // content-free #moderatedPost, even after the edit reaches the index. p.Holds(t, "the unadmitted edited post to read as notFound anonymously with its newly added image blocked", func() (bool, error) { hidden, err := notFoundAnonymously() if err != nil || !hidden { @@ -453,8 +545,7 @@ func TestModerationPostRemovalContract(t *testing.T) { } return pathBlocked(editImagePath) }) - awaitStatus(t, p, uri, community.DID, "pending_reacceptance", - "the author's edit to invalidate the old community acceptance") + requireModeratedAuthorPost(title, content, editedTitle, editedContent) // The community account can update its acceptance at the same subject rkey, // so restore can be checked through the public thread rather than an @@ -497,6 +588,8 @@ func TestModerationPostRemovalContract(t *testing.T) { return ok && post["uri"] == uri && post["$type"] == nil && record["title"] == editedTitle, nil }) p.FreshReadQuota(t, "restored-post-thread") + // As copies the client IP, and FreshReadQuota just replaced it. + authorView = p.AppView.As(authorToken) p.Await(t, "the restored post's comment thread to serve publicly", func() (bool, error) { thread, err := p.Thread(context.Background(), uri, nil) if err != nil { @@ -505,6 +598,15 @@ func TestModerationPostRemovalContract(t *testing.T) { node, found := thread.find(commentURI) return thread.Post.URI == uri && found && node.Comment.Record["content"] == commentText, nil }, withReadCadence()) + restoredAuthorPost, _, err := readAuthorPost() + require.NoError(t, err) + require.Equal(t, uri, restoredAuthorPost["uri"]) + require.NotContains(t, restoredAuthorPost, "$type", "restoration must serve a post view, not a moderated placeholder") + restoredRecord, ok := restoredAuthorPost["record"].(map[string]any) + require.True(t, ok, "restored author's post must carry a record") + require.Equal(t, editedTitle, restoredRecord["title"]) + require.Equal(t, editedContent, restoredRecord["content"]) + require.NotContains(t, restoredAuthorPost, "moderation") require.Contains(t, communityFeedURIs(t, p, community.DID), uri, "restoration must return the re-accepted post to the community feed") // Positive control for the blocked-path checks above: both paths are real, diff --git a/tests/e2e/user_contract_test.go b/tests/e2e/user_contract_test.go index 4543f96..647d36c 100644 --- a/tests/e2e/user_contract_test.go +++ b/tests/e2e/user_contract_test.go @@ -125,7 +125,7 @@ func (p *pipeline) profileWithImages(ctx context.Context, actor string) (Profile // uploaded. The proxy re-encodes and resizes by preset, so a byte comparison // would be asserting the image pipeline's output rather than its reachability, // and would fail the day a preset changed. -func requireServesImage(t *testing.T, p *pipeline, kind, rawURL string) { +func requireServesImage(t *testing.T, p *pipeline, kind, rawURL string) testkit.BinaryResponse { t.Helper() appview, err := url.Parse(testkit.Endpoints().AppView.BaseURL) @@ -148,6 +148,7 @@ func requireServesImage(t *testing.T, p *pipeline, kind, rawURL string) { require.Truef(t, strings.HasPrefix(resp.ContentType, "image/"), "the %s URL served content type %q rather than an image/*: a client will not render it, "+ "and an HTML error page returned with a 200 looks exactly like this", kind, resp.ContentType) + return resp } // blobRefValue renders a testkit blob reference the way a record embeds it. diff --git a/tests/testkit/appview.go b/tests/testkit/appview.go index d81161d..736a93c 100644 --- a/tests/testkit/appview.go +++ b/tests/testkit/appview.go @@ -64,6 +64,8 @@ type StatusError struct { XRPCShaped bool // Body is the raw response, truncated to maxErrorBody. Body string + // Header is a copy of the response headers, including cache policy on errors. + Header http.Header } func (e *StatusError) Error() string { @@ -333,12 +335,13 @@ func (c *XRPCClient) Get(ctx context.Context, path string, out any) error { return c.do(req, path, out) } -// BinaryResponse is a non-JSON response: what was served, and enough about it -// to assert the service really served content rather than merely not failing. +// BinaryResponse is a non-JSON response: what was served, including its headers, +// and enough about it to assert the service served content rather than merely not failing. type BinaryResponse struct { Status int ContentType string Body []byte + Header http.Header } // GetBinary fetches a plain path and returns the raw response. @@ -354,11 +357,10 @@ type BinaryResponse struct { // against — a proxy that cannot reach the blob store — is upstream of the // status code the proxy chooses to report. // -// So this returns the three facts a caller needs to make the real claim -// (status, content type, bytes) rather than folding them into a bool. The body -// is bounded: a test asserting an image is non-empty does not need to buffer an -// arbitrarily large one, and an unbounded read here would make a runaway -// response a hang instead of a failure. +// So this returns status, content type, bytes and response headers rather than +// folding them into a bool. The body is bounded: a test asserting an image is +// non-empty does not need to buffer an arbitrarily large one, and an unbounded +// read here would make a runaway response a hang instead of a failure. // // Unlike Get, a non-2xx is returned as a StatusError, so callers keep the // familiar testkit.IsStatus handling. @@ -399,6 +401,7 @@ func (c *XRPCClient) GetBinary(ctx context.Context, path string) (BinaryResponse Status: resp.StatusCode, ContentType: resp.Header.Get("Content-Type"), Body: body, + Header: resp.Header.Clone(), }, nil } @@ -459,6 +462,7 @@ func newStatusError(nsid string, resp *http.Response) *StatusError { Method: nsid, StatusCode: resp.StatusCode, Body: strings.TrimSpace(string(body)), + Header: resp.Header.Clone(), } if readErr != nil { // Said rather than swallowed. An empty Body reads as "the service diff --git a/tests/testkit/appview_test.go b/tests/testkit/appview_test.go index 395a161..6d28258 100644 --- a/tests/testkit/appview_test.go +++ b/tests/testkit/appview_test.go @@ -181,6 +181,38 @@ func TestStatusError_HandlesANonXRPCBody(t *testing.T) { assert.Contains(t, err.Error(), "upstream is down") } +func TestXRPCClient_GetBinaryExposesSuccessAndErrorHeaders(t *testing.T) { + stub := newStubService(t, func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/img/blocked" { + w.Header().Set("Cache-Control", "no-store") + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte("blob not found")) + return + } + w.Header().Set("Cache-Control", "public, max-age=86400") + w.Header().Set("ETag", `"avatar-blob"`) + w.Header().Set("Content-Type", "image/jpeg") + _, _ = w.Write([]byte("image bytes")) + }) + client := NewXRPCClient(stub.URL) + + served, err := client.GetBinary(t.Context(), "/img/served") + require.NoError(t, err) + assert.Equal(t, http.StatusOK, served.Status) + assert.Equal(t, "public, max-age=86400", served.Header.Get("Cache-Control")) + assert.Equal(t, `"avatar-blob"`, served.Header.Get("ETag")) + assert.Equal(t, "image/jpeg", served.ContentType) + assert.Equal(t, []byte("image bytes"), served.Body) + + _, err = client.GetBinary(t.Context(), "/img/blocked") + require.Error(t, err) + var statusError *StatusError + require.ErrorAs(t, err, &statusError) + assert.Equal(t, http.StatusNotFound, statusError.StatusCode) + assert.Equal(t, "no-store", statusError.Header.Get("Cache-Control")) + assert.Equal(t, "blob not found", statusError.Body) +} + func TestStatusOf_IsZeroForTransportFailures(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {})) address := server.URL