diff --git a/appview/db/db.go b/appview/db/db.go index 466be3da..c5e51750 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -1451,6 +1451,30 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return err }) + orm.RunMigration(conn, logger, "unify-pds-record-migration-table", func(tx *sql.Tx) error { + _, err := tx.Exec(` + insert into pds_migration ( + name, + did, + collection, + rkey, + status, + updated_at + ) + select + 'add-repo-did', + user_did, + record_nsid, + record_rkey, + status, + updated_at + from pds_rewrite_status; + + drop table pds_rewrite_status; + `) + return err + }) + return &DB{ db, logger, diff --git a/appview/db/repos.go b/appview/db/repos.go index 965b31c5..7b0f3ce4 100644 --- a/appview/db/repos.go +++ b/appview/db/repos.go @@ -1,6 +1,7 @@ package db import ( + "context" "database/sql" "errors" "fmt" @@ -10,6 +11,7 @@ import ( "time" "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" "tangled.org/core/appview/models" "tangled.org/core/appview/pagination" "tangled.org/core/orm" @@ -589,20 +591,22 @@ func GetRepoByDid(e Execer, repoDid string) (*models.Repo, error) { return GetRepo(e, orm.FilterEq("repo_did", repoDid)) } +// TODO: just queue every legacy records regardless of target repo has a DID or not. +// doable after we have `repo_did` column in db for each tables. func EnqueuePdsRewritesForRepo(tx *sql.Tx, repoDid, repoAtUri string) error { type record struct { userDidCol string table string - nsid string + nsid syntax.NSID fkCol string } sources := []record{ - {"did", "repos", "sh.tangled.repo", "at_uri"}, - {"did", "issues", "sh.tangled.repo.issue", "repo_at"}, - {"owner_did", "pulls", "sh.tangled.repo.pull", "repo_at"}, - {"did", "collaborators", "sh.tangled.repo.collaborator", "repo_at"}, - {"did", "artifacts", "sh.tangled.repo.artifact", "repo_at"}, - {"did", "stars", "sh.tangled.feed.star", "subject_at"}, + {"did", "repos", tangled.RepoNSID, "at_uri"}, + {"did", "issues", tangled.RepoIssueNSID, "repo_at"}, + {"owner_did", "pulls", tangled.RepoPullNSID, "repo_at"}, + {"did", "collaborators", tangled.RepoCollaboratorNSID, "repo_at"}, + {"did", "artifacts", tangled.RepoArchiveNSID, "repo_at"}, + {"did", "stars", tangled.FeedStarNSID, "subject_at"}, } for _, src := range sources { @@ -629,7 +633,7 @@ func EnqueuePdsRewritesForRepo(tx *sql.Tx, repoDid, repoAtUri string) error { } for _, p := range pairs { - if err := EnqueuePdsRewrite(tx, p.did, repoDid, src.nsid, p.rkey, repoAtUri); err != nil { + if err := EnqueuePdsRecordMigration(context.Background(), tx, "add-repo-did", syntax.DID(p.did), src.nsid, syntax.RecordKey(p.rkey)); err != nil { return fmt.Errorf("enqueue pds rewrite for %s/%s: %w", src.table, p.rkey, err) } } @@ -657,7 +661,7 @@ func EnqueuePdsRewritesForRepo(tx *sql.Tx, repoDid, repoAtUri string) error { } for _, d := range profileDids { - if err := EnqueuePdsRewrite(tx, d, repoDid, "sh.tangled.actor.profile", "self", repoAtUri); err != nil { + if err := EnqueuePdsRecordMigration(context.Background(), tx, "add-repo-did", syntax.DID(d), tangled.ActorProfileNSID, "self"); err != nil { return fmt.Errorf("enqueue pds rewrite for profile/%s: %w", d, err) } } @@ -665,60 +669,6 @@ func EnqueuePdsRewritesForRepo(tx *sql.Tx, repoDid, repoAtUri string) error { return nil } -type PdsRewrite struct { - Id int - RepoDid string - RecordNsid string - RecordRkey string - OldRepoAt string -} - -func GetPendingPdsRewrites(e Execer, userDid string) ([]PdsRewrite, error) { - rows, err := e.Query( - `SELECT id, repo_did, record_nsid, record_rkey, old_repo_at - FROM pds_rewrite_status - WHERE user_did = ? AND status = 'pending'`, - userDid, - ) - if err != nil { - return nil, err - } - defer rows.Close() - - var rewrites []PdsRewrite - for rows.Next() { - var r PdsRewrite - if err := rows.Scan(&r.Id, &r.RepoDid, &r.RecordNsid, &r.RecordRkey, &r.OldRepoAt); err != nil { - return nil, err - } - rewrites = append(rewrites, r) - } - return rewrites, rows.Err() -} - -func CompletePdsRewrite(e Execer, id int) error { - _, err := e.Exec( - `UPDATE pds_rewrite_status SET status = 'done', updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?`, - id, - ) - return err -} - -func EnqueuePdsRewrite(e Execer, userDid, repoDid, recordNsid, recordRkey, oldRepoAt string) error { - _, err := e.Exec( - `INSERT INTO pds_rewrite_status - (user_did, repo_did, record_nsid, record_rkey, old_repo_at, status) - VALUES (?, ?, ?, ?, ?, 'pending') - ON CONFLICT(user_did, record_nsid, record_rkey) DO UPDATE SET - status = 'pending', - repo_did = excluded.repo_did, - old_repo_at = excluded.old_repo_at, - updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')`, - userDid, repoDid, recordNsid, recordRkey, oldRepoAt, - ) - return err -} - func CascadeRepoDid(tx *sql.Tx, repoAtUri, repoDid string) error { _, err := tx.Exec( `UPDATE repos SET repo_did = ? WHERE at_uri = ?`, diff --git a/appview/ingester.go b/appview/ingester.go index 5be37ab3..812ff583 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -70,11 +70,11 @@ func (i *Ingester) Ingest() processFunc { case tangled.GraphVouchNSID: err = i.ingestVouch(ctx, e) case tangled.FeedStarNSID: - err = i.ingestStar(e) + err = i.ingestStar(ctx, e) case tangled.PublicKeyNSID: err = i.ingestPublicKey(e) case tangled.RepoArtifactNSID: - err = i.ingestArtifact(e) + err = i.ingestArtifact(ctx, e) case tangled.ActorProfileNSID: err = i.ingestProfile(ctx, e) case tangled.SpindleMemberNSID: @@ -114,7 +114,7 @@ func (i *Ingester) Ingest() processFunc { } } -func (i *Ingester) ingestStar(e *jmodels.Event) error { +func (i *Ingester) ingestStar(ctx context.Context, e *jmodels.Event) error { var err error did := e.Did @@ -154,7 +154,7 @@ func (i *Ingester) ingestStar(e *jmodels.Event) error { star.RepoAt = subjectUri repo, repoErr := db.GetRepoByAtUri(i.Db, subjectUri.String()) if repoErr == nil && repo.RepoDid != "" { - if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.FeedStarNSID, e.Commit.RKey, *record.Subject); enqErr != nil { + if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.FeedStarNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { l.Warn("failed to enqueue PDS rewrite for star", "err", enqErr, "did", did, "repoDid", repo.RepoDid) } } @@ -325,7 +325,7 @@ func (i *Ingester) ingestPublicKey(e *jmodels.Event) error { return nil } -func (i *Ingester) ingestArtifact(e *jmodels.Event) error { +func (i *Ingester) ingestArtifact(ctx context.Context, e *jmodels.Event) error { did := e.Did var err error @@ -373,7 +373,7 @@ func (i *Ingester) ingestArtifact(e *jmodels.Event) error { repoDid = *record.RepoDid } if repoDid != "" && (record.RepoDid == nil || *record.RepoDid == "") && record.Repo != nil { - if enqErr := db.EnqueuePdsRewrite(i.Db, did, repoDid, tangled.RepoArtifactNSID, e.Commit.RKey, *record.Repo); enqErr != nil { + if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoArtifactNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { l.Warn("failed to enqueue PDS rewrite for artifact", "err", enqErr, "did", did, "repoDid", repoDid) } } @@ -1011,7 +1011,7 @@ func (i *Ingester) ingestIssue(ctx context.Context, e *jmodels.Event) error { if record.Repo != nil { repo, repoErr := db.GetRepoByAtUri(i.Db, *record.Repo) if repoErr == nil && repo.RepoDid != "" { - if enqErr := db.EnqueuePdsRewrite(i.Db, did, repo.RepoDid, tangled.RepoIssueNSID, rkey, *record.Repo); enqErr != nil { + if enqErr := db.EnqueuePdsRecordMigration(ctx, i.Db, "add-repo-did", syntax.DID(did), syntax.NSID(tangled.RepoIssueNSID), syntax.RecordKey(e.Commit.RKey)); enqErr != nil { l.Warn("failed to enqueue PDS rewrite for issue", "err", enqErr, "did", did, "repoDid", repo.RepoDid) } } diff --git a/appview/oauth/handler.go b/appview/oauth/handler.go index f08836b2..8323046b 100644 --- a/appview/oauth/handler.go +++ b/appview/oauth/handler.go @@ -13,7 +13,6 @@ import ( "time" comatproto "github.com/bluesky-social/indigo/api/atproto" - atpclient "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/auth/oauth" lexutil "github.com/bluesky-social/indigo/lex/util" xrpc "github.com/bluesky-social/indigo/xrpc" @@ -273,153 +272,6 @@ func (o *OAuth) ensureTangledProfile(sessData *oauth.ClientSessionData) { l.Debug("successfully created empty Tangled profile on PDS and DB") } -func (o *OAuth) PdsRewriteMiddleware(next http.Handler) http.Handler { - return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - defer next.ServeHTTP(w, r) - - sess, err := o.ResumeSession(r) - if err != nil { - return - } - - go o.drainPdsRewrites(sess.Data) - }) -} - -func (o *OAuth) drainPdsRewrites(sessData *oauth.ClientSessionData) { - ctx := context.Background() - did := sessData.AccountDID.String() - l := o.Logger.With("did", did, "handler", "drainPdsRewrites") - - rewrites, err := db.GetPendingPdsRewrites(o.Db, did) - if err != nil { - l.Error("failed to get pending rewrites", "err", err) - return - } - if len(rewrites) == 0 { - return - } - - l.Info("draining pending PDS rewrites", "count", len(rewrites)) - - sess, err := o.ClientApp.ResumeSession(ctx, sessData.AccountDID, sessData.SessionID) - if err != nil { - l.Error("failed to resume session for PDS rewrites", "err", err) - return - } - client := sess.APIClient() - - for _, rw := range rewrites { - if err := o.rewritePdsRecord(ctx, client, did, rw); err != nil { - l.Error("failed to rewrite PDS record", - "nsid", rw.RecordNsid, - "rkey", rw.RecordRkey, - "repo_did", rw.RepoDid, - "err", err) - continue - } - - if err := db.CompletePdsRewrite(o.Db, rw.Id); err != nil { - l.Error("failed to mark rewrite complete", "id", rw.Id, "err", err) - } - } -} - -func (o *OAuth) rewritePdsRecord(ctx context.Context, client *atpclient.APIClient, userDid string, rw db.PdsRewrite) error { - ex, err := comatproto.RepoGetRecord(ctx, client, "", rw.RecordNsid, userDid, rw.RecordRkey) - if err != nil { - return fmt.Errorf("get record: %w", err) - } - - val := ex.Value.Val - repoDid := rw.RepoDid - - switch rw.RecordNsid { - case tangled.RepoNSID: - rec, ok := val.(*tangled.Repo) - if !ok { - return fmt.Errorf("unexpected type for repo record") - } - rec.RepoDid = &repoDid - - case tangled.RepoIssueNSID: - rec, ok := val.(*tangled.RepoIssue) - if !ok { - return fmt.Errorf("unexpected type for issue record") - } - rec.RepoDid = &repoDid - - case tangled.RepoPullNSID: - rec, ok := val.(*tangled.RepoPull) - if !ok { - return fmt.Errorf("unexpected type for pull record") - } - if rec.Target != nil { - rec.Target.RepoDid = &repoDid - } - if rec.Source != nil && rec.Source.Repo != nil && *rec.Source.Repo == rw.OldRepoAt { - rec.Source.RepoDid = &repoDid - } - - case tangled.RepoCollaboratorNSID: - rec, ok := val.(*tangled.RepoCollaborator) - if !ok { - return fmt.Errorf("unexpected type for collaborator record") - } - rec.RepoDid = &repoDid - - case tangled.RepoArtifactNSID: - rec, ok := val.(*tangled.RepoArtifact) - if !ok { - return fmt.Errorf("unexpected type for artifact record") - } - rec.RepoDid = &repoDid - - case tangled.FeedStarNSID: - rec, ok := val.(*tangled.FeedStar) - if !ok { - return fmt.Errorf("unexpected type for star record") - } - rec.SubjectDid = &repoDid - - case tangled.ActorProfileNSID: - rec, ok := val.(*tangled.ActorProfile) - if !ok { - return fmt.Errorf("unexpected type for profile record") - } - rewritten := make([]string, 0, len(rec.PinnedRepositories)) - for _, pin := range rec.PinnedRepositories { - if strings.HasPrefix(pin, "did:") { - rewritten = append(rewritten, pin) - continue - } - repo, repoErr := db.GetRepoByAtUri(o.Db, pin) - if repoErr != nil || repo.RepoDid == "" { - rewritten = append(rewritten, pin) - continue - } - rewritten = append(rewritten, repo.RepoDid) - } - rec.PinnedRepositories = rewritten - - default: - return fmt.Errorf("unsupported NSID for PDS rewrite: %s", rw.RecordNsid) - } - - _, err = comatproto.RepoPutRecord(ctx, client, &comatproto.RepoPutRecord_Input{ - Collection: rw.RecordNsid, - Repo: userDid, - Rkey: rw.RecordRkey, - SwapRecord: ex.Cid, - Record: &lexutil.LexiconTypeDecoder{Val: val}, - }) - if err != nil { - return fmt.Errorf("put record: %w", err) - } - - return nil -} - // create a AppPasswordSession using apppasswords type AppPasswordSession struct { AccessJwt string `json:"accessJwt"` diff --git a/appview/state/router.go b/appview/state/router.go index 34437605..8f89458c 100644 --- a/appview/state/router.go +++ b/appview/state/router.go @@ -38,9 +38,6 @@ func (s *State) Router() http.Handler { s.logger, ) - // TODO(boltless): merge this into BackgroundMigrationMiddleware - router.Use(s.oauth.PdsRewriteMiddleware) - m := migration.NewMigration(s.db, s.oauth, s.idResolver.Directory(), s.logger) router.Use(m.BackgroundMigrationMiddleware)