internal/deps/worker.go
245 lines · 7106 bytes
1package deps
2
3import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "os"
9 "sort"
10 "strings"
11 "sync"
12 "time"
13
14 "gitbay.org/gitbay/internal/config"
15 "gitbay.org/gitbay/internal/gitutil"
16 "gitbay.org/gitbay/internal/store"
17)
18
19// manifestLimit bounds one manifest read out of the object store.
20const manifestLimit = 1 << 20
21
22// lookups is how many registry requests one repository has in flight.
23const lookups = 4
24
25// Worker sweeps repositories that have opted in, comparing their manifests
26// against the registries and maintaining one issue per repository.
27type Worker struct {
28 St *store.Store
29 Cfg config.Config
30 RepoDir func(owner, name string) string
31 Client *Client
32 Tick time.Duration
33}
34
35func New(st *store.Store, cfg config.Config, repoDir func(owner, name string) string, version string) *Worker {
36 tick := time.Hour
37 if v := os.Getenv("GITBAY_DEPS_TICK"); v != "" {
38 if d, err := time.ParseDuration(v); err == nil {
39 tick = d
40 }
41 }
42 return &Worker{St: st, Cfg: cfg, RepoDir: repoDir, Client: NewClient(version), Tick: tick}
43}
44
45// Run sweeps until ctx ends.
46func (w *Worker) Run(ctx context.Context) {
47 t := time.NewTicker(w.Tick)
48 defer t.Stop()
49 for {
50 select {
51 case <-ctx.Done():
52 return
53 case <-t.C:
54 w.Sweep(ctx)
55 }
56 }
57}
58
59// Sweep checks every repository whose interval has elapsed. Split from the
60// ticker for tests.
61func (w *Worker) Sweep(ctx context.Context) {
62 interval := w.Cfg.Deps.CheckIntervalHours
63 if interval <= 0 {
64 interval = 24
65 }
66 due, err := w.St.DueDepChecks(interval * 3600)
67 if err != nil {
68 slog.Error("deps: listing due repos", "err", err)
69 return
70 }
71 for _, repo := range due {
72 if ctx.Err() != nil {
73 return
74 }
75 msg := ""
76 if err := w.check(ctx, repo); err != nil {
77 slog.Warn("deps check failed", "repo", repo.Path(), "err", err)
78 msg = err.Error()
79 }
80 // Stamped either way, so a repo that fails every time is retried on
81 // the interval rather than on every sweep.
82 w.St.SetDepCheckResult(repo.ID, msg)
83 }
84}
85
86// check compares one repository against the registries and reconciles its
87// issue with the result.
88func (w *Worker) check(ctx context.Context, repo store.Repo) error {
89 dir := w.RepoDir(repo.OwnerName, repo.Name)
90 sha, err := gitutil.ResolveRef(dir, "refs/heads/"+repo.DefaultBranch)
91 if err != nil {
92 return nil // empty repo, or no default branch yet
93 }
94 found := Scan(func(path string) ([]byte, error) {
95 return gitutil.ReadBlob(dir, sha, path, manifestLimit)
96 })
97 if len(found) == 0 {
98 return w.reconcile(repo, nil)
99 }
100 behind, err := w.behind(ctx, found)
101 if err != nil {
102 return err
103 }
104 return w.reconcile(repo, behind)
105}
106
107// behind queries the registries and keeps the dependencies whose latest
108// release is greater than what the manifest declares. A lookup that fails
109// is dropped rather than failing the sweep: one unreachable package should
110// not silence the rest, and the next sweep tries again.
111func (w *Worker) behind(ctx context.Context, found []Dep) ([]store.DepReport, error) {
112 var (
113 mu sync.Mutex
114 out []store.DepReport
115 errs []string
116 wg sync.WaitGroup
117 tickets = make(chan struct{}, lookups)
118 )
119 for _, d := range found {
120 wg.Add(1)
121 go func(d Dep) {
122 defer wg.Done()
123 tickets <- struct{}{}
124 defer func() { <-tickets }()
125 latest, err := w.Client.Latest(ctx, d.Ecosystem, d.Name)
126 mu.Lock()
127 defer mu.Unlock()
128 switch {
129 case err != nil:
130 errs = append(errs, fmt.Sprintf("%s %s: %v", d.Ecosystem, d.Name, err))
131 case latest != "" && Newer(d.Current, latest):
132 out = append(out, store.DepReport{
133 Ecosystem: d.Ecosystem, Name: d.Name, Current: d.Current, Latest: latest})
134 }
135 }(d)
136 }
137 wg.Wait()
138 if ctx.Err() != nil {
139 return nil, ctx.Err()
140 }
141 // Every lookup failing means the registries are unreachable, which is
142 // worth recording; a few failing is ordinary.
143 if len(errs) == len(found) && len(errs) > 0 {
144 return nil, errors.New(errs[0])
145 }
146 sort.Slice(out, func(i, j int) bool {
147 if out[i].Ecosystem != out[j].Ecosystem {
148 return out[i].Ecosystem < out[j].Ecosystem
149 }
150 return out[i].Name < out[j].Name
151 })
152 return out, nil
153}
154
155// reconcile brings the repository's issue in line with what is behind:
156// opened when something first falls behind, rewritten when the set changes,
157// closed when nothing is behind any more.
158func (w *Worker) reconcile(repo store.Repo, behind []store.DepReport) error {
159 check, err := w.St.DepCheckFor(repo.ID)
160 if err != nil {
161 return err
162 }
163 previous, err := w.St.ReportedDeps(repo.ID)
164 if err != nil {
165 return err
166 }
167 issue, hasIssue := w.openIssue(repo, check.IssueNumber)
168 if len(behind) == 0 {
169 if hasIssue {
170 w.St.SetIssueState(issue.ID, "closed")
171 w.St.SetDepIssue(repo.ID, 0)
172 }
173 if len(previous) > 0 {
174 return w.St.ReplaceDepReports(repo.ID, nil)
175 }
176 return nil
177 }
178 // Nothing changed since the last report: leave it be, whether the issue
179 // is still open or the maintainer has closed it. Reopening on an
180 // unchanged set would make closing the issue pointless.
181 if same(previous, behind) {
182 return nil
183 }
184 if err := w.St.ReplaceDepReports(repo.ID, behind); err != nil {
185 return err
186 }
187 body := Body(repo.DefaultBranch, behind)
188 if hasIssue {
189 if err := w.St.UpdateIssueText(issue.ID, nil, &body, nil); err != nil {
190 return err
191 }
192 w.notify(repo, issue.Number, fmt.Sprintf("updated issue #%d", issue.Number), body)
193 return nil
194 }
195 author, err := w.St.UserByUsername(store.BotUsername)
196 if err != nil {
197 return fmt.Errorf("loading %s: %w", store.BotUsername, err)
198 }
199 number, err := w.St.CreateIssue(repo.ID, author.ID, IssueTitle, body, "md")
200 if err != nil {
201 return err
202 }
203 if err := w.St.SetDepIssue(repo.ID, number); err != nil {
204 return err
205 }
206 w.St.RecordEvent(repo.ID, author.ID, "issue.created", fmt.Sprintf(`{"number":%d}`, number))
207 w.notify(repo, number, fmt.Sprintf("opened issue #%d", number), body)
208 return nil
209}
210
211// openIssue loads the issue this worker maintains, if there still is one.
212// A closed issue counts as gone: reopening one the maintainer closed would
213// be arguing with them, so the next change opens a fresh issue.
214func (w *Worker) openIssue(repo store.Repo, number int64) (store.Issue, bool) {
215 if number == 0 {
216 return store.Issue{}, false
217 }
218 issue, err := w.St.IssueByNumber(repo.ID, number)
219 if err != nil || issue.State != "open" {
220 return store.Issue{}, false
221 }
222 return issue, true
223}
224
225// notify mails the repo's owners, the same targets and shape as an issue
226// filed over SSH. Best-effort, like every other notification.
227func (w *Worker) notify(repo store.Repo, number int64, action, body string) {
228 if w.Cfg.Mail.SMTPHost == "" {
229 return
230 }
231 targets, err := w.St.RepoNotifyTargets(repo)
232 if err != nil {
233 return
234 }
235 subject := fmt.Sprintf("[%s] #%d: %s", repo.Path(), number, IssueTitle)
236 text := fmt.Sprintf("%s %s\n\n%s\n%s/%s/issues/%d\n", store.BotUsername, action, body,
237 strings.TrimSuffix(w.Cfg.Server.SiteURL, "/"), repo.Path(), number)
238 for _, id := range targets {
239 email, err := w.St.PrimaryVerifiedEmail(id)
240 if err != nil || email == "" {
241 continue
242 }
243 w.St.EnqueueMail(email, subject, text)
244 }
245}