From a6fc4cf285945375cf9adf1769d6da8f4f567250 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Oscar=20S=C3=B6derlund?= Date: Tue, 4 Aug 2026 10:54:50 +0200 Subject: [PATCH] feat(capture): dynamic sink-level audio capture --- internal/audio/pcm/pcm.go | 29 ++ internal/audio/pcm/pcm_test.go | 90 ++++ internal/capture/capture.go | 11 +- internal/capture/parec.go | 167 ++++--- internal/capture/parec_test.go | 583 +++++++++++++++++++++---- internal/capture/reader.go | 74 ++++ internal/capture/reconcile.go | 218 +++++++++ internal/protocol/parec/client.go | 61 ++- internal/protocol/parec/client_test.go | 92 +++- internal/recorder/capture.go | 11 +- 10 files changed, 1110 insertions(+), 226 deletions(-) create mode 100644 internal/capture/reader.go create mode 100644 internal/capture/reconcile.go diff --git a/internal/audio/pcm/pcm.go b/internal/audio/pcm/pcm.go index 0d6df4a..5f506eb 100644 --- a/internal/audio/pcm/pcm.go +++ b/internal/audio/pcm/pcm.go @@ -45,3 +45,32 @@ func TrimTrailingFrames(pcm []byte, frames, frameBytes int) []byte { } return pcm } + +// Mix combines s16le mono PCM frames via saturating addition. Silent or idle +// inputs (all-zero) don't attenuate active ones the way averaging would. +// Frames are expected to be FrameBytes long; shorter frames contribute only +// their available samples. Mixing zero frames returns FrameBytes of silence. +func Mix(frames ...[]byte) []byte { + out := make([]byte, FrameBytes) + for _, f := range frames { + n := min(len(f), FrameBytes) / 2 + for i := range n { + a := int16(binary.LittleEndian.Uint16(out[i*2 : i*2+2])) + b := int16(binary.LittleEndian.Uint16(f[i*2 : i*2+2])) + binary.LittleEndian.PutUint16(out[i*2:i*2+2], uint16(saturatingAdd(a, b))) + } + } + return out +} + +func saturatingAdd(a, b int16) int16 { + sum := int32(a) + int32(b) + switch { + case sum > math.MaxInt16: + return math.MaxInt16 + case sum < math.MinInt16: + return math.MinInt16 + default: + return int16(sum) + } +} diff --git a/internal/audio/pcm/pcm_test.go b/internal/audio/pcm/pcm_test.go index 006b86f..7d68c39 100644 --- a/internal/audio/pcm/pcm_test.go +++ b/internal/audio/pcm/pcm_test.go @@ -1,8 +1,10 @@ package pcm import ( + "bytes" "encoding/binary" "math" + "slices" "testing" ) @@ -71,6 +73,94 @@ func TestFrameCount(t *testing.T) { } } +func sampleFrame(t *testing.T, samples ...int16) []byte { + t.Helper() + buf := make([]byte, FrameBytes) + for i, s := range samples { + binary.LittleEndian.PutUint16(buf[i*2:], uint16(s)) + } + return buf +} + +func readSamples(t *testing.T, data []byte, n int) []int16 { + t.Helper() + out := make([]int16, n) + for i := range n { + out[i] = int16(binary.LittleEndian.Uint16(data[i*2:])) + } + return out +} + +func TestMix_NoInputs(t *testing.T) { + got := Mix() + if len(got) != FrameBytes { + t.Fatalf("len = %d, want %d", len(got), FrameBytes) + } + for _, b := range got { + if b != 0 { + t.Fatalf("expected silence, got non-zero byte") + } + } +} + +func TestMix_SingleInput(t *testing.T) { + f := sampleFrame(t, 100, -200, 300) + got := Mix(f) + if !bytes.Equal(got, f) { + t.Errorf("single input should be unchanged") + } +} + +func TestMix_TwoInputsAdd(t *testing.T) { + a := sampleFrame(t, 100, -200, 300) + b := sampleFrame(t, 50, -50, -100) + got := Mix(a, b) + want := []int16{150, -250, 200} + if s := readSamples(t, got, len(want)); !slices.Equal(s, want) { + t.Errorf("Mix() = %v, want %v", s, want) + } +} + +func TestMix_SilentDoesNotAttenuate(t *testing.T) { + active := sampleFrame(t, 1000, -1000, 500) + silent := make([]byte, FrameBytes) + got := Mix(active, silent) + if !bytes.Equal(got, active) { + t.Errorf("mixing with silence should leave active frame unchanged") + } +} + +func TestMix_PositiveOverflowClamps(t *testing.T) { + a := sampleFrame(t, math.MaxInt16) + b := sampleFrame(t, math.MaxInt16) + got := Mix(a, b) + want := []int16{math.MaxInt16} + if s := readSamples(t, got, 1); !slices.Equal(s, want) { + t.Errorf("Mix() = %v, want %v", s, want) + } +} + +func TestMix_NegativeOverflowClamps(t *testing.T) { + a := sampleFrame(t, math.MinInt16) + b := sampleFrame(t, math.MinInt16) + got := Mix(a, b) + want := []int16{math.MinInt16} + if s := readSamples(t, got, 1); !slices.Equal(s, want) { + t.Errorf("Mix() = %v, want %v", s, want) + } +} + +func TestMix_InputBuffersUnchanged(t *testing.T) { + a := sampleFrame(t, 100, -200) + b := sampleFrame(t, 50, -50) + aCopy := append([]byte(nil), a...) + bCopy := append([]byte(nil), b...) + Mix(a, b) + if !bytes.Equal(a, aCopy) || !bytes.Equal(b, bCopy) { + t.Errorf("Mix() must not mutate its input buffers") + } +} + func TestTrimTrailingFrames(t *testing.T) { pcm := make([]byte, FrameBytes*3) tests := []struct { diff --git a/internal/capture/capture.go b/internal/capture/capture.go index 21b79a8..46a0a9a 100644 --- a/internal/capture/capture.go +++ b/internal/capture/capture.go @@ -4,12 +4,19 @@ import ( "context" "github.com/odsod/recorder/internal/audio/frame" + "github.com/odsod/recorder/internal/protocol/parec" ) // Source abstracts dual-channel audio capture (system + microphone). type Source interface { Start(ctx context.Context) (<-chan frame.Dual, error) Stop() error - MonitorSource() string - MicSource() string +} + +// sinkClient is the subset of *parec.Client that capture depends on. +// Narrowed for testability; *parec.Client satisfies it structurally. +type sinkClient interface { + ListSinks(ctx context.Context, req parec.ListSinksRequest) (parec.ListSinksResponse, error) + GetDefaultSource(ctx context.Context, req parec.GetDefaultSourceRequest) (parec.GetDefaultSourceResponse, error) + StartCapture(ctx context.Context, req parec.StartCaptureRequest) (*parec.CaptureStream, error) } diff --git a/internal/capture/parec.go b/internal/capture/parec.go index 2383baf..ebfca4c 100644 --- a/internal/capture/parec.go +++ b/internal/capture/parec.go @@ -2,119 +2,116 @@ package capture import ( "context" - "errors" - "fmt" - "io" "sync" + "time" "github.com/odsod/recorder/internal/audio/frame" "github.com/odsod/recorder/internal/audio/pcm" "github.com/odsod/recorder/internal/protocol/parec" ) -// Parec implements Source using PulseAudio's parec command. +// tickInterval is the mix loop's emission period, matching pcm.FrameBytes' +// one second of audio. A package-level var so tests can shrink it. +var tickInterval = time.Second + +// Parec implements Source by dynamically monitoring every PulseAudio sink +// plus the default microphone, mixing simultaneous sink audio together. type Parec struct { - client *parec.Client - monitor string - mic string - stop func() -} + client sinkClient -// NewParec creates a Parec source using the given parec protocol client. -func NewParec(client *parec.Client) *Parec { - return &Parec{client: client} -} + mu sync.Mutex + sinks map[string]*reader // keyed by sink name + mic *reader + micName string + + sinkBackoff backoff + sinkLastTry time.Time + micBackoff backoff + micLastTry time.Time -// MonitorSource returns the system audio monitor source name. -func (c *Parec) MonitorSource() string { - return c.monitor + cancel context.CancelFunc + wg sync.WaitGroup + stopOnce sync.Once } -// MicSource returns the microphone source name. -func (c *Parec) MicSource() string { - return c.mic +// NewParec creates a Parec source using the given parec protocol client. +func NewParec(client *parec.Client) *Parec { + return &Parec{client: client, sinks: make(map[string]*reader)} } -// Start begins capturing system and microphone audio, returning a channel of frames. +// Start begins dynamic sink and microphone capture, returning a channel of +// mixed frames. Discovery and stream startup happen in the background; +// outages (no sinks, unreachable PulseAudio, a dead parec process) surface +// as silent frames rather than failing Start or closing the channel — the +// channel only closes once ctx is done. func (c *Parec) Start(ctx context.Context) (<-chan frame.Dual, error) { - sinkResp, err := c.client.GetDefaultSink(ctx, parec.GetDefaultSinkRequest{}) - if err != nil { - return nil, err - } - sourceResp, err := c.client.GetDefaultSource(ctx, parec.GetDefaultSourceRequest{}) - if err != nil { - return nil, err - } - c.monitor = sinkResp.MonitorSource - c.mic = sourceResp.Source - - sysStream, err := c.client.StartCapture(ctx, parec.StartCaptureRequest{ - Device: c.monitor, SampleRate: pcm.SampleRate, - }) - if err != nil { - return nil, fmt.Errorf("start sys parec: %w", err) - } - micStream, err := c.client.StartCapture(ctx, parec.StartCaptureRequest{ - Device: c.mic, SampleRate: pcm.SampleRate, - }) - if err != nil { - _ = sysStream.Close() - return nil, fmt.Errorf("start mic parec: %w", err) - } + runCtx, cancel := context.WithCancel(ctx) + c.cancel = cancel frames := make(chan frame.Dual, 2) - done := make(chan struct{}) - var once sync.Once - c.stop = func() { - once.Do(func() { - close(done) - _ = sysStream.Close() - _ = micStream.Close() - }) - } - go func() { - defer close(frames) - defer c.stop() - for { - select { - case <-ctx.Done(): - return - case <-done: - return - default: - } + c.wg.Go(func() { c.reconcileLoop(runCtx) }) + c.wg.Go(func() { c.mixLoop(runCtx, frames) }) - sysData, err := frame.Read(sysStream, pcm.FrameBytes) - if err != nil { - if !errors.Is(err, io.EOF) { - return - } - return - } - micData, err := frame.Read(micStream, pcm.FrameBytes) - if err != nil { - micData = frame.Silent(pcm.FrameBytes) - } + return frames, nil +} - f := frame.Dual{Sys: sysData, Mic: micData} +func (c *Parec) mixLoop(ctx context.Context, out chan<- frame.Dual) { + defer close(out) + ticker := time.NewTicker(tickInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + f := c.mixTick() select { - case frames <- f: + case out <- f: case <-ctx.Done(): return - case <-done: - return } } - }() + } +} - return frames, nil +func (c *Parec) mixTick() frame.Dual { + c.mu.Lock() + sinkFrames := make([][]byte, 0, len(c.sinks)) + for _, r := range c.sinks { + if data, ok := r.take(); ok { + sinkFrames = append(sinkFrames, data) + } + } + micReader := c.mic + c.mu.Unlock() + + mic := frame.Silent(pcm.FrameBytes) + if micReader != nil { + if data, ok := micReader.take(); ok { + mic = data + } + } + + return frame.Dual{Sys: pcm.Mix(sinkFrames...), Mic: mic} } -// Stop terminates the capture processes. +// Stop terminates all capture streams. Idempotent. func (c *Parec) Stop() error { - if c.stop != nil { - c.stop() - } + c.stopOnce.Do(func() { + if c.cancel != nil { + c.cancel() + } + c.wg.Wait() + + c.mu.Lock() + for _, r := range c.sinks { + _ = r.stop() + } + if c.mic != nil { + _ = c.mic.stop() + } + c.mu.Unlock() + }) return nil } diff --git a/internal/capture/parec_test.go b/internal/capture/parec_test.go index df4346a..c6da817 100644 --- a/internal/capture/parec_test.go +++ b/internal/capture/parec_test.go @@ -4,145 +4,548 @@ import ( "bytes" "context" "io" + "sync" "testing" + "time" - "github.com/odsod/recorder/internal/audio/frame" "github.com/odsod/recorder/internal/audio/pcm" "github.com/odsod/recorder/internal/protocol/parec" ) +// fakeRunner backs a real *parec.Client so fakeSinkClient can produce real +// *parec.CaptureStream values, matching the existing hand-rolled-fake +// convention used by internal/protocol/parec's own tests. type fakeRunner struct { - outputFn func(ctx context.Context, name string, args ...string) ([]byte, error) - startFn func(ctx context.Context, name string, args ...string) (io.ReadCloser, func() error, error) + mu sync.Mutex + streams map[string]func() (io.ReadCloser, func() error, error) +} + +func newFakeRunner() *fakeRunner { + return &fakeRunner{streams: make(map[string]func() (io.ReadCloser, func() error, error))} } func (f *fakeRunner) Output(ctx context.Context, name string, args ...string) ([]byte, error) { - return f.outputFn(ctx, name, args...) + panic("fakeRunner.Output should not be called; capture talks to fakeSinkClient directly") } func (f *fakeRunner) Start(ctx context.Context, name string, args ...string) (io.ReadCloser, func() error, error) { - return f.startFn(ctx, name, args...) -} - -func TestParec_Start(t *testing.T) { - sysData := bytes.Repeat([]byte{0x01}, pcm.FrameBytes*2) - micData := bytes.Repeat([]byte{0x02}, pcm.FrameBytes*2) - startCount := 0 - - runner := &fakeRunner{ - outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { - if args[0] == "get-default-sink" { - return []byte("sink\n"), nil - } - return []byte("mic\n"), nil - }, - startFn: func(ctx context.Context, name string, args ...string) (io.ReadCloser, func() error, error) { - startCount++ - if startCount == 1 { - return io.NopCloser(bytes.NewReader(sysData)), func() error { return nil }, nil - } - return io.NopCloser(bytes.NewReader(micData)), func() error { return nil }, nil - }, - } - - src := NewParec(parec.New(runner)) + device := deviceArg(args) + f.mu.Lock() + fn, ok := f.streams[device] + f.mu.Unlock() + if !ok { + return io.NopCloser(bytes.NewReader(nil)), func() error { return nil }, nil + } + return fn() +} + +func deviceArg(args []string) string { + for _, a := range args { + if len(a) > len("--device=") && a[:len("--device=")] == "--device=" { + return a[len("--device="):] + } + } + return "" +} + +// setStream registers repeating frame data for a device: each StartCapture +// call for that device gets a fresh reader over data, repeated indefinitely +// so tests don't need to size data to an exact number of ticks. +func (f *fakeRunner) setStream(device string, data []byte) { + f.mu.Lock() + defer f.mu.Unlock() + f.streams[device] = func() (io.ReadCloser, func() error, error) { + return io.NopCloser(&repeatingReader{data: data}), func() error { return nil }, nil + } +} + +// setStreamOnce registers one-shot frame data (EOF after data is exhausted), +// used to simulate a reader dying. +func (f *fakeRunner) setStreamOnce(device string, data []byte) { + f.mu.Lock() + defer f.mu.Unlock() + f.streams[device] = func() (io.ReadCloser, func() error, error) { + return io.NopCloser(bytes.NewReader(data)), func() error { return nil }, nil + } +} + +// repeatingReader loops over data forever, so a fake capture stream never +// looks like it died from the reader's perspective. +type repeatingReader struct { + data []byte + pos int +} + +func (r *repeatingReader) Read(p []byte) (int, error) { + if len(r.data) == 0 { + return 0, io.EOF + } + n := copy(p, r.data[r.pos:]) + r.pos += n + if r.pos >= len(r.data) { + r.pos = 0 + } + return n, nil +} + +// fakeSinkClient implements sinkClient directly (no JSON/CommandRunner +// involved), for scripting reconciliation scenarios plainly. +type fakeSinkClient struct { + mu sync.Mutex + runner *fakeRunner + client *parec.Client + sinks []parec.Sink + sinksErr error + source string + sourceErr error + startErrFn func(device string) error +} + +func newFakeSinkClient() *fakeSinkClient { + runner := newFakeRunner() + return &fakeSinkClient{runner: runner, client: parec.New(runner)} +} + +func (f *fakeSinkClient) ListSinks(ctx context.Context, _ parec.ListSinksRequest) (parec.ListSinksResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + if f.sinksErr != nil { + return parec.ListSinksResponse{}, f.sinksErr + } + return parec.ListSinksResponse{Sinks: f.sinks}, nil +} + +func (f *fakeSinkClient) GetDefaultSource( + ctx context.Context, + _ parec.GetDefaultSourceRequest, +) (parec.GetDefaultSourceResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + if f.sourceErr != nil { + return parec.GetDefaultSourceResponse{}, f.sourceErr + } + return parec.GetDefaultSourceResponse{Source: f.source}, nil +} + +func (f *fakeSinkClient) StartCapture( + ctx context.Context, + req parec.StartCaptureRequest, +) (*parec.CaptureStream, error) { + f.mu.Lock() + startErrFn := f.startErrFn + f.mu.Unlock() + if startErrFn != nil { + if err := startErrFn(req.Device); err != nil { + return nil, err + } + } + return f.client.StartCapture(ctx, req) +} + +func (f *fakeSinkClient) setSinks(sinks ...parec.Sink) { + f.mu.Lock() + defer f.mu.Unlock() + f.sinks = sinks +} + +func (f *fakeSinkClient) setSinksErr(err error) { + f.mu.Lock() + defer f.mu.Unlock() + f.sinksErr = err +} + +func (f *fakeSinkClient) setSource(name string) { + f.mu.Lock() + defer f.mu.Unlock() + f.source = name +} + +func (f *fakeSinkClient) setSourceErr(err error) { + f.mu.Lock() + defer f.mu.Unlock() + f.sourceErr = err +} + +// waitFor polls cond until it returns true or the timeout elapses. +func waitFor(t *testing.T, timeout time.Duration, cond func() bool) { + t.Helper() + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(time.Millisecond) + } + if !cond() { + t.Fatal("condition not met before timeout") + } +} + +func withFastPolling(t *testing.T) { + t.Helper() + origSink, origMic, origTick := sinkPollInterval, micPollInterval, tickInterval + origInitial, origMax := backoffInitial, backoffMax + sinkPollInterval = 5 * time.Millisecond + micPollInterval = 5 * time.Millisecond + tickInterval = 5 * time.Millisecond + backoffInitial = 5 * time.Millisecond + backoffMax = 20 * time.Millisecond + t.Cleanup(func() { + sinkPollInterval, micPollInterval, tickInterval = origSink, origMic, origTick + backoffInitial, backoffMax = origInitial, origMax + }) +} + +func TestParec_MultiSinkMixing(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("a.monitor", bytes.Repeat([]byte{0x01, 0x00}, pcm.FrameBytes/2)) + client.runner.setStream("b.monitor", bytes.Repeat([]byte{0x02, 0x00}, pcm.FrameBytes/2)) + client.setSinks( + parec.Sink{Name: "a", MonitorSource: "a.monitor"}, + parec.Sink{Name: "b", MonitorSource: "b.monitor"}, + ) + + src := &Parec{client: client, sinks: make(map[string]*reader)} frames, err := src.Start(context.Background()) if err != nil { t.Fatal(err) } + defer func() { _ = src.Stop() }() - frame1, ok := <-frames - if !ok { - t.Fatal("expected first frame") + f := <-frames + for f.Sys[0] == 0 { + f = <-frames } - if !bytes.Equal(frame1.Sys, sysData[:pcm.FrameBytes]) { - t.Error("sys frame mismatch") + if f.Sys[0] != 0x03 { + t.Errorf("expected mixed sys byte 0x03, got %#x", f.Sys[0]) } - if !bytes.Equal(frame1.Mic, micData[:pcm.FrameBytes]) { - t.Error("mic frame mismatch") +} + +func TestParec_SinkAppears(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("a.monitor", bytes.Repeat([]byte{0x05, 0x00}, pcm.FrameBytes/2)) + + src := &Parec{client: client, sinks: make(map[string]*reader)} + frames, err := src.Start(context.Background()) + if err != nil { + t.Fatal(err) } + defer func() { _ = src.Stop() }() - frame2, ok := <-frames - if !ok { - t.Fatal("expected second frame") + client.setSinks(parec.Sink{Name: "a", MonitorSource: "a.monitor"}) + + found := false + for range 200 { + f := <-frames + if f.Sys[0] == 0x05 { + found = true + break + } } - if !bytes.Equal(frame2.Sys, sysData[pcm.FrameBytes:]) { - t.Error("second sys frame mismatch") + if !found { + t.Error("expected sink audio to appear after discovery") } +} - if _, ok := <-frames; ok { - t.Error("expected channel closed after streams exhausted") - } +func TestParec_SinkDisappears(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("a.monitor", bytes.Repeat([]byte{0x05, 0x00}, pcm.FrameBytes/2)) + client.runner.setStream("b.monitor", bytes.Repeat([]byte{0x06, 0x00}, pcm.FrameBytes/2)) + client.setSinks( + parec.Sink{Name: "a", MonitorSource: "a.monitor"}, + parec.Sink{Name: "b", MonitorSource: "b.monitor"}, + ) - if src.MonitorSource() != "sink.monitor" { - t.Errorf("MonitorSource = %q, want sink.monitor", src.MonitorSource()) - } - if src.MicSource() != "mic" { - t.Errorf("MicSource = %q, want mic", src.MicSource()) + src := &Parec{client: client, sinks: make(map[string]*reader)} + if _, err := src.Start(context.Background()); err != nil { + t.Fatal(err) } + defer func() { _ = src.Stop() }() + + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + return len(src.sinks) == 2 + }) + + client.setSinks(parec.Sink{Name: "b", MonitorSource: "b.monitor"}) + + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, hasA := src.sinks["a"] + _, hasB := src.sinks["b"] + return !hasA && hasB + }) } -func TestParec_MicFallbackToSilence(t *testing.T) { - sysData := bytes.Repeat([]byte{0x01}, pcm.FrameBytes) - startCount := 0 +func TestParec_SinkReappearsNewIndex(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("a.monitor", bytes.Repeat([]byte{0x05, 0x00}, pcm.FrameBytes/2)) + client.setSinks(parec.Sink{Name: "a", MonitorSource: "a.monitor"}) - runner := &fakeRunner{ - outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { - if args[0] == "get-default-sink" { - return []byte("sink\n"), nil - } - return []byte("mic\n"), nil - }, - startFn: func(ctx context.Context, name string, args ...string) (io.ReadCloser, func() error, error) { - startCount++ - if startCount == 1 { - return io.NopCloser(bytes.NewReader(sysData)), func() error { return nil }, nil - } - return io.NopCloser(bytes.NewReader(nil)), func() error { return nil }, nil - }, + src := &Parec{client: client, sinks: make(map[string]*reader)} + _, err := src.Start(context.Background()) + if err != nil { + t.Fatal(err) } + defer func() { _ = src.Stop() }() + + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, ok := src.sinks["a"] + return ok + }) + + client.setSinks() + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, ok := src.sinks["a"] + return !ok + }) + + client.setSinks(parec.Sink{Name: "a", MonitorSource: "a.monitor"}) + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, ok := src.sinks["a"] + return ok + }) +} + +func TestParec_MicHotSwapSuccess(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("mic1", bytes.Repeat([]byte{0x07, 0x00}, pcm.FrameBytes/2)) + client.runner.setStream("mic2", bytes.Repeat([]byte{0x08, 0x00}, pcm.FrameBytes/2)) + client.setSource("mic1") - src := NewParec(parec.New(runner)) + src := &Parec{client: client, sinks: make(map[string]*reader)} frames, err := src.Start(context.Background()) if err != nil { t.Fatal(err) } + defer func() { _ = src.Stop() }() - f, ok := <-frames - if !ok { - t.Fatal("expected frame") + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + return src.micName == "mic1" + }) + + client.setSource("mic2") + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + return src.micName == "mic2" + }) + + found := false + for range 200 { + f := <-frames + if f.Mic[0] == 0x08 { + found = true + break + } } - if !bytes.Equal(f.Mic, frame.Silent(pcm.FrameBytes)) { - t.Error("expected silent mic fallback") + if !found { + t.Error("expected mic audio from new source after hot-swap") } } -func TestParec_Stop(t *testing.T) { - closed := false - runner := &fakeRunner{ - outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { - if args[0] == "get-default-sink" { - return []byte("sink\n"), nil - } - return []byte("mic\n"), nil - }, - startFn: func(ctx context.Context, name string, args ...string) (io.ReadCloser, func() error, error) { - return io.NopCloser(bytes.NewReader(nil)), func() error { - closed = true - return nil - }, nil - }, +func TestParec_MicHotSwapFailureRetainsOld(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStream("mic1", bytes.Repeat([]byte{0x07, 0x00}, pcm.FrameBytes/2)) + client.setSource("mic1") + + src := &Parec{client: client, sinks: make(map[string]*reader)} + _, err := src.Start(context.Background()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = src.Stop() }() + + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + return src.micName == "mic1" + }) + + client.mu.Lock() + client.startErrFn = func(device string) error { + if device == "mic2" { + return io.ErrClosedPipe + } + return nil } + client.mu.Unlock() + client.setSource("mic2") - src := NewParec(parec.New(runner)) + time.Sleep(100 * time.Millisecond) + + src.mu.Lock() + name := src.micName + mic := src.mic + src.mu.Unlock() + if name != "mic1" { + t.Errorf("expected mic to remain mic1 after failed swap, got %q", name) + } + if mic == nil { + t.Fatal("expected mic reader to still be present") + } + select { + case <-mic.dead(): + t.Error("old mic reader should still be alive after a failed swap") + default: + } +} + +func TestParec_ReaderExitTriggersRestartWithoutKillingOthers(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.runner.setStreamOnce("a.monitor", bytes.Repeat([]byte{0x05, 0x00}, pcm.FrameBytes/2)) + client.runner.setStream("b.monitor", bytes.Repeat([]byte{0x06, 0x00}, pcm.FrameBytes/2)) + client.setSinks( + parec.Sink{Name: "a", MonitorSource: "a.monitor"}, + parec.Sink{Name: "b", MonitorSource: "b.monitor"}, + ) + + src := &Parec{client: client, sinks: make(map[string]*reader)} _, err := src.Start(context.Background()) if err != nil { t.Fatal(err) } + defer func() { _ = src.Stop() }() + + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, hasA := src.sinks["a"] + _, hasB := src.sinks["b"] + return hasA && hasB + }) + + src.mu.Lock() + bReaderBefore := src.sinks["b"] + src.mu.Unlock() + + // Once "a"'s one-shot stream is exhausted, its reader dies; the next + // sink reconciliation pass must restart it in place without touching "b". + waitFor(t, time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + a, ok := src.sinks["a"] + if !ok { + return false + } + select { + case <-a.dead(): + return false // not yet restarted + default: + return true + } + }) + + src.mu.Lock() + bReaderAfter := src.sinks["b"] + src.mu.Unlock() + if bReaderBefore != bReaderAfter { + t.Error("sink b's reader should not have been restarted") + } +} + +func TestParec_PulseAudioUnreachableThenRecovers(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + client.setSinksErr(io.ErrClosedPipe) + client.setSourceErr(io.ErrClosedPipe) + + src := &Parec{client: client, sinks: make(map[string]*reader)} + frames, err := src.Start(context.Background()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = src.Stop() }() + + for range 3 { + f := <-frames + if !isSilent(f.Sys) || !isSilent(f.Mic) { + t.Error("expected silence while PulseAudio is unreachable") + } + } + + client.runner.setStream("a.monitor", bytes.Repeat([]byte{0x09, 0x00}, pcm.FrameBytes/2)) + client.setSinksErr(nil) + client.setSinks(parec.Sink{Name: "a", MonitorSource: "a.monitor"}) + client.setSourceErr(nil) + client.setSource("mic1") + client.runner.setStream("mic1", bytes.Repeat([]byte{0x0a, 0x00}, pcm.FrameBytes/2)) + + waitFor(t, 2*time.Second, func() bool { + src.mu.Lock() + defer src.mu.Unlock() + _, hasSink := src.sinks["a"] + return hasSink && src.micName == "mic1" + }) +} + +func isSilent(data []byte) bool { + for _, b := range data { + if b != 0 { + return false + } + } + return true +} + +func TestParec_Stop_Idempotent(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + src := &Parec{client: client, sinks: make(map[string]*reader)} + if _, err := src.Start(context.Background()); err != nil { + t.Fatal(err) + } + if err := src.Stop(); err != nil { + t.Fatal(err) + } if err := src.Stop(); err != nil { t.Fatal(err) } - if !closed { - t.Error("expected capture streams to close on Stop") +} + +func TestParec_Backpressure_SlowConsumerDoesNotDropFrames(t *testing.T) { + withFastPolling(t) + client := newFakeSinkClient() + src := &Parec{client: client, sinks: make(map[string]*reader)} + frames, err := src.Start(context.Background()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = src.Stop() }() + + // Don't read frames for a while; mixLoop must block on send rather than + // panic or busy-loop, and Stop() must still terminate cleanly afterward. + time.Sleep(50 * time.Millisecond) + + drained := 0 + timeout := time.After(time.Second) +loop: + for drained < 3 { + select { + case <-frames: + drained++ + case <-timeout: + break loop + } + } + if drained < 3 { + t.Errorf("expected to drain buffered/backlogged frames, got %d", drained) } } diff --git a/internal/capture/reader.go b/internal/capture/reader.go new file mode 100644 index 0000000..b6dc0ec --- /dev/null +++ b/internal/capture/reader.go @@ -0,0 +1,74 @@ +package capture + +import ( + "sync" + + "github.com/odsod/recorder/internal/audio/frame" + "github.com/odsod/recorder/internal/audio/pcm" + "github.com/odsod/recorder/internal/protocol/parec" +) + +// reader continuously reads frames from a parec.CaptureStream and makes the +// most recent frame available without blocking capture. A slow or absent +// consumer never backs up the underlying parec process. +type reader struct { + stream *parec.CaptureStream + + frameCh chan []byte + done chan struct{} + + stopOnce sync.Once +} + +// startReader begins reading 1-second frames from stream in the background. +func startReader(stream *parec.CaptureStream) *reader { + r := &reader{ + stream: stream, + frameCh: make(chan []byte, 1), + done: make(chan struct{}), + } + go r.run() + return r +} + +func (r *reader) run() { + defer close(r.done) + for { + data, err := frame.Read(r.stream, pcm.FrameBytes) + if err != nil { + return + } + select { + case r.frameCh <- data: + default: + <-r.frameCh + r.frameCh <- data + } + } +} + +// take returns the most recently read frame if one is available. It never +// blocks: if no frame has arrived since the last take, ok is false. +func (r *reader) take() (data []byte, ok bool) { + select { + case data := <-r.frameCh: + return data, true + default: + return nil, false + } +} + +// dead returns a channel that is closed when the read loop has exited +// (stream EOF or error). +func (r *reader) dead() <-chan struct{} { + return r.done +} + +// stop terminates the underlying capture stream. Idempotent. +func (r *reader) stop() error { + var err error + r.stopOnce.Do(func() { + err = r.stream.Close() + }) + return err +} diff --git a/internal/capture/reconcile.go b/internal/capture/reconcile.go new file mode 100644 index 0000000..f337adf --- /dev/null +++ b/internal/capture/reconcile.go @@ -0,0 +1,218 @@ +package capture + +import ( + "context" + "log/slog" + "maps" + "time" + + "github.com/odsod/recorder/internal/audio/pcm" + "github.com/odsod/recorder/internal/protocol/parec" +) + +// Poll intervals for sink and microphone discovery. Sinks and default +// sources change on human timescales (plugging in headphones, changing +// output device in a settings applet); periodic polling is sufficient and +// avoids a pactl subscribe event subsystem. +var ( + sinkPollInterval = 2 * time.Second + micPollInterval = 2 * time.Second +) + +// backoff tracks a reset-on-success retry delay for one recovering +// operation. Not goroutine-safe by design: each failure axis (sinks, mic) +// owns its own instance, so failures on one never throttle the other. +type backoff struct { + cur time.Duration +} + +var ( + backoffInitial = 500 * time.Millisecond + backoffMax = 30 * time.Second +) + +func (b *backoff) next() time.Duration { + if b.cur == 0 { + b.cur = backoffInitial + } else { + b.cur = min(b.cur*2, backoffMax) + } + return b.cur +} + +func (b *backoff) reset() { + b.cur = 0 +} + +func (c *Parec) reconcileLoop(ctx context.Context) { + sinkTicker := time.NewTicker(sinkPollInterval) + defer sinkTicker.Stop() + micTicker := time.NewTicker(micPollInterval) + defer micTicker.Stop() + + c.reconcileSinks(ctx) + c.reconcileMic(ctx) + + for { + select { + case <-ctx.Done(): + return + case <-sinkTicker.C: + c.reconcileSinks(ctx) + case <-micTicker.C: + c.reconcileMic(ctx) + } + } +} + +// reconcileSinks lists the current sinks and diffs them against active +// readers by name. A transient ListSinks failure leaves existing readers +// running untouched. +// +// Known limitation: virtual/loopback sink chains (e.g. an EasyEffects +// post-processing sink chained after a raw device sink) will be monitored +// as separate sinks and double-counted when mixed. Not solved here. +func (c *Parec) reconcileSinks(ctx context.Context) { + if time.Since(c.sinkLastTry) < c.sinkBackoff.cur { + return + } + c.sinkLastTry = time.Now() + + resp, err := c.client.ListSinks(ctx, parec.ListSinksRequest{}) + if err != nil { + slog.WarnContext(ctx, "list sinks failed", + "err", err, + "retryIn", c.sinkBackoff.next(), + ) + return + } + c.sinkBackoff.reset() + + wanted := make(map[string]parec.Sink, len(resp.Sinks)) + for _, s := range resp.Sinks { + wanted[s.Name] = s + } + + c.mu.Lock() + current := make(map[string]*reader, len(c.sinks)) + maps.Copy(current, c.sinks) + c.mu.Unlock() + + for name, r := range current { + if _, ok := wanted[name]; !ok { + c.stopSink(name, r, "sink disappeared") + continue + } + select { + case <-r.dead(): + c.stopSink(name, r, "sink reader exited") + default: + } + } + + for name, sink := range wanted { + c.mu.Lock() + _, active := c.sinks[name] + c.mu.Unlock() + if active { + continue + } + + stream, err := c.client.StartCapture(ctx, parec.StartCaptureRequest{ + Device: sink.MonitorSource, SampleRate: pcm.SampleRate, + }) + if err != nil { + slog.WarnContext(ctx, "sink capture start failed", + "sink", name, + "err", err, + ) + continue + } + + r := startReader(stream) + c.mu.Lock() + c.sinks[name] = r + c.mu.Unlock() + slog.InfoContext(ctx, "sink capture started", "sink", name) + } +} + +func (c *Parec) stopSink(name string, r *reader, reason string) { + c.mu.Lock() + delete(c.sinks, name) + c.mu.Unlock() + _ = r.stop() + slog.InfoContext(context.Background(), "sink capture stopped", + "sink", name, + "reason", reason, + ) +} + +// reconcileMic resolves the default microphone source and swaps to it only +// after the replacement capture starts successfully, so a failed swap never +// drops below one working mic reader. +func (c *Parec) reconcileMic(ctx context.Context) { + if time.Since(c.micLastTry) < c.micBackoff.cur { + return + } + c.micLastTry = time.Now() + + resp, err := c.client.GetDefaultSource(ctx, parec.GetDefaultSourceRequest{}) + if err != nil { + slog.WarnContext(ctx, "get default source failed", + "err", err, + "retryIn", c.micBackoff.next(), + ) + return + } + c.micBackoff.reset() + + c.mu.Lock() + name := resp.Source + current := c.mic + currentName := c.micName + c.mu.Unlock() + + dead := false + if current != nil { + select { + case <-current.dead(): + dead = true + default: + } + } + + if current != nil && !dead && name == currentName { + return + } + + stream, err := c.client.StartCapture(ctx, parec.StartCaptureRequest{ + Device: name, SampleRate: pcm.SampleRate, + }) + if err != nil { + slog.WarnContext(ctx, "mic capture start failed", + "source", name, + "err", err, + ) + return + } + + r := startReader(stream) + c.mu.Lock() + c.mic = r + c.micName = name + c.mu.Unlock() + + if current != nil { + _ = current.stop() + } + + if currentName == "" { + slog.InfoContext(ctx, "mic capture started", "source", name) + } else { + slog.InfoContext(ctx, "default microphone changed", + "old", currentName, + "new", name, + ) + } +} diff --git a/internal/protocol/parec/client.go b/internal/protocol/parec/client.go index 4727c3b..5db0338 100644 --- a/internal/protocol/parec/client.go +++ b/internal/protocol/parec/client.go @@ -4,13 +4,14 @@ // command for streaming raw PCM audio. The CommandRunner interface allows // injecting a fake for testing without real PulseAudio. // -// Query operations (GetDefaultSink, GetDefaultSource) follow the standard +// Query operations (ListSinks, GetDefaultSource) follow the standard // request/response pattern. StartCapture returns a CaptureStream that // implements io.Reader for continuous audio data and io.Closer to stop capture. package parec import ( "context" + "encoding/json" "fmt" "io" "os/exec" @@ -70,27 +71,6 @@ func NewDefault() *Client { return &Client{runner: ExecRunner{}} } -// GetDefaultSinkRequest is empty; the default sink is a system-global query. -type GetDefaultSinkRequest struct{} - -// GetDefaultSinkResponse contains the monitor source for the default output device. -type GetDefaultSinkResponse struct { - // MonitorSource is the PulseAudio source name for capturing system audio - // (the default sink name with ".monitor" appended). - MonitorSource string -} - -// GetDefaultSink queries the system's default audio output device. -func (c *Client) GetDefaultSink(ctx context.Context, _ GetDefaultSinkRequest) (GetDefaultSinkResponse, error) { - out, err := c.runner.Output(ctx, "pactl", "get-default-sink") - if err != nil { - return GetDefaultSinkResponse{}, fmt.Errorf("pactl get-default-sink: %w", err) - } - return GetDefaultSinkResponse{ - MonitorSource: strings.TrimSpace(string(out)) + ".monitor", - }, nil -} - // GetDefaultSourceRequest is empty; the default source is a system-global query. type GetDefaultSourceRequest struct{} @@ -111,6 +91,43 @@ func (c *Client) GetDefaultSource(ctx context.Context, _ GetDefaultSourceRequest }, nil } +// ListSinksRequest is empty; sinks are a system-global query. +type ListSinksRequest struct{} + +// Sink describes one PulseAudio sink. Name is its stable identity across +// polls; numeric indexes are reused by PipeWire across device churn. +type Sink struct { + // Name is the sink's PulseAudio name. + Name string + // MonitorSource is the PulseAudio source name for capturing audio played + // through this sink (Name with ".monitor" appended). + MonitorSource string +} + +// ListSinksResponse contains all currently known sinks. +type ListSinksResponse struct { + Sinks []Sink +} + +// ListSinks enumerates all PulseAudio sinks. +func (c *Client) ListSinks(ctx context.Context, _ ListSinksRequest) (ListSinksResponse, error) { + out, err := c.runner.Output(ctx, "pactl", "--format=json", "list", "sinks") + if err != nil { + return ListSinksResponse{}, fmt.Errorf("pactl list sinks: %w", err) + } + var wire []struct { + Name string `json:"name"` + } + if err := json.Unmarshal(out, &wire); err != nil { + return ListSinksResponse{}, fmt.Errorf("parse pactl sinks json: %w", err) + } + sinks := make([]Sink, 0, len(wire)) + for _, s := range wire { + sinks = append(sinks, Sink{Name: s.Name, MonitorSource: s.Name + ".monitor"}) + } + return ListSinksResponse{Sinks: sinks}, nil +} + // StartCaptureRequest specifies the audio device and format for streaming capture. type StartCaptureRequest struct { // Device is the PulseAudio source name to capture from. diff --git a/internal/protocol/parec/client_test.go b/internal/protocol/parec/client_test.go index 5e4448c..5510642 100644 --- a/internal/protocol/parec/client_test.go +++ b/internal/protocol/parec/client_test.go @@ -24,70 +24,122 @@ func (f *fakeRunner) Start(ctx context.Context, name string, args ...string) (io return f.startFn(ctx, name, args...) } -func TestGetDefaultSink(t *testing.T) { +func TestGetDefaultSource(t *testing.T) { runner := &fakeRunner{ outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { if name != "pactl" { t.Errorf("expected pactl, got %s", name) } - if len(args) != 1 || args[0] != "get-default-sink" { + if len(args) != 1 || args[0] != "get-default-source" { t.Errorf("unexpected args: %v", args) } - return []byte("alsa_output.pci\n"), nil + return []byte("alsa_input.usb\n"), nil }, } client := parec.New(runner) - resp, err := client.GetDefaultSink(context.Background(), parec.GetDefaultSinkRequest{}) + resp, err := client.GetDefaultSource(context.Background(), parec.GetDefaultSourceRequest{}) if err != nil { t.Fatal(err) } - if resp.MonitorSource != "alsa_output.pci.monitor" { - t.Errorf("expected 'alsa_output.pci.monitor', got %q", resp.MonitorSource) + if resp.Source != "alsa_input.usb" { + t.Errorf("expected 'alsa_input.usb', got %q", resp.Source) } } -func TestGetDefaultSource(t *testing.T) { +func TestGetDefaultSource_Error(t *testing.T) { + runner := &fakeRunner{ + outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { + return nil, errors.New("command not found") + }, + } + + client := parec.New(runner) + _, err := client.GetDefaultSource(context.Background(), parec.GetDefaultSourceRequest{}) + if err == nil { + t.Fatal("expected error") + } + if !strings.Contains(err.Error(), "pactl get-default-source") { + t.Errorf("expected wrapped error, got: %v", err) + } +} + +func TestListSinks(t *testing.T) { runner := &fakeRunner{ outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { if name != "pactl" { t.Errorf("expected pactl, got %s", name) } - if len(args) != 1 || args[0] != "get-default-source" { - t.Errorf("unexpected args: %v", args) + expectedArgs := []string{"--format=json", "list", "sinks"} + if len(args) != len(expectedArgs) { + t.Fatalf("expected %d args, got %d: %v", len(expectedArgs), len(args), args) } - return []byte("alsa_input.usb\n"), nil + for i, exp := range expectedArgs { + if args[i] != exp { + t.Errorf("arg %d: expected %q, got %q", i, exp, args[i]) + } + } + return []byte(`[ + {"name": "alsa_output.pci.hdmi"}, + {"name": "bluez_output.usb"} + ]`), nil }, } client := parec.New(runner) - resp, err := client.GetDefaultSource(context.Background(), parec.GetDefaultSourceRequest{}) + resp, err := client.ListSinks(context.Background(), parec.ListSinksRequest{}) if err != nil { t.Fatal(err) } - if resp.Source != "alsa_input.usb" { - t.Errorf("expected 'alsa_input.usb', got %q", resp.Source) + want := []parec.Sink{ + {Name: "alsa_output.pci.hdmi", MonitorSource: "alsa_output.pci.hdmi.monitor"}, + {Name: "bluez_output.usb", MonitorSource: "bluez_output.usb.monitor"}, + } + if len(resp.Sinks) != len(want) { + t.Fatalf("expected %d sinks, got %d: %v", len(want), len(resp.Sinks), resp.Sinks) + } + for i, w := range want { + if resp.Sinks[i] != w { + t.Errorf("sink %d: expected %+v, got %+v", i, w, resp.Sinks[i]) + } } } -func TestGetDefaultSink_Error(t *testing.T) { +func TestListSinks_Empty(t *testing.T) { runner := &fakeRunner{ outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { - return nil, errors.New("command not found") + return []byte(`[]`), nil }, } client := parec.New(runner) - _, err := client.GetDefaultSink(context.Background(), parec.GetDefaultSinkRequest{}) + resp, err := client.ListSinks(context.Background(), parec.ListSinksRequest{}) + if err != nil { + t.Fatal(err) + } + if len(resp.Sinks) != 0 { + t.Errorf("expected no sinks, got %v", resp.Sinks) + } +} + +func TestListSinks_InvalidJSON(t *testing.T) { + runner := &fakeRunner{ + outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { + return []byte(`not json`), nil + }, + } + + client := parec.New(runner) + _, err := client.ListSinks(context.Background(), parec.ListSinksRequest{}) if err == nil { t.Fatal("expected error") } - if !strings.Contains(err.Error(), "pactl get-default-sink") { + if !strings.Contains(err.Error(), "parse pactl sinks json") { t.Errorf("expected wrapped error, got: %v", err) } } -func TestGetDefaultSource_Error(t *testing.T) { +func TestListSinks_CommandError(t *testing.T) { runner := &fakeRunner{ outputFn: func(ctx context.Context, name string, args ...string) ([]byte, error) { return nil, errors.New("command not found") @@ -95,11 +147,11 @@ func TestGetDefaultSource_Error(t *testing.T) { } client := parec.New(runner) - _, err := client.GetDefaultSource(context.Background(), parec.GetDefaultSourceRequest{}) + _, err := client.ListSinks(context.Background(), parec.ListSinksRequest{}) if err == nil { t.Fatal("expected error") } - if !strings.Contains(err.Error(), "pactl get-default-source") { + if !strings.Contains(err.Error(), "pactl list sinks") { t.Errorf("expected wrapped error, got: %v", err) } } diff --git a/internal/recorder/capture.go b/internal/recorder/capture.go index 11ffebf..c385d56 100644 --- a/internal/recorder/capture.go +++ b/internal/recorder/capture.go @@ -25,13 +25,6 @@ func (r *Recorder) captureLoop(ctx context.Context, chunkCh chan<- AudioChunk) { } defer func() { _ = r.svc.Capture.Stop() }() - slog.InfoContext(ctx, "system source configured", - "source", r.svc.Capture.MonitorSource(), - ) - slog.InfoContext(ctx, "mic source configured", - "source", r.svc.Capture.MicSource(), - ) - accum := chunk.New(chunk.DefaultConfig()) audioGate := gate.Default() var wasSpeech bool @@ -47,6 +40,10 @@ func (r *Recorder) captureLoop(ctx context.Context, chunkCh chan<- AudioChunk) { return case frame, ok := <-frames: if !ok { + // capture.Source only closes this channel on ctx.Done(); + // outages (no sinks, dead parec, unreachable PulseAudio) + // surface as silent frame.Dual values instead, so reaching + // here means real shutdown. if out, ok := accum.Flush(); ok { r.emitChunk(ctx, out.SysPCM, out.MicPCM, out.StartTime, audioGate, chunkCh) }