Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions docs/MODEL_RPM.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# Local per-model completion limits

Add `modelRPM` to the existing Zero config file:

```json
{"modelRPM": {"gpt-4.1": 30, "claude-sonnet-4.5": 15}}
```

Keys resolve through the model registry; custom IDs match exactly after trimming
whitespace. Missing keys and zero mean unlimited. Negative limits and conflicting
aliases are configuration errors. Project config may add or tighten a user cap,
but cannot relax it. Changes take effect on the next Zero process start.

For `zero exec` and interactive agent turns, admission happens immediately before
the wrapped turn session's `Stream`. N admissions fit in any sliding 60-second
window; N+1 returns a local rate-limit error before entering that method. The
error reports when a slot expires. Exec uses its existing provider-error exit
code plus a local-limit hint. There is no sleep or hidden retry. Failed admitted
requests still consume a slot.

Sessions and model switches share a limiter in one process. Separate processes
have independent windows; restarting clears them. This is not an account-wide
cap. It counts Stream admissions, not all HTTP requests: setup/prewarm,
compaction, provider-internal retries, discovery and direct calls outside that
boundary are not counted. No TPM, cost budgets, persistence, cross-process
coordination or new runtime/Python dependency is introduced.
1 change: 1 addition & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ Zero.
## User Docs

- [Install](INSTALL.md)
- [Local per-model completion limits](MODEL_RPM.md)
- [Update flow](UPDATE.md)
- [OAuth logins and subscription-backed providers](oauth-subscriptions.md)

Expand Down
4 changes: 3 additions & 1 deletion internal/cli/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -959,7 +959,9 @@ func runInteractiveTUIWithSetup(stderr io.Writer, deps appDeps, permissionMode a
// notice when project hooks/plugins were dropped for an untrusted workspace.
hookDispatcher, hookSkip := newHookDispatcherWithExtra(workspaceRoot, pluginActivation.hooks, trustRoot, executionRunner)
emitTrustNotice(stderr, hookSkip, pluginActivation.trustSkip, mcpSkip)
modelRPM := zeroruntime.NewModelRPMLimiter(resolved.ModelRPM)
return deps.runTUI(context.Background(), tui.Options{
ModelRPM: modelRPM,
Cwd: workspaceRoot,
Version: version,
Theme: theme,
Expand All @@ -982,7 +984,7 @@ func runInteractiveTUIWithSetup(stderr io.Writer, deps appDeps, permissionMode a
if provider == nil {
return nil
}
if optimized, ok := providers.OptimizedTurnSessions(profile, provider, providers.Options{}); ok {
if optimized, ok := providers.ConfiguredTurnSessions(profile, provider, providers.Options{ModelRPM: modelRPM}); ok {
return optimized
}
return providers.DefaultTurnSessions(profile, provider, providers.Options{})
Expand Down
10 changes: 8 additions & 2 deletions internal/cli/app_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -343,8 +343,10 @@ func TestRunNoArgsFallsBackToUsableProviderWhenNoneMarkedActive(t *testing.T) {
// config.json has providers configured but none marked active (e.g. a
// blank/stale activeProvider field) — Resolve returns the successfully
// normalized providers list alongside ErrNoActiveProvider.
return config.ResolvedConfig{Providers: []config.ProviderProfile{usable}},
fmt.Errorf("%w: active provider %q not found", config.ErrNoActiveProvider, "")
return config.ResolvedConfig{
Providers: []config.ProviderProfile{usable},
ModelRPM: map[string]int{"gpt-test": 1},
}, fmt.Errorf("%w: active provider %q not found", config.ErrNoActiveProvider, "")
},
newProvider: func(profile config.ProviderProfile) (zeroruntime.Provider, error) {
providerProfile = profile
Expand Down Expand Up @@ -374,6 +376,10 @@ func TestRunNoArgsFallsBackToUsableProviderWhenNoneMarkedActive(t *testing.T) {
}
if providerProfile.Name != "work" {
t.Fatalf("provider used = %q, want fallback to the usable saved provider %q", providerProfile.Name, "work")

}
if launchedOptions.ModelRPM == nil {
t.Fatal("ModelRPM = nil, want limiter preserved through provider fallback")
}
}

Expand Down
12 changes: 6 additions & 6 deletions internal/cli/exec.go
Original file line number Diff line number Diff line change
Expand Up @@ -452,20 +452,20 @@ func runExec(args []string, stdout io.Writer, stderr io.Writer, deps appDeps) in
// escalation so post-switch turns are attributed to the escalated model.
currentModel := resolved.Provider.Model
// Optimized OpenAI turn sessions (ZERO_OPENAI_TURN_SESSION, default off). nil
// when gated off or the profile is ineligible: agent.Run then wraps the
// provider in its default adapter. This is the run's STARTING session provider
// and is used whether or not escalation is enabled.
turnSessions, _ := providers.OptimizedTurnSessions(resolved.Provider, provider, providers.Options{})
// when gated off or ineligible, unless a configured RPM limiter needs to
// wrap the default adapter. The limiter is shared across model switches.
sessionOptions := providers.Options{ModelRPM: zeroruntime.NewModelRPMLimiter(resolved.ModelRPM)}
turnSessions, _ := providers.ConfiguredTurnSessions(resolved.Provider, provider, sessionOptions)
// Both switchers come from one shared builder so exec and the interactive TUI
// cannot drift on the nil contracts the agent loop depends on. The session
// switcher is nil unless this run STARTED optimized, which is what keeps a
// default-adapter run on the default adapter.
// switcher preserves the starting transport and shares the RPM limiter.
var modelSwitcher func(context.Context, string) (agent.Provider, error)
var modelSessionSwitcher func(context.Context, string) (zeroruntime.TurnSessionProvider, error)
if options.allowEscalation {
modelSwitcher, modelSessionSwitcher = providers.EscalationSwitchers(
resolved.Provider, provider, deps.newProvider,
func(modelID string) { currentModel = modelID },
sessionOptions,
)
}

Expand Down
83 changes: 83 additions & 0 deletions internal/cli/model_rpm_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
package cli

import (
"bytes"
"context"
"errors"
"github.com/Gitlawb/zero/internal/config"
"github.com/Gitlawb/zero/internal/mcp"
"github.com/Gitlawb/zero/internal/tools"
"github.com/Gitlawb/zero/internal/tui"
"github.com/Gitlawb/zero/internal/zeroruntime"
"strings"
"testing"
)

type rpmExecProvider struct{ calls int }

func (p *rpmExecProvider) StreamCompletion(context.Context, zeroruntime.CompletionRequest) (<-chan zeroruntime.StreamEvent, error) {
p.calls++
ch := make(chan zeroruntime.StreamEvent, 4)
if p.calls == 1 {
ch <- zeroruntime.StreamEvent{Type: zeroruntime.StreamEventToolCallStart, ToolCallID: "fixture", ToolName: "read_file"}
ch <- zeroruntime.StreamEvent{Type: zeroruntime.StreamEventToolCallDelta, ToolCallID: "fixture", ArgumentsFragment: `{"path":"missing-fixture.txt"}`}
ch <- zeroruntime.StreamEvent{Type: zeroruntime.StreamEventToolCallEnd, ToolCallID: "fixture"}
} else {
ch <- zeroruntime.StreamEvent{Type: zeroruntime.StreamEventText, Content: "unexpected second completion"}
}
ch <- zeroruntime.StreamEvent{Type: zeroruntime.StreamEventDone}
close(ch)
return ch, nil
}
func rpmTestDeps(t *testing.T, p *rpmExecProvider) appDeps {
t.Helper()
root := t.TempDir()
for _, k := range []string{"HOME", "USERPROFILE", "APPDATA", "LOCALAPPDATA", "XDG_CONFIG_HOME", "XDG_CACHE_HOME", "XDG_DATA_HOME", "XDG_STATE_HOME"} {
t.Setenv(k, root)
}
return appDeps{getwd: func() (string, error) { return root, nil }, resolveConfig: func(string, config.Overrides) (config.ResolvedConfig, error) {
c := execResolvedConfig()
c.ModelRPM = map[string]int{c.Provider.Model: 1}
return c, nil
}, newProvider: func(config.ProviderProfile) (zeroruntime.Provider, error) { return p, nil }, registerMCPTools: func(context.Context, *tools.Registry, config.MCPConfig, mcp.RegisterOptions) (mcpToolRuntime, error) {
return noopMCPRuntime{}, nil
}}
}
func TestRunExecRPMRefusesSecondCompletion(t *testing.T) {
p := &rpmExecProvider{}
deps := rpmTestDeps(t, p)
var out, err bytes.Buffer
code := runWithDeps([]string{"exec", "fixture"}, &out, &err, deps)
if code != exitProvider || p.calls != 1 || !strings.Contains(err.String(), "Local modelRPM cap reached") {
t.Fatalf("exit=%d calls=%d stdout=%s stderr=%s", code, p.calls, out.String(), err.String())
}
}
func TestInteractiveRPMFactorySharesWindow(t *testing.T) {
p := &rpmExecProvider{}
deps := rpmTestDeps(t, p)
launched := false
deps.runTUI = func(ctx context.Context, o tui.Options) int {
launched = true
if o.ModelRPM == nil {
t.Fatal("TUI escalation lacks limiter")
}
for i := 0; i < 2; i++ {
ss := o.NewTurnSessionProvider(o.ProviderProfile, p)
s, e := ss.OpenTurnSession(ctx)
if e != nil {
t.Fatal(e)
}
_, e = s.Stream(ctx, zeroruntime.CompletionRequest{})
var hit *zeroruntime.RPMLimitError
if (i == 0 && e != nil) || (i == 1 && !errors.As(e, &hit)) {
t.Fatalf("run=%d err=%v", i, e)
}
}
return exitSuccess
}
var out, err bytes.Buffer
code := runWithDeps(nil, &out, &err, deps)
if !launched || code != exitSuccess || p.calls != 1 {
t.Fatalf("launched=%v exit=%d calls=%d stderr=%s", launched, code, p.calls, err.String())
}
}
48 changes: 48 additions & 0 deletions internal/config/model_rpm.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
package config

import (
"fmt"
"github.com/Gitlawb/zero/internal/modelregistry"
"strings"
)

func normalizeModelRPM(limits map[string]int) (map[string]int, error) {
if len(limits) == 0 {
return nil, nil
}
registry, err := modelregistry.DefaultRegistry()
if err != nil {
return nil, err
}
out := make(map[string]int, len(limits))
for key, n := range limits {
model := strings.TrimSpace(key)
if model == "" || n < 0 {
return nil, fmt.Errorf("invalid modelRPM entry %q: model must be non-empty and limit >= 0", key)
}
if id, ok := registry.ResolveID(model); ok {
model = id
}
if old, exists := out[model]; exists && old != n {
return nil, fmt.Errorf("conflicting modelRPM entries resolve to %q", model)
}
out[model] = n
}
return out, nil
}

// Project config can tighten a user cap but cannot relax it.
func mergeModelRPM(dst *map[string]int, src map[string]int, tighten bool) {
if len(src) == 0 {
return
}
if *dst == nil {
*dst = make(map[string]int)
}
for model, n := range src {
old := (*dst)[model]
if !tighten || old == 0 || (n > 0 && n < old) {
(*dst)[model] = n
}
}
}
54 changes: 54 additions & 0 deletions internal/config/model_rpm_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package config

import (
"encoding/json"
"reflect"
"testing"
)

func TestModelRPMRoundTrip(t *testing.T) {
var c FileConfig
if e := json.Unmarshal([]byte(`{"modelRPM":{" OPENAI:GPT-4.1 ":2,"custom":0},"future":true}`), &c); e != nil {
t.Fatal(e)
}
want := map[string]int{"gpt-4.1": 2, "custom": 0}
if !reflect.DeepEqual(c.ModelRPM, want) {
t.Fatal(c.ModelRPM)
}
b, e := json.Marshal(c)
if e != nil {
t.Fatal(e)
}
var d FileConfig
if e = json.Unmarshal(b, &d); e != nil {
t.Fatal(e)
}
if !reflect.DeepEqual(d.ModelRPM, want) || string(d.Extra["future"]) != "true" {
t.Fatalf("round trip=%s", b)
}
}
func TestModelRPMInvalid(t *testing.T) {
for _, s := range []string{`{"modelRPM":{"":1}}`, `{"modelRPM":{"a":-1}}`, `{"modelRPM":{"a":1.5}}`, `{"modelRPM":{"a":"2"}}`, `{"modelRPM":{"gpt-4.1":1,"openai:gpt-4.1":2}}`} {
var c FileConfig
if e := json.Unmarshal([]byte(s), &c); e == nil {
t.Fatalf("accepted %s", s)
}
}
}
func TestModelRPMProjectCanOnlyTighten(t *testing.T) {
user := writeConfig(t, `{"modelRPM":{"gpt-4.1":2},"activeProvider":"test","providers":[{"name":"test","provider":"openai","model":"gpt-4.1"}]}`)
for _, n := range []string{"0", "1", "5"} {
project := writeConfig(t, `{"modelRPM":{"openai:gpt-4.1":`+n+`}}`)
c, e := Resolve(ResolveOptions{UserConfigPath: user, ProjectConfigPath: project, Env: map[string]string{}})
if e != nil {
t.Fatal(e)
}
want := 2
if n == "1" {
want = 1
}
if c.ModelRPM["gpt-4.1"] != want {
t.Fatal(c.ModelRPM)
}
}
}
11 changes: 10 additions & 1 deletion internal/config/resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,9 @@ func Resolve(options ResolveOptions) (ResolvedConfig, error) {
// trusted global-config merge, which must keep honouring the setting.
commandConfig.Sandbox.Enabled = nil
commandConfig.CrossSessionInbound = ""
// External provider commands cannot relax the user cap.
mergeModelRPM(&cfg.ModelRPM, commandConfig.ModelRPM, true)
commandConfig.ModelRPM = nil
mergeConfig(&cfg, commandConfig)
}

Expand Down Expand Up @@ -158,10 +161,14 @@ func Resolve(options ResolveOptions) (ResolvedConfig, error) {
// normalized (but active-less) profile list — keep it so a caller can fall
// back to an already-configured usable provider instead of treating this
// like a config with nothing set up at all.
return ResolvedConfig{Providers: providers}, err
return ResolvedConfig{
ModelRPM: cfg.ModelRPM,
Providers: providers,
}, err
}

return ResolvedConfig{
ModelRPM: cfg.ModelRPM,
ActiveProvider: active.Name,
Providers: providers,
Provider: active,
Expand Down Expand Up @@ -230,6 +237,7 @@ func loadConfigFile(path string) (FileConfig, error) {
}

func mergeConfig(dst *FileConfig, src FileConfig) {
mergeModelRPM(&dst.ModelRPM, src.ModelRPM, false)
if activeProvider := strings.TrimSpace(src.ActiveProvider); activeProvider != "" {
dst.ActiveProvider = activeProvider
}
Expand Down Expand Up @@ -291,6 +299,7 @@ func mergeConfig(dst *FileConfig, src FileConfig) {
}

func mergeProjectConfig(dst *FileConfig, src FileConfig) error {
mergeModelRPM(&dst.ModelRPM, src.ModelRPM, true)
if activeProvider := strings.TrimSpace(src.ActiveProvider); activeProvider != "" {
dst.ActiveProvider = activeProvider
}
Expand Down
14 changes: 10 additions & 4 deletions internal/config/resolver_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1161,10 +1161,13 @@ func TestResolveKeepsNormalizedProvidersWhenNoneMarkedActive(t *testing.T) {
// activeProvider is blank/stale — a caller like the interactive TUI still
// needs the normalized list to fall back to an already-usable provider
// instead of forcing a full re-onboarding wizard.
path := writeConfig(t, `{"providers":[
{"name":"work","provider_kind":"openai","apiKey":"sk-test","model":"gpt-test"},
{"name":"other","provider_kind":"openai","apiKey":"sk-other","model":"gpt-test"}
]}`)
path := writeConfig(t, `{
"modelRPM":{"gpt-test":1},
"providers":[
{"name":"work","provider_kind":"openai","apiKey":"sk-test","model":"gpt-test"},
{"name":"other","provider_kind":"openai","apiKey":"sk-other","model":"gpt-test"}
]
}`)

resolved, err := Resolve(ResolveOptions{ProjectConfigPath: path, Env: map[string]string{}})
if !errors.Is(err, ErrNoActiveProvider) {
Expand All @@ -1173,6 +1176,9 @@ func TestResolveKeepsNormalizedProvidersWhenNoneMarkedActive(t *testing.T) {
if len(resolved.Providers) != 2 {
t.Fatalf("Providers = %#v, want the 2 normalized profiles preserved despite the error", resolved.Providers)
}
if got := resolved.ModelRPM["gpt-test"]; got != 1 {
t.Fatalf("ModelRPM[gpt-test] = %d, want 1 preserved despite ErrNoActiveProvider", got)
}
}

func TestResolveTrimsProviderProfileAliasesBeforeFallback(t *testing.T) {
Expand Down
Loading