fix(realtime): keep VAD-commit transcription alive across barge-in, cancel it at teardown, order commits

Barge-in (new speech onset) cancels the turn's SourceVAD response context
(realtime_turncoord.go: respSink.cancel(SourceVAD)). The VAD commit body
runs under that same context, so an in-flight Whisper STT call was aborted
with 'context canceled' whenever the caller kept talking while the first
chunk was transcribing. The user's turn was lost: no transcript, no
LLM/TTS response.

v1 of this fix ran the transcription with context.WithoutCancel(ctx). The
review correctly pointed out two correctness gaps:

1. Teardown lost its cancellation. WithoutCancel detaches from every
   cancellation, so a transcription in flight at session close outlived
   the session and blocked respSink.shutdown (which joins the response
   goroutines) until the backend finished the job.
2. Out-of-order commits. Consecutive commits run in parallel goroutines,
   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].

Changes (core/http/endpoints/openai/):
- Session gains a session-lifetime context (sessionCtx), cancelled by
  conncoord's Teardown BEFORE respSink.shutdown joins the response
  goroutines. The transcription (and the voice-gate resolution) run under
  it: they survive barge-in (which cancels only the per-response context)
  but are cancelled with the session.
- Commit slots order the user-item appends in speech order:
  Session.nextCommitSlot() is claimed at commit issue time (VAD CommitTurn
  / client commit), a commit's item append waits on the previous slot's
  done (aborts on the session context), and every exit closes the slot so
  a failed or torn-down commit never blocks the next. Transcriptions stay
  parallel; only the appends are ordered.
- 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 (appendUserItem, split out of generateResponse) so
  the LLM context keeps the full user input, but no response is generated
  for the superseded turn; the newer speech triggers its own response on
  the complete history.
- Regression tests (realtime_commit_order_test.go) cover both review
  schedules — teardown during an in-flight transcription, and
  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.
- docs/design/realtime-state-machines.md: implementation-status entry for
  the committed-turn pipeline (transcription lifetime + commit order).

Fixes #12445

Validated: builds, go vet clean, all openai specs + respcoord/turncoord/
conncoord suites pass under -race (incl. the 3 new regression specs).
Production A/B (call-center voice agent, SIP, silero-vad +
whisper-large-turbo + LLM + TTS, server_vad ~600 ms) on LocalAI v4.11.0:
unpatched — 'transcription_failed: context canceled', first part of the
utterance lost, agent answers only the remainder; patched — full
transcript committed, agent answers the complete utterance, barge-in
still cancels the in-flight assistant TTS response as intended, and
teardown cancels the in-flight transcription instead of waiting for the
backend.

Signed-off-by: nexxtmobile.de <kai@nexxtmobile.de>
This commit is contained in:
nexxtmobile.de committed 2026-10-03 16:41:11 +00:00
1 parent 372b1f8983
commit c8f991be14
8 files changed
+421 -38

No files matched your search

+131 -10
View File
@@ -203,6 +203,56 @@ type Session struct {
// decision is serialized through respcoord.Coordinator, guaranteeing at most // decision is serialized through respcoord.Coordinator, guaranteeing at most
// one live response. See realtime_respcoord.go. // one live response. See realtime_respcoord.go.
respSink *responseSink 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()) { 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). // into two overlapping responses (see realtime_respcoord.go).
session.respSink = newResponseSink() 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 // Create a default conversation
conversationID := generateConversationID() conversationID := generateConversationID()
conversation := &Conversation{ conversation := &Conversation{
@@ -946,8 +1001,11 @@ func runRealtimeSession(application *application.Application, t Transport, model
ItemID: generateItemID(), 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) { 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: case types.InputAudioBufferClearEvent:
@@ -1824,8 +1882,8 @@ func vadScanWindowSec(sv *types.RealtimeSessionSemanticVad, silenceThreshold flo
return window return window
} }
func commitUtterance(ctx context.Context, utt []byte, session *Session, conv *Conversation, t Transport) { func commitUtterance(ctx context.Context, utt []byte, session *Session, conv *Conversation, t Transport, slot *commitSlot) {
commitUtteranceWithTranscript(ctx, utt, nil, nil, "", session, conv, t) commitUtteranceWithTranscript(ctx, utt, nil, nil, "", session, conv, t, slot)
} }
// commitUtteranceWithTranscript commits one user turn. live carries the // 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. // 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 // itemID is the turn's conversation item id ("" mints a fresh one); it must
// match the id any live deltas were sent under. // 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 { if len(utt) == 0 {
return 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") f, err := os.CreateTemp("", "realtime-audio-chunk-*.wav")
if err != nil { if err != nil {
xlog.Error("failed to create temp file", "error", err) 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) resolveCh = make(chan resolveOutcome, 1)
wavPath := f.Name() wavPath := f.Name()
go func() { 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} resolveCh <- resolveOutcome{res: r, err: rerr}
}() }()
} }
@@ -1932,7 +2006,19 @@ func commitUtteranceWithTranscript(ctx context.Context, utt []byte, live *liveUt
// emitTranscription streams transcript deltas when // emitTranscription streams transcript deltas when
// pipeline.streaming.transcription is set, otherwise emits a single // pipeline.streaming.transcription is set, otherwise emits a single
// completed event; either way it returns the final transcript text. // 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 { if err != nil {
// Drain the gate goroutine before returning so its in-flight read of // Drain the gate goroutine before returning so its in-flight read of
// the temp WAV finishes before the deferred os.Remove fires. // 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 // sound-detection-only session (no transcription) has no LLM stage, so it
// stops here after emitting the sound-detection event. // stops here after emitting the sound-detection event.
if session.InputAudioTranscription != nil && !session.TranscriptionOnly && strings.TrimSpace(transcript) != "" { 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) 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 // 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) { // appendUserItem adds the committed user turn to the conversation and
xlog.Debug("Generating realtime response...") // notifies the client, returning the item. Split out of generateResponse so a
// turn whose response was superseded (barge-in during transcription) can still
// Create user message item // 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{ item := types.MessageItemUnion{
User: &types.MessageItemUser{ User: &types.MessageItemUser{
ID: generateItemID(), ID: generateItemID(),
@@ -2214,6 +2325,16 @@ func generateResponse(ctx context.Context, session *Session, utt []byte, transcr
sendEvent(t, types.ConversationItemAddedEvent{ sendEvent(t, types.ConversationItemAddedEvent{
Item: item, 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 // Surface the recognized speaker to the client. Skip the event for an
// unidentified speaker unless announce_unknown is set. // unidentified speaker unless announce_unknown is set.
@@ -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"}))
})
})
@@ -102,12 +102,21 @@ func (s *connSink) Perform(e conncoord.Effect) {
} }
s.wg.Wait() 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 // and join all response goroutines (which also closes their TTS
// pipelines, M5). After this no response can start. // pipelines, M5). After this no response can start.
s.session.respSink.shutdown() 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. // join any in-flight summarize+evict so it cannot outlive the session.
for _, conv := range s.session.Conversations { for _, conv := range s.session.Conversations {
if conv.compaction != nil { if conv.compaction != nil {
@@ -377,8 +377,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() {
session := newTranscriptionOnlySession(m, true) session := newTranscriptionOnlySession(m, true)
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, nil, commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, nil, &schema.TranscriptionResult{Text: "batch text", Eou: true}, "item_turn", session, &Conversation{}, tr, session.nextCommitSlot())
&schema.TranscriptionResult{Text: "batch text", Eou: true}, "item_turn", session, &Conversation{}, tr)
Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0))
Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1))
@@ -389,8 +388,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() {
session := newTranscriptionOnlySession(m, true) session := newTranscriptionOnlySession(m, true)
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, &liveUtterance{Text: "hello"}, nil, "item_turn", session, &Conversation{}, tr, session.nextCommitSlot())
&liveUtterance{Text: "hello"}, nil, "item_turn", session, &Conversation{}, tr)
Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionDelta)).To(Equal(0))
Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1))
@@ -411,8 +409,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() {
session := newTranscriptionOnlySession(m, false) session := newTranscriptionOnlySession(m, false)
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, commitUtteranceWithTranscript(context.Background(), []byte{1, 2}, &liveUtterance{}, nil, "", session, &Conversation{}, tr, session.nextCommitSlot())
&liveUtterance{}, nil, "", session, &Conversation{}, tr)
Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1)) Expect(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1))
}) })
@@ -423,7 +420,7 @@ var _ = Describe("commitUtteranceWithTranscript", func() {
tr := &fakeTransport{} tr := &fakeTransport{}
conv := &Conversation{} 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(tr.countEvents(types.ServerEventTypeConversationItemInputAudioTranscriptionCompleted)).To(Equal(1))
Expect(conv.Items).To(BeEmpty()) Expect(conv.Items).To(BeEmpty())
@@ -167,7 +167,7 @@ var _ = Describe("commitUtterance (sound-detection-only session)", func() {
tr := &fakeTransport{} tr := &fakeTransport{}
utt := make([]byte, 32) // non-empty PCM so commitUtterance proceeds 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)) Expect(tr.countEvents(types.ServerEventTypeConversationItemSoundDetection)).To(Equal(1))
// No transcription happened. // No transcription happened.
@@ -125,8 +125,11 @@ func (s *turnSink) Perform(e turncoord.Effect) {
audio := s.commitAudio audio := s.commitAudio
gated := s.commitGated gated := s.commitGated
conv := s.conv 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) { 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: case turncoord.DiscardTurn:
// No-op if the stream was never open (server_vad / already idle). // No-op if the stream was never open (server_vad / already idle).
@@ -83,7 +83,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() {
config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) config.VoiceGateWhenEvery, config.VoiceGateRejectEvent))
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtterance(context.Background(), utt, session, &Conversation{}, tr) commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse())
// The LLM/TTS pipeline ran to completion. // The LLM/TTS pipeline ran to completion.
@@ -98,7 +98,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() {
config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) config.VoiceGateWhenEvery, config.VoiceGateRejectEvent))
tr := &fakeTransport{} 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. // Hard barrier: the LLM/TTS pipeline never ran.
Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0))
@@ -114,7 +114,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() {
config.VoiceGateWhenEvery, config.VoiceGateRejectEvent)) config.VoiceGateWhenEvery, config.VoiceGateRejectEvent))
tr := &fakeTransport{} 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(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0))
Expect(hasSpeakerNotAuthorized(tr)).To(BeTrue()) Expect(hasSpeakerNotAuthorized(tr)).To(BeTrue())
@@ -125,7 +125,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() {
config.VoiceGateWhenEvery, config.VoiceGateRejectSilent)) config.VoiceGateWhenEvery, config.VoiceGateRejectSilent))
tr := &fakeTransport{} 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(tr.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0))
Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse())
@@ -138,7 +138,7 @@ var _ = Describe("realtime voice gate integration (commitUtterance)", func() {
// First utterance: authorized, marks the session verified. // First utterance: authorized, marks the session verified.
tr1 := &fakeTransport{} tr1 := &fakeTransport{}
commitUtterance(context.Background(), utt, session, &Conversation{}, tr1) commitUtterance(context.Background(), utt, session, &Conversation{}, tr1, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse())
Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) 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. // Second utterance still proceeds because when:first skips re-verification.
tr2 := &fakeTransport{} tr2 := &fakeTransport{}
commitUtterance(context.Background(), utt, session, &Conversation{}, tr2) commitUtterance(context.Background(), utt, session, &Conversation{}, tr2, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr2)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr2)).To(BeFalse())
Expect(tr2.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) 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} session.voiceGate.cfg.Identity = &config.VoiceIdentityConfig{Announce: true}
tr := &fakeTransport{} 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)) Expect(tr.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1))
}) })
@@ -184,12 +184,12 @@ var _ = Describe("realtime speaker surfacing (commitUtterance)", func() {
session, _ := itSession(gate) session, _ := itSession(gate)
tr := &fakeTransport{} 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)) Expect(tr.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(0))
gate.cfg.Identity.AnnounceUnknown = true gate.cfg.Identity.AnnounceUnknown = true
tr2 := &fakeTransport{} 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)) 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) session.voiceGate.cfg.Enforce = boolPtr(false)
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtterance(context.Background(), utt, session, &Conversation{}, tr) commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr)).To(BeFalse())
Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) Expect(tr.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1))
@@ -227,7 +227,7 @@ var _ = Describe("realtime speaker personalization (triggerResponseAtTurn)", fun
session.Instructions = "You are helpful." session.Instructions = "You are helpful."
tr := &fakeTransport{} tr := &fakeTransport{}
commitUtterance(context.Background(), utt, session, &Conversation{}, tr) commitUtterance(context.Background(), utt, session, &Conversation{}, tr, session.nextCommitSlot())
user := findRole(m.lastMessages, "user") user := findRole(m.lastMessages, "user")
Expect(user).ToNot(BeNil()) Expect(user).ToNot(BeNil())
@@ -257,12 +257,12 @@ var _ = Describe("realtime speaker personalization (triggerResponseAtTurn)", fun
} }
s1, m1 := base() 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")) Expect(findRole(m1.lastMessages, "system").StringContent).ToNot(ContainSubstring("unknown"))
s2, m2 := base() s2, m2 := base()
s2.voiceGate.cfg.Identity.NoteUnknown = true 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.")) 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. // Turn 1: authorized; identity resolved, speaker surfaced, response runs.
tr1 := &fakeTransport{} tr1 := &fakeTransport{}
commitUtterance(context.Background(), utt, session, conv, tr1) commitUtterance(context.Background(), utt, session, conv, tr1, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse())
Expect(tr1.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1)) Expect(tr1.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1))
Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 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 // enforces, the turn is dropped fail-closed rather than riding on the
// cached first verification. // cached first verification.
tr2 := &fakeTransport{} tr2 := &fakeTransport{}
commitUtterance(context.Background(), utt, session, conv, tr2) commitUtterance(context.Background(), utt, session, conv, tr2, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr2)).To(BeTrue()) Expect(hasSpeakerNotAuthorized(tr2)).To(BeTrue())
Expect(tr2.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0)) Expect(tr2.countEvents(types.ServerEventTypeResponseDone)).To(Equal(0))
}) })
@@ -326,7 +326,7 @@ var _ = Describe("realtime when:first with identity (commitUtterance)", func() {
conv := &Conversation{} conv := &Conversation{}
tr1 := &fakeTransport{} tr1 := &fakeTransport{}
commitUtterance(context.Background(), utt, session, conv, tr1) commitUtterance(context.Background(), utt, session, conv, tr1, session.nextCommitSlot())
Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse()) Expect(hasSpeakerNotAuthorized(tr1)).To(BeFalse())
Expect(tr1.countEvents(types.ServerEventTypeResponseDone)).To(BeNumerically(">=", 1)) 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 // speaker event still fires and the per-message name is set, proving the
// per-turn re-resolution (not the cached first verification) drove it. // per-turn re-resolution (not the cached first verification) drove it.
tr2 := &fakeTransport{} 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)) Expect(tr2.countEvents(types.ServerEventTypeConversationItemSpeaker)).To(Equal(1))
var lastUser *schema.Message var lastUser *schema.Message
for i := range m.lastMessages { for i := range m.lastMessages {
+37
View File
@@ -458,6 +458,43 @@ property-test oracles, and FizzBee invariants:
M5's by its existing `Closed`; the persistent coordinators (M3/M4) carry the M5's by its existing `Closed`; the persistent coordinators (M3/M4) carry the
explicit `Terminated` state. 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) ## 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 Researched against live GitHub/pkg.go.dev data. **Verdict: hand-roll a typed transition