From a590c735cfa4ec55ba6b32e297600d83e1a3e46a Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Fri, 10 Jul 2026 18:03:29 +0800
Subject: [PATCH] feat: stream realtime asr from helper

---
 docs/runtime-contract.md                |    1 
 src/audio.rs                            |   10 
 fixtures/asr-realtime-happy.ndjson      |    5 
 src/asr_realtime.rs                     |  584 ++++++++++++++++++++++++++++++++++++
 Cargo.toml                              |    1 
 src/main.rs                             |  183 +++++++++++
 tools/validate-asr-realtime-fixture.mjs |  116 +++++++
 Cargo.lock                              |    1 
 src/service.rs                          |   15 
 README.md                               |    5 
 10 files changed, 907 insertions(+), 14 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock
index 03c9ec2..deae385 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -240,6 +240,7 @@
  "reqwest",
  "serde",
  "serde_json",
+ "sha2",
  "tokio",
  "tracing",
  "tracing-subscriber",
diff --git a/Cargo.toml b/Cargo.toml
index a6d823d..c048bb8 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -15,6 +15,7 @@
 reqwest = { version = "0.12.24", default-features = false, features = ["rustls-tls", "json", "stream"] }
 serde = { version = "1.0.228", features = ["derive"] }
 serde_json = "1.0.145"
+sha2 = "0.10.9"
 tokio = { version = "1.48.0", features = ["macros", "rt-multi-thread", "signal", "time"] }
 tracing = "0.1.41"
 tracing-subscriber = { version = "0.3.20", features = ["env-filter", "fmt"] }
diff --git a/README.md b/README.md
index d81dfdc..e9ba507 100644
--- a/README.md
+++ b/README.md
@@ -54,9 +54,12 @@
 ```bash
 node tools/validate-turn-stream-fixture.mjs fixtures/turn-stream-happy.ndjson
 node tools/validate-turn-stream-fixture.mjs fixtures/turn-stream-mp3-chunks.ndjson
+node tools/validate-asr-realtime-fixture.mjs fixtures/asr-realtime-happy.ndjson
 ```
 
 这条快速线只校验 NDJSON contract、事件顺序和 `pcm_s16le` / `mp3` chunk 基本约束,用于提前发现 `replyPlaybackMode`、`reply_state`、`reply_audio_chunk`、`turn_completed` 等字段破坏。它不能替代 Docker 镜像构建、真实 LiveKit smoke 或 iPhone 真机验收。
+
+ASR realtime fixture 额外校验 `session_start -> audio_chunk -> vad_speech_end -> finish`、连续 `chunkSeq`、`pcm_s16le / 16000Hz / mono` 与 `24KiB` 单行上限。它只验证 helper 到 Java 的 wire contract,不代表 provider partial/final 已通过。
 
 ## 本机运行
 
@@ -98,6 +101,8 @@
 
 `sessions/start` 收到 Java 传入的 LiveKit bot 入房材料后,会拉起现有 worker 子进程承接媒体链路。helper service 本身不做 ASR / LLM / TTS / 消息 / 计费,也不持久化业务数据。
 
+当 Java 在 `sessions/start.runtime` 下发 `asrRealtimeEnabled=true`、`asrRealtimeUrl` 与 `asrRealtimeChunkDurationMs=200` 时,worker 会从 VAD speech start 起接收 20ms 用户音频帧,聚合成 `16kHz / mono / 200ms` NDJSON chunk 持续上传给 Java;最后一块允许短于 200ms。helper 只负责音频传输和 fallback 编排,ASR provider session、partial/final 归一化、activity 和 `asrResultRef` 仍由 Java 管理;realtime 失败时只回退一次既有 final-only ASR stream。
+
 鉴权口径:
 
 - Java 调 helper 控制面使用 `Authorization: Bearer {helperAuthToken}`。
diff --git a/docs/runtime-contract.md b/docs/runtime-contract.md
index d521b07..aa48970 100644
--- a/docs/runtime-contract.md
+++ b/docs/runtime-contract.md
@@ -27,6 +27,7 @@
 ```bash
 node tools/validate-turn-stream-fixture.mjs fixtures/turn-stream-happy.ndjson
 node tools/validate-turn-stream-fixture.mjs fixtures/turn-stream-mp3-chunks.ndjson
+node tools/validate-asr-realtime-fixture.mjs fixtures/asr-realtime-happy.ndjson
 ```
 
 The fixture validates only the protocol shape and fast-path invariants:
diff --git a/fixtures/asr-realtime-happy.ndjson b/fixtures/asr-realtime-happy.ndjson
new file mode 100644
index 0000000..50b560b
--- /dev/null
+++ b/fixtures/asr-realtime-happy.ndjson
@@ -0,0 +1,5 @@
+{"event":"session_start","callId":"call_fixture","traceId":"trace_fixture","turnId":"turn-0001","runtimeSessionNonceHash":"d7c778da6f43","audio":{"format":"pcm_s16le","sampleRate":16000,"channels":1}}
+{"event":"audio_chunk","chunkSeq":1,"audioBase64":"AAAAAAAAAAA=","durationMs":200}
+{"event":"audio_chunk","chunkSeq":2,"audioBase64":"AAAAAAAAAAA=","durationMs":200}
+{"event":"vad_speech_end","speechDurationMs":400,"endReason":"silence"}
+{"event":"finish"}
diff --git a/src/asr_realtime.rs b/src/asr_realtime.rs
new file mode 100644
index 0000000..db28b8c
--- /dev/null
+++ b/src/asr_realtime.rs
@@ -0,0 +1,584 @@
+use std::{io, time::Duration};
+
+use anyhow::{Context, Result, anyhow};
+use base64::{Engine as _, engine::general_purpose};
+use futures_util::stream;
+use reqwest::Client;
+use serde::Deserialize;
+use serde_json::json;
+use sha2::{Digest, Sha256};
+use tokio::{
+    sync::mpsc,
+    task::JoinHandle,
+    time::{Instant, timeout},
+};
+
+const SAMPLE_RATE_16K: u32 = 16_000;
+const CHANNELS_MONO: u32 = 1;
+const MIN_CHUNK_DURATION_MS: u64 = 20;
+const MAX_CHUNK_DURATION_MS: u64 = 1_000;
+const UPLOAD_QUEUE_CAPACITY: usize = 64;
+const MAX_NDJSON_LINE_BYTES: usize = 24 * 1024;
+const FINISH_RESPONSE_TIMEOUT: Duration = Duration::from_secs(15);
+const CANCEL_RESPONSE_TIMEOUT: Duration = Duration::from_secs(3);
+
+#[derive(Clone)]
+pub(crate) struct RealtimeAsrConfig {
+    pub(crate) enabled: bool,
+    pub(crate) url: Option<String>,
+    pub(crate) runtime_token: Option<String>,
+    pub(crate) runtime_session_nonce: Option<String>,
+    pub(crate) chunk_duration_ms: u64,
+}
+
+impl RealtimeAsrConfig {
+    pub(crate) fn is_ready(&self) -> bool {
+        self.enabled
+            && non_blank(&self.url)
+            && non_blank(&self.runtime_token)
+            && non_blank(&self.runtime_session_nonce)
+            && self.chunk_duration_ms >= MIN_CHUNK_DURATION_MS
+            && self.chunk_duration_ms <= MAX_CHUNK_DURATION_MS
+            && self.chunk_duration_ms % MIN_CHUNK_DURATION_MS == 0
+    }
+}
+
+pub(crate) struct RealtimeAsrUpload {
+    sender: Option<mpsc::Sender<Vec<u8>>>,
+    task: JoinHandle<Result<RealtimeAsrOutcome>>,
+    chunker: Pcm16kChunker,
+}
+
+pub(crate) struct RealtimeAsrOutcome {
+    pub(crate) status: String,
+    pub(crate) asr_result_ref: Option<String>,
+    pub(crate) provider_alias: Option<String>,
+    pub(crate) partial_count: u64,
+    pub(crate) fallback_reason: Option<String>,
+    pub(crate) fallback_stage: Option<String>,
+    pub(crate) chunk_count: u64,
+    pub(crate) audio_bytes: u64,
+    pub(crate) wall_ms: u64,
+}
+
+impl RealtimeAsrUpload {
+    pub(crate) fn start(
+        http: Client,
+        config: RealtimeAsrConfig,
+        call_id: &str,
+        trace_id: &str,
+        turn_id: &str,
+        initial_samples_48k: &[i16],
+    ) -> Result<Self> {
+        if !config.is_ready() {
+            return Err(anyhow!("realtime asr config is not ready"));
+        }
+        let nonce = config.runtime_session_nonce.as_deref().unwrap_or_default();
+        let session_line = session_start_line(call_id, trace_id, turn_id, nonce)?;
+        let (sender, receiver) = mpsc::channel(UPLOAD_QUEUE_CAPACITY);
+        let request_call_id = call_id.to_string();
+        let request_trace_id = trace_id.to_string();
+        let request_turn_id = turn_id.to_string();
+        let task = tokio::spawn(run_upload(
+            http,
+            config.clone(),
+            request_call_id,
+            request_trace_id,
+            request_turn_id,
+            receiver,
+        ));
+        let mut upload = Self {
+            sender: Some(sender),
+            task,
+            chunker: Pcm16kChunker::new(config.chunk_duration_ms),
+        };
+        if let Err(error) = upload
+            .try_send_line(session_line)
+            .and_then(|_| upload.push_48k_samples(initial_samples_48k))
+        {
+            upload.task.abort();
+            return Err(error);
+        }
+        Ok(upload)
+    }
+
+    pub(crate) fn push_48k_samples(&mut self, samples: &[i16]) -> Result<()> {
+        for chunk in self.chunker.push_48k(samples) {
+            let line = audio_chunk_line(chunk.seq, &chunk.samples)?;
+            self.try_send_line(line)?;
+        }
+        Ok(())
+    }
+
+    pub(crate) async fn finish(
+        mut self,
+        speech_duration_ms: u64,
+        end_reason: &str,
+    ) -> Result<RealtimeAsrOutcome> {
+        if let Some(chunk) = self.chunker.flush() {
+            let line = audio_chunk_line(chunk.seq, &chunk.samples)?;
+            self.try_send_line(line)?;
+        }
+        self.try_send_line(vad_speech_end_line(speech_duration_ms, end_reason)?)?;
+        self.try_send_line(finish_line()?)?;
+        self.sender.take();
+
+        let chunk_count = self.chunker.chunk_count;
+        let audio_bytes = self.chunker.audio_bytes;
+        let mut task = self.task;
+        let mut outcome = match timeout(FINISH_RESPONSE_TIMEOUT, &mut task).await {
+            Ok(joined) => joined.context("realtime asr upload task failed")??,
+            Err(_) => {
+                task.abort();
+                return Err(anyhow!("realtime asr finish response timeout"));
+            }
+        };
+        outcome.chunk_count = chunk_count;
+        outcome.audio_bytes = audio_bytes;
+        Ok(outcome)
+    }
+
+    pub(crate) async fn cancel(mut self, reason: &str) {
+        if let Ok(line) = cancel_line(reason) {
+            let _ = self.try_send_line(line);
+        }
+        self.sender.take();
+        if timeout(CANCEL_RESPONSE_TIMEOUT, &mut self.task)
+            .await
+            .is_err()
+        {
+            self.task.abort();
+        }
+    }
+
+    fn try_send_line(&self, line: Vec<u8>) -> Result<()> {
+        self.sender
+            .as_ref()
+            .ok_or_else(|| anyhow!("realtime asr upload is closed"))?
+            .try_send(line)
+            .map_err(|error| anyhow!("realtime asr upload backpressure: {error}"))
+    }
+}
+
+async fn run_upload(
+    http: Client,
+    config: RealtimeAsrConfig,
+    call_id: String,
+    trace_id: String,
+    turn_id: String,
+    receiver: mpsc::Receiver<Vec<u8>>,
+) -> Result<RealtimeAsrOutcome> {
+    let started_at = Instant::now();
+    let body_stream = stream::unfold(receiver, |mut receiver| async move {
+        receiver
+            .recv()
+            .await
+            .map(|chunk| (Ok::<Vec<u8>, io::Error>(chunk), receiver))
+    });
+    let response = http
+        .post(config.url.as_deref().unwrap_or_default())
+        .header("Content-Type", "application/x-ndjson")
+        .header(
+            "X-CV-Runtime-Token",
+            config.runtime_token.as_deref().unwrap_or_default(),
+        )
+        .header("X-CV-Call-Id", &call_id)
+        .header("X-CV-Trace-Id", &trace_id)
+        .header(
+            "X-CV-Runtime-Session-Nonce",
+            config.runtime_session_nonce.as_deref().unwrap_or_default(),
+        )
+        .body(reqwest::Body::wrap_stream(body_stream))
+        .send()
+        .await
+        .context("failed to post realtime asr stream")?;
+    let status = response.status();
+    if !status.is_success() {
+        return Err(anyhow!(
+            "realtime asr http failed status={} turn={}",
+            status.as_u16(),
+            turn_id
+        ));
+    }
+    let body: RuntimeCommonResult<RuntimeRealtimeAsrResp> = response
+        .json()
+        .await
+        .context("failed to decode realtime asr response")?;
+    if body.code != 0 {
+        return Err(anyhow!(
+            "realtime asr common result failed code={}",
+            body.code
+        ));
+    }
+    let data = body
+        .data
+        .ok_or_else(|| anyhow!("realtime asr response data missing"))?;
+    Ok(RealtimeAsrOutcome {
+        status: data.status.unwrap_or_else(|| "unknown".to_string()),
+        asr_result_ref: data.asr_result_ref,
+        provider_alias: data.provider_alias,
+        partial_count: data.partial_count.unwrap_or_default(),
+        fallback_reason: data.fallback_reason,
+        fallback_stage: data.fallback_stage,
+        chunk_count: 0,
+        audio_bytes: 0,
+        wall_ms: started_at.elapsed().as_millis() as u64,
+    })
+}
+
+struct Pcm16kChunker {
+    pending: Vec<i16>,
+    samples_per_chunk: usize,
+    next_seq: u64,
+    chunk_count: u64,
+    audio_bytes: u64,
+}
+
+impl Pcm16kChunker {
+    fn new(chunk_duration_ms: u64) -> Self {
+        Self {
+            pending: Vec::new(),
+            samples_per_chunk: SAMPLE_RATE_16K as usize * chunk_duration_ms as usize / 1000,
+            next_seq: 0,
+            chunk_count: 0,
+            audio_bytes: 0,
+        }
+    }
+
+    fn push_48k(&mut self, samples: &[i16]) -> Vec<PcmChunk> {
+        self.pending.extend(samples.iter().step_by(3).copied());
+        let mut chunks = Vec::new();
+        while self.pending.len() >= self.samples_per_chunk {
+            let samples = self.pending.drain(..self.samples_per_chunk).collect();
+            chunks.push(self.record_chunk(samples));
+        }
+        chunks
+    }
+
+    fn flush(&mut self) -> Option<PcmChunk> {
+        if self.pending.is_empty() {
+            return None;
+        }
+        let samples = std::mem::take(&mut self.pending);
+        Some(self.record_chunk(samples))
+    }
+
+    fn record_chunk(&mut self, samples: Vec<i16>) -> PcmChunk {
+        self.next_seq = self.next_seq.saturating_add(1);
+        self.chunk_count = self.chunk_count.saturating_add(1);
+        self.audio_bytes = self
+            .audio_bytes
+            .saturating_add((samples.len() * size_of::<i16>()) as u64);
+        PcmChunk {
+            seq: self.next_seq,
+            samples,
+        }
+    }
+}
+
+struct PcmChunk {
+    seq: u64,
+    samples: Vec<i16>,
+}
+
+fn session_start_line(
+    call_id: &str,
+    trace_id: &str,
+    turn_id: &str,
+    runtime_session_nonce: &str,
+) -> Result<Vec<u8>> {
+    encode_line(json!({
+        "event": "session_start",
+        "callId": call_id,
+        "traceId": trace_id,
+        "turnId": turn_id,
+        "runtimeSessionNonceHash": sha12(runtime_session_nonce),
+        "audio": {
+            "format": "pcm_s16le",
+            "sampleRate": SAMPLE_RATE_16K,
+            "channels": CHANNELS_MONO,
+        }
+    }))
+}
+
+fn audio_chunk_line(chunk_seq: u64, samples: &[i16]) -> Result<Vec<u8>> {
+    let mut bytes = Vec::with_capacity(samples.len() * size_of::<i16>());
+    for sample in samples {
+        bytes.extend_from_slice(&sample.to_le_bytes());
+    }
+    encode_line(json!({
+        "event": "audio_chunk",
+        "chunkSeq": chunk_seq,
+        "audioBase64": general_purpose::STANDARD.encode(bytes),
+        "durationMs": ((samples.len() as u64) * 1000 / u64::from(SAMPLE_RATE_16K)).max(1),
+    }))
+}
+
+fn vad_speech_end_line(speech_duration_ms: u64, end_reason: &str) -> Result<Vec<u8>> {
+    encode_line(json!({
+        "event": "vad_speech_end",
+        "speechDurationMs": speech_duration_ms,
+        "endReason": end_reason,
+    }))
+}
+
+fn finish_line() -> Result<Vec<u8>> {
+    encode_line(json!({"event": "finish"}))
+}
+
+fn cancel_line(reason: &str) -> Result<Vec<u8>> {
+    encode_line(json!({"event": "cancel", "reason": reason}))
+}
+
+fn encode_line(value: serde_json::Value) -> Result<Vec<u8>> {
+    let mut line = serde_json::to_vec(&value)?;
+    line.push(b'\n');
+    if line.len() > MAX_NDJSON_LINE_BYTES {
+        return Err(anyhow!("realtime asr ndjson line exceeds 24KiB"));
+    }
+    Ok(line)
+}
+
+fn sha12(value: &str) -> String {
+    let digest = Sha256::digest(value.as_bytes());
+    format!("{digest:x}")[..12].to_string()
+}
+
+fn non_blank(value: &Option<String>) -> bool {
+    value.as_ref().is_some_and(|value| !value.trim().is_empty())
+}
+
+#[derive(Deserialize)]
+struct RuntimeCommonResult<T> {
+    code: i64,
+    data: Option<T>,
+}
+
+#[derive(Deserialize)]
+#[serde(rename_all = "camelCase")]
+struct RuntimeRealtimeAsrResp {
+    status: Option<String>,
+    asr_result_ref: Option<String>,
+    provider_alias: Option<String>,
+    partial_count: Option<u64>,
+    fallback_reason: Option<String>,
+    fallback_stage: Option<String>,
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use std::{
+        io::{Read, Write},
+        net::TcpListener,
+        sync::mpsc as std_mpsc,
+        thread,
+    };
+
+    #[test]
+    fn chunker_aggregates_20ms_frames_into_200ms_chunks_and_flushes_tail() {
+        let mut chunker = Pcm16kChunker::new(200);
+        let first = chunker.push_48k(&vec![7; 9_600]);
+        let second = chunker.push_48k(&vec![9; 4_800]);
+        let tail = chunker.flush();
+
+        assert_eq!(1, first.len());
+        assert_eq!(1, first[0].seq);
+        assert_eq!(3_200, first[0].samples.len());
+        assert!(second.is_empty());
+        assert_eq!(2, tail.as_ref().map(|chunk| chunk.seq).unwrap_or_default());
+        assert_eq!(
+            1_600,
+            tail.as_ref()
+                .map(|chunk| chunk.samples.len())
+                .unwrap_or_default()
+        );
+        assert_eq!(2, chunker.chunk_count);
+        assert_eq!(9_600, chunker.audio_bytes);
+    }
+
+    #[test]
+    fn session_start_uses_canonical_nonce_hash_and_audio_contract() {
+        let line = session_start_line("call-001", "trace-001", "turn-0001", "nonce-001")
+            .expect("session start line");
+        let value: serde_json::Value = serde_json::from_slice(&line).expect("valid json");
+
+        assert_eq!("session_start", value["event"]);
+        assert_eq!(sha12("nonce-001"), value["runtimeSessionNonceHash"]);
+        assert_eq!("pcm_s16le", value["audio"]["format"]);
+        assert_eq!(16000, value["audio"]["sampleRate"]);
+        assert_eq!(1, value["audio"]["channels"]);
+    }
+
+    #[test]
+    fn maximum_audio_chunk_stays_within_java_line_limit() {
+        let line = audio_chunk_line(1, &vec![0; 8_000]).expect("maximum chunk line");
+        assert!(line.len() <= MAX_NDJSON_LINE_BYTES);
+    }
+
+    #[tokio::test]
+    async fn upload_streams_canonical_ndjson_and_reads_final_response() {
+        let response = json!({
+            "code": 0,
+            "data": {
+                "status": "final",
+                "asrResultRef": "asr_rt_fixture",
+                "providerAlias": "fixture",
+                "partialCount": 2
+            }
+        })
+        .to_string();
+        let (url, captured, server) = spawn_http_fixture(response);
+        let mut upload = RealtimeAsrUpload::start(
+            Client::new(),
+            fixture_config(url),
+            "call-001",
+            "trace-001",
+            "turn-0001",
+            &vec![1; 9_600],
+        )
+        .expect("start upload");
+        upload
+            .push_48k_samples(&vec![2; 9_600])
+            .expect("push audio");
+
+        let outcome = upload.finish(400, "silence").await.expect("finish upload");
+        let request = captured.recv().expect("captured request");
+        server.join().expect("fixture server");
+
+        assert_eq!("final", outcome.status);
+        assert_eq!(Some("asr_rt_fixture"), outcome.asr_result_ref.as_deref());
+        assert_eq!(2, outcome.chunk_count);
+        assert_eq!(12_800, outcome.audio_bytes);
+        assert!(
+            request
+                .headers
+                .contains("x-cv-runtime-session-nonce: nonce-001")
+        );
+        let events = request
+            .body
+            .lines()
+            .map(|line| serde_json::from_str::<serde_json::Value>(line).expect("event json"))
+            .collect::<Vec<_>>();
+        assert_eq!(5, events.len());
+        assert_eq!("session_start", events[0]["event"]);
+        assert_eq!(1, events[1]["chunkSeq"]);
+        assert_eq!(2, events[2]["chunkSeq"]);
+        assert_eq!("vad_speech_end", events[3]["event"]);
+        assert_eq!("finish", events[4]["event"]);
+    }
+
+    #[tokio::test]
+    async fn cancel_ends_stream_without_finish_event() {
+        let response = json!({"code": 0, "data": {"status": "cancelled"}}).to_string();
+        let (url, captured, server) = spawn_http_fixture(response);
+        let upload = RealtimeAsrUpload::start(
+            Client::new(),
+            fixture_config(url),
+            "call-002",
+            "trace-002",
+            "turn-0002",
+            &vec![1; 9_600],
+        )
+        .expect("start upload");
+
+        upload.cancel("call_end").await;
+        let request = captured.recv().expect("captured request");
+        server.join().expect("fixture server");
+
+        assert!(request.body.contains("\"event\":\"cancel\""));
+        assert!(!request.body.contains("\"event\":\"finish\""));
+    }
+
+    fn fixture_config(url: String) -> RealtimeAsrConfig {
+        RealtimeAsrConfig {
+            enabled: true,
+            url: Some(url),
+            runtime_token: Some("token-001".to_string()),
+            runtime_session_nonce: Some("nonce-001".to_string()),
+            chunk_duration_ms: 200,
+        }
+    }
+
+    struct CapturedRequest {
+        headers: String,
+        body: String,
+    }
+
+    fn spawn_http_fixture(
+        response_body: String,
+    ) -> (
+        String,
+        std_mpsc::Receiver<CapturedRequest>,
+        thread::JoinHandle<()>,
+    ) {
+        let listener = TcpListener::bind("127.0.0.1:0").expect("bind fixture server");
+        let address = listener.local_addr().expect("fixture address");
+        let (sender, receiver) = std_mpsc::channel();
+        let server = thread::spawn(move || {
+            let (mut stream, _) = listener.accept().expect("accept request");
+            stream
+                .set_read_timeout(Some(Duration::from_secs(3)))
+                .expect("read timeout");
+            let mut request = Vec::new();
+            let mut buffer = [0u8; 4096];
+            loop {
+                let read = stream.read(&mut buffer).expect("read request");
+                if read == 0 {
+                    break;
+                }
+                request.extend_from_slice(&buffer[..read]);
+                if request.windows(5).any(|window| window == b"0\r\n\r\n") {
+                    break;
+                }
+            }
+            let header_end = request
+                .windows(4)
+                .position(|window| window == b"\r\n\r\n")
+                .expect("request headers");
+            let headers = String::from_utf8_lossy(&request[..header_end]).to_ascii_lowercase();
+            let body = decode_chunked_body(&request[header_end + 4..]);
+            sender
+                .send(CapturedRequest {
+                    headers,
+                    body: String::from_utf8(body).expect("utf8 ndjson"),
+                })
+                .expect("capture request");
+            let response = format!(
+                "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
+                response_body.len(),
+                response_body
+            );
+            stream
+                .write_all(response.as_bytes())
+                .expect("write response");
+        });
+        (
+            format!("http://{address}/runtime/asr/realtime"),
+            receiver,
+            server,
+        )
+    }
+
+    fn decode_chunked_body(mut input: &[u8]) -> Vec<u8> {
+        let mut body = Vec::new();
+        loop {
+            let line_end = input
+                .windows(2)
+                .position(|window| window == b"\r\n")
+                .expect("chunk size line");
+            let size = usize::from_str_radix(
+                std::str::from_utf8(&input[..line_end]).expect("chunk size utf8"),
+                16,
+            )
+            .expect("chunk size");
+            input = &input[line_end + 2..];
+            if size == 0 {
+                break;
+            }
+            body.extend_from_slice(&input[..size]);
+            input = &input[size + 2..];
+        }
+        body
+    }
+}
diff --git a/src/audio.rs b/src/audio.rs
index 46613eb..cec22f0 100644
--- a/src/audio.rs
+++ b/src/audio.rs
@@ -1012,10 +1012,7 @@
         assert_eq!(5, frames.len());
         assert_eq!(
             4_800,
-            frames
-                .iter()
-                .map(|frame| frame.samples.len())
-                .sum::<usize>()
+            frames.iter().map(|frame| frame.data.len()).sum::<usize>()
         );
         assert!(frames.iter().all(|frame| frame.sample_rate == 48_000));
         assert!(frames.iter().all(|frame| frame.num_channels == 1));
@@ -1066,10 +1063,7 @@
         assert!(frames.iter().all(|frame| frame.samples_per_channel == 320));
         assert_eq!(
             1_600,
-            frames
-                .iter()
-                .map(|frame| frame.samples.len())
-                .sum::<usize>()
+            frames.iter().map(|frame| frame.data.len()).sum::<usize>()
         );
     }
 }
diff --git a/src/main.rs b/src/main.rs
index 1e969cc..62107a9 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -1,3 +1,4 @@
+mod asr_realtime;
 mod audio;
 mod service;
 
@@ -12,6 +13,7 @@
 };
 
 use anyhow::{Context, Result, anyhow};
+use asr_realtime::{RealtimeAsrConfig, RealtimeAsrOutcome, RealtimeAsrUpload};
 use audio::{AudioDiagnostics, load_pre_recorded_frames};
 use base64::{Engine as _, engine::general_purpose};
 use futures_util::StreamExt;
@@ -272,6 +274,9 @@
     runtime_turn_bridge_mode: String,
     runtime_asr_stream_enabled: bool,
     runtime_asr_stream_url: Option<String>,
+    runtime_asr_realtime_enabled: bool,
+    runtime_asr_realtime_url: Option<String>,
+    runtime_asr_realtime_chunk_duration_ms: u64,
     runtime_turn_artifact_dir: Option<String>,
     runtime_session_nonce: Option<String>,
     user_audio_observer_enabled: bool,
@@ -336,6 +341,9 @@
     bridge_mode: String,
     asr_stream_enabled: bool,
     asr_stream_url: Option<String>,
+    asr_realtime_enabled: bool,
+    asr_realtime_url: Option<String>,
+    asr_realtime_chunk_duration_ms: u64,
     artifact_dir: Option<String>,
     runtime_session_nonce: Option<String>,
     audio_debug_dump_dir: Option<String>,
@@ -349,6 +357,9 @@
             bridge_mode: config.runtime_turn_bridge_mode.clone(),
             asr_stream_enabled: config.runtime_asr_stream_enabled,
             asr_stream_url: config.runtime_asr_stream_url.clone(),
+            asr_realtime_enabled: config.runtime_asr_realtime_enabled,
+            asr_realtime_url: config.runtime_asr_realtime_url.clone(),
+            asr_realtime_chunk_duration_ms: config.runtime_asr_realtime_chunk_duration_ms,
             artifact_dir: config.runtime_turn_artifact_dir.clone(),
             runtime_session_nonce: config.runtime_session_nonce.clone(),
             audio_debug_dump_dir: config.audio_debug_dump_dir.clone(),
@@ -396,6 +407,16 @@
                 .as_ref()
                 .is_some_and(|value| !value.is_empty())
     }
+
+    fn realtime_asr_config(&self) -> RealtimeAsrConfig {
+        RealtimeAsrConfig {
+            enabled: self.asr_realtime_enabled,
+            url: self.asr_realtime_url.clone(),
+            runtime_token: self.bridge_token.clone(),
+            runtime_session_nonce: self.runtime_session_nonce.clone(),
+            chunk_duration_ms: self.asr_realtime_chunk_duration_ms,
+        }
+    }
 }
 
 impl Config {
@@ -419,6 +440,12 @@
                 .unwrap_or_else(|_| "json".to_string()),
             runtime_asr_stream_enabled: bool_env("CV_RUNTIME_ASR_STREAM_ENABLED", false),
             runtime_asr_stream_url: optional_env("CV_RUNTIME_ASR_STREAM_URL"),
+            runtime_asr_realtime_enabled: bool_env("CV_RUNTIME_ASR_REALTIME_ENABLED", false),
+            runtime_asr_realtime_url: optional_env("CV_RUNTIME_ASR_REALTIME_URL"),
+            runtime_asr_realtime_chunk_duration_ms: u64_env(
+                "CV_RUNTIME_ASR_REALTIME_CHUNK_DURATION_MS",
+                200,
+            ),
             runtime_turn_artifact_dir: optional_env("CV_RUNTIME_TURN_ARTIFACT_DIR"),
             runtime_session_nonce: optional_env("CV_RUNTIME_SESSION_NONCE"),
             user_audio_observer_enabled: bool_env("CV_ENABLE_USER_AUDIO_OBSERVER", true),
@@ -775,6 +802,7 @@
     call_id: &str,
     trace_id: &str,
     turn: FinishedSpeechTurn,
+    realtime_asr_result_ref: Option<String>,
 ) {
     let turn_pipeline_started_at = Instant::now();
     if !bridge_config.is_ready() {
@@ -834,8 +862,10 @@
                     "endReason": turn.end_reason.as_str(),
                 }),
             );
-            let asr_result_ref =
-                request_asr_result_ref(http, bridge_config, call_id, trace_id, &turn).await;
+            let asr_result_ref = match realtime_asr_result_ref {
+                Some(value) => Some(value),
+                None => request_asr_result_ref(http, bridge_config, call_id, trace_id, &turn).await,
+            };
             if bridge_config.is_stream_mode() {
                 match request_turn_bridge_stream(
                     http,
@@ -2446,6 +2476,7 @@
         } else {
             None
         };
+        let mut realtime_asr_upload: Option<RealtimeAsrUpload> = None;
 
         while let Some(frame) = stream.next().await {
             frame_count += 1;
@@ -2478,7 +2509,8 @@
 
             if let Some(vad) = simple_vad.as_mut() {
                 if vad_enabled_gate.load(Ordering::Acquire) {
-                    if let Some(turn) = vad.observe_frame(
+                    let was_in_speech = vad.in_speech;
+                    let turn = vad.observe_frame(
                         &call_id,
                         &trace_id,
                         &participant_alias,
@@ -2486,7 +2518,65 @@
                         frame_count,
                         elapsed_ms,
                         &frame,
-                    ) {
+                    );
+                    let is_in_speech = vad.in_speech;
+
+                    if !was_in_speech && is_in_speech {
+                        let turn_id = format!("turn-{:04}", vad.turn_index);
+                        match RealtimeAsrUpload::start(
+                            http.clone(),
+                            turn_bridge_config.realtime_asr_config(),
+                            &call_id,
+                            &trace_id,
+                            &turn_id,
+                            &vad.speech_samples,
+                        ) {
+                            Ok(upload) => {
+                                info!(
+                                    call_id = %call_id,
+                                    trace_id = %trace_id,
+                                    turn_id = %turn_id,
+                                    "runtime helper asr_realtime_session_started"
+                                );
+                                realtime_asr_upload = Some(upload);
+                            }
+                            Err(error) if turn_bridge_config.asr_realtime_enabled => {
+                                warn!(
+                                    call_id = %call_id,
+                                    trace_id = %trace_id,
+                                    turn_id = %turn_id,
+                                    error = %safe_error(&error.to_string()),
+                                    "runtime helper asr_realtime_start_failed_fallback"
+                                );
+                            }
+                            Err(_) => {}
+                        }
+                    } else if was_in_speech {
+                        let push_failed = realtime_asr_upload
+                            .as_mut()
+                            .and_then(|upload| upload.push_48k_samples(frame.data.as_ref()).err());
+                        if let Some(error) = push_failed {
+                            warn!(
+                                call_id = %call_id,
+                                trace_id = %trace_id,
+                                error = %safe_error(&error.to_string()),
+                                "runtime helper asr_realtime_upload_failed_fallback"
+                            );
+                            if let Some(upload) = realtime_asr_upload.take() {
+                                tokio::spawn(async move {
+                                    upload.cancel("upload_backpressure").await;
+                                });
+                            }
+                        }
+                    }
+
+                    if let Some(turn) = turn {
+                        let realtime_asr_result_ref = match realtime_asr_upload.take() {
+                            Some(upload) => {
+                                finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+                            }
+                            None => None,
+                        };
                         handle_finished_turn(
                             &http,
                             &turn_bridge_config,
@@ -2494,10 +2584,22 @@
                             &call_id,
                             &trace_id,
                             turn,
+                            realtime_asr_result_ref,
                         )
                         .await;
+                    } else if was_in_speech && !is_in_speech {
+                        if let Some(upload) = realtime_asr_upload.take() {
+                            tokio::spawn(async move {
+                                upload.cancel("speech_too_short").await;
+                            });
+                        }
                     }
                 } else {
+                    if let Some(upload) = realtime_asr_upload.take() {
+                        tokio::spawn(async move {
+                            upload.cancel("vad_disabled").await;
+                        });
+                    }
                     vad.observe_disabled_frame(
                         &call_id,
                         &trace_id,
@@ -2518,9 +2620,26 @@
                 &track_sid_alias,
                 started_at.elapsed().as_millis() as u64,
             ) {
-                handle_finished_turn(&http, &turn_bridge_config, &sink, &call_id, &trace_id, turn)
-                    .await;
+                let realtime_asr_result_ref = match realtime_asr_upload.take() {
+                    Some(upload) => {
+                        finish_realtime_asr_upload(upload, &call_id, &trace_id, &turn).await
+                    }
+                    None => None,
+                };
+                handle_finished_turn(
+                    &http,
+                    &turn_bridge_config,
+                    &sink,
+                    &call_id,
+                    &trace_id,
+                    turn,
+                    realtime_asr_result_ref,
+                )
+                .await;
             }
+        }
+        if let Some(upload) = realtime_asr_upload.take() {
+            upload.cancel("stream_end").await;
         }
 
         info!(
@@ -2536,6 +2655,58 @@
     })
 }
 
+async fn finish_realtime_asr_upload(
+    upload: RealtimeAsrUpload,
+    call_id: &str,
+    trace_id: &str,
+    turn: &FinishedSpeechTurn,
+) -> Option<String> {
+    match upload.finish(turn.duration_ms, &turn.end_reason).await {
+        Ok(RealtimeAsrOutcome {
+            status,
+            asr_result_ref,
+            provider_alias,
+            partial_count,
+            fallback_reason,
+            fallback_stage,
+            chunk_count,
+            audio_bytes,
+            wall_ms,
+        }) => {
+            info!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                status = %status,
+                provider_alias = ?provider_alias,
+                partial_count,
+                chunk_count,
+                audio_bytes,
+                wall_ms,
+                asr_result_ref_present = asr_result_ref.is_some(),
+                fallback_reason = ?fallback_reason,
+                fallback_stage = ?fallback_stage,
+                "runtime helper asr_realtime_finished"
+            );
+            if status == "final" {
+                asr_result_ref
+            } else {
+                None
+            }
+        }
+        Err(error) => {
+            warn!(
+                call_id = %call_id,
+                trace_id = %trace_id,
+                turn_id = %turn.turn_id,
+                error = %safe_error(&error.to_string()),
+                "runtime helper asr_realtime_failed_fallback"
+            );
+            None
+        }
+    }
+}
+
 #[derive(Clone)]
 struct SimpleVadConfig {
     rms_threshold: f64,
diff --git a/src/service.rs b/src/service.rs
index ed7ba42..bb64683 100644
--- a/src/service.rs
+++ b/src/service.rs
@@ -410,6 +410,18 @@
         if let Some(value) = &runtime.asr_stream_url {
             command.env("CV_RUNTIME_ASR_STREAM_URL", value);
         }
+        if let Some(value) = runtime.asr_realtime_enabled {
+            command.env("CV_RUNTIME_ASR_REALTIME_ENABLED", value.to_string());
+        }
+        if let Some(value) = &runtime.asr_realtime_url {
+            command.env("CV_RUNTIME_ASR_REALTIME_URL", value);
+        }
+        if let Some(value) = runtime.asr_realtime_chunk_duration_ms {
+            command.env(
+                "CV_RUNTIME_ASR_REALTIME_CHUNK_DURATION_MS",
+                value.to_string(),
+            );
+        }
         if let Some(true) = runtime.audio_debug_dump_enabled {
             if let Some(value) = &config.audio_debug_dump_dir {
                 command.env("CV_AUDIO_DEBUG_DUMP_DIR", value);
@@ -976,6 +988,9 @@
     audio_debug_dump_enabled: Option<bool>,
     asr_streaming_enabled: Option<bool>,
     asr_stream_url: Option<String>,
+    asr_realtime_enabled: Option<bool>,
+    asr_realtime_url: Option<String>,
+    asr_realtime_chunk_duration_ms: Option<u64>,
 }
 
 #[derive(Default, Deserialize)]
diff --git a/tools/validate-asr-realtime-fixture.mjs b/tools/validate-asr-realtime-fixture.mjs
new file mode 100755
index 0000000..64414b1
--- /dev/null
+++ b/tools/validate-asr-realtime-fixture.mjs
@@ -0,0 +1,116 @@
+#!/usr/bin/env node
+
+import fs from 'node:fs'
+import path from 'node:path'
+
+const file = process.argv[2] || 'fixtures/asr-realtime-happy.ndjson'
+const absolute = path.resolve(file)
+const raw = fs.readFileSync(absolute, 'utf8')
+const lines = raw.split(/\r?\n/).map((line) => line.trim()).filter(Boolean)
+const allowedEvents = new Set(['session_start', 'audio_chunk', 'vad_speech_end', 'finish', 'cancel'])
+const forbiddenActivityEvents = new Set([
+  'asr_realtime_session_started',
+  'asr_realtime_audio_chunk',
+  'asr_realtime_speech_end',
+  'asr_realtime_finish_received',
+  'asr_realtime_cancelled',
+])
+
+let state = 'new'
+let lastChunkSeq = 0
+let audioChunkCount = 0
+let audioBytes = 0
+
+function fail(message) {
+  console.error(`asr realtime fixture invalid: ${message}`)
+  process.exit(1)
+}
+
+for (const [index, line] of lines.entries()) {
+  if (Buffer.byteLength(`${line}\n`, 'utf8') > 24 * 1024) {
+    fail(`line ${index + 1} exceeds 24KiB`)
+  }
+  let event
+  try {
+    event = JSON.parse(line)
+  } catch {
+    fail(`line ${index + 1} is not JSON`)
+  }
+  if (forbiddenActivityEvents.has(event.event)) {
+    fail(`line ${index + 1} uses activity event ${event.event} as a wire event`)
+  }
+  if (!allowedEvents.has(event.event)) {
+    fail(`line ${index + 1} has unsupported event ${event.event}`)
+  }
+
+  if (state === 'new') {
+    if (event.event !== 'session_start') {
+      fail('first event must be session_start')
+    }
+    if (!event.callId || !event.traceId || !event.turnId || !event.runtimeSessionNonceHash) {
+      fail('session_start missing binding fields')
+    }
+    if (event.audio?.format !== 'pcm_s16le' || event.audio?.sampleRate !== 16000 || event.audio?.channels !== 1) {
+      fail('session_start audio must be pcm_s16le/16000Hz/mono')
+    }
+    state = 'streaming'
+    continue
+  }
+
+  if (state === 'terminal') {
+    fail(`event ${event.event} arrived after terminal event`)
+  }
+  if (event.event === 'audio_chunk') {
+    if (state !== 'streaming') {
+      fail('audio_chunk arrived after vad_speech_end')
+    }
+    if (!Number.isInteger(event.chunkSeq) || event.chunkSeq !== lastChunkSeq + 1) {
+      fail('audio_chunk chunkSeq must start at 1 and be contiguous')
+    }
+    if (!Number.isInteger(event.durationMs) || event.durationMs < 1 || event.durationMs > 1000) {
+      fail('audio_chunk durationMs must be 1..1000')
+    }
+    const payload = Buffer.from(event.audioBase64 || '', 'base64')
+    if (payload.length === 0 || payload.length > 16000 || payload.length % 2 !== 0) {
+      fail('audio_chunk payload must be non-empty, <=16000 bytes and 16-bit aligned')
+    }
+    lastChunkSeq = event.chunkSeq
+    audioChunkCount += 1
+    audioBytes += payload.length
+    continue
+  }
+  if (event.event === 'vad_speech_end') {
+    if (state !== 'streaming' || audioChunkCount === 0) {
+      fail('vad_speech_end requires at least one audio_chunk')
+    }
+    state = 'speech_ended'
+    continue
+  }
+  if (event.event === 'finish') {
+    if (state !== 'speech_ended') {
+      fail('finish must follow vad_speech_end')
+    }
+    state = 'terminal'
+    continue
+  }
+  if (event.event === 'cancel') {
+    if (!event.reason) {
+      fail('cancel missing reason')
+    }
+    state = 'terminal'
+    continue
+  }
+  fail(`duplicate session_start at line ${index + 1}`)
+}
+
+if (state !== 'terminal') {
+  fail('fixture must end with finish or cancel')
+}
+
+console.log(JSON.stringify({
+  ok: true,
+  file: path.relative(process.cwd(), absolute),
+  eventCount: lines.length,
+  audioChunkCount,
+  audioBytes,
+}))

--
Gitblit v1.9.3