diff --git a/appview/db/db.go b/appview/db/db.go index c999b280..466be3da 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -1429,6 +1429,28 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return err }) + orm.RunMigration(conn, logger, "add-pds-migration", func(tx *sql.Tx) error { + _, err := tx.Exec(` + create table if not exists pds_migration ( + name text not null, + + -- record at_uri + did text not null, + collection text not null, + rkey text not null, + + status text not null default 'pending', + error_msg text, + retry_count integer not null default 0, + retry_after integer not null default 0, + updated_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + unique(name, did, collection, rkey) + ); + `) + return err + }) + return &DB{ db, logger, diff --git a/appview/db/migration.go b/appview/db/migration.go new file mode 100644 index 00000000..f75c49cd --- /dev/null +++ b/appview/db/migration.go @@ -0,0 +1,89 @@ +package db + +import ( + "context" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/appview/models" +) + +// "migration" for records stored in user's PDS, not AppView DB + +// ListPendingPdsRecordMigrations queries list of pending PDS migrations for given user. +// Only pending migrations whose `retry_after` has elapsed are returned. +func ListPendingPdsRecordMigrations(ctx context.Context, e Execer, user syntax.DID) ([]*models.PDSMigration, error) { + rows, err := e.QueryContext(ctx, + `with picked as ( + select rowid + from pds_migration + where did = ? + and status = 'pending' + and retry_after < ? + ) + update pds_migration + set status = ? + where rowid in (select rowid from picked) + returning name, did, collection, rkey, status, error_msg, retry_count, retry_after`, + user, + time.Now().Unix(), + models.PDSMigrationStatusRunning, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var migrations []*models.PDSMigration + for rows.Next() { + var migration models.PDSMigration + if err := rows.Scan( + &migration.Name, + &migration.Did, + &migration.Collection, + &migration.Rkey, + &migration.Status, + &migration.ErrorMsg, + &migration.RetryCount, + &migration.RetryAfter, + ); err != nil { + return nil, err + } + migrations = append(migrations, &migration) + } + if err := rows.Err(); err != nil { + return nil, err + } + + return migrations, nil +} + +func EnqueuePdsRecordMigration(ctx context.Context, e Execer, name string, did syntax.DID, collection syntax.NSID, rkey syntax.RecordKey) error { + _, err := e.ExecContext(ctx, + `insert into pds_migration (name, did, collection, rkey) + values (?, ?, ?, ?)`, + name, did, collection, rkey, + ) + return err +} + +func UpdatePdsRecordMigration(ctx context.Context, e Execer, migration *models.PDSMigration) error { + _, err := e.ExecContext(ctx, + `update pds_migration + set status = ?, + error_msg = ?, + retry_count = ?, + retry_after = ?, + updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') + where name = ? and did = ? and collection = ? and rkey = ?`, + migration.Status, + migration.ErrorMsg, + migration.RetryCount, + migration.RetryAfter, + migration.Name, + migration.Did, + migration.Collection, + migration.Rkey, + ) + return err +} diff --git a/appview/migration/migrate_add_repo_did.go b/appview/migration/migrate_add_repo_did.go new file mode 100644 index 00000000..02cea4da --- /dev/null +++ b/appview/migration/migrate_add_repo_did.go @@ -0,0 +1,151 @@ +package migration + +import ( + "context" + "fmt" + "strings" + + comatproto "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/atclient" + "github.com/bluesky-social/indigo/atproto/syntax" + lexutil "github.com/bluesky-social/indigo/lex/util" + "tangled.org/core/api/tangled" + "tangled.org/core/appview/db" +) + +func (s *Migration) migrateAddRepoDid(ctx context.Context, client *atclient.APIClient, did syntax.DID, record syntax.ATURI) error { + // TODO: use agnostic.RepoGetRecord instead + ex, err := comatproto.RepoGetRecord(ctx, client, "", record.Collection().String(), did.String(), record.RecordKey().String()) + if err != nil { + return fmt.Errorf("pds: %w", err) + } + + val := ex.Value.Val + + switch record.Collection() { + case tangled.RepoNSID: + rec, ok := val.(*tangled.Repo) + if !ok { + return fmt.Errorf("unexpected type for repo record") + } + repo, err := db.GetRepoByAtUri(s.db, record.String()) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.RepoDid = &repo.RepoDid + + case tangled.RepoIssueNSID: + rec, ok := val.(*tangled.RepoIssue) + if !ok { + return fmt.Errorf("unexpected type for issue record") + } + if rec.Repo != nil { + repoAt := *rec.Repo + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.RepoDid = &repo.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.Repo != nil { + repoAt := *rec.Target.Repo + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.Target.RepoDid = &repo.RepoDid + } + if rec.Source != nil && rec.Source.Repo != nil { + repoAt := *rec.Source.Repo + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.Source.RepoDid = &repo.RepoDid + } + + case tangled.RepoCollaboratorNSID: + rec, ok := val.(*tangled.RepoCollaborator) + if !ok { + return fmt.Errorf("unexpected type for collaborator record") + } + if rec.Repo != nil { + repoAt := *rec.Repo + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.RepoDid = &repo.RepoDid + } + + case tangled.RepoArtifactNSID: + rec, ok := val.(*tangled.RepoArtifact) + if !ok { + return fmt.Errorf("unexpected type for artifact record") + } + if rec.Repo != nil { + repoAt := *rec.Repo + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.RepoDid = &repo.RepoDid + } + + case tangled.FeedStarNSID: + rec, ok := val.(*tangled.FeedStar) + if !ok { + return fmt.Errorf("unexpected type for star record") + } + if rec.Subject != nil { + repoAt := *rec.Subject + repo, err := db.GetRepoByAtUri(s.db, repoAt) + if err != nil { + return fmt.Errorf("db: failed to query repo: %w", err) + } + rec.SubjectDid = &repo.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(s.db, pin) + if repoErr != nil || repo.RepoDid == "" { + rewritten = append(rewritten, pin) + continue + } + rewritten = append(rewritten, repo.RepoDid) + } + rec.PinnedRepositories = rewritten + + default: + return fmt.Errorf("unexpected collection: '%s'", record.Collection()) + } + + _, err = comatproto.RepoPutRecord(ctx, client, &comatproto.RepoPutRecord_Input{ + Repo: did.String(), + Collection: record.Collection().String(), + Rkey: record.RecordKey().String(), + SwapRecord: ex.Cid, + Record: &lexutil.LexiconTypeDecoder{Val: val}, + }) + if err != nil { + return fmt.Errorf("put record: %w", err) + } + + return nil +} diff --git a/appview/migration/migration.go b/appview/migration/migration.go new file mode 100644 index 00000000..4c9eb135 --- /dev/null +++ b/appview/migration/migration.go @@ -0,0 +1,100 @@ +package migration + +import ( + "context" + "fmt" + "log/slog" + "net/http" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/atclient" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" + + "tangled.org/core/appview/db" + "tangled.org/core/appview/models" + "tangled.org/core/appview/oauth" +) + +type Migration struct { + db *db.DB + oauth *oauth.OAuth + dir identity.Directory + logger *slog.Logger +} + +func NewMigration(db *db.DB, oauth *oauth.OAuth, dir identity.Directory, logger *slog.Logger) *Migration { + return &Migration{ + db, oauth, dir, logger, + } +} + +func (s *Migration) BackgroundMigrationMiddleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + defer next.ServeHTTP(w, r) + + client, err := s.oauth.AuthorizedClient(r) + if err != nil { + return + } + if client.AccountDID == nil { + return + } + + go s.runPendingMigrations(context.Background(), *client.AccountDID, client) + }) +} + +func (s *Migration) runPendingMigrations(ctx context.Context, did syntax.DID, client *atclient.APIClient) { + l := s.logger.With("did", did) + migrations, err := db.ListPendingPdsRecordMigrations(ctx, s.db, did) + if err != nil { + l.Error("failed to query pending migrations", "err", err) + return + } + + for _, migration := range migrations { + if err := s.migrate(ctx, client, migration); err != nil { + l.Error("migration failed", "err", err) + } + } +} + +func (s *Migration) migrate(ctx context.Context, client *atclient.APIClient, migration *models.PDSMigration) error { + l := s.logger.With( + "name", migration.Name, + "aturi", migration.RecordAtUri(), + ) + + var err error + switch migration.Name { + case "add-repo-did": + err = s.migrateAddRepoDid(ctx, client, migration.Did, migration.RecordAtUri()) + default: + return fmt.Errorf("unexpected migration name %s", migration.Name) + } + + if err == nil { + l.Info("migrated") + migration.Status = models.PDSMigrationStatusDone + } else { + l.Warn("failed to migrate", "err", err) + + errMsg := err.Error() + var retryCount = migration.RetryCount + 1 + var retryAfter = time.Now().Add(3 * time.Second).Unix() + + // remove null bytes + errMsg = strings.ReplaceAll(errMsg, "\x00", "") + + migration.Status = models.PDSMigrationStatusPending + migration.ErrorMsg = &errMsg + migration.RetryCount = retryCount + migration.RetryAfter = retryAfter + } + if err := db.UpdatePdsRecordMigration(ctx, s.db, migration); err != nil { + return fmt.Errorf("failed to update migration status: %w", err) + } + return nil +} diff --git a/appview/models/migration.go b/appview/models/migration.go new file mode 100644 index 00000000..28bc5069 --- /dev/null +++ b/appview/models/migration.go @@ -0,0 +1,37 @@ +package models + +import ( + "fmt" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type PDSRecordMigration struct { + Did syntax.DID + Name string // name of the migration + Records []syntax.ATURI // records that need a migration + ErrorMsg *string // error message from previous attempt +} + +type PDSMigration struct { + Name string // name of the migration + Did syntax.DID // record owner + Collection syntax.NSID // record collection + Rkey syntax.RecordKey // record rkey + Status PDSMigrationStatus + ErrorMsg *string // error message from previous attempt + RetryCount int + RetryAfter int64 // Unix timestamp (seconds) +} + +type PDSMigrationStatus string + +const ( + PDSMigrationStatusPending PDSMigrationStatus = "pending" + PDSMigrationStatusRunning PDSMigrationStatus = "running" + PDSMigrationStatusDone PDSMigrationStatus = "done" +) + +func (m *PDSMigration) RecordAtUri() syntax.ATURI { + return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", m.Did, m.Collection, m.Rkey)) +} diff --git a/appview/state/router.go b/appview/state/router.go index 827abc13..34437605 100644 --- a/appview/state/router.go +++ b/appview/state/router.go @@ -12,6 +12,7 @@ import ( "tangled.org/core/appview/knots" "tangled.org/core/appview/labels" "tangled.org/core/appview/middleware" + "tangled.org/core/appview/migration" "tangled.org/core/appview/notifications" "tangled.org/core/appview/pipelines" "tangled.org/core/appview/pulls" @@ -37,8 +38,12 @@ 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) + router.Get("/pwa-manifest.json", s.WebAppManifest) router.Get("/robots.txt", s.RobotsTxt) router.Get("/.well-known/security.txt", s.SecurityTxt)