This repository has no description
9.1 kB
356 lines
1package knotmirror
2
3import (
4 "context"
5 "database/sql"
6 "errors"
7 "fmt"
8 "log/slog"
9 "math/rand"
10 "net/http"
11 "net/url"
12 "strings"
13 "sync"
14 "time"
15
16 "github.com/bluesky-social/indigo/atproto/syntax"
17 "tangled.org/core/knotmirror/config"
18 "tangled.org/core/knotmirror/db"
19 "tangled.org/core/knotmirror/models"
20 "tangled.org/core/log"
21)
22
23type Resyncer struct {
24 logger *slog.Logger
25 db *sql.DB
26 gitm GitMirrorManager
27
28 claimJobMu sync.Mutex
29
30 runningJobs map[syntax.ATURI]context.CancelFunc
31 runningJobsMu sync.Mutex
32
33 repoFetchTimeout time.Duration
34 manualResyncTimeout time.Duration
35 parallelism int
36
37 knotBackoff map[string]time.Time
38 knotBackoffMu sync.RWMutex
39}
40
41func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, cfg *config.Config) *Resyncer {
42 return &Resyncer{
43 logger: log.SubLogger(l, "resyncer"),
44 db: db,
45 gitm: gitm,
46
47 runningJobs: make(map[syntax.ATURI]context.CancelFunc),
48
49 repoFetchTimeout: cfg.GitRepoFetchTimeout,
50 manualResyncTimeout: 30 * time.Minute,
51 parallelism: cfg.ResyncParallelism,
52
53 knotBackoff: make(map[string]time.Time),
54 }
55}
56
57func (r *Resyncer) Start(ctx context.Context) {
58 for i := 0; i < r.parallelism; i++ {
59 go r.runResyncWorker(ctx, i)
60 }
61}
62
63func (r *Resyncer) runResyncWorker(ctx context.Context, workerID int) {
64 l := r.logger.With("worker", workerID)
65 for {
66 select {
67 case <-ctx.Done():
68 l.Info("resync worker shutting down", "error", ctx.Err())
69 return
70 default:
71 }
72 repoAt, found, err := r.claimResyncJob(ctx)
73 if err != nil {
74 l.Error("failed to claim resync job", "error", err)
75 time.Sleep(time.Second)
76 continue
77 }
78 if !found {
79 time.Sleep(time.Second)
80 continue
81 }
82 l.Info("processing resync", "aturi", repoAt)
83 if err := r.resyncRepo(ctx, repoAt); err != nil {
84 l.Error("resync failed", "aturi", repoAt, "error", err)
85 }
86 }
87}
88
89func (r *Resyncer) registerRunning(repo syntax.ATURI, cancel context.CancelFunc) {
90 r.runningJobsMu.Lock()
91 defer r.runningJobsMu.Unlock()
92
93 if _, exists := r.runningJobs[repo]; exists {
94 return
95 }
96 r.runningJobs[repo] = cancel
97}
98
99func (r *Resyncer) unregisterRunning(repo syntax.ATURI) {
100 r.runningJobsMu.Lock()
101 defer r.runningJobsMu.Unlock()
102
103 delete(r.runningJobs, repo)
104}
105
106func (r *Resyncer) CancelResyncJob(repo syntax.ATURI) {
107 r.runningJobsMu.Lock()
108 defer r.runningJobsMu.Unlock()
109
110 cancel, ok := r.runningJobs[repo]
111 if !ok {
112 return
113 }
114 delete(r.runningJobs, repo)
115 cancel()
116}
117
118// TriggerResyncJob manually triggers the resync job
119func (r *Resyncer) TriggerResyncJob(ctx context.Context, repoAt syntax.ATURI) error {
120 repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt)
121 if err != nil {
122 return fmt.Errorf("failed to get repo: %w", err)
123 }
124 if repo == nil {
125 return fmt.Errorf("repo not found: %s", repoAt)
126 }
127
128 if repo.State == models.RepoStateResyncing {
129 return fmt.Errorf("repo already resyncing")
130 }
131
132 repo.State = models.RepoStatePending
133 repo.RetryAfter = -1 // resyncer will prioritize this
134
135 if err := db.UpsertRepo(ctx, r.db, repo); err != nil {
136 return fmt.Errorf("updating repo state to pending %w", err)
137 }
138 return nil
139}
140
141func (r *Resyncer) claimResyncJob(ctx context.Context) (syntax.ATURI, bool, error) {
142 // use mutex to prevent duplicated jobs
143 r.claimJobMu.Lock()
144 defer r.claimJobMu.Unlock()
145
146 var repoAt syntax.ATURI
147 now := time.Now().Unix()
148 if err := r.db.QueryRowContext(ctx,
149 `update repos
150 set state = $1
151 where at_uri = (
152 select at_uri from repos
153 where state in ($2, $3, $4)
154 and (retry_after = -1 or retry_after = 0 or retry_after < $5)
155 order by
156 (retry_after = -1) desc,
157 (retry_after = 0) desc,
158 retry_after
159 limit 1
160 )
161 returning at_uri
162 `,
163 models.RepoStateResyncing,
164 models.RepoStatePending, models.RepoStateDesynchronized, models.RepoStateError,
165 now,
166 ).Scan(&repoAt); err != nil {
167 if errors.Is(err, sql.ErrNoRows) {
168 return "", false, nil
169 }
170 return "", false, err
171 }
172
173 return repoAt, true, nil
174}
175
176func (r *Resyncer) resyncRepo(ctx context.Context, repoAt syntax.ATURI) error {
177 // ctx, span := tracer.Start(ctx, "resyncRepo")
178 // span.SetAttributes(attribute.String("aturi", repoAt))
179 // defer span.End()
180
181 resyncsStarted.Inc()
182 startTime := time.Now()
183
184 jobCtx, cancel := context.WithCancel(ctx)
185 r.registerRunning(repoAt, cancel)
186 defer r.unregisterRunning(repoAt)
187
188 success, err := r.doResync(jobCtx, repoAt)
189 if !success {
190 resyncsFailed.Inc()
191 resyncDuration.Observe(time.Since(startTime).Seconds())
192 return r.handleResyncFailure(ctx, repoAt, err)
193 }
194
195 resyncsCompleted.Inc()
196 resyncDuration.Observe(time.Since(startTime).Seconds())
197 return nil
198}
199
200func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, error) {
201 // ctx, span := tracer.Start(ctx, "doResync")
202 // span.SetAttributes(attribute.String("aturi", repoAt))
203 // defer span.End()
204
205 repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt)
206 if err != nil {
207 return false, fmt.Errorf("failed to get repo: %w", err)
208 }
209 if repo == nil { // untracked repo, skip
210 return false, nil
211 }
212
213 r.knotBackoffMu.RLock()
214 backoffUntil, inBackoff := r.knotBackoff[repo.KnotDomain]
215 r.knotBackoffMu.RUnlock()
216 if inBackoff && time.Now().Before(backoffUntil) {
217 return false, nil
218 }
219
220 // HACK: check knot reachability with short timeout before running actual fetch.
221 // This is crucial as git-cli doesn't support http connection timeout.
222 // `http.lowSpeedTime` is only applied _after_ the connection.
223 if err := r.checkKnotReachability(ctx, repo); err != nil {
224 if isRateLimitError(err) {
225 r.knotBackoffMu.Lock()
226 r.knotBackoff[repo.KnotDomain] = time.Now().Add(10 * time.Second)
227 r.knotBackoffMu.Unlock()
228 return false, nil
229 }
230 // TODO: suspend repo on 404. KnotStream updates will change the repo state back online
231 return false, fmt.Errorf("knot unreachable: %w", err)
232 }
233
234 timeout := r.repoFetchTimeout
235 if repo.RetryAfter == -1 {
236 timeout = r.manualResyncTimeout
237 }
238 fetchCtx, cancel := context.WithTimeout(ctx, timeout)
239 defer cancel()
240
241 if err := r.gitm.Sync(fetchCtx, repo); err != nil {
242 return false, err
243 }
244
245 // repo.GitRev = <processed git.refUpdate revision>
246 // repo.RepoSha = <sha256 sum of git refs>
247 repo.State = models.RepoStateActive
248 repo.ErrorMsg = ""
249 repo.RetryCount = 0
250 repo.RetryAfter = 0
251 if err := db.UpsertRepo(ctx, r.db, repo); err != nil {
252 return false, fmt.Errorf("updating repo state to active %w", err)
253 }
254 return true, nil
255}
256
257type knotStatusError struct {
258 StatusCode int
259}
260
261func (ke *knotStatusError) Error() string {
262 return fmt.Sprintf("request failed with status code (HTTP %d)", ke.StatusCode)
263}
264
265func isRateLimitError(err error) bool {
266 var knotErr *knotStatusError
267 if errors.As(err, &knotErr) {
268 return knotErr.StatusCode == http.StatusTooManyRequests
269 }
270 return false
271}
272
273// checkKnotReachability checks if Knot is reachable and is valid git remote server
274func (r *Resyncer) checkKnotReachability(ctx context.Context, repo *models.Repo) error {
275 repoUrl, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), true)
276 if err != nil {
277 return err
278 }
279
280 repoUrl += "/info/refs?service=git-upload-pack"
281
282 client := http.Client{
283 Timeout: 30 * time.Second,
284 }
285 req, err := http.NewRequestWithContext(ctx, "GET", repoUrl, nil)
286 if err != nil {
287 return err
288 }
289 req.Header.Set("User-Agent", "git/2.x")
290 req.Header.Set("Accept", "*/*")
291
292 resp, err := client.Do(req)
293 if err != nil {
294 var uerr *url.Error
295 if errors.As(err, &uerr) {
296 return fmt.Errorf("request failed: %w", uerr.Unwrap())
297 }
298 return fmt.Errorf("request failed: %w", err)
299 }
300 defer resp.Body.Close()
301
302 if resp.StatusCode != http.StatusOK {
303 return &knotStatusError{resp.StatusCode}
304 }
305
306 // check if target is git server
307 ct := resp.Header.Get("Content-Type")
308 if !strings.Contains(ct, "application/x-git-upload-pack-advertisement") {
309 return fmt.Errorf("unexpected content-type: %s", ct)
310 }
311
312 return nil
313}
314
315func (r *Resyncer) handleResyncFailure(ctx context.Context, repoAt syntax.ATURI, err error) error {
316 r.logger.Debug("handleResyncFailure", "at_uri", repoAt, "err", err)
317 var state models.RepoState
318 var errMsg string
319 if err == nil {
320 state = models.RepoStateDesynchronized
321 errMsg = ""
322 } else {
323 state = models.RepoStateError
324 errMsg = err.Error()
325 }
326
327 repo, err := db.GetRepoByAtUri(ctx, r.db, repoAt)
328 if err != nil {
329 return fmt.Errorf("failed to get repo: %w", err)
330 }
331 if repo == nil {
332 return fmt.Errorf("failed to get repo. repo '%s' doesn't exist in db", repoAt)
333 }
334
335 // start a 1 min & go up to 1 hr between retries
336 var retryCount = repo.RetryCount + 1
337 var retryAfter = time.Now().Add(backoff(retryCount, 60) * 60).Unix()
338
339 // remove null bytes
340 errMsg = strings.ReplaceAll(errMsg, "\x00", "")
341
342 repo.State = state
343 repo.ErrorMsg = errMsg
344 repo.RetryCount = retryCount
345 repo.RetryAfter = retryAfter
346 if err := db.UpsertRepo(ctx, r.db, repo); err != nil {
347 return fmt.Errorf("failed to update repo state: %w", err)
348 }
349 return nil
350}
351
352func backoff(retries int, max int) time.Duration {
353 dur := min(1<<retries, max)
354 jitter := time.Millisecond * time.Duration(rand.Intn(1000))
355 return time.Second*time.Duration(dur) + jitter
356}