use super::a2dp::{A2dpError, AospA2dpSink, TARGET_RATE};
use super::edge_tts::TtsEvent;
use super::pipeline;
use std::time::{Duration, Instant};
use thiserror::Error;
const EDGE_TICKS_PER_SECOND: f64 = 10_000_000.0;
const BYTES_PER_STEREO_FRAME: u64 = 4;
const DRAIN_SAFETY_PAD: f64 = 0.5;
#[derive(Debug, Clone)]
pub struct WordMark {
pub time_s: f64,
pub text: String,
}
#[derive(Debug, Clone, Default)]
pub struct Prepared {
pub stereo: Vec<i16>,
pub bounds: Vec<(u64, String)>,
}
pub const CHUNK_FRAMES: usize = 2_205;
pub const LEAD_FRAMES: u64 = 8_820;
pub const LEAD_IN_FRAMES: usize = TARGET_RATE / 2;
pub const SENTENCE_GAP_FRAMES: usize = TARGET_RATE / 1000 * 400;
pub const PARA_GAP_FRAMES: usize = TARGET_RATE / 1000 * 700;
#[derive(Debug, Error)]
pub enum PlayerError {
#[error("a2dp: {0}")]
A2dp(#[from] A2dpError),
#[error("pipeline: {0}")]
Pipeline(String),
}
pub struct Player {
sink: AospA2dpSink,
started_at: Option<Instant>,
frames_written: u64,
clock_base: u64,
clock_anchor: Option<Instant>,
paused_since: Option<Instant>,
chunk_buf: Vec<u8>,
silence_buf: Vec<u8>,
}
impl Player {
pub async fn open() -> Result<Self, PlayerError> {
let mut sink = AospA2dpSink::open().await?;
sink.start().await?;
Ok(Player {
sink,
started_at: None,
frames_written: 0,
clock_base: 0,
clock_anchor: None,
paused_since: None,
chunk_buf: vec![0u8; CHUNK_FRAMES * BYTES_PER_STEREO_FRAME as usize],
silence_buf: vec![0u8; CHUNK_FRAMES * BYTES_PER_STEREO_FRAME as usize],
})
}
pub fn prepare(events: &[TtsEvent]) -> Result<Prepared, PlayerError> {
let mut mp3: Vec<u8> = Vec::new();
let mut bounds: Vec<(u64, String)> = Vec::new();
for ev in events {
match ev {
TtsEvent::Audio(b) => mp3.extend_from_slice(b),
TtsEvent::WordBoundary { offset, text, .. } => bounds.push((*offset, text.clone())),
TtsEvent::TurnEnd => {}
}
}
if mp3.is_empty() {
return Ok(Prepared {
stereo: Vec::new(),
bounds,
});
}
let mono = pipeline::decode_mp3(&mp3).map_err(PlayerError::Pipeline)?;
let resampled = pipeline::resample_mono(&mono).map_err(PlayerError::Pipeline)?;
let stereo = pipeline::mono_to_stereo(&resampled);
Ok(Prepared { stereo, bounds })
}
pub fn next_utt_start_s(&self) -> f64 {
self.frames_written as f64 / TARGET_RATE as f64
}
pub fn marks(prepared: &Prepared, utt_start_s: f64) -> Vec<WordMark> {
prepared
.bounds
.iter()
.map(|(ticks, text)| WordMark {
time_s: utt_start_s + *ticks as f64 / EDGE_TICKS_PER_SECOND,
text: text.clone(),
})
.collect()
}
pub async fn write_chunk(&mut self, stereo: &[i16]) -> Result<(), PlayerError> {
if self.started_at.is_none() {
self.started_at = Some(Instant::now());
}
let n = stereo.len() * 2;
self.chunk_buf.resize(n, 0);
for (i, s) in stereo.iter().enumerate() {
self.chunk_buf[i * 2] = (*s & 0xFF) as u8;
self.chunk_buf[i * 2 + 1] = (*s >> 8) as u8;
}
self.sink.write_pcm(&self.chunk_buf[..n]).await?;
self.frames_written += (stereo.len() / 2) as u64;
Ok(())
}
pub async fn keepalive(&mut self) -> Result<(), PlayerError> {
self.silence_buf.fill(0);
self.sink.write_pcm(&self.silence_buf).await?;
Ok(())
}
pub fn playback_frames(&self) -> u64 {
let span = match self.clock_anchor {
Some(a) => (TARGET_RATE as f64 * a.elapsed().as_secs_f64()) as u64,
None => 0,
};
self.clock_base + span
}
pub fn lead_frames(&self) -> u64 {
self.frames_written.saturating_sub(self.playback_frames())
}
pub fn socket_buffered_frames(&self) -> u64 {
self.sink.unsent_bytes() as u64 / 4 }
pub fn frames_written(&self) -> u64 {
self.frames_written
}
pub fn resume_clock(&mut self) {
if self.clock_anchor.is_none() {
self.clock_anchor = Some(Instant::now());
}
self.paused_since = None;
}
pub fn pause_clock(&mut self) {
self.clock_base = self.frames_written;
self.clock_anchor = None;
self.paused_since = Some(Instant::now());
}
pub fn paused_since(&self) -> Option<Duration> {
self.paused_since.map(|t| t.elapsed())
}
pub async fn play_utterance(
&mut self,
events: &[TtsEvent],
) -> Result<Vec<WordMark>, PlayerError> {
let prepared = Self::prepare(events)?;
if prepared.stereo.is_empty() {
return Ok(Vec::new());
}
if self.started_at.is_none() {
self.started_at = Some(Instant::now());
}
let utt_start_s = self.frames_written as f64 / TARGET_RATE as f64;
let bytes: Vec<u8> = prepared
.stereo
.iter()
.flat_map(|s| s.to_le_bytes())
.collect();
self.sink.write_pcm(&bytes).await?;
self.frames_written += (prepared.stereo.len() / 2) as u64;
Ok(Self::marks(&prepared, utt_start_s))
}
fn buffered_frames(&self) -> u64 {
self.sink.unsent_bytes() as u64 / BYTES_PER_STEREO_FRAME
}
pub async fn drain_and_stop(mut self) -> Result<(), PlayerError> {
let drain_secs = self.buffered_frames() as f64 / TARGET_RATE as f64 + DRAIN_SAFETY_PAD;
if drain_secs > 0.0 {
tokio::time::sleep(Duration::from_secs_f64(drain_secs)).await;
}
self.sink.stop().await?;
Ok(())
}
}