From 432171096774d34d7f039e30bf886701d02025ee Mon Sep 17 00:00:00 2001 From: Sergei Bronnikov Date: Fri, 25 Sep 2026 09:33:30 +0100 Subject: [PATCH 1/3] ModifyResponse --- internal/server/web/proxy/x_custom.go | 68 +++++++++++++++++++++++++++ 1 file changed, 68 insertions(+) diff --git a/internal/server/web/proxy/x_custom.go b/internal/server/web/proxy/x_custom.go index c0132f6..cf665e7 100644 --- a/internal/server/web/proxy/x_custom.go +++ b/internal/server/web/proxy/x_custom.go @@ -1,6 +1,7 @@ package proxy import ( + "bytes" "context" "errors" "fmt" @@ -9,12 +10,53 @@ import ( "github.com/bricks-cloud/bricksllm/internal/telemetry" "github.com/bricks-cloud/bricksllm/internal/util" "github.com/gin-gonic/gin" + "io" "net/http" "net/http/httputil" "net/url" "strings" ) +// xCustomMaxCapturedResponseBytes caps how much of a response body gets +// buffered for later processing (c.Set). Responses larger than this are +// still proxied in full, they're just not captured. +const xCustomMaxCapturedResponseBytes = 5 * 1024 * 1024 // 5MB + +// xCustomCapturingBody wraps a response body so its bytes keep flowing to +// the client exactly as they arrive (no buffering delay, streaming stays +// real-time), while also being copied into an in-memory buffer for later +// use. onClose runs once the upstream body has been fully read/closed. +type xCustomCapturingBody struct { + io.ReadCloser + buf bytes.Buffer + exceeded bool + onClose func(data []byte, exceeded bool) +} + +func (b *xCustomCapturingBody) Read(p []byte) (int, error) { + n, err := b.ReadCloser.Read(p) + if n > 0 && !b.exceeded { + if b.buf.Len()+n > xCustomMaxCapturedResponseBytes { + b.exceeded = true + b.buf.Reset() + } else { + b.buf.Write(p[:n]) + } + } + + return n, err +} + +func (b *xCustomCapturingBody) Close() error { + err := b.ReadCloser.Close() + + if b.onClose != nil { + b.onClose(b.buf.Bytes(), b.exceeded) + } + + return err +} + func getXCustomHandler(prod bool) gin.HandlerFunc { return func(c *gin.Context) { log := util.GetLogFromCtx(c) @@ -68,6 +110,32 @@ func getXCustomHandler(prod bool) gin.HandlerFunc { r.Out.URL.Path, r.Out.URL.RawPath = target.Path, target.RawPath r.Out.WithContext(ctx) }, + ModifyResponse: func(res *http.Response) error { + if res.Body == nil { + return nil + } + + isStreaming := strings.Contains(res.Header.Get("Content-Type"), "text/event-stream") + + res.Body = &xCustomCapturingBody{ + ReadCloser: res.Body, + onClose: func(data []byte, exceeded bool) { + if exceeded || len(data) == 0 { + return + } + + if isStreaming { + c.Set("content", string(data)) + c.Set("streaming_response", data) + return + } + + c.Set("response", data) + }, + } + + return nil + }, } proxy.ServeHTTP(c.Writer, c.Request) } From 881a879da9d8769885fd736412a0fa0c3a14da26 Mon Sep 17 00:00:00 2001 From: Sergei Bronnikov Date: Fri, 25 Sep 2026 10:15:08 +0100 Subject: [PATCH 2/3] fixes --- internal/server/web/proxy/x_custom.go | 32 ++++++++++++++++++++++----- 1 file changed, 27 insertions(+), 5 deletions(-) diff --git a/internal/server/web/proxy/x_custom.go b/internal/server/web/proxy/x_custom.go index cf665e7..abb8606 100644 --- a/internal/server/web/proxy/x_custom.go +++ b/internal/server/web/proxy/x_custom.go @@ -11,6 +11,7 @@ import ( "github.com/bricks-cloud/bricksllm/internal/util" "github.com/gin-gonic/gin" "io" + "mime" "net/http" "net/http/httputil" "net/url" @@ -30,7 +31,11 @@ type xCustomCapturingBody struct { io.ReadCloser buf bytes.Buffer exceeded bool - onClose func(data []byte, exceeded bool) + // failed is set when the upstream body read ends in an error other than + // io.EOF (aborted/truncated response), so we don't capture partial data + // as if it were the complete response. + failed bool + onClose func(data []byte, exceeded bool) } func (b *xCustomCapturingBody) Read(p []byte) (int, error) { @@ -38,19 +43,26 @@ func (b *xCustomCapturingBody) Read(p []byte) (int, error) { if n > 0 && !b.exceeded { if b.buf.Len()+n > xCustomMaxCapturedResponseBytes { b.exceeded = true - b.buf.Reset() + // drop the reference so the already-buffered bytes can be + // garbage collected instead of being held for the rest of + // a possibly long-lived stream. + b.buf = bytes.Buffer{} } else { b.buf.Write(p[:n]) } } + if err != nil && err != io.EOF { + b.failed = true + } + return n, err } func (b *xCustomCapturingBody) Close() error { err := b.ReadCloser.Close() - if b.onClose != nil { + if b.onClose != nil && !b.failed { b.onClose(b.buf.Bytes(), b.exceeded) } @@ -109,13 +121,23 @@ func getXCustomHandler(prod bool) gin.HandlerFunc { r.SetURL(target) r.Out.URL.Path, r.Out.URL.RawPath = target.Path, target.RawPath r.Out.WithContext(ctx) + + // Let the transport negotiate and transparently decompress + // the upstream response itself; otherwise a forwarded + // client Accept-Encoding disables that and ModifyResponse + // would capture raw compressed bytes instead of text. + r.Out.Header.Del("Accept-Encoding") }, ModifyResponse: func(res *http.Response) error { - if res.Body == nil { + if res.Body == nil || res.StatusCode == http.StatusSwitchingProtocols { + // A 101 response's Body is an io.ReadWriteCloser used + // for bidirectional upgrade proxying (e.g. WebSocket); + // wrapping it would strip that and break the upgrade. return nil } - isStreaming := strings.Contains(res.Header.Get("Content-Type"), "text/event-stream") + mediaType, _, _ := mime.ParseMediaType(res.Header.Get("Content-Type")) + isStreaming := mediaType == "text/event-stream" res.Body = &xCustomCapturingBody{ ReadCloser: res.Body, From 637c55087ba812487f69899b5bca2b188a0557ac Mon Sep 17 00:00:00 2001 From: Sergei Bronnikov Date: Fri, 25 Sep 2026 11:18:18 +0100 Subject: [PATCH 3/3] stream --- internal/server/web/proxy/x_custom.go | 1 + 1 file changed, 1 insertion(+) diff --git a/internal/server/web/proxy/x_custom.go b/internal/server/web/proxy/x_custom.go index abb8606..a61df56 100644 --- a/internal/server/web/proxy/x_custom.go +++ b/internal/server/web/proxy/x_custom.go @@ -149,6 +149,7 @@ func getXCustomHandler(prod bool) gin.HandlerFunc { if isStreaming { c.Set("content", string(data)) c.Set("streaming_response", data) + c.Set("stream", true) return }