diff --git a/core/http/endpoints/openai/realtime.go b/core/http/endpoints/openai/realtime.go index c152d2c22..178756060 100644 --- a/core/http/endpoints/openai/realtime.go +++ b/core/http/endpoints/openai/realtime.go @@ -203,6 +203,56 @@ 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. + 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: without the gate a fast second transcription +// would commit ahead of a slow first one, and the assistant response — and +// the conversation history — would see the user input out of order or missing +// (issue #12445, review schedule 2). +type commitSlot struct { + // prevDone is the previous slot's done channel (nil for the first commit); + // this commit's item append waits on it. + prevDone chan struct{} + // done is closed when this commit has appended its user item — or aborted + // (error, teardown) — so the next commit is never blocked. + 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. +func (session *Session) nextCommitSlot() *commitSlot { + session.commitOrderMu.Lock() + var prevDone chan struct{} + if session.commitTail != nil { + prevDone = session.commitTail.done + } + slot := &commitSlot{prevDone: prevDone, done: make(chan struct{})} + session.commitTail = slot + session.commitOrderMu.Unlock() + return slot } func (s *Session) installVoiceBinding(voice string, params map[string]string, release func()) { @@ -614,6 +664,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 +1001,11 @@ func runRealtimeSession(application *application.Application, t Transport, model ItemID: generateItemID(), }) + // Claim the commit order before the (parallel) transcription starts, + // so user items commit in speech order (issue #12445). + slot := session.nextCommitSlot() session.respSink.issue(context.Background(), respcoord.SourceClient, func(ctx context.Context) { - commitUtterance(ctx, allAudio, session, conversation, t) + commitUtterance(ctx, allAudio, session, conversation, t, slot) }) case types.InputAudioBufferClearEvent: @@ -1824,8 +1882,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,11 +1895,24 @@ 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) { + // Close the slot on EVERY exit — before the empty-utt early return too — + // so the next commit's ordering wait is never blocked, even when this + // turn errors out or the session is torn down (issue #12445). + defer close(slot.done) if len(utt) == 0 { return } + // 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() + } + f, err := os.CreateTemp("", "realtime-audio-chunk-*.wav") if err != nil { xlog.Error("failed to create temp file", "error", err) @@ -1895,7 +1966,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 +2006,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 +2103,29 @@ 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: 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). 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 +2298,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 +2325,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..f66db2e97 --- /dev/null +++ b/core/http/endpoints/openai/realtime_commit_order_test.go @@ -0,0 +1,216 @@ +package openai + +import ( + "context" + "sync/atomic" + "sync" + + . "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/types" + "github.com/mudler/LocalAI/core/schema" +) + +// These specs are the regression tests for the two schedules from the review +// of issue #12445 (PR fix/realtime-bargein-vad-commit-cancel): +// +// 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]. +// +// They drive the REAL commit path (commitUtteranceWithTranscript + +// responseSink/respcoord) 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]. started / +// finished announce entry/exit of call n so specs can sequence the schedules. +type heldTranscribeModel struct { + *fakeModel + n int32 + texts []string + 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() + } + } + 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, 4), + finished: make(chan int32, 4), + } + } + + 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}}, + }, + } + session.sessionCtx, session.sessionCancel = context.WithCancel(context.Background()) + return session + } + + 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{} + + done := make(chan struct{}) + go func() { + commitUtteranceWithTranscript(context.Background(), utt, nil, nil, "", session, conv, tr, session.nextCommitSlot()) + close(done) + }() + + // The first transcription is in flight (held by the double). + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Teardown cancels the session context (conncoord does this BEFORE + // respSink.shutdown joins the response goroutines). The commit must + // return promptly — not wait for the backend to finish the job. + session.sessionCancel() + + Eventually(done, "2s").Should(BeClosed()) + Expect(userTexts(conv)).To(BeEmpty(), "a cancelled transcription commits no user item") + }) + + It("commits user items in speech order even when transcriptions finish out of order", func() { + m := newModel("first", "second") + session := newSession(m) + tr := &fakeTransport{} + conv := &Conversation{} + var wg sync.WaitGroup + + // Commit 1 (issued first = speech order) — its transcription is held. + wg.Add(1) + go func() { + defer wg.Done() + commitUtteranceWithTranscript(context.Background(), utt, nil, nil, "", session, conv, tr, session.nextCommitSlot()) + }() + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Commit 2 (issued second) — its transcription finishes FIRST. + wg.Add(1) + go func() { + defer wg.Done() + commitUtteranceWithTranscript(context.Background(), utt, nil, nil, "", session, conv, tr, session.nextCommitSlot()) + }() + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + // Now release commit 1; both commits complete. + close(m.release) + wg.Wait() + + // The conversation is in speech order — NOT [second first]. + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + // The second response's LLM context saw both user turns, in order. + 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{} + var wg sync.WaitGroup + + ctx1, cancel1 := context.WithCancel(context.Background()) + wg.Add(1) + go func() { + defer wg.Done() + commitUtteranceWithTranscript(ctx1, utt, nil, nil, "", session, conv, tr, session.nextCommitSlot()) + }() + Eventually(m.started, "2s").Should(Receive(Equal(int32(1)))) + + // Barge-in: the new speech onset cancels the in-flight turn's response + // context while its transcription is still running. + cancel1() + + wg.Add(1) + go func() { + defer wg.Done() + commitUtteranceWithTranscript(context.Background(), utt, nil, nil, "", session, conv, tr, session.nextCommitSlot()) + }() + Eventually(m.finished, "2s").Should(Receive(Equal(int32(2)))) + + // The held transcription of the barged-in turn survives the barge-in + // (session context) and completes. + close(m.release) + wg.Wait() + + // The barged-in turn's user item is committed — the LLM context keeps + // the full user input — but it gets NO response of its own. + Expect(userTexts(conv)).To(Equal([]string{"first", "second"})) + Expect(tr.countEvents(types.ServerEventTypeResponseCreated)).To(Equal(1)) + // The response (for the barge-in turn) was built on the complete, + // ordered history. + 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..8e24a54a7 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 + // Claim the commit order before the (parallel) transcription starts, + // so user items commit in speech order (issue #12445). + slot := s.session.nextCommitSlot() s.session.respSink.issue(s.vadContext, respcoord.SourceVAD, func(ctx context.Context) { - commitUtteranceWithTranscript(ctx, audio, live, gated, itemID, s.session, conv, s.transport) + 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..cb81f1b69 100644 --- a/docs/design/realtime-state-machines.md +++ b/docs/design/realtime-state-machines.md @@ -458,6 +458,43 @@ 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]`. Fix: a + per-session **commit slot chain** (`Session.nextCommitSlot`, claimed at + commit *issue* time so slot order == speech order). A commit's user-item + append waits on the previous slot's `done` (aborts on the session + context); every exit closes its own slot, so a failed or torn-down commit + never blocks the next. Transcriptions stay parallel; only the item appends + are ordered. + Regression tests: `realtime_commit_order_test.go` (both review schedules — + teardown during an in-flight transcription; held-first/finished-second + out-of-order completion — plus the barge-in-during-transcription item + survival), driving the real commit path 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