Merge branch 'release/dev-974'
This commit is contained in:
commit
fd41adac77
7 changed files with 521 additions and 54 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)."
|
||||
|
|
|
|||
Loading…
Reference in a new issue