From ba5ab795a2de8e8dbd9bfda357a7434a65637beb Mon Sep 17 00:00:00 2001 From: Luis Pater Date: Tue, 11 Aug 2026 04:32:20 +0800 Subject: [PATCH] feat(plugin): add schema-v3 stream chunk contract to omit payload request bodies - Bump plugin schema to version 3 and introduce `SchemaVersionStreamChunkOmitRequestBody`. - Treat missing plugin schema versions as legacy during RPC registration (`0 -> 1`) and expose schema on plugin descriptors. - In stream interception, keep request headers/bodies on header-init chunk and stop re-sending them on payload chunks for schema-v3+ plugins, with per-chunk cloning for legacy plugins. Closes: #4876 --- internal/pluginhost/adapters_interceptors.go | 34 ++++++- internal/pluginhost/adapters_test.go | 79 +++++++++++++++++ internal/pluginhost/rpc_client.go | 8 +- internal/pluginhost/rpc_schema_test.go | 9 +- sdk/api/handlers/handlers_interceptors.go | 19 ++++ .../handlers/handlers_interceptors_test.go | 88 +++++++++++++++++-- sdk/api/handlers/handlers_stream.go | 63 +++++++++---- sdk/pluginabi/types.go | 8 +- sdk/pluginabi/types_test.go | 7 +- sdk/pluginapi/types.go | 15 +++- 10 files changed, 297 insertions(+), 33 deletions(-) diff --git a/internal/pluginhost/adapters_interceptors.go b/internal/pluginhost/adapters_interceptors.go index 91d0227f5..7239d4d4b 100644 --- a/internal/pluginhost/adapters_interceptors.go +++ b/internal/pluginhost/adapters_interceptors.go @@ -10,6 +10,7 @@ import ( "strings" coreauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" + "github.com/router-for-me/CLIProxyAPI/v7/sdk/pluginabi" "github.com/router-for-me/CLIProxyAPI/v7/sdk/pluginapi" log "github.com/sirupsen/logrus" ) @@ -211,8 +212,15 @@ func (h *Host) InterceptStreamChunkExcept(ctx context.Context, req pluginapi.Str nextReq := req nextReq.RequestHeaders = cloneHeader(req.RequestHeaders) nextReq.ResponseHeaders = cloneHeader(current.Headers) - nextReq.OriginalRequest = bytes.Clone(req.OriginalRequest) - nextReq.RequestBody = bytes.Clone(req.RequestBody) + // Schema v3+ omits request bodies on payload chunks to avoid re-sending multi-MB + // prompts across cgo/JSON for every frame. Legacy plugins still receive them. + if req.ChunkIndex != pluginapi.StreamChunkHeaderInitIndex && streamChunkOmitsRequestBodies(record.plugin.SchemaVersion) { + nextReq.OriginalRequest = nil + nextReq.RequestBody = nil + } else { + nextReq.OriginalRequest = bytes.Clone(req.OriginalRequest) + nextReq.RequestBody = bytes.Clone(req.RequestBody) + } nextReq.Body = bytes.Clone(current.Body) nextReq.HistoryChunks = cloneByteSlices(req.HistoryChunks) nextReq.Metadata = cloneInterceptorMetadata(req.Metadata) @@ -244,6 +252,28 @@ func (h *Host) HasStreamInterceptors() bool { return false } +// StreamChunkPayloadIncludesRequestBody reports whether any active stream chunk +// interceptor still requires OriginalRequest/RequestBody on payload chunks +// (schema_version < SchemaVersionStreamChunkOmitRequestBody). +func (h *Host) StreamChunkPayloadIncludesRequestBody() bool { + if h == nil { + return false + } + for _, record := range h.activeRecords() { + if h.isPluginFused(record.id) || record.plugin.Capabilities.StreamChunkInterceptor == nil { + continue + } + if !streamChunkOmitsRequestBodies(record.plugin.SchemaVersion) { + return true + } + } + return false +} + +func streamChunkOmitsRequestBodies(schemaVersion uint32) bool { + return schemaVersion >= pluginabi.SchemaVersionStreamChunkOmitRequestBody +} + func (h *Host) HasRequestInterceptors() bool { if h == nil { return false diff --git a/internal/pluginhost/adapters_test.go b/internal/pluginhost/adapters_test.go index 9540040f2..de62918ec 100644 --- a/internal/pluginhost/adapters_test.go +++ b/internal/pluginhost/adapters_test.go @@ -19,6 +19,7 @@ import ( coreauth "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/auth" coreexecutor "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/executor" coreusage "github.com/router-for-me/CLIProxyAPI/v7/sdk/cliproxy/usage" + "github.com/router-for-me/CLIProxyAPI/v7/sdk/pluginabi" "github.com/router-for-me/CLIProxyAPI/v7/sdk/pluginapi" sdktranslator "github.com/router-for-me/CLIProxyAPI/v7/sdk/translator" ) @@ -1691,6 +1692,84 @@ func TestHasStreamInterceptorsReflectsActiveStreamInterceptors(t *testing.T) { } } +func TestStreamChunkRequestBodyPolicyBySchemaVersion(t *testing.T) { + var legacyGot, modernGot pluginapi.StreamChunkInterceptRequest + host := newHostWithRecords( + capabilityRecord{ + id: "legacy", + plugin: pluginapi.Plugin{ + SchemaVersion: 2, + Capabilities: pluginapi.Capabilities{ + StreamChunkInterceptor: responseInterceptorFunc{ + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) (pluginapi.StreamChunkInterceptResponse, error) { + legacyGot = req + return pluginapi.StreamChunkInterceptResponse{Body: req.Body}, nil + }, + }, + }, + }, + }, + capabilityRecord{ + id: "modern", + plugin: pluginapi.Plugin{ + SchemaVersion: pluginabi.SchemaVersionStreamChunkOmitRequestBody, + Capabilities: pluginapi.Capabilities{ + StreamChunkInterceptor: responseInterceptorFunc{ + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) (pluginapi.StreamChunkInterceptResponse, error) { + modernGot = req + return pluginapi.StreamChunkInterceptResponse{Body: req.Body}, nil + }, + }, + }, + }, + }, + ) + if !host.StreamChunkPayloadIncludesRequestBody() { + t.Fatal("StreamChunkPayloadIncludesRequestBody() = false, want true when legacy stream interceptor is active") + } + + _ = host.InterceptStreamChunk(context.Background(), pluginapi.StreamChunkInterceptRequest{ + OriginalRequest: []byte("original"), + RequestBody: []byte("request"), + Body: []byte("chunk"), + ChunkIndex: 0, + }) + if string(legacyGot.OriginalRequest) != "original" || string(legacyGot.RequestBody) != "request" { + t.Fatalf("legacy payload bodies = original:%q body:%q, want preserved", legacyGot.OriginalRequest, legacyGot.RequestBody) + } + if len(modernGot.OriginalRequest) != 0 || len(modernGot.RequestBody) != 0 { + t.Fatalf("modern payload bodies = original:%q body:%q, want omitted", modernGot.OriginalRequest, modernGot.RequestBody) + } + + legacyGot = pluginapi.StreamChunkInterceptRequest{} + modernGot = pluginapi.StreamChunkInterceptRequest{} + _ = host.InterceptStreamChunk(context.Background(), pluginapi.StreamChunkInterceptRequest{ + OriginalRequest: []byte("original"), + RequestBody: []byte("request"), + ChunkIndex: pluginapi.StreamChunkHeaderInitIndex, + }) + if string(legacyGot.OriginalRequest) != "original" || string(modernGot.OriginalRequest) != "original" { + t.Fatalf("header-init bodies not preserved: legacy=%q modern=%q", legacyGot.OriginalRequest, modernGot.OriginalRequest) + } + + modernOnly := newHostWithRecords(capabilityRecord{ + id: "modern-only", + plugin: pluginapi.Plugin{ + SchemaVersion: pluginabi.SchemaVersionStreamChunkOmitRequestBody, + Capabilities: pluginapi.Capabilities{ + StreamChunkInterceptor: responseInterceptorFunc{ + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) (pluginapi.StreamChunkInterceptResponse, error) { + return pluginapi.StreamChunkInterceptResponse{}, nil + }, + }, + }, + }, + }) + if modernOnly.StreamChunkPayloadIncludesRequestBody() { + t.Fatal("StreamChunkPayloadIncludesRequestBody() = true, want false for schema v3+ only") + } +} + func TestHasRequestInterceptorsReflectsActiveRequestInterceptors(t *testing.T) { responseOnly := newHostWithRecords(capabilityRecord{ id: "response", diff --git a/internal/pluginhost/rpc_client.go b/internal/pluginhost/rpc_client.go index 01319fd78..881f232cf 100644 --- a/internal/pluginhost/rpc_client.go +++ b/internal/pluginhost/rpc_client.go @@ -68,8 +68,14 @@ func registerRPCPlugin(ctx context.Context, host *Host, id string, client plugin return pluginapi.Plugin{}, fmt.Errorf("plugin schema version %d is not supported", resp.SchemaVersion) } adapter := &rpcPluginAdapter{id: id, host: host, client: client} + schemaVersion := resp.SchemaVersion + if schemaVersion == 0 { + // Missing schema_version is treated as the original contract. + schemaVersion = 1 + } plugin := pluginapi.Plugin{ - Metadata: resp.Metadata, + Metadata: resp.Metadata, + SchemaVersion: schemaVersion, Capabilities: pluginapi.Capabilities{ FrontendAuthProviderExclusive: resp.Capabilities.FrontendAuthProvider && resp.Capabilities.FrontendAuthProviderExclusive, ExecutorModelScope: resp.Capabilities.ExecutorModelScope, diff --git a/internal/pluginhost/rpc_schema_test.go b/internal/pluginhost/rpc_schema_test.go index 1746b66a8..6b5256635 100644 --- a/internal/pluginhost/rpc_schema_test.go +++ b/internal/pluginhost/rpc_schema_test.go @@ -108,12 +108,16 @@ func TestRegisterRPCPluginSendsHostSchemaVersion(t *testing.T) { registerResult: validTestPlugin("schema"), }) - if _, errRegister := registerRPCPlugin(context.Background(), nil, "schema", lookup, pluginabi.MethodPluginRegister, []byte("mode: test")); errRegister != nil { + registered, errRegister := registerRPCPlugin(context.Background(), nil, "schema", lookup, pluginabi.MethodPluginRegister, []byte("mode: test")) + if errRegister != nil { t.Fatalf("registerRPCPlugin() error = %v", errRegister) } if lookup.lastLifecycle.SchemaVersion != pluginabi.SchemaVersion { t.Fatalf("lifecycle schema_version = %d, want %d", lookup.lastLifecycle.SchemaVersion, pluginabi.SchemaVersion) } + if registered.SchemaVersion != pluginabi.SchemaVersion { + t.Fatalf("registered SchemaVersion = %d, want %d", registered.SchemaVersion, pluginabi.SchemaVersion) + } if string(lookup.lastLifecycle.ConfigYAML) != "mode: test" { t.Fatalf("lifecycle config = %q, want input config", lookup.lastLifecycle.ConfigYAML) } @@ -146,6 +150,9 @@ func TestRegisterRPCPluginAcceptsModelRouterOnSchema1(t *testing.T) { if registered.Capabilities.ModelRouter == nil { t.Fatal("ModelRouter = nil, want adapter") } + if registered.SchemaVersion != 1 { + t.Fatalf("registered SchemaVersion = %d, want 1", registered.SchemaVersion) + } } func TestRPCModelRouteUsesAdapter(t *testing.T) { diff --git a/sdk/api/handlers/handlers_interceptors.go b/sdk/api/handlers/handlers_interceptors.go index 2c9cb282c..a8b356045 100644 --- a/sdk/api/handlers/handlers_interceptors.go +++ b/sdk/api/handlers/handlers_interceptors.go @@ -32,6 +32,25 @@ type streamInterceptorDetector interface { HasStreamInterceptors() bool } +// streamChunkRequestBodyPolicy reports whether payload stream-chunk interceptors +// still require OriginalRequest/RequestBody (legacy schema_version < 3). +type streamChunkRequestBodyPolicy interface { + StreamChunkPayloadIncludesRequestBody() bool +} + +// streamChunkPayloadIncludesRequestBody returns true when at least one active +// stream interceptor needs per-chunk request bodies. Evaluated per call so +// mid-stream plugin reloads stay correct. Unknown hosts default to true. +func streamChunkPayloadIncludesRequestBody(host PluginInterceptorHost) bool { + if host == nil { + return false + } + if policy, ok := host.(streamChunkRequestBodyPolicy); ok { + return policy.StreamChunkPayloadIncludesRequestBody() + } + return true +} + type requestInterceptorDetector interface { HasRequestInterceptors() bool } diff --git a/sdk/api/handlers/handlers_interceptors_test.go b/sdk/api/handlers/handlers_interceptors_test.go index 738b1a776..0b328f3e6 100644 --- a/sdk/api/handlers/handlers_interceptors_test.go +++ b/sdk/api/handlers/handlers_interceptors_test.go @@ -26,6 +26,8 @@ type handlerInterceptorTestHost struct { interceptResponse func(context.Context, pluginapi.ResponseInterceptRequest) pluginapi.ResponseInterceptResponse interceptStreamChunk func(context.Context, pluginapi.StreamChunkInterceptRequest) pluginapi.StreamChunkInterceptResponse completeRequest func(context.Context, pluginapi.RequestCompletion) + // includeStreamChunkRequestBodies simulates legacy schema_version < 3 plugins. + includeStreamChunkRequestBodies bool } type handlerInterceptorNoStreamTestHost struct { @@ -82,6 +84,15 @@ func (h *handlerInterceptorTestHost) CompleteRequest(ctx context.Context, comple } } +// StreamChunkPayloadIncludesRequestBody implements streamChunkRequestBodyPolicy. +// Default false simulates schema_version >= 3 (omit request bodies on payload chunks). +func (h *handlerInterceptorTestHost) StreamChunkPayloadIncludesRequestBody() bool { + if h == nil { + return false + } + return h.includeStreamChunkRequestBodies +} + type interceptorCaptureExecutor struct { provider string @@ -905,17 +916,23 @@ func TestHandlerStreamInterceptorRewritesAndDropsChunks(t *testing.T) { if req.RequestHeaders.Get("X-Stage") != "after" { t.Fatalf("stream request headers = %#v, want after-auth header", req.RequestHeaders) } - if string(req.OriginalRequest) != `{"stage":"after-stream"}` { - t.Fatalf("stream original request = %q, want after-auth body", req.OriginalRequest) - } - if string(req.RequestBody) != `{"stage":"after-stream"}` { - t.Fatalf("stream request body = %q, want after-auth body", req.RequestBody) - } if req.ChunkIndex == pluginapi.StreamChunkHeaderInitIndex { + if string(req.OriginalRequest) != `{"stage":"after-stream"}` { + t.Fatalf("stream original request = %q, want after-auth body", req.OriginalRequest) + } + if string(req.RequestBody) != `{"stage":"after-stream"}` { + t.Fatalf("stream request body = %q, want after-auth body", req.RequestBody) + } headers := cloneHeader(req.ResponseHeaders) headers.Set("X-Stream", "plugin") return pluginapi.StreamChunkInterceptResponse{Headers: headers} } + if len(req.OriginalRequest) != 0 { + t.Fatalf("payload chunk OriginalRequest = %q, want omitted for schema v3+", req.OriginalRequest) + } + if len(req.RequestBody) != 0 { + t.Fatalf("payload chunk RequestBody = %q, want omitted for schema v3+", req.RequestBody) + } if req.ResponseHeaders.Get("X-Upstream") != "stream" { t.Fatalf("stream response headers = %#v, want upstream header", req.ResponseHeaders) } @@ -957,6 +974,65 @@ func TestHandlerStreamInterceptorRewritesAndDropsChunks(t *testing.T) { } } +func TestHandlerStreamInterceptorLegacySchemaClonesRequestBodiesOnPayloadChunks(t *testing.T) { + model := "handler-interceptor-stream-legacy-clone-model" + executor := &interceptorCaptureExecutor{ + stream: func(ctx context.Context, auth *coreauth.Auth, req coreexecutor.Request, opts coreexecutor.Options) (*coreexecutor.StreamResult, error) { + chunks := make(chan coreexecutor.StreamChunk, 2) + chunks <- coreexecutor.StreamChunk{Payload: []byte("first")} + chunks <- coreexecutor.StreamChunk{Payload: []byte("second")} + close(chunks) + return &coreexecutor.StreamResult{ + Headers: http.Header{"X-Upstream": []string{"stream"}}, + Chunks: chunks, + }, nil + }, + } + handler := newInterceptorHandler(t, model, executor, &sdkconfig.SDKConfig{PassthroughHeaders: true}) + var payloadBodies [][]byte + handler.SetPluginHost(&handlerInterceptorTestHost{ + includeStreamChunkRequestBodies: true, + interceptRequestAfterAuth: func(ctx context.Context, req pluginapi.RequestInterceptRequest) pluginapi.RequestInterceptResponse { + return pluginapi.RequestInterceptResponse{Body: []byte(`{"stage":"legacy-stream"}`)} + }, + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) pluginapi.StreamChunkInterceptResponse { + if req.ChunkIndex == pluginapi.StreamChunkHeaderInitIndex { + if string(req.OriginalRequest) != `{"stage":"legacy-stream"}` || string(req.RequestBody) != `{"stage":"legacy-stream"}` { + t.Fatalf("header-init bodies = original:%q body:%q", req.OriginalRequest, req.RequestBody) + } + // Mutate delivered slices; later chunks must not observe this mutation. + req.OriginalRequest[0] = 'X' + req.RequestBody[0] = 'Y' + return pluginapi.StreamChunkInterceptResponse{} + } + if string(req.OriginalRequest) != `{"stage":"legacy-stream"}` { + t.Fatalf("payload OriginalRequest = %q, want isolated clone of after-auth body", req.OriginalRequest) + } + if string(req.RequestBody) != `{"stage":"legacy-stream"}` { + t.Fatalf("payload RequestBody = %q, want isolated clone of after-auth body", req.RequestBody) + } + payloadBodies = append(payloadBodies, req.OriginalRequest) + req.OriginalRequest[0] = 'Z' + return pluginapi.StreamChunkInterceptResponse{Body: req.Body} + }, + }) + + dataChan, _, errChan := handler.ExecuteStreamWithAuthManager(context.Background(), "openai", model, []byte(fmt.Sprintf(`{"model":%q}`, model)), "") + for range dataChan { + } + for msg := range errChan { + if msg != nil { + t.Fatalf("unexpected stream error: %+v", msg) + } + } + if len(payloadBodies) != 2 { + t.Fatalf("payload body deliveries = %d, want 2", len(payloadBodies)) + } + if &payloadBodies[0][0] == &payloadBodies[1][0] { + t.Fatal("payload OriginalRequest slices alias across chunks; want fresh clones") + } +} + func TestHandlerStreamInterceptorInitializesHeadersBeforeReturn(t *testing.T) { model := "handler-interceptor-stream-header-before-return-model" initStarted := make(chan struct{}) diff --git a/sdk/api/handlers/handlers_stream.go b/sdk/api/handlers/handlers_stream.go index 6669e0833..7be8243f5 100644 --- a/sdk/api/handlers/handlers_stream.go +++ b/sdk/api/handlers/handlers_stream.go @@ -81,19 +81,28 @@ func (h *BaseAPIHandler) streamWithPluginExecutor(ctx context.Context, entryProt streamInterceptorsActive := streamInterceptorsEnabled(interceptorHost) rawStreamHeaders := cloneHeader(streamResult.Headers) baseStreamHeaders := cloneHeader(streamResult.Headers) + // Request headers and request bodies are stream-invariant. Keep a private snapshot + // and clone into each interceptor call so plugins cannot mutate shared storage. + // Schema v3+ payload chunks omit these bodies (host also strips per plugin). + var streamRequestHeaders http.Header + var streamOriginalRequest []byte + var streamRequestBody []byte applyStreamHeaders := func(headers http.Header) { rawStreamHeaders = finalInterceptorHeaders(rawStreamHeaders, headers) } if streamInterceptorsActive { + streamRequestHeaders = cloneHeader(opts.Headers) + streamOriginalRequest = cloneBytes(opts.OriginalRequest) + streamRequestBody = cloneBytes(req.Payload) intercepted := interceptStreamChunk(ctx, interceptorHost, pluginapi.StreamChunkInterceptRequest{ RequestID: lifecycle.requestID(), SourceFormat: responseProtocol, Model: modelName, RequestedModel: originalRequestedModel, - RequestHeaders: cloneHeader(opts.Headers), + RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), - OriginalRequest: cloneBytes(opts.OriginalRequest), - RequestBody: cloneBytes(req.Payload), + OriginalRequest: cloneBytes(streamOriginalRequest), + RequestBody: cloneBytes(streamRequestBody), ChunkIndex: pluginapi.StreamChunkHeaderInitIndex, Metadata: opts.Metadata, }, execOptions.SkipInterceptorPluginID) @@ -161,20 +170,25 @@ func (h *BaseAPIHandler) streamWithPluginExecutor(ctx context.Context, entryProt } payload := cloneBytes(chunk.Payload) if streamInterceptorsActive { - intercepted := interceptStreamChunk(ctx, interceptorHost, pluginapi.StreamChunkInterceptRequest{ + chunkReq := pluginapi.StreamChunkInterceptRequest{ RequestID: lifecycle.requestID(), SourceFormat: responseProtocol, Model: modelName, RequestedModel: originalRequestedModel, - RequestHeaders: cloneHeader(opts.Headers), + RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), - OriginalRequest: cloneBytes(opts.OriginalRequest), - RequestBody: cloneBytes(req.Payload), Body: payload, HistoryChunks: cloneByteSlices(historyChunks), ChunkIndex: chunkIndex, Metadata: opts.Metadata, - }, execOptions.SkipInterceptorPluginID) + } + // Re-evaluate each chunk so mid-stream plugin reloads stay correct. + // Schema v3+ omits bodies here (one header-init clone only). + if streamChunkPayloadIncludesRequestBody(interceptorHost) { + chunkReq.OriginalRequest = cloneBytes(streamOriginalRequest) + chunkReq.RequestBody = cloneBytes(streamRequestBody) + } + intercepted := interceptStreamChunk(ctx, interceptorHost, chunkReq, execOptions.SkipInterceptorPluginID) applyStreamHeaders(intercepted.Headers) if len(intercepted.Body) > 0 { payload = cloneBytes(intercepted.Body) @@ -323,6 +337,12 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context streamClosedBeforeRead := false streamCanceledBeforeRead := false streamHeaderInitialized := false + // Request headers/bodies are stream-invariant after after-auth capture. Keep a private + // snapshot and clone into each interceptor call so plugins cannot mutate shared storage. + // Schema v3+ payload chunks omit these bodies (host also strips per plugin). + var streamRequestHeaders http.Header + var streamOriginalRequest []byte + var streamRequestBody []byte applyStreamHeaders := func(headers http.Header) { rawStreamHeaders = finalInterceptorHeaders(rawStreamHeaders, headers) @@ -333,15 +353,18 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context return } executedReq, executedOpts := executedRequest() + streamRequestHeaders = cloneHeader(executedOpts.Headers) + streamOriginalRequest = cloneBytes(executedOpts.OriginalRequest) + streamRequestBody = cloneBytes(executedReq.Payload) intercepted := interceptStreamChunk(ctx, interceptorHost, pluginapi.StreamChunkInterceptRequest{ RequestID: lifecycle.requestID(), SourceFormat: responseProtocol, Model: normalizedModel, RequestedModel: originalRequestedModel, - RequestHeaders: cloneHeader(executedOpts.Headers), + RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), - OriginalRequest: cloneBytes(executedOpts.OriginalRequest), - RequestBody: cloneBytes(executedReq.Payload), + OriginalRequest: cloneBytes(streamOriginalRequest), + RequestBody: cloneBytes(streamRequestBody), ChunkIndex: pluginapi.StreamChunkHeaderInitIndex, Metadata: executedOpts.Metadata, }, execOptions.SkipInterceptorPluginID) @@ -353,21 +376,25 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context applyStreamHeaderInit() payload = cloneBytes(payload) if streamInterceptorsActive { - executedReq, executedOpts := executedRequest() - intercepted := interceptStreamChunk(ctx, interceptorHost, pluginapi.StreamChunkInterceptRequest{ + chunkReq := pluginapi.StreamChunkInterceptRequest{ RequestID: lifecycle.requestID(), SourceFormat: responseProtocol, Model: normalizedModel, RequestedModel: originalRequestedModel, - RequestHeaders: cloneHeader(executedOpts.Headers), + RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), - OriginalRequest: cloneBytes(executedOpts.OriginalRequest), - RequestBody: cloneBytes(executedReq.Payload), Body: payload, HistoryChunks: cloneByteSlices(historyChunks), ChunkIndex: *chunkIndex, - Metadata: executedOpts.Metadata, - }, execOptions.SkipInterceptorPluginID) + Metadata: opts.Metadata, + } + // Re-evaluate each chunk so mid-stream plugin reloads stay correct. + // Schema v3+ omits bodies here (one header-init clone only). + if streamChunkPayloadIncludesRequestBody(interceptorHost) { + chunkReq.OriginalRequest = cloneBytes(streamOriginalRequest) + chunkReq.RequestBody = cloneBytes(streamRequestBody) + } + intercepted := interceptStreamChunk(ctx, interceptorHost, chunkReq, execOptions.SkipInterceptorPluginID) applyStreamHeaders(intercepted.Headers) if len(intercepted.Body) > 0 { payload = cloneBytes(intercepted.Body) diff --git a/sdk/pluginabi/types.go b/sdk/pluginabi/types.go index 916eb4eb0..97c41a136 100644 --- a/sdk/pluginabi/types.go +++ b/sdk/pluginabi/types.go @@ -7,7 +7,13 @@ const ( ABIVersion uint32 = 1 // SchemaVersion tracks the RPC JSON contract exchanged at plugin.register. // Version 2 adds request lifecycle completion and active request termination. - SchemaVersion uint32 = 2 + // Version 3 omits OriginalRequest/RequestBody on payload stream chunks + // (ChunkIndex >= 0); those fields remain on StreamChunkHeaderInitIndex only. + // Plugins that still need per-chunk request bodies should keep schema_version < 3. + SchemaVersion uint32 = 3 + // SchemaVersionStreamChunkOmitRequestBody is the first schema version that omits + // request bodies on payload stream-chunk interceptor calls. + SchemaVersionStreamChunkOmitRequestBody uint32 = 3 ) const ( diff --git a/sdk/pluginabi/types_test.go b/sdk/pluginabi/types_test.go index c594546d7..8fa63542e 100644 --- a/sdk/pluginabi/types_test.go +++ b/sdk/pluginabi/types_test.go @@ -27,8 +27,11 @@ func TestEnvelopeRoundTrip(t *testing.T) { } func TestMethodNamesAreStable(t *testing.T) { - if SchemaVersion != 2 { - t.Fatalf("SchemaVersion = %d, want 2", SchemaVersion) + if SchemaVersion != 3 { + t.Fatalf("SchemaVersion = %d, want 3", SchemaVersion) + } + if SchemaVersionStreamChunkOmitRequestBody != 3 { + t.Fatalf("SchemaVersionStreamChunkOmitRequestBody = %d, want 3", SchemaVersionStreamChunkOmitRequestBody) } if MethodPluginRegister != "plugin.register" { t.Fatalf("MethodPluginRegister = %q", MethodPluginRegister) diff --git a/sdk/pluginapi/types.go b/sdk/pluginapi/types.go index b1edde46b..6add5d694 100644 --- a/sdk/pluginapi/types.go +++ b/sdk/pluginapi/types.go @@ -15,6 +15,9 @@ type Plugin struct { Metadata Metadata // Capabilities declares the optional integration points implemented by the plugin. Capabilities Capabilities + // SchemaVersion is the plugin contract version negotiated at registration. + // Zero means unset (treated as legacy by the host). + SchemaVersion uint32 } // Metadata describes a plugin for registry, logging, and diagnostics. @@ -1087,9 +1090,17 @@ type StreamChunkInterceptRequest struct { RequestedModel string RequestHeaders http.Header ResponseHeaders http.Header + // OriginalRequest contains the raw client request body. + // Always populated on header-init (ChunkIndex == StreamChunkHeaderInitIndex), as a fresh clone. + // On payload chunks (ChunkIndex >= 0): + // - schema_version >= 3: omitted (nil); cache from header-init or request intercept hooks + // - schema_version < 3: populated as a fresh clone each call (legacy compatibility) + // Callers must treat this slice as read-only; hosts clone before delivery to keep snapshots isolated. OriginalRequest []byte - RequestBody []byte - Body []byte + // RequestBody contains the provider/executed request payload. + // Same population / cloning / schema-version rules as OriginalRequest. + RequestBody []byte + Body []byte // HistoryChunks contains a bounded recent history of chunks already delivered downstream. // The host currently retains at most 64 chunks and 1 MiB total history bytes. HistoryChunks [][]byte