diff --git a/core/http/endpoints/openai/realtime.go b/core/http/endpoints/openai/realtime.go index c152d2c22..ebe8b5241 100644 --- a/core/http/endpoints/openai/realtime.go +++ b/core/http/endpoints/openai/realtime.go @@ -203,6 +203,91 @@ type Session struct { // decision is serialized through respcoord.Coordinator, guaranteeing at most // one live response. See realtime_respcoord.go. respSink *responseSink + + // sessionCtx is the session-lifetime context: cancelled at teardown + // (conncoord Teardown, BEFORE respSink.shutdown joins the response + // goroutines) but NOT by barge-in / response supersede (those cancel + // respSink's per-response contexts only). Committed-turn work that must + // outlive a barge-in but must never outlive the session — the + // transcription of a VAD commit (issue #12445) — runs under it. + sessionCtx context.Context + sessionCancel context.CancelFunc + + // commitOrderMu guards commitTail AND is the shared ordering boundary + // between the two commit producers (see issueCommit). + commitOrderMu sync.Mutex + // commitTail is the newest commit slot. Slots form a chain that orders + // user-item appends in speech order, not transcription-completion order + // (see commitSlot). + commitTail *commitSlot +} + +// commitSlot orders committed user turns in the conversation. The +// transcriptions of consecutive turns run in parallel (each under the session +// context, so a barge-in on a later turn does not cancel an earlier +// transcription), but a turn may append its user item only after the previous +// turn has appended its own — and the slot releases (done closes) only after +// the previous commit has FINISHED, on every exit path (success, error, empty +// transcript, teardown). Without the append gate a fast second transcription +// would commit ahead of a slow first one; without the release gate a failed +// middle commit would release the third turn before the first appended. In +// both cases the assistant response — and the conversation history — would +// see the user input out of order or missing (issue #12445 + review follow-up). +type commitSlot struct { + // prevDone is the previous slot's done channel (nil for the first commit); + // this commit's item append waits on it, and its slot release waits on it + // too (on every exit path). + prevDone chan struct{} + // done is closed when this commit has fully finished — appended its user + // item or aborted (error, empty transcript, teardown) — AND the + // predecessor has finished, so the next commit is released in order. + done chan struct{} +} + +// nextCommitSlot claims the next position in the commit order. It is called +// at commit ISSUE time (VAD CommitTurn / client input_audio_buffer.commit) — +// synchronously on the issuing goroutine, in the order the turns were +// detected — before the (parallel) transcriptions start, so slot order == +// speech order even if the commit goroutines schedule out of order. The +// production commit paths go through issueCommit (slot claim + issue under +// one lock); this primitive is also used by tests that drive a commit body +// directly. +func (session *Session) nextCommitSlot() *commitSlot { + session.commitOrderMu.Lock() + defer session.commitOrderMu.Unlock() + return session.nextCommitSlotLocked() +} + +// nextCommitSlotLocked claims the next slot; commitOrderMu must be held. +func (session *Session) nextCommitSlotLocked() *commitSlot { + var prevDone chan struct{} + if session.commitTail != nil { + prevDone = session.commitTail.done + } + slot := &commitSlot{prevDone: prevDone, done: make(chan struct{})} + session.commitTail = slot + return slot +} + +// issueCommit is the single ordering boundary both commit producers — the VAD +// CommitTurn (realtime_turncoord.go) and the client input_audio_buffer.commit +// (read loop) — must go through. Claiming the next commit slot and issuing +// the commit body happen under ONE lock, so slot order == issue order: the +// coordinator's supersession (a newer issue cancels the in-flight response) +// can only cancel a response issued LATER in speech order, never an earlier +// turn whose slot a later issue would overtake (issue #12445, review +// follow-up: "slot reservation and response issuance need a shared ordering +// boundary across both producers"). respSink.issue is non-blocking — it +// registers the body and applies the coordinator start; the commit work runs +// in the spawned goroutine — so holding the lock across it does not stall +// VAD/barge-in handling. +func (session *Session) issueCommit(parent context.Context, source respcoord.Source, run func(ctx context.Context, slot *commitSlot)) { + session.commitOrderMu.Lock() + defer session.commitOrderMu.Unlock() + slot := session.nextCommitSlotLocked() + session.respSink.issue(parent, source, func(ctx context.Context) { + run(ctx, slot) + }) } func (s *Session) installVoiceBinding(voice string, params map[string]string, release func()) { @@ -614,6 +699,11 @@ func runRealtimeSession(application *application.Application, t Transport, model // into two overlapping responses (see realtime_respcoord.go). session.respSink = newResponseSink() + // Session-lifetime context for committed-turn work that must outlive a + // barge-in but not the session (transcription, commit ordering). Cancelled + // at teardown BEFORE respSink.shutdown joins the response goroutines. + session.sessionCtx, session.sessionCancel = context.WithCancel(context.Background()) + // Create a default conversation conversationID := generateConversationID() conversation := &Conversation{ @@ -946,8 +1036,11 @@ func runRealtimeSession(application *application.Application, t Transport, model ItemID: generateItemID(), }) - session.respSink.issue(context.Background(), respcoord.SourceClient, func(ctx context.Context) { - commitUtterance(ctx, allAudio, session, conversation, t) + // Issue through the shared commit-ordering boundary (slot claim + + // issue under one lock, shared with the VAD commit path) so slot + // order == issue order across both producers (issue #12445). + session.issueCommit(context.Background(), respcoord.SourceClient, func(ctx context.Context, slot *commitSlot) { + commitUtterance(ctx, allAudio, session, conversation, t, slot) }) case types.InputAudioBufferClearEvent: @@ -1824,8 +1917,8 @@ func vadScanWindowSec(sv *types.RealtimeSessionSemanticVad, silenceThreshold flo return window } -func commitUtterance(ctx context.Context, utt []byte, session *Session, conv *Conversation, t Transport) { - commitUtteranceWithTranscript(ctx, utt, nil, nil, "", session, conv, t) +func commitUtterance(ctx context.Context, utt []byte, session *Session, conv *Conversation, t Transport, slot *commitSlot) { + commitUtteranceWithTranscript(ctx, utt, nil, nil, "", session, conv, t, slot) } // commitUtteranceWithTranscript commits one user turn. live carries the @@ -1837,7 +1930,33 @@ func commitUtterance(ctx context.Context, utt []byte, session *Session, conv *Co // is written to a temp WAV and transcribed via the file path as before. // itemID is the turn's conversation item id ("" mints a fresh one); it must // match the id any live deltas were sent under. -func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUtterance, gated *schema.TranscriptionResult, itemID string, session *Session, conv *Conversation, t Transport) { +func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUtterance, gated *schema.TranscriptionResult, itemID string, session *Session, conv *Conversation, t Transport, slot *commitSlot) { + // sctx is the session-lifetime context (see Session.sessionCtx): it + // survives a barge-in (which cancels only the per-response ctx) but is + // cancelled at teardown. Unit tests that build a bare Session leave it + // nil; fall back like responseSink.Perform does for a nil parent. + sctx := session.sessionCtx + if sctx == nil { + sctx = context.Background() + } + + // Release the slot on EVERY exit — including the empty-utt early return, + // transcription errors, gate rejections, empty transcripts and teardown — + // but ONLY after the predecessor has finished: a failed or skipped middle + // commit must not release the next commit before the earlier one appended + // its item, or the next response would be built on a history missing the + // earlier turn (issue #12445, review follow-up schedule 1). The session + // context can still stop the wait, so teardown never blocks on a + // never-finishing predecessor. + defer func() { + if slot.prevDone != nil { + select { + case <-slot.prevDone: + case <-sctx.Done(): + } + } + close(slot.done) + }() if len(utt) == 0 { return } @@ -1895,7 +2014,10 @@ func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUt resolveCh = make(chan resolveOutcome, 1) wavPath := f.Name() go func() { - r, rerr := session.voiceGate.Resolve(ctx, wavPath) + // Turn work like the transcription: must outlive a barge-in (the + // decision belongs to THIS utterance), so run under the session + // context (issue #12445). + r, rerr := session.voiceGate.Resolve(sctx, wavPath) resolveCh <- resolveOutcome{res: r, err: rerr} }() } @@ -1932,7 +2054,19 @@ func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUt // emitTranscription streams transcript deltas when // pipeline.streaming.transcription is set, otherwise emits a single // completed event; either way it returns the final transcript text. - transcript, err = emitTranscription(ctx, t, session, itemID, f.Name()) + // + // The transcription runs under the SESSION context, not the turn's + // response context: a barge-in (new speech onset) or a newer commit + // cancels that context while Whisper is still in flight, and + // cancelling the transcription would lose the user's input for this + // turn ("transcription_failed: context canceled", issue #12445). The + // session context is cancelled at teardown instead (conncoord + // Teardown, before respSink.shutdown joins this goroutine), so the + // transcription can never outlive the session. The response to this + // turn is still cancelled — see the ctx check below — only the + // transcript survives, so the next response is built on the full + // utterance. + transcript, err = emitTranscription(sctx, t, session, itemID, f.Name()) if err != nil { // Drain the gate goroutine before returning so its in-flight read of // the temp WAV finishes before the deferred os.Remove fires. @@ -2017,6 +2151,31 @@ func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUt // sound-detection-only session (no transcription) has no LLM stage, so it // stops here after emitting the sound-detection event. if session.InputAudioTranscription != nil && !session.TranscriptionOnly && strings.TrimSpace(transcript) != "" { + // Commit ordering (append gate): wait until the previous committed + // turn has appended its user item, so the conversation history — and + // the response built on it — sees the turns in speech order, not + // transcription-completion order (issue #12445, review schedule 2). + // The deferred slot release (release gate) enforces the same order on + // every other exit path. Aborts on teardown. + if slot.prevDone != nil { + select { + case <-slot.prevDone: + case <-sctx.Done(): + return + } + } + // If the turn's response context was cancelled while the (detached) + // transcription ran — barge-in, superseded by a newer commit — the + // user item still commits, so the LLM context keeps the full user + // input, but no response is generated for this (superseded) turn: + // the newer speech triggers its own response on the complete history. + // (Teardown is not a barge-in: the slot wait above already returned + // when the session context was cancelled.) + if ctx.Err() != nil { + xlog.Debug("skipping response: turn context cancelled during transcription (barge-in); committing user item only") + appendUserItem(conv, t, utt, transcript, speaker) + return + } generateResponse(ctx, session, utt, transcript, speaker, conv, t) } } @@ -2189,10 +2348,12 @@ func speakerNote(s *types.Speaker, noteUnknown bool) string { } // Function to generate a response based on the conversation -func generateResponse(ctx context.Context, session *Session, utt []byte, transcript string, speaker *types.Speaker, conv *Conversation, t Transport) { - xlog.Debug("Generating realtime response...") - - // Create user message item +// appendUserItem adds the committed user turn to the conversation and +// notifies the client, returning the item. Split out of generateResponse so a +// turn whose response was superseded (barge-in during transcription) can still +// commit its item — the LLM context must keep the full user input — without +// generating a response for it (issue #12445). +func appendUserItem(conv *Conversation, t Transport, utt []byte, transcript string, speaker *types.Speaker) types.MessageItemUnion { item := types.MessageItemUnion{ User: &types.MessageItemUser{ ID: generateItemID(), @@ -2214,6 +2375,16 @@ func generateResponse(ctx context.Context, session *Session, utt []byte, transcr sendEvent(t, types.ConversationItemAddedEvent{ Item: item, }) + return item +} + +// generateResponse creates the user message item for a committed turn and +// triggers the assistant response (LLM + TTS) for it. +func generateResponse(ctx context.Context, session *Session, utt []byte, transcript string, speaker *types.Speaker, conv *Conversation, t Transport) { + xlog.Debug("Generating realtime response...") + + // Create user message item + item := appendUserItem(conv, t, utt, transcript, speaker) // Surface the recognized speaker to the client. Skip the event for an // unidentified speaker unless announce_unknown is set. diff --git a/core/http/endpoints/openai/realtime_commit_order_test.go b/core/http/endpoints/openai/realtime_commit_order_test.go new file mode 100644 index 000000000..a5b024c18 --- /dev/null +++ b/core/http/endpoints/openai/realtime_commit_order_test.go @@ -0,0 +1,346 @@ +package openai + +import ( + "context" + "errors" + "sync/atomic" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "github.com/mudler/LocalAI/core/backend" + "github.com/mudler/LocalAI/core/config" + "github.com/mudler/LocalAI/core/http/endpoints/openai/respcoord" + "github.com/mudler/LocalAI/core/http/endpoints/openai/types" + "github.com/mudler/LocalAI/core/schema" +) + +// These specs are the regression tests for the schedules from the review of +// issue #12445 (PR fix/realtime-bargein-vad-commit-cancel), v1–v3: +// +// 1. teardown during an in-flight transcription: the transcription must be +// cancelled with the session — a detached context (WithoutCancel) would +// keep the backend job (and the respSink.shutdown join) blocked until the +// backend was explicitly released; +// 2. out-of-order completion: with the first transcription held and the +// second completed, the conversation must NOT become [second first] and +// the second response must not see only [second]; +// 3a. a FAILED middle commit (transcription error) must not release the +// third turn before the first one finishes — response history must not be +// [third], final history must not be [third first]; +// 3b. an EMPTY middle commit (empty transcript / gate rejection behave the +// same) must not release the third turn early either; +// 3c. the two producers (VAD CommitTurn, client input_audio_buffer.commit) +// must not be able to reserve a slot and issue it in opposite order — +// supersession may only cancel a response issued LATER in speech order. +// +// They drive the REAL issue path — Session.issueCommit into the real +// responseSink/respcoord (so coordinator supersession and the spawned +// response goroutines are exercised, not a bare +// commitUtteranceWithTranscript call) — with a transcription double whose +// wait is controlled by the spec and which honours context cancellation like +// a real backend job. + +// heldTranscribeModel scripts the commit-path transcription: call n blocks +// until release (or ctx cancellation) and then returns texts[n-1] (or +// errs[n-1], when set). started / finished announce entry/exit of call n so +// specs can sequence the schedules. +type heldTranscribeModel struct { + *fakeModel + n int32 + texts []string + errs []error + release chan struct{} + started chan int32 + finished chan int32 +} + +func (m *heldTranscribeModel) Transcribe(ctx context.Context, _ string, _ string, _ bool, _ bool, _ string) (*schema.TranscriptionResult, error) { + n := atomic.AddInt32(&m.n, 1) + select { + case m.started <- n: + default: + } + if n == 1 { + select { + case <-m.release: + case <-ctx.Done(): + return nil, ctx.Err() + } + } + if int(n) <= len(m.errs) && m.errs[n-1] != nil { + select { + case m.finished <- n: + default: + } + return nil, m.errs[n-1] + } + select { + case m.finished <- n: + default: + } + return &schema.TranscriptionResult{Text: m.texts[n-1]}, nil +} + +var _ = Describe("realtime commit order and transcription lifetime (issue #12445)", func() { + utt := []byte{1, 2, 3, 4} + + newModel := func(texts ...string) *heldTranscribeModel { + return &heldTranscribeModel{ + fakeModel: &fakeModel{ + cfg: &config.ModelConfig{}, + predictTokens: []string{"ok"}, + predictResp: backend.LLMResponse{Response: "ok"}, + ttsStreamChunks: [][]byte{{1}}, + ttsStreamRate: 24000, + }, + texts: texts, + release: make(chan struct{}), + started: make(chan int32, 8), + finished: make(chan int32, 8), + } + } + + newSession := func(m *heldTranscribeModel) *Session { + on := true + session := &Session{ + OutputSampleRate: 24000, + InputAudioTranscription: &types.AudioTranscription{}, + ModelInterface: m, + ModelConfig: &config.ModelConfig{ + Pipeline: config.Pipeline{Streaming: config.PipelineStreaming{LLM: &on, TTS: &on}}, + }, + respSink: newResponseSink(), + } + session.sessionCtx, session.sessionCancel = context.WithCancel(context.Background()) + return session + } + + // issue drives one commit through the REAL issue path: the shared + // ordering boundary (slot claim + respSink.issue under one lock) and the + // real respcoord supersession — NOT a direct + // commitUtteranceWithTranscript call. + issue := func(session *Session, parent context.Context, source respcoord.Source, conv *Conversation, tr *fakeTransport) { + session.issueCommit(parent, source, func(ctx context.Context, slot *commitSlot) { + commitUtteranceWithTranscript(ctx, utt, nil, nil, "", session, conv, tr, slot) + }) + } + + userTexts := func(conv *Conversation) []string { + var out []string + conv.Lock.Lock() + defer conv.Lock.Unlock() + for _, it := range conv.Items { + if it.User != nil && len(it.User.Content) > 0 { + out = append(out, it.User.Content[0].Transcript) + } + } + return out + } + + userMsgTexts := func(msgs schema.Messages) []string { + var out []string + for _, msg := range msgs { + if msg.Role == string(types.MessageRoleUser) { + out = append(out, msg.StringContent) + } + } + return out + } + + It("cancels an in-flight transcription at teardown instead of outliving the session", func() { + m := newModel("first") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + // The first transcription is in flight (held by the double). + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Teardown (conncoord): cancel the session context, then shut down the + // response sink — which JOINS the response goroutines. The join must + // return promptly: the transcription is cancelled with the session, not + // allowed to outlive it (a WithoutCancel context would block here until + // the backend was released). + done := make(chan struct{}) + go func() { + session.sessionCancel() + session.respSink.shutdown() + close(done) + }() + + Eventually(done, "2s").Should(BeClosed()) + Expect(userTexts(conv)).To(BeEmpty(), "a cancelled transcription commits no user item") + }) + + It("commits user items in speech order when the first transcription finishes last", func() { + m := newModel("first", "second") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // Commit 1 (issued first = speech order) — its transcription is held. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Commit 2 (issued second) supersedes commit 1's response, and its + // transcription finishes FIRST. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + // Now release commit 1; both commit bodies complete. + close(m.release) + session.respSink.wait() + + // The conversation is in speech order — NOT [second first]. + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + // The superseded turn got no response of its own; the surviving + // (second) response was built on the complete, ordered history. + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "second"})) + }) + + It("keeps the user item when a barge-in cancels the turn context during transcription", func() { + m := newModel("first", "second") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // The first turn is in flight (transcription held). + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Barge-in: the new speech onset cancels the in-flight turn's response + // context (turncoord: respSink.cancel(SourceVAD)) while its + // transcription is still running. + session.respSink.cancel(respcoord.SourceVAD) + + // The held transcription of the barged-in turn survives the barge-in + // (session context) and completes; the turn commits its item but gets + // no response of its own. + close(m.release) + session.respSink.wait() + Expect(userTexts(conv)).To(Equal([]string{"first"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(BeZero()) + + // The barge-in turn then commits and responds on the complete history. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + session.respSink.wait() + + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + // The response was built on the complete, ordered history. + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "second"})) + }) + + It("a failed middle commit does not release later turns before earlier ones finish", func() { + m := newModel("first", "second", "third") + m.errs = []error{nil, errors.New("stt backend down"), nil} + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // Turn A (issued first) — its transcription is held. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Turn B (issued second, supersedes A) — its transcription FAILS. + // Without the release gate, B would close its slot on the error return + // without waiting for A, releasing turn C early. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + // Turn C (issued third, supersedes B) — its transcription completes... + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(3)))) + // ...but C's item append must still wait for A (through B's slot), so + // no response has been created yet. + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(BeZero()) + + // Release A; everything drains in order. + close(m.release) + session.respSink.wait() + + // B contributes no item; C's response saw the earlier turn — NOT + // [third] — and the conversation is in speech order — NOT + // [third first]. + Expect(userTexts(conv)).To(Equal([]string{"first", "third"})) + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "third"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + }) + + It("an empty middle commit does not release later turns before earlier ones finish", func() { + m := newModel("first", "", "third") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // Turn A (issued first) — its transcription is held. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Turn B (issued second, supersedes A) — its transcription returns an + // EMPTY transcript: no item, no response, early return. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + // Turn C (issued third) — completes, but must still wait for A. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(3)))) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(BeZero()) + + close(m.release) + session.respSink.wait() + + Expect(userTexts(conv)).To(Equal([]string{"first", "third"})) + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "third"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + }) + + It("orders commits across the VAD and client producers (VAD first)", func() { + m := newModel("first", "second") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // VAD issues turn A (held). The client then issues turn B — through + // the same shared boundary, so B can only be issued AFTER A's slot is + // claimed and A's response issued; B supersedes A (the later issue + // wins), never the other way around. + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + issue(session, context.Background(), respcoord.SourceClient, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + close(m.release) + session.respSink.wait() + + // A's item survives (committed, no response); B's response is built on + // the complete, ordered history. + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "second"})) + }) + + It("orders commits across the VAD and client producers (client first)", func() { + m := newModel("first", "second") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + + // Client issues turn A (held); VAD then issues turn B — same shared + // boundary, opposite producer order. + issue(session, context.Background(), respcoord.SourceClient, conv, tr) + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + issue(session, context.Background(), respcoord.SourceVAD, conv, tr) + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + close(m.release) + session.respSink.wait() + + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + Expect(userMsgTexts(m.lastMessages)).To(Equal([]string{"first", "second"})) + }) +}) diff --git a/core/http/endpoints/openai/realtime_conncoord.go b/core/http/endpoints/openai/realtime_conncoord.go index 0dc6016bf..d94c6132e 100644 --- a/core/http/endpoints/openai/realtime_conncoord.go +++ b/core/http/endpoints/openai/realtime_conncoord.go @@ -102,12 +102,21 @@ func (s *connSink) Perform(e conncoord.Effect) { } s.wg.Wait() - // 2. Terminate the response coordinator (M3): cancel the in-flight response + // 2. Cancel the session-lifetime context FIRST: committed turns' + // transcriptions run under it (issue #12445), and respSink.shutdown + // below joins the response goroutines — a transcription that + // outlived the session would block the join (and the teardown) + // until the backend finished the job. + if s.session.sessionCancel != nil { + s.session.sessionCancel() + } + + // 3. Terminate the response coordinator (M3): cancel the in-flight response // and join all response goroutines (which also closes their TTS // pipelines, M5). After this no response can start. s.session.respSink.shutdown() - // 3. Terminate every conversation's compaction coordinator (M4): cancel + + // 4. Terminate every conversation's compaction coordinator (M4): cancel + // join any in-flight summarize+evict so it cannot outlive the session. for _, conv := range s.session.Conversations { if conv.compaction != nil { diff --git a/core/http/endpoints/openai/realtime_semantic_vad_test.go b/core/http/endpoints/openai/realtime_semantic_vad_test.go index 5a92e3b15..4621901b2 100644 --- a/core/http/endpoints/openai/realtime_semantic_vad_test.go +++ b/core/http/endpoints/openai/realtime_semantic_vad_test.go @@ -377,8 +377,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() { session := newTranscriptionOnlySession(m, true) tr := &fakeTransport{} - commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, nil, - &schema.TranscriptionResult{Text: "batch text", Eou: true}, "item_turn", session, &Conversation{}, tr) + commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, nil, &schema.TranscriptionResult{Text: "batch text", Eou: true}, "item_turn", session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) @@ -389,8 +388,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() { session := newTranscriptionOnlySession(m, true) tr := &fakeTransport{} - commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, - &liveUtterance{Text: "hello"}, nil, "item_turn", session, &Conversation{}, tr) + commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, &liveUtterance{Text: "hello"}, nil, "item_turn", session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) @@ -411,8 +409,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() { session := newTranscriptionOnlySession(m, false) tr := &fakeTransport{} - commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, - &liveUtterance{}, nil, "", session, &Conversation{}, tr) + commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, &liveUtterance{}, nil, "", session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) }) @@ -423,7 +420,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() { tr := &fakeTransport{} conv := &Conversation{} - commitUtterance(context.Background(), []byte{1, 2}, session, conv, tr) + commitUtterance(context.Background(), []byte{1, 2}, session, conv, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) Expect(conv.Items).To(BeEmpty()) diff --git a/core/http/endpoints/openai/realtime_sound_detection_test.go b/core/http/endpoints/openai/realtime_sound_detection_test.go index 94406e7cc..a67d19c1a 100644 --- a/core/http/endpoints/openai/realtime_sound_detection_test.go +++ b/core/http/endpoints/openai/realtime_sound_detection_test.go @@ -167,7 +167,7 @@ var _ = Describe("commitUtterance (sound-detection-only session)", func() { tr := &fakeTransport{} utt := make([]byte, 32) // non-empty PCM so commitUtterance proceeds - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemSoundDetection)).To(Equal(1)) // No transcription happened. diff --git a/core/http/endpoints/openai/realtime_turncoord.go b/core/http/endpoints/openai/realtime_turncoord.go index f0d599f7e..8faf12479 100644 --- a/core/http/endpoints/openai/realtime_turncoord.go +++ b/core/http/endpoints/openai/realtime_turncoord.go @@ -125,8 +125,11 @@ func (s *turnSink) Perform(e turncoord.Effect) { audio := s.commitAudio gated := s.commitGated conv := s.conv - s.session.respSink.issue(s.vadContext, respcoord.SourceVAD, func(ctx context.Context) { - commitUtteranceWithTranscript(ctx, audio, live, gated, itemID, s.session, conv, s.transport) + // Issue through the shared commit-ordering boundary (slot claim + + // issue under one lock, shared with the client commit path) so slot + // order == issue order across both producers (issue #12445). + s.session.issueCommit(s.vadContext, respcoord.SourceVAD, func(ctx context.Context, slot *commitSlot) { + commitUtteranceWithTranscript(ctx, audio, live, gated, itemID, s.session, conv, s.transport, slot) }) case turncoord.DiscardTurn: // No-op if the stream was never open (server_vad / already idle). diff --git a/core/http/endpoints/openai/realtime_voicegate_integration_test.go b/core/http/endpoints/openai/realtime_voicegate_integration_test.go index 4da774c76..07a5d20f9 100644 --- a/core/http/endpoints/openai/realtime_voicegate_integration_test.go +++ b/core/http/endpoints/openai/realtime_voicegate_integration_test.go @@ -83,7 +83,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) // The LLM/TTS pipeline ran to completion. @@ -98,7 +98,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) // Hard barrier: the LLM/TTS pipeline never ran. Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) @@ -114,7 +114,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) Expect(hasSpeakerNotAuthorized(tr)).To(BeTrue()) @@ -125,7 +125,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { config.VoiceGateWhenEvery, config.VoiceGateRejectSilent)) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) @@ -138,7 +138,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { // First utterance: authorized, marks the session verified. tr1 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr1) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr1, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) @@ -147,7 +147,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() { // Second utterance still proceeds because when:first skips re-verification. tr2 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr2) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr2, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr2)).To(BeFalse()) Expect(tr2.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) }) @@ -162,7 +162,7 @@ var _ = Describe("realtime speaker surfacing (commitUtterance)", func() { session.voiceGate.cfg.Identity = &config.VoiceIdentityConfig{Announce: true} tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1)) }) @@ -184,12 +184,12 @@ var _ = Describe("realtime speaker surfacing (commitUtterance)", func() { session, _ := itSession(gate) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(tr.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(0)) gate.cfg.Identity.AnnounceUnknown = true tr2 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr2) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr2, session.nextCommitSlot()) Expect(tr2.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1)) }) @@ -199,7 +199,7 @@ var _ = Describe("realtime speaker surfacing (commitUtterance)", func() { session.voiceGate.cfg.Enforce = boolPtr(false) tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) @@ -227,7 +227,7 @@ var _ = Describe("realtime speaker personalization (triggerResponseAtTurn)", fun session.Instructions = "You are helpful." tr := &fakeTransport{} - commitUtterance(context.Background(), utt, session, &Conversation{}, tr) + commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot()) user := findRole(m.lastMessages, "user") Expect(user).ToNot(BeNil()) @@ -257,12 +257,12 @@ var _ = Describe("realtime speaker personalization (triggerResponseAtTurn)", fun } s1, m1 := base() - commitUtterance(context.Background(), utt, s1, &Conversation{}, &fakeTransport{}) + commitUtterance(context.Background(), utt, s1, &Conversation{}, &fakeTransport{}, s1.nextCommitSlot()) Expect(findRole(m1.lastMessages, "system").StringContent).ToNot(ContainSubstring("unknown")) s2, m2 := base() s2.voiceGate.cfg.Identity.NoteUnknown = true - commitUtterance(context.Background(), utt, s2, &Conversation{}, &fakeTransport{}) + commitUtterance(context.Background(), utt, s2, &Conversation{}, &fakeTransport{}, s2.nextCommitSlot()) Expect(findRole(m2.lastMessages, "system").StringContent).To(ContainSubstring("The current speaker is unknown.")) }) }) @@ -305,7 +305,7 @@ var _ = Describe("realtime when:first with identity (commitUtterance)", func() { // Turn 1: authorized; identity resolved, speaker surfaced, response runs. tr1 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, conv, tr1) + commitUtterance(context.Background(), utt, session, conv, tr1, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(tr1.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1)) Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) @@ -315,7 +315,7 @@ var _ = Describe("realtime when:first with identity (commitUtterance)", func() { // enforces, the turn is dropped fail-closed rather than riding on the // cached first verification. tr2 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, conv, tr2) + commitUtterance(context.Background(), utt, session, conv, tr2, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr2)).To(BeTrue()) Expect(tr2.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) }) @@ -326,7 +326,7 @@ var _ = Describe("realtime when:first with identity (commitUtterance)", func() { conv := &Conversation{} tr1 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, conv, tr1) + commitUtterance(context.Background(), utt, session, conv, tr1, session.nextCommitSlot()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) @@ -334,7 +334,7 @@ var _ = Describe("realtime when:first with identity (commitUtterance)", func() { // speaker event still fires and the per-message name is set, proving the // per-turn re-resolution (not the cached first verification) drove it. tr2 := &fakeTransport{} - commitUtterance(context.Background(), utt, session, conv, tr2) + commitUtterance(context.Background(), utt, session, conv, tr2, session.nextCommitSlot()) Expect(tr2.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1)) var lastUser *schema.Message for i := range m.lastMessages { diff --git a/docs/design/realtime-state-machines.md b/docs/design/realtime-state-machines.md index 700929783..01a4229a6 100644 --- a/docs/design/realtime-state-machines.md +++ b/docs/design/realtime-state-machines.md @@ -458,6 +458,61 @@ property-test oracles, and FizzBee invariants: M5's by its existing `Closed`; the persistent coordinators (M3/M4) carry the explicit `Terminated` state. +- **Committed-turn pipeline: transcription lifetime + commit order (issue #12445, + done).** Two cross-cutting defects in the VAD commit path, neither of which a + single machine owned: + - *Transcription lifetime.* `commitUtteranceWithTranscript` ran the + utterance transcription under the turn's **response** context (M3). A + barge-in (new speech onset) or a superseding commit cancels that context + while Whisper is still in flight, so the in-flight transcription died + ("transcription_failed: context canceled") and the user's input was lost + from the conversation — the next response answered the second half of a + two-part utterance. Detaching with `context.WithoutCancel` fixed the + barge-in but broke teardown: the transcription then outlived the session, + and `respSink.shutdown` (which joins the response goroutines) blocked until + the backend finished the job. Fix: the transcription (and the voice-gate + resolution) now run under a **session-lifetime context** + (`Session.sessionCtx`) — cancelled at teardown by `conncoord`'s `Teardown` + *before* `respSink.shutdown()` joins, untouched by barge-in. The + turn's response context still cancels the *response*: when it was + cancelled during the (now detached) transcription, the user item is + committed (`appendUserItem`, split out of `generateResponse`) but no + response is generated for the superseded turn — the newer speech triggers + its own response on the complete history. + - *Commit order.* Consecutive commits run in parallel goroutines (M3 spawns + one per `issue`), so a fast second transcription could append its user + item before a slow first one — the conversation became + `[second, first]` and the second response saw only `[second]`. Two more + schedules broke the naive fix: a FAILED middle commit (error, empty + transcript, gate rejection) closed its slot without waiting for the + earlier one, releasing the third turn before the first appended + (`[third first]`); and the two producers (VAD `CommitTurn`, client + `input_audio_buffer.commit`) reserved the slot and called `respSink.issue` + separately, so a pause between the two let the other producer reserve AND + issue first — the later issue then superseded the EARLIER turn's response + (response history `[first]`, final history `[first second]` with the + second turn un-answered). Fix: a per-session **commit slot chain** with + two gates, and a **shared issue boundary**: + - `Session.issueCommit` claims the next slot and issues the commit body + under ONE lock (`commitOrderMu`), so slot order == issue order across + both producers; `respSink.issue` is non-blocking, so the lock never + stalls VAD/barge-in handling. + - *Append gate:* a commit's user-item append waits on the previous slot's + `done` (aborts on the session context). + - *Release gate:* the slot releases (`done` closes) only after the + predecessor has finished — on EVERY exit path, including errors, empty + transcripts and teardown (session context can still stop the wait), so + a failed or skipped middle commit never releases the next turn early. + Transcriptions stay parallel; only the item appends (and the slot + releases) are ordered. + Regression tests: `realtime_commit_order_test.go` (teardown during an + in-flight transcription; held-first/finished-second out-of-order completion; + barge-in-during-transcription item survival; a failed AND an empty middle + commit; interleaved VAD/client producers), driving the REAL issue path — + `issueCommit` + `responseSink`/`respcoord` supersession — with a + transcription double that honours context cancellation. Verified: builds, + openai specs under `-race`. + ## Part 5 — Library vs hand-rolled (Go ecosystem, verified 2026-06) Researched against live GitHub/pkg.go.dev data. **Verdict: hand-roll a typed transition