package openai import ( "bytes" "context" "encoding/json" "fmt" "net/http" "sort" "strings" "sync" "time" "go.uber.org/zap" edgeservice "iop/apps/edge/internal/service" "iop/packages/go/streamgate" ) type openAIResponsesAttemptResult struct { text string reasoning string toolCalls []any usage *openAIUsage dispatch edgeservice.RunDispatch collectErr error } type openAIResponsesResultHolder struct { mu sync.Mutex result openAIResponsesAttemptResult set bool } func (h *openAIResponsesResultHolder) beginAttempt() { h.mu.Lock() h.result = openAIResponsesAttemptResult{} h.set = false h.mu.Unlock() } func (h *openAIResponsesResultHolder) store(result openAIResponsesAttemptResult) { h.mu.Lock() h.result = result h.set = true h.mu.Unlock() } func (h *openAIResponsesResultHolder) get() (openAIResponsesAttemptResult, bool) { h.mu.Lock() defer h.mu.Unlock() return h.result, h.set } type openAIResponsesAttemptContext struct { mu sync.Mutex dc *responsesDispatchContext } // newOpenAIResponsesPoolTunnelDispatchContext creates the request-local // Responses runtime context for an initial provider tunnel. No public // Responses normalization is attempted here: the initial body remains a raw // provider passthrough. On recovery, the admission builder derives the // normalized Run and tunnel body from the private rebuilt body instead. func newOpenAIResponsesPoolTunnelDispatchContext(requestCtx *responsesRequestContext, pool edgeservice.ProviderPoolDispatchRequest) *responsesDispatchContext { return &responsesDispatchContext{ responsesRequestContext: requestCtx, req: responsesRequest{ Model: requestCtx.envelope.Model, Stream: requestCtx.envelope.Stream, }, runMetadata: cloneMetadata(pool.Run.Metadata), submitReq: pool.Run, poolDispatch: &pool, } } func (s *openAIResponsesAttemptContext) set(dc *responsesDispatchContext) { s.mu.Lock() s.dc = dc s.mu.Unlock() } func (s *openAIResponsesAttemptContext) get() *responsesDispatchContext { s.mu.Lock() defer s.mu.Unlock() return s.dc } type openAIResponsesEventSource struct { dc *responsesDispatchContext handle edgeservice.RunResult holder *openAIResponsesResultHolder usage *openAIStreamGateUsageHolder attempt *openAIAttemptUsage mu sync.Mutex started bool loaded bool pending []streamgate.NormalizedEvent } func newOpenAIResponsesEventSource(dc *responsesDispatchContext, handle edgeservice.RunResult, holder *openAIResponsesResultHolder, usage *openAIStreamGateUsageHolder, attempts ...*openAIAttemptUsage) *openAIResponsesEventSource { source := &openAIResponsesEventSource{dc: dc, handle: handle, holder: holder, usage: usage} if len(attempts) > 0 { source.attempt = attempts[0] } return source } func (s *openAIResponsesEventSource) NextEvent(ctx context.Context) (streamgate.NormalizedEvent, error) { s.mu.Lock() if !s.started { s.started = true s.holder.beginAttempt() s.mu.Unlock() return streamgate.NewResponseStartEvent(streamGateChannelDefault, http.StatusOK, map[string]string{"Content-Type": "application/json"}, time.Now()) } if len(s.pending) > 0 { event := s.pending[0] s.pending = s.pending[1:] s.mu.Unlock() return event, nil } if s.loaded { s.mu.Unlock() return newOpenAIProviderErrorEvent(streamGateErrorStreamClosed) } s.loaded = true s.mu.Unlock() text, reasoning, _, toolCalls, usage, _, err := collectRunResult(ctx, s.handle.Stream(), s.handle.WaitTimeout()) if err != nil { s.holder.store(openAIResponsesAttemptResult{dispatch: s.handle.Dispatch(), collectErr: err}) return newOpenAIProviderErrorEvent(streamGateErrorRunFailed) } text, reasoning, _ = normalizeCompletionOutput(s.dc.outputPolicy, text, reasoning, false) result := openAIResponsesAttemptResult{text: text, reasoning: reasoning, toolCalls: toolCalls, usage: usage, dispatch: s.handle.Dispatch()} s.holder.store(result) obs := usageObservationFromOpenAIUsage(usage, len(reasoning)) s.attempt.observe(obs) if s.usage != nil { s.usage.set(obs) } var events []streamgate.NormalizedEvent if text != "" { event, eventErr := streamgate.NewTextDeltaEvent(streamGateChannelDefault, text, time.Now()) if eventErr != nil { return streamgate.NormalizedEvent{}, eventErr } events = append(events, event) } if reasoning != "" { event, eventErr := streamgate.NewReasoningDeltaEvent(streamGateChannelDefault, reasoning, time.Now()) if eventErr != nil { return streamgate.NormalizedEvent{}, eventErr } events = append(events, event) } for i, raw := range toolCalls { call, ok := decodeResponsesToolCall(raw, i) if !ok || call.Arguments == "" { continue } event, eventErr := streamgate.NewToolCallFragmentEvent(streamGateChannelDefault, call.CallID, call.Name, call.Arguments, time.Now()) if eventErr != nil { return streamgate.NormalizedEvent{}, eventErr } events = append(events, event) } terminal, terminalErr := streamgate.NewTerminalEvent(streamGateChannelDefault, time.Now()) if terminalErr != nil { return streamgate.NormalizedEvent{}, terminalErr } events = append(events, terminal) s.mu.Lock() s.pending = append(s.pending, events...) event := s.pending[0] s.pending = s.pending[1:] s.mu.Unlock() return event, nil } var _ streamgate.NormalizedEventSource = (*openAIResponsesEventSource)(nil) type openAIResponsesToolCall struct { ID string CallID string Name string Arguments string } func decodeResponsesToolCall(raw any, index int) (openAIResponsesToolCall, bool) { encoded, err := json.Marshal(raw) if err != nil { return openAIResponsesToolCall{}, false } var value struct { ID string `json:"id"` CallID string `json:"call_id"` Name string `json:"name"` Arguments string `json:"arguments"` Function struct { Name string `json:"name"` Arguments string `json:"arguments"` } `json:"function"` } if json.Unmarshal(encoded, &value) != nil { return openAIResponsesToolCall{}, false } name := value.Name if name == "" { name = value.Function.Name } args := value.Arguments if args == "" { args = value.Function.Arguments } callID := value.CallID if callID == "" { callID = value.ID } if callID == "" { callID = fmt.Sprintf("call-%d", index) } id := value.ID if id == "" { id = fmt.Sprintf("fc-%d", index) } if name == "" { name = "function" } return openAIResponsesToolCall{ID: id, CallID: callID, Name: name, Arguments: args}, true } func responsesOutputItems(text string, toolCalls []any) []responsesOutputItem { items := []responsesOutputItem{{ Type: "message", Role: "assistant", Content: []responsesContentItem{{Type: "output_text", Text: text}}, }} for i, raw := range toolCalls { call, ok := decodeResponsesToolCall(raw, i) if !ok { continue } items = append(items, responsesOutputItem{ Type: "function_call", ID: call.ID, CallID: call.CallID, Name: call.Name, Arguments: call.Arguments, }) } return items } type openAIResponsesReleaseSink struct { server *Server w http.ResponseWriter req responsesRequest holder *openAIResponsesResultHolder recoveryAdmission *openAIRecoveryAdmissionState mu sync.Mutex terminalCommitted bool terminalSuccess bool } func (s *openAIResponsesReleaseSink) setRecoveryAdmissionState(state *openAIRecoveryAdmissionState) { s.mu.Lock() s.recoveryAdmission = state s.mu.Unlock() } func newOpenAIResponsesReleaseSink(server *Server, w http.ResponseWriter, dc *responsesDispatchContext, holder *openAIResponsesResultHolder) *openAIResponsesReleaseSink { return &openAIResponsesReleaseSink{server: server, w: w, req: dc.req, holder: holder} } func (s *openAIResponsesReleaseSink) terminalStatus() (bool, bool) { s.mu.Lock() defer s.mu.Unlock() return s.terminalCommitted, s.terminalSuccess } func (s *openAIResponsesReleaseSink) CommitResponseStart(context.Context, streamgate.ResponseStart) (streamgate.CommitState, error) { return streamgate.CommitStateStreamOpen, nil } func (s *openAIResponsesReleaseSink) Release(_ context.Context, event streamgate.ReleaseEvent) (streamgate.CommitState, error) { switch event.Kind() { case streamgate.EventKindTextDelta, streamgate.EventKindReasoningDelta, streamgate.EventKindToolCallFragment: return streamgate.CommitStateStreamOpen, nil default: return streamgate.CommitStateStreamOpen, fmt.Errorf("openai stream gate: responses sink does not support %q", event.Kind()) } } func (s *openAIResponsesReleaseSink) CommitTerminal(_ context.Context, terminal streamgate.TerminalResult) (streamgate.CommitState, error) { s.mu.Lock() defer s.mu.Unlock() s.terminalCommitted = true s.terminalSuccess = terminal.Success() result, ok := s.holder.get() if !terminal.Success() && s.recoveryAdmission.rejected() { writeError(s.w, http.StatusBadRequest, "invalid_request_error", openAIStreamGateCandidateRejectedMessage) return streamgate.CommitStateTerminalCommitted, nil } if !terminal.Success() || !ok || result.collectErr != nil { message := openAIStreamGateErrorMessage(terminal) status := http.StatusBadGateway if ok && result.collectErr != nil { message = result.collectErr.Error() status = httpStatusForRunError(result.collectErr) } writeError(s.w, status, "run_error", message) return streamgate.CommitStateTerminalCommitted, nil } var usage openAIUsage if result.usage != nil { usage = *result.usage } s.server.logger.Info("openai responses output", zap.String("run_id", result.dispatch.RunID), zap.Int("content_len", len(result.text)), zap.Int("reasoning_len", len(result.reasoning)), ) writeJSON(s.w, http.StatusOK, responsesResponse{ ID: "resp-" + result.dispatch.RunID, Object: "response", CreatedAt: time.Now().Unix(), Model: responseModel(s.req.Model, result.dispatch.Target), OutputText: result.text, Output: responsesOutputItems(result.text, result.toolCalls), Usage: usage, }) return streamgate.CommitStateTerminalCommitted, nil } var _ openAIStreamGateSink = (*openAIResponsesReleaseSink)(nil) // openAIResponsesPoolReleaseSink keeps the public Responses shape owned by // the caller rather than by a recovery attempt. A streaming caller may begin // on an SSE tunnel, while a private continuation is deliberately stream:false // and can subsequently be admitted to either provider path. In that case the // continuation's raw body must still be consumed from the tunnel codec, but // its semantic event is rendered as Responses SSE instead of being appended // after the already committed caller stream. type openAIResponsesPoolReleaseSink struct { w http.ResponseWriter flusher http.Flusher holder *openAIResponsesResultHolder selector *openAIStreamGateCodecSelector codec *openAITunnelCodecState mu sync.Mutex attemptStreaming bool wroteHeader bool terminalCommitted bool terminalSuccess bool recoveryAdmission *openAIRecoveryAdmissionState usage *openAIStreamGateUsageHolder responseState openAIResponsesSSEState } // openAIResponsesSSEState is deliberately owned by one caller stream. A pool // recovery can change its provider transport, but it must not reset the public // Responses event identity or sequence that the caller has already observed. type openAIResponsesSSEState struct { responseID string responseOpened bool model string sequenceNumber int64 nextOutput int nextFallback int message openAIResponsesSSEItem reasoning *openAIResponsesSSEItem functionCalls []*openAIResponsesSSEItem functionIndex map[string]int } // openAIResponsesSSEItem tracks one output item across a raw provider prefix and // a synthetic continuation. Each lifecycle transition is an independent flag so // terminal completion can emit exactly the transitions a raw prefix has not // already released, in protocol order, and never conflates a content-part close // with a text/arguments close or an item close. type openAIResponsesSSEItem struct { id string itemType string role string callID string name string outputIndex int outputIndexAssigned bool contentIndex int itemOpened bool contentOpened bool textDone bool argumentsDone bool contentDone bool itemDone bool value strings.Builder } func newOpenAIResponsesPoolReleaseSink(w http.ResponseWriter, holder *openAIResponsesResultHolder, selector *openAIStreamGateCodecSelector) *openAIResponsesPoolReleaseSink { flusher, _ := w.(http.Flusher) sink := &openAIResponsesPoolReleaseSink{ w: w, flusher: flusher, holder: holder, selector: selector, codec: &openAITunnelCodecState{}, } // These stable fallbacks are request-local. Raw provider events replace them // as soon as the provider supplies canonical Responses identifiers. sink.responseState.responseID = "resp-streamgate-1" sink.responseState.message = openAIResponsesSSEItem{id: "msg-streamgate-1", itemType: "message", role: "assistant", contentIndex: 0} sink.responseState.functionIndex = make(map[string]int) return sink } func (s *openAIResponsesPoolReleaseSink) setRecoveryAdmissionState(state *openAIRecoveryAdmissionState) { s.mu.Lock() s.recoveryAdmission = state s.mu.Unlock() } func (s *openAIResponsesPoolReleaseSink) bindAttempt(streaming bool, dispatch edgeservice.RunDispatch) { s.mu.Lock() s.attemptStreaming = streaming s.responseState.model = actualOpenAIModel(dispatch) s.mu.Unlock() } func (s *openAIResponsesPoolReleaseSink) setUsageHolder(usage *openAIStreamGateUsageHolder) { s.mu.Lock() s.usage = usage s.mu.Unlock() } func (s *openAIResponsesPoolReleaseSink) terminalStatus() (bool, bool) { s.mu.Lock() defer s.mu.Unlock() return s.terminalCommitted, s.terminalSuccess } func (s *openAIResponsesPoolReleaseSink) resolvedCodec() openAIStreamGateCodec { return s.selector.get() } func (s *openAIResponsesPoolReleaseSink) useRawTunnelWireLocked() bool { return s.selector.get() == openAIStreamGateCodecTunnel && s.attemptStreaming } func (s *openAIResponsesPoolReleaseSink) commitSSEHeaderLocked(status int) { if s.wroteHeader { return } if status == 0 { status = http.StatusOK } s.w.Header().Set("Content-Type", "text/event-stream") s.w.Header().Set("Cache-Control", "no-cache") s.w.Header().Set("Connection", "keep-alive") s.w.WriteHeader(status) s.wroteHeader = true } func (s *openAIResponsesPoolReleaseSink) writeSSELocked(value any) error { payload, err := json.Marshal(value) if err != nil { return err } if _, err := fmt.Fprintf(s.w, "data: %s\n\n", payload); err != nil { return err } if s.flusher != nil { s.flusher.Flush() } return nil } func (s *openAIResponsesPoolReleaseSink) writeDoneLocked() error { if _, err := fmt.Fprint(s.w, "data: [DONE]\n\n"); err != nil { return err } if s.flusher != nil { s.flusher.Flush() } return nil } func (s *openAIResponsesPoolReleaseSink) nextSequenceLocked() int64 { s.responseState.sequenceNumber++ return s.responseState.sequenceNumber } func (s *openAIResponsesPoolReleaseSink) observeRawResponsesWireLocked(payload []byte) { for len(payload) > 0 { frame, rest, ok := takeOpenAISSEFrame(payload) if !ok { return } payload = rest data := openAISSEData(frame) if data == "" || data == "[DONE]" { continue } var event map[string]any if json.Unmarshal([]byte(data), &event) == nil { s.observeRawResponsesEventLocked(event) } } } func (s *openAIResponsesPoolReleaseSink) observeRawResponsesEventLocked(event map[string]any) { state := &s.responseState eventType, _ := event["type"].(string) if sequence, ok := event["sequence_number"].(float64); ok && int64(sequence) > state.sequenceNumber { state.sequenceNumber = int64(sequence) } if response, ok := event["response"].(map[string]any); ok { if id, ok := response["id"].(string); ok && id != "" { state.responseID = id } if model, ok := response["model"].(string); ok && model != "" { state.model = model } } if eventType == "response.created" { state.responseOpened = true } itemID, _ := event["item_id"].(string) outputIndex, hasOutputIndex := event["output_index"].(float64) contentIndex, hasContentIndex := event["content_index"].(float64) // Resolve exactly one typed item from event-specific evidence. An output-item // event carries the authoritative id and type; delta and done events reference // an already-registered item by item_id plus the event family, so an untyped // non-message item can never fall through and overwrite the message identity. itemState := s.rawEventItemLocked(eventType, itemID, event["item"], int(outputIndex)) if itemState == nil { return } if hasOutputIndex { itemState.outputIndex = int(outputIndex) itemState.outputIndexAssigned = true if itemState.outputIndex >= state.nextOutput { state.nextOutput = itemState.outputIndex + 1 } } if hasContentIndex { itemState.contentIndex = int(contentIndex) } if item, ok := event["item"].(map[string]any); ok { if role, ok := item["role"].(string); ok && role != "" { itemState.role = role } if itemState.itemType == "function_call" { itemState.callID = nonEmptyString(item, "call_id", itemState.callID) itemState.name = nonEmptyString(item, "name", itemState.name) } } switch eventType { case "response.output_item.added": itemState.itemOpened = true case "response.output_item.done": itemState.itemDone = true case "response.content_part.added": itemState.contentOpened = true case "response.content_part.done": itemState.contentDone = true case "response.output_text.delta": if delta, ok := event["delta"].(string); ok { itemState.value.WriteString(delta) } case "response.output_text.done": itemState.textDone = true case "response.reasoning_text.delta", "response.reasoning_summary_text.delta": if delta, ok := event["delta"].(string); ok { itemState.value.WriteString(delta) } case "response.reasoning_text.done", "response.reasoning_summary_text.done": itemState.textDone = true case "response.function_call_arguments.delta": if delta, ok := event["delta"].(string); ok { itemState.value.WriteString(delta) } // Some providers echo function metadata on the delta; keep it only as a // fallback for the item-level identity, never the canonical source. itemState.callID = nonEmptyString(event, "call_id", itemState.callID) itemState.name = nonEmptyString(event, "name", itemState.name) case "response.function_call_arguments.done": itemState.argumentsDone = true } } // rawEventItemLocked resolves the single item a raw Responses event refers to. // Output-item events carry the authoritative id and type inside item; other // events reference an already-registered item by item_id, with the event family // implying a function or reasoning type so an untyped lookup never rewrites the // message identity with a function or reasoning item. func (s *openAIResponsesPoolReleaseSink) rawEventItemLocked(eventType, itemID string, rawItem any, outputIndex int) *openAIResponsesSSEItem { if item, ok := rawItem.(map[string]any); ok { if id, ok := item["id"].(string); ok && id != "" { itemID = id } itemType, _ := item["type"].(string) return s.itemForRawLocked(itemID, itemType, outputIndex) } itemType := "" if strings.HasPrefix(eventType, "response.function_call_arguments.") { itemType = "function_call" } else if strings.HasPrefix(eventType, "response.reasoning_") { itemType = "reasoning" } return s.itemForRawLocked(itemID, itemType, outputIndex) } func nonEmptyString(m map[string]any, key, fallback string) string { if value, ok := m[key].(string); ok && value != "" { return value } return fallback } func (s *openAIResponsesPoolReleaseSink) itemForRawLocked(itemID, itemType string, outputIndex int) *openAIResponsesSSEItem { state := &s.responseState if itemType == "function_call" { if itemID == "" { itemID = s.nextFallbackIDLocked("fc") } if index, ok := state.functionIndex[itemID]; ok { return state.functionCalls[index] } item := &openAIResponsesSSEItem{id: itemID, itemType: "function_call", outputIndex: outputIndex} state.functionIndex[itemID] = len(state.functionCalls) state.functionCalls = append(state.functionCalls, item) return item } if itemType == "reasoning" || (state.reasoning != nil && itemID != "" && state.reasoning.id == itemID) { if state.reasoning == nil { state.reasoning = &openAIResponsesSSEItem{id: itemID, itemType: "reasoning", outputIndex: outputIndex, contentIndex: 0} } return state.reasoning } if itemID != "" { state.message.id = itemID } return &state.message } func (s *openAIResponsesPoolReleaseSink) recordFunctionCallLocked(itemID, callID, name, arguments string, outputIndex int) *openAIResponsesSSEItem { if itemID == "" { itemID = callID } if itemID == "" { itemID = s.nextFallbackIDLocked("fc") } item := s.itemForRawLocked(itemID, "function_call", outputIndex) item.value.WriteString(arguments) if item.callID == "" { item.callID = callID } if item.name == "" { item.name = name } return item } func (s *openAIResponsesPoolReleaseSink) nextFallbackIDLocked(prefix string) string { s.responseState.nextFallback++ return fmt.Sprintf("%s-streamgate-%d", prefix, s.responseState.nextFallback) } // orderedItemsLocked returns every tracked item ordered by its preserved // provider output index so terminal completion and the terminal response // object present the items in the same schema-stable order. func (s *openAIResponsesPoolReleaseSink) orderedItemsLocked() []*openAIResponsesSSEItem { state := &s.responseState items := []*openAIResponsesSSEItem{&state.message} if state.reasoning != nil { items = append(items, state.reasoning) } items = append(items, state.functionCalls...) sort.SliceStable(items, func(i, j int) bool { return items[i].outputIndex < items[j].outputIndex }) return items } func (s *openAIResponsesPoolReleaseSink) responseObjectLocked() map[string]any { state := &s.responseState usage := usageObservation{} if s.usage != nil { usage = s.usage.get() } output := make([]any, 0, len(state.functionCalls)+2) for _, item := range s.orderedItemsLocked() { if !item.itemOpened { continue } output = append(output, s.completedItemObjectLocked(item)) } return map[string]any{ "id": state.responseID, "object": "response", "created_at": time.Now().Unix(), "model": state.model, "status": "completed", "output_text": state.message.value.String(), "output": output, "usage": map[string]any{"input_tokens": usage.inputTokens, "output_tokens": usage.outputTokens, "total_tokens": usage.inputTokens + usage.outputTokens}, } } func responsesContentPartType(itemType string) string { if itemType == "reasoning" { return "reasoning_text" } return "output_text" } // responsesContentPartObject builds one content part. Output-text parts carry // the required annotations/logprobs arrays; reasoning-text parts do not. func responsesContentPartObject(itemType, text string) map[string]any { part := map[string]any{"type": responsesContentPartType(itemType), "text": text} if itemType != "reasoning" { part["annotations"] = []any{} part["logprobs"] = []any{} } return part } // openingItemObjectLocked builds the in-progress item object for // response.output_item.added with empty opening content or arguments. func (s *openAIResponsesPoolReleaseSink) openingItemObjectLocked(item *openAIResponsesSSEItem) map[string]any { if item.itemType == "function_call" { return map[string]any{"id": item.id, "type": "function_call", "status": "in_progress", "call_id": item.callID, "name": item.name, "arguments": ""} } object := map[string]any{"id": item.id, "type": item.itemType, "status": "in_progress", "content": []any{}} if item.role != "" { object["role"] = item.role } return object } // completedItemObjectLocked builds the completed item object shared by // response.output_item.done and the terminal response output array. func (s *openAIResponsesPoolReleaseSink) completedItemObjectLocked(item *openAIResponsesSSEItem) map[string]any { if item.itemType == "function_call" { return map[string]any{"id": item.id, "type": "function_call", "status": "completed", "call_id": item.callID, "name": item.name, "arguments": item.value.String()} } object := map[string]any{"id": item.id, "type": item.itemType, "status": "completed", "content": []any{responsesContentPartObject(item.itemType, item.value.String())}} if item.role != "" { object["role"] = item.role } return object } func (s *openAIResponsesPoolReleaseSink) commitProviderErrorLocked(providerErr openAITunnelErrorResponse) (streamgate.CommitState, error) { for key, value := range providerErr.headers { s.w.Header().Set(key, value) } status := providerErr.status if status == 0 { status = http.StatusBadGateway } s.w.WriteHeader(status) s.wroteHeader = true if len(providerErr.body) > 0 { if _, err := s.w.Write(providerErr.body); err != nil { return streamgate.CommitStateTerminalCommitted, err } } if s.flusher != nil { s.flusher.Flush() } return streamgate.CommitStateTerminalCommitted, nil } func (s *openAIResponsesPoolReleaseSink) commitRawTunnelTerminalLocked() (streamgate.CommitState, error) { if payload, ok := s.codec.popTerminal(); ok && len(payload) > 0 { s.observeRawResponsesWireLocked(payload) if _, err := s.w.Write(payload); err != nil { return streamgate.CommitStateTerminalCommitted, err } if s.flusher != nil { s.flusher.Flush() } } return streamgate.CommitStateTerminalCommitted, nil } func (s *openAIResponsesPoolReleaseSink) commitResponsesErrorTerminalLocked(message string) (streamgate.CommitState, error) { s.commitSSEHeaderLocked(http.StatusOK) if err := s.writeSSELocked(map[string]any{ "type": "error", "code": "run_error", "message": message, "sequence_number": s.nextSequenceLocked(), }); err != nil { return streamgate.CommitStateTerminalCommitted, err } return streamgate.CommitStateTerminalCommitted, s.writeDoneLocked() } func (s *openAIResponsesPoolReleaseSink) CommitResponseStart(_ context.Context, rs streamgate.ResponseStart) (streamgate.CommitState, error) { s.mu.Lock() defer s.mu.Unlock() if s.wroteHeader { return streamgate.CommitStateStreamOpen, nil } if s.useRawTunnelWireLocked() { for key, value := range rs.Headers() { s.w.Header().Set(key, value) } status := rs.Status() if status == 0 { status = http.StatusOK } s.w.WriteHeader(status) s.wroteHeader = true if s.flusher != nil { s.flusher.Flush() } return streamgate.CommitStateStreamOpen, nil } s.commitSSEHeaderLocked(rs.Status()) return streamgate.CommitStateStreamOpen, nil } func (s *openAIResponsesPoolReleaseSink) ensureResponseOpenedLocked() error { state := &s.responseState if state.responseOpened { return nil } if err := s.writeSSELocked(map[string]any{"type": "response.created", "response": map[string]any{"id": state.responseID, "object": "response", "model": state.model, "status": "in_progress"}, "sequence_number": s.nextSequenceLocked()}); err != nil { return err } state.responseOpened = true return nil } func (s *openAIResponsesPoolReleaseSink) ensureItemOpenedLocked(item *openAIResponsesSSEItem) error { if err := s.ensureResponseOpenedLocked(); err != nil { return err } if !item.itemOpened { if !item.outputIndexAssigned { item.outputIndex = s.responseState.nextOutput item.outputIndexAssigned = true s.responseState.nextOutput++ } if err := s.writeSSELocked(map[string]any{"type": "response.output_item.added", "output_index": item.outputIndex, "item": s.openingItemObjectLocked(item), "sequence_number": s.nextSequenceLocked()}); err != nil { return err } item.itemOpened = true } if item.itemType == "function_call" || item.contentOpened { return nil } if err := s.writeSSELocked(map[string]any{"type": "response.content_part.added", "item_id": item.id, "output_index": item.outputIndex, "content_index": item.contentIndex, "part": responsesContentPartObject(item.itemType, ""), "sequence_number": s.nextSequenceLocked()}); err != nil { return err } item.contentOpened = true return nil } // completeItemLocked closes an opened item by emitting only the lifecycle // transitions it has not already released, in protocol order, and marks each // flag only after a successful write. A raw prefix that already emitted a // transition is never re-emitted, and an omitted one is always filled. func (s *openAIResponsesPoolReleaseSink) completeItemLocked(item *openAIResponsesSSEItem) error { if !item.itemOpened || item.itemDone { return nil } if item.itemType == "function_call" { if !item.argumentsDone { if err := s.writeSSELocked(map[string]any{"type": "response.function_call_arguments.done", "item_id": item.id, "output_index": item.outputIndex, "name": item.name, "arguments": item.value.String(), "sequence_number": s.nextSequenceLocked()}); err != nil { return err } item.argumentsDone = true } } else { if !item.textDone { doneType := "response.output_text.done" if item.itemType == "reasoning" { doneType = "response.reasoning_text.done" } payload := map[string]any{"type": doneType, "item_id": item.id, "output_index": item.outputIndex, "content_index": item.contentIndex, "text": item.value.String(), "sequence_number": s.nextSequenceLocked()} if item.itemType != "reasoning" { payload["logprobs"] = []any{} } if err := s.writeSSELocked(payload); err != nil { return err } item.textDone = true } if !item.contentDone { if err := s.writeSSELocked(map[string]any{"type": "response.content_part.done", "item_id": item.id, "output_index": item.outputIndex, "content_index": item.contentIndex, "part": responsesContentPartObject(item.itemType, item.value.String()), "sequence_number": s.nextSequenceLocked()}); err != nil { return err } item.contentDone = true } } if err := s.writeSSELocked(map[string]any{"type": "response.output_item.done", "output_index": item.outputIndex, "item": s.completedItemObjectLocked(item), "sequence_number": s.nextSequenceLocked()}); err != nil { return err } item.itemDone = true return nil } func (s *openAIResponsesPoolReleaseSink) completeOpenItemsLocked() error { for _, item := range s.orderedItemsLocked() { if err := s.completeItemLocked(item); err != nil { return err } } return nil } func (s *openAIResponsesPoolReleaseSink) Release(_ context.Context, event streamgate.ReleaseEvent) (streamgate.CommitState, error) { s.mu.Lock() defer s.mu.Unlock() if s.useRawTunnelWireLocked() { payload, ok := s.codec.popRelease() if !ok { return streamgate.CommitStateStreamOpen, fmt.Errorf("openai stream gate: Responses tunnel codec lost release wire") } if len(payload) > 0 { s.observeRawResponsesWireLocked(payload) if _, err := s.w.Write(payload); err != nil { return streamgate.CommitStateStreamOpen, err } if s.flusher != nil { s.flusher.Flush() } } return streamgate.CommitStateStreamOpen, nil } // A private non-stream tunnel can still queue a raw JSON response. Drain it // in lockstep, then keep the public stream valid by serializing the same // semantic release event as a Responses SSE payload. if s.selector.get() == openAIStreamGateCodecTunnel { _, _ = s.codec.popRelease() } state := &s.responseState var payload any switch event.Kind() { case streamgate.EventKindTextDelta: text, err := event.AsTextDelta() if err != nil { return streamgate.CommitStateStreamOpen, err } if err := s.ensureItemOpenedLocked(&state.message); err != nil { return streamgate.CommitStateStreamOpen, err } state.message.value.WriteString(text) payload = map[string]any{"type": "response.output_text.delta", "item_id": state.message.id, "output_index": state.message.outputIndex, "content_index": state.message.contentIndex, "delta": text, "sequence_number": s.nextSequenceLocked()} case streamgate.EventKindReasoningDelta: reasoning, err := event.AsReasoningDelta() if err != nil { return streamgate.CommitStateStreamOpen, err } if state.reasoning == nil { state.reasoning = &openAIResponsesSSEItem{id: s.nextFallbackIDLocked("reasoning"), itemType: "reasoning", contentIndex: 0} } if err := s.ensureItemOpenedLocked(state.reasoning); err != nil { return streamgate.CommitStateStreamOpen, err } state.reasoning.value.WriteString(reasoning) payload = map[string]any{"type": "response.reasoning_text.delta", "item_id": state.reasoning.id, "output_index": state.reasoning.outputIndex, "content_index": state.reasoning.contentIndex, "delta": reasoning, "sequence_number": s.nextSequenceLocked()} case streamgate.EventKindToolCallFragment: call, err := event.AsToolCallFragment() if err != nil { return streamgate.CommitStateStreamOpen, err } item := s.recordFunctionCallLocked(call.ID, call.ID, call.Name, call.Arguments, state.nextOutput) if err := s.ensureItemOpenedLocked(item); err != nil { return streamgate.CommitStateStreamOpen, err } // The delta event carries only its documented identity/index/value; the // canonical call_id and name stay on the item and its done payload. payload = map[string]any{"type": "response.function_call_arguments.delta", "item_id": item.id, "output_index": item.outputIndex, "delta": call.Arguments, "sequence_number": s.nextSequenceLocked()} default: return streamgate.CommitStateStreamOpen, fmt.Errorf("openai stream gate: Responses pool sink does not support %q", event.Kind()) } if err := s.writeSSELocked(payload); err != nil { return streamgate.CommitStateStreamOpen, err } return streamgate.CommitStateStreamOpen, nil } func (s *openAIResponsesPoolReleaseSink) CommitTerminal(_ context.Context, terminal streamgate.TerminalResult) (streamgate.CommitState, error) { s.mu.Lock() defer s.mu.Unlock() s.terminalCommitted = true s.terminalSuccess = terminal.Success() if providerErr, ok := s.codec.popErrorResponse(); ok && !s.wroteHeader { return s.commitProviderErrorLocked(providerErr) } if s.useRawTunnelWireLocked() && terminal.Success() { return s.commitRawTunnelTerminalLocked() } if s.selector.get() == openAIStreamGateCodecTunnel { _, _ = s.codec.popTerminal() } s.commitSSEHeaderLocked(http.StatusOK) if terminal.Success() { if err := s.completeOpenItemsLocked(); err != nil { return streamgate.CommitStateTerminalCommitted, err } if err := s.writeSSELocked(map[string]any{"type": "response.completed", "response": s.responseObjectLocked(), "sequence_number": s.nextSequenceLocked()}); err != nil { return streamgate.CommitStateTerminalCommitted, err } return streamgate.CommitStateTerminalCommitted, s.writeDoneLocked() } message := openAIStreamGateErrorMessage(terminal) if s.recoveryAdmission != nil && s.recoveryAdmission.rejected() { message = openAIStreamGateCandidateRejectedMessage } return s.commitResponsesErrorTerminalLocked(message) } var _ openAIStreamGateSink = (*openAIResponsesPoolReleaseSink)(nil) func newOpenAIResponsesRecoveryAdmissionBuilder(server *Server, initial *responsesDispatchContext, state *openAIResponsesAttemptContext) openAIAttemptAdmissionBuilder { return func(ctx context.Context, request streamgate.RebuiltRequest, body []byte) (openAIAttemptAdmission, error) { // Continuation repair has a private Responses item array, while exact // replay/schema repair retains the original public request body. Admit // the private form only through its exact decoder; preserving the normal // path here keeps existing lossless recovery semantics intact. resume, resumeErr := decodeOpenAIResponsesResumeRequest(body) var dc *responsesDispatchContext var err error if resumeErr == nil { dc, err = server.newResponsesResumeDispatchContext(initial.responsesRequestContext, resume) } else { var req responsesRequest if err = decodeResponsesRequest(json.NewDecoder(bytes.NewReader(body)), &req); err == nil { dc, err = server.newResponsesDispatchContext(initial.responsesRequestContext, req) } } if err != nil { return openAIAttemptAdmission{}, err } state.set(dc) if initial.poolDispatch == nil { return openAIAttemptAdmission{kind: openAIAdmissionRun, run: dc.submitReq}, nil } pool := *initial.poolDispatch pool.Run = dc.submitReq pool.Run.ProviderPool = true // A continuation is a private non-streaming Responses request. Keep the // provider-selection and auth hooks from the initial template, but make // every request-owned tunnel field agree with the admitted replacement // context rather than the caller's initial streaming tunnel. pool.Tunnel.Stream = dc.req.Stream pool.Tunnel.Metadata = cloneMetadata(dc.runMetadata) pool.Tunnel.EstimatedInputTokens = dc.submitReq.EstimatedInputTokens pool.Tunnel.ContextClass = dc.submitReq.ContextClass pool.Tunnel.BuildBody = func(target string) ([]byte, error) { return rewriteResponsesModel(body, target) } pool.PrepareRun = func(runReq edgeservice.SubmitRunRequest) (edgeservice.SubmitRunRequest, error) { runReq.Prompt = dc.submitReq.Prompt runReq.Input = dc.submitReq.Input runReq.Metadata = dc.submitReq.Metadata runReq.EstimatedInputTokens = dc.submitReq.EstimatedInputTokens runReq.ContextClass = dc.submitReq.ContextClass return runReq, nil } return openAIAttemptAdmission{kind: openAIAdmissionPool, pool: pool}, nil } } // buildOpenAIResponsesStreamGateRuntime builds the normalized initial-attempt // variant used by direct and initially-normalized provider-pool Responses // dispatches. Provider-pool tunnel attempts use the attempt form below so a // recovery can safely switch codecs before anything is committed downstream. func (s *Server) buildOpenAIResponsesStreamGateRuntime(dc *responsesDispatchContext, handle edgeservice.RunResult, sink openAIStreamGateSink, registry streamgate.FilterRegistrySnapshot) (*streamgate.RequestRuntime, *openAIStreamGateUsageHolder, error) { return s.buildOpenAIResponsesStreamGateRuntimeFromAttempt( dc, openAIAttemptTransport{path: openAIAdmissionRun, run: handle}, handle.Dispatch(), handle.Close, sink, registry, ) } // buildOpenAIResponsesStreamGateRuntimeFromAttempt keeps the initial // provider-pool tunnel in the same Responses-specific runtime used for every // recovery. That gives the admission builder one private resume body from // which both normalized and tunnel replacement requests are derived, rather // than retaining caller-derived Run/PrepareRun state from a generic tunnel // runtime. func (s *Server) buildOpenAIResponsesStreamGateRuntimeFromAttempt(dc *responsesDispatchContext, initial openAIAttemptTransport, dispatch edgeservice.RunDispatch, closeInitial func(), sink openAIStreamGateSink, registry streamgate.FilterRegistrySnapshot) (*streamgate.RequestRuntime, *openAIStreamGateUsageHolder, error) { holderSink, ok := sink.(*openAIResponsesReleaseSink) if !ok { if composite, compositeOK := sink.(*openAICompositeReleaseSink); compositeOK { holderSink, _ = composite.normalized.(*openAIResponsesReleaseSink) } } var holder *openAIResponsesResultHolder if holderSink != nil { holder = holderSink.holder } if poolSink, poolOK := sink.(*openAIResponsesPoolReleaseSink); poolOK { holder = poolSink.holder } if holder == nil { return nil, nil, fmt.Errorf("openai responses stream gate: normalized sink is required") } usage := &openAIStreamGateUsageHolder{} if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok { poolSink.setUsageHolder(usage) } state := &openAIResponsesAttemptContext{dc: dc} recoverySource := newOpenAIRecoverySourceStore(dc.ingress) rebuilder, err := newOpenAIRequestRebuilder(dc.ingress, openAIRebuildEndpointResponses, recoverySource, s.openAIResumeContextWindowTokens(dc.req.Model)) if err != nil { return nil, nil, err } selector := newOpenAIStreamGateCodecSelector(openAIStreamGateCodecNormalized) if composite, ok := sink.(*openAICompositeReleaseSink); ok { selector = composite.selector } if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok { selector = poolSink.selector } factory := func(transport openAIAttemptTransport) (streamgate.NormalizedEventSource, error) { var src streamgate.NormalizedEventSource switch transport.path { case openAIAdmissionRun: selector.set(openAIStreamGateCodecNormalized) openAIResponsesTunnelCodecStateForSink(sink).reset() attemptDC := state.get() if attemptDC == nil || transport.run == nil { return nil, fmt.Errorf("openai responses normalized attempt is incomplete") } if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok { poolSink.bindAttempt(attemptDC.req.Stream, transport.run.Dispatch()) } src = newOpenAIResponsesEventSource(attemptDC, transport.run, holder, usage, transport.usage) case openAIAdmissionTunnel: selector.set(openAIStreamGateCodecTunnel) codecState := openAIResponsesTunnelCodecStateForSink(sink) codecState.reset() attemptDC := state.get() if attemptDC == nil || transport.tunnel == nil { return nil, fmt.Errorf("openai responses tunnel attempt is incomplete") } if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok { poolSink.bindAttempt(attemptDC.req.Stream, transport.tunnel.Dispatch()) } assembler := &providerChatAssembler{streaming: attemptDC.req.Stream} rewriter := newProviderModelRewriter(attemptDC.req.Stream, "") tunnelSource := newOpenAITunnelEndpointEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, openAIRebuildEndpointResponses, codecState) src = &openAIStreamGateUsageTrackingTunnelSource{openAITunnelEventSource: tunnelSource, usage: usage, attempt: transport.usage} default: return nil, fmt.Errorf("openai responses unsupported attempt path %q", transport.path) } return newOpenAIRecoverySourceEventSource(src, recoverySource), nil } dispatcher, err := newOpenAIAttemptDispatcher(s.service, rebuilder.RebuiltStore(), newOpenAIResponsesRecoveryAdmissionBuilder(s, dc, state), factory, dc.usage) if err != nil { return nil, nil, err } bindOpenAIRecoveryAdmissionState(sink, dispatcher.admissionState()) initial.bindUsage(dispatch) initialSource, err := factory(initial) if err != nil { return nil, nil, err } controller := &openAIAttemptController{ service: s.service, dispatch: dispatch, closeTransport: closeInitial, usageRecorder: dc.usage, usageBinding: initial.usageBinding, usage: initial.usage, } binding, err := streamgate.NewAttemptBinding( openAIStreamGateSafeToken("attempt", dispatch.RunID), actualOpenAIModel(dispatch), actualOpenAIProvider(dispatch), actualOpenAIExecutionPath(dispatch, initial.path), initialSource, controller, ) if err != nil { return nil, nil, err } opts, err := s.streamGateRuntimeOptions() if err != nil { return nil, nil, err } snapRef, err := dc.ingress.recoveryRef() if err != nil { return nil, nil, err } snapshot, err := streamgate.NewRequestRuntimeSnapshot( openAIStreamGateSafeToken("req", dispatch.RunID), streamGateConfigGeneration, s.streamGateConfig().EffectiveEnvironment(), openAIRebuildEndpointResponses, openAIRebuildFamily, opts, registry, nil, snapRef, dispatcher, rebuilder, nil, nil, sink, ) if err != nil { return nil, nil, err } snapshot = snapshot.WithObservationSink(s.observationSink()) modelGroup := strings.TrimSpace(dc.req.Model) if modelGroup == "" { modelGroup = actualOpenAIModel(dispatch) } runtime, err := streamgate.NewRequestRuntime(snapshot, modelGroup, binding) if err != nil { return nil, nil, err } return runtime, usage, nil } func openAIResponsesTunnelCodecStateForSink(sink openAIStreamGateSink) *openAITunnelCodecState { if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok { return poolSink.codec } return openAITunnelCodecStateForSink(sink) } func (s *Server) runOpenAIResponsesStreamGate(w http.ResponseWriter, dc *responsesDispatchContext, handle edgeservice.RunResult) { s.runOpenAIResponsesStreamGateAttempt( w, dc, openAIAttemptTransport{path: openAIAdmissionRun, run: handle}, handle.Dispatch(), handle.Close, ) } // runOpenAIResponsesPoolStreamGate starts an initial provider-pool tunnel in // the Responses runtime. The initial provider body remains passthrough, but a // pre-commit recovery may now select either provider path and commits only the // successful replacement codec. func (s *Server) runOpenAIResponsesPoolStreamGate(w http.ResponseWriter, requestCtx *responsesRequestContext, pool edgeservice.ProviderPoolDispatchRequest, handle edgeservice.ProviderTunnelResult) { dc := newOpenAIResponsesPoolTunnelDispatchContext(requestCtx, pool) s.runOpenAIResponsesStreamGateAttempt( w, dc, openAIAttemptTransport{path: openAIAdmissionTunnel, tunnel: handle}, handle.Dispatch(), handle.Close, ) } func (s *Server) runOpenAIResponsesStreamGateAttempt(w http.ResponseWriter, dc *responsesDispatchContext, initial openAIAttemptTransport, dispatch edgeservice.RunDispatch, closeInitial func()) { holder := &openAIResponsesResultHolder{} normalized := newOpenAIResponsesReleaseSink(s, w, dc, holder) selector := newOpenAIStreamGateCodecSelector(openAIStreamGateCodecNormalized) var sink openAIStreamGateSink = normalized if dc.poolDispatch != nil { if dc.responsesRequestContext.envelope.Stream { sink = newOpenAIResponsesPoolReleaseSink(w, holder, selector) } else { tunnel := newOpenAIBufferedTunnelReleaseSink(w, nil, "") sink = newOpenAICompositeReleaseSink(selector, normalized, tunnel) } } fctx, err := s.openAIResponsesOutputFilterContext(dc.responsesRequestContext) if err != nil { closeInitial() dc.finishUsageRequest(usageStatusError, openAIAttemptResponseMode(initial.path)) writeError(w, http.StatusInternalServerError, "run_error", "stream gate runtime unavailable") return } registry, err := openAIStreamGateRegistrySnapshotFor(s.streamGateConfig(), fctx) if err != nil { closeInitial() dc.finishUsageRequest(usageStatusError, openAIAttemptResponseMode(initial.path)) writeError(w, http.StatusInternalServerError, "run_error", "stream gate runtime unavailable") return } runtime, _, err := s.buildOpenAIResponsesStreamGateRuntimeFromAttempt(dc, initial, dispatch, closeInitial, sink, registry) if err != nil { closeInitial() dc.finishUsageRequest(usageStatusError, openAIAttemptResponseMode(initial.path)) writeError(w, http.StatusInternalServerError, "run_error", "stream gate runtime unavailable") return } runErr := runtime.Run(dc.r.Context()) committed, success := sink.terminalStatus() _ = runtime.CloseRequestResources(context.Background(), runErr == nil && committed && success) responseMode := responseModeNormalized if composite, ok := sink.(*openAICompositeReleaseSink); ok && composite.resolvedCodec() == openAIStreamGateCodecTunnel { responseMode = responseModePassthrough } if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok && poolSink.resolvedCodec() == openAIStreamGateCodecTunnel { responseMode = responseModePassthrough } status := streamGateUsageStatus(runErr, committed, success) dc.finishUsageRequest(status, responseMode) }