Build & push image / build-and-push (push) Successful in 33s
Makes gadfly a consumer of executus (run.Executor compaction/bounding/budget/critic + fanout) and fixes the large-PR token burn in size-gated layers: paginated get_diff, downshift above GADFLY_HUGE_DIFF_BYTES, and a swarm-wide GADFLY_PR_BUDGET_SECS backstop. Small PRs untouched; advisory-only and the static binary preserved. Dogfood swarm reviewed it (6 models, 21 real findings graded + folded in). Co-authored-by: Steve Dudenhoeffer <[email protected]> Co-committed-by: Steve Dudenhoeffer <[email protected]>
372 lines
14 KiB
Go
372 lines
14 KiB
Go
package main
|
||
|
||
// executus.go wires gadfly's agentic review path onto the executus run kernel
|
||
// (gitea.stevedudenhoeffer.com/steve/executus), layered above majordomo. The
|
||
// majordomoEngine no longer drives majordomo's agent loop directly; it builds a
|
||
// run.Executor that gives gadfly, for free:
|
||
//
|
||
// - context compaction (executus/compact): once the transcript a step would
|
||
// SEND crosses a token threshold derived from the model's real context
|
||
// window, the runaway middle is folded into a one-paragraph summary by a
|
||
// cheap summarizer model — so a big diff + accumulating read_file/grep
|
||
// results can't balloon every re-sent step (the large-PR burn).
|
||
// - run bounding + a per-PR spend budget (executus/run Ports.Budget): a hard
|
||
// token/seconds ceiling so a pathological PR can't drain the usage block.
|
||
// - the wrap-up nudge, re-expressed as an executus Critic (Ports.Critic): the
|
||
// steer that tells a step-hungry model to stop investigating and write its
|
||
// answer is now the critic seam, not a bespoke RunOption.
|
||
//
|
||
// Everything degrades to today's behavior when unconfigured: nil summarizer or a
|
||
// 0 context window disables compaction; nil budget disables the ceiling; the
|
||
// claude-code engine shells out and is unaffected by any of this.
|
||
//
|
||
// gadfly keeps its own model.go resolution (so GADFLY_ENDPOINT_<NAME> http
|
||
// aliases, failover chains, and the claude-code engine all survive) — the
|
||
// run.Executor is handed gadfly's already-resolved model via a trivial resolver,
|
||
// not routed through executus's tier table.
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"os"
|
||
"strconv"
|
||
"strings"
|
||
"sync"
|
||
"sync/atomic"
|
||
"time"
|
||
|
||
llm "gitea.stevedudenhoeffer.com/steve/majordomo/llm"
|
||
|
||
"gitea.stevedudenhoeffer.com/steve/executus/compact"
|
||
"gitea.stevedudenhoeffer.com/steve/executus/model"
|
||
exrun "gitea.stevedudenhoeffer.com/steve/executus/run"
|
||
exectool "gitea.stevedudenhoeffer.com/steve/executus/tool"
|
||
)
|
||
|
||
const (
|
||
// defaultCompactRatio is the fraction of the model's context window at which
|
||
// compaction fires. It is deliberately LOWER than executus's own 0.7 default:
|
||
// on the large-PR burn the per-step transcript is the embedded diff (~17K) plus
|
||
// accumulating read_file results, which rarely reaches 0.7×262K≈183K — so a
|
||
// 0.7 threshold never bites. ~0.45×262K≈118K folds the runaway middle while a
|
||
// transcript is still well under the cap. Override with GADFLY_COMPACT_RATIO.
|
||
defaultCompactRatio = 0.45
|
||
// defaultCompactKeepRecent / defaultCompactSummaryWords mirror executus's own
|
||
// compactor defaults; surfaced as gadfly env knobs for tuning.
|
||
defaultCompactKeepRecent = 8
|
||
defaultCompactSummaryWords = 200
|
||
// contextTokenLookupTimeout bounds the one-shot /api/show call that resolves a
|
||
// cloud model's context window at executor-build time. Kept short so a slow or
|
||
// unreachable endpoint adds at most this to startup before degrading to
|
||
// no-compaction (rather than the provider cache's default 15s).
|
||
contextTokenLookupTimeout = 5 * time.Second
|
||
)
|
||
|
||
// runSeq mints a unique-per-process RunID suffix for each executor run so audit
|
||
// and the run kernel can tell one pass from another within a binary process.
|
||
var runSeq atomic.Uint64
|
||
|
||
// wrappedTool adapts an already-built majordomo llm.Tool (gadfly's sandboxed
|
||
// read_file/grep/get_diff/… closures over the repoFS) to executus's tool.Tool
|
||
// interface so the run kernel can build a toolbox from them by name. gadfly's
|
||
// tools need no caller/channel identity, so BuildLLM ignores the Invocation and
|
||
// returns the pre-built tool; Permission is the zero value (private, ungated).
|
||
type wrappedTool struct{ t llm.Tool }
|
||
|
||
func (w wrappedTool) Name() string { return w.t.Name }
|
||
func (w wrappedTool) Description() string { return w.t.Description }
|
||
func (w wrappedTool) Permission() exectool.Permission { return exectool.Permission{} }
|
||
func (w wrappedTool) BuildLLM(_ exectool.Invocation) llm.Tool { return w.t }
|
||
|
||
// gadflyToolRegistry registers the repo's read-only tools (plus the optional
|
||
// delegate_investigation worker tool) in a fresh executus tool.Registry and
|
||
// returns it along with the tool names for RunnableAgent.LowLevelTools.
|
||
func gadflyToolRegistry(fs *repoFS) (exectool.Registry, []string, error) {
|
||
reg := exectool.NewRegistry()
|
||
tools := fs.allTools()
|
||
names := make([]string, 0, len(tools))
|
||
for _, t := range tools {
|
||
if err := reg.Register(wrappedTool{t: t}); err != nil {
|
||
return nil, nil, fmt.Errorf("register tool %q: %w", t.Name, err)
|
||
}
|
||
names = append(names, t.Name)
|
||
}
|
||
return reg, names, nil
|
||
}
|
||
|
||
// gadflyBudget is gadfly's per-PR spend ceiling, satisfying run.Ports.Budget.
|
||
// It gates a run BEFORE it makes any model call (Check) once the process has
|
||
// spent its token or wall-clock allowance on this PR. Tokens are fed in
|
||
// out-of-band via addUsage (the Budget interface's Commit only carries seconds);
|
||
// the engine calls addUsage after each pass with run.Result.Usage. A nil
|
||
// *gadflyBudget is never installed — caps of 0 mean "unlimited", so the port is
|
||
// only wired when at least one cap is set.
|
||
//
|
||
// The guard is PASS-granular: Check runs before each pass, so it stops the NEXT
|
||
// pass once the budget is spent but cannot abort a single runaway pass mid-flight.
|
||
// The swarm-wide GADFLY_PR_BUDGET_SECS wall-clock backstop (entrypoint.sh) is what
|
||
// bounds a mid-pass runaway.
|
||
type gadflyBudget struct {
|
||
mu sync.Mutex
|
||
maxTokens int64
|
||
maxSeconds float64
|
||
tokens int64
|
||
seconds float64
|
||
}
|
||
|
||
// newPRBudget builds the per-PR budget from env, or nil when neither cap is set
|
||
// (the default — the swarm-wide ceiling lives in entrypoint.sh; this is the
|
||
// per-process belt to its suspenders).
|
||
func newPRBudget() *gadflyBudget {
|
||
toks := envInt("GADFLY_PR_TOKEN_BUDGET", 0)
|
||
secs := envInt("GADFLY_PR_TIME_BUDGET_SECS", 0)
|
||
if toks <= 0 && secs <= 0 {
|
||
return nil
|
||
}
|
||
return &gadflyBudget{maxTokens: int64(toks), maxSeconds: float64(secs)}
|
||
}
|
||
|
||
func (b *gadflyBudget) Check(_ context.Context, _ string) error {
|
||
if b == nil {
|
||
return nil
|
||
}
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
if b.maxTokens > 0 && b.tokens >= b.maxTokens {
|
||
return fmt.Errorf("gadfly: per-PR token budget exhausted (%d/%d)", b.tokens, b.maxTokens)
|
||
}
|
||
if b.maxSeconds > 0 && b.seconds >= b.maxSeconds {
|
||
return fmt.Errorf("gadfly: per-PR time budget exhausted (%.0f/%.0fs)", b.seconds, b.maxSeconds)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (b *gadflyBudget) Commit(_ context.Context, _ string, runtimeSeconds float64) {
|
||
if b == nil {
|
||
return
|
||
}
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.seconds += runtimeSeconds
|
||
}
|
||
|
||
// addUsage records a finished pass's token spend toward the budget. Safe on nil.
|
||
func (b *gadflyBudget) addUsage(u llm.Usage) {
|
||
if b == nil {
|
||
return
|
||
}
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.tokens += int64(u.InputTokens) + int64(u.OutputTokens)
|
||
}
|
||
|
||
// wrapUpCritic re-expresses gadfly's wrap-up nudge as an executus run.Critic:
|
||
// once a run comes within wrapUpReserve steps of its cap, Steer() injects the
|
||
// "stop calling tools and write your final answer" message so a thorough model
|
||
// spends its last steps finalizing instead of hard-failing empty. It sets no
|
||
// hard deadline (Deadline()==zero) and never raises the step ceiling
|
||
// (MaxSteps()==0, defer to the run's MaxIterations) — it is purely the nudge.
|
||
type wrapUpCritic struct{ reserve int }
|
||
|
||
func (c *wrapUpCritic) Monitor(_ context.Context, info exrun.RunInfo, _ time.Duration) exrun.CriticHandle {
|
||
return &wrapUpHandle{maxSteps: info.MaxIterations, reserve: c.reserve}
|
||
}
|
||
|
||
type wrapUpHandle struct {
|
||
mu sync.Mutex
|
||
maxSteps int
|
||
reserve int
|
||
done int // steps completed so far
|
||
nudged bool
|
||
}
|
||
|
||
func (h *wrapUpHandle) RecordStep(iter int, _ *llm.Response) {
|
||
h.mu.Lock()
|
||
h.done = iter + 1
|
||
h.mu.Unlock()
|
||
}
|
||
func (h *wrapUpHandle) RecordToolStart(string, string) {}
|
||
func (h *wrapUpHandle) Steer() []llm.Message {
|
||
h.mu.Lock()
|
||
defer h.mu.Unlock()
|
||
at := h.maxSteps - h.reserve
|
||
if at < 1 {
|
||
at = 1
|
||
}
|
||
if !h.nudged && h.maxSteps > 0 && h.done >= at {
|
||
h.nudged = true
|
||
return []llm.Message{llm.UserText(wrapUpInstruction)}
|
||
}
|
||
return nil
|
||
}
|
||
func (h *wrapUpHandle) Deadline() time.Time { return time.Time{} }
|
||
func (h *wrapUpHandle) MaxSteps() int { return 0 }
|
||
func (h *wrapUpHandle) KillCause() error { return nil }
|
||
func (h *wrapUpHandle) Stop() {}
|
||
|
||
// reviewExecutor bundles a run.Executor with the per-run wiring the engine needs
|
||
// for each pass (the tool names to expose, the model spec to report as the tier,
|
||
// the per-PR caller id, and the budget to feed token usage into).
|
||
type reviewExecutor struct {
|
||
ex *exrun.Executor
|
||
toolNames []string
|
||
modelSpec string
|
||
callerID string
|
||
budget *gadflyBudget
|
||
}
|
||
|
||
// newReviewExecutor builds the run.Executor for the in-process majordomo review
|
||
// path. mdl is gadfly's already-resolved review model; summarizer is the cheap
|
||
// model the compactor uses (nil disables compaction). Compaction also needs the
|
||
// model's context window (resolved once here, not per pass); a 0 window likewise
|
||
// disables it. The budget (may be nil) becomes the run.Ports.Budget gate.
|
||
func newReviewExecutor(fs *repoFS, mdl, summarizer llm.Model, modelSpec string, budget *gadflyBudget) (*reviewExecutor, error) {
|
||
reg, names, err := gadflyToolRegistry(fs)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
// gadfly resolves exactly one model per process; the run kernel's resolver
|
||
// just hands that model back regardless of the tier string it is asked for.
|
||
modelsResolver := func(ctx context.Context, _ string) (context.Context, llm.Model, error) {
|
||
return ctx, mdl, nil
|
||
}
|
||
|
||
var compactor compact.CompactorFactory
|
||
var ctxTokens func(string) int
|
||
if summarizer != nil && compactionEnabled() {
|
||
if window := resolveContextTokens(modelSpec); window > 0 {
|
||
sumResolver := func(ctx context.Context, _ string) (context.Context, llm.Model, error) {
|
||
return ctx, summarizer, nil
|
||
}
|
||
compactor = compact.NewCompactor(compact.CompactorConfig{
|
||
Models: sumResolver,
|
||
KeepRecent: envInt("GADFLY_COMPACT_KEEP_RECENT", defaultCompactKeepRecent),
|
||
SummaryWordCap: envInt("GADFLY_COMPACT_SUMMARY_WORDS", defaultCompactSummaryWords),
|
||
})
|
||
ctxTokens = func(string) int { return window } // memoized: one window per process
|
||
}
|
||
}
|
||
|
||
var ports exrun.Ports
|
||
ports.Critic = &wrapUpCritic{reserve: wrapUpReserve()}
|
||
if budget != nil {
|
||
ports.Budget = budget
|
||
}
|
||
|
||
cfg := exrun.Config{
|
||
Registry: reg,
|
||
Models: modelsResolver,
|
||
Compactor: compactor,
|
||
ContextTokens: ctxTokens,
|
||
Defaults: exrun.Defaults{
|
||
// MaxIterations/MaxRuntime are intentionally omitted: every pass sets its
|
||
// own per-run cap on the RunnableAgent below (the review and recheck caps
|
||
// differ), so a Defaults value here would always be overridden — dead. This
|
||
// leaves only the cross-pass guards + the compaction ratio.
|
||
MaxConsecutiveToolErrors: 4,
|
||
MaxSameToolCallRepeats: 4,
|
||
CompactionThresholdRatio: compactRatio(),
|
||
FallbackTier: modelSpec,
|
||
},
|
||
Ports: ports,
|
||
}
|
||
return &reviewExecutor{
|
||
ex: exrun.New(cfg),
|
||
toolNames: names,
|
||
modelSpec: modelSpec,
|
||
callerID: prCallerID(),
|
||
budget: budget,
|
||
}, nil
|
||
}
|
||
|
||
// run executes one agent pass (review or recheck) through the run kernel and
|
||
// returns the model's final text. An empty answer with no error is reported as
|
||
// an error so the caller (reviewWithSpecialist) renders the advisory "reviewer
|
||
// failed to complete" notice rather than a blank section.
|
||
func (r *reviewExecutor) run(ctx context.Context, system, task string, maxSteps int) (string, error) {
|
||
res := r.ex.Run(ctx, exrun.RunnableAgent{
|
||
Name: "gadfly-review",
|
||
SystemPrompt: system,
|
||
ModelTier: r.modelSpec,
|
||
MaxIterations: maxSteps,
|
||
MaxRuntime: reviewTimeout(),
|
||
LowLevelTools: r.toolNames,
|
||
Critic: exrun.CriticConfig{Enabled: true},
|
||
}, exectool.Invocation{
|
||
RunID: fmt.Sprintf("gadfly-%d", runSeq.Add(1)),
|
||
CallerID: r.callerID,
|
||
}, task)
|
||
|
||
// Feed token spend toward the per-PR budget out-of-band (Commit carries only
|
||
// seconds; the executor already called it). Safe on a nil budget.
|
||
r.budget.addUsage(res.Usage)
|
||
|
||
if res.Err != nil {
|
||
return "", res.Err
|
||
}
|
||
if out := strings.TrimSpace(res.Output); out != "" {
|
||
return out, nil
|
||
}
|
||
return "", errors.New("agent produced no output")
|
||
}
|
||
|
||
// prCallerID is the budget/audit caller key: the repo + PR, so a budget keyed on
|
||
// it is naturally per-PR. Falls back to "local" for an out-of-CI run.
|
||
func prCallerID() string {
|
||
repo := strings.TrimSpace(os.Getenv("GADFLY_REPO"))
|
||
pr := strings.TrimSpace(os.Getenv("GADFLY_PR"))
|
||
if repo == "" && pr == "" {
|
||
return "local"
|
||
}
|
||
return repo + "#" + pr
|
||
}
|
||
|
||
// compactionEnabled reports whether context compaction should be wired. On
|
||
// unless GADFLY_COMPACT is explicitly falsey.
|
||
func compactionEnabled() bool { return envBool("GADFLY_COMPACT", true) }
|
||
|
||
// compactRatio is the compaction threshold as a fraction of the model context
|
||
// window (GADFLY_COMPACT_RATIO), clamped to (0,1]; default defaultCompactRatio.
|
||
func compactRatio() float64 {
|
||
v := strings.TrimSpace(os.Getenv("GADFLY_COMPACT_RATIO"))
|
||
if v == "" {
|
||
return defaultCompactRatio
|
||
}
|
||
f, err := strconv.ParseFloat(v, 64)
|
||
if err != nil || f <= 0 || f > 1 {
|
||
return defaultCompactRatio
|
||
}
|
||
return f
|
||
}
|
||
|
||
// resolveContextTokens returns the review model's context window in tokens, used
|
||
// to set the compaction threshold. GADFLY_MODEL_CONTEXT_TOKENS overrides it
|
||
// (needed for custom/self-hosted endpoints executus can't introspect); otherwise
|
||
// it asks executus/model, which knows the static catalog and can fetch an
|
||
// Ollama Cloud model's limit via /api/show (one call, at executor-build time).
|
||
// Returns 0 — disabling compaction — for an unknown model, mirroring executus's
|
||
// "unknown ⇒ don't budget" contract.
|
||
func resolveContextTokens(modelSpec string) int {
|
||
if v := envInt("GADFLY_MODEL_CONTEXT_TOKENS", 0); v > 0 {
|
||
return v
|
||
}
|
||
key := strings.TrimSpace(os.Getenv("GADFLY_API_KEY"))
|
||
if key == "" {
|
||
key = strings.TrimSpace(os.Getenv("OLLAMA_API_KEY"))
|
||
}
|
||
cache := model.NewCloudOllamaLimitCache("", key, nil)
|
||
ctx, cancel := context.WithTimeout(context.Background(), contextTokenLookupTimeout)
|
||
defer cancel()
|
||
if n, ok := model.MaxContextTokensResolving(ctx, modelSpec, cache); ok {
|
||
return n
|
||
}
|
||
// Unknown model or a failed lookup (e.g. no key / unreachable endpoint): don't
|
||
// guess — compaction is disabled. Log it so a misconfiguration is debuggable
|
||
// rather than silently dropping the protection. Set GADFLY_MODEL_CONTEXT_TOKENS
|
||
// to force a window for an endpoint executus can't introspect.
|
||
fmt.Fprintf(os.Stderr, "gadfly: no context window resolved for %q; compaction disabled (set GADFLY_MODEL_CONTEXT_TOKENS to enable it)\n", modelSpec)
|
||
return 0
|
||
}
|