rtx-csm: Phase 6f.trace — per-phase timing in handle_connection

Adds explicit phase boundaries inside the WS connection loop so we can
see where voice-loop wall-time actually goes. Each turn emits one info
log line plus four new gauges on /metrics:

  recv_phase_ms          turn_start  -> EOT received (audio_recv + parallel STT)
  stt_post_phase_ms      EOT received -> conv_fut start (post-EOT flush)
  llm_to_first_audio_ms  conv_fut start -> first PCM byte (LLM + 1st sentence TTS)
  conv_total_ms          conv_fut start -> all sentences TTS'd
  total_turn_ms          turn_start -> done event

First measured turn (mock LLM, FP CSM, M-series Metal, 10.43s LibriSpeech):
  recv=4740ms  stt_post=29ms  llm_to_first_audio=4126ms
  conv_total=12735ms  total=17505ms

Two findings worth keeping:
  1. stt_post=29ms confirms the Phase 6c.3d parallel-STT optimization is
     working — the post-EOT flush is effectively free, all the heavy
     lifting happened during receive.
  2. Earlier "12s STT" estimate from the bench tool's transcript_ms was
     measuring the wrong thing (its clock includes bench-side audio_send
     that dumps frames at full speed; the server finishes receive in
     ~4.7s of which most is parallel STT). The instrumentation now
     attributes time correctly.

Co-Authored-By: Claude Opus 4.7 (1M context) <[email protected]>
This commit is contained in:
osobh
2026-04-27 04:04:27 -07:00
co-authored by Claude Opus 4.7
parent c56cc68cdc
commit f385999a18
@@ -241,6 +241,20 @@ struct Metrics {
/// End-to-end turn latency (audio-in → first-audio-out): sum + count. /// End-to-end turn latency (audio-in → first-audio-out): sum + count.
e2e_first_ms_sum: AtomicU64, e2e_first_ms_sum: AtomicU64,
e2e_first_ms_count: AtomicU64, e2e_first_ms_count: AtomicU64,
/// Receive phase: turn start → EOT (or VAD auto-EOT) received.
recv_phase_ms_sum: AtomicU64,
recv_phase_ms_count: AtomicU64,
/// LLM-to-first-audio: from conv_fut start (after STT done) to first
/// PCM chunk reaching the pump. Bundles LLM stream + first sentence
/// TTS gen.
llm_to_first_audio_ms_sum: AtomicU64,
llm_to_first_audio_ms_count: AtomicU64,
/// Full conversation phase: conv_fut start → all sentences TTS'd.
conv_total_ms_sum: AtomicU64,
conv_total_ms_count: AtomicU64,
/// Total turn wall time: turn start → done event sent.
total_turn_ms_sum: AtomicU64,
total_turn_ms_count: AtomicU64,
} }
struct Shared { struct Shared {
@@ -519,6 +533,10 @@ async fn metrics_handler(State(shared): State<Arc<Shared>>) -> impl IntoResponse
let stt_count = m.stt_latency_ms_count.load(Ordering::Relaxed).max(1); let stt_count = m.stt_latency_ms_count.load(Ordering::Relaxed).max(1);
let tts_count = m.tts_latency_ms_count.load(Ordering::Relaxed).max(1); let tts_count = m.tts_latency_ms_count.load(Ordering::Relaxed).max(1);
let e2e_count = m.e2e_first_ms_count.load(Ordering::Relaxed).max(1); let e2e_count = m.e2e_first_ms_count.load(Ordering::Relaxed).max(1);
let recv_count = m.recv_phase_ms_count.load(Ordering::Relaxed).max(1);
let llm_to_audio_count = m.llm_to_first_audio_ms_count.load(Ordering::Relaxed).max(1);
let conv_count = m.conv_total_ms_count.load(Ordering::Relaxed).max(1);
let total_count = m.total_turn_ms_count.load(Ordering::Relaxed).max(1);
let body = format!( let body = format!(
"# TYPE rtx_csm_turns_total counter\n\ "# TYPE rtx_csm_turns_total counter\n\
rtx_csm_turns_total {}\n\ rtx_csm_turns_total {}\n\
@@ -533,7 +551,15 @@ async fn metrics_handler(State(shared): State<Arc<Shared>>) -> impl IntoResponse
# TYPE rtx_csm_tts_latency_ms_avg gauge\n\ # TYPE rtx_csm_tts_latency_ms_avg gauge\n\
rtx_csm_tts_latency_ms_avg {}\n\ rtx_csm_tts_latency_ms_avg {}\n\
# TYPE rtx_csm_e2e_first_audio_ms_avg gauge\n\ # TYPE rtx_csm_e2e_first_audio_ms_avg gauge\n\
rtx_csm_e2e_first_audio_ms_avg {}\n", rtx_csm_e2e_first_audio_ms_avg {}\n\
# TYPE rtx_csm_recv_phase_ms_avg gauge\n\
rtx_csm_recv_phase_ms_avg {}\n\
# TYPE rtx_csm_llm_to_first_audio_ms_avg gauge\n\
rtx_csm_llm_to_first_audio_ms_avg {}\n\
# TYPE rtx_csm_conv_total_ms_avg gauge\n\
rtx_csm_conv_total_ms_avg {}\n\
# TYPE rtx_csm_total_turn_ms_avg gauge\n\
rtx_csm_total_turn_ms_avg {}\n",
m.turns_total.load(Ordering::Relaxed), m.turns_total.load(Ordering::Relaxed),
m.errors_total.load(Ordering::Relaxed), m.errors_total.load(Ordering::Relaxed),
m.connections_total.load(Ordering::Relaxed), m.connections_total.load(Ordering::Relaxed),
@@ -541,6 +567,10 @@ async fn metrics_handler(State(shared): State<Arc<Shared>>) -> impl IntoResponse
m.stt_latency_ms_sum.load(Ordering::Relaxed) / stt_count, m.stt_latency_ms_sum.load(Ordering::Relaxed) / stt_count,
m.tts_latency_ms_sum.load(Ordering::Relaxed) / tts_count, m.tts_latency_ms_sum.load(Ordering::Relaxed) / tts_count,
m.e2e_first_ms_sum.load(Ordering::Relaxed) / e2e_count, m.e2e_first_ms_sum.load(Ordering::Relaxed) / e2e_count,
m.recv_phase_ms_sum.load(Ordering::Relaxed) / recv_count,
m.llm_to_first_audio_ms_sum.load(Ordering::Relaxed) / llm_to_audio_count,
m.conv_total_ms_sum.load(Ordering::Relaxed) / conv_count,
m.total_turn_ms_sum.load(Ordering::Relaxed) / total_count,
); );
( (
StatusCode::OK, StatusCode::OK,
@@ -720,6 +750,18 @@ async fn handle_connection(mut socket: WebSocket, shared: Arc<Shared>) {
} }
} }
} }
// Mark the end of the audio-receive phase. Everything after this
// is post-EOT processing (STT flush + LLM + TTS).
let recv_done_t = std::time::Instant::now();
let recv_phase_ms = turn_start.elapsed().as_millis() as u64;
shared
.metrics
.recv_phase_ms_sum
.fetch_add(recv_phase_ms, Ordering::Relaxed);
shared
.metrics
.recv_phase_ms_count
.fetch_add(1, Ordering::Relaxed);
// Drain the asr_delay buffer so any trailing words flush. // Drain the asr_delay buffer so any trailing words flush.
let final_events = { let final_events = {
let mut stt = shared.stt.lock().await; let mut stt = shared.stt.lock().await;
@@ -810,6 +852,11 @@ async fn handle_connection(mut socket: WebSocket, shared: Arc<Shared>) {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>(); let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
let history_clone = history.clone(); let history_clone = history.clone();
let metrics_for_tts = shared.clone(); let metrics_for_tts = shared.clone();
// Capture the moment conv_fut starts so we can derive
// llm_to_first_audio (LLM stream + first-sentence TTS) and
// conv_total (whole conversation) latencies.
let conv_start_t = std::time::Instant::now();
let mut llm_to_first_audio_ms: Option<u64> = None;
let conv_fut = async { let conv_fut = async {
let mut g = shared.generator.lock().await; let mut g = shared.generator.lock().await;
let mut conv = Converse::new(&shared.llm, &mut *g); let mut conv = Converse::new(&shared.llm, &mut *g);
@@ -869,6 +916,20 @@ async fn handle_connection(mut socket: WebSocket, shared: Arc<Shared>) {
.metrics .metrics
.e2e_first_ms_count .e2e_first_ms_count
.fetch_add(1, Ordering::Relaxed); .fetch_add(1, Ordering::Relaxed);
// llm_to_first_audio: from conv_fut
// start to the first PCM chunk arriving
// here. Bundles LLM stream + first
// sentence's TTS gen.
let lta = conv_start_t.elapsed().as_millis() as u64;
llm_to_first_audio_ms = Some(lta);
shared
.metrics
.llm_to_first_audio_ms_sum
.fetch_add(lta, Ordering::Relaxed);
shared
.metrics
.llm_to_first_audio_ms_count
.fetch_add(1, Ordering::Relaxed);
} }
if socket.send(Message::Binary(buf)).await.is_err() { if socket.send(Message::Binary(buf)).await.is_err() {
return PumpResult::Disconnected; return PumpResult::Disconnected;
@@ -938,6 +999,15 @@ async fn handle_connection(mut socket: WebSocket, shared: Arc<Shared>) {
} }
}; };
let (conv_res, pump_res) = tokio::join!(conv_fut, pump_fut); let (conv_res, pump_res) = tokio::join!(conv_fut, pump_fut);
let conv_total_ms = conv_start_t.elapsed().as_millis() as u64;
shared
.metrics
.conv_total_ms_sum
.fetch_add(conv_total_ms, Ordering::Relaxed);
shared
.metrics
.conv_total_ms_count
.fetch_add(1, Ordering::Relaxed);
let mut barge_in_audio: Option<Vec<f32>> = None; let mut barge_in_audio: Option<Vec<f32>> = None;
match pump_res { match pump_res {
PumpResult::Done => {} PumpResult::Done => {}
@@ -973,6 +1043,31 @@ async fn handle_connection(mut socket: WebSocket, shared: Arc<Shared>) {
let done_msg = serde_json::json!({"event":"done","assistant":&assistant_text}); let done_msg = serde_json::json!({"event":"done","assistant":&assistant_text});
send_text(&mut socket, &done_msg.to_string()).await; send_text(&mut socket, &done_msg.to_string()).await;
} }
let total_turn_ms = turn_start.elapsed().as_millis() as u64;
shared
.metrics
.total_turn_ms_sum
.fetch_add(total_turn_ms, Ordering::Relaxed);
shared
.metrics
.total_turn_ms_count
.fetch_add(1, Ordering::Relaxed);
// Per-turn timing summary. recv + stt_post + conv_total ≈ total
// (small slack for mutex acquire, channel close, log formatting).
// llm_to_first_audio is a sub-phase of conv_total: it's the time
// from conv_fut start to the first PCM byte hitting the pump.
let stt_post_ms = conv_start_t
.saturating_duration_since(recv_done_t)
.as_millis() as u64;
tracing::info!(
"turn timing: recv={}ms stt_post={}ms llm_to_first_audio={}ms \
conv_total={}ms total={}ms",
recv_phase_ms,
stt_post_ms,
llm_to_first_audio_ms.unwrap_or(0),
conv_total_ms,
total_turn_ms,
);
// Stash any barge-in audio so the next turn picks up where the user started. // Stash any barge-in audio so the next turn picks up where the user started.
carry_over = barge_in_audio; carry_over = barge_in_audio;
} }