Files
LocalAI/backend/go/parakeet-cpp/live.go
T
localai-org-maint-botandEttore Di Giacinto 2ae6cae70d feat(parakeet-cpp): name speakers from the shared voice registry (#12382)
* feat(voice): list registered voices and record which encoder made them

The voice registry could register, identify and forget but not list, and
it did not remember which speaker encoder produced an embedding. Add
Metadata.Model and Registry.List, answered from the index the store
registry already keeps for Forget. Needed so a backend can be given the
registered voices that match its own speaker encoder.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(voice): store the encoder model with a registered voice

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(voice): pick the registered voices that match a speaker model

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(proto): carry known voices and speaker names on diarize and live messages

Assisted-by: Claude:claude-haiku-4-5 [Claude Code]

* feat(diarization): name speakers from the voice registry

When a diarization model has a speaker_model option, the endpoint sends
the registered voices made by that encoder to the backend. The backend's
name and name_score come back as extra fields next to the normalized
SPEAKER_NN speaker, and the speakers summary carries the first name seen
for each speaker. RTTM output and results without names are unchanged.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(live): pass registered voices to a live session and surface speaker names

Live sessions now send the registered voices that match the model's
speaker_model to the backend, and each speaker segment carries the name
the backend matched. The realtime segment event gains an optional
speaker_name field.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(parakeet-cpp): load a speaker model and build per-request voice registries

Adds the speaker bindings (ABI v9 and v10, probed separately), the
speaker_model, speaker_threshold and speaker_margin options, and a
per-request registry builder over the known voices.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(parakeet-cpp): name the speakers in Diarize from the known voices

Diarize builds a per-request speaker registry from the known voices when a
speaker model is loaded, calls the named C functions, and puts each slot's
registered name and score on the segments. The registry is freed on every
path. A library without ABI 10 reports Unimplemented instead of dropping
the names.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(parakeet-cpp): name speakers in the live scene stream

The live scene stream now begins with a known-voice registry when a
speaker model is loaded and the live config carries voices, and each
closed speaker segment takes its slot's current name from the feed's
names map. A segment that closes before its slot is identified has an
empty name. The registry is freed after the stream, on every path.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* feat(gallery): speaker naming entries and docs for parakeet-cpp

Add three gallery entries that load the WeSpeaker ResNet34 speaker model
next to the diarization or realtime scene models, and document speaker
names in the voice recognition, diarization, audio to text and realtime
pages.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* fix(parakeet-cpp): skip an unusable registered voice instead of failing the request

A registered voice with the wrong embedding size, or one the C side
refused, failed the whole diarization request, so one legacy voice broke
the model for every user. Skip such voices with a warning that does not
carry the voice name, and take the plain path when none is left.

Also map an exact 0 speaker threshold or margin to a tiny positive value,
since the C side reads 0 as "use the default", and fix a stale comment
about which contexts Free() walks.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* fix(diarization): warn once per model about voices from another encoder; document the privacy limit

The different-encoder warning fired on every request. Log it once per
feature and speaker model, then at debug level. Document that the global
voice registry lets any caller of a speaker_model model learn matching
names, and that skipped wrong-sized voices are logged.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

* chore(parakeet-cpp): bump parakeet.cpp to 8c8cec0 (C-API v10) and check speaker naming against the real library

The pin moves from 623a968 to 8c8cec0, which brings in everything merged
in parakeet.cpp since: the voice identification change (C-API v9, #78) and
raw-embedding enroll plus diarize-only speaker naming (C-API v10, #79).

New real-library specs (gated on PARAKEET_BACKEND_TEST_SPEAKER_MODEL,
_DIAR_MODEL, _WAV and, for the live path, _STREAM_MODEL) name the two
speakers of two_speakers.wav from a committed pair of WeSpeaker embeddings,
with the voices passed in reversed order. They also check that the float32
threshold reaches C through purego. The shared test loader now registers
the v9/v10 and scene symbols as main.go does.

The rebase onto origin/master had no conflicts.

Assisted-by: Claude:claude-sonnet-5-5 [Claude Code]

---------

Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
2026-10-01 08:25:52 +02:00

245 lines
8.9 KiB
Go

package main
import (
"strings"
"time"
"github.com/mudler/LocalAI/pkg/grpc/grpcerrors"
pb "github.com/mudler/LocalAI/pkg/grpc/proto"
"github.com/mudler/xlog"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
// liveSampleRate is the only PCM rate the parakeet C streaming API accepts.
const liveSampleRate = 16000
// AudioTranscriptionLive drives one cache-aware streaming session over audio
// fed incrementally by the caller (the realtime API's semantic_vad turn
// detection). Contract:
//
// - the first request must carry a Config; a Config mid-stream resets the
// decode session (free + begin) and drops accumulated transcript state;
// - a Ready ack is sent right after a successful stream_begin so callers
// can degrade synchronously when the model has no streaming support
// (LiveTranscriptionUnsupported, codes.Unimplemented);
// - every feed that produced output is forwarded as {delta, eou, words};
// the <EOU>/<EOB> flag is the model's own utterance boundary and the
// decoder auto-resets after it, so one session spans many utterances;
// - closing the send side finalizes: the held-back tail chunk is flushed
// (the last ~2 encoder frames of words only appear here) and a terminal
// FinalResult carries the full transcript Text only. Per-utterance
// segments, duration, and the terminal <EOU> flag are NOT produced here —
// the realtime core consumes the streamed per-feed tokens and the final
// Text; those batch fields are the file path's concern (see
// AudioTranscriptionStream).
//
// Engine access is serialized per C call (streamBegin/streamFeed*/streamFree
// take engineMu internally), never for the session lifetime — unary
// transcription keeps flowing between feeds.
func (p *ParakeetCpp) AudioTranscriptionLive(in <-chan *pb.TranscriptLiveRequest, out chan<- *pb.TranscriptLiveResponse) error {
defer close(out)
if p.ctxPtr == 0 {
if err := p.notASRError(); err != nil {
return err
}
return grpcerrors.ModelNotLoaded("parakeet-cpp")
}
first, ok := <-in
if !ok {
return nil // caller closed without sending anything
}
cfg := first.GetConfig()
if cfg == nil {
return status.Error(codes.InvalidArgument, "parakeet-cpp: first live message must carry a config")
}
if err := validateLiveConfig(cfg); err != nil {
return err
}
stream, err := p.streamBegin(cfg.GetLanguage())
if err != nil {
return err
}
if stream == 0 {
return grpcerrors.LiveTranscriptionUnsupported("parakeet-cpp",
"loaded model is not a cache-aware streaming model")
}
// stream is reassigned on a mid-stream Config reset; free whatever is
// current when the RPC unwinds.
defer func() { p.streamFree(stream) }()
// scene runs a no-ASR scene stream (diarization/sound only) beside the
// ASR session when a diarization_model:/sound_model: companion is loaded
// (see scene.go). A zero handle means scene events are disabled: no
// companions, or the begin/a later feed call failed (logged below / in
// feedSlicesScene), in which case live transcription continues ASR-only.
// Reassigned on a mid-stream Config reset alongside stream, which also
// brings back a scene stream the session had disabled after an earlier
// scene error.
var scene sceneStreamHandle
if p.sceneWanted() {
scene = p.sceneBegin(cfg.GetKnownVoices())
if scene.s == 0 {
xlog.Warn("parakeet-cpp: scene stream begin failed; live continues without speaker/sound events")
}
}
defer func() { p.sceneFree(scene) }()
out <- &pb.TranscriptLiveResponse{Ready: true}
var (
full strings.Builder
fedSecs float64
// behindSec accumulates how far decode wall time has fallen behind
// the audio it was fed. A live caller feeds in real time, so a
// persistent positive backlog means every downstream signal —
// including the <EOU> the turn detector waits on — arrives that many
// seconds late. Warned once per session; reset by a Config reset.
behindSec float64
behindWarned bool
)
// emit sends one decode increment as its own response when it carries
// anything: either the ASR side (delta/eou/eob/words, accumulated into
// the running transcript for the closing FinalResult) or the scene
// side (closed speaker/sound events), never both at once — the live
// audio loop below calls it once for the ASR result right after the ASR
// feed and, separately, once more for the scene document after the
// scene feed (see feedSlicesScene), so a slice with both produces two
// responses, ASR first. No segmentation or boundary latch here — the
// live consumer reads only the streamed tokens and the final Text;
// per-utterance segments and the terminal <EOU> flag are an
// offline-path concern (see AudioTranscriptionStream / boundary.go).
emit := func(r streamFeedResult, sceneDoc sceneFeedJSON) error {
if r.Delta != "" {
full.WriteString(r.Delta)
}
speakers := liveSpeakersToProto(sceneDoc.Speakers, sceneDoc.Names)
sounds := liveSoundsToProto(sceneDoc.Sounds)
if r.Delta != "" || r.Eou || r.Eob || len(r.Words) > 0 || len(speakers) > 0 || len(sounds) > 0 {
out <- &pb.TranscriptLiveResponse{
Delta: r.Delta,
Eou: r.Eou,
Eob: r.Eob,
Words: liveWordsToProto(r.Words),
Speakers: speakers,
Sounds: sounds,
}
}
return nil
}
for req := range in {
switch payload := req.GetPayload().(type) {
case *pb.TranscriptLiveRequest_Config:
if err := validateLiveConfig(payload.Config); err != nil {
return err
}
// Reset: a fresh decode session, dropping accumulated state.
p.streamFree(stream)
stream, err = p.streamBegin(payload.Config.GetLanguage())
if err != nil {
return err
}
if stream == 0 {
return grpcerrors.LiveTranscriptionUnsupported("parakeet-cpp",
"loaded model is not a cache-aware streaming model")
}
// The scene stream is freed and begun again alongside the ASR
// session, mirroring the reset above.
p.sceneFree(scene)
scene = sceneStreamHandle{}
if p.sceneWanted() {
scene = p.sceneBegin(payload.Config.GetKnownVoices())
if scene.s == 0 {
xlog.Warn("parakeet-cpp: scene stream begin failed; live continues without speaker/sound events")
}
}
full.Reset()
fedSecs = 0
behindSec = 0
behindWarned = false
case *pb.TranscriptLiveRequest_Audio:
pcm := payload.Audio.GetPcm()
audioSec := float64(len(pcm)) / liveSampleRate
fedSecs += audioSec
start := time.Now()
// nil ctx: a live session is bounded by this request channel, not a
// context — cancellation is the caller closing the stream.
var asrWall, sceneWall time.Duration
scene, asrWall, sceneWall, err = p.feedSlicesScene(nil, stream, scene, pcm, emit)
if err != nil {
return err
}
wallSec := time.Since(start).Seconds()
behindSec += wallSec - audioSec
if behindSec < 0 {
behindSec = 0
}
xlog.Debug("parakeet-cpp: live feed",
"audio_ms", int(audioSec*1000), "wall_ms", int(wallSec*1000),
"asr_wall_ms", int(asrWall.Seconds()*1000), "scene_wall_ms", int(sceneWall.Seconds()*1000),
"behind_ms", int(behindSec*1000), "fed_s", fedSecs)
if behindSec > 1 && !behindWarned {
behindWarned = true
xlog.Warn("parakeet-cpp: live decode is falling behind real time; "+
"end-of-utterance signals will arrive late",
"behind_s", behindSec, "fed_s", fedSecs)
}
}
}
// Send side closed: flush the streaming tail and emit the final transcript.
// The live FinalResult carries only Text — the authoritative full-turn
// transcript the realtime core commits. Per-utterance segments, duration,
// and the terminal <EOU> flag are not produced on the live path.
if err := p.flushTail(stream, func(r streamFeedResult) error {
return emit(r, sceneFeedJSON{})
}); err != nil {
return err
}
// The scene stream gets its own is_last flush (it consumes no new audio
// here, so it is not part of flushTail above); its remaining events go
// out before the terminal FinalResult, then the stream is released by the
// deferred sceneFree above.
if scene.s != 0 {
doc, err := p.sceneFeed(scene, nil, true)
if err != nil {
xlog.Warn("parakeet-cpp: live scene finalize failed", "err", err)
} else if err := emit(streamFeedResult{}, doc); err != nil {
return err
}
}
out <- &pb.TranscriptLiveResponse{
FinalResult: &pb.TranscriptResult{Text: strings.TrimSpace(full.String())},
}
return nil
}
func validateLiveConfig(cfg *pb.TranscriptLiveConfig) error {
if sr := cfg.GetSampleRate(); sr != 0 && sr != liveSampleRate {
return status.Errorf(codes.InvalidArgument,
"parakeet-cpp: unsupported live sample_rate %d (only %d)", sr, liveSampleRate)
}
return nil
}
func liveWordsToProto(words []transcriptWord) []*pb.TranscriptWord {
if len(words) == 0 {
return nil
}
out := make([]*pb.TranscriptWord, len(words))
for i, w := range words {
out[i] = &pb.TranscriptWord{
Start: secondsToNanos(w.Start),
End: secondsToNanos(w.End),
Text: w.W,
}
}
return out
}