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"
>
-
+