iop/apps/edge/internal/openai/responses_stream_gate.go
toki d1e32b6e06 feat(openai): 반복 출력 복구를 구현한다
출력 반복을 요청 단위로 감지하고 안전한 continuation lifecycle을 보장하기 위해 Chat/Responses codec과 stream gate evidence를 함께 정렬한다.
2026-07-29 18:40:36 +09:00

1266 lines
46 KiB
Go

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
mu sync.Mutex
started bool
loaded bool
pending []streamgate.NormalizedEvent
}
func newOpenAIResponsesEventSource(dc *responsesDispatchContext, handle edgeservice.RunResult, holder *openAIResponsesResultHolder, usage *openAIStreamGateUsageHolder) *openAIResponsesEventSource {
return &openAIResponsesEventSource{dc: dc, handle: handle, holder: holder, usage: usage}
}
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)
if s.usage != nil {
s.usage.set(usageObservationFromOpenAIUsage(usage, len(reasoning)))
}
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)
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}
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)
if err != nil {
return nil, nil, err
}
bindOpenAIRecoveryAdmissionState(sink, dispatcher.admissionState())
initialSource, err := factory(initial)
if err != nil {
return nil, nil, err
}
controller := &openAIAttemptController{service: s.service, dispatch: dispatch, closeTransport: closeInitial}
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()
writeError(w, http.StatusInternalServerError, "run_error", "stream gate runtime unavailable")
return
}
registry, err := openAIStreamGateRegistrySnapshotFor(s.streamGateConfig(), fctx)
if err != nil {
closeInitial()
writeError(w, http.StatusInternalServerError, "run_error", "stream gate runtime unavailable")
return
}
runtime, usage, err := s.buildOpenAIResponsesStreamGateRuntimeFromAttempt(dc, initial, dispatch, closeInitial, sink, registry)
if err != nil {
closeInitial()
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)
labels := dc.usageLabels(s, responseModeNormalized)
if composite, ok := sink.(*openAICompositeReleaseSink); ok && composite.resolvedCodec() == openAIStreamGateCodecTunnel {
labels = dc.usageLabels(s, responseModePassthrough)
}
if poolSink, ok := sink.(*openAIResponsesPoolReleaseSink); ok && poolSink.resolvedCodec() == openAIStreamGateCodecTunnel {
labels = dc.usageLabels(s, responseModePassthrough)
}
status := streamGateUsageStatus(runErr, committed, success)
if status == usageStatusSuccess {
emitUsageMetrics(labels, status, usage.get())
return
}
emitUsageMetrics(labels, status, usageObservation{})
}