This commit is contained in:
iamdoubz
2026-06-30 23:18:29 -05:00
parent 849fa79647
commit f9b172aff4
24 changed files with 1867 additions and 177 deletions
+127
View File
@@ -344,6 +344,66 @@ impl Resampler {
}
}
/// Read a WAV file written by this app's own `WavWriter` (Float32 or Int16;
/// see `write_wav_bytes`) and return mono `f32` samples at 16kHz — the batch
/// counterpart to the live capture path's downmix/resample, used for crash
/// recovery (`Transcriber::transcribe_file`, T2.8).
///
/// A working WAV left behind by an abnormal exit never had `finalize()` patch
/// its RIFF/`data` chunk sizes, so they read back as 0 — exactly the file
/// recovery has to work from (FR-REL-1). Trusting hound's own duration/sample
/// iterator would silently see zero frames despite real PCM bytes on disk (and
/// feeding whisper.cpp zero samples hangs rather than erroring — verified).
/// So this reads the `data` payload as raw bytes directly off disk instead of
/// through hound's size-aware iterator; hound is only used to parse the
/// format/`fmt ` chunk.
pub fn read_wav_mono_16k(path: &Path) -> Result<Vec<f32>, AudioError> {
let reader =
hound::WavReader::open(path).map_err(|e| AudioError::Capture(format!("open wav: {e}")))?;
let spec = reader.spec();
let channels = spec.channels.max(1) as usize;
drop(reader);
let file = std::fs::read(path).map_err(|e| AudioError::Capture(format!("read wav: {e}")))?;
let data_marker = file
.windows(4)
.position(|w| w == b"data")
.ok_or_else(|| AudioError::Capture("wav has no data chunk".into()))?;
let payload_start = data_marker + 8; // "data" + 4-byte (possibly unfinalized) size field
let bytes = file.get(payload_start..).unwrap_or(&[]);
let mono: Vec<f32> = match (spec.sample_format, spec.bits_per_sample) {
(SampleFormat::Float, 32) => bytes
.chunks_exact(4 * channels)
.map(|frame| {
frame
.chunks_exact(4)
.map(|c| f32::from_le_bytes(c.try_into().unwrap()))
.sum::<f32>()
/ channels as f32
})
.collect(),
(SampleFormat::Int, 16) => bytes
.chunks_exact(2 * channels)
.map(|frame| {
frame
.chunks_exact(2)
.map(|c| i16::from_le_bytes(c.try_into().unwrap()) as f32 / i16::MAX as f32)
.sum::<f32>()
/ channels as f32
})
.collect(),
(fmt, bits) => {
return Err(AudioError::Device(format!(
"unsupported wav format: {fmt:?} {bits}-bit"
)))
}
};
let mut resampler = Resampler::new(spec.sample_rate);
Ok(resampler.process(&mono))
}
#[cfg(test)]
mod tests {
use super::*;
@@ -381,6 +441,73 @@ mod tests {
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn read_wav_mono_16k_reads_a_normally_finalized_file() {
let dir = std::env::temp_dir().join(format!("wa-test-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("audio.wav");
let spec = WavSpec {
channels: 1,
sample_rate: 16_000,
bits_per_sample: 32,
sample_format: SampleFormat::Float,
};
let mut writer = WavWriter::create(&path, spec).unwrap();
for i in 0..1600 {
writer.write_sample((i as f32 / 1600.0) - 0.5).unwrap();
}
writer.finalize().unwrap();
let mono = read_wav_mono_16k(&path).unwrap();
assert_eq!(mono.len(), 1600); // already 16kHz mono, no resampling needed
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn read_wav_mono_16k_recovers_samples_from_an_unfinalized_crash_file() {
// hound's own `WavWriter` patches the RIFF/data sizes in both
// `finalize()` and `flush()`, and its `Drop` impl finalizes too — so
// there's no way to get it to leave an unfinalized file on disk short
// of `mem::forget`-ing mid-write, which depends on `BufWriter`'s
// internal (unspecified) flush threshold. Building the crashed file
// by hand instead gives a deterministic repro of exactly what a real
// crash leaves: a valid `fmt ` chunk with the RIFF/data sizes at 0.
let dir = std::env::temp_dir().join(format!("wa-test-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("audio.wav");
let samples: Vec<f32> = (0..1600).map(|i| (i as f32 / 1600.0) - 0.5).collect();
let pcm: Vec<u8> = samples.iter().flat_map(|s| s.to_le_bytes()).collect();
let mut file = Vec::new();
file.extend_from_slice(b"RIFF");
file.extend_from_slice(&0u32.to_le_bytes()); // unfinalized RIFF size
file.extend_from_slice(b"WAVE");
file.extend_from_slice(b"fmt ");
file.extend_from_slice(&16u32.to_le_bytes()); // fmt chunk size
file.extend_from_slice(&3u16.to_le_bytes()); // WAVE_FORMAT_IEEE_FLOAT
file.extend_from_slice(&1u16.to_le_bytes()); // channels
file.extend_from_slice(&16_000u32.to_le_bytes()); // sample rate
file.extend_from_slice(&(16_000 * 4u32).to_le_bytes()); // byte rate
file.extend_from_slice(&4u16.to_le_bytes()); // block align
file.extend_from_slice(&32u16.to_le_bytes()); // bits per sample
file.extend_from_slice(b"data");
file.extend_from_slice(&0u32.to_le_bytes()); // unfinalized data size
file.extend_from_slice(&pcm);
std::fs::write(&path, &file).unwrap();
let mono = read_wav_mono_16k(&path).unwrap();
assert_eq!(
mono.len(),
1600,
"should recover real samples despite the zeroed header"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn decode_mono_f32_downmixes_stereo() {
let format = WaveFormat::new(32, 32, &SampleType::Float, 48_000, 2, None);
+238 -44
View File
@@ -8,10 +8,14 @@
use crate::audio::{AudioCapture, WasapiCapture};
use crate::error::WaResult;
use crate::models::*;
use crate::notes::NotesRenderer;
use crate::paths::{meeting_dir, settings_path, wa_root, whisper_model_path};
use crate::storage::{FinalizeMeeting, Meeting, NewMeeting};
use crate::transcription::{run_streaming_worker, Transcriber, WhisperTranscriber};
use crate::{error::WaError, AppState, RecordingSession};
use serde::Deserialize;
use std::path::PathBuf;
use std::sync::{Arc, Mutex as StdMutex};
use tauri::{AppHandle, Emitter, State};
#[derive(Deserialize)]
@@ -23,30 +27,9 @@ pub struct StartRecordingArgs {
pub record: bool,
}
// ---- Phase 1 file-backed settings + storage paths ----
// Full SQLite persistence is Phase 2 (`storage::Store`, still `todo!()`); Phase
// 1 only needs `settings.json` (consent/default-retention, FR-REC-1/2) and the
// `meetings/<id>/` folder layout from `docs/03-data-model.md`.
fn local_appdata() -> PathBuf {
PathBuf::from(std::env::var("LOCALAPPDATA").unwrap_or_else(|_| ".".to_string()))
}
fn wa_root() -> PathBuf {
local_appdata().join("WhispAssist")
}
fn settings_path() -> PathBuf {
wa_root().join("settings.json")
}
fn meetings_dir() -> PathBuf {
wa_root().join("meetings")
}
fn whisper_model_path() -> PathBuf {
wa_root().join("models").join("ggml-base.en-q5_1.bin")
}
// ---- File-backed settings (consent/default-retention/storage policy) ----
// `settings.json` per `docs/03-data-model.md`; meeting rows/files go through
// `AppState.store` (Phase 2, `storage::SqliteStore`).
fn default_settings() -> Settings {
Settings {
@@ -60,10 +43,12 @@ fn default_settings() -> Settings {
default_record: false,
consent_acknowledged: false,
sync_enabled: false,
retention_max_age_days: None,
retention_max_size_gb: None,
}
}
fn load_settings() -> Settings {
pub(crate) fn load_settings() -> Settings {
std::fs::read_to_string(settings_path())
.ok()
.and_then(|s| serde_json::from_str(&s).ok())
@@ -114,10 +99,17 @@ pub async fn start_recording(
));
}
let meeting_id = uuid::Uuid::new_v4().to_string();
let meeting_dir = meetings_dir().join(&meeting_id);
std::fs::create_dir_all(&meeting_dir).map_err(|e| WaError::new("storage", e.to_string()))?;
let wav_path = meeting_dir.join("audio.wav");
let meeting_id = state
.store
.create_meeting(NewMeeting {
title: args
.meeting_title
.unwrap_or_else(|| "Untitled meeting".to_string()),
calendar_event_id: args.calendar_event_id,
})
.await
.map_err(|e| WaError::new("storage", e.to_string()))?;
let wav_path = meeting_dir(&meeting_id).join("audio.wav");
// Bounded so a slow/stalled transcription worker can never back up the capture thread.
let (frame_tx, frame_rx) = std::sync::mpsc::sync_channel::<Vec<f32>>(8);
@@ -125,6 +117,8 @@ pub async fn start_recording(
.start(&wav_path, frame_tx)
.map_err(|e| WaError::new("audio", e.to_string()))?;
let segments: Arc<StdMutex<Vec<TranscriptSegment>>> = Arc::new(StdMutex::new(Vec::new()));
let segments_for_worker = segments.clone();
let app_for_worker = app.clone();
let meeting_id_for_worker = meeting_id.clone();
let transcription_worker = std::thread::Builder::new()
@@ -132,6 +126,9 @@ pub async fn start_recording(
.spawn(
move || match WhisperTranscriber::load(&model_path, BackendId::Cpu) {
Ok(transcriber) => run_streaming_worker(&transcriber, frame_rx, |segment| {
if let Ok(mut buf) = segments_for_worker.lock() {
buf.push(segment.clone());
}
let _ = app_for_worker.emit(
"transcript://segment",
serde_json::json!({ "meetingId": meeting_id_for_worker, "segment": segment }),
@@ -149,6 +146,7 @@ pub async fn start_recording(
wav_path,
started_at: std::time::Instant::now(),
transcription_worker,
segments,
});
drop(guard);
@@ -193,13 +191,53 @@ pub async fn stop_recording(
// guarantee it has fully drained the audio before we act on retention.
let _ = session.transcription_worker.join();
let segments = session
.segments
.lock()
.map(|g| g.clone())
.unwrap_or_default();
let segment_count = segments.len();
// No diarization until Phase 4 — every segment is provisionally "S1".
let speakers = vec![SpeakerInfo {
label: "S1".to_string(),
display_name: None,
participant_id: None,
}];
let model_used = whisper_model_path()
.file_stem()
.and_then(|s| s.to_str())
.map(str::to_string);
// T2.10: persist transcript.json (via finalize_meeting) and notes.md
// *before* touching the working WAV, so a crash here still leaves a
// recoverable, regenerable meeting.
state
.store
.finalize_meeting(
&meeting_id,
FinalizeMeeting {
segments: segments.clone(),
speakers: speakers.clone(),
duration_secs: (summary.duration_ms / 1000) as i64,
recorded: session.retention,
language: None,
backend_used: Some(BackendId::Cpu.as_str().to_string()),
model_used,
},
)
.await
.map_err(|e| WaError::new("storage", e.to_string()))?;
let notes_md = crate::notes::MarkdownNotes.to_markdown(&segments, &speakers, None);
let _ = std::fs::write(meeting_dir(&meeting_id).join("notes.md"), notes_md);
let _ = app.emit(
"transcript://finalized",
serde_json::json!({ "meetingId": meeting_id, "segmentCount": serde_json::Value::Null }),
serde_json::json!({ "meetingId": meeting_id, "segmentCount": segment_count }),
);
// ADR-0009: delete the working WAV only after the transcript worker has
// finished reading it, and only when retention is off.
// ADR-0009: delete the working WAV only after the transcript is finalized
// above, and only when retention is off.
if !session.retention {
let _ = std::fs::remove_file(&session.wav_path);
}
@@ -299,35 +337,191 @@ pub async fn hardware_status() -> WaResult<Vec<BackendInfo>> {
todo!("Phase 3 — hardware_status")
}
/// Re-run transcription from a `recovering` meeting's working `audio.wav`
/// (T2.8, FR-REL-1). CPU-bound, so it runs on a blocking task rather than
/// tying up an async worker.
#[tauri::command]
pub async fn resume_transcription(
app: AppHandle,
state: State<'_, AppState>,
meeting_id: MeetingId,
) -> WaResult<()> {
let model_path = whisper_model_path();
if !model_path.exists() {
return Err(WaError::new(
"transcription",
format!(
"whisper model not found at {}; run `npm run download-models` first",
model_path.display()
),
));
}
let wav_path = meeting_dir(&meeting_id).join("audio.wav");
if !wav_path.exists() {
return Err(WaError::new(
"recovery",
"no working audio.wav to recover from",
));
}
// Note: deliberately doesn't emit `recording://state` — that event drives
// the main header's single live-session Record/Stop toggle, and reusing
// it here would make the header show a phantom "Stop" button for a
// recovery run that isn't a `RecordingSession` at all. The "recovering"
// badge in the meetings list is this operation's own progress signal.
let segments = tauri::async_runtime::spawn_blocking(move || {
WhisperTranscriber::load(&model_path, BackendId::Cpu)
.and_then(|transcriber| transcriber.transcribe_file(&wav_path))
})
.await
.map_err(|e| WaError::new("transcription", e.to_string()))?
.map_err(|e| WaError::new("transcription", e.to_string()))?;
let speakers = vec![SpeakerInfo {
label: "S1".to_string(),
display_name: None,
participant_id: None,
}];
let duration_secs = segments
.last()
.map(|s| (s.end_ms / 1000) as i64)
.unwrap_or(0);
let model_used = whisper_model_path()
.file_stem()
.and_then(|s| s.to_str())
.map(str::to_string);
state
.store
.finalize_meeting(
&meeting_id,
FinalizeMeeting {
segments: segments.clone(),
speakers: speakers.clone(),
duration_secs,
// The working WAV surviving a crash is the only signal we have
// left about intent; keep it rather than silently discard it.
recorded: true,
language: None,
backend_used: Some(BackendId::Cpu.as_str().to_string()),
model_used,
},
)
.await
.map_err(|e| WaError::new("storage", e.to_string()))?;
let notes_md = crate::notes::MarkdownNotes.to_markdown(&segments, &speakers, None);
let _ = std::fs::write(meeting_dir(&meeting_id).join("notes.md"), notes_md);
let _ = app.emit(
"transcript://finalized",
serde_json::json!({ "meetingId": meeting_id, "segmentCount": segments.len() }),
);
Ok(())
}
// ---- Meetings / storage (Phase 2) ----
#[tauri::command]
pub async fn list_meetings(_query: Option<String>) -> WaResult<Vec<MeetingListItem>> {
todo!("Phase 2 — list_meetings")
pub async fn list_meetings(
state: State<'_, AppState>,
query: Option<String>,
) -> WaResult<Vec<MeetingListItem>> {
state
.store
.list_meetings(query)
.await
.map_err(|e| WaError::new("storage", e.to_string()))
}
#[tauri::command]
pub async fn get_meeting(_meeting_id: MeetingId) -> WaResult<MeetingListItem> {
todo!("Phase 2 — get_meeting (returns full Meeting incl. transcript)")
pub async fn get_meeting(state: State<'_, AppState>, meeting_id: MeetingId) -> WaResult<Meeting> {
state
.store
.get_meeting(&meeting_id)
.await
.map_err(|e| WaError::new("storage", e.to_string()))
}
#[tauri::command]
pub async fn delete_meeting(_meeting_id: MeetingId) -> WaResult<()> {
todo!("Phase 2 — delete_meeting")
pub async fn delete_meeting(state: State<'_, AppState>, meeting_id: MeetingId) -> WaResult<()> {
let guard = state.session.lock().await;
if guard.as_ref().is_some_and(|s| s.meeting_id == meeting_id) {
return Err(WaError::new(
"recording",
"cannot delete a meeting that is currently recording",
));
}
drop(guard);
state
.store
.delete_meeting(&meeting_id)
.await
.map_err(|e| WaError::new("storage", e.to_string()))
}
#[tauri::command]
pub async fn update_notes(_meeting_id: MeetingId, _markdown: String) -> WaResult<()> {
todo!("Phase 2 — update_notes")
pub async fn update_notes(
state: State<'_, AppState>,
meeting_id: MeetingId,
markdown: String,
) -> WaResult<()> {
state
.store
.update_notes(&meeting_id, &markdown)
.await
.map_err(|e| WaError::new("storage", e.to_string()))
}
/// `format` ∈ `md | bundle` (PDF/Word stay Phase 8). `dest` is a file path for
/// `md`, a destination folder for `bundle` (the frontend gets it from a native
/// Save/choose-folder dialog).
#[tauri::command]
pub async fn export_meeting(
_meeting_id: MeetingId,
_dest: String,
_format: String,
state: State<'_, AppState>,
meeting_id: MeetingId,
dest: String,
format: String,
) -> WaResult<String> {
todo!("Phase 2/8 — export_meeting (md|pdf|docx|bundle)")
let meeting = state
.store
.get_meeting(&meeting_id)
.await
.map_err(|e| WaError::new("storage", e.to_string()))?;
let dest_path = PathBuf::from(&dest);
match format.as_str() {
"md" => {
crate::notes::MarkdownNotes
.export(
&meeting.notes_markdown,
&dest_path,
crate::notes::ExportFormat::Md,
)
.map_err(|e| WaError::new("export", e.to_string()))?;
}
"bundle" => {
std::fs::create_dir_all(&dest_path)
.map_err(|e| WaError::new("export", e.to_string()))?;
let source_dir = meeting_dir(&meeting_id);
for name in ["audio.wav", "transcript.json"] {
let src = source_dir.join(name);
if src.exists() {
std::fs::copy(&src, dest_path.join(name))
.map_err(|e| WaError::new("export", e.to_string()))?;
}
}
std::fs::write(dest_path.join("notes.md"), &meeting.notes_markdown)
.map_err(|e| WaError::new("export", e.to_string()))?;
}
other => {
return Err(WaError::new(
"export",
format!("unsupported export format: {other}"),
))
}
}
Ok(dest)
}
// ---- LLM (Phase 5) ----
+42 -6
View File
@@ -15,11 +15,13 @@ pub mod llm;
pub mod mcp;
pub mod models;
pub mod notes;
pub mod paths;
pub mod storage;
pub mod sync;
pub mod transcription;
use std::path::PathBuf;
use std::sync::{Arc, Mutex as StdMutex};
use std::thread::JoinHandle;
use tauri::tray::TrayIcon;
use tauri::Manager;
@@ -27,12 +29,11 @@ use tokio::sync::Mutex;
/// Shared application state handed to every Tauri command via `tauri::State`.
///
/// Phase 1 only needs the in-flight recording session. `hardware`/`store`/
/// `llm`/`sync`/`mcp` service handles join this struct once their concrete
/// implementations exist (Phase 2/3/5/9/10) — wiring them in ahead of that
/// isn't possible yet (several don't have a concrete impl to construct) and
/// would just be unused plumbing until then.
/// `hardware`/`llm`/`sync`/`mcp` service handles join this struct once their
/// concrete implementations exist (Phase 3/5/9/10) — several don't have one to
/// construct yet, so wiring them in early would just be unused plumbing.
pub struct AppState {
pub store: Arc<dyn storage::Store>,
/// The single in-flight recording, if any. Guarded so start/stop/pause are atomic.
pub session: Mutex<Option<RecordingSession>>,
}
@@ -46,6 +47,10 @@ pub struct RecordingSession {
pub wav_path: PathBuf,
pub started_at: std::time::Instant,
pub transcription_worker: JoinHandle<()>,
/// Segments produced so far, appended by the same closure that emits
/// `transcript://segment` — lets `stop_recording` persist the full
/// transcript without a second round-trip through the event channel.
pub segments: Arc<StdMutex<Vec<models::TranscriptSegment>>>,
}
/// Wraps the tray icon so it can be looked up from commands to update its
@@ -56,17 +61,47 @@ pub struct TrayHandle(pub TrayIcon);
pub fn run() {
tracing_subscriber::fmt().with_env_filter("info").init();
let store: Arc<dyn storage::Store> = Arc::new(
tauri::async_runtime::block_on(storage::SqliteStore::connect())
.expect("failed to initialize storage (wa.db)"),
);
let store_for_setup = store.clone();
tauri::Builder::default()
.plugin(tauri_plugin_dialog::init())
.manage(AppState {
store,
session: Mutex::new(None),
})
.setup(|app| {
.setup(move |app| {
let icon = tauri::image::Image::from_bytes(include_bytes!("../icons/tray.png"))?;
let tray = tauri::tray::TrayIconBuilder::new()
.icon(icon)
.tooltip("WhispAssist — idle")
.build(app)?;
app.manage(TrayHandle(tray));
// Startup recovery + retention pass (FR-REL-1, FR-STORE-2). Spawned so
// it never blocks the window from showing (NFR-PERF-4); nothing here
// repeats on a timer (NFR-RES-1).
let store = store_for_setup;
tauri::async_runtime::spawn(async move {
match store.recover_scan().await {
Ok(recovered) if !recovered.is_empty() => {
tracing::info!("recovered {} interrupted meeting(s)", recovered.len());
}
Ok(_) => {}
Err(e) => tracing::error!("startup recovery scan failed: {e}"),
}
let settings = commands::load_settings();
let policy = storage::Retention {
max_age_days: settings.retention_max_age_days,
max_size_gb: settings.retention_max_size_gb,
};
if let Err(e) = store.enforce_retention(policy).await {
tracing::error!("startup retention enforcement failed: {e}");
}
});
Ok(())
})
.invoke_handler(tauri::generate_handler![
@@ -76,6 +111,7 @@ pub fn run() {
commands::resume_recording,
commands::set_recording_retention,
commands::acknowledge_recording_consent,
commands::resume_transcription,
commands::hardware_status,
commands::list_meetings,
commands::get_meeting,
+41 -1
View File
@@ -15,6 +15,18 @@ pub enum BackendId {
Cpu,
}
impl BackendId {
pub fn as_str(&self) -> &'static str {
match self {
BackendId::Npu => "npu",
BackendId::Nvidia => "nvidia",
BackendId::Amd => "amd",
BackendId::Intel => "intel",
BackendId::Cpu => "cpu",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BackendInfo {
pub id: BackendId,
@@ -25,7 +37,7 @@ pub struct BackendInfo {
pub vram_mb: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum MeetingStatus {
Recording,
@@ -35,6 +47,31 @@ pub enum MeetingStatus {
Error,
}
impl MeetingStatus {
pub fn as_str(&self) -> &'static str {
match self {
MeetingStatus::Recording => "recording",
MeetingStatus::Transcribing => "transcribing",
MeetingStatus::Ready => "ready",
MeetingStatus::Recovering => "recovering",
MeetingStatus::Error => "error",
}
}
/// Parses a `meetings.status` DB value. Unrecognized values map to `Error`
/// rather than panicking — a stored value should always be one written by
/// this app, but a row is not worth crashing the app over.
pub fn parse(s: &str) -> Self {
match s {
"recording" => MeetingStatus::Recording,
"transcribing" => MeetingStatus::Transcribing,
"ready" => MeetingStatus::Ready,
"recovering" => MeetingStatus::Recovering,
_ => MeetingStatus::Error,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TranscriptSegment {
pub id: u64,
@@ -112,6 +149,9 @@ pub struct Settings {
pub consent_acknowledged: bool,
// Sync master switch (ADR-0010). Default OFF. Target rows live in the DB; secrets in OS keychain.
pub sync_enabled: bool,
// Storage retention policy (FR-STORE-2). None = no cap on that dimension.
pub retention_max_age_days: Option<u32>,
pub retention_max_size_gb: Option<u32>,
}
// ---- Sync (ADR-0010) ----
+121 -10
View File
@@ -42,20 +42,131 @@ pub struct MarkdownNotes;
impl NotesRenderer for MarkdownNotes {
fn to_markdown(
&self,
_segments: &[TranscriptSegment],
_speakers: &[SpeakerInfo],
_summary_md: Option<&str>,
segments: &[TranscriptSegment],
speakers: &[SpeakerInfo],
summary_md: Option<&str>,
) -> String {
// T2.4: resolve speaker names, group dialogue, prepend summary section.
todo!("Phase 2 — assemble Markdown")
let mut out = String::new();
if let Some(summary) = summary_md {
let summary = summary.trim();
if !summary.is_empty() {
out.push_str(summary);
out.push_str("\n\n---\n\n");
}
}
let name_for = |label: &str| -> String {
speakers
.iter()
.find(|s| s.label == label)
.and_then(|s| s.display_name.clone())
.unwrap_or_else(|| label.to_string())
};
// Group consecutive segments from the same speaker into one paragraph
// (matters once Phase 4 diarization produces more than one speaker).
let mut current_speaker: Option<&str> = None;
let mut buffer = String::new();
for seg in segments {
let text = seg.text.trim();
if text.is_empty() {
continue;
}
if current_speaker != Some(seg.speaker.as_str()) {
if let Some(speaker) = current_speaker {
out.push_str(&format!("**{}:** {}\n\n", name_for(speaker), buffer.trim()));
}
current_speaker = Some(seg.speaker.as_str());
buffer.clear();
}
if !buffer.is_empty() {
buffer.push(' ');
}
buffer.push_str(text);
}
if let Some(speaker) = current_speaker {
out.push_str(&format!("**{}:** {}\n\n", name_for(speaker), buffer.trim()));
}
out.trim_end().to_string()
}
fn export(
&self,
_markdown: &str,
_dest: &Path,
_fmt: ExportFormat,
markdown: &str,
dest: &Path,
fmt: ExportFormat,
) -> Result<PathBuf, NotesError> {
// T2.6 (.md/bundle); T8.4 (pdf/docx via local conversion).
todo!("Phase 2/8 — export notes")
match fmt {
ExportFormat::Md => {
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(dest, markdown)?;
Ok(dest.to_path_buf())
}
// Bundling pulls in audio.wav/transcript.json alongside notes.md —
// that's meeting-folder orchestration, not Markdown rendering, so
// `commands::export_meeting` handles it directly rather than here.
ExportFormat::Bundle => Err(NotesError::Export(
"bundle export is assembled in commands::export_meeting, not NotesRenderer".into(),
)),
ExportFormat::Pdf | ExportFormat::Docx => Err(NotesError::Export(
"PDF/Word export lands in Phase 8".into(),
)),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn seg(id: u64, speaker: &str, text: &str) -> TranscriptSegment {
TranscriptSegment {
id,
start_ms: id * 1000,
end_ms: id * 1000 + 900,
speaker: speaker.to_string(),
text: text.to_string(),
confidence: Some(0.9),
interim: false,
}
}
#[test]
fn groups_consecutive_same_speaker_segments_into_one_paragraph() {
let segments = vec![
seg(0, "S1", "Hello"),
seg(1, "S1", "world."),
seg(2, "S2", "Hi there."),
];
let md = MarkdownNotes.to_markdown(&segments, &[], None);
assert_eq!(md, "**S1:** Hello world.\n\n**S2:** Hi there.");
}
#[test]
fn resolves_display_name_and_falls_back_to_label() {
let segments = vec![seg(0, "S1", "Hi")];
let speakers = [SpeakerInfo {
label: "S1".to_string(),
display_name: Some("Alex".to_string()),
participant_id: None,
}];
assert_eq!(
MarkdownNotes.to_markdown(&segments, &speakers, None),
"**Alex:** Hi"
);
assert_eq!(
MarkdownNotes.to_markdown(&segments, &[], None),
"**S1:** Hi"
);
}
#[test]
fn skips_blank_segments_and_prepends_summary() {
let segments = vec![seg(0, "S1", " "), seg(1, "S1", "Real text.")];
let md = MarkdownNotes.to_markdown(&segments, &[], Some("## Summary\nDone."));
assert_eq!(md, "## Summary\nDone.\n\n---\n\n**S1:** Real text.");
}
}
+34
View File
@@ -0,0 +1,34 @@
//! Shared on-disk layout helpers (`docs/03-data-model.md`). Used by `storage`,
//! `commands`, and `lib` so the root/meetings/models/settings paths are computed
//! exactly one way.
use crate::models::MeetingId;
use std::path::PathBuf;
pub fn local_appdata() -> PathBuf {
PathBuf::from(std::env::var("LOCALAPPDATA").unwrap_or_else(|_| ".".to_string()))
}
pub fn wa_root() -> PathBuf {
local_appdata().join("WhispAssist")
}
pub fn db_path() -> PathBuf {
wa_root().join("wa.db")
}
pub fn settings_path() -> PathBuf {
wa_root().join("settings.json")
}
pub fn meetings_dir() -> PathBuf {
wa_root().join("meetings")
}
pub fn meeting_dir(id: &MeetingId) -> PathBuf {
meetings_dir().join(id)
}
pub fn whisper_model_path() -> PathBuf {
wa_root().join("models").join("ggml-base.en-q5_1.bin")
}
+357 -18
View File
@@ -4,8 +4,14 @@
//! Invariant: audio is the source of truth. Derived artifacts (transcript/notes/
//! summary) are regenerable; retention never touches an in-progress meeting.
use crate::models::{MeetingId, MeetingListItem};
use crate::models::{MeetingId, MeetingListItem, MeetingStatus, SpeakerInfo, TranscriptSegment};
use crate::paths;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use sqlx::sqlite::{SqliteConnectOptions, SqlitePool, SqlitePoolOptions};
use sqlx::Row;
use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, thiserror::Error)]
pub enum StoreError {
@@ -17,11 +23,49 @@ pub enum StoreError {
NotFound(String),
}
impl From<sqlx::Error> for StoreError {
fn from(e: sqlx::Error) -> Self {
StoreError::Db(e.to_string())
}
}
pub struct NewMeeting {
pub title: String,
pub calendar_event_id: Option<String>,
}
/// What actually happened during the session — richer than the doc's bare
/// `&MeetingId` because `finalize_meeting` has to persist what was captured.
pub struct FinalizeMeeting {
pub segments: Vec<TranscriptSegment>,
pub speakers: Vec<SpeakerInfo>,
pub duration_secs: i64,
pub recorded: bool,
pub language: Option<String>,
pub backend_used: Option<String>,
pub model_used: Option<String>,
}
/// Full meeting detail: DB row + transcript + speakers + notes (`get_meeting`'s
/// documented contract — `docs/04-api-contracts.md` — which the old
/// `MeetingListItem`-only stub couldn't represent).
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Meeting {
pub id: MeetingId,
pub title: String,
pub started_at: i64,
pub ended_at: Option<i64>,
pub duration_secs: Option<i64>,
pub status: MeetingStatus,
pub recorded: bool,
pub language: Option<String>,
pub backend_used: Option<String>,
pub model_used: Option<String>,
pub segments: Vec<TranscriptSegment>,
pub speakers: Vec<SpeakerInfo>,
pub notes_markdown: String,
}
pub struct Retention {
pub max_age_days: Option<u32>,
pub max_size_gb: Option<u32>,
@@ -30,51 +74,346 @@ pub struct Retention {
#[async_trait]
pub trait Store: Send + Sync {
async fn create_meeting(&self, m: NewMeeting) -> Result<MeetingId, StoreError>;
async fn finalize_meeting(&self, id: &MeetingId) -> Result<(), StoreError>;
async fn finalize_meeting(&self, id: &MeetingId, s: FinalizeMeeting) -> Result<(), StoreError>;
async fn list_meetings(
&self,
query: Option<String>,
) -> Result<Vec<MeetingListItem>, StoreError>;
async fn get_meeting(&self, id: &MeetingId) -> Result<Meeting, StoreError>;
async fn delete_meeting(&self, id: &MeetingId) -> Result<(), StoreError>;
async fn update_notes(&self, id: &MeetingId, markdown: &str) -> Result<(), StoreError>;
/// Full-text search across transcripts + notes (Phase 8, FR-SEARCH-1).
async fn search(&self, query: &str) -> Result<Vec<MeetingListItem>, StoreError>;
/// Startup reconcile: meetings with audio but no finalized transcript (FR-REL-1).
async fn recover_scan(&self) -> Result<Vec<MeetingId>, StoreError>;
/// Enforce retention; returns count removed. Skips in-progress meetings.
/// Enforce retention; returns count removed. Skips in-progress meetings
/// (anything whose `status` isn't a terminal `ready`/`error`).
async fn enforce_retention(&self, policy: Retention) -> Result<u32, StoreError>;
}
/// SQLite-backed store. Migrations live in `migrations/` (sqlx::migrate!).
pub struct SqliteStore;
/// SQLite-backed store. Migrations live in `migrations/` (`sqlx::migrate!`).
pub struct SqliteStore {
pool: SqlitePool,
}
impl SqliteStore {
pub async fn connect() -> Result<Self, StoreError> {
let root = paths::wa_root();
std::fs::create_dir_all(&root)?;
let opts = SqliteConnectOptions::new()
.filename(paths::db_path())
.create_if_missing(true)
.foreign_keys(true);
let pool = SqlitePoolOptions::new()
.max_connections(5)
.connect_with(opts)
.await?;
sqlx::migrate!("./migrations")
.run(&pool)
.await
.map_err(|e| StoreError::Db(e.to_string()))?;
Ok(Self { pool })
}
}
fn now_unix() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
/// On-disk shape of `transcript.json` (`docs/03-data-model.md`).
#[derive(Debug, Serialize, Deserialize, Default)]
struct TranscriptFile {
schema: u32,
meeting_id: MeetingId,
language: Option<String>,
model: Option<String>,
backend: Option<String>,
segments: Vec<TranscriptSegment>,
speakers: Vec<TranscriptSpeaker>,
}
#[derive(Debug, Serialize, Deserialize)]
struct TranscriptSpeaker {
label: String,
display_name: Option<String>,
}
fn row_to_list_item(row: &sqlx::sqlite::SqliteRow) -> MeetingListItem {
MeetingListItem {
id: row.get("id"),
title: row.get("title"),
started_at: row.get("started_at"),
duration_secs: row.get("duration_secs"),
status: MeetingStatus::parse(row.get::<String, _>("status").as_str()),
}
}
#[async_trait]
impl Store for SqliteStore {
async fn create_meeting(&self, _m: NewMeeting) -> Result<MeetingId, StoreError> {
todo!("Phase 2 — create meeting row + folder")
async fn create_meeting(&self, m: NewMeeting) -> Result<MeetingId, StoreError> {
let id = uuid::Uuid::new_v4().to_string();
let folder = paths::meeting_dir(&id);
std::fs::create_dir_all(&folder)?;
let audio_path = folder.join("audio.wav");
let now = now_unix();
sqlx::query(
"INSERT INTO meetings (id, title, started_at, folder_path, audio_path, status, calendar_event_id, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, 'recording', ?, ?, ?)",
)
.bind(&id)
.bind(&m.title)
.bind(now)
.bind(folder.display().to_string())
.bind(audio_path.display().to_string())
.bind(&m.calendar_event_id)
.bind(now)
.bind(now)
.execute(&self.pool)
.await?;
Ok(id)
}
async fn finalize_meeting(&self, _id: &MeetingId) -> Result<(), StoreError> {
todo!("Phase 2 — persist transcript.json + metadata")
async fn finalize_meeting(&self, id: &MeetingId, s: FinalizeMeeting) -> Result<(), StoreError> {
let now = now_unix();
sqlx::query(
"UPDATE meetings SET status = 'ready', ended_at = ?, duration_secs = ?, recorded = ?,
language = ?, backend_used = ?, model_used = ?, updated_at = ? WHERE id = ?",
)
.bind(now)
.bind(s.duration_secs)
.bind(s.recorded as i64)
.bind(&s.language)
.bind(&s.backend_used)
.bind(&s.model_used)
.bind(now)
.bind(id)
.execute(&self.pool)
.await?;
for speaker in &s.speakers {
sqlx::query(
"INSERT INTO speakers (id, meeting_id, label, display_name)
VALUES (?, ?, ?, ?)
ON CONFLICT(meeting_id, label) DO UPDATE SET display_name = excluded.display_name",
)
.bind(uuid::Uuid::new_v4().to_string())
.bind(id)
.bind(&speaker.label)
.bind(&speaker.display_name)
.execute(&self.pool)
.await?;
}
let transcript = TranscriptFile {
schema: 1,
meeting_id: id.clone(),
language: s.language.clone(),
model: s.model_used.clone(),
backend: s.backend_used.clone(),
segments: s.segments,
speakers: s
.speakers
.iter()
.map(|sp| TranscriptSpeaker {
label: sp.label.clone(),
display_name: sp.display_name.clone(),
})
.collect(),
};
let json =
serde_json::to_string_pretty(&transcript).map_err(|e| StoreError::Db(e.to_string()))?;
std::fs::write(paths::meeting_dir(id).join("transcript.json"), json)?;
Ok(())
}
async fn list_meetings(
&self,
_query: Option<String>,
query: Option<String>,
) -> Result<Vec<MeetingListItem>, StoreError> {
todo!("Phase 2 — list meetings")
let rows = match query.filter(|q| !q.trim().is_empty()) {
Some(q) => {
let like = format!("%{}%", q.trim());
sqlx::query("SELECT id, title, started_at, duration_secs, status FROM meetings WHERE title LIKE ? ORDER BY started_at DESC")
.bind(like)
.fetch_all(&self.pool)
.await?
}
None => {
sqlx::query("SELECT id, title, started_at, duration_secs, status FROM meetings ORDER BY started_at DESC")
.fetch_all(&self.pool)
.await?
}
};
Ok(rows.iter().map(row_to_list_item).collect())
}
async fn delete_meeting(&self, _id: &MeetingId) -> Result<(), StoreError> {
todo!("Phase 2 — delete meeting + folder")
async fn get_meeting(&self, id: &MeetingId) -> Result<Meeting, StoreError> {
let row = sqlx::query(
"SELECT id, title, started_at, ended_at, duration_secs, status, recorded, language, backend_used, model_used
FROM meetings WHERE id = ?",
)
.bind(id)
.fetch_optional(&self.pool)
.await?
.ok_or_else(|| StoreError::NotFound(id.clone()))?;
let speaker_rows = sqlx::query(
"SELECT label, display_name FROM speakers WHERE meeting_id = ? ORDER BY label",
)
.bind(id)
.fetch_all(&self.pool)
.await?;
let speakers: Vec<SpeakerInfo> = speaker_rows
.iter()
.map(|r| SpeakerInfo {
label: r.get("label"),
display_name: r.get("display_name"),
participant_id: None,
})
.collect();
let folder = paths::meeting_dir(id);
let segments = std::fs::read_to_string(folder.join("transcript.json"))
.ok()
.and_then(|s| serde_json::from_str::<TranscriptFile>(&s).ok())
.map(|t| t.segments)
.unwrap_or_default();
let notes_markdown = std::fs::read_to_string(folder.join("notes.md")).unwrap_or_default();
Ok(Meeting {
id: row.get("id"),
title: row.get("title"),
started_at: row.get("started_at"),
ended_at: row.get("ended_at"),
duration_secs: row.get("duration_secs"),
status: MeetingStatus::parse(row.get::<String, _>("status").as_str()),
recorded: row.get::<i64, _>("recorded") != 0,
language: row.get("language"),
backend_used: row.get("backend_used"),
model_used: row.get("model_used"),
segments,
speakers,
notes_markdown,
})
}
async fn update_notes(&self, _id: &MeetingId, _markdown: &str) -> Result<(), StoreError> {
todo!("Phase 2 — write notes.md + FTS")
async fn delete_meeting(&self, id: &MeetingId) -> Result<(), StoreError> {
sqlx::query("DELETE FROM meetings WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
let folder = paths::meeting_dir(id);
if folder.exists() {
std::fs::remove_dir_all(folder)?;
}
Ok(())
}
async fn update_notes(&self, id: &MeetingId, markdown: &str) -> Result<(), StoreError> {
std::fs::write(paths::meeting_dir(id).join("notes.md"), markdown)?;
sqlx::query("UPDATE meetings SET updated_at = ? WHERE id = ?")
.bind(now_unix())
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn search(&self, _query: &str) -> Result<Vec<MeetingListItem>, StoreError> {
todo!("Phase 8 — FTS5 search")
}
async fn recover_scan(&self) -> Result<Vec<MeetingId>, StoreError> {
todo!("Phase 2 — recovery scan")
let rows = sqlx::query(
"SELECT id, audio_path FROM meetings WHERE status IN ('recording', 'transcribing')",
)
.fetch_all(&self.pool)
.await?;
let mut recovered = Vec::new();
for row in rows {
let id: MeetingId = row.get("id");
let audio_path: String = row.get("audio_path");
if PathBuf::from(&audio_path).exists() {
sqlx::query(
"UPDATE meetings SET status = 'recovering', updated_at = ? WHERE id = ?",
)
.bind(now_unix())
.bind(&id)
.execute(&self.pool)
.await?;
recovered.push(id);
} else {
// No working audio and never finalized: nothing to recover from.
sqlx::query("UPDATE meetings SET status = 'error', updated_at = ? WHERE id = ?")
.bind(now_unix())
.bind(&id)
.execute(&self.pool)
.await?;
}
}
Ok(recovered)
}
async fn enforce_retention(&self, _policy: Retention) -> Result<u32, StoreError> {
todo!("Phase 2 — retention enforcement")
async fn enforce_retention(&self, policy: Retention) -> Result<u32, StoreError> {
if policy.max_age_days.is_none() && policy.max_size_gb.is_none() {
return Ok(0);
}
let candidates = sqlx::query(
"SELECT id, folder_path, started_at FROM meetings WHERE status IN ('ready', 'error') ORDER BY started_at ASC",
)
.fetch_all(&self.pool)
.await?;
let mut removed = 0u32;
let now = now_unix();
let max_age_secs = policy.max_age_days.map(|d| d as i64 * 86_400);
// Age-based pruning first (cheap, no disk walk needed).
let mut kept = Vec::new();
for row in &candidates {
let id: MeetingId = row.get("id");
let started_at: i64 = row.get("started_at");
if let Some(max_age) = max_age_secs {
if now - started_at > max_age {
self.delete_meeting(&id).await?;
removed += 1;
continue;
}
}
kept.push((id, row.get::<String, _>("folder_path")));
}
// Size-based pruning, oldest first, until under the cap.
if let Some(max_gb) = policy.max_size_gb {
let cap_bytes = max_gb as u64 * 1_073_741_824;
let sizes: Vec<(MeetingId, u64)> = kept
.iter()
.map(|(id, folder)| (id.clone(), dir_size(folder)))
.collect();
let mut total: u64 = sizes.iter().map(|(_, s)| s).sum();
let mut idx = 0;
while total > cap_bytes && idx < sizes.len() {
let (id, size) = &sizes[idx];
self.delete_meeting(id).await?;
total = total.saturating_sub(*size);
removed += 1;
idx += 1;
}
}
Ok(removed)
}
}
fn dir_size(path: &str) -> u64 {
std::fs::read_dir(path)
.into_iter()
.flatten()
.flatten()
.filter_map(|entry| entry.metadata().ok())
.map(|m| m.len())
.sum()
}
+94 -51
View File
@@ -31,7 +31,8 @@ pub trait Transcriber: Send + Sync {
Self: Sized;
/// Streaming: emit interim + final segments for a window (FR-TRX-2).
fn transcribe_stream(&self, audio: AudioWindow, out: SegmentSink) -> Result<(), TrxError>;
/// Batch: one-shot, higher accuracy (FR-TRX-3).
/// Batch: one-shot over a whole file, higher accuracy — also the crash-recovery
/// path (FR-TRX-3, T2.8): re-run over the working `audio.wav` from scratch.
fn transcribe_file(&self, wav: &Path) -> Result<Vec<TranscriptSegment>, TrxError>;
}
@@ -42,6 +43,88 @@ pub struct WhisperTranscriber {
next_id: std::sync::atomic::AtomicU64,
}
#[cfg(feature = "cpu-transcription")]
impl WhisperTranscriber {
/// Shared decode path for both streaming windows and batch files.
///
/// `single_segment` forces whisper.cpp to treat the whole buffer as one
/// segment — needed for the short (~4s) streaming windows (see
/// `transcribe_stream`'s doc comment for why), but wrong for batch/file
/// transcription: a multi-minute recording should get whisper's normal
/// speech-boundary-driven multi-segment output, and its internal 30s-chunk
/// seek loop handles audio longer than one chunk on its own.
fn run_full(
&self,
samples: &[f32],
offset_ms: u64,
single_segment: bool,
) -> Result<Vec<TranscriptSegment>, TrxError> {
// whisper.cpp hangs (not errors) on an empty buffer — verified against
// a crash-recovered WAV whose header hadn't been finalized. Belt and
// suspenders alongside `read_wav_mono_16k` reading raw bytes directly.
if samples.is_empty() {
return Ok(Vec::new());
}
let mut state = self
.ctx
.create_state()
.map_err(|e| TrxError::Inference(e.to_string()))?;
// best_of: 1 — fastest greedy decode, matching the CPU/typical-laptop
// baseline (NFR-PERF-3).
let mut params =
whisper_rs::FullParams::new(whisper_rs::SamplingStrategy::Greedy { best_of: 1 });
params.set_n_threads(available_threads());
params.set_translate(false);
params.set_print_progress(false);
params.set_print_realtime(false);
params.set_print_timestamps(false);
params.set_suppress_blank(true);
params.set_single_segment(single_segment);
state
.full(params, samples)
.map_err(|e| TrxError::Inference(e.to_string()))?;
// A forced single_segment's reported end_timestamp() reflects
// whisper.cpp's internal 30s-padded mel frame, not the real window
// length — confirmed even with `duration_ms` set, so don't trust it.
// We already know the window's true length from `samples`, which is
// the more reliable source for `single_segment` windows; whisper's own
// multi-segment timestamps (batch/`transcribe_file`) stay authoritative.
let window_end_ms = offset_ms + (samples.len() as u64 * 1000) / 16_000;
let mut out = Vec::new();
for seg in state.as_iter() {
let text = seg.to_str_lossy().unwrap_or_default().trim().to_string();
if text.is_empty() {
continue;
}
// Whisper timestamps are centiseconds (10ms units).
let start_ms = offset_ms + seg.start_timestamp().max(0) as u64 * 10;
let end_ms = if single_segment {
window_end_ms.max(start_ms)
} else {
offset_ms + seg.end_timestamp().max(0) as u64 * 10
};
out.push(TranscriptSegment {
id: self
.next_id
.fetch_add(1, std::sync::atomic::Ordering::SeqCst),
start_ms,
end_ms,
// Diarization lands in Phase 4; every segment is provisionally "S1" until then.
speaker: "S1".to_string(),
text,
confidence: Some((1.0 - seg.no_speech_probability()).clamp(0.0, 1.0)),
interim: false,
});
}
Ok(out)
}
}
#[cfg(feature = "cpu-transcription")]
impl Transcriber for WhisperTranscriber {
fn load(model: &Path, _backend: BackendId) -> Result<Self, TrxError> {
@@ -57,55 +140,13 @@ impl Transcriber for WhisperTranscriber {
}
fn transcribe_stream(&self, audio: AudioWindow, out: SegmentSink) -> Result<(), TrxError> {
let mut state = self
.ctx
.create_state()
.map_err(|e| TrxError::Inference(e.to_string()))?;
// best_of: 1 — fastest greedy decode, matching the CPU/typical-laptop
// baseline (NFR-PERF-3); real-time windowed transcription favors low
// latency per window over the extra accuracy of a wider search.
let mut params =
whisper_rs::FullParams::new(whisper_rs::SamplingStrategy::Greedy { best_of: 1 });
params.set_n_threads(available_threads());
params.set_translate(false);
params.set_print_progress(false);
params.set_print_realtime(false);
params.set_print_timestamps(false);
params.set_suppress_blank(true);
// Each call already gets exactly one fixed-size window (not a full 30s
// clip); without this, whisper.cpp's segment-boundary heuristic hits
// "single timestamp ending" on short buffers and silently drops the
// whole window even when decoding succeeded (verified against real
// audio: text decoded correctly but zero segments came out until this
// was set — see whisper.cpp's own `stream` example, which sets this
// for the same reason).
params.set_single_segment(true);
state
.full(params, &audio.samples)
.map_err(|e| TrxError::Inference(e.to_string()))?;
for seg in state.as_iter() {
let text = seg.to_str_lossy().unwrap_or_default().trim().to_string();
if text.is_empty() {
continue;
}
// Whisper timestamps are centiseconds (10ms units).
let start_ms = audio.offset_ms + seg.start_timestamp().max(0) as u64 * 10;
let end_ms = audio.offset_ms + seg.end_timestamp().max(0) as u64 * 10;
let segment = TranscriptSegment {
id: self
.next_id
.fetch_add(1, std::sync::atomic::Ordering::SeqCst),
start_ms,
end_ms,
// Diarization lands in Phase 4; every segment is provisionally "S1" until then.
speaker: "S1".to_string(),
text,
confidence: Some((1.0 - seg.no_speech_probability()).clamp(0.0, 1.0)),
interim: false,
};
// clip); without `single_segment: true`, whisper.cpp's segment-boundary
// heuristic hits "single timestamp ending" on short buffers and drops
// the whole window even when decoding succeeded (verified against real
// audio — see whisper.cpp's own `stream` example, which sets this for
// the same reason).
for segment in self.run_full(&audio.samples, audio.offset_ms, true)? {
if out.send(segment).is_err() {
break; // receiver gone (meeting stopped) — nothing left to do
}
@@ -113,8 +154,10 @@ impl Transcriber for WhisperTranscriber {
Ok(())
}
fn transcribe_file(&self, _wav: &Path) -> Result<Vec<TranscriptSegment>, TrxError> {
todo!("Phase 3 — batch transcription")
fn transcribe_file(&self, wav: &Path) -> Result<Vec<TranscriptSegment>, TrxError> {
let samples =
crate::audio::read_wav_mono_16k(wav).map_err(|e| TrxError::Load(e.to_string()))?;
self.run_full(&samples, 0, false)
}
}