From d86d802ad51f6c0d304136dedabdacdbd9cd7e77 Mon Sep 17 00:00:00 2001 From: Bretton Date: Thu, 24 Sep 2026 02:42:54 -0700 Subject: [PATCH] feat(moderation): add instance-admin authority and getSubjectState MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Introduce the internal/core/moderation domain with an operator-managed admin DID allowlist, DID-verified admin authorization, and the getSubjectState read over content the AppView already indexes. No mutation storage in this chunk. - MODERATION_ADMINS allowlist config (comma-separated DIDs; empty grants nobody; a non-DID entry fails Load naming the variable), added to the test clear list, with prod compose and .env.prod.example entries. - RequireInstanceAdmin middleware accepting OAuth sealed sessions and indigo-validated PDS service JWTs (aud INSTANCE_DID, lxm bound to the request NSID). Answers 401/403/503 without logging tokens or IPs; caller-supplied actor/authority fields are never read. - Shared authenticateSealedSession helper that fails closed. - moderation.Service.GetSubjectState and the SubjectReader over the existing post and comment repos (present/deleted/unavailable record state, opaque v0 version token). - Sentinels for every §14.5 error code and their HTTP mapping; authorization precedes subject errors. - social.coves.moderation.getSubjectState route, registered as authRequired; moderation NSIDs stay AppView-served and out of the PDS OAuth scope list. - CI bootstrap of two moderation admin PDS accounts whose DIDs are passed to the AppView as MODERATION_ADMINS. - T2 subject-state contract, excluded from test-e2e-dev. Co-Authored-By: Claude Opus 5.5 (1M context) --- .env.ci | 10 + .env.dev | 8 + .env.prod.example | 7 + Makefile | 14 +- cmd/server/moderation_wiring.go | 16 + cmd/server/moderation_wiring_test.go | 99 ++++ cmd/server/oauth_scopes_test.go | 6 + cmd/server/routes.go | 1 + cmd/server/wiring.go | 20 +- docker-compose.ci.yml | 8 +- docker-compose.prod.yml | 4 + internal/api/handlers/moderation/errors.go | 25 + .../handlers/moderation/get_subject_state.go | 88 ++++ .../moderation/get_subject_state_test.go | 170 +++++++ internal/api/middleware/admin_auth.go | 123 +++++ internal/api/middleware/admin_auth_test.go | 455 ++++++++++++++++++ internal/api/middleware/auth.go | 120 +++-- internal/api/routes/harness_test.go | 14 + internal/api/routes/moderation.go | 18 + ...deration_subject_state_integration_test.go | 184 +++++++ internal/api/routes/registration_test.go | 4 + internal/config/config.go | 19 + internal/config/moderation_admins_test.go | 53 ++ internal/config/testing.go | 1 + internal/core/moderation/authority.go | 19 + internal/core/moderation/authority_test.go | 44 ++ internal/core/moderation/errors.go | 23 + internal/core/moderation/harness_test.go | 14 + internal/core/moderation/interfaces.go | 34 ++ internal/core/moderation/service.go | 49 ++ internal/core/moderation/service_test.go | 121 +++++ internal/core/moderation/subject_reader.go | 53 ++ .../subject_reader_integration_test.go | 153 ++++++ internal/core/moderation/types.go | 78 +++ internal/core/moderation/version.go | 5 + scripts/ci-bootstrap.sh | 124 +++-- scripts/lib/ci-stack.sh | 24 +- .../moderation_subject_state_contract_test.go | 89 ++++ tests/testkit/pds.go | 48 ++ 39 files changed, 2259 insertions(+), 86 deletions(-) create mode 100644 cmd/server/moderation_wiring.go create mode 100644 cmd/server/moderation_wiring_test.go create mode 100644 internal/api/handlers/moderation/errors.go create mode 100644 internal/api/handlers/moderation/get_subject_state.go create mode 100644 internal/api/handlers/moderation/get_subject_state_test.go create mode 100644 internal/api/middleware/admin_auth.go create mode 100644 internal/api/middleware/admin_auth_test.go create mode 100644 internal/api/routes/harness_test.go create mode 100644 internal/api/routes/moderation.go create mode 100644 internal/api/routes/moderation_subject_state_integration_test.go create mode 100644 internal/config/moderation_admins_test.go create mode 100644 internal/core/moderation/authority.go create mode 100644 internal/core/moderation/authority_test.go create mode 100644 internal/core/moderation/errors.go create mode 100644 internal/core/moderation/harness_test.go create mode 100644 internal/core/moderation/interfaces.go create mode 100644 internal/core/moderation/service.go create mode 100644 internal/core/moderation/service_test.go create mode 100644 internal/core/moderation/subject_reader.go create mode 100644 internal/core/moderation/subject_reader_integration_test.go create mode 100644 internal/core/moderation/types.go create mode 100644 internal/core/moderation/version.go create mode 100644 tests/e2e/moderation_subject_state_contract_test.go diff --git a/.env.ci b/.env.ci index edbefac..1a64fba 100644 --- a/.env.ci +++ b/.env.ci @@ -67,6 +67,16 @@ PDS_DID_PLC_URL=http://localhost:3002 PDS_INSTANCE_HANDLE=testuser123.local.coves.dev PDS_INSTANCE_PASSWORD=test-password-123 +# Instance moderation admins. scripts/ci-bootstrap.sh creates both accounts on +# the PDS before the AppView boots and writes their DIDs as MODERATION_ADMINS +# into .ci-out/moderation-admins-.env, which the appview +# service loads (docker-compose.ci.yml). tests/testkit ModerationAdmin logs in +# with these credentials. MODERATION_ADMINS itself is deliberately not set here. +CI_MODERATION_ADMIN_ONE_HANDLE=modadmin-one.local.coves.dev +CI_MODERATION_ADMIN_ONE_PASSWORD=moderation-admin-one-password +CI_MODERATION_ADMIN_TWO_HANDLE=modadmin-two.local.coves.dev +CI_MODERATION_ADMIN_TWO_PASSWORD=moderation-admin-two-password + # ============================================================================= # The federated PDS and the relay (CI only — no dev-stack equivalent) # ============================================================================= diff --git a/.env.dev b/.env.dev index f5247b6..00da2a7 100644 --- a/.env.dev +++ b/.env.dev @@ -206,6 +206,14 @@ PDS_INSTANCE_PASSWORD=test-password-123 # - did:plc:igjbg5cex7poojsniebvmafb = test-aggregator.local.coves.dev (dev) TRUSTED_AGGREGATOR_DIDS=did:plc:yyf34padpfjknejyutxtionr,did:plc:igjbg5cex7poojsniebvmafb,did:plc:jn4tlbpkdms5tahfrylct5g7 +# ============================================================================= +# Instance Moderation Admins +# ============================================================================= +# Comma-separated DIDs allowed to call the social.coves.moderation.* admin +# endpoints. Empty grants nobody. Every entry must be a DID; a handle fails +# startup. Authorization never uses handles or COMMUNITY_CREATORS. +MODERATION_ADMINS= + # ============================================================================= # Development Settings # ============================================================================= diff --git a/.env.prod.example b/.env.prod.example index aaa868c..546af69 100644 --- a/.env.prod.example +++ b/.env.prod.example @@ -141,6 +141,13 @@ CURSOR_SECRET=CHANGE_ME_CURSOR_SECRET # Comma-separated list. If not set, any authenticated user can create communities. # COMMUNITY_CREATORS=did:plc:abc123,did:plc:def456 +# Instance moderation admins. Comma-separated list of DIDs only (handles are +# rejected at boot). Empty grants admin authority to nobody. Independent of +# COMMUNITY_CREATORS: being allowed to create communities grants no admin +# authority, and vice versa. Use dedicated admin accounts that grant no generic +# scopes to third-party apps. +# MODERATION_ADMINS= + # Optional: trust virtual bridge PDS hosts for bridged identity discovery and # bridge-asserted bridgedStats. Required when deploying Tidepool; use its # public scheme+host URL (comma-separated if multiple bridges are trusted). diff --git a/Makefile b/Makefile index c0af29f..c563cfb 100644 --- a/Makefile +++ b/Makefile @@ -244,6 +244,8 @@ test-e2e-dev: ## T2 against the long-lived DEV stack (debugging only - not how C @echo "$(YELLOW) minus the federation contracts: they need the SECOND PDS and the relay,$(RESET)" @echo "$(YELLOW) which exist only in the hermetic stack (docker-compose.ci.yml). The dev$(RESET)" @echo "$(YELLOW) stack has one PDS and Jetstream wired straight to it.$(RESET)" + @echo "$(YELLOW) minus the moderation contracts: they need the bootstrap admin accounts$(RESET)" + @echo "$(YELLOW) that scripts/ci-bootstrap.sh provisions only in the hermetic stack.$(RESET)" @# run_pipeline_tier is the ONE definition of how T2 is invoked — the gate, @# 'make test-e2e' and this hatch all call it, so the flags cannot drift @# apart. Sourced here rather than copied for exactly that reason. @@ -263,8 +265,16 @@ test-e2e-dev: ## T2 against the long-lived DEV stack (debugging only - not how C @# A federation contract named without it is not a silent pass either: @# testkit.NewFederatedPDS fatals on the spot, naming PDS2_URL and this @# hatch. - @bash -c 'source ./scripts/lib/runner-ready.sh && run_pipeline_tier -skip "^TestReliability|Federat"' - @echo "$(GREEN)✓ Pipeline tier complete (against the dev stack; no reliability suite, no federation contracts)$(RESET)" + @# + @# "^TestModeration" catches every instance-moderation contract. They act as + @# the two bootstrap admin accounts (testkit.ModerationAdmin), which + @# scripts/ci-bootstrap.sh provisions and writes into MODERATION_ADMINS only + @# for the hermetic stack. The dev stack has neither the accounts nor the + @# allowlist (.env.dev sets MODERATION_ADMINS empty), so these contracts + @# would fatal here on every run. Prefix rather than a list, as with + @# "Federat", so a new moderation contract is excluded without an edit. + @bash -c 'source ./scripts/lib/runner-ready.sh && run_pipeline_tier -skip "^TestReliability|Federat|^TestModeration"' + @echo "$(GREEN)✓ Pipeline tier complete (against the dev stack; no reliability suite, no federation contracts, no moderation contracts)$(RESET)" test-db-reset: ## Reset test database @echo "$(GREEN)Resetting test database...$(RESET)" diff --git a/cmd/server/moderation_wiring.go b/cmd/server/moderation_wiring.go new file mode 100644 index 0000000..9c75cef --- /dev/null +++ b/cmd/server/moderation_wiring.go @@ -0,0 +1,16 @@ +package main + +import ( + oauthlib "github.com/bluesky-social/indigo/atproto/auth/oauth" + + "Coves/internal/api/middleware" + "Coves/internal/config" + "Coves/internal/core/moderation" +) + +// buildInstanceAdminMiddleware gates the moderation endpoints on the +// operator-managed MODERATION_ADMINS allowlist. +func buildInstanceAdminMiddleware(cfg *config.Config, unsealer middleware.SessionUnsealer, store oauthlib.ClientAuthStore, validator middleware.ServiceAuthValidator) *middleware.InstanceAdminMiddleware { + return middleware.NewInstanceAdminMiddleware( + unsealer, store, validator, moderation.NewAllowlistAuthority(cfg.Moderation.Admins)) +} diff --git a/cmd/server/moderation_wiring_test.go b/cmd/server/moderation_wiring_test.go new file mode 100644 index 0000000..f007c61 --- /dev/null +++ b/cmd/server/moderation_wiring_test.go @@ -0,0 +1,99 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" + + "Coves/internal/api/middleware" + "Coves/internal/atproto/oauth" + "Coves/internal/config" + indigooauth "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type wiringSessionUnsealer map[string]*oauth.SealedSession + +func (unsealer wiringSessionUnsealer) UnsealSession(token string) (*oauth.SealedSession, error) { + if session, ok := unsealer[token]; ok { + return session, nil + } + return nil, fmt.Errorf("unknown session") +} + +type wiringOAuthStore struct { + indigooauth.ClientAuthStore + sessions map[string]*indigooauth.ClientSessionData +} + +func (store *wiringOAuthStore) GetSession(_ context.Context, did syntax.DID, sessionID string) (*indigooauth.ClientSessionData, error) { + if session, ok := store.sessions[did.String()+":"+sessionID]; ok { + return session, nil + } + return nil, oauth.ErrSessionNotFound +} + +func TestBuildInstanceAdminMiddlewareUsesModerationAdminsOnly(t *testing.T) { + const ( + creatorDID = "did:plc:creator" + adminDID = "did:plc:admin" + path = "/xrpc/social.coves.moderation.getSubjectState" + ) + for _, test := range []struct { + name string + admins []string + caller string + wantStatus int + wantCode string + wantCalled bool + }{ + {name: "community creator is not admin", admins: []string{adminDID}, caller: creatorDID, wantStatus: http.StatusForbidden, wantCode: "Forbidden"}, + {name: "listed admin is authorized", admins: []string{adminDID}, caller: adminDID, wantStatus: http.StatusNoContent, wantCalled: true}, + {name: "empty admin list grants nobody", caller: adminDID, wantStatus: http.StatusForbidden, wantCode: "Forbidden"}, + } { + t.Run(test.name, func(t *testing.T) { + cfg := &config.Config{ + Instance: config.InstanceConfig{AllowedCommunityCreators: []string{creatorDID}}, + Moderation: config.ModerationConfig{Admins: test.admins}, + } + unsealer := wiringSessionUnsealer{} + store := &wiringOAuthStore{sessions: map[string]*indigooauth.ClientSessionData{}} + for _, did := range []string{creatorDID, adminDID} { + token := "session-for-" + did + unsealer[token] = &oauth.SealedSession{DID: did, SessionID: "browser", ExpiresAt: time.Now().Add(time.Hour).Unix()} + store.sessions[did+":browser"] = &indigooauth.ClientSessionData{AccountDID: syntax.DID(did), SessionID: "browser"} + } + + gate := buildInstanceAdminMiddleware(cfg, unsealer, store, nil) + if !assert.NotNil(t, gate, "the server must construct the instance-admin gate") { + return + } + called := false + var authenticatedDID string + handler := gate.RequireInstanceAdmin(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + called = true + authenticatedDID = middleware.GetUserDID(r) + w.WriteHeader(http.StatusNoContent) + })) + request := httptest.NewRequest(http.MethodGet, path, nil) + request.Header.Set("Authorization", "Bearer session-for-"+test.caller) + response := httptest.NewRecorder() + handler.ServeHTTP(response, request) + assert.Equal(t, test.wantStatus, response.Code) + assert.Equal(t, test.wantCalled, called) + if test.wantCalled { + assert.Equal(t, test.caller, authenticatedDID) + } else if assert.NotEmpty(t, response.Body.Bytes(), "refusal must include a JSON error") { + var body map[string]any + require.NoError(t, json.Unmarshal(response.Body.Bytes(), &body)) + assert.Equal(t, test.wantCode, body["error"]) + } + }) + } +} diff --git a/cmd/server/oauth_scopes_test.go b/cmd/server/oauth_scopes_test.go index 45bb0d8..9113d79 100644 --- a/cmd/server/oauth_scopes_test.go +++ b/cmd/server/oauth_scopes_test.go @@ -77,3 +77,9 @@ func TestOAuthScopes_RetainCommunityPostThroughTheDrain(t *testing.T) { assert.Containsf(t, post, "action=delete", "the retained community.post grant must include action=delete — the drain's whole job is deleting those records") } + +func TestOAuthScopes_ExcludeAppViewModeration(t *testing.T) { + for _, scope := range oauthScopes() { + assert.NotContains(t, scope, "social.coves.moderation", "moderation endpoints are AppView-served, not PDS-proxied") + } +} diff --git a/cmd/server/routes.go b/cmd/server/routes.go index 3a19299..7a92fee 100644 --- a/cmd/server/routes.go +++ b/cmd/server/routes.go @@ -87,6 +87,7 @@ func registerXRPCRoutes(r chi.Router, app *application) { routes.RegisterUserBlockRoutes(r, app.userBlockService, app.authMiddleware) routes.RegisterCommentRoutes(r, app.commentService, app.authMiddleware) routes.RegisterAdminReportRoutes(r, app.adminReportService, app.authMiddleware) + routes.RegisterModerationRoutes(r, app.moderationService, app.instanceAdminAuth) routes.RegisterCommunitySuggestionRoutes(r, app.communitySuggestionService, app.authMiddleware, app.cfg.Instance.AllowedCommunityCreators) diff --git a/cmd/server/wiring.go b/cmd/server/wiring.go index 5117bee..347c0ad 100644 --- a/cmd/server/wiring.go +++ b/cmd/server/wiring.go @@ -26,6 +26,7 @@ import ( "Coves/internal/core/communitysuggestions" "Coves/internal/core/discover" "Coves/internal/core/imageproxy" + "Coves/internal/core/moderation" "Coves/internal/core/posts" "Coves/internal/core/timeline" "Coves/internal/core/unfurl" @@ -90,12 +91,14 @@ type application struct { credentialCipher *credentialcipher.Cipher // Identity and authentication - identityResolver identity.Resolver - oauthClient *oauth.OAuthClient - oauthStore *oauth.MobileAwareStoreWrapper - oauthHandler *oauth.OAuthHandler - authMiddleware *middleware.OAuthAuthMiddleware - dualAuth *middleware.DualAuthMiddleware + identityResolver identity.Resolver + oauthClient *oauth.OAuthClient + oauthStore *oauth.MobileAwareStoreWrapper + oauthHandler *oauth.OAuthHandler + authMiddleware *middleware.OAuthAuthMiddleware + dualAuth *middleware.DualAuthMiddleware + serviceAuthValidator middleware.ServiceAuthValidator + instanceAdminAuth *middleware.InstanceAdminMiddleware // Repositories reused outside their own service (Jetstream consumers, // route options). @@ -135,6 +138,7 @@ type application struct { commentService comments.Service userBlockService userblocks.Service adminReportService adminreports.Service + moderationService moderation.Service communitySuggestionService communitysuggestions.Service feedService communityFeeds.Service timelineService timeline.Service @@ -434,6 +438,9 @@ func (a *application) buildServices(ctx context.Context) error { a.apiKeyService = aggregators.NewAPIKeyService(a.aggregatorRepo, a.oauthClient.ClientApp) a.buildDualAuth() + a.instanceAdminAuth = buildInstanceAdminMiddleware(a.cfg, a.oauthClient, a.oauthStore, a.serviceAuthValidator) + a.moderationService = moderation.NewService(moderation.NewRepositorySubjectReader(a.postRepo, a.commentRepo)) + slog.Info("instance moderation admins configured", "admin_count", len(a.cfg.Moderation.Admins)) // The SSRF hatch is open only in dev, where the links a developer pastes and // the fixtures the test suite serves both live on the developer's own @@ -765,6 +772,7 @@ func (a *application) buildDualAuth() { Dir: identityDir, TimestampLeeway: 30 * time.Second, } + a.serviceAuthValidator = serviceValidator a.dualAuth = middleware.NewDualAuthMiddleware( a.oauthClient, // SessionUnsealer for OAuth diff --git a/docker-compose.ci.yml b/docker-compose.ci.yml index 3446506..bc2c562 100644 --- a/docker-compose.ci.yml +++ b/docker-compose.ci.yml @@ -90,8 +90,8 @@ # itself. # # Usage: driven by scripts/ci.sh, which stages startup so the relay is crawling -# both PDSes and the instance PDS account exists before the AppView boots. Not -# intended for direct `up`. +# both PDSes and the instance and moderation admin accounts exist before the +# AppView boots. Not intended for direct `up`. services: # Owns the network namespace every other service joins. Does nothing else; @@ -548,6 +548,10 @@ services: network_mode: "service:netns" env_file: - .env.ci + # Bootstrap writes this into the runner's /src/.ci-out bind mount before + # AppView starts. It is absent during the initial infrastructure bring-up. + - path: .ci-out/moderation-admins-${COVES_CI_PROJECT:-coves-ci}.env + required: false environment: # .env.ci silences logs so the go test -json stream stays readable. The # AppView is exempt: when an E2E test fails, these logs are the only diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index bb92d5e..2382a7b 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -128,6 +128,10 @@ services: # Comma-separated list of DIDs TRUSTED_AGGREGATOR_DIDS: ${TRUSTED_AGGREGATOR_DIDS:-} + # Instance moderation admins (comma-separated DIDs). Empty grants nobody + # admin authority. Independent of COMMUNITY_CREATORS. + MODERATION_ADMINS: ${MODERATION_ADMINS:-} + # Trusted virtual bridge PDS hosts. Required for Tidepool-authored # profiles/posts to be indexed and for bridge-asserted bridgedStats to # be accepted. Example: https://tdpl.io diff --git a/internal/api/handlers/moderation/errors.go b/internal/api/handlers/moderation/errors.go new file mode 100644 index 0000000..c7fa7e0 --- /dev/null +++ b/internal/api/handlers/moderation/errors.go @@ -0,0 +1,25 @@ +package moderation + +import ( + "errors" + "log/slog" + "net/http" + + "Coves/internal/api/xrpc" + "Coves/internal/core/moderation" +) + +func writeSubjectStateError(w http.ResponseWriter, err error) { + switch { + case errors.Is(err, moderation.ErrInvalidSubject): + xrpc.WriteError(w, http.StatusBadRequest, "InvalidSubject", "Invalid subject URI") + case errors.Is(err, moderation.ErrModerationUnavailable): + slog.Error("moderation subject state unavailable", "error", err) + xrpc.WriteError(w, http.StatusServiceUnavailable, "ModerationUnavailable", "Moderation state temporarily unavailable") + default: + // Subject AT-URIs, DIDs and database errors carry no credentials, so the + // full error is logged; the response body stays generic. + slog.Error("unexpected moderation subject state failure", "error", err) + xrpc.WriteError(w, http.StatusInternalServerError, "InternalServerError", "An internal error occurred") + } +} diff --git a/internal/api/handlers/moderation/get_subject_state.go b/internal/api/handlers/moderation/get_subject_state.go new file mode 100644 index 0000000..b1a0847 --- /dev/null +++ b/internal/api/handlers/moderation/get_subject_state.go @@ -0,0 +1,88 @@ +// Package moderation serves the social.coves.moderation.* XRPC endpoints. +package moderation + +import ( + "net/http" + + "Coves/internal/api/xrpc" + "Coves/internal/core/moderation" +) + +// GetSubjectStateHandler serves social.coves.moderation.getSubjectState. +type GetSubjectStateHandler struct { + service moderation.Service +} + +// NewGetSubjectStateHandler builds the handler over the moderation service. +func NewGetSubjectStateHandler(service moderation.Service) *GetSubjectStateHandler { + return &GetSubjectStateHandler{service: service} +} + +// HandleGetSubjectState answers the query. +func (h *GetSubjectStateHandler) HandleGetSubjectState(w http.ResponseWriter, r *http.Request) { + subject := r.URL.Query().Get("subject") + if subject == "" { + xrpc.WriteError(w, http.StatusBadRequest, "InvalidSubject", "subject is required") + return + } + state, err := h.service.GetSubjectState(r.Context(), subject) + if err != nil { + writeSubjectStateError(w, err) + return + } + if state == nil { + writeSubjectStateError(w, nil) + return + } + + view := subjectStateView{ + Subject: state.Subject, + Version: state.Version, + Moderation: moderationView{State: state.Moderation.State}, + RecordState: state.RecordState, + } + if state.CurrentSubject != nil { + view.CurrentSubject = &strongRefView{URI: state.CurrentSubject.URI, CID: state.CurrentSubject.CID} + } + if state.LocalRemoval != nil { + view.LocalRemoval = &actionRefView{ServiceDID: state.LocalRemoval.ServiceDID, ActionID: state.LocalRemoval.ActionID} + } + for _, label := range state.LocalLabels { + view.LocalLabels = append(view.LocalLabels, localLabelView{ + Value: label.Value, + Action: actionRefView{ServiceDID: label.Action.ServiceDID, ActionID: label.Action.ActionID}, + }) + } + xrpc.WriteJSON(w, http.StatusOK, struct { + State subjectStateView `json:"state"` + }{State: view}) +} + +type subjectStateView struct { + Subject string `json:"subject"` + Version string `json:"version"` + Moderation moderationView `json:"moderation"` + RecordState moderation.RecordState `json:"recordState"` + CurrentSubject *strongRefView `json:"currentSubject,omitempty"` + LocalRemoval *actionRefView `json:"localRemoval,omitempty"` + LocalLabels []localLabelView `json:"localLabels,omitempty"` +} + +type moderationView struct { + State string `json:"state"` +} + +type strongRefView struct { + URI string `json:"uri"` + CID string `json:"cid"` +} + +type actionRefView struct { + ServiceDID string `json:"serviceDid"` + ActionID string `json:"actionId"` +} + +type localLabelView struct { + Value string `json:"value"` + Action actionRefView `json:"action"` +} diff --git a/internal/api/handlers/moderation/get_subject_state_test.go b/internal/api/handlers/moderation/get_subject_state_test.go new file mode 100644 index 0000000..1183a13 --- /dev/null +++ b/internal/api/handlers/moderation/get_subject_state_test.go @@ -0,0 +1,170 @@ +package moderation + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + "net/http" + "net/http/httptest" + "net/url" + "testing" + + "Coves/internal/core/moderation" + "Coves/internal/validation" + "github.com/bluesky-social/indigo/atproto/atdata" + "github.com/bluesky-social/indigo/atproto/lexicon" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type subjectStateServiceFake struct { + state *moderation.SubjectState + err error + subjects []string +} + +func (service *subjectStateServiceFake) GetSubjectState(_ context.Context, subject string) (*moderation.SubjectState, error) { + service.subjects = append(service.subjects, subject) + return service.state, service.err +} + +func requestGetSubjectState(service *subjectStateServiceFake, subject string, includeSubject bool) *httptest.ResponseRecorder { + target := "/xrpc/social.coves.moderation.getSubjectState" + if includeSubject { + target += "?" + url.Values{"subject": {subject}}.Encode() + } + response := httptest.NewRecorder() + NewGetSubjectStateHandler(service).HandleGetSubjectState(response, httptest.NewRequest(http.MethodGet, target, nil)) + return response +} + +func assertSubjectStateError(t *testing.T, response *httptest.ResponseRecorder, status int, code string) { + t.Helper() + assert.Equal(t, status, response.Code) + if assert.NotEmpty(t, response.Body.Bytes(), "XRPC error must include a JSON body") { + var body map[string]any + require.NoError(t, json.Unmarshal(response.Body.Bytes(), &body)) + assert.Equal(t, code, body["error"]) + } +} + +func TestGetSubjectStateHandlerMissingSubject(t *testing.T) { + service := &subjectStateServiceFake{} + response := requestGetSubjectState(service, "", false) + assertSubjectStateError(t, response, http.StatusBadRequest, "InvalidSubject") + assert.Empty(t, service.subjects, "missing subject must be rejected before calling the service") +} + +func TestGetSubjectStateHandlerServiceErrors(t *testing.T) { + const subject = "at://did:plc:subject/social.coves.community.postv2/3kabc" + for _, test := range []struct { + name string + err error + status int + code string + }{ + {name: "invalid subject", err: fmt.Errorf("subject: %w", moderation.ErrInvalidSubject), status: http.StatusBadRequest, code: "InvalidSubject"}, + {name: "unavailable", err: fmt.Errorf("reading state: %w", moderation.ErrModerationUnavailable), status: http.StatusServiceUnavailable, code: "ModerationUnavailable"}, + } { + t.Run(test.name, func(t *testing.T) { + service := &subjectStateServiceFake{err: test.err} + response := requestGetSubjectState(service, subject, true) + assertSubjectStateError(t, response, test.status, test.code) + assert.Equal(t, []string{subject}, service.subjects) + }) + } +} + +// TestGetSubjectStateHandlerLogsServerFailures swaps the process-global slog +// default, so it must not run in parallel. +func TestGetSubjectStateHandlerLogsServerFailures(t *testing.T) { + const subject = "at://did:plc:subject/social.coves.community.postv2/3kabc" + for _, test := range []struct { + name string + err error + status int + code string + }{ + { + name: "unavailable", + err: fmt.Errorf("%w: %w", moderation.ErrModerationUnavailable, errors.New("pg: connection refused")), + status: http.StatusServiceUnavailable, code: "ModerationUnavailable", + }, + { + name: "unexpected", + err: fmt.Errorf("reading state: %w", errors.New("scan: column count mismatch")), + status: http.StatusInternalServerError, code: "InternalServerError", + }, + } { + t.Run(test.name, func(t *testing.T) { + var logged bytes.Buffer + previous := slog.Default() + slog.SetDefault(slog.New(slog.NewTextHandler(&logged, &slog.HandlerOptions{Level: slog.LevelError}))) + t.Cleanup(func() { slog.SetDefault(previous) }) + + response := requestGetSubjectState(&subjectStateServiceFake{err: test.err}, subject, true) + + assertSubjectStateError(t, response, test.status, test.code) + assert.Contains(t, logged.String(), test.err.Error(), "the underlying error must reach the server log") + assert.NotContains(t, response.Body.String(), test.err.Error(), "the underlying error must not reach the response body") + }) + } +} + +func TestGetSubjectStateHandlerExactLexiconOutput(t *testing.T) { + const ( + subject = "at://did:plc:subject/social.coves.community.postv2/3kabc" + cid = "bafyreib6tbnql2ux3whnfysbzabthaj2vvck53nimhbi5g5a7jgvgr5eqm" + ) + for _, test := range []struct { + name string + recordState moderation.RecordState + current *moderation.StrongRef + wantCurrent bool + }{ + {name: "present", recordState: moderation.RecordStatePresent, current: &moderation.StrongRef{URI: subject, CID: cid}, wantCurrent: true}, + {name: "deleted", recordState: moderation.RecordStateDeleted}, + {name: "unavailable", recordState: moderation.RecordStateUnavailable}, + } { + t.Run(test.name, func(t *testing.T) { + catalog := lexicon.NewBaseCatalog() + require.NoError(t, catalog.LoadDirectory("../../../atproto/lexicon")) + state := &moderation.SubjectState{ + Subject: subject, Version: "v0", Moderation: moderation.ModerationView{State: "clear"}, + RecordState: test.recordState, CurrentSubject: test.current, + } + service := &subjectStateServiceFake{state: state} + response := requestGetSubjectState(service, subject, true) + + assert.Equal(t, []string{subject}, service.subjects, "handler must pass the subject through unchanged") + assert.Equal(t, http.StatusOK, response.Code) + assert.Equal(t, "application/json", response.Header().Get("Content-Type")) + require.NotEmpty(t, response.Body.Bytes(), "successful query must serialize the subject state") + + var got map[string]any + require.NoError(t, json.Unmarshal(response.Body.Bytes(), &got)) + wantState := map[string]any{ + "subject": subject, "version": "v0", "moderation": map[string]any{"state": "clear"}, + "recordState": string(test.recordState), + } + if test.wantCurrent { + wantState["currentSubject"] = map[string]any{"uri": subject, "cid": cid} + } + assert.Equal(t, map[string]any{"state": wantState}, got, "extra fields and null optional fields are forbidden") + + decoded, err := atdata.UnmarshalJSON(response.Body.Bytes()) + require.NoError(t, err) + assert.NoError(t, validation.ValidateData(catalog, decoded, "social.coves.moderation.getSubjectState#output", 0)) + }) + } +} + +func TestGetSubjectStateHandlerPassesExactQuerySubject(t *testing.T) { + const subject = "at://did:plc:verbatim/social.coves.community.comment/3kquery" + service := &subjectStateServiceFake{} + _ = requestGetSubjectState(service, subject, true) + assert.Equal(t, []string{subject}, service.subjects) +} diff --git a/internal/api/middleware/admin_auth.go b/internal/api/middleware/admin_auth.go new file mode 100644 index 0000000..5354a10 --- /dev/null +++ b/internal/api/middleware/admin_auth.go @@ -0,0 +1,123 @@ +package middleware + +import ( + "context" + "encoding/json" + "log/slog" + "net/http" + "strings" + + oauthlib "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/bluesky-social/indigo/atproto/syntax" +) + +// InstanceAdminAuthority answers whether an authenticated DID is on the +// operator-managed instance admin allowlist. +type InstanceAdminAuthority interface { + IsInstanceAdmin(did string) bool +} + +// InstanceAdminMiddleware authenticates a caller by PDS service JWT or OAuth +// sealed session and then requires the DID to be an instance admin. +type InstanceAdminMiddleware struct { + unsealer SessionUnsealer + store oauthlib.ClientAuthStore + validator ServiceAuthValidator + authority InstanceAdminAuthority +} + +// NewInstanceAdminMiddleware builds the admin gate over the existing +// authentication primitives. +func NewInstanceAdminMiddleware(unsealer SessionUnsealer, store oauthlib.ClientAuthStore, validator ServiceAuthValidator, authority InstanceAdminAuthority) *InstanceAdminMiddleware { + return &InstanceAdminMiddleware{unsealer: unsealer, store: store, validator: validator, authority: authority} +} + +// RequireInstanceAdmin refuses any request whose caller is not a verified +// instance admin. +func (m *InstanceAdminMiddleware) RequireInstanceAdmin(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + var token string + if header := r.Header.Get("Authorization"); header != "" { + var ok bool + token, ok = extractBearerToken(header) + if !ok { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Invalid Authorization header", "invalid_authorization_header", "", nil) + return + } + } else if cookie, err := r.Cookie("coves_session"); err == nil { + token = cookie.Value + } + if token == "" { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Missing authentication", "missing_authentication", "", nil) + return + } + + var did, authMethod string + ctx := r.Context() + if isJWTFormat(token) { + methodName, ok := strings.CutPrefix(r.URL.Path, "/xrpc/") + if !ok || m.validator == nil { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Invalid service authentication", "missing_service_validator_or_method", "", nil) + return + } + method, err := syntax.ParseNSID(methodName) + if err != nil { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Invalid service method", "invalid_service_method", "", err) + return + } + issuer, err := m.validator.Validate(ctx, token, &method) + if err != nil { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Invalid or expired service JWT", "invalid_service_jwt", "", err) + return + } + did, authMethod = issuer.String(), AuthMethodServiceJWT + } else { + auth := authenticateSealedSession(ctx, token, m.unsealer, m.store) + if auth.failure == sealedSessionStoreFailure { + refuseInstanceAdmin(w, r, http.StatusServiceUnavailable, "ModerationUnavailable", "Session lookup temporarily unavailable", auth.failure.reason(), auth.did, auth.err) + return + } + if auth.failure != sealedSessionNoFailure { + refuseInstanceAdmin(w, r, http.StatusUnauthorized, "AuthRequired", "Invalid or expired session", auth.failure.reason(), auth.did, nil) + return + } + did, authMethod = auth.did, AuthMethodOAuth + ctx = context.WithValue(ctx, OAuthSessionKey, auth.session) + ctx = context.WithValue(ctx, UserAccessToken, auth.session.AccessToken) + } + + if m.authority == nil || !m.authority.IsInstanceAdmin(did) { + refuseInstanceAdmin(w, r, http.StatusForbidden, "Forbidden", "Instance admin authority required", "not_instance_admin", did, nil) + return + } + ctx = context.WithValue(ctx, UserDIDKey, did) + ctx = context.WithValue(ctx, AuthMethodKey, authMethod) + ctx = context.WithValue(ctx, IsAggregatorAuthKey, false) + next.ServeHTTP(w, r.WithContext(ctx)) + }) +} + +// refuseInstanceAdmin logs and writes a refusal. A server-side failure (5xx) +// logs at Error; a refused caller logs at Warn. The client IP is not logged. +func refuseInstanceAdmin(w http.ResponseWriter, r *http.Request, status int, code, message, reason, did string, cause error) { + attributes := []any{"method", r.Method, "path", r.URL.Path, "reason", reason} + if did != "" { + attributes = append(attributes, "did", did) + } + if cause != nil { + attributes = append(attributes, "error", cause) + } + level := slog.LevelWarn + if status >= http.StatusInternalServerError { + level = slog.LevelError + } + slog.Log(r.Context(), level, "instance admin request refused", attributes...) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + if err := json.NewEncoder(w).Encode(struct { + Error string `json:"error"` + Message string `json:"message"` + }{Error: code, Message: message}); err != nil { + slog.Warn("failed to write instance admin error response", "error", err) + } +} diff --git a/internal/api/middleware/admin_auth_test.go b/internal/api/middleware/admin_auth_test.go new file mode 100644 index 0000000..42bc705 --- /dev/null +++ b/internal/api/middleware/admin_auth_test.go @@ -0,0 +1,455 @@ +package middleware + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "log" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "Coves/internal/atproto/oauth" + + "github.com/bluesky-social/indigo/atproto/atcrypto" + indigoauth "github.com/bluesky-social/indigo/atproto/auth" + oauthlib "github.com/bluesky-social/indigo/atproto/auth/oauth" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const ( + instanceAdminDID = "did:plc:instanceadmin" + otherActorDID = "did:plc:ordinaryuser" + adminMethodPath = "/xrpc/social.coves.moderation.getSubjectState" +) + +type instanceAdminAuthorityFake map[string]bool + +func (a instanceAdminAuthorityFake) IsInstanceAdmin(did string) bool { return a[did] } + +type instanceAdminFixture struct { + client *mockOAuthClient + unsealer SessionUnsealer + sessions *mockOAuthStore + store oauthlib.ClientAuthStore + validator ServiceAuthValidator + authority instanceAdminAuthorityFake +} + +func newInstanceAdminFixture() *instanceAdminFixture { + sessions := newMockOAuthStore() + client := newMockOAuthClient() + return &instanceAdminFixture{ + client: client, + unsealer: client, + sessions: sessions, + store: sessions, + authority: instanceAdminAuthorityFake{instanceAdminDID: true}, + } +} + +func (f *instanceAdminFixture) sealedToken(t *testing.T, did string) string { + t.Helper() + const sessionID = "browser" + require.NoError(t, f.sessions.SaveSession(t.Context(), oauthlib.ClientSessionData{ + AccountDID: syntax.DID(did), SessionID: sessionID, + })) + return f.client.createTestToken(did, sessionID, time.Hour) +} + +func instanceAdminBearerRequest(token string) *http.Request { + request := httptest.NewRequest(http.MethodGet, adminMethodPath, nil) + request.Header.Set("Authorization", "Bearer "+token) + return request +} + +type instanceAdminResult struct { + response *httptest.ResponseRecorder + called bool + did string +} + +func (f *instanceAdminFixture) serve(request *http.Request) instanceAdminResult { + result := instanceAdminResult{response: httptest.NewRecorder()} + gate := NewInstanceAdminMiddleware(f.unsealer, f.store, f.validator, f.authority) + handler := gate.RequireInstanceAdmin(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + result.called = true + result.did = GetUserDID(r) + w.WriteHeader(http.StatusNoContent) + })) + handler.ServeHTTP(result.response, request) + return result +} + +func assertInstanceAdminResult(t *testing.T, got instanceAdminResult, wantStatus int, wantCode string, wantCalled bool, wantDID string) { + t.Helper() + assert.Equal(t, wantStatus, got.response.Code) + assert.Equal(t, wantCalled, got.called, "next handler must run only for an authorized admin") + if wantCalled { + assert.Equal(t, wantDID, got.did, "authenticated DID in next handler") + } + if wantCode != "" { + if assert.NotEmpty(t, got.response.Body.Bytes(), "refusal must include a JSON error") { + var body struct { + Error string `json:"error"` + Message string `json:"message"` + } + require.NoError(t, json.Unmarshal(got.response.Body.Bytes(), &body)) + assert.Equal(t, wantCode, body.Error) + assert.NotEmpty(t, body.Message) + } + } +} + +func TestRequireInstanceAdminOAuthRefusals(t *testing.T) { + for _, test := range []struct { + name string + request func(t *testing.T, fixture *instanceAdminFixture) *http.Request + wantStatus int + wantCode string + }{ + { + name: "no credential", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + request: func(_ *testing.T, _ *instanceAdminFixture) *http.Request { + return httptest.NewRequest(http.MethodGet, adminMethodPath, nil) + }, + }, + { + name: "basic authentication", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + request: func(_ *testing.T, _ *instanceAdminFixture) *http.Request { + request := httptest.NewRequest(http.MethodGet, adminMethodPath, nil) + request.Header.Set("Authorization", "Basic abc") + return request + }, + }, + { + name: "authenticated non-admin", wantStatus: http.StatusForbidden, wantCode: "Forbidden", + request: func(t *testing.T, f *instanceAdminFixture) *http.Request { + return instanceAdminBearerRequest(f.sealedToken(t, otherActorDID)) + }, + }, + { + name: "unsealable token", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + request: func(_ *testing.T, _ *instanceAdminFixture) *http.Request { + return instanceAdminBearerRequest("garbage-token") + }, + }, + { + name: "session missing", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + request: func(_ *testing.T, f *instanceAdminFixture) *http.Request { + return instanceAdminBearerRequest(f.client.createTestToken(instanceAdminDID, "missing", time.Hour)) + }, + }, + { + name: "session DID mismatch", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + request: func(_ *testing.T, f *instanceAdminFixture) *http.Request { + f.sessions.sessions[instanceAdminDID+":browser"] = &oauthlib.ClientSessionData{ + AccountDID: syntax.DID(otherActorDID), SessionID: "browser", + } + return instanceAdminBearerRequest(f.client.createTestToken(instanceAdminDID, "browser", time.Hour)) + }, + }, + { + name: "session store unavailable", wantStatus: http.StatusServiceUnavailable, wantCode: "ModerationUnavailable", + request: func(_ *testing.T, f *instanceAdminFixture) *http.Request { + f.store = &authenticationFailureStore{failure: errors.New("database unavailable")} + return instanceAdminBearerRequest(f.client.createTestToken(instanceAdminDID, "browser", time.Hour)) + }, + }, + } { + t.Run(test.name, func(t *testing.T) { + fixture := newInstanceAdminFixture() + request := test.request(t, fixture) + assertInstanceAdminResult(t, fixture.serve(request), test.wantStatus, test.wantCode, false, "") + }) + } +} + +func TestRequireInstanceAdminOAuthSuccess(t *testing.T) { + for _, useCookie := range []bool{false, true} { + name := "bearer header" + if useCookie { + name = "session cookie" + } + t.Run(name, func(t *testing.T) { + fixture := newInstanceAdminFixture() + token := fixture.sealedToken(t, instanceAdminDID) + request := httptest.NewRequest(http.MethodGet, adminMethodPath, nil) + if useCookie { + request.AddCookie(&http.Cookie{Name: "coves_session", Value: token}) + } else { + request.Header.Set("Authorization", "Bearer "+token) + } + assertInstanceAdminResult(t, fixture.serve(request), http.StatusNoContent, "", true, instanceAdminDID) + }) + } +} + +func TestRequireInstanceAdminIgnoresCallerActorAndAuthority(t *testing.T) { + t.Run("non-admin cannot claim admin in query or body", func(t *testing.T) { + fixture := newInstanceAdminFixture() + request := httptest.NewRequest(http.MethodPost, + adminMethodPath+"?actor="+instanceAdminDID+"&authority="+instanceAdminDID, + strings.NewReader(`{"actor":"`+instanceAdminDID+`","authority":"`+instanceAdminDID+`"}`)) + request.Header.Set("Content-Type", "application/json") + request.Header.Set("Authorization", "Bearer "+fixture.sealedToken(t, otherActorDID)) + assertInstanceAdminResult(t, fixture.serve(request), http.StatusForbidden, "Forbidden", false, "") + }) + t.Run("admin remains authenticated despite other actor", func(t *testing.T) { + fixture := newInstanceAdminFixture() + request := httptest.NewRequest(http.MethodGet, adminMethodPath+"?actor="+otherActorDID, nil) + request.Header.Set("Authorization", "Bearer "+fixture.sealedToken(t, instanceAdminDID)) + assertInstanceAdminResult(t, fixture.serve(request), http.StatusNoContent, "", true, instanceAdminDID) + }) +} + +type recordingInstanceAdminValidator struct { + mockServiceAuthValidator + called bool + lexMethod *syntax.NSID +} + +func (v *recordingInstanceAdminValidator) Validate(ctx context.Context, token string, lexMethod *syntax.NSID) (syntax.DID, error) { + v.called = true + if lexMethod != nil { + method := *lexMethod + v.lexMethod = &method + } + return v.mockServiceAuthValidator.Validate(ctx, token, lexMethod) +} + +func TestRequireInstanceAdminServiceJWTWithMockValidator(t *testing.T) { + for _, test := range []struct { + name string + validator *recordingInstanceAdminValidator + wantStatus int + wantCode string + wantCalled bool + }{ + { + name: "listed DID and endpoint-scoped lexMethod", validator: &recordingInstanceAdminValidator{ + mockServiceAuthValidator: mockServiceAuthValidator{returnDID: syntax.DID(instanceAdminDID)}, + }, wantStatus: http.StatusNoContent, wantCalled: true, + }, + { + name: "validator rejects token", validator: &recordingInstanceAdminValidator{ + mockServiceAuthValidator: mockServiceAuthValidator{shouldFail: true}, + }, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + }, + { + name: "unlisted issuer", validator: &recordingInstanceAdminValidator{ + mockServiceAuthValidator: mockServiceAuthValidator{returnDID: syntax.DID(otherActorDID)}, + }, wantStatus: http.StatusForbidden, wantCode: "Forbidden", + }, + { + name: "nil validator", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + }, + } { + t.Run(test.name, func(t *testing.T) { + fixture := newInstanceAdminFixture() + if test.validator != nil { + fixture.validator = test.validator + } + assertInstanceAdminResult(t, fixture.serve(instanceAdminBearerRequest("header.payload.signature")), + test.wantStatus, test.wantCode, test.wantCalled, instanceAdminDID) + if test.validator != nil { + assert.True(t, test.validator.called, "JWT must be validated, not treated as a sealed session") + if assert.NotNil(t, test.validator.lexMethod, "JWT validation must bind the called endpoint") { + assert.Equal(t, "social.coves.moderation.getSubjectState", test.validator.lexMethod.String()) + } + } + }) + } +} + +// mapSessionUnsealer unseals exactly the tokens it holds, so a JWT-shaped +// string can also be a valid sealed session. +type mapSessionUnsealer map[string]*oauth.SealedSession + +func (u mapSessionUnsealer) UnsealSession(token string) (*oauth.SealedSession, error) { + if sealed, ok := u[token]; ok { + return sealed, nil + } + return nil, errors.New("unknown sealed token") +} + +func TestRequireInstanceAdminServiceJWTRefusalDoesNotFallBackToSealedSession(t *testing.T) { + const jwtShapedToken = "a.b.c" + require.True(t, isJWTFormat(jwtShapedToken), "fixture token must route as a service JWT") + + for _, test := range []struct { + name string + validator ServiceAuthValidator + }{ + {name: "validator rejects token", validator: &mockServiceAuthValidator{shouldFail: true}}, + {name: "nil validator"}, + } { + t.Run(test.name, func(t *testing.T) { + fixture := newInstanceAdminFixture() + fixture.unsealer = mapSessionUnsealer{jwtShapedToken: { + DID: instanceAdminDID, SessionID: "browser", ExpiresAt: time.Now().Add(time.Hour).Unix(), + }} + require.NoError(t, fixture.sessions.SaveSession(t.Context(), oauthlib.ClientSessionData{ + AccountDID: syntax.DID(instanceAdminDID), SessionID: "browser", + })) + sealed := authenticateSealedSession(t.Context(), jwtShapedToken, fixture.unsealer, fixture.store) + require.Equal(t, sealedSessionNoFailure, sealed.failure, "the JWT-shaped token must also be a valid admin sealed session") + fixture.validator = test.validator + + assertInstanceAdminResult(t, fixture.serve(instanceAdminBearerRequest(jwtShapedToken)), + http.StatusUnauthorized, "AuthRequired", false, "") + }) + } +} + +func TestRequireInstanceAdminRefusalLogs(t *testing.T) { + const remoteAddress = "203.0.113.77:4242" + for _, test := range []struct { + name string + path string + credential func(t *testing.T, fixture *instanceAdminFixture) string + wantStatus int + wantCode string + wantLevel string + wantError string + }{ + { + name: "unsealable token", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", wantLevel: "WARN", + credential: func(_ *testing.T, _ *instanceAdminFixture) string { return "secret-token-value" }, + }, + { + name: "non-admin", wantStatus: http.StatusForbidden, wantCode: "Forbidden", wantLevel: "WARN", + credential: func(t *testing.T, f *instanceAdminFixture) string { return f.sealedToken(t, otherActorDID) }, + }, + { + name: "service JWT rejected", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", wantLevel: "WARN", + wantError: "mock validation failure", + credential: func(_ *testing.T, f *instanceAdminFixture) string { + f.validator = &mockServiceAuthValidator{shouldFail: true} + return "distinctive-jwt-header.distinctive-jwt-payload.distinctive-jwt-signature" + }, + }, + { + name: "service JWT for invalid method", path: "/xrpc/not-an-nsid", wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired", + wantLevel: "WARN", wantError: "NSID syntax didn't validate via regex", + credential: func(_ *testing.T, f *instanceAdminFixture) string { + f.validator = &mockServiceAuthValidator{returnDID: syntax.DID(instanceAdminDID)} + return "method-jwt-header.method-jwt-payload.method-jwt-signature" + }, + }, + { + name: "session store unavailable", wantStatus: http.StatusServiceUnavailable, wantCode: "ModerationUnavailable", wantLevel: "ERROR", + wantError: "database unavailable", + credential: func(_ *testing.T, f *instanceAdminFixture) string { + f.store = &authenticationFailureStore{failure: errors.New("database unavailable")} + return f.client.createTestToken(instanceAdminDID, "browser", time.Hour) + }, + }, + } { + t.Run(test.name, func(t *testing.T) { + var logged bytes.Buffer + previousSlog := slog.Default() + previousLog := log.Writer() + slog.SetDefault(slog.New(slog.NewTextHandler(&logged, nil))) + log.SetOutput(&logged) + t.Cleanup(func() { + slog.SetDefault(previousSlog) + log.SetOutput(previousLog) + }) + + fixture := newInstanceAdminFixture() + token := test.credential(t, fixture) + request := instanceAdminBearerRequest(token) + if test.path != "" { + request.URL.Path = test.path + } + request.RemoteAddr = remoteAddress + assertInstanceAdminResult(t, fixture.serve(request), test.wantStatus, test.wantCode, false, "") + output := logged.String() + assert.Contains(t, output, "level="+test.wantLevel, "refusal log level") + if test.wantError != "" { + assert.Contains(t, output, test.wantError, "refusal log must carry the underlying error") + } + assert.NotContains(t, output, token, "credentials must not appear in refusal logs") + assert.NotContains(t, output, "203.0.113.77", "client IP must not appear in refusal logs") + }) + } +} + +func TestRequireInstanceAdminRealServiceJWT(t *testing.T) { + method := syntax.NSID("social.coves.moderation.getSubjectState") + otherMethod := syntax.NSID("social.coves.moderation.removeContent") + const audience = "did:web:appview.test" + + for _, test := range []struct { + name string + issuer syntax.DID + audience string + lexMethod *syntax.NSID + wrongKey bool + expired bool + wantStatus int + wantCode string + }{ + {name: "valid admin", issuer: syntax.DID(instanceAdminDID), audience: audience, lexMethod: &method, wantStatus: http.StatusNoContent}, + {name: "wrong audience", issuer: syntax.DID(instanceAdminDID), audience: "did:web:other.test", lexMethod: &method, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired"}, + {name: "wrong lexMethod", issuer: syntax.DID(instanceAdminDID), audience: audience, lexMethod: &otherMethod, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired"}, + {name: "missing lexMethod", issuer: syntax.DID(instanceAdminDID), audience: audience, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired"}, + {name: "wrong signing key", issuer: syntax.DID(instanceAdminDID), audience: audience, lexMethod: &method, wrongKey: true, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired"}, + {name: "expired token", issuer: syntax.DID(instanceAdminDID), audience: audience, lexMethod: &method, expired: true, wantStatus: http.StatusUnauthorized, wantCode: "AuthRequired"}, + {name: "valid unlisted issuer", issuer: syntax.DID(otherActorDID), audience: audience, lexMethod: &method, wantStatus: http.StatusForbidden, wantCode: "Forbidden"}, + } { + t.Run(test.name, func(t *testing.T) { + privateKey, err := atcrypto.GeneratePrivateKeyP256() + require.NoError(t, err) + publicKey, err := privateKey.PublicKey() + require.NoError(t, err) + directory := identity.NewMockDirectory() + directory.Insert(identity.Identity{ + DID: test.issuer, + Keys: map[string]identity.VerificationMethod{ + "atproto": {Type: "Multikey", PublicKeyMultibase: publicKey.Multibase()}, + }, + }) + validator := &indigoauth.ServiceAuthValidator{ + Audience: audience, Dir: directory, TimestampLeeway: 30 * time.Second, + } + signingKey := atcrypto.PrivateKey(privateKey) + if test.wrongKey { + signingKey, err = atcrypto.GeneratePrivateKeyP256() + require.NoError(t, err) + } + lifetime := time.Minute + if test.expired { + lifetime = -2 * time.Minute + } + token, err := indigoauth.SignServiceAuth(test.issuer, test.audience, lifetime, test.lexMethod, signingKey) + require.NoError(t, err) + verifiedDID, validationErr := validator.Validate(t.Context(), token, &method) + if test.wantStatus == http.StatusUnauthorized { + require.Error(t, validationErr, "the signed JWT fixture must be rejected by indigo") + } else { + require.NoError(t, validationErr, "the signed JWT fixture must be verifiable by indigo") + require.Equal(t, test.issuer, verifiedDID) + } + + fixture := newInstanceAdminFixture() + fixture.validator = validator + assertInstanceAdminResult(t, fixture.serve(instanceAdminBearerRequest(token)), + test.wantStatus, test.wantCode, test.wantStatus == http.StatusNoContent, test.issuer.String()) + }) + } +} + +func TestSealedSessionAuthenticationZeroValueIsNotAuthenticated(t *testing.T) { + var zero sealedSessionAuthentication + assert.NotEqual(t, sealedSessionNoFailure, zero.failure, "an unset result must not read as authenticated") + assert.NotEmpty(t, zero.failure.reason(), "an unset result must have a refusal reason") +} diff --git a/internal/api/middleware/auth.go b/internal/api/middleware/auth.go index 5247931..31e1181 100644 --- a/internal/api/middleware/auth.go +++ b/internal/api/middleware/auth.go @@ -592,52 +592,47 @@ func (m *DualAuthMiddleware) handleAPIKeyAuth(w http.ResponseWriter, r *http.Req // handleOAuthAuth handles authentication using OAuth sealed session tokens (existing logic) func (m *DualAuthMiddleware) handleOAuthAuth(w http.ResponseWriter, r *http.Request, next http.Handler, token string) { - // Authenticate using sealed token - sealedSession, err := m.unsealer.UnsealSession(token) - if err != nil { + auth := authenticateSealedSession(r.Context(), token, m.unsealer, m.store) + switch auth.failure { + case sealedSessionUnsealFailed: log.Printf("[AUTH_FAILURE] type=unseal_failed ip=%s method=%s path=%s error=%v", - r.RemoteAddr, r.Method, r.URL.Path, err) + r.RemoteAddr, r.Method, r.URL.Path, auth.err) writeAuthError(w, "Invalid or expired token") return - } - - // Parse DID - did, err := syntax.ParseDID(sealedSession.DID) - if err != nil { + case sealedSessionInvalidDID: log.Printf("[AUTH_FAILURE] type=invalid_did ip=%s method=%s path=%s did=%s error=%v", - r.RemoteAddr, r.Method, r.URL.Path, sealedSession.DID, err) + r.RemoteAddr, r.Method, r.URL.Path, auth.did, auth.err) writeAuthError(w, "Invalid DID in token") return - } - - // Load full OAuth session from database - session, err := m.store.GetSession(r.Context(), did, sealedSession.SessionID) - if err != nil && !errors.Is(err, oauth.ErrSessionNotFound) { - writeSessionStoreError(w, r, sealedSession.DID, sealedSession.SessionID, err) + case sealedSessionStoreFailure: + writeSessionStoreError(w, r, auth.did, auth.sessionID, auth.err) return - } - if err != nil { + case sealedSessionNotFound: log.Printf("[AUTH_FAILURE] type=session_not_found ip=%s method=%s path=%s did=%s session_id=%s error=%v", - r.RemoteAddr, r.Method, r.URL.Path, sealedSession.DID, sealedSession.SessionID, err) + r.RemoteAddr, r.Method, r.URL.Path, auth.did, auth.sessionID, auth.err) writeAuthError(w, "Session not found or expired") return - } - - // Verify session DID matches token DID - if session.AccountDID.String() != sealedSession.DID { + case sealedSessionDIDMismatch: log.Printf("[AUTH_FAILURE] type=did_mismatch ip=%s method=%s path=%s token_did=%s session_did=%s", - r.RemoteAddr, r.Method, r.URL.Path, sealedSession.DID, session.AccountDID.String()) + r.RemoteAddr, r.Method, r.URL.Path, auth.did, auth.session.AccountDID.String()) writeAuthError(w, "Session DID mismatch") return + case sealedSessionNoFailure: + // Authenticated; continue below. + default: + log.Printf("[AUTH_FAILURE] type=%s ip=%s method=%s path=%s", + auth.failure.reason(), r.RemoteAddr, r.Method, r.URL.Path) + writeAuthError(w, "Invalid or expired token") + return } log.Printf("[AUTH_SUCCESS] type=oauth ip=%s method=%s path=%s did=%s session_id=%s", - r.RemoteAddr, r.Method, r.URL.Path, sealedSession.DID, sealedSession.SessionID) + r.RemoteAddr, r.Method, r.URL.Path, auth.did, auth.sessionID) // Inject user info and session into context - ctx := context.WithValue(r.Context(), UserDIDKey, sealedSession.DID) - ctx = context.WithValue(ctx, OAuthSessionKey, session) - ctx = context.WithValue(ctx, UserAccessToken, session.AccessToken) + ctx := context.WithValue(r.Context(), UserDIDKey, auth.did) + ctx = context.WithValue(ctx, OAuthSessionKey, auth.session) + ctx = context.WithValue(ctx, UserAccessToken, auth.session.AccessToken) ctx = context.WithValue(ctx, IsAggregatorAuthKey, false) ctx = context.WithValue(ctx, AuthMethodKey, AuthMethodOAuth) @@ -645,6 +640,75 @@ func (m *DualAuthMiddleware) handleOAuthAuth(w http.ResponseWriter, r *http.Requ next.ServeHTTP(w, r.WithContext(ctx)) } +type sealedSessionFailure uint8 + +// The zero value is sealedSessionUnauthenticated so that an unset result can +// never read as authenticated; success must be set explicitly. +const ( + sealedSessionUnauthenticated sealedSessionFailure = iota + sealedSessionNoFailure + sealedSessionUnsealFailed + sealedSessionInvalidDID + sealedSessionNotFound + sealedSessionStoreFailure + sealedSessionDIDMismatch +) + +func (failure sealedSessionFailure) reason() string { + switch failure { + case sealedSessionUnsealFailed: + return "unseal_failed" + case sealedSessionInvalidDID: + return "invalid_did" + case sealedSessionNotFound: + return "session_not_found" + case sealedSessionStoreFailure: + return "session_store_failure" + case sealedSessionDIDMismatch: + return "did_mismatch" + case sealedSessionNoFailure: + return "" + default: + return "unauthenticated" + } +} + +type sealedSessionAuthentication struct { + did string + sessionID string + session *oauthlib.ClientSessionData + failure sealedSessionFailure + err error +} + +func authenticateSealedSession(ctx context.Context, token string, unsealer SessionUnsealer, store oauthlib.ClientAuthStore) sealedSessionAuthentication { + sealed, err := unsealer.UnsealSession(token) + if err != nil { + return sealedSessionAuthentication{failure: sealedSessionUnsealFailed, err: err} + } + result := sealedSessionAuthentication{did: sealed.DID, sessionID: sealed.SessionID} + did, err := syntax.ParseDID(sealed.DID) + if err != nil { + result.failure, result.err = sealedSessionInvalidDID, err + return result + } + result.session, err = store.GetSession(ctx, did, sealed.SessionID) + if errors.Is(err, oauth.ErrSessionNotFound) { + result.failure, result.err = sealedSessionNotFound, err + return result + } + if err != nil { + result.failure, result.err = sealedSessionStoreFailure, err + return result + } + if result.session.AccountDID.String() != sealed.DID { + result.failure = sealedSessionDIDMismatch + return result + } + result.failure = sealedSessionNoFailure + return result +} + // isJWTFormat checks if a token has JWT format (three parts separated by dots). // NOTE: This is a format heuristic for routing, not security validation. // Actual JWT signature verification happens in ServiceAuthValidator.Validate(). diff --git a/internal/api/routes/harness_test.go b/internal/api/routes/harness_test.go new file mode 100644 index 0000000..1bb702c --- /dev/null +++ b/internal/api/routes/harness_test.go @@ -0,0 +1,14 @@ +//go:build integration + +package routes_test + +import ( + "os" + "testing" + + "Coves/tests/testkit" +) + +func TestMain(m *testing.M) { + os.Exit(testkit.Main(m, testkit.RequirePostgres)) +} diff --git a/internal/api/routes/moderation.go b/internal/api/routes/moderation.go new file mode 100644 index 0000000..cf87ad8 --- /dev/null +++ b/internal/api/routes/moderation.go @@ -0,0 +1,18 @@ +package routes + +import ( + handler "Coves/internal/api/handlers/moderation" + "Coves/internal/api/middleware" + "Coves/internal/core/moderation" + + "github.com/go-chi/chi/v5" +) + +// RegisterModerationRoutes registers the social.coves.moderation.* endpoints. +// Moderation NSIDs are AppView-served, never PDS-proxied, so the PDS OAuth +// scopes in cmd/server remain unchanged (see oauth_scopes_test.go). +func RegisterModerationRoutes(r chi.Router, service moderation.Service, adminAuth *middleware.InstanceAdminMiddleware) { + getSubjectState := handler.NewGetSubjectStateHandler(service) + r.With(adminAuth.RequireInstanceAdmin).Get( + "/xrpc/social.coves.moderation.getSubjectState", getSubjectState.HandleGetSubjectState) +} diff --git a/internal/api/routes/moderation_subject_state_integration_test.go b/internal/api/routes/moderation_subject_state_integration_test.go new file mode 100644 index 0000000..4edc26c --- /dev/null +++ b/internal/api/routes/moderation_subject_state_integration_test.go @@ -0,0 +1,184 @@ +//go:build integration + +package routes_test + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "net/url" + "testing" + "time" + + "Coves/internal/api/middleware" + "Coves/internal/api/routes" + "Coves/internal/core/moderation" + "Coves/internal/db/postgres" + "Coves/internal/validation" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + "github.com/bluesky-social/indigo/atproto/atdata" + "github.com/bluesky-social/indigo/atproto/lexicon" + "github.com/go-chi/chi/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const getSubjectStatePath = "/xrpc/social.coves.moderation.getSubjectState" + +type subjectStateHTTPResponse struct { + status int + body []byte +} + +func requestSubjectState(t *testing.T, client *http.Client, serverURL, subject, token string) subjectStateHTTPResponse { + t.Helper() + + target := serverURL + getSubjectStatePath + "?" + url.Values{"subject": {subject}}.Encode() + request, err := http.NewRequestWithContext(t.Context(), http.MethodGet, target, nil) + require.NoError(t, err) + if token != "" { + request.Header.Set("Authorization", "Bearer "+token) + } + + response, err := client.Do(request) + require.NoError(t, err) + defer response.Body.Close() + + body, err := io.ReadAll(response.Body) + require.NoError(t, err) + return subjectStateHTTPResponse{status: response.StatusCode, body: body} +} + +func requireXRPCError(t *testing.T, response subjectStateHTTPResponse, status int, code string) { + t.Helper() + + require.Equalf(t, status, response.status, "status = %d, want %d (body: %s)", response.status, status, response.body) + var body map[string]any + require.NoErrorf(t, json.Unmarshal(response.body, &body), "decoding XRPC error body %q", response.body) + assert.Equal(t, code, body["error"]) +} + +func TestGetSubjectState(t *testing.T) { + db := testkit.DB(t) + postRepo := postgres.NewPostRepository(db) + commentRepo := postgres.NewCommentRepository(db) + reader := moderation.NewRepositorySubjectReader(postRepo, commentRepo) + service := moderation.NewService(reader) + + adminDID := fixtures.DID(testkit.UniqueIDWithPrefix(t, "admin")) + nonAdminDID := fixtures.DID(testkit.UniqueIDWithPrefix(t, "nonadmin")) + authorName := testkit.UniqueIDWithPrefix(t, "subjectauthor") + authorDID := fixtures.DID(authorName) + fixtures.User(t, db, authorName+".test", authorDID) + communityName := testkit.UniqueIDWithPrefix(t, "subjectcommunity") + communityDID, err := fixtures.Community(t.Context(), db, communityName, "owner"+communityName) + require.NoError(t, err) + + const postTitle = "distinctive title" + uri := fixtures.Post(t, db, communityDID, authorDID, postTitle, 0, time.Now()) + indexedPost, err := postRepo.GetRawIndexedRow(t.Context(), uri) + require.NoError(t, err) + + const ( + adminToken = "admin-sealed-session" + adminSessionID = "admin-session" + nonAdminToken = "non-admin-sealed-session" + nonAdminSessionID = "non-admin-session" + ) + unsealer := fixtures.NewSessionUnsealer() + store := fixtures.NewOAuthStore() + unsealer.AddSession(adminToken, adminDID, adminSessionID) + store.AddSession(adminDID, adminSessionID, "admin-access-token") + unsealer.AddSession(nonAdminToken, nonAdminDID, nonAdminSessionID) + store.AddSession(nonAdminDID, nonAdminSessionID, "non-admin-access-token") + + authority := moderation.NewAllowlistAuthority([]string{adminDID}) + adminAuth := middleware.NewInstanceAdminMiddleware(unsealer, store, nil, authority) + mux := chi.NewRouter() + routes.RegisterModerationRoutes(mux, service, adminAuth) + server := httptest.NewServer(mux) + t.Cleanup(server.Close) + + t.Run("given an indexed post when an admin requests its state then the response is content-free", func(t *testing.T) { + response := requestSubjectState(t, server.Client(), server.URL, uri, adminToken) + require.Equalf(t, http.StatusOK, response.status, "status = %d, want 200 (body: %s)", response.status, response.body) + assert.NotContains(t, string(response.body), postTitle, "the moderation state must not retain indexed post content") + + data, err := atdata.UnmarshalJSON(response.body) + require.NoError(t, err, "the response must decode as atProto JSON") + catalog := lexicon.NewBaseCatalog() + require.NoError(t, catalog.LoadDirectory("../../atproto/lexicon")) + require.NoError(t, validation.ValidateData( + catalog, + data, + "social.coves.moderation.getSubjectState#output", + 0, + ), "the response must satisfy the getSubjectState output lexicon") + + var body map[string]any + require.NoError(t, json.Unmarshal(response.body, &body)) + state, ok := body["state"].(map[string]any) + require.Truef(t, ok, "state must be an object; got %#v", body["state"]) + assert.Equal(t, uri, state["subject"]) + assert.Equal(t, "present", state["recordState"]) + assert.Equal(t, "v0", state["version"]) + assert.NotContains(t, state, "localRemoval") + assert.NotContains(t, state, "localLabels") + + moderationView, ok := state["moderation"].(map[string]any) + require.Truef(t, ok, "moderation must be an object; got %#v", state["moderation"]) + assert.Equal(t, "clear", moderationView["state"]) + + currentSubject, ok := state["currentSubject"].(map[string]any) + require.Truef(t, ok, "currentSubject must be an object; got %#v", state["currentSubject"]) + assert.Equal(t, uri, currentSubject["uri"]) + assert.Equal(t, indexedPost.CID, currentSubject["cid"]) + }) + + t.Run("given a non-admin session when state is requested then access is forbidden", func(t *testing.T) { + response := requestSubjectState(t, server.Client(), server.URL, uri, nonAdminToken) + requireXRPCError(t, response, http.StatusForbidden, "Forbidden") + }) + + t.Run("given no credential when state is requested then authentication is required", func(t *testing.T) { + response := requestSubjectState(t, server.Client(), server.URL, uri, "") + requireXRPCError(t, response, http.StatusUnauthorized, "AuthRequired") + }) + + t.Run("given no credential and a malformed subject when requested then authentication is checked first", func(t *testing.T) { + response := requestSubjectState(t, server.Client(), server.URL, "not a uri", "") + requireXRPCError(t, response, http.StatusUnauthorized, "AuthRequired") + }) + + t.Run("given a non-admin and a malformed subject when requested then authorization is checked first", func(t *testing.T) { + response := requestSubjectState(t, server.Client(), server.URL, "not a uri", nonAdminToken) + requireXRPCError(t, response, http.StatusForbidden, "Forbidden") + }) + + t.Run("given a never-indexed supported subject when an admin requests it then unavailable state is returned", func(t *testing.T) { + const neverIndexedURI = "at://did:plc:neverindexed/social.coves.community.postv2/3kabc" + response := requestSubjectState(t, server.Client(), server.URL, neverIndexedURI, adminToken) + require.Equalf(t, http.StatusOK, response.status, "status = %d, want 200 (body: %s)", response.status, response.body) + + var body map[string]any + require.NoError(t, json.Unmarshal(response.body, &body)) + state, ok := body["state"].(map[string]any) + require.Truef(t, ok, "state must be an object; got %#v", body["state"]) + assert.Equal(t, neverIndexedURI, state["subject"]) + assert.Equal(t, "unavailable", state["recordState"]) + assert.Equal(t, "v0", state["version"]) + assert.NotContains(t, state, "currentSubject") + moderationView, ok := state["moderation"].(map[string]any) + require.Truef(t, ok, "moderation must be an object; got %#v", state["moderation"]) + assert.Equal(t, "clear", moderationView["state"]) + }) + + t.Run("given an unsupported subject collection when an admin requests it then the subject is invalid", func(t *testing.T) { + const unsupportedURI = "at://did:plc:x/social.coves.actor.profile/self" + response := requestSubjectState(t, server.Client(), server.URL, unsupportedURI, adminToken) + requireXRPCError(t, response, http.StatusBadRequest, "InvalidSubject") + }) +} diff --git a/internal/api/routes/registration_test.go b/internal/api/routes/registration_test.go index e7a82c2..a928cb7 100644 --- a/internal/api/routes/registration_test.go +++ b/internal/api/routes/registration_test.go @@ -247,6 +247,9 @@ var declaredRoutes = []declaredRoute{ // RegisterAdminReportRoutes — social.coves.admin.* {http.MethodPost, "/xrpc/social.coves.admin.submitReport", authRequired, 10, false}, + // RegisterModerationRoutes — instance-admin-only state query. + {http.MethodGet, "/xrpc/social.coves.moderation.getSubjectState", authRequired, 0, false}, + // RegisterCommunitySuggestionRoutes — social.coves.community.suggestion.* {http.MethodGet, "/xrpc/social.coves.community.suggestion.list", authOptional, 0, false}, {http.MethodGet, "/xrpc/social.coves.community.suggestion.get", authOptional, 0, false}, @@ -425,6 +428,7 @@ var theRouter = sync.OnceValue(func() builtRouter { RegisterUserBlockRoutes(mux, unreachableUserBlockService{}, auth) RegisterCommentRoutes(mux, nil, auth) RegisterAdminReportRoutes(mux, nil, auth) + RegisterModerationRoutes(mux, nil, middleware.NewInstanceAdminMiddleware(unsealer, nil, nil, nil)) RegisterCommunitySuggestionRoutes(mux, nil, auth, nil) RegisterCommunityFeedRoutes(mux, nil, nil, nil, auth) RegisterTimelineRoutes(mux, nil, nil, nil, auth) diff --git a/internal/config/config.go b/internal/config/config.go index 0d6785c..e03e8d5 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -14,6 +14,8 @@ import ( "Coves/internal/core/bridgedvotes" "Coves/internal/core/imageproxy" + + "github.com/bluesky-social/indigo/atproto/syntax" ) // devCursorSecret is the placeholder HMAC key used for pagination cursors when @@ -122,6 +124,9 @@ type Config struct { // Submissions bounds what one author may post into one community. Submissions SubmissionsConfig + // Moderation holds the instance moderation settings. + Moderation ModerationConfig + // CursorSecret is the HMAC key that signs pagination cursors, preventing // clients from forging or tampering with them. CursorSecret string @@ -137,6 +142,13 @@ type Config struct { EncryptionKeyGenerated bool } +// ModerationConfig holds the instance moderation settings. +type ModerationConfig struct { + // Admins is the operator-managed allowlist of instance admin DIDs + // (MODERATION_ADMINS). An empty list grants nobody admin authority. + Admins []string +} + // DatabaseConfig holds the AppView PostgreSQL connection and pool settings. // // The pool bounds matter: database/sql defaults to an unlimited number of open @@ -515,6 +527,7 @@ func Load() (*Config, error) { if err := cfg.loadSubmissions(); err != nil { return nil, err } + cfg.Moderation.Admins = csvVar("MODERATION_ADMINS") cfg.PDS = PDSConfig{ URL: stringVar("PDS_URL", "http://localhost:3001"), @@ -1011,6 +1024,12 @@ func (c *Config) loadSubmissions() error { // instead of one restart per mistake. func (c *Config) Validate() error { var problems []string + for _, did := range c.Moderation.Admins { + if _, err := syntax.ParseDID(did); err != nil { + problems = append(problems, "MODERATION_ADMINS must contain only valid DIDs") + break + } + } if c.Database.URL == "" { problems = append(problems, "DATABASE_URL is required") diff --git a/internal/config/moderation_admins_test.go b/internal/config/moderation_admins_test.go new file mode 100644 index 0000000..7e70f2f --- /dev/null +++ b/internal/config/moderation_admins_test.go @@ -0,0 +1,53 @@ +package config + +import ( + "os" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestLoadModerationAdmins(t *testing.T) { + tests := []struct { + name string + value string + want []string + }{ + {name: "trims and drops empty entries", value: " did:plc:a , did:plc:b ,, ", want: []string{"did:plc:a", "did:plc:b"}}, + {name: "empty list grants nobody", value: ""}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + clearEnv(t) + prodEnv(t) + t.Setenv("MODERATION_ADMINS", test.value) + + cfg, err := Load() + require.NoError(t, err) + if test.want == nil { + assert.Empty(t, cfg.Moderation.Admins) + } else { + assert.Equal(t, test.want, cfg.Moderation.Admins) + } + }) + } +} + +func TestLoadModerationAdminsRejectsNonDID(t *testing.T) { + clearEnv(t) + prodEnv(t) + t.Setenv("MODERATION_ADMINS", "did:plc:a,alice.test") + + _, err := Load() + require.Error(t, err) + assert.ErrorContains(t, err, "MODERATION_ADMINS") +} + +func TestClearEnvForTestClearsModerationAdmins(t *testing.T) { + t.Setenv("MODERATION_ADMINS", "did:plc:a") + + ClearEnvForTest(t) + + assert.Empty(t, os.Getenv("MODERATION_ADMINS")) +} diff --git a/internal/config/testing.go b/internal/config/testing.go index f8f54db..a765010 100644 --- a/internal/config/testing.go +++ b/internal/config/testing.go @@ -18,6 +18,7 @@ var loadedEnvVars = []string{ "OAUTH_CLIENT_PRIVATE_KEY", "OAUTH_CLIENT_KEY_ID", "INSTANCE_DID", "INSTANCE_DOMAIN", "COMMUNITY_CREATORS", "TRUSTED_BRIDGE_PDS_HOSTS", "SKIP_DID_WEB_VERIFICATION", + "MODERATION_ADMINS", "BRIDGED_VOTE_POLL_INTERVAL", "BRIDGED_VOTE_POLL_LOOKBACK", "BRIDGED_VOTE_POLL_SWEEP_CAP", "PDS_URL", "PDS_INSTANCE_HANDLE", "PDS_INSTANCE_PASSWORD", "PDS_ADMIN_PASSWORD", "JETSTREAM_FEEDS", "REDRIVE_INTERVAL", "IDENTITY_NEGATIVE_CACHE_TTL", diff --git a/internal/core/moderation/authority.go b/internal/core/moderation/authority.go new file mode 100644 index 0000000..7455b57 --- /dev/null +++ b/internal/core/moderation/authority.go @@ -0,0 +1,19 @@ +package moderation + +// NewAllowlistAuthority builds an Authority from the configured admin DIDs. +func NewAllowlistAuthority(dids []string) Authority { + allowed := make(map[string]struct{}, len(dids)) + for _, did := range dids { + allowed[did] = struct{}{} + } + return allowlistAuthority{allowed: allowed} +} + +type allowlistAuthority struct { + allowed map[string]struct{} +} + +func (a allowlistAuthority) IsInstanceAdmin(did string) bool { + _, ok := a.allowed[did] + return ok +} diff --git a/internal/core/moderation/authority_test.go b/internal/core/moderation/authority_test.go new file mode 100644 index 0000000..49e3502 --- /dev/null +++ b/internal/core/moderation/authority_test.go @@ -0,0 +1,44 @@ +package moderation_test + +import ( + "testing" + + "Coves/internal/core/moderation" + "github.com/stretchr/testify/assert" +) + +func TestAllowlistAuthorityDefaultsToNobody(t *testing.T) { + for _, test := range []struct { + name string + dids []string + }{ + {name: "nil list"}, + {name: "empty list", dids: []string{}}, + } { + t.Run(test.name, func(t *testing.T) { + authority := moderation.NewAllowlistAuthority(test.dids) + for _, did := range []string{"", "did:plc:a"} { + assert.False(t, authority.IsInstanceAdmin(did), "DID %q must not be an admin", did) + } + }) + } +} + +func TestAllowlistAuthorityMatchesExactDID(t *testing.T) { + authority := moderation.NewAllowlistAuthority([]string{"did:plc:alice"}) + for _, test := range []struct { + name string + did string + want bool + }{ + {name: "listed DID", did: "did:plc:alice", want: true}, + {name: "unlisted DID", did: "did:plc:bob"}, + {name: "different case", did: "did:plc:Alice"}, + {name: "trailing whitespace", did: "did:plc:alice "}, + {name: "prefix", did: "did:plc:alic"}, + } { + t.Run(test.name, func(t *testing.T) { + assert.Equal(t, test.want, authority.IsInstanceAdmin(test.did)) + }) + } +} diff --git a/internal/core/moderation/errors.go b/internal/core/moderation/errors.go new file mode 100644 index 0000000..78ca2a8 --- /dev/null +++ b/internal/core/moderation/errors.go @@ -0,0 +1,23 @@ +package moderation + +import "errors" + +// One sentinel per error code the moderation lexicons declare (PRD §14.5), +// plus the reader-level ErrSubjectNotIndexed which never reaches the wire. +var ( + ErrAuthRequired = errors.New("moderation: authentication required") + ErrForbidden = errors.New("moderation: forbidden") + ErrInvalidRequest = errors.New("moderation: invalid request") + ErrInvalidSubject = errors.New("moderation: invalid subject") + ErrSubjectNotFound = errors.New("moderation: subject not found") + ErrDecisionNotFound = errors.New("moderation: decision not found") + ErrInvalidDecision = errors.New("moderation: invalid decision") + ErrContentChanged = errors.New("moderation: content changed") + ErrStateConflict = errors.New("moderation: state conflict") + ErrIdempotencyConflict = errors.New("moderation: idempotency conflict") + ErrUnsupportedReason = errors.New("moderation: unsupported reason") + ErrUnsupportedLabel = errors.New("moderation: unsupported label") + ErrInvalidCursor = errors.New("moderation: invalid cursor") + ErrModerationUnavailable = errors.New("moderation: temporarily unavailable") + ErrSubjectNotIndexed = errors.New("moderation: subject never indexed") +) diff --git a/internal/core/moderation/harness_test.go b/internal/core/moderation/harness_test.go new file mode 100644 index 0000000..855d292 --- /dev/null +++ b/internal/core/moderation/harness_test.go @@ -0,0 +1,14 @@ +//go:build integration + +package moderation_test + +import ( + "os" + "testing" + + "Coves/tests/testkit" +) + +func TestMain(m *testing.M) { + os.Exit(testkit.Main(m, testkit.RequirePostgres)) +} diff --git a/internal/core/moderation/interfaces.go b/internal/core/moderation/interfaces.go new file mode 100644 index 0000000..919ea8b --- /dev/null +++ b/internal/core/moderation/interfaces.go @@ -0,0 +1,34 @@ +package moderation + +import ( + "context" + + "Coves/internal/core/comments" + "Coves/internal/core/posts" +) + +// Authority answers whether a DID may act as an instance admin. +type Authority interface { + IsInstanceAdmin(did string) bool +} + +// SubjectReader reports what the AppView has indexed for a subject URI. +// It returns ErrSubjectNotIndexed for a URI the AppView has never seen. +type SubjectReader interface { + ReadSubject(ctx context.Context, uri string) (*IndexedRecord, error) +} + +// PostReader is the ungated post lookup the reader adapter needs. +type PostReader interface { + GetRawIndexedRow(ctx context.Context, uri string) (*posts.Post, error) +} + +// CommentReader is the admission-blind comment lookup the reader adapter needs. +type CommentReader interface { + GetByURI(ctx context.Context, uri string) (*comments.Comment, error) +} + +// Service is the moderation domain's read surface. +type Service interface { + GetSubjectState(ctx context.Context, subject string) (*SubjectState, error) +} diff --git a/internal/core/moderation/service.go b/internal/core/moderation/service.go new file mode 100644 index 0000000..4113e69 --- /dev/null +++ b/internal/core/moderation/service.go @@ -0,0 +1,49 @@ +package moderation + +import ( + "context" + "errors" + "fmt" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +// NewService builds the moderation read service over a SubjectReader. +func NewService(reader SubjectReader) Service { + return &service{reader: reader} +} + +type service struct { + reader SubjectReader +} + +func (s *service) GetSubjectState(ctx context.Context, subject string) (*SubjectState, error) { + uri, err := syntax.ParseATURI(subject) + if err != nil || !uri.Authority().IsDID() || !IsSubjectCollection(uri.Collection().String()) || uri.RecordKey().String() == "" { + return nil, fmt.Errorf("%w: expected a record URI with a DID authority and supported collection", ErrInvalidSubject) + } + + state := &SubjectState{ + Subject: subject, + Version: InitialVersion, + Moderation: ModerationView{State: ModerationStateClear}, + } + record, err := s.reader.ReadSubject(ctx, subject) + if errors.Is(err, ErrSubjectNotIndexed) { + state.RecordState = RecordStateUnavailable + return state, nil + } + if err != nil { + return nil, fmt.Errorf("%w: %w", ErrModerationUnavailable, err) + } + if record == nil { + return nil, fmt.Errorf("%w: subject reader returned no record", ErrModerationUnavailable) + } + if record.Deleted { + state.RecordState = RecordStateDeleted + } else { + state.RecordState = RecordStatePresent + state.CurrentSubject = &StrongRef{URI: record.URI, CID: record.CID} + } + return state, nil +} diff --git a/internal/core/moderation/service_test.go b/internal/core/moderation/service_test.go new file mode 100644 index 0000000..e7e61e3 --- /dev/null +++ b/internal/core/moderation/service_test.go @@ -0,0 +1,121 @@ +package moderation_test + +import ( + "context" + "errors" + "testing" + + "Coves/internal/core/moderation" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type fakeSubjectReader struct { + calls []string + record *moderation.IndexedRecord + err error +} + +func (reader *fakeSubjectReader) ReadSubject(_ context.Context, uri string) (*moderation.IndexedRecord, error) { + reader.calls = append(reader.calls, uri) + return reader.record, reader.err +} + +func assertInitialSubjectState(t *testing.T, state *moderation.SubjectState, uri string, recordState moderation.RecordState, current *moderation.StrongRef) { + t.Helper() + require.NotNil(t, state) + assert.Equal(t, uri, state.Subject) + assert.Equal(t, recordState, state.RecordState) + assert.Equal(t, moderation.InitialVersion, state.Version) + assert.Equal(t, moderation.ModerationStateClear, state.Moderation.State) + assert.Equal(t, current, state.CurrentSubject) + assert.Nil(t, state.LocalRemoval) + assert.Empty(t, state.LocalLabels) +} + +func TestGetSubjectStateRejectsInvalidSubjectsBeforeReading(t *testing.T) { + for _, subject := range []string{ + "", + "not a uri", + "at://alice.test/social.coves.community.post/3kabc", + "at://did:plc:x/social.coves.actor.profile/self", + "at://did:plc:x/social.coves.community.comment", + "at://did:plc:x", + "https://example.com/x", + } { + t.Run(subject, func(t *testing.T) { + reader := &fakeSubjectReader{} + state, err := moderation.NewService(reader).GetSubjectState(t.Context(), subject) + assert.ErrorIs(t, err, moderation.ErrInvalidSubject) + assert.Nil(t, state) + assert.Empty(t, reader.calls, "invalid subjects must not reach the reader") + }) + } +} + +func TestGetSubjectStateNeverIndexed(t *testing.T) { + uri := "at://did:plc:neverindexed/social.coves.community.postv2/3kabc" + reader := &fakeSubjectReader{err: moderation.ErrSubjectNotIndexed} + + state, err := moderation.NewService(reader).GetSubjectState(t.Context(), uri) + require.NoError(t, err) + assertInitialSubjectState(t, state, uri, moderation.RecordStateUnavailable, nil) + assert.Equal(t, []string{uri}, reader.calls) +} + +func TestGetSubjectStateIndexedPresent(t *testing.T) { + for _, test := range []struct { + name string + collection string + }{ + {name: "comment", collection: moderation.CommentCollection}, + {name: "postv2", collection: moderation.PostV2Collection}, + {name: "legacy post", collection: moderation.LegacyPostCollection}, + } { + t.Run(test.name, func(t *testing.T) { + uri := "at://did:plc:x/" + test.collection + "/3kabc" + cid := "bafy...distinct" + reader := &fakeSubjectReader{record: &moderation.IndexedRecord{URI: uri, CID: cid}} + + state, err := moderation.NewService(reader).GetSubjectState(t.Context(), uri) + require.NoError(t, err) + assertInitialSubjectState(t, state, uri, moderation.RecordStatePresent, &moderation.StrongRef{URI: uri, CID: cid}) + assert.Equal(t, []string{uri}, reader.calls) + }) + } +} + +func TestGetSubjectStateAuthorDeleted(t *testing.T) { + uri := "at://did:plc:x/social.coves.community.comment/3kabc" + reader := &fakeSubjectReader{record: &moderation.IndexedRecord{URI: uri, CID: "bafy...distinct", Deleted: true}} + + state, err := moderation.NewService(reader).GetSubjectState(t.Context(), uri) + require.NoError(t, err) + assertInitialSubjectState(t, state, uri, moderation.RecordStateDeleted, nil) + assert.Equal(t, []string{uri}, reader.calls) +} + +func TestGetSubjectStateReaderFailure(t *testing.T) { + uri := "at://did:plc:x/social.coves.community.postv2/3kabc" + reader := &fakeSubjectReader{err: errors.New("db down")} + + state, err := moderation.NewService(reader).GetSubjectState(t.Context(), uri) + assert.ErrorIs(t, err, moderation.ErrModerationUnavailable) + assert.Nil(t, state) + assert.Equal(t, []string{uri}, reader.calls) +} + +func TestGetSubjectStateNeverIndexedVersionIsStable(t *testing.T) { + uri := "at://did:plc:neverindexed/social.coves.community.postv2/3kabc" + reader := &fakeSubjectReader{err: moderation.ErrSubjectNotIndexed} + service := moderation.NewService(reader) + + first, err := service.GetSubjectState(t.Context(), uri) + require.NoError(t, err) + require.NotNil(t, first) + second, err := service.GetSubjectState(t.Context(), uri) + require.NoError(t, err) + require.NotNil(t, second) + assert.Equal(t, first.Version, second.Version) + assert.Equal(t, []string{uri, uri}, reader.calls) +} diff --git a/internal/core/moderation/subject_reader.go b/internal/core/moderation/subject_reader.go new file mode 100644 index 0000000..5c043c5 --- /dev/null +++ b/internal/core/moderation/subject_reader.go @@ -0,0 +1,53 @@ +package moderation + +import ( + "context" + "errors" + "fmt" + + "Coves/internal/core/comments" + "Coves/internal/core/posts" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +// NewRepositorySubjectReader reads subjects from the existing post and +// comment repositories, dispatching on the URI's collection. +func NewRepositorySubjectReader(postReader PostReader, commentReader CommentReader) SubjectReader { + return &repositorySubjectReader{postReader: postReader, commentReader: commentReader} +} + +type repositorySubjectReader struct { + postReader PostReader + commentReader CommentReader +} + +func (r *repositorySubjectReader) ReadSubject(ctx context.Context, uri string) (*IndexedRecord, error) { + parsed, err := syntax.ParseATURI(uri) + if err != nil || !parsed.Authority().IsDID() || parsed.RecordKey().String() == "" { + return nil, fmt.Errorf("%w: expected a record URI with a DID authority", ErrInvalidSubject) + } + + switch parsed.Collection().String() { + case PostV2Collection, LegacyPostCollection: + post, err := r.postReader.GetRawIndexedRow(ctx, uri) + if errors.Is(err, posts.ErrNotFound) { + return nil, ErrSubjectNotIndexed + } + if err != nil { + return nil, fmt.Errorf("reading post: %w", err) + } + return &IndexedRecord{URI: post.URI, CID: post.CID, Deleted: post.DeletedAt != nil}, nil + case CommentCollection: + comment, err := r.commentReader.GetByURI(ctx, uri) + if errors.Is(err, comments.ErrCommentNotFound) { + return nil, ErrSubjectNotIndexed + } + if err != nil { + return nil, fmt.Errorf("reading comment: %w", err) + } + return &IndexedRecord{URI: comment.URI, CID: comment.CID, Deleted: comment.DeletedAt != nil}, nil + default: + return nil, fmt.Errorf("%w: unsupported collection", ErrInvalidSubject) + } +} diff --git a/internal/core/moderation/subject_reader_integration_test.go b/internal/core/moderation/subject_reader_integration_test.go new file mode 100644 index 0000000..b6f7764 --- /dev/null +++ b/internal/core/moderation/subject_reader_integration_test.go @@ -0,0 +1,153 @@ +//go:build integration + +package moderation_test + +import ( + "testing" + "time" + + "Coves/internal/core/moderation" + "Coves/internal/db/postgres" + "Coves/tests/fixtures" + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestRepositorySubjectReaderIndexedRecords(t *testing.T) { + db := testkit.DB(t) + ctx := t.Context() + reader := moderation.NewRepositorySubjectReader(postgres.NewPostRepository(db), postgres.NewCommentRepository(db)) + + authorName := testkit.UniqueIDWithPrefix(t, "subauthor") + authorDID := fixtures.DID(authorName) + fixtures.User(t, db, authorName+".test", authorDID) + communityName := testkit.UniqueIDWithPrefix(t, "subcommunity") + communityDID, err := fixtures.Community(ctx, db, communityName, "owner"+communityName) + require.NoError(t, err) + + legacyURI := fixtures.Post(t, db, communityDID, authorDID, "legacy", 0, time.Now()) + deletedLegacyURI := fixtures.Post(t, db, communityDID, authorDID, "deleted legacy", 0, time.Now()) + const legacyCID = "bafytest" + + insertPostV2 := func(title, cid string) string { + t.Helper() + rkey := testkit.TID() + uri := "at://" + authorDID + "/" + moderation.PostV2Collection + "/" + rkey + _, err := db.ExecContext(ctx, ` + INSERT INTO posts (uri, cid, rkey, author_did, community_did, title, created_at) + VALUES ($1, $2, $3, $4, $5, $6, NOW()) + `, uri, cid, rkey, authorDID, communityDID, title) + require.NoError(t, err) + return uri + } + postV2CID := "bafyreipresent" + testkit.UniqueIDWithPrefix(t, "cid") + postV2URI := insertPostV2("present postv2", postV2CID) + deletedPostV2URI := insertPostV2("deleted postv2", "bafyreideletedpostv2") + _, err = db.ExecContext(ctx, `UPDATE posts SET deleted_at = NOW() WHERE uri IN ($1, $2)`, deletedPostV2URI, deletedLegacyURI) + require.NoError(t, err) + + pendingURI := insertPostV2("unadmitted postv2", "bafyreipendingpostv2") + _, err = db.ExecContext(ctx, ` + INSERT INTO community_post_admissions (community_did, post_uri, status, evaluated_cid) + VALUES ($1, $2, 'pending', $3) + `, communityDID, pendingURI, "bafyreipendingpostv2") + require.NoError(t, err) + + insertComment := func(cid string) string { + t.Helper() + rkey := testkit.TID() + uri := "at://" + authorDID + "/" + moderation.CommentCollection + "/" + rkey + _, err := db.ExecContext(ctx, ` + INSERT INTO comments (uri, cid, rkey, commenter_did, root_uri, root_cid, parent_uri, parent_cid, content, created_at) + VALUES ($1, $2, $3, $4, $5, $6, $5, $6, $7, NOW()) + `, uri, cid, rkey, authorDID, postV2URI, postV2CID, "a comment") + require.NoError(t, err) + return uri + } + commentURI := insertComment("bafyreipresentcomment") + deletedCommentURI := insertComment("bafyreideletedcomment") + _, err = db.ExecContext(ctx, ` + UPDATE comments SET deleted_at = NOW(), deletion_reason = 'author', deleted_by = $1 WHERE uri = $2 + `, authorDID, deletedCommentURI) + require.NoError(t, err) + + for _, test := range []struct { + name string + uri string + cid string + deleted bool + }{ + {"indexed postv2", postV2URI, postV2CID, false}, + {"indexed legacy post", legacyURI, legacyCID, false}, + {"indexed comment", commentURI, "bafyreipresentcomment", false}, + {"author-deleted comment", deletedCommentURI, "bafyreideletedcomment", true}, + {"soft-deleted postv2", deletedPostV2URI, "bafyreideletedpostv2", true}, + {"soft-deleted legacy post", deletedLegacyURI, legacyCID, true}, + {"pending postv2", pendingURI, "bafyreipendingpostv2", false}, + } { + t.Run(test.name, func(t *testing.T) { + record, err := reader.ReadSubject(t.Context(), test.uri) + require.NoError(t, err) + assert.Equal(t, &moderation.IndexedRecord{URI: test.uri, CID: test.cid, Deleted: test.deleted}, record) + }) + } + + for _, test := range []struct { + name string + uri string + }{ + {"never-indexed postv2", "at://" + authorDID + "/" + moderation.PostV2Collection + "/" + testkit.TID()}, + {"never-indexed legacy post", "at://" + communityDID + "/" + moderation.LegacyPostCollection + "/" + testkit.TID()}, + {"never-indexed comment", "at://" + authorDID + "/" + moderation.CommentCollection + "/" + testkit.TID()}, + } { + t.Run(test.name, func(t *testing.T) { + record, err := reader.ReadSubject(t.Context(), test.uri) + assert.Nil(t, record) + assert.ErrorIs(t, err, moderation.ErrSubjectNotIndexed) + }) + } + + service := moderation.NewService(reader) + t.Run("service author-deleted comment", func(t *testing.T) { + state, err := service.GetSubjectState(t.Context(), deletedCommentURI) + require.NoError(t, err) + require.NotNil(t, state) + assert.Equal(t, moderation.RecordStateDeleted, state.RecordState) + assert.Nil(t, state.CurrentSubject) + }) + t.Run("service indexed postv2", func(t *testing.T) { + state, err := service.GetSubjectState(t.Context(), postV2URI) + require.NoError(t, err) + require.NotNil(t, state) + assert.Equal(t, moderation.RecordStatePresent, state.RecordState) + assert.Equal(t, &moderation.StrongRef{URI: postV2URI, CID: postV2CID}, state.CurrentSubject) + }) + t.Run("service never-indexed version", func(t *testing.T) { + uri := "at://" + authorDID + "/" + moderation.PostV2Collection + "/" + testkit.TID() + first, err := service.GetSubjectState(t.Context(), uri) + require.NoError(t, err) + second, err := service.GetSubjectState(t.Context(), uri) + require.NoError(t, err) + require.NotNil(t, first) + require.NotNil(t, second) + assert.Equal(t, moderation.RecordStateUnavailable, first.RecordState) + assert.Equal(t, "v0", first.Version) + assert.Equal(t, first.Version, second.Version) + }) +} + +func TestRepositorySubjectReaderDatabaseFailureIsNotMissing(t *testing.T) { + // This is a dedicated clone: testkit closes its pool during cleanup, and + // closing the pool here must not interrupt another test's fixture database. + db := testkit.DB(t) + reader := moderation.NewRepositorySubjectReader(postgres.NewPostRepository(db), postgres.NewCommentRepository(db)) + require.NoError(t, db.Close()) + uri := "at://" + fixtures.DID(testkit.UniqueIDWithPrefix(t, "dbfailure")) + "/" + moderation.PostV2Collection + "/" + testkit.TID() + + record, err := reader.ReadSubject(t.Context(), uri) + assert.Nil(t, record) + require.Error(t, err) + assert.NotErrorIs(t, err, moderation.ErrSubjectNotIndexed) +} diff --git a/internal/core/moderation/types.go b/internal/core/moderation/types.go new file mode 100644 index 0000000..45e60b2 --- /dev/null +++ b/internal/core/moderation/types.go @@ -0,0 +1,78 @@ +// Package moderation is the instance-admin moderation domain: who may act +// (Authority) and what the AppView knows about a subject (SubjectState). +package moderation + +// RecordState is the repository availability of a subject record as the +// AppView has indexed it. +type RecordState string + +const ( + RecordStatePresent RecordState = "present" + RecordStateDeleted RecordState = "deleted" + RecordStateUnavailable RecordState = "unavailable" +) + +// ModerationStateClear is the effective removal state of a subject with no +// applicable removal decision. +const ModerationStateClear = "clear" + +// Subject collections getSubjectState and the mutations accept. +const ( + CommentCollection = "social.coves.community.comment" + PostV2Collection = "social.coves.community.postv2" + LegacyPostCollection = "social.coves.community.post" +) + +// SubjectCollections is the set of NSIDs a moderation subject may live in. +var SubjectCollections = map[string]struct{}{ + CommentCollection: {}, + PostV2Collection: {}, + LegacyPostCollection: {}, +} + +// IsSubjectCollection reports whether nsid is a supported subject collection. +func IsSubjectCollection(nsid string) bool { + _, ok := SubjectCollections[nsid] + return ok +} + +// StrongRef is a com.atproto.repo.strongRef. +type StrongRef struct { + URI string + CID string +} + +// ActionRef identifies an action in a moderation service's log. +type ActionRef struct { + ServiceDID string + ActionID string +} + +// LocalLabel is a content label established by a local action. +type LocalLabel struct { + Value string + Action ActionRef +} + +// ModerationView is the effective public removal state of a subject. +type ModerationView struct { + State string +} + +// SubjectState is the versioned moderation and repository state of a subject. +type SubjectState struct { + Subject string + Version string + Moderation ModerationView + RecordState RecordState + CurrentSubject *StrongRef + LocalRemoval *ActionRef + LocalLabels []LocalLabel +} + +// IndexedRecord is what the AppView has indexed for a subject URI. +type IndexedRecord struct { + URI string + CID string + Deleted bool +} diff --git a/internal/core/moderation/version.go b/internal/core/moderation/version.go new file mode 100644 index 0000000..44404b7 --- /dev/null +++ b/internal/core/moderation/version.go @@ -0,0 +1,5 @@ +package moderation + +// InitialVersion is the opaque state token of a subject with no moderation +// rows. Later chunks advance the encoding; readers treat it as opaque. +const InitialVersion = "v0" diff --git a/scripts/ci-bootstrap.sh b/scripts/ci-bootstrap.sh index a0f2fd6..f3b8639 100755 --- a/scripts/ci-bootstrap.sh +++ b/scripts/ci-bootstrap.sh @@ -1,6 +1,6 @@ #!/usr/bin/env bash # Seeds the CI stack: registers both PDSes with the relay, then creates the -# account the AppView authenticates as. +# instance account and two moderation admin accounts before the AppView boots. # # Both halves must happen after the infrastructure is healthy and before the # AppView boots, which is why they share a stage in scripts/lib/ci-stack.sh @@ -29,6 +29,11 @@ RELAY_ADMIN_KEY=${RELAY_ADMIN_KEY:-ci-relay-admin-key} HANDLE=${PDS_INSTANCE_HANDLE:?PDS_INSTANCE_HANDLE must be set (see .env.ci)} PASSWORD=${PDS_INSTANCE_PASSWORD:?PDS_INSTANCE_PASSWORD must be set (see .env.ci)} EMAIL=${COVES_CI_INSTANCE_EMAIL:-instance@local.coves.dev} +ADMIN_ONE_HANDLE=${CI_MODERATION_ADMIN_ONE_HANDLE:?CI_MODERATION_ADMIN_ONE_HANDLE must be set (see .env.ci)} +ADMIN_ONE_PASSWORD=${CI_MODERATION_ADMIN_ONE_PASSWORD:?CI_MODERATION_ADMIN_ONE_PASSWORD must be set (see .env.ci)} +ADMIN_TWO_HANDLE=${CI_MODERATION_ADMIN_TWO_HANDLE:?CI_MODERATION_ADMIN_TWO_HANDLE must be set (see .env.ci)} +ADMIN_TWO_PASSWORD=${CI_MODERATION_ADMIN_TWO_PASSWORD:?CI_MODERATION_ADMIN_TWO_PASSWORD must be set (see .env.ci)} +CI_PROJECT=${COVES_CI_PROJECT:?COVES_CI_PROJECT must be set by the CI runner} # --------------------------------------------------------------------------- # The relay: raise the crawl limit, then announce both PDSes @@ -80,53 +85,84 @@ announce_host "the AppView's PDS" "${PDS_URL#http://}" announce_host "the federated PDS" "${PDS2_URL#http://}" # --------------------------------------------------------------------------- -# The instance account +# PDS accounts # --------------------------------------------------------------------------- -echo "▶ Creating the instance PDS account ($HANDLE)..." +provision_account() { + local label=$1 handle=$2 email=$3 password=$4 response body session_code account_did + echo "▶ Creating the $label PDS account ($handle)..." -# The PDS runs with PDS_INVITE_REQUIRED=false, so no invite code is needed. -response=$( - curl -sS -o /tmp/createAccount.out -w '%{http_code}' \ - -X POST "$PDS_URL/xrpc/com.atproto.server.createAccount" \ - -H 'Content-Type: application/json' \ - -d "{\"handle\":\"$HANDLE\",\"email\":\"$EMAIL\",\"password\":\"$PASSWORD\"}" -) -body=$(cat /tmp/createAccount.out) + # The PDS runs with PDS_INVITE_REQUIRED=false, so no invite code is needed. + response=$( + curl -sS -o /tmp/createAccount.out -w '%{http_code}' \ + -X POST "$PDS_URL/xrpc/com.atproto.server.createAccount" \ + -H 'Content-Type: application/json' \ + -d "{\"handle\":\"$handle\",\"email\":\"$email\",\"password\":\"$password\"}" + ) + body=$(cat /tmp/createAccount.out) -case "$response" in -200 | 201) - echo " ✓ created" - ;; -400) - # Idempotency: a retried run against a stack that was kept alive - # (COVES_CI_KEEP_STACK=1) will find the handle already taken. Anything else - # in the 400 is a real configuration problem and must not be swallowed. - if grep -qi 'handle.*taken\|already.*exists\|AlreadyExists' <<<"$body"; then - echo " ✓ already exists (reusing)" - else - echo " ✗ PDS rejected the account creation: $body" >&2 - exit 1 + case "$response" in + 200 | 201) + echo " ✓ created" + ;; + 400) + # A kept stack can already have this handle. Any other 400 is a real + # configuration error, not a successful retry. + if grep -qi 'handle.*taken\|already.*exists\|AlreadyExists' <<<"$body"; then + echo " ✓ already exists (reusing)" + else + echo " ✗ PDS rejected the account creation: $body" >&2 + return 1 + fi + ;; + *) + echo " ✗ unexpected HTTP $response from the PDS: $body" >&2 + return 1 + ;; + esac + + # Creating an account and authenticating as it are different claims. + echo "▶ Verifying the $label credentials authenticate..." + session_code=$( + curl -sS -o /tmp/createSession.out -w '%{http_code}' \ + -X POST "$PDS_URL/xrpc/com.atproto.server.createSession" \ + -H 'Content-Type: application/json' \ + -d "{\"identifier\":\"$handle\",\"password\":\"$password\"}" + ) + if [[ $session_code != 200 ]]; then + echo " ✗ could not authenticate as $handle (HTTP $session_code): $(cat /tmp/createSession.out)" >&2 + return 1 fi - ;; -*) - echo " ✗ unexpected HTTP $response from the PDS: $body" >&2 - exit 1 - ;; -esac + echo " ✓ credentials valid" -# Prove the credentials the AppView will use actually work. Creating the account -# and being able to authenticate as it are different claims, and the AppView -# failing to log in surfaces much later and much more confusingly — as community -# writes failing deep inside a test. -echo "▶ Verifying the instance credentials authenticate..." -session_code=$( - curl -sS -o /tmp/createSession.out -w '%{http_code}' \ - -X POST "$PDS_URL/xrpc/com.atproto.server.createSession" \ - -H 'Content-Type: application/json' \ - -d "{\"identifier\":\"$HANDLE\",\"password\":\"$PASSWORD\"}" -) -if [[ $session_code != 200 ]]; then - echo " ✗ could not authenticate as $HANDLE (HTTP $session_code): $(cat /tmp/createSession.out)" >&2 + # The runner has GNU grep but no jq. Do not print the createSession body: + # it contains access and refresh tokens alongside the DID. + if ! account_did=$(grep -oE '"did"[[:space:]]*:[[:space:]]*"did:[a-z]+:[^"]+"' /tmp/createSession.out | cut -d '"' -f 4) || + [[ ! $account_did =~ ^did:[a-z]+:[a-zA-Z0-9._:%-]+$ ]]; then + echo " ✗ PDS session for $handle did not contain a valid DID" >&2 + return 1 + fi + PROVISIONED_DID=$account_did +} + +provision_account "instance" "$HANDLE" "$EMAIL" "$PASSWORD" +provision_account "moderation admin one" "$ADMIN_ONE_HANDLE" "${CI_MODERATION_ADMIN_ONE_EMAIL:-modadmin-one@local.coves.dev}" "$ADMIN_ONE_PASSWORD" +admin_one_did=$PROVISIONED_DID +provision_account "moderation admin two" "$ADMIN_TWO_HANDLE" "${CI_MODERATION_ADMIN_TWO_EMAIL:-modadmin-two@local.coves.dev}" "$ADMIN_TWO_PASSWORD" +admin_two_did=$PROVISIONED_DID + +if [[ -z $admin_one_did || -z $admin_two_did || $admin_one_did == "$admin_two_did" ]]; then + echo " ✗ moderation admin accounts must have distinct, non-empty DIDs" >&2 exit 1 fi -echo " ✓ credentials valid" + +admins_dir=/src/.ci-out +admins_file="$admins_dir/moderation-admins-${CI_PROJECT}.env" +mkdir -p "$admins_dir" +admins_tmp=$(mktemp "$admins_file.tmp.XXXXXX") +trap 'rm -f "$admins_tmp"' EXIT +printf 'MODERATION_ADMINS=%s,%s\n' "$admin_one_did" "$admin_two_did" >"$admins_tmp" +# The runner writes as root, but host-side Compose must be able to read it. +chmod 644 "$admins_tmp" +mv "$admins_tmp" "$admins_file" +trap - EXIT +echo " ✓ moderation admin DIDs saved for the AppView" diff --git a/scripts/lib/ci-stack.sh b/scripts/lib/ci-stack.sh index 1212aec..ef49d55 100644 --- a/scripts/lib/ci-stack.sh +++ b/scripts/lib/ci-stack.sh @@ -13,8 +13,8 @@ # # So the bring-up lives here once, and both callers get the identical stack: # same images built from the working tree, same staging (the AppView cannot -# start until its PDS account exists), same egress-blocked network, same cache -# volumes. +# start until its PDS account and moderation admins exist), same egress-blocked +# network, same cache volumes. # # CONTRACT FOR CALLERS # @@ -261,6 +261,7 @@ stack_discard_previous() { fail "remove them by hand and retry: $(compose_cmdline down -v --remove-orphans)" return 1 fi + rm -f "$OUT_DIR/moderation-admins-${PROJECT}.env" ok "clean slate" } @@ -302,7 +303,7 @@ stack_prefetch_modules() { stack_start() { # Staged deliberately: the AppView is NOT started with the rest. It # authenticates to the PDS as PDS_INSTANCE_HANDLE, and in a fresh PDS that - # account does not exist yet, so it has to be created in between. + # account and the moderation admins do not exist yet. step "Starting infrastructure (Postgres ×4, PLC, PDS ×2, relay, Turnstile stub)" # Jetstream is deliberately absent from this --wait list. `up --wait` fails # outright — "has no healthcheck configured" — for any service it is asked @@ -327,10 +328,25 @@ stack_start() { compose up -d jetstream ok "jetstream started (readiness gated by the runner)" - step "Seeding the stack (relay crawl announcements, instance PDS account)" + step "Seeding the stack (relay crawl announcements, instance and moderation admin accounts)" # --no-deps so this does not drag the AppView up before its account exists. compose run --rm --no-deps --entrypoint bash runner /src/scripts/ci-bootstrap.sh + local admins_file="$OUT_DIR/moderation-admins-${PROJECT}.env" admins_line admins + if [[ ! -f $admins_file ]] || ! IFS= read -r admins_line <"$admins_file"; then + fail "bootstrap did not write $admins_file with MODERATION_ADMINS; refusing to start AppView" + return 1 + fi + if [[ ! $admins_line =~ ^MODERATION_ADMINS=did:[a-z]+:[a-zA-Z0-9._:%-]+,did:[a-z]+:[a-zA-Z0-9._:%-]+$ ]]; then + fail "invalid MODERATION_ADMINS in $admins_file: expected two DIDs; refusing to start AppView" + return 1 + fi + admins=${admins_line#MODERATION_ADMINS=} + if [[ ${admins%,*} == "${admins#*,}" ]]; then + fail "duplicate moderation admin DIDs in $admins_file; refusing to start AppView" + return 1 + fi + step "Starting the AppView" compose up -d --wait appview ok "appview healthy on :8081" diff --git a/tests/e2e/moderation_subject_state_contract_test.go b/tests/e2e/moderation_subject_state_contract_test.go new file mode 100644 index 0000000..96030f1 --- /dev/null +++ b/tests/e2e/moderation_subject_state_contract_test.go @@ -0,0 +1,89 @@ +//go:build e2e + +package e2e + +import ( + "net/http" + "net/url" + "testing" + + "Coves/tests/testkit" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const subjectStateMethod = "social.coves.moderation.getSubjectState" + +type moderationSubjectStateResponse struct { + State struct { + Subject string `json:"subject"` + RecordState string `json:"recordState"` + Version string `json:"version"` + Moderation struct { + State string `json:"state"` + } `json:"moderation"` + CurrentSubject struct { + URI string `json:"uri"` + CID string `json:"cid"` + } `json:"currentSubject"` + } `json:"state"` +} + +func TestModerationSubjectStateAPIContract(t *testing.T) { + p := newPipeline(t) + author := p.IndexedAccount(t, "modsubject") + community := indexedCommunity(t, p, "ms", author.DID) + rkey := testkit.TID() + post := author.PutRecord(t, postV2Collection, rkey, + postV2Record(community.DID, "moderation subject "+testkit.UniqueID(t), "a post awaiting admission")) + awaitStatus(t, p, post.URI, community.DID, "pending", + "the author's post to reach the AppView index through the consumers") + + params := url.Values{"subject": {post.URI}} + admins := []struct { + name string + number int + }{ + {name: "first bootstrap admin", number: 1}, + {name: "second bootstrap admin", number: 2}, + } + for _, admin := range admins { + t.Run(admin.name, func(t *testing.T) { + account := testkit.ModerationAdmin(t, admin.number) + token := account.ServiceAuth(t, communityInstanceDID, subjectStateMethod) + var response moderationSubjectStateResponse + err := p.AppView.As(token).Query(t.Context(), subjectStateMethod, params, &response) + require.NoError(t, err, "the bootstrapped admin must be allowed to read the indexed subject") + assert.Equal(t, post.URI, response.State.Subject) + assert.Equal(t, "present", response.State.RecordState) + assert.Equal(t, "v0", response.State.Version) + assert.Equal(t, "clear", response.State.Moderation.State) + assert.Equal(t, post.URI, response.State.CurrentSubject.URI) + assert.Equal(t, post.CID, response.State.CurrentSubject.CID) + }) + } + + nonAdmin := p.IndexedAccount(t, "modoutsider") + t.Run("non-admin service JWT is forbidden", func(t *testing.T) { + token := nonAdmin.ServiceAuth(t, communityInstanceDID, subjectStateMethod) + err := p.AppView.As(token).Query(t.Context(), subjectStateMethod, params, nil) + requireXRPCRefusal(t, err, http.StatusForbidden, "Forbidden", "a non-admin's valid service JWT") + }) + t.Run("missing credential requires authentication", func(t *testing.T) { + err := p.AppView.Query(t.Context(), subjectStateMethod, params, nil) + requireXRPCRefusal(t, err, http.StatusUnauthorized, "AuthRequired", "an unauthenticated request") + }) + t.Run("service JWT for another method is rejected", func(t *testing.T) { + account := testkit.ModerationAdmin(t, 1) + token := account.ServiceAuth(t, communityInstanceDID, "social.coves.community.post.create") + err := p.AppView.As(token).Query(t.Context(), subjectStateMethod, params, nil) + requireXRPCRefusal(t, err, http.StatusUnauthorized, "AuthRequired", "a service JWT bound to another method") + }) + t.Run("service JWT for another audience is rejected", func(t *testing.T) { + account := testkit.ModerationAdmin(t, 1) + token := account.ServiceAuth(t, "did:web:other.test", subjectStateMethod) + err := p.AppView.As(token).Query(t.Context(), subjectStateMethod, params, nil) + requireXRPCRefusal(t, err, http.StatusUnauthorized, "AuthRequired", "a service JWT addressed to another service") + }) +} diff --git a/tests/testkit/pds.go b/tests/testkit/pds.go index 4ea9d0e..1421f9b 100644 --- a/tests/testkit/pds.go +++ b/tests/testkit/pds.go @@ -6,6 +6,7 @@ import ( "encoding/binary" "fmt" "net/url" + "os" "strings" "sync" "time" @@ -367,6 +368,29 @@ func (p *PDS) Login(t TestingT, identifier, password string) *Account { return acct } +// ModerationAdmin logs into one of the two accounts provisioned by the CI +// bootstrap before the AppView starts. Only testkit reads their credentials. +func ModerationAdmin(t TestingT, n int) *Account { + t.Helper() + if n != 1 && n != 2 { + t.Fatalf("testkit.ModerationAdmin: account number must be 1 or 2, got %d", n) + return nil + } + suffix := "ONE" + if n == 2 { + suffix = "TWO" + } + handleVariable := "CI_MODERATION_ADMIN_" + suffix + "_HANDLE" + passwordVariable := "CI_MODERATION_ADMIN_" + suffix + "_PASSWORD" + handle, password := os.Getenv(handleVariable), os.Getenv(passwordVariable) + if handle == "" || password == "" { + t.Fatalf("testkit.ModerationAdmin: %s and %s must be set; run 'make ci' to bootstrap moderation admin accounts", + handleVariable, passwordVariable) + return nil + } + return NewPDS(t).Login(t, handle, password) +} + // sessionResponse is the body com.atproto.server.createAccount and // com.atproto.server.createSession both return. type sessionResponse struct { @@ -463,6 +487,30 @@ func (a *Account) XRPC() *XRPCClient { return a.client } +// ServiceAuth asks this account's PDS to sign a service JWT for one AppView +// method and audience. Neither the token nor its signing credentials are logged. +func (a *Account) ServiceAuth(t TestingT, audience, lexiconMethod string) string { + t.Helper() + var out struct { + Token string `json:"token"` + } + ctx, cancel := context.WithTimeout(context.Background(), defaultXRPCTimeout) + defer cancel() + err := a.XRPC().Query(ctx, "com.atproto.server.getServiceAuth", url.Values{ + "aud": {audience}, + "lxm": {lexiconMethod}, + }, &out) + if err != nil { + t.Fatalf("testkit: minting service auth for %s via com.atproto.server.getServiceAuth: %v", a.DID, err) + return "" + } + if out.Token == "" { + t.Fatalf("testkit: com.atproto.server.getServiceAuth answered 200 without a token for %s", a.DID) + return "" + } + return out.Token +} + // handleLabel returns the part of a handle before the first dot. func handleLabel(handle string) string { label, _, _ := strings.Cut(handle, ".") -- 2.51.2