fix(voice): address PR review findings — 3 critical, 5 warning, 4 suggestion
Critical fixes: - Pause mechanism: workers now hold requests during pause instead of dropping them. Queue and worker pool share the same AtomicBool flag via VoiceQueue::paused_flag(). Submit() rejects while paused. - Seed type: sr-voice accepts u64 seeds over IPC (explicit u32 truncation for llama.cpp sampler, documented). Warning fixes: - HashMap → BTreeMap in cache.rs and worker.rs (D-010 determinism mandate). Added Ord derives to CacheKey, ContentType, TellCategory. - VoicePipe::generate() watchdog kills child after 120s timeout to prevent indefinite blocking on read_line. - VoiceCacheStore Drop impl calls save_all() on shutdown. - trait-modifiers.ron: fixed 3 wrong trait names (Impulsive→Compassionate, Methodical→Incurious, Stubborn→Ruthless) to match PersonalityTrait enum. Suggestion fixes: - Worker spawn: log error + reduce pool instead of panic on thread failure. - on_battery(): added macOS detection via pmset. - Epistemic markers: lowercased constants, removed redundant to_lowercase(). - cache.rs: documented non-atomic write tradeoff. - queue.rs: reprioritize() bypasses pause check (it runs during pause). Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -11,7 +11,8 @@
|
||||
// - Behavioral, not emotional. Never label a feeling.
|
||||
// - Each clause targets a distinct speech dimension so traits stack cleanly.
|
||||
// Bold = delivery force. Cautious = word selection. Curious = sentence shape.
|
||||
// etc. A Bold+Guarded NPC speaks with force but picks words carefully — no conflict.
|
||||
// Compassionate = cadence (engaged). Incurious = cadence (flat). Ruthless = position.
|
||||
// A Bold+Compassionate NPC speaks with force but gives the answer room to land — no conflict.
|
||||
// - Kept to 1-2 sentences. LLM context is limited and these share space with
|
||||
// persona, tell-state, and epistemic marker instructions.
|
||||
|
||||
@@ -31,11 +32,11 @@
|
||||
// Framing: no softening, no omission. The whole thing, bluntly.
|
||||
"Honest": "This character gives the whole answer, including the part that doesn't reflect well. Nothing is softened to spare anyone.",
|
||||
|
||||
// Cadence: fast entry, no ramp-up. The response is already in motion.
|
||||
"Impulsive": "This character starts talking before the thought is finished. The sentence catches up to itself mid-way.",
|
||||
// Cadence: open, engaged. Responses arrive with momentum and care.
|
||||
"Compassionate": "This character gives the answer room to land. There's no rush past the difficult part — it gets said, plainly, without looking away.",
|
||||
|
||||
// Cadence: deliberate pacing. Sequence matters. One thing before the next.
|
||||
"Methodical": "This character takes things in order. Sentences build, one piece at a time, and don't arrive at the point before the steps do.",
|
||||
// Cadence: flat, unengaged. The response does its job and stops.
|
||||
"Incurious": "This character answers what was asked. Nothing extra surfaces — no follow-up, no interest, no second look at what was just said.",
|
||||
|
||||
// Volume and reach: minimal, inward-facing. Not interested in being heard widely.
|
||||
"Reclusive": "This character answers what was asked and closes the door. There is no invitation for follow-up.",
|
||||
@@ -43,6 +44,6 @@
|
||||
// Texture: open, inclusive. Others are assumed to be present and welcome.
|
||||
"Social": "This character addresses the conversation, not just the question. A word or two lands that wasn't strictly necessary — the kind that keeps things warm.",
|
||||
|
||||
// Position: established, not revisable. Sentences don't open up once they're out.
|
||||
"Stubborn": "This character doesn't revise mid-sentence. What comes out is what they meant, and nothing in the delivery invites renegotiation.",
|
||||
// Position: sharp, economical. Nothing is offered that doesn't serve the speaker.
|
||||
"Ruthless": "This character cuts to the useful part. Courtesy is absent, not hostile — just unnecessary. The sentence ends when the point is made.",
|
||||
}
|
||||
|
||||
@@ -12,7 +12,7 @@ const TOP_P: f32 = 0.9;
|
||||
#[derive(serde::Deserialize)]
|
||||
struct GenerateRequest {
|
||||
prompt: String,
|
||||
seed: Option<u32>,
|
||||
seed: Option<u64>,
|
||||
}
|
||||
|
||||
pub fn run_server(
|
||||
@@ -70,7 +70,7 @@ fn handle_generate(engine: &InferenceEngine, mut request: tiny_http::Request) {
|
||||
};
|
||||
|
||||
eprintln!(" generate: {} chars", req.prompt.len());
|
||||
match engine.generate(&req.prompt, MAX_TOKENS, TEMPERATURE, TOP_P, req.seed) {
|
||||
match engine.generate(&req.prompt, MAX_TOKENS, TEMPERATURE, TOP_P, req.seed.map(|s| s as u32)) {
|
||||
Ok(result) => {
|
||||
eprintln!(" -> {} tokens, {:.1} t/s", result.tokens_generated, result.tokens_per_sec);
|
||||
respond(request, 200, &serde_json::to_string(&result).unwrap());
|
||||
|
||||
@@ -19,7 +19,7 @@ const TOP_P: f32 = 0.9;
|
||||
#[derive(serde::Deserialize)]
|
||||
struct StdioRequest {
|
||||
prompt: String,
|
||||
seed: Option<u32>,
|
||||
seed: Option<u64>,
|
||||
}
|
||||
|
||||
pub fn run_stdio(engine: InferenceEngine) -> Result<(), Box<dyn std::error::Error>> {
|
||||
@@ -42,7 +42,9 @@ pub fn run_stdio(engine: InferenceEngine) -> Result<(), Box<dyn std::error::Erro
|
||||
let response = match serde_json::from_str::<StdioRequest>(&line) {
|
||||
Ok(req) => {
|
||||
eprintln!(" stdio: {} chars", req.prompt.len());
|
||||
match engine.generate(&req.prompt, MAX_TOKENS, TEMPERATURE, TOP_P, req.seed) {
|
||||
// Truncate u64 seed to u32 for llama.cpp sampler (deterministic
|
||||
// within same process, seed space reduction is acceptable).
|
||||
match engine.generate(&req.prompt, MAX_TOKENS, TEMPERATURE, TOP_P, req.seed.map(|s| s as u32)) {
|
||||
Ok(result) => {
|
||||
eprintln!(
|
||||
" -> {} tokens, {:.1} t/s",
|
||||
|
||||
@@ -37,7 +37,7 @@ use crate::simulation::tier::ActiveSim;
|
||||
///
|
||||
/// Derived each tick from NPC simulation state — not authored per NPC.
|
||||
/// Five categories correspond to the D-024 tell taxonomy.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
|
||||
pub enum TellCategory {
|
||||
/// NPC exhibits nervous behaviour: Major secret + stress exceeds half of threshold.
|
||||
Nervous,
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
//! the same directory. No separate baked path.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
use std::collections::BTreeMap;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::PathBuf;
|
||||
@@ -17,7 +17,7 @@ use crate::npc::tell_state::TellCategory;
|
||||
use crate::voice::prompt_builder::ContentType;
|
||||
|
||||
/// Cache key for a single voiced line.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
|
||||
pub struct CacheKey {
|
||||
pub culture_id: String,
|
||||
pub npc_stable_id: u64,
|
||||
@@ -35,7 +35,7 @@ pub struct ZoneVoiceCache {
|
||||
/// Injector version hash — cache miss if this doesn't match.
|
||||
pub injector_version: String,
|
||||
/// Cached voiced lines.
|
||||
pub entries: HashMap<CacheKey, String>,
|
||||
pub entries: BTreeMap<CacheKey, String>,
|
||||
}
|
||||
|
||||
impl ZoneVoiceCache {
|
||||
@@ -43,7 +43,7 @@ impl ZoneVoiceCache {
|
||||
Self {
|
||||
model_version,
|
||||
injector_version,
|
||||
entries: HashMap::new(),
|
||||
entries: BTreeMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -79,7 +79,7 @@ pub struct VoiceCacheStore {
|
||||
/// Current injector version hash.
|
||||
injector_version: String,
|
||||
/// Loaded zone caches.
|
||||
zones: HashMap<u32, ZoneVoiceCache>,
|
||||
zones: BTreeMap<u32, ZoneVoiceCache>,
|
||||
}
|
||||
|
||||
impl VoiceCacheStore {
|
||||
@@ -95,7 +95,7 @@ impl VoiceCacheStore {
|
||||
world_seed,
|
||||
model_version,
|
||||
injector_version,
|
||||
zones: HashMap::new(),
|
||||
zones: BTreeMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -126,6 +126,10 @@ impl VoiceCacheStore {
|
||||
}
|
||||
|
||||
/// Persist a zone's cache to disk as MessagePack.
|
||||
///
|
||||
/// Not atomic (no rename-into-place) — a crash mid-write can produce a
|
||||
/// partial file. This is acceptable: worst case is a cache miss on next
|
||||
/// load, triggering re-inference or base text fallback.
|
||||
pub fn save_zone(&self, zone_id: u32) -> io::Result<()> {
|
||||
let Some(cache) = self.zones.get(&zone_id) else {
|
||||
return Ok(());
|
||||
@@ -175,6 +179,14 @@ impl VoiceCacheStore {
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for VoiceCacheStore {
|
||||
fn drop(&mut self) {
|
||||
if let Err(e) = self.save_all() {
|
||||
tracing::warn!(error = %e, "failed to save voice cache on shutdown");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Determine which tell states should be cached for a given base text.
|
||||
///
|
||||
/// Length-gated variant count (Spike 1 finding):
|
||||
|
||||
@@ -343,32 +343,53 @@ fn compute_max_slots(free_mb: u64, running_instances: usize) -> usize {
|
||||
}
|
||||
}
|
||||
|
||||
/// Check if the system is on battery power (Linux).
|
||||
/// Check if the system is on battery power.
|
||||
///
|
||||
/// Linux: reads `/sys/class/power_supply/` sysfs entries.
|
||||
/// macOS: runs `pmset -g batt` and checks for "Battery Power".
|
||||
/// Other platforms: returns false (assumes AC power).
|
||||
fn on_battery() -> bool {
|
||||
let Ok(entries) = std::fs::read_dir("/sys/class/power_supply/") else {
|
||||
return false;
|
||||
};
|
||||
|
||||
for entry in entries.flatten() {
|
||||
let type_path = entry.path().join("type");
|
||||
let status_path = entry.path().join("status");
|
||||
|
||||
let Ok(supply_type) = std::fs::read_to_string(&type_path) else {
|
||||
continue;
|
||||
#[cfg(target_os = "linux")]
|
||||
{
|
||||
let Ok(entries) = std::fs::read_dir("/sys/class/power_supply/") else {
|
||||
return false;
|
||||
};
|
||||
if supply_type.trim() != "Battery" {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Ok(status) = std::fs::read_to_string(&status_path) {
|
||||
let status = status.trim();
|
||||
if status == "Discharging" {
|
||||
return true;
|
||||
for entry in entries.flatten() {
|
||||
let type_path = entry.path().join("type");
|
||||
let status_path = entry.path().join("status");
|
||||
|
||||
let Ok(supply_type) = std::fs::read_to_string(&type_path) else {
|
||||
continue;
|
||||
};
|
||||
if supply_type.trim() != "Battery" {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Ok(status) = std::fs::read_to_string(&status_path) {
|
||||
if status.trim() == "Discharging" {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
false
|
||||
}
|
||||
|
||||
false
|
||||
#[cfg(target_os = "macos")]
|
||||
{
|
||||
Command::new("pmset")
|
||||
.args(["-g", "batt"])
|
||||
.output()
|
||||
.ok()
|
||||
.map(|o| String::from_utf8_lossy(&o.stdout).contains("Battery Power"))
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
|
||||
{
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -16,7 +16,7 @@ use crate::npc::blueprint::CultureProfile;
|
||||
use crate::npc::tell_state::TellCategory;
|
||||
|
||||
/// Content type determines the task verb in the prompt.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
|
||||
pub enum ContentType {
|
||||
/// Spoken dialogue — re-voiced with "Re-voice".
|
||||
Dialogue,
|
||||
@@ -84,11 +84,13 @@ fn tell_injector(category: TellCategory) -> &'static str {
|
||||
/// converting "I heard the night crew stopped the line" to
|
||||
/// "Line tripped twice. What's the plan?" — losing the epistemic framing
|
||||
/// that is semantically load-bearing for the perception system.
|
||||
/// Note: all entries must be lowercase — matched against `base_text.to_lowercase()`.
|
||||
/// The original-case version is reconstructed from the base text for the PRESERVE clause.
|
||||
const EPISTEMIC_MARKERS: &[&str] = &[
|
||||
"I heard",
|
||||
"I think",
|
||||
"I saw",
|
||||
"I noticed",
|
||||
"i heard",
|
||||
"i think",
|
||||
"i saw",
|
||||
"i noticed",
|
||||
"someone told me",
|
||||
"they say",
|
||||
"apparently",
|
||||
@@ -118,7 +120,7 @@ fn extract_epistemic_markers(base_text: &str) -> Vec<&'static str> {
|
||||
let lower = base_text.to_lowercase();
|
||||
EPISTEMIC_MARKERS
|
||||
.iter()
|
||||
.filter(|marker| lower.contains(&marker.to_lowercase()))
|
||||
.filter(|marker| lower.contains(*marker))
|
||||
.copied()
|
||||
.collect()
|
||||
}
|
||||
@@ -461,7 +463,7 @@ mod tests {
|
||||
42,
|
||||
);
|
||||
assert!(result.prompt.contains("PRESERVE:"));
|
||||
assert!(result.prompt.contains("I heard"));
|
||||
assert!(result.prompt.contains("i heard"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
+40
-14
@@ -5,7 +5,8 @@
|
||||
|
||||
use std::cmp::Ordering;
|
||||
use std::collections::BinaryHeap;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
|
||||
use std::sync::Arc;
|
||||
|
||||
use crossbeam_channel::{Receiver, Sender, TrySendError};
|
||||
|
||||
@@ -83,7 +84,7 @@ impl Ord for VoiceRequest {
|
||||
pub struct VoiceQueue {
|
||||
sender: Sender<VoiceRequest>,
|
||||
receiver: Receiver<VoiceRequest>,
|
||||
paused: Arc<Mutex<bool>>,
|
||||
paused: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
impl VoiceQueue {
|
||||
@@ -92,12 +93,16 @@ impl VoiceQueue {
|
||||
Self {
|
||||
sender,
|
||||
receiver,
|
||||
paused: Arc::new(Mutex::new(false)),
|
||||
paused: Arc::new(AtomicBool::new(false)),
|
||||
}
|
||||
}
|
||||
|
||||
/// Submit a voice request. Returns `false` if the queue is full (backpressure).
|
||||
/// Submit a voice request. Returns `false` if the queue is full or paused.
|
||||
pub fn submit(&self, request: VoiceRequest) -> bool {
|
||||
if self.is_paused() {
|
||||
tracing::trace!("voice queue paused — dropping submission");
|
||||
return false;
|
||||
}
|
||||
match self.sender.try_send(request) {
|
||||
Ok(()) => true,
|
||||
Err(TrySendError::Full(_)) => {
|
||||
@@ -116,6 +121,15 @@ impl VoiceQueue {
|
||||
self.receiver.clone()
|
||||
}
|
||||
|
||||
/// Get the shared pause flag for worker threads.
|
||||
///
|
||||
/// Workers check this flag to avoid processing requests during zone
|
||||
/// transitions. The same `Arc<AtomicBool>` is shared between the queue
|
||||
/// and the worker pool.
|
||||
pub fn paused_flag(&self) -> Arc<AtomicBool> {
|
||||
Arc::clone(&self.paused)
|
||||
}
|
||||
|
||||
/// Current number of pending requests in the channel.
|
||||
pub fn pending_count(&self) -> usize {
|
||||
self.sender.len()
|
||||
@@ -123,21 +137,17 @@ impl VoiceQueue {
|
||||
|
||||
/// Pause the queue (zone transition start).
|
||||
pub fn pause(&self) {
|
||||
if let Ok(mut p) = self.paused.lock() {
|
||||
*p = true;
|
||||
}
|
||||
self.paused.store(true, AtomicOrdering::SeqCst);
|
||||
}
|
||||
|
||||
/// Resume the queue (zone transition complete).
|
||||
pub fn resume(&self) {
|
||||
if let Ok(mut p) = self.paused.lock() {
|
||||
*p = false;
|
||||
}
|
||||
self.paused.store(false, AtomicOrdering::SeqCst);
|
||||
}
|
||||
|
||||
/// Check if the queue is paused.
|
||||
pub fn is_paused(&self) -> bool {
|
||||
self.paused.lock().map(|p| *p).unwrap_or(false)
|
||||
self.paused.load(AtomicOrdering::SeqCst)
|
||||
}
|
||||
|
||||
/// Reprioritize all pending requests after a zone change.
|
||||
@@ -162,11 +172,16 @@ impl VoiceQueue {
|
||||
let count = pending.len();
|
||||
let mut resubmitted = 0;
|
||||
|
||||
// Re-tag and re-submit
|
||||
// Re-tag and re-submit (bypass pause check — reprioritize is called
|
||||
// while paused and needs to refill the channel).
|
||||
for mut req in pending {
|
||||
req.priority = classify(&req);
|
||||
if self.submit(req) {
|
||||
resubmitted += 1;
|
||||
match self.sender.try_send(req) {
|
||||
Ok(()) => resubmitted += 1,
|
||||
Err(TrySendError::Full(_)) => {
|
||||
tracing::trace!("voice queue full during reprioritize — dropping request");
|
||||
}
|
||||
Err(TrySendError::Disconnected(_)) => break,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -287,6 +302,17 @@ mod tests {
|
||||
assert!(!queue.is_paused());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn submit_rejected_when_paused() {
|
||||
let queue = VoiceQueue::new();
|
||||
queue.pause();
|
||||
assert!(!queue.submit(make_request(Priority::High, "should be rejected")));
|
||||
assert_eq!(queue.pending_count(), 0);
|
||||
queue.resume();
|
||||
assert!(queue.submit(make_request(Priority::High, "should be accepted")));
|
||||
assert_eq!(queue.pending_count(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reprioritize_reshuffles_on_zone_change() {
|
||||
let queue = VoiceQueue::new();
|
||||
|
||||
+70
-18
@@ -4,7 +4,7 @@
|
||||
//! 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::collections::BTreeMap;
|
||||
use std::io::{BufRead, BufReader, Write as IoWrite};
|
||||
use std::process::{Child, ChildStdin, ChildStdout, Command, Stdio};
|
||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||
@@ -22,6 +22,11 @@ use crate::voice::queue::VoiceRequest;
|
||||
/// Minimum token count for a valid response. Below this, retry once.
|
||||
const MIN_TOKENS: usize = 4;
|
||||
|
||||
/// Maximum time to wait for a single inference response before killing
|
||||
/// the sr-voice child process. Gemma 2B at ~16 t/s should complete 64
|
||||
/// tokens in ~4s; 120s covers extreme slow hardware with margin.
|
||||
const INFERENCE_TIMEOUT: Duration = Duration::from_secs(120);
|
||||
|
||||
/// Worker pool manages inference worker threads.
|
||||
pub struct WorkerPool {
|
||||
workers: Vec<WorkerHandle>,
|
||||
@@ -38,7 +43,7 @@ struct WorkerHandle {
|
||||
pub struct WorkerContext {
|
||||
pub receiver: Receiver<VoiceRequest>,
|
||||
pub cache: Arc<Mutex<VoiceCacheStore>>,
|
||||
pub cultures: Arc<HashMap<String, CultureProfile>>,
|
||||
pub cultures: Arc<BTreeMap<String, CultureProfile>>,
|
||||
pub shutdown: Arc<AtomicBool>,
|
||||
pub active_count: Arc<AtomicUsize>,
|
||||
pub paused: Arc<AtomicBool>,
|
||||
@@ -55,16 +60,19 @@ pub struct VoiceProcessConfig {
|
||||
|
||||
impl WorkerPool {
|
||||
/// Spawn `count` worker threads, each owning a piped sr-voice child process.
|
||||
///
|
||||
/// `paused` should come from `VoiceQueue::paused_flag()` so workers share
|
||||
/// the same pause signal as the queue.
|
||||
pub fn spawn(
|
||||
count: usize,
|
||||
config: &VoiceProcessConfig,
|
||||
receiver: Receiver<VoiceRequest>,
|
||||
cache: Arc<Mutex<VoiceCacheStore>>,
|
||||
cultures: Arc<HashMap<String, CultureProfile>>,
|
||||
cultures: Arc<BTreeMap<String, CultureProfile>>,
|
||||
paused: Arc<AtomicBool>,
|
||||
) -> Self {
|
||||
let shutdown = Arc::new(AtomicBool::new(false));
|
||||
let active_count = Arc::new(AtomicUsize::new(0));
|
||||
let paused = Arc::new(AtomicBool::new(false));
|
||||
|
||||
let mut workers = Vec::with_capacity(count);
|
||||
|
||||
@@ -79,15 +87,20 @@ impl WorkerPool {
|
||||
};
|
||||
let cfg = config.clone();
|
||||
|
||||
let thread = thread::Builder::new()
|
||||
match thread::Builder::new()
|
||||
.name(format!("voice-worker-{}", id))
|
||||
.spawn(move || worker_loop(id, cfg, ctx))
|
||||
.expect("failed to spawn voice worker thread");
|
||||
|
||||
workers.push(WorkerHandle {
|
||||
thread: Some(thread),
|
||||
id,
|
||||
});
|
||||
{
|
||||
Ok(thread) => {
|
||||
workers.push(WorkerHandle {
|
||||
thread: Some(thread),
|
||||
id,
|
||||
});
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!(id, error = %e, "failed to spawn voice worker — reducing pool");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::info!(count, "voice worker pool started");
|
||||
@@ -165,6 +178,9 @@ impl VoicePipe {
|
||||
}
|
||||
|
||||
/// Send a prompt and read the response (JSONL: one JSON object per line).
|
||||
///
|
||||
/// Spawns a watchdog thread that kills the child process after
|
||||
/// `INFERENCE_TIMEOUT` to prevent indefinite blocking on `read_line`.
|
||||
fn generate(&mut self, prompt: &str, seed: Option<u64>) -> Result<String, String> {
|
||||
let request = serde_json::json!({
|
||||
"prompt": prompt,
|
||||
@@ -175,6 +191,11 @@ impl VoicePipe {
|
||||
.map_err(|e| format!("failed to serialize request: {}", e))?;
|
||||
line.push('\n');
|
||||
|
||||
// Check if child has already exited before writing
|
||||
if let Some(status) = self.child.try_wait().ok().flatten() {
|
||||
return Err(format!("sr-voice process exited with {}", status));
|
||||
}
|
||||
|
||||
self.stdin
|
||||
.write_all(line.as_bytes())
|
||||
.map_err(|e| format!("failed to write to sr-voice stdin: {}", e))?;
|
||||
@@ -182,13 +203,33 @@ impl VoicePipe {
|
||||
.flush()
|
||||
.map_err(|e| format!("failed to flush sr-voice stdin: {}", e))?;
|
||||
|
||||
// Watchdog: kill the child if it doesn't respond within the timeout.
|
||||
// This unblocks the read_line below (stdout closes → read returns empty).
|
||||
let child_id = self.child.id();
|
||||
let cancel = Arc::new(AtomicBool::new(false));
|
||||
let cancel_clone = Arc::clone(&cancel);
|
||||
let watchdog = std::thread::spawn(move || {
|
||||
std::thread::sleep(INFERENCE_TIMEOUT);
|
||||
if !cancel_clone.load(Ordering::Relaxed) {
|
||||
tracing::warn!(pid = child_id, "sr-voice inference timeout — killing child");
|
||||
let _ = Command::new("kill")
|
||||
.arg("-9")
|
||||
.arg(child_id.to_string())
|
||||
.output();
|
||||
}
|
||||
});
|
||||
|
||||
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))?;
|
||||
let read_result = self.reader.read_line(&mut response_line);
|
||||
|
||||
// Cancel the watchdog — response arrived (or EOF).
|
||||
cancel.store(true, Ordering::Relaxed);
|
||||
let _ = watchdog.join();
|
||||
|
||||
read_result.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());
|
||||
return Err("sr-voice process closed stdout (timeout or crash)".to_string());
|
||||
}
|
||||
|
||||
let body: serde_json::Value = serde_json::from_str(&response_line)
|
||||
@@ -239,9 +280,20 @@ fn worker_loop(id: usize, config: VoiceProcessConfig, ctx: WorkerContext) {
|
||||
Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
|
||||
};
|
||||
|
||||
// Skip while paused (zone transition)
|
||||
if ctx.paused.load(Ordering::Relaxed) {
|
||||
continue;
|
||||
// During zone transitions the queue is paused — wait rather than
|
||||
// process (priorities may be stale). Sleep briefly and re-check.
|
||||
if ctx.paused.load(Ordering::SeqCst) {
|
||||
// Don't drop the request — sleep and retry the pause check.
|
||||
// The request stays in our local variable until we can process it.
|
||||
while ctx.paused.load(Ordering::SeqCst) {
|
||||
if ctx.shutdown.load(Ordering::SeqCst) {
|
||||
break;
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(50));
|
||||
}
|
||||
if ctx.shutdown.load(Ordering::SeqCst) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
ctx.active_count.fetch_add(1, Ordering::Relaxed);
|
||||
|
||||
Reference in New Issue
Block a user