From 817fd58d718e4a56bde60180ba1139bafea0d81f Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Mon, 18 May 2026 10:25:59 +0900 Subject: [PATCH] knotmirror: performant language indexer `git.listLanguages` has been one of the method that fails most often. Opening multiple git repos simultaneously can eaily cause OOM and language indexing itself usually takes super long. So several changes: - use `gitea.CatFileBatch` instead of go-git to avoid OOM - skip files larger than 16KB - sync HEAD ref language stats in knotmirror db - cache language stats info by commits (30d TTL) When syncing HEAD ref language stats, we do indexing on background. KnotMirror maintains internal "repo_stats_update" queue and right after `doResync` is done, enqueue the language stat indexing job so we can pre-index the language stats of HEAD ref. It's ok to spam this queue because all later events will be eventually ignored as we are resolving HEAD lazily. Signed-off-by: Seongmin Lee --- appview/state/knotstream.go | 1 + knotmirror/db/db.go | 14 ++ knotmirror/db/repo_index.go | 68 +++++++ knotmirror/knotmirror.go | 9 +- knotmirror/knotstream/scheduler.go | 14 +- knotmirror/knotstream/slurper.go | 4 +- knotmirror/repoindexer/language.go | 249 ++++++++++++++++++++++++++ knotmirror/resyncer.go | 24 ++- knotmirror/xrpc/git_list_languages.go | 98 +++------- knotmirror/xrpc/gitea/batch.go | 14 +- knotmirror/xrpc/xrpc.go | 5 +- 11 files changed, 398 insertions(+), 102 deletions(-) create mode 100644 knotmirror/db/repo_index.go create mode 100644 knotmirror/repoindexer/language.go diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go index 49877e3e..3df5985b 100644 --- a/appview/state/knotstream.go +++ b/appview/state/knotstream.go @@ -98,6 +98,7 @@ func knotIngester(d *db.DB, enforcer *rbac.Enforcer, posthog posthog.Client, not } } +// TODO(boltless): remove this. knotmirror should do all sort of indexing func ingestRefUpdate(ctx context.Context, d *db.DB, enforcer *rbac.Enforcer, pc posthog.Client, notifier notify.Notifier, dev bool, c *config.Config, cfClient *cloudflare.Client, source ec.Source, msg ec.Message) error { logger := log.FromContext(ctx) diff --git a/knotmirror/db/db.go b/knotmirror/db/db.go index 24540c24..5e6914a8 100644 --- a/knotmirror/db/db.go +++ b/knotmirror/db/db.go @@ -69,10 +69,24 @@ func Make(ctx context.Context, dbUrl string, maxConns int) (*sql.DB, error) { constraint hosts_pkey primary key (hostname) ); + -- repo language stats at HEAD + create table if not exists repo_head_languages ( + repo text not null, -- repo identifier (did) + commit text not null, -- commit id (oid) + language text not null, + size integer not null check (size >= 0), + + constraint repo_head_languages_pkey + primary key (repo, commit, language) + ); + create index if not exists idx_repos_aturi on repos (at_uri); create index if not exists idx_repos_db_updated_at on repos (db_updated_at desc); create index if not exists idx_hosts_db_updated_at on hosts (db_updated_at desc); + create index if not exists idx_repo_head_languages_repo_commit + on repo_head_languages (repo, commit); + create or replace function set_updated_at() returns trigger as $$ begin diff --git a/knotmirror/db/repo_index.go b/knotmirror/db/repo_index.go new file mode 100644 index 00000000..e67f5f25 --- /dev/null +++ b/knotmirror/db/repo_index.go @@ -0,0 +1,68 @@ +package db + +import ( + "context" + "database/sql" + "fmt" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/go-git/go-git/v5/plumbing" +) + +func IsLanguageIndexed(ctx context.Context, e *sql.DB, repoId syntax.DID, commitId plumbing.Hash) (bool, error) { + var exists bool + err := e.QueryRow(`select exists(select 1 from repo_head_languages where repo = $1)`, repoId).Scan(&exists) + return exists, err +} + +func InsertLanguages(ctx context.Context, e *sql.DB, repoId syntax.DID, commitId plumbing.Hash, langs map[string]int64) error { + tx, err := e.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("BeginTx: %w", err) + } + defer tx.Rollback() + + if _, err := tx.Exec(`delete from repo_head_languages where repo = $1`, repoId); err != nil { + return fmt.Errorf("deleting old languages: %w", err) + } + + for lang, size := range langs { + if _, err := tx.Exec( + `insert into repo_head_languages (repo, commit, language, size) + values ($1, $2, $3, $4)`, + repoId, commitId.String(), lang, size, + ); err != nil { + return fmt.Errorf("inserting language: %w", err) + } + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("tx.Commit: %w", err) + } + return nil +} + +func ListLanguages(ctx context.Context, e *sql.DB, repoId syntax.DID, commitId plumbing.Hash) (map[string]int64, error) { + sizes := make(map[string]int64) + + rows, err := e.Query(`select language, size from repo_head_languages where repo = $1 and commit = $2`, repoId, commitId.String()) + if err != nil { + return nil, err + } + defer rows.Close() + + for rows.Next() { + var lang string + var size int64 + if err := rows.Scan(&lang, &size); err != nil { + return nil, err + } + sizes[lang] = size + } + + if err := rows.Err(); err != nil { + return nil, err + } + + return sizes, nil +} diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go index 6c004c98..902baccf 100644 --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -14,6 +14,7 @@ import ( "tangled.org/core/knotmirror/config" "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/knotstream" + "tangled.org/core/knotmirror/repoindexer" "tangled.org/core/knotmirror/models" "tangled.org/core/knotmirror/xrpc" "tangled.org/core/log" @@ -56,10 +57,14 @@ func Run(ctx context.Context, cfg *config.Config) error { } logger.Info(fmt.Sprintf("clearing resyning states: %d records updated", rows)) + indexer := repoindexer.NewIndexer(logger, cfg, rdb) + indexScheduler := repoindexer.NewBackgroundIndexScheduler(logger, cfg, db, indexer) + indexScheduler.Start(ctx) + knotstream := knotstream.NewKnotStream(logger, db, cfg) crawler := NewCrawler(logger, db) - resyncer := NewResyncer(logger, db, gitm, cfg) - xrpc := xrpc.New(logger, cfg, db, rdb, resolver, knotstream) + resyncer := NewResyncer(logger, db, gitm, indexScheduler, cfg) + xrpc := xrpc.New(logger, cfg, db, rdb, indexer, resolver, knotstream) adminpage := NewAdminServer(logger, db, resyncer, xrpc, resolver) // maintain repository list with tap diff --git a/knotmirror/knotstream/scheduler.go b/knotmirror/knotstream/scheduler.go index 95103808..b3fd537b 100644 --- a/knotmirror/knotstream/scheduler.go +++ b/knotmirror/knotstream/scheduler.go @@ -24,7 +24,7 @@ type ParallelScheduler struct { } type Task struct { - key string + Key string message []byte } @@ -46,13 +46,13 @@ func (s *ParallelScheduler) Start(ctx context.Context) { func (s *ParallelScheduler) AddTask(ctx context.Context, task *Task) { s.lk.Lock() - if st, ok := s.scheduled[task.key]; ok { + if st, ok := s.scheduled[task.Key]; ok { // schedule task - s.scheduled[task.key] = append(st, task) + s.scheduled[task.Key] = append(st, task) s.lk.Unlock() return } - s.scheduled[task.key] = []*Task{} + s.scheduled[task.Key] = []*Task{} s.lk.Unlock() select { @@ -77,16 +77,16 @@ func (s *ParallelScheduler) ForEach(ctx context.Context, fn func(context.Context s.lk.Lock() func() { - rem, ok := s.scheduled[task.key] + rem, ok := s.scheduled[task.Key] if !ok { s.logger.Error("should always have an 'active' entry if a worker is processing a job") } if len(rem) == 0 { - delete(s.scheduled, task.key) + delete(s.scheduled, task.Key) task = nil } else { task = rem[0] - s.scheduled[task.key] = rem[1:] + s.scheduled[task.Key] = rem[1:] } // TODO: update seq from received message diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go index 539d8e6e..36f4318d 100644 --- a/knotmirror/knotstream/slurper.go +++ b/knotmirror/knotstream/slurper.go @@ -257,7 +257,7 @@ func (s *KnotSlurper) handleConnection(ctx context.Context, conn *websocket.Conn } sub.scheduler.AddTask(ctx, &Task{ - key: sub.hostname, // TODO: replace to repository AT-URI for better concurrency + Key: sub.hostname, // TODO: replace to repository AT-URI for better concurrency message: msg, }) } @@ -281,7 +281,7 @@ func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { return fmt.Errorf("unmarshaling message: %w", err) } - if err := s.ProcessLegacyGitRefUpdate(ctx, task.key, &legacyMessage); err != nil { + if err := s.ProcessLegacyGitRefUpdate(ctx, task.Key, &legacyMessage); err != nil { return fmt.Errorf("processing gitRefUpdate: %w", err) } return nil diff --git a/knotmirror/repoindexer/language.go b/knotmirror/repoindexer/language.go new file mode 100644 index 00000000..50bc0742 --- /dev/null +++ b/knotmirror/repoindexer/language.go @@ -0,0 +1,249 @@ +package repoindexer + +import ( + "bufio" + "context" + "database/sql" + "encoding/json" + "fmt" + "io" + "log/slog" + "path/filepath" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/go-enry/go-enry/v2" + "github.com/go-git/go-git/v5/plumbing" + "github.com/go-git/go-git/v5/plumbing/filemode" + "github.com/go-git/go-git/v5/plumbing/object" + "github.com/redis/go-redis/v9" + "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/knotstream" + "tangled.org/core/knotmirror/xrpc/gitea" + "tangled.org/core/log" +) + +const ( + fileSizeLimit = 16 * 1024 // read up to 16 KiB for language detection + bigFileSize = 1024 * 1024 // skip content read for blobs over 1 MiB + langIndexRepoCommit = "lang_index:%s:%s" // lang_index:{did}:{oid} + langIndexRepoCommitTTL = 30 * 24 * time.Hour +) + +// Language indexing strategy: +// +// git.refUpdate HEAD -> store to db +// git.refUpdate other -> on-demand calculation, cache +// +// NOTE: currently all event we have is git.refUpdate, so background indexing +// job will always be triggered. +// TODO(boltless): don't queue indexing job on "sync" type event while repo is +// not active. + +type Indexer struct { + logger *slog.Logger + cfg *config.Config + rdb *redis.Client +} + +func NewIndexer(l *slog.Logger, cfg *config.Config, rdb *redis.Client) *Indexer { + indexer := &Indexer{ + logger: log.SubLogger(l, "indexer"), + cfg: cfg, + rdb: rdb, + } + return indexer +} + +func NewBackgroundIndexScheduler(l *slog.Logger, cfg *config.Config, e *sql.DB, indexer *Indexer) *knotstream.ParallelScheduler { + return knotstream.NewParallelScheduler( + 4, + "repo_stats_update", // NOTE: this is unused + func(ctx context.Context, t *knotstream.Task) error { + start := time.Now() + repoId := syntax.DID(t.Key) + + // resolve HEAD to commitId + commit, err := gitea.GetCommit(ctx, indexer.repoPath(repoId), "HEAD") + if err != nil { + return fmt.Errorf("failed to resolve HEAD: %w", err) + } + + l := l.With("repo", repoId, "hash", commit.Hash) + + // check if (did,oid) is already indexed + indexed, err := db.IsLanguageIndexed(ctx, e, repoId, commit.Hash) + if err != nil { + l.Error("failed to query langs", "err", err) + indexed = false + // continue + } + if indexed { + return nil + } + + langs, err := indexer.IndexLanguages(ctx, repoId, commit.Hash) + if err != nil { + return fmt.Errorf("indexing langs: %w", err) + } + + l.Info("pre-indexed language stats", "duration", time.Since(start)) + + if err := db.InsertLanguages(ctx, e, repoId, commit.Hash, langs); err != nil { + return fmt.Errorf("failed to insert langs into db: %w", err) + } + + return nil + }, + ) +} + +func (i *Indexer) repoPath(repo syntax.DID) string { + return filepath.Join(i.cfg.GitRepoBasePath, repo.String()) +} + +// IndexLanguages index the repository language stats at given commit +func (i *Indexer) IndexLanguages(ctx context.Context, repoId syntax.DID, commitId plumbing.Hash) (map[string]int64, error) { + if i.rdb != nil { + if val, err := i.rdb.Get(ctx, fmt.Sprintf(langIndexRepoCommit, repoId, commitId.String())).Result(); err == nil { + i.logger.Debug("serve from cache") + var sizes map[string]int64 + if err := json.Unmarshal([]byte(val), &sizes); err == nil { + return sizes, nil + } + } + } + + sizes, err := IndexLanguagesInner(ctx, i.repoPath(repoId), commitId) + if err != nil { + return nil, err + } + + if i.rdb != nil { + if encoded, err := json.Marshal(sizes); err == nil { + i.logger.Debug("cache language") + if err := i.rdb.Set(ctx, + fmt.Sprintf(langIndexRepoCommit, repoId, commitId.String()), + encoded, + langIndexRepoCommitTTL, + ).Err(); err != nil { + i.logger.Error("failed to cache languages", "err", err) + } + } + } + return sizes, nil +} + +func IndexLanguagesInner(ctx context.Context, repoPath string, commitId plumbing.Hash) (map[string]int64, error) { + tree, err := gitea.GetTree(ctx, repoPath, commitId.String()+"^{tree}") + if err != nil { + return nil, err + } + + bw, br, close := gitea.CatFileBatch(ctx, repoPath) + defer close() + + sizes, err := batchAnalyzeTree(ctx, bw, br, tree) + if err != nil { + return nil, err + } + return sizes, nil +} + +func batchAnalyzeTree(ctx context.Context, bw io.WriteCloser, br *bufio.Reader, tree *object.Tree) (map[string]int64, error) { + sizes := make(map[string]int64) + for _, entry := range tree.Entries { + select { + case <-ctx.Done(): + return nil, ctx.Err() + default: + } + + switch entry.Mode { + case filemode.Dir: + subTree, err := gitea.BatchGetTree(bw, br, entry.Hash.String()) + if err != nil { + return nil, err + } + subTreeSizes, err := batchAnalyzeTree(ctx, bw, br, subTree) + if err != nil { + return nil, err + } + for name, size := range subTreeSizes { + sizes[name] += size + } + case filemode.Symlink, filemode.Submodule: + // skip symlink/submodule + default: + _, err := bw.Write([]byte(entry.Hash.String() + "\n")) + if err != nil { + return nil, err + } + _, _, size, err := gitea.ReadBatchLine(br) + if err != nil { + return nil, err + } + // skip large file + if size > bigFileSize { + if err := gitea.DiscardFull(br, size+1); err != nil { + return nil, err + } + continue + } + + sizeToRead := size + discard := int64(1) + if size > fileSizeLimit { + sizeToRead = fileSizeLimit + discard = size - fileSizeLimit + 1 + } + content, err := io.ReadAll(io.LimitReader(br, sizeToRead)) + if err != nil { + return nil, err + } + if err := gitea.DiscardFull(br, discard); err != nil { + return nil, err + } + + language, noskip := analyzeLanguage(entry.Name, content) + if noskip { + sizes[language] += size + } + } + } + return sizes, nil +} + +func analyzeLanguage(fileName string, content []byte) (string, bool) { + // skip generated file + // TODO: follow gitattributes, lazyily read content (filter by filename first) + if enry.IsGenerated(fileName, content) || enry.IsBinary(content) || strings.HasSuffix(fileName, "bun.lock") { + return "", false + } + + language := func(fileName string, content []byte) string { + language, ok := enry.GetLanguageByExtension(fileName) + if ok { + return language + } + language, ok = enry.GetLanguageByFilename(fileName) + if ok { + return language + } + if len(content) == 0 { + return enry.OtherLanguage + } + return enry.GetLanguage(fileName, content) + }(fileName, content) + if group := enry.GetLanguageGroup(language); group != "" { + language = group + } + + langType := enry.GetLanguageType(language) + if langType != enry.Programming && langType != enry.Markup && langType != enry.Unknown { + return "", false + } + return language, true +} diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 594ef5ad..9818e2ba 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -16,15 +16,17 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/knotmirror/config" "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/knotstream" "tangled.org/core/knotmirror/models" "tangled.org/core/log" ) type Resyncer struct { - logger *slog.Logger - db *sql.DB - gitm GitMirrorManager - cfg *config.Config + logger *slog.Logger + db *sql.DB + gitm GitMirrorManager + cfg *config.Config + indexer *knotstream.ParallelScheduler claimJobMu sync.Mutex @@ -41,12 +43,13 @@ type Resyncer struct { httpClient *http.Client } -func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, cfg *config.Config) *Resyncer { +func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, indexer *knotstream.ParallelScheduler, cfg *config.Config) *Resyncer { return &Resyncer{ - logger: log.SubLogger(l, "resyncer"), - db: db, - gitm: gitm, - cfg: cfg, + logger: log.SubLogger(l, "resyncer"), + db: db, + gitm: gitm, + cfg: cfg, + indexer: indexer, runningJobs: make(map[syntax.ATURI]context.CancelFunc), @@ -248,6 +251,9 @@ func (r *Resyncer) doResync(ctx context.Context, repoAt syntax.ATURI) (bool, err return false, err } + // queue repo_stats_update job + r.indexer.AddTask(context.TODO(), &knotstream.Task{Key: repo.RepoDid.String()}) + // repo.GitRev = // repo.RepoSha = repo.State = models.RepoStateActive diff --git a/knotmirror/xrpc/git_list_languages.go b/knotmirror/xrpc/git_list_languages.go index c63455e8..263eaa6b 100644 --- a/knotmirror/xrpc/git_list_languages.go +++ b/knotmirror/xrpc/git_list_languages.go @@ -2,21 +2,14 @@ package xrpc import ( "context" - "encoding/json" "fmt" - "math" "net/http" "time" "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" - "tangled.org/core/knotserver/git" -) - -const ( - RepoLanguagesByDid = "git_list_languages:repo:%s:%s" - RepoLanguagesTTL = 24 * time.Hour + "tangled.org/core/knotmirror/xrpc/gitea" ) func (x *Xrpc) ListLanguages(w http.ResponseWriter, r *http.Request) { @@ -34,86 +27,39 @@ func (x *Xrpc) ListLanguages(w http.ResponseWriter, r *http.Request) { return } - if val, err := x.rdb.Get(r.Context(), fmt.Sprintf(RepoLanguagesByDid, repo, ref)).Result(); err == nil { - l.Debug("served from cache") - var langs []*tangled.GitTempListLanguages_Language - err = json.Unmarshal([]byte(val), &langs) - if err == nil { - writeJson(w, http.StatusOK, &tangled.GitTempListLanguages_Output{ - Ref: ref, - Languages: langs, - }) - return - } - } - - var out *tangled.GitTempListLanguages_Output - out, err = x.listLanguages(r.Context(), repo, ref) - if err != nil { - l.Warn("local mirror failed, trying proxy", "err", err) - if x.proxyToKnot(w, r, repo) { - return - } - writeErr(w, err) - return - } - - go func() { - ctx := context.Background() - encoded, err := json.Marshal(out.Languages) - if err != nil { - return - } - x.rdb.Set(ctx, fmt.Sprintf(RepoLanguagesByDid, repo, ref), encoded, RepoLanguagesTTL) - }() - - writeJson(w, http.StatusOK, out) -} + ctx := r.Context() -func (x *Xrpc) listLanguages(ctx context.Context, repo syntax.DID, ref string) (*tangled.GitTempListLanguages_Output, error) { repoPath, err := x.makeRepoPath(ctx, repo) if err != nil { - return nil, fmt.Errorf("resolving repo did: %w", err) + l.Error("failed to make repo path", "err", err) + writeJson(w, http.StatusNotFound, atclient.ErrorBody{Name: "RepoNotFound", Message: fmt.Sprintf("unknown repository: %s", repo)}) + return } - gr, err := git.Open(repoPath, ref) + commit, err := gitea.GetCommit(ctx, repoPath, ref) if err != nil { - return nil, &atclient.APIError{StatusCode: http.StatusNotFound, Name: "RepoNotFound", Message: "failed to find git repo"} + l.Error("failed to get commit", "err", err) + writeJson(w, http.StatusNotFound, atclient.ErrorBody{Name: "RefNotFound", Message: fmt.Sprintf("unknown git ref: %s", repo)}) + return } - ctx, cancel := context.WithTimeout(ctx, 1*time.Second) + indexCtx, cancel := context.WithTimeout(ctx, 1 * time.Second) defer cancel() - - sizes, err := gr.AnalyzeLanguages(ctx) + sizes, err := x.indexer.IndexLanguages(indexCtx, repo, commit.Hash) if err != nil { - return nil, fmt.Errorf("analyzing languages: %w", err) - } - - return &tangled.GitTempListLanguages_Output{ - Ref: ref, - Languages: sizesToLanguages(sizes), - }, nil -} - -func sizesToLanguages(sizes git.LangBreakdown) []*tangled.GitTempListLanguages_Language { - var apiLanguages []*tangled.GitTempListLanguages_Language - var totalSize int64 - for _, size := range sizes { - totalSize += size + l.Error("failed to serve languages", "err", err) + writeJson(w, http.StatusNotFound, atclient.ErrorBody{Name: "InternalServerError", Message: "failed to serve languages"}) + return } - for name, size := range sizes { - percentagef64 := float64(size) / float64(totalSize) * 100 - percentage := math.Round(percentagef64) - - lang := &tangled.GitTempListLanguages_Language{ - Name: name, - Size: size, - Percentage: int64(percentage), - } - - apiLanguages = append(apiLanguages, lang) + var out tangled.GitTempListLanguages_Output + for lang, size := range sizes { + out.Total += size + out.Languages = append(out.Languages, &tangled.GitTempListLanguages_Language{ + Name: lang, + Size: size, + }) } - return apiLanguages + writeJson(w, http.StatusOK, &out) } diff --git a/knotmirror/xrpc/gitea/batch.go b/knotmirror/xrpc/gitea/batch.go index 017802f5..fe49ba6e 100644 --- a/knotmirror/xrpc/gitea/batch.go +++ b/knotmirror/xrpc/gitea/batch.go @@ -49,24 +49,28 @@ func GetCommit(ctx context.Context, repoPath, rev string) (*object.Commit, error } func GetTree(ctx context.Context, repoPath, rev string) (*object.Tree, error) { - wr, rd, cancel := CatFileBatch(ctx, repoPath) + bw, br, cancel := CatFileBatch(ctx, repoPath) defer cancel() - if _, err := wr.Write([]byte(rev + "\n")); err != nil { + return BatchGetTree(bw, br, rev) +} + +func BatchGetTree(bw io.WriteCloser, br *bufio.Reader, rev string) (*object.Tree, error) { + if _, err := bw.Write([]byte(rev + "\n")); err != nil { return nil, fmt.Errorf("write rev: %w", err) } - sha, typ, size, err := ReadBatchLine(rd) + sha, typ, size, err := ReadBatchLine(br) if err != nil { return nil, fmt.Errorf("resolve %s: %w", rev, err) } if typ != "tree" { - if err := DiscardFull(rd, size+1); err != nil { + if err := DiscardFull(br, size+1); err != nil { return nil, err } return nil, fmt.Errorf("unexpected type: %s for tree: %s", typ, rev) } - entries, err := catBatchParseTreeEntries(rd, size) + entries, err := catBatchParseTreeEntries(br, size) if err != nil { return nil, fmt.Errorf("read tree %s: %w", rev, err) } diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go index a2dbcc13..9e90f188 100644 --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -15,6 +15,7 @@ import ( "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/repoindexer" "tangled.org/core/knotmirror/knotstream" "tangled.org/core/log" ) @@ -23,6 +24,7 @@ type Xrpc struct { cfg *config.Config db *sql.DB rdb *redis.Client + indexer *repoindexer.Indexer resolver *idresolver.Resolver ks *knotstream.KnotStream logger *slog.Logger @@ -30,7 +32,7 @@ type Xrpc struct { inflight *inflightTracker } -func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { +func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, indexer *repoindexer.Indexer, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { httpClient := &http.Client{ Timeout: 30 * time.Second, } @@ -41,6 +43,7 @@ func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, cfg: cfg, db: db, rdb: rdb, + indexer: indexer, resolver: resolver, ks: ks, logger: log.SubLogger(logger, "xrpc"), -- 2.51.2