Parent directory

recorder.go

4367 bytes
  1package stt
  2
  3import (
  4	"context"
  5	"errors"
  6	"fmt"
  7	"os"
  8	"os/exec"
  9	"strings"
 10	"sync"
 11	"syscall"
 12	"time"
 13)
 14
 15// stopGrace is how long the recorder waits for a graceful SIGINT shutdown
 16// before escalating to SIGKILL.
 17const stopGrace = 500 * time.Millisecond
 18
 19// Recorder wraps the configured capture command for push-to-talk recording. It
 20// is an HTTP/process client only and never manages any model service. The
 21// command string is split on whitespace; the temp WAV output path is appended
 22// as the final argument at Start time.
 23type Recorder struct {
 24	command   []string
 25	maxRecord time.Duration
 26}
 27
 28func NewRecorder(command string, maxRecord time.Duration) *Recorder {
 29	return &Recorder{command: strings.Fields(command), maxRecord: maxRecord}
 30}
 31
 32// Recording is a live push-to-talk capture. The recording ends when the caller
 33// calls Stop, the start context is cancelled, or the max duration elapses; in
 34// all cases Done closes and the temp WAV path becomes available from Stop.
 35type Recording struct {
 36	cmd       *exec.Cmd
 37	program   string
 38	wavPath   string
 39	maxRecord time.Duration
 40
 41	stopCh      chan struct{}
 42	doneCh      chan struct{}
 43	cancel      context.CancelFunc
 44	stopReqOnce sync.Once
 45	stopOnce    sync.Once
 46	weStopped   bool
 47	resultPath  string
 48	resultErr   error
 49}
 50
 51// Start launches the capture command with a fresh temp WAV path appended as the
 52// final argument. The provided context is respected: cancelling it stops the
 53// recording, and maxRecord caps how long it may run.
 54func (r *Recorder) Start(ctx context.Context) (*Recording, error) {
 55	if len(r.command) == 0 {
 56		return nil, errors.New("stt: empty recorder command")
 57	}
 58
 59	f, err := os.CreateTemp("", "kaiwari-record-*.wav")
 60	if err != nil {
 61		return nil, fmt.Errorf("stt: create temp wav: %w", err)
 62	}
 63	wavPath := f.Name()
 64	_ = f.Close()
 65
 66	args := make([]string, 0, len(r.command)+1)
 67	args = append(args, r.command...)
 68	args = append(args, wavPath)
 69
 70	cmd := exec.Command(args[0], args[1:]...)
 71	if err := cmd.Start(); err != nil {
 72		_ = os.Remove(wavPath)
 73		return nil, fmt.Errorf("stt: start recorder %q: %w", r.command[0], err)
 74	}
 75
 76	recCtx, cancel := context.WithCancel(ctx)
 77	rec := &Recording{
 78		cmd:       cmd,
 79		program:   r.command[0],
 80		wavPath:   wavPath,
 81		maxRecord: r.maxRecord,
 82		stopCh:    make(chan struct{}),
 83		doneCh:    make(chan struct{}),
 84		cancel:    cancel,
 85	}
 86	go rec.watchdog(recCtx)
 87	return rec, nil
 88}
 89
 90// Stop requests a graceful stop (SIGINT, then SIGKILL after a short grace
 91// period), waits for the process to be reaped, and returns the temp WAV path.
 92// The file is left on disk for the ASR step, which deletes it after
 93// uploading. If the recorder command exited on its own with an error, the temp
 94// file is removed and that error is returned. Stop is safe to call more than
 95// once.
 96func (rec *Recording) Stop() (string, error) {
 97	if rec == nil || rec.cmd == nil {
 98		return "", errors.New("stt: recording not started")
 99	}
100	rec.stopReqOnce.Do(func() { close(rec.stopCh) })
101	<-rec.doneCh
102	rec.cancel()
103	return rec.resultPath, rec.resultErr
104}
105
106func (rec *Recording) watchdog(ctx context.Context) {
107	var timerC <-chan time.Time
108	if rec.maxRecord > 0 {
109		t := time.NewTimer(rec.maxRecord)
110		defer t.Stop()
111		timerC = t.C
112	}
113	select {
114	case <-ctx.Done():
115	case <-timerC:
116	case <-rec.stopCh:
117	}
118
119	rec.terminate()
120	waitErr := rec.cmd.Wait()
121	switch {
122	case rec.weStopped:
123		rec.resultPath = rec.wavPath
124	case waitErr != nil:
125		_ = os.Remove(rec.wavPath)
126		rec.resultErr = fmt.Errorf("stt: recorder %q exited: %w", rec.program, waitErr)
127	default:
128		rec.resultPath = rec.wavPath
129	}
130	close(rec.doneCh)
131}
132
133func (rec *Recording) terminate() {
134	rec.stopOnce.Do(func() {
135		if gracefulStop(rec.cmd.Process) {
136			rec.weStopped = true
137		}
138	})
139}
140
141// gracefulStop signals the process with SIGINT, waits up to stopGrace for it to
142// exit, then escalates to SIGKILL. It reports whether it stopped a live
143// process (as opposed to finding an already-dead one).
144func gracefulStop(p *os.Process) bool {
145	if p == nil {
146		return false
147	}
148	if err := p.Signal(syscall.SIGINT); err != nil {
149		return false
150	}
151	deadline := time.Now().Add(stopGrace)
152	for time.Now().Before(deadline) {
153		if !processRunning(p) {
154			return true
155		}
156		time.Sleep(20 * time.Millisecond)
157	}
158	_ = p.Kill()
159	return true
160}
161
162func processRunning(p *os.Process) bool {
163	return p.Signal(syscall.Signal(0)) == nil
164}