From 58f2054dcffd158a2b4ff0f7d6ccfd518b573ee2 Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 11 Sep 2026 05:58:44 +0300 Subject: [PATCH] knotmirror/xrpc: limit concurrent git requests and return 429 when busy Signed-off-by: dawn --- knotmirror/knotmirror.go | 8 ++++++- knotmirror/xrpc/git_get_merge_base.go | 2 +- knotmirror/xrpc/xrpc.go | 31 +++++++++++++++++++++++++++ 3 files changed, 39 insertions(+), 2 deletions(-) diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go index 2157f64ca..56c6cb1e3 100644 --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -102,7 +102,13 @@ func Run(ctx context.Context, cfg *config.Config) error { } mux.Mount("/xrpc", xrpc.Router()) - if err := http.ListenAndServe(cfg.Listen, mux); err != nil { + srv := &http.Server{ + Addr: cfg.Listen, + Handler: mux, + ReadHeaderTimeout: 5 * time.Second, + IdleTimeout: 60 * time.Second, + } + if err := srv.ListenAndServe(); err != nil { logger.Error("xrpc server failed", "error", err) } }() diff --git a/knotmirror/xrpc/git_get_merge_base.go b/knotmirror/xrpc/git_get_merge_base.go index a7a1b7308..a460af74b 100644 --- a/knotmirror/xrpc/git_get_merge_base.go +++ b/knotmirror/xrpc/git_get_merge_base.go @@ -46,7 +46,7 @@ func (x *Xrpc) GetMergeBase(w http.ResponseWriter, r *http.Request) { setCacheControl(w, base, head) - cmd := exec.Command("git", "-C", repoPath, "merge-base", head, base) + cmd := exec.CommandContext(ctx, "git", "-C", repoPath, "merge-base", head, base) out, err := cmd.Output() if err != nil { writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalError", Message: "failed to compute merge-base"}) diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go index 15ca53377..5dd4c5c37 100644 --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -16,6 +16,8 @@ import ( "github.com/hashicorp/golang-lru/v2/expirable" "github.com/redis/go-redis/v9" "github.com/samber/lo" + "context" + "golang.org/x/sync/semaphore" "golang.org/x/sync/singleflight" "tangled.org/core/api/tangled" "tangled.org/core/gitutil" @@ -42,6 +44,7 @@ type Xrpc struct { v2WriteClient *http.Client committers *expirable.LRU[string, []string] committerGroup singleflight.Group + gitSem *semaphore.Weighted inflight *inflightTracker serviceSigner *serviceauth.Signer } @@ -66,6 +69,7 @@ func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, v2Client: &http.Client{Timeout: 10 * time.Second}, v2WriteClient: &http.Client{Timeout: 5 * time.Minute}, committers: expirable.NewLRU[string, []string](committerCacheCapacity, nil, committerCacheFreshness), + gitSem: semaphore.NewWeighted(64), inflight: newInflightTracker(), } } @@ -82,6 +86,7 @@ func (x *Xrpc) Router() http.Handler { r.Group(func(r chi.Router) { r.Use(x.inflight.middleware) r.Use(x.forwardSuspended) + r.Use(x.limitGitConcurrency) r.Get("/"+tangled.GitTempGetArchiveNSID, x.GetArchive) r.Get("/"+tangled.GitTempGetBlobNSID, x.GetBlob) @@ -98,6 +103,12 @@ func (x *Xrpc) Router() http.Handler { r.Get("/"+tangled.GitTempListCommitsNSID, x.ListCommits) r.Get("/"+tangled.GitTempListLanguagesNSID, x.ListLanguages) r.Get("/"+tangled.GitTempListTagsNSID, x.ListTags) + }) + + r.Group(func(r chi.Router) { + r.Use(x.inflight.middleware) + r.Use(x.forwardSuspended) + r.Post("/"+tangled.SyncRequestCrawlNSID, x.RequestCrawl) r.Get("/sh.tangled.git.temp2.getBlame", x.proxyV2) @@ -196,3 +207,23 @@ func setCacheControl(w http.ResponseWriter, refs ...string) { w.Header().Set("Cache-Control", "public, no-cache") } } + +func (x *Xrpc) limitGitConcurrency(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second) + defer cancel() + + if err := x.gitSem.Acquire(ctx, 1); err != nil { + x.logger.Warn("git concurrency limit reached, shedding request", "path", r.URL.Path, "err", err) + w.Header().Set("Retry-After", "1") + writeJson(w, http.StatusTooManyRequests, atclient.ErrorBody{ + Name: "RateLimitExceeded", + Message: "server busy, please retry shortly", + }) + return + } + defer x.gitSem.Release(1) + + next.ServeHTTP(w, r) + }) +} -- 2.51.2