1use std::fs;
24use std::hash::{BuildHasher, Hasher, RandomState};
25use std::io::{BufRead, BufReader, Read, Write};
26use std::net::{Ipv4Addr, TcpListener, TcpStream};
27use std::path::{Path, PathBuf};
28use std::process::Command;
29use std::time::{Duration, Instant, SystemTime};
30
31use serde_json::{Value, json};
32
33use crate::error::{FttsError, FttsExitCode};
34use crate::synth::{LoadedModel, ModelBundle, SynthesizedAudio};
35use ftts_core::{NormalizationMode, NormalizationOptions, SynthesisRequest};
36
37const DEFAULT_IDLE: Duration = Duration::from_secs(600);
40
41const DEFAULT_CLIENT_READ_TIMEOUT: Duration = Duration::from_secs(600);
46
47fn client_read_timeout() -> Duration {
48 std::env::var("FTTS_RESIDENT_CLIENT_TIMEOUT_SECS")
49 .ok()
50 .and_then(|value| value.parse::<u64>().ok())
51 .filter(|&seconds| seconds > 0)
55 .map_or(DEFAULT_CLIENT_READ_TIMEOUT, Duration::from_secs)
56}
57
58const DEFAULT_SPAWN_WAIT: Duration = Duration::from_secs(30);
64
65fn spawn_wait() -> Duration {
66 std::env::var("FTTS_RESIDENT_SPAWN_WAIT_SECS")
67 .ok()
68 .and_then(|value| value.parse::<u64>().ok())
69 .map_or(DEFAULT_SPAWN_WAIT, Duration::from_secs)
70}
71
72const PROTOCOL: u64 = 1;
73
74pub struct WireRequest<'a> {
77 pub text: &'a str,
78 pub normalize: &'a str,
80 pub trace: bool,
81 pub speaker: &'a [f32],
82 pub seed: u64,
83}
84
85fn parse_normalize(label: &str) -> Option<NormalizationMode> {
86 match label {
87 "verbatim" => Some(NormalizationMode::Verbatim),
88 "conservative" => Some(NormalizationMode::Conservative),
89 "locale-aware" => Some(NormalizationMode::LocaleAware),
90 _ => None,
91 }
92}
93
94fn idle_period() -> Duration {
95 std::env::var("FTTS_RESIDENT_IDLE_SECS")
96 .ok()
97 .and_then(|value| value.parse::<u64>().ok())
98 .map_or(DEFAULT_IDLE, Duration::from_secs)
99}
100
101pub fn enabled(no_resident_flag: bool) -> bool {
104 if no_resident_flag {
105 return false;
106 }
107 !matches!(
108 std::env::var("FTTS_NO_RESIDENT").ok().as_deref(),
109 Some("1") | Some("true")
110 )
111}
112
113fn resident_dir() -> Option<PathBuf> {
116 if let Ok(dir) = std::env::var("FTTS_RESIDENT_DIR") {
117 return Some(PathBuf::from(dir));
118 }
119 #[allow(deprecated)] std::env::home_dir().map(|home| home.join(".cache/franken_tts"))
121}
122
123fn root_digest(root: &Path) -> u64 {
132 let canonical = fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf());
133 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
134 for byte in canonical.to_string_lossy().as_bytes() {
135 hash ^= u64::from(*byte);
136 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
137 }
138 hash
139}
140
141fn state_path(root: &Path) -> Option<PathBuf> {
142 resident_dir().map(|dir| dir.join(format!("resident-{:016x}.json", root_digest(root))))
143}
144
145struct DaemonState {
146 port: u16,
147 token: String,
148}
149
150fn read_state(root: &Path) -> Option<DaemonState> {
151 let path = state_path(root)?;
152 let raw = fs::read_to_string(path).ok()?;
153 let value: Value = serde_json::from_str(&raw).ok()?;
154 Some(DaemonState {
155 port: u16::try_from(value.get("port")?.as_u64()?).ok()?,
156 token: value.get("token")?.as_str()?.to_owned(),
157 })
158}
159
160fn write_state(root: &Path, port: u16, token: &str) -> std::io::Result<PathBuf> {
161 let path = state_path(root).ok_or_else(|| {
162 std::io::Error::other("no home directory and no FTTS_RESIDENT_DIR; cannot go resident")
163 })?;
164 if let Some(parent) = path.parent() {
165 fs::create_dir_all(parent)?;
166 }
167 let body = json!({
168 "port": port,
169 "token": token,
170 "pid": std::process::id(),
171 "version": env!("CARGO_PKG_VERSION"),
172 "bundle_root": root.to_string_lossy(),
173 })
174 .to_string();
175 let staging = path.with_extension("json.tmp");
179 #[cfg(unix)]
180 {
181 use std::io::Write as _;
182 use std::os::unix::fs::OpenOptionsExt;
183 let mut file = fs::OpenOptions::new()
184 .write(true)
185 .create(true)
186 .truncate(true)
187 .mode(0o600)
188 .open(&staging)?;
189 file.write_all(body.as_bytes())?;
190 }
191 #[cfg(not(unix))]
192 fs::write(&staging, &body)?;
193 fs::rename(&staging, &path)?;
194 Ok(path)
195}
196
197fn fresh_token() -> String {
200 let mut token = String::with_capacity(32);
201 for _ in 0..2 {
202 let mut hasher = RandomState::new().build_hasher();
203 hasher.write_u128(std::time::UNIX_EPOCH.elapsed().map_or(0, |d| d.as_nanos()));
204 hasher.write_u32(std::process::id());
205 token.push_str(&format!("{:016x}", hasher.finish()));
206 }
207 token
208}
209
210fn artifact_stamp(bundle: &ModelBundle) -> (u64, u64) {
213 let path = bundle.canonical_main.as_deref().unwrap_or(&bundle.main);
214 let Ok(meta) = fs::metadata(path) else {
215 return (0, 0);
216 };
217 let mtime = meta
218 .modified()
219 .ok()
220 .and_then(|time| time.duration_since(SystemTime::UNIX_EPOCH).ok())
221 .map_or(0, |d| d.as_secs());
222 (mtime, meta.len())
223}
224
225pub fn try_synthesize(
234 bundle: &ModelBundle,
235 request: &WireRequest<'_>,
236) -> Result<Option<SynthesizedAudio>, FttsError> {
237 if request.speaker.iter().any(|value| !value.is_finite()) {
241 return Ok(None);
242 }
243 match connect(bundle) {
244 Some(stream) => roundtrip(stream, bundle, request),
245 None => Ok(None),
246 }
247}
248
249fn connect(bundle: &ModelBundle) -> Option<TcpStream> {
250 for attempt in 0..3 {
254 if attempt > 0 {
255 std::thread::sleep(Duration::from_millis(300));
256 }
257 if let Some(stream) = connect_once(bundle) {
258 return Some(stream);
259 }
260 if read_state(&bundle.root).is_none() {
261 break; }
263 }
264 if let Some(path) = state_path(&bundle.root)
266 && path.exists()
267 {
268 let _ = fs::remove_file(&path);
269 }
270 spawn_daemon(&bundle.root)?;
271 let deadline = Instant::now() + spawn_wait();
272 while Instant::now() < deadline {
273 if let Some(stream) = connect_once(bundle) {
274 return Some(stream);
275 }
276 std::thread::sleep(Duration::from_millis(50));
277 }
278 None
279}
280
281fn connect_once(bundle: &ModelBundle) -> Option<TcpStream> {
282 let state = read_state(&bundle.root)?;
283 let stream = TcpStream::connect_timeout(
284 &(Ipv4Addr::LOCALHOST, state.port).into(),
285 Duration::from_millis(1000),
286 )
287 .ok()?;
288 stream.set_read_timeout(Some(client_read_timeout())).ok()?;
289 stream.set_nodelay(true).ok();
290 Some(stream)
291}
292
293fn spawn_daemon(root: &Path) -> Option<()> {
294 let exe = std::env::current_exe().ok()?;
295 let mut command = Command::new(exe);
296 command
297 .arg("resident-daemon")
298 .arg("--bundle-root")
299 .arg(root)
300 .stdin(std::process::Stdio::null())
301 .stdout(std::process::Stdio::null())
302 .stderr(std::process::Stdio::null());
303 if let Ok(log) = std::env::var("FTTS_RESIDENT_LOG")
306 && let Ok(file) = fs::OpenOptions::new().create(true).append(true).open(&log)
307 && let Ok(err) = file.try_clone()
308 {
309 command.stdout(file).stderr(err);
310 }
311 #[cfg(windows)]
312 {
313 use std::os::windows::process::CommandExt;
314 command.creation_flags(0x0000_0008 | 0x0800_0000);
316 }
317 command.spawn().ok().map(|_child| ())
318}
319
320fn roundtrip(
321 mut stream: TcpStream,
322 bundle: &ModelBundle,
323 request: &WireRequest<'_>,
324) -> Result<Option<SynthesizedAudio>, FttsError> {
325 let state = match read_state(&bundle.root) {
326 Some(state) => state,
327 None => return Ok(None),
328 };
329 let header = json!({
330 "protocol": PROTOCOL,
331 "op": "synthesize",
332 "token": state.token,
333 "version": env!("CARGO_PKG_VERSION"),
334 "bundle_root": bundle.root.to_string_lossy(),
335 "text": request.text,
336 "normalize": request.normalize,
337 "trace": request.trace,
338 "speaker": request.speaker,
339 "seed": request.seed,
340 });
341 if stream
342 .write_all(format!("{header}\n").as_bytes())
343 .and_then(|()| stream.flush())
344 .is_err()
345 {
346 return Ok(None);
347 }
348
349 let mut reader = BufReader::new(stream);
350 let mut line = String::new();
351 if reader.read_line(&mut line).is_err() || line.trim().is_empty() {
352 return Ok(None);
353 }
354 let Ok(reply) = serde_json::from_str::<Value>(&line) else {
355 return Ok(None);
356 };
357 if reply.get("ok").and_then(Value::as_bool) == Some(true) {
358 const MAX_WIRE_SAMPLES: u64 = 200_000_000;
363 let samples = reply
364 .get("samples")
365 .and_then(Value::as_u64)
366 .filter(|&n| n <= MAX_WIRE_SAMPLES)
367 .and_then(|n| usize::try_from(n).ok())
368 .unwrap_or(0);
369 let mut bytes = vec![0u8; samples * 4];
370 if reader.read_exact(&mut bytes).is_err() {
371 return Ok(None);
372 }
373 let (chunks, _remainder) = bytes.as_chunks::<4>();
374 let pcm = chunks
375 .iter()
376 .map(|chunk| f32::from_le_bytes(*chunk))
377 .collect();
378 return Ok(Some(SynthesizedAudio {
379 pcm,
380 frames: reply.get("frames").and_then(Value::as_u64).unwrap_or(0),
381 prepared_token_count: reply
382 .get("prepared_token_count")
383 .and_then(Value::as_u64)
384 .unwrap_or(0) as usize,
385 ttfa: reply
386 .get("ttfa_ms")
387 .and_then(Value::as_u64)
388 .map(Duration::from_millis),
389 }));
390 }
391 match reply.get("kind").and_then(Value::as_str) {
393 Some("synthesis") => {
394 let code = reply
395 .get("exit_code")
396 .and_then(Value::as_u64)
397 .unwrap_or(FttsExitCode::Generic.as_u8().into());
398 let message = reply
399 .get("message")
400 .and_then(Value::as_str)
401 .unwrap_or("resident synthesis failed")
402 .to_owned();
403 Err(wire_error(code, message))
404 }
405 _ => Ok(None),
406 }
407}
408
409fn wire_error(exit_code: u64, message: String) -> FttsError {
410 match exit_code {
411 3 => FttsError::ModelNotFound(message),
412 4 => FttsError::Input(message),
413 5 => FttsError::BudgetTimeout(message),
414 7 => FttsError::ArtifactFormat(message),
415 8 => FttsError::EnrollmentQualityRefusal(message),
416 _ => FttsError::Generic(message),
417 }
418}
419
420pub fn run_daemon(bundle_root: &Path) -> Result<(), FttsError> {
428 let bundle = ModelBundle::resolve(bundle_root)?;
429 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0))
430 .map_err(|error| FttsError::Generic(format!("resident daemon cannot bind: {error}")))?;
431 let port = listener
432 .local_addr()
433 .map_err(|error| FttsError::Generic(format!("resident daemon has no address: {error}")))?
434 .port();
435 let token = fresh_token();
436 let state_file = write_state(&bundle.root, port, &token)
437 .map_err(|error| FttsError::Generic(format!("cannot write resident state: {error}")))?;
438 listener
439 .set_nonblocking(true)
440 .map_err(|error| FttsError::Generic(format!("resident daemon socket mode: {error}")))?;
441 eprintln!(
442 "resident daemon serving {} on 127.0.0.1:{port}",
443 bundle.root.display()
444 );
445
446 let idle = idle_period();
447 let mut resident: Option<(LoadedModel, ftts_core::TtsEngine, (u64, u64))> = None;
448 let mut deadline = Instant::now() + idle;
449
450 loop {
451 match listener.accept() {
452 Ok((stream, _peer)) => {
453 let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
465 handle_connection(stream, &bundle, &token, port, &mut resident);
466 }));
467 if outcome.is_err() {
468 eprintln!("resident daemon: request panicked; connection dropped, model kept");
470 }
471 deadline = Instant::now() + idle;
472 }
473 Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
474 if Instant::now() >= deadline {
475 eprintln!("resident daemon idle exit");
476 remove_state_if_ours(&state_file, port);
477 return Ok(());
478 }
479 std::thread::sleep(Duration::from_millis(100));
480 }
481 Err(_) => {
482 remove_state_if_ours(&state_file, port);
483 return Ok(());
484 }
485 }
486 }
487}
488
489fn remove_state_if_ours(state_file: &Path, port: u16) {
498 let ours = fs::read_to_string(state_file)
499 .ok()
500 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
501 .and_then(|value| value.get("port")?.as_u64())
502 .is_some_and(|recorded| recorded == u64::from(port));
503 if ours {
504 let _ = fs::remove_file(state_file);
505 }
506}
507
508#[allow(clippy::type_complexity)]
509fn handle_connection(
510 stream: TcpStream,
511 bundle: &ModelBundle,
512 token: &str,
513 port: u16,
514 resident: &mut Option<(LoadedModel, ftts_core::TtsEngine, (u64, u64))>,
515) {
516 let _ = stream.set_nonblocking(false);
517 let _ = stream.set_read_timeout(Some(Duration::from_secs(10)));
518 let _ = stream.set_write_timeout(Some(Duration::from_secs(60)));
524 let _ = stream.set_nodelay(true);
525
526 const MAX_REQUEST_BYTES: u64 = 1024 * 1024;
534 let mut reader = BufReader::new(stream);
535 let mut line = String::new();
536 if (&mut reader)
537 .take(MAX_REQUEST_BYTES)
538 .read_line(&mut line)
539 .is_err()
540 {
541 return;
542 }
543 let mut stream = reader.into_inner();
544 if line.len() as u64 >= MAX_REQUEST_BYTES {
547 let reply = json!({ "ok": false, "kind": "request", "message": "request too large" });
548 let _ = stream.write_all(format!("{reply}\n").as_bytes());
549 return;
550 }
551 let Ok(request) = serde_json::from_str::<Value>(&line) else {
552 return;
553 };
554
555 let refuse = |stream: &mut TcpStream, kind: &str, message: &str| {
556 let reply = json!({ "ok": false, "kind": kind, "message": message });
557 let _ = stream.write_all(format!("{reply}\n").as_bytes());
558 };
559
560 if request.get("token").and_then(Value::as_str) != Some(token) {
562 refuse(&mut stream, "auth", "bad token");
563 return;
564 }
565 if request.get("protocol").and_then(Value::as_u64) != Some(PROTOCOL)
566 || request.get("version").and_then(Value::as_str) != Some(env!("CARGO_PKG_VERSION"))
567 {
568 refuse(&mut stream, "version", "resident daemon version mismatch");
571 if let Some(state) = state_path(&bundle.root) {
574 remove_state_if_ours(&state, port);
575 }
576 std::process::exit(0);
577 }
578
579 if let Some(wire_root) = request.get("bundle_root").and_then(Value::as_str) {
583 let wire_canonical =
584 fs::canonicalize(wire_root).unwrap_or_else(|_| PathBuf::from(wire_root));
585 let own_canonical = fs::canonicalize(&bundle.root).unwrap_or_else(|_| bundle.root.clone());
586 if wire_canonical != own_canonical {
587 refuse(
588 &mut stream,
589 "request",
590 "resident daemon serves a different bundle root",
591 );
592 return;
593 }
594 }
595
596 let stamp_now = artifact_stamp(bundle);
598 if let Some((_, _, loaded_stamp)) = resident.as_ref()
599 && *loaded_stamp != stamp_now
600 {
601 refuse(&mut stream, "stale", "model artifact changed since load");
602 if let Some(state) = state_path(&bundle.root) {
603 remove_state_if_ours(&state, port);
604 }
605 std::process::exit(0);
606 }
607
608 let text = request.get("text").and_then(Value::as_str).unwrap_or("");
609 let normalize = request
610 .get("normalize")
611 .and_then(Value::as_str)
612 .and_then(parse_normalize);
613 let trace = request
614 .get("trace")
615 .and_then(Value::as_bool)
616 .unwrap_or(false);
617 let seed = request.get("seed").and_then(Value::as_u64).unwrap_or(0);
618 let speaker: Vec<f32> = match request.get("speaker").and_then(Value::as_array) {
629 Some(values) => {
630 let mut vector = Vec::with_capacity(values.len());
631 for value in values {
632 let Some(number) = value.as_f64() else {
633 refuse(
634 &mut stream,
635 "request",
636 "speaker vector contains a non-numeric entry",
637 );
638 return;
639 };
640 #[allow(clippy::cast_possible_truncation)]
641 let narrowed = number as f32;
642 if !narrowed.is_finite() {
643 refuse(
644 &mut stream,
645 "request",
646 "speaker vector contains a non-finite value",
647 );
648 return;
649 }
650 vector.push(narrowed);
651 }
652 vector
653 }
654 None => Vec::new(),
655 };
656 let Some(mode) = normalize else {
657 refuse(&mut stream, "request", "unknown normalize mode");
658 return;
659 };
660 if text.is_empty() || speaker.is_empty() {
661 refuse(&mut stream, "request", "empty text or speaker");
662 return;
663 }
664
665 if resident.is_none() {
667 let loaded = match LoadedModel::load(bundle) {
668 Ok(loaded) => loaded,
669 Err(error) => {
670 let reply = json!({
671 "ok": false,
672 "kind": "synthesis",
673 "exit_code": error.exit_code().as_u8(),
674 "message": error.to_string(),
675 });
676 let _ = stream.write_all(format!("{reply}\n").as_bytes());
677 return;
678 }
679 };
680 let engine = match ftts_core::TtsEngine::from_process_environment() {
681 Ok(engine) => engine,
682 Err(error) => {
683 refuse(
684 &mut stream,
685 "engine",
686 &format!("engine start failed: {error}"),
687 );
688 return;
689 }
690 };
691 *resident = Some((loaded, engine, stamp_now));
692 }
693 let (loaded, engine, _) = resident.as_ref().expect("hydrated just above");
694
695 let synthesis_request = SynthesisRequest::new(text.to_owned())
696 .with_normalization_options(NormalizationOptions {
697 mode,
698 ..NormalizationOptions::default()
699 })
700 .with_normalization_trace(trace);
701 let cancellation = ftts_core::CancellationToken::new();
702 let observer = |_event: ftts_core::SynthesisEvent| {};
703 match crate::synth::synthesize(
704 loaded,
705 engine,
706 &synthesis_request,
707 &speaker,
708 seed,
709 &cancellation,
710 &observer,
711 ) {
712 Ok(audio) => {
713 let header = json!({
714 "ok": true,
715 "samples": audio.pcm.len(),
716 "frames": audio.frames,
717 "prepared_token_count": audio.prepared_token_count,
718 "ttfa_ms": audio.ttfa.map(|d| u64::try_from(d.as_millis()).unwrap_or(u64::MAX)),
719 });
720 if stream.write_all(format!("{header}\n").as_bytes()).is_err() {
721 return;
722 }
723 let mut bytes = Vec::with_capacity(audio.pcm.len() * 4);
724 for sample in &audio.pcm {
725 bytes.extend_from_slice(&sample.to_le_bytes());
726 }
727 let _ = stream.write_all(&bytes);
728 let _ = stream.flush();
729 }
730 Err(error) => {
731 let reply = json!({
732 "ok": false,
733 "kind": "synthesis",
734 "exit_code": error.exit_code().as_u8(),
735 "message": error.to_string(),
736 });
737 let _ = stream.write_all(format!("{reply}\n").as_bytes());
738 }
739 }
740}
741
742#[cfg(test)]
743mod tests {
744 use super::*;
745
746 #[test]
747 fn normalize_labels_round_trip() {
748 for label in ["verbatim", "conservative", "locale-aware"] {
749 assert!(parse_normalize(label).is_some(), "{label}");
750 }
751 assert!(parse_normalize("aggressive").is_none());
752 }
753
754 #[test]
758 fn an_endless_request_line_is_bounded_not_fatal() {
759 use std::io::Write as _;
760 use std::net::TcpListener;
761
762 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).unwrap();
763 let port = listener.local_addr().unwrap().port();
764
765 let flood = std::thread::spawn(move || {
767 let Ok(mut stream) = TcpStream::connect((Ipv4Addr::LOCALHOST, port)) else {
768 return;
769 };
770 let block = vec![b'x'; 64 * 1024];
771 while stream.write_all(&block).is_ok() {}
773 });
774
775 let (stream, _) = listener.accept().unwrap();
776 let mut reader = BufReader::new(stream);
777 let mut line = String::new();
778 const MAX: u64 = 1024 * 1024;
779 let read = (&mut reader).take(MAX).read_line(&mut line).unwrap_or(0);
780
781 assert!(
782 read as u64 <= MAX,
783 "read {read} bytes, past the {MAX}-byte cap"
784 );
785 assert!(
786 line.len() as u64 <= MAX,
787 "buffered {} bytes, past the cap",
788 line.len()
789 );
790 drop(reader);
791 let _ = flood.join();
792 }
793
794 #[test]
801 fn speaker_vectors_are_validated_rather_than_salvaged() {
802 fn parse(values: &[Value]) -> Result<Vec<f32>, &'static str> {
804 let mut vector = Vec::with_capacity(values.len());
805 for value in values {
806 let number = value.as_f64().ok_or("non-numeric")?;
807 let narrowed = number as f32;
808 if !narrowed.is_finite() {
809 return Err("non-finite");
810 }
811 vector.push(narrowed);
812 }
813 Ok(vector)
814 }
815
816 assert_eq!(parse(&[json!(1.0), json!(-2.5)]).unwrap(), vec![1.0, -2.5]);
817 assert_eq!(
818 parse(&[json!(1.0), json!("x"), json!(2.0)]),
819 Err("non-numeric"),
820 "a non-numeric entry must be refused, not dropped"
821 );
822 assert_eq!(parse(&[json!(null)]), Err("non-numeric"));
823 assert_eq!(
825 parse(&[json!(1e300)]),
826 Err("non-finite"),
827 "an f64 that overflows f32 becomes infinity and must be refused"
828 );
829 }
830
831 #[test]
832 fn wire_errors_keep_their_exit_class() {
833 let cases = [
834 (3u64, FttsExitCode::ModelNotFound),
835 (4, FttsExitCode::Input),
836 (5, FttsExitCode::BudgetTimeout),
837 (7, FttsExitCode::ArtifactFormat),
838 (8, FttsExitCode::EnrollmentQualityRefusal),
839 (1, FttsExitCode::Generic),
840 (99, FttsExitCode::Generic),
841 ];
842 for (code, expected) in cases {
843 assert_eq!(wire_error(code, String::new()).exit_code(), expected);
844 }
845 }
846
847 #[test]
848 fn root_digest_distinguishes_roots_and_is_stable() {
849 let a = root_digest(Path::new("/tmp/model-a"));
850 let b = root_digest(Path::new("/tmp/model-b"));
851 assert_ne!(a, b);
852 assert_eq!(a, root_digest(Path::new("/tmp/model-a")));
853 }
854
855 #[test]
856 fn tokens_are_distinct_and_hex() {
857 let one = fresh_token();
858 let two = fresh_token();
859 assert_eq!(one.len(), 32);
860 assert!(one.bytes().all(|b| b.is_ascii_hexdigit()));
861 assert_ne!(one, two, "two RandomState-seeded tokens collided");
862 }
863
864 #[test]
865 fn state_file_round_trips_and_respects_dir_override() {
866 let dir = std::env::temp_dir().join(format!("ftts-resident-test-{}", std::process::id()));
867 let root = Path::new("/tmp/some-model-root");
870 let path = dir.join(format!("resident-{:016x}.json", root_digest(root)));
871 fs::create_dir_all(&dir).unwrap();
872 let body =
873 json!({"port": 45123, "token": "abc123", "pid": 1, "version": "x", "bundle_root": "y"});
874 fs::write(&path, body.to_string()).unwrap();
875 let raw = fs::read_to_string(&path).unwrap();
876 let value: Value = serde_json::from_str(&raw).unwrap();
877 assert_eq!(value.get("port").and_then(Value::as_u64), Some(45123));
878 assert_eq!(value.get("token").and_then(Value::as_str), Some("abc123"));
879 let _ = fs::remove_file(&path);
880 }
881
882 #[test]
885 fn daemon_refuses_bad_token() {
886 let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).unwrap();
887 let port = listener.local_addr().unwrap().port();
888 let handle = std::thread::spawn(move || {
889 let (stream, _) = listener.accept().unwrap();
890 let bundle = ModelBundle {
891 root: PathBuf::from("/nonexistent"),
892 main: PathBuf::from("/nonexistent/model.safetensors"),
893 canonical_main: None,
894 codec: PathBuf::from("/nonexistent/codec"),
895 };
896 let mut resident = None;
897 handle_connection(stream, &bundle, "right-token", 0, &mut resident);
898 });
899 let mut stream = TcpStream::connect((Ipv4Addr::LOCALHOST, port)).unwrap();
900 let request = json!({
901 "protocol": PROTOCOL,
902 "op": "synthesize",
903 "token": "wrong-token",
904 "version": env!("CARGO_PKG_VERSION"),
905 "text": "hi",
906 "normalize": "verbatim",
907 "speaker": [0.0],
908 "seed": 0,
909 });
910 stream.write_all(format!("{request}\n").as_bytes()).unwrap();
911 let mut reply = String::new();
912 BufReader::new(&mut stream).read_line(&mut reply).unwrap();
913 let value: Value = serde_json::from_str(&reply).unwrap();
914 assert_eq!(value.get("ok").and_then(Value::as_bool), Some(false));
915 assert_eq!(value.get("kind").and_then(Value::as_str), Some("auth"));
916 handle.join().unwrap();
917 }
918}