mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-05 12:34:43 -04:00
* 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>
380 lines
11 KiB
Go
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
|
|
}
|