From ef53db1193b8361c16629a044174708f0569eb7a Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 10 Jun 2025 16:42:37 -0700 Subject: [PATCH] model: implement database migrations (#279) * model: implement database migrations * media: shut down streams when keys are revoked --- pkg/atproto/atproto.go | 47 ++------------- pkg/atproto/firehose.go | 20 +++++++ pkg/atproto/locks.go | 47 +++++++++++++++ pkg/atproto/sync.go | 17 +++--- pkg/cmd/streamplace.go | 4 ++ pkg/media/key_revocation.go | 34 +++++++++++ pkg/media/media_signer.go | 12 ++++ pkg/media/media_signer_ext.go | 11 ++++ pkg/media/mkv_ingest.go | 2 + pkg/media/webrtc_ingest.go | 19 +++--- pkg/model/model.go | 2 + pkg/model/repo.go | 9 +++ pkg/model/resync.go | 1 + pkg/model/signing_key.go | 32 ++++++++-- pkg/resync/resync.go | 110 ++++++++++++++++++++++++++++++++++ 15 files changed, 306 insertions(+), 61 deletions(-) create mode 100644 pkg/atproto/locks.go create mode 100644 pkg/media/key_revocation.go create mode 100644 pkg/model/resync.go create mode 100644 pkg/resync/resync.go diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index 125b7290..f14dd40c 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -4,12 +4,9 @@ import ( "bytes" "context" "fmt" - "strings" - "sync" comatproto "github.com/bluesky-social/indigo/api/atproto" _ "github.com/bluesky-social/indigo/api/bsky" - atcrypto "github.com/bluesky-social/indigo/atproto/crypto" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/repo" @@ -17,35 +14,12 @@ import ( "github.com/ipfs/go-cid" "go.opentelemetry.io/otel" "stream.place/streamplace/pkg/aqhttp" - "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/model" ) var SyncGetRepo = comatproto.SyncGetRepo -// handleLocks provides per-handle synchronization -var handleLocks = struct { - sync.Mutex - locks map[string]*sync.Mutex -}{ - locks: make(map[string]*sync.Mutex), -} - -// getHandleLock returns a mutex for the given handle -func getHandleLock(handle string) *sync.Mutex { - handleLocks.Lock() - defer handleLocks.Unlock() - - if lock, exists := handleLocks.locks[handle]; exists { - return lock - } - - lock := &sync.Mutex{} - handleLocks.locks[handle] = lock - return lock -} - func (atsync *ATProtoSynchronizer) SyncBlueskyRepoCached(ctx context.Context, handle string, mod model.Model) (*model.Repo, error) { ctx, span := otel.Tracer("signer").Start(ctx, "SyncBlueskyRepoCached") defer span.End() @@ -66,14 +40,13 @@ type mstNode struct { } func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle string, mod model.Model) (*model.Repo, error) { - ctx = log.WithLogValues(ctx, "func", "SyncBlueskyRepo", "handle", handle) - // Get handle-specific lock and ensure synchronized access - ident, err := ResolveIdent(ctx, handle) if err != nil { return nil, fmt.Errorf("failed to resolve Bluesky handle %s: %w", handle, err) } + ctx = log.WithLogValues(ctx, "did", ident.DID.String(), "handle", ident.Handle.String()) + handleLock := getHandleLock(ident.DID.String()) handleLock.Lock() defer handleLock.Unlock() @@ -102,6 +75,7 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s } log.Log(ctx, "resolved bluesky identity", "did", ident.DID, "handle", ident.Handle, "pds", ident.PDSEndpoint()) + pdsLock := getPDSLock(ident.PDSEndpoint()) xrpcc := xrpc.Client{ Host: ident.PDSEndpoint(), Client: &aqhttp.Client, @@ -109,7 +83,9 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s if xrpcc.Host == "" { return nil, fmt.Errorf("no PDS endpoint found for Bluesky identity %s", handle) } + pdsLock.Lock() repoBytes, err := SyncGetRepo(ctx, &xrpcc, ident.DID.String(), rev) + pdsLock.Unlock() if err != nil { return nil, fmt.Errorf("failed to fetch repo for %s from PDS %s: %w", ident.DID.String(), xrpcc.Host, err) } @@ -124,7 +100,7 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s // return nil, fmt.Errorf("failed to write encoded repo bytes to file: %w", err) // } - log.Log(ctx, "got diff", "bytes", len(repoBytes)) + log.Debug(ctx, "got diff", "bytes", len(repoBytes)) r, err := repo.ReadRepoFromCar(ctx, bytes.NewReader(repoBytes)) if err != nil { @@ -179,17 +155,6 @@ func (atsync *ATProtoSynchronizer) SyncBlueskyRepo(ctx context.Context, handle s return &newRepo, nil } -func parseSigningKey(ctx context.Context, key string) error { - if !strings.HasPrefix(key, constants.DID_KEY_PREFIX) { - return fmt.Errorf("invalid key format for DID key: %s", key) - } - _, err := atcrypto.ParsePublicDIDKey(key) - if err != nil { - return fmt.Errorf("failed to parse multibase key %s: %w", key, err) - } - return nil -} - var ResolveIdent = resolveIdent func resolveIdent(ctx context.Context, arg string) (*identity.Identity, error) { diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index e522300c..dbe1c382 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -246,6 +246,26 @@ func (atsync *ATProtoSynchronizer) handleCommitEventOps(ctx context.Context, evt } } + if collection.String() == constants.PLACE_STREAM_KEY { + log.Warn(ctx, "revoking stream key", "userDID", evt.Repo, "rkey", rkey.String()) + key, err := atsync.Model.GetSigningKeyByRKey(ctx, rkey.String()) + if err != nil { + log.Error(ctx, "failed to get signing key", "err", err) + continue + } + if key == nil { + log.Warn(ctx, "no signing key found for stream key", "userDID", evt.Repo, "rkey", rkey.String()) + continue + } + now := time.Now() + key.RevokedAt = &now + err = atsync.Model.UpdateSigningKey(key) + if err != nil { + log.Error(ctx, "failed to revoke signing key", "err", err) + } + atsync.Bus.Publish(evt.Repo, key) + } + default: log.Error(ctx, "unexpected record op kind") } diff --git a/pkg/atproto/locks.go b/pkg/atproto/locks.go new file mode 100644 index 00000000..8304f399 --- /dev/null +++ b/pkg/atproto/locks.go @@ -0,0 +1,47 @@ +package atproto + +import "sync" + +// handleLocks provides per-handle synchronization +var handleLocks = struct { + sync.Mutex + locks map[string]*sync.Mutex +}{ + locks: make(map[string]*sync.Mutex), +} + +// getHandleLock returns a mutex for the given handle +func getHandleLock(handle string) *sync.Mutex { + handleLocks.Lock() + defer handleLocks.Unlock() + + if lock, exists := handleLocks.locks[handle]; exists { + return lock + } + + lock := &sync.Mutex{} + handleLocks.locks[handle] = lock + return lock +} + +// pdsLocks provides per-pds synchronization +var pdsLocks = struct { + sync.Mutex + locks map[string]*sync.Mutex +}{ + locks: make(map[string]*sync.Mutex), +} + +// getpdsLock returns a mutex for the given pds +func getPDSLock(pds string) *sync.Mutex { + pdsLocks.Lock() + defer pdsLocks.Unlock() + + if lock, exists := pdsLocks.locks[pds]; exists { + return lock + } + + lock := &sync.Mutex{} + pdsLocks.locks[pds] = lock + return lock +} diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index 2b5f448a..588197c5 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -87,17 +87,19 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD if err != nil { return fmt.Errorf("failed to sync bluesky repo: %w", err) } - streamerRepo, err := atsync.SyncBlueskyRepoCached(ctx, rec.Streamer, atsync.Model) + + _, err = atsync.SyncBlueskyRepoCached(ctx, rec.Streamer, atsync.Model) if err != nil { - return fmt.Errorf("failed to sync bluesky repo: %w", err) + log.Error(ctx, "failed to sync bluesky repo", "err", err) } + log.Debug(ctx, "streamplace.ChatMessage detected", "message", rec.Text, "repo", repo.Handle) - block, err := atsync.Model.GetUserBlock(ctx, streamerRepo.DID, userDID) + block, err := atsync.Model.GetUserBlock(ctx, rec.Streamer, userDID) if err != nil { return fmt.Errorf("failed to get user block: %w", err) } if block != nil { - log.Debug(ctx, "excluding message from blocked user", "userDID", userDID, "subjectDID", streamerRepo.DID) + log.Debug(ctx, "excluding message from blocked user", "userDID", userDID, "subjectDID", rec.Streamer) return nil } mcm := &model.ChatMessage{ @@ -107,7 +109,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD ChatMessage: recCBOR, RepoDID: userDID, Repo: repo, - StreamerRepoDID: streamerRepo.DID, + StreamerRepoDID: rec.Streamer, IndexedAt: &now, } if rec.Reply != nil && rec.Reply.Parent != nil && rec.Reply.Root != nil { @@ -129,11 +131,11 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD if err != nil { log.Error(ctx, "failed to convert chat message to streamplace message view", "err", err) } - go atsync.Bus.Publish(streamerRepo.DID, scm) + go atsync.Bus.Publish(rec.Streamer, scm) if !isUpdate && !isFirstSync { for _, webhook := range atsync.CLI.DiscordWebhooks { - if webhook.DID == streamerRepo.DID && webhook.Type == "chat" { + if webhook.DID == rec.Streamer && webhook.Type == "chat" { go func() { err := discord.SendChat(ctx, webhook, r, scm) if err != nil { @@ -362,6 +364,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD } key := model.SigningKey{ DID: rec.SigningKey, + RKey: rkey.String(), CreatedAt: time.Time(), RepoDID: userDID, } diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 8f545096..67fd9921 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -29,6 +29,7 @@ import ( "stream.place/streamplace/pkg/notifications" "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/replication/boring" + "stream.place/streamplace/pkg/resync" "stream.place/streamplace/pkg/rtmps" v0 "stream.place/streamplace/pkg/schema/v0" "stream.place/streamplace/pkg/spmetrics" @@ -206,6 +207,9 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { spmetrics.Version.WithLabelValues(build.Version).Inc() aqhttp.UserAgent = fmt.Sprintf("streamplace/%s", build.Version) + if len(os.Args) > 1 && os.Args[1] == "resync" { + return resync.Resync(ctx, &cli) + } err = os.MkdirAll(cli.DataDir, os.ModePerm) if err != nil { diff --git a/pkg/media/key_revocation.go b/pkg/media/key_revocation.go new file mode 100644 index 00000000..c8486543 --- /dev/null +++ b/pkg/media/key_revocation.go @@ -0,0 +1,34 @@ +package media + +import ( + "context" + "fmt" + + "github.com/go-gst/go-gst/gst" + "stream.place/streamplace/pkg/model" +) + +// Handle shutting down a pipeline when a signing key is revoked +func (mm *MediaManager) HandleKeyRevocation(ctx context.Context, ms MediaSigner, pipeline *gst.Pipeline) { + sub := mm.bus.Subscribe(ms.Streamer()) + defer mm.bus.Unsubscribe(ms.Streamer(), sub) + for { + select { + case <-ctx.Done(): + return + case msg := <-sub: + signingKey, ok := msg.(*model.SigningKey) + if !ok { + continue + } + if signingKey.RevokedAt == nil { + continue + } + if signingKey.DID == ms.DID() { + err := fmt.Errorf("signing key revoked, ending stream: %s", signingKey.RKey) + pipeline.Error(err.Error(), err) + return + } + } + } +} diff --git a/pkg/media/media_signer.go b/pkg/media/media_signer.go index c8e0946b..05f26a4e 100644 --- a/pkg/media/media_signer.go +++ b/pkg/media/media_signer.go @@ -15,6 +15,7 @@ import ( "go.opentelemetry.io/otel" "stream.place/streamplace/pkg/aqio" "stream.place/streamplace/pkg/aqtime" + "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/crypto/aqpub" "stream.place/streamplace/pkg/crypto/signers" @@ -26,6 +27,7 @@ type MediaSigner interface { SignMP4(ctx context.Context, input io.ReadSeeker, start int64) ([]byte, error) Pub() aqpub.Pub Streamer() string + DID() string } type MediaSignerLocal struct { @@ -34,6 +36,7 @@ type MediaSignerLocal struct { AQPub aqpub.Pub Cert []byte TAURL string + did string } func prepareCert(ctx context.Context, cli *config.CLI, signer crypto.Signer) ([]byte, string, error) { @@ -77,12 +80,17 @@ func MakeMediaSigner(ctx context.Context, cli *config.CLI, streamer string, sign if err != nil { return nil, err } + did, err := atproto.ParsePubKey(signer.Public().(*ecdsa.PublicKey)) + if err != nil { + return nil, err + } return &MediaSignerLocal{ Signer: signer, Cert: cert, StreamerName: streamer, TAURL: cli.TAURL, AQPub: pub, + did: did.DIDKey(), }, nil } @@ -173,3 +181,7 @@ func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, st func (ms *MediaSignerLocal) Pub() aqpub.Pub { return ms.AQPub } + +func (ms *MediaSignerLocal) DID() string { + return ms.did +} diff --git a/pkg/media/media_signer_ext.go b/pkg/media/media_signer_ext.go index 2b5a7471..66fd7aa6 100644 --- a/pkg/media/media_signer_ext.go +++ b/pkg/media/media_signer_ext.go @@ -14,6 +14,7 @@ import ( "github.com/decred/dcrd/dcrec/secp256k1" "github.com/mr-tron/base58" "go.opentelemetry.io/otel" + "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/crypto/aqpub" "stream.place/streamplace/pkg/spmetrics" @@ -27,6 +28,7 @@ type MediaSignerExt struct { streamer string keyBs []byte taURL string + did string } func MakeMediaSignerExt(ctx context.Context, cli *config.CLI, streamer string, keyBs []byte) (MediaSigner, error) { @@ -43,6 +45,10 @@ func MakeMediaSignerExt(ctx context.Context, cli *config.CLI, streamer string, k if err != nil { return nil, err } + did, err := atproto.ParsePubKey(signer.Public().(*ecdsa.PublicKey)) + if err != nil { + return nil, err + } return &MediaSignerExt{ // cli: cli, signer: signer, @@ -51,6 +57,7 @@ func MakeMediaSignerExt(ctx context.Context, cli *config.CLI, streamer string, k pub: pub, keyBs: keyBs, taURL: cli.TAURL, + did: did.DIDKey(), }, nil } @@ -115,3 +122,7 @@ func (ms *MediaSignerExt) Pub() aqpub.Pub { func (ms *MediaSignerExt) Streamer() string { return ms.streamer } + +func (ms *MediaSignerExt) DID() string { + return ms.did +} diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go index 00b3a629..66f3ce9e 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mkv_ingest.go @@ -68,6 +68,8 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media busErr <- err }() + go mm.HandleKeyRevocation(ctx, ms, pipeline) + err = pipeline.SetState(gst.StatePlaying) if err != nil { return err diff --git a/pkg/media/webrtc_ingest.go b/pkg/media/webrtc_ingest.go index a1259ecb..a73cb87f 100644 --- a/pkg/media/webrtc_ingest.go +++ b/pkg/media/webrtc_ingest.go @@ -99,13 +99,6 @@ func (mm *MediaManager) WebRTCIngest(ctx context.Context, offer *webrtc.SessionD } audioSrc := app.SrcFromElement(audioSrcElem) - go func() { - <-ctx.Done() - if cErr := peerConnection.Close(); cErr != nil { - log.Log(ctx, "cannot close peerConnection: %v\n", cErr) - } - }() - // Set the remote SessionDescription if err = peerConnection.SetRemoteDescription(*offer); err != nil { return nil, fmt.Errorf("failed to set remote description: %w", err) @@ -145,11 +138,21 @@ func (mm *MediaManager) WebRTCIngest(ctx context.Context, offer *webrtc.SessionD go func() { if err := HandleBusMessages(ctx, pipeline); err != nil { - log.Log(ctx, "pipeilne error", "error", err) + log.Log(ctx, "pipeline error", "error", err) } cancel() }() + // subscription to bus messages for key revocation + go mm.HandleKeyRevocation(ctx, signer, pipeline) + + go func() { + <-ctx.Done() + if cErr := peerConnection.Close(); cErr != nil { + log.Log(ctx, "cannot close peerConnection: %v\n", cErr) + } + }() + log.Debug(ctx, "starting pipeline") // Start the pipeline diff --git a/pkg/model/model.go b/pkg/model/model.go index 3c17ecca..ed3f3311 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -49,10 +49,12 @@ type Model interface { GetRepoByHandle(handle string) (*Repo, error) GetRepoByHandleOrDID(arg string) (*Repo, error) GetRepoBySigningKey(signingKey string) (*Repo, error) + GetAllRepos() ([]Repo, error) UpdateRepo(repo *Repo) error UpdateSigningKey(key *SigningKey) error GetSigningKey(ctx context.Context, did, repoDID string) (*SigningKey, error) + GetSigningKeyByRKey(ctx context.Context, rkey string) (*SigningKey, error) GetSigningKeysForRepo(repoDID string) ([]SigningKey, error) CreateFollow(ctx context.Context, userDID, rev string, follow *bsky.GraphFollow) error diff --git a/pkg/model/repo.go b/pkg/model/repo.go index fe20f9d8..64ca807c 100644 --- a/pkg/model/repo.go +++ b/pkg/model/repo.go @@ -30,6 +30,15 @@ func (m *DBModel) GetRepo(did string) (*Repo, error) { return &repoModel, nil } +func (m *DBModel) GetAllRepos() ([]Repo, error) { + var repos []Repo + res := m.DB.Find(&repos) + if res.Error != nil { + return nil, res.Error + } + return repos, nil +} + func (m *DBModel) GetRepoByHandle(handle string) (*Repo, error) { var repoModel Repo res := m.DB.Where("handle = ?", handle).First(&repoModel) diff --git a/pkg/model/resync.go b/pkg/model/resync.go new file mode 100644 index 00000000..8b537907 --- /dev/null +++ b/pkg/model/resync.go @@ -0,0 +1 @@ +package model diff --git a/pkg/model/signing_key.go b/pkg/model/signing_key.go index ae5c1bea..2e2a2f23 100644 --- a/pkg/model/signing_key.go +++ b/pkg/model/signing_key.go @@ -3,6 +3,7 @@ package model import ( "context" "errors" + "fmt" "time" "go.opentelemetry.io/otel" @@ -10,11 +11,12 @@ import ( ) type SigningKey struct { - DID string `gorm:"primaryKey;column:did" json:"did"` - RepoDID string `gorm:"primaryKey;column:repo_did" json:"repoDID"` - Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:RepoDID;references:DID"` - CreatedAt time.Time `json:"createdAt"` - RevokedAt time.Time `json:"revokedAt"` + DID string `gorm:"primaryKey;column:did" json:"did"` + RepoDID string `gorm:"primaryKey;column:repo_did" json:"repoDID"` + RKey string `gorm:"column:rkey;index" json:"rkey"` + Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:RepoDID;references:DID"` + CreatedAt time.Time `json:"createdAt"` + RevokedAt *time.Time `json:"revokedAt"` } func (SigningKey) TableName() string { @@ -33,6 +35,26 @@ func (m *DBModel) GetSigningKey(ctx context.Context, did, repoDID string) (*Sign if errors.Is(res.Error, gorm.ErrRecordNotFound) { return nil, nil } + if key.RevokedAt != nil { + return nil, fmt.Errorf("signing key revoked") + } + if res.Error != nil { + return nil, res.Error + } + return &key, nil +} + +func (m *DBModel) GetSigningKeyByRKey(ctx context.Context, rkey string) (*SigningKey, error) { + _, span := otel.Tracer("signer").Start(ctx, "GetSigningKeyByRKey") + defer span.End() + var key SigningKey + res := m.DB.Model(SigningKey{}).Where("rkey = ?", rkey).First(&key) + if errors.Is(res.Error, gorm.ErrRecordNotFound) { + return nil, nil + } + if key.RevokedAt != nil { + return nil, fmt.Errorf("signing key revoked") + } if res.Error != nil { return nil, res.Error } diff --git a/pkg/resync/resync.go b/pkg/resync/resync.go new file mode 100644 index 00000000..5bb2029c --- /dev/null +++ b/pkg/resync/resync.go @@ -0,0 +1,110 @@ +package resync + +import ( + "context" + "fmt" + "time" + + "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/atproto" + "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/model" +) + +// resync a fresh database from the PDSses, copying over the few pieces of local state +// that we have +func Resync(ctx context.Context, cli *config.CLI) error { + oldMod, err := model.MakeDB(cli.DBPath) + if err != nil { + return err + } + tempDBPath := cli.DBPath + ".temp." + fmt.Sprintf("%d", time.Now().UnixNano()) + newMod, err := model.MakeDB(tempDBPath) + if err != nil { + return err + } + repos, err := oldMod.GetAllRepos() + if err != nil { + return err + } + + atsync := &atproto.ATProtoSynchronizer{ + CLI: cli, + Model: newMod, + Noter: nil, + Bus: bus.NewBus(), + } + + doneMap := make(map[string]bool) + + g, ctx := errgroup.WithContext(ctx) + + doneChan := make(chan string) + go func() { + for { + select { + case <-ctx.Done(): + return + case did := <-doneChan: + doneMap[did] = true + case <-time.After(10 * time.Second): + for _, repo := range repos { + if !doneMap[repo.DID] { + log.Warn(ctx, "remaining repos to sync", "did", repo.DID, "handle", repo.Handle, "pds", repo.PDS) + } + } + } + } + }() + + for _, repo := range repos { + repo := repo // capture range variable + doneMap[repo.DID] = false + g.Go(func() error { + log.Warn(ctx, "syncing repo", "did", repo.DID, "handle", repo.Handle) + ctx := log.WithLogValues(ctx, "resyncDID", repo.DID, "resyncHandle", repo.Handle) + _, err := atsync.SyncBlueskyRepoCached(ctx, repo.Handle, newMod) + if err != nil { + log.Error(ctx, "failed to sync repo", "did", repo.DID, "handle", repo.Handle, "err", err) + return nil + } + log.Log(ctx, "synced repo", "did", repo.DID, "handle", repo.Handle) + doneChan <- repo.DID + return nil + }) + } + + if err := g.Wait(); err != nil { + return err + } + + oauthSessions, err := oldMod.ListOAuthSessions() + if err != nil { + return err + } + for _, session := range oauthSessions { + err := newMod.CreateOAuthSession(session.DownstreamDPoPJKT, &session) + if err != nil { + return fmt.Errorf("failed to create oauth session: %w", err) + } + } + log.Log(ctx, "migrated oauth sessions", "count", len(oauthSessions)) + + notificationTokens, err := oldMod.ListNotifications() + if err != nil { + return err + } + for _, token := range notificationTokens { + err := newMod.CreateNotification(token.Token, token.RepoDID) + if err != nil { + return fmt.Errorf("failed to create notification: %w", err) + } + } + log.Log(ctx, "migrated notification tokens", "count", len(notificationTokens)) + + log.Log(ctx, "resync complete!", "newDBPath", tempDBPath) + + return nil +} -- 2.51.2