feat(concurrency): provider-wide lens budget, drop the model cap
Concurrency was two multiplicative gates in two processes: entrypoint.sh capped MODELS-at-once per provider (GADFLY_PROVIDER_CONCURRENCY) while each model's binary separately capped its own lenses (GADFLY_LENS_CONCURRENCY). A model therefore held its whole model-slot until its LAST lens finished, stalling the next model even with idle lens capacity. Collapse to one throttle: a provider-wide lens budget shared across all of that provider's models. entrypoint now runs every model in a lane at once and seeds a single cross-process permit pool per lane (a dir of N flock files, sized by GADFLY_PROVIDER_LENS_CONCURRENCY -> GADFLY_LENS_CONCURRENCY). Each lens pass (review+recheck) acquires a permit before it runs and releases it after, so a model winding down immediately yields its freed permits to another model's queued lenses. flock auto-releases on process death, so a killed/crashed model can't leak budget. - cmd/gadfly/lenssem.go: the flock permit pool (+ lenssem_test.go). - main.go: runSpecialists holds a shared permit per lens; fanout sized to the budget so a lone model can use all of it. Falls back to the in-process limit when no pool is set (local runs, tests). - entrypoint.sh: drop provider_cap/DEFAULT_CONC; run_lane runs all models and seeds the per-lane pool. - GADFLY_PROVIDER_CONCURRENCY / GADFLY_CONCURRENCY are now ignored; the reusable workflow marks provider_concurrency deprecated and stops forwarding it. Docs (README, CLAUDE.md, examples) updated per the maintenance rule. Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
@@ -0,0 +1,86 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"syscall"
|
||||
"time"
|
||||
)
|
||||
|
||||
// lensSem is the PROVIDER-WIDE lens-permit pool. Historically each model's binary
|
||||
// throttled its own lenses in-process (GADFLY_LENS_CONCURRENCY) while entrypoint.sh
|
||||
// separately capped how many MODELS ran at once — two multiplicative gates in
|
||||
// different processes. That let a model hold its whole slot until its last lens
|
||||
// finished, stalling the next model even with idle lens capacity.
|
||||
//
|
||||
// Instead, entrypoint.sh now runs all of a provider's models concurrently and
|
||||
// seeds ONE directory of N permit files per provider (N = the provider's lens
|
||||
// budget). Every model process in that lane draws from the same pool: a lens pass
|
||||
// (review + recheck) acquires a permit before it runs and releases it after, so a
|
||||
// model winding down immediately yields its freed permits to another model's
|
||||
// queued lenses. Permits are held with flock, which the kernel drops when the
|
||||
// holding process exits — so a killed/crashed model frees its permits for free.
|
||||
type lensSem struct {
|
||||
dir string
|
||||
size int
|
||||
}
|
||||
|
||||
// lensSemPollInterval is how often a blocked acquirer re-sweeps the permit files.
|
||||
// Lens passes run for many seconds to minutes, so a coarse poll adds negligible
|
||||
// latency while keeping the mechanism a few lines of stdlib (no IPC primitives).
|
||||
const lensSemPollInterval = 150 * time.Millisecond
|
||||
|
||||
// activeLensSem returns the shared semaphore configured by entrypoint.sh, or nil
|
||||
// when it isn't set (standalone/local runs, tests, or an older entrypoint) — in
|
||||
// which case runSpecialists falls back to the in-process fanout limit alone, i.e.
|
||||
// the pre-existing per-model behavior.
|
||||
func activeLensSem() *lensSem {
|
||||
dir := os.Getenv("GADFLY_LENS_SEM_DIR")
|
||||
if dir == "" {
|
||||
return nil
|
||||
}
|
||||
n := envInt("GADFLY_LENS_SEM_SIZE", 0)
|
||||
if n < 1 {
|
||||
return nil
|
||||
}
|
||||
return &lensSem{dir: dir, size: n}
|
||||
}
|
||||
|
||||
// acquire blocks until a permit is free (or ctx is done) and returns a release
|
||||
// func. On cancellation it returns a no-op release plus ctx.Err(), so callers can
|
||||
// surface the lens as "did not run" rather than leaking a permit.
|
||||
func (s *lensSem) acquire(ctx context.Context) (func(), error) {
|
||||
for {
|
||||
if release, ok := s.tryAcquire(); ok {
|
||||
return release, nil
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return func() {}, ctx.Err()
|
||||
case <-time.After(lensSemPollInterval):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// tryAcquire makes one non-blocking sweep over the permit files, returning a
|
||||
// release func for the first one it locks. The permit files are created lazily
|
||||
// (entrypoint.sh only guarantees the directory exists), so a fresh pool needs no
|
||||
// seeding step beyond mkdir. Closing the *os.File in the returned func drops the
|
||||
// flock (the lock lives on the open file description).
|
||||
func (s *lensSem) tryAcquire() (func(), bool) {
|
||||
for i := 0; i < s.size; i++ {
|
||||
path := filepath.Join(s.dir, fmt.Sprintf("permit.%d", i))
|
||||
f, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE, 0o600)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
if err := syscall.Flock(int(f.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil {
|
||||
f.Close()
|
||||
continue
|
||||
}
|
||||
return func() { f.Close() }, true
|
||||
}
|
||||
return nil, false
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// flock permits are per open-file-description, so two separate opens of the same
|
||||
// permit file conflict even within one process — the exhaustion, blocking, and
|
||||
// max-in-flight tests below therefore exercise the same semantics a real
|
||||
// multi-process lane would see.
|
||||
|
||||
func TestActiveLensSemUnset(t *testing.T) {
|
||||
t.Setenv("GADFLY_LENS_SEM_DIR", "")
|
||||
if s := activeLensSem(); s != nil {
|
||||
t.Fatalf("expected nil sem when GADFLY_LENS_SEM_DIR unset, got %+v", s)
|
||||
}
|
||||
// A dir with a zero/blank size is inert too (must not divide the pool by 0).
|
||||
t.Setenv("GADFLY_LENS_SEM_DIR", t.TempDir())
|
||||
t.Setenv("GADFLY_LENS_SEM_SIZE", "0")
|
||||
if s := activeLensSem(); s != nil {
|
||||
t.Fatalf("expected nil sem when size < 1, got %+v", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLensSemExhaustionAndRelease(t *testing.T) {
|
||||
s := &lensSem{dir: t.TempDir(), size: 2}
|
||||
|
||||
r1, ok := s.tryAcquire()
|
||||
if !ok {
|
||||
t.Fatal("first acquire should succeed")
|
||||
}
|
||||
r2, ok := s.tryAcquire()
|
||||
if !ok {
|
||||
t.Fatal("second acquire should succeed")
|
||||
}
|
||||
if _, ok := s.tryAcquire(); ok {
|
||||
t.Fatal("third acquire should fail: pool of 2 is exhausted")
|
||||
}
|
||||
|
||||
r1() // free one permit
|
||||
r3, ok := s.tryAcquire()
|
||||
if !ok {
|
||||
t.Fatal("acquire should succeed after a release")
|
||||
}
|
||||
r2()
|
||||
r3()
|
||||
}
|
||||
|
||||
func TestLensSemAcquireBlocksThenCancels(t *testing.T) {
|
||||
s := &lensSem{dir: t.TempDir(), size: 1}
|
||||
|
||||
// Immediate success while a permit is free.
|
||||
release, err := s.acquire(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("acquire on a free pool: %v", err)
|
||||
}
|
||||
|
||||
// With the only permit held, acquire must block until ctx expires.
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
start := time.Now()
|
||||
if _, err := s.acquire(ctx); err == nil {
|
||||
t.Fatal("acquire on an exhausted pool should return ctx error, not a permit")
|
||||
}
|
||||
if waited := time.Since(start); waited < 50*time.Millisecond {
|
||||
t.Fatalf("acquire returned too fast (%v); it should have blocked on the full pool", waited)
|
||||
}
|
||||
release()
|
||||
}
|
||||
|
||||
func TestLensSemNeverExceedsSize(t *testing.T) {
|
||||
const size = 3
|
||||
s := &lensSem{dir: t.TempDir(), size: size}
|
||||
|
||||
var inFlight, peak int64
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < 12; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
release, err := s.acquire(context.Background())
|
||||
if err != nil {
|
||||
t.Errorf("acquire: %v", err)
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
n := atomic.AddInt64(&inFlight, 1)
|
||||
for {
|
||||
p := atomic.LoadInt64(&peak)
|
||||
if n <= p || atomic.CompareAndSwapInt64(&peak, p, n) {
|
||||
break
|
||||
}
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
atomic.AddInt64(&inFlight, -1)
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if peak > size {
|
||||
t.Fatalf("max in-flight lens permits = %d, exceeds budget %d", peak, size)
|
||||
}
|
||||
if peak == 0 {
|
||||
t.Fatal("no permits were ever acquired")
|
||||
}
|
||||
}
|
||||
+46
-23
@@ -49,15 +49,18 @@
|
||||
// GADFLY_RECHECK set to 0/false to skip the recheck pass (optional, default on).
|
||||
// GADFLY_RECHECK_MAX_STEPS recheck-pass step cap (optional, default 16).
|
||||
// GADFLY_TIMEOUT_SECS overall deadline in seconds, shared by both passes (optional, default 300).
|
||||
// GADFLY_LENS_CONCURRENCY how many specialist lenses run concurrently within this
|
||||
// model (optional, default 1 = sequential). Total in-flight
|
||||
// model requests ≈ this × entrypoint.sh's per-provider model
|
||||
// concurrency, so keep the product within the backend's budget.
|
||||
// GADFLY_LENS_CONCURRENCY how many specialist lenses run concurrently (optional,
|
||||
// default 1 = sequential). Under entrypoint.sh this is the
|
||||
// PROVIDER-WIDE lens budget, shared across all of that
|
||||
// provider's models via a permit pool (GADFLY_LENS_SEM_DIR),
|
||||
// so it bounds total lens passes in flight per provider.
|
||||
// GADFLY_PROVIDER_LENS_CONCURRENCY per-provider override for the above, as a
|
||||
// "provider=N,provider=N" map keyed by the SAME provider
|
||||
// lanes as GADFLY_PROVIDER_CONCURRENCY (e.g.
|
||||
// "ollama-cloud=3,m1=1"). Wins over GADFLY_LENS_CONCURRENCY
|
||||
// for the model's provider; falls back to it otherwise.
|
||||
// "provider=N,provider=N" map keyed by the provider lanes
|
||||
// (e.g. "ollama-cloud=3,m1=1"). Wins over
|
||||
// GADFLY_LENS_CONCURRENCY for the model's provider.
|
||||
// GADFLY_LENS_SEM_DIR / GADFLY_LENS_SEM_SIZE set by entrypoint.sh: the shared
|
||||
// cross-process lens-permit pool (dir of N flock files)
|
||||
// and its size. Unset => in-process lensConcurrency only.
|
||||
// GADFLY_MAX_DIFF_CHARS diff chars embedded in the review prompt (optional, default 60000;
|
||||
// the full diff is reachable via the paginated get_diff tool).
|
||||
//
|
||||
@@ -248,19 +251,29 @@ func run() error {
|
||||
// own per-lens timeout (reviewWithSpecialist) and the lenses only read the
|
||||
// immutable repoFS, so concurrency simply overlaps independent passes.
|
||||
//
|
||||
// Caution: this fans out WITHIN one model. It multiplies with entrypoint.sh's
|
||||
// per-provider model concurrency, so total concurrent backend requests ≈
|
||||
// (models at once) × (lenses at once). To fan lenses out without oversubscribing
|
||||
// the backend, run models one at a time (provider lane cap 1) and raise this.
|
||||
// Provider-wide throttling: when entrypoint.sh runs several of a provider's
|
||||
// models at once it seeds a shared lens-permit pool (activeLensSem) that every
|
||||
// model's lenses draw from, so the real cap is total lens passes in flight per
|
||||
// provider — not (models at once) × (lenses at once). Absent that pool (local
|
||||
// runs, tests) this falls back to the in-process lensConcurrency() limit alone.
|
||||
func runSpecialists(eng reviewEngine, base string, specialists []Specialist, task, diff string) []specialistResult {
|
||||
// Optional live status board: publishes this model's per-lens progress to a
|
||||
// file the entrypoint board renders. Inert (no-op) unless GADFLY_STATUS_FILE
|
||||
// is set, so plain runs are unaffected.
|
||||
sw := newStatusWriter(os.Getenv("GADFLY_MODEL"), modelProvider(), specialists)
|
||||
|
||||
// The cross-process pool (if any) is the real ceiling; size the in-process
|
||||
// fanout to it so a lone model in its lane can use the whole provider budget,
|
||||
// while extra goroutines simply block in sem.acquire until a permit frees.
|
||||
sem := activeLensSem()
|
||||
maxConcurrent := lensConcurrency()
|
||||
if sem != nil {
|
||||
maxConcurrent = sem.size
|
||||
}
|
||||
|
||||
fanResults := fanout.Run(context.Background(), specialists, fanout.Options[Specialist]{
|
||||
MaxConcurrent: lensConcurrency(),
|
||||
}, func(_ context.Context, sp Specialist) (res specialistResult, _ error) {
|
||||
MaxConcurrent: maxConcurrent,
|
||||
}, func(ctx context.Context, sp Specialist) (res specialistResult, _ error) {
|
||||
// A panic in one lens must not crash the whole binary (which would kill
|
||||
// every other lens's output) or leave this lens stuck at "running" on the
|
||||
// status board. fanout does not recover fn panics, so we do it here:
|
||||
@@ -271,6 +284,17 @@ func runSpecialists(eng reviewEngine, base string, specialists []Specialist, tas
|
||||
sw.set(sp.Name, lensFinished, "", true)
|
||||
}
|
||||
}()
|
||||
// Hold a provider-wide permit for this lens's whole review+recheck pass.
|
||||
// While waiting the lens stays "queued" on the board; cancellation before
|
||||
// a permit frees surfaces as a did-not-run lens rather than a leaked slot.
|
||||
if sem != nil {
|
||||
release, err := sem.acquire(ctx)
|
||||
if err != nil {
|
||||
sw.set(sp.Name, lensFinished, "", true)
|
||||
return specialistResult{spec: sp, out: fmt.Sprintf("⚠️ This reviewer did not run: %v", err), verdict: verdictUnknown, errored: true}, nil
|
||||
}
|
||||
defer release()
|
||||
}
|
||||
sw.set(sp.Name, lensRunning, "", false)
|
||||
out, errored := reviewWithSpecialist(eng, base, sp, task, diff)
|
||||
v := parseVerdict(out)
|
||||
@@ -293,14 +317,13 @@ func runSpecialists(eng reviewEngine, base string, specialists []Specialist, tas
|
||||
return results
|
||||
}
|
||||
|
||||
// lensConcurrency resolves how many specialist lenses run at once for THIS run's
|
||||
// model. It mirrors entrypoint.sh's per-provider MODEL concurrency: a
|
||||
// per-provider override in GADFLY_PROVIDER_LENS_CONCURRENCY ("provider=N,...")
|
||||
// wins for the model's provider, otherwise the GADFLY_LENS_CONCURRENCY scalar
|
||||
// (default 1). The provider is resolved by modelProvider() — the SAME lane rule
|
||||
// entrypoint uses for GADFLY_PROVIDER_CONCURRENCY — so e.g.
|
||||
// "ollama-cloud=3,m1=1" fans cloud lenses out while keeping a slow local box
|
||||
// serial, exactly the way the model map does for whole models.
|
||||
// lensConcurrency resolves the lens budget for THIS run's provider: a per-provider
|
||||
// override in GADFLY_PROVIDER_LENS_CONCURRENCY ("provider=N,...") wins for the
|
||||
// model's provider (resolved by modelProvider()), otherwise the
|
||||
// GADFLY_LENS_CONCURRENCY scalar (default 1). Under entrypoint.sh the SAME value
|
||||
// seeds the shared cross-process permit pool (activeLensSem), so it is the
|
||||
// provider-wide budget rather than a per-model one; standalone it caps the single
|
||||
// model's in-process fanout.
|
||||
func lensConcurrency() int {
|
||||
if n, ok := providerOverride("GADFLY_PROVIDER_LENS_CONCURRENCY", modelProvider()); ok {
|
||||
return n
|
||||
@@ -310,7 +333,7 @@ func lensConcurrency() int {
|
||||
|
||||
// providerOverride parses a "provider=N,provider=N" env map and returns the
|
||||
// value for provider when present and valid (>0). Mirrors entrypoint.sh's
|
||||
// provider_cap lookup so the two concurrency maps share one syntax.
|
||||
// provider_lens_cap lookup so the two share one syntax.
|
||||
func providerOverride(envName, provider string) (int, bool) {
|
||||
for _, item := range strings.Split(os.Getenv(envName), ",") {
|
||||
k, v, ok := strings.Cut(item, "=")
|
||||
|
||||
+2
-2
@@ -162,8 +162,8 @@ func buildSpec(provider, model string) string {
|
||||
// entrypoint.sh's provider_of: the segment before the first "/" in GADFLY_MODEL,
|
||||
// else GADFLY_PROVIDER, else the default (ollama-cloud). The binary reviews one
|
||||
// model per invocation, so this is that model's provider — used to resolve
|
||||
// per-provider policy (e.g. lens concurrency) against the SAME provider keys
|
||||
// entrypoint uses for GADFLY_PROVIDER_CONCURRENCY.
|
||||
// per-provider policy (e.g. the lens budget) against the SAME provider keys
|
||||
// entrypoint uses for GADFLY_PROVIDER_LENS_CONCURRENCY.
|
||||
func modelProvider() string {
|
||||
model := strings.TrimSpace(os.Getenv("GADFLY_MODEL"))
|
||||
if pfx, _, ok := strings.Cut(model, "/"); ok {
|
||||
|
||||
Reference in New Issue
Block a user