diff --git a/internal/pluginhost/adapters_interceptors.go b/internal/pluginhost/adapters_interceptors.go index 12f73112a..df68f4c18 100644 --- a/internal/pluginhost/adapters_interceptors.go +++ b/internal/pluginhost/adapters_interceptors.go @@ -109,9 +109,11 @@ func (h *Host) InterceptRequestAfterAuthExcept(ctx context.Context, req pluginap func (h *Host) interceptRequest(ctx context.Context, req pluginapi.RequestInterceptRequest, method string, invoke func(pluginapi.RequestInterceptor, context.Context, pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error), skipPluginID string) pluginapi.RequestInterceptResponse { current := pluginapi.RequestInterceptResponse{ Headers: cloneHeader(req.Headers), - Body: bytes.Clone(req.Body), } skipPluginID = strings.TrimSpace(skipPluginID) + currentBase := req.Body + bodyModified := false + for _, record := range h.activeRecords() { interceptor := record.plugin.Capabilities.RequestInterceptor if h.isPluginFused(record.id) || interceptor == nil || record.id == skipPluginID { @@ -119,14 +121,19 @@ func (h *Host) interceptRequest(ctx context.Context, req pluginapi.RequestInterc } nextReq := req nextReq.Headers = cloneHeader(current.Headers) - nextReq.Body = bytes.Clone(current.Body) + if len(currentBase) > 0 { + nextReq.Body = bytes.Clone(currentBase) + } else { + nextReq.Body = nil + } nextReq.Metadata = cloneInterceptorMetadata(req.Metadata) if resp, ok := h.callRequestInterceptor(ctx, record, method, func(callCtx context.Context, callReq pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { return invoke(interceptor, callCtx, callReq) }, nextReq); ok { current.Headers = mergeHeaders(current.Headers, resp.Headers, resp.ClearHeaders) if len(resp.Body) > 0 { - current.Body = bytes.Clone(resp.Body) + currentBase = bytes.Clone(resp.Body) + bodyModified = true } if resp.Terminate { current.Terminate = true @@ -137,6 +144,9 @@ func (h *Host) interceptRequest(ctx context.Context, req pluginapi.RequestInterc } } } + if bodyModified { + current.Body = currentBase + } return current } diff --git a/internal/pluginhost/adapters_test.go b/internal/pluginhost/adapters_test.go index 3a8878fc0..5fc5e9a0c 100644 --- a/internal/pluginhost/adapters_test.go +++ b/internal/pluginhost/adapters_test.go @@ -3785,3 +3785,232 @@ func (failingReadCloser) Read(p []byte) (int, error) { func (failingReadCloser) Close() error { return nil } + +func TestInterceptRequest_ReadOnlyInterceptorDoesNotReplaceBody_Issue6101(t *testing.T) { + // A read-only interceptor (like cpa-account-config-manager) returns an unmodified body (empty/nil Body). + // Per RequestInterceptResponse documentation, "Body replaces the current request body only when non-empty." + // When an interceptor does not mutate the body, host.InterceptRequestBeforeAuth and InterceptRequestAfterAuth + // must return an empty/nil Body so downstream callers know the request body was not replaced, + // rather than returning a cloned copy of the input body. + host := newHostWithRecords(capabilityRecord{ + id: "read-only-observer", + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + return pluginapi.RequestInterceptResponse{ + Headers: http.Header{"X-Observed": []string{"true"}}, + }, nil + }), + }}, + }) + + inputBody := []byte("large-request-payload-for-codex-session") + gotBefore := host.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{ + Body: inputBody, + }) + if len(gotBefore.Body) != 0 { + t.Fatalf("InterceptRequestBeforeAuth returned non-empty body %q, want empty body for read-only interceptor", gotBefore.Body) + } + if gotBefore.Headers.Get("X-Observed") != "true" { + t.Fatalf("expected header X-Observed to be preserved, got %#v", gotBefore.Headers) + } + + gotAfter := host.InterceptRequestAfterAuth(context.Background(), pluginapi.RequestInterceptRequest{ + Body: inputBody, + }) + if len(gotAfter.Body) != 0 { + t.Fatalf("InterceptRequestAfterAuth returned non-empty body %q, want empty body for read-only interceptor", gotAfter.Body) + } +} + +func TestInterceptRequest_FirstPluginInPlaceMutationDiscardedIfEmptyOrFailed(t *testing.T) { + // If plugin 1 mutates req.Body in-place but returns an empty body (or fails), + // plugin 2 must still receive the uncorrupted original body from currentBase. + var plugin2ReceivedBody []byte + host := newHostWithRecords( + capabilityRecord{ + id: "plugin-1-mutator", + priority: 20, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + // In-place mutation on nextReq.Body + req.Body[0] = 'X' + // Returns empty body (mutation is not returned as a replacement) + return pluginapi.RequestInterceptResponse{ + Headers: http.Header{"X-Plugin-1": []string{"ran"}}, + }, nil + }), + }}, + }, + capabilityRecord{ + id: "plugin-2-observer", + priority: 10, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + plugin2ReceivedBody = bytes.Clone(req.Body) + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }, + ) + + inputBody := []byte("clean-payload") + got := host.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{ + Body: inputBody, + }) + + if string(inputBody) != "clean-payload" { + t.Fatalf("caller inputBody was corrupted: %q", string(inputBody)) + } + if string(plugin2ReceivedBody) != "clean-payload" { + t.Fatalf("plugin 2 received corrupted body: %q, want clean-payload", string(plugin2ReceivedBody)) + } + if len(got.Body) != 0 { + t.Fatalf("got.Body should be empty when no plugin returned non-empty body, got: %q", string(got.Body)) + } +} + +func TestInterceptRequest_FirstPluginInPlaceMutationDiscardedOnError(t *testing.T) { + // If plugin 1 mutates req.Body in-place but returns an error, + // plugin 2 must still receive the uncorrupted original body from currentBase. + var plugin2ReceivedBody []byte + host := newHostWithRecords( + capabilityRecord{ + id: "plugin-1-error", + priority: 20, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + req.Body[0] = 'E' + return pluginapi.RequestInterceptResponse{}, errors.New("interceptor failed") + }), + }}, + }, + capabilityRecord{ + id: "plugin-2-observer", + priority: 10, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + plugin2ReceivedBody = bytes.Clone(req.Body) + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }, + ) + + inputBody := []byte("clean-payload") + got := host.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{ + Body: inputBody, + }) + + if string(inputBody) != "clean-payload" { + t.Fatalf("caller inputBody was corrupted: %q", string(inputBody)) + } + if string(plugin2ReceivedBody) != "clean-payload" { + t.Fatalf("plugin 2 received corrupted body: %q, want clean-payload", string(plugin2ReceivedBody)) + } + if len(got.Body) != 0 { + t.Fatalf("got.Body should be empty when no plugin returned non-empty body, got: %q", string(got.Body)) + } +} + +func TestInterceptRequest_FirstPluginInPlaceMutationDiscardedOnPanic(t *testing.T) { + var plugin2ReceivedBody []byte + host := newHostWithRecords( + capabilityRecord{ + id: "plugin-1-panic", + priority: 20, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + req.Body[0] = 'Z' + panic("boom") + }), + }}, + }, + capabilityRecord{ + id: "plugin-2-observer", + priority: 10, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + plugin2ReceivedBody = bytes.Clone(req.Body) + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }, + ) + + inputBody := []byte("clean-payload") + got := host.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{ + Body: inputBody, + }) + + if string(inputBody) != "clean-payload" { + t.Fatalf("caller inputBody was corrupted: %q", string(inputBody)) + } + if string(plugin2ReceivedBody) != "clean-payload" { + t.Fatalf("plugin 2 received corrupted body: %q, want clean-payload", string(plugin2ReceivedBody)) + } + if len(got.Body) != 0 { + t.Fatalf("got.Body should be empty when no plugin returned non-empty body, got: %q", string(got.Body)) + } +} + +func BenchmarkHostRequestInterceptors_ReadOnly(b *testing.B) { + sizes := []struct { + name string + bytes int + }{ + {name: "1MiB", bytes: 1 << 20}, + {name: "8MiB", bytes: 8 << 20}, + } + // 1 read-only plugin (matching cpa-account-config-manager) + hostSingle := newHostWithRecords(capabilityRecord{ + id: "read-only-observer", + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }) + // 2 read-only plugins + hostMulti := newHostWithRecords( + capabilityRecord{ + id: "observer-1", priority: 20, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }, + capabilityRecord{ + id: "observer-2", priority: 10, + plugin: pluginapi.Plugin{Capabilities: pluginapi.Capabilities{ + RequestInterceptor: requestInterceptorFunc(func(ctx context.Context, req pluginapi.RequestInterceptRequest) (pluginapi.RequestInterceptResponse, error) { + return pluginapi.RequestInterceptResponse{}, nil + }), + }}, + }, + ) + + for _, size := range sizes { + payload := make([]byte, size.bytes) + b.Run("single/"+size.name, func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for range b.N { + got := hostSingle.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{Body: payload}) + if len(got.Body) != 0 { + b.Fatal("expected empty body for read-only interceptor") + } + } + }) + b.Run("multi/"+size.name, func(b *testing.B) { + b.ReportAllocs() + b.ResetTimer() + for range b.N { + got := hostMulti.InterceptRequestBeforeAuth(context.Background(), pluginapi.RequestInterceptRequest{Body: payload}) + if len(got.Body) != 0 { + b.Fatal("expected empty body for read-only interceptor") + } + } + }) + } +} diff --git a/sdk/api/handlers/handlers_interceptors.go b/sdk/api/handlers/handlers_interceptors.go index 61bbe3867..28d257dfb 100644 --- a/sdk/api/handlers/handlers_interceptors.go +++ b/sdk/api/handlers/handlers_interceptors.go @@ -375,7 +375,7 @@ func (c *requestAfterAuthCapture) record(req coreexecutor.RequestAfterAuthInterc return } headers := mergeRequestInterceptorHeaders(req.Headers, resp.Headers, resp.ClearHeaders) - body := cloneBytes(req.Body) + var body []byte var originalRequest []byte originalRequestReplaced := false if len(resp.Body) > 0 { @@ -402,11 +402,11 @@ func (c *requestAfterAuthCapture) apply(req coreexecutor.Request, opts coreexecu if !c.set { return req, opts } - req.Payload = cloneBytes(c.body) - opts.Headers = cloneHeader(c.headers) if c.originalRequestReplaced { + req.Payload = cloneBytes(c.body) opts.OriginalRequest = cloneBytes(c.originalRequest) } + opts.Headers = cloneHeader(c.headers) return req, opts } @@ -479,7 +479,7 @@ func (h *BaseAPIHandler) applyRequestInterceptorsBeforeAuth(ctx context.Context, RequestedModel: requestedModel, Stream: opts.Stream, Headers: cloneHeader(opts.Headers), - Body: cloneBytes(req.Payload), + Body: req.Payload, Metadata: opts.Metadata, }, skipPluginID) opts.Headers = finalInterceptorHeaders(opts.Headers, resp.Headers) @@ -561,7 +561,7 @@ func (h *BaseAPIHandler) applyRequestInterceptorsAfterAuth(ctx context.Context, RequestedModel: req.RequestedModel, Stream: req.Stream, Headers: cloneHeader(req.Headers), - Body: cloneBytes(req.Body), + Body: req.Body, Metadata: req.Metadata, }, skipPluginID) return coreexecutor.RequestAfterAuthInterceptResponse{ diff --git a/sdk/api/handlers/handlers_interceptors_test.go b/sdk/api/handlers/handlers_interceptors_test.go index 626c2f024..310a69c13 100644 --- a/sdk/api/handlers/handlers_interceptors_test.go +++ b/sdk/api/handlers/handlers_interceptors_test.go @@ -10,6 +10,7 @@ import ( "sync" "testing" "time" + "unsafe" "github.com/gin-gonic/gin" "github.com/router-for-me/CLIProxyAPI/v7/internal/interfaces" @@ -56,7 +57,6 @@ func (h *handlerInterceptorTestHost) InterceptRequestBeforeAuth(ctx context.Cont } return pluginapi.RequestInterceptResponse{ Headers: cloneHeader(req.Headers), - Body: cloneBytes(req.Body), } } @@ -66,7 +66,6 @@ func (h *handlerInterceptorTestHost) InterceptRequestAfterAuth(ctx context.Conte } return pluginapi.RequestInterceptResponse{ Headers: cloneHeader(req.Headers), - Body: cloneBytes(req.Body), } } @@ -1695,3 +1694,63 @@ func TestWriteModelListResponse_ExposesResponseToPluginInterceptors(t *testing.T t.Fatalf("API_RESPONSE not recorded properly: %#v", apiResp) } } + +func TestRequestAfterAuthCapture_ReadOnlyInterceptorDoesNotReallocatePayload_Issue6101(t *testing.T) { + capture := &requestAfterAuthCapture{} + payload := []byte(`{"messages":[{"role":"user","content":"hello"}]}`) + req := coreexecutor.Request{ + Model: "gpt-5.4", + Payload: payload, + } + opts := coreexecutor.Options{ + OriginalRequest: payload, + } + + // Read-only after-auth response: Body is empty + capture.record(coreexecutor.RequestAfterAuthInterceptRequest{ + Body: payload, + }, coreexecutor.RequestAfterAuthInterceptResponse{ + Body: nil, + }) + + appliedReq, appliedOpts := capture.apply(req, opts) + if unsafe.SliceData(appliedReq.Payload) != unsafe.SliceData(req.Payload) { + t.Fatal("requestAfterAuthCapture.apply reallocated req.Payload when interceptor did not mutate body") + } + if unsafe.SliceData(appliedOpts.OriginalRequest) != unsafe.SliceData(opts.OriginalRequest) { + t.Fatal("requestAfterAuthCapture.apply reallocated opts.OriginalRequest when interceptor did not mutate body") + } +} + +func TestApplyRequestInterceptors_ReadOnlyInterceptorDoesNotReallocatePayload_Issue6101(t *testing.T) { + // A read-only interceptor returns an empty/nil Body. + host := &handlerInterceptorTestHost{ + interceptRequestBeforeAuth: func(ctx context.Context, req pluginapi.RequestInterceptRequest) pluginapi.RequestInterceptResponse { + return pluginapi.RequestInterceptResponse{ + Headers: req.Headers, + } + }, + } + handler := NewBaseAPIHandlers(&sdkconfig.SDKConfig{}, nil) + handler.SetPluginHost(host) + + payload := []byte(`{"messages":[{"role":"user","content":"hello"}]}`) + req := coreexecutor.Request{ + Model: "gpt-5.4", + Payload: payload, + } + opts := coreexecutor.Options{ + OriginalRequest: payload, + } + + gotReq, gotOpts, err := handler.applyRequestInterceptorsBeforeAuth(context.Background(), "openai", "gpt-5.4", "req-1", req, opts, "") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if unsafe.SliceData(gotReq.Payload) != unsafe.SliceData(req.Payload) { + t.Fatal("applyRequestInterceptorsBeforeAuth reallocated req.Payload when interceptor did not mutate body") + } + if unsafe.SliceData(gotOpts.OriginalRequest) != unsafe.SliceData(opts.OriginalRequest) { + t.Fatal("applyRequestInterceptorsBeforeAuth reallocated opts.OriginalRequest when interceptor did not mutate body") + } +} diff --git a/sdk/cliproxy/auth/conductor_execution.go b/sdk/cliproxy/auth/conductor_execution.go index b365f5d85..7a5aafbcf 100644 --- a/sdk/cliproxy/auth/conductor_execution.go +++ b/sdk/cliproxy/auth/conductor_execution.go @@ -339,7 +339,7 @@ func applyRequestAfterAuthInterceptor(ctx context.Context, executor ProviderExec RequestedModel: requestedModel, Stream: opts.Stream, Headers: cloneRequestHeaders(opts.Headers), - Body: bytes.Clone(req.Payload), + Body: req.Payload, Metadata: opts.Metadata, }) opts.Headers = mergeRequestHeaders(opts.Headers, resp.Headers, resp.ClearHeaders) diff --git a/sdk/cliproxy/executor/types.go b/sdk/cliproxy/executor/types.go index 2d1f76580..79d874d43 100644 --- a/sdk/cliproxy/executor/types.go +++ b/sdk/cliproxy/executor/types.go @@ -115,7 +115,7 @@ type RequestAfterAuthInterceptRequest struct { Stream bool // Headers contains the current upstream request headers. Headers http.Header - // Body contains the current request payload. + // Body contains the current request payload. Treat it as read-only; modifications must be returned in RequestAfterAuthInterceptResponse.Body. Body []byte // Metadata is a best-effort cloned context snapshot. Treat it as read-only and JSON-like. Metadata map[string]any diff --git a/sdk/cliproxy/session/identity.go b/sdk/cliproxy/session/identity.go index ed35c2613..12c3fc068 100644 --- a/sdk/cliproxy/session/identity.go +++ b/sdk/cliproxy/session/identity.go @@ -2,7 +2,6 @@ package session import ( - "bytes" "crypto/sha256" "encoding/hex" "encoding/json" @@ -227,10 +226,13 @@ func DerivedID(metadata map[string]any) string { } // Enrich derives a session identity once and places it in both request and option metadata. +// When opts.OriginalRequest is unset, it shares req.Payload as the read-only original request +// baseline to avoid multi-megabyte allocations on large payloads. Callers that mutate req.Payload +// in-place after Enrich must explicitly provide an independent opts.OriginalRequest. func Enrich(req cliproxyexecutor.Request, opts cliproxyexecutor.Options) (cliproxyexecutor.Request, cliproxyexecutor.Options) { payload := opts.OriginalRequest if len(payload) == 0 && len(req.Payload) > 0 { - opts.OriginalRequest = bytes.Clone(req.Payload) + opts.OriginalRequest = req.Payload payload = opts.OriginalRequest } executionID := firstNormalizedMetadataID(cliproxyexecutor.ExecutionSessionMetadataKey, opts.Metadata, req.Metadata) diff --git a/sdk/cliproxy/session/identity_test.go b/sdk/cliproxy/session/identity_test.go index 4063109b9..e632a59e5 100644 --- a/sdk/cliproxy/session/identity_test.go +++ b/sdk/cliproxy/session/identity_test.go @@ -4,6 +4,7 @@ import ( "net/http" "strings" "testing" + "unsafe" cliproxyexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" sdktranslator "github.com/router-for-me/CLIProxyAPI/v7/sdk/translator" @@ -621,3 +622,29 @@ func TestNormalizeToCanonicalUUID(t *testing.T) { t.Fatalf("NormalizeToCanonicalUUID(%q) = %q, want 01a07e72-c84d-7fd3-8207-d217b41cc649", derivedUUID, got) } } + +func TestEnrich_DoesNotClonePayloadWhenPopulatingOriginalRequest_Issue6101(t *testing.T) { + req := cliproxyexecutor.Request{ + Model: "gpt-5.4", + Payload: []byte(`{"messages":[{"role":"user","content":"test"}]}`), + } + opts := cliproxyexecutor.Options{} + + _, enrichedOpts := Enrich(req, opts) + if len(enrichedOpts.OriginalRequest) == 0 { + t.Fatal("expected OriginalRequest to be populated") + } + if unsafe.SliceData(enrichedOpts.OriginalRequest) != unsafe.SliceData(req.Payload) { + t.Fatal("Enrich allocated a new copy of Payload instead of reusing the slice reference for OriginalRequest") + } + + // When caller explicitly provides an independent OriginalRequest, Enrich preserves it + independentOriginal := []byte(`{"messages":[{"role":"user","content":"original"}]}`) + optsWithOriginal := cliproxyexecutor.Options{ + OriginalRequest: independentOriginal, + } + _, enrichedOpts2 := Enrich(req, optsWithOriginal) + if unsafe.SliceData(enrichedOpts2.OriginalRequest) != unsafe.SliceData(independentOriginal) { + t.Fatal("Enrich replaced explicitly provided OriginalRequest") + } +} diff --git a/sdk/pluginapi/types.go b/sdk/pluginapi/types.go index 0cab10fcb..ca9958783 100644 --- a/sdk/pluginapi/types.go +++ b/sdk/pluginapi/types.go @@ -1068,7 +1068,7 @@ type RequestInterceptRequest struct { Stream bool // Headers contains the current upstream request headers. Headers http.Header - // Body contains the current request payload. + // Body contains the current request payload. Treat it as read-only; modifications must be returned in RequestInterceptResponse.Body. Body []byte // Metadata is a best-effort cloned context snapshot. Treat it as read-only and JSON-like. Metadata map[string]any