This repository has no description
7.6 kB
274 lines
1package repoindexer
2
3import (
4 "bufio"
5 "context"
6 "database/sql"
7 "encoding/json"
8 "fmt"
9 "io"
10 "log/slog"
11 "net/http"
12 "path/filepath"
13 "strings"
14 "time"
15
16 "github.com/bluesky-social/indigo/atproto/syntax"
17 "github.com/go-enry/go-enry/v2"
18 "github.com/go-git/go-git/v5/plumbing"
19 "github.com/go-git/go-git/v5/plumbing/filemode"
20 "github.com/go-git/go-git/v5/plumbing/object"
21 "github.com/redis/go-redis/v9"
22 "tangled.org/core/knotmirror/config"
23 "tangled.org/core/knotmirror/db"
24 "tangled.org/core/knotmirror/knotstream"
25 "tangled.org/core/knotmirror/xrpc/gitea"
26 "tangled.org/core/log"
27)
28
29const (
30 fileSizeLimit = 16 * 1024 // read up to 16 KiB for language detection
31 bigFileSize = 1024 * 1024 // skip content read for blobs over 1 MiB
32 langIndexRepoCommit = "lang_index:%s:%s" // lang_index:{did}:{oid}
33 langIndexRepoCommitTTL = 30 * 24 * time.Hour
34)
35
36// Language indexing strategy:
37//
38// git.refUpdate HEAD -> store to db
39// git.refUpdate other -> on-demand calculation, cache
40//
41// NOTE: currently all event we have is git.refUpdate, so background indexing
42// job will always be triggered.
43// TODO(boltless): don't queue indexing job on "sync" type event while repo is
44// not active.
45
46type Indexer struct {
47 logger *slog.Logger
48 cfg *config.Config
49 rdb *redis.Client
50}
51
52func NewIndexer(l *slog.Logger, cfg *config.Config, rdb *redis.Client) *Indexer {
53 indexer := &Indexer{
54 logger: log.SubLogger(l, "indexer"),
55 cfg: cfg,
56 rdb: rdb,
57 }
58 return indexer
59}
60
61func NewBackgroundIndexScheduler(l *slog.Logger, cfg *config.Config, e *sql.DB, indexer *Indexer) *knotstream.ParallelScheduler {
62 return knotstream.NewParallelScheduler(
63 4,
64 "repo_stats_update", // NOTE: this is unused
65 func(ctx context.Context, t *knotstream.Task) error {
66 start := time.Now()
67 repoId := syntax.DID(t.Key)
68
69 // resolve HEAD to commitId
70 commit, err := gitea.GetCommit(ctx, indexer.repoPath(repoId), "HEAD")
71 if err != nil {
72 return fmt.Errorf("failed to resolve HEAD: %w", err)
73 }
74
75 l := l.With("repo", repoId, "hash", commit.Hash)
76
77 // check if (did,oid) is already indexed
78 indexed, err := db.IsLanguageIndexed(ctx, e, repoId, commit.Hash)
79 if err != nil {
80 l.Error("failed to query langs", "err", err)
81 indexed = false
82 // continue
83 }
84 if indexed {
85 return nil
86 }
87
88 langs, err := indexer.IndexLanguages(ctx, repoId, commit.Hash)
89 if err != nil {
90 return fmt.Errorf("indexing langs: %w", err)
91 }
92
93 l.Info("pre-indexed language stats", "duration", time.Since(start))
94
95 if err := db.InsertLanguages(ctx, e, repoId, commit.Hash, langs); err != nil {
96 return fmt.Errorf("failed to insert langs into db: %w", err)
97 }
98
99 // HACK(boltless): ping appview to update language cache.
100 go func() {
101 url := fmt.Sprintf("%s/%s", cfg.AppviewUrl, repoId.String())
102 pingCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
103 defer cancel()
104 req, err := http.NewRequestWithContext(pingCtx, http.MethodGet, url, nil)
105 if err != nil {
106 l.Warn("appview ping: build request failed", "err", err)
107 return
108 }
109 resp, err := http.DefaultClient.Do(req)
110 if err != nil {
111 l.Warn("appview ping failed", "url", url, "err", err)
112 return
113 }
114 defer resp.Body.Close()
115 // drain body to ensure the appview completes rendering
116 if _, err := io.Copy(io.Discard, resp.Body); err != nil {
117 l.Warn("appview ping: drain response failed", "url", url, "err", err)
118 return
119 }
120 l.Info("appview pinged", "url", url, "status", resp.StatusCode)
121 }()
122
123 return nil
124 },
125 )
126}
127
128func (i *Indexer) repoPath(repo syntax.DID) string {
129 return filepath.Join(i.cfg.GitRepoBasePath, repo.String())
130}
131
132// IndexLanguages index the repository language stats at given commit
133func (i *Indexer) IndexLanguages(ctx context.Context, repoId syntax.DID, commitId plumbing.Hash) (map[string]int64, error) {
134 if i.rdb != nil {
135 if val, err := i.rdb.Get(ctx, fmt.Sprintf(langIndexRepoCommit, repoId, commitId.String())).Result(); err == nil {
136 i.logger.Debug("serve from cache")
137 var sizes map[string]int64
138 if err := json.Unmarshal([]byte(val), &sizes); err == nil {
139 return sizes, nil
140 }
141 }
142 }
143
144 sizes, err := IndexLanguagesInner(ctx, i.repoPath(repoId), commitId)
145 if err != nil {
146 return nil, err
147 }
148
149 if i.rdb != nil {
150 if encoded, err := json.Marshal(sizes); err == nil {
151 i.logger.Debug("cache language")
152 if err := i.rdb.Set(ctx,
153 fmt.Sprintf(langIndexRepoCommit, repoId, commitId.String()),
154 encoded,
155 langIndexRepoCommitTTL,
156 ).Err(); err != nil {
157 i.logger.Error("failed to cache languages", "err", err)
158 }
159 }
160 }
161 return sizes, nil
162}
163
164func IndexLanguagesInner(ctx context.Context, repoPath string, commitId plumbing.Hash) (map[string]int64, error) {
165 tree, err := gitea.GetTree(ctx, repoPath, commitId.String()+"^{tree}")
166 if err != nil {
167 return nil, err
168 }
169
170 bw, br, close := gitea.CatFileBatch(ctx, repoPath)
171 defer close()
172
173 sizes, err := batchAnalyzeTree(ctx, bw, br, tree)
174 if err != nil {
175 return nil, err
176 }
177 return sizes, nil
178}
179
180func batchAnalyzeTree(ctx context.Context, bw io.WriteCloser, br *bufio.Reader, tree *object.Tree) (map[string]int64, error) {
181 sizes := make(map[string]int64)
182 for _, entry := range tree.Entries {
183 select {
184 case <-ctx.Done():
185 return nil, ctx.Err()
186 default:
187 }
188
189 switch entry.Mode {
190 case filemode.Dir:
191 subTree, err := gitea.BatchGetTree(bw, br, entry.Hash.String())
192 if err != nil {
193 return nil, err
194 }
195 subTreeSizes, err := batchAnalyzeTree(ctx, bw, br, subTree)
196 if err != nil {
197 return nil, err
198 }
199 for name, size := range subTreeSizes {
200 sizes[name] += size
201 }
202 case filemode.Symlink, filemode.Submodule:
203 // skip symlink/submodule
204 default:
205 _, err := bw.Write([]byte(entry.Hash.String() + "\n"))
206 if err != nil {
207 return nil, err
208 }
209 _, _, size, err := gitea.ReadBatchLine(br)
210 if err != nil {
211 return nil, err
212 }
213 // skip large file
214 if size > bigFileSize {
215 if err := gitea.DiscardFull(br, size+1); err != nil {
216 return nil, err
217 }
218 continue
219 }
220
221 sizeToRead := size
222 discard := int64(1)
223 if size > fileSizeLimit {
224 sizeToRead = fileSizeLimit
225 discard = size - fileSizeLimit + 1
226 }
227 content, err := io.ReadAll(io.LimitReader(br, sizeToRead))
228 if err != nil {
229 return nil, err
230 }
231 if err := gitea.DiscardFull(br, discard); err != nil {
232 return nil, err
233 }
234
235 language, noskip := analyzeLanguage(entry.Name, content)
236 if noskip {
237 sizes[language] += size
238 }
239 }
240 }
241 return sizes, nil
242}
243
244func analyzeLanguage(fileName string, content []byte) (string, bool) {
245 // skip generated file
246 // TODO: follow gitattributes, lazyily read content (filter by filename first)
247 if enry.IsGenerated(fileName, content) || enry.IsBinary(content) || strings.HasSuffix(fileName, "bun.lock") {
248 return "", false
249 }
250
251 language := func(fileName string, content []byte) string {
252 language, ok := enry.GetLanguageByExtension(fileName)
253 if ok {
254 return language
255 }
256 language, ok = enry.GetLanguageByFilename(fileName)
257 if ok {
258 return language
259 }
260 if len(content) == 0 {
261 return enry.OtherLanguage
262 }
263 return enry.GetLanguage(fileName, content)
264 }(fileName, content)
265 if group := enry.GetLanguageGroup(language); group != "" {
266 language = group
267 }
268
269 langType := enry.GetLanguageType(language)
270 if langType != enry.Programming && langType != enry.Markup && langType != enry.Unknown {
271 return "", false
272 }
273 return language, true
274}