refactor(voice): replace HTTP with stdin/stdout IPC for sr-voice workers
Gemma 2 T&C compliance: exposed HTTP ports allow mods or external code to reach the model, complicating license enforcement. Switch to piped stdin/stdout (JSONL protocol) so the model is only reachable through the game server's internal queue. - worker.rs: VoicePipe owns Child + piped stdin/stdout, VoiceProcessConfig replaces port-based config, workers spawn their own sr-voice child - hardware.rs: remove VoiceInstanceManager (port/process lifecycle), replace with evaluate_scaling() free function + HardwareProbe::voice_config() - sr-voice: add --stdio flag to serve command, new stdio.rs JSONL mode - Remove ureq dependency from server crate (no longer needed) Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
+115
-215
@@ -1,54 +1,40 @@
|
||||
//! Hardware detection + dynamic sr-voice instance management (D-138, Spike 2).
|
||||
//! Hardware detection + dynamic scaling decisions (D-138, Spike 2).
|
||||
//!
|
||||
//! Determines how many parallel LLM workers the system can sustain and manages
|
||||
//! sr-voice process lifecycle. The scaling ceiling is conservative:
|
||||
//! Determines how many parallel LLM workers the system can sustain.
|
||||
//! The scaling ceiling is conservative:
|
||||
//!
|
||||
//! max_new = (free_resource - existing_llm_usage) / 2 / PER_INSTANCE_COST
|
||||
//!
|
||||
//! "free_resource" is VRAM when a GPU is detected (nvidia-smi / AMD sysfs),
|
||||
//! or system RAM otherwise. This ensures the voice pipeline never takes more
|
||||
//! than half the available headroom after accounting for its own instances.
|
||||
//!
|
||||
//! Process lifecycle is handled by `WorkerPool` / `VoicePipe` in `worker.rs`.
|
||||
//! This module only probes hardware and advises on scaling — it does not
|
||||
//! spawn or stop sr-voice processes.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::PathBuf;
|
||||
use std::process::{Child, Command};
|
||||
use std::sync::atomic::AtomicBool;
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
use std::process::Command;
|
||||
use std::time::Duration;
|
||||
|
||||
use sysinfo::System;
|
||||
|
||||
/// RAM budget per sr-voice instance (Gemma 2B Q4_K_M ≈ 1.5 GB resident).
|
||||
use crate::voice::worker::VoiceProcessConfig;
|
||||
|
||||
/// RAM/VRAM budget per sr-voice instance (Gemma 2B Q4_K_M ≈ 1.5 GB resident).
|
||||
const PER_INSTANCE_RAM_MB: u64 = 1536;
|
||||
|
||||
/// Minimum free RAM to allow any voice instance at all.
|
||||
/// Minimum free resource to allow any voice instance at all.
|
||||
const MIN_FREE_RAM_MB: u64 = 1536;
|
||||
|
||||
/// Base port for sr-voice instances. Worker N listens on BASE_PORT + N.
|
||||
const BASE_PORT: u16 = 8321;
|
||||
|
||||
/// Context window size for sr-voice instances.
|
||||
const CTX_SIZE: u32 = 512;
|
||||
|
||||
/// How often the scaler thread checks for scale-up/down opportunities.
|
||||
/// How often the scaler thread checks for scale-up/down opportunities.
|
||||
/// Used by the runtime scaler loop (not yet implemented).
|
||||
const _SCALE_CHECK_INTERVAL: Duration = Duration::from_secs(10);
|
||||
|
||||
/// Queue depth threshold — sustained above this triggers scale-up consideration.
|
||||
const QUEUE_DEPTH_SCALE_UP: usize = 32;
|
||||
|
||||
/// Worker idle duration before scale-down.
|
||||
const IDLE_BEFORE_SCALE_DOWN: Duration = Duration::from_secs(60);
|
||||
|
||||
/// Managed sr-voice process instances.
|
||||
#[derive(Debug)]
|
||||
struct VoiceInstance {
|
||||
process: Child,
|
||||
port: u16,
|
||||
spawned_at: Instant,
|
||||
}
|
||||
|
||||
/// Whether the scaling resource is GPU VRAM or system RAM.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum ResourceMode {
|
||||
@@ -75,18 +61,16 @@ pub struct HardwareProbe {
|
||||
pub threads_per_instance: u32,
|
||||
}
|
||||
|
||||
/// Manages sr-voice process lifecycle and dynamic scaling.
|
||||
pub struct VoiceInstanceManager {
|
||||
/// Path to sr-voice binary.
|
||||
binary_path: PathBuf,
|
||||
/// Path to model file.
|
||||
model_path: PathBuf,
|
||||
/// Running instances keyed by worker ID.
|
||||
instances: HashMap<usize, VoiceInstance>,
|
||||
/// Hardware probe from startup.
|
||||
probe: HardwareProbe,
|
||||
/// Shutdown signal.
|
||||
shutdown: Arc<AtomicBool>,
|
||||
impl HardwareProbe {
|
||||
/// Build a `VoiceProcessConfig` from probe results and user paths.
|
||||
pub fn voice_config(&self, binary_path: String, model_path: String) -> VoiceProcessConfig {
|
||||
VoiceProcessConfig {
|
||||
binary_path,
|
||||
model_path,
|
||||
threads: self.threads_per_instance,
|
||||
ctx_size: CTX_SIZE,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Detect GPU type. Run once at install/first-run, persist to settings.
|
||||
@@ -176,6 +160,70 @@ pub fn probe_hardware(mode: Option<ResourceMode>) -> HardwareProbe {
|
||||
probe
|
||||
}
|
||||
|
||||
/// Evaluate whether the worker pool should scale up or down.
|
||||
///
|
||||
/// Called periodically by the voice pipeline coordinator. Does not take
|
||||
/// action itself — returns a `ScalingDecision` for the caller to act on.
|
||||
pub fn evaluate_scaling(
|
||||
mode: ResourceMode,
|
||||
running_workers: usize,
|
||||
queue_depth: usize,
|
||||
active_workers: usize,
|
||||
max_idle_duration: Option<Duration>,
|
||||
) -> ScalingDecision {
|
||||
// Battery → scale to 1
|
||||
if on_battery() && running_workers > 1 {
|
||||
return ScalingDecision::ScaleDown {
|
||||
reason: "battery power detected".into(),
|
||||
};
|
||||
}
|
||||
|
||||
// Re-probe free resource for current conditions
|
||||
let current_free_mb = probe_current_free(mode);
|
||||
let max_slots = compute_max_slots(current_free_mb, running_workers);
|
||||
|
||||
// Scale up: queue pressure + capacity available
|
||||
if queue_depth >= QUEUE_DEPTH_SCALE_UP
|
||||
&& running_workers < max_slots
|
||||
&& current_free_mb >= MIN_FREE_RAM_MB
|
||||
{
|
||||
return ScalingDecision::ScaleUp {
|
||||
reason: format!(
|
||||
"queue depth {} >= {}, {} slots available",
|
||||
queue_depth, QUEUE_DEPTH_SCALE_UP, max_slots
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
// Scale down: worker idle too long + more than 1 worker
|
||||
if running_workers > 1 {
|
||||
if let Some(idle) = max_idle_duration {
|
||||
if idle >= IDLE_BEFORE_SCALE_DOWN && active_workers < running_workers {
|
||||
return ScalingDecision::ScaleDown {
|
||||
reason: format!("worker idle for {}s", idle.as_secs()),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ScalingDecision::Hold
|
||||
}
|
||||
|
||||
/// Result of a scaling evaluation.
|
||||
#[derive(Debug)]
|
||||
pub enum ScalingDecision {
|
||||
/// No change needed.
|
||||
Hold,
|
||||
/// Spawn an additional worker (caller decides which).
|
||||
ScaleUp { reason: String },
|
||||
/// Stop an idle worker (caller picks the most idle).
|
||||
ScaleDown { reason: String },
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// GPU / resource probing
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Probe NVIDIA GPU VRAM via nvidia-smi.
|
||||
/// Returns (total_mb, free_mb) for the first GPU, or None.
|
||||
fn probe_nvidia_vram() -> Option<(u64, u64)> {
|
||||
@@ -323,180 +371,6 @@ fn on_battery() -> bool {
|
||||
false
|
||||
}
|
||||
|
||||
impl VoiceInstanceManager {
|
||||
/// Create a new manager. Does not spawn any instances yet.
|
||||
pub fn new(
|
||||
binary_path: PathBuf,
|
||||
model_path: PathBuf,
|
||||
probe: HardwareProbe,
|
||||
shutdown: Arc<AtomicBool>,
|
||||
) -> Self {
|
||||
Self {
|
||||
binary_path,
|
||||
model_path,
|
||||
instances: HashMap::new(),
|
||||
probe,
|
||||
shutdown,
|
||||
}
|
||||
}
|
||||
|
||||
/// Number of currently running instances.
|
||||
pub fn instance_count(&self) -> usize {
|
||||
self.instances.len()
|
||||
}
|
||||
|
||||
/// Base port for worker connections.
|
||||
pub fn base_port(&self) -> u16 {
|
||||
BASE_PORT
|
||||
}
|
||||
|
||||
/// The hardware probe from startup.
|
||||
pub fn probe(&self) -> &HardwareProbe {
|
||||
&self.probe
|
||||
}
|
||||
|
||||
/// Spawn a sr-voice instance for the given worker ID.
|
||||
/// Returns the port it's listening on, or an error.
|
||||
pub fn spawn_instance(&mut self, worker_id: usize) -> Result<u16, String> {
|
||||
let port = BASE_PORT + worker_id as u16;
|
||||
|
||||
if self.instances.contains_key(&worker_id) {
|
||||
return Ok(port); // already running
|
||||
}
|
||||
|
||||
let child = Command::new(&self.binary_path)
|
||||
.arg("serve")
|
||||
.arg("--model")
|
||||
.arg(&self.model_path)
|
||||
.arg("--port")
|
||||
.arg(port.to_string())
|
||||
.arg("--threads")
|
||||
.arg(self.probe.threads_per_instance.to_string())
|
||||
.arg("--ctx-size")
|
||||
.arg(CTX_SIZE.to_string())
|
||||
.spawn()
|
||||
.map_err(|e| format!("failed to spawn sr-voice on port {}: {}", port, e))?;
|
||||
|
||||
tracing::info!(worker_id, port, "spawned sr-voice instance");
|
||||
|
||||
self.instances.insert(worker_id, VoiceInstance {
|
||||
process: child,
|
||||
port,
|
||||
spawned_at: Instant::now(),
|
||||
});
|
||||
|
||||
Ok(port)
|
||||
}
|
||||
|
||||
/// Stop a sr-voice instance for the given worker ID.
|
||||
pub fn stop_instance(&mut self, worker_id: usize) {
|
||||
if let Some(mut instance) = self.instances.remove(&worker_id) {
|
||||
let _ = instance.process.kill();
|
||||
let _ = instance.process.wait();
|
||||
tracing::info!(worker_id, port = instance.port, "stopped sr-voice instance");
|
||||
}
|
||||
}
|
||||
|
||||
/// Evaluate whether to scale up or down based on current conditions.
|
||||
///
|
||||
/// Returns (should_scale_up, should_scale_down_worker_id).
|
||||
pub fn evaluate_scaling(
|
||||
&self,
|
||||
queue_depth: usize,
|
||||
active_workers: usize,
|
||||
worker_idle_durations: &HashMap<usize, Duration>,
|
||||
) -> ScalingDecision {
|
||||
// Battery → scale to 1
|
||||
if on_battery() && self.instances.len() > 1 {
|
||||
return ScalingDecision::ScaleDown {
|
||||
reason: "battery power detected".into(),
|
||||
};
|
||||
}
|
||||
|
||||
// Re-probe free resource (VRAM or RAM) for current conditions
|
||||
let current_free_mb = probe_current_free(self.probe.mode);
|
||||
let max_slots = compute_max_slots(current_free_mb, self.instances.len());
|
||||
|
||||
// Scale up: queue pressure + capacity available
|
||||
if queue_depth >= QUEUE_DEPTH_SCALE_UP
|
||||
&& self.instances.len() < max_slots
|
||||
&& current_free_mb >= MIN_FREE_RAM_MB
|
||||
{
|
||||
return ScalingDecision::ScaleUp {
|
||||
reason: format!(
|
||||
"queue depth {} >= {}, {} slots available",
|
||||
queue_depth, QUEUE_DEPTH_SCALE_UP, max_slots
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
// Scale down: worker idle too long + more than 1 instance
|
||||
if self.instances.len() > 1 {
|
||||
for (&worker_id, &idle_time) in worker_idle_durations {
|
||||
if idle_time >= IDLE_BEFORE_SCALE_DOWN && active_workers < self.instances.len() {
|
||||
return ScalingDecision::ScaleDown {
|
||||
reason: format!(
|
||||
"worker {} idle for {}s",
|
||||
worker_id,
|
||||
idle_time.as_secs()
|
||||
),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
ScalingDecision::Hold
|
||||
}
|
||||
|
||||
/// Shut down all sr-voice instances.
|
||||
pub fn shutdown_all(&mut self) {
|
||||
let ids: Vec<usize> = self.instances.keys().copied().collect();
|
||||
for id in ids {
|
||||
self.stop_instance(id);
|
||||
}
|
||||
}
|
||||
|
||||
/// Wait for a sr-voice instance to become healthy (responds to /health).
|
||||
/// Returns true if healthy within timeout, false otherwise.
|
||||
pub fn wait_for_healthy(&self, port: u16, timeout: Duration) -> bool {
|
||||
let url = format!("http://127.0.0.1:{}/health", port);
|
||||
let deadline = Instant::now() + timeout;
|
||||
|
||||
while Instant::now() < deadline {
|
||||
let result = ureq::AgentBuilder::new()
|
||||
.timeout(Duration::from_secs(1))
|
||||
.build()
|
||||
.get(&url)
|
||||
.call();
|
||||
|
||||
if result.is_ok() {
|
||||
return true;
|
||||
}
|
||||
|
||||
std::thread::sleep(Duration::from_millis(500));
|
||||
}
|
||||
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for VoiceInstanceManager {
|
||||
fn drop(&mut self) {
|
||||
self.shutdown_all();
|
||||
}
|
||||
}
|
||||
|
||||
/// Result of a scaling evaluation.
|
||||
#[derive(Debug)]
|
||||
pub enum ScalingDecision {
|
||||
/// No change needed.
|
||||
Hold,
|
||||
/// Spawn an additional instance.
|
||||
ScaleUp { reason: String },
|
||||
/// Stop an instance (pick the most idle worker).
|
||||
ScaleDown { reason: String },
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Tests
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -544,8 +418,7 @@ mod tests {
|
||||
let slots_0 = compute_max_slots(8192, 0); // 8GB free, 0 running
|
||||
let slots_2 = compute_max_slots(8192, 2); // 8GB free, 2 running
|
||||
|
||||
// With 2 running, effective headroom is smaller so ceiling is lower per-new-instance
|
||||
// but total (running + new) can still be higher
|
||||
// Both should be non-zero with 8GB
|
||||
assert!(slots_0 > 0);
|
||||
assert!(slots_2 > 0);
|
||||
// The key property: free RAM measured at probe time is the same,
|
||||
@@ -561,9 +434,36 @@ mod tests {
|
||||
assert!(probe.threads_per_instance >= 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn voice_config_from_probe() {
|
||||
let probe = HardwareProbe {
|
||||
mode: ResourceMode::Cpu,
|
||||
total_mb: 16384,
|
||||
free_mb: 8192,
|
||||
max_slots: 2,
|
||||
cpu_threads: 8,
|
||||
threads_per_instance: 2,
|
||||
};
|
||||
let config = probe.voice_config("/usr/bin/sr-voice".into(), "/models/gemma.gguf".into());
|
||||
assert_eq!(config.threads, 2);
|
||||
assert_eq!(config.ctx_size, CTX_SIZE);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn battery_detection_does_not_crash() {
|
||||
// Just verify it doesn't panic — result depends on hardware
|
||||
let _ = on_battery();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn evaluate_scaling_hold_when_calm() {
|
||||
let decision = evaluate_scaling(
|
||||
ResourceMode::Cpu,
|
||||
2, // running
|
||||
5, // queue depth (low)
|
||||
1, // active
|
||||
None,
|
||||
);
|
||||
assert!(matches!(decision, ScalingDecision::Hold));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,8 +9,8 @@
|
||||
//! - `prompt_builder` — composition engine: NPC data + culture + tell state → prompt string
|
||||
//! - `cache` — MessagePack voice cache (store/retrieve, length-gated variants)
|
||||
//! - `queue` — crossbeam work queue with priority + backpressure
|
||||
//! - `worker` — inference worker pool (dynamic scaling, owns sr-voice HTTP clients)
|
||||
//! - `hardware` — hardware detection + dynamic sr-voice instance management
|
||||
//! - `worker` — inference worker pool (each worker owns a piped sr-voice child process)
|
||||
//! - `hardware` — hardware detection + dynamic scaling decisions
|
||||
|
||||
pub mod cache;
|
||||
pub mod hardware;
|
||||
|
||||
+123
-57
@@ -1,8 +1,12 @@
|
||||
//! Inference worker pool (D-138, Spike 2).
|
||||
//!
|
||||
//! Dynamic pool of worker threads, each owning an HTTP client to its own
|
||||
//! sr-voice instance. Pool size controlled by hardware detection.
|
||||
//! Dynamic pool of worker threads, each owning a piped stdin/stdout connection
|
||||
//! to its own sr-voice child process. No network ports — the model is only
|
||||
//! reachable through the game server's queue (Gemma 2 T&C compliance).
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::io::{BufRead, BufReader, Write as IoWrite};
|
||||
use std::process::{Child, ChildStdin, ChildStdout, Command, Stdio};
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::thread::{self, JoinHandle};
|
||||
@@ -18,12 +22,6 @@ use crate::voice::queue::VoiceRequest;
|
||||
/// Minimum token count for a valid response. Below this, retry once.
|
||||
const MIN_TOKENS: usize = 4;
|
||||
|
||||
/// How long to wait before retrying connection to sr-voice.
|
||||
const _RECONNECT_INTERVAL: Duration = Duration::from_secs(30);
|
||||
|
||||
/// HTTP request timeout for inference calls.
|
||||
const INFERENCE_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
|
||||
/// Worker pool manages inference worker threads.
|
||||
pub struct WorkerPool {
|
||||
workers: Vec<WorkerHandle>,
|
||||
@@ -46,14 +44,20 @@ pub struct WorkerContext {
|
||||
pub paused: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
use std::collections::HashMap;
|
||||
/// Configuration for spawning sr-voice child processes.
|
||||
#[derive(Clone)]
|
||||
pub struct VoiceProcessConfig {
|
||||
pub binary_path: String,
|
||||
pub model_path: String,
|
||||
pub threads: u32,
|
||||
pub ctx_size: u32,
|
||||
}
|
||||
|
||||
impl WorkerPool {
|
||||
/// Spawn `count` worker threads, each connecting to sr-voice on
|
||||
/// `base_port + worker_id`.
|
||||
/// Spawn `count` worker threads, each owning a piped sr-voice child process.
|
||||
pub fn spawn(
|
||||
count: usize,
|
||||
base_port: u16,
|
||||
config: &VoiceProcessConfig,
|
||||
receiver: Receiver<VoiceRequest>,
|
||||
cache: Arc<Mutex<VoiceCacheStore>>,
|
||||
cultures: Arc<HashMap<String, CultureProfile>>,
|
||||
@@ -73,11 +77,11 @@ impl WorkerPool {
|
||||
active_count: Arc::clone(&active_count),
|
||||
paused: Arc::clone(&paused),
|
||||
};
|
||||
let port = base_port + id as u16;
|
||||
let cfg = config.clone();
|
||||
|
||||
let thread = thread::Builder::new()
|
||||
.name(format!("voice-worker-{}", id))
|
||||
.spawn(move || worker_loop(id, port, ctx))
|
||||
.spawn(move || worker_loop(id, cfg, ctx))
|
||||
.expect("failed to spawn voice worker thread");
|
||||
|
||||
workers.push(WorkerHandle {
|
||||
@@ -86,7 +90,7 @@ impl WorkerPool {
|
||||
});
|
||||
}
|
||||
|
||||
tracing::info!(count, base_port, "voice worker pool started");
|
||||
tracing::info!(count, "voice worker pool started");
|
||||
|
||||
Self {
|
||||
workers,
|
||||
@@ -123,13 +127,105 @@ impl Drop for WorkerPool {
|
||||
}
|
||||
}
|
||||
|
||||
/// Main worker loop: receive requests, build prompts, call sr-voice, cache results.
|
||||
fn worker_loop(id: usize, port: u16, ctx: WorkerContext) {
|
||||
let base_url = format!("http://127.0.0.1:{}", port);
|
||||
tracing::debug!(id, port, "voice worker started");
|
||||
/// A piped connection to a sr-voice child process.
|
||||
struct VoicePipe {
|
||||
child: Child,
|
||||
stdin: ChildStdin,
|
||||
reader: BufReader<ChildStdout>,
|
||||
}
|
||||
|
||||
// Below-normal thread priority is handled at the OS level by the
|
||||
// sr-voice process itself (nice value). Worker threads inherit it.
|
||||
impl VoicePipe {
|
||||
/// Spawn a sr-voice child process with piped stdin/stdout.
|
||||
fn spawn(config: &VoiceProcessConfig) -> Result<Self, String> {
|
||||
let mut child = Command::new(&config.binary_path)
|
||||
.arg("serve")
|
||||
.arg("--model")
|
||||
.arg(&config.model_path)
|
||||
.arg("--threads")
|
||||
.arg(config.threads.to_string())
|
||||
.arg("--ctx-size")
|
||||
.arg(config.ctx_size.to_string())
|
||||
.arg("--stdio")
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::inherit())
|
||||
.spawn()
|
||||
.map_err(|e| format!("failed to spawn sr-voice: {}", e))?;
|
||||
|
||||
let stdin = child.stdin.take()
|
||||
.ok_or_else(|| "failed to capture sr-voice stdin".to_string())?;
|
||||
let stdout = child.stdout.take()
|
||||
.ok_or_else(|| "failed to capture sr-voice stdout".to_string())?;
|
||||
|
||||
Ok(Self {
|
||||
child,
|
||||
stdin,
|
||||
reader: BufReader::new(stdout),
|
||||
})
|
||||
}
|
||||
|
||||
/// Send a prompt and read the response (JSONL: one JSON object per line).
|
||||
fn generate(&mut self, prompt: &str, seed: Option<u64>) -> Result<String, String> {
|
||||
let request = serde_json::json!({
|
||||
"prompt": prompt,
|
||||
"seed": seed,
|
||||
});
|
||||
|
||||
let mut line = serde_json::to_string(&request)
|
||||
.map_err(|e| format!("failed to serialize request: {}", e))?;
|
||||
line.push('\n');
|
||||
|
||||
self.stdin
|
||||
.write_all(line.as_bytes())
|
||||
.map_err(|e| format!("failed to write to sr-voice stdin: {}", e))?;
|
||||
self.stdin
|
||||
.flush()
|
||||
.map_err(|e| format!("failed to flush sr-voice stdin: {}", e))?;
|
||||
|
||||
let mut response_line = String::new();
|
||||
self.reader
|
||||
.read_line(&mut response_line)
|
||||
.map_err(|e| format!("failed to read from sr-voice stdout: {}", e))?;
|
||||
|
||||
if response_line.is_empty() {
|
||||
return Err("sr-voice process closed stdout".to_string());
|
||||
}
|
||||
|
||||
let body: serde_json::Value = serde_json::from_str(&response_line)
|
||||
.map_err(|e| format!("failed to parse response JSON: {}", e))?;
|
||||
|
||||
if let Some(err) = body.get("error") {
|
||||
return Err(format!("sr-voice error: {}", err));
|
||||
}
|
||||
|
||||
body["text"]
|
||||
.as_str()
|
||||
.map(|s| s.trim().to_string())
|
||||
.ok_or_else(|| "response missing 'text' field".to_string())
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for VoicePipe {
|
||||
fn drop(&mut self) {
|
||||
let _ = self.child.kill();
|
||||
let _ = self.child.wait();
|
||||
}
|
||||
}
|
||||
|
||||
/// Main worker loop: spawn sr-voice child, receive requests, process them.
|
||||
fn worker_loop(id: usize, config: VoiceProcessConfig, ctx: WorkerContext) {
|
||||
tracing::debug!(id, "voice worker starting sr-voice child process");
|
||||
|
||||
let mut pipe = match VoicePipe::spawn(&config) {
|
||||
Ok(p) => {
|
||||
tracing::info!(id, "voice worker connected to sr-voice via stdio");
|
||||
p
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!(id, error = %e, "voice worker failed to spawn sr-voice — exiting");
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
loop {
|
||||
if ctx.shutdown.load(Ordering::SeqCst) {
|
||||
@@ -145,23 +241,22 @@ fn worker_loop(id: usize, port: u16, ctx: WorkerContext) {
|
||||
|
||||
// Skip while paused (zone transition)
|
||||
if ctx.paused.load(Ordering::Relaxed) {
|
||||
// Re-queue the request — it wasn't consumed
|
||||
let _ = ctx.receiver.clone(); // can't re-send, just drop during pause
|
||||
continue;
|
||||
}
|
||||
|
||||
ctx.active_count.fetch_add(1, Ordering::Relaxed);
|
||||
process_request(id, &base_url, &request, &ctx);
|
||||
process_request(id, &mut pipe, &request, &ctx);
|
||||
ctx.active_count.fetch_sub(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
tracing::debug!(id, "voice worker stopped");
|
||||
// VoicePipe::drop kills the child process
|
||||
}
|
||||
|
||||
/// Process a single voice request: build prompt → infer → validate → cache.
|
||||
fn process_request(
|
||||
worker_id: usize,
|
||||
base_url: &str,
|
||||
pipe: &mut VoicePipe,
|
||||
request: &VoiceRequest,
|
||||
ctx: &WorkerContext,
|
||||
) {
|
||||
@@ -186,8 +281,8 @@ fn process_request(
|
||||
request.seed,
|
||||
);
|
||||
|
||||
// Call sr-voice
|
||||
let result = call_sr_voice(base_url, &built.prompt);
|
||||
// Call sr-voice via stdio pipe
|
||||
let result = pipe.generate(&built.prompt, Some(request.seed));
|
||||
|
||||
match result {
|
||||
Ok(text) if text.split_whitespace().count() >= MIN_TOKENS => {
|
||||
@@ -207,12 +302,11 @@ fn process_request(
|
||||
request.tell_state,
|
||||
request.seed.wrapping_add(1),
|
||||
);
|
||||
match call_sr_voice(base_url, &retry_built.prompt) {
|
||||
match pipe.generate(&retry_built.prompt, Some(request.seed.wrapping_add(1))) {
|
||||
Ok(text) if text.split_whitespace().count() >= MIN_TOKENS => {
|
||||
cache_result(request, &text, ctx);
|
||||
}
|
||||
_ => {
|
||||
// Graceful degradation: cache base text
|
||||
tracing::debug!(
|
||||
worker_id,
|
||||
npc = request.npc_stable_id,
|
||||
@@ -233,34 +327,6 @@ fn process_request(
|
||||
}
|
||||
}
|
||||
|
||||
/// POST to sr-voice /generate endpoint and return the generated text.
|
||||
fn call_sr_voice(base_url: &str, prompt: &str) -> Result<String, String> {
|
||||
let url = format!("{}/generate", base_url);
|
||||
|
||||
let payload = serde_json::json!({ "prompt": prompt });
|
||||
|
||||
let response = ureq::AgentBuilder::new()
|
||||
.timeout(INFERENCE_TIMEOUT)
|
||||
.build()
|
||||
.post(&url)
|
||||
.send_json(payload);
|
||||
|
||||
match response {
|
||||
Ok(resp) => {
|
||||
let body_str = resp
|
||||
.into_string()
|
||||
.map_err(|e| format!("failed to read response: {}", e))?;
|
||||
let body: serde_json::Value = serde_json::from_str(&body_str)
|
||||
.map_err(|e| format!("failed to parse JSON: {}", e))?;
|
||||
body["text"]
|
||||
.as_str()
|
||||
.map(|s| s.trim().to_string())
|
||||
.ok_or_else(|| "response missing 'text' field".to_string())
|
||||
}
|
||||
Err(e) => Err(format!("HTTP error: {}", e)),
|
||||
}
|
||||
}
|
||||
|
||||
/// Cache the inference result.
|
||||
fn cache_result(request: &VoiceRequest, text: &str, ctx: &WorkerContext) {
|
||||
let key = cache_key_from_request(request);
|
||||
|
||||
Reference in New Issue
Block a user