Skip to content
Merged
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
6 changes: 6 additions & 0 deletions .github/scripts/fork-check.sh
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ fork_packages=(
./internal/api/...
./internal/config/...
./internal/registry/...
./internal/runtime/executor/...
./internal/watcher/...
./sdk/api/handlers/...
./sdk/cliproxy/...
Expand Down Expand Up @@ -53,6 +54,11 @@ require_hook sdk/cliproxy/auth/conductor_selection.go 'availableAuthsForPoolSele
require_hook sdk/cliproxy/auth/conductor_selection.go 'HasAccountPools(ctx)' 2
require_hook sdk/api/handlers/handlers_execution.go 'h.accountPoolExecutionContext(ctx)' 2
require_hook sdk/api/handlers/handlers_stream.go 'h.accountPoolExecutionContext(ctx)'
require_hook internal/runtime/executor/meta_executor_stream.go 'helps.ObserveMetaSubscriptionUsage(ctx, eventData)'
require_hook internal/runtime/executor/meta_executor_execute.go 'helps.ObserveMetaSubscriptionUsageSSE(ctx, data)'
require_hook sdk/cliproxy/auth/quota_signals.go 'metaQuotaSignalHeaders[name]'
require_hook sdk/cliproxy/auth/conductor_lifecycle.go 'auth.Quota = mergeQuotaObservation(auth.Quota, existing.Quota)'
require_hook sdk/api/handlers/openai/openai_responses_handlers.go 'responsesSSETrailingFrame(frame)'

echo "== fork tests"
go test -count=1 -run '^TestFork' "${fork_packages[@]}"
Expand Down
29 changes: 29 additions & 0 deletions FORK.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ Everything upstream documents still applies. This page covers only what the fork
| `client.native-model-lists` | `false` | Claude Code is only shown Claude models; Codex is only shown Codex models. |
| `client.key-scopes` | none | Limit a client API key to the credentials of the providers it lists. |
| (no setting) | on | Serve the Muse Code client endpoints (`/muse-code/models` and friends) so the Muse CLI can use the proxy. |
| (no setting) | on | Record Meta's subscription usage on the credential and pass its event through to Responses clients. |

`docker/fork/config.default.yaml`, the configuration the image writes on first start, has every
one of them turned on. `config.example.yaml` is left exactly as upstream ships it, so it does not
Expand Down Expand Up @@ -185,6 +186,30 @@ to an endpoint named there, not to one passed with `--base-url`) and give it a c
export META_API_KEY="sk-muse-77b0..."
```

### Meta subscription usage

Meta has no usage endpoint. It reports the account's subscription windows with every Responses
stream, in a `response.subscription_usage` event after `response.completed`:

```json
{"type":"response.subscription_usage","subscription":{"tier":"…",
"window":{"used_percent":12,"window_duration_mins":300,"resets_at":1791465557},
"weekly":{"used_percent":5,"resets_at":1791763200}}}
```

The fork records that event on the Meta credential as quota signals, the same way Claude and Codex
rate-limit headers are recorded: `X-Meta-Window-Used-Percent`, `X-Meta-Window-Minutes`,
`X-Meta-Window-Reset-At`, `X-Meta-Weekly-Used-Percent`, `X-Meta-Weekly-Reset-At` and `X-Meta-Tier`
(percentages verbatim, reset times in unix seconds). They appear under `quota` in the management
`auth-files` list, which is what T3 Code's usage view reads, and the soonest-reset order uses
them. A credential reports nothing until something has used it through the proxy.

Responses clients also receive the event, although it arrives after the terminal event; the Muse
CLI reads it to show its quota. Every other frame after a terminal event is still dropped.

Reloading a credential file from disk keeps its last observation, since the file does not store
it. This applies to every provider.

## Container image

`ghcr.io/ntindle/cliproxyapi:latest` is built by `.github/workflows/fork-image.yml` from
Expand Down Expand Up @@ -255,5 +280,9 @@ The fork changes these upstream files, each by a few lines, so these are where c
| `sdk/api/handlers/handlers_execution.go`, `sdk/api/handlers/handlers_stream.go` | Apply key scopes to the providers resolved for a request. |
| `sdk/api/handlers/handlers_interceptors.go` | Applies key scopes to every model list. |
| `internal/api/server_routes.go` | Registers the Muse Code endpoints. |
| `internal/runtime/executor/meta_executor_stream.go`, `internal/runtime/executor/meta_executor_execute.go` | Record Meta's subscription usage event. |
| `sdk/cliproxy/auth/quota_signals.go` | Accepts the `X-Meta-*` quota signals. |
| `sdk/cliproxy/auth/conductor_lifecycle.go` | Keeps the quota observation when a credential file is reloaded. |
| `sdk/api/handlers/openai/openai_responses_handlers.go` | Passes the subscription usage event after the terminal event. |

Everything else the fork adds lives in files upstream does not have.
113 changes: 113 additions & 0 deletions internal/runtime/executor/helps/meta_quota.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
package helps

import (
"bytes"
"context"
"net/http"
"strconv"
"strings"

"github.com/router-for-me/CLIProxyAPI/v8/internal/logging"
"github.com/tidwall/gjson"
)

// MetaSubscriptionUsageEventType is the Responses stream event Meta sends after
// response.completed with the account's subscription usage.
const MetaSubscriptionUsageEventType = "response.subscription_usage"

// Quota signal names for Meta subscription usage. They follow the X-Codex-*
// pattern so the passive quota snapshot and the management API treat both alike.
const (
MetaQuotaTierHeader = "X-Meta-Tier"
MetaQuotaWindowUsedHeader = "X-Meta-Window-Used-Percent"
MetaQuotaWindowMinutesHeader = "X-Meta-Window-Minutes"
MetaQuotaWindowResetHeader = "X-Meta-Window-Reset-At"
MetaQuotaWeeklyUsedHeader = "X-Meta-Weekly-Used-Percent"
MetaQuotaWeeklyResetHeader = "X-Meta-Weekly-Reset-At"
)

// ParseMetaSubscriptionUsageHeaders converts one Meta response.subscription_usage
// event into the bounded header representation used for passive quota
// observations. Meta sends it once per response, after response.completed:
//
// {"type":"response.subscription_usage","subscription":{"tier":"…",
// "window":{"used_percent":0,"window_duration_mins":300,"resets_at":1791465557},
// "weekly":{"used_percent":5,"resets_at":1791763200}}}
//
// window is the five-hour-class block and weekly the rolling week. Percentages
// pass through verbatim (above 100 means over quota) and reset times stay unix
// seconds. A window without a percentage or reset time is dropped, and an event
// with no usable window returns nil.
func ParseMetaSubscriptionUsageHeaders(payload []byte) http.Header {
if len(payload) == 0 || gjson.GetBytes(payload, "type").String() != MetaSubscriptionUsageEventType {
return nil
}
return metaSubscriptionUsageHeaders(gjson.GetBytes(payload, "subscription"))
}

// ParseMetaSubscriptionUsageSSE finds the last response.subscription_usage
// event in a buffered Responses event stream.
func ParseMetaSubscriptionUsageSSE(body []byte) http.Header {
var headers http.Header
for _, line := range bytes.Split(body, []byte("\n")) {
line = bytes.TrimSpace(line)
if !bytes.HasPrefix(line, []byte("data:")) {
continue
}
if parsed := ParseMetaSubscriptionUsageHeaders(bytes.TrimSpace(line[len("data:"):])); parsed != nil {
headers = parsed
}
}
return headers
}

// ObserveMetaSubscriptionUsage records a subscription usage event on the request
// context. The conductor stores it on the credential when the request finishes.
func ObserveMetaSubscriptionUsage(ctx context.Context, payload []byte) {
logging.MergeResponseHeaders(ctx, ParseMetaSubscriptionUsageHeaders(payload))
}

// ObserveMetaSubscriptionUsageSSE records the subscription usage carried by a
// buffered Responses event stream.
func ObserveMetaSubscriptionUsageSSE(ctx context.Context, body []byte) {
logging.MergeResponseHeaders(ctx, ParseMetaSubscriptionUsageSSE(body))
}

func metaSubscriptionUsageHeaders(subscription gjson.Result) http.Header {
if !subscription.IsObject() {
return nil
}
headers := make(http.Header)
window := subscription.Get("window")
minutes := window.Get("window_duration_mins")
if used, reset, ok := metaQuotaWindow(window); ok && minutes.Type == gjson.Number && minutes.Int() > 0 {
headers.Set(MetaQuotaWindowUsedHeader, used)
headers.Set(MetaQuotaWindowMinutesHeader, strconv.FormatInt(minutes.Int(), 10))
headers.Set(MetaQuotaWindowResetHeader, reset)
}
if used, reset, ok := metaQuotaWindow(subscription.Get("weekly")); ok {
headers.Set(MetaQuotaWeeklyUsedHeader, used)
headers.Set(MetaQuotaWeeklyResetHeader, reset)
}
if len(headers) == 0 {
return nil
}
// The tier is an opaque plan id. It only labels the snapshot, so a value that
// could forge a request-log line is dropped rather than failing the event.
if tier := strings.TrimSpace(subscription.Get("tier").String()); tier != "" && len(tier) <= 128 && !strings.ContainsAny(tier, "\r\n") {
headers.Set(MetaQuotaTierHeader, tier)
}
return headers
}

func metaQuotaWindow(window gjson.Result) (used, reset string, ok bool) {
if !window.IsObject() {
return "", "", false
}
usedPercent := window.Get("used_percent")
resetsAt := window.Get("resets_at")
if usedPercent.Type != gjson.Number || usedPercent.Float() < 0 || resetsAt.Type != gjson.Number || resetsAt.Int() <= 0 {
return "", "", false
}
return strings.TrimSpace(usedPercent.Raw), strconv.FormatInt(resetsAt.Int(), 10), true
}
93 changes: 93 additions & 0 deletions internal/runtime/executor/helps/meta_quota_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
package helps

import (
"context"
"net/http"
"testing"

"github.com/router-for-me/CLIProxyAPI/v8/internal/logging"
)

// Recorded from api.meta.ai on 2026-10-08; only the tier id is replaced.
const metaSubscriptionUsageEvent = `{"subscription":{"tier":"tier-1","weekly":{"resets_at":1791763200,"used_percent":5},"window":{"resets_at":1791465557,"used_percent":0,"window_duration_mins":300}},"type":"response.subscription_usage"}`

func TestForkParseMetaSubscriptionUsageHeaders(t *testing.T) {
headers := ParseMetaSubscriptionUsageHeaders([]byte(metaSubscriptionUsageEvent))
want := map[string]string{
MetaQuotaTierHeader: "tier-1",
MetaQuotaWindowUsedHeader: "0",
MetaQuotaWindowMinutesHeader: "300",
MetaQuotaWindowResetHeader: "1791465557",
MetaQuotaWeeklyUsedHeader: "5",
MetaQuotaWeeklyResetHeader: "1791763200",
}
if len(headers) != len(want) {
t.Fatalf("headers = %#v, want %d entries", headers, len(want))
}
for name, value := range want {
if got := headers.Get(name); got != value {
t.Fatalf("%s = %q, want %q", name, got, value)
}
}
}

func TestForkParseMetaSubscriptionUsageKeepsOverQuotaAndDropsInvalidWindows(t *testing.T) {
headers := ParseMetaSubscriptionUsageHeaders([]byte(`{"type":"response.subscription_usage","subscription":{` +
`"window":{"used_percent":104.5,"window_duration_mins":300,"resets_at":1791465557},` +
`"weekly":{"used_percent":-1,"resets_at":1791763200}}}`))
if got := headers.Get(MetaQuotaWindowUsedHeader); got != "104.5" {
t.Fatalf("window used = %q, want the over-quota value verbatim", got)
}
if headers.Get(MetaQuotaWeeklyUsedHeader) != "" || headers.Get(MetaQuotaWeeklyResetHeader) != "" {
t.Fatalf("negative weekly percentage was kept: %#v", headers)
}
if headers.Get(MetaQuotaTierHeader) != "" {
t.Fatalf("missing tier produced a header: %#v", headers)
}

for _, payload := range []string{
``,
`{"type":"response.completed","subscription":{"window":{"used_percent":1,"window_duration_mins":300,"resets_at":1}}}`,
`{"type":"response.subscription_usage"}`,
`{"type":"response.subscription_usage","subscription":{"window":{"used_percent":1,"window_duration_mins":0,"resets_at":1791465557}}}`,
`{"type":"response.subscription_usage","subscription":{"weekly":{"used_percent":"5","resets_at":1791763200}}}`,
`{"type":"response.subscription_usage","subscription":{"tier":"t","weekly":{"used_percent":5}}}`,
} {
if headers := ParseMetaSubscriptionUsageHeaders([]byte(payload)); headers != nil {
t.Fatalf("payload %s produced %#v, want nil", payload, headers)
}
}
}

func TestForkParseMetaSubscriptionUsageDropsUnsafeTier(t *testing.T) {
headers := ParseMetaSubscriptionUsageHeaders([]byte(`{"type":"response.subscription_usage","subscription":{"tier":"a\nb","weekly":{"used_percent":5,"resets_at":1791763200}}}`))
if headers.Get(MetaQuotaWeeklyUsedHeader) != "5" || headers.Get(MetaQuotaTierHeader) != "" {
t.Fatalf("headers = %#v, want the weekly window without the tier", headers)
}
}

func TestForkParseMetaSubscriptionUsageSSEUsesTheLastEvent(t *testing.T) {
body := "event: response.completed\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_1\"}}\n\n" +
"event: response.subscription_usage\n" +
"data: {\"type\":\"response.subscription_usage\",\"subscription\":{\"weekly\":{\"used_percent\":5,\"resets_at\":1791763200}}}\n\n" +
"event: response.subscription_usage\n" +
"data: " + metaSubscriptionUsageEvent + "\n\n"
headers := ParseMetaSubscriptionUsageSSE([]byte(body))
if headers.Get(MetaQuotaWindowMinutesHeader) != "300" || headers.Get(MetaQuotaTierHeader) != "tier-1" {
t.Fatalf("headers = %#v, want the last event", headers)
}
if ParseMetaSubscriptionUsageSSE([]byte("data: {\"type\":\"response.completed\"}\n\n")) != nil {
t.Fatal("a stream without subscription usage produced headers")
}
}

func TestForkObserveMetaSubscriptionUsageMergesIntoResponseHeaders(t *testing.T) {
ctx := logging.WithResponseHeadersHolder(context.Background())
logging.SetResponseHeaders(ctx, http.Header{"X-Request-Id": []string{"req-1"}})
ObserveMetaSubscriptionUsage(ctx, []byte(metaSubscriptionUsageEvent))
headers := logging.GetResponseHeaders(ctx)
if headers.Get("X-Request-Id") != "req-1" || headers.Get(MetaQuotaWeeklyUsedHeader) != "5" {
t.Fatalf("response headers = %#v", headers)
}
}
3 changes: 3 additions & 0 deletions internal/runtime/executor/meta_executor_execute.go
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,9 @@ func (e *MetaExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, req
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
return resp, wrapMetaUpstreamError(httpResp.StatusCode, data)
}
// Meta sends subscription usage after response.completed, which is where the
// translation below stops reading.
helps.ObserveMetaSubscriptionUsageSSE(ctx, data)

out, errCompleted := e.translateMetaCompleted(ctx, req, prepared, data)
if errCompleted != nil {
Expand Down
94 changes: 94 additions & 0 deletions internal/runtime/executor/meta_executor_quota_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
package executor

import (
"context"
"net/http"
"net/http/httptest"
"strings"
"testing"

"github.com/router-for-me/CLIProxyAPI/v8/internal/config"
"github.com/router-for-me/CLIProxyAPI/v8/internal/logging"
"github.com/router-for-me/CLIProxyAPI/v8/internal/runtime/executor/helps"
cliproxyauth "github.com/router-for-me/CLIProxyAPI/v8/sdk/cliproxy/auth"
cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v8/sdk/cliproxy/executor"
sdktranslator "github.com/router-for-me/CLIProxyAPI/v8/sdk/translator"
)

// Meta sends the subscription usage after response.completed (recorded 2026-10-08).
const metaSubscriptionUsageSSE = "event: response.output_item.done\n" +
"data: {\"type\":\"response.output_item.done\",\"output_index\":0,\"item\":{\"id\":\"item_0\",\"type\":\"message\",\"role\":\"assistant\",\"content\":[{\"type\":\"output_text\",\"text\":\"OK\"}]}}\n\n" +
"event: response.completed\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_1\",\"status\":\"completed\",\"model\":\"muse-spark-1.3\",\"usage\":{\"input_tokens\":3,\"output_tokens\":1,\"total_tokens\":4}}}\n\n" +
"event: response.subscription_usage\n" +
"data: {\"subscription\":{\"tier\":\"tier-1\",\"weekly\":{\"resets_at\":1791763200,\"used_percent\":5},\"window\":{\"resets_at\":1791465557,\"used_percent\":12,\"window_duration_mins\":300}},\"type\":\"response.subscription_usage\"}\n\n"

func newMetaSubscriptionUsageServer(t *testing.T) *cliproxyauth.Auth {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
_, _ = w.Write([]byte(metaSubscriptionUsageSSE))
}))
t.Cleanup(server.Close)
return &cliproxyauth.Auth{
Provider: "meta",
Attributes: map[string]string{"api_key": "meta-token", "base_url": server.URL},
}
}

func assertMetaSubscriptionUsageObserved(t *testing.T, ctx context.Context) {
t.Helper()
headers := logging.GetResponseHeaders(ctx)
if headers.Get(helps.MetaQuotaWindowUsedHeader) != "12" ||
headers.Get(helps.MetaQuotaWindowMinutesHeader) != "300" ||
headers.Get(helps.MetaQuotaWeeklyResetHeader) != "1791763200" ||
headers.Get(helps.MetaQuotaTierHeader) != "tier-1" {
t.Fatalf("observed response headers = %#v, want the Meta subscription usage", headers)
}
}

func TestForkMetaExecutorStreamRecordsAndForwardsSubscriptionUsage(t *testing.T) {
auth := newMetaSubscriptionUsageServer(t)
ctx := logging.WithResponseHeadersHolder(context.Background())
result, err := NewMetaExecutor(&config.Config{}).ExecuteStream(ctx, auth, cliproxyexecutor.Request{
Model: "muse-spark-1.3",
Payload: []byte(`{"model":"muse-spark-1.3","input":"hello","stream":true}`),
}, cliproxyexecutor.Options{
Stream: true,
SourceFormat: sdktranslator.FromString("openai-response"),
})
if err != nil {
t.Fatalf("ExecuteStream() error = %v", err)
}
var forwarded strings.Builder
for chunk := range result.Chunks {
if chunk.Err != nil {
t.Fatalf("stream chunk error = %v", chunk.Err)
}
forwarded.Write(chunk.Payload)
forwarded.WriteByte('\n')
}
// The Muse CLI reads the event after response.completed, so the client must still get it.
if !strings.Contains(forwarded.String(), `"type":"response.subscription_usage"`) {
t.Fatalf("forwarded stream lacks the subscription usage event:\n%s", forwarded.String())
}
assertMetaSubscriptionUsageObserved(t, ctx)
}

func TestForkMetaExecutorExecuteRecordsSubscriptionUsage(t *testing.T) {
auth := newMetaSubscriptionUsageServer(t)
ctx := logging.WithResponseHeadersHolder(context.Background())
resp, err := NewMetaExecutor(&config.Config{}).Execute(ctx, auth, cliproxyexecutor.Request{
Model: "muse-spark-1.3",
Payload: []byte(`{"model":"muse-spark-1.3","input":"hello"}`),
}, cliproxyexecutor.Options{
SourceFormat: sdktranslator.FromString("openai-response"),
})
if err != nil {
t.Fatalf("Execute() error = %v", err)
}
if !strings.Contains(string(resp.Payload), "resp_1") {
t.Fatalf("Execute() payload = %s, want the completed response", resp.Payload)
}
assertMetaSubscriptionUsageObserved(t, ctx)
}
2 changes: 2 additions & 0 deletions internal/runtime/executor/meta_executor_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,8 @@ func (e *MetaExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Aut
streamUsage.Observe(detail, true)
}
eventData = patchCodexCompletedOutput(eventData, outputItemsByIndex, outputItemsFallback)
case helps.MetaSubscriptionUsageEventType:
helps.ObserveMetaSubscriptionUsage(ctx, eventData)
}
if !emitTranslatedLine(append([]byte("data: "), eventData...)) {
return
Expand Down
Loading