fix(realtime): detach VAD-commit transcription from barge-in cancellation (#12446)

* 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>

* fix(realtime): release commit slots in order on every exit; share slot+issue boundary

Follow-up to the review of c8f991be (issue #12445): two schedules still
broke commit ordering.

1. A failed middle commit released later turns before earlier turns
   finished. slot.done closed on every return, but the wait on
   slot.prevDone happened only on the successful nonempty-transcript
   path. Hold transcription A, let B fail (or return an empty
   transcript / be rejected by the voice gate), then complete C: B
   closed its channel without waiting for A, so C appended and started
   its response without A (response history [third], final history
   [third first]). Fix: the slot now releases (done closes) only AFTER
   the predecessor has finished — on EVERY exit path, including errors,
   empty transcripts, gate rejections and teardown (the session context
   can still stop the wait, so teardown never blocks on a
   never-finishing predecessor). The success path keeps its append gate
   (wait before appending the user item); the deferred release gate
   enforces the same order on every other exit.

2. Slot order and response issue order could disagree between the two
   producers. The VAD CommitTurn and the client
   input_audio_buffer.commit reserved the slot and called
   respSink.issue separately; a pause between the two let the other
   producer reserve AND issue first, so the later issue superseded the
   EARLIER turn's response (response history [first], final history
   [first second], second turn un-answered). Fix: both producers now go
   through Session.issueCommit, which claims the slot and issues the
   body under one lock (commitOrderMu) — slot order == issue order.
   respSink.issue is non-blocking, so the lock never stalls
   VAD/barge-in handling.

Regression tests (realtime_commit_order_test.go) now drive the REAL
issue path — Session.issueCommit into the real responseSink/respcoord,
so coordinator supersession and the spawned response goroutines are
exercised — and cover: teardown during an in-flight transcription;
held-first/finished-second out-of-order completion; barge-in
(respSink.cancel) item survival; a FAILED middle commit; an EMPTY
middle commit; interleaved VAD/client producers in both directions.
The failed/empty middle specs fail deterministically without the
release gate (verified against the pre-fix code).

docs/design/realtime-state-machines.md: implementation-status entry
updated (append gate + release gate + shared issue boundary).

Validated: builds, go vet clean, all openai specs + respcoord/turncoord/
conncoord suites pass under -race (476 specs, incl. the 4 new ones).

Fixes #12445

Signed-off-by: nexxtmobile.de <kai@nexxtmobile.de>

---------

Signed-off-by: nexxtmobile.de <kai@nexxtmobile.de>
Co-authored-by: nexxtmobile.de <kai@nexxtmobile.de>
This commit is contained in:
nexxtmobile.deandnexxtmobile.de authored and GitHub committed 2026-10-05 01:27:40 +02:00
1 parent d66383c0e2
commit 4d681b7f6d
8 files changed
+621 -40

No files matched your search

+182 -11
View File
@@ -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.
@@ -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"}))
})
})
@@ -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 {
@@ -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())
@@ -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.
@@ -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).
@@ -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 {
+55
View File
@@ -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