Half-duplex capture and barge-in in mavwaked (#287) #76
+33
-92
@@ -12,8 +12,15 @@
|
|||||||
// (30ms frames, 16kHz PCM) matches silero-vad's input interface exactly, so
|
// (30ms frames, 16kHz PCM) matches silero-vad's input interface exactly, so
|
||||||
// swapping energy-threshold for ONNX-inference is a local change in vad.go.
|
// swapping energy-threshold for ONNX-inference is a local change in vad.go.
|
||||||
//
|
//
|
||||||
|
// 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
|
||||||
|
// herself. -barge-in punches one hole in that gate — sustained energy well
|
||||||
|
// above the speaker's leak level cuts playback so he can talk over her. It is
|
||||||
|
// off by default because the threshold is room-specific; see playback.go.
|
||||||
|
//
|
||||||
// usage:
|
// usage:
|
||||||
// mavwaked # default ALSA device, 127.0.0.1:9100
|
// mavwaked # default ALSA device, 127.0.0.1:9100
|
||||||
|
// mavwaked -barge-in # let him interrupt her mid-reply
|
||||||
// mavwaked -device hw:1,0 -addr 10.42.0.1:9100
|
// mavwaked -device hw:1,0 -addr 10.42.0.1:9100
|
||||||
// mavwaked -test file.wav # read from file, no arecord
|
// mavwaked -test file.wav # read from file, no arecord
|
||||||
package main
|
package main
|
||||||
@@ -60,6 +67,9 @@ func run(args []string) error {
|
|||||||
silenceMs := flag.Int("silence-ms", defaultSilenceMs, "silence ms to end utterance")
|
silenceMs := flag.Int("silence-ms", defaultSilenceMs, "silence ms to end utterance")
|
||||||
maxMs := flag.Int("max-ms", defaultMaxMs, "max utterance ms")
|
maxMs := flag.Int("max-ms", defaultMaxMs, "max utterance ms")
|
||||||
testFile := flag.String("test", "", "read PCM from file instead of arecord (testing only)")
|
testFile := flag.String("test", "", "read PCM from file instead of arecord (testing only)")
|
||||||
|
bargeIn := flag.Bool("barge-in", false, "cut Maven off when he talks over her (needs a room-tuned -barge-in-rms)")
|
||||||
|
bargeRMS := flag.Int("barge-in-rms", defaultBargeRMS, "RMS x10000 a frame must clear to count as barge-in")
|
||||||
|
bargeFrames := flag.Int("barge-in-frames", defaultBargeFrames, "consecutive frames over -barge-in-rms before playback is cut")
|
||||||
flag.CommandLine.Parse(args)
|
flag.CommandLine.Parse(args)
|
||||||
|
|
||||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
|
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM, syscall.SIGHUP)
|
||||||
@@ -117,14 +127,21 @@ func run(args []string) error {
|
|||||||
|
|
||||||
defer src.Close()
|
defer src.Close()
|
||||||
|
|
||||||
return captureLoop(ctx, src, vad, vc, *lang)
|
var barge bargeInConfig
|
||||||
|
if *bargeIn {
|
||||||
|
barge = bargeInConfig{RMS: float64(*bargeRMS) / 10000.0, Frames: *bargeFrames}
|
||||||
|
log.Printf("mavwaked: barge-in on (rms %.4f x %d frames)", barge.RMS, barge.Frames)
|
||||||
|
}
|
||||||
|
sess := newSession(vad, newAplayPlayer(), &voiceSender{vc: vc}, *lang, barge)
|
||||||
|
|
||||||
|
return captureLoop(ctx, src, sess)
|
||||||
}
|
}
|
||||||
|
|
||||||
// captureLoop reads PCM from src, runs VAD, and sends complete utterances to
|
// captureLoop reads PCM from src and hands whole frames to the session.
|
||||||
// the voice server. Returns when ctx is done or src is exhausted.
|
// Returns when ctx is done or src is exhausted.
|
||||||
func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client, lang string) error {
|
func captureLoop(ctx context.Context, src io.Reader, sess *session) error {
|
||||||
br := bufio.NewReaderSize(src, defaultReadSize)
|
br := bufio.NewReaderSize(src, defaultReadSize)
|
||||||
frameBytes := vad.FrameSamples() * 2 // 480 samples × 2 bytes = 960 bytes per 30ms
|
frameBytes := sess.vad.FrameSamples() * 2 // 480 samples × 2 bytes = 960 bytes per 30ms
|
||||||
|
|
||||||
log.Printf("mavwaked: capture loop starting (frame=%d bytes, %dms)",
|
log.Printf("mavwaked: capture loop starting (frame=%d bytes, %dms)",
|
||||||
frameBytes, defaultFrameMs)
|
frameBytes, defaultFrameMs)
|
||||||
@@ -147,7 +164,7 @@ func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client,
|
|||||||
// Flush partial frame.
|
// Flush partial frame.
|
||||||
partial = append(partial, buf[:n]...)
|
partial = append(partial, buf[:n]...)
|
||||||
if len(partial) >= frameBytes {
|
if len(partial) >= frameBytes {
|
||||||
if err := processFrame(partial[:frameBytes], vad, vc, lang); err != nil {
|
if err := sess.feed(ctx, partial[:frameBytes]); err != nil {
|
||||||
log.Printf("mavwaked: process frame: %v", err)
|
log.Printf("mavwaked: process frame: %v", err)
|
||||||
}
|
}
|
||||||
partial = partial[frameBytes:]
|
partial = partial[frameBytes:]
|
||||||
@@ -165,107 +182,31 @@ func captureLoop(ctx context.Context, src io.Reader, vad *VAD, vc *voice.Client,
|
|||||||
partial = nil
|
partial = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := processFrame(full, vad, vc, lang); err != nil {
|
if err := sess.feed(ctx, full); err != nil {
|
||||||
log.Printf("mavwaked: process frame: %v", err)
|
log.Printf("mavwaked: process frame: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// processFrame feeds one 30ms PCM frame to the VAD and sends any completed
|
// voiceSender is the production utteranceSender: one PushToTalk round-trip
|
||||||
// utterance to the voice server.
|
// over the voice wire. SurfaceVoice (not the default SurfacePCClient that
|
||||||
func processFrame(frame []byte, vad *VAD, vc *voice.Client, lang string) error {
|
// c.PushToTalk uses) caps everything at L0, which is what makes an accidental
|
||||||
samples := PCMToI16(frame)
|
// VAD trigger safe.
|
||||||
utt, state := vad.Feed(samples)
|
type voiceSender struct{ vc *voice.Client }
|
||||||
|
|
||||||
if state == StateSpeech {
|
func (s *voiceSender) Send(ctx context.Context, utt audio.Audio, lang string) (audio.Audio, error) {
|
||||||
// Speech is in progress; nothing to send yet.
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
if utt.Bytes == nil {
|
|
||||||
// Still in silence, or short speech that didn't trigger.
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// We have a complete utterance — send it to the voice server.
|
|
||||||
return sendUtterance(context.Background(), utt, vc, lang)
|
|
||||||
}
|
|
||||||
|
|
||||||
// sendUtterance sends audio to the voice server and plays the reply.
|
|
||||||
func sendUtterance(ctx context.Context, utt audio.Audio, vc *voice.Client, lang string) error {
|
|
||||||
dur := utt.Duration()
|
|
||||||
log.Printf("mavwaked: utterance complete (%.2fs, %d bytes), sending...",
|
|
||||||
dur, len(utt.Bytes))
|
|
||||||
|
|
||||||
// Use SendRequest directly so we can set SurfaceVoice instead of the
|
|
||||||
// default SurfacePCClient that c.PushToTalk uses.
|
|
||||||
var resp voice.PushToTalkResp
|
var resp voice.PushToTalkResp
|
||||||
err := vc.SendRequest(ctx, voice.MethodPushToTalk, voice.PushToTalkReq{
|
err := s.vc.SendRequest(ctx, voice.MethodPushToTalk, voice.PushToTalkReq{
|
||||||
Audio: utt,
|
Audio: utt,
|
||||||
Lang: lang,
|
Lang: lang,
|
||||||
Surface: voice.SurfaceVoice,
|
Surface: voice.SurfaceVoice,
|
||||||
}, &resp)
|
}, &resp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("push-to-talk: %w", err)
|
return audio.Audio{}, fmt.Errorf("push-to-talk: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Printf("mavwaked: reply: %q (%.2fs audio)", resp.ReplyText, resp.ReplyAudio.Duration())
|
log.Printf("mavwaked: reply: %q (%.2fs audio)", resp.ReplyText, resp.ReplyAudio.Duration())
|
||||||
|
|
||||||
// Play the reply audio.
|
|
||||||
if len(resp.ReplyAudio.Bytes) > 0 {
|
|
||||||
go playAudio(resp.ReplyAudio)
|
|
||||||
} else {
|
|
||||||
log.Printf("mavwaked: empty reply audio (text only)")
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(resp.RoutedChannels) > 0 {
|
if len(resp.RoutedChannels) > 0 {
|
||||||
log.Printf("mavwaked: also routed to: %v", resp.RoutedChannels)
|
log.Printf("mavwaked: also routed to: %v", resp.RoutedChannels)
|
||||||
}
|
}
|
||||||
|
return resp.ReplyAudio, nil
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// playAudio pipes PCM audio to aplay(1) for playback. Runs in a goroutine.
|
|
||||||
func playAudio(a audio.Audio) {
|
|
||||||
// Build WAV header for aplay (or pipe raw PCM with the right format flags).
|
|
||||||
cmd := exec.Command("aplay",
|
|
||||||
"-f", "S16_LE",
|
|
||||||
"-r", fmt.Sprintf("%d", a.Format.SampleRate),
|
|
||||||
"-c", fmt.Sprintf("%d", a.Format.Channels),
|
|
||||||
"-t", "raw",
|
|
||||||
)
|
|
||||||
|
|
||||||
stdin, err := cmd.StdinPipe()
|
|
||||||
if err != nil {
|
|
||||||
log.Printf("mavwaked: aplay stdin pipe: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := cmd.Start(); err != nil {
|
|
||||||
log.Printf("mavwaked: start aplay: %v", err)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// Write audio to aplay's stdin.
|
|
||||||
if _, err := stdin.Write(a.Bytes); err != nil {
|
|
||||||
log.Printf("mavwaked: write to aplay: %v", err)
|
|
||||||
}
|
|
||||||
_ = stdin.Close()
|
|
||||||
|
|
||||||
// Wait for playback to finish (with a timeout).
|
|
||||||
done := make(chan error, 1)
|
|
||||||
go func() {
|
|
||||||
done <- cmd.Wait()
|
|
||||||
}()
|
|
||||||
|
|
||||||
select {
|
|
||||||
case err := <-done:
|
|
||||||
if err != nil {
|
|
||||||
log.Printf("mavwaked: aplay: %v", err)
|
|
||||||
}
|
|
||||||
case <-time.After(30 * time.Second):
|
|
||||||
log.Printf("mavwaked: aplay timeout, killing")
|
|
||||||
_ = cmd.Process.Kill()
|
|
||||||
<-done
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,139 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
// Reply playback, and the half-duplex gate around it (Vikunja #287).
|
||||||
|
//
|
||||||
|
// Before this, playback was `go playAudio(reply)` — fire and forget, with no
|
||||||
|
// handle on the running aplay. Two things fell out of that, and both are
|
||||||
|
// audible:
|
||||||
|
//
|
||||||
|
// 1. Self-trigger. The capture loop keeps feeding the VAD while the speaker
|
||||||
|
// is playing, so Maven's own reply comes back in through the mic, trips
|
||||||
|
// the VAD, and is sent to the daemon as a fresh utterance. She answers
|
||||||
|
// herself. There is no acoustic echo canceller in this pipeline, so the
|
||||||
|
// only correct fix is half-duplex: while she is speaking, the capture
|
||||||
|
// side is muted.
|
||||||
|
//
|
||||||
|
// 2. No barge-in. Talking over her did nothing — there was nothing to
|
||||||
|
// cancel, because nobody held the process handle.
|
||||||
|
//
|
||||||
|
// The two are the same mechanism seen from opposite sides, so they live
|
||||||
|
// together here. Echo suppression is unconditional (it fixes a bug). Barge-in
|
||||||
|
// is off unless -barge-in is passed, because it needs a room-specific energy
|
||||||
|
// threshold: with no echo canceller, the only way to tell "he is talking over
|
||||||
|
// her" from "the mic is hearing her" is that he is louder, and how much
|
||||||
|
// louder depends on where the mic sits relative to the speaker.
|
||||||
|
|
||||||
|
import (
|
||||||
|
"log"
|
||||||
|
"os/exec"
|
||||||
|
"strconv"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/audio"
|
||||||
|
)
|
||||||
|
|
||||||
|
// player plays one reply at a time and can be cut off mid-utterance.
|
||||||
|
type player interface {
|
||||||
|
// Play starts playback of a, replacing anything already playing, and
|
||||||
|
// returns immediately.
|
||||||
|
Play(a audio.Audio)
|
||||||
|
// Stop ends playback now. A no-op when nothing is playing.
|
||||||
|
Stop()
|
||||||
|
// Playing reports whether audio is currently going out of the speaker.
|
||||||
|
Playing() bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// aplayPlayer pipes raw PCM to aplay(1). Stop kills the child, which is what
|
||||||
|
// makes barge-in instant rather than "instant at the end of the sentence".
|
||||||
|
type aplayPlayer struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
cmd *exec.Cmd
|
||||||
|
playing bool
|
||||||
|
// gen rises on every Play/Stop so a finishing playback cannot clear the
|
||||||
|
// playing flag of the one that replaced it.
|
||||||
|
gen uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
func newAplayPlayer() *aplayPlayer { return &aplayPlayer{} }
|
||||||
|
|
||||||
|
func (p *aplayPlayer) Play(a audio.Audio) {
|
||||||
|
if len(a.Bytes) == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
p.Stop()
|
||||||
|
|
||||||
|
cmd := exec.Command("aplay",
|
||||||
|
"-f", "S16_LE",
|
||||||
|
"-r", strconv.Itoa(a.Format.SampleRate),
|
||||||
|
"-c", strconv.Itoa(a.Format.Channels),
|
||||||
|
"-t", "raw",
|
||||||
|
)
|
||||||
|
stdin, err := cmd.StdinPipe()
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("mavwaked: aplay stdin pipe: %v", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := cmd.Start(); err != nil {
|
||||||
|
log.Printf("mavwaked: start aplay: %v", err)
|
||||||
|
_ = stdin.Close()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
p.mu.Lock()
|
||||||
|
p.gen++
|
||||||
|
gen := p.gen
|
||||||
|
p.cmd = cmd
|
||||||
|
p.playing = true
|
||||||
|
p.mu.Unlock()
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
if _, err := stdin.Write(a.Bytes); err != nil {
|
||||||
|
// Broken pipe is the expected outcome of Stop().
|
||||||
|
log.Printf("mavwaked: write to aplay: %v", err)
|
||||||
|
}
|
||||||
|
_ = stdin.Close()
|
||||||
|
|
||||||
|
done := make(chan error, 1)
|
||||||
|
go func() { done <- cmd.Wait() }()
|
||||||
|
select {
|
||||||
|
case err := <-done:
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("mavwaked: aplay: %v", err)
|
||||||
|
}
|
||||||
|
case <-time.After(30 * time.Second):
|
||||||
|
log.Printf("mavwaked: aplay timeout, killing")
|
||||||
|
if pr := cmd.Process; pr != nil {
|
||||||
|
_ = pr.Kill()
|
||||||
|
}
|
||||||
|
<-done
|
||||||
|
}
|
||||||
|
|
||||||
|
p.mu.Lock()
|
||||||
|
if p.gen == gen {
|
||||||
|
p.playing = false
|
||||||
|
p.cmd = nil
|
||||||
|
}
|
||||||
|
p.mu.Unlock()
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *aplayPlayer) Stop() {
|
||||||
|
p.mu.Lock()
|
||||||
|
cmd := p.cmd
|
||||||
|
if cmd != nil {
|
||||||
|
p.gen++
|
||||||
|
p.playing = false
|
||||||
|
p.cmd = nil
|
||||||
|
}
|
||||||
|
p.mu.Unlock()
|
||||||
|
if cmd != nil && cmd.Process != nil {
|
||||||
|
_ = cmd.Process.Kill()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *aplayPlayer) Playing() bool {
|
||||||
|
p.mu.Lock()
|
||||||
|
defer p.mu.Unlock()
|
||||||
|
return p.playing
|
||||||
|
}
|
||||||
@@ -0,0 +1,33 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/audio"
|
||||||
|
)
|
||||||
|
|
||||||
|
// The real player must be safe to poke when nothing is playing — the capture
|
||||||
|
// loop calls Playing() on every 30ms frame, and Stop() lands on an idle
|
||||||
|
// player whenever a barge-in races the end of a reply. Neither may need
|
||||||
|
// aplay(1) to be installed.
|
||||||
|
func TestAplayPlayerIdleIsSafe(t *testing.T) {
|
||||||
|
p := newAplayPlayer()
|
||||||
|
if p.Playing() {
|
||||||
|
t.Fatal("a fresh player reports playing")
|
||||||
|
}
|
||||||
|
p.Stop()
|
||||||
|
p.Stop()
|
||||||
|
if p.Playing() {
|
||||||
|
t.Fatal("playing after Stop on an idle player")
|
||||||
|
}
|
||||||
|
// Empty audio is a text-only turn: nothing to play, no process to spawn.
|
||||||
|
p.Play(audio.Audio{Format: audio.PCM16kMono})
|
||||||
|
if p.Playing() {
|
||||||
|
t.Fatal("empty audio started playback")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestAplayPlayerSatisfiesPlayer(t *testing.T) {
|
||||||
|
var _ player = newAplayPlayer()
|
||||||
|
var _ player = &fakePlayer{}
|
||||||
|
}
|
||||||
@@ -0,0 +1,123 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
// The capture session: what happens to one 30ms frame, given whether Maven is
|
||||||
|
// currently speaking. Split out of main.go's processFrame so the decision is
|
||||||
|
// testable without a mic, a speaker, or a daemon (Vikunja #287).
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"log"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/audio"
|
||||||
|
)
|
||||||
|
|
||||||
|
// utteranceSender ships one complete utterance to the voice server and
|
||||||
|
// returns the reply audio to play. The real one round-trips over the voice
|
||||||
|
// wire; tests substitute a recorder.
|
||||||
|
type utteranceSender interface {
|
||||||
|
Send(ctx context.Context, utt audio.Audio, lang string) (audio.Audio, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// bargeInConfig holds the two numbers barge-in needs. Zero Frames disables
|
||||||
|
// barge-in entirely — the half-duplex gate still runs.
|
||||||
|
type bargeInConfig struct {
|
||||||
|
// RMS is the normalised energy a frame must exceed to count as him
|
||||||
|
// talking over her rather than the mic hearing her. It is deliberately
|
||||||
|
// far above the VAD's own floor: the speaker leaks into the mic at
|
||||||
|
// roughly ambient level, a person talking at the mic does not.
|
||||||
|
RMS float64
|
||||||
|
// Frames is how many consecutive frames must clear RMS before playback
|
||||||
|
// is cut. One loud frame is a door closing; five in a row is a voice.
|
||||||
|
Frames int
|
||||||
|
}
|
||||||
|
|
||||||
|
// Enabled reports whether barge-in should be attempted at all.
|
||||||
|
func (c bargeInConfig) Enabled() bool { return c.Frames > 0 && c.RMS > 0 }
|
||||||
|
|
||||||
|
// session is the per-client capture state machine.
|
||||||
|
type session struct {
|
||||||
|
vad *VAD
|
||||||
|
player player
|
||||||
|
sender utteranceSender
|
||||||
|
lang string
|
||||||
|
barge bargeInConfig
|
||||||
|
|
||||||
|
// loudFrames counts consecutive over-threshold frames seen while she is
|
||||||
|
// speaking. Reset whenever a frame falls back under the threshold, and
|
||||||
|
// whenever playback ends.
|
||||||
|
loudFrames int
|
||||||
|
|
||||||
|
// counters, read by tests and logged on the way out.
|
||||||
|
suppressed int // frames dropped because she was speaking
|
||||||
|
bargeIns int // times playback was cut because he spoke over her
|
||||||
|
sent int // utterances shipped to the daemon
|
||||||
|
}
|
||||||
|
|
||||||
|
func newSession(vad *VAD, p player, s utteranceSender, lang string, barge bargeInConfig) *session {
|
||||||
|
return &session{vad: vad, player: p, sender: s, lang: lang, barge: barge}
|
||||||
|
}
|
||||||
|
|
||||||
|
// feed processes one 30ms PCM frame.
|
||||||
|
//
|
||||||
|
// While the player is running the capture side is muted: the VAD is not fed
|
||||||
|
// and no utterance can be produced, so Maven's own reply cannot come back in
|
||||||
|
// as a new command. The one thing that gets through is barge-in — sustained
|
||||||
|
// energy well above the speaker's leak level cuts playback, and capture
|
||||||
|
// resumes on the very next frame with a clean VAD.
|
||||||
|
func (s *session) feed(ctx context.Context, frame []byte) error {
|
||||||
|
if s.player.Playing() {
|
||||||
|
s.suppressed++
|
||||||
|
if !s.barge.Enabled() {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if frameRMS(PCMToI16(frame)) < s.barge.RMS {
|
||||||
|
s.loudFrames = 0
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
s.loudFrames++
|
||||||
|
if s.loudFrames < s.barge.Frames {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
// He is talking over her. Cut her off, drop the VAD state that
|
||||||
|
// accumulated from the echo, and start listening for real.
|
||||||
|
s.player.Stop()
|
||||||
|
s.bargeIns++
|
||||||
|
s.loudFrames = 0
|
||||||
|
s.vad.Reset()
|
||||||
|
log.Printf("mavwaked: barge-in — stopped playback")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Not speaking. If we just stopped, make sure no echo-era state leaks
|
||||||
|
// into the next utterance.
|
||||||
|
if s.loudFrames != 0 {
|
||||||
|
s.loudFrames = 0
|
||||||
|
s.vad.Reset()
|
||||||
|
}
|
||||||
|
|
||||||
|
utt, state := s.vad.Feed(PCMToI16(frame))
|
||||||
|
if state == StateSpeech || utt.Bytes == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return s.dispatch(ctx, utt)
|
||||||
|
}
|
||||||
|
|
||||||
|
// dispatch ships a complete utterance and plays whatever comes back.
|
||||||
|
func (s *session) dispatch(ctx context.Context, utt audio.Audio) error {
|
||||||
|
log.Printf("mavwaked: utterance complete (%.2fs, %d bytes), sending...", utt.Duration(), len(utt.Bytes))
|
||||||
|
reply, err := s.sender.Send(ctx, utt, s.lang)
|
||||||
|
s.sent++
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if len(reply.Bytes) == 0 {
|
||||||
|
log.Printf("mavwaked: empty reply audio (text only)")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
// The VAD has been accumulating from the buffered mic stream while the
|
||||||
|
// round-trip blocked. None of it is a command — reset before the
|
||||||
|
// speaker opens, so the first post-reply frame starts clean.
|
||||||
|
s.vad.Reset()
|
||||||
|
s.player.Play(reply)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,282 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"math"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/audio"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakePlayer records Play/Stop instead of shelling out to aplay.
|
||||||
|
type fakePlayer struct {
|
||||||
|
playing bool
|
||||||
|
plays int
|
||||||
|
stops int
|
||||||
|
last audio.Audio
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *fakePlayer) Play(a audio.Audio) { p.playing = true; p.plays++; p.last = a }
|
||||||
|
func (p *fakePlayer) Stop() { p.playing = false; p.stops++ }
|
||||||
|
func (p *fakePlayer) Playing() bool { return p.playing }
|
||||||
|
|
||||||
|
// fakeSender records what was shipped and hands back a canned reply.
|
||||||
|
type fakeSender struct {
|
||||||
|
sent []audio.Audio
|
||||||
|
reply audio.Audio
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *fakeSender) Send(_ context.Context, utt audio.Audio, _ string) (audio.Audio, error) {
|
||||||
|
s.sent = append(s.sent, utt)
|
||||||
|
return s.reply, s.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func replyAudio() audio.Audio {
|
||||||
|
return audio.Audio{Format: audio.PCM16kMono, Bytes: make([]byte, 16000)}
|
||||||
|
}
|
||||||
|
|
||||||
|
// frameAt returns a 30ms frame whose RMS is approximately rms.
|
||||||
|
func frameAt(rms float64) []byte {
|
||||||
|
amp := rms * math.Sqrt2 * 32768
|
||||||
|
f := make([]int16, frameSamples)
|
||||||
|
for i := range f {
|
||||||
|
f[i] = int16(amp * math.Sin(2*math.Pi*440*float64(i)/16000))
|
||||||
|
}
|
||||||
|
return pcmBytes(f)
|
||||||
|
}
|
||||||
|
|
||||||
|
func silentBytes() []byte { return make([]byte, frameSamples*2) }
|
||||||
|
|
||||||
|
// newTestSession wires a session with fakes and a default VAD.
|
||||||
|
func newTestSession(barge bargeInConfig) (*session, *fakePlayer, *fakeSender) {
|
||||||
|
p := &fakePlayer{}
|
||||||
|
s := &fakeSender{reply: replyAudio()}
|
||||||
|
return newSession(NewVAD(0, 0, 0, 0), p, s, "ru", barge), p, s
|
||||||
|
}
|
||||||
|
|
||||||
|
// speakThenPause drives a full utterance through the session: enough loud
|
||||||
|
// frames to trigger, then enough silence to end it.
|
||||||
|
func speakThenPause(t *testing.T, sess *session) {
|
||||||
|
t.Helper()
|
||||||
|
speechFrames := (defaultSpeechMs + defaultFrameMs - 1) / defaultFrameMs
|
||||||
|
silenceFrames := (defaultSilenceMs+defaultFrameMs-1)/defaultFrameMs + 2
|
||||||
|
loud := frameAt(0.35)
|
||||||
|
for i := 0; i < speechFrames+5; i++ {
|
||||||
|
if err := sess.feed(context.Background(), loud); err != nil {
|
||||||
|
t.Fatalf("feed loud frame %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for i := 0; i < silenceFrames; i++ {
|
||||||
|
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||||
|
t.Fatalf("feed silent frame %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSessionSendsUtteranceAndPlaysReply(t *testing.T) {
|
||||||
|
sess, p, snd := newTestSession(bargeInConfig{})
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
|
||||||
|
if len(snd.sent) != 1 {
|
||||||
|
t.Fatalf("sent %d utterances, want 1", len(snd.sent))
|
||||||
|
}
|
||||||
|
if snd.sent[0].Format != audio.PCM16kMono {
|
||||||
|
t.Errorf("utterance format = %+v, want canonical", snd.sent[0].Format)
|
||||||
|
}
|
||||||
|
if p.plays != 1 {
|
||||||
|
t.Errorf("plays = %d, want 1", p.plays)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The bug this whole file exists for: while the speaker is running, the mic
|
||||||
|
// hears Maven and the old code shipped that back as a fresh command.
|
||||||
|
func TestSessionDoesNotHearItselfWhilePlaying(t *testing.T) {
|
||||||
|
sess, p, snd := newTestSession(bargeInConfig{})
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
if !p.Playing() {
|
||||||
|
t.Fatal("expected playback to be running after the reply")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Feed a long stretch of loud audio — Maven's own voice coming back in.
|
||||||
|
base := sess.suppressed
|
||||||
|
loud := frameAt(0.35)
|
||||||
|
for i := 0; i < 200; i++ {
|
||||||
|
if err := sess.feed(context.Background(), loud); err != nil {
|
||||||
|
t.Fatalf("feed echo frame %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(snd.sent) != 1 {
|
||||||
|
t.Fatalf("sent %d utterances, want 1 — her own reply was captured as a command", len(snd.sent))
|
||||||
|
}
|
||||||
|
if got := sess.suppressed - base; got != 200 {
|
||||||
|
t.Errorf("suppressed %d of the 200 echo frames, want all of them", got)
|
||||||
|
}
|
||||||
|
if p.stops != 0 {
|
||||||
|
t.Errorf("stops = %d, want 0 — barge-in is off, nothing should cut her off", p.stops)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// With barge-in off, no amount of noise stops playback.
|
||||||
|
func TestSessionBargeInDisabledByDefault(t *testing.T) {
|
||||||
|
sess, p, _ := newTestSession(bargeInConfig{})
|
||||||
|
if sess.barge.Enabled() {
|
||||||
|
t.Fatal("zero bargeInConfig must be disabled")
|
||||||
|
}
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
veryLoud := frameAt(0.6)
|
||||||
|
for i := 0; i < 50; i++ {
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
}
|
||||||
|
if p.stops != 0 || sess.bargeIns != 0 {
|
||||||
|
t.Fatalf("stops = %d, bargeIns = %d, want 0 with barge-in off", p.stops, sess.bargeIns)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSessionBargeInCutsPlayback(t *testing.T) {
|
||||||
|
barge := bargeInConfig{RMS: 0.12, Frames: 5}
|
||||||
|
sess, p, _ := newTestSession(barge)
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
if !p.Playing() {
|
||||||
|
t.Fatal("expected playback after the reply")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Four loud frames must not be enough — a door closing is not a voice.
|
||||||
|
veryLoud := frameAt(0.35)
|
||||||
|
for i := 0; i < 4; i++ {
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
}
|
||||||
|
if p.stops != 0 {
|
||||||
|
t.Fatalf("playback cut after 4 frames, want it to hold until %d", barge.Frames)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The fifth cuts her off.
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
if p.stops != 1 || sess.bargeIns != 1 {
|
||||||
|
t.Fatalf("stops = %d, bargeIns = %d, want 1 and 1", p.stops, sess.bargeIns)
|
||||||
|
}
|
||||||
|
if p.Playing() {
|
||||||
|
t.Fatal("still playing after barge-in")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A burst that falls back under the threshold resets the counter, so noise
|
||||||
|
// spread over a whole reply never accumulates into a false barge-in.
|
||||||
|
func TestSessionBargeInNeedsConsecutiveFrames(t *testing.T) {
|
||||||
|
sess, p, _ := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
|
||||||
|
veryLoud := frameAt(0.35)
|
||||||
|
quiet := frameAt(0.02)
|
||||||
|
for i := 0; i < 20; i++ {
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
_ = sess.feed(context.Background(), quiet)
|
||||||
|
}
|
||||||
|
if p.stops != 0 || sess.bargeIns != 0 {
|
||||||
|
t.Fatalf("stops = %d, bargeIns = %d, want 0 — two-frame bursts must not accumulate", p.stops, sess.bargeIns)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Speaker leak sits near the room floor; it must never reach the barge-in bar.
|
||||||
|
func TestSessionEchoLevelAudioNeverBargesIn(t *testing.T) {
|
||||||
|
sess, p, _ := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
|
||||||
|
base := sess.suppressed
|
||||||
|
leak := frameAt(0.05) // loud enough for the VAD, far under the barge bar
|
||||||
|
for i := 0; i < 300; i++ {
|
||||||
|
_ = sess.feed(context.Background(), leak)
|
||||||
|
}
|
||||||
|
if p.stops != 0 {
|
||||||
|
t.Fatalf("stops = %d, want 0 — speaker leak must not read as barge-in", p.stops)
|
||||||
|
}
|
||||||
|
if got := sess.suppressed - base; got != 300 {
|
||||||
|
t.Errorf("suppressed %d of the 300 leak frames, want all of them", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// After barge-in the VAD must start clean, so the interrupting speech is
|
||||||
|
// captured as a whole utterance rather than joined onto echo state.
|
||||||
|
func TestSessionCapturesTheInterruptingUtterance(t *testing.T) {
|
||||||
|
sess, p, snd := newTestSession(bargeInConfig{RMS: 0.12, Frames: 5})
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
|
||||||
|
veryLoud := frameAt(0.35)
|
||||||
|
for i := 0; i < 5; i++ {
|
||||||
|
_ = sess.feed(context.Background(), veryLoud)
|
||||||
|
}
|
||||||
|
if p.stops != 1 {
|
||||||
|
t.Fatalf("expected barge-in, stops = %d", p.stops)
|
||||||
|
}
|
||||||
|
|
||||||
|
// He keeps talking; that is a new command.
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
if len(snd.sent) != 2 {
|
||||||
|
t.Fatalf("sent %d utterances, want 2 — the interruption itself must be heard", len(snd.sent))
|
||||||
|
}
|
||||||
|
if p.plays != 2 {
|
||||||
|
t.Errorf("plays = %d, want 2", p.plays)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A failed round-trip must surface as an error and must not start playback.
|
||||||
|
func TestSessionSendErrorDoesNotPlay(t *testing.T) {
|
||||||
|
p := &fakePlayer{}
|
||||||
|
snd := &fakeSender{err: errors.New("boom")}
|
||||||
|
sess := newSession(NewVAD(0, 0, 0, 0), p, snd, "ru", bargeInConfig{})
|
||||||
|
|
||||||
|
speechFrames := (defaultSpeechMs + defaultFrameMs - 1) / defaultFrameMs
|
||||||
|
silenceFrames := (defaultSilenceMs+defaultFrameMs-1)/defaultFrameMs + 2
|
||||||
|
loud := frameAt(0.35)
|
||||||
|
var lastErr error
|
||||||
|
for i := 0; i < speechFrames+5; i++ {
|
||||||
|
_ = sess.feed(context.Background(), loud)
|
||||||
|
}
|
||||||
|
for i := 0; i < silenceFrames; i++ {
|
||||||
|
if err := sess.feed(context.Background(), silentBytes()); err != nil {
|
||||||
|
lastErr = err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if lastErr == nil {
|
||||||
|
t.Fatal("send error was swallowed")
|
||||||
|
}
|
||||||
|
if p.plays != 0 || p.Playing() {
|
||||||
|
t.Fatalf("plays = %d, playing = %v, want no playback on a failed round-trip", p.plays, p.Playing())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// An empty reply (text-only turn) must leave the capture side open.
|
||||||
|
func TestSessionEmptyReplyLeavesCaptureOpen(t *testing.T) {
|
||||||
|
p := &fakePlayer{}
|
||||||
|
snd := &fakeSender{reply: audio.Audio{Format: audio.PCM16kMono}}
|
||||||
|
sess := newSession(NewVAD(0, 0, 0, 0), p, snd, "ru", bargeInConfig{})
|
||||||
|
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
if p.plays != 0 {
|
||||||
|
t.Fatalf("plays = %d, want 0 for an empty reply", p.plays)
|
||||||
|
}
|
||||||
|
speakThenPause(t, sess)
|
||||||
|
if len(snd.sent) != 2 {
|
||||||
|
t.Fatalf("sent %d, want 2 — capture must stay open when there is no audio reply", len(snd.sent))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestBargeInConfigEnabled(t *testing.T) {
|
||||||
|
cases := []struct {
|
||||||
|
c bargeInConfig
|
||||||
|
want bool
|
||||||
|
}{
|
||||||
|
{bargeInConfig{}, false},
|
||||||
|
{bargeInConfig{RMS: 0.12}, false},
|
||||||
|
{bargeInConfig{Frames: 5}, false},
|
||||||
|
{bargeInConfig{RMS: 0.12, Frames: 5}, true},
|
||||||
|
}
|
||||||
|
for _, tc := range cases {
|
||||||
|
if got := tc.c.Enabled(); got != tc.want {
|
||||||
|
t.Errorf("%+v.Enabled() = %v, want %v", tc.c, got, tc.want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -23,6 +23,15 @@ const (
|
|||||||
defaultSilenceMs = 800 // silence hold before declaring end-of-utterance
|
defaultSilenceMs = 800 // silence hold before declaring end-of-utterance
|
||||||
defaultMaxMs = 10000 // cap single utterance at 10s
|
defaultMaxMs = 10000 // cap single utterance at 10s
|
||||||
defaultMinRMS = 0.01 // RMS floor (same as mavsttd)
|
defaultMinRMS = 0.01 // RMS floor (same as mavsttd)
|
||||||
|
|
||||||
|
// Barge-in thresholds. Only used when -barge-in is passed. The RMS is
|
||||||
|
// x10000 like -min-rms, and sits an order of magnitude above the VAD's
|
||||||
|
// own floor on purpose: with no acoustic echo canceller, a frame only
|
||||||
|
// counts as "he is talking over her" if it is far louder than what the
|
||||||
|
// speaker leaks back into the mic. 5 frames is 150ms — long enough that
|
||||||
|
// a door or a cough does not cut her off mid-sentence.
|
||||||
|
defaultBargeRMS = 1200 // 0.12 normalised RMS
|
||||||
|
defaultBargeFrames = 5
|
||||||
)
|
)
|
||||||
|
|
||||||
// frameSamples — samples per 30ms frame at 16kHz.
|
// frameSamples — samples per 30ms frame at 16kHz.
|
||||||
|
|||||||
Reference in New Issue
Block a user