This commit is contained in:
iamdoubz
2026-06-30 22:17:30 -05:00
parent aae22a0e23
commit 02834e9cd2
38 changed files with 15517 additions and 203 deletions
+11 -2
View File
@@ -32,13 +32,22 @@ pub struct RunOutcome {
/// Spawns a local coding-agent CLI headless against a repo (claude -p / codex exec / opencode run / copilot).
#[async_trait]
pub trait AgentRunner: Send + Sync {
async fn run(&self, brief: &FeatureBrief, repo: &Path, out: LineSink) -> Result<RunOutcome, AgentError>;
async fn run(
&self,
brief: &FeatureBrief,
repo: &Path,
out: LineSink,
) -> Result<RunOutcome, AgentError>;
}
/// Creates a work item from a brief (e.g. a GitHub issue, optionally assigned to Copilot's cloud agent).
#[async_trait]
pub trait IssueTracker: Send + Sync {
async fn create_issue(&self, brief: &FeatureBrief, assign_copilot: bool) -> Result<String /* url */, AgentError>;
async fn create_issue(
&self,
brief: &FeatureBrief,
assign_copilot: bool,
) -> Result<String /* url */, AgentError>;
}
// Concrete impls (CliAgentRunner, GitHubIssueTracker) are added in Phase 10c behind a feature flag.
+379 -17
View File
@@ -1,13 +1,30 @@
//! Audio capture service — WASAPI loopback (Phase 1, FR-CAP-1/2).
//!
//! Design notes (see docs/02-architecture.md):
//! - WASAPI loopback CANNOT use event-callback mode, so capture POLLS on a
//! dedicated thread (research finding, docs/07).
//! - Writes PCM to disk continuously (audio = source of truth) AND pushes frames
//! into a bounded ring buffer consumed by the transcription worker.
//! - Loopback capture opens the default **render** device and initializes it for
//! `Direction::Capture` (see `initialize_client` below). The `wasapi` crate wires
//! `AUDCLNT_STREAMFLAGS_EVENTCALLBACK | AUDCLNT_STREAMFLAGS_LOOPBACK` together for
//! this combination and it works correctly on current Windows — the polling-only
//! claim in `docs/07-research-findings.md` was written before verifying against the
//! crate's actual source and doesn't hold for this dependency; the capture loop
//! below waits on the WASAPI event handle (with a short timeout so it can also
//! notice a stop/pause request) rather than busy-polling.
//! - Writes PCM to disk continuously, in the device's native format, so `audio.wav`
//! is a byte-accurate capture (audio = source of truth). A downmixed/resampled
//! 16kHz-mono copy is pushed into the bounded `FrameSink` for live transcription;
//! a full consumer drops preview frames (`try_send`) without blocking the disk
//! write or the capture thread.
//! - Does no inference itself.
use hound::{SampleFormat, WavSpec, WavWriter};
use std::fs::File;
use std::io::BufWriter;
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::SyncSender;
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use wasapi::{Direction, SampleType, ShareMode, WaveFormat};
#[derive(Debug, thiserror::Error)]
pub enum AudioError {
@@ -20,10 +37,15 @@ pub enum AudioError {
}
/// Opaque handle to a running capture, returned by `start` and consumed by `stop`.
pub struct CaptureHandle;
pub struct CaptureHandle {
running: Arc<AtomicBool>,
paused: Arc<AtomicBool>,
thread: JoinHandle<Result<CaptureSummary, AudioError>>,
}
/// Where captured frames are delivered for live transcription.
pub type FrameSink = std::sync::mpsc::Sender<Vec<f32>>;
/// Where captured frames are delivered for live transcription: mono f32 @ 16kHz,
/// bounded so a slow/absent consumer can never stall the capture thread.
pub type FrameSink = SyncSender<Vec<f32>>;
pub struct CaptureSummary {
pub duration_ms: u64,
@@ -44,18 +66,358 @@ pub struct WasapiCapture;
#[cfg(feature = "audio")]
impl AudioCapture for WasapiCapture {
fn start(&self, _wav_path: &Path, _sink: FrameSink) -> Result<CaptureHandle, AudioError> {
// T1.2: open default render device in loopback mode, spawn polling thread,
// write WAV via `hound`, push frames to `sink`.
todo!("Phase 1 — WASAPI loopback capture")
fn start(&self, wav_path: &Path, sink: FrameSink) -> Result<CaptureHandle, AudioError> {
let running = Arc::new(AtomicBool::new(true));
let paused = Arc::new(AtomicBool::new(false));
let running_th = running.clone();
let paused_th = paused.clone();
let wav_path = wav_path.to_path_buf();
let thread = thread::Builder::new()
.name("wa-audio-capture".into())
.spawn(move || capture_loop(&wav_path, &sink, &running_th, &paused_th))
.map_err(|e| AudioError::Capture(format!("spawn failed: {e}")))?;
Ok(CaptureHandle {
running,
paused,
thread,
})
}
fn pause(&self, _h: &CaptureHandle) -> Result<(), AudioError> {
todo!("Phase 1 — pause capture")
fn pause(&self, h: &CaptureHandle) -> Result<(), AudioError> {
h.paused.store(true, Ordering::SeqCst);
Ok(())
}
fn resume(&self, _h: &CaptureHandle) -> Result<(), AudioError> {
todo!("Phase 1 — resume capture")
fn resume(&self, h: &CaptureHandle) -> Result<(), AudioError> {
h.paused.store(false, Ordering::SeqCst);
Ok(())
}
fn stop(&self, _h: CaptureHandle) -> Result<CaptureSummary, AudioError> {
todo!("Phase 1 — stop capture, finalize WAV")
fn stop(&self, h: CaptureHandle) -> Result<CaptureSummary, AudioError> {
h.running.store(false, Ordering::SeqCst);
h.thread
.join()
.map_err(|_| AudioError::Capture("capture thread panicked".into()))?
}
}
/// Runs on a dedicated OS thread for the lifetime of a `CaptureHandle`. Owns the
/// WASAPI client and the WAV writer; exits (and finalizes the WAV) once `running`
/// is cleared.
#[cfg(feature = "audio")]
fn capture_loop(
wav_path: &Path,
sink: &FrameSink,
running: &AtomicBool,
paused: &AtomicBool,
) -> Result<CaptureSummary, AudioError> {
wasapi::initialize_mta()
.ok()
.map_err(|e| AudioError::Device(format!("COM init failed: {e}")))?;
let device = wasapi::get_default_device(&Direction::Render)
.map_err(|e| AudioError::Device(format!("no default render device: {e}")))?;
let mut audio_client = device
.get_iaudioclient()
.map_err(|e| AudioError::Device(format!("activate IAudioClient failed: {e}")))?;
let format = audio_client
.get_mixformat()
.map_err(|e| AudioError::Device(format!("GetMixFormat failed: {e}")))?;
let (_default_period, min_period) = audio_client
.get_periods()
.map_err(|e| AudioError::Device(format!("GetDevicePeriod failed: {e}")))?;
// `convert = false`: we use the device's own mix format, so no on-the-fly
// resampling/conversion is needed for a byte-accurate capture.
audio_client
.initialize_client(
&format,
min_period,
&Direction::Capture,
&ShareMode::Shared,
false,
)
.map_err(|e| AudioError::Capture(format!("Initialize failed: {e}")))?;
let event_handle = audio_client
.set_get_eventhandle()
.map_err(|e| AudioError::Capture(format!("SetEventHandle failed: {e}")))?;
let capture_client = audio_client
.get_audiocaptureclient()
.map_err(|e| AudioError::Capture(format!("GetService(IAudioCaptureClient) failed: {e}")))?;
let spec = wav_spec_for(&format)?;
let mut writer = WavWriter::create(wav_path, spec).map_err(|e| {
AudioError::Capture(format!("could not create {}: {e}", wav_path.display()))
})?;
audio_client
.start_stream()
.map_err(|e| AudioError::Capture(format!("start_stream failed: {e}")))?;
let mut resampler = Resampler::new(format.get_samplespersec());
let mut queue: std::collections::VecDeque<u8> = std::collections::VecDeque::new();
let mut frames_written: u64 = 0;
while running.load(Ordering::Relaxed) {
// Short timeout so we periodically re-check `running` even with no data.
let _ = event_handle.wait_for_event(100);
capture_client
.read_from_device_to_deque(&mut queue)
.map_err(|e| AudioError::Capture(format!("GetBuffer failed: {e}")))?;
if queue.is_empty() {
continue;
}
let bytes: Vec<u8> = queue.drain(..).collect();
// Must keep pulling WASAPI buffers even while paused (required to avoid
// overrun); only skip persisting/forwarding the audio.
if paused.load(Ordering::Relaxed) {
continue;
}
frames_written += write_wav_bytes(&mut writer, &bytes, &format)?;
let mono = decode_mono_f32(&bytes, &format)?;
let resampled = resampler.process(&mono);
if !resampled.is_empty() {
let _ = sink.try_send(resampled); // drop on backpressure; disk write is unaffected
}
}
audio_client.stop_stream().ok();
writer
.finalize()
.map_err(|e| AudioError::Capture(format!("wav finalize failed: {e}")))?;
let sample_rate = format.get_samplespersec();
Ok(CaptureSummary {
duration_ms: (frames_written * 1000) / sample_rate.max(1) as u64,
sample_rate,
channels: format.get_nchannels(),
})
}
fn wav_spec_for(format: &WaveFormat) -> Result<WavSpec, AudioError> {
let sample_format = match format
.get_subformat()
.map_err(|e| AudioError::Device(format!("unrecognized mix format: {e}")))?
{
SampleType::Float => SampleFormat::Float,
SampleType::Int => SampleFormat::Int,
};
Ok(WavSpec {
channels: format.get_nchannels(),
sample_rate: format.get_samplespersec(),
bits_per_sample: format.get_bitspersample(),
sample_format,
})
}
/// Write raw WASAPI capture bytes to the WAV writer in their native format.
/// Returns the number of frames written. Only the two mix formats WASAPI shared
/// mode actually produces in practice (32-bit float, 16-bit PCM) are supported;
/// anything else is a loud error rather than a silently corrupt recording.
fn write_wav_bytes(
writer: &mut WavWriter<BufWriter<File>>,
bytes: &[u8],
format: &WaveFormat,
) -> Result<u64, AudioError> {
let sample_type = format
.get_subformat()
.map_err(|e| AudioError::Device(format!("unrecognized mix format: {e}")))?;
let channels = format.get_nchannels() as u64;
match (sample_type, format.get_bitspersample()) {
(SampleType::Float, 32) => {
for chunk in bytes.chunks_exact(4) {
let v = f32::from_le_bytes(chunk.try_into().unwrap());
writer
.write_sample(v)
.map_err(|e| AudioError::Capture(format!("wav write: {e}")))?;
}
Ok(bytes.len() as u64 / 4 / channels.max(1))
}
(SampleType::Int, 16) => {
for chunk in bytes.chunks_exact(2) {
let v = i16::from_le_bytes(chunk.try_into().unwrap());
writer
.write_sample(v)
.map_err(|e| AudioError::Capture(format!("wav write: {e}")))?;
}
Ok(bytes.len() as u64 / 2 / channels.max(1))
}
(st, bits) => Err(AudioError::Device(format!(
"unsupported capture format: {st} {bits}-bit"
))),
}
}
/// Downmix raw WASAPI capture bytes to mono `f32` in `[-1.0, 1.0]`, at the
/// device's native sample rate (resampling to 16kHz happens separately).
fn decode_mono_f32(bytes: &[u8], format: &WaveFormat) -> Result<Vec<f32>, AudioError> {
let sample_type = format
.get_subformat()
.map_err(|e| AudioError::Device(format!("unrecognized mix format: {e}")))?;
let channels = format.get_nchannels() as usize;
if channels == 0 {
return Ok(Vec::new());
}
let mut mono = Vec::new();
match (sample_type, format.get_bitspersample()) {
(SampleType::Float, 32) => {
for frame in bytes.chunks_exact(4 * channels) {
let sum: f32 = frame
.chunks_exact(4)
.map(|c| f32::from_le_bytes(c.try_into().unwrap()))
.sum();
mono.push(sum / channels as f32);
}
}
(SampleType::Int, 16) => {
for frame in bytes.chunks_exact(2 * channels) {
let sum: f32 = frame
.chunks_exact(2)
.map(|c| i16::from_le_bytes(c.try_into().unwrap()) as f32 / i16::MAX as f32)
.sum();
mono.push(sum / channels as f32);
}
}
(st, bits) => {
return Err(AudioError::Device(format!(
"unsupported capture format: {st} {bits}-bit"
)))
}
}
Ok(mono)
}
const TARGET_SAMPLE_RATE: u32 = 16_000;
/// Streaming linear-interpolation resampler, mono f32 in -> mono f32 @ 16kHz out.
/// Carries fractional position and unconsumed tail samples across calls so
/// chunk boundaries don't introduce phase discontinuities.
/// `// ponytail: linear resampler, ceiling is STT-grade quality; swap for a sinc
/// resampler if transcription accuracy complaints trace back to aliasing.`
struct Resampler {
src_rate: u32,
pos: f64,
carry: Vec<f32>,
}
impl Resampler {
fn new(src_rate: u32) -> Self {
Self {
src_rate,
pos: 0.0,
carry: Vec::new(),
}
}
fn process(&mut self, mono_in: &[f32]) -> Vec<f32> {
if mono_in.is_empty() {
return Vec::new();
}
if self.src_rate == TARGET_SAMPLE_RATE {
return mono_in.to_vec();
}
let mut buf = std::mem::take(&mut self.carry);
buf.extend_from_slice(mono_in);
if buf.len() < 2 {
self.carry = buf;
return Vec::new();
}
let ratio = self.src_rate as f64 / TARGET_SAMPLE_RATE as f64;
let mut out = Vec::new();
while (self.pos as usize) + 1 < buf.len() {
let i = self.pos as usize;
let frac = (self.pos - i as f64) as f32;
out.push(buf[i] * (1.0 - frac) + buf[i + 1] * frac);
self.pos += ratio;
}
let keep_from = (self.pos as usize).min(buf.len() - 1);
self.pos -= keep_from as f64;
self.carry = buf[keep_from..].to_vec();
out
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn wav_writer_roundtrip_produces_valid_header_and_duration() {
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: 2,
sample_rate: 48_000,
bits_per_sample: 32,
sample_format: SampleFormat::Float,
};
let mut writer = WavWriter::create(&path, spec).unwrap();
let samples: Vec<f32> = (0..2000).map(|i| (i as f32 / 2000.0) - 0.5).collect();
let bytes: Vec<u8> = samples.iter().flat_map(|s| s.to_le_bytes()).collect();
// format.get_subformat()/get_bitspersample() need a real WaveFormat; build one
// via WaveFormat::new instead of a live device for this pure unit test.
let format = WaveFormat::new(32, 32, &SampleType::Float, 48_000, 2, None);
write_wav_bytes(&mut writer, &bytes, &format).unwrap();
writer.finalize().unwrap();
let reader = hound::WavReader::open(&path).unwrap();
let read_spec = reader.spec();
assert_eq!(read_spec.channels, 2);
assert_eq!(read_spec.sample_rate, 48_000);
assert_eq!(read_spec.bits_per_sample, 32);
let frame_count = reader.duration();
assert_eq!(frame_count as usize, samples.len() / 2);
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);
// One frame: left=1.0, right=-1.0 -> mono average 0.0.
let mut bytes = Vec::new();
bytes.extend_from_slice(&1.0f32.to_le_bytes());
bytes.extend_from_slice(&(-1.0f32).to_le_bytes());
let mono = decode_mono_f32(&bytes, &format).unwrap();
assert_eq!(mono, vec![0.0]);
}
#[test]
fn resampler_exact_integer_ratio_picks_aligned_samples() {
// 48kHz -> 16kHz is an exact 3:1 ratio, so linear interpolation should land
// exactly on the original samples at indices 0, 3, 6, ...
let mut r = Resampler::new(48_000);
let input: Vec<f32> = (0..30).map(|i| i as f32).collect();
let out = r.process(&input);
let expected: Vec<f32> = (0..10).map(|i| (i * 3) as f32).collect();
assert_eq!(out, expected);
}
#[test]
fn resampler_continues_correctly_across_chunk_boundaries() {
let mut r = Resampler::new(48_000);
let first: Vec<f32> = (0..30).map(|i| i as f32).collect();
let second: Vec<f32> = (30..60).map(|i| i as f32).collect();
let mut out = r.process(&first);
out.extend(r.process(&second));
let expected: Vec<f32> = (0..20).map(|i| (i * 3) as f32).collect();
assert_eq!(out, expected);
}
#[test]
fn resampler_passthrough_when_already_target_rate() {
let mut r = Resampler::new(16_000);
let input = vec![0.1f32, 0.2, 0.3];
assert_eq!(r.process(&input), input);
}
}
+292 -24
View File
@@ -2,12 +2,17 @@
//! Contract: `docs/04-api-contracts.md`. Commands return promptly; long work
//! is spawned and reported via events ("recording://*", "transcript://*", …).
//!
//! These are typed stubs. Each `todo!()` maps to a roadmap task and must be
//! replaced with a real implementation that depends on the service traits.
//! Recording/settings commands (Phase 1) are real implementations; everything
//! else is still a typed `todo!()` stub mapped to its roadmap task.
use crate::audio::{AudioCapture, WasapiCapture};
use crate::error::WaResult;
use crate::models::*;
use crate::transcription::{run_streaming_worker, Transcriber, WhisperTranscriber};
use crate::{error::WaError, AppState, RecordingSession};
use serde::Deserialize;
use std::path::PathBuf;
use tauri::{AppHandle, Emitter, State};
#[derive(Deserialize)]
pub struct StartRecordingArgs {
@@ -18,40 +23,273 @@ 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")
}
fn default_settings() -> Settings {
Settings {
theme: "system".into(),
storage_root: wa_root().display().to_string(),
llm_provider: "ollama".into(),
llm_endpoint: "http://localhost:11434".into(),
llm_model: "llama3".into(),
preferred_backend: "auto".into(),
low_overhead: false,
default_record: false,
consent_acknowledged: false,
sync_enabled: false,
}
}
fn load_settings() -> Settings {
std::fs::read_to_string(settings_path())
.ok()
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_else(default_settings)
}
fn save_settings(settings: &Settings) -> Result<(), WaError> {
let path = settings_path();
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| WaError::new("settings", e.to_string()))?;
}
let json = serde_json::to_string_pretty(settings)
.map_err(|e| WaError::new("settings", e.to_string()))?;
std::fs::write(path, json).map_err(|e| WaError::new("settings", e.to_string()))
}
// ---- Recording lifecycle (Phase 1) ----
#[tauri::command]
pub async fn start_recording(_args: StartRecordingArgs) -> WaResult<MeetingId> {
// T1.2/T1.3: create meeting row, start WASAPI capture, spawn transcription worker.
todo!("Phase 1 — start_recording")
pub async fn start_recording(
app: AppHandle,
state: State<'_, AppState>,
args: StartRecordingArgs,
) -> WaResult<MeetingId> {
let mut guard = state.session.lock().await;
if guard.is_some() {
return Err(WaError::new(
"recording",
"a meeting is already in progress",
));
}
if args.record && !load_settings().consent_acknowledged {
return Err(WaError::new(
"consent",
"recording consent has not been acknowledged yet",
));
}
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 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");
// 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);
let capture = WasapiCapture
.start(&wav_path, frame_tx)
.map_err(|e| WaError::new("audio", e.to_string()))?;
let app_for_worker = app.clone();
let meeting_id_for_worker = meeting_id.clone();
let transcription_worker = std::thread::Builder::new()
.name("wa-transcription".into())
.spawn(
move || match WhisperTranscriber::load(&model_path, BackendId::Cpu) {
Ok(transcriber) => run_streaming_worker(&transcriber, frame_rx, |segment| {
let _ = app_for_worker.emit(
"transcript://segment",
serde_json::json!({ "meetingId": meeting_id_for_worker, "segment": segment }),
);
}),
Err(e) => tracing::error!("failed to load whisper model: {e}"),
},
)
.map_err(|e| WaError::new("transcription", e.to_string()))?;
*guard = Some(RecordingSession {
meeting_id: meeting_id.clone(),
capture,
retention: args.record,
wav_path,
started_at: std::time::Instant::now(),
transcription_worker,
});
drop(guard);
crate::update_tray_tooltip(&app, "WhispAssist — recording");
let _ = app.emit(
"recording://state",
serde_json::json!({ "meetingId": meeting_id, "state": "recording", "elapsedMs": 0 }),
);
Ok(meeting_id)
}
#[tauri::command]
pub async fn stop_recording(_meeting_id: MeetingId) -> WaResult<()> {
// T1.3 + T4: finalize WAV, run diarization, persist, optional summary.
todo!("Phase 1 — stop_recording")
pub async fn stop_recording(
app: AppHandle,
state: State<'_, AppState>,
meeting_id: MeetingId,
) -> WaResult<()> {
let mut guard = state.session.lock().await;
let session = match guard.take() {
Some(s) if s.meeting_id == meeting_id => s,
Some(s) => {
*guard = Some(s);
return Err(WaError::new(
"recording",
"meeting_id does not match the active recording",
));
}
None => {
return Err(WaError::new(
"recording",
"no meeting is currently recording",
))
}
};
drop(guard);
let summary = WasapiCapture
.stop(session.capture)
.map_err(|e| WaError::new("audio", e.to_string()))?;
// The capture thread has dropped its `FrameSink`; the transcription worker's
// `recv()` now returns `Err` and the worker exits on its own — join to
// guarantee it has fully drained the audio before we act on retention.
let _ = session.transcription_worker.join();
let _ = app.emit(
"transcript://finalized",
serde_json::json!({ "meetingId": meeting_id, "segmentCount": serde_json::Value::Null }),
);
// ADR-0009: delete the working WAV only after the transcript worker has
// finished reading it, and only when retention is off.
if !session.retention {
let _ = std::fs::remove_file(&session.wav_path);
}
crate::update_tray_tooltip(&app, "WhispAssist — idle");
let _ = app.emit(
"recording://state",
serde_json::json!({ "meetingId": meeting_id, "state": "stopped", "elapsedMs": summary.duration_ms }),
);
Ok(())
}
#[tauri::command]
pub async fn pause_recording(_meeting_id: MeetingId) -> WaResult<()> {
todo!("Phase 1 — pause_recording")
pub async fn pause_recording(
app: AppHandle,
state: State<'_, AppState>,
meeting_id: MeetingId,
) -> WaResult<()> {
let guard = state.session.lock().await;
let session = guard
.as_ref()
.filter(|s| s.meeting_id == meeting_id)
.ok_or_else(|| WaError::new("recording", "no matching active recording"))?;
WasapiCapture
.pause(&session.capture)
.map_err(|e| WaError::new("audio", e.to_string()))?;
drop(guard);
let _ = app.emit(
"recording://state",
serde_json::json!({ "meetingId": meeting_id, "state": "paused", "elapsedMs": 0 }),
);
Ok(())
}
#[tauri::command]
pub async fn resume_recording(_meeting_id: MeetingId) -> WaResult<()> {
todo!("Phase 1 — resume_recording")
pub async fn resume_recording(
app: AppHandle,
state: State<'_, AppState>,
meeting_id: MeetingId,
) -> WaResult<()> {
let guard = state.session.lock().await;
let session = guard
.as_ref()
.filter(|s| s.meeting_id == meeting_id)
.ok_or_else(|| WaError::new("recording", "no matching active recording"))?;
WasapiCapture
.resume(&session.capture)
.map_err(|e| WaError::new("audio", e.to_string()))?;
drop(guard);
let _ = app.emit(
"recording://state",
serde_json::json!({ "meetingId": meeting_id, "state": "recording", "elapsedMs": 0 }),
);
Ok(())
}
/// Toggle audio retention mid-meeting (ADR-0009, FR-REC-1).
#[tauri::command]
pub async fn set_recording_retention(_meeting_id: MeetingId, _record: bool) -> WaResult<()> {
todo!("Phase 1 — set_recording_retention")
pub async fn set_recording_retention(
app: AppHandle,
state: State<'_, AppState>,
meeting_id: MeetingId,
record: bool,
) -> WaResult<()> {
if record && !load_settings().consent_acknowledged {
return Err(WaError::new(
"consent",
"recording consent has not been acknowledged yet",
));
}
let mut guard = state.session.lock().await;
let session = guard
.as_mut()
.filter(|s| s.meeting_id == meeting_id)
.ok_or_else(|| WaError::new("recording", "no matching active recording"))?;
session.retention = record;
drop(guard);
let _ = app.emit(
"recording://retention",
serde_json::json!({ "meetingId": meeting_id, "record": record }),
);
Ok(())
}
/// Record the one-time recording-consent acknowledgment (FR-REC-2).
#[tauri::command]
pub async fn acknowledge_recording_consent() -> WaResult<()> {
todo!("Phase 1 — acknowledge_recording_consent")
let mut settings = load_settings();
settings.consent_acknowledged = true;
save_settings(&settings)
}
// ---- Hardware (Phase 3) ----
@@ -107,7 +345,10 @@ pub async fn set_llm_provider(_config: serde_json::Value) -> WaResult<()> {
}
#[tauri::command]
pub async fn generate_summary(_meeting_id: MeetingId, _template_id: Option<String>) -> WaResult<()> {
pub async fn generate_summary(
_meeting_id: MeetingId,
_template_id: Option<String>,
) -> WaResult<()> {
// Streams via "llm://token" / "llm://done".
todo!("Phase 5 — generate_summary")
}
@@ -179,13 +420,18 @@ pub async fn sync_status(_meeting_id: Option<MeetingId>) -> WaResult<Vec<SyncJob
// ---- Feature briefs + MCP server (Phase 10b, ADR-0011) ----
#[tauri::command]
pub async fn create_feature_brief(_meeting_id: MeetingId, _target_repo: Option<String>) -> WaResult<FeatureBrief> {
pub async fn create_feature_brief(
_meeting_id: MeetingId,
_target_repo: Option<String>,
) -> WaResult<FeatureBrief> {
// T10.6: distill transcript → structured brief via the configured LlmProvider.
todo!("Phase 10b — create_feature_brief")
}
#[tauri::command]
pub async fn list_feature_briefs(_meeting_id: Option<MeetingId>) -> WaResult<Vec<FeatureBriefInfo>> {
pub async fn list_feature_briefs(
_meeting_id: Option<MeetingId>,
) -> WaResult<Vec<FeatureBriefInfo>> {
todo!("Phase 10b — list_feature_briefs")
}
@@ -207,7 +453,11 @@ pub async fn mcp_status() -> WaResult<serde_json::Value> {
/// Enable/disable the loopback MCP server; returns endpoint + token on enable (FR-MCP-1/6).
#[tauri::command]
pub async fn set_mcp_enabled(_enabled: bool, _transport: Option<String>, _port: Option<u16>) -> WaResult<serde_json::Value> {
pub async fn set_mcp_enabled(
_enabled: bool,
_transport: Option<String>,
_port: Option<u16>,
) -> WaResult<serde_json::Value> {
todo!("Phase 10b — set_mcp_enabled (loopback only, token)")
}
@@ -225,12 +475,20 @@ pub async fn mcp_access_log(_limit: Option<u32>) -> WaResult<Vec<McpAccessEntry>
// ---- Agent push / task-tracker handoff (Phase 10c, later) ----
#[tauri::command]
pub async fn run_agent(_brief_id: String, _tool: String, _repo_path: String) -> WaResult<serde_json::Value> {
pub async fn run_agent(
_brief_id: String,
_tool: String,
_repo_path: String,
) -> WaResult<serde_json::Value> {
todo!("Phase 10c — run_agent (claude|codex|opencode|copilot)")
}
#[tauri::command]
pub async fn create_issue_from_brief(_brief_id: String, _tracker: String, _assign_copilot: Option<bool>) -> WaResult<serde_json::Value> {
pub async fn create_issue_from_brief(
_brief_id: String,
_tracker: String,
_assign_copilot: Option<bool>,
) -> WaResult<serde_json::Value> {
todo!("Phase 10c — create_issue_from_brief")
}
@@ -238,12 +496,22 @@ pub async fn create_issue_from_brief(_brief_id: String, _tracker: String, _assig
#[tauri::command]
pub async fn get_settings() -> WaResult<Settings> {
todo!("Phase 2 — get_settings")
Ok(load_settings())
}
#[tauri::command]
pub async fn update_settings(_patch: serde_json::Value) -> WaResult<Settings> {
todo!("Phase 2 — update_settings")
pub async fn update_settings(patch: serde_json::Value) -> WaResult<Settings> {
let mut current = serde_json::to_value(load_settings())
.map_err(|e| WaError::new("settings", e.to_string()))?;
if let (Some(base), Some(patch)) = (current.as_object_mut(), patch.as_object()) {
for (key, value) in patch {
base.insert(key.clone(), value.clone());
}
}
let merged: Settings =
serde_json::from_value(current).map_err(|e| WaError::new("settings", e.to_string()))?;
save_settings(&merged)?;
Ok(merged)
}
/// Reports current egress + LLM endpoint so the UI can prove local-only handling (FR-SEC-2).
+4 -1
View File
@@ -13,7 +13,10 @@ pub struct WaError {
impl WaError {
pub fn new(kind: impl Into<String>, message: impl Into<String>) -> Self {
Self { kind: kind.into(), message: message.into() }
Self {
kind: kind.into(),
message: message.into(),
}
}
}
+39 -10
View File
@@ -19,17 +19,20 @@ pub mod storage;
pub mod sync;
pub mod transcription;
use std::sync::Arc;
use std::path::PathBuf;
use std::thread::JoinHandle;
use tauri::tray::TrayIcon;
use tauri::Manager;
use tokio::sync::Mutex;
/// Shared application state handed to every Tauri command via `tauri::State`.
/// Service handles are `Arc`-wrapped so they can be cloned into worker tasks.
///
/// 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.
pub struct AppState {
pub hardware: Arc<dyn hardware::HardwareDetector>,
pub store: Arc<dyn storage::Store>,
pub llm: Arc<dyn llm::LlmProvider>,
pub sync: Arc<dyn sync::SyncManager>,
pub mcp: Arc<dyn mcp::McpServer>,
/// The single in-flight recording, if any. Guarded so start/stop/pause are atomic.
pub session: Mutex<Option<RecordingSession>>,
}
@@ -37,16 +40,35 @@ pub struct AppState {
/// Tracks the currently-recording meeting and its worker handles.
pub struct RecordingSession {
pub meeting_id: models::MeetingId,
// capture handle, transcription worker join handle, etc. — populated in Phase 1.
pub capture: audio::CaptureHandle,
/// Audio retention for this meeting (ADR-0009); toggle-able mid-meeting.
pub retention: bool,
pub wav_path: PathBuf,
pub started_at: std::time::Instant,
pub transcription_worker: JoinHandle<()>,
}
/// Wraps the tray icon so it can be looked up from commands to update its
/// tooltip on recording state changes (FR-CAP-4).
pub struct TrayHandle(pub TrayIcon);
/// Build state, register commands/events, and run the app.
pub fn run() {
tracing_subscriber::fmt().with_env_filter("info").init();
// NOTE: concrete service implementations are constructed here once they exist.
// The skeleton wires the command surface; `todo!()`s mark per-phase tasks.
tauri::Builder::default()
.manage(AppState {
session: Mutex::new(None),
})
.setup(|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));
Ok(())
})
.invoke_handler(tauri::generate_handler![
commands::start_recording,
commands::stop_recording,
@@ -91,3 +113,10 @@ pub fn run() {
.run(tauri::generate_context!())
.expect("error while running WhispAssist");
}
/// Used by `commands.rs` to keep the tray tooltip honest about capture state (FR-CAP-4).
pub(crate) fn update_tray_tooltip(app: &tauri::AppHandle, text: &str) {
if let Some(tray) = app.try_state::<TrayHandle>() {
let _ = tray.0.set_tooltip(Some(text));
}
}
+21 -5
View File
@@ -76,7 +76,11 @@ pub trait McpServer: Send + Sync {
/// Distills a transcript into an agent-ready spec (FR-MCP-4) using the configured LlmProvider.
#[async_trait]
pub trait FeatureBriefBuilder: Send + Sync {
async fn build(&self, meeting_id: &MeetingId, target_repo: Option<&str>) -> Result<FeatureBrief, BriefError>;
async fn build(
&self,
meeting_id: &MeetingId,
target_repo: Option<&str>,
) -> Result<FeatureBrief, BriefError>;
}
/// Default rmcp-backed server (feature `mcp`).
@@ -97,10 +101,22 @@ impl McpServer for RmcpServer {
fn tools(&self) -> Vec<McpToolDescriptor> {
// T10.5: the tools an agent can call. Tools-first for Copilot compatibility.
vec![
McpToolDescriptor { name: "list_recent_meetings", description: "Recent meetings (scoped)." },
McpToolDescriptor { name: "get_transcript", description: "Transcript for a meeting (scoped)." },
McpToolDescriptor { name: "get_action_items", description: "Action items for a meeting." },
McpToolDescriptor { name: "get_feature_brief", description: "Agent-ready spec distilled from a meeting." },
McpToolDescriptor {
name: "list_recent_meetings",
description: "Recent meetings (scoped).",
},
McpToolDescriptor {
name: "get_transcript",
description: "Transcript for a meeting (scoped).",
},
McpToolDescriptor {
name: "get_action_items",
description: "Action items for a meeting.",
},
McpToolDescriptor {
name: "get_feature_brief",
description: "Agent-ready spec distilled from a meeting.",
},
]
}
}
+3 -3
View File
@@ -49,7 +49,7 @@ pub struct TranscriptSegment {
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SpeakerInfo {
pub label: String, // "S1"
pub label: String, // "S1"
pub display_name: Option<String>,
pub participant_id: Option<String>,
}
@@ -100,9 +100,9 @@ pub struct Participant {
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Settings {
pub theme: String, // system|light|dark
pub theme: String, // system|light|dark
pub storage_root: String,
pub llm_provider: String, // ollama|custom|off
pub llm_provider: String, // ollama|custom|off
pub llm_endpoint: String,
pub llm_model: String,
pub preferred_backend: String, // auto|npu|nvidia|amd|intel|cpu
+8 -2
View File
@@ -33,7 +33,8 @@ pub trait NotesRenderer: Send + Sync {
summary_md: Option<&str>,
) -> String;
fn export(&self, markdown: &str, dest: &Path, fmt: ExportFormat) -> Result<PathBuf, NotesError>;
fn export(&self, markdown: &str, dest: &Path, fmt: ExportFormat)
-> Result<PathBuf, NotesError>;
}
pub struct MarkdownNotes;
@@ -48,7 +49,12 @@ impl NotesRenderer for MarkdownNotes {
// T2.4: resolve speaker names, group dialogue, prepend summary section.
todo!("Phase 2 — assemble Markdown")
}
fn export(&self, _markdown: &str, _dest: &Path, _fmt: ExportFormat) -> Result<PathBuf, NotesError> {
fn export(
&self,
_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")
}
+8 -2
View File
@@ -31,7 +31,10 @@ pub struct Retention {
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 list_meetings(&self, query: Option<String>) -> Result<Vec<MeetingListItem>, StoreError>;
async fn list_meetings(
&self,
query: Option<String>,
) -> Result<Vec<MeetingListItem>, 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).
@@ -53,7 +56,10 @@ impl Store for SqliteStore {
async fn finalize_meeting(&self, _id: &MeetingId) -> Result<(), StoreError> {
todo!("Phase 2 — persist transcript.json + metadata")
}
async fn list_meetings(&self, _query: Option<String>) -> Result<Vec<MeetingListItem>, StoreError> {
async fn list_meetings(
&self,
_query: Option<String>,
) -> Result<Vec<MeetingListItem>, StoreError> {
todo!("Phase 2 — list meetings")
}
async fn delete_meeting(&self, _id: &MeetingId) -> Result<(), StoreError> {
+20 -6
View File
@@ -41,7 +41,12 @@ pub trait SyncTarget: Send + Sync {
/// True if a file with this hash already exists remotely (skip-if-unchanged).
async fn exists(&self, remote_path: &str, sha256: &str) -> Result<bool, SyncError>;
/// Upload a file; resumable/chunked for large artifacts.
async fn put(&self, local: &Path, remote_path: &str, prog: ProgressSink) -> Result<(), SyncError>;
async fn put(
&self,
local: &Path,
remote_path: &str,
prog: ProgressSink,
) -> Result<(), SyncError>;
}
/// Owns the durable queue, retry/backoff, credential resolution, TLS enforcement,
@@ -49,7 +54,11 @@ pub trait SyncTarget: Send + Sync {
#[async_trait]
pub trait SyncManager: Send + Sync {
/// Enqueue a meeting's selected artifacts for one or all enabled targets.
async fn enqueue_meeting(&self, meeting_id: &MeetingId, target_id: Option<&str>) -> Result<(), SyncError>;
async fn enqueue_meeting(
&self,
meeting_id: &MeetingId,
target_id: Option<&str>,
) -> Result<(), SyncError>;
/// Drive pending jobs (called on finalize, on startup, and on a low-frequency timer).
async fn pump(&self) -> Result<(), SyncError>;
async fn status(&self, meeting_id: Option<&MeetingId>) -> Result<Vec<SyncJobInfo>, SyncError>;
@@ -61,11 +70,11 @@ pub trait SyncManager: Send + Sync {
/// WebDAV provider — primary targets + Synology (feature `sync`).
#[cfg(feature = "sync")]
pub struct WebDavTarget {
pub base_url: String, // https://host/remote.php/dav/files/<user>/ , /seafdav , /dav , …
pub base_url: String, // https://host/remote.php/dav/files/<user>/ , /seafdav , /dav , …
pub remote_base_path: String,
pub username: String,
pub credential_ref: String, // key into the OS credential store — resolved at use, not stored here
pub third_party: bool, // false for self-hosted
pub credential_ref: String, // key into the OS credential store — resolved at use, not stored here
pub third_party: bool, // false for self-hosted
pub allow_plaintext_lan: bool,
}
@@ -89,7 +98,12 @@ impl SyncTarget for WebDavTarget {
async fn exists(&self, _remote_path: &str, _sha256: &str) -> Result<bool, SyncError> {
todo!("Phase 9 — WebDAV exists / hash compare")
}
async fn put(&self, _local: &Path, _remote_path: &str, _prog: ProgressSink) -> Result<(), SyncError> {
async fn put(
&self,
_local: &Path,
_remote_path: &str,
_prog: ProgressSink,
) -> Result<(), SyncError> {
// T9.2/T9.5: PUT (chunked for large files); stream progress.
todo!("Phase 9 — WebDAV upload")
}
+137 -7
View File
@@ -6,6 +6,7 @@
use crate::models::{BackendId, TranscriptSegment};
use std::path::Path;
use std::sync::mpsc::{Receiver, Sender};
#[derive(Debug, thiserror::Error)]
pub enum TrxError {
@@ -22,7 +23,7 @@ pub struct AudioWindow {
}
/// Where produced segments are delivered (interim then final).
pub type SegmentSink = std::sync::mpsc::Sender<TranscriptSegment>;
pub type SegmentSink = Sender<TranscriptSegment>;
pub trait Transcriber: Send + Sync {
fn load(model: &Path, backend: BackendId) -> Result<Self, TrxError>
@@ -36,22 +37,151 @@ pub trait Transcriber: Send + Sync {
/// whisper.cpp-backed transcriber (CPU baseline; GPU via Cargo features).
#[cfg(feature = "cpu-transcription")]
pub struct WhisperTranscriber;
pub struct WhisperTranscriber {
ctx: whisper_rs::WhisperContext,
next_id: std::sync::atomic::AtomicU64,
}
#[cfg(feature = "cpu-transcription")]
impl Transcriber for WhisperTranscriber {
fn load(_model: &Path, _backend: BackendId) -> Result<Self, TrxError> {
// T1.5 / T3.3: init whisper-rs with the backend's acceleration features.
todo!("Phase 1 — load whisper model")
fn load(model: &Path, _backend: BackendId) -> Result<Self, TrxError> {
let ctx = whisper_rs::WhisperContext::new_with_params(
model,
whisper_rs::WhisperContextParameters::default(),
)
.map_err(|e| TrxError::Load(e.to_string()))?;
Ok(Self {
ctx,
next_id: std::sync::atomic::AtomicU64::new(0),
})
}
fn transcribe_stream(&self, _audio: AudioWindow, _out: SegmentSink) -> Result<(), TrxError> {
todo!("Phase 1 — streaming transcription")
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,
};
if out.send(segment).is_err() {
break; // receiver gone (meeting stopped) — nothing left to do
}
}
Ok(())
}
fn transcribe_file(&self, _wav: &Path) -> Result<Vec<TranscriptSegment>, TrxError> {
todo!("Phase 3 — batch transcription")
}
}
#[cfg(feature = "cpu-transcription")]
fn available_threads() -> std::ffi::c_int {
std::thread::available_parallelism()
.map(|n| n.get() as std::ffi::c_int)
.unwrap_or(4)
}
/// ONNX Runtime + DirectML transcriber for the NPU tier (Phase 3).
#[cfg(feature = "directml")]
pub struct OnnxNpuTranscriber;
/// Streaming window worker (Phase 1, T1.5/T1.6): accumulates raw 16kHz-mono
/// chunks from the `audio` service into fixed-size, **non-overlapping** windows
/// and runs one `transcribe_stream` pass per window as it fills, forwarding
/// each produced segment to `on_segment` (e.g. a Tauri event emit).
///
/// `transcribe_stream` takes a `SegmentSink` per the `Transcriber` trait (so a
/// future async/threaded engine can push mid-inference), but whisper.cpp's
/// `full()` call is synchronous — by the time a window's `transcribe_stream`
/// call returns, every segment it produced is already sitting in a fresh
/// per-window channel, so this drains it inline rather than needing a second
/// long-lived thread just to bridge segments out.
///
/// True incremental/partial-word streaming (and window overlap for continuity)
/// are out of scope for Phase 1: whisper.cpp transcribes each window from
/// scratch, so an overlapping window would re-emit the overlapped words a
/// second time with no stitching logic to merge them — a worse rough edge for
/// a live transcript than the occasional word clipped at a window boundary.
/// "Near real time" (FR-TRX-2) is met by short (~4s) windows; `interim` stays
/// `false` for every segment produced here.
#[cfg(feature = "cpu-transcription")]
pub fn run_streaming_worker<T, F>(transcriber: &T, frame_rx: Receiver<Vec<f32>>, mut on_segment: F)
where
T: Transcriber,
F: FnMut(TranscriptSegment),
{
const SAMPLE_RATE: usize = 16_000;
const WINDOW_SECS: f32 = 4.0;
let window_len = (WINDOW_SECS * SAMPLE_RATE as f32) as usize;
let mut buf: Vec<f32> = Vec::new();
let mut offset_ms: u64 = 0;
let run_window = |transcriber: &T, samples: Vec<f32>, offset_ms: u64, on_segment: &mut F| {
let (tx, rx) = std::sync::mpsc::channel();
let window = AudioWindow { samples, offset_ms };
if let Err(e) = transcriber.transcribe_stream(window, tx) {
tracing::warn!("transcription window failed: {e}");
}
while let Ok(segment) = rx.try_recv() {
on_segment(segment);
}
};
while let Ok(chunk) = frame_rx.recv() {
buf.extend_from_slice(&chunk);
while buf.len() >= window_len {
let samples: Vec<f32> = buf.drain(..window_len).collect();
run_window(transcriber, samples, offset_ms, &mut on_segment);
offset_ms += (window_len as u64 * 1000) / SAMPLE_RATE as u64;
}
}
// Final partial window on stop, if there's enough audio to be worth a pass.
if buf.len() > SAMPLE_RATE / 2 {
run_window(transcriber, buf, offset_ms, &mut on_segment);
}
}