package accountmigration import ( "context" "fmt" "log/slog" "sync" "time" 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" "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/appview/oauth" "tangled.org/core/xrpc/xrpcclient" "tangled.org/core/orm" ) const defaultConcurrency = 4 type Worker struct { db *db.DB oauth *oauth.OAuth dev bool logger *slog.Logger concurrency int startMu sync.Mutex started bool } func NewWorker(d *db.DB, o *oauth.OAuth, dev bool, logger *slog.Logger) *Worker { return &Worker{ db: d, oauth: o, dev: dev, logger: logger, concurrency: defaultConcurrency, } } func (w *Worker) Start(ctx context.Context) { w.startMu.Lock() defer w.startMu.Unlock() if w.started { return } w.started = true for i := 0; i < w.concurrency; i++ { go w.runWorker(ctx, i) } } func (w *Worker) runWorker(ctx context.Context, id int) { l := w.logger.With("worker", id) for { select { case <-ctx.Done(): return default: } row, ok, err := db.ClaimNextPending(ctx, w.db) if err != nil { l.Error("claim failed", "err", err) time.Sleep(time.Second) continue } if !ok { time.Sleep(time.Second) continue } l.Info("migrating repo", "row", row) w.importRepo(ctx, l, row) } } func (w *Worker) importRepo(ctx context.Context, l *slog.Logger, row *db.GitRepoMigration) { l = l.With("id", row.ID, "owner", row.OwnerDid, "name", row.Name, "knot", row.Knot) err := w.doImportRepo(ctx, row) if err == nil { if err := db.MarkGitRepoMigrationDone(ctx, w.db, row.ID); err != nil { l.Error("mark done failed", "err", err) } return } l.Warn("job failed", "err", err) if uerr := db.MarkGitRepoMigrationFailed(ctx, w.db, row.ID, err.Error()); uerr != nil { l.Error("mark failed failed", "err", uerr) } } func (w *Worker) doImportRepo(ctx context.Context, row *db.GitRepoMigration) error { atpClient, err := w.oauth.ClientForDidSession(ctx, row.OwnerDid, row.SessionID) if err != nil { return fmt.Errorf("resume session: %w", err) } existing, _ := db.GetRepo(w.db, orm.FilterEq("did", row.OwnerDid), orm.FilterEq("rkey", row.Name), ) if existing != nil || rkeyOccupied(ctx, atpClient, row.OwnerDid.String(), row.Name) { return fmt.Errorf("rkey %q already exists", row.Name) } sc, err := w.oauth.ServiceClientForDidSession(ctx, row.OwnerDid, row.SessionID, oauth.WithService(row.Knot), oauth.WithLxm(tangled.RepoCreateNSID), oauth.WithDev(w.dev), oauth.WithTimeout(10*time.Minute), ) if err != nil { return fmt.Errorf("knot service client: %w", err) } defaultBranch := "main" createResp, err := tangled.RepoCreate(ctx, sc, &tangled.RepoCreate_Input{ Rkey: row.Name, Name: row.Name, DefaultBranch: &defaultBranch, Source: &row.CloneUrl, }) if xrpcerr := xrpcclient.HandleXrpcErr(err); xrpcerr != nil { return fmt.Errorf("knot RepoCreate: %w", err) } if err != nil { return fmt.Errorf("knot RepoCreate: %w", err) } var repoDid string if createResp != nil && createResp.RepoDid != nil { repoDid = *createResp.RepoDid } if repoDid == "" { return fmt.Errorf("knot returned empty repo DID") } repo := &models.Repo{ Did: row.OwnerDid.String(), Name: row.Name, Knot: row.Knot, Rkey: row.Name, Description: row.Description, Created: time.Now(), RepoDid: repoDid, } record := repo.AsRecord() rkey := row.Name _, err = comatproto.RepoCreateRecord(ctx, atpClient, &comatproto.RepoCreateRecord_Input{ Collection: tangled.RepoNSID, Repo: row.OwnerDid.String(), Rkey: &rkey, Record: &lexutil.LexiconTypeDecoder{ Val: &record, }, }) if err != nil { w.cleanupKnotRepo(ctx, sc, row.OwnerDid, row.Knot, row.Name) return fmt.Errorf("PDS PutRecord: %w", err) } return nil } func (w *Worker) cleanupKnotRepo(ctx context.Context, sc *xrpc.Client, did syntax.DID, knot, name string) { cctx, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() if err := tangled.RepoDelete(cctx, sc, &tangled.RepoDelete_Input{ Did: did.String(), Name: name, Rkey: name, }); err != nil { w.logger.Error("cleanup: RepoDelete failed", "knot", knot, "name", name, "err", err) } } func rkeyOccupied(ctx context.Context, client *atclient.APIClient, did, rkey string) bool { probeCtx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() resp, err := comatproto.RepoGetRecord(probeCtx, client, "", tangled.RepoNSID, did, rkey) return err == nil && resp != nil }