diff --git a/cmd/knotmirror/main.go b/cmd/knotmirror/main.go index 749f479a..9867d68d 100644 --- a/cmd/knotmirror/main.go +++ b/cmd/knotmirror/main.go @@ -11,6 +11,8 @@ import ( "github.com/urfave/cli/v3" "tangled.org/core/knotmirror" "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/migrate" "tangled.org/core/log" ) @@ -42,6 +44,12 @@ func run(args []string) error { Action: runKnotMirror, Flags: []cli.Flag{}, }, + { + Name: "migrate-disk", + Usage: "rename mirror dirs from {did}/{rkey} -> {repo_did}; daemon must be stopped", + Action: runMigrateDisk, + Flags: []cli.Flag{}, + }, } return app.Run(ctx, args) } @@ -56,3 +64,20 @@ func runKnotMirror(ctx context.Context, cmd *cli.Command) error { logger.Debug("config loaded:", "config", cfg) return knotmirror.Run(ctx, cfg) } + +func runMigrateDisk(ctx context.Context, cmd *cli.Command) error { + logger := log.FromContext(ctx) + cfg, err := config.Load(ctx) + if err != nil { + return err + } + database, err := db.Make(ctx, cfg.DbUrl, 4) + if err != nil { + return err + } + defer database.Close() + + stats, err := migrate.RenameDisk(ctx, cfg.GitRepoBasePath, database, logger) + logger.Info("migrate-disk complete", "stats", stats.String()) + return err +} diff --git a/knotmirror/db/db.go b/knotmirror/db/db.go index 22eda56e..24540c24 100644 --- a/knotmirror/db/db.go +++ b/knotmirror/db/db.go @@ -7,6 +7,7 @@ import ( "time" _ "github.com/jackc/pgx/v5/stdlib" + "tangled.org/core/log" ) func Make(ctx context.Context, dbUrl string, maxConns int) (*sql.DB, error) { @@ -96,5 +97,9 @@ func Make(ctx context.Context, dbUrl string, maxConns int) (*sql.DB, error) { return nil, fmt.Errorf("initializing db schema: %w", err) } + if err := RunMigrations(ctx, conn, log.FromContext(ctx), Migrations); err != nil { + return nil, fmt.Errorf("running migrations: %w", err) + } + return db, nil } diff --git a/knotmirror/db/migrations.go b/knotmirror/db/migrations.go new file mode 100644 index 00000000..8a0f2879 --- /dev/null +++ b/knotmirror/db/migrations.go @@ -0,0 +1,72 @@ +package db + +import ( + "context" + "database/sql" + "fmt" + "log/slog" +) + +type MigrationFn = func(context.Context, *sql.Tx) error + +type Migration struct { + Name string + Fn MigrationFn +} + +func ensureMigrationsTable(ctx context.Context, conn *sql.Conn) error { + _, err := conn.ExecContext(ctx, ` + create table if not exists migrations ( + name text primary key, + applied_at timestamptz not null default now() + ); + `) + return err +} + +func RunMigration(ctx context.Context, conn *sql.Conn, logger *slog.Logger, m Migration) error { + logger = logger.With("migration", m.Name) + + tx, err := conn.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin migration tx: %w", err) + } + defer tx.Rollback() + + var exists bool + if err := tx.QueryRowContext(ctx, `select exists (select 1 from migrations where name = $1)`, m.Name).Scan(&exists); err != nil { + return fmt.Errorf("checking migration state: %w", err) + } + if exists { + logger.Debug("migration already applied") + return nil + } + + if err := m.Fn(ctx, tx); err != nil { + logger.Error("migration failed", "err", err) + return fmt.Errorf("running migration %s: %w", m.Name, err) + } + + if _, err := tx.ExecContext(ctx, `insert into migrations (name) values ($1)`, m.Name); err != nil { + return fmt.Errorf("recording migration: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("commit migration: %w", err) + } + + logger.Info("migration applied") + return nil +} + +func RunMigrations(ctx context.Context, conn *sql.Conn, logger *slog.Logger, ms []Migration) error { + if err := ensureMigrationsTable(ctx, conn); err != nil { + return fmt.Errorf("ensuring migrations table: %w", err) + } + for _, m := range ms { + if err := RunMigration(ctx, conn, logger, m); err != nil { + return err + } + } + return nil +} diff --git a/knotmirror/db/migrations_list.go b/knotmirror/db/migrations_list.go new file mode 100644 index 00000000..256ec524 --- /dev/null +++ b/knotmirror/db/migrations_list.go @@ -0,0 +1,56 @@ +package db + +import ( + "context" + "database/sql" + "fmt" + + "tangled.org/core/log" +) + +var Migrations = []Migration{ + { + Name: "repos_pk_to_repo_did", + Fn: reposPkToRepoDid, + }, +} + +func reposPkToRepoDid(ctx context.Context, tx *sql.Tx) error { + if _, err := tx.ExecContext(ctx, + `alter table repos add column if not exists repo_did text`, + ); err != nil { + return fmt.Errorf("adding repo_did column: %w", err) + } + var bad int + if err := tx.QueryRowContext(ctx, + `select count(*) from repos where repo_did is null or repo_did = ''`, + ).Scan(&bad); err != nil { + return fmt.Errorf("counting rows with null repo_did: %w", err) + } + if bad > 0 { + log.FromContext(ctx).Warn( + "dropping repos with null repo_did; their on-disk dirs will be orphaned. re-crawl via tap to restore", + "count", bad, + ) + if _, err := tx.ExecContext(ctx, + `delete from repos where repo_did is null or repo_did = ''`, + ); err != nil { + return fmt.Errorf("deleting null repo_did rows: %w", err) + } + } + return execAll(ctx, tx, + `alter table repos alter column repo_did set not null`, + `alter table repos drop constraint if exists repos_pkey`, + `alter table repos add constraint repos_pkey primary key (repo_did)`, + ) +} + +func execAll(ctx context.Context, tx *sql.Tx, stmts ...string) error { + if len(stmts) == 0 { + return nil + } + if _, err := tx.ExecContext(ctx, stmts[0]); err != nil { + return err + } + return execAll(ctx, tx, stmts[1:]...) +} diff --git a/knotmirror/db/repos.go b/knotmirror/db/repos.go index 36313b1e..29565d4e 100644 --- a/knotmirror/db/repos.go +++ b/knotmirror/db/repos.go @@ -11,22 +11,16 @@ import ( "tangled.org/core/knotmirror/models" ) -func AddRepo(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.RecordKey, cid syntax.CID, name, knot string) error { - if _, err := e.ExecContext(ctx, - `insert into repos (did, rkey, cid, name, knot_domain) - values ($1, $2, $3, $4, $5)`, - did, rkey, cid, name, knot, - ); err != nil { - return fmt.Errorf("inserting repo: %w", err) - } - return nil -} - func UpsertRepo(ctx context.Context, e *sql.DB, repo *models.Repo) error { + if repo.RepoDid == "" { + return fmt.Errorf("upsert repo: repo_did is required") + } if _, err := e.ExecContext(ctx, - `insert into repos (did, rkey, cid, name, knot_domain, git_rev, repo_sha, state, error_msg, retry_count, retry_after) - values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) - on conflict(did, rkey) do update set + `insert into repos (did, rkey, cid, name, knot_domain, repo_did, git_rev, repo_sha, state, error_msg, retry_count, retry_after) + values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12) + on conflict(repo_did) do update set + did = excluded.did, + rkey = excluded.rkey, cid = excluded.cid, name = excluded.name, knot_domain = excluded.knot_domain, @@ -36,12 +30,12 @@ func UpsertRepo(ctx context.Context, e *sql.DB, repo *models.Repo) error { error_msg = excluded.error_msg, retry_count = excluded.retry_count, retry_after = excluded.retry_after`, - // where repos.cid != excluded.cid`, repo.Did, repo.Rkey, repo.Cid, repo.Name, repo.KnotDomain, + repo.RepoDid, repo.GitRev, repo.RepoSha, repo.State, @@ -54,13 +48,13 @@ func UpsertRepo(ctx context.Context, e *sql.DB, repo *models.Repo) error { return nil } -func UpdateRepoState(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.RecordKey, state models.RepoState) error { +func UpdateRepoState(ctx context.Context, e *sql.DB, repoDid syntax.DID, state models.RepoState) error { if _, err := e.ExecContext(ctx, `update repos set state = $1 - where did = $2 and rkey = $3`, + where repo_did = $2`, state, - did, rkey, + repoDid, ); err != nil { return fmt.Errorf("updating repo: %w", err) } @@ -78,31 +72,29 @@ func DeleteRepo(ctx context.Context, e *sql.DB, did syntax.DID, rkey syntax.Reco return nil } -func GetRepoByName(ctx context.Context, e *sql.DB, did syntax.DID, name string) (*models.Repo, error) { +const repoColumns = ` + did, + rkey, + cid, + name, + knot_domain, + repo_did, + git_rev, + repo_sha, + state, + error_msg, + retry_count, + retry_after` + +func scanRepo(row interface{ Scan(...any) error }) (*models.Repo, error) { var repo models.Repo - if err := e.QueryRowContext(ctx, - `select - did, - rkey, - cid, - name, - knot_domain, - git_rev, - repo_sha, - state, - error_msg, - retry_count, - retry_after - from repos - where did = $1 and name = $2`, - did, - name, - ).Scan( + if err := row.Scan( &repo.Did, &repo.Rkey, &repo.Cid, &repo.Name, &repo.KnotDomain, + &repo.RepoDid, &repo.GitRev, &repo.RepoSha, &repo.State, @@ -110,51 +102,43 @@ func GetRepoByName(ctx context.Context, e *sql.DB, did syntax.DID, name string) &repo.RetryCount, &repo.RetryAfter, ); err != nil { + return nil, err + } + return &repo, nil +} + +func GetRepoByRepoDid(ctx context.Context, e *sql.DB, repoDid syntax.DID) (*models.Repo, error) { + row := e.QueryRowContext(ctx, + `select`+repoColumns+` + from repos + where repo_did = $1`, + repoDid, + ) + repo, err := scanRepo(row) + if err != nil { if errors.Is(err, sql.ErrNoRows) { return nil, nil } return nil, fmt.Errorf("querying repo: %w", err) } - return &repo, nil + return repo, nil } func GetRepoByAtUri(ctx context.Context, e *sql.DB, aturi syntax.ATURI) (*models.Repo, error) { - var repo models.Repo - if err := e.QueryRowContext(ctx, - `select - did, - rkey, - cid, - name, - knot_domain, - git_rev, - repo_sha, - state, - error_msg, - retry_count, - retry_after + row := e.QueryRowContext(ctx, + `select`+repoColumns+` from repos where at_uri = $1`, aturi, - ).Scan( - &repo.Did, - &repo.Rkey, - &repo.Cid, - &repo.Name, - &repo.KnotDomain, - &repo.GitRev, - &repo.RepoSha, - &repo.State, - &repo.ErrorMsg, - &repo.RetryCount, - &repo.RetryAfter, - ); err != nil { + ) + repo, err := scanRepo(row) + if err != nil { if errors.Is(err, sql.ErrNoRows) { return nil, nil } return nil, fmt.Errorf("querying repo: %w", err) } - return &repo, nil + return repo, nil } func ListRepos(ctx context.Context, e *sql.DB, page pagination.Page, did, knot, state string) ([]models.Repo, error) { @@ -188,18 +172,7 @@ func ListRepos(ctx context.Context, e *sql.DB, page pagination.Page, did, knot, } query := ` - select - did, - rkey, - cid, - name, - knot_domain, - git_rev, - repo_sha, - state, - error_msg, - retry_count, - retry_after + select` + repoColumns + ` from repos ` + whereClause + pageClause rows, err := e.QueryContext(ctx, query, args...) @@ -210,23 +183,11 @@ func ListRepos(ctx context.Context, e *sql.DB, page pagination.Page, did, knot, var repos []models.Repo for rows.Next() { - var repo models.Repo - if err := rows.Scan( - &repo.Did, - &repo.Rkey, - &repo.Cid, - &repo.Name, - &repo.KnotDomain, - &repo.GitRev, - &repo.RepoSha, - &repo.State, - &repo.ErrorMsg, - &repo.RetryCount, - &repo.RetryAfter, - ); err != nil { + repo, err := scanRepo(rows) + if err != nil { return nil, fmt.Errorf("scanning row: %w", err) } - repos = append(repos, repo) + repos = append(repos, *repo) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("scanning rows: %w ", err) diff --git a/knotmirror/git.go b/knotmirror/git.go index 6bee91d0..b6081be8 100644 --- a/knotmirror/git.go +++ b/knotmirror/git.go @@ -19,8 +19,6 @@ import ( type GitMirrorManager interface { Exist(repo *models.Repo) (bool, error) - // RemoteSetUrl updates git repository 'origin' remote - RemoteSetUrl(ctx context.Context, repo *models.Repo) error // Clone clones the repository as a mirror Clone(ctx context.Context, repo *models.Repo) error // Fetch fetches the repository @@ -44,33 +42,16 @@ func NewCliGitMirrorManager(repoBasePath string, knotUseSSL bool) *CliGitMirrorM var _ GitMirrorManager = new(CliGitMirrorManager) func (c *CliGitMirrorManager) makeRepoPath(repo *models.Repo) string { - return filepath.Join(c.repoBasePath, repo.Did.String(), repo.Rkey.String()) + return filepath.Join(c.repoBasePath, repo.RepoDid.String()) } func (c *CliGitMirrorManager) Exist(repo *models.Repo) (bool, error) { return isDir(c.makeRepoPath(repo)) } -func (c *CliGitMirrorManager) RemoteSetUrl(ctx context.Context, repo *models.Repo) error { - path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) - if err != nil { - return fmt.Errorf("constructing repo remote url: %w", err) - } - cmd := exec.CommandContext(ctx, "git", "-C", path, "remote", "set-url", "origin", url) - if out, err := cmd.CombinedOutput(); err != nil { - if ctx.Err() != nil { - return ctx.Err() - } - msg := string(out) - return fmt.Errorf("running 'git remote set-url origin %s': %w\n%s", url, err, msg) - } - return nil -} - func (c *CliGitMirrorManager) Clone(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) if err != nil { return fmt.Errorf("constructing repo remote url: %w", err) } @@ -94,12 +75,15 @@ func (c *CliGitMirrorManager) clone(ctx context.Context, path, url string) error func (c *CliGitMirrorManager) Fetch(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - return c.fetch(ctx, path) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) + if err != nil { + return fmt.Errorf("constructing repo remote url: %w", err) + } + return c.fetch(ctx, path, url) } -func (c *CliGitMirrorManager) fetch(ctx context.Context, path string) error { - // TODO: use `repo.Knot` instead of depending on origin - cmd := exec.CommandContext(ctx, "git", "-C", path, "fetch", "--prune", "origin") +func (c *CliGitMirrorManager) fetch(ctx context.Context, path, url string) error { + cmd := exec.CommandContext(ctx, "git", "-C", path, "fetch", "--prune", url, "+refs/*:refs/*") if out, err := cmd.CombinedOutput(); err != nil { if ctx.Err() != nil { return ctx.Err() @@ -111,7 +95,7 @@ func (c *CliGitMirrorManager) fetch(ctx context.Context, path string) error { func (c *CliGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) if err != nil { return fmt.Errorf("constructing repo remote url: %w", err) } @@ -125,7 +109,7 @@ func (c *CliGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error return fmt.Errorf("cloning repo: %w", err) } } else { - if err := c.fetch(ctx, path); err != nil { + if err := c.fetch(ctx, path, url); err != nil { return fmt.Errorf("fetching repo: %w", err) } } @@ -191,20 +175,16 @@ func NewGoGitMirrorClient(repoBasePath string, knotUseSSL bool) *GoGitMirrorMana var _ GitMirrorManager = new(GoGitMirrorManager) func (c *GoGitMirrorManager) makeRepoPath(repo *models.Repo) string { - return filepath.Join(c.repoBasePath, repo.Did.String(), repo.Rkey.String()) + return filepath.Join(c.repoBasePath, repo.RepoDid.String()) } func (c *GoGitMirrorManager) Exist(repo *models.Repo) (bool, error) { return isDir(c.makeRepoPath(repo)) } -func (c *GoGitMirrorManager) RemoteSetUrl(ctx context.Context, repo *models.Repo) error { - panic("unimplemented") -} - func (c *GoGitMirrorManager) Clone(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) if err != nil { return fmt.Errorf("constructing repo remote url: %w", err) } @@ -224,7 +204,7 @@ func (c *GoGitMirrorManager) clone(ctx context.Context, path, url string) error func (c *GoGitMirrorManager) Fetch(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) if err != nil { return fmt.Errorf("constructing repo remote url: %w", err) } @@ -250,7 +230,7 @@ func (c *GoGitMirrorManager) fetch(ctx context.Context, path, url string) error func (c *GoGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error { path := c.makeRepoPath(repo) - url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), c.knotUseSSL) + url, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), c.knotUseSSL) if err != nil { return fmt.Errorf("constructing repo remote url: %w", err) } @@ -271,7 +251,7 @@ func (c *GoGitMirrorManager) Sync(ctx context.Context, repo *models.Repo) error return nil } -func makeRepoRemoteUrl(knot, didSlashRepo string, knotUseSSL bool) (string, error) { +func makeRepoRemoteUrl(knot, repoIdentifier string, knotUseSSL bool) (string, error) { if !strings.Contains(knot, "://") { if knotUseSSL { knot = "https://" + knot @@ -289,7 +269,7 @@ func makeRepoRemoteUrl(knot, didSlashRepo string, knotUseSSL bool) (string, erro return "", fmt.Errorf("unsupported scheme: %s", u.Scheme) } - u = u.JoinPath(didSlashRepo) + u = u.JoinPath(repoIdentifier) return u.String(), nil } diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go index b28ca1dc..10e50c84 100644 --- a/knotmirror/knotstream/slurper.go +++ b/knotmirror/knotstream/slurper.go @@ -15,7 +15,6 @@ import ( "github.com/bluesky-social/indigo/util/ssrf" "github.com/carlmjohnson/versioninfo" "github.com/gorilla/websocket" - "tangled.org/core/api/tangled" "tangled.org/core/knotmirror/config" "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/models" @@ -262,10 +261,15 @@ func (s *KnotSlurper) handleConnection(ctx context.Context, conn *websocket.Conn } } +type legacyGitRefUpdate struct { + OwnerDid *string `json:"ownerDid,omitempty"` + RepoDid *string `json:"repo,omitempty"` +} + type LegacyGitEvent struct { Rkey string Nsid string - Event tangled.GitRefUpdate + Event legacyGitRefUpdate } func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { @@ -280,33 +284,36 @@ func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { return nil } +// lookupRepoForRefUpdate resolves the local repo row for an incoming refUpdate +// via the stable RepoDid join. Returns (nil, "", nil) when the event has no +// repoDid (unjoinable) and (nil, key, nil) on a clean miss. +func (s *KnotSlurper) lookupRepoForRefUpdate(ctx context.Context, evt *LegacyGitEvent) (*models.Repo, string, error) { + if evt.Event.RepoDid == nil || *evt.Event.RepoDid == "" { + return nil, "", nil + } + repoDid := syntax.DID(*evt.Event.RepoDid) + curr, err := db.GetRepoByRepoDid(ctx, s.db, repoDid) + return curr, repoDid.String(), err +} + func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source string, evt *LegacyGitEvent) error { knotstreamEventsReceived.Inc() l := s.logger.With("src", source) - ownerDid := "" - if evt.Event.OwnerDid != nil { - ownerDid = *evt.Event.OwnerDid - } else { - // handle legacy event - if evt.Event.RepoDid != nil { - ownerDid = *evt.Event.RepoDid - } - } - curr, err := db.GetRepoByName(ctx, s.db, syntax.DID(ownerDid), evt.Event.RepoName) + curr, lookupKey, err := s.lookupRepoForRefUpdate(ctx, evt) if err != nil { - return fmt.Errorf("failed to get repo '%s': %w", ownerDid+"/"+evt.Event.RepoName, err) + return fmt.Errorf("failed to get repo '%s': %w", lookupKey, err) } if curr == nil { - // if repo doesn't exist in DB, just ignore the event. That repo is unknown. - // - // Normally did+name is already enough to perform git-fetch as that's - // what needed to fetch the repository. - // But we want to store that in did/rkey in knot-mirror. - // Therefore, we should ignore when the repository is unknown. - // Hopefully crawler will sync it later. - l.Warn("skipping event from unknown repo", "did/name", ownerDid+"/"+evt.Event.RepoName) + if lookupKey == "" { + l.Warn("skipping gitRefUpdate: event has no fields to join on", + "repo_did", evt.Event.RepoDid) + } else { + // if repo doesn't exist in DB, just ignore the event. That repo is unknown. + // Hopefully crawler/tap will sync it later. + l.Warn("skipping event from unknown repo", "key", lookupKey) + } knotstreamEventsSkipped.Inc() return nil } @@ -325,13 +332,8 @@ func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source stri return nil } - // if curr.State == models.RepoStateResyncing { - // firehoseEventsSkipped.Inc() - // return fp.events.addToResyncBuffer(ctx, commit) - // } - // can't skip anything, update repo state - if err := db.UpdateRepoState(ctx, s.db, curr.Did, curr.Rkey, models.RepoStateDesynchronized); err != nil { + if err := db.UpdateRepoState(ctx, s.db, curr.RepoDid, models.RepoStateDesynchronized); err != nil { return err } diff --git a/knotmirror/migrate/migrate.go b/knotmirror/migrate/migrate.go new file mode 100644 index 00000000..2fed459c --- /dev/null +++ b/knotmirror/migrate/migrate.go @@ -0,0 +1,135 @@ +package migrate + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log/slog" + "os" + "path/filepath" + "strings" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/db" +) + +type Stats struct { + Renamed int + Skipped int + Orphaned int + OwnerDirsRm int + AlreadyExists int +} + +func (s Stats) String() string { + return fmt.Sprintf( + "renamed=%d skipped=%d orphaned=%d owner_dirs_removed=%d target_existed=%d", + s.Renamed, s.Skipped, s.Orphaned, s.OwnerDirsRm, s.AlreadyExists, + ) +} + +func RenameDisk(ctx context.Context, base string, database *sql.DB, logger *slog.Logger) (Stats, error) { + entries, err := os.ReadDir(base) + if err != nil { + return Stats{}, fmt.Errorf("reading base path: %w", err) + } + return reduceEntries(ctx, entries, 0, Stats{}, ownerStep(base, database, logger)) +} + +type stepFn func(context.Context, os.DirEntry, Stats) (Stats, error) + +func reduceEntries(ctx context.Context, entries []os.DirEntry, idx int, acc Stats, fn stepFn) (Stats, error) { + if idx >= len(entries) { + return acc, nil + } + if err := ctx.Err(); err != nil { + return acc, err + } + next, err := fn(ctx, entries[idx], acc) + if err != nil { + return next, err + } + return reduceEntries(ctx, entries, idx+1, next, fn) +} + +func ownerStep(base string, database *sql.DB, logger *slog.Logger) stepFn { + return func(ctx context.Context, entry os.DirEntry, acc Stats) (Stats, error) { + if !entry.IsDir() || !strings.HasPrefix(entry.Name(), "did:") { + return acc, nil + } + ownerPath := filepath.Join(base, entry.Name()) + if _, err := os.Stat(filepath.Join(ownerPath, "HEAD")); err == nil { + return acc, nil + } + subEntries, err := os.ReadDir(ownerPath) + if err != nil { + logger.Error("reading owner dir", "ownerPath", ownerPath, "err", err) + return acc, nil + } + next, err := reduceEntries(ctx, subEntries, 0, acc, rkeyStep(base, database, logger, syntax.DID(entry.Name()), ownerPath)) + if err != nil { + return next, err + } + remaining, err := os.ReadDir(ownerPath) + if err == nil && len(remaining) == 0 { + if rmErr := os.Remove(ownerPath); rmErr == nil { + next.OwnerDirsRm++ + logger.Info("removed empty owner dir", "ownerPath", ownerPath) + } else { + logger.Warn("failed to remove empty owner dir", "ownerPath", ownerPath, "err", rmErr) + } + } + return next, nil + } +} + +func rkeyStep(base string, database *sql.DB, logger *slog.Logger, ownerDid syntax.DID, ownerPath string) stepFn { + return func(ctx context.Context, sub os.DirEntry, acc Stats) (Stats, error) { + if !sub.IsDir() { + return acc, nil + } + rkey := sub.Name() + subPath := filepath.Join(ownerPath, rkey) + l := logger.With("did", ownerDid, "rkey", rkey, "subPath", subPath) + + if _, err := os.Stat(filepath.Join(subPath, "HEAD")); err != nil { + l.Warn("skipping non-repo subdir") + acc.Skipped++ + return acc, nil + } + + aturi := syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", ownerDid, tangled.RepoNSID, rkey)) + repo, err := db.GetRepoByAtUri(ctx, database, aturi) + if err != nil { + return acc, fmt.Errorf("looking up repo by aturi %s: %w", aturi, err) + } + if repo == nil { + l.Warn("orphan disk repo, no DB row; leaving in place") + acc.Orphaned++ + return acc, nil + } + if repo.RepoDid == "" { + l.Warn("DB row has empty repo_did; leaving in place") + acc.Orphaned++ + return acc, nil + } + + target := filepath.Join(base, repo.RepoDid.String()) + if _, err := os.Stat(target); err == nil { + l.Warn("target path already exists; leaving source in place", "target", target) + acc.AlreadyExists++ + return acc, nil + } else if !errors.Is(err, os.ErrNotExist) { + return acc, fmt.Errorf("stat target %s: %w", target, err) + } + + if err := os.Rename(subPath, target); err != nil { + return acc, fmt.Errorf("rename %s -> %s: %w", subPath, target, err) + } + acc.Renamed++ + l.Info("renamed", "target", target) + return acc, nil + } +} diff --git a/knotmirror/models/models.go b/knotmirror/models/models.go index 8dbc27e9..088fb700 100644 --- a/knotmirror/models/models.go +++ b/knotmirror/models/models.go @@ -14,6 +14,7 @@ type Repo struct { // content of tangled.Repo Name string KnotDomain string + RepoDid syntax.DID GitRev syntax.TID // last processed git.refUpdate revision RepoSha string // sha256 sum of git refs (to avoid no-op git fetch) @@ -27,8 +28,8 @@ func (r *Repo) AtUri() syntax.ATURI { return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", r.Did, tangled.RepoNSID, r.Rkey)) } -func (r *Repo) DidSlashRepo() string { - return fmt.Sprintf("%s/%s", r.Did, r.Name) +func (r *Repo) RepoIdentifier() string { + return r.RepoDid.String() } type RepoState string diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 97023799..594ef5ad 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -278,7 +278,7 @@ func isRateLimitError(err error) bool { // checkKnotReachability checks if Knot is reachable and is valid git remote server func (r *Resyncer) checkKnotReachability(ctx context.Context, repo *models.Repo) error { - repoUrl, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), r.cfg.KnotUseSSL) + repoUrl, err := makeRepoRemoteUrl(repo.KnotDomain, repo.RepoIdentifier(), r.cfg.KnotUseSSL) if err != nil { return err } diff --git a/knotmirror/tapclient.go b/knotmirror/tapclient.go index 30d9237e..190d08b2 100644 --- a/knotmirror/tapclient.go +++ b/knotmirror/tapclient.go @@ -11,6 +11,7 @@ import ( "strings" "time" + "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" "tangled.org/core/knotmirror/config" "tangled.org/core/knotmirror/db" @@ -104,32 +105,23 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error errMsg = "suspending non-public knot" } + if record.RepoDid == nil || *record.RepoDid == "" { + t.logger.Warn("dropping repo record without repo_did", "did", evt.Did, "rkey", evt.Rkey) + return nil + } repo := &models.Repo{ Did: evt.Did, Rkey: evt.Rkey, Cid: evt.CID, - Name: record.Name, + Name: evt.Rkey.String(), KnotDomain: knotUrl, + RepoDid: syntax.DID(*record.RepoDid), State: status, ErrorMsg: errMsg, RetryAfter: 0, // clear retry info RetryCount: 0, } - if evt.Action == tapc.RecordUpdateAction { - exist, err := t.gitm.Exist(repo) - if err != nil { - return fmt.Errorf("checking git repo existence: %w", err) - } - if exist { - // update git repo remote url - if err := t.gitm.RemoteSetUrl(ctx, repo); err != nil { - return fmt.Errorf("updating git repo remote url: %w", err) - } - } - } - - t.logger.Debug("tap: upserting repo with knot", "knot", repo.KnotDomain) if err := db.UpsertRepo(ctx, t.db, repo); err != nil { return fmt.Errorf("upserting repo to db: %w", err) } diff --git a/knotmirror/xrpc/git_list_branches.go b/knotmirror/xrpc/git_list_branches.go index 28604972..28158911 100644 --- a/knotmirror/xrpc/git_list_branches.go +++ b/knotmirror/xrpc/git_list_branches.go @@ -9,6 +9,7 @@ import ( "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/knotmirror/db" "tangled.org/core/knotserver/git" "tangled.org/core/types" ) @@ -82,14 +83,15 @@ func (x *Xrpc) listBranches(ctx context.Context, repo syntax.ATURI, limit int, c } func (x *Xrpc) makeRepoPath(ctx context.Context, repo syntax.ATURI) (string, error) { - id, err := x.resolver.ResolveIdent(ctx, repo.Authority().String()) + r, err := db.GetRepoByAtUri(ctx, x.db, repo) if err != nil { - return "", err + return "", fmt.Errorf("looking up repo: %w", err) } - - return filepath.Join( - x.cfg.GitRepoBasePath, - id.DID.String(), - repo.RecordKey().String(), - ), nil + if r == nil { + return "", fmt.Errorf("repo not found: %s", repo) + } + if r.RepoDid == "" { + return "", fmt.Errorf("repo missing repo_did: %s", repo) + } + return filepath.Join(x.cfg.GitRepoBasePath, r.RepoDid.String()), nil } diff --git a/knotmirror/xrpc/proxy.go b/knotmirror/xrpc/proxy.go index 8694c5c8..e167c0b2 100644 --- a/knotmirror/xrpc/proxy.go +++ b/knotmirror/xrpc/proxy.go @@ -13,6 +13,7 @@ import ( indigoxrpc "github.com/bluesky-social/indigo/xrpc" "tangled.org/core/api/tangled" "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/models" ) var mirrorToKnotNSID = map[string]string{ @@ -40,8 +41,8 @@ var hopByHopHeaders = map[string]bool{ } type knotInfo struct { - baseURL string - didSlashRepo string + baseURL string + repoIdentifier string } func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, error) { @@ -60,7 +61,7 @@ func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, } } } - return &knotInfo{baseURL: knotURL, didSlashRepo: repo.DidSlashRepo()}, nil + return &knotInfo{baseURL: knotURL, repoIdentifier: repo.RepoIdentifier()}, nil } owner, err := x.resolver.ResolveIdent(ctx, repoAt.Authority().String()) @@ -75,6 +76,9 @@ func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, } record := out.Value.Val.(*tangled.Repo) + if record.RepoDid == nil || *record.RepoDid == "" { + return nil, fmt.Errorf("repo record has no repo_did") + } knotURL := record.Knot if !strings.Contains(record.Knot, "://") { if host, _ := db.GetHost(ctx, x.db, record.Knot); host != nil { @@ -89,9 +93,27 @@ func (x *Xrpc) resolveKnot(ctx context.Context, repoAt syntax.ATURI) (*knotInfo, } } + rkey := repoAt.RecordKey().String() + repoDid := syntax.DID(*record.RepoDid) + go func() { + bgCtx := context.Background() + pending := &models.Repo{ + Did: owner.DID, + Rkey: repoAt.RecordKey(), + Cid: (*syntax.CID)(out.Cid), + Name: rkey, + KnotDomain: knotURL, + RepoDid: repoDid, + State: models.RepoStatePending, + } + if upsertErr := db.UpsertRepo(bgCtx, x.db, pending); upsertErr != nil { + x.logger.Error("failed to upsert repo after proxy resolution", "err", upsertErr) + } + }() + return &knotInfo{ - baseURL: knotURL, - didSlashRepo: fmt.Sprintf("%s/%s", owner.DID, record.Name), + baseURL: knotURL, + repoIdentifier: repoDid.String(), }, nil } @@ -112,7 +134,7 @@ func (x *Xrpc) proxyToKnot(w http.ResponseWriter, r *http.Request, repoAt syntax for k, v := range r.URL.Query() { params[k] = v } - params.Set("repo", knot.didSlashRepo) + params.Set("repo", knot.repoIdentifier) target := fmt.Sprintf("%s/xrpc/%s?%s", knot.baseURL, knotNSID, params.Encode()) diff --git a/knotmirror/xrpc/sync_request_crawl.go b/knotmirror/xrpc/sync_request_crawl.go index 61978dba..900b8beb 100644 --- a/knotmirror/xrpc/sync_request_crawl.go +++ b/knotmirror/xrpc/sync_request_crawl.go @@ -71,12 +71,19 @@ func (x *Xrpc) RequestCrawl(w http.ResponseWriter, r *http.Request) { } } + if record.RepoDid == nil || *record.RepoDid == "" { + l.Warn("dropping repo crawl request without repo_did", "did", owner.DID, "rkey", repoAt.RecordKey()) + writeErr(w, fmt.Errorf("repo record missing repo_did")) + return + } + repo := &models.Repo{ Did: owner.DID, Rkey: repoAt.RecordKey(), Cid: (*syntax.CID)(out.Cid), - Name: record.Name, + Name: repoAt.RecordKey().String(), KnotDomain: knotUrl, + RepoDid: syntax.DID(*record.RepoDid), State: models.RepoStatePending, ErrorMsg: "", RetryAfter: 0,