internal/deps/worker.go

v1.31.0
gitbay/internal/deps/worker.go history · blame · raw

260 lines · 7746 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		// Migration 0028 leaves the account uncreated when the name was
198		// already taken, which is the one case worth spelling out: the
199		// feature is stuck until an operator frees the name.
200		return fmt.Errorf("no %s account to author the issue (the name was taken when this instance upgraded): %w",
201			store.BotUsername, err)
202	}
203	number, err := w.St.CreateIssue(repo.ID, author.ID, IssueTitle, body, "md")
204	if err != nil {
205		return err
206	}
207	if err := w.St.SetDepIssue(repo.ID, number); err != nil {
208		return err
209	}
210	w.St.RecordEvent(repo.ID, author.ID, "issue.created", fmt.Sprintf(`{"number":%d}`, number))
211	w.notify(repo, number, fmt.Sprintf("opened issue #%d", number), body)
212	return nil
213}
214
215// openIssue loads the issue this worker maintains, if there still is one.
216// A closed issue counts as gone: reopening one the maintainer closed would
217// be arguing with them, so the next change opens a fresh issue.
218func (w *Worker) openIssue(repo store.Repo, number int64) (store.Issue, bool) {
219	if number == 0 {
220		return store.Issue{}, false
221	}
222	issue, err := w.St.IssueByNumber(repo.ID, number)
223	if err != nil || issue.State != "open" {
224		return store.Issue{}, false
225	}
226	return issue, true
227}
228
229// notify files an inbox row for the repo's owners and watchers and mails
230// them, the same targets and shape as an issue opened over SSH.
231// Best-effort, like every other notification.
232func (w *Worker) notify(repo store.Repo, number int64, action, body string) {
233	targets, err := w.St.RepoNotifyTargets(repo)
234	if err != nil {
235		return
236	}
237	author, err := w.St.UserByUsername(store.BotUsername)
238	if err != nil {
239		return
240	}
241	recipients, err := w.St.NotifyRecipients(repo.ID, author.ID, targets, true)
242	if err != nil {
243		return
244	}
245	subject := fmt.Sprintf("[%s] #%d: %s", repo.Path(), number, IssueTitle)
246	text := fmt.Sprintf("%s %s\n\n%s\n%s/%s/issues/%d\n", store.BotUsername, action, body,
247		strings.TrimSuffix(w.Cfg.Server.SiteURL, "/"), repo.Path(), number)
248	path := fmt.Sprintf("%s/issues/%d", repo.Path(), number)
249	for _, id := range recipients {
250		w.St.AddNotice(id, repo.ID, "issue", store.BotUsername, action, path)
251		if w.Cfg.Mail.SMTPHost == "" {
252			continue
253		}
254		email, err := w.St.ActivityMailAddress(id)
255		if err != nil || email == "" {
256			continue
257		}
258		w.St.EnqueueMail(email, subject, text)
259	}
260}