1#![doc(html_favicon_url = "https://zng-ui.github.io/res/zng-logo-icon.png")]
2#![doc(html_logo_url = "https://zng-ui.github.io/res/zng-logo.png")]
3#![doc = include_str!(concat!("../", std::env!("CARGO_PKG_README")))]
15#![warn(unused_extern_crates)]
16#![warn(missing_docs)]
17
18use std::{path::PathBuf, pin::Pin};
19
20use zng_app::{
21 update::UPDATES,
22 view_process::{
23 VIEW_PROCESS, VIEW_PROCESS_INITED_EVENT, ViewAudioHandle,
24 raw_events::{RAW_AUDIO_DECODE_ERROR_EVENT, RAW_AUDIO_DECODED_EVENT, RAW_AUDIO_METADATA_DECODED_EVENT},
25 },
26};
27use zng_app_context::app_local;
28use zng_clone_move::clmv;
29use zng_task::channel::{IpcBytes, IpcReadHandle};
30use zng_txt::ToTxt;
31use zng_unique_id::{IdEntry, IdMap};
32use zng_unit::ByteLength;
33use zng_var::{ResponseVar, Var, VarHandle, response_var, var};
34use zng_view_api::audio::{AudioDecoded, AudioMetadata, AudioRequest};
35
36mod types;
37pub use types::*;
38
39mod output;
40pub use output::*;
41
42app_local! {
43 static AUDIOS_SV: AudiosService = AudiosService::new();
44 static AUDIOS_EXTENSIONS: Vec<Box<dyn AudiosExtension>> = vec![];
45}
46
47struct AudiosService {
48 load_in_headless: Var<bool>,
49 limits: Var<AudioLimits>,
50
51 cache: IdMap<AudioHash, AudioVar>,
52 outputs: IdMap<AudioOutputId, WeakAudioOutput>,
53 perm_outputs: IdMap<AudioOutputId, AudioOutput>,
54}
55impl AudiosService {
56 pub fn new() -> Self {
57 Self {
58 load_in_headless: var(false),
59 limits: var(AudioLimits::default()),
60
61 cache: IdMap::new(),
62 outputs: IdMap::new(),
63 perm_outputs: IdMap::new(),
64 }
65 }
66}
67
68pub struct AUDIOS;
76impl AUDIOS {
77 pub fn load_in_headless(&self) -> Var<bool> {
87 AUDIOS_SV.read().load_in_headless.clone()
88 }
89
90 pub fn limits(&self) -> Var<AudioLimits> {
92 AUDIOS_SV.read().limits.clone()
93 }
94
95 pub fn read(&self, path: impl Into<PathBuf>) -> AudioVar {
101 self.audio_impl(path.into().into(), AudioOptions::cache(), None)
102 }
103
104 #[cfg(feature = "http")]
113 pub fn download<U>(&self, uri: U, accept: Option<zng_txt::Txt>) -> AudioVar
114 where
115 U: TryInto<zng_task::http::Uri>,
116 <U as TryInto<zng_task::http::Uri>>::Error: ToTxt,
117 {
118 match uri.try_into() {
119 Ok(uri) => self.audio_impl(AudioSource::Download(uri, accept), AudioOptions::cache(), None),
120 Err(e) => zng_var::const_var(AudioTrack::new_error(e.to_txt())),
121 }
122 }
123
124 pub fn from_static(&self, data: &'static [u8], format: impl Into<AudioDataFormat>) -> AudioVar {
144 self.audio_impl((data, format.into()).into(), AudioOptions::cache(), None)
145 }
146
147 pub fn from_data(&self, data: IpcBytes, format: impl Into<AudioDataFormat>) -> AudioVar {
155 self.audio_impl((data, format.into()).into(), AudioOptions::cache(), None)
156 }
157
158 pub fn audio(&self, source: impl Into<AudioSource>, options: AudioOptions, limits: Option<AudioLimits>) -> AudioVar {
167 self.audio_impl(source.into(), options, limits)
168 }
169 fn audio_impl(&self, source: AudioSource, options: AudioOptions, limits: Option<AudioLimits>) -> AudioVar {
170 let r = var(AudioTrack::new_loading());
171 let ri = r.read_only();
172 UPDATES.once_update("AUDIOS.audio", move || {
173 audio(source, options, limits, r);
174 });
175 ri
176 }
177
178 pub fn audio_task<F>(&self, source: impl IntoFuture<IntoFuture = F>, options: AudioOptions, limits: Option<AudioLimits>) -> AudioVar
190 where
191 F: Future<Output = AudioSource> + Send + 'static,
192 {
193 self.audio_task_impl(Box::pin(source.into_future()), options, limits)
194 }
195 fn audio_task_impl(
196 &self,
197 source: Pin<Box<dyn Future<Output = AudioSource> + Send + 'static>>,
198 options: AudioOptions,
199 limits: Option<AudioLimits>,
200 ) -> AudioVar {
201 let r = var(AudioTrack::new_loading());
202 let ri = r.read_only();
203 zng_task::spawn(async move {
204 let source = source.await;
205 audio(source, options, limits, r);
206 });
207 ri
208 }
209
210 pub fn register(&self, key: Option<AudioHash>, audio: (ViewAudioHandle, AudioMetadata, AudioDecoded)) -> AudioVar {
218 let r = var(AudioTrack::new_loading());
219 let rr = r.read_only();
220 UPDATES.once_update("AUDIOS.register", move || {
221 audio_view(key, audio.0, audio.1, audio.2, None, r);
222 });
223 rr
224 }
225
226 pub fn clean(&self, key: AudioHash) {
231 UPDATES.once_update("AUDIOS.clean", move || {
232 if let IdEntry::Occupied(e) = AUDIOS_SV.write().cache.entry(key)
233 && e.get().strong_count() == 1
234 {
235 e.remove();
236 }
237 });
238 }
239
240 pub fn purge(&self, key: AudioHash) {
245 UPDATES.once_update("AUDIOS.purge", move || {
246 AUDIOS_SV.write().cache.remove(&key);
247 });
248 }
249
250 pub fn cache_key(&self, audio: &AudioTrack) -> Option<AudioHash> {
252 let key = audio.cache_key?;
253 if AUDIOS_SV.read().cache.contains_key(&key) {
254 Some(key)
255 } else {
256 None
257 }
258 }
259
260 pub fn is_cached(&self, audio: &AudioTrack) -> bool {
262 match &audio.cache_key {
263 Some(k) => AUDIOS_SV.read().cache.contains_key(k),
264 None => false,
265 }
266 }
267
268 pub fn clean_all(&self) {
270 UPDATES.once_update("AUDIOS.clean_all", || {
271 AUDIOS_SV.write().cache.retain(|_, v| v.strong_count() > 1);
272 });
273 }
274
275 pub fn purge_all(&self) {
280 UPDATES.once_update("AUDIOS.purge_all", || {
281 AUDIOS_SV.write().cache.clear();
282 });
283 }
284
285 pub fn extend(&self, extension: Box<dyn AudiosExtension>) {
289 UPDATES.once_update("AUDIOS.extend", move || {
290 AUDIOS_EXTENSIONS.write().push(extension);
291 });
292 }
293
294 pub fn available_formats(&self) -> Vec<AudioFormat> {
296 let mut formats = VIEW_PROCESS.info().audio.clone();
297
298 for ext in AUDIOS_EXTENSIONS.read().iter() {
299 ext.available_formats(&mut formats);
300 }
301
302 formats
303 }
304
305 #[cfg(feature = "http")]
306 fn http_accept(&self) -> zng_txt::Txt {
307 let mut s = String::new();
308 let mut sep = "";
309 for f in self.available_formats() {
310 for f in f.media_type_suffixes_iter() {
311 s.push_str(sep);
312 s.push_str("audio/");
313 s.push_str(f);
314 sep = ",";
315 }
316 }
317 s.into()
318 }
319}
320
321impl AUDIOS {
322 pub fn open_output(
333 &self,
334 id: impl Into<AudioOutputId>,
335 init: impl FnOnce(&mut AudioOutputOptions) + Send + 'static,
336 ) -> ResponseVar<AudioOutput> {
337 self.open_audio_out(id.into(), Box::new(init))
338 }
339
340 fn open_audio_out(&self, id: AudioOutputId, init: Box<dyn FnOnce(&mut AudioOutputOptions) + Send>) -> ResponseVar<AudioOutput> {
341 let (responder, r) = response_var();
342 UPDATES.once_update("AUDIOS.open_output", move || match AUDIOS_SV.write().outputs.entry(id) {
343 IdEntry::Occupied(mut e) => match e.get().upgrade() {
344 Some(r) => responder.respond(r),
345 None => {
346 let mut opt = AudioOutputOptions::default();
347 init(&mut opt);
348 let r = AudioOutput::open(id, opt);
349 e.insert(r.downgrade());
350 responder.respond(r);
351 }
352 },
353 IdEntry::Vacant(e) => {
354 let mut opt = AudioOutputOptions::default();
355 init(&mut opt);
356 let r = AudioOutput::open(id, opt);
357 e.insert(r.downgrade());
358 responder.respond(r);
359 }
360 });
361 r
362 }
363
364 pub(crate) fn perm_output(&self, output: &AudioOutput) {
365 AUDIOS_SV.write().perm_outputs.entry(output.id()).or_insert_with(|| output.clone());
366 }
367}
368
369fn audio(mut source: AudioSource, mut options: AudioOptions, limits: Option<AudioLimits>, r: Var<AudioTrack>) {
370 let limits = limits.unwrap_or_else(|| AUDIOS_SV.read().limits.get());
371
372 {
374 let mut exts = AUDIOS_EXTENSIONS.write();
375 if !exts.is_empty() {
376 tracing::trace!("process audio with {} extensions", exts.len());
377 }
378 for ext in exts.iter_mut() {
379 ext.audio(&limits, &mut source, &mut options);
380 }
381 }
382
383 let mut s = AUDIOS_SV.write();
385
386 if let AudioSource::Audio(var) = source {
387 var.set_bind(&r).perm();
389 r.hold(var).perm();
390 return;
391 }
392
393 if !VIEW_PROCESS.is_available() && !s.load_in_headless.get() {
394 tracing::debug!("ignoring audio request due headless mode");
395 return;
396 }
397
398 let key = source.hash128(&options).unwrap();
399
400 match options.cache_mode {
402 AudioCacheMode::Ignore => (),
403 AudioCacheMode::Cache => {
404 match s.cache.entry(key) {
405 IdEntry::Occupied(e) => {
406 let var = e.get();
408 var.set_bind(&r).perm();
409 r.hold(var.clone()).perm();
410 return;
411 }
412 IdEntry::Vacant(e) => {
413 e.insert(r.clone());
415 }
416 }
417 }
418 AudioCacheMode::Retry => {
419 match s.cache.entry(key) {
420 IdEntry::Occupied(mut e) => {
421 let var = e.get();
422 if var.with(AudioTrack::is_error) {
423 r.set_bind(var).perm();
428 var.hold(r.clone()).perm();
429
430 e.insert(r.clone());
432 } else {
433 var.set_bind(&r).perm();
435 r.hold(var.clone()).perm();
436 return;
437 }
438 }
439 IdEntry::Vacant(e) => {
440 e.insert(r.clone());
442 }
443 }
444 }
445 AudioCacheMode::Reload => {
446 match s.cache.entry(key) {
447 IdEntry::Occupied(mut e) => {
448 let var = e.get();
449 r.set_bind(var).perm();
450 var.hold(r.clone()).perm();
451
452 e.insert(r.clone());
453 }
454 IdEntry::Vacant(e) => {
455 e.insert(r.clone());
457 }
458 }
459 }
460 }
461 drop(s);
462
463 match source {
464 AudioSource::Read(path) => {
465 fn read(path: &PathBuf, limit: (&AudioSourceFilter<PathBuf>, ByteLength)) -> std::io::Result<IpcReadHandle> {
466 if !limit.0.allows(path) {
467 return Err(std::io::Error::new(
468 std::io::ErrorKind::PermissionDenied,
469 "file path no allowed by limit",
470 ));
471 }
472 let file = std::fs::File::open(path)?;
473 if file.metadata()?.len() > limit.1.bytes() {
474 return Err(std::io::Error::new(std::io::ErrorKind::InvalidData, "file length exceeds limit"));
475 }
476 IpcReadHandle::best_read_blocking(file)
477 }
478 let data_format = match path.extension() {
479 Some(ext) => AudioDataFormat::FileExtension(ext.to_string_lossy().to_txt()),
480 None => AudioDataFormat::Unknown,
481 };
482 zng_task::spawn_wait(move || match read(&path, (&limits.allow_path, limits.max_encoded_len)) {
483 Ok(data) => audio_data(false, Some(key), data_format, data, options, limits, r),
484 Err(e) => {
485 r.set(AudioTrack::new_error(e.to_txt()));
486 }
487 });
488 }
489 #[cfg(feature = "http")]
490 AudioSource::Download(uri, accept) => {
491 let accept = accept.unwrap_or_else(|| AUDIOS.http_accept());
492
493 use zng_task::http::*;
494 async fn download(
495 uri: Uri,
496 accept: zng_txt::Txt,
497 limit: (AudioSourceFilter<Uri>, ByteLength),
498 ) -> Result<(AudioDataFormat, IpcBytes), Error> {
499 if !limit.0.allows(&uri) {
500 return Err(Box::new(std::io::Error::new(
501 std::io::ErrorKind::PermissionDenied,
502 "uri no allowed by limit",
503 )));
504 }
505 let request = Request::get(uri)?.max_length(limit.1).header(header::ACCEPT, accept.as_str())?;
506 let mut response = send(request).await?;
507 response.error().await?;
508 let data_format = match response.header().get(&header::CONTENT_TYPE).and_then(|m| m.to_str().ok()) {
509 Some(m) => AudioDataFormat::MimeType(m.to_txt()),
510 None => AudioDataFormat::Unknown,
511 };
512 let data = response.body().await?;
513
514 Ok((data_format, data))
515 }
516
517 zng_task::spawn(async move {
518 match download(uri, accept, (limits.allow_uri.clone(), limits.max_encoded_len)).await {
519 Ok((fmt, data)) => {
520 audio_data(false, Some(key), fmt, data.into(), options, limits, r);
521 }
522 Err(e) => r.set(AudioTrack::new_error(e.to_txt())),
523 }
524 });
525 }
526 AudioSource::Data(_, data, format) => audio_data(false, Some(key), format, data.into(), options, limits, r),
527 _ => unreachable!(),
528 }
529}
530
531fn audio_data(
533 is_respawn: bool,
534 cache_key: Option<AudioHash>,
535 format: AudioDataFormat,
536 data: IpcReadHandle,
537 options: AudioOptions,
538 limits: AudioLimits,
539 r: Var<AudioTrack>,
540) {
541 if !is_respawn && let Some(key) = cache_key {
542 let mut exts = AUDIOS_EXTENSIONS.write();
543 if !exts.is_empty() {
544 tracing::trace!("process audio_data with {} extensions", exts.len());
545 }
546 for ext in exts.iter_mut() {
547 if let Some(replacement) = ext.audio_data(limits.max_decoded_len, &key, &data, &format, &options) {
548 replacement.set_bind(&r).perm();
549 r.hold(replacement).perm();
550
551 tracing::trace!("extension replaced audio_data");
552 return;
553 }
554 }
555 }
556
557 if !VIEW_PROCESS.is_available() {
558 tracing::debug!("ignoring audio view request after test load due to headless mode");
559 return;
560 }
561
562 let mut data = data;
563 let data_clone = match data.duplicate_or_read_blocking() {
564 Ok(d) => d,
565 Err(e) => {
566 tracing::error!("audio data lost, {e}");
567 return;
568 }
569 };
570 let data = data;
571 let mut request = AudioRequest::new(format.clone(), data_clone, limits.max_decoded_len.bytes());
572 request.tracks = options.tracks;
573
574 let try_gen = VIEW_PROCESS.generation();
575
576 match VIEW_PROCESS.add_audio(request) {
577 Ok(view_img) => audio_view(
578 cache_key,
579 view_img,
580 AudioMetadata::default(),
581 AudioDecoded::default(),
582 Some((format, data, options, limits)),
583 r,
584 ),
585 Err(_) => {
586 tracing::debug!("audio view request failed, will retry on respawn");
587
588 zng_task::spawn(async move {
589 VIEW_PROCESS_INITED_EVENT.wait_match(move |a| a.generation != try_gen).await;
590 audio_data(true, cache_key, format, data, options, limits, r);
591 });
592 }
593 }
594}
595fn audio_view(
597 cache_key: Option<AudioHash>,
598 handle: ViewAudioHandle,
599 meta: AudioMetadata,
600 decoded: AudioDecoded,
601 respawn_data: Option<(AudioDataFormat, IpcReadHandle, AudioOptions, AudioLimits)>,
602 r: Var<AudioTrack>,
603) {
604 let aud = AudioTrack::new(cache_key, handle, meta, decoded);
605 let is_loaded = aud.is_loaded();
606 let is_dummy = aud.view_handle().is_dummy();
607 r.set(aud);
608
609 if is_loaded {
610 audio_decoded(r);
611 return;
612 }
613
614 if is_dummy {
615 tracing::error!("tried to register dummy handle");
616 return;
617 }
618
619 let decoding_respawn_handle = if respawn_data.is_some() {
621 let r_weak = r.downgrade();
622 let mut respawn_data = respawn_data;
623 VIEW_PROCESS_INITED_EVENT.hook(move |_| {
624 if let Some(r) = r_weak.upgrade() {
625 let (format, data, options, limits) = respawn_data.take().unwrap();
626 audio_data(true, cache_key, format, data, options, limits, r);
627 }
628 false
629 })
630 } else {
631 VarHandle::dummy()
633 };
634
635 let r_weak = r.downgrade();
637 let decode_error_handle = RAW_AUDIO_DECODE_ERROR_EVENT.hook(move |args| match r_weak.upgrade() {
638 Some(r) => {
639 if let Some(handle) = args.handle.upgrade()
640 && r.with(|aud| aud.view_handle() == &handle)
641 {
642 r.set(AudioTrack::new_error(args.error.clone()));
643 false
644 } else {
645 r.with(AudioTrack::is_loading)
646 }
647 }
648 None => false,
649 });
650
651 let r_weak = r.downgrade();
653 let decode_meta_handle = RAW_AUDIO_METADATA_DECODED_EVENT.hook(move |args| match r_weak.upgrade() {
654 Some(r) => {
655 let handle = match args.handle.upgrade() {
656 Some(h) => h,
657 None => return r.with(AudioTrack::is_loading),
658 };
659 if r.with(|aud| aud.view_handle() == &handle) {
660 let meta = args.meta.clone();
661 r.modify(move |i| i.meta = meta);
662 } else if let Some(p) = &args.meta.parent
663 && p.parent == r.with(|aud| aud.view_handle().audio_id())
664 {
665 let mut decoded = AudioDecoded::default();
667 decoded.id = args.meta.id;
668 let track = var(AudioTrack::new(None, handle.clone(), args.meta.clone(), decoded.clone()));
669 r.modify(clmv!(track, |i| i.insert_track(track)));
670 audio_view(None, handle, args.meta.clone(), decoded, None, track);
671 }
672 r.with(AudioTrack::is_loading)
673 }
674 None => false,
675 });
676
677 let r_weak = r.downgrade();
679 RAW_AUDIO_DECODED_EVENT
680 .hook(move |args| {
681 let _hold = [&decoding_respawn_handle, &decode_error_handle, &decode_meta_handle];
682 match r_weak.upgrade() {
683 Some(r) => {
684 if let Some(handle) = args.handle.upgrade()
685 && r.with(|aud| aud.view_handle() == &handle)
686 {
687 let data = args.audio.upgrade().unwrap();
688 let is_loading = !data.is_full;
689 r.modify(move |i| i.data = (*data.0).clone());
690 if !is_loading {
691 audio_decoded(r);
692 }
693 is_loading
694 } else {
695 r.with(AudioTrack::is_loading)
696 }
697 }
698 None => false,
699 }
700 })
701 .perm();
702}
703fn audio_decoded(r: Var<AudioTrack>) {
705 let r_weak = r.downgrade();
706 VIEW_PROCESS_INITED_EVENT
707 .hook(move |_| {
708 if let Some(r) = r_weak.upgrade() {
709 let aud = r.get();
710 if !aud.is_loaded() {
711 return false;
713 }
714
715 let options = AudioOptions::none();
717 let format = AudioDataFormat::InterleavedF32 {
718 channel_count: aud.channel_count(),
719 sample_rate: aud.sample_rate(),
720 total_duration: aud.total_duration(),
721 };
722 audio_data(
723 true,
724 aud.cache_key,
725 format,
726 aud.chunk().into_inner().into(),
727 options,
728 AudioLimits::none(),
729 r,
730 );
731 }
732 false
733 })
734 .perm();
735}