- 1
//! The voice surface (docs/design/49-live-voice.md): `/voice/speak`, - 2
//! `/voice/providers`, `WS /voice/session`, and the transcription of voice - 3
//! notes the gateway admits from channels. - 4
//! - 5
//! Every route resolves its provider through [`vak_voice::VoiceProvider`] and - 6
//! its credentials through the canonical secret chain, and every paid call - 7
//! draws on one shared per-minute [`RequestWindow`]. Transcription has one - 8
//! implementation, [`transcribe`], used by the socket and by channel voice - 9
//! notes alike. A channel note is transcribed only after allowlist admission, - 10
//! through its chat's bot → chat voice tiers ([`narrowed`]). - 11
- 12
use crate::AppState; - 13
use axum::{ - 14
Json, - 15
extract::{ - 16
Query, State, - 17
ws::{Message, WebSocket, WebSocketUpgrade}, - 18
}, - 19
http::{HeaderValue, StatusCode, header}, - 20
response::{IntoResponse, Response}, - 21
}; - 22
use futures::SinkExt; - 23
use serde::Deserialize; - 24
use std::collections::HashSet; - 25
use std::sync::Arc; - 26
use std::sync::Mutex; - 27
use std::sync::atomic::{AtomicUsize, Ordering}; - 28
use std::time::{Duration, Instant}; - 29
use tokio_util::sync::CancellationToken; - 30
use vak_voice::protocol::{ClientControl, DiscardReason, ServerControl, VOICE_PROTOCOL_VERSION}; - 31
use vak_voice::{SpeakFormat, VoiceError, VoiceProvider, audio, vad::SpeechEvidence}; - 32
- 33
/// Rolling one-minute budget of paid voice requests for this process. - 34
pub(crate) struct RequestWindow { - 35
window: Mutex<(Instant, usize)>, - 36
} - 37
- 38
impl RequestWindow { - 39
pub(crate) fn new() -> Self { - 40
Self { - 41
window: Mutex::new((Instant::now(), 0)), - 42
} - 43
} - 44
- 45
pub(crate) fn admit(&self, limit: usize) -> bool { - 46
self.admit_at(limit, Instant::now()) - 47
} - 48
- 49
fn admit_at(&self, limit: usize, now: Instant) -> bool { - 50
let mut window = self - 51
.window - 52
.lock() - 53
.unwrap_or_else(std::sync::PoisonError::into_inner); - 54
if now.saturating_duration_since(window.0) >= Duration::from_secs(60) { - 55
*window = (now, 0); - 56
} - 57
if window.1 >= limit { - 58
return false; - 59
} - 60
window.1 += 1; - 61
true - 62
} - 63
} - 64
- 65
/// Why a voice route could not serve a request. - 66
#[derive(Debug)] - 67
pub(crate) enum RouteError { - 68
Disabled, - 69
Config(String), - 70
NoCredential(VoiceProvider), - 71
Provider(vak_llm::LlmError), - 72
Local(VoiceError), - 73
} - 74
- 75
impl RouteError { - 76
fn status(&self) -> StatusCode { - 77
match self { - 78
Self::Disabled => StatusCode::CONFLICT, - 79
Self::Config(_) => StatusCode::BAD_REQUEST, - 80
Self::NoCredential(_) => StatusCode::SERVICE_UNAVAILABLE, - 81
Self::Provider(error) => match error { - 82
vak_llm::LlmError::Auth(_) | vak_llm::LlmError::InvalidRequest(_) => { - 83
StatusCode::BAD_REQUEST - 84
} - 85
vak_llm::LlmError::RateLimit { .. } => StatusCode::TOO_MANY_REQUESTS, - 86
vak_llm::LlmError::Overloaded(_) => StatusCode::SERVICE_UNAVAILABLE, - 87
_ => StatusCode::BAD_GATEWAY, - 88
}, - 89
Self::Local(VoiceError::InvalidRequest(_)) => StatusCode::BAD_REQUEST, - 90
Self::Local(VoiceError::Unavailable(_)) => StatusCode::SERVICE_UNAVAILABLE, - 91
Self::Local(_) => StatusCode::BAD_GATEWAY, - 92
} - 93
} - 94
- 95
fn message(&self) -> String { - 96
match self { - 97
Self::Disabled => "Voice is disabled in workspace settings".into(), - 98
Self::Config(message) => message.clone(), - 99
Self::NoCredential(provider) => format!( - 100
"No credential is configured for the {provider} voice provider ({})", - 101
provider.credential_vars().join(" or ") - 102
), - 103
Self::Provider(error) => error.to_string(), - 104
Self::Local(error) => error.to_string(), - 105
} - 106
} - 107
- 108
fn receipt_outcome(&self) -> (vak_llm::FailureDomain, vak_llm::Settlement) { - 109
match self { - 110
Self::Provider(error) => vak_llm::work::classify_error(error), - 111
Self::Local(VoiceError::Cancelled) => ( - 112
vak_llm::FailureDomain::Unknown, - 113
vak_llm::Settlement::Cancelled, - 114
), - 115
_ => (vak_llm::FailureDomain::Request, vak_llm::Settlement::Failed), - 116
} - 117
} - 118
} - 119
- 120
impl IntoResponse for RouteError { - 121
fn into_response(self) -> Response { - 122
( - 123
self.status(), - 124
Json(serde_json::json!({ "error": self.message() })), - 125
) - 126
.into_response() - 127
} - 128
} - 129
- 130
fn error_response(status: StatusCode, message: impl Into<String>) -> Response { - 131
(status, Json(serde_json::json!({ "error": message.into() }))).into_response() - 132
} - 133
- 134
fn route(settings: &vak_config::VoiceSettings) -> Result<VoiceProvider, RouteError> { - 135
if !settings.enabled { - 136
return Err(RouteError::Disabled); - 137
} - 138
VoiceProvider::resolve(settings.provider.as_deref()).map_err(RouteError::Config) - 139
} - 140
- 141
/// The operator-installed executable for a local engine, if configured. - 142
fn local_engine(var: &str) -> Option<std::path::PathBuf> { - 143
vak_config::get_var(var) - 144
.filter(|path| !path.trim().is_empty()) - 145
.map(std::path::PathBuf::from) - 146
} - 147
- 148
fn credential(provider: VoiceProvider) -> Result<String, RouteError> { - 149
provider - 150
.credential_vars() - 151
.iter() - 152
.find_map(|name| vak_config::get_var(name).filter(|value| !value.trim().is_empty())) - 153
.map(|value| value.trim().to_string()) - 154
.ok_or(RouteError::NoCredential(provider)) - 155
} - 156
- 157
/// A hosted provider's model for one operation, from the workspace pin. - 158
/// Model ids come from discovery and configuration, never from source - 159
/// (invariant 9), so a missing pin is a configuration gap. - 160
fn pinned_model(pin: Option<&String>, operation: &str) -> Result<String, RouteError> { - 161
pin.map(|model| model.trim().to_string()) - 162
.filter(|model| !model.is_empty()) - 163
.ok_or_else(|| { - 164
RouteError::Config(format!( - 165
"Voice {operation} needs a model; choose one in Voice settings" - 166
)) - 167
}) - 168
} - 169
- 170
fn openai_config(api_key: String) -> vak_llm::openai::OpenAiConfig { - 171
vak_llm::openai::OpenAiConfig { - 172
api_key, - 173
base_url: vak_llm::openai::OPENAI_DEFAULT_BASE_URL.into(), - 174
cache_key: false, - 175
openrouter: false, - 176
} - 177
} - 178
- 179
/// Transcribe encoded audio through the effective voice route. The one - 180
/// transcription implementation for every surface. - 181
pub(crate) async fn transcribe( - 182
settings: &vak_config::VoiceSettings, - 183
audio: &[u8], - 184
mime: &str, - 185
cancel: &CancellationToken, - 186
) -> Result<String, RouteError> { - 187
match route(settings)? { - 188
VoiceProvider::Local => { - 189
vak_voice::LocalTranscriber::new(local_engine(vak_voice::TRANSCRIBER_VAR)) - 190
.transcribe(audio, cancel) - 191
.await - 192
.map_err(RouteError::Local) - 193
} - 194
VoiceProvider::OpenAi => { - 195
let model = pinned_model(settings.transcription_model.as_ref(), "transcription")?; - 196
let key = credential(VoiceProvider::OpenAi)?; - 197
vak_llm::openai::transcribe(&openai_config(key), audio, mime, &model, cancel) - 198
.await - 199
.map_err(RouteError::Provider) - 200
} - 201
VoiceProvider::Gemini => { - 202
let model = pinned_model(settings.transcription_model.as_ref(), "transcription")?; - 203
let key = credential(VoiceProvider::Gemini)?; - 204
let config = vak_llm::google_live::GoogleLiveConfig::new(key, model); - 205
vak_llm::google_live::transcribe(&config, audio, mime, cancel) - 206
.await - 207
.map_err(RouteError::Provider) - 208
} - 209
} - 210
.map(|text| text.trim().to_string()) - 211
} - 212
- 213
#[derive(Deserialize)] - 214
pub(crate) struct SpeakBody { - 215
text: String, - 216
/// `wav` (default), `pcm16`, `ogg_opus` or `mp3`, within what the - 217
/// provider can produce. - 218
#[serde(default)] - 219
format: Option<String>, - 220
/// A voice being auditioned (the admin console's Preview): the narrowest - 221
/// tier, above everything the conversation resolves to. - 222
#[serde(default)] - 223
voice_override: Option<vak_config::VoiceConfig>, - 224
/// The conversation the answer belongs to. A gateway chat's session - 225
/// brings that chat's bot → chat voice tiers; any session brings its - 226
/// Agent's voice style and personality. - 227
#[serde(default)] - 228
session_id: Option<String>, - 229
} - 230
- 231
/// `POST /voice/speak`: synthesize `text` through the effective voice route - 232
/// and return the encoded audio, with a `x-vak-work-receipt` header recording - 233
/// the dispatch whatever its outcome. - 234
pub(crate) async fn voice_speak( - 235
State(state): State<AppState>, - 236
Json(body): Json<SpeakBody>, - 237
) -> Response { - 238
if body.text.trim().is_empty() { - 239
return error_response(StatusCode::BAD_REQUEST, "text must not be empty"); - 240
} - 241
let voice = resolve_speaker(&state, &body); - 242
let settings = voice.settings; - 243
let provider = match route(&settings) { - 244
Ok(provider) => provider, - 245
Err(error) => return error.into_response(), - 246
}; - 247
if body.text.chars().count() > settings.max_text_chars { - 248
return error_response( - 249
StatusCode::PAYLOAD_TOO_LARGE, - 250
format!( - 251
"text exceeds voice.max_text_chars ({})", - 252
settings.max_text_chars - 253
), - 254
); - 255
} - 256
let requested = body.format.as_deref().unwrap_or("wav"); - 257
let Some(format) = - 258
SpeakFormat::parse(requested).filter(|format| provider.speak_formats().contains(format)) - 259
else { - 260
return error_response( - 261
StatusCode::BAD_REQUEST, - 262
format!("the {provider} voice provider cannot produce '{requested}' audio"), - 263
); - 264
}; - 265
if !state.voice_requests.admit(settings.max_requests_per_minute) { - 266
return error_response( - 267
StatusCode::TOO_MANY_REQUESTS, - 268
"voice request rate limit exceeded", - 269
); - 270
} - 271
let (voice_name, persona) = (voice.voice_name, voice.persona); - 272
let model = match provider { - 273
VoiceProvider::Local => None, - 274
_ => match pinned_model(settings.synthesis_model.as_ref(), "synthesis") { - 275
Ok(model) => Some(model), - 276
Err(error) => return error.into_response(), - 277
}, - 278
}; - 279
let mut receipt = vak_llm::WorkReceipt::new( - 280
vak_llm::WorkPurpose::VoiceSynthesis, - 281
provider.as_str(), - 282
model.as_deref().unwrap_or("local"), - 283
); - 284
let cancel = CancellationToken::new(); - 285
let started = Instant::now(); - 286
let result = synthesize( - 287
provider, - 288
model.as_deref(), - 289
&body.text, - 290
format, - 291
voice_name.as_deref(), - 292
persona.as_deref(), - 293
&cancel, - 294
) - 295
.await; - 296
let latency_ms = started.elapsed().as_millis() as u64; - 297
let mut response = match result { - 298
Ok(audio) => { - 299
receipt.record( - 300
vak_llm::AttemptReason::Initial, - 301
vak_llm::FailureDomain::Unknown, - 302
vak_llm::Settlement::Ok, - 303
latency_ms, - 304
None, - 305
None, - 306
); - 307
([(header::CONTENT_TYPE, format.mime())], audio).into_response() - 308
} - 309
Err(error) => { - 310
let (domain, settlement) = error.receipt_outcome(); - 311
receipt.record( - 312
vak_llm::AttemptReason::Initial, - 313
domain, - 314
settlement, - 315
latency_ms, - 316
None, - 317
Some(error.message()), - 318
); - 319
error.into_response() - 320
} - 321
}; - 322
if let Ok(encoded) = serde_json::to_string(&receipt) - 323
&& let Ok(value) = HeaderValue::try_from(encoded) - 324
{ - 325
response.headers_mut().insert("x-vak-work-receipt", value); - 326
} - 327
response - 328
} - 329
- 330
async fn synthesize( - 331
provider: VoiceProvider, - 332
model: Option<&str>, - 333
text: &str, - 334
format: SpeakFormat, - 335
voice_name: Option<&str>, - 336
persona: Option<&str>, - 337
cancel: &CancellationToken, - 338
) -> Result<Vec<u8>, RouteError> { - 339
match provider { - 340
VoiceProvider::Local => vak_voice::LocalTtsSpeaker::new(local_engine(vak_voice::TTS_VAR)) - 341
.speak(text, format, cancel) - 342
.await - 343
.map_err(RouteError::Local), - 344
VoiceProvider::OpenAi => { - 345
let key = credential(provider)?; - 346
let response_format = match format { - 347
SpeakFormat::Wav => "wav", - 348
SpeakFormat::Pcm16 => "pcm", - 349
SpeakFormat::OggOpus => "opus", - 350
SpeakFormat::Mp3 => "mp3", - 351
}; - 352
vak_llm::openai::speak( - 353
&openai_config(key), - 354
text, - 355
model.unwrap_or_default(), - 356
voice_name.or(Some("alloy")), - 357
response_format, - 358
cancel, - 359
) - 360
.await - 361
.map_err(RouteError::Provider) - 362
} - 363
VoiceProvider::Gemini => { - 364
let key = credential(provider)?; - 365
let config = - 366
vak_llm::google_live::GoogleLiveConfig::new(key, model.unwrap_or_default()); - 367
vak_llm::google_live::speak(&config, text, persona, voice_name, cancel) - 368
.await - 369
.map_err(RouteError::Provider) - 370
} - 371
} - 372
} - 373
- 374
/// Workspace voice settings narrowed by a voice tier: its provider and - 375
/// model pins win where set, like every other route tier (invariant 23). - 376
pub(crate) fn narrowed( - 377
mut settings: vak_config::VoiceSettings, - 378
tier: Option<&vak_config::VoiceConfig>, - 379
) -> vak_config::VoiceSettings { - 380
let Some(tier) = tier else { - 381
return settings; - 382
}; - 383
let pin = |value: &Option<String>| value.clone().filter(|v| !v.trim().is_empty()); - 384
if let Some(provider) = pin(&tier.provider) { - 385
settings.provider = Some(provider); - 386
} - 387
if let Some(model) = pin(&tier.transcription_model) { - 388
settings.transcription_model = Some(model); - 389
} - 390
if let Some(model) = pin(&tier.synthesis_model) { - 391
settings.synthesis_model = Some(model); - 392
} - 393
settings - 394
} - 395
- 396
/// Whether a bot or chat voice tier can be stored: a pinned provider must be - 397
/// one that exists, and every pin must be a plausible identifier. Checked at - 398
/// write time so a typo is refused where it is made, not discovered when a - 399
/// voice note fails. - 400
pub(crate) fn check_tier(tier: &vak_config::VoiceConfig) -> Result<(), String> { - 401
if let Some(provider) = tier.provider.as_deref().filter(|p| !p.trim().is_empty()) { - 402
provider.trim().parse::<VoiceProvider>()?; - 403
} - 404
for (name, value) in [ - 405
("voice_name", &tier.voice_name), - 406
("transcription_model", &tier.transcription_model), - 407
("synthesis_model", &tier.synthesis_model), - 408
] { - 409
if value.as_deref().is_some_and(|v| v.chars().count() > 256) { - 410
return Err(format!("{name} must be at most 256 characters")); - 411
} - 412
} - 413
Ok(()) - 414
} - 415
- 416
/// The gateway chat a session is bound to, if any. - 417
fn chat_for_session(state: &AppState, session_id: &str) -> Option<String> { - 418
state - 419
.gateway - 420
.bindings_snapshot() - 421
.into_iter() - 422
.find(|(_, binding)| binding.session_id.as_deref() == Some(session_id)) - 423
.map(|(key, _)| key) - 424
} - 425
- 426
struct Speaker { - 427
settings: vak_config::VoiceSettings, - 428
voice_name: Option<String>, - 429
persona: Option<String>, - 430
} - 431
- 432
/// Route, voice and persona for a synthesis request, narrowest first: the - 433
/// auditioned override > the bound chat's bot → chat tiers (and its - 434
/// workspace) > the conversation's Agent identity > the workspace. - 435
fn resolve_speaker(state: &AppState, body: &SpeakBody) -> Speaker { - 436
let chat = body - 437
.session_id - 438
.as_deref() - 439
.and_then(|id| chat_for_session(state, id)); - 440
let core = chat - 441
.as_deref() - 442
.and_then(|key| state.gateway.core_for_entry(&state.core, key).ok()) - 443
.unwrap_or_else(|| state.core.clone()); - 444
let chat_voice = chat - 445
.as_deref() - 446
.and_then(|key| state.gateway.resolve_voice(key)); - 447
let settings = narrowed( - 448
narrowed(core.effective_voice(), chat_voice.as_ref()), - 449
body.voice_override.as_ref(), - 450
); - 451
let non_empty = |value: Option<String>| value.filter(|v| !v.trim().is_empty()); - 452
let voice_name = non_empty( - 453
body.voice_override - 454
.as_ref() - 455
.and_then(|voice| voice.voice_name.clone()), - 456
) - 457
.or_else(|| { - 458
non_empty( - 459
chat_voice - 460
.as_ref() - 461
.and_then(|voice| voice.voice_name.clone()), - 462
) - 463
}); - 464
// The persona comes from the bot/chat `identity` prompt block - 465
// (docs/design/45-prompt-layers.md); an explicit override still wins. - 466
let persona = non_empty( - 467
body.voice_override - 468
.as_ref() - 469
.and_then(|voice| voice.persona.clone()), - 470
) - 471
.or_else(|| { - 472
chat.as_deref() - 473
.and_then(|key| state.gateway.resolve_persona(key)) - 474
}) - 475
.or_else(|| { - 476
let header = crate::read_historical_header(state, body.session_id.as_deref()?, None)?; - 477
let agent = header.agent?; - 478
let style = match agent.voice.as_str() { - 479
"calm" => "Speak calmly, warmly, and at an unhurried pace.", - 480
"bright" => "Speak with clear, friendly energy.", - 481
"quiet" => "Speak gently, evenly, and without theatrical emphasis.", - 482
_ => "Speak naturally and clearly.", - 483
}; - 484
Some(if agent.personality.trim().is_empty() { - 485
style.to_string() - 486
} else { - 487
format!("{style} {}", agent.personality) - 488
}) - 489
}); - 490
Speaker { - 491
settings, - 492
voice_name, - 493
persona, - 494
} - 495
} - 496
- 497
/// What the model is told about one channel voice note. - 498
#[derive(Debug, PartialEq, Eq)] - 499
pub(crate) enum VoiceNote { - 500
Heard(String), - 501
Unheard(String), - 502
} - 503
- 504
impl VoiceNote { - 505
/// The note's place in the user message: its words, or an honest account - 506
/// of why there are none. - 507
pub(crate) fn prompt_line(&self) -> String { - 508
match self { - 509
Self::Heard(text) => text.clone(), - 510
Self::Unheard(reason) => format!("[voice note not transcribed: {reason}]"), - 511
} - 512
} - 513
} - 514
- 515
/// Transcribe one voice note from an admitted gateway chat through that - 516
/// chat's resolved route and the shared per-minute budget. Called only after - 517
/// allowlist admission, so an unknown or pending chat never spends a provider - 518
/// call. - 519
pub(crate) async fn transcribe_voice_note( - 520
state: &AppState, - 521
core: &vak_core::Core, - 522
key: &str, - 523
mime: &str, - 524
data_base64: &str, - 525
rejected: Option<&str>, - 526
) -> VoiceNote { - 527
use base64::Engine as _; - 528
if let Some(reason) = rejected { - 529
return VoiceNote::Unheard(reason.to_string()); - 530
} - 531
let Ok(audio) = base64::engine::general_purpose::STANDARD.decode(data_base64.trim()) else { - 532
return VoiceNote::Unheard("the audio could not be decoded".into()); - 533
}; - 534
let settings = narrowed( - 535
core.effective_voice(), - 536
state.gateway.resolve_voice(key).as_ref(), - 537
); - 538
if audio.is_empty() || audio.len() as u64 > settings.max_audio_bytes { - 539
return VoiceNote::Unheard(format!( - 540
"the audio must be 1 to {} bytes", - 541
settings.max_audio_bytes - 542
)); - 543
} - 544
if let Err(error) = route(&settings) { - 545
return VoiceNote::Unheard(error.message()); - 546
} - 547
if !state.voice_requests.admit(settings.max_requests_per_minute) { - 548
return VoiceNote::Unheard("this minute's voice request budget is spent".into()); - 549
} - 550
match transcribe(&settings, &audio, mime, &CancellationToken::new()).await { - 551
Ok(text) if text.is_empty() => VoiceNote::Unheard("no speech was recognized".into()), - 552
Ok(text) => VoiceNote::Heard(text), - 553
Err(error) => VoiceNote::Unheard(error.message()), - 554
} - 555
} - 556
- 557
/// `GET /voice/providers`: the static provider catalogue with credential or - 558
/// engine readiness. Never includes credential values or model ids. - 559
pub(crate) async fn voice_providers() -> Json<serde_json::Value> { - 560
let providers: Vec<_> = VoiceProvider::ALL - 561
.into_iter() - 562
.zip(vak_voice::catalogue()) - 563
.map(|(provider, descriptor)| { - 564
let (configured, readiness) = if provider == VoiceProvider::Local { - 565
let listen = vak_voice::engine_readiness( - 566
local_engine(vak_voice::TRANSCRIBER_VAR).as_deref(), - 567
); - 568
let speak = - 569
vak_voice::engine_readiness(local_engine(vak_voice::TTS_VAR).as_deref()); - 570
// Listening is what a conversation needs; without speech the - 571
// answer stays text. - 572
let readiness = serde_json::json!({ - 573
"ready": listen.ready, - 574
"detail": format!("transcriber {}; speech {}", listen.detail, speak.detail), - 575
}); - 576
(listen.configured || speak.configured, Some(readiness)) - 577
} else { - 578
(credential(provider).is_ok(), None) - 579
}; - 580
serde_json::json!({ - 581
"name": descriptor.name, - 582
"formats": descriptor.formats, - 583
"credential_vars": descriptor.credential_vars, - 584
"configured": configured, - 585
"readiness": readiness, - 586
}) - 587
}) - 588
.collect(); - 589
Json(serde_json::json!({ "providers": providers })) - 590
} - 591
- 592
#[derive(Debug, Deserialize)] - 593
pub(crate) struct SessionQuery { - 594
session_id: String, - 595
} - 596
- 597
/// `WS /voice/session`: a governed spoken conversation bound to one Agent - 598
/// session. Final transcripts become ordinary turns on that session. - 599
pub(crate) async fn voice_socket( - 600
State(state): State<AppState>, - 601
Query(query): Query<SessionQuery>, - 602
headers: header::HeaderMap, - 603
upgrade: WebSocketUpgrade, - 604
) -> Response { - 605
// An upgrade is a GET, so the router's mutation Origin check does not - 606
// cover it; a microphone relay that spends provider money must check. - 607
let origin = headers.get(header::ORIGIN).and_then(|v| v.to_str().ok()); - 608
if !crate::origin_is_trusted(origin, &state.core.config().server.trusted_hosts) { - 609
return error_response(StatusCode::FORBIDDEN, "voice session origin is not trusted"); - 610
} - 611
upgrade.on_upgrade(move |socket| drive(socket, state, query.session_id)) - 612
} - 613
- 614
async fn send(socket: &mut WebSocket, control: &ServerControl) -> bool { - 615
match control.encode() { - 616
Ok(text) => socket.send(Message::Text(text.into())).await.is_ok(), - 617
Err(_) => false, - 618
} - 619
} - 620
- 621
async fn send_error(socket: &mut WebSocket, message: impl Into<String>, remedy: &str) -> bool { - 622
send( - 623
socket, - 624
&ServerControl::Error { - 625
message: message.into(), - 626
remedy: Some(remedy.into()), - 627
}, - 628
) - 629
.await - 630
} - 631
- 632
struct Lease(Arc<AtomicUsize>); - 633
- 634
impl Lease { - 635
fn acquire(active: &Arc<AtomicUsize>, limit: usize) -> Option<Self> { - 636
active - 637
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| { - 638
(current < limit).then_some(current + 1) - 639
}) - 640
.ok() - 641
.map(|_| Self(Arc::clone(active))) - 642
} - 643
} - 644
- 645
impl Drop for Lease { - 646
fn drop(&mut self) { - 647
self.0.fetch_sub(1, Ordering::AcqRel); - 648
} - 649
} - 650
- 651
/// Bound on a spoken answer's text, keeping its frame far below the - 652
/// protocol's control-frame limit whatever the configured synthesis cap. - 653
const MAX_REPLY_CHARS: usize = 8_000; - 654
- 655
/// Protocol state of one socket: at most one open utterance, and an - 656
/// utterance id is never reused once it has been closed. - 657
#[derive(Default)] - 658
struct Utterances { - 659
open: Option<String>, - 660
audio: Vec<u8>, - 661
closed: HashSet<String>, - 662
} - 663
- 664
impl Utterances { - 665
fn start(&mut self, id: String) -> Result<(), &'static str> { - 666
if id.trim().is_empty() || self.closed.contains(&id) { - 667
return Err("utterance ids must be non-empty and never reused"); - 668
} - 669
// A new start abandons an unfinished utterance: stopping voice - 670
// mid-sentence never submits it. - 671
self.audio.clear(); - 672
self.open = Some(id); - 673
Ok(()) - 674
} - 675
- 676
fn stop(&mut self, id: &str) -> Result<Vec<u8>, &'static str> { - 677
if self.open.as_deref() != Some(id) { - 678
return Err("speech_stopped does not name the open utterance"); - 679
} - 680
self.open = None; - 681
self.closed.insert(id.to_string()); - 682
Ok(std::mem::take(&mut self.audio)) - 683
} - 684
- 685
/// Audio outside an open utterance is dropped, never buffered. - 686
fn push(&mut self, bytes: &[u8]) { - 687
if self.open.is_some() { - 688
self.audio.extend_from_slice(bytes); - 689
} - 690
} - 691
} - 692
- 693
/// What a closed utterance became. - 694
enum Outcome { - 695
Discarded(DiscardReason), - 696
Transcript(String), - 697
Failed(RouteError), - 698
} - 699
- 700
async fn close_utterance(state: &AppState, pcm: &[u8]) -> Outcome { - 701
// Measured here, not trusted from the client: steady room noise that a - 702
// client detector let through never reaches a provider. - 703
if !SpeechEvidence::measure(pcm, audio::SOCKET_SAMPLE_RATE_HZ).is_speech() { - 704
return Outcome::Discarded(DiscardReason::InsufficientSpeech); - 705
} - 706
// Re-read per utterance so a settings change applies without reconnecting. - 707
let settings = state.core.effective_voice(); - 708
if let Err(error) = route(&settings) { - 709
return Outcome::Failed(error); - 710
} - 711
if !state.voice_requests.admit(settings.max_requests_per_minute) { - 712
return Outcome::Discarded(DiscardReason::RateLimited); - 713
} - 714
let wav = match audio::wrap_wav(pcm, audio::PcmSpec::SOCKET) { - 715
Ok(wav) => wav, - 716
Err(error) => return Outcome::Failed(RouteError::Local(error)), - 717
}; - 718
match transcribe(&settings, &wav, "audio/wav", &CancellationToken::new()).await { - 719
Ok(text) if text.is_empty() => Outcome::Discarded(DiscardReason::EmptyTranscript), - 720
Ok(text) => Outcome::Transcript(text), - 721
Err(error) => Outcome::Failed(error), - 722
} - 723
} - 724
- 725
async fn drive(mut socket: WebSocket, state: AppState, session_id: String) { - 726
let settings = state.core.effective_voice(); - 727
// Refuse before the client asks for the microphone when no utterance - 728
// could be served: disabled, or no usable provider route. - 729
if let Err(error) = route(&settings) { - 730
send_error(&mut socket, error.message(), "Open Voice settings").await; - 731
return; - 732
} - 733
let handle = state - 734
.sessions - 735
.lock() - 736
.ok() - 737
.and_then(|sessions| sessions.get(&session_id).cloned()); - 738
let Some(handle) = handle else { - 739
send_error( - 740
&mut socket, - 741
"Voice conversation is unavailable", - 742
"Reopen the Agent conversation and try again", - 743
) - 744
.await; - 745
return; - 746
}; - 747
let Some(_lease) = Lease::acquire(&state.voice_active, settings.max_concurrent) else { - 748
send_error( - 749
&mut socket, - 750
"Voice session concurrency limit reached", - 751
"Wait for an active voice session to finish or increase the configured limit", - 752
) - 753
.await; - 754
return; - 755
}; - 756
let deadline = tokio::time::Instant::now() + Duration::from_secs(settings.max_session_secs); - 757
let ready = ServerControl::Ready { - 758
protocol_version: VOICE_PROTOCOL_VERSION, - 759
sample_rate_hz: audio::SOCKET_SAMPLE_RATE_HZ, - 760
channels: audio::PcmSpec::SOCKET.channels, - 761
}; - 762
if !send(&mut socket, &ready).await { - 763
return; - 764
} - 765
let mut utterances = Utterances::default(); - 766
let mut received_bytes: u64 = 0; - 767
let (completed_tx, mut completed_rx) = - 768
tokio::sync::mpsc::unbounded_channel::<(String, String)>(); - 769
loop { - 770
let message = tokio::select! { - 771
inbound = futures::StreamExt::next(&mut socket) => match inbound { - 772
Some(Ok(message)) => message, - 773
_ => break, - 774
}, - 775
Some((utterance_id, reply)) = completed_rx.recv() => { - 776
let text: String = reply.chars().take(MAX_REPLY_CHARS).collect(); - 777
if !send(&mut socket, &ServerControl::TurnCompleted { utterance_id, text }).await { - 778
break; - 779
} - 780
continue; - 781
} - 782
() = tokio::time::sleep_until(deadline) => { - 783
send_error( - 784
&mut socket, - 785
"This voice session reached its time limit", - 786
"Press Voice to start a new session", - 787
) - 788
.await; - 789
break; - 790
} - 791
}; - 792
match message { - 793
Message::Binary(bytes) => { - 794
received_bytes = received_bytes.saturating_add(bytes.len() as u64); - 795
if received_bytes > settings.max_audio_bytes { - 796
send_error( - 797
&mut socket, - 798
format!( - 799
"Voice audio budget exceeded ({} bytes)", - 800
settings.max_audio_bytes - 801
), - 802
"Start a new session or increase the configured budget", - 803
) - 804
.await; - 805
break; - 806
} - 807
if vak_voice::protocol::validate_audio(&bytes).is_err() { - 808
send_error( - 809
&mut socket, - 810
"Invalid PCM audio frame", - 811
"Send mono 16-bit PCM frames", - 812
) - 813
.await; - 814
break; - 815
} - 816
utterances.push(&bytes); - 817
} - 818
Message::Text(text) => { - 819
let control = match ClientControl::decode(text.as_str()) { - 820
Ok(control) => control, - 821
Err(error) => { - 822
send_error(&mut socket, error.to_string(), "Update the voice client").await; - 823
break; - 824
} - 825
}; - 826
match control { - 827
ClientControl::SpeechStarted { utterance_id } => { - 828
if let Err(reason) = utterances.start(utterance_id) { - 829
send_error(&mut socket, reason, "Update the voice client").await; - 830
break; - 831
} - 832
} - 833
ClientControl::SpeechStopped { utterance_id } => { - 834
let pcm = match utterances.stop(&utterance_id) { - 835
Ok(pcm) => pcm, - 836
Err(reason) => { - 837
send_error(&mut socket, reason, "Update the voice client").await; - 838
break; - 839
} - 840
}; - 841
let frame = match close_utterance(&state, &pcm).await { - 842
Outcome::Discarded(reason) => ServerControl::Discarded { - 843
utterance_id, - 844
reason, - 845
}, - 846
Outcome::Failed(error) => ServerControl::Error { - 847
message: format!("Voice transcription failed: {}", error.message()), - 848
remedy: Some( - 849
"Check Voice settings or use the text composer".into(), - 850
), - 851
}, - 852
Outcome::Transcript(text) => { - 853
if let Ok(mut log) = handle.session.lock() - 854
&& let Some(log) = log.as_mut() - 855
{ - 856
let _ = log.append_voice_transcript( - 857
format!("voice:{utterance_id}"), - 858
text.clone(), - 859
true, - 860
); - 861
} - 862
// A final transcript is an ordinary user turn - 863
// through the same governed runner as typed - 864
// input: permissions, intent, budgets and - 865
// receipts are one contract. - 866
let (reply_tx, reply_rx) = tokio::sync::oneshot::channel(); - 867
crate::gateway::start_turn_chain_with_gateway( - 868
state.gateway.clone(), - 869
&handle.core, - 870
handle.clone(), - 871
crate::gateway::compose_voice_prompt(&text), - 872
Some(reply_tx), - 873
); - 874
let completed = completed_tx.clone(); - 875
let answered = utterance_id.clone(); - 876
tokio::spawn(async move { - 877
if let Ok(reply) = reply_rx.await { - 878
let _ = completed.send((answered, reply.text)); - 879
} - 880
}); - 881
ServerControl::Transcript { utterance_id, text } - 882
} - 883
}; - 884
let failed = matches!(frame, ServerControl::Error { .. }); - 885
if !send(&mut socket, &frame).await || failed { - 886
break; - 887
} - 888
} - 889
ClientControl::Playback { - 890
utterance_id, - 891
emitted_ms, - 892
interrupted, - 893
} => { - 894
if let Ok(mut log) = handle.session.lock() - 895
&& let Some(log) = log.as_mut() - 896
{ - 897
let _ = log.append_voice_playback( - 898
format!("voice:{utterance_id}"), - 899
emitted_ms, - 900
interrupted, - 901
); - 902
} - 903
} - 904
} - 905
} - 906
Message::Close(_) => break, - 907
_ => {} - 908
} - 909
} - 910
let _ = socket.close().await; - 911
} - 912
- 913
#[cfg(test)] - 914
#[allow(clippy::unwrap_used, clippy::expect_used)] - 915
mod tests { - 916
use super::*; - 917
- 918
#[test] - 919
fn request_window_enforces_its_limit_and_rolls_over() { - 920
let window = RequestWindow::new(); - 921
let start = Instant::now(); - 922
assert!(window.admit_at(2, start)); - 923
assert!(window.admit_at(2, start)); - 924
assert!(!window.admit_at(2, start)); - 925
assert!(window.admit_at(2, start + Duration::from_secs(61))); - 926
} - 927
- 928
#[test] - 929
fn a_zero_limit_admits_nothing() { - 930
assert!(!RequestWindow::new().admit(0)); - 931
} - 932
- 933
#[test] - 934
fn lease_admission_is_atomic_and_released_on_drop() { - 935
let active = Arc::new(AtomicUsize::new(0)); - 936
let first = Lease::acquire(&active, 1).expect("first lease"); - 937
assert!(Lease::acquire(&active, 1).is_none()); - 938
drop(first); - 939
assert!(Lease::acquire(&active, 1).is_some()); - 940
assert_eq!(active.load(Ordering::Acquire), 0); - 941
} - 942
- 943
#[test] - 944
fn utterance_boundaries_are_strict() { - 945
let mut utterances = Utterances::default(); - 946
utterances.push(&[1, 0]); - 947
assert!(utterances.start("u1".into()).is_ok()); - 948
utterances.push(&[2, 0]); - 949
assert!(utterances.stop("u2").is_err()); - 950
assert_eq!( - 951
utterances.stop("u1").unwrap(), - 952
vec![2, 0], - 953
"audio before the start is dropped" - 954
); - 955
assert!(utterances.start("u1".into()).is_err()); - 956
assert!(utterances.start(" ".into()).is_err()); - 957
} - 958
- 959
#[test] - 960
fn a_restart_abandons_the_unfinished_utterance() { - 961
let mut utterances = Utterances::default(); - 962
utterances.start("u1".into()).unwrap(); - 963
utterances.push(&[9, 9]); - 964
utterances.start("u2".into()).unwrap(); - 965
assert!(utterances.stop("u2").unwrap().is_empty()); - 966
assert!(utterances.stop("u1").is_err()); - 967
} - 968
- 969
#[test] - 970
fn an_unset_or_unknown_provider_is_a_configuration_error() { - 971
let mut settings = vak_config::VoiceSettings { - 972
enabled: true, - 973
..Default::default() - 974
}; - 975
assert!(matches!(route(&settings), Err(RouteError::Config(_)))); - 976
settings.provider = Some("google".into()); - 977
assert!(matches!(route(&settings), Err(RouteError::Config(_)))); - 978
settings.provider = Some("openai".into()); - 979
assert_eq!(route(&settings).unwrap(), VoiceProvider::OpenAi); - 980
settings.enabled = false; - 981
assert!(matches!(route(&settings), Err(RouteError::Disabled))); - 982
} - 983
- 984
#[test] - 985
fn a_chat_tier_narrows_the_workspace_route_and_blanks_inherit() { - 986
let workspace = vak_config::VoiceSettings { - 987
enabled: true, - 988
provider: Some("gemini".into()), - 989
transcription_model: Some("stt-workspace".into()), - 990
synthesis_model: Some("tts-workspace".into()), - 991
..Default::default() - 992
}; - 993
let chat = vak_config::VoiceConfig { - 994
provider: Some("openai".into()), - 995
synthesis_model: Some("tts-chat".into()), - 996
transcription_model: Some(" ".into()), - 997
..Default::default() - 998
}; - 999
let effective = narrowed(workspace.clone(), Some(&chat)); - 1000
assert_eq!(effective.provider.as_deref(), Some("openai"));
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.