This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knotmirror / repoindexer / language.go
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}