fix(executor): map message_too_big WebSocket errors to structured API responses

- Added `mapCodexWebsocketReadError` to handle `CloseMessageTooBig` errors with proper status and error code mapping.
- Updated error propagation and logging in Codex WebSocket executor.
- Introduced corresponding unit test to verify `message_too_big` error mapping in streamed responses.

Closes: #4017
This commit is contained in:
Luis Pater
2026-07-08 03:07:55 +08:00
parent 078ed1787b
commit 4f157fbdff
2 changed files with 82 additions and 6 deletions

View File

@@ -5,6 +5,7 @@ package executor
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"net"
@@ -347,8 +348,9 @@ func (e *CodexWebsocketsExecutor) Execute(ctx context.Context, auth *cliproxyaut
}
msgType, payload, errRead := readCodexWebsocketMessage(ctx, sess, conn, readCh)
if errRead != nil {
helps.RecordAPIWebsocketError(ctx, e.cfg, "read", errRead)
return resp, errRead
mappedErr := mapCodexWebsocketReadError(errRead)
helps.RecordAPIWebsocketError(ctx, e.cfg, "read", mappedErr)
return resp, mappedErr
}
if msgType != websocket.TextMessage {
if msgType == websocket.BinaryMessage {
@@ -608,11 +610,12 @@ func (e *CodexWebsocketsExecutor) ExecuteStream(ctx context.Context, auth *clipr
_ = send(cliproxyexecutor.StreamChunk{Err: ctx.Err()})
return
}
mappedErr := mapCodexWebsocketReadError(errRead)
terminateReason = "read_error"
terminateErr = errRead
helps.RecordAPIWebsocketError(ctx, e.cfg, "read", errRead)
reporter.PublishFailure(ctx, errRead)
_ = send(cliproxyexecutor.StreamChunk{Err: errRead})
terminateErr = mappedErr
helps.RecordAPIWebsocketError(ctx, e.cfg, "read", mappedErr)
reporter.PublishFailure(ctx, mappedErr)
_ = send(cliproxyexecutor.StreamChunk{Err: mappedErr})
return
}
if msgType != websocket.TextMessage {
@@ -724,6 +727,17 @@ func writeCodexWebsocketMessage(sess *codexWebsocketSession, conn *websocket.Con
return conn.WriteMessage(websocket.TextMessage, payload)
}
func mapCodexWebsocketReadError(err error) error {
if err == nil {
return nil
}
var closeErr *websocket.CloseError
if errors.As(err, &closeErr) && closeErr.Code == websocket.CloseMessageTooBig {
return statusErr{code: http.StatusRequestEntityTooLarge, msg: `{"error":{"message":"upstream websocket message too big","type":"invalid_request_error","code":"message_too_big"}}`}
}
return err
}
func buildCodexWebsocketRequestBody(body []byte) []byte {
if len(body) == 0 {
return nil

View File

@@ -227,6 +227,68 @@ func TestCodexWebsocketsExecuteStreamPropagatesUpstreamErrorForDownstreamWebsock
}
}
func TestCodexWebsocketsExecuteStreamMapsMessageTooBigClose(t *testing.T) {
upgrader := websocket.Upgrader{CheckOrigin: func(*http.Request) bool { return true }}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
t.Errorf("upgrade websocket: %v", err)
return
}
defer func() { _ = conn.Close() }()
if _, _, errRead := conn.ReadMessage(); errRead != nil {
t.Errorf("read upstream websocket message: %v", errRead)
return
}
deadline := time.Now().Add(time.Second)
closeMessage := websocket.FormatCloseMessage(websocket.CloseMessageTooBig, "message too big")
if errWrite := conn.WriteControl(websocket.CloseMessage, closeMessage, deadline); errWrite != nil {
t.Errorf("write close websocket message: %v", errWrite)
return
}
}))
defer server.Close()
exec := NewCodexWebsocketsExecutor(&config.Config{SDKConfig: config.SDKConfig{DisableImageGeneration: config.DisableImageGenerationAll}})
auth := &cliproxyauth.Auth{Attributes: map[string]string{"api_key": "sk-test", "base_url": server.URL}}
req := cliproxyexecutor.Request{
Model: "gpt-5-codex",
Payload: []byte(`{"model":"gpt-5-codex","input":[{"type":"message","role":"user","content":"hello"}]}`),
}
opts := cliproxyexecutor.Options{
SourceFormat: sdktranslator.FromString("openai-response"),
ResponseFormat: sdktranslator.FromString("openai-response"),
}
result, err := exec.ExecuteStream(context.Background(), auth, req, opts)
if err != nil {
t.Fatalf("ExecuteStream() error = %v", err)
}
select {
case chunk, ok := <-result.Chunks:
if !ok {
t.Fatal("stream closed before error chunk")
}
if chunk.Err == nil {
t.Fatal("error chunk Err = nil, want message-too-big error")
}
statusErr, ok := chunk.Err.(interface{ StatusCode() int })
if !ok {
t.Fatalf("error type %T does not expose StatusCode", chunk.Err)
}
if got := statusErr.StatusCode(); got != http.StatusRequestEntityTooLarge {
t.Fatalf("status = %d, want %d", got, http.StatusRequestEntityTooLarge)
}
if got := gjson.Get(chunk.Err.Error(), "error.code").String(); got != "message_too_big" {
t.Fatalf("error code = %q, want message_too_big; err=%v", got, chunk.Err)
}
case <-time.After(5 * time.Second):
t.Fatal("timed out waiting for error stream chunk")
}
}
func TestCodexWebsocketsUpstreamDisconnectChanSignalsOnInvalidate(t *testing.T) {
upgrader := websocket.Upgrader{CheckOrigin: func(*http.Request) bool { return true }}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {