package db import ( "context" "database/sql" "errors" "time" "github.com/bluesky-social/indigo/atproto/syntax" ) type GitRepoMigrationStatus string const ( GitRepoMigrationStatusPending GitRepoMigrationStatus = "pending" GitRepoMigrationStatusRunning GitRepoMigrationStatus = "running" GitRepoMigrationStatusDone GitRepoMigrationStatus = "done" GitRepoMigrationStatusFailed GitRepoMigrationStatus = "failed" ) const ( GitRepoMigrationSourceGitHub = "github" ) type GitRepoMigrations []GitRepoMigration func (m GitRepoMigrations) AnyActive() bool { for _, r := range m { if r.Status == GitRepoMigrationStatusPending || r.Status == GitRepoMigrationStatusRunning { return true } } return false } type GitRepoMigration struct { ID int64 OwnerDid syntax.DID SourceKind string CloneUrl string Name string Knot string Description string Website string Topics []string SessionID string Status GitRepoMigrationStatus ErrorMsg string UpdatedAt time.Time } func InsertGitRepoMigrations(ctx context.Context, e *DB, rows []GitRepoMigration) error { if len(rows) == 0 { return nil } txx, err := e.BeginTx(ctx, nil) if err != nil { return err } defer txx.Rollback() stmt, err := txx.PrepareContext(ctx, ` insert into gitrepo_migrations (owner_did, source_kind, clone_url, name, knot, description, session_id, status, error_msg, updated_at) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) on conflict(owner_did, name) do update set source_kind = excluded.source_kind, clone_url = excluded.clone_url, knot = excluded.knot, description = excluded.description, session_id = excluded.session_id, status = 'pending', error_msg = '', updated_at = excluded.updated_at `) if err != nil { return err } defer stmt.Close() now := time.Now().UTC().Format(time.RFC3339) for _, r := range rows { status := r.Status if status == "" { status = GitRepoMigrationStatusPending } if _, err := stmt.ExecContext(ctx, r.OwnerDid, r.SourceKind, r.CloneUrl, r.Name, r.Knot, r.Description, r.SessionID, status, r.ErrorMsg, now, ); err != nil { return err } } return txx.Commit() } func ClaimNextPending(ctx context.Context, e Execer) (*GitRepoMigration, bool, error) { var ( m GitRepoMigration updatedAt string ) err := e.QueryRowContext(ctx, ` update gitrepo_migrations set status = 'running', updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') where id = ( select id from gitrepo_migrations gm where status = 'pending' and not exists ( select 1 from gitrepo_migrations gm2 where gm2.owner_did = gm.owner_did and gm2.status = 'running' ) order by id limit 1 ) returning id, owner_did, source_kind, clone_url, name, knot, description, session_id, status, error_msg, updated_at `).Scan( &m.ID, &m.OwnerDid, &m.SourceKind, &m.CloneUrl, &m.Name, &m.Knot, &m.Description, &m.SessionID, &m.Status, &m.ErrorMsg, &updatedAt, ) if errors.Is(err, sql.ErrNoRows) { return nil, false, nil } if err != nil { return nil, false, err } if t, err := time.Parse(time.RFC3339, updatedAt); err == nil { m.UpdatedAt = t } return &m, true, nil } func MarkGitRepoMigrationDone(ctx context.Context, e Execer, id int64) error { _, err := e.ExecContext(ctx, ` update gitrepo_migrations set status = 'done', error_msg = '', updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') where id = ? `, id) return err } func MarkGitRepoMigrationFailed(ctx context.Context, e Execer, id int64, errMsg string) error { _, err := e.ExecContext(ctx, ` update gitrepo_migrations set status = 'failed', error_msg = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') where id = ? `, errMsg, id) return err } func ListGitRepoMigrationsForOwner(ctx context.Context, e Execer, owner string) (GitRepoMigrations, error) { rows, err := e.QueryContext(ctx, ` select id, owner_did, source_kind, clone_url, name, knot, description, session_id, status, coalesce(error_msg, ''), updated_at from gitrepo_migrations where owner_did = ? order by id desc `, owner) if err != nil { return nil, err } defer rows.Close() var out GitRepoMigrations for rows.Next() { var ( m GitRepoMigration updatedAt string ) if err := rows.Scan( &m.ID, &m.OwnerDid, &m.SourceKind, &m.CloneUrl, &m.Name, &m.Knot, &m.Description, &m.SessionID, &m.Status, &m.ErrorMsg, &updatedAt, ); err != nil { return nil, err } if t, err := time.Parse(time.RFC3339, updatedAt); err == nil { m.UpdatedAt = t } out = append(out, m) } return out, rows.Err() } func ListEnqueuedGitRepoNames(ctx context.Context, e Execer, owner string) (map[string]struct{}, error) { rows, err := e.QueryContext(ctx, ` select name from gitrepo_migrations where owner_did = ? and status in ('pending', 'running', 'done') `, owner) if err != nil { return nil, err } defer rows.Close() out := map[string]struct{}{} for rows.Next() { var n string if err := rows.Scan(&n); err != nil { return nil, err } out[n] = struct{}{} } return out, rows.Err() } func ReapStaleRunningGitRepoMigrations(ctx context.Context, e Execer) error { _, err := e.ExecContext(ctx, ` update gitrepo_migrations set status = 'pending', updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') where status = 'running' `) return err }