From 2fcc1093c7629ab94b460d083520ffd9c8064815 Mon Sep 17 00:00:00 2001 From: toki Date: Thu, 13 Aug 2026 11:41:23 +0900 Subject: [PATCH] =?UTF-8?q?fix(openai):=20=EC=8A=A4=ED=8A=B8=EB=A6=BC=20?= =?UTF-8?q?=EC=A0=95=EC=A7=80=20=EC=A2=85=EB=A3=8C=EB=A5=BC=20=EB=B3=B5?= =?UTF-8?q?=EA=B5=AC=ED=95=9C=EB=8B=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit semantic filter 비활성 경로도 endpoint codec과 terminal sink를 유지해야 논리적 finish 뒤 정지가 열린 응답으로 남지 않는다. --- .../internal/openai/responses_stream_gate.go | 18 +- .../openai/stream_gate_pipeline_test.go | 90 +++++++++- .../internal/openai/stream_gate_runtime.go | 113 +++++++++--- .../openai/stream_gate_stall_recovery_test.go | 166 ++++++++++++++++++ .../internal/node/liveness_watchdog_test.go | 73 ++++++++ docs/edge-local-dev-guide.md | 12 ++ scripts/e2e-openai-lemonade.sh | 103 +++++++++-- 7 files changed, 521 insertions(+), 54 deletions(-) diff --git a/apps/edge/internal/openai/responses_stream_gate.go b/apps/edge/internal/openai/responses_stream_gate.go index 7368000c..aa6c0ead 100644 --- a/apps/edge/internal/openai/responses_stream_gate.go +++ b/apps/edge/internal/openai/responses_stream_gate.go @@ -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) diff --git a/apps/edge/internal/openai/stream_gate_pipeline_test.go b/apps/edge/internal/openai/stream_gate_pipeline_test.go index 893286e5..0f4ac289 100644 --- a/apps/edge/internal/openai/stream_gate_pipeline_test.go +++ b/apps/edge/internal/openai/stream_gate_pipeline_test.go @@ -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, diff --git a/apps/edge/internal/openai/stream_gate_runtime.go b/apps/edge/internal/openai/stream_gate_runtime.go index f8f30538..5a71f37d 100644 --- a/apps/edge/internal/openai/stream_gate_runtime.go +++ b/apps/edge/internal/openai/stream_gate_runtime.go @@ -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) diff --git a/apps/edge/internal/openai/stream_gate_stall_recovery_test.go b/apps/edge/internal/openai/stream_gate_stall_recovery_test.go index 37d0bd8b..14b8d04a 100644 --- a/apps/edge/internal/openai/stream_gate_stall_recovery_test.go +++ b/apps/edge/internal/openai/stream_gate_stall_recovery_test.go @@ -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 diff --git a/apps/node/internal/node/liveness_watchdog_test.go b/apps/node/internal/node/liveness_watchdog_test.go index 2f28e4a3..6db38171 100644 --- a/apps/node/internal/node/liveness_watchdog_test.go +++ b/apps/node/internal/node/liveness_watchdog_test.go @@ -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") diff --git a/docs/edge-local-dev-guide.md b/docs/edge-local-dev-guide.md index d0d32c9a..71edea1b 100644 --- a/docs/edge-local-dev-guide.md +++ b/docs/edge-local-dev-guide.md @@ -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 diff --git a/scripts/e2e-openai-lemonade.sh b/scripts/e2e-openai-lemonade.sh index ee760d8d..de6ea806 100755 --- a/scripts/e2e-openai-lemonade.sh +++ b/scripts/e2e-openai-lemonade.sh @@ -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" < "$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)."