use super::ui::{PickerCandidate, PickerMessage};
use crate::config::{Config, Provider};
use crate::telemetry::{TelemetrySink, TranscriptionEvent as PipelineEvent};
use crate::transcription::{
OrderedItemTranscript, RealtimeTranscriber, TranscriptSegment, TranscriptionEvent,
TranscriptionMetadata,
};
use std::path::PathBuf;
use std::sync::Mutex;
pub(super) struct PickerStatusSink {
tx: std::sync::mpsc::Sender<PickerMessage>,
provider: Provider,
model: String,
streaming: bool,
last: Mutex<Option<String>>,
}
impl PickerStatusSink {
pub(super) fn new(
tx: std::sync::mpsc::Sender<PickerMessage>,
provider: Provider,
model: String,
streaming: bool,
) -> Self {
Self {
tx,
provider,
model,
streaming,
last: Mutex::new(None),
}
}
fn push(&self, status_text: String) {
if let Ok(mut last) = self.last.lock() {
if last.as_deref() == Some(status_text.as_str()) {
return;
}
*last = Some(status_text.clone());
}
let _ = self.tx.send(PickerMessage::CandidateStatus {
provider: self.provider,
model: self.model.clone(),
streaming: self.streaming,
status_text,
});
}
}
impl TelemetrySink for PickerStatusSink {
fn emit(&self, event: PipelineEvent) {
match event {
PipelineEvent::PreflightStarted { .. } => {
self.push("pre-validating model…".to_string())
}
PipelineEvent::PreflightCompleted { .. } => {}
PipelineEvent::RequestStarted { .. } => self.push("connecting…".to_string()),
PipelineEvent::ConnectionEstablished { .. } => self.push("uploading…".to_string()),
PipelineEvent::UploadComplete { .. } => self.push("waiting for server…".to_string()),
PipelineEvent::ResponseHeaders { status, .. } if (200..300).contains(&status) => {
self.push("transcribing…".to_string())
}
PipelineEvent::RetryScheduled {
kind,
attempt,
max,
delay,
..
} => {
let label = match kind {
crate::telemetry::RetryKind::Connection => "connect retry",
crate::telemetry::RetryKind::Data => "server retry",
};
if delay.is_zero() {
self.push(format!("{} {}/{}…", label, attempt, max))
} else {
self.push(format!(
"{} {}/{} in {}s…",
label,
attempt,
max,
delay.as_secs()
))
}
}
PipelineEvent::Status { message, .. } => self.push(message),
_ => {}
}
}
}
const REALTIME_FEED_CHUNK: usize = 480;
#[derive(Default)]
struct PickerTranscriptAccumulator {
generic_text: String,
item_text: OrderedItemTranscript,
timed_segments: Vec<TranscriptSegment>,
}
impl PickerTranscriptAccumulator {
fn apply(&mut self, event: &TranscriptionEvent) -> Option<String> {
match event {
TranscriptionEvent::TextDelta { text } => {
self.generic_text.push_str(text);
Some(self.final_text())
}
TranscriptionEvent::SegmentDelta { text, start, end } => {
if text.is_empty() {
return None;
}
let trimmed = text.trim().to_string();
if let (Some(start), Some(end)) = (start, end) {
self.timed_segments.push(TranscriptSegment {
start: *start,
end: *end,
text: trimmed.clone(),
});
}
if !self.generic_text.is_empty() {
self.generic_text.push(' ');
}
self.generic_text.push_str(&trimmed);
Some(self.final_text())
}
TranscriptionEvent::ItemCreated {
item_id,
previous_item_id,
} => {
self.item_text
.item_created(item_id, previous_item_id.as_deref());
None
}
TranscriptionEvent::ItemTextDelta {
item_id,
content_index,
text,
} => {
self.item_text.append_delta(item_id, *content_index, text);
Some(self.final_text())
}
TranscriptionEvent::ItemTextCompleted {
item_id,
content_index,
transcript,
} => {
self.item_text.complete(item_id, *content_index, transcript);
Some(self.final_text())
}
_ => None,
}
}
fn final_text(&self) -> String {
if self.item_text.is_empty() {
self.generic_text.trim().to_string()
} else {
self.item_text.snapshot().trim().to_string()
}
}
}
pub(super) async fn run_realtime_transcription(
mut transcriber: Box<dyn RealtimeTranscriber>,
samples: std::sync::Arc<Vec<i16>>,
tx: std::sync::mpsc::Sender<PickerMessage>,
provider: Provider,
model: String,
cancel_token: tokio_util::sync::CancellationToken,
) {
let status_sink: std::sync::Arc<dyn TelemetrySink> = std::sync::Arc::new(
PickerStatusSink::new(tx.clone(), provider, model.clone(), true),
);
transcriber.set_sink(status_sink);
transcriber.set_cancel_token(cancel_token);
let (audio_tx, audio_rx) = tokio::sync::mpsc::channel::<Vec<i16>>(100);
let event_rx = match transcriber.transcribe_realtime(audio_rx).await {
Ok(rx) => rx,
Err(e) => {
log::warn!("realtime {}:{} connect failed: {}", provider, model, e);
let _ = tx.send(PickerMessage::Candidate(Box::new(PickerCandidate::error(
provider,
model,
format!("{e}"),
true,
))));
return;
}
};
let feeder_samples = samples.clone();
tokio::spawn(async move {
let data = &*feeder_samples;
let mut offset = 0;
while offset < data.len() {
let end = (offset + REALTIME_FEED_CHUNK).min(data.len());
let chunk = data[offset..end].to_vec();
if audio_tx.send(chunk).await.is_err() {
break;
}
offset = end;
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
});
let mut event_rx = event_rx;
let mut transcript = PickerTranscriptAccumulator::default();
loop {
match event_rx.recv().await {
Some(event @ TranscriptionEvent::TextDelta { .. })
| Some(event @ TranscriptionEvent::SegmentDelta { .. })
| Some(event @ TranscriptionEvent::ItemCreated { .. })
| Some(event @ TranscriptionEvent::ItemTextDelta { .. })
| Some(event @ TranscriptionEvent::ItemTextCompleted { .. }) => {
if let Some(accumulated_text) = transcript.apply(&event) {
let _ = tx.send(PickerMessage::StreamUpdate {
provider,
model: model.clone(),
accumulated_text,
});
}
}
Some(TranscriptionEvent::Done) => {
let final_text = transcript.final_text();
let _ = tx.send(PickerMessage::Candidate(Box::new(
PickerCandidate::success(
provider,
model,
final_text,
true,
if transcript.timed_segments.is_empty() {
None
} else {
Some(transcript.timed_segments)
},
TranscriptionMetadata::default(),
),
)));
return;
}
Some(TranscriptionEvent::Error { message }) => {
let _ = tx.send(PickerMessage::Candidate(Box::new(PickerCandidate::error(
provider, model, message, true,
))));
return;
}
None => {
let final_text = transcript.final_text();
let _ = tx.send(PickerMessage::Candidate(Box::new(
PickerCandidate::success(
provider,
model,
final_text,
true,
if transcript.timed_segments.is_empty() {
None
} else {
Some(transcript.timed_segments)
},
TranscriptionMetadata::default(),
),
)));
return;
}
_ => {
}
}
}
}
pub(super) fn spawn_transcription(
tasks: &mut tokio::task::JoinSet<()>,
tx: std::sync::mpsc::Sender<PickerMessage>,
audio: PathBuf,
provider: Provider,
model: String,
config: std::sync::Arc<Config>,
_ignored_cancel: tokio_util::sync::CancellationToken,
) {
let sink: std::sync::Arc<dyn TelemetrySink> = std::sync::Arc::new(PickerStatusSink::new(
tx.clone(),
provider,
model.clone(),
false,
));
tasks.spawn(async move {
let local_job =
match crate::transcription::jobs::register_local(&audio, provider, &model, false) {
Ok(j) => Some(j),
Err(e) => {
log::warn!(
"picker: failed to register job for {}:{}: {}",
provider,
model,
e
);
None
}
};
let cancel_token = local_job
.as_ref()
.map(|j| j.cancel_token())
.unwrap_or_default();
let msg = match crate::transcription::transcribe_audio(
&audio,
config.as_ref(),
provider,
Some(&model),
false,
crate::transcription::TranscribeOptions {
allow_api: true,
policy: crate::transcription::RequestTimeoutPolicy::UserAttended,
cancel_token: Some(cancel_token),
skip_legacy_lock: local_job.is_some(),
},
&sink,
)
.await
{
Ok(res) => {
let text = res.text.trim().to_string();
PickerMessage::Candidate(Box::new(PickerCandidate::success(
provider,
model,
text,
false,
res.segments,
res.metadata,
)))
}
Err(e) => {
log::warn!("candidate {}:{} failed: {}", provider, model, e);
PickerMessage::Candidate(Box::new(PickerCandidate::error(
provider,
model,
format!("{e}"),
false,
)))
}
};
drop(local_job);
let _ = tx.send(msg);
});
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Instant;
#[test]
fn picker_reconciles_item_snapshots_and_final_text() {
let mut transcript = PickerTranscriptAccumulator::default();
assert_eq!(
transcript.apply(&TranscriptionEvent::ItemCreated {
item_id: "item-1".to_string(),
previous_item_id: None,
}),
None
);
assert_eq!(
transcript.apply(&TranscriptionEvent::ItemTextDelta {
item_id: "item-1".to_string(),
content_index: 0,
text: "Hello world".to_string(),
}),
Some("Hello world".to_string())
);
assert_eq!(
transcript.apply(&TranscriptionEvent::ItemTextCompleted {
item_id: "item-1".to_string(),
content_index: 0,
transcript: "Hello, corrected world.".to_string(),
}),
Some("Hello, corrected world.".to_string())
);
assert_eq!(transcript.final_text(), "Hello, corrected world.");
}
#[test]
fn picker_generic_segments_remain_additive() {
let mut transcript = PickerTranscriptAccumulator::default();
transcript.apply(&TranscriptionEvent::TextDelta {
text: "prefix".to_string(),
});
transcript.apply(&TranscriptionEvent::SegmentDelta {
text: "first".to_string(),
start: None,
end: None,
});
transcript.apply(&TranscriptionEvent::SegmentDelta {
text: "second".to_string(),
start: None,
end: None,
});
assert_eq!(transcript.final_text(), "prefix first second");
}
fn drain_status(
rx: &std::sync::mpsc::Receiver<PickerMessage>,
) -> Vec<Option<(Provider, String, String)>> {
let mut out = Vec::new();
while let Ok(msg) = rx.try_recv() {
out.push(match msg {
PickerMessage::CandidateStatus {
provider,
model,
status_text,
..
} => Some((provider, model, status_text)),
_ => None,
});
}
out
}
#[test]
fn picker_status_sink_translates_phase_transitions() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(
tx,
Provider::Mistral,
"voxtral-mini-2602".to_string(),
false,
);
let now = Instant::now();
sink.emit(PipelineEvent::PreflightStarted { t: now });
sink.emit(PipelineEvent::PreflightCompleted {
success: true,
t: now,
});
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://example.invalid".into(),
t: now,
});
sink.emit(PipelineEvent::ConnectionEstablished { t: now });
sink.emit(PipelineEvent::UploadProgress {
bytes_sent: 1024,
total: 8192,
t: now,
});
sink.emit(PipelineEvent::UploadProgress {
bytes_sent: 4096,
total: 8192,
t: now,
});
sink.emit(PipelineEvent::UploadComplete {
total: 8192,
t: now,
});
sink.emit(PipelineEvent::ResponseHeaders {
status: 200,
t: now,
});
sink.emit(PipelineEvent::DownloadProgress {
bytes_received: 64,
total: None,
t: now,
});
sink.emit(PipelineEvent::ResponseComplete { total: 64, t: now });
sink.emit(PipelineEvent::RequestCompleted {
success: true,
t: now,
});
let drained = drain_status(&rx);
let strings: Vec<&str> = drained
.iter()
.filter_map(|o| o.as_ref().map(|(_, _, s)| s.as_str()))
.collect();
assert_eq!(
strings,
vec![
"pre-validating model…",
"connecting…",
"uploading…",
"waiting for server…",
"transcribing…",
],
"phase-transition strings must arrive in the documented order"
);
}
#[test]
fn picker_status_sink_cache_hit_starts_at_connecting() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(tx, Provider::Mistral, "voxtral".to_string(), false);
let now = Instant::now();
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://x".into(),
t: now,
});
sink.emit(PipelineEvent::ConnectionEstablished { t: now });
sink.emit(PipelineEvent::UploadComplete {
total: 1024,
t: now,
});
sink.emit(PipelineEvent::ResponseHeaders {
status: 200,
t: now,
});
let drained = drain_status(&rx);
let strings: Vec<&str> = drained
.iter()
.filter_map(|o| o.as_ref().map(|(_, _, s)| s.as_str()))
.collect();
assert_eq!(
strings,
vec![
"connecting…",
"uploading…",
"waiting for server…",
"transcribing…",
],
"cache-hit path must skip the pre-validating status"
);
}
#[test]
fn picker_status_sink_preflight_completed_is_silent() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(tx, Provider::Mistral, "voxtral".to_string(), false);
let now = Instant::now();
sink.emit(PipelineEvent::PreflightStarted { t: now });
let _ = drain_status(&rx);
sink.emit(PipelineEvent::PreflightCompleted {
success: true,
t: now,
});
sink.emit(PipelineEvent::PreflightCompleted {
success: false,
t: now,
});
let drained = drain_status(&rx);
assert!(
drained.is_empty(),
"PreflightCompleted must not push any status; got {:?}",
drained
);
}
#[test]
fn picker_status_sink_dedups_repeated_phase() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(tx, Provider::Mistral, "voxtral".to_string(), false);
let now = Instant::now();
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://x".into(),
t: now,
});
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://x".into(),
t: now,
});
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://x".into(),
t: now,
});
let drained = drain_status(&rx);
assert_eq!(drained.len(), 1, "dedup must collapse repeated phases");
}
#[test]
fn picker_status_sink_formats_retry_attempts() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(tx, Provider::OpenAI, "whisper-1".to_string(), false);
let now = Instant::now();
sink.emit(PipelineEvent::RetryScheduled {
kind: crate::telemetry::RetryKind::Connection,
attempt: 2,
max: 5,
reason: "transient".into(),
delay: std::time::Duration::ZERO,
t: now,
});
sink.emit(PipelineEvent::RetryScheduled {
kind: crate::telemetry::RetryKind::Data,
attempt: 1,
max: 3,
reason: "5xx".into(),
delay: std::time::Duration::from_secs(30),
t: now,
});
let drained = drain_status(&rx);
let strings: Vec<&str> = drained
.iter()
.filter_map(|o| o.as_ref().map(|(_, _, s)| s.as_str()))
.collect();
assert_eq!(
strings,
vec!["connect retry 2/5…", "server retry 1/3 in 30s…"]
);
}
#[test]
fn picker_status_sink_drops_non_2xx_response_headers() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(tx, Provider::Mistral, "voxtral".to_string(), false);
let now = Instant::now();
sink.emit(PipelineEvent::ResponseHeaders {
status: 401,
t: now,
});
sink.emit(PipelineEvent::ResponseHeaders {
status: 500,
t: now,
});
let drained = drain_status(&rx);
assert!(
drained.is_empty(),
"non-2xx response headers must not emit a status; got: {:?}",
drained
);
}
#[test]
fn picker_status_sink_preserves_row_identity() {
let (tx, rx) = std::sync::mpsc::channel::<PickerMessage>();
let sink = PickerStatusSink::new(
tx,
Provider::Mistral,
"voxtral-mini-2602".to_string(),
false,
);
let now = Instant::now();
sink.emit(PipelineEvent::RequestStarted {
endpoint: "https://x".into(),
t: now,
});
let drained = drain_status(&rx);
assert_eq!(drained.len(), 1);
let (p, m, _) = drained[0].as_ref().expect("status message");
assert_eq!(*p, Provider::Mistral);
assert_eq!(m, "voxtral-mini-2602");
}
}