From 07299d8677c1a511bf6469d3dc7770f84920313d Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Sun, 31 May 2026 06:57:02 +0000 Subject: [PATCH] knotmirror: use repoDid for resync identifier Signed-off-by: Seongmin Lee --- knotmirror/adminpage.go | 8 ++++---- knotmirror/resyncer.go | 60 ++++++++++++++++++++++++++++++------------------------------ knotmirror/tapclient.go | 14 +------------- knotmirror/templates/repos.html | 2 +- 4 file(s) changed, 36 insertion(s)(+), 48 deletion(s)(-) diff --git a/knotmirror/adminpage.go b/knotmirror/adminpage.go --- a/knotmirror/adminpage.go +++ b/knotmirror/adminpage.go @@ -182,8 +182,8 @@ func (s *AdminServer) handleRepoResyncTrigger() http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var repoQuery = r.FormValue("repo") - repo, err := syntax.ParseATURI(repoQuery) - if err != nil || repo.RecordKey() == "" { + repo, err := syntax.ParseDID(repoQuery) + if err != nil { writeNotif(w, http.StatusBadRequest, fmt.Sprintf("repo parameter invalid: %s", repoQuery)) return } @@ -201,8 +201,8 @@ func (s *AdminServer) handleRepoResyncCancel() http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { var repoQuery = r.FormValue("repo") - repo, err := syntax.ParseATURI(repoQuery) - if err != nil || repo.RecordKey() == "" { + repo, err := syntax.ParseDID(repoQuery) + if err != nil { writeNotif(w, http.StatusBadRequest, fmt.Sprintf("repo parameter invalid: %s", repoQuery)) return } diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -30,7 +30,7 @@ indexer *knotstream.ParallelScheduler claimJobMu sync.Mutex - runningJobs map[syntax.ATURI]context.CancelFunc + runningJobs map[syntax.DID]context.CancelFunc runningJobsMu sync.Mutex repoFetchTimeout time.Duration @@ -51,7 +51,7 @@ gitm: gitm, cfg: cfg, indexer: indexer, - runningJobs: make(map[syntax.ATURI]context.CancelFunc), + runningJobs: make(map[syntax.DID]context.CancelFunc), repoFetchTimeout: cfg.GitRepoFetchTimeout, manualResyncTimeout: 30 * time.Minute, @@ -78,7 +78,7 @@ l.Info("resync worker shutting down", "error", ctx.Err()) return default: } - repoAt, found, err := r.claimResyncJob(ctx) + repoDid, found, err := r.claimResyncJob(ctx) if err != nil { l.Error("failed to claim resync job", "error", err) time.Sleep(time.Second) @@ -88,14 +88,14 @@ if !found { time.Sleep(time.Second) continue } - l.Info("processing resync", "aturi", repoAt) - if err := r.resyncRepo(ctx, repoAt); err != nil { - l.Error("resync failed", "aturi", repoAt, "error", err) + l.Info("processing resync", "did", repoDid) + if err := r.resyncRepo(ctx, repoDid); err != nil { + l.Error("resync failed", "did", repoDid, "error", err) } } } -func (r *Resyncer) registerRunning(repo syntax.ATURI, cancel context.CancelFunc) { +func (r *Resyncer) registerRunning(repo syntax.DID, cancel context.CancelFunc) { r.runningJobsMu.Lock() defer r.runningJobsMu.Unlock() @@ -105,14 +105,14 @@ } r.runningJobs[repo] = cancel } -func (r *Resyncer) unregisterRunning(repo syntax.ATURI) { +func (r *Resyncer) unregisterRunning(repo syntax.DID) { r.runningJobsMu.Lock() defer r.runningJobsMu.Unlock() delete(r.runningJobs, repo) } -func (r *Resyncer) CancelResyncJob(repo syntax.ATURI) { +func (r *Resyncer) CancelResyncJob(repo syntax.DID) { r.runningJobsMu.Lock() defer r.runningJobsMu.Unlock() @@ -125,13 +125,13 @@ cancel() } // TriggerResyncJob manually triggers the resync job -func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoAt syntax.ATURI) error { - repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) +func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoDid syntax.DID) error { + repo, err := db.GetRepoByRepoDid(ctx, r.db, repoDid) if err != nil { return fmt.Errorf("failed to get repo: %w", err) } if repo == nil { - return fmt.Errorf("repo not found: %s", repoAt) + return fmt.Errorf("repo not found: %s", repoDid) } if repo.State == models.RepoStateResyncing { @@ -147,18 +147,18 @@ } return nil } -func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.ATURI, bool, error) { +func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.DID, bool, error) { // use mutex to prevent duplicated jobs r.claimJobMu.Lock() defer r.claimJobMu.Unlock() - var repoAt syntax.ATURI + var repoDid syntax.DID now := time.Now().Unix() if err := r.db.QueryRowContext(ctx, `update repos set state = $1 - where at_uri = ( - select at_uri from repos + where repo_did = ( + select repo_did from repos where state in ($2, $3, $4) and (retry_after = -1 or retry_after = 0 or retry_after < $5) order by @@ -167,22 +167,22 @@ (retry_after = 0) desc, retry_after limit 1 ) - returning at_uri + returning repo_did `, models.RepoStateResyncing, models.RepoStatePending, models.RepoStateDesynchronized, models.RepoStateError, now, - ).Scan(&repoAt); err != nil { + ).Scan(&repoDid); err != nil { if errors.Is(err, sql.ErrNoRows) { return "", false, nil } return "", false, err } - return repoAt, true, nil + return repoDid, true, nil } -func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error { +func (r *Resyncer) resyncRepo(ctx context.Context, repoDid syntax.DID) error { // ctx, span := tracer.Start(ctx, "resyncRepo") // span.SetAttributes(attribute.String("aturi", repoAt)) // defer span.End() @@ -191,14 +191,14 @@ resyncsStarted.Inc() startTime := time.Now() jobCtx, cancel := context.WithCancel(ctx) - r.registerRunning(repoAt, cancel) - defer r.unregisterRunning(repoAt) + r.registerRunning(repoDid, cancel) + defer r.unregisterRunning(repoDid) - success, err := r.doResync(jobCtx, repoAt) + success, err := r.doResync(jobCtx, repoDid) if !success { resyncsFailed.Inc() resyncDuration.Observe(time.Since(startTime).Seconds()) - return r.handleResyncFailure(ctx, repoAt, err) + return r.handleResyncFailure(ctx, repoDid, err) } resyncsCompleted.Inc() @@ -206,12 +206,12 @@ resyncDuration.Observe(time.Since(startTime).Seconds()) return nil } -func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, error) { +func (r *Resyncer) doResync(ctx context.Context, repoDid syntax.DID) (bool, error) { // ctx, span := tracer.Start(ctx, "doResync") // span.SetAttributes(attribute.String("aturi", repoAt)) // defer span.End() - repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) + repo, err := db.GetRepoByRepoDid(ctx, r.db, repoDid) if err != nil { return false, fmt.Errorf("failed to get repo: %w", err) } @@ -323,8 +323,8 @@ return nil } -func (r *Resyncer) handleResyncFailure(ctx context.Context, repoAt syntax.ATURI, err error) error { - r.logger.Debug("handleResyncFailure", "at_uri", repoAt, "err", err) +func (r *Resyncer) handleResyncFailure(ctx context.Context, repoDid syntax.DID, err error) error { + r.logger.Debug("handleResyncFailure", "at_uri", repoDid, "err", err) var state models.RepoState var errMsg string if err == nil { @@ -335,12 +335,12 @@ state = models.RepoStateError errMsg = err.Error() } - repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt) + repo, err := db.GetRepoByRepoDid(ctx, r.db, repoDid) if err != nil { return fmt.Errorf("failed to get repo: %w", err) } if repo == nil { - return fmt.Errorf("failed to get repo. repo '%s' doesn't exist in db", repoAt) + return fmt.Errorf("failed to get repo. repo '%s' doesn't exist in db", repoDid) } // start a 1 min & go up to 1 hr between retries diff --git a/knotmirror/tapclient.go b/knotmirror/tapclient.go --- a/knotmirror/tapclient.go +++ b/knotmirror/tapclient.go @@ -133,19 +133,7 @@ } } case tapc.RecordDeleteAction: - aturi := syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", evt.Did, tangled.RepoNSID, evt.Rkey)) - repo, err := db.GetRepoByAtUri(ctx, t.db, aturi) - if err != nil { - return fmt.Errorf("looking up repo before delete: %w", err) - } - if repo != nil { - if err := t.gitm.Delete(repo); err != nil { - return fmt.Errorf("removing mirror dir: %w", err) - } - } - if err := db.DeleteRepo(ctx, t.db, evt.Did, evt.Rkey); err != nil { - return fmt.Errorf("deleting repo from db: %w", err) - } + // no-op. deletion of sh.tangled.repo record doesn't mean repository deletion } return nil } diff --git a/knotmirror/templates/repos.html b/knotmirror/templates/repos.html --- a/knotmirror/templates/repos.html +++ b/knotmirror/templates/repos.html @@ -65,7 +65,7 @@ {{- end }} hx-swap="none" hx-disabled-elt="find button" > - + -- tangled.sh