Files
LocalAI/core/backend/transcript_live.go
T
mudler-agentandEttore Di Giacinto 9eb5a9e61d feat(audio): remember speakers from diarization (#12414)
* feat(schema): validate portable speaker profiles

Add the versioned profile schema for explicit speaker enrollment.
Validate compatibility against separately supplied loaded-encoder metadata.
Reject unusable speakers, invalid vectors, and inconsistent clean spans.

This slice does not change HTTP routes, backend integration, or the UI.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(parakeet): export profiles with transcripts

Export opt-in speaker profiles and trusted encoder metadata.
Replay registrations by ID so duplicate display names keep independent
vectors.

Use one profile-capable diarization for slots, names, and clean spans.
Assign timestamped ASR words to those slots without a second diarization.
Preserve legacy opt-out and no-ASR behavior, and propagate failures.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(audio): enroll portable speaker profiles

Gate profile exports with voice-recognition permission and validate
registration against metadata from the loaded encoder. Preserve audio
enrollment and independent registrations with duplicate display names.

Exclude diarization and registration exchanges before API trace capture
so persisted traces cannot retain profile vectors or JSON audio.

Defer candidate dimensions to trusted loaded metadata. Sort candidates
by registration ID so incompatible profiles cannot suppress legacy voices
through registry iteration order. Keep portable identity checks closed
when trusted metadata is unavailable.

Test persisted traces, explicit slot zero, and selection through offline
and live transport. Document privacy and the ephemeral registry lifecycle.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(ui): remember speakers from diarization

Add a Studio page for diarization and opt-in speaker profiles. Preview
clean intervals from the original recording before explicit registration.

Join profiles by raw speaker labels, preserve duplicate names, and relabel
turns only after a successful save. Discard stale results when the model
or recording changes. Share registration metadata with voice management
without storing vectors or recordings from this flow.

Document permissions and the global, ephemeral registry. Cover enrollment,
permissions, previews, and asynchronous races with mocked Playwright tests.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: clarify HTTP speaker enrollment support

Replace the stale enrollment limitation with the current HTTP workflow.
Distinguish native transport from explicit registration and link its docs.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* chore(parakeet): pin merged speaker profile support

Use the merged commit from mudler/parakeet.cpp#80.
Its tree matches the previously accepted native pin.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: add diarization enrollment setup example

Connect the existing gallery modes to the speaker enrollment workflow.
Show installation, private profile export, explicit raw-slot registration,
and later recognition without another export.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs(blog): explain diarization speaker profiles

Put the diarization walkthrough on the LocalAI website in the feature PR.
Cover the three gallery modes, explicit enrollment, and privacy limits.
Link setup instructions and keep availability conditional on feature support.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs(blog): focus diarization on everyday use

Explain what users can do with recordings before the setup steps.
Replace the technical walkthrough with a short Studio guide and link
readers to the existing reference for model names and developer use.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs(blog): lead with speaker capabilities

Present speaker recognition through everyday uses and a short UI flow.
Keep technical reference details in the existing documentation.

Assisted-by: OpenAI:unknown
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(diarization): satisfy Go lint checks

Avoid copying protobuf message state when extending backend status, check the multipart reader close result, and document the focused testing.T lint exemptions.

Assisted-by: nib:gpt-5.6-sol

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

---------

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
2026-10-02 08:14:00 +02:00

380 lines
11 KiB
Go

package backend
import (
"context"
"errors"
"fmt"
"io"
"maps"
"sync"
"time"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/schema"
"github.com/mudler/LocalAI/core/services/voicerecognition"
"github.com/mudler/LocalAI/core/trace"
grpcPkg "github.com/mudler/LocalAI/pkg/grpc"
"github.com/mudler/LocalAI/pkg/grpc/proto"
"github.com/mudler/LocalAI/pkg/model"
"github.com/mudler/LocalAI/pkg/sound"
"github.com/mudler/xlog"
)
// LiveTranscriptionEvent is one streamed event from a live (bidirectional)
// transcription session. Delta/Eou/Eob/Words arrive as the user speaks; Final
// is set exactly once, on the terminal event after Close flushes the decode
// tail. Eou means the model judged the user yielded the turn; Eob means a
// backchannel ("uh-huh") ended — callers must NOT treat Eob as a turn
// boundary.
type LiveTranscriptionEvent struct {
Delta string
Eou bool
Eob bool
Words []schema.TranscriptionWord
Speakers []LiveSpeakerSegment
Sounds []LiveSoundEvent
Final *schema.TranscriptionResult
}
// LiveSpeakerSegment is one closed speaker segment from a companion
// diarization/scene stream running alongside live transcription. Start/End
// are stream-relative seconds (mapped from the backend's nanoseconds).
type LiveSpeakerSegment struct {
Speaker string
// Name is the registered speaker name the backend matched, empty when the
// speaker is unknown.
Name string
Start float64
End float64
}
// LiveSoundEvent is one closed sound event from a companion sound/scene
// stream running alongside live transcription. Start/End are stream-relative
// seconds (mapped from the backend's nanoseconds).
type LiveSoundEvent struct {
Label string
Index int
Peak float32
Start float64
End float64
}
// LiveTranscriptionSession is a handle on an open live transcription stream.
// Feed pushes 16 kHz mono float PCM; Close signals end-of-audio, waits for
// the backend's terminal Final event to be delivered, and releases the
// stream.
type LiveTranscriptionSession interface {
Feed(pcm []float32) error
Close() error
}
// liveCloseDrainTimeout bounds how long Close waits for the backend to flush
// the decode tail before force-cancelling the stream. Finalize is one short
// engine call; seconds here means the backend is wedged.
const liveCloseDrainTimeout = 10 * time.Second
type liveTranscriptionSession struct {
stream grpcPkg.AudioTranscriptionLiveClient
cancel context.CancelFunc
recvDone chan struct{}
recvErr error // written by the recv goroutine before recvDone closes
closeOnce sync.Once
closeErr error
trace *liveTraceState // nil when tracing was disabled at open
release func()
}
func (s *liveTranscriptionSession) Feed(pcm []float32) error {
s.trace.addPCM(pcm)
return s.stream.Send(&proto.TranscriptLiveRequest{
Payload: &proto.TranscriptLiveRequest_Audio{Audio: &proto.TranscriptLiveAudio{Pcm: pcm}},
})
}
func (s *liveTranscriptionSession) Close() error {
s.closeOnce.Do(func() {
err := s.stream.CloseSend()
select {
case <-s.recvDone:
case <-time.After(liveCloseDrainTimeout):
xlog.Warn("live transcription: backend did not finalize in time; cancelling stream")
s.cancel()
<-s.recvDone
}
s.cancel()
if err == nil {
err = s.recvErr
}
s.closeErr = err
s.trace.record(err)
s.release()
})
return s.closeErr
}
// liveSampleRate is the PCM rate of a live transcription session, fixed by
// the session config sent in ModelTranscriptionLive.
const liveSampleRate = 16000
// liveTraceState accumulates what the per-turn backend trace needs while a
// live session runs: a bounded copy of the fed PCM for the audio snippet,
// the decode outputs, and timing. One trace is recorded at Close — the live
// path never touches the unary transcription wrapper, so without this a
// streaming-only pipeline produced no transcription traces at all. Feed and
// the recv goroutine run concurrently; mu guards the accumulators.
type liveTraceState struct {
appConfig *config.ApplicationConfig
modelName string
backend string
language string
started time.Time
traceID string
mu sync.Mutex
pcm []byte // first trace.MaxSnippetSeconds of fed audio, int16 LE
fedSamples int // ALL samples fed, beyond the snippet cap
deltaEvents int
eouEvents int
eobEvents int
finalText string
}
func newLiveTraceState(modelConfig config.ModelConfig, appConfig *config.ApplicationConfig, language string) *liveTraceState {
if !appConfig.EnableTracing {
return nil
}
trace.InitBackendTracingIfEnabled(appConfig.TracingMaxItems, appConfig.TracingMaxBodyBytes)
started := time.Now()
return &liveTraceState{
appConfig: appConfig,
modelName: modelConfig.Name,
backend: modelConfig.Backend,
language: language,
started: started,
traceID: trace.BeginBackendTrace(trace.BackendTrace{Timestamp: started, Type: trace.BackendTraceTranscription, ModelName: modelConfig.Name, Backend: modelConfig.Backend, Summary: "live transcription"}),
}
}
func (ts *liveTraceState) addPCM(pcm []float32) {
if ts == nil {
return
}
ts.mu.Lock()
defer ts.mu.Unlock()
ts.fedSamples += len(pcm)
maxBytes := trace.MaxSnippetSeconds * liveSampleRate * 2
if room := (maxBytes - len(ts.pcm)) / 2; room > 0 {
if len(pcm) > room {
pcm = pcm[:room]
}
ts.pcm = append(ts.pcm, sound.Float32sToInt16LEBytes(pcm)...)
}
}
func (ts *liveTraceState) observe(ev LiveTranscriptionEvent) {
if ts == nil {
return
}
ts.mu.Lock()
defer ts.mu.Unlock()
if ev.Delta != "" {
ts.deltaEvents++
}
if ev.Eou {
ts.eouEvents++
}
if ev.Eob {
ts.eobEvents++
}
if ev.Final != nil {
ts.finalText = ev.Final.Text
}
}
func (ts *liveTraceState) record(closeErr error) {
if ts == nil || !ts.appConfig.EnableTracing {
return
}
ts.mu.Lock()
data := map[string]any{
"source": "live_stream",
"language": ts.language,
"result_text": ts.finalText,
"eou_events": ts.eouEvents,
"eob_events": ts.eobEvents,
"delta_events": ts.deltaEvents,
}
if snippet := trace.AudioSnippetFromPCM(ts.pcm, liveSampleRate, ts.fedSamples*2, ts.appConfig.TracingMaxBodyBytes); snippet != nil {
maps.Copy(data, snippet)
}
summary := "live -> " + ts.finalText
ts.mu.Unlock()
bt := trace.BackendTrace{
ID: ts.traceID,
Timestamp: ts.started,
Duration: time.Since(ts.started),
Type: trace.BackendTraceTranscription,
ModelName: ts.modelName,
Backend: ts.backend,
Summary: trace.TruncateString(summary, 200),
Data: data,
}
if closeErr != nil {
bt.Error = closeErr.Error()
}
trace.RecordBackendTrace(bt)
}
// LiveOption tunes a live transcription session.
type LiveOption func(*liveOptions)
type liveOptions struct {
knownVoices []voicerecognition.KnownVoice
}
// WithKnownVoices gives the backend the registered voices it may use to name
// the speakers it detects. Backends without speaker identification ignore them.
func WithKnownVoices(v []voicerecognition.KnownVoice) LiveOption {
return func(o *liveOptions) { o.knownVoices = v }
}
// liveConfigProto builds the first message of a live session.
func liveConfigProto(language string, o liveOptions) *proto.TranscriptLiveConfig {
cfg := &proto.TranscriptLiveConfig{Language: language, SampleRate: liveSampleRate}
for _, v := range o.knownVoices {
cfg.KnownVoices = append(cfg.KnownVoices, &proto.KnownVoice{Id: v.ID, Name: v.Name, Embedding: v.Embedding, Model: v.Model})
}
return cfg
}
// ModelTranscriptionLive loads the transcription backend, opens the
// bidirectional AudioTranscriptionLive RPC, sends the session config, and
// BLOCKS until the backend's ready ack. A grpcerrors.
// IsLiveTranscriptionUnsupported error means the backend (or the loaded
// model) cannot do live transcription and the caller should degrade to the
// unary/file path. After a successful return, onEvent is invoked from a
// background goroutine — in order, one event at a time — for every response
// the backend streams, ending with the Final event triggered by Close.
func ModelTranscriptionLive(ctx context.Context, language string,
ml *model.ModelLoader, modelConfig config.ModelConfig, appConfig *config.ApplicationConfig,
onEvent func(LiveTranscriptionEvent), opts ...LiveOption) (LiveTranscriptionSession, error) {
lo := liveOptions{}
for _, f := range opts {
f(&lo)
}
transcriptionModel, err := loadTranscriptionModel(ctx, ml, modelConfig, appConfig)
if err != nil {
return nil, err
}
lo.knownVoices = compatiblePortableVoices(ctx, transcriptionModel, lo.knownVoices)
release, err := AcquireGlobalBackendSlot()
if err != nil {
return nil, err
}
// The derived cancel out-lives this call inside the session: Close uses
// it to unwind the stream (and, in embed mode, the server-side recv
// pump, which only stops on send-close or context cancellation).
streamCtx, cancel := context.WithCancel(ctx)
stream, err := transcriptionModel.AudioTranscriptionLive(streamCtx)
if err != nil {
cancel()
release()
return nil, err
}
fail := func(err error) (LiveTranscriptionSession, error) {
_ = stream.CloseSend()
cancel()
release()
return nil, err
}
if err := stream.Send(&proto.TranscriptLiveRequest{
Payload: &proto.TranscriptLiveRequest_Config{Config: liveConfigProto(language, lo)},
}); err != nil {
return fail(err)
}
// Ready-ack contract: the backend answers a successful open with a
// {ready:true} response before any transcript data; unsupported
// backends surface Unimplemented here instead.
ack, err := stream.Recv()
if err != nil {
return fail(err)
}
if !ack.GetReady() {
return fail(fmt.Errorf("live transcription: backend %q broke the ready-ack contract (first response carried data)", modelConfig.Backend))
}
s := &liveTranscriptionSession{
stream: stream,
cancel: cancel,
recvDone: make(chan struct{}),
trace: newLiveTraceState(modelConfig, appConfig, language),
release: release,
}
go func() {
defer close(s.recvDone)
for {
resp, err := stream.Recv()
if err != nil {
if !errors.Is(err, io.EOF) && streamCtx.Err() == nil {
xlog.Warn("live transcription stream ended unexpectedly", "error", err)
s.recvErr = err
}
return
}
ev := liveEventFromProto(resp)
if ev.Delta == "" && !ev.Eou && !ev.Eob && len(ev.Words) == 0 && ev.Final == nil {
continue // duplicate ready ack / keep-alive: nothing to deliver
}
s.trace.observe(ev)
onEvent(ev)
}
}()
return s, nil
}
func liveEventFromProto(r *proto.TranscriptLiveResponse) LiveTranscriptionEvent {
ev := LiveTranscriptionEvent{
Delta: r.GetDelta(),
Eou: r.GetEou(),
Eob: r.GetEob(),
}
for _, w := range r.GetWords() {
ev.Words = append(ev.Words, schema.TranscriptionWord{
Start: time.Duration(w.Start),
End: time.Duration(w.End),
Text: w.Text,
Speaker: w.Speaker,
})
}
for _, s := range r.GetSpeakers() {
ev.Speakers = append(ev.Speakers, LiveSpeakerSegment{
Speaker: s.GetSpeaker(),
Name: s.GetName(),
Start: time.Duration(s.GetStart()).Seconds(),
End: time.Duration(s.GetEnd()).Seconds(),
})
}
for _, s := range r.GetSounds() {
ev.Sounds = append(ev.Sounds, LiveSoundEvent{
Label: s.GetLabel(),
Index: int(s.GetIndex()),
Peak: s.GetPeak(),
Start: time.Duration(s.GetStart()).Seconds(),
End: time.Duration(s.GetEnd()).Seconds(),
})
}
if r.GetFinalResult() != nil {
ev.Final = transcriptResultFromProto(r.GetFinalResult())
}
return ev
}