Compare commits

...

9 Commits

Author SHA1 Message Date
claude ce6a6821a9 Test the gate without three ONNX files (V-487)
keywordGate is an interface so the decision that ships an utterance can be
exercised with a fake that fires on demand. A gate that can only be tested
with a model file is a gate nobody tests.

The three that carry the fixed-when criterion: keywordless speech never
reaches STT, the keyword does, and barge-in still cuts her off mid-sentence.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 14:19:01 +04:00
claude 479b0c4475 Speech without the keyword no longer reaches STT (V-487)
Until now every utterance near the microphone became a turn. SurfaceVoice caps
acts at L0, which made that safe rather than expensive, but L0 does not cap
reading: the room could still hear his facts read back.

The gate sits at dispatch, not at the VAD. The keyword opens a window, the VAD
closes the utterance when he stops, and dispatch asks whether the window was
open. That ordering is what lets him say "Мэйвен" and then a sentence: the
window has to outlive the word by the length of what follows it.

One keyword buys 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 exists to end.

Her own voice cannot wake her. Every path above the gate returns while the
player is running, so no frame of her reply is ever scored, and the streaming
state is cleared when playback ends.

Nil is a working value. Without -wake-model the gate is open and this is
yesterday's mavwaked, which is what an operator with a missing file should get
rather than a daemon that refuses to listen.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 14:19:01 +04:00
claude 877b1fd4f8 Score the keyword every 80ms without re-reading old audio (V-487)
melContext is 480 because melspectrogram.onnx returns N/160-3 frames and frame
i covers [i*160, i*160+400). With 480 samples of history the buffer is 8
frames and the oldest continues exactly one hop after the previous call's
newest. Less history leaves a gap.

Feed reports the threshold 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.

Nil is the CLOSED gate rather than the open one. A nil that answers "yes,
keyword" reads as a working wake word in every log line it produces.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 14:18:46 +04:00
claude 21a42cb3e6 Load openWakeWord's three models and run their tensors (V-487)
The two feature models are frozen and pretrained; only the 100KB head was
trained here. The shapes were measured rather than assumed: 2.0s of 16kHz
audio gives 197 mel frames, and 76-frame windows at stride 8 give exactly the
16 embeddings the head was fitted on.

This file knows tensors and nothing about the 80ms cadence, which is why the
scaling openWakeWord applies between the two feature models lives here.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 14:18:46 +04:00
claude b8279f6a22 Merge pull request 'mavwaked registers as a voice consumer it cannot honor, so a spoken turn silences every nudge' (#218) from task/671-mavwaked-registers-as-a-voice-consumer-i into master 2026-08-09 11:45:35 +02:00
claude 8c30971a96 Say that mavwaked now holds the conn from startup (V-671)
The lazy-connect note is no longer true and the trap it described was the
opposite way round: the session existed and the audio was discarded.

diff-budget.sh blocks the branch at 615 changed lines. This commit is
markdown only, which the repo's own pre-commit hook exempts, and it
corrects a line the code in this branch has just falsified.
2026-08-09 13:45:20 +04:00
claude d0ea927ac3 Pin the five things a nudge must do at the speaker (V-671)
It reaches the player, but not from the push goroutine. An unusable push
is dropped and does not wedge the next one. It waits for a reply to
finish. It resets the VAD, so the frames before it are not spliced onto
what he says after. And a second nudge replaces an unspoken first.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 13:45:04 +04:00
claude 9c7bafd5b1 Let mavwaked hear the nudges it was already being sent (V-671)
It wired no PushHandler, and SendRequest discards a push frame when there
is none. That was not a missing feature but a silent one. mavend routes a
nudge to the voice session that spoke most recently, so once mavwaked had
spoken once it WAS that session. PushToMostRecent succeeded, the
dispatcher counted the nudge delivered and stopped rerouting to the away
channels, and mavwaked threw the audio away. He heard nothing, anywhere.

It now connects at startup rather than at the first utterance, because
the dispatcher has to tell "he is not at the machine" from "he is, and
she has nothing to say". The receiver redials on its own clock, since
mavend restarts on every deploy.

A nudge is queued, not played where it arrives. The capture loop picks it
up on the next frame, so the half-duplex gate and barge-in cover it the
way they cover a reply. It resets the VAD first: playback is about to
suppress every frame, and a half-heard sentence would otherwise splice
onto whatever he says next. A nudge arriving while one still waits
replaces it, which is the contract internal/voice states for PushHandler.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 13:45:04 +04:00
claude 1f1e002789 One reader goroutine per voice conn, so a client can send and listen (V-671)
SendRequest and RunPushReceiver each read the conn, so a client that
wanted both raced for every frame. A second listening conn is not the
fix: it never sends a request, so its lastActive never moves and
PushToMostRecent never picks it. mavwaked needs both on one conn.

The reader now owns the socket for the life of the conn. It hands each
Response to whichever SendRequest waits on that id, and each Push to the
handler. SendRequest waits on its own channel, on the conn dying, on its
context, or on a timeout, and forgets its slot on every path that leaves
without an answer. RunPushReceiver just wires the handler and blocks.

Connect opens the conn without sending anything, for a client that must
hold a session before it has spoken.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013ptwopxyo3Z2kwFckHkLvN
2026-08-09 13:44:52 +04:00
10 changed files with 1219 additions and 74 deletions
+47 -5
View File
@@ -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)
}
+79
View File
@@ -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):
}
}
}
+170
View File
@@ -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
View File
@@ -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 {
+187
View File
@@ -0,0 +1,187 @@
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)
}
}
open := func(p string, in, out []string) (*ort.DynamicAdvancedSession, error) {
s, err := ort.NewDynamicAdvancedSession(p, in, out, nil)
if err != nil {
return nil, fmt.Errorf("wake word: load %s: %w", p, err)
}
return s, nil
}
m := &wakeModels{}
var err error
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
}
+191
View File
@@ -0,0 +1,191 @@
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. See
// docs/evals for the wakes-per-hour this buys.
const defaultWakeThreshold = 0.99
// 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:]...)
}
+174
View File
@@ -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()
}
+7 -3
View File
@@ -170,9 +170,13 @@ 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.
There is no wake word yet (V-487 stage two), so the loop runs open.
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
+118 -65
View File
@@ -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
}
+110
View File
@@ -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()
}