internal/symbols/worker.go
247 lines · 7483 bytes
1package symbols
2
3import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "os"
9 "time"
10
11 "gitbay.org/gitbay/internal/gitutil"
12 "gitbay.org/gitbay/internal/store"
13)
14
15// Bounds on one index run, and how it writes.
16const (
17 // DefaultMaxSymbols is where a repository's index stops growing.
18 DefaultMaxSymbols = 200_000
19 // DefaultMaxBytes bounds the names, keys and paths one index holds, so
20 // a tree of long names cannot fill the database within the count.
21 DefaultMaxBytes = 32 << 20
22 // DefaultMaxTime is how long one run may take.
23 DefaultMaxTime = 2 * time.Minute
24 // DefaultChunkRows is how many symbols one write transaction carries.
25 DefaultChunkRows = 5000
26 // DefaultBackoff is how long a failed tree waits before it is tried
27 // again.
28 DefaultBackoff = time.Hour
29)
30
31// Worker builds the index for each repository that has asked for one:
32// post-receive and the merge path ask after the default branch moves, and
33// `admin symbols reindex` asks with force. One repository at a time.
34type Worker struct {
35 St *store.Store
36 RepoDir func(owner, name string) string
37 Tick time.Duration
38 MaxSymbols int
39 MaxBytes int
40 MaxTime time.Duration
41 ChunkRows int
42 Backoff time.Duration
43 // chunkHook runs after each chunk is written; tests read the index
44 // mid-build through it.
45 chunkHook func(indexID int64)
46}
47
48func New(st *store.Store, repoDir func(owner, name string) string) *Worker {
49 tick := 5 * time.Second
50 if v := os.Getenv("GITBAY_SYMBOLS_TICK"); v != "" {
51 if d, err := time.ParseDuration(v); err == nil {
52 tick = d
53 }
54 }
55 return NewWith(st, repoDir, tick)
56}
57
58// NewWith is New with the tick given, and every bound at its default.
59func NewWith(st *store.Store, repoDir func(owner, name string) string, tick time.Duration) *Worker {
60 return &Worker{St: st, RepoDir: repoDir, Tick: tick,
61 MaxSymbols: DefaultMaxSymbols, MaxBytes: DefaultMaxBytes, MaxTime: DefaultMaxTime,
62 ChunkRows: DefaultChunkRows, Backoff: DefaultBackoff}
63}
64
65// Run sweeps until ctx ends.
66func (w *Worker) Run(ctx context.Context) {
67 t := time.NewTicker(w.Tick)
68 defer t.Stop()
69 for {
70 select {
71 case <-ctx.Done():
72 return
73 case <-t.C:
74 w.Sweep(ctx)
75 }
76 }
77}
78
79// Sweep handles every due request once. A request is cleared whatever
80// the outcome, except that a first failure is kept for one retry after
81// Backoff: a failure is retried once, not in a loop, and again only when
82// something asks.
83func (w *Worker) Sweep(ctx context.Context) {
84 reqs, err := w.St.SymbolRequests()
85 if err != nil {
86 slog.Error("symbols: listing requests", "err", err)
87 return
88 }
89 for _, req := range reqs {
90 if ctx.Err() != nil {
91 return
92 }
93 failed, err := w.Index(ctx, req.RepoID, req.Force)
94 if ctx.Err() != nil {
95 return // shutting down: the request stays for the next start
96 }
97 if err != nil {
98 slog.Warn("symbols: indexing", "repo", req.RepoID, "err", err)
99 }
100 if failed && req.Attempts == 0 {
101 w.St.DeferSymbolRequest(req, int(w.Backoff.Seconds()))
102 continue
103 }
104 w.St.DoneSymbolRequest(req)
105 }
106}
107
108// Index brings one repository's index up to its default branch's head.
109// A head whose tree is already indexed is left alone unless force is set,
110// and so is one whose tree failed within Backoff. The new index is
111// written in chunks while no read can see it, then published in one
112// short transaction. A run cut short by a bound publishes what it found
113// as partial; one that cannot build records the failure and leaves the
114// current index in place, reporting failed. The error is for the log.
115func (w *Worker) Index(ctx context.Context, repoID int64, force bool) (failed bool, err error) {
116 repo, err := w.St.RepoByID(repoID)
117 if errors.Is(err, store.ErrNotFound) {
118 return false, nil
119 } else if err != nil {
120 return false, err
121 }
122 // A run that stopped part way left its building index behind.
123 if err := w.St.PurgeSymbolIndexes(repo.ID); err != nil {
124 return false, err
125 }
126 dir := w.RepoDir(repo.OwnerName, repo.Name)
127 commit, err := gitutil.ResolveRef(dir, "refs/heads/"+repo.DefaultBranch)
128 if err != nil {
129 return false, nil // no default branch yet: nothing to index
130 }
131 tree, err := gitutil.ResolveTree(dir, commit)
132 if err != nil {
133 return false, err
134 }
135 if !force {
136 if cur, err := w.St.SymbolIndexFor(repo.ID); err == nil && cur.Tree == tree {
137 return false, nil
138 }
139 if recent, err := w.St.SymbolFailureRecent(repo.ID, tree, int(w.Backoff.Seconds())); err != nil || recent {
140 return false, err
141 }
142 }
143 x := store.SymbolIndex{RepoID: repo.ID, Commit: commit, Tree: tree, State: "ok"}
144 syms, files, runErr := w.collect(ctx, dir, tree)
145 if ctx.Err() != nil {
146 return false, ctx.Err()
147 }
148 x.Files = files
149 switch {
150 case errors.Is(runErr, errSymbolCap):
151 x.State, x.Note = "partial", fmt.Sprintf("stopped at %d symbols", w.MaxSymbols)
152 case errors.Is(runErr, errByteBudget):
153 x.State, x.Note = "partial", fmt.Sprintf("stopped at %d bytes of names and paths", w.MaxBytes)
154 case errors.Is(runErr, context.DeadlineExceeded):
155 x.State, x.Note = "partial", fmt.Sprintf("stopped after %s", w.MaxTime)
156 case runErr != nil:
157 if err := w.St.RecordSymbolFailure(repo.ID, tree, runErr.Error()); err != nil {
158 return true, err
159 }
160 return true, fmt.Errorf("%s failed: %v", repo.Path(), runErr)
161 }
162 if x.ID, err = w.St.BeginSymbolIndex(repo.ID, commit, tree); err != nil {
163 return false, err
164 }
165 for len(syms) > 0 {
166 if ctx.Err() != nil {
167 return false, ctx.Err() // the next run purges what was written
168 }
169 n := min(len(syms), w.ChunkRows)
170 if err := w.St.AddSymbols(x.ID, syms[:n]); err != nil {
171 return false, err
172 }
173 syms = syms[n:]
174 if w.chunkHook != nil {
175 w.chunkHook(x.ID)
176 }
177 }
178 if err := w.St.PublishSymbolIndex(x); err != nil {
179 return false, err
180 }
181 if err := w.St.PurgeSymbolIndexes(repo.ID); err != nil {
182 return false, err
183 }
184 if x.State != "ok" {
185 return false, fmt.Errorf("%s %s: %s", repo.Path(), x.State, x.Note)
186 }
187 return false, nil
188}
189
190var (
191 errSymbolCap = errors.New("symbol cap reached")
192 errByteBudget = errors.New("byte budget reached")
193)
194
195// collect reads every indexable blob in tree and extracts its symbols,
196// stopping at the symbol cap, the byte budget or the time bound with what
197// it has.
198func (w *Worker) collect(ctx context.Context, dir, tree string) ([]store.SymbolRow, int, error) {
199 ctx, cancel := context.WithTimeout(ctx, w.MaxTime)
200 defer cancel()
201 blobs, err := gitutil.ListBlobs(ctx, dir, tree)
202 if err != nil {
203 if ctx.Err() != nil {
204 return nil, 0, ctx.Err()
205 }
206 return nil, 0, err
207 }
208 var paths, shas []string
209 for _, b := range blobs {
210 if b.Mode == "120000" || b.Mode == "160000" || Skip(b.Name, b.Size) {
211 continue
212 }
213 paths = append(paths, b.Name)
214 shas = append(shas, b.SHA)
215 }
216 var out []store.SymbolRow
217 files, used := 0, 0
218 var stop error
219 err = gitutil.CatBlobs(ctx, dir, shas, func(i int, data []byte) bool {
220 if gitutil.IsBinary(data) {
221 return true
222 }
223 files++
224 for _, s := range Extract(paths[i], data) {
225 if len(out) >= w.MaxSymbols {
226 stop = errSymbolCap
227 return false
228 }
229 // The key is counted too: it is stored beside the name.
230 n := len(s.Name) + len(s.Key) + len(paths[i])
231 if used+n > w.MaxBytes {
232 stop = errByteBudget
233 return false
234 }
235 used += n
236 out = append(out, store.SymbolRow{Name: s.Name, Key: s.Key, Kind: s.Kind, Path: paths[i], Line: s.Line})
237 }
238 return true
239 })
240 if stop != nil {
241 return out, files, stop
242 }
243 if err != nil && ctx.Err() == nil {
244 return nil, files, err
245 }
246 return out, files, err
247}