Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 99e73ea653 | |||
| ab1784f5e1 | |||
| 2c73493bf8 | |||
| ff202c0c35 | |||
| 02d96e611d | |||
| 62eef01c18 | |||
| 1a8aed35b8 | |||
| ce6a6821a9 | |||
| 479b0c4475 | |||
| 877b1fd4f8 | |||
| 21a42cb3e6 | |||
| b8279f6a22 | |||
| 8c30971a96 | |||
| d0ea927ac3 | |||
| 9c7bafd5b1 | |||
| 1f1e002789 | |||
| ce91d20ac8 |
+47
-5
@@ -12,10 +12,17 @@
|
||||
// samples and the capture frame is 480, so silero.go re-chunks. This comment
|
||||
// used to say the two matched, which was true of silero v4.
|
||||
//
|
||||
// There is still no wake-word model, so anything spoken near the microphone
|
||||
// becomes a turn (V-487 stage two). The SurfaceVoice auth layer caps all
|
||||
// commands at L0 (no destructive acts), which is what makes an accidental
|
||||
// trigger safe rather than expensive.
|
||||
// The keyword is "Мэйвен" and it is required, when -wake-model points at the
|
||||
// head (V-487 stage two). Without it anything spoken near the microphone
|
||||
// becomes a turn, which the SurfaceVoice auth layer makes safe rather than
|
||||
// expensive: it caps all commands at L0, no destructive acts. It does not cap
|
||||
// reading, so an open gate still lets the room hear his facts read back.
|
||||
// wakeword.go holds the cadence and wakefeatures.go the three models.
|
||||
//
|
||||
// The conn carries both directions. mavwaked sends utterances and receives
|
||||
// proactive nudges on it, and it is opened at startup rather than at the first
|
||||
// utterance, because mavend registers a voice session on accept. See nudge.go
|
||||
// for why a nudge that is not heard is worse than one that is not delivered.
|
||||
//
|
||||
// While a reply is playing the capture side is muted (half-duplex): without
|
||||
// it, Maven's own voice comes back in through the mic and she answers
|
||||
@@ -57,6 +64,12 @@ const (
|
||||
defaultAddr = "127.0.0.1:9100"
|
||||
defaultLang = "ru"
|
||||
defaultReadSize = 4096 // max PCM bytes per read from arecord (fits multiple frames)
|
||||
|
||||
// defaultWakeWindowMs — how long the keyword stays good for. He says
|
||||
// "Мэйвен" and then a sentence, and the VAD does not close the utterance
|
||||
// until he stops, so this has to outlive the word by the length of what
|
||||
// follows it. It is spent on dispatch: one keyword, one turn.
|
||||
defaultWakeWindowMs = 8000
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -81,12 +94,18 @@ func run(args []string) error {
|
||||
vadModel := flag.String("vad-model", "", "silero-vad onnx file; empty runs the energy threshold instead")
|
||||
vadThreshold := flag.Float64("vad-threshold", defaultSileroThreshold, "speech probability a frame must clear")
|
||||
onnxLib := flag.String("onnx-lib", os.Getenv("MAVEN_ONNX_LIB"), "libonnxruntime.so, needed with -vad-model")
|
||||
wakeModel := flag.String("wake-model", "", "keyword head onnx; empty ships every utterance, as before V-487")
|
||||
wakeMel := flag.String("wake-mel", "", "melspectrogram.onnx, required with -wake-model")
|
||||
wakeEmbed := flag.String("wake-embed", "", "embedding_model.onnx, required with -wake-model")
|
||||
wakeThreshold := flag.Float64("wake-threshold", defaultWakeThreshold, "score the keyword must clear")
|
||||
wakeWindowMs := flag.Int("wake-window-ms", defaultWakeWindowMs, "ms an utterance may still start after the keyword")
|
||||
flag.CommandLine.Parse(args)
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
|
||||
defer stop()
|
||||
|
||||
// Voice client — reused across utterances; SendRequest reconnects on error.
|
||||
// Voice client — one conn carrying both directions. SendRequest reconnects
|
||||
// on error, and the push receiver redials on its own clock.
|
||||
vc := voice.Dial(*addr)
|
||||
defer vc.Close()
|
||||
|
||||
@@ -165,6 +184,29 @@ func run(args []string) error {
|
||||
}
|
||||
sess := newSession(vad, newAplayPlayer(), &voiceSender{vc: vc}, *lang, barge)
|
||||
|
||||
// Keyword gate. A model that will not load is logged and not fatal, for
|
||||
// the same reason silero's is not: an open gate is the daemon he had
|
||||
// yesterday, and a daemon that refuses to start is not.
|
||||
if *wakeModel != "" {
|
||||
w, err := newWakeWord(*wakeMel, *wakeEmbed, *wakeModel, *onnxLib, *wakeThreshold)
|
||||
if err != nil {
|
||||
log.Printf("mavwaked: wake word unavailable, every utterance is a turn: %v", err)
|
||||
} else {
|
||||
defer w.Close()
|
||||
sess.UseWakeWord(w, time.Duration(*wakeWindowMs)*time.Millisecond)
|
||||
log.Printf("mavwaked: wake word from %s, threshold %.3f, window %dms",
|
||||
*wakeModel, *wakeThreshold, *wakeWindowMs)
|
||||
}
|
||||
}
|
||||
|
||||
// Listen for nudges alongside capture. Connect eagerly so mavend has a
|
||||
// voice session before he has said anything: without one, a nudge routed
|
||||
// to voice finds nobody home and goes to the away channels instead.
|
||||
if err := vc.Connect(ctx); err != nil {
|
||||
log.Printf("mavwaked: voice server not reachable yet, retrying in background: %v", err)
|
||||
}
|
||||
go runNudgeReceiver(ctx, vc, sess)
|
||||
|
||||
return captureLoop(ctx, src, sess)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
package main
|
||||
|
||||
// The receiving half of the voice reach (V-671).
|
||||
//
|
||||
// mavwaked used to send and never listen. It wired no PushHandler, and
|
||||
// SendRequest discards a push frame when there is none. The consequence was
|
||||
// not a missing feature but a silent one: mavend routes a nudge to the voice
|
||||
// session that spoke most recently, and once mavwaked had spoken once it WAS
|
||||
// that session. PushToMostRecent succeeded, the dispatcher counted the nudge
|
||||
// delivered and stopped rerouting to telegram and ntfy, and mavwaked threw the
|
||||
// audio away. He heard nothing, anywhere.
|
||||
//
|
||||
// So the connection is opened at startup rather than at the first utterance,
|
||||
// and it is held open. A client that has never connected has no session, and
|
||||
// the dispatcher must be able to tell "he is not at the machine" from "he is,
|
||||
// and she has nothing to say".
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/voice"
|
||||
)
|
||||
|
||||
// nudgeRetry is how long to wait before dialling again after the conn ends.
|
||||
// mavend restarts on every deploy, and a listener that gives up then is a
|
||||
// listener that is deaf until the next reboot.
|
||||
const nudgeRetry = 5 * time.Second
|
||||
|
||||
// nudgeHandler decodes a push and hands the audio to the session, which
|
||||
// speaks it through the same player the reply path uses. It does not play
|
||||
// anything itself: the half-duplex gate and barge-in live on the capture
|
||||
// loop, and a nudge has to sit under both.
|
||||
type nudgeHandler struct{ sess *session }
|
||||
|
||||
func (h *nudgeHandler) OnPush(p voice.Push) {
|
||||
if p.Kind != voice.PushKindAudioNudge {
|
||||
log.Printf("mavwaked: ignoring push of unknown kind %q", p.Kind)
|
||||
return
|
||||
}
|
||||
var ap voice.AudioNudgePush
|
||||
if err := json.Unmarshal(p.Params, &ap); err != nil {
|
||||
log.Printf("mavwaked: nudge: decode: %v", err)
|
||||
return
|
||||
}
|
||||
log.Printf("mavwaked: nudge from rule %q (severity %d): %q (%.2fs audio)",
|
||||
ap.RuleName, ap.Severity, ap.Text, ap.Audio.Duration())
|
||||
if len(ap.Audio.Bytes) == 0 {
|
||||
// mavttsd was down or the text was empty. Say so rather than going
|
||||
// quiet: the dispatcher already counted this one as delivered.
|
||||
log.Printf("mavwaked: nudge %q carried no audio, nothing to speak", ap.RuleName)
|
||||
return
|
||||
}
|
||||
h.sess.Nudge(ap.Audio)
|
||||
}
|
||||
|
||||
// runNudgeReceiver keeps a push handler wired for as long as ctx lives,
|
||||
// redialling whenever the conn ends. Returns when ctx is cancelled.
|
||||
func runNudgeReceiver(ctx context.Context, vc *voice.Client, sess *session) {
|
||||
h := &nudgeHandler{sess: sess}
|
||||
for {
|
||||
err := vc.RunPushReceiver(ctx, h)
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
log.Printf("mavwaked: nudge receiver: %v", err)
|
||||
} else {
|
||||
log.Printf("mavwaked: voice connection ended, reconnecting in %s", nudgeRetry)
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(nudgeRetry):
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
package main
|
||||
|
||||
// The receiving half: a nudge pushed by mavend has to reach the speaker, and
|
||||
// it has to obey the same two gates a reply obeys (V-671).
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/audio"
|
||||
"github.com/kami/maven/internal/voice"
|
||||
)
|
||||
|
||||
func nudgeAudio() audio.Audio {
|
||||
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 8000)}
|
||||
}
|
||||
|
||||
// pushFrame builds the frame mavend's voicesink sends.
|
||||
func pushFrame(t *testing.T, a audio.Audio) voice.Push {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(voice.AudioNudgePush{
|
||||
RuleName: "test-rule",
|
||||
Severity: 3,
|
||||
Audio: a,
|
||||
Text: "пора пить воду",
|
||||
Ts: time.Unix(0, 0),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("marshal push: %v", err)
|
||||
}
|
||||
return voice.Push{Kind: voice.PushKindAudioNudge, Params: body}
|
||||
}
|
||||
|
||||
// The defect itself: the push arrived and nothing came out of the speaker.
|
||||
func TestNudgeReachesThePlayer(t *testing.T) {
|
||||
sess, p, snd := newTestSession(bargeInConfig{})
|
||||
(&nudgeHandler{sess: sess}).OnPush(pushFrame(t, nudgeAudio()))
|
||||
|
||||
if p.plays != 0 {
|
||||
t.Fatal("nudge played from the push goroutine; it must wait for the capture loop")
|
||||
}
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != 1 {
|
||||
t.Fatalf("plays = %d, want 1", p.plays)
|
||||
}
|
||||
if len(p.last.Bytes) != 8000 {
|
||||
t.Errorf("played %d bytes, want the nudge audio", len(p.last.Bytes))
|
||||
}
|
||||
if sess.nudges != 1 {
|
||||
t.Errorf("nudges = %d, want 1", sess.nudges)
|
||||
}
|
||||
if len(snd.sent) != 0 {
|
||||
t.Errorf("a nudge must not be shipped back to the daemon as an utterance")
|
||||
}
|
||||
}
|
||||
|
||||
// A push of some other kind, or one carrying no audio, must not reach the
|
||||
// player and must not wedge the one that follows.
|
||||
func TestNudgeIgnoresUnusablePushes(t *testing.T) {
|
||||
sess, p, _ := newTestSession(bargeInConfig{})
|
||||
h := &nudgeHandler{sess: sess}
|
||||
|
||||
h.OnPush(voice.Push{Kind: "something-else", Params: json.RawMessage(`{}`)})
|
||||
h.OnPush(voice.Push{Kind: voice.PushKindAudioNudge, Params: json.RawMessage(`not json`)})
|
||||
h.OnPush(pushFrame(t, audio.Audio{Format: audio.PCM16kMono}))
|
||||
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != 0 {
|
||||
t.Fatalf("plays = %d, want 0", p.plays)
|
||||
}
|
||||
|
||||
h.OnPush(pushFrame(t, nudgeAudio()))
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != 1 {
|
||||
t.Fatalf("plays after a usable nudge = %d, want 1", p.plays)
|
||||
}
|
||||
}
|
||||
|
||||
// The half-duplex gate covers a nudge exactly as it covers a reply: she does
|
||||
// not start one over herself, and the mic stays muted while it runs.
|
||||
func TestNudgeWaitsForTheReplyToFinish(t *testing.T) {
|
||||
sess, p, _ := newTestSession(bargeInConfig{})
|
||||
speakThenPause(t, sess)
|
||||
if !p.Playing() {
|
||||
t.Fatal("expected the reply to be playing")
|
||||
}
|
||||
plays := p.plays
|
||||
|
||||
(&nudgeHandler{sess: sess}).OnPush(pushFrame(t, nudgeAudio()))
|
||||
for i := 0; i < 20; i++ {
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
}
|
||||
if p.plays != plays {
|
||||
t.Fatalf("nudge cut across the reply: plays = %d, want %d", p.plays, plays)
|
||||
}
|
||||
|
||||
p.playing = false
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != plays+1 {
|
||||
t.Fatalf("nudge never played after the reply ended: plays = %d", p.plays)
|
||||
}
|
||||
}
|
||||
|
||||
// Speaking a nudge must not leave half a sentence in the VAD. The frames
|
||||
// captured before it are pre-nudge speech, and splicing them onto whatever he
|
||||
// says afterwards ships one utterance that is two.
|
||||
func TestNudgeResetsTheVAD(t *testing.T) {
|
||||
sess, p, snd := newTestSession(bargeInConfig{})
|
||||
loud := frameAt(0.35)
|
||||
speechFrames := (defaultSpeechMs + defaultFrameMs - 1) / defaultFrameMs
|
||||
for i := 0; i < speechFrames+5; i++ {
|
||||
if err := sess.feed(context.Background(), loud); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
(&nudgeHandler{sess: sess}).OnPush(pushFrame(t, nudgeAudio()))
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != 1 {
|
||||
t.Fatalf("nudge did not play: plays = %d", p.plays)
|
||||
}
|
||||
|
||||
// Playback ends, silence follows. The half-formed utterance must be gone
|
||||
// rather than closing on the first quiet frame.
|
||||
p.playing = false
|
||||
silenceFrames := (defaultSilenceMs+defaultFrameMs-1)/defaultFrameMs + 2
|
||||
for i := 0; i < silenceFrames; i++ {
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
}
|
||||
if len(snd.sent) != 0 {
|
||||
t.Fatalf("sent %d utterances after a nudge, want 0", len(snd.sent))
|
||||
}
|
||||
}
|
||||
|
||||
// Two nudges queued back to back: the newer one is what he hears. The
|
||||
// PushHandler contract in internal/voice says the next nudge replaces the
|
||||
// stale one rather than dogpiling on it.
|
||||
func TestNudgeReplacesAnUnspokenOne(t *testing.T) {
|
||||
sess, p, _ := newTestSession(bargeInConfig{})
|
||||
h := &nudgeHandler{sess: sess}
|
||||
|
||||
h.OnPush(pushFrame(t, audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 4000)}))
|
||||
h.OnPush(pushFrame(t, audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 12000)}))
|
||||
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if p.plays != 1 {
|
||||
t.Fatalf("plays = %d, want 1", p.plays)
|
||||
}
|
||||
if len(p.last.Bytes) != 12000 {
|
||||
t.Errorf("played %d bytes, want the newer nudge", len(p.last.Bytes))
|
||||
}
|
||||
}
|
||||
+136
-1
@@ -7,6 +7,7 @@ package main
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/audio"
|
||||
@@ -19,6 +20,15 @@ type utteranceSender interface {
|
||||
Send(ctx context.Context, utt audio.Audio, lang string) (audio.Audio, error)
|
||||
}
|
||||
|
||||
// keywordGate answers whether the keyword has just been spoken. The
|
||||
// production one is wakeWord; tests substitute a recorder, because a gate that
|
||||
// can only be exercised with three ONNX files is a gate nobody tests.
|
||||
type keywordGate interface {
|
||||
Feed(frame []int16) bool
|
||||
Reset()
|
||||
Score() float64
|
||||
}
|
||||
|
||||
// bargeInConfig holds the two numbers barge-in needs. Zero Frames disables
|
||||
// barge-in entirely — the half-duplex gate still runs.
|
||||
type bargeInConfig struct {
|
||||
@@ -62,11 +72,29 @@ type session struct {
|
||||
// whenever playback ends.
|
||||
loudFrames int
|
||||
|
||||
// wake is the keyword gate, or nil when no model was loaded. wakeUntil is
|
||||
// how long a keyword stays good for: he says "Мэйвен" and then a sentence,
|
||||
// and the VAD does not close the utterance until he stops, so the window
|
||||
// has to outlive the word by the length of what follows it.
|
||||
wake keywordGate
|
||||
wakeWindow time.Duration
|
||||
wakeUntil time.Time
|
||||
|
||||
// pending holds a nudge the push receiver handed over, waiting for the
|
||||
// capture loop to speak it. It is the one field written from another
|
||||
// goroutine, hence the mutex; everything else in this struct belongs to
|
||||
// the capture loop alone.
|
||||
nudgeMu sync.Mutex
|
||||
pending *audio.Audio
|
||||
|
||||
// counters, read by tests and logged on the way out.
|
||||
suppressed int // frames dropped because she was speaking
|
||||
dropped int // frames dropped as round-trip backlog
|
||||
bargeIns int // times playback was cut because he spoke over her
|
||||
sent int // utterances shipped to the daemon
|
||||
nudges int // proactive pushes spoken through the speaker
|
||||
wakes int // times the keyword opened the gate
|
||||
ignored int // complete utterances dropped because the keyword was absent
|
||||
|
||||
// loudSum and loudSeen accumulate the energy of suppressed frames, so
|
||||
// the operator can read what the room actually measures and set
|
||||
@@ -79,6 +107,12 @@ func newSession(vad *VAD, p player, s utteranceSender, lang string, barge bargeI
|
||||
return &session{vad: vad, player: p, sender: s, lang: lang, barge: barge, now: time.Now}
|
||||
}
|
||||
|
||||
// UseWakeWord puts the keyword gate in front of dispatch. Without it every
|
||||
// utterance is shipped, which is what mavwaked did before V-487 stage two.
|
||||
func (s *session) UseWakeWord(w keywordGate, window time.Duration) {
|
||||
s.wake, s.wakeWindow = w, window
|
||||
}
|
||||
|
||||
// frameDuration is the wall time one captured frame represents.
|
||||
const frameDuration = defaultFrameMs * time.Millisecond
|
||||
|
||||
@@ -138,6 +172,7 @@ func (s *session) feed(ctx context.Context, frame []byte) error {
|
||||
s.bargeIns++
|
||||
s.loudFrames = 0
|
||||
s.vad.Reset()
|
||||
s.resetWake()
|
||||
log.Printf("mavwaked: barge-in — stopped playback")
|
||||
s.replayRecent()
|
||||
return nil
|
||||
@@ -148,15 +183,101 @@ func (s *session) feed(ctx context.Context, frame []byte) error {
|
||||
if s.loudFrames != 0 {
|
||||
s.loudFrames = 0
|
||||
s.vad.Reset()
|
||||
// The wake word saw nothing during playback, so what it holds is from
|
||||
// before she spoke. Judging what he says next on it would score a
|
||||
// sentence that ended a reply ago.
|
||||
s.resetWake()
|
||||
}
|
||||
|
||||
utt, state := s.vad.Feed(PCMToI16(frame))
|
||||
if s.startPendingNudge() {
|
||||
return nil
|
||||
}
|
||||
|
||||
// The keyword is scored on the same frames the VAD sees, and only on the
|
||||
// ones that reach here: every path above returns while she is speaking, so
|
||||
// her own voice saying "Мэйвен" cannot wake her.
|
||||
pcm := PCMToI16(frame)
|
||||
if s.wake != nil && s.wake.Feed(pcm) {
|
||||
s.wakes++
|
||||
s.wakeUntil = s.now().Add(s.wakeWindow)
|
||||
log.Printf("mavwaked: keyword heard (score %.3f), listening for %s",
|
||||
s.wake.Score(), s.wakeWindow)
|
||||
}
|
||||
|
||||
utt, state := s.vad.Feed(pcm)
|
||||
if state == StateSpeech || utt.Bytes == nil {
|
||||
return nil
|
||||
}
|
||||
return s.dispatch(ctx, utt)
|
||||
}
|
||||
|
||||
// Nudge hands proactive audio to the session, to be spoken as soon as the
|
||||
// capture loop finds a quiet moment. Safe to call from the push receiver
|
||||
// goroutine; nothing else here is.
|
||||
//
|
||||
// A nudge arriving while one is already waiting REPLACES it. That is the
|
||||
// contract internal/voice states for PushHandler: the next nudge replaces the
|
||||
// stale one in his attention rather than dogpiling on it.
|
||||
func (s *session) Nudge(a audio.Audio) {
|
||||
if len(a.Bytes) == 0 {
|
||||
return
|
||||
}
|
||||
s.nudgeMu.Lock()
|
||||
if s.pending != nil {
|
||||
log.Printf("mavwaked: nudge replaced one still waiting to be spoken")
|
||||
}
|
||||
s.pending = &a
|
||||
s.nudgeMu.Unlock()
|
||||
}
|
||||
|
||||
// takeNudge removes and returns the waiting nudge, or nil.
|
||||
func (s *session) takeNudge() *audio.Audio {
|
||||
s.nudgeMu.Lock()
|
||||
defer s.nudgeMu.Unlock()
|
||||
a := s.pending
|
||||
s.pending = nil
|
||||
return a
|
||||
}
|
||||
|
||||
// startPendingNudge speaks a waiting nudge and reports whether it started
|
||||
// one. It runs on the capture loop, past the half-duplex gate, so a nudge
|
||||
// never cuts across a reply and never plays into a backlog drain.
|
||||
//
|
||||
// The VAD is reset first. Playback is about to suppress every frame until it
|
||||
// ends, and a half-heard sentence left in the VAD would splice onto whatever
|
||||
// he says afterwards. Barge-in needs no special case: it reads the player,
|
||||
// and the player does not care which audio it is playing.
|
||||
func (s *session) startPendingNudge() bool {
|
||||
a := s.takeNudge()
|
||||
if a == nil {
|
||||
return false
|
||||
}
|
||||
s.vad.Reset()
|
||||
s.nudges++
|
||||
log.Printf("mavwaked: speaking nudge (%.2fs audio)", a.Duration())
|
||||
s.player.Play(*a)
|
||||
return true
|
||||
}
|
||||
|
||||
// awake reports whether an utterance ending now was addressed to her.
|
||||
//
|
||||
// With no wake word loaded every utterance is, which is exactly what mavwaked
|
||||
// did before this gate existed. An operator with no model file gets the old
|
||||
// daemon rather than a daemon that refuses to hear anything.
|
||||
func (s *session) awake() bool {
|
||||
if s.wake == nil {
|
||||
return true
|
||||
}
|
||||
return s.now().Before(s.wakeUntil)
|
||||
}
|
||||
|
||||
// resetWake drops the gate's streaming state when there is a gate.
|
||||
func (s *session) resetWake() {
|
||||
if s.wake != nil {
|
||||
s.wake.Reset()
|
||||
}
|
||||
}
|
||||
|
||||
// keepRecent stores a copy of one barge-in trigger frame, keeping at most
|
||||
// barge.Frames of them.
|
||||
func (s *session) keepRecent(frame []byte) {
|
||||
@@ -197,6 +318,19 @@ func (s *session) replayRecent() {
|
||||
// whole backlog straight into the VAD, and a Send error did the same on every
|
||||
// failed turn, so a dead socket drove a retry loop off nothing but backlog.
|
||||
func (s *session) dispatch(ctx context.Context, utt audio.Audio) error {
|
||||
if !s.awake() {
|
||||
s.ignored++
|
||||
log.Printf("mavwaked: utterance ignored, keyword not heard (%.2fs, %d ignored so far)",
|
||||
utt.Duration(), s.ignored)
|
||||
s.vad.Reset()
|
||||
s.resetWake()
|
||||
return nil
|
||||
}
|
||||
// One keyword, one turn. A window that renewed itself on every reply would
|
||||
// leave the microphone open for as long as he kept talking, which is the
|
||||
// state this gate exists to end.
|
||||
s.wakeUntil = time.Time{}
|
||||
|
||||
log.Printf("mavwaked: utterance complete (%.2fs, %d bytes), sending...", utt.Duration(), len(utt.Bytes))
|
||||
start := s.now()
|
||||
reply, err := s.sender.Send(ctx, utt, s.lang)
|
||||
@@ -225,6 +359,7 @@ func (s *session) dispatch(ctx context.Context, utt audio.Audio) error {
|
||||
// recorded before she started speaking.
|
||||
func (s *session) dropBacklog(start time.Time) {
|
||||
s.vad.Reset()
|
||||
s.resetWake()
|
||||
s.loudFrames = 0
|
||||
s.recent = s.recent[:0]
|
||||
if elapsed := s.now().Sub(start); elapsed > 0 {
|
||||
|
||||
@@ -0,0 +1,202 @@
|
||||
package main
|
||||
|
||||
// The three models behind the wake word (V-487 stage two).
|
||||
//
|
||||
// openWakeWord's pipeline, run in a row:
|
||||
//
|
||||
// audio -> melspectrogram.onnx -> 32-bin mel frames, one per 10ms
|
||||
// 76 frames -> embedding_model.onnx -> one 96-dim embedding per 80ms
|
||||
// 16 embeds -> maven_wakeword.onnx -> one score
|
||||
//
|
||||
// The first two are frozen and pretrained. Only the last was trained here,
|
||||
// which is why it is 100KB and the other two are megabytes. The shapes are
|
||||
// not guesses: 2.0s of 16kHz audio measures 197 mel frames, and 76-frame
|
||||
// windows at stride 8 give exactly the 16 embeddings the head was fitted on.
|
||||
//
|
||||
// This file knows ONNX and nothing about the 80ms cadence. wakeword.go knows
|
||||
// the cadence and nothing about tensors.
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
ort "github.com/yalue/onnxruntime_go"
|
||||
)
|
||||
|
||||
const (
|
||||
// melHop — samples per mel frame. 10ms at 16kHz.
|
||||
melHop = 160
|
||||
// melBins — mel bins per frame, fixed by melspectrogram.onnx.
|
||||
melBins = 32
|
||||
// embedFrames — mel frames one embedding is computed over, 760ms.
|
||||
embedFrames = 76
|
||||
// embedStride — mel frames between embeddings, 80ms.
|
||||
embedStride = 8
|
||||
// embedDim — the embedding width.
|
||||
embedDim = 96
|
||||
// headWindow — embeddings the head scores at once, 1.28s of audio.
|
||||
headWindow = 16
|
||||
|
||||
// melContext — samples of history prepended to each incremental mel
|
||||
// call, chosen so the eight frames this call yields continue exactly
|
||||
// where the previous call's eight stopped.
|
||||
//
|
||||
// melspectrogram.onnx returns N/160-3 frames for N samples, and frame i
|
||||
// covers [i*160, i*160+400). With 480 samples of history the buffer is
|
||||
// 1760 samples, which is 8 frames, and the oldest of them starts one hop
|
||||
// after the newest of the previous call. Less history leaves a gap: the
|
||||
// first frames of a bare chunk would be computed against silence.
|
||||
melContext = 480
|
||||
|
||||
// chunkSamples — audio per embedding step, 80ms.
|
||||
chunkSamples = embedStride * melHop
|
||||
)
|
||||
|
||||
// wakeModels holds the three ONNX sessions. It runs on CPU threads beside
|
||||
// silero and never touches the GPU. That is a rule, not a result: a wake word
|
||||
// that waits on card admission is not a wake word.
|
||||
type wakeModels struct {
|
||||
mel *ort.DynamicAdvancedSession
|
||||
emb *ort.DynamicAdvancedSession
|
||||
head *ort.DynamicAdvancedSession
|
||||
}
|
||||
|
||||
// newWakeModels loads all three. melPath and embedPath are openWakeWord's
|
||||
// frozen feature models; headPath is the keyword head trained for "Мэйвен".
|
||||
func newWakeModels(melPath, embedPath, headPath, libPath string) (*wakeModels, error) {
|
||||
if !ort.IsInitialized() {
|
||||
if libPath != "" {
|
||||
ort.SetSharedLibraryPath(libPath)
|
||||
}
|
||||
if err := ort.InitializeEnvironment(); err != nil {
|
||||
return nil, fmt.Errorf("wake word: onnx runtime: %w", err)
|
||||
}
|
||||
}
|
||||
// One thread per session, not the default of every core. Measured on
|
||||
// workpc: the default took mavwaked from 68% of one core to 335% of
|
||||
// three, for three graphs that each run in well under 80ms single
|
||||
// threaded. An always-on gate that eats a quarter of the workstation is
|
||||
// not a gate he will leave running.
|
||||
opts, err := ort.NewSessionOptions()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("wake word: session options: %w", err)
|
||||
}
|
||||
defer opts.Destroy()
|
||||
if err := opts.SetIntraOpNumThreads(1); err != nil {
|
||||
return nil, fmt.Errorf("wake word: intra-op threads: %w", err)
|
||||
}
|
||||
if err := opts.SetInterOpNumThreads(1); err != nil {
|
||||
return nil, fmt.Errorf("wake word: inter-op threads: %w", err)
|
||||
}
|
||||
open := func(p string, in, out []string) (*ort.DynamicAdvancedSession, error) {
|
||||
s, err := ort.NewDynamicAdvancedSession(p, in, out, opts)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("wake word: load %s: %w", p, err)
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
m := &wakeModels{}
|
||||
if m.mel, err = open(melPath, []string{"input"}, []string{"output"}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if m.emb, err = open(embedPath, []string{"input_1"}, []string{"conv2d_19"}); err != nil {
|
||||
m.Close()
|
||||
return nil, err
|
||||
}
|
||||
if m.head, err = open(headPath, []string{"embeddings"}, []string{"score"}); err != nil {
|
||||
m.Close()
|
||||
return nil, err
|
||||
}
|
||||
return m, nil
|
||||
}
|
||||
|
||||
// Close releases the three sessions.
|
||||
func (m *wakeModels) Close() {
|
||||
if m == nil {
|
||||
return
|
||||
}
|
||||
for _, s := range []*ort.DynamicAdvancedSession{m.mel, m.emb, m.head} {
|
||||
if s != nil {
|
||||
s.Destroy()
|
||||
}
|
||||
}
|
||||
m.mel, m.emb, m.head = nil, nil, nil
|
||||
}
|
||||
|
||||
// melFrames runs one buffer of samples and returns the mel frames it yielded.
|
||||
func (m *wakeModels) melFrames(buf []float32) ([][melBins]float32, error) {
|
||||
in, err := ort.NewTensor(ort.NewShape(1, int64(len(buf))), buf)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer in.Destroy()
|
||||
|
||||
n := int64(len(buf)/melHop - 3)
|
||||
if n < 1 {
|
||||
return nil, fmt.Errorf("wake word: %d samples yield no mel frames", len(buf))
|
||||
}
|
||||
out, err := ort.NewEmptyTensor[float32](ort.NewShape(1, 1, n, melBins))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer out.Destroy()
|
||||
|
||||
if err := m.mel.Run([]ort.Value{in}, []ort.Value{out}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
data := out.GetData()
|
||||
frames := make([][melBins]float32, n)
|
||||
for i := range frames {
|
||||
for j := 0; j < melBins; j++ {
|
||||
// The scaling openWakeWord applies between the two feature
|
||||
// models, and the head was fitted on its output.
|
||||
frames[i][j] = data[i*melBins+j]/10.0 + 2.0
|
||||
}
|
||||
}
|
||||
return frames, nil
|
||||
}
|
||||
|
||||
// embedding runs embedFrames mel frames through the frozen embedder.
|
||||
func (m *wakeModels) embedding(mels [][melBins]float32) ([embedDim]float32, error) {
|
||||
var e [embedDim]float32
|
||||
flat := make([]float32, 0, embedFrames*melBins)
|
||||
for _, f := range mels {
|
||||
flat = append(flat, f[:]...)
|
||||
}
|
||||
in, err := ort.NewTensor(ort.NewShape(1, embedFrames, melBins, 1), flat)
|
||||
if err != nil {
|
||||
return e, err
|
||||
}
|
||||
defer in.Destroy()
|
||||
out, err := ort.NewEmptyTensor[float32](ort.NewShape(1, 1, 1, embedDim))
|
||||
if err != nil {
|
||||
return e, err
|
||||
}
|
||||
defer out.Destroy()
|
||||
if err := m.emb.Run([]ort.Value{in}, []ort.Value{out}); err != nil {
|
||||
return e, err
|
||||
}
|
||||
copy(e[:], out.GetData())
|
||||
return e, nil
|
||||
}
|
||||
|
||||
// score runs the trained head over headWindow embeddings.
|
||||
func (m *wakeModels) score(embeds [][embedDim]float32) (float64, error) {
|
||||
flat := make([]float32, 0, headWindow*embedDim)
|
||||
for _, e := range embeds {
|
||||
flat = append(flat, e[:]...)
|
||||
}
|
||||
in, err := ort.NewTensor(ort.NewShape(1, headWindow, embedDim), flat)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer in.Destroy()
|
||||
out, err := ort.NewEmptyTensor[float32](ort.NewShape(1, 1))
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer out.Destroy()
|
||||
if err := m.head.Run([]ort.Value{in}, []ort.Value{out}); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return float64(out.GetData()[0]), nil
|
||||
}
|
||||
@@ -0,0 +1,195 @@
|
||||
package main
|
||||
|
||||
// The wake word, "Мэйвен" (V-487 stage two).
|
||||
//
|
||||
// Silero answers "is this frame speech". It does not answer "was this said to
|
||||
// her", and until this file existed nothing did: every utterance near the
|
||||
// microphone became a turn. What made that safe rather than expensive was
|
||||
// SurfaceVoice capping acts at L0, and L0 does not cap reading, so the room
|
||||
// could still hear his facts read back.
|
||||
//
|
||||
// This file owns the 80ms cadence and the three rings of state between the
|
||||
// models. wakefeatures.go owns the tensors.
|
||||
//
|
||||
// Nil is a working value, and it is the CLOSED gate rather than the open one.
|
||||
// Feed on a nil receiver reports no keyword; session.go asks separately
|
||||
// whether a gate exists at all. That split is deliberate: a nil that answers
|
||||
// "yes, keyword" reads as a working wake word in every log line it produces.
|
||||
|
||||
import (
|
||||
"log"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// defaultWakeThreshold — score above which the keyword was said.
|
||||
//
|
||||
// Picked from the false-accept rate on held-out Russian speech, not from
|
||||
// accuracy: a miss costs him a repeat, a false accept costs a turn nobody
|
||||
// asked for. Over 65 minutes of Common Voice, 0.99 woke her three times and
|
||||
// 0.999 once, and the difference in recall was one render out of 126. So the
|
||||
// default is the strict one. `docs/evals/2026-08-09-wake-word.md` has both
|
||||
// tables.
|
||||
const defaultWakeThreshold = 0.999
|
||||
|
||||
// wakeWord is the streaming state around wakeModels. It is fed the same
|
||||
// capture frames the VAD sees and answers whether the keyword has just been
|
||||
// spoken.
|
||||
type wakeWord struct {
|
||||
mu sync.Mutex
|
||||
m *wakeModels
|
||||
|
||||
threshold float64
|
||||
|
||||
// pending holds captured samples not yet part of a full 80ms chunk, and
|
||||
// history holds the melContext samples before them.
|
||||
pending []float32
|
||||
history []float32
|
||||
|
||||
// mels is the newest embedFrames mel frames, oldest first.
|
||||
mels [][melBins]float32
|
||||
// embeds is the newest headWindow embeddings, oldest first.
|
||||
embeds [][embedDim]float32
|
||||
|
||||
last float64 // most recent score, held between chunks
|
||||
}
|
||||
|
||||
// newWakeWord loads the models and wraps them in the streaming gate.
|
||||
func newWakeWord(melPath, embedPath, headPath, libPath string, threshold float64) (*wakeWord, error) {
|
||||
m, err := newWakeModels(melPath, embedPath, headPath, libPath)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if threshold <= 0 {
|
||||
threshold = defaultWakeThreshold
|
||||
}
|
||||
return &wakeWord{m: m, threshold: threshold}, nil
|
||||
}
|
||||
|
||||
// Close releases the models.
|
||||
func (w *wakeWord) Close() {
|
||||
if w == nil {
|
||||
return
|
||||
}
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
w.m.Close()
|
||||
w.m = nil
|
||||
}
|
||||
|
||||
// Feed takes one capture frame and reports whether the keyword was heard on
|
||||
// it. A nil wakeWord hears nothing.
|
||||
func (w *wakeWord) Feed(frame []int16) bool {
|
||||
if w == nil {
|
||||
return false
|
||||
}
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
for _, v := range frame {
|
||||
w.pending = append(w.pending, float32(v)/32768.0)
|
||||
}
|
||||
fired := false
|
||||
for len(w.pending) >= chunkSamples {
|
||||
chunk := w.pending[:chunkSamples]
|
||||
if w.step(chunk) {
|
||||
fired = true
|
||||
}
|
||||
w.history = append(w.history[:0], tailFloat32(append(w.history, chunk...), melContext)...)
|
||||
// Slide the remainder to the front rather than reslicing. This runs
|
||||
// every 80ms for as long as the daemon lives.
|
||||
w.pending = append(w.pending[:0], w.pending[chunkSamples:]...)
|
||||
}
|
||||
return fired
|
||||
}
|
||||
|
||||
// Reset drops the streaming state, so a fresh utterance is not judged on audio
|
||||
// from before it. Called after every dispatch and after barge-in, for the same
|
||||
// reason silero is: echo-era history must not score the next sentence, and her
|
||||
// own voice saying the keyword must not wake her.
|
||||
func (w *wakeWord) Reset() {
|
||||
if w == nil {
|
||||
return
|
||||
}
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
w.pending, w.history = w.pending[:0], w.history[:0]
|
||||
w.mels, w.embeds = nil, nil
|
||||
w.last = 0
|
||||
}
|
||||
|
||||
// Score returns the most recent score, for the operator to read out of the
|
||||
// journal when picking a threshold for his room.
|
||||
func (w *wakeWord) Score() float64 {
|
||||
if w == nil {
|
||||
return 0
|
||||
}
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
return w.last
|
||||
}
|
||||
|
||||
// step runs one 80ms chunk through all three models. It returns true when the
|
||||
// score crosses the threshold on this chunk.
|
||||
func (w *wakeWord) step(chunk []float32) bool {
|
||||
buf := make([]float32, 0, melContext+len(chunk))
|
||||
if pad := melContext - len(w.history); pad > 0 {
|
||||
buf = append(buf, make([]float32, pad)...)
|
||||
}
|
||||
buf = append(buf, tailFloat32(w.history, melContext)...)
|
||||
buf = append(buf, chunk...)
|
||||
|
||||
frames, err := w.m.melFrames(buf)
|
||||
if err != nil {
|
||||
// A failed inference must not silence the microphone. Hold the last
|
||||
// score and let the next chunk try again.
|
||||
log.Printf("mavwaked: wake word: mel: %v", err)
|
||||
return false
|
||||
}
|
||||
w.mels = tailMel(append(w.mels, frames...), embedFrames)
|
||||
if len(w.mels) < embedFrames {
|
||||
return false
|
||||
}
|
||||
e, err := w.m.embedding(w.mels)
|
||||
if err != nil {
|
||||
log.Printf("mavwaked: wake word: embedding: %v", err)
|
||||
return false
|
||||
}
|
||||
w.embeds = tailEmbed(append(w.embeds, e), headWindow)
|
||||
if len(w.embeds) < headWindow {
|
||||
return false
|
||||
}
|
||||
score, err := w.m.score(w.embeds)
|
||||
if err != nil {
|
||||
log.Printf("mavwaked: wake word: head: %v", err)
|
||||
return false
|
||||
}
|
||||
// Report the crossing, not the state. A keyword held above the threshold
|
||||
// for a second is one wake, and firing on every chunk of it would make the
|
||||
// gate look open when it is merely slow to fall.
|
||||
crossed := score >= w.threshold && w.last < w.threshold
|
||||
w.last = score
|
||||
return crossed
|
||||
}
|
||||
|
||||
// The three rings. Each keeps the newest n entries and nothing older.
|
||||
|
||||
func tailFloat32(s []float32, n int) []float32 {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return s[len(s)-n:]
|
||||
}
|
||||
|
||||
func tailMel(s [][melBins]float32, n int) [][melBins]float32 {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return append(s[:0], s[len(s)-n:]...)
|
||||
}
|
||||
|
||||
func tailEmbed(s [][embedDim]float32, n int) [][embedDim]float32 {
|
||||
if len(s) <= n {
|
||||
return s
|
||||
}
|
||||
return append(s[:0], s[len(s)-n:]...)
|
||||
}
|
||||
@@ -0,0 +1,174 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeGate fires on demand instead of running three ONNX models. The gate's
|
||||
// own arithmetic is measured on real audio in docs/evals; what these tests
|
||||
// cover is the thing that decides whether an utterance is shipped.
|
||||
type fakeGate struct {
|
||||
fireOn int // fire when this many frames have been fed, 0 never fires
|
||||
fed int
|
||||
resets int
|
||||
}
|
||||
|
||||
func (g *fakeGate) Feed(_ []int16) bool {
|
||||
g.fed++
|
||||
return g.fireOn > 0 && g.fed == g.fireOn
|
||||
}
|
||||
func (g *fakeGate) Reset() { g.resets++ }
|
||||
func (g *fakeGate) Score() float64 { return 1 }
|
||||
|
||||
// wakingSession wires a session whose gate fires on the first frame it sees.
|
||||
func wakingSession(fireOn int, window time.Duration) (*session, *fakePlayer, *fakeSender, *fakeGate) {
|
||||
sess, p, snd := newTestSession(bargeInConfig{})
|
||||
g := &fakeGate{fireOn: fireOn}
|
||||
sess.UseWakeWord(g, window)
|
||||
return sess, p, snd, g
|
||||
}
|
||||
|
||||
func TestKeywordlessSpeechNeverReachesSTT(t *testing.T) {
|
||||
sess, p, snd, g := wakingSession(0, 8*time.Second)
|
||||
speakThenPause(t, sess)
|
||||
|
||||
if len(snd.sent) != 0 {
|
||||
t.Fatalf("sent %d utterances, want 0 — this is the whole point of V-487", len(snd.sent))
|
||||
}
|
||||
if sess.ignored != 1 {
|
||||
t.Errorf("ignored = %d, want 1", sess.ignored)
|
||||
}
|
||||
if p.plays != 0 {
|
||||
t.Errorf("plays = %d, want 0", p.plays)
|
||||
}
|
||||
if g.fed == 0 {
|
||||
t.Error("the gate was never fed a frame")
|
||||
}
|
||||
}
|
||||
|
||||
func TestKeywordOpensTheGate(t *testing.T) {
|
||||
sess, p, snd, _ := wakingSession(1, 8*time.Second)
|
||||
speakThenPause(t, sess)
|
||||
|
||||
if len(snd.sent) != 1 {
|
||||
t.Fatalf("sent %d utterances, want 1", len(snd.sent))
|
||||
}
|
||||
if sess.wakes != 1 {
|
||||
t.Errorf("wakes = %d, want 1", sess.wakes)
|
||||
}
|
||||
if sess.ignored != 0 {
|
||||
t.Errorf("ignored = %d, want 0", sess.ignored)
|
||||
}
|
||||
if p.plays != 1 {
|
||||
t.Errorf("plays = %d, want 1", p.plays)
|
||||
}
|
||||
}
|
||||
|
||||
// One keyword buys one turn. Without this the microphone stays open for as
|
||||
// long as he keeps talking, which is the state the gate exists to end.
|
||||
func TestOneKeywordBuysOneTurn(t *testing.T) {
|
||||
sess, p, snd, _ := wakingSession(1, 8*time.Second)
|
||||
speakThenPause(t, sess)
|
||||
p.Stop() // she finished her reply
|
||||
sess.discard = 0 // the backlog drain is not what this measures
|
||||
speakThenPause(t, sess)
|
||||
|
||||
if len(snd.sent) != 1 {
|
||||
t.Fatalf("sent %d utterances, want 1: the second had no keyword", len(snd.sent))
|
||||
}
|
||||
if sess.ignored != 1 {
|
||||
t.Errorf("ignored = %d, want 1", sess.ignored)
|
||||
}
|
||||
}
|
||||
|
||||
// The keyword is heard, then he says nothing for longer than the window. What
|
||||
// he says after that is not addressed to her.
|
||||
func TestTheKeywordExpires(t *testing.T) {
|
||||
sess, _, snd, _ := wakingSession(1, 500*time.Millisecond)
|
||||
now := time.Unix(1750000000, 0)
|
||||
sess.now = func() time.Time { return now }
|
||||
|
||||
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
if sess.wakes != 1 {
|
||||
t.Fatalf("wakes = %d, want 1", sess.wakes)
|
||||
}
|
||||
now = now.Add(2 * time.Second)
|
||||
speakThenPause(t, sess)
|
||||
|
||||
if len(snd.sent) != 0 {
|
||||
t.Fatalf("sent %d utterances, want 0 — the keyword had expired", len(snd.sent))
|
||||
}
|
||||
}
|
||||
|
||||
// Barge-in cuts her off whether or not the keyword was heard. What he says
|
||||
// after cutting her off still has to carry it.
|
||||
func TestBargeInStillInterruptsHer(t *testing.T) {
|
||||
sess, p, _, g := wakingSession(0, 8*time.Second)
|
||||
sess.barge = bargeInConfig{RMS: 0.2, Frames: 3}
|
||||
p.playing = true
|
||||
loud := frameAt(0.35)
|
||||
for i := 0; i < 4; i++ {
|
||||
if err := sess.feed(context.Background(), loud); err != nil {
|
||||
t.Fatalf("feed %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
if sess.bargeIns != 1 {
|
||||
t.Fatalf("bargeIns = %d, want 1", sess.bargeIns)
|
||||
}
|
||||
if p.stops != 1 {
|
||||
t.Errorf("stops = %d, want 1", p.stops)
|
||||
}
|
||||
if g.resets == 0 {
|
||||
t.Error("barge-in left pre-playback audio in the gate")
|
||||
}
|
||||
}
|
||||
|
||||
// Her own reply must not wake her. Frames captured while the player runs never
|
||||
// reach the gate, and the gate is cleared when playback ends.
|
||||
func TestHerOwnVoiceNeverReachesTheGate(t *testing.T) {
|
||||
sess, p, _, g := wakingSession(1, 8*time.Second)
|
||||
p.playing = true
|
||||
for i := 0; i < 10; i++ {
|
||||
if err := sess.feed(context.Background(), frameAt(0.35)); err != nil {
|
||||
t.Fatalf("feed: %v", err)
|
||||
}
|
||||
}
|
||||
if g.fed != 0 {
|
||||
t.Fatalf("gate was fed %d frames while she was speaking, want 0", g.fed)
|
||||
}
|
||||
if sess.wakes != 0 {
|
||||
t.Errorf("wakes = %d, want 0", sess.wakes)
|
||||
}
|
||||
}
|
||||
|
||||
// No model, no gate: the daemon behaves exactly as it did before V-487 stage
|
||||
// two. An operator with a missing file gets yesterday's mavwaked, not one that
|
||||
// refuses to hear anything.
|
||||
func TestNoGateShipsEveryUtterance(t *testing.T) {
|
||||
sess, _, snd := newTestSession(bargeInConfig{})
|
||||
speakThenPause(t, sess)
|
||||
|
||||
if len(snd.sent) != 1 {
|
||||
t.Fatalf("sent %d utterances, want 1", len(snd.sent))
|
||||
}
|
||||
if sess.ignored != 0 {
|
||||
t.Errorf("ignored = %d, want 0", sess.ignored)
|
||||
}
|
||||
}
|
||||
|
||||
// A nil *wakeWord is the closed gate, not a crash and not an open one.
|
||||
func TestNilWakeWordHearsNothing(t *testing.T) {
|
||||
var w *wakeWord
|
||||
if w.Feed([]int16{0, 0, 0}) {
|
||||
t.Error("a nil wake word reported the keyword")
|
||||
}
|
||||
if w.Score() != 0 {
|
||||
t.Error("a nil wake word reported a score")
|
||||
}
|
||||
w.Reset()
|
||||
w.Close()
|
||||
}
|
||||
+37
-14
@@ -4,11 +4,22 @@
|
||||
# user unit because it needs his ALSA session and his ssh agent, and because
|
||||
# it should stop when he logs out.
|
||||
#
|
||||
# THERE IS NO WAKE WORD YET (V-487 stage two). Anything spoken near the fifine
|
||||
# becomes a turn. What makes that safe rather than expensive is voiceSender:
|
||||
# it sends Surface=SurfaceVoice, which caps every command at L0, so no
|
||||
# accidental trigger runs a destructive act. It does not stop her answering
|
||||
# out loud, so this unit is his to stop when the room is not his alone.
|
||||
# The keyword is "Мэйвен" and the three -wake- flags are what require it
|
||||
# (V-487 stage two). Without them anything spoken near the fifine becomes a
|
||||
# turn, which voiceSender makes safe rather than expensive: it sends
|
||||
# Surface=SurfaceVoice, capping every command at L0. That does not stop her
|
||||
# answering out loud, which is the whole reason the keyword exists.
|
||||
#
|
||||
# The threshold is 0.999 and it is the binary's default, so it is not passed.
|
||||
# It came from 65 minutes of held-out Russian speech through this same binary:
|
||||
# 0.9 false wakes an hour against 2.8 at 0.99, for one lost render out of 126
|
||||
# (docs/evals/2026-08-09-wake-word.md). If the room proves noisier than the
|
||||
# corpus, read the scores out of this unit's journal and pass -wake-threshold.
|
||||
# Do not lower it by guessing.
|
||||
#
|
||||
# A keyword shorter than 1.32s can be heard too late to be used, because the
|
||||
# head scores 1.28s of audio and the VAD has closed the utterance by then.
|
||||
# "Мэйвен, <request>" is unaffected. A bare "Мэйвен" is the case that fails.
|
||||
#
|
||||
# -vad-model is passed on purpose. Silero answers "is this frame speech" where
|
||||
# the energy floor answers "is this frame loud". It declines white noise at
|
||||
@@ -24,27 +35,39 @@
|
||||
# systemctl --user enable --now mavwaked.service
|
||||
|
||||
[Unit]
|
||||
Description=Maven always-on listening (VAD, no wake word yet)
|
||||
Description=Maven always-on listening (silero VAD, "Мэйвен" keyword)
|
||||
# The tunnel is the only path to mavend and the only thing authenticating it.
|
||||
Requires=maven-voice-tunnel.service
|
||||
After=maven-voice-tunnel.service
|
||||
|
||||
[Service]
|
||||
# card 0 is the fifine USB microphone. Named, and not "default", because the
|
||||
# default device follows whatever pipewire last decided and this daemon should
|
||||
# not change ears when he plugs in a headset.
|
||||
# The Scarlett Solo 4th Gen, and not the fifine. The fifine was the device
|
||||
# here for three days and mavwaked never logged one utterance in them, because
|
||||
# it returns RMS 0.00004 with its capture switch on and its ALSA volume at the
|
||||
# full 496 of 496. That silence is in the hardware, so no flag reaches it.
|
||||
#
|
||||
# Named CARD=Gen and not card 4, because a USB card number moves when
|
||||
# something else is replugged and this daemon must not change ears quietly.
|
||||
# Not "default" either: that follows whatever pipewire last decided.
|
||||
#
|
||||
# plughw and not hw. mavwaked asks arecord for 16kHz mono, which is what the
|
||||
# whole pipeline is canonical in. The fifine offers 2 channels at 44100 or
|
||||
# 48000 and nothing else, so bare hw:0,0 dies on "Channels count non
|
||||
# available" before a frame is read. plughw puts ALSA's downmix and resampler
|
||||
# in front. Any replacement microphone wants the same treatment.
|
||||
# whole pipeline is canonical in. Neither microphone offers it, so bare hw
|
||||
# dies on "Channels count non available" before a frame is read. plughw puts
|
||||
# ALSA's downmix and resampler in front. Any replacement wants the same.
|
||||
#
|
||||
# The Scarlett measured RMS 0.003 against 0.14 on the onboard input, so its
|
||||
# front-panel gain is the thing to raise if she mishears. That is a knob, not
|
||||
# a control ALSA exposes. The two loud devices, the onboard ALC897 and the
|
||||
# camera, both clip at peak 1.0 and are worse candidates, not better ones.
|
||||
Environment=LD_LIBRARY_PATH=%h/.local/lib
|
||||
ExecStart=%h/.local/bin/mavwaked \
|
||||
-device plughw:0,0 \
|
||||
-device plughw:CARD=Gen,DEV=0 \
|
||||
-addr 127.0.0.1:9100 \
|
||||
-lang ru \
|
||||
-vad-model %h/.local/share/maven/models/silero_vad.onnx \
|
||||
-wake-model %h/.local/share/maven/models/maven_wakeword.onnx \
|
||||
-wake-mel %h/.local/share/maven/models/melspectrogram.onnx \
|
||||
-wake-embed %h/.local/share/maven/models/embedding_model.onnx \
|
||||
-onnx-lib %h/.local/lib/libonnxruntime.so
|
||||
Restart=on-failure
|
||||
RestartSec=5
|
||||
|
||||
+19
-3
@@ -170,9 +170,25 @@ The device is `plughw:0,0` and not `hw:0,0`. The fifine offers 2 channels at
|
||||
44100 or 48000 and nothing else, and mavwaked asks arecord for 16kHz mono. Bare
|
||||
`hw` dies on "Channels count non available" before a frame is read.
|
||||
|
||||
There is no wake word yet (V-487 stage two), so the loop runs open. mavwaked
|
||||
connects lazily, so `voicesink` cannot push a nudge to it until it has sent one
|
||||
utterance.
|
||||
The three `-wake-` flags require the keyword "Мэйвен" (V-487 stage two). Drop
|
||||
them and the loop runs open, which is what it did before. The threshold is the
|
||||
binary's default of 0.999 and is not passed. Over 65 minutes of held-out
|
||||
Russian speech it woke her 0.9 times an hour against 2.8 at 0.99. That cost one
|
||||
lost render out of 126 (`docs/evals/2026-08-09-wake-word.md`).
|
||||
|
||||
The three sessions are pinned to one thread each. onnxruntime otherwise sizes
|
||||
its pool to every core and spins between runs, which took mavwaked from 68% of
|
||||
one core to 335%. With the cap it sits at 81%, so the gate costs about 13%.
|
||||
|
||||
A keyword shorter than 1.32s can be heard too late to be used. The head scores
|
||||
1.28s of audio, and the VAD has closed the utterance by then.
|
||||
"Мэйвен, <request>" is unaffected. A bare "Мэйвен" is the case that fails.
|
||||
|
||||
mavwaked connects at startup and holds the conn, so a nudge routed to voice
|
||||
reaches the speaker before he has said anything (V-671). It used to connect
|
||||
lazily, which made the failure silent rather than absent: after one utterance
|
||||
the session existed, `PushToMostRecent` succeeded, the dispatcher stopped
|
||||
rerouting to telegram and ntfy, and mavwaked discarded the audio.
|
||||
|
||||
**Passwords are read from files, never taken as flag values.** `mavcaldav` uses
|
||||
`-pass-file` and `-render-pass-file`. `mavpoll` and `mavmaild` follow the same
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
# The "Мэйвен" wake word: what it hears and what it invents
|
||||
|
||||
*Measured 2026-08-09 on workpc and homesrv. V-487, stage two of two.*
|
||||
|
||||
Stage one gave mavwaked silero-vad, which answers "is this frame speech".
|
||||
Nothing answered "was this said to her", so every utterance near the
|
||||
microphone became a turn. SurfaceVoice caps acts at L0, which made that safe
|
||||
rather than expensive. L0 does not cap reading, so the room could still hear
|
||||
his facts read back.
|
||||
|
||||
The keyword is "Мэйвен". openWakeWord's two frozen feature models do the
|
||||
hearing and a 100KB head trained here draws the boundary. It runs on CPU
|
||||
beside silero and never touches the GPU.
|
||||
|
||||
## Why a per-window accuracy is not a number anyone can act on
|
||||
|
||||
The gate scores every 80ms. A 1.7% false-accept rate per window sounds small
|
||||
and means a wake every few seconds. The useful question is how many times an hour
|
||||
it wakes on speech that was not the keyword. So every table below counts
|
||||
threshold crossings over whole clips and divides by the audio duration.
|
||||
|
||||
A crossing, not a window above the threshold. A keyword held high for half a
|
||||
second is one wake, not six.
|
||||
|
||||
## The data
|
||||
|
||||
Positives are 600 silero TTS renders of three stressings of the keyword, six
|
||||
speakers, ten trailing phrases, augmented eight ways each. Hard negatives are
|
||||
560 renders of confusable Russian words. Real speech is Common Voice ru and Golos.
|
||||
The 74257 Common Voice clips were already on workpc from the CrisperWhisper
|
||||
work. The 200 Golos clips came from the CW2 WER eval.
|
||||
|
||||
Splits are by source file. Augmented copies of one render on both sides of a
|
||||
split would measure memorisation.
|
||||
|
||||
Golos was never trained on at any stage, so it answers the harder question:
|
||||
does this survive a change of speakers and rooms.
|
||||
|
||||
## Three heads
|
||||
|
||||
Each row is a full retrain. The false-accept column is 8.89 hours of Common
|
||||
Voice that no stage of training had seen.
|
||||
|
||||
| trained on | recall (window) | false wakes/hour @0.99 |
|
||||
|---|---|---|
|
||||
| TTS + 13.7 min of Golos | 0.869 | not measurable |
|
||||
| + 4000 Common Voice clips | 0.836 | 21.9 |
|
||||
| + 3837 mined hard negatives | 0.784 | 4.2 |
|
||||
| + 753 more mined | 0.810 | 3.4 |
|
||||
|
||||
The first row is why the second exists. Thirteen minutes of held-out speech
|
||||
cannot measure a rate for a gate that scores twelve times a second. A head
|
||||
trained only against TTS learns to tell TTS from not-TTS.
|
||||
|
||||
Mining is the whole story after that. Random negatives teach the head what
|
||||
most speech sounds like. They do not teach it the few syllable sequences that
|
||||
score high, because 4000 clips barely contain them. So the current head was
|
||||
run over 20000 fresh clips, keeping every window it scored above 0.05. That
|
||||
found 3837 windows in 855512. Repeating those ten times in the next training
|
||||
run cut the rate five-fold.
|
||||
|
||||
The second round found 753 in 852240, a fifth of the yield, and bought a
|
||||
further 20%. It also recovered recall, which the first round had cost. Whether
|
||||
a third round is worth 25 minutes of workpc is untested.
|
||||
|
||||
## Where the threshold came from
|
||||
|
||||
Both columns are held out. Positives are the 126 renders in the test split.
|
||||
Speech is 65.1 minutes of Common Voice, disjoint from every training and
|
||||
mining pool. Both were run through the built `mavwaked` binary reading PCM from a
|
||||
file, not through the python that trained the head.
|
||||
|
||||
| threshold | renders shipped | false wakes/hour |
|
||||
|---|---|---|
|
||||
| 0.99 | 116 / 126 | 2.8 |
|
||||
| 0.999 | 115 / 126 | 0.9 |
|
||||
|
||||
One render against a third of the false wakes. `defaultWakeThreshold` is
|
||||
0.999.
|
||||
|
||||
Golos disagrees. It gave 2 wakes in 14 minutes at every threshold, which is
|
||||
8.7 per hour. Two events is not a rate. What it does say is that a handful of real utterances score above 0.999
|
||||
and no threshold will move them.
|
||||
|
||||
## What it costs him
|
||||
|
||||
Ten of the 126 held-out renders were heard and still dropped, and every one
|
||||
was an utterance shorter than 1.32s. The head scores 16 embeddings, or 1.28s of
|
||||
audio. The score therefore peaks up to a second after a short keyword ends.
|
||||
By then the VAD has closed the utterance and dispatch has already asked.
|
||||
|
||||
Real commands are "Мэйвен, <request>" and run past two seconds, which gives
|
||||
the head the whole request to peak during. A bare "Мэйвен" with nothing after
|
||||
it is the case that fails. One fix would hold an ignored utterance for a grace
|
||||
period and ship it if the keyword lands late. It is not built.
|
||||
|
||||
## What it costs the workstation
|
||||
|
||||
Under systemd on workpc, mavwaked sat at 335% of a core with the gate on and
|
||||
68% with only silero. onnxruntime sizes its thread pool to every core and spins
|
||||
between runs, and this gate runs three graphs twelve times a second. Pinning
|
||||
all three sessions to one thread brought it to 81%, so the keyword costs about
|
||||
13% of one core. The three graphs each finish in well under 80ms that way.
|
||||
|
||||
## What was not measured
|
||||
|
||||
No room recordings. Every negative above is a clean corpus clip. This gate
|
||||
will live among a television, a fan and the far side of a kitchen. None of
|
||||
those are in these numbers.
|
||||
|
||||
No measurement of him. Training on his voice means copying his transcripts off
|
||||
homesrv, which is his call and has not been asked.
|
||||
+118
-65
@@ -9,15 +9,18 @@
|
||||
//
|
||||
// The wire is symmetric: a Request from the client is answered by a
|
||||
// Response with a matching ID, OR a server-initiated Push frame (no ID)
|
||||
// may arrive interleaved. SendRequest loops reading frames, drops Push
|
||||
// frames to the harness if a receiver is running (or silently if not),
|
||||
// and returns the first Response with the matching ID.
|
||||
// may arrive interleaved. One reader goroutine per connection owns the
|
||||
// socket. It hands each Response to whichever SendRequest is waiting on
|
||||
// that ID and each Push to the handler, so a client may send and listen
|
||||
// at the same time on one conn. mavwaked needs exactly that: it speaks
|
||||
// utterances and it must hear nudges, and the server routes a nudge to
|
||||
// the session that spoke most recently, so a second listening conn would
|
||||
// never be picked (V-671).
|
||||
package voice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
@@ -41,13 +44,22 @@ type PushHandler interface {
|
||||
// Client — one connection to the voice.Server.
|
||||
type Client struct {
|
||||
addr string
|
||||
mu sync.Mutex
|
||||
c net.Conn
|
||||
nextID atomic.Uint64
|
||||
|
||||
// pushCh fan-out: a reader goroutine (started by RunPushReceiver)
|
||||
// writes Push frames here; SendRequest also drains it when no reader
|
||||
// is running (drops the frame in that case).
|
||||
mu sync.Mutex
|
||||
c net.Conn
|
||||
// pending holds one channel per in-flight request, keyed by frame id.
|
||||
// The reader goroutine delivers the Response here and deletes the entry.
|
||||
pending map[uint64]chan *Response
|
||||
// dead is closed by the reader goroutine when this conn ends, so a
|
||||
// waiting SendRequest fails at once instead of at its own deadline.
|
||||
dead chan struct{}
|
||||
|
||||
// wmu serialises writes. Frames must not interleave on the wire.
|
||||
wmu sync.Mutex
|
||||
|
||||
// pushH is set by RunPushReceiver and survives a reconnect, because the
|
||||
// client that wants pushes wants them on whatever conn it ends up with.
|
||||
pushMu sync.Mutex
|
||||
pushH PushHandler
|
||||
}
|
||||
@@ -55,6 +67,15 @@ type Client struct {
|
||||
// Dial returns a Client that will connect to addr on first use.
|
||||
func Dial(addr string) *Client { return &Client{addr: addr} }
|
||||
|
||||
// Connect opens the conn now rather than on the first request. mavwaked calls
|
||||
// it at startup: the server registers a session on accept, and a client that
|
||||
// has never connected cannot be sent a nudge.
|
||||
func (c *Client) Connect(ctx context.Context) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
return c.ensureConnLocked(ctx)
|
||||
}
|
||||
|
||||
// Close releases the conn. Idempotent.
|
||||
func (c *Client) Close() error {
|
||||
c.mu.Lock()
|
||||
@@ -75,12 +96,14 @@ func (c *Client) PushToTalk(ctx context.Context, a audio.Audio, lang string) (Pu
|
||||
return out, err
|
||||
}
|
||||
|
||||
// requestTimeout bounds a round-trip with no deadline on its context. It is
|
||||
// generous because the far end runs speech-to-text, a router and a voice.
|
||||
const requestTimeout = 120 * time.Second
|
||||
|
||||
// SendRequest sends one Request frame and waits for the matching Response.
|
||||
// Push frames received while waiting are dropped on the floor UNLESS a
|
||||
// PushHandler has been wired via RunPushReceiver, in which case the handler
|
||||
// is invoked inline (still synchronous with the SendRequest caller's
|
||||
// read). For sanity, the reference client runs either one-shot (no
|
||||
// receiver) or interactive (RunPushReceiver, no concurrent SendRequest).
|
||||
// Push frames arriving meanwhile go to the handler on the reader goroutine,
|
||||
// so listening and sending on one Client is supported rather than merely
|
||||
// tolerated.
|
||||
func (c *Client) SendRequest(ctx context.Context, m Method, params any, out any) error {
|
||||
body, err := marshalParams(params)
|
||||
if err != nil {
|
||||
@@ -94,58 +117,69 @@ func (c *Client) SendRequest(ctx context.Context, m Method, params any, out any)
|
||||
c.mu.Unlock()
|
||||
return err
|
||||
}
|
||||
conn := c.c
|
||||
conn, dead := c.c, c.dead
|
||||
ch := make(chan *Response, 1)
|
||||
c.pending[id] = ch
|
||||
c.mu.Unlock()
|
||||
|
||||
if dl, ok := ctx.Deadline(); ok {
|
||||
_ = conn.SetDeadline(dl)
|
||||
} else {
|
||||
_ = conn.SetDeadline(time.Now().Add(120 * time.Second))
|
||||
}
|
||||
defer conn.SetDeadline(time.Time{})
|
||||
|
||||
if err := writeFrame(conn, &req); err != nil {
|
||||
c.teardown()
|
||||
c.wmu.Lock()
|
||||
err = writeFrame(conn, &req)
|
||||
c.wmu.Unlock()
|
||||
if err != nil {
|
||||
c.forget(id)
|
||||
c.teardownConn(conn)
|
||||
return err
|
||||
}
|
||||
for {
|
||||
resp, push, err := readOneFrame(conn)
|
||||
if err != nil {
|
||||
c.teardown()
|
||||
return err
|
||||
}
|
||||
if push != nil {
|
||||
c.deliverPush(*push)
|
||||
continue
|
||||
}
|
||||
if resp.ID != id {
|
||||
continue // not ours; ignore (singleplex ⇒ shouldn't happen)
|
||||
}
|
||||
if resp.Error != nil {
|
||||
return hydrate(resp.Error)
|
||||
}
|
||||
if out != nil {
|
||||
if err := json.Unmarshal(resp.Result, out); err != nil {
|
||||
return fmt.Errorf("voice: unmarshal result: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
|
||||
timer := time.NewTimer(requestTimeout)
|
||||
defer timer.Stop()
|
||||
|
||||
var resp *Response
|
||||
select {
|
||||
case resp = <-ch:
|
||||
case <-dead:
|
||||
c.forget(id)
|
||||
return fmt.Errorf("voice: connection closed before reply")
|
||||
case <-ctx.Done():
|
||||
c.forget(id)
|
||||
return ctx.Err()
|
||||
case <-timer.C:
|
||||
c.forget(id)
|
||||
c.teardownConn(conn)
|
||||
return fmt.Errorf("voice: no reply within %s", requestTimeout)
|
||||
}
|
||||
|
||||
if resp.Error != nil {
|
||||
return hydrate(resp.Error)
|
||||
}
|
||||
if out != nil {
|
||||
if err := json.Unmarshal(resp.Result, out); err != nil {
|
||||
return fmt.Errorf("voice: unmarshal result: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// RunPushReceiver spawns a reader goroutine that delivers Push frames to h
|
||||
// until the conn closes or Close is called. Today's reference client uses
|
||||
// this in -listen mode (proactive voice playback). SendRequest and
|
||||
// RunPushReceiver SHOULD NOT be used concurrently on the same Client — the
|
||||
// wire is singleplex at the reference client's scale; production picks one
|
||||
// mode per conn. Returns when the goroutine ends (ctx cancel or conn close).
|
||||
// forget drops an abandoned request so a late Response is discarded rather
|
||||
// than delivered to nobody.
|
||||
func (c *Client) forget(id uint64) {
|
||||
c.mu.Lock()
|
||||
delete(c.pending, id)
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// RunPushReceiver wires h and blocks until the conn ends or ctx is
|
||||
// cancelled. Frames are read by the per-conn reader goroutine, so a client
|
||||
// may call SendRequest on the same Client while this is running. Returns nil
|
||||
// when the conn ended, so a caller that wants to stay reachable reconnects
|
||||
// and calls it again.
|
||||
func (c *Client) RunPushReceiver(ctx context.Context, h PushHandler) error {
|
||||
c.mu.Lock()
|
||||
if err := c.ensureConnLocked(ctx); err != nil {
|
||||
c.mu.Unlock()
|
||||
return err
|
||||
}
|
||||
conn := c.c
|
||||
dead := c.dead
|
||||
c.mu.Unlock()
|
||||
|
||||
c.pushMu.Lock()
|
||||
@@ -158,21 +192,34 @@ func (c *Client) RunPushReceiver(ctx context.Context, h PushHandler) error {
|
||||
c.pushMu.Unlock()
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-dead:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// readLoop owns conn for its whole life. It ends on any read error, which is
|
||||
// how a closed conn, a killed server and a cancelled dial all arrive here.
|
||||
func (c *Client) readLoop(conn net.Conn, dead chan struct{}) {
|
||||
defer close(dead)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
default:
|
||||
}
|
||||
_, push, err := readOneFrame(conn)
|
||||
resp, push, err := readOneFrame(conn)
|
||||
if err != nil {
|
||||
if errors.Is(err, io.EOF) || errors.Is(err, net.ErrClosed) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
c.teardownConn(conn)
|
||||
return
|
||||
}
|
||||
if push != nil {
|
||||
c.deliverPush(*push)
|
||||
continue
|
||||
}
|
||||
c.mu.Lock()
|
||||
ch := c.pending[resp.ID]
|
||||
delete(c.pending, resp.ID)
|
||||
c.mu.Unlock()
|
||||
if ch != nil {
|
||||
ch <- resp
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -196,13 +243,19 @@ func (c *Client) ensureConnLocked(ctx context.Context) error {
|
||||
return fmt.Errorf("voice: dial %s: %w", c.addr, err)
|
||||
}
|
||||
c.c = conn
|
||||
c.pending = make(map[uint64]chan *Response)
|
||||
c.dead = make(chan struct{})
|
||||
go c.readLoop(conn, c.dead)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Client) teardown() {
|
||||
// teardownConn closes conn and forgets it, but only if it is still the live
|
||||
// one. A reconnect may already have replaced it, and closing the new conn
|
||||
// because the old one died takes the client down on every hiccup.
|
||||
func (c *Client) teardownConn(conn net.Conn) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if c.c != nil {
|
||||
if c.c != nil && c.c == conn {
|
||||
_ = c.c.Close()
|
||||
c.c = nil
|
||||
}
|
||||
|
||||
@@ -221,3 +221,113 @@ func TestClientListenModeReceivesPush(t *testing.T) {
|
||||
type pushHandlerFunc func(Push)
|
||||
|
||||
func (f pushHandlerFunc) OnPush(p Push) { f(p) }
|
||||
|
||||
// The whole point of the per-conn reader (V-671): mavwaked speaks utterances
|
||||
// and must hear nudges, and the server routes a nudge to the session that
|
||||
// spoke most recently. A second listening conn would never be picked, so both
|
||||
// directions have to share one conn.
|
||||
func TestClientSendsAndListensOnOneConn(t *testing.T) {
|
||||
l := newTestListener(t)
|
||||
sess := NewSessions()
|
||||
h := &stubHandler{}
|
||||
srv := NewServer(l.Addr().String(), h, sess)
|
||||
if err := srv.Listen(); err != nil {
|
||||
t.Fatalf("listen: %v", err)
|
||||
}
|
||||
defer func() {
|
||||
_ = srv.Close()
|
||||
waitPort()
|
||||
}()
|
||||
go func() { _ = srv.Serve() }()
|
||||
|
||||
c := Dial(l.Addr().String())
|
||||
defer c.Close()
|
||||
|
||||
got := make(chan audio.Audio, 4)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go func() {
|
||||
_ = c.RunPushReceiver(ctx, pushHandlerFunc(func(p Push) {
|
||||
var ap AudioNudgePush
|
||||
if err := json.Unmarshal(p.Params, &ap); err == nil {
|
||||
got <- ap.Audio
|
||||
}
|
||||
}))
|
||||
}()
|
||||
|
||||
for i := 0; i < 100 && sess.Active() < 1; i++ {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
if sess.Active() != 1 {
|
||||
t.Fatalf("active sessions = %d, want exactly 1", sess.Active())
|
||||
}
|
||||
|
||||
// A round-trip while the receiver is running. Before the reader owned the
|
||||
// conn, this and the receiver raced for every frame.
|
||||
resp, err := c.PushToTalk(context.Background(), audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("hello")}, "ru")
|
||||
if err != nil {
|
||||
t.Fatalf("PushToTalk with a receiver running: %v", err)
|
||||
}
|
||||
if resp.ReplyText != "got it" {
|
||||
t.Fatalf("ReplyText = %q, want %q", resp.ReplyText, "got it")
|
||||
}
|
||||
|
||||
// And the nudge still arrives, on the session that just spoke.
|
||||
err = sess.PushToMostRecent(context.Background(), AudioNudgePush{
|
||||
RuleName: "after-speaking",
|
||||
Audio: audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("proactive")},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("PushToMostRecent: %v", err)
|
||||
}
|
||||
select {
|
||||
case a := <-got:
|
||||
if string(a.Bytes) != "proactive" {
|
||||
t.Fatalf("received %q, want the nudge audio", string(a.Bytes))
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("nudge never reached the handler after the client had spoken")
|
||||
}
|
||||
}
|
||||
|
||||
// A request abandoned by its context must not leave its slot behind, or a
|
||||
// long-running client leaks one channel per timeout.
|
||||
func TestClientForgetsAbandonedRequests(t *testing.T) {
|
||||
l := newTestListener(t)
|
||||
sess := NewSessions()
|
||||
srv := NewServer(l.Addr().String(), &blockingHandler{}, sess)
|
||||
if err := srv.Listen(); err != nil {
|
||||
t.Fatalf("listen: %v", err)
|
||||
}
|
||||
defer func() {
|
||||
_ = srv.Close()
|
||||
waitPort()
|
||||
}()
|
||||
go func() { _ = srv.Serve() }()
|
||||
|
||||
c := Dial(l.Addr().String())
|
||||
defer c.Close()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
_, err := c.PushToTalk(ctx, audio.Audio{Format: audio.PCM16kMono, Bytes: []byte("x")}, "ru")
|
||||
if err == nil {
|
||||
t.Fatal("expected the round-trip to fail on its context")
|
||||
}
|
||||
|
||||
c.mu.Lock()
|
||||
n := len(c.pending)
|
||||
c.mu.Unlock()
|
||||
if n != 0 {
|
||||
t.Fatalf("pending = %d after an abandoned request, want 0", n)
|
||||
}
|
||||
}
|
||||
|
||||
// blockingHandler never answers, so the client's context is what ends the
|
||||
// round-trip.
|
||||
type blockingHandler struct{}
|
||||
|
||||
func (blockingHandler) HandlePushToTalk(ctx context.Context, _ PushToTalkReq, _ uint64) (PushToTalkResp, error) {
|
||||
<-ctx.Done()
|
||||
return PushToTalkResp{}, ctx.Err()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user