use std::time::Duration;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use crate::app::events::AppEvent;
use crate::entities::message::{Message, MessageRole};
use crate::features::tts_command::TtsScope;
use crate::shared::config::SecretSlot;
use crate::shared::i18n::Locale;
use crate::shared::tts::{TtsEngine, TtsSetupError, engines_from_config, playback::Playback};
use super::Orchestrator;
const QUEUE_AHEAD: usize = 2;
const POLL_INTERVAL: Duration = Duration::from_millis(80);
const DRAIN_TAIL: Duration = Duration::from_millis(300);
impl Orchestrator {
pub(super) fn handle_tts(&mut self, scope: TtsScope) {
self.stop_tts();
let Some((profile_id, messages)) = self.active_id.and_then(|id| match self.view(id)? {
super::ChatView::Top(chat) => Some((chat.profile_id, chat.messages.clone())),
super::ChatView::Child { parent, run } => {
Some((parent.profile_id, run.messages.clone()))
}
}) else {
self.fail_tts(self.ui_locale().t("ui.err.tts_no_active_chat"));
return;
};
let speech_loc = self.profile_locale(profile_id);
let chunks =
match build_utterances(&messages, scope, self.config.tts.speak_roles, speech_loc) {
Some(text) => text,
None => {
self.fail_tts(self.ui_locale().t("ui.err.tts_nothing_to_speak"));
return;
}
};
let stored = self.config.tts.secret_key().and_then(|k| {
crate::shared::secrets::stored_key(&self.config.api_keys, &k.storage_name())
});
let (engine, user_engine) = match engines_from_config(&self.config.tts, stored) {
Ok(pair) => pair,
Err(err) => {
self.fail_tts(self.ui_locale().t(setup_error_key(err)));
return;
}
};
let chunks = chunk_utterances(&chunks, engine.max_input_chars());
if chunks.is_empty() {
self.fail_tts(self.ui_locale().t("ui.err.tts_nothing_to_speak"));
return;
}
let playback = match Playback::open() {
Ok(p) => std::sync::Arc::new(p),
Err(err) => {
self.fail_tts(
&self
.ui_locale()
.tf("ui.err.tts_no_audio", &[("err", &err.to_string())]),
);
return;
}
};
let cancel = CancellationToken::new();
let task_id = Uuid::new_v4();
self.tts_cancel = Some(cancel.clone());
self.tts_gen = Some(task_id);
self.tts_playback = Some(playback.clone());
let _ = self.evt_tx.send(AppEvent::TtsActive(true));
spawn_tts(TtsTask {
engine,
user_engine,
chunks,
cancel,
task_id,
playback,
loc: self.ui_locale(),
evt_tx: self.evt_tx.clone(),
done_tx: self.tts_done_tx.clone(),
});
}
pub(super) fn handle_tts_pause(&mut self) {
if let Some(pb) = &self.tts_playback {
pb.pause();
}
}
pub(super) fn handle_tts_resume(&mut self) {
if let Some(pb) = &self.tts_playback {
pb.resume();
}
}
pub(super) fn stop_tts(&mut self) {
if let Some(token) = self.tts_cancel.take() {
token.cancel();
}
if let Some(pb) = self.tts_playback.take() {
pb.resume();
}
if self.tts_gen.take().is_some() {
let _ = self.evt_tx.send(AppEvent::TtsActive(false));
}
}
pub(super) fn handle_tts_done(&mut self, task_id: Uuid) {
if self.tts_gen == Some(task_id) {
self.tts_gen = None;
self.tts_cancel = None;
self.tts_playback = None;
let _ = self.evt_tx.send(AppEvent::TtsActive(false));
}
}
fn fail_tts(&self, msg: &str) {
let _ = self.evt_tx.send(AppEvent::Error(msg.to_string()));
}
}
fn setup_error_key(err: TtsSetupError) -> &'static str {
match err {
TtsSetupError::Model => "ui.err.tts_no_model",
TtsSetupError::ApiKey => "ui.err.tts_no_api_key",
TtsSetupError::Url => "ui.err.tts_no_url",
}
}
pub(super) fn build_utterances(
messages: &[Message],
scope: TtsScope,
speak_roles: bool,
loc: &'static Locale,
) -> Option<Vec<(MessageRole, String)>> {
let spoken: Vec<&Message> = messages
.iter()
.filter(|m| matches!(m.role, MessageRole::User | MessageRole::Assistant))
.filter(|m| !m.text.trim().is_empty())
.collect();
let take = match scope {
TtsScope::Last => 1,
TtsScope::Recent(n) => n,
TtsScope::All => spoken.len(),
};
let start = spoken.len().saturating_sub(take);
let out: Vec<(MessageRole, String)> = spoken[start..]
.iter()
.filter_map(|m| {
let text = crate::shared::markdown::speakable_text(&m.text, loc);
if text.is_empty() {
return None;
}
let text = if speak_roles {
let role = match m.role {
MessageRole::User => loc.t("speak.role.user"),
_ => loc.t("speak.role.assistant"),
};
format!("{role} {text}")
} else {
text
};
Some((m.role, text))
})
.collect();
(!out.is_empty()).then_some(out)
}
pub(super) fn chunk_utterances(
utterances: &[(MessageRole, String)],
max_chars: usize,
) -> Vec<(MessageRole, String)> {
let max = max_chars.max(1);
let mut out = Vec::new();
for (role, utterance) in utterances {
for block in utterance.lines() {
let block = block.trim();
if block.is_empty() {
continue;
}
let mut chunks = Vec::new();
pack_sentences(block, max, &mut chunks);
out.extend(chunks.into_iter().map(|c| (*role, c)));
}
}
out
}
fn pack_sentences(block: &str, max: usize, out: &mut Vec<String>) {
let mut cur = String::new();
for sentence in crate::features::tools::rag::split_sentences(block) {
append_pieces(&sentence, max, &mut cur, out);
}
if !cur.is_empty() {
out.push(cur);
}
}
fn append_pieces(sentence: &str, max: usize, cur: &mut String, out: &mut Vec<String>) {
for piece in split_long(sentence, max) {
let piece = piece.trim();
if piece.is_empty() {
continue;
}
let extra = if cur.is_empty() { 0 } else { 1 };
if !cur.is_empty() && cur.chars().count() + extra + piece.chars().count() > max {
out.push(std::mem::take(cur));
}
if !cur.is_empty() {
cur.push(' ');
}
cur.push_str(piece);
}
}
fn split_long(sentence: &str, max: usize) -> Vec<String> {
if sentence.chars().count() <= max {
return vec![sentence.to_string()];
}
let mut out = Vec::new();
let mut cur = String::new();
for word in sentence.split_whitespace() {
let wlen = word.chars().count();
if wlen > max {
split_giant_word(word, max, &mut cur, &mut out);
continue;
}
let extra = if cur.is_empty() { 0 } else { 1 };
if cur.chars().count() + extra + wlen > max {
out.push(std::mem::take(&mut cur));
}
if !cur.is_empty() {
cur.push(' ');
}
cur.push_str(word);
}
if !cur.is_empty() {
out.push(cur);
}
out
}
fn split_giant_word(word: &str, max: usize, cur: &mut String, out: &mut Vec<String>) {
if !cur.is_empty() {
out.push(std::mem::take(cur));
}
let chars: Vec<char> = word.chars().collect();
for part in chars.chunks(max) {
out.push(part.iter().collect());
}
}
struct TtsTask {
engine: Box<dyn TtsEngine>,
user_engine: Option<Box<dyn TtsEngine>>,
chunks: Vec<(MessageRole, String)>,
cancel: CancellationToken,
task_id: Uuid,
playback: std::sync::Arc<Playback>,
loc: &'static Locale,
evt_tx: tokio::sync::mpsc::UnboundedSender<AppEvent>,
done_tx: tokio::sync::mpsc::UnboundedSender<Uuid>,
}
fn spawn_tts(task: TtsTask) {
tokio::spawn(async move {
for (role, chunk) in &task.chunks {
if task.cancel.is_cancelled() {
break;
}
if !synth_chunk(&task, *role, chunk).await {
break;
}
}
while !task.playback.is_drained() && !task.cancel.is_cancelled() {
tokio::time::sleep(POLL_INTERVAL).await;
}
if task.cancel.is_cancelled() {
task.playback.stop();
} else {
tokio::time::sleep(DRAIN_TAIL).await;
}
let _ = task.done_tx.send(task.task_id);
});
}
async fn synth_chunk(task: &TtsTask, role: MessageRole, chunk: &str) -> bool {
let fail = |msg: String| {
let _ = task.evt_tx.send(AppEvent::Error(msg));
};
while task.playback.queued() >= QUEUE_AHEAD && !task.cancel.is_cancelled() {
tokio::time::sleep(POLL_INTERVAL).await;
}
if task.cancel.is_cancelled() {
return false;
}
let active: &dyn TtsEngine = match role {
MessageRole::User => task.user_engine.as_deref().unwrap_or(task.engine.as_ref()),
_ => task.engine.as_ref(),
};
match active.synthesize(chunk, &task.cancel).await {
Ok(clip) if !clip.is_empty() => {
if let Err(err) = task.playback.enqueue(clip) {
fail(
task.loc
.tf("ui.err.tts_playback", &[("err", &err.to_string())]),
);
return false;
}
true
}
Ok(_) => true,
Err(err) => {
if !task.cancel.is_cancelled() {
fail(
task.loc
.tf("ui.err.tts_synth", &[("err", &err.to_string())]),
);
}
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::shared::i18n::{Lang, locale};
fn ru() -> &'static Locale {
locale(Lang::Ru)
}
fn chat_messages() -> Vec<Message> {
vec![
Message::user("первое от пользователя"),
Message::assistant("первый ответ"),
Message::user("второе от пользователя"),
Message::assistant("второй ответ"),
]
}
fn texts(v: Vec<(MessageRole, String)>) -> Vec<String> {
v.into_iter().map(|(_, s)| s).collect()
}
fn asst(s: &str) -> (MessageRole, String) {
(MessageRole::Assistant, s.to_string())
}
fn chunk_texts(chunks: &[(MessageRole, String)]) -> Vec<String> {
chunks.iter().map(|(_, s)| s.clone()).collect()
}
#[test]
fn last_scope_takes_only_final_message() {
let out = build_utterances(&chat_messages(), TtsScope::Last, false, ru()).unwrap();
assert_eq!(
out,
vec![(MessageRole::Assistant, "второй ответ.".to_string())]
);
}
#[test]
fn recent_scope_takes_tail_in_chronological_order() {
let out = build_utterances(&chat_messages(), TtsScope::Recent(3), false, ru()).unwrap();
assert_eq!(out.len(), 3);
assert_eq!(out[0].0, MessageRole::Assistant);
assert!(out[0].1.contains("первый ответ"), "order: {out:?}");
assert_eq!(out[1].0, MessageRole::User);
assert!(out[2].1.contains("второй ответ"));
let all = build_utterances(&chat_messages(), TtsScope::Recent(99), false, ru()).unwrap();
assert_eq!(all.len(), 4);
}
#[test]
fn all_scope_takes_everything_spoken() {
let out = build_utterances(&chat_messages(), TtsScope::All, false, ru()).unwrap();
assert_eq!(out.len(), 4);
}
#[test]
fn service_messages_and_empty_text_are_skipped() {
let messages = vec![
Message::new(MessageRole::System, "системная инструкция"),
Message::user(" "),
Message::new(MessageRole::Tool, "результат инструмента"),
Message::assistant("настоящий ответ"),
];
let out = build_utterances(&messages, TtsScope::All, false, ru()).unwrap();
assert_eq!(out.len(), 1, "only user/assistant get spoken: {out:?}");
assert!(out[0].1.contains("настоящий ответ"));
assert!(build_utterances(&[], TtsScope::All, false, ru()).is_none());
assert!(
build_utterances(
&[Message::new(MessageRole::System, "только системное")],
TtsScope::All,
false,
ru()
)
.is_none()
);
}
#[test]
fn role_prefixes_apply_to_every_scope_including_single() {
let one = texts(build_utterances(&chat_messages(), TtsScope::Last, true, ru()).unwrap());
assert!(
one[0].starts_with(ru().t("speak.role.assistant")),
"the role prefix appears on a single message too: {one:?}"
);
let many =
texts(build_utterances(&chat_messages(), TtsScope::Recent(2), true, ru()).unwrap());
assert!(many[0].starts_with(ru().t("speak.role.user")));
assert!(many[1].starts_with(ru().t("speak.role.assistant")));
}
#[test]
fn message_whose_text_is_all_skippable_is_dropped() {
let messages = vec![Message::assistant("```rust\nfn main() {}\n```")];
let out = build_utterances(&messages, TtsScope::All, false, ru()).unwrap();
assert_eq!(out.len(), 1);
assert!(out[0].1.contains(ru().t("speak.skip.code")), "{out:?}");
}
#[test]
fn chunk_inherits_role_of_source_message() {
let src = vec![
(MessageRole::User, "Раз. Два.".to_string()),
(MessageRole::Assistant, "Три. Четыре.".to_string()),
];
let chunks = chunk_utterances(&src, 6);
assert!(chunks.iter().all(|(_, c)| c.chars().count() <= 6));
assert_eq!(chunks.first().unwrap().0, MessageRole::User);
assert_eq!(chunks.last().unwrap().0, MessageRole::Assistant);
}
#[test]
fn chunking_respects_limit_and_sentence_boundaries() {
let text = "Первое предложение. Второе предложение! Третье предложение?";
let chunks = chunk_utterances(&[asst(text)], 25);
let ch = chunk_texts(&chunks);
assert!(
ch.iter().all(|c| c.chars().count() <= 25),
"the limit is respected: {ch:?}"
);
assert!(
ch.iter()
.all(|c| c.ends_with('.') || c.ends_with('!') || c.ends_with('?')),
"a chunk ends at a sentence boundary: {ch:?}"
);
assert_eq!(ch.join(" ").replace(" ", " "), text);
}
#[test]
fn short_text_stays_single_chunk() {
let chunks = chunk_utterances(&[asst("Коротко.")], 4096);
assert_eq!(chunk_texts(&chunks), vec!["Коротко.".to_string()]);
}
#[test]
fn utterances_are_not_merged_across_messages() {
let chunks = chunk_utterances(&[asst("Раз."), asst("Два.")], 4096);
assert_eq!(
chunk_texts(&chunks),
vec!["Раз.".to_string(), "Два.".to_string()]
);
}
#[test]
fn overlong_sentence_is_split_by_words_then_chars() {
let long_words = "слово ".repeat(20);
let chunks = chunk_utterances(&[asst(long_words.trim())], 20);
let ch = chunk_texts(&chunks);
assert!(ch.iter().all(|c| c.chars().count() <= 20), "{ch:?}");
assert!(ch.len() > 1);
let giant = "я".repeat(50);
let chunks = chunk_utterances(&[asst(&giant)], 20);
let ch = chunk_texts(&chunks);
assert!(ch.iter().all(|c| c.chars().count() <= 20), "{ch:?}");
assert_eq!(ch.concat().chars().count(), giant.chars().count());
}
#[test]
fn setup_error_keys_exist_in_bundle() {
for err in [
TtsSetupError::Model,
TtsSetupError::ApiKey,
TtsSetupError::Url,
] {
let key = setup_error_key(err);
assert!(ru().has_key(key), "key {key} must be in the bundle");
}
}
}