mirror of
https://github.com/Kpa-clawbot/meshcore-analyzer.git
synced 2026-10-10 16:57:28 +00:00
fix(api): cap concurrent cold reach scans at 2, answer 429 beyond that (#2127)
`GET /api/nodes/{pubkey}/reach` is unauthenticated and a cold-cache
request runs a full scan. `singleflight` only collapses identical keys,
so distinct `(pubkey, days)` requests each start a scan, and they share
the 4-connection SQLite pool with every other handler. A handful of
requests for different nodes can hold every reader and stall the whole
API.
**Fix:** a small semaphore allows two cold scans at once. A cold request
that finds both slots busy gets `429` with `Retry-After: 5` instead of
queueing. Cached answers are unaffected, and the two in-flight scans
still complete normally.
**Tests:** `node_reach_concurrency_test.go` — with both slots held a
cold request returns 429 with `Retry-After`; after release the same
request runs normally. Full server suite passes.
**Note:** our instance runs a fork that already bounds reach work
differently (an async job queue), so this exact change is not what we
run in production. The test suite is the verification here.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-authored-by: Claude Mythos 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Mythos 5.1
parent
0a6870c508
commit
caa66356bd
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log"
|
||||
"net/http"
|
||||
"sort"
|
||||
@@ -331,6 +332,12 @@ type reachCacheEntry struct {
|
||||
type reachState struct {
|
||||
cacheMu sync.RWMutex
|
||||
cache map[string]reachCacheEntry
|
||||
// buildSem bounds how many cold-cache reach scans run at the same time.
|
||||
// singleflight only collapses identical keys; distinct (pubkey, days)
|
||||
// keys each start a scan, and the endpoint is unauthenticated, so
|
||||
// without a cap a handful of requests can hold every SQLite reader.
|
||||
buildSemOnce sync.Once
|
||||
buildSem chan struct{}
|
||||
// sf dedups concurrent cold-cache requests for the same key so N
|
||||
// simultaneous callers run the scan + attribution once, not N times.
|
||||
sf singleflight.Group
|
||||
@@ -492,6 +499,11 @@ func (s *Server) handleNodeReach(w http.ResponseWriter, r *http.Request) {
|
||||
if raw, ok := s.reachCacheGet(cacheKey); ok {
|
||||
return raw, nil
|
||||
}
|
||||
release, ok := s.reachAcquireBuildSlot()
|
||||
if !ok {
|
||||
return nil, errReachBusy
|
||||
}
|
||||
defer release()
|
||||
resp, ok, cErr := s.computeNodeReach(r.Context(), pubkey, days)
|
||||
if cErr != nil {
|
||||
// Real backend failure (e.g. DB scan exploded) — propagate so the
|
||||
@@ -510,6 +522,11 @@ func (s *Server) handleNodeReach(w http.ResponseWriter, r *http.Request) {
|
||||
s.reachCachePut(cacheKey, raw)
|
||||
return raw, nil
|
||||
})
|
||||
if errors.Is(err, errReachBusy) {
|
||||
w.Header().Set("Retry-After", strconv.Itoa(reachBusyRetryAfterSeconds))
|
||||
writeError(w, http.StatusTooManyRequests, "too many reach computations in progress, retry shortly")
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
writeError(w, 500, "reach computation failed")
|
||||
return
|
||||
@@ -792,3 +809,26 @@ func (s *Server) scanReachRows(ctx context.Context, tokens map[string]bool, sinc
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// reachMaxConcurrentBuilds is the number of cold-cache reach scans allowed at
|
||||
// once. Two keeps half of the default 4-connection SQLite pool free for
|
||||
// every other handler while a scan runs.
|
||||
const reachMaxConcurrentBuilds = 2
|
||||
|
||||
// reachBusyRetryAfterSeconds is the Retry-After we send with a 429.
|
||||
const reachBusyRetryAfterSeconds = 5
|
||||
|
||||
// errReachBusy is returned inside the singleflight when no build slot is free.
|
||||
var errReachBusy = errors.New("reach: too many concurrent builds")
|
||||
|
||||
// reachAcquireBuildSlot takes a build slot without blocking. It returns a
|
||||
// release func and true, or nil and false when all slots are busy.
|
||||
func (s *Server) reachAcquireBuildSlot() (func(), bool) {
|
||||
s.reach.buildSemOnce.Do(func() { s.reach.buildSem = make(chan struct{}, reachMaxConcurrentBuilds) })
|
||||
select {
|
||||
case s.reach.buildSem <- struct{}{}:
|
||||
return func() { <-s.reach.buildSem }, true
|
||||
default:
|
||||
return nil, false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// GET /api/nodes/{pubkey}/reach is unauthenticated and a cold-cache request
|
||||
// starts a scan. Only reachMaxConcurrentBuilds scans may run at once; the
|
||||
// next cold request gets 429 with Retry-After instead of queueing on the
|
||||
// shared SQLite pool.
|
||||
func TestNodeReachConcurrentBuildCap(t *testing.T) {
|
||||
srv, router := setupTestServerWithAPIKey(t, "")
|
||||
const key = "ab00000000000000000000000000000000000000000000000000000000000001"
|
||||
|
||||
get := func() *httptest.ResponseRecorder {
|
||||
req := httptest.NewRequest("GET", "/api/nodes/"+key+"/reach?days=7", nil)
|
||||
w := httptest.NewRecorder()
|
||||
router.ServeHTTP(w, req)
|
||||
return w
|
||||
}
|
||||
|
||||
// Fill every build slot as if scans were in flight.
|
||||
var releases []func()
|
||||
for i := 0; i < reachMaxConcurrentBuilds; i++ {
|
||||
rel, ok := srv.reachAcquireBuildSlot()
|
||||
if !ok {
|
||||
t.Fatalf("slot %d should be free", i)
|
||||
}
|
||||
releases = append(releases, rel)
|
||||
}
|
||||
if w := get(); w.Code != http.StatusTooManyRequests {
|
||||
t.Fatalf("expected 429 while all build slots are busy, got %d (body: %s)", w.Code, w.Body.String())
|
||||
} else if w.Header().Get("Retry-After") == "" {
|
||||
t.Fatalf("429 must carry Retry-After")
|
||||
}
|
||||
|
||||
// Free the slots: the same request now runs the scan. The key is not a
|
||||
// known node, so the normal answer is 404 — the point is that it is no
|
||||
// longer 429.
|
||||
for _, rel := range releases {
|
||||
rel()
|
||||
}
|
||||
if w := get(); w.Code == http.StatusTooManyRequests {
|
||||
t.Fatalf("still 429 after slots were released (body: %s)", w.Body.String())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user