mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-10 07:47:29 -04:00
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:
1 parent
372b1f8983
commit
c8f991be14
8 files changed
+421
-38
No files matched your search
@@ -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 {
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in new issue
Block a user