diff --git a/internal/pluginhost/adapters_interceptors.go b/internal/pluginhost/adapters_interceptors.go index ffc8c337..12f73112 100644 --- a/internal/pluginhost/adapters_interceptors.go +++ b/internal/pluginhost/adapters_interceptors.go @@ -277,7 +277,13 @@ func (h *Host) InterceptStreamChunkExcept(ctx context.Context, req pluginapi.Str nextReq.RequestBody = bytes.Clone(req.RequestBody) } nextReq.Body = bytes.Clone(current.Body) - nextReq.HistoryChunks = cloneByteSlices(req.HistoryChunks) + // Schema v5+ omits HistoryChunks on payload chunks to avoid cloning and re-encoding + // recent chunk windows across cgo/JSON for every frame. Legacy plugins still receive them. + if req.ChunkIndex != pluginapi.StreamChunkHeaderInitIndex && streamChunkOmitsHistory(record.plugin.SchemaVersion) { + nextReq.HistoryChunks = nil + } else { + nextReq.HistoryChunks = cloneByteSlices(req.HistoryChunks) + } nextReq.Metadata = cloneInterceptorMetadata(req.Metadata) if resp, ok := h.callStreamChunkInterceptor(ctx, record, interceptor, nextReq); ok { current.Headers = mergeHeaders(current.Headers, resp.Headers, resp.ClearHeaders) @@ -329,6 +335,28 @@ func streamChunkOmitsRequestBodies(schemaVersion uint32) bool { return schemaVersion >= pluginabi.SchemaVersionStreamChunkOmitRequestBody } +// StreamChunkPayloadIncludesHistory reports whether any active stream chunk +// interceptor still requires HistoryChunks on payload chunks +// (schema_version < SchemaVersionStreamChunkOmitHistory). +func (h *Host) StreamChunkPayloadIncludesHistory() bool { + if h == nil { + return false + } + for _, record := range h.activeRecords() { + if h.isPluginFused(record.id) || record.plugin.Capabilities.StreamChunkInterceptor == nil { + continue + } + if !streamChunkOmitsHistory(record.plugin.SchemaVersion) { + return true + } + } + return false +} + +func streamChunkOmitsHistory(schemaVersion uint32) bool { + return schemaVersion >= pluginabi.SchemaVersionStreamChunkOmitHistory +} + 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 5b40f84c..17fd0965 100644 --- a/internal/pluginhost/adapters_test.go +++ b/internal/pluginhost/adapters_test.go @@ -1770,6 +1770,86 @@ func TestStreamChunkRequestBodyPolicyBySchemaVersion(t *testing.T) { } } +func TestStreamChunkHistoryPolicyBySchemaVersion(t *testing.T) { + var legacyGot, modernGot pluginapi.StreamChunkInterceptRequest + host := newHostWithRecords( + capabilityRecord{ + id: "legacy", + plugin: pluginapi.Plugin{ + SchemaVersion: 4, + 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.SchemaVersionStreamChunkOmitHistory, + 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.StreamChunkPayloadIncludesHistory() { + t.Fatal("StreamChunkPayloadIncludesHistory() = false, want true when legacy stream interceptor is active") + } + + _ = host.InterceptStreamChunk(context.Background(), pluginapi.StreamChunkInterceptRequest{ + HistoryChunks: [][]byte{[]byte("first")}, + Body: []byte("chunk"), + ChunkIndex: 0, + }) + if len(legacyGot.HistoryChunks) != 1 || string(legacyGot.HistoryChunks[0]) != "first" { + t.Fatalf("legacy payload history = %#v, want preserved", legacyGot.HistoryChunks) + } + if len(modernGot.HistoryChunks) != 0 { + t.Fatalf("modern payload history = %#v, want omitted", modernGot.HistoryChunks) + } + + legacyGot = pluginapi.StreamChunkInterceptRequest{} + modernGot = pluginapi.StreamChunkInterceptRequest{} + _ = host.InterceptStreamChunk(context.Background(), pluginapi.StreamChunkInterceptRequest{ + HistoryChunks: [][]byte{[]byte("first")}, + Body: []byte("chunk"), + ChunkIndex: pluginapi.StreamChunkHeaderInitIndex, + }) + if len(legacyGot.HistoryChunks) != 1 || string(legacyGot.HistoryChunks[0]) != "first" { + t.Fatalf("legacy init history = %#v, want preserved", legacyGot.HistoryChunks) + } + if len(modernGot.HistoryChunks) != 1 || string(modernGot.HistoryChunks[0]) != "first" { + t.Fatalf("modern init history = %#v, want preserved", modernGot.HistoryChunks) + } + + modernOnly := newHostWithRecords(capabilityRecord{ + id: "modern-only", + plugin: pluginapi.Plugin{ + SchemaVersion: pluginabi.SchemaVersionStreamChunkOmitHistory, + Capabilities: pluginapi.Capabilities{ + StreamChunkInterceptor: responseInterceptorFunc{ + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) (pluginapi.StreamChunkInterceptResponse, error) { + return pluginapi.StreamChunkInterceptResponse{}, nil + }, + }, + }, + }, + }) + if modernOnly.StreamChunkPayloadIncludesHistory() { + t.Fatal("StreamChunkPayloadIncludesHistory() = true, want false for schema v5+ only") + } +} + func TestHasRequestInterceptorsReflectsActiveRequestInterceptors(t *testing.T) { responseOnly := newHostWithRecords(capabilityRecord{ id: "response", diff --git a/sdk/api/handlers/handlers_interceptors.go b/sdk/api/handlers/handlers_interceptors.go index d1b5d442..1f1cc6ad 100644 --- a/sdk/api/handlers/handlers_interceptors.go +++ b/sdk/api/handlers/handlers_interceptors.go @@ -51,6 +51,25 @@ func streamChunkPayloadIncludesRequestBody(host PluginInterceptorHost) bool { return true } +// streamChunkHistoryPolicy reports whether payload stream-chunk interceptors +// still require HistoryChunks (legacy schema_version < 5). +type streamChunkHistoryPolicy interface { + StreamChunkPayloadIncludesHistory() bool +} + +// streamChunkPayloadIncludesHistory returns true when at least one active +// stream interceptor needs per-chunk history chunks. Evaluated per call so +// mid-stream plugin reloads stay correct. Unknown hosts default to true. +func streamChunkPayloadIncludesHistory(host PluginInterceptorHost) bool { + if host == nil { + return false + } + if policy, ok := host.(streamChunkHistoryPolicy); ok { + return policy.StreamChunkPayloadIncludesHistory() + } + 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 42db2bcb..a13e0ab0 100644 --- a/sdk/api/handlers/handlers_interceptors_test.go +++ b/sdk/api/handlers/handlers_interceptors_test.go @@ -29,6 +29,8 @@ type handlerInterceptorTestHost struct { completeRequest func(context.Context, pluginapi.RequestCompletion) // includeStreamChunkRequestBodies simulates legacy schema_version < 3 plugins. includeStreamChunkRequestBodies bool + // includeStreamChunkHistory simulates legacy schema_version < 5 plugins. + includeStreamChunkHistory bool } type handlerInterceptorNoStreamTestHost struct { @@ -108,6 +110,15 @@ func (h *handlerInterceptorTestHost) StreamChunkPayloadIncludesRequestBody() boo return h.includeStreamChunkRequestBodies } +// StreamChunkPayloadIncludesHistory implements streamChunkHistoryPolicy. +// Default false simulates schema_version >= 5 (omit history on payload chunks). +func (h *handlerInterceptorTestHost) StreamChunkPayloadIncludesHistory() bool { + if h == nil { + return false + } + return h.includeStreamChunkHistory +} + type interceptorCaptureExecutor struct { provider string @@ -975,6 +986,7 @@ func TestHandlerStreamInterceptorRewritesAndDropsChunks(t *testing.T) { handler := newInterceptorHandler(t, model, executor, &sdkconfig.SDKConfig{PassthroughHeaders: true}) var streamCalls int handler.SetPluginHost(&handlerInterceptorTestHost{ + includeStreamChunkHistory: true, interceptRequestBeforeAuth: func(ctx context.Context, req pluginapi.RequestInterceptRequest) pluginapi.RequestInterceptResponse { headers := cloneHeader(req.Headers) if headers == nil { @@ -1122,6 +1134,97 @@ func TestHandlerStreamInterceptorLegacySchemaClonesRequestBodiesOnPayloadChunks( } } +func TestHandlerStreamInterceptorModernSchemaOmitsHistoryOnPayloadChunks(t *testing.T) { + model := "handler-interceptor-stream-modern-omit-history-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 observedHistories [][][]byte + handler.SetPluginHost(&handlerInterceptorTestHost{ + includeStreamChunkHistory: false, + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) pluginapi.StreamChunkInterceptResponse { + if req.ChunkIndex == pluginapi.StreamChunkHeaderInitIndex { + return pluginapi.StreamChunkInterceptResponse{} + } + observedHistories = append(observedHistories, cloneByteSlices(req.HistoryChunks)) + 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(observedHistories) != 2 { + t.Fatalf("payload chunk count = %d, want 2", len(observedHistories)) + } + for i, h := range observedHistories { + if len(h) != 0 { + t.Fatalf("chunk %d history = %#v, want omitted for schema v5+", i, h) + } + } +} + +func TestHandlerStreamInterceptorLegacySchemaClonesHistoryChunksOnPayloadChunks(t *testing.T) { + model := "handler-interceptor-stream-legacy-history-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 observedHistories [][][]byte + handler.SetPluginHost(&handlerInterceptorTestHost{ + includeStreamChunkHistory: true, + interceptStreamChunk: func(ctx context.Context, req pluginapi.StreamChunkInterceptRequest) pluginapi.StreamChunkInterceptResponse { + if req.ChunkIndex == pluginapi.StreamChunkHeaderInitIndex { + return pluginapi.StreamChunkInterceptResponse{} + } + observedHistories = append(observedHistories, cloneByteSlices(req.HistoryChunks)) + 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(observedHistories) != 2 { + t.Fatalf("payload chunk count = %d, want 2", len(observedHistories)) + } + if len(observedHistories[0]) != 0 { + t.Fatalf("chunk 0 history = %#v, want empty", observedHistories[0]) + } + if len(observedHistories[1]) != 1 || string(observedHistories[1][0]) != "first" { + t.Fatalf("chunk 1 history = %#v, want ['first']", observedHistories[1]) + } +} + 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 f1cc74eb..1f36ac4f 100644 --- a/sdk/api/handlers/handlers_stream.go +++ b/sdk/api/handlers/handlers_stream.go @@ -224,11 +224,14 @@ func (h *BaseAPIHandler) streamWithPluginExecutor(ctx context.Context, entryProt RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), Body: payload, - HistoryChunks: cloneByteSlices(historyChunks), ChunkIndex: chunkIndex, Metadata: opts.Metadata, } // Re-evaluate each chunk so mid-stream plugin reloads stay correct. + // Schema v5+ omits history here. + if streamChunkPayloadIncludesHistory(interceptorHost) { + chunkReq.HistoryChunks = cloneByteSlices(historyChunks) + } // Schema v3+ omits bodies here (one header-init clone only). if streamChunkPayloadIncludesRequestBody(interceptorHost) { chunkReq.OriginalRequest = cloneBytes(streamOriginalRequest) @@ -270,7 +273,7 @@ func (h *BaseAPIHandler) streamWithPluginExecutor(ctx context.Context, entryProt } select { case dataChan <- payload: - if streamInterceptorsActive { + if streamInterceptorsActive && streamChunkPayloadIncludesHistory(interceptorHost) { historyChunks = appendStreamInterceptorHistory(historyChunks, payload) } case <-done: @@ -441,11 +444,14 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context RequestHeaders: cloneHeader(streamRequestHeaders), ResponseHeaders: cloneHeader(rawStreamHeaders), Body: payload, - HistoryChunks: cloneByteSlices(historyChunks), ChunkIndex: *chunkIndex, Metadata: opts.Metadata, } // Re-evaluate each chunk so mid-stream plugin reloads stay correct. + // Schema v5+ omits history here. + if streamChunkPayloadIncludesHistory(interceptorHost) { + chunkReq.HistoryChunks = cloneByteSlices(historyChunks) + } // Schema v3+ omits bodies here (one header-init clone only). if streamChunkPayloadIncludesRequestBody(interceptorHost) { chunkReq.OriginalRequest = cloneBytes(streamOriginalRequest) @@ -658,7 +664,7 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context } return } - if streamInterceptorsActive { + if streamInterceptorsActive && streamChunkPayloadIncludesHistory(interceptorHost) { historyChunks = appendStreamInterceptorHistory(historyChunks, bootstrapPayload) } } @@ -722,7 +728,7 @@ func (h *BaseAPIHandler) executeStreamWithAuthManagerFormats(ctx context.Context } return } - if streamInterceptorsActive { + if streamInterceptorsActive && streamChunkPayloadIncludesHistory(interceptorHost) { historyChunks = appendStreamInterceptorHistory(historyChunks, payload) } } diff --git a/sdk/pluginabi/types.go b/sdk/pluginabi/types.go index 0df98f49..17790432 100644 --- a/sdk/pluginabi/types.go +++ b/sdk/pluginabi/types.go @@ -11,13 +11,19 @@ const ( // (ChunkIndex >= 0); those fields remain on StreamChunkHeaderInitIndex only. // Plugins that still need per-chunk request bodies should keep schema_version < 3. // Version 4 adds upstream WebSocket response event observation. - SchemaVersion uint32 = 4 + // Version 5 omits HistoryChunks on payload stream chunks (ChunkIndex >= 0); + // those fields remain on StreamChunkHeaderInitIndex only. Plugins that still need + // per-chunk history chunks should keep schema_version < 5. + SchemaVersion uint32 = 5 // SchemaVersionStreamChunkOmitRequestBody is the first schema version that omits // request bodies on payload stream-chunk interceptor calls. SchemaVersionStreamChunkOmitRequestBody uint32 = 3 // SchemaVersionWebSocketResponseObserver is the first schema version that supports // upstream WebSocket response event observation. SchemaVersionWebSocketResponseObserver uint32 = 4 + // SchemaVersionStreamChunkOmitHistory is the first schema version that omits + // history chunks on payload stream-chunk interceptor calls. + SchemaVersionStreamChunkOmitHistory uint32 = 5 ) const ( diff --git a/sdk/pluginabi/types_test.go b/sdk/pluginabi/types_test.go index 60489e72..5e96820c 100644 --- a/sdk/pluginabi/types_test.go +++ b/sdk/pluginabi/types_test.go @@ -27,8 +27,8 @@ func TestEnvelopeRoundTrip(t *testing.T) { } func TestMethodNamesAreStable(t *testing.T) { - if SchemaVersion != 4 { - t.Fatalf("SchemaVersion = %d, want 4", SchemaVersion) + if SchemaVersion != 5 { + t.Fatalf("SchemaVersion = %d, want 5", SchemaVersion) } if SchemaVersionWebSocketResponseObserver != 4 { t.Fatalf("SchemaVersionWebSocketResponseObserver = %d, want 4", SchemaVersionWebSocketResponseObserver) @@ -36,6 +36,9 @@ func TestMethodNamesAreStable(t *testing.T) { if SchemaVersionStreamChunkOmitRequestBody != 3 { t.Fatalf("SchemaVersionStreamChunkOmitRequestBody = %d, want 3", SchemaVersionStreamChunkOmitRequestBody) } + if SchemaVersionStreamChunkOmitHistory != 5 { + t.Fatalf("SchemaVersionStreamChunkOmitHistory = %d, want 5", SchemaVersionStreamChunkOmitHistory) + } if MethodPluginRegister != "plugin.register" { t.Fatalf("MethodPluginRegister = %q", MethodPluginRegister) } diff --git a/sdk/pluginapi/types.go b/sdk/pluginapi/types.go index adcf1ff1..161ca440 100644 --- a/sdk/pluginapi/types.go +++ b/sdk/pluginapi/types.go @@ -1110,6 +1110,11 @@ type StreamChunkInterceptRequest struct { 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. + // Always preserved on header-init (ChunkIndex == StreamChunkHeaderInitIndex) when non-empty. + // On payload chunks (ChunkIndex >= 0): + // - schema_version >= 5: omitted (nil) to avoid per-chunk cloning and serialization + // - schema_version < 5: populated as a fresh clone each call (legacy compatibility) + // Callers must treat these slices as read-only; hosts clone before delivery to keep snapshots isolated. HistoryChunks [][]byte // ChunkIndex starts at 0 for payload chunks. StreamChunkHeaderInitIndex marks the header-only initialization call. ChunkIndex int