diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index 6ac2bd7..85fd8bb 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -149,6 +149,9 @@ task documents and git history rather than this list. serving the repos (crawl imports state, not per-op events). Live records heal organically (any bridgedStats vote sweep or upstream edit re-emits the record as an update commit, which the AppView upserts), but quiet - posts stay invisible. Consider an admin "re-emit" endpoint that walks a - repo and appends no-op update commits for records older than the relay - subscription, for operators onboarding a relay after content exists. + posts stay invisible. RESOLVED: `POST /admin/reemit` walks a repo (or + all repos) and re-emits every record as a delete+create commit pair — + honest diffs, unchanged at-uris/CIDs (a same-value "touch" update was + rejected as the design: its op list would not match the empty MST diff, + which sync-v1.1-validating relays may refuse). Remaining follow-up: the + e2e harness has no scenario covering it. diff --git a/README.md b/README.md index 49214cf..aff4ca9 100644 --- a/README.md +++ b/README.md @@ -302,6 +302,16 @@ curl -X DELETE localhost:8091/admin/communities \ curl localhost:8091/admin/metrics \ -H "Authorization: Bearer dev-admin-token" # expvar counters + +# re-emit a repo's records onto the firehose as delete+create commit pairs +# (identical values, so at-uris and CIDs are unchanged). For the relay +# cold-start gap: records committed before a relay first subscribed never +# re-emit on their own, so a Jetstream-fed AppView cannot index them. +# {} (or empty body) re-emits every active repo; tombstoned repos are +# always skipped. +curl -X POST localhost:8091/admin/reemit \ + -H "Authorization: Bearer dev-admin-token" \ + -d '{"did":"did:plc:aaa..."}' # (includes tidepool_lexicon_validation_failures — non-zero in production, # where validation failures log-and-write, means investigate) ``` @@ -399,6 +409,10 @@ it as a `subscribeRepos` upstream: `listRepos` (paginated), `getRepoStatus` - `com.atproto.server.describeServer`, `/xrpc/_health` - `com.atproto.identity.resolveHandle` + `/.well-known/atproto-did` (task 03) +- `/.well-known/did.json` — the DID document for the bridge's own derived + `did:web:` service identity (the `hostedBy` of every + bridged community profile; consumers verifying that claim resolve it + here). 404 when a non-did:web `BRIDGE_SERVICE_DID` is provisioned. Repos whose actor revoked consent (tombstoned) report `RepoDeactivated` / `active: false` with status `deleted`, and their content endpoints stop diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index b2fecf2..c7162e5 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -123,17 +123,33 @@ func run(logger *slog.Logger) error { _, _ = w.Write([]byte("ok")) }) + // The bridge's own DID (community.profile createdBy/hostedBy, the bare + // hostname's handle, and /.well-known/did.json). An operator may + // pre-provision one; otherwise the bridge identifies as did:web on its + // own hostname. Derived HERE, before the resolver, so the bare-hostname + // handle resolves to it and the did:web document is served — consumers + // (the Coves AppView's hostedBy verification) resolve that document and + // fail closed when it is missing. + serviceDID := cfg.BridgeServiceDID + if serviceDID == "" { + serviceDID = "did:web:" + cfg.BridgeHostname + logger.Info("BRIDGE_SERVICE_DID not set, deriving from hostname", "did", serviceDID) + } + // Handle resolution for the bridged handle space (task 03). Bridged // handles are subdomains of BRIDGE_HOSTNAME; wildcard DNS routes them // all here (see README, "Handle resolution & DNS"). actors := store.NewBridgedActors(database) - resolver := identity.NewStoreResolver(actors, cfg.BridgeHostname, cfg.BridgeServiceDID) + resolver := identity.NewStoreResolver(actors, cfg.BridgeHostname, serviceDID) router.Get("/xrpc/com.atproto.identity.resolveHandle", identity.ResolveHandleHandler(resolver, logger)) router.Get("/.well-known/atproto-did", identity.WellKnownDIDHandler(resolver, logger)) // Cert-issuance gate for TLS-terminating proxies with on-demand // issuance (production Caddy asks here before requesting a cert for a // bridged-handle subdomain; see docker-compose.prod.yml header). router.Get("/.well-known/tidepool-tls-ask", identity.TLSAskHandler(resolver, logger)) + // The bridge's own did:web document (404 when a non-did:web + // BRIDGE_SERVICE_DID is provisioned). + router.Get("/.well-known/did.json", identity.DIDWebHandler(serviceDID, cfg.BridgeHostname)) // The sync surface (task 04): com.atproto.sync.* + subscribeRepos, // describeServer, _health — everything a relay or Jetstream needs to @@ -242,15 +258,6 @@ func run(logger *slog.Logger) error { return err } - // The bridge's own DID (community.profile createdBy/hostedBy). An - // operator may pre-provision one; otherwise the bridge identifies as - // did:web on its own hostname. - serviceDID := cfg.BridgeServiceDID - if serviceDID == "" { - serviceDID = "did:web:" + cfg.BridgeHostname - logger.Info("BRIDGE_SERVICE_DID not set, deriving from hostname", "did", serviceDID) - } - objects := store.NewAPObjects(database) communities := store.NewCommunities(database) tombstones := store.NewTombstones(database) @@ -394,6 +401,7 @@ func run(logger *slog.Logger) error { Communities: communities, Service: serviceActor, Backfill: backfill, + Repos: repoManager, Logger: logger, }) if err != nil { diff --git a/internal/identity/didweb.go b/internal/identity/didweb.go new file mode 100644 index 0000000..0432cbf --- /dev/null +++ b/internal/identity/didweb.go @@ -0,0 +1,35 @@ +package identity + +import ( + "net/http" + "strings" +) + +// DIDWebHandler serves GET /.well-known/did.json — the DID document for the +// bridge's own did:web service identity. The bridge identifies as +// did:web: when no BRIDGE_SERVICE_DID is provisioned, and +// that DID appears as `hostedBy` in every bridged community.profile record; +// consumers verifying the claim (the Coves AppView's bidirectional check +// above all) resolve it right here and require `id` to equal the DID and +// `alsoKnownAs` to contain at://. Without this document the +// service DID is unresolvable and every hostedBy verification fails closed. +// +// Only the did:web method needs the endpoint: when the operator provisions +// a BRIDGE_SERVICE_DID of any other method (did:plc), the document lives in +// that method's directory instead and this handler answers 404. +func DIDWebHandler(serviceDID, hostname string) http.HandlerFunc { + hostname = strings.ToLower(hostname) + serves := serviceDID == "did:web:"+hostname + + return func(w http.ResponseWriter, r *http.Request) { + if !serves { + http.Error(w, "no did:web identity on this host", http.StatusNotFound) + return + } + writeJSON(w, http.StatusOK, map[string]any{ + "@context": []string{"https://www.w3.org/ns/did/v1"}, + "id": serviceDID, + "alsoKnownAs": []string{"at://" + hostname}, + }) + } +} diff --git a/internal/identity/didweb_test.go b/internal/identity/didweb_test.go new file mode 100644 index 0000000..38b07ea --- /dev/null +++ b/internal/identity/didweb_test.go @@ -0,0 +1,51 @@ +package identity + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDIDWebHandler(t *testing.T) { + get := func(h http.HandlerFunc) *httptest.ResponseRecorder { + rec := httptest.NewRecorder() + h(rec, httptest.NewRequest(http.MethodGet, "/.well-known/did.json", nil)) + return rec + } + + t.Run("serves the document for the derived did:web identity", func(t *testing.T) { + rec := get(DIDWebHandler("did:web:tidepool.example", "tidepool.example")) + require.Equal(t, http.StatusOK, rec.Code) + var doc struct { + Context []string `json:"@context"` + ID string `json:"id"` + AlsoKnownAs []string `json:"alsoKnownAs"` + } + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &doc)) + // The exact fields the Coves AppView's bidirectional hostedBy + // verification requires: id equal to the DID, and alsoKnownAs + // containing at://. + assert.Equal(t, "did:web:tidepool.example", doc.ID) + assert.Contains(t, doc.AlsoKnownAs, "at://tidepool.example") + assert.Contains(t, doc.Context, "https://www.w3.org/ns/did/v1") + }) + + t.Run("hostname is case-normalized", func(t *testing.T) { + rec := get(DIDWebHandler("did:web:tidepool.example", "Tidepool.Example")) + assert.Equal(t, http.StatusOK, rec.Code) + }) + + t.Run("404 when a non-did:web service DID is provisioned", func(t *testing.T) { + rec := get(DIDWebHandler("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "tidepool.example")) + assert.Equal(t, http.StatusNotFound, rec.Code) + }) + + t.Run("404 when the did:web is for a different host", func(t *testing.T) { + rec := get(DIDWebHandler("did:web:other.example", "tidepool.example")) + assert.Equal(t, http.StatusNotFound, rec.Code) + }) +} diff --git a/internal/ingest/follow.go b/internal/ingest/follow.go index 4bca716..86355b9 100644 --- a/internal/ingest/follow.go +++ b/internal/ingest/follow.go @@ -43,7 +43,10 @@ type AdminOptions struct { Service *ap.ServiceActor // Backfill serves the on-demand backfill endpoint (optional). Backfill Backfiller - Logger *slog.Logger + // Repos serves POST /admin/reemit (optional; the endpoint answers 501 + // when nil). See reemit.go for what re-emission is for. + Repos RepoReemitter + Logger *slog.Logger } // Admin is the operator API driving the community subscription lifecycle: @@ -62,6 +65,7 @@ type Admin struct { communities store.Communities service *ap.ServiceActor backfill Backfiller + repos RepoReemitter logger *slog.Logger // reconciler serves POST /admin/communities/reconcile; nil (the // endpoint answers 501) unless a follow list is configured. Set once @@ -102,6 +106,7 @@ func NewAdmin(opts AdminOptions) (*Admin, error) { communities: opts.Communities, service: opts.Service, backfill: opts.Backfill, + repos: opts.Repos, logger: logger, }, nil } @@ -120,6 +125,7 @@ func (a *Admin) Routes(r chi.Router) { r.Get("/communities", a.handleList) r.Post("/communities/backfill", a.handleBackfill) r.Post("/communities/reconcile", a.handleReconcile) + r.Post("/reemit", a.handleReemit) r.Method(http.MethodGet, "/metrics", http.HandlerFunc(scopedMetrics)) }) } diff --git a/internal/ingest/reemit.go b/internal/ingest/reemit.go new file mode 100644 index 0000000..8105a16 --- /dev/null +++ b/internal/ingest/reemit.go @@ -0,0 +1,153 @@ +package ingest + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "net/http" + "sort" + + "tidepool/internal/errors" + "tidepool/internal/repo" +) + +// RepoReemitter is the slice of the repo manager the admin re-emit endpoint +// drives. An interface so the handler is testable without postgres. +type RepoReemitter interface { + ListDIDs(ctx context.Context) ([]string, error) + ListRecords(ctx context.Context, did string) ([]repo.RecordEntry, error) + DeleteRecord(ctx context.Context, did, collection, rkey string) (*repo.CommitResult, error) + PutRecord(ctx context.Context, did, collection, rkey string, record map[string]any) (*repo.CommitResult, error) +} + +// reemitCollectionRank orders a repo's records so downstream indexers see +// identity-bearing records before the content that references them: an +// AppView consuming the re-emitted stream needs a community's profile +// indexed before its posts arrive (and an author's profile before their +// comments), or the content lands with dangling references. Cross-repo +// order is relay-dependent regardless (see README), so this is best-effort +// within one repo, not a delivery guarantee. +func reemitCollectionRank(collection string) int { + switch collection { + case "social.coves.community.profile", "social.coves.actor.profile": + return 0 + case "social.coves.community.post": + return 1 + default: + return 2 + } +} + +// ReemitResult reports one repo's re-emission. +type ReemitResult struct { + DID string `json:"did"` + Records int `json:"records"` + Reemited int `json:"reemitted"` + Error string `json:"error,omitempty"` +} + +// reemitRepo re-emits every record of one DID onto the firehose as a +// delete commit followed by a create commit with the identical value. +// +// Why delete+create and not a "touch" update: re-putting an identical +// record is an idempotent no-op by design (no commit, no event), and a +// forced update op whose value is unchanged would produce a commit whose +// op list does not match its (empty) MST diff — which sync-v1.1-validating +// relays are entitled to reject. The delete/create pair makes two honest +// commits with real diffs. Deterministic rkeys and identical bytes mean +// the record returns under the same at-uri with the same CID, so existing +// strongRefs stay valid; consumers see a transient delete, which they +// already tolerate (out-of-order tombstones are part of the AP diet). +// +// This exists for the relay cold-start gap (FOLLOWUPS): records committed +// before a relay's first subscription never re-emit on their own, so a +// Jetstream-fed AppView can never index them. +func reemitRepo(ctx context.Context, repos RepoReemitter, did string, logger *slog.Logger) ReemitResult { + res := ReemitResult{DID: did} + entries, err := repos.ListRecords(ctx, did) + if err != nil { + res.Error = err.Error() + return res + } + res.Records = len(entries) + sort.SliceStable(entries, func(i, j int) bool { + return reemitCollectionRank(entries[i].Collection) < reemitCollectionRank(entries[j].Collection) + }) + for _, e := range entries { + if ctx.Err() != nil { + res.Error = ctx.Err().Error() + return res + } + if _, err := repos.DeleteRecord(ctx, did, e.Collection, e.Rkey); err != nil { + // A vanished record (raced by a concurrent scrub) is fine to + // skip; anything else aborts this repo so the operator sees it. + if errors.IsNotFound(err) { + continue + } + res.Error = fmt.Sprintf("delete %s/%s: %v", e.Collection, e.Rkey, err) + return res + } + if _, err := repos.PutRecord(ctx, did, e.Collection, e.Rkey, e.Value); err != nil { + // The record is now deleted but not recreated — surface loudly; + // re-running the endpoint cannot restore it (the value is gone + // from the tree), so log the full value for manual recovery. + logger.Error("reemit: recreate failed after delete — record dropped from repo", + "did", did, "collection", e.Collection, "rkey", e.Rkey, "error", err) + res.Error = fmt.Sprintf("recreate %s/%s (RECORD DROPPED, see logs): %v", e.Collection, e.Rkey, err) + return res + } + res.Reemited++ + } + return res +} + +// handleReemit serves POST /admin/reemit: {"did":"did:plc:..."} re-emits one +// repo, an empty body (or {}) re-emits every active repo. Answers 200 with +// per-repo results; a body-level "failed" count > 0 means at least one repo +// aborted early and the response details why. +func (a *Admin) handleReemit(w http.ResponseWriter, r *http.Request) { + if a.repos == nil { + http.Error(w, "re-emit not configured", http.StatusNotImplemented) + return + } + var req struct { + DID string `json:"did"` + } + if r.ContentLength > 0 { + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + http.Error(w, `body must be {} or {"did":"did:..."}`, http.StatusBadRequest) + return + } + } + + ctx := r.Context() + var dids []string + if req.DID != "" { + dids = []string{req.DID} + } else { + all, err := a.repos.ListDIDs(ctx) + if err != nil { + a.logger.Error("reemit: list dids", "error", err) + http.Error(w, "listing repos failed", http.StatusInternalServerError) + return + } + dids = all + } + + results := make([]ReemitResult, 0, len(dids)) + failed := 0 + for _, did := range dids { + res := reemitRepo(ctx, a.repos, did, a.logger) + if res.Error != "" { + failed++ + } + results = append(results, res) + } + a.logger.Info("reemit complete", "repos", len(results), "failed", failed) + writeJSON(w, http.StatusOK, map[string]any{ + "repos": len(results), + "failed": failed, + "result": results, + }) +} diff --git a/internal/ingest/reemit_test.go b/internal/ingest/reemit_test.go new file mode 100644 index 0000000..1cf894b --- /dev/null +++ b/internal/ingest/reemit_test.go @@ -0,0 +1,112 @@ +package ingest + +import ( + "context" + "fmt" + "log/slog" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/errors" + "tidepool/internal/repo" +) + +// fakeReemitter records the delete/put sequence the re-emit core drives. +type fakeReemitter struct { + dids []string + records map[string][]repo.RecordEntry + // calls logs "delete did coll/rkey" / "put did coll/rkey" in order. + calls []string + // failPut, when matching "coll/rkey", fails that PutRecord. + failPut string + // missingDelete, when matching "coll/rkey", makes DeleteRecord NotFound. + missingDelete string +} + +func (f *fakeReemitter) ListDIDs(context.Context) ([]string, error) { return f.dids, nil } + +func (f *fakeReemitter) ListRecords(_ context.Context, did string) ([]repo.RecordEntry, error) { + entries, ok := f.records[did] + if !ok { + return nil, errors.NewNotFoundError("repo", did) + } + return entries, nil +} + +func (f *fakeReemitter) DeleteRecord(_ context.Context, did, collection, rkey string) (*repo.CommitResult, error) { + if collection+"/"+rkey == f.missingDelete { + return nil, errors.NewNotFoundError("record", rkey) + } + f.calls = append(f.calls, fmt.Sprintf("delete %s %s/%s", did, collection, rkey)) + return &repo.CommitResult{}, nil +} + +func (f *fakeReemitter) PutRecord(_ context.Context, did, collection, rkey string, _ map[string]any) (*repo.CommitResult, error) { + if collection+"/"+rkey == f.failPut { + return nil, fmt.Errorf("boom") + } + f.calls = append(f.calls, fmt.Sprintf("put %s %s/%s", did, collection, rkey)) + return &repo.CommitResult{}, nil +} + +func entry(collection, rkey string) repo.RecordEntry { + return repo.RecordEntry{Collection: collection, Rkey: rkey, Value: map[string]any{"$type": collection}} +} + +func TestReemitRepo(t *testing.T) { + logger := slog.Default() + + t.Run("emits profiles before posts before comments, delete then put each", func(t *testing.T) { + f := &fakeReemitter{records: map[string][]repo.RecordEntry{ + "did:plc:c": { + entry("social.coves.community.comment", "c1"), + entry("social.coves.community.post", "p1"), + entry("social.coves.community.profile", "self"), + }, + }} + res := reemitRepo(t.Context(), f, "did:plc:c", logger) + require.Empty(t, res.Error) + assert.Equal(t, 3, res.Records) + assert.Equal(t, 3, res.Reemited) + assert.Equal(t, []string{ + "delete did:plc:c social.coves.community.profile/self", + "put did:plc:c social.coves.community.profile/self", + "delete did:plc:c social.coves.community.post/p1", + "put did:plc:c social.coves.community.post/p1", + "delete did:plc:c social.coves.community.comment/c1", + "put did:plc:c social.coves.community.comment/c1", + }, f.calls) + }) + + t.Run("a record vanished mid-run is skipped, not fatal", func(t *testing.T) { + f := &fakeReemitter{ + records: map[string][]repo.RecordEntry{ + "did:plc:c": {entry("social.coves.community.post", "p1"), entry("social.coves.community.post", "p2")}, + }, + missingDelete: "social.coves.community.post/p1", + } + res := reemitRepo(t.Context(), f, "did:plc:c", logger) + require.Empty(t, res.Error) + assert.Equal(t, 1, res.Reemited, "the vanished record is skipped, the other re-emits") + }) + + t.Run("recreate failure aborts loudly", func(t *testing.T) { + f := &fakeReemitter{ + records: map[string][]repo.RecordEntry{ + "did:plc:c": {entry("social.coves.community.post", "p1"), entry("social.coves.community.post", "p2")}, + }, + failPut: "social.coves.community.post/p1", + } + res := reemitRepo(t.Context(), f, "did:plc:c", logger) + require.Contains(t, res.Error, "RECORD DROPPED") + assert.Equal(t, 0, res.Reemited) + }) + + t.Run("unknown repo reports the error", func(t *testing.T) { + f := &fakeReemitter{records: map[string][]repo.RecordEntry{}} + res := reemitRepo(t.Context(), f, "did:plc:nope", logger) + assert.NotEmpty(t, res.Error) + }) +} diff --git a/internal/repo/list.go b/internal/repo/list.go new file mode 100644 index 0000000..7752b81 --- /dev/null +++ b/internal/repo/list.go @@ -0,0 +1,102 @@ +package repo + +import ( + "context" + "database/sql" + "fmt" + "strings" + + "github.com/bluesky-social/indigo/atproto/atdata" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/ipfs/go-cid" +) + +// RecordEntry is one record surfaced by ListRecords: its path plus the +// decoded value, ready to re-put. +type RecordEntry struct { + Collection string + Rkey string + CID string + Value map[string]any +} + +// ListDIDs returns every DID with a repo (a repo_state row), excluding +// tombstoned actors: a consent_state of 'deleted' freezes the bridged +// identity (task-01 semantics), so bulk operations like the admin re-emit +// must never touch those repos. Actors under 'nobridge' suppression stay +// listed — their records are already scrubbed, so walking them is a no-op. +func (m *Manager) ListDIDs(ctx context.Context) ([]string, error) { + rows, err := m.db.QueryContext(ctx, ` + SELECT rs.did + FROM repo_state rs + LEFT JOIN bridged_actors ba ON ba.did = rs.did + WHERE ba.consent_state IS DISTINCT FROM 'deleted' + ORDER BY rs.did`) + if err != nil { + return nil, fmt.Errorf("repo: list dids: %w", err) + } + defer func() { _ = rows.Close() }() + var dids []string + for rows.Next() { + var did string + if err := rows.Scan(&did); err != nil { + return nil, fmt.Errorf("repo: scan did: %w", err) + } + dids = append(dids, did) + } + return dids, rows.Err() +} + +// ListRecords walks the DID's current MST and returns every record with its +// decoded value, in tree (path) order. Missing repo is an error satisfying +// errors.IsNotFound. The walk runs on one REPEATABLE READ snapshot for the +// same consistency reasons documented on GetRecord. +func (m *Manager) ListRecords(ctx context.Context, did string) ([]RecordEntry, error) { + if _, err := syntax.ParseDID(did); err != nil { + return nil, fmt.Errorf("repo: invalid did %q: %w", did, err) + } + + tx, err := m.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true, Isolation: sql.LevelRepeatableRead}) + if err != nil { + return nil, fmt.Errorf("repo: begin read tx: %w", err) + } + defer func() { _ = tx.Rollback() }() + + state, err := readRepoState(ctx, tx, did, false) + if err != nil { + return nil, err + } + tree, _, err := loadTree(ctx, tx, did, state.headCID) + if err != nil { + return nil, err + } + + src := &txBlockSource{tx: tx, did: did} + entries := []RecordEntry{} + walkErr := tree.Walk(func(key []byte, val cid.Cid) error { + path := string(key) + collection, rkey, ok := strings.Cut(path, "/") + if !ok { + return fmt.Errorf("repo: malformed MST key %q in %s", path, did) + } + blk, err := src.Get(ctx, val) + if err != nil { + return fmt.Errorf("repo: read record block %s: %w", val, err) + } + value, err := atdata.UnmarshalCBOR(blk.RawData()) + if err != nil { + return fmt.Errorf("repo: decode record %s: %w", val, err) + } + entries = append(entries, RecordEntry{ + Collection: collection, + Rkey: rkey, + CID: val.String(), + Value: value, + }) + return nil + }) + if walkErr != nil { + return nil, walkErr + } + return entries, nil +} diff --git a/internal/repo/list_test.go b/internal/repo/list_test.go new file mode 100644 index 0000000..9bfa872 --- /dev/null +++ b/internal/repo/list_test.go @@ -0,0 +1,103 @@ +package repo + +import ( + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/errors" + "tidepool/internal/testutil" +) + +func TestListRecords(t *testing.T) { + manager, _, _ := testManager(t) + ctx := t.Context() + + want := map[string]string{} // path -> CID + for i := 0; i < 5; i++ { + rkey := testRKey(i) + res, err := manager.PutRecord(ctx, testDID, testCollection, rkey, testRecord(fmt.Sprintf("post %d", i))) + require.NoError(t, err) + want[testCollection+"/"+rkey] = res.RecordCID + } + + entries, err := manager.ListRecords(ctx, testDID) + require.NoError(t, err) + require.Len(t, entries, len(want)) + for _, e := range entries { + path := e.Collection + "/" + e.Rkey + assert.Equal(t, want[path], e.CID, "CID for %s", path) + assert.Equal(t, testCollection, e.Collection) + require.NotNil(t, e.Value) + assert.Equal(t, testCollection, e.Value["$type"]) + assert.Contains(t, e.Value["text"], "post ") + } + + t.Run("missing repo is NotFound", func(t *testing.T) { + _, err := manager.ListRecords(ctx, "did:plc:doesnotexistanywhere1") + require.Error(t, err) + assert.True(t, errors.IsNotFound(err), "want NotFound, got %v", err) + }) + + t.Run("invalid did is an error", func(t *testing.T) { + _, err := manager.ListRecords(ctx, "not-a-did") + require.Error(t, err) + }) +} + +func TestListRecords_DeleteThenRecreateKeepsCID(t *testing.T) { + // The admin re-emit contract: delete + re-put of the identical value + // must produce two real commits and land the record back under the + // same CID (strongRefs stay valid). + manager, _, _ := testManager(t) + ctx := t.Context() + + rkey := testRKey(1) + orig, err := manager.PutRecord(ctx, testDID, testCollection, rkey, testRecord("hello")) + require.NoError(t, err) + + entries, err := manager.ListRecords(ctx, testDID) + require.NoError(t, err) + require.Len(t, entries, 1) + + del, err := manager.DeleteRecord(ctx, testDID, testCollection, rkey) + require.NoError(t, err) + require.False(t, del.NoOp) + require.Greater(t, del.Seq, orig.Seq) + + back, err := manager.PutRecord(ctx, testDID, testCollection, rkey, entries[0].Value) + require.NoError(t, err) + require.False(t, back.NoOp, "recreate must be a real commit") + require.Greater(t, back.Seq, del.Seq) + assert.Equal(t, orig.RecordCID, back.RecordCID, "identical value must round-trip to the identical CID") +} + +func TestListDIDs(t *testing.T) { + manager, database, _ := testManager(t) + testutil.Truncate(t, database, "bridged_actors") + ctx := t.Context() + + const ( + activeDID = "did:plc:activeactor11111111111a" + deletedDID = "did:plc:deletedactor1111111111a" + bareDID = "did:plc:norowactor111111111111a" // repo without a bridged_actors row + ) + for i, did := range []string{activeDID, deletedDID, bareDID} { + _, err := manager.PutRecord(ctx, did, testCollection, testRKey(i), testRecord("x")) + require.NoError(t, err) + } + _, err := database.ExecContext(ctx, ` + INSERT INTO bridged_actors (did, ap_actor_id, handle, actor_type, consent_state) + VALUES ($1, 'https://x/u/a', 'a.x.test', 'person', 'ok'), + ($2, 'https://x/u/b', 'b.x.test', 'person', 'deleted')`, + activeDID, deletedDID) + require.NoError(t, err) + + dids, err := manager.ListDIDs(ctx) + require.NoError(t, err) + assert.Contains(t, dids, activeDID) + assert.Contains(t, dids, bareDID, "repos without an actor row (e.g. legacy) must list") + assert.NotContains(t, dids, deletedDID, "tombstoned actors are frozen and must never re-emit") +}