fix(openai): 스트림 정지 종료를 복구한다

semantic filter 비활성 경로도 endpoint codec과 terminal sink를 유지해야 논리적 finish 뒤 정지가 열린 응답으로 남지 않는다.
This commit is contained in:
toki 2026-08-13 11:41:23 +09:00
parent eaf8fbaf21
commit 2fcc1093c7
7 changed files with 521 additions and 54 deletions

View file

@ -1112,7 +1112,6 @@ func (s *Server) buildOpenAIResponsesStreamGateRuntime(dc *responsesDispatchCont
// 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, stallStates ...*openAIStallRecoveryState) (*streamgate.RequestRuntime, *openAIStreamGateUsageHolder, error) {
semanticEnabled := s.streamGateSemanticEnabled()
var stallState *openAIStallRecoveryState
if len(stallStates) > 0 {
stallState = stallStates[0]
@ -1177,12 +1176,10 @@ func (s *Server) buildOpenAIResponsesStreamGateRuntimeFromAttempt(dc *responsesD
}
assembler := &providerChatAssembler{streaming: attemptDC.req.Stream}
rewriter := newProviderModelRewriter(attemptDC.req.Stream, "")
var tunnelSource *openAITunnelEventSource
if semanticEnabled {
tunnelSource = newOpenAITunnelEndpointEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, openAIRebuildEndpointResponses, codecState)
} else {
tunnelSource = newOpenAITunnelEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, codecState)
}
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)
@ -1278,12 +1275,7 @@ func (s *Server) runOpenAIResponsesStreamGateAttempt(w http.ResponseWriter, dc *
var sink openAIStreamGateSink = normalized
if dc.poolDispatch != nil {
if dc.responsesRequestContext.envelope.Stream {
if s.streamGateSemanticEnabled() {
sink = newOpenAIResponsesPoolReleaseSink(w, holder, selector)
} else {
flusher, _ := w.(http.Flusher)
sink = newOpenAICompositeReleaseSink(selector, normalized, newOpenAITunnelReleaseSink(w, flusher))
}
sink = newOpenAIResponsesPoolReleaseSink(w, holder, selector)
} else {
tunnel := newOpenAIBufferedTunnelReleaseSink(w, nil, "")
sink = newOpenAICompositeReleaseSink(selector, normalized, tunnel)

View file

@ -387,7 +387,7 @@ func TestOpenAITunnelCodecTerminalWire(t *testing.T) {
endpoint: openAIRebuildEndpointChat,
frames: [][]byte{
[]byte("data: {\"choices\":[{\"delta\":{\"content\":\"answer\"}}]}\n\n"),
[]byte("data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n"),
[]byte("data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}],\"usage\":{\"prompt_tokens\":2,\"completion_tokens\":1}}\n\n"),
[]byte("data: [DONE]\n\n"),
},
},
@ -413,12 +413,19 @@ func TestOpenAITunnelCodecTerminalWire(t *testing.T) {
state := &openAITunnelCodecState{}
codec := newOpenAITunnelEndpointCodec(tc.endpoint, state)
var events []streamgate.NormalizedEvent
for _, frame := range tc.frames {
for index, frame := range tc.frames {
decoded, err := codec.decode(frame, false)
if err != nil {
t.Fatalf("decode: %v", err)
}
events = append(events, decoded...)
if index < len(tc.frames)-1 {
for _, event := range events {
if event.Kind() == streamgate.EventKindTerminal {
t.Fatalf("protocol finish created a Core terminal before [DONE]/END: frame=%d events=%v", index, eventKinds(events))
}
}
}
}
decoded, err := codec.finishTransport(nil, false)
if err != nil {
@ -810,6 +817,85 @@ func TestOpenAITunnelHTTPErrorLifecycle(t *testing.T) {
}
}
// TestOpenAIDirectResponsesSemanticDisabledStallTerminal drives the generic
// direct Responses tunnel runtime (not the provider-pool Responses wrapper).
// Once raw Responses SSE is visible, a typed stall must discard a staged
// response.completed and append exactly one native error plus [DONE].
func TestOpenAIDirectResponsesSemanticDisabledStallTerminal(t *testing.T) {
srv := NewServer(config.EdgeOpenAIConf{TimeoutSec: 5}, &providerFakeRunService{}, nil)
rawBody := []byte(`{"model":"client-model","stream":true,"input":"hi"}`)
requestCtx := newTestRequestContext(t, routeDispatch{Adapter: "openai-compat", Target: "served-model", TimeoutSec: 5}, rawBody)
req := openAITunnelStreamGateRequest{
route: routeDispatch{Adapter: "openai-compat", Target: "served-model", TimeoutSec: 5},
ingress: requestCtx.ingress, endpoint: openAIRebuildEndpointResponses,
method: http.MethodPost, path: "/v1/responses", operation: string(config.OperationResponses),
stream: true, modelGroupKey: "client-model", semanticSet: true, semanticEnabled: false,
authorize: func(context.Context) (map[string]string, error) { return nil, nil },
rewriteBody: func(body []byte, _ string) ([]byte, error) { return body, nil },
}
frames := bufferedTunnelFrames(
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_RESPONSE_START, StatusCode: http.StatusOK, Headers: map[string]string{"Content-Type": "text/event-stream"}},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_BODY, Body: []byte(
"data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp-direct\",\"object\":\"response\",\"status\":\"in_progress\"},\"sequence_number\":1}\n\n" +
"data: {\"type\":\"response.output_text.delta\",\"item_id\":\"msg-direct\",\"output_index\":0,\"content_index\":0,\"delta\":\"direct-prefix\",\"sequence_number\":2}\n\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"raw-provider-terminal-must-stay-hidden\",\"object\":\"response\",\"status\":\"completed\"},\"sequence_number\":3}\n\n",
)},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR, Failure: confirmedStallFailure("available")},
)
closeCount := 0
handle := &countingDirectTunnelHandle{
fakeTunnelHandle: fakeTunnelHandle{dispatch: edgeservice.RunDispatch{
RunID: "direct-responses-stall", ModelGroupKey: "client-model", Adapter: "openai-compat",
Target: "served-model", ProviderID: "provider-a", ExecutionPath: string(edgeservice.ProviderPoolPathTunnel),
}, frames: frames},
closed: &closeCount,
}
fctx, err := srv.openAITunnelOutputFilterContext(req)
if err != nil {
t.Fatalf("openAITunnelOutputFilterContext: %v", err)
}
stallState, registration, err := openAIStallRecoveryRegistration(fctx)
if err != nil {
t.Fatalf("openAIStallRecoveryRegistration: %v", err)
}
registry, err := openAIStreamGateRegistrySnapshotFor(srv.streamGateConfig(), fctx, registration)
if err != nil {
t.Fatalf("openAIStreamGateRegistrySnapshotFor: %v", err)
}
w := newRecordingResponseWriter()
sink := newOpenAIResponsesPoolReleaseSink(w, &openAIResponsesResultHolder{}, newOpenAIStreamGateCodecSelector(openAIStreamGateCodecTunnel))
runtime, _, err := srv.buildOpenAITunnelStreamGateRuntime(req, handle, sink, registry, stallState)
if err != nil {
t.Fatalf("buildOpenAITunnelStreamGateRuntime: %v", err)
}
runErr := runtime.Run(t.Context())
committed, success := sink.terminalStatus()
if closeErr := runtime.CloseRequestResources(t.Context(), false); closeErr != nil {
t.Fatalf("CloseRequestResources: %v", closeErr)
}
body := w.body.String()
if runErr != nil || !committed || success || w.code != http.StatusOK || closeCount != 1 {
t.Fatalf("runtime=(err=%v committed=%v success=%v status=%d closes=%d body=%q)", runErr, committed, success, w.code, closeCount, body)
}
if !strings.Contains(body, "direct-prefix") || strings.Contains(body, "raw-provider-terminal-must-stay-hidden") || strings.Contains(body, "provider body") || strings.Contains(body, "raw provider metadata") {
t.Fatalf("direct Responses terminal leaked or lost wire: %q", body)
}
if strings.Count(body, `"type":"error"`) != 1 || strings.Count(body, "data: [DONE]") != 1 || strings.Contains(body, `"type":"response.completed"`) {
t.Fatalf("direct Responses terminal count mismatch: %q", body)
}
}
type countingDirectTunnelHandle struct {
fakeTunnelHandle
closed *int
}
func (h *countingDirectTunnelHandle) Close() {
if h.closed != nil {
*h.closed = *h.closed + 1
}
}
// TestOpenAITunnelHTTPErrorRawPassthroughRuntime closes the production
// source -> Core -> release-sink gap for unmatched provider errors. A non-2xx
// body stays opaque even when it looks like endpoint success/tool evidence,

View file

@ -647,6 +647,48 @@ func newOpenAITunnelEndpointEventSource(stream edgeservice.ProviderTunnelStream,
return source
}
// decodeEndpointTunnelBody keeps endpoint parsing active for known Chat and
// Responses frames while preserving the legacy behavior of opaque SSE data.
// The latter remains a release event per complete provider frame; recognized
// endpoint finish wire stays staged until [DONE], END, or a typed error owns
// the terminal.
func decodeEndpointTunnelBody(codec *openAITunnelEndpointCodec, body []byte) ([]streamgate.NormalizedEvent, error) {
if codec == nil || codec.terminal {
return nil, nil
}
codec.pending = append(codec.pending, body...)
var out []streamgate.NormalizedEvent
for {
frame, rest, ok := takeOpenAISSEFrame(codec.pending)
if !ok {
break
}
codec.pending = rest
data := openAISSEData(frame)
if strings.TrimSpace(data) != "" && strings.TrimSpace(data) != "[DONE]" && !json.Valid([]byte(data)) {
payload := append(append([]byte(nil), codec.stagedWire...), frame...)
codec.stagedWire = nil
codec.state.pushRelease(payload)
event, err := streamgate.NewTextDeltaEvent(streamGateChannelDefault, string(frame), time.Now())
if err != nil {
return nil, err
}
out = append(out, event)
continue
}
events, err := codec.decodeFrame(frame)
if err != nil {
return nil, err
}
out = append(out, events...)
if codec.terminal {
codec.pending = nil
break
}
}
return out, nil
}
func (s *openAITunnelEventSource) NextEvent(ctx context.Context) (streamgate.NormalizedEvent, error) {
s.mu.Lock()
if len(s.pending) > 0 {
@ -756,7 +798,7 @@ func (s *openAITunnelEventSource) translateFrame(frame *iop.ProviderTunnelFrame)
return events, nil
}
if s.codec != nil {
decoded, err := s.codec.decode(rewritten, false)
decoded, err := decodeEndpointTunnelBody(s.codec, rewritten)
if err != nil {
return nil, err
}
@ -775,12 +817,18 @@ func (s *openAITunnelEventSource) translateFrame(frame *iop.ProviderTunnelFrame)
return nil, nil
case iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR:
if s.codec == nil && (frame.GetFailure() == nil || frame.GetFailure().GetCode() != openAIStallFailureCode) {
if frame.GetFailure() == nil || frame.GetFailure().GetCode() != openAIStallFailureCode {
message := frame.GetError()
if message == "" {
message = "provider tunnel failed"
}
s.compatState.setCompatibilityProviderTerminal(message)
// The generic provider-error compatibility path owns its legacy
// pre/post-start behavior and must not acquire an endpoint-authored
// SSE terminal merely because endpoint parsing is always enabled.
s.compatState.mu.Lock()
s.compatState.endpoint = ""
s.compatState.mu.Unlock()
}
ev, err := newOpenAIProviderErrorEventFromFailure(frame.GetFailure(), streamGateErrorTunnelFailed)
if err != nil {
@ -1238,12 +1286,10 @@ func (s *Server) newOpenAIChatAttemptEventSourceFactory(
rewriter := newProviderModelRewriter(dc.req.Stream, dc.req.Model)
state := openAITunnelCodecStateForSink(cfg.sink)
state.reset()
var tunnelSrc *openAITunnelEventSource
if cfg.semanticEnabled {
tunnelSrc = newOpenAITunnelEndpointEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, openAIRebuildEndpointChat, state)
} else {
tunnelSrc = newOpenAITunnelEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, state)
}
tunnelSrc := newOpenAITunnelEndpointEventSource(
transport.tunnel.Stream(), transport.tunnel.WaitTimeout(),
rewriter, assembler, openAIRebuildEndpointChat, state,
)
dispatch := transport.tunnel.Dispatch()
tunnelSrc.onTerminal = func(obs *providerAssembledObservation, bodyBytes int) {
s.logger.Info("openai chat completion passthrough closed",
@ -1806,13 +1852,14 @@ func (s *Server) buildOpenAITunnelStreamGateRuntime(
if len(stallStates) > 0 {
stallState = stallStates[0]
}
semanticEnabled := req.semanticEnabled
if !req.semanticSet {
// Direct runtime fixtures predate the product-level semantic switch.
// Product callers always set semanticSet explicitly.
semanticEnabled = true
}
usage := &openAIStreamGateUsageHolder{}
var responseSink *openAIResponsesPoolReleaseSink
if req.endpoint == openAIRebuildEndpointResponses {
responseSink, _ = sink.(*openAIResponsesPoolReleaseSink)
if responseSink != nil {
responseSink.setUsageHolder(usage)
}
}
recoverySource := newOpenAIRecoverySourceStore(req.ingress)
rebuilder, err := newOpenAIRequestRebuilder(req.ingress, req.endpoint, recoverySource, s.openAIResumeContextWindowTokens(req.modelGroupKey))
@ -1832,13 +1879,15 @@ func (s *Server) buildOpenAITunnelStreamGateRuntime(
assembler := &providerChatAssembler{streaming: req.stream}
rewriter := newProviderModelRewriter(req.stream, req.requestModel)
state := openAITunnelCodecStateForSink(sink)
state.reset()
var src *openAITunnelEventSource
if semanticEnabled {
src = newOpenAITunnelEndpointEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, req.endpoint, state)
} else {
src = newOpenAITunnelEventSource(transport.tunnel.Stream(), transport.tunnel.WaitTimeout(), rewriter, assembler, state)
if responseSink != nil {
state = openAIResponsesTunnelCodecStateForSink(responseSink)
responseSink.bindAttempt(req.stream, transport.tunnel.Dispatch())
}
state.reset()
src := newOpenAITunnelEndpointEventSource(
transport.tunnel.Stream(), transport.tunnel.WaitTimeout(),
rewriter, assembler, req.endpoint, state,
)
tracking := &openAIStreamGateUsageTrackingTunnelSource{openAITunnelEventSource: src, usage: usage, attempt: transport.usage}
return newOpenAIRecoverySourceEventSource(tracking, recoverySource), nil
}
@ -1863,13 +1912,15 @@ func (s *Server) buildOpenAITunnelStreamGateRuntime(
initialAssembler := &providerChatAssembler{streaming: req.stream}
initialRewriter := newProviderModelRewriter(req.stream, req.requestModel)
initialState := openAITunnelCodecStateForSink(sink)
initialState.reset()
var initialEventSource *openAITunnelEventSource
if semanticEnabled {
initialEventSource = newOpenAITunnelEndpointEventSource(handle.Stream(), handle.WaitTimeout(), initialRewriter, initialAssembler, req.endpoint, initialState)
} else {
initialEventSource = newOpenAITunnelEventSource(handle.Stream(), handle.WaitTimeout(), initialRewriter, initialAssembler, initialState)
if responseSink != nil {
initialState = openAIResponsesTunnelCodecStateForSink(responseSink)
responseSink.bindAttempt(req.stream, dispatch)
}
initialState.reset()
initialEventSource := newOpenAITunnelEndpointEventSource(
handle.Stream(), handle.WaitTimeout(), initialRewriter, initialAssembler,
req.endpoint, initialState,
)
initialSource := &openAIStreamGateUsageTrackingTunnelSource{
openAITunnelEventSource: initialEventSource,
usage: usage,
@ -1948,8 +1999,14 @@ func (s *Server) runOpenAITunnelStreamGate(w http.ResponseWriter, r *http.Reques
req.semanticEnabled = s.streamGateSemanticEnabled()
req.semanticSet = true
flusher, _ := w.(http.Flusher)
var sink *openAITunnelReleaseSink
if req.stream {
var sink openAIStreamGateSink
if req.stream && req.endpoint == openAIRebuildEndpointResponses {
sink = newOpenAIResponsesPoolReleaseSink(
w,
&openAIResponsesResultHolder{},
newOpenAIStreamGateCodecSelector(openAIStreamGateCodecTunnel),
)
} else if req.stream {
sink = newOpenAITunnelReleaseSink(w, flusher)
} else {
sink = newOpenAIBufferedTunnelReleaseSink(w, flusher, req.requestModel)

View file

@ -244,6 +244,172 @@ func assertStallAttemptClosedOnce(t *testing.T, service *scriptedPoolRunService,
}
}
func logicalFinishStallAttempt(runID, content string, committed bool) scriptedPoolAttempt {
wire := []byte(fmt.Sprintf("data: {\"id\":\"logical-finish\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{\"content\":%q},\"finish_reason\":\"stop\"}]}\n\n", content))
if committed {
// The first content frame opens the stream. The following finish-only
// frame has no semantic event, so its identifiable wire remains pending
// until [DONE]/END and must be discarded by the later stall terminal.
wire = append(
[]byte("data: {\"id\":\"logical-finish\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"committed-prefix\"},\"finish_reason\":null}]}\n\n"),
[]byte(fmt.Sprintf("data: {\"id\":%q,\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n", content))...,
)
}
return scriptedPoolAttempt{
path: string(edgeservice.ProviderPoolPathTunnel), runID: runID, provider: "provider-a", target: "served-a",
frames: bufferedTunnelFrames(
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_RESPONSE_START, StatusCode: http.StatusOK, Headers: map[string]string{"Content-Type": "text/event-stream"}},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_BODY, Body: wire},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR, Failure: confirmedStallFailure("available")},
),
}
}
func logicalFinishMatrixServer(service runService, holdRunes int) *Server {
budget := 1
srv := NewServer(config.EdgeOpenAIConf{
TimeoutSec: 5,
StreamEvidenceGate: config.StreamEvidenceGateConf{
Enabled: true, MaxRequestFaultRecovery: &budget,
Filters: []config.StreamGateFilterPolicyConf{{
Filter: config.StreamGateFilterRepeatGuard, Enforcement: config.StreamGateFilterEnforcementBlocking,
Priority: 10, HoldEvidenceRunes: holdRunes,
}},
},
}, service, nil)
srv.SetModelCatalog([]config.ModelCatalogEntry{{
ID: "matrix-model", ContextWindowTokens: 8192,
Providers: map[string]string{"provider-a": "served-a", "provider-b": "served-b"},
}})
return srv
}
// TestOpenAIStallAfterLogicalFinishMatrix drives the observed shape through
// the production tunnel codec/Core/sink: finish_reason is pending protocol
// wire, while only the typed stall terminal decides replay or sanitized close.
func TestOpenAIStallAfterLogicalFinishMatrix(t *testing.T) {
t.Run("pre-commit pending finish replays once", func(t *testing.T) {
const pending = "pending-finish-must-stay-hidden"
const recovered = "recovered-output-only"
service := newScriptedPoolRunService(
logicalFinishStallAttempt("logical-pre-a", pending, false),
stallMatrixSuccessAttempt(openAIRebuildEndpointChat, string(edgeservice.ProviderPoolPathTunnel), true, "logical-pre-b", "provider-b", recovered),
)
w := runStallMatrixHandler(t, logicalFinishMatrixServer(service, 500), openAIRebuildEndpointChat, true, nil)
body := w.Body.String()
if w.Code != http.StatusOK || service.poolSubmits() != 2 || !strings.Contains(body, recovered) || strings.Contains(body, pending) {
t.Fatalf("pre-commit recovery=(status=%d dispatches=%d body=%q)", w.Code, service.poolSubmits(), body)
}
if strings.Contains(body, "provider body") || strings.Contains(body, "raw provider metadata") {
t.Fatalf("raw failure data leaked: %q", body)
}
requests := stallPoolRequests(service)
if len(requests) != 2 || requests[1].AvoidProviderID != "provider-a" || !requests[1].AllowAvoidedProviderFallback {
t.Fatalf("pre-commit replay requests=%+v", requests)
}
if strings.Count(body, recovered) != 1 {
t.Fatalf("recovered output count is not one: %q", body)
}
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "logical-pre-a")
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "logical-pre-b")
})
t.Run("post-commit finish emits one sanitized terminal", func(t *testing.T) {
const pending = "pending-finish-after-commit"
service := newScriptedPoolRunService(
logicalFinishStallAttempt("logical-post-a", pending, true),
stallMatrixSuccessAttempt(openAIRebuildEndpointChat, string(edgeservice.ProviderPoolPathTunnel), true, "forbidden", "provider-b", "must-not-replay"),
)
w := runStallMatrixHandler(t, logicalFinishMatrixServer(service, 1), openAIRebuildEndpointChat, true, nil)
body := w.Body.String()
if w.Code != http.StatusOK || service.poolSubmits() != 1 || !strings.Contains(body, "committed-prefix") || strings.Contains(body, "must-not-replay") {
t.Fatalf("post-commit recovery=(status=%d dispatches=%d body=%q)", w.Code, service.poolSubmits(), body)
}
if strings.Count(body, `"type":"run_error"`) != 1 || strings.Count(body, "data: [DONE]") != 1 {
t.Fatalf("post-commit terminal count mismatch: %q", body)
}
if strings.Contains(body, "provider body") || strings.Contains(body, "raw provider metadata") || strings.Contains(body, pending) {
t.Fatalf("pending/raw bytes leaked after commit: %q", body)
}
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "logical-post-a")
})
}
func semanticDisabledCommittedStallAttempt(endpoint, runID, pending string) scriptedPoolAttempt {
body := []byte(
"data: {\"id\":\"disabled-chat\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"disabled-committed-prefix\"},\"finish_reason\":null}]}\n\n" +
fmt.Sprintf("data: {\"id\":%q,\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0,\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n", pending),
)
if endpoint == openAIRebuildEndpointResponses {
body = []byte(
"data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp-disabled\",\"object\":\"response\",\"status\":\"in_progress\"},\"sequence_number\":1}\n\n" +
"data: {\"type\":\"response.output_text.delta\",\"item_id\":\"msg-disabled\",\"output_index\":0,\"content_index\":0,\"delta\":\"disabled-committed-prefix\",\"sequence_number\":2}\n\n" +
fmt.Sprintf("data: {\"type\":\"response.completed\",\"response\":{\"id\":%q,\"object\":\"response\",\"status\":\"completed\"},\"sequence_number\":3}\n\n", pending),
)
}
return scriptedPoolAttempt{
path: string(edgeservice.ProviderPoolPathTunnel), runID: runID, provider: "provider-a", target: "served-a",
frames: bufferedTunnelFrames(
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_RESPONSE_START, StatusCode: http.StatusOK, Headers: map[string]string{"Content-Type": "text/event-stream"}},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_BODY, Body: body},
&iop.ProviderTunnelFrame{Kind: iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR, Failure: confirmedStallFailure("available")},
),
}
}
// TestOpenAISemanticGateDisabledStallTerminalMatrix proves that endpoint
// framing and the private typed-stall owner remain active when semantic filters
// are disabled. Pre-commit stalls replay once; an opened Chat/Responses stream
// never replays and receives one sanitized endpoint-native terminal.
func TestOpenAISemanticGateDisabledStallTerminalMatrix(t *testing.T) {
for _, endpoint := range []string{openAIRebuildEndpointChat, openAIRebuildEndpointResponses} {
t.Run(endpoint+"/pre-commit", func(t *testing.T) {
const recovered = "semantic-disabled-recovered"
service := newScriptedPoolRunService(
stallMatrixFailureAttempt(string(edgeservice.ProviderPoolPathTunnel), "disabled-pre-a", "provider-a", "available"),
stallMatrixSuccessAttempt(endpoint, string(edgeservice.ProviderPoolPathTunnel), true, "disabled-pre-b", "provider-b", recovered),
)
w := runStallMatrixHandler(t, stallMatrixServer(service, false, 1), endpoint, true, nil)
body := w.Body.String()
if w.Code != http.StatusOK || service.poolSubmits() != 2 || strings.Count(body, recovered) != 1 {
t.Fatalf("pre-commit response=(status=%d dispatches=%d body=%q)", w.Code, service.poolSubmits(), body)
}
if strings.Contains(body, "provider body") || strings.Contains(body, "raw provider metadata") {
t.Fatalf("raw failure data leaked: %q", body)
}
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "disabled-pre-a")
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "disabled-pre-b")
})
t.Run(endpoint+"/post-commit", func(t *testing.T) {
const pending = "semantic-disabled-pending-terminal"
service := newScriptedPoolRunService(
semanticDisabledCommittedStallAttempt(endpoint, "disabled-post-a", pending),
stallMatrixSuccessAttempt(endpoint, string(edgeservice.ProviderPoolPathTunnel), true, "forbidden", "provider-b", "must-not-replay"),
)
w := runStallMatrixHandler(t, stallMatrixServer(service, false, 1), endpoint, true, nil)
body := w.Body.String()
if w.Code != http.StatusOK || service.poolSubmits() != 1 || !strings.Contains(body, "disabled-committed-prefix") || strings.Contains(body, "must-not-replay") {
t.Fatalf("post-commit response=(status=%d dispatches=%d body=%q)", w.Code, service.poolSubmits(), body)
}
if strings.Contains(body, pending) || strings.Contains(body, "provider body") || strings.Contains(body, "raw provider metadata") {
t.Fatalf("pending/raw failure data leaked: %q", body)
}
if strings.Count(body, "data: [DONE]") != 1 {
t.Fatalf("DONE terminal count mismatch: %q", body)
}
if endpoint == openAIRebuildEndpointChat {
if strings.Count(body, `"type":"run_error"`) != 1 {
t.Fatalf("Chat error terminal count mismatch: %q", body)
}
} else if strings.Count(body, `"type":"error"`) != 1 || strings.Contains(body, `"type":"response.completed"`) {
t.Fatalf("Responses error terminal mismatch: %q", body)
}
assertStallAttemptClosedOnce(t, service, string(edgeservice.ProviderPoolPathTunnel), "disabled-post-a")
})
}
}
// TestOpenAIStallRecoveryMatrix proves S05 through the supported production
// handlers and the production runtime adapter. It covers every endpoint/path/
// semantic-policy recovery product, then exercises the shared budget and every

View file

@ -453,6 +453,79 @@ func TestTunnelSinkStallClaimSerializesAcceptedFrame(t *testing.T) {
}
}
// TestTunnelWatchdogFinishFrameWithoutEndStallsOnce reproduces a provider
// stream that has emitted a logical OpenAI finish frame but has not closed its
// HTTP response. The finish frame is body progress, not a transport terminal:
// the watchdog must expire once and fence every later usage/END frame.
func TestTunnelWatchdogFinishFrameWithoutEndStallsOnce(t *testing.T) {
clock := newManualAttemptClock()
adapter := newControlledWatchdogAdapter("tunnel-finish-without-end")
n := newWatchdogNode(t, adapter, clock)
pipe := newWatchdogPipe(t)
done := make(chan error, 1)
go func() {
done <- n.OnProviderTunnelRequest(context.Background(), pipe.sess, &iop.ProviderTunnelRequest{
RunId: "tunnel-finish-without-end", TunnelId: "tunnel", Adapter: adapter.Name(), Target: "target", ResponseStallTimeoutMs: 1000,
})
}()
call := <-adapter.tunnelCalls
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{
RunID: "tunnel-finish-without-end", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindResponseStart,
StatusCode: 200, Headers: map[string]string{"Content-Type": "text/event-stream"},
}); err != nil {
t.Fatal(err)
}
logicalFinish := []byte("data: {\"choices\":[{\"delta\":{\"content\":\"finished body\"},\"finish_reason\":\"stop\"}]}\n\n")
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{
RunID: "tunnel-finish-without-end", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindBody, Body: logicalFinish,
}); err != nil {
t.Fatal(err)
}
if frame := waitTunnelFrame(t, pipe.frames); frame.GetKind() != iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_RESPONSE_START {
t.Fatalf("response start = %+v", frame)
}
if frame := waitTunnelFrame(t, pipe.frames); frame.GetKind() != iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_BODY || string(frame.GetBody()) != string(logicalFinish) {
t.Fatalf("logical finish body = %+v", frame)
}
stallTimer := clock.waitTimer(t, 0)
requireTimerDurations(t, stallTimer, time.Second, time.Second, time.Second)
stallTimer.fire()
waitContextCanceled(t, call.ctx)
grace := clock.waitTimer(t, 1)
adapter.tunnelReturn <- nil
if err := <-done; err != errProviderResponseStalled {
t.Fatalf("tunnel result = %v", err)
}
terminal := waitTunnelFrame(t, pipe.frames)
if terminal.GetKind() != iop.ProviderTunnelFrameKind_PROVIDER_TUNNEL_FRAME_KIND_ERROR ||
terminal.GetFailure().GetCode() != string(runtime.FailureCodeResponseStalled) ||
terminal.GetMetadata()["failure_code"] != string(runtime.FailureCodeResponseStalled) {
t.Fatalf("stall terminal = %+v", terminal)
}
requireTimerDurations(t, grace, defaultAttemptCloseGrace)
// These are the exact late frames an adapter can race after cancellation.
// Both must be accepted as no-ops by the fenced sink.
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{
RunID: "tunnel-finish-without-end", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindUsage,
Usage: &runtime.UsageStats{InputTokens: 1, OutputTokens: 1},
}); err != nil {
t.Fatal(err)
}
if err := call.sink.EmitTunnelFrame(context.Background(), runtime.ProviderTunnelFrame{
RunID: "tunnel-finish-without-end", TunnelID: "tunnel", Kind: runtime.ProviderTunnelFrameKindEnd, End: true,
}); err != nil {
t.Fatal(err)
}
select {
case frame := <-pipe.frames:
t.Fatalf("late frame escaped the terminal fence: %+v", frame)
default:
}
}
func TestRunWatchdogStaleExpiryYieldsToProgress(t *testing.T) {
clock := newManualAttemptClock()
adapter := newControlledWatchdogAdapter("run-stale-expiry")

View file

@ -72,6 +72,18 @@ nodes:
`models[]`를 사용할 경우 `openai.adapter`/`openai.target`은 하향 호환 fallback이다. 실제 routing은 `nodes[].providers[].id`를 기준으로 provider-pool에서 선택된다.
### Provider response-stall timeout ownership
dev-runtime에서 긴 응답의 정지 판정은 아래 순서를 유지한다.
```text
Node provider response_stall_timeout_ms = 120000
external caller boundary ≈ 180000
Edge request hard timeout > 180000
```
`response_stall_timeout_ms` 변경은 restart-required다. `config check``config refresh --mode dry-run`에서 이를 확인한 뒤 Edge와 Node를 같은 source ref로 rebuild/restart하고, 각 binary의 build identity와 실행 중인 process identity를 다시 대조한다. tracked 검증 근거에는 source/build/config 식별자, 단조 시간, terminal 개수와 결과 분류만 남기며 prompt, output, token, credential 원문은 기록하지 않는다.
확인:
```bash

View file

@ -43,8 +43,14 @@ cleanup() {
exec 3>&- 2>/dev/null || true
fi
if [ -n "$EDGE_PID" ]; then wait "$EDGE_PID" 2>/dev/null || true; fi
if [ -n "$NODE_PID" ]; then kill "$NODE_PID" 2>/dev/null || true; fi
if [ -n "$LEMONADE_PID" ]; then kill "$LEMONADE_PID" 2>/dev/null || true; fi
if [ -n "$NODE_PID" ]; then
kill "$NODE_PID" 2>/dev/null || true
wait "$NODE_PID" 2>/dev/null || true
fi
if [ -n "$LEMONADE_PID" ]; then
kill "$LEMONADE_PID" 2>/dev/null || true
wait "$LEMONADE_PID" 2>/dev/null || true
fi
if [ "${IOP_LEMONADE_KEEP_TMP:-0}" = "1" ]; then
echo "[openai-lemonade] keeping temp dir: $TMP_DIR"
return
@ -116,6 +122,8 @@ authorization := os.Args[3]
var modelRequests atomic.Int64
var chatRequests atomic.Int64
var responseRequests atomic.Int64
var hangRequests atomic.Int64
var hangCancellations atomic.Int64
mux := http.NewServeMux()
mux.HandleFunc("/v1/models", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet || r.RequestURI != "/v1/models" {
@ -128,7 +136,7 @@ return
}
modelRequests.Add(1)
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"object":"list","data":[{"id":"` + model + `"}]}`))
_, _ = w.Write([]byte(`{"object":"list","data":[{"id":"` + model + `"},{"id":"lemonade-profile-alias"}]}`))
})
mux.HandleFunc("/v1/chat/completions", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost || r.RequestURI != "/v1/chat/completions" {
@ -160,6 +168,15 @@ if f, ok := w.(http.Flusher); ok {
f.Flush()
}
}
if len(req.Messages) == 1 && req.Messages[0].Content == "IOP_HANG_AFTER_FINISH" {
hangRequests.Add(1)
_, _ = w.Write([]byte("data: " + `{"choices":[{"delta":{"content":"IOP_HANG_VISIBLE_PREFIX"},"finish_reason":"stop"}]}` + "\n\n"))
flush()
<-r.Context().Done()
hangCancellations.Add(1)
log.Print("IOP_HANG_PROVIDER_RAW_DIAGNOSTIC")
return
}
_, _ = w.Write([]byte("data: " + `{"choices":[{"delta":{"reasoning_content":"thinking..."}}]}` + "\n\n"))
flush()
_, _ = w.Write([]byte("data: " + `{"choices":[{"delta":{"content":"IOP_OPENAI_"}}]}` + "\n\n"))
@ -199,19 +216,23 @@ return
models := modelRequests.Load()
chats := chatRequests.Load()
responses := responseRequests.Load()
if models < 1 || chats != 2 || responses != 1 {
http.Error(w, fmt.Sprintf("unexpected request counts: models=%d chats=%d responses=%d", models, chats, responses), http.StatusInternalServerError)
hangs := hangRequests.Load()
cancellations := hangCancellations.Load()
if models < 1 || chats != 4 || responses != 1 || hangs != 1 || cancellations != 1 {
http.Error(w, fmt.Sprintf("unexpected request counts: models=%d chats=%d responses=%d hangs=%d cancellations=%d", models, chats, responses, hangs, cancellations), http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
_, _ = fmt.Fprintf(w, `{"models":%d,"chats":%d,"responses":%d}`, models, chats, responses)
_, _ = fmt.Fprintf(w, `{"models":%d,"chats":%d,"responses":%d,"hangs":%d,"cancellations":%d}`, models, chats, responses, hangs, cancellations)
})
log.Fatal(http.ListenAndServe(os.Args[1], mux))
}
EOF
FAKE_AUTHORIZATION="Bearer fake-lemonade-token"
go run "$FAKE_LEMONADE" "127.0.0.1:$LEMONADE_PORT" "$MODEL" "$FAKE_AUTHORIZATION" > "$TMP_DIR/fake_lemonade.out" 2>&1 &
FAKE_LEMONADE_BIN="$TMP_DIR/fake-lemonade"
go build -o "$FAKE_LEMONADE_BIN" "$FAKE_LEMONADE"
"$FAKE_LEMONADE_BIN" "127.0.0.1:$LEMONADE_PORT" "$MODEL" "$FAKE_AUTHORIZATION" > "$TMP_DIR/fake_lemonade.out" 2>&1 &
LEMONADE_PID=$!
wait_port 127.0.0.1 "$LEMONADE_PORT" "fake lemonade"
fi
@ -282,6 +303,7 @@ $(cat "$PROVIDER_HEADER_BLOCK")
- "lemonade-profile-alias"
health: available
capacity: 4
response_stall_timeout_ms: 200
EOF
cat > "$NODE_CONFIG" <<EOF
@ -359,6 +381,45 @@ else
fi
grep -q 'data: \[DONE\]' "$STREAM_OUT"
if [ "$MODE" = "fake" ]; then
# Flush a logical finish and deliberately keep the provider response open.
# Node's 200ms typed watchdog owns the failure before Edge's 30s request
# timeout. The committed stream must receive one sanitized terminal only.
HANG_OUT="$TMP_DIR/hang-after-finish.txt"
HANG_STARTED=$SECONDS
curl -fsS -N \
-H "Content-Type: application/json" \
-d '{"model":"client-request-model","stream":true,"max_tokens":32,"messages":[{"role":"user","content":"IOP_HANG_AFTER_FINISH"}]}' \
"http://127.0.0.1:$OPENAI_PORT/v1/chat/completions" > "$HANG_OUT"
HANG_ELAPSED=$((SECONDS - HANG_STARTED))
if (( HANG_ELAPSED >= 10 )); then
echo "[openai-lemonade] hang fixture exceeded the bounded IOP liveness window"
exit 1
fi
HANG_PREFIX_COUNT=$(grep -c 'IOP_HANG_VISIBLE_PREFIX' "$HANG_OUT" || true)
HANG_FINISH_COUNT=$(grep -c '"finish_reason":"stop"' "$HANG_OUT" || true)
HANG_ERROR_COUNT=$(grep -c '"type":"run_error"' "$HANG_OUT" || true)
HANG_DONE_COUNT=$(grep -c '^data: \[DONE\]$' "$HANG_OUT" || true)
if [ "$HANG_PREFIX_COUNT" -ne 1 ] || [ "$HANG_FINISH_COUNT" -ne 1 ] || [ "$HANG_ERROR_COUNT" -ne 1 ] || [ "$HANG_DONE_COUNT" -ne 1 ]; then
echo "[openai-lemonade] hang fixture terminal counts prefix=$HANG_PREFIX_COUNT finish=$HANG_FINISH_COUNT run_error=$HANG_ERROR_COUNT done=$HANG_DONE_COUNT"
exit 1
fi
if grep -q 'IOP_HANG_PROVIDER_RAW_DIAGNOSTIC' "$HANG_OUT"; then
echo "[openai-lemonade] raw provider diagnostic leaked to the client"
exit 1
fi
# A normal request after cancellation proves the provider handler and the
# Node/Edge path remain usable after the fenced attempt.
POST_CANCEL_OUT="$TMP_DIR/post-cancel-stream.txt"
curl -fsS -N \
-H "Content-Type: application/json" \
-d '{"model":"client-request-model","stream":true,"max_tokens":32,"messages":[{"role":"user","content":"Reply after cancellation"}]}' \
"http://127.0.0.1:$OPENAI_PORT/v1/chat/completions" > "$POST_CANCEL_OUT"
grep -q '"content":"IOP_OPENAI_' "$POST_CANCEL_OUT"
[ "$(grep -c '^data: \[DONE\]$' "$POST_CANCEL_OUT")" -eq 1 ]
fi
# Command-path smoke via the edge CLI.
SMOKE_OUT="$TMP_DIR/iop-edge-smoke-openai.txt"
(cd "$REPO_ROOT" && go run ./apps/edge/cmd/edge smoke openai \
@ -371,15 +432,35 @@ grep -q "IOP Edge OpenAI Smoke Test SUCCESS!" "$SMOKE_OUT"
if [ "$MODE" = "fake" ]; then
FAKE_ASSERT_OUT="$TMP_DIR/fake_assert.json"
curl -fsS "http://127.0.0.1:$LEMONADE_PORT/assert" > "$FAKE_ASSERT_OUT"
grep -q '"chats":2' "$FAKE_ASSERT_OUT"
grep -q '"chats":4' "$FAKE_ASSERT_OUT"
grep -q '"responses":1' "$FAKE_ASSERT_OUT"
grep -q '"hangs":1' "$FAKE_ASSERT_OUT"
grep -q '"cancellations":1' "$FAKE_ASSERT_OUT"
fi
if grep -i -E "node reported error|error run_id=|\[[^]]+-evt\] error|panic:" "$EDGE_OUT" "$NODE_OUT" >/dev/null; then
FAILURE_MATCHES="$TMP_DIR/failure-matches.txt"
grep -i -E "node reported error|error run_id=|\[[^]]+-evt\] error|panic:" "$EDGE_OUT" "$NODE_OUT" > "$FAILURE_MATCHES" || true
if [ "$MODE" = "fake" ]; then
UNEXPECTED_FAILURES="$TMP_DIR/unexpected-failures.txt"
grep -vi -E "response_stalled|provider response stalled" "$FAILURE_MATCHES" > "$UNEXPECTED_FAILURES" || true
else
UNEXPECTED_FAILURES="$FAILURE_MATCHES"
fi
if [ -s "$UNEXPECTED_FAILURES" ]; then
echo "[openai-lemonade] detected failure marker"
echo "=== EDGE OUTPUT ==="; cat "$EDGE_OUT"
echo "=== NODE OUTPUT ==="; cat "$NODE_OUT"
echo "[openai-lemonade] raw runtime output retained only in the temporary artifact directory"
exit 1
fi
if [ "$MODE" = "fake" ]; then
grep -qi -E "response_stalled|provider response stalled" "$EDGE_OUT" "$NODE_OUT"
kill "$LEMONADE_PID"
wait "$LEMONADE_PID" 2>/dev/null || true
if kill -0 "$LEMONADE_PID" 2>/dev/null; then
echo "[openai-lemonade] fake provider survived teardown"
exit 1
fi
LEMONADE_PID=""
fi
echo "[openai-lemonade] OpenAI-compatible Lemonade serving test PASSED (mode=$MODE)."