359ae81d1f
piper_handler: close stdin/stdout pipes on Start() failure and on WriteString error instead of leaking fds. Propagate WriteString error. worker/client: log SetDeadline errors instead of discarding them. voice/session: pushAudio marshals params inline and returns the marshal error instead of swallowing it via mustParams (removed). tool/matcher: log ListTools errors instead of silently returning an empty allowlist that refuses every act. config: applyDefaults now sets RouterThreshold and ToolTimeout defaults so consumers self-contained defaults are belt-and-suspenders.
191 lines
5.4 KiB
Go
191 lines
5.4 KiB
Go
// worker/client.go — the core-side dialer.
|
|
//
|
|
// Core is the CLIENT of its stt/tts worker modules: it dials them, ships
|
|
// audio bytes (Transcribe) and reply text (Synthesize), reads back the
|
|
// result. One Client ⇒ one conn ⇒ one mutex ⇒ no frame interleaving by
|
|
// construction (a module needing parallelism opens N clients, but stt/tts
|
|
// jobs are serial at single-user scale — the model is the bottleneck, not
|
|
// the wire). ctx cancel ⇒ conn close (a half-sent frame desyncs the stream;
|
|
// teardown is the clean recovery, a fresh Dial is the caller's job on next
|
|
// call). Same instinct as internal/ipc/client.go.
|
|
//
|
|
// The daemon holds two Clients — one for the stt module, one for tts — each
|
|
// behind a stt.Transcriber / tts.Synthesizer interface (see internal/stt,
|
|
// internal/tts). The interface IS the swap seam: a Stub impl satisfies the
|
|
// same interface in-process; the Remote wraps this Client. Core never sees
|
|
// a difference.
|
|
package worker
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"net"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Client — one connection to one worker module. NOT goroutine-safe for
|
|
// concurrent calls on the same conn: the mutex serializes writes per
|
|
// request so a frame doesn't interleave with another, but parallel calls
|
|
// block on each other. Open N clients for N parallel jobs.
|
|
type Client struct {
|
|
path string
|
|
|
|
mu sync.Mutex
|
|
c net.Conn
|
|
dial func() (net.Conn, error)
|
|
}
|
|
|
|
// Dial opens a Client to the worker socket at path. The first call lazily
|
|
// dials; subsequent calls reuse the conn (a fresh dial happens on next call
|
|
// after a teardown). Lazy dial keeps a worker that's restarting from
|
|
// blocking core's startup; core attempts the dial on first use.
|
|
func Dial(path string) *Client {
|
|
return &Client{
|
|
path: path,
|
|
dial: func() (net.Conn, error) {
|
|
return net.Dial("unix", path)
|
|
},
|
|
}
|
|
}
|
|
|
|
// Close releases the conn. Idempotent; subsequent calls re-dial on demand.
|
|
func (c *Client) Close() error {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
if c.c == nil {
|
|
return nil
|
|
}
|
|
err := c.c.Close()
|
|
c.c = nil
|
|
return err
|
|
}
|
|
|
|
// Transcribe calls the transcribe verb on the connected worker module. Raises
|
|
// ErrUnknownMethod when the worker serves synthesize only; the daemon wiring
|
|
// has pointed this Client at the wrong socket, surfaced as a clean error.
|
|
func (c *Client) Transcribe(ctx context.Context, req TranscribeReq) (TranscribeResp, error) {
|
|
var resp TranscribeResp
|
|
err := c.call(ctx, MethodTranscribe, req, &resp)
|
|
return resp, err
|
|
}
|
|
|
|
// Synthesize calls the synthesize verb.
|
|
func (c *Client) Synthesize(ctx context.Context, req SynthesizeReq) (SynthesizeResp, error) {
|
|
var resp SynthesizeResp
|
|
err := c.call(ctx, MethodSynthesize, req, &resp)
|
|
return resp, err
|
|
}
|
|
|
|
func (c *Client) call(ctx context.Context, m Method, params any, out any) error {
|
|
body, err := marshalParams(params)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req := Request{Method: m, Params: body}
|
|
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
|
|
if err := c.ensureConnLocked(ctx); err != nil {
|
|
return err
|
|
}
|
|
// ctx cancel ⇒ close conn so a stalled peer doesn't hang the caller; a
|
|
// half-sent frame would desync the stream. Teardown is the clean
|
|
// recovery; the next call re-dials.
|
|
if dl, ok := ctx.Deadline(); ok {
|
|
if err := c.c.SetDeadline(dl); err != nil {
|
|
log.Printf("worker: set deadline: %v", err)
|
|
}
|
|
} else {
|
|
if err := c.c.SetDeadline(time.Now().Add(defaultCallTimeout)); err != nil {
|
|
log.Printf("worker: set deadline: %v", err)
|
|
}
|
|
}
|
|
defer c.c.SetDeadline(time.Time{})
|
|
|
|
if err := writeFrame(c.c, req); err != nil {
|
|
c.teardownLocked()
|
|
return err
|
|
}
|
|
var resp Response
|
|
if err := readFrame(c.c, &resp); err != nil {
|
|
c.teardownLocked()
|
|
return err
|
|
}
|
|
if resp.Error != nil {
|
|
return hydrate(resp.Error)
|
|
}
|
|
if err := json.Unmarshal(resp.Result, out); err != nil {
|
|
return fmt.Errorf("worker: unmarshal result: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c *Client) ensureConnLocked(ctx context.Context) error {
|
|
if c.c != nil {
|
|
return nil
|
|
}
|
|
connCh := make(chan net.Conn, 1)
|
|
errCh := make(chan error, 1)
|
|
go func() {
|
|
nc, err := c.dial()
|
|
if err != nil {
|
|
errCh <- err
|
|
return
|
|
}
|
|
connCh <- nc
|
|
}()
|
|
select {
|
|
case nc := <-connCh:
|
|
c.c = nc
|
|
return nil
|
|
case err := <-errCh:
|
|
return fmt.Errorf("worker: dial %s: %w", c.path, err)
|
|
case <-ctx.Done():
|
|
// the dial goroutine will finish or not; the conn if any is leaked
|
|
// to GC. acceptable — dial failures are rare, and a leaked closed
|
|
// socket is harmless.
|
|
return ctx.Err()
|
|
}
|
|
}
|
|
|
|
func (c *Client) teardownLocked() {
|
|
if c.c != nil {
|
|
_ = c.c.Close()
|
|
c.c = nil
|
|
}
|
|
}
|
|
|
|
// defaultCallTimeout — a worker that hangs > 60s is dead or stuck on a
|
|
// model forward pass that's gone off the rails. Surface as a timeout
|
|
// instead of hanging the daemon's tick / voice path. Production may override
|
|
// via ctx (a long meeting-record transcription, post-MVP).
|
|
const defaultCallTimeout = 60 * time.Second
|
|
|
|
func hydrate(e *RpcError) error {
|
|
switch e.Code {
|
|
case codeUnknownMethod:
|
|
return fmt.Errorf("%w: %s", ErrUnknownMethod, e.Message)
|
|
case codeBadParams:
|
|
return fmt.Errorf("%w: %s", ErrBadParams, e.Message)
|
|
default:
|
|
if e.Message != "" {
|
|
return fmt.Errorf("worker: %s: %s", e.Code, e.Message)
|
|
}
|
|
return fmt.Errorf("worker: %s", e.Code)
|
|
}
|
|
}
|
|
|
|
func marshalParams(v any) (json.RawMessage, error) {
|
|
if v == nil {
|
|
return nil, nil
|
|
}
|
|
b, err := json.Marshal(v)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("worker: marshal params: %w", err)
|
|
}
|
|
return b, nil
|
|
} |