diff --git a/cmd/rematerialize-posts/main.go b/cmd/rematerialize-posts/main.go new file mode 100644 index 0000000..aaaaa69 --- /dev/null +++ b/cmd/rematerialize-posts/main.go @@ -0,0 +1,357 @@ +// Command rematerialize-posts is the operator-invoked cutover tool of +// docs/PRD_AUTHOR_OWNED_POSTS.md §11: it moves every legacy +// social.coves.community.post (written into a COMMUNITY's repo under the +// community's credentials) to an author-owned social.coves.community.postv2 in +// the AUTHOR's repo, plus the community's acceptance that pins it, and then — +// and only then — deletes the old record. +// +// It is a THIN WRAPPER. All of the safety logic lives in posts.Rematerializer; +// this file only wires the production seams the state machine drives: +// +// - the ledger (migration 037), so the run is resumable and idempotent; +// - the author-repo factory, so each postv2 is signed by its own author (an +// aggregator's stored session for a non-interactive author, never a forged +// admin signature); +// - the DIRECT community acceptance writer, so a since-banned author's live +// post is preserved rather than re-adjudicated; +// - a LegacySource over the real community repos: listRecords to discover the +// deprecated posts, deleteRecord to remove them once verified. +// +// It is run by hand during the deploy window (§11 step 4), against a database +// whose migrations are already applied, and it reports a census that refuses to +// declare the migration complete while any post was left as legacy — the gate on +// the separate, irreversible legacy-removal follow-up. +package main + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "flag" + "fmt" + "log" + "os" + "time" + + "Coves/internal/atproto/oauth" + "Coves/internal/atproto/pds" + "Coves/internal/config" + "Coves/internal/core/aggregators" + "Coves/internal/core/blobs" + "Coves/internal/core/communities" + "Coves/internal/core/posts" + postgresRepo "Coves/internal/db/postgres" + + _ "github.com/lib/pq" +) + +// legacyPostCollection is the deprecated community-repo post collection the tool +// drains. It is the lexicon NSID, spelled here rather than imported because the +// posts package keeps its copy private; the two must agree, and there is exactly +// one correct string. +const legacyPostCollection = "social.coves.community.post" + +// listPageSize bounds each communities/records page so an instance with a large +// catalogue is enumerated in bounded queries rather than one unbounded read. +const listPageSize = 100 + +func main() { + communityFilter := flag.String("community", "", + "restrict the run to a single community DID (a staged rollout); empty means every hosted community") + flag.Parse() + + cfg, err := config.Load() + if err != nil { + log.Fatalf("rematerialize-posts: loading config: %v", err) + } + + db, err := openDatabase(cfg) + if err != nil { + log.Fatalf("rematerialize-posts: %v", err) + } + defer func() { _ = db.Close() }() + + ctx := context.Background() + + // The OAuth client is the credential seam for BOTH kinds of author: a human's + // browser session is not available in a batch tool, so the only authors this + // run can re-materialize are non-interactive ones (aggregators) whose stored + // session it resumes. Any human-authored post surfaces as a no-creds fallback + // and is left as legacy — never forged. + oauthClient, err := buildOAuthClient(cfg, db) + if err != nil { + log.Fatalf("rematerialize-posts: %v", err) + } + + blobService := blobs.NewBlobService(cfg.PDS.URL) + provisioner := communities.NewPDSAccountProvisioner(cfg.Instance.Domain, cfg.PDS.URL) + communityService := communities.NewCommunityService( + postgresRepo.NewCommunityRepository(db), + cfg.PDS.URL, + cfg.Instance.DID, + cfg.Instance.Domain, + provisioner, + oauthClient, + blobService, + ) + + // The DIRECT acceptance writer over the production community-repo factory — + // the same credential-presence hosting test the acceptance engine uses. The + // tool holds this writer and no decider, so it cannot re-run admission. + repoFactory := posts.NewCommunityRepoFactory(communityService) + writer := posts.NewCommunityRecordWriter(repoFactory, time.Now) + + authorFactory := posts.NewAuthorRepoFactory(oauthClient.ClientApp, aggregators.DefaultSessionID) + + source := &realLegacySource{ + communities: communityService, + creds: communityService, + communityFilter: *communityFilter, + } + + tool := &posts.Rematerializer{ + Source: source, + Ledger: postgresRepo.NewRematerializeLedger(db), + AuthorRepos: authorFactory, + Acceptances: writer, + } + + report, err := tool.Run(ctx) + if err != nil { + log.Fatalf("rematerialize-posts: the run failed: %v", err) + } + + logCensus(report) + + // A surviving fallback means at least one post still lives only as a legacy + // record, so the operator must NOT proceed to the legacy-removal step. Exiting + // non-zero makes that a machine-checkable gate rather than a line of output an + // operator might skim past. + if !report.Complete { + log.Printf("rematerialize-posts: INCOMPLETE — %d post(s) left as legacy; do not run the legacy-removal step", report.Fallbacks) + os.Exit(1) + } + log.Printf("rematerialize-posts: complete — every discovered post was re-materialized") +} + +// realLegacySource enumerates the deprecated community.post records across the +// hosted communities and deletes them from their community repos. +// +// It reaches the PDS through a full pds.Client rather than the narrowed +// CommunityRepo the acceptance writer uses, because discovery needs listRecords +// and deletion needs deleteRecord — neither of which the write-narrowed surface +// carries. The credentials come from the same source the acceptance writer's +// factory uses, so the two never disagree about which communities are hosted. +type realLegacySource struct { + communities communities.Service + creds posts.CommunityCredentialSource + communityFilter string +} + +// ListLegacyPosts walks every hosted community and lists its remaining +// social.coves.community.post records as LegacyPosts. +func (s *realLegacySource) ListLegacyPosts(ctx context.Context) ([]posts.LegacyPost, error) { + dids, err := s.hostedCommunityDIDs(ctx) + if err != nil { + return nil, err + } + + var legacy []posts.LegacyPost + for _, did := range dids { + client, err := s.communityClient(ctx, did) + if err != nil { + return nil, fmt.Errorf("opening the repo of %s to list legacy posts: %w", did, err) + } + + cursor := "" + for { + page, err := client.ListRecords(ctx, legacyPostCollection, listPageSize, cursor) + if err != nil { + return nil, fmt.Errorf("listing %s in %s: %w", legacyPostCollection, did, err) + } + for _, entry := range page.Records { + post, err := legacyPostFromEntry(did, entry) + if err != nil { + return nil, fmt.Errorf("decoding legacy record %s: %w", entry.URI, err) + } + legacy = append(legacy, post) + } + if page.Cursor == "" { + break + } + cursor = page.Cursor + } + } + return legacy, nil +} + +// DeleteLegacyPost removes the old community.post from its community repo. A +// delete of an already-gone record is success — it is the step a crash after the +// migrated checkpoint retries, so idempotence is the contract. +func (s *realLegacySource) DeleteLegacyPost(ctx context.Context, legacy posts.LegacyPost) error { + client, err := s.communityClient(ctx, legacy.CommunityDID) + if err != nil { + return fmt.Errorf("opening the repo of %s to delete %s: %w", legacy.CommunityDID, legacy.URI, err) + } + rkey := legacy.URI[lastSlash(legacy.URI)+1:] + if err := client.DeleteRecord(ctx, legacyPostCollection, rkey); err != nil { + if errors.Is(err, pds.ErrNotFound) { + return nil + } + return fmt.Errorf("deleting %s: %w", legacy.URI, err) + } + return nil +} + +// hostedCommunityDIDs returns the DIDs of the communities this AppView can sign +// for — the only ones whose posts it can re-materialize and whose old records it +// can delete. A --community filter narrows the run to one for a staged rollout. +func (s *realLegacySource) hostedCommunityDIDs(ctx context.Context) ([]string, error) { + if s.communityFilter != "" { + return []string{s.communityFilter}, nil + } + + var dids []string + offset := 0 + for { + page, err := s.communities.ListCommunities(ctx, communities.ListCommunitiesRequest{ + Limit: listPageSize, + Offset: offset, + }) + if err != nil { + return nil, fmt.Errorf("listing communities: %w", err) + } + for _, community := range page { + // Hosting is credential presence, never a claimed profile field: only a + // community whose refresh token this AppView holds can be written to, and + // the repo factory would refuse the rest with ErrCommunityNotHosted. + if community.PDSRefreshToken != "" { + dids = append(dids, community.DID) + } + } + if len(page) < listPageSize { + break + } + offset += listPageSize + } + return dids, nil +} + +// communityClient opens a full PDS client bound to one community's repo, over +// freshly-renewed stored credentials. +func (s *realLegacySource) communityClient(ctx context.Context, did string) (pds.Client, error) { + community, err := s.creds.GetByDID(ctx, did) + if err != nil { + return nil, fmt.Errorf("reading the credentials of %s: %w", did, err) + } + if community == nil { + return nil, fmt.Errorf("reading the credentials of %s: no such community is indexed", did) + } + fresh, err := s.creds.EnsureFreshToken(ctx, community) + if err != nil { + return nil, fmt.Errorf("renewing the credentials of %s: %w", did, err) + } + if fresh == nil || fresh.PDSAccessToken == "" { + return nil, fmt.Errorf("renewing the credentials of %s: no access token came back", did) + } + return pds.NewFromAccessToken(fresh.PDSURL, fresh.DID, fresh.PDSAccessToken) +} + +// legacyPostFromEntry decodes one listRecords entry into a LegacyPost, carrying +// the decoded body forward so the tool re-materializes the record's ACTUAL +// content rather than a re-fetch that might have changed under it. +func legacyPostFromEntry(communityDID string, entry pds.RecordEntry) (posts.LegacyPost, error) { + raw, err := json.Marshal(entry.Value) + if err != nil { + return posts.LegacyPost{}, fmt.Errorf("re-encoding record value: %w", err) + } + var record posts.PostRecord + if err := json.Unmarshal(raw, &record); err != nil { + return posts.LegacyPost{}, fmt.Errorf("decoding record value: %w", err) + } + if record.Author == "" { + return posts.LegacyPost{}, fmt.Errorf("record %s carries no author field to re-author under", entry.URI) + } + return posts.LegacyPost{ + URI: entry.URI, + CID: entry.CID, + CommunityDID: communityDID, + AuthorDID: record.Author, + Record: record, + }, nil +} + +// logCensus prints the run's per-state tally. +func logCensus(report posts.RematerializeReport) { + log.Printf("rematerialize-posts: census — discovered=%d done=%d fallbacks=%d complete=%v", + report.Discovered, report.Done, report.Fallbacks, report.Complete) + for state, n := range report.ByState { + log.Printf(" %-22s %d", state, n) + } +} + +func lastSlash(s string) int { + for i := len(s) - 1; i >= 0; i-- { + if s[i] == '/' { + return i + } + } + return -1 +} + +// openDatabase opens the AppView Postgres the ledger and community catalogue +// live in, with the pool bounds from config. +func openDatabase(cfg *config.Config) (*sql.DB, error) { + db, err := sql.Open("postgres", cfg.Database.URL) + if err != nil { + return nil, fmt.Errorf("opening the database: %w", err) + } + db.SetMaxOpenConns(cfg.Database.MaxOpenConns) + db.SetMaxIdleConns(cfg.Database.MaxIdleConns) + if err := db.Ping(); err != nil { + _ = db.Close() + return nil, fmt.Errorf("pinging the database: %w", err) + } + return db, nil +} + +// buildOAuthClient assembles the OAuth client the author-repo factory resumes +// aggregator sessions through — the same wiring cmd/server uses, so a session +// this tool resumes is byte-identical to one the write path would. +func buildOAuthClient(cfg *config.Config, db *sql.DB) (*oauth.OAuthClient, error) { + store := oauth.NewMobileAwareStoreWrapper(oauth.NewPostgresOAuthStore(db, 0)) + client, err := oauth.NewOAuthClient(&oauth.OAuthConfig{ + PublicURL: cfg.OAuth.PublicURL, + SealSecret: cfg.OAuth.SealSecret, + Scopes: oauthScopes(), + DevMode: cfg.IsDevEnv, + AllowPrivateIPs: cfg.IsDevEnv, + PLCURL: cfg.Identity.PLCURL, + PDSURL: cfg.PDS.URL, + ClientPrivateKeyMultibase: cfg.OAuth.ClientPrivateKeyMultibase, + ClientKeyID: cfg.OAuth.ClientKeyID, + }, store) + if err != nil { + return nil, fmt.Errorf("initializing the OAuth client: %w", err) + } + return client, nil +} + +// oauthScopes is the granted scope set an aggregator's resumed session must +// carry to write a postv2. It mirrors cmd/server's list, which is the authority; +// the two must agree, so a divergence here shows up as a scope the resumed +// session lacks at the first write. +func oauthScopes() []string { + return []string{ + "atproto", + "blob:*/*", + "repo:social.coves.community.post?action=create&action=update&action=delete", + "repo:social.coves.community.comment?action=create&action=update&action=delete", + "repo:social.coves.community.profile?action=create&action=update&action=delete", + "repo:social.coves.community.subscription?action=create&action=update&action=delete", + "repo:social.coves.actor.profile?action=create&action=update&action=delete", + "repo:social.coves.feed.vote?action=create&action=delete", + "repo:social.coves.actor.block?action=create&action=delete", + } +} diff --git a/internal/core/posts/rematerialize.go b/internal/core/posts/rematerialize.go index f459854..53ace98 100644 --- a/internal/core/posts/rematerialize.go +++ b/internal/core/posts/rematerialize.go @@ -3,6 +3,7 @@ package posts import ( "context" "errors" + "fmt" "time" ) @@ -212,9 +213,6 @@ type Rematerializer struct { Acceptances CommunityRecordWriter } -// errRematerializeNotImplemented is the RED sentinel every stub body returns. -var errRematerializeNotImplemented = errors.New("rematerialize: not implemented (RED stub)") - // RematerializeRkey is the postv2 record key the tool writes a legacy record at. // // IT IS A PURE, STABLE FUNCTION OF THE OLD RECORD'S URI — the single @@ -225,17 +223,173 @@ var errRematerializeNotImplemented = errors.New("rematerialize: not implemented // re-run would draw a different key and duplicate the post. // // The scheme is the SubjectRkey digest scheme applied to the OLD URI: a total, -// collision-free (SHA-256) function that the write path already trusts. -func RematerializeRkey(legacyPostURI string) string { return "" } +// collision-free (SHA-256) function that the write path already trusts. It is +// SubjectRkey verbatim, not a re-derivation, so the tool and the write path +// cannot come to disagree about what key a subject hashes to. +func RematerializeRkey(legacyPostURI string) string { return SubjectRkey(legacyPostURI) } // Run discovers every legacy record and drives each to a terminal state, // returning the census. +// +// A per-record error FAILS THE RUN rather than being logged and skipped: the +// safety properties are all ordering ones, and continuing past a record the tool +// could not verify would let the operator read a "done"-heavy census as +// permission to run the irreversible legacy-removal step while a record sits +// half-migrated. A no-creds fallback is NOT such an error — it is an expected +// terminal outcome the census counts — so it lets the run continue while still +// holding Complete false. func (r *Rematerializer) Run(ctx context.Context) (RematerializeReport, error) { - return RematerializeReport{}, errRematerializeNotImplemented + legacies, err := r.Source.ListLegacyPosts(ctx) + if err != nil { + return RematerializeReport{}, fmt.Errorf("enumerating legacy posts: %w", err) + } + + for _, legacy := range legacies { + if _, err := r.RematerializeOne(ctx, legacy); err != nil { + return RematerializeReport{}, fmt.Errorf("re-materializing %s: %w", legacy.URI, err) + } + } + + byState, err := r.Ledger.CountByState(ctx) + if err != nil { + return RematerializeReport{}, fmt.Errorf("taking the census: %w", err) + } + + report := RematerializeReport{ByState: byState} + for state, n := range byState { + report.Discovered += n + if state == RematerializeDone { + report.Done += n + } + if IsFallback(state) { + report.Fallbacks += n + } + } + // The gate on the separate, irreversible legacy-removal follow-up (§11 step + // 6): a surviving fallback is a post still living only as a legacy record, so + // the run must not tell the operator the migration is finished. + report.Complete = report.Fallbacks == 0 + + return report, nil } // RematerializeOne drives a single legacy record from wherever its ledger row // stands to a terminal state, and returns the state it reached. +// +// The steps are guarded on the ledger state each moves FROM, so a resumed run +// re-enters at exactly the step its predecessor stopped before and re-does none +// of the completed ones. The load-bearing ordering is VERIFY BEFORE DELETE: the +// old record is deleted only after the postv2 and its acceptance are confirmed to +// pin the same CID, and the migrated checkpoint is persisted BEFORE the delete so +// a crash there retries only the delete. func (r *Rematerializer) RematerializeOne(ctx context.Context, legacy LegacyPost) (RematerializeState, error) { - return "", errRematerializeNotImplemented + row, err := r.Ledger.Discover(ctx, legacy.URI, legacy.AuthorDID) + if err != nil { + return "", err + } + + // A row already in a terminal fallback state is left exactly as it stands: the + // credential census reached its verdict on a prior pass, and re-opening it + // would be the one thing §11 step 3 forbids — a second chance to forge. + if IsFallback(row.State) { + return row.State, nil + } + + // Step 1 — postv2_written. Write the author-owned postv2 at the deterministic + // rkey. createAuthorRecord is create-only and converges by read, so a resume + // that re-enters here (it will not, but the guard is honest) would find its own + // first attempt rather than mint a second post. + if row.State == RematerializeDiscovered { + repo, err := r.AuthorRepos(ctx, legacy.AuthorDID, nil) + if err != nil { + // NO CREDENTIALS IS A TERMINAL FALLBACK, NEVER A FORGERY. An author whose + // repo cannot be restored is left as legacy — the postv2 is not written and + // the old record survives — because re-authoring under any other identity + // reintroduces the §2 impersonation the whole flip removes. + if errors.Is(err, ErrNoAuthorCredentials) { + reason := fmt.Sprintf("author %s has no restorable repo credentials: %v", legacy.AuthorDID, err) + if markErr := r.Ledger.MarkFallback(ctx, legacy.URI, RematerializeFallbackLeftLegacy, reason); markErr != nil { + return row.State, markErr + } + return RematerializeFallbackLeftLegacy, nil + } + return row.State, fmt.Errorf("opening the author repo of %s: %w", legacy.AuthorDID, err) + } + + rkey := RematerializeRkey(legacy.URI) + newURI, newCID, _, err := createAuthorRecord(ctx, repo, rkey, postV2From(legacy.Record)) + if err != nil { + return row.State, fmt.Errorf("writing the postv2 for %s: %w", legacy.URI, err) + } + + if err := r.Ledger.RecordPostV2Written(ctx, legacy.URI, newURI, newCID, rkey); err != nil { + return row.State, err + } + row.State = RematerializePostV2Written + row.NewURI, row.NewCID, row.NewRkey = newURI, newCID, rkey + } + + // Step 2 — verified. Write the community's acceptance DIRECT (never through the + // engine — see the type's doc) pinning the NEW postv2 CID, then RE-READ the + // postv2 and confirm it still pins that CID. Verification reads the standing + // record rather than trusting the write's returned CID, so a concurrent edit + // landing in the write→verify window is caught here — before anything is + // deleted. + if row.State == RematerializePostV2Written { + if _, err := r.Acceptances.WriteAcceptance(ctx, CommunityWriteCommand{ + CommunityDID: legacy.CommunityDID, + PostURI: row.NewURI, + PostCID: row.NewCID, + }); err != nil { + return row.State, fmt.Errorf("writing the acceptance for %s: %w", row.NewURI, err) + } + + repo, err := r.AuthorRepos(ctx, legacy.AuthorDID, nil) + if err != nil { + return row.State, fmt.Errorf("re-opening the author repo of %s to verify: %w", legacy.AuthorDID, err) + } + standing, err := repo.GetRecord(ctx, PostV2Collection, row.NewRkey) + if err != nil { + return row.State, fmt.Errorf("verifying the postv2 for %s: %w", row.NewURI, err) + } + if standing.CID != row.NewCID { + // VERIFY BEFORE DELETE fails closed: the acceptance now pins a CID the + // postv2 no longer carries, so deleting the old record would destroy the + // only copy of a post whose new attestation points at content that no + // longer stands. No checkpoint, no delete — the row stays at + // postv2_written for a later pass to re-verify. + return row.State, fmt.Errorf( + "verifying the postv2 for %s: the standing record pins %s but the acceptance pinned %s (a concurrent edit landed mid-verify)", + row.NewURI, standing.CID, row.NewCID) + } + + if err := r.Ledger.MarkVerified(ctx, legacy.URI); err != nil { + return row.State, err + } + row.State = RematerializeVerified + } + + // Step 3 — migrated. The checkpoint BEFORE the delete: postv2 and acceptance + // verified, old record still present. Persisting it as its own state is what + // lets a crash on the delete retry ONLY the delete. + if row.State == RematerializeVerified { + if err := r.Ledger.MarkMigrated(ctx, legacy.URI); err != nil { + return row.State, err + } + row.State = RematerializeMigrated + } + + // Step 4 — done. Delete the old community.post; a delete of an already-gone + // record is success (the source's contract), so a resumed delete is idempotent. + if row.State == RematerializeMigrated { + if err := r.Source.DeleteLegacyPost(ctx, legacy); err != nil { + return row.State, fmt.Errorf("deleting the old record %s: %w", legacy.URI, err) + } + if err := r.Ledger.MarkDone(ctx, legacy.URI); err != nil { + return row.State, err + } + row.State = RematerializeDone + } + + return row.State, nil } diff --git a/internal/db/migrations/037_create_post_rematerialization_ledger.sql b/internal/db/migrations/037_create_post_rematerialization_ledger.sql new file mode 100644 index 0000000..4e80b79 --- /dev/null +++ b/internal/db/migrations/037_create_post_rematerialization_ledger.sql @@ -0,0 +1,73 @@ +-- +goose Up +-- The re-materialization ledger: one row per legacy social.coves.community.post +-- record, tracking how far the cutover tool moved it (docs/PRD_AUTHOR_OWNED_POSTS.md +-- §11 the rev-2.8 deploy runbook). +-- +-- WHY THIS EXISTS. The tool moves every legacy community-repo post to an +-- author-owned postv2 plus a community acceptance, and then — and ONLY then — +-- deletes the old record. Getting that order wrong makes a live post vanish or +-- relaunders a removed one, so the tool has to be resumable and idempotent, and +-- this table is what makes resume possible: a crash reads the row back and picks +-- up exactly where it stopped rather than re-doing work or skipping the delete. +-- +-- THE DISTINCTION THIS TABLE HAS TO MAKE REPRESENTABLE is migrated ≠ done. +-- `migrated` means "verified safe to delete, old record STILL PRESENT"; `done` +-- means "old record deleted". A crash between them resumes by retrying ONLY the +-- delete, and a delete of an already-gone record is success — which is coherent +-- only if the checkpoint BEFORE the delete is its own persisted state. Collapsing +-- the two would either re-write the postv2 on every resume or skip the delete a +-- crash owes. +-- +-- THE TWO FALLBACK STATES are the credential census (§11 step 3): a record whose +-- author credentials cannot be restored is left as legacy — never re-authored +-- under a forged signature, which would reintroduce the §2 impersonation the flip +-- exists to remove — and the run refuses to report "complete" while any such row +-- survives, gating the operator's separate, irreversible legacy-removal step. +CREATE TABLE post_rematerialization_ledger ( + -- The OLD community.post AT-URI is the whole key: one row per legacy record, + -- so a re-run UPDATES the row it resumes from rather than accumulating a + -- second for the same post, and the deterministic postv2 rkey is derived from + -- this exact string. + old_uri TEXT PRIMARY KEY, + + -- The state machine's cursor. NOT NULL — a row with no state cannot be + -- resumed from — and CHECK-closed so a typo'd transition is refused at the + -- schema, where every writer meets it, rather than sitting in the table as an + -- unresumable row nothing would ever advance. migrated and done are BOTH + -- listed and DISTINCT on purpose (see the header). + state TEXT NOT NULL CHECK (state IN ( + 'discovered', + 'postv2_written', + 'verified', + 'migrated', + 'done', + 'fallback_left_legacy', + 'fallback_no_creds' + )), + + -- Who the postv2 is re-authored under: the legacy record's `author` field. + -- Nullable is acceptable — a fallback row may never resolve a repo to write + -- into — but it is the audit trail for which repo the tool wrote to. + author_did TEXT, + + -- The postv2 coordinates, populated at the postv2_written transition and read + -- back (never recomputed) on resume, so a resumed run converges on the record + -- its first attempt wrote rather than deriving a fresh CID. NULL until then. + new_uri TEXT, + new_cid TEXT, + new_rkey TEXT, + + -- The human-readable note on a fallback row (why it was left as legacy). + reason TEXT, + + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); + +COMMENT ON TABLE post_rematerialization_ledger IS 'One row per legacy community.post: the cutover tool''s resumable, idempotent progress ledger (PRD_AUTHOR_OWNED_POSTS 11)'; +COMMENT ON COLUMN post_rematerialization_ledger.old_uri IS 'The OLD community.post AT-URI; the primary key and the material the deterministic postv2 rkey is derived from'; +COMMENT ON COLUMN post_rematerialization_ledger.state IS 'The state-machine cursor; CHECK-closed. migrated (safe to delete, old record present) is DISTINCT from done (old record deleted)'; +COMMENT ON COLUMN post_rematerialization_ledger.new_uri IS 'The postv2 URI written at the postv2_written transition; read back on resume, never recomputed'; + +-- +goose Down +DROP TABLE IF EXISTS post_rematerialization_ledger; diff --git a/internal/db/postgres/rematerialize_ledger.go b/internal/db/postgres/rematerialize_ledger.go index 0f771e4..17bf2ed 100644 --- a/internal/db/postgres/rematerialize_ledger.go +++ b/internal/db/postgres/rematerialize_ledger.go @@ -3,7 +3,7 @@ package postgres import ( "context" "database/sql" - "errors" + "fmt" "Coves/internal/core/posts" ) @@ -12,12 +12,12 @@ import ( // per legacy social.coves.community.post record, tracking its progress through // the cutover state machine (docs/PRD_AUTHOR_OWNED_POSTS.md §11). // -// # THIS FILE IS A RED STUB -// -// It exists so the state-machine tests compile and so cmd/rematerialize-posts -// has a concrete ledger to wire. Every method returns the not-implemented -// sentinel; GREEN implements the single-statement upserts and the census query -// against the migration-037 table. +// Every transition is a single parameterized statement. The state-advancing +// updates are GUARDED on the state they move FROM, and a guard that matches no +// row is an error rather than a silent no-op: the tool advances the machine in +// strict order, so a transition finding the row in an unexpected state means the +// ledger and the tool have diverged, and that is a fault to surface, not to +// swallow into a false success. // rematerializeLedger is the migration-037-backed posts.RematerializeLedger. type rematerializeLedger struct { @@ -29,36 +29,160 @@ func NewRematerializeLedger(db *sql.DB) posts.RematerializeLedger { return &rematerializeLedger{db: db} } -var errRematerializeLedgerNotImplemented = errors.New("rematerialize ledger: not implemented (RED stub)") - +// Discover upserts the row for oldURI in state discovered, idempotently: a +// re-run finds the existing row (whatever state it stands in) rather than +// resetting it, then reads it back so the caller resumes from where it stopped. func (l *rematerializeLedger) Discover(ctx context.Context, oldURI, authorDID string) (posts.RematerializeLedgerRow, error) { - return posts.RematerializeLedgerRow{}, errRematerializeLedgerNotImplemented + // ON CONFLICT DO NOTHING keeps a resumed row untouched; a plain INSERT would + // reset an in-flight row back to discovered and re-do the whole migration. + _, err := l.db.ExecContext(ctx, ` + INSERT INTO post_rematerialization_ledger (old_uri, state, author_did, created_at, updated_at) + VALUES ($1, $2, $3, NOW(), NOW()) + ON CONFLICT (old_uri) DO NOTHING + `, oldURI, string(posts.RematerializeDiscovered), nullString(authorDID)) + if err != nil { + return posts.RematerializeLedgerRow{}, fmt.Errorf("discovering %s: %w", oldURI, err) + } + + row, found, err := l.Get(ctx, oldURI) + if err != nil { + return posts.RematerializeLedgerRow{}, err + } + if !found { + return posts.RematerializeLedgerRow{}, fmt.Errorf("discovering %s: the row vanished between upsert and read", oldURI) + } + return row, nil } +// Get reads one row. found is false when the URI has never been discovered. func (l *rematerializeLedger) Get(ctx context.Context, oldURI string) (posts.RematerializeLedgerRow, bool, error) { - return posts.RematerializeLedgerRow{}, false, errRematerializeLedgerNotImplemented + var ( + row posts.RematerializeLedgerRow + state string + authorDID sql.NullString + newURI sql.NullString + newCID sql.NullString + newRkey sql.NullString + reason sql.NullString + ) + err := l.db.QueryRowContext(ctx, ` + SELECT old_uri, state, author_did, new_uri, new_cid, new_rkey, reason, created_at, updated_at + FROM post_rematerialization_ledger + WHERE old_uri = $1 + `, oldURI).Scan(&row.OldURI, &state, &authorDID, &newURI, &newCID, &newRkey, &reason, &row.CreatedAt, &row.UpdatedAt) + if err == sql.ErrNoRows { + return posts.RematerializeLedgerRow{}, false, nil + } + if err != nil { + return posts.RematerializeLedgerRow{}, false, fmt.Errorf("reading ledger row %s: %w", oldURI, err) + } + + row.State = posts.RematerializeState(state) + row.AuthorDID = authorDID.String + row.NewURI = newURI.String + row.NewCID = newCID.String + row.NewRkey = newRkey.String + row.Reason = reason.String + return row, true, nil } +// RecordPostV2Written moves discovered → postv2_written and records the postv2 +// coordinates the resume path reads back. func (l *rematerializeLedger) RecordPostV2Written(ctx context.Context, oldURI, newURI, newCID, newRkey string) error { - return errRematerializeLedgerNotImplemented + return l.guardedTransition(ctx, ` + UPDATE post_rematerialization_ledger + SET state = $2, new_uri = $3, new_cid = $4, new_rkey = $5, updated_at = NOW() + WHERE old_uri = $1 AND state = $6 + `, "postv2_written", oldURI, + string(posts.RematerializePostV2Written), newURI, newCID, newRkey, string(posts.RematerializeDiscovered)) } +// MarkVerified moves postv2_written → verified. func (l *rematerializeLedger) MarkVerified(ctx context.Context, oldURI string) error { - return errRematerializeLedgerNotImplemented + return l.guardedTransition(ctx, ` + UPDATE post_rematerialization_ledger + SET state = $2, updated_at = NOW() + WHERE old_uri = $1 AND state = $3 + `, "verified", oldURI, + string(posts.RematerializeVerified), string(posts.RematerializePostV2Written)) } +// MarkMigrated moves verified → migrated — the checkpoint before delete. func (l *rematerializeLedger) MarkMigrated(ctx context.Context, oldURI string) error { - return errRematerializeLedgerNotImplemented + return l.guardedTransition(ctx, ` + UPDATE post_rematerialization_ledger + SET state = $2, updated_at = NOW() + WHERE old_uri = $1 AND state = $3 + `, "migrated", oldURI, + string(posts.RematerializeMigrated), string(posts.RematerializeVerified)) } +// MarkDone moves migrated → done, after the old record is deleted. func (l *rematerializeLedger) MarkDone(ctx context.Context, oldURI string) error { - return errRematerializeLedgerNotImplemented + return l.guardedTransition(ctx, ` + UPDATE post_rematerialization_ledger + SET state = $2, updated_at = NOW() + WHERE old_uri = $1 AND state = $3 + `, "done", oldURI, + string(posts.RematerializeDone), string(posts.RematerializeMigrated)) } +// MarkFallback moves a discovered row to a terminal fallback state with a reason. +// The from-state guard is intentionally broad — a fallback is only ever reached +// from discovered in cycle 1 — but the reason is always recorded for the census. func (l *rematerializeLedger) MarkFallback(ctx context.Context, oldURI string, state posts.RematerializeState, reason string) error { - return errRematerializeLedgerNotImplemented + if !posts.IsFallback(state) { + return fmt.Errorf("marking %s as fallback: %q is not a fallback state", oldURI, state) + } + return l.guardedTransition(ctx, ` + UPDATE post_rematerialization_ledger + SET state = $2, reason = $3, updated_at = NOW() + WHERE old_uri = $1 AND state = $4 + `, string(state), oldURI, + string(state), reason, string(posts.RematerializeDiscovered)) } +// CountByState is the census: how many rows sit in each state, so the run can +// refuse "complete" while any fallback survives. func (l *rematerializeLedger) CountByState(ctx context.Context) (map[posts.RematerializeState]int, error) { - return nil, errRematerializeLedgerNotImplemented + rows, err := l.db.QueryContext(ctx, ` + SELECT state, COUNT(*) FROM post_rematerialization_ledger GROUP BY state + `) + if err != nil { + return nil, fmt.Errorf("counting ledger rows by state: %w", err) + } + defer func() { _ = rows.Close() }() + + counts := map[posts.RematerializeState]int{} + for rows.Next() { + var state string + var n int + if err := rows.Scan(&state, &n); err != nil { + return nil, fmt.Errorf("scanning census row: %w", err) + } + counts[posts.RematerializeState(state)] = n + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterating census rows: %w", err) + } + return counts, nil +} + +// guardedTransition runs a from-state-guarded UPDATE and treats a no-op as the +// error it is: the tool only ever fires a transition when the row stands in the +// expected state, so matching no row means the ledger and the tool diverged. +func (l *rematerializeLedger) guardedTransition(ctx context.Context, query, toState, oldURI string, args ...any) error { + full := append([]any{oldURI}, args...) + res, err := l.db.ExecContext(ctx, query, full...) + if err != nil { + return fmt.Errorf("transitioning %s to %s: %w", oldURI, toState, err) + } + affected, err := res.RowsAffected() + if err != nil { + return fmt.Errorf("transitioning %s to %s: reading rows affected: %w", oldURI, toState, err) + } + if affected == 0 { + return fmt.Errorf("transitioning %s to %s: no row in the expected prior state (the ledger and the tool have diverged)", oldURI, toState) + } + return nil }