Files
OwnCord/Server/ws/livekit_process.go
T
J3vbandClaude Opus 5 0594a130cb fix(ws): give the LiveKit health check its own HTTP transport (#1356)
NewLiveKitProcess built its health-check http.Client without a Transport,
so it fell back to the process-wide http.DefaultTransport.

httptest.Server.Close calls CloseIdleConnections on http.DefaultTransport
by design ("assume most users of httptest.Server will be using the standard
transport, so help them out"), and ws is full of t.Parallel tests that each
defer srv.Close(). Any one of them finishing while a health check held a
pooled connection severed that request:

  livekit_test.go:978: HealthCheck: livekit health check failed:
    Get "http://127.0.0.1:41343": net/http: HTTP/1.x transport connection
    broken: http: CloseIdleConnections called

That surfaced as an unrelated-looking CI failure on a TypeScript lint bump
(#1341). It is not purely a test artifact: in production the health check
also shared one connection pool with every other DefaultTransport user in
the server process.

Cloning DefaultTransport keeps its tuned defaults (proxy, dial and TLS
timeouts, HTTP/2) while giving the client a private pool.

Locked by TestHealthCheckClientOwnsItsTransport, which fails on the
unfixed constructor.

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-11 19:56:54 +02:00

411 lines
12 KiB
Go

// Package ws provides the LiveKit companion process manager.
//
// LiveKitProcess manages the lifecycle of a livekit-server binary running
// alongside chatserver. It auto-generates a minimal livekit.yaml config,
// starts the process, monitors health, and restarts on crash.
package ws
import (
"context"
"fmt"
"log/slog"
"net/http"
"os"
"os/exec"
"path/filepath"
"strings"
"time"
"github.com/owncord/server/config"
"github.com/owncord/server/syncutil"
)
// LiveKitProcess manages the companion livekit-server binary.
type LiveKitProcess struct {
cfg *config.VoiceConfig
tlsCfg *config.TLSConfig
dataDir string
httpClient *http.Client // for health checks — no redirect following
mu syncutil.Mutex
cmd *exec.Cmd
cancel context.CancelFunc
stopped bool
runDone chan struct{} // closed by runLoop when cmd.Wait() returns
loopDone chan struct{} // closed when runLoop exits entirely
}
// NewLiveKitProcess creates a new process manager. It does not start the
// process — call Start() to launch the LiveKit server.
// The tlsCfg is used to configure TLS on the LiveKit server using the same
// certs as OwnCord (avoids mixed-content blocks in WebView2).
func NewLiveKitProcess(cfg *config.VoiceConfig, tlsCfg *config.TLSConfig, dataDir string) *LiveKitProcess {
return &LiveKitProcess{
cfg: cfg,
tlsCfg: tlsCfg,
dataDir: dataDir,
httpClient: &http.Client{
// Own the transport rather than inheriting http.DefaultTransport:
// its pool is shared process-wide, and httptest.Server.Close closes
// its idle connections, which severs in-flight health checks.
Transport: http.DefaultTransport.(*http.Transport).Clone(),
CheckRedirect: func(*http.Request, []*http.Request) error {
return http.ErrUseLastResponse
},
},
}
}
// autoGeneratedMarker identifies a livekit.yaml written by OwnCord. A file
// without this marker is treated as user-managed and never overwritten.
const autoGeneratedMarker = "Auto-generated by OwnCord"
// generateConfig writes a minimal livekit.yaml for the companion process.
// If the file exists and lacks the auto-generated marker, it is treated as
// user-managed and left untouched so operators can set LiveKit options
// OwnCord does not model (ips.includes, interfaces, stun_servers, ...).
func (p *LiveKitProcess) generateConfig() (string, error) {
cfgPath := filepath.Join(p.dataDir, "livekit.yaml")
if existing, err := os.ReadFile(cfgPath); err == nil {
// An empty/whitespace-only file is a truncated leftover, not a
// user-managed config — regenerate it rather than wedging LiveKit.
if len(strings.TrimSpace(string(existing))) > 0 &&
!strings.Contains(string(existing), autoGeneratedMarker) {
slog.Info("livekit: livekit.yaml is user-managed (auto-generated marker absent), not overwriting", "path", cfgPath)
slog.Warn("livekit: ensure the keys entry in your livekit.yaml matches voice.livekit_api_key / voice.livekit_api_secret")
return cfgPath, nil
}
} else if !os.IsNotExist(err) {
return "", fmt.Errorf("reading existing livekit config: %w", err)
}
// No TURN TLS config — LiveKit signaling is proxied through OwnCord's
// HTTPS server at /livekit/*, so no separate TLS is needed on LiveKit.
// Sanitize credentials for safe YAML interpolation: reject strings
// containing characters that could break YAML structure.
//
// The error names the offending config field but never echoes the
// offending byte: this error is wrapped by Start() and logged by the
// caller, so quoting a character from the key or secret would write a
// byte of a credential to the server log in clear text.
const unsafeYAML = ":#{}\n\r\"\\"
for _, cred := range []struct {
field string
value string
}{
{"voice.livekit_api_key", p.cfg.LiveKitAPIKey},
{"voice.livekit_api_secret", p.cfg.LiveKitAPISecret},
} {
if strings.ContainsAny(cred.value, unsafeYAML) {
return "", fmt.Errorf(`%s contains a character that cannot be safely written to livekit.yaml (one of : # { } " \ CR LF)`, cred.field)
}
}
// Build node_ip line only when configured (required for remote users behind NAT).
// Validate: must be a plain IP address (no YAML-breaking chars).
nodeIPLine := ""
if p.cfg.NodeIP != "" {
for _, ch := range p.cfg.NodeIP {
if ch == '"' || ch == '\\' || ch == '\n' || ch == '\r' || ch == '#' || ch == '{' || ch == '}' {
return "", fmt.Errorf("node_ip contains unsafe character %q", string(ch))
}
}
nodeIPLine = fmt.Sprintf("\n node_ip: %q", p.cfg.NodeIP)
}
// Advertise LAN host candidates alongside the external mapping so clients
// on the local network can reach a dual-homed (LAN + public IP) server.
advertiseInternalLine := ""
if p.cfg.AdvertiseInternalIP {
advertiseInternalLine = "\n advertise_internal_ip: true"
}
content := fmt.Sprintf(`# Auto-generated by OwnCord — regenerated on every server start.
# To manage this file yourself (custom rtc options, multiple interfaces, etc.),
# delete the first line above; OwnCord will then leave the file untouched.
# Your keys entry must still match voice.livekit_api_key /
# voice.livekit_api_secret in config.yaml.
port: 7880
rtc:
port_range_start: 50000
port_range_end: 60000
use_external_ip: true%s%s
pli_throttle:
low_quality: 500ms
mid_quality: 1s
high_quality: 1s
keys:
"%s": "%s"
logging:
level: info
`, nodeIPLine, advertiseInternalLine, p.cfg.LiveKitAPIKey, p.cfg.LiveKitAPISecret)
if err := os.MkdirAll(p.dataDir, 0o750); err != nil {
return "", fmt.Errorf("creating data dir: %w", err)
}
if err := os.WriteFile(cfgPath, []byte(content), 0o600); err != nil {
return "", fmt.Errorf("writing livekit config: %w", err)
}
return cfgPath, nil
}
// Start launches the livekit-server binary. If LiveKitBinaryPath is empty
// and auto-download is disabled, this is a no-op (assumes LiveKit is managed
// externally). With auto-download enabled, the binary is fetched (and
// checksum-verified) in the background before the process starts, so server
// startup is never blocked on the download.
func (p *LiveKitProcess) Start() error {
if p.cfg.LiveKitBinaryPath == "" && !p.cfg.AutoDownloadLiveKit {
slog.Info("livekit: no binary path configured, assuming externally managed")
return nil
}
p.mu.Lock()
defer p.mu.Unlock()
if p.cmd != nil {
return fmt.Errorf("livekit process already running")
}
cfgPath, err := p.generateConfig()
if err != nil {
return fmt.Errorf("generating livekit config: %w", err)
}
ctx, cancel := context.WithCancel(context.Background())
p.cancel = cancel
p.loopDone = make(chan struct{})
go func() {
defer func() {
p.mu.Lock()
if p.loopDone != nil {
close(p.loopDone)
}
p.mu.Unlock()
}()
binPath := p.cfg.LiveKitBinaryPath
if binPath == "" {
resolved, dlErr := p.resolveBinary(ctx)
if dlErr != nil {
if ctx.Err() == nil {
slog.Error("livekit: auto-download failed — voice stays offline until livekit-server is available",
"error", dlErr,
"hint", "check the network, or set voice.livekit_binary in config.yaml to a binary you provide")
}
return
}
binPath = resolved
}
p.runLoop(ctx, cfgPath, binPath)
}()
return nil
}
// resolveBinary downloads (or reuses a cached copy of) the pinned
// livekit-server release, retrying a few times so a transient network hiccup
// at boot does not permanently disable voice until the next restart.
func (p *LiveKitProcess) resolveBinary(ctx context.Context) (string, error) {
const (
attempts = 3
retryDelay = 15 * time.Second
)
var lastErr error
for i := range attempts {
if i > 0 {
slog.Warn("livekit: retrying download", "attempt", i+1, "error", lastErr)
select {
case <-time.After(retryDelay):
case <-ctx.Done():
return "", ctx.Err()
}
}
attemptCtx, cancel := context.WithTimeout(ctx, 10*time.Minute)
path, err := EnsureLiveKitBinary(attemptCtx, p.dataDir, p.cfg.LiveKitVersion)
cancel()
if err == nil {
return path, nil
}
lastErr = err
}
return "", lastErr
}
// runLoop starts and restarts the process until stopped or context cancelled.
// Uses exponential backoff (3s → 6s → 12s … up to 60s) and stops after 10
// consecutive rapid failures (process exits within 30 seconds). loopDone is
// closed by the Start goroutine that calls this.
func (p *LiveKitProcess) runLoop(ctx context.Context, cfgPath, binPath string) {
const (
baseDelay = 3 * time.Second
maxDelay = 60 * time.Second
maxRetries = 10
stableAfter = 30 * time.Second // reset counter if process runs longer than this
)
rapidFailures := 0
delay := baseDelay
for {
if ctx.Err() != nil {
return
}
cmd := exec.CommandContext(ctx, binPath, "--config", cfgPath) //nolint:gosec // G204: binary path from trusted server config or verified download
cmd.Stdout = os.Stdout
cmd.Stderr = os.Stderr
cmd.WaitDelay = 6 * time.Second // bound Wait to prevent goroutine leak on Windows
slog.Info("livekit: starting process",
"binary", binPath,
"config", cfgPath,
"rapid_failures", rapidFailures)
startTime := time.Now()
p.mu.Lock()
if p.stopped {
p.mu.Unlock()
return
}
// Start inside the critical section that publishes p.cmd: Start is
// what writes cmd.Process, which IsRunning and Stop read under p.mu —
// started after the unlock, that write races every such read.
err := cmd.Start()
if err == nil {
p.cmd = cmd
p.runDone = make(chan struct{})
}
p.mu.Unlock()
if err == nil {
err = cmd.Wait()
}
p.mu.Lock()
p.cmd = nil
if p.runDone != nil {
close(p.runDone)
p.runDone = nil
}
stopped := p.stopped
p.mu.Unlock()
if stopped || ctx.Err() != nil {
slog.Info("livekit: process stopped")
return
}
// If the process ran for a while, it was stable — reset backoff.
if time.Since(startTime) > stableAfter {
rapidFailures = 0
delay = baseDelay
} else {
rapidFailures++
}
if err != nil {
slog.Error("livekit: process exited unexpectedly",
"error", err,
"rapid_failures", rapidFailures,
"restart_delay", delay)
}
if rapidFailures >= maxRetries {
slog.Error("livekit: too many rapid failures, giving up",
"rapid_failures", rapidFailures)
return
}
select {
case <-time.After(delay):
slog.Info("livekit: restarting process")
case <-ctx.Done():
return
}
// Exponential backoff capped at maxDelay.
delay *= 2
if delay > maxDelay {
delay = maxDelay
}
}
}
// IsRunning returns true if the companion process is currently running.
func (p *LiveKitProcess) IsRunning() bool {
p.mu.Lock()
defer p.mu.Unlock()
return p.cmd != nil && p.cmd.Process != nil
}
// HealthCheck probes the LiveKit HTTP endpoint to verify it is accepting
// connections. Returns true if the server responds (any status code).
func (p *LiveKitProcess) HealthCheck(ctx context.Context) (bool, error) {
httpURL := wsToHTTP(p.cfg.LiveKitURL)
ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, httpURL, nil)
if err != nil {
return false, fmt.Errorf("creating health check request: %w", err)
}
resp, err := p.httpClient.Do(req)
if err != nil {
return false, fmt.Errorf("livekit health check failed: %w", err)
}
_ = resp.Body.Close()
return true, nil
}
// Stop gracefully stops the companion process.
// It cancels the context (which signals runLoop) and waits up to 5 seconds
// for the process to exit. The actual cmd.Wait() is done by runLoop — we only
// monitor the process here to avoid calling exec.Cmd.Wait() twice (which has
// undefined behavior).
func (p *LiveKitProcess) Stop() {
p.mu.Lock()
p.stopped = true
cancel := p.cancel
cmd := p.cmd
done := p.runDone
loopDone := p.loopDone
p.mu.Unlock()
if cancel != nil {
cancel()
}
// Wait for runLoop's cmd.Wait() to return (which closes runDone).
// This avoids calling cmd.Wait() or cmd.Process.Wait() from a second
// goroutine, which is unsafe on Windows.
if done != nil {
select {
case <-done:
slog.Info("livekit: process exited cleanly")
case <-time.After(5 * time.Second):
slog.Warn("livekit: process did not exit in time, killing")
if cmd != nil && cmd.Process != nil {
_ = cmd.Process.Kill()
}
}
}
// Wait for the entire runLoop goroutine to finish, ensuring no
// new iteration can start after Stop returns.
if loopDone != nil {
select {
case <-loopDone:
case <-time.After(5 * time.Second):
}
}
}