diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index 4158cd4..ddd7d27 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -481,10 +481,25 @@ func run(logger *slog.Logger) error { // wired end to end must not start accumulating it. var consumerDone <-chan struct{} if cfg.ConsumerEnabled { - consumerDone, err = startConsumer(ctx, cfg, database, repoManager, personasService, apClient, personasService, logger) + var acceptEngine *accept.Engine + consumerDone, acceptEngine, err = startConsumer(ctx, cfg, database, repoManager, personasService, apClient, personasService, logger) if err != nil { return err } + // The acceptance-engine admin surface (task 16): list admissions with + // their reasons + force re-admit. It shares the /admin bearer and mounts + // only WITH the consumer, because a force re-admit needs the engine. A + // deployment with the consumer off has no admissions to inspect. + acceptAdmin, err := accept.NewAdmin(accept.AdminOptions{ + Token: cfg.AdminToken, + Admissions: accept.NewAdmissions(database), + Engine: acceptEngine, + Logger: logger, + }) + if err != nil { + return err + } + acceptAdmin.Routes(router) } // Host routing wraps everything: the chi router keeps answering for the @@ -606,7 +621,7 @@ func startConsumer( apClient *ap.Client, signers outbound.SignerProvider, logger *slog.Logger, -) (<-chan struct{}, error) { +) (<-chan struct{}, *accept.Engine, error) { // The most SSRF-exposed egress in the bridge: the well-known host comes // from a DID document a stranger controls, so it shares the AP client's // guard rather than using a bare http.Client. @@ -620,7 +635,7 @@ func startConsumer( Logger: logger, }) if err != nil { - return nil, fmt.Errorf("consumer: handle resolver: %w", err) + return nil, nil, fmt.Errorf("consumer: handle resolver: %w", err) } // The outbound delivery pipe (task 15). Because this function runs ONLY when @@ -640,7 +655,7 @@ func startConsumer( Logger: logger, }) if err != nil { - return nil, fmt.Errorf("consumer: outbound enqueuer: %w", err) + return nil, nil, fmt.Errorf("consumer: outbound enqueuer: %w", err) } var worker *outbound.Worker @@ -662,7 +677,7 @@ func startConsumer( Logger: logger, }) if err != nil { - return nil, fmt.Errorf("consumer: outbound worker: %w", err) + return nil, nil, fmt.Errorf("consumer: outbound worker: %w", err) } } @@ -685,7 +700,7 @@ func startConsumer( Logger: logger, }) if err != nil { - return nil, fmt.Errorf("consumer: acceptance engine: %w", err) + return nil, nil, fmt.Errorf("consumer: acceptance engine: %w", err) } dispatcher, err := consume.NewDispatcher(consume.Options{ @@ -701,7 +716,7 @@ func startConsumer( Logger: logger, }) if err != nil { - return nil, fmt.Errorf("consumer: dispatcher: %w", err) + return nil, nil, fmt.Errorf("consumer: dispatcher: %w", err) } // The collection filter is load-bearing: without wantedCollections this @@ -709,7 +724,7 @@ func startConsumer( // by record. subscribeURL, err := consume.SubscribeURL(cfg.JetstreamURL, consume.WantedCollections()) if err != nil { - return nil, fmt.Errorf("consumer: %w", err) + return nil, nil, fmt.Errorf("consumer: %w", err) } state := consume.NewPostgresStateStore(database, consume.CursorSchemaVersion) @@ -760,5 +775,5 @@ func startConsumer( close(done) }() logger.Info("jetstream consumer started", "url", subscribeURL) - return done, nil + return done, engine, nil } diff --git a/internal/accept/admin.go b/internal/accept/admin.go new file mode 100644 index 0000000..29860bf --- /dev/null +++ b/internal/accept/admin.go @@ -0,0 +1,179 @@ +package accept + +import ( + "crypto/subtle" + "encoding/json" + stderrors "errors" + "log/slog" + "net/http" + + "github.com/go-chi/chi/v5" + + "tidepool/internal/errors" +) + +// AdminOptions wires the admissions admin surface. +type AdminOptions struct { + // Token is the bearer token protecting /admin (config.AdminToken), matching + // the ingest.Admin pattern. + Token string + // Admissions is the decision ledger the list endpoint reads. + Admissions *Admissions + // Engine serves the force re-admit (Readmit). + Engine *Engine + // Logger receives rejection/authz warnings. Nil uses slog.Default(). + Logger *slog.Logger +} + +// Admin is the operator API for the acceptance engine's admission ledger: +// +// GET /admin/admissions ?status=&community= (list decisions + reasons) +// POST /admin/admissions/readmit {"post":"at://..."} (force re-admit one post) +// +// Both endpoints require "Authorization: Bearer $ADMIN_TOKEN". post.getStatus +// reads a post's admission state from the firehose-visible acceptance/removal +// records; THIS surface is the bridge operator's own window on the WHY of every +// rejection, which those records cannot carry. +type Admin struct { + token string + admissions *Admissions + engine *Engine + logger *slog.Logger +} + +// NewAdmin validates options and builds the admin API. +func NewAdmin(opts AdminOptions) (*Admin, error) { + switch { + case opts.Token == "": + return nil, errors.NewValidationError("token", "must not be empty") + case opts.Admissions == nil: + return nil, errors.NewValidationError("admissions", "must not be nil") + case opts.Engine == nil: + return nil, errors.NewValidationError("engine", "must not be nil") + } + logger := opts.Logger + if logger == nil { + logger = slog.Default() + } + return &Admin{token: opts.Token, admissions: opts.Admissions, engine: opts.Engine, logger: logger}, nil +} + +// Routes mounts the admissions admin API, bearer-protected like the rest of +// /admin. It registers its full /admin/... paths inside a middleware GROUP +// rather than a Route("/admin") subrouter, so it composes onto a router that +// already mounts a /admin subtree (ingest.Admin) instead of colliding with it. +func (a *Admin) Routes(r chi.Router) { + r.Group(func(r chi.Router) { + r.Use(adminBearer(a.token, a.logger)) + r.Get("/admin/admissions", a.handleList) + r.Post("/admin/admissions/readmit", a.handleReadmit) + }) +} + +// adminListItem is one admission on the wire: the operator's triage view of a +// decision. The evaluated_snapshot is deliberately omitted (it can be large and +// is only readmit's input). +type adminListItem struct { + Post string `json:"post"` + Community string `json:"community"` + Status string `json:"status"` + DecisionCode string `json:"decisionCode"` + EvaluatedCID string `json:"evaluatedCid"` +} + +// handleList lists admissions with their status, decision_code and evaluated_cid, +// filterable by ?status= and/or ?community=. +func (a *Admin) handleList(w http.ResponseWriter, r *http.Request) { + admissions, err := a.admissions.List(r.Context(), AdmissionFilter{ + Status: r.URL.Query().Get("status"), + Community: r.URL.Query().Get("community"), + }) + if err != nil { + a.logger.Error("admin list admissions failed", "error", err) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "internal error"}) + return + } + items := make([]adminListItem, 0, len(admissions)) + for _, adm := range admissions { + items = append(items, adminListItem{ + Post: adm.PostURI, + Community: adm.CommunityDID, + Status: adm.Status, + DecisionCode: adm.DecisionCode, + EvaluatedCID: adm.EvaluatedCID, + }) + } + writeJSON(w, http.StatusOK, map[string]any{"admissions": items}) +} + +// readmitRequest is the POST body naming the post to force re-admit. +type readmitRequest struct { + Post string `json:"post"` +} + +// readmitResult is the wire outcome of a force re-admit. +type readmitResult struct { + Post string `json:"post"` + Status string `json:"status"` + DecisionCode string `json:"decisionCode"` + Enqueued bool `json:"enqueued"` +} + +// handleReadmit force re-runs admission for one post (Engine.Readmit). A post +// that passes now is accepted + enqueued; one that still fails returns the +// current rejection (a reported outcome, HTTP 200, not an error). An unknown post +// is 404; a post whose stored snapshot did not survive is 422 (unrecoverable +// without a task-18 getRecord fetch). +func (a *Admin) handleReadmit(w http.ResponseWriter, r *http.Request) { + var body readmitRequest + if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.Post == "" { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "body must be {\"post\":\"at://...\"}"}) + return + } + result, err := a.engine.Readmit(r.Context(), body.Post) + switch { + case err == nil: + writeJSON(w, http.StatusOK, readmitResult{ + Post: result.PostURI, + Status: result.Status, + DecisionCode: result.DecisionCode, + Enqueued: result.Enqueued, + }) + case errors.IsNotFound(err): + writeJSON(w, http.StatusNotFound, map[string]string{"error": "no admission for " + body.Post}) + case stderrors.Is(err, ErrUnrecoverableReadmit): + // Surfaced, never a silent no-op: the record body the post was decided + // against did not survive (legacy row, pre-migration-022), so there is + // nothing to re-run admission from until a getRecord fetch lands (task 18). + writeJSON(w, http.StatusUnprocessableEntity, map[string]string{ + "error": "readmit unrecoverable: no stored record snapshot", "post": body.Post}) + default: + a.logger.Error("admin readmit failed", "post", body.Post, "error", err) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "internal error"}) + } +} + +// adminBearer is the constant-time bearer guard, mirroring ingest.requireBearer +// (kept local so this surface stays self-contained). +func adminBearer(token string, logger *slog.Logger) func(http.Handler) http.Handler { + return func(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + const prefix = "Bearer " + header := r.Header.Get("Authorization") + if len(header) <= len(prefix) || header[:len(prefix)] != prefix || + subtle.ConstantTimeCompare([]byte(header[len(prefix):]), []byte(token)) != 1 { + logger.Warn("admin request rejected: bad bearer token", "path", r.URL.Path) + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + next.ServeHTTP(w, r) + }) + } +} + +// writeJSON is the shared JSON responder GREEN's handlers use. +func writeJSON(w http.ResponseWriter, status int, body any) { + w.Header().Set("Content-Type", "application/json; charset=utf-8") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(body) +} diff --git a/internal/accept/admin_test.go b/internal/accept/admin_test.go new file mode 100644 index 0000000..50bcb7f --- /dev/null +++ b/internal/accept/admin_test.go @@ -0,0 +1,244 @@ +package accept + +import ( + "context" + "database/sql" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/go-chi/chi/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// Round 3: the admin/debug surface — the last task-16 deliverable. Coves' +// post.getStatus reads a post's admission state from the firehose-visible +// acceptance/removal records; THIS surface is the bridge operator's own window +// on the WHY of every rejection, which those records cannot carry, plus a force +// re-admit (the mod-override seam task 17 reuses for restore). + +const adminToken = "admin-secret-token" + +// adminRequest issues one admin request and returns the recorder. An empty token +// omits the Authorization header (the 401 path). +func adminRequest(t *testing.T, router chi.Router, method, path, token, body string) *httptest.ResponseRecorder { + t.Helper() + var r io.Reader + if body != "" { + r = strings.NewReader(body) + } + req := httptest.NewRequest(method, "https://bridge.example"+path, r) + if token != "" { + req.Header.Set("Authorization", "Bearer "+token) + } + if body != "" { + req.Header.Set("Content-Type", "application/json") + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req) + return rec +} + +// newAdminRouter mounts an Admin over the given engine + this DB's admissions. +func newAdminRouter(t *testing.T, conn *sql.DB, engine *Engine) chi.Router { + t.Helper() + admin, err := NewAdmin(AdminOptions{ + Token: adminToken, + Admissions: NewAdmissions(conn), + Engine: engine, + }) + require.NoError(t, err) + router := chi.NewRouter() + admin.Routes(router) + return router +} + +// listItem is the wire contract this surface pins for one listed admission. +type listItem struct { + Post string `json:"post"` + Community string `json:"community"` + Status string `json:"status"` + DecisionCode string `json:"decisionCode"` + EvaluatedCID string `json:"evaluatedCid"` +} + +type listResponse struct { + Admissions []listItem `json:"admissions"` +} + +// readmitResponse is the wire contract for a force re-admit. +type readmitResponse struct { + Post string `json:"post"` + Status string `json:"status"` + DecisionCode string `json:"decisionCode"` + Enqueued bool `json:"enqueued"` +} + +// --------------------------------------------------------------------------- +// A1 — GET /admin/admissions: list decisions + reasons, filterable by status +// --------------------------------------------------------------------------- + +func TestAdminList_FiltersByStatusAndReturnsReasons(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + adm := NewAdmissions(conn) + + // Four decisions in the ledger: one accepted, two rejected (distinct + // reasons), one removed. + seed := []Admission{ + {CommunityDID: acCommunityDID, PostURI: acPostURI + "-acc", AuthorDID: acAuthorDID, Status: StatusAccepted, EvaluatedCID: acPostCID, AcceptanceRKey: "rk", AcceptedCID: acPostCID}, + {CommunityDID: acCommunityDID, PostURI: acPostURI + "-rej1", AuthorDID: acAuthorDID, Status: StatusRejected, DecisionCode: DecisionTitleRequired, EvaluatedCID: acPostCID}, + {CommunityDID: acCommunityDID, PostURI: acPostURI + "-rej2", AuthorDID: acAuthorDID, Status: StatusRejected, DecisionCode: DecisionRateLimit, EvaluatedCID: acPostCID}, + {CommunityDID: acCommunityDID, PostURI: acPostURI + "-rem", AuthorDID: acAuthorDID, Status: StatusRemoved, DecisionCode: DecisionTitleRequired, EvaluatedCID: acPostCID}, + } + for _, a := range seed { + require.NoError(t, adm.Record(ctx, a)) + } + + engine := engineWith(t, conn, newRepos(t, conn), realEnqueuer(t, conn)) + router := newAdminRouter(t, conn, engine) + + rec := adminRequest(t, router, http.MethodGet, "/admin/admissions?status=rejected", adminToken, "") + require.Equal(t, http.StatusOK, rec.Code, "listing rejected admissions must succeed") + + var resp listResponse + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &resp), + "the list body is {\"admissions\":[{post,community,status,decisionCode,evaluatedCid}]}") + + codes := map[string]string{} + for _, item := range resp.Admissions { + assert.Equal(t, StatusRejected, item.Status, "status=rejected must return only rejected rows") + codes[item.Post] = item.DecisionCode + } + require.Len(t, resp.Admissions, 2, + "status=rejected returns exactly the two rejected admissions, not the accepted or removed ones") + assert.Equal(t, DecisionTitleRequired, codes[acPostURI+"-rej1"], + "each rejection carries its distinct machine-readable reason — the whole point of this surface") + assert.Equal(t, DecisionRateLimit, codes[acPostURI+"-rej2"]) +} + +// --------------------------------------------------------------------------- +// A2 — POST /admin/admissions/readmit: force re-admit one post +// --------------------------------------------------------------------------- + +// A rejected post whose cause has cleared (opted-out → re-enabled) is re-admitted +// from stored state: acceptance written, Create{Page} enqueued, ledger accepted. +func TestAdminReadmit_ReAdmitsAPostWhoseRejectionCauseCleared(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + engine := engineWith(t, conn, repos, enq) + dispatcher := wireDispatcher(t, conn, engine, enq) + router := newAdminRouter(t, conn, engine) + + // Opted out → the post is rejected (and its record snapshot is stored on the + // ledger row, so readmit can re-run from stored state). + _, err := store.NewFederationPrefs(conn).Upsert(ctx, store.FederationPref{ + DID: acAuthorDID, Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + status, code := admissionOf(t, conn, acCommunityDID, acPostURI) + require.Equal(t, StatusRejected, status) + require.Equal(t, DecisionOptedOut, code) + + // The author re-enables federation. + require.NoError(t, store.NewFederationPrefs(conn).Delete(ctx, acAuthorDID)) + + // Force re-admit. + rec := adminRequest(t, router, http.MethodPost, "/admin/admissions/readmit", adminToken, + `{"post":"`+acPostURI+`"}`) + require.Equal(t, http.StatusOK, rec.Code, "readmit of a now-eligible post succeeds") + + var resp readmitResponse + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &resp)) + assert.Equal(t, StatusAccepted, resp.Status, "the re-admission now passes") + assert.True(t, resp.Enqueued, "a newly accepted post enqueues its Create{Page}") + + // The acceptance now stands, the Create{Page} is enqueued, and the ledger is + // accepted. + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.True(t, ok, "readmit wrote the acceptance record") + assert.Equal(t, 1, activityKindCount(t, conn, "Create"), + "readmit enqueued exactly one Create{Page}") + status, _ = admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusAccepted, status, "the ledger row flips to accepted") +} + +// A post that STILL fails admission (opted-out, never re-enabled) returns the +// failure: still rejected, nothing enqueued, no acceptance. +func TestAdminReadmit_StillFailingPostStaysRejected(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + engine := engineWith(t, conn, repos, enq) + dispatcher := wireDispatcher(t, conn, engine, enq) + router := newAdminRouter(t, conn, engine) + + _, err := store.NewFederationPrefs(conn).Upsert(ctx, store.FederationPref{ + DID: acAuthorDID, Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + + // Readmit WITHOUT clearing the opt-out. + rec := adminRequest(t, router, http.MethodPost, "/admin/admissions/readmit", adminToken, + `{"post":"`+acPostURI+`"}`) + require.Equal(t, http.StatusOK, rec.Code, + "a readmit that still fails is a reported outcome, not an HTTP error") + + var resp readmitResponse + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &resp)) + assert.Equal(t, StatusRejected, resp.Status, "the post still fails admission") + assert.Equal(t, DecisionOptedOut, resp.DecisionCode, "and the response carries why") + assert.False(t, resp.Enqueued, "nothing is enqueued for a post that still fails") + + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.False(t, ok, "no acceptance is written") + assert.Zero(t, countRows(t, conn, "outbound_activities"), "and nothing is enqueued") +} + +// --------------------------------------------------------------------------- +// A3 — authz: both endpoints require the admin bearer +// --------------------------------------------------------------------------- + +func TestAdminAdmissions_RequireBearer(t *testing.T) { + conn := acceptanceDB(t) + engine := engineWith(t, conn, newRepos(t, conn), realEnqueuer(t, conn)) + router := newAdminRouter(t, conn, engine) + + for _, tc := range []struct { + name, method, path, body string + }{ + {"list", http.MethodGet, "/admin/admissions", ""}, + {"readmit", http.MethodPost, "/admin/admissions/readmit", `{"post":"at://x/y/z"}`}, + } { + t.Run(tc.name, func(t *testing.T) { + // No token. + rec := adminRequest(t, router, tc.method, tc.path, "", tc.body) + assert.Equal(t, http.StatusUnauthorized, rec.Code, "no bearer → 401") + + // Wrong token. + rec = adminRequest(t, router, tc.method, tc.path, "wrong-token", tc.body) + assert.Equal(t, http.StatusUnauthorized, rec.Code, "wrong bearer → 401") + + // Valid token: auth passes (the handler may still be unimplemented, + // but it is NOT a 401). + rec = adminRequest(t, router, tc.method, tc.path, adminToken, tc.body) + assert.NotEqual(t, http.StatusUnauthorized, rec.Code, + "a valid bearer must pass the guard") + }) + } +} diff --git a/internal/accept/admissions.go b/internal/accept/admissions.go index 01f5e16..4933aa6 100644 --- a/internal/accept/admissions.go +++ b/internal/accept/admissions.go @@ -33,6 +33,12 @@ type Admission struct { AcceptanceRKey string AcceptedCID string Redrivable bool + // EvaluatedSnapshot is the postv2 record (plus resolved context) this + // decision was made against, stored on EVERY decision so /admin/admissions/ + // readmit can re-run admission from stored state (migration 022). Empty + // ('{}') means the body did not survive — readmit is unrecoverable without a + // getRecord fetch from the author's PDS (task 18). + EvaluatedSnapshot []byte } // Admissions persists the decision ledger. @@ -69,11 +75,17 @@ func (a *Admissions) record(ctx context.Context, ex execer, adm Admission) error if adm.Status == "" { return errors.NewValidationError("admission.status", "must be set") } + // A JSONB NOT NULL column rejects a NULL, so an unset snapshot coalesces to + // the empty object the DEFAULT would have used. + snapshot := adm.EvaluatedSnapshot + if len(snapshot) == 0 { + snapshot = []byte("{}") + } _, err := ex.ExecContext(ctx, ` INSERT INTO admissions (community_did, post_uri, author_did, status, decision_code, evaluated_cid, - acceptance_rkey, accepted_cid, redrivable) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) + acceptance_rkey, accepted_cid, redrivable, evaluated_snapshot) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) ON CONFLICT (community_did, post_uri) DO UPDATE SET author_did = EXCLUDED.author_did, status = EXCLUDED.status, @@ -82,9 +94,10 @@ func (a *Admissions) record(ctx context.Context, ex execer, adm Admission) error acceptance_rkey = EXCLUDED.acceptance_rkey, accepted_cid = EXCLUDED.accepted_cid, redrivable = EXCLUDED.redrivable, + evaluated_snapshot = EXCLUDED.evaluated_snapshot, updated_at = now()`, adm.CommunityDID, adm.PostURI, adm.AuthorDID, adm.Status, adm.DecisionCode, adm.EvaluatedCID, - adm.AcceptanceRKey, adm.AcceptedCID, adm.Redrivable) + adm.AcceptanceRKey, adm.AcceptedCID, adm.Redrivable, snapshot) if err != nil { return fmt.Errorf("accept: record admission %s/%s: %w", adm.CommunityDID, adm.PostURI, err) } @@ -97,11 +110,11 @@ func (a *Admissions) Get(ctx context.Context, communityDID, postURI string) (*Ad var adm Admission err := a.db.QueryRowContext(ctx, ` SELECT community_did, post_uri, author_did, status, decision_code, evaluated_cid, - acceptance_rkey, accepted_cid, redrivable + acceptance_rkey, accepted_cid, redrivable, evaluated_snapshot FROM admissions WHERE community_did = $1 AND post_uri = $2`, communityDID, postURI).Scan( &adm.CommunityDID, &adm.PostURI, &adm.AuthorDID, &adm.Status, &adm.DecisionCode, &adm.EvaluatedCID, - &adm.AcceptanceRKey, &adm.AcceptedCID, &adm.Redrivable) + &adm.AcceptanceRKey, &adm.AcceptedCID, &adm.Redrivable, &adm.EvaluatedSnapshot) if stderrors.Is(err, sql.ErrNoRows) { return nil, errors.NewNotFoundError("admission", communityDID+"/"+postURI) } @@ -111,6 +124,61 @@ func (a *Admissions) Get(ctx context.Context, communityDID, postURI string) (*Ad return &adm, nil } +// GetByPostURI returns the admission for a post at-uri alone — the readmit and +// admin-list path, which knows the post but not necessarily its community. The +// post_uri is globally unique (it embeds the author repo), so at most one row +// matches. A miss satisfies errors.IsNotFound. +func (a *Admissions) GetByPostURI(ctx context.Context, postURI string) (*Admission, error) { + var adm Admission + err := a.db.QueryRowContext(ctx, ` + SELECT community_did, post_uri, author_did, status, decision_code, evaluated_cid, + acceptance_rkey, accepted_cid, redrivable, evaluated_snapshot + FROM admissions WHERE post_uri = $1`, postURI).Scan( + &adm.CommunityDID, &adm.PostURI, &adm.AuthorDID, &adm.Status, &adm.DecisionCode, &adm.EvaluatedCID, + &adm.AcceptanceRKey, &adm.AcceptedCID, &adm.Redrivable, &adm.EvaluatedSnapshot) + if stderrors.Is(err, sql.ErrNoRows) { + return nil, errors.NewNotFoundError("admission", postURI) + } + if err != nil { + return nil, fmt.Errorf("accept: get admission %s: %w", postURI, err) + } + return &adm, nil +} + +// AdmissionFilter narrows a List. Empty fields are wildcards. +type AdmissionFilter struct { + Status string + Community string +} + +// List returns admissions matching the filter, newest first — the admin +// surface's read of pending/rejected/removed decisions with their reasons. The +// snapshot is deliberately NOT returned (it can be large and the listing is a +// triage view); readmit reads it through GetByPostURI. +func (a *Admissions) List(ctx context.Context, filter AdmissionFilter) ([]Admission, error) { + rows, err := a.db.QueryContext(ctx, ` + SELECT community_did, post_uri, author_did, status, decision_code, evaluated_cid, + acceptance_rkey, accepted_cid, redrivable + FROM admissions + WHERE ($1 = '' OR status = $1) + AND ($2 = '' OR community_did = $2) + ORDER BY updated_at DESC`, filter.Status, filter.Community) + if err != nil { + return nil, fmt.Errorf("accept: list admissions: %w", err) + } + defer func() { _ = rows.Close() }() + var out []Admission + for rows.Next() { + var adm Admission + if err := rows.Scan(&adm.CommunityDID, &adm.PostURI, &adm.AuthorDID, &adm.Status, + &adm.DecisionCode, &adm.EvaluatedCID, &adm.AcceptanceRKey, &adm.AcceptedCID, &adm.Redrivable); err != nil { + return nil, fmt.Errorf("accept: scan admission: %w", err) + } + out = append(out, adm) + } + return out, rows.Err() +} + // CountAccepted reports how many posts one author currently has ACCEPTED in one // community, excluding one post_uri (the post being decided — a repin must not // count against its own author). It backs the per-author-per-community flood diff --git a/internal/accept/engine.go b/internal/accept/engine.go index 961a336..3d3577d 100644 --- a/internal/accept/engine.go +++ b/internal/accept/engine.go @@ -17,8 +17,10 @@ import ( "context" "database/sql" "encoding/json" + stderrors "errors" "fmt" "log/slog" + "strings" "time" "github.com/bluesky-social/indigo/atproto/atdata" @@ -254,6 +256,10 @@ func (e *Engine) AdmitPost(ctx context.Context, did string, commit *consume.Comm Status: StatusRejected, DecisionCode: code, EvaluatedCID: commit.CID, + // The record body + context this decision was made against, so a later + // force re-admit can re-run admission from stored state (a rejection + // writes no outbound_objects, and the author's PDS is not local). + EvaluatedSnapshot: evaluatedSnapshot(commit), }) } @@ -405,13 +411,14 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, return err } return e.admissions.RecordTx(sctx, tx, Admission{ - AuthorDID: did, - CommunityDID: communityDID, - PostURI: postURI, - Status: StatusAccepted, - EvaluatedCID: commit.CID, - AcceptanceRKey: rkey, - AcceptedCID: commit.CID, + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postURI, + Status: StatusAccepted, + EvaluatedCID: commit.CID, + AcceptanceRKey: rkey, + AcceptedCID: commit.CID, + EvaluatedSnapshot: evaluatedSnapshot(commit), }) } @@ -450,12 +457,13 @@ func (e *Engine) removeAccepted(ctx context.Context, did, communityDID, postURI return err } return e.admissions.RecordTx(sctx, tx, Admission{ - AuthorDID: did, - CommunityDID: communityDID, - PostURI: postURI, - Status: StatusRemoved, - DecisionCode: code, - EvaluatedCID: commit.CID, + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postURI, + Status: StatusRemoved, + DecisionCode: code, + EvaluatedCID: commit.CID, + EvaluatedSnapshot: evaluatedSnapshot(commit), }) } @@ -586,3 +594,151 @@ func publishedAtOf(record map[string]any) time.Time { } return time.Time{} } + +// ReadmitResult reports the outcome of a force re-admit (A2). Enqueued is true +// when the re-admission newly wrote/repinned the acceptance and enqueued a +// Create/Update{Page}; false when the post STILL fails admission (the result +// then carries the current rejection Status and DecisionCode). +type ReadmitResult struct { + PostURI string + Status string + DecisionCode string + Enqueued bool +} + +// ErrUnrecoverableReadmit is returned by Readmit when the post's admissions row +// carries no stored record snapshot ('{}' — a legacy row, or a decision made +// before migration 022). Admission cannot be re-run from nothing, and the postv2 +// lives in the author's native PDS Tidepool does not host, so this surfaces as a +// distinct error (HTTP 422) rather than a silent no-op. A task-18 +// com.atproto.repo.getRecord fetch would recover it. +var ErrUnrecoverableReadmit = stderrors.New("accept: no stored record snapshot to re-run admission from") + +// Readmit force re-runs admission for ONE post (the mod-override seam task 17 +// reuses for restore). It re-runs the SAME decide() path against the record +// snapshot stored on the post's admissions row (evaluated_snapshot, migration +// 022) — task 16 re-admits from STORED STATE, because the postv2 lives in the +// author's native PDS that Tidepool does not host, so there is no local record to +// re-read. Passes now → the acceptance is written/repinned and a Create/Update +// {Page} enqueued (Enqueued=true); still fails → the current rejection/removal is +// reported and nothing is enqueued. +func (e *Engine) Readmit(ctx context.Context, postATURI string) (*ReadmitResult, error) { + adm, err := e.admissions.GetByPostURI(ctx, postATURI) + if err != nil { + return nil, err // NotFound flows through; the handler maps it to 404. + } + + commit, did, ok := rebuildCommit(postATURI, adm.EvaluatedSnapshot) + if !ok { + return nil, ErrUnrecoverableReadmit + } + communityDID, _ := commit.Record["community"].(string) + if communityDID == "" { + communityDID = adm.CommunityDID + } + + // Re-run the exact admission pipeline AdmitPost uses — no duplicated policy. + prior, priorBound, err := e.priorBinding(ctx, postATURI) + if err != nil { + return nil, err + } + code, discard, err := e.decide(ctx, did, commit, communityDID, prior, priorBound) + if err != nil { + return nil, err + } + if discard { + // The stored community no longer matches the event's — nothing is written; + // reported as still-failing with the immutability cause. + return &ReadmitResult{PostURI: postATURI, Status: adm.Status, DecisionCode: DecisionCommunityImmutable}, nil + } + if code != "" { + priorAccepted := priorBound && !prior.IsTombstoned() + if priorAccepted { + // Was accepted, now fails: this is a removal, exactly as AdmitPost would. + if err := e.removeAccepted(ctx, did, communityDID, postATURI, commit, prior, code); err != nil { + return nil, err + } + return &ReadmitResult{PostURI: postATURI, Status: StatusRemoved, DecisionCode: code}, nil + } + // Still rejected: refresh the ledger with the current cause (idempotent), + // enqueue nothing. + if err := e.admissions.Record(ctx, Admission{ + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postATURI, + Status: StatusRejected, + DecisionCode: code, + EvaluatedCID: commit.CID, + EvaluatedSnapshot: evaluatedSnapshot(commit), + }); err != nil { + return nil, err + } + return &ReadmitResult{PostURI: postATURI, Status: StatusRejected, DecisionCode: code}, nil + } + + // Passes now: write/repin the acceptance and enqueue the Page (reusing accept()). + if err := e.accept(ctx, did, communityDID, postATURI, commit); err != nil { + return nil, err + } + return &ReadmitResult{PostURI: postATURI, Status: StatusAccepted, Enqueued: true}, nil +} + +// evaluatedSnapshot serializes the postv2 record and the context a Readmit needs +// to rebuild the CommitEvent it re-runs admission against. It is stored on EVERY +// decision (accept, reject, remove). A marshal failure yields nil, which the +// store coalesces to '{}' — an unrecoverable readmit, never a wrong one. +func evaluatedSnapshot(commit *consume.CommitEvent) []byte { + b, err := json.Marshal(map[string]any{ + "record": commit.Record, + "cid": commit.CID, + "rev": commit.Rev, + "operation": commit.Operation, + "collection": commit.Collection, + }) + if err != nil { + return nil + } + return b +} + +// rebuildCommit reconstructs the CommitEvent (and its author DID) a Readmit +// re-runs admission against, from the post at-uri and the stored evaluated +// snapshot. ok=false means the snapshot did not survive (legacy '{}' or a +// malformed at-uri): the caller surfaces that as ErrUnrecoverableReadmit. +func rebuildCommit(postURI string, snapshot []byte) (commit *consume.CommitEvent, authorDID string, ok bool) { + trimmed := strings.TrimPrefix(postURI, "at://") + parts := strings.SplitN(trimmed, "/", 3) + if len(parts) != 3 || parts[0] == "" || parts[1] == "" || parts[2] == "" { + return nil, "", false + } + did, collection, rkey := parts[0], parts[1], parts[2] + + if len(snapshot) == 0 { + return nil, "", false + } + var snap struct { + Record map[string]any `json:"record"` + CID string `json:"cid"` + Rev string `json:"rev"` + Operation string `json:"operation"` + } + if err := json.Unmarshal(snapshot, &snap); err != nil { + return nil, "", false + } + if len(snap.Record) == 0 { + // '{}' — a decision recorded before the snapshot column existed. + return nil, "", false + } + op := snap.Operation + if op == "" { + op = "create" + } + return &consume.CommitEvent{ + Rev: snap.Rev, + Operation: op, + Collection: collection, + RKey: rkey, + CID: snap.CID, + Record: snap.Record, + }, did, true +} diff --git a/internal/db/migrations/022_admission_snapshot.sql b/internal/db/migrations/022_admission_snapshot.sql new file mode 100644 index 0000000..76926e6 --- /dev/null +++ b/internal/db/migrations/022_admission_snapshot.sql @@ -0,0 +1,21 @@ +-- +goose Up +-- Task 16 admin surface (readmit recoverability). +-- +-- POST /admin/admissions/readmit re-runs admission for ONE post. The primary +-- use case is re-admitting a post that was REJECTED (opted-out, paused, …) once +-- the cause clears. But a rejection writes NO outbound_objects row — that state +-- is created only on ACCEPT — so the postv2 record body a rejected post was +-- decided against does not survive anywhere. The author's repo is a native PDS +-- Tidepool does not host, so the body cannot be read back locally. +-- +-- evaluated_snapshot stores the record (plus resolved context) the engine +-- decided against, on EVERY decision (accept AND reject), so readmit can re-run +-- from stored state alone. A fresh com.atproto.repo.getRecord fetch from the +-- author's PDS (fresher content) is deferred to task 18; this column is the +-- task-16 baseline. Empty ('{}') marks a decision made before this column +-- existed, or a legacy row — readmit surfaces that as an unrecoverable readmit +-- rather than silently re-admitting nothing. +ALTER TABLE admissions ADD COLUMN evaluated_snapshot JSONB NOT NULL DEFAULT '{}'::jsonb; + +-- +goose Down +ALTER TABLE admissions DROP COLUMN IF EXISTS evaluated_snapshot;