use std::cell::{Cell, RefCell};
use std::collections::VecDeque;
use std::rc::{Rc, Weak};
use std::time::Duration;
use dioxus::core::{Runtime as DioxusRuntime, ScopeId};
use dioxus::prelude::*;
use js_sys::{Array, Reflect, Uint8Array};
use wasm_bindgen::JsCast;
use wasm_bindgen::prelude::*;
use web_sys::{
AudioContext, Blob, BlobEvent, BlobPropertyBag, MediaRecorder, MediaRecorderOptions,
MediaStream, MediaStreamAudioSourceNode, MediaStreamTrack, MediaStreamTrackState,
};
use super::*;
use crate::analysis::{AudioAnalyser, AudioAnalyserControl, peak_amplitude};
use crate::devices::web::{audio_error_from_js, media_devices, stop_stream};
use crate::{AudioData, AudioErrorKind, RecordingChunk};
enum RecordingSourceRequest {
Acquire(Option<AudioInputId>),
Supplied(RecordingSource),
}
enum RecordingInput {
Acquire(Option<AudioInputId>),
Supplied(PreparedSource),
}
enum AcceptedRecordingInput {
Acquire(Option<AudioInputId>),
Supplied(PendingCapture),
}
pub(super) fn use_web_audio_recorder(
options: RecorderOptions,
selected_input: ReadSignal<Option<AudioInputId>>,
) -> AudioRecorder {
let mut status = use_signal(|| RecorderStatus::Idle);
let mut completed = use_signal(|| None::<RecordedAudio>);
let mut analyser = use_signal(|| None::<AudioAnalyser>);
let mut elapsed = use_signal(|| Duration::ZERO);
let mut source_availability = use_signal(|| None::<RecordingSourceAvailability>);
let mut requested_constraints = use_signal(|| None::<RecordingConstraints>);
let mut constraint_capabilities = use_signal(|| None::<RecorderConstraintCapabilities>);
let mut settings = use_signal(|| None::<RecordingSourceSettings>);
let mut media_type = use_signal(|| None::<String>);
let mut outcome = use_signal(|| None::<RecordingOutcome>);
let mut chunk_delivery_failure = use_signal(|| None::<RecordingChunkDeliveryFailure>);
let mut microphone = use_signal(|| MicrophoneStatus {
permission: MicrophonePermission::Unknown,
recorder: RecorderStatus::Idle,
input_device: None,
muted: false,
});
let runtime = use_hook(|| Rc::new(RefCell::new(Runtime::default())));
let dioxus_runtime = DioxusRuntime::current();
let dioxus_scope = dioxus_runtime.current_scope_id();
use_effect(move || {
constraint_capabilities.set(read_constraint_capabilities());
});
{
let runtime = Rc::downgrade(&runtime);
use_hook(|| Rc::new(UnmountGuard(runtime)));
}
let runtime_for_start = runtime.clone();
let begin_start = Rc::new(RefCell::new(move |source: RecordingSourceRequest| {
let input_device = match &source {
RecordingSourceRequest::Acquire(input_device) => input_device.clone(),
RecordingSourceRequest::Supplied(_) => None,
};
if let Err(error) = options.validate() {
let accepted = runtime_for_start
.borrow_mut()
.lifecycle
.configuration_failed(error.clone());
if accepted {
status.set(RecorderStatus::Failed(error.clone()));
microphone.set(MicrophoneStatus {
permission: MicrophonePermission::Unknown,
recorder: RecorderStatus::Failed(error),
input_device,
muted: false,
});
}
return Err(command_error("invalid recorder options"));
}
let (input, requested, initial_permission, acquires_device) = match source {
RecordingSourceRequest::Acquire(input_device) => (
RecordingInput::Acquire(input_device),
Some(options.constraints.clone()),
MicrophonePermission::Prompt,
true,
),
RecordingSourceRequest::Supplied(source) => (
RecordingInput::Supplied(prepare_supplied_source(&source)?),
None,
MicrophonePermission::Unknown,
false,
),
};
let initial_source_availability = match &input {
RecordingInput::Acquire(_) => None,
RecordingInput::Supplied(source) => Some(source.availability()),
};
let recording_id = runtime_for_start.borrow_mut().lifecycle.start()?;
let input = match input {
RecordingInput::Acquire(input_device) => AcceptedRecordingInput::Acquire(input_device),
RecordingInput::Supplied(source) => {
AcceptedRecordingInput::Supplied(PendingCapture::new(source))
}
};
outcome.set(None);
chunk_delivery_failure.set(None);
requested_constraints.set(requested);
settings.set(None);
media_type.set(None);
{
let mut runtime = runtime_for_start.borrow_mut();
runtime.elapsed_ms = 0.0;
runtime.segment_started_at = None;
runtime.last_peak_at = 0.0;
runtime.peaks.clear();
runtime.selected_device = input_device.clone();
runtime.muted = false;
runtime.terminal_error = None;
runtime.acquires_device = acquires_device;
runtime.microphone_permission = initial_permission;
}
analyser.set(None);
elapsed.set(Duration::ZERO);
source_availability.set(initial_source_availability);
publish_status(&runtime_for_start, &mut status, &mut microphone);
let runtime = runtime_for_start.clone();
let options = options.clone();
let dioxus_runtime = dioxus_runtime.clone();
wasm_bindgen_futures::spawn_local(async move {
let result = start_recording(
recording_id,
input,
options.clone(),
&runtime,
status,
completed,
analyser,
elapsed,
source_availability,
microphone,
settings,
media_type,
outcome,
chunk_delivery_failure,
dioxus_runtime.clone(),
dioxus_scope,
)
.await;
if !runtime.borrow().mounted {
return;
}
dioxus_runtime.in_scope(dioxus_scope, || match result {
Ok(())
if matches!(
runtime.borrow().lifecycle.status(),
RecorderStatus::Recording
) =>
{
let runtime_for_timer = runtime.clone();
spawn(async move {
run_timer(
recording_id,
options.peak_interval,
runtime_for_timer,
elapsed,
)
.await;
});
}
Ok(()) => {}
Err(error) => {
let mut signals = RecorderTerminalSignals {
status,
analyser,
source_availability,
microphone,
outcome,
};
fail_start(recording_id, error, &runtime, &mut signals);
}
});
});
Ok(())
}));
let begin_device_start = begin_start.clone();
let start: Callback<(), Result<(), RecorderCommandError>> = use_callback(move |()| {
(begin_device_start.borrow_mut())(RecordingSourceRequest::Acquire(selected_input()))
});
let begin_supplied_start = begin_start.clone();
let start_with_source: Callback<RecordingSource, Result<(), RecorderCommandError>> =
use_callback(move |source| {
(begin_supplied_start.borrow_mut())(RecordingSourceRequest::Supplied(source))
});
let runtime_for_pause = runtime.clone();
let pause: Callback<(), Result<(), RecorderCommandError>> = use_callback(move |()| {
let recorder = runtime_for_pause
.borrow()
.recording
.as_ref()
.map(|recording| recording.recorder.clone())
.ok_or_else(|| command_error("no active recorder"))?;
recorder
.pause()
.map_err(|_| command_error("browser rejected pause"))?;
let mut runtime = runtime_for_pause.borrow_mut();
runtime.lifecycle.pause()?;
runtime.accumulate_elapsed();
elapsed.set(duration_from_ms(runtime.elapsed_ms));
drop(runtime);
publish_status(&runtime_for_pause, &mut status, &mut microphone);
Ok(())
});
let runtime_for_resume = runtime.clone();
let resume: Callback<(), Result<(), RecorderCommandError>> = use_callback(move |()| {
let recorder = runtime_for_resume
.borrow()
.recording
.as_ref()
.map(|recording| recording.recorder.clone())
.ok_or_else(|| command_error("no active recorder"))?;
recorder
.resume()
.map_err(|_| command_error("browser rejected resume"))?;
let mut runtime = runtime_for_resume.borrow_mut();
runtime.lifecycle.resume()?;
runtime.segment_started_at = Some(now_ms());
drop(runtime);
publish_status(&runtime_for_resume, &mut status, &mut microphone);
Ok(())
});
let runtime_for_boundary = runtime.clone();
let request_chunk_boundary: Callback<(), Result<(), RecorderCommandError>> =
use_callback(move |()| {
let recorder = {
let runtime = runtime_for_boundary.borrow();
runtime.lifecycle.request_chunk_boundary()?;
let recording = runtime
.recording
.as_ref()
.ok_or_else(|| command_error("no active recorder"))?;
if recording.chunk_delivery.is_none() {
return Err(command_error(
"Recording Chunk delivery was not enabled when recording started",
));
}
recording.recorder.clone()
};
recorder
.request_data()
.map_err(|_| command_error("browser rejected the chunk boundary request"))
});
let runtime_for_stop = runtime.clone();
let stop: Callback<(), Result<(), RecorderCommandError>> = use_callback(move |()| {
let mut signals = RecorderTerminalSignals {
status,
analyser,
source_availability,
microphone,
outcome,
};
stop_or_cancel(false, &runtime_for_stop, &mut elapsed, &mut signals)
});
let runtime_for_cancel = runtime.clone();
let cancel: Callback<(), Result<(), RecorderCommandError>> = use_callback(move |()| {
let mut signals = RecorderTerminalSignals {
status,
analyser,
source_availability,
microphone,
outcome,
};
stop_or_cancel(true, &runtime_for_cancel, &mut elapsed, &mut signals)
});
let runtime_for_clear = runtime.clone();
let clear_completed = use_callback(move |()| {
completed.set(None);
runtime_for_clear.borrow_mut().lifecycle.clear_completed();
});
let runtime_for_take = runtime.clone();
let take_completed = use_callback(move |()| {
let value = completed.write().take();
if value.is_some() {
runtime_for_take.borrow_mut().lifecycle.clear_completed();
}
value
});
AudioRecorder {
status: status.into(),
completed: completed.into(),
analyser: analyser.into(),
elapsed: elapsed.into(),
source_availability: source_availability.into(),
microphone: microphone.into(),
requested_constraints: requested_constraints.into(),
constraint_capabilities: constraint_capabilities.into(),
settings: settings.into(),
media_type: media_type.into(),
outcome: outcome.into(),
chunk_delivery_failure: chunk_delivery_failure.into(),
start,
start_with_source,
pause,
resume,
request_chunk_boundary,
stop,
cancel,
take_completed,
clear_completed,
}
}
struct Runtime {
lifecycle: RecorderLifecycle,
recording: Option<WebRecording>,
elapsed_ms: f64,
segment_started_at: Option<f64>,
last_peak_at: f64,
peaks: Vec<u8>,
selected_device: Option<AudioInputId>,
muted: bool,
terminal_error: Option<AudioError>,
acquires_device: bool,
microphone_permission: MicrophonePermission,
mounted: bool,
}
struct RecorderTerminalSignals {
status: Signal<RecorderStatus>,
analyser: Signal<Option<AudioAnalyser>>,
source_availability: Signal<Option<RecordingSourceAvailability>>,
microphone: Signal<MicrophoneStatus>,
outcome: Signal<Option<RecordingOutcome>>,
}
impl Default for Runtime {
fn default() -> Self {
Self {
lifecycle: RecorderLifecycle::default(),
recording: None,
elapsed_ms: 0.0,
segment_started_at: None,
last_peak_at: 0.0,
peaks: Vec::new(),
selected_device: None,
muted: false,
terminal_error: None,
acquires_device: false,
microphone_permission: MicrophonePermission::Unknown,
mounted: true,
}
}
}
impl Runtime {
fn accumulate_elapsed(&mut self) {
if let Some(started_at) = self.segment_started_at.take() {
self.elapsed_ms += (now_ms() - started_at).max(0.0);
}
}
}
#[derive(Clone)]
struct ChunkDelivery {
state: Rc<RefCell<ChunkDeliveryState>>,
}
struct ChunkDeliveryState {
configuration: RecordingChunkDelivery,
dioxus_runtime: Rc<DioxusRuntime>,
dioxus_scope: ScopeId,
failure_signal: Signal<Option<RecordingChunkDeliveryFailure>>,
pending: VecDeque<PendingChunk>,
worker_active: bool,
enabled: bool,
drain_waiters: Vec<js_sys::Function>,
}
struct PendingChunk {
recording_id: RecordingId,
sequence: u64,
blob: Blob,
media_type: String,
}
impl ChunkDelivery {
fn new(
configuration: RecordingChunkDelivery,
dioxus_runtime: Rc<DioxusRuntime>,
dioxus_scope: ScopeId,
failure_signal: Signal<Option<RecordingChunkDeliveryFailure>>,
) -> Self {
Self {
state: Rc::new(RefCell::new(ChunkDeliveryState {
configuration,
dioxus_runtime,
dioxus_scope,
failure_signal,
pending: VecDeque::new(),
worker_active: false,
enabled: true,
drain_waiters: Vec::new(),
})),
}
}
fn enqueue(&self, chunk: PendingChunk) {
let should_start = {
let mut state = self.state.borrow_mut();
if !state.enabled {
return;
}
state.pending.push_back(chunk);
if state.worker_active {
false
} else {
state.worker_active = true;
true
}
};
if should_start {
wasm_bindgen_futures::spawn_local(deliver_chunks(self.state.clone()));
}
}
fn invalidate(&self) {
let waiters = {
let mut state = self.state.borrow_mut();
state.enabled = false;
state.pending.clear();
std::mem::take(&mut state.drain_waiters)
};
resolve_waiters(waiters);
}
async fn wait_until_drained(&self) {
let state = self.state.clone();
let promise = js_sys::Promise::new(&mut move |resolve, _| {
let resolve_now = {
let mut state = state.borrow_mut();
if !state.enabled || (!state.worker_active && state.pending.is_empty()) {
true
} else {
state.drain_waiters.push(resolve.clone());
false
}
};
if resolve_now {
let _ = resolve.call0(&JsValue::UNDEFINED);
}
});
let _ = wasm_bindgen_futures::JsFuture::from(promise).await;
}
}
async fn deliver_chunks(state: Rc<RefCell<ChunkDeliveryState>>) {
loop {
let pending = {
let mut state = state.borrow_mut();
if state.enabled {
state.pending.pop_front()
} else {
state.pending.clear();
None
}
};
let Some(pending) = pending else {
let waiters = {
let mut state = state.borrow_mut();
state.worker_active = false;
std::mem::take(&mut state.drain_waiters)
};
resolve_waiters(waiters);
return;
};
let bytes = match blob_to_owned_bytes(&pending.blob).await {
Ok(bytes) if !bytes.is_empty() => bytes,
result => {
let error = result.err().unwrap_or_else(|| {
AudioError::new(
AudioErrorKind::Backend,
"Recording Chunk conversion produced no data",
)
});
let publication = {
let mut state = state.borrow_mut();
if !state.enabled {
None
} else {
state.enabled = false;
state.pending.clear();
Some((
state.dioxus_runtime.clone(),
state.dioxus_scope,
state.failure_signal,
std::mem::take(&mut state.drain_waiters),
))
}
};
if let Some((dioxus_runtime, dioxus_scope, mut failure_signal, waiters)) =
publication
{
resolve_waiters(waiters);
dioxus_runtime.in_scope(dioxus_scope, || {
failure_signal.set(Some(RecordingChunkDeliveryFailure {
recording_id: pending.recording_id,
failed_sequence: pending.sequence,
error,
}));
});
}
state.borrow_mut().worker_active = false;
return;
}
};
let callback = {
let state = state.borrow();
state.enabled.then(|| {
(
state.configuration.clone(),
state.dioxus_runtime.clone(),
state.dioxus_scope,
)
})
};
if let Some((callback, dioxus_runtime, dioxus_scope)) = callback {
dioxus_runtime.in_scope(dioxus_scope, || {
callback.call(RecordingChunk {
recording_id: pending.recording_id,
sequence: pending.sequence,
bytes,
media_type: pending.media_type,
});
});
}
}
}
fn resolve_waiters(waiters: Vec<js_sys::Function>) {
for resolve in waiters {
let _ = resolve.call0(&JsValue::UNDEFINED);
}
}
async fn blob_to_owned_bytes(blob: &Blob) -> Result<Vec<u8>, AudioError> {
let buffer = wasm_bindgen_futures::JsFuture::from(blob.array_buffer())
.await
.map_err(audio_error_from_js)?;
Ok(Uint8Array::new(&buffer).to_vec())
}
struct UnmountGuard(Weak<RefCell<Runtime>>);
impl Drop for UnmountGuard {
fn drop(&mut self) {
if let Some(runtime) = self.0.upgrade() {
let mut runtime = runtime.borrow_mut();
runtime.mounted = false;
runtime.recording.take();
runtime.lifecycle.abandon();
}
}
}
struct PreparedSource {
stream: MediaStream,
track: MediaStreamTrack,
track_cleanup: TrackCleanup,
settings: Option<RecordingSourceSettings>,
}
impl PreparedSource {
fn availability(&self) -> RecordingSourceAvailability {
track_availability(&self.track)
}
}
fn track_availability(track: &MediaStreamTrack) -> RecordingSourceAvailability {
if track.muted() {
RecordingSourceAvailability::Interrupted
} else {
RecordingSourceAvailability::Live
}
}
#[derive(Clone, Copy)]
enum TrackCleanup {
Preserve,
Stop,
}
impl TrackCleanup {
fn apply(self, track: &MediaStreamTrack) {
if matches!(self, Self::Stop) {
track.stop();
}
}
}
struct PendingCapture {
stream: Option<MediaStream>,
track: Option<MediaStreamTrack>,
track_cleanup: TrackCleanup,
settings: Option<RecordingSourceSettings>,
context: Option<AudioContext>,
}
impl PendingCapture {
fn new(source: PreparedSource) -> Self {
Self {
stream: Some(source.stream),
track: Some(source.track),
track_cleanup: source.track_cleanup,
settings: source.settings,
context: None,
}
}
fn track(&self) -> &MediaStreamTrack {
self.track.as_ref().expect("pending capture owns its track")
}
fn availability(&self) -> RecordingSourceAvailability {
track_availability(self.track())
}
fn settings(&self) -> Option<RecordingSourceSettings> {
self.settings.clone()
}
fn into_parts(mut self) -> (MediaStream, MediaStreamTrack, TrackCleanup, AudioContext) {
(
self.stream.take().expect("pending capture owns its stream"),
self.track.take().expect("pending capture owns its track"),
self.track_cleanup,
self.context
.take()
.expect("pending capture owns its context"),
)
}
}
impl Drop for PendingCapture {
fn drop(&mut self) {
if let Some(track) = self.track.take() {
self.track_cleanup.apply(&track);
}
if let Some(context) = self.context.take() {
settle_audio_promise(context.close());
}
}
}
struct WebRecording {
recorder: MediaRecorder,
start_resolver: Rc<RefCell<Option<js_sys::Function>>>,
_stream: MediaStream,
context: AudioContext,
_source: MediaStreamAudioSourceNode,
analyser: AudioAnalyserControl,
chunks: Rc<RefCell<Vec<Blob>>>,
chunk_delivery: Option<ChunkDelivery>,
_on_data: Closure<dyn FnMut(BlobEvent)>,
_on_stop: Closure<dyn FnMut()>,
_on_error: Closure<dyn FnMut()>,
track: MediaStreamTrack,
track_cleanup: TrackCleanup,
on_mute: Closure<dyn FnMut()>,
on_unmute: Closure<dyn FnMut()>,
on_ended: Closure<dyn FnMut()>,
}
impl Drop for WebRecording {
fn drop(&mut self) {
settle_start(&self.start_resolver);
if let Some(delivery) = &self.chunk_delivery {
delivery.invalidate();
}
self.analyser.set_available(false);
self.recorder.set_onstart(None);
self.recorder.set_ondataavailable(None);
self.recorder.set_onstop(None);
self.recorder.set_onerror(None);
let _ = self.recorder.stop();
let _ = self
.track
.remove_event_listener_with_callback("mute", self.on_mute.as_ref().unchecked_ref());
let _ = self
.track
.remove_event_listener_with_callback("unmute", self.on_unmute.as_ref().unchecked_ref());
let _ = self
.track
.remove_event_listener_with_callback("ended", self.on_ended.as_ref().unchecked_ref());
self.track_cleanup.apply(&self.track);
settle_audio_promise(self.context.close());
}
}
#[allow(clippy::too_many_arguments)]
async fn start_recording(
recording_id: RecordingId,
input: AcceptedRecordingInput,
options: RecorderOptions,
runtime: &Rc<RefCell<Runtime>>,
mut status: Signal<RecorderStatus>,
mut completed: Signal<Option<RecordedAudio>>,
mut analyser_signal: Signal<Option<AudioAnalyser>>,
mut elapsed: Signal<Duration>,
mut source_availability_signal: Signal<Option<RecordingSourceAvailability>>,
mut microphone: Signal<MicrophoneStatus>,
mut settings_signal: Signal<Option<RecordingSourceSettings>>,
mut media_type_signal: Signal<Option<String>>,
mut outcome_signal: Signal<Option<RecordingOutcome>>,
chunk_delivery_failure: Signal<Option<RecordingChunkDeliveryFailure>>,
dioxus_runtime: Rc<DioxusRuntime>,
dioxus_scope: ScopeId,
) -> Result<(), AudioError> {
let mut pending = match input {
AcceptedRecordingInput::Acquire(input_device) => PendingCapture::new(
prepare_acquired_source(input_device.as_ref(), &options.constraints).await?,
),
AcceptedRecordingInput::Supplied(pending) => pending,
};
let initial_availability = pending.availability();
if !runtime.borrow().mounted
|| runtime.borrow().lifecycle.active_recording != Some(recording_id)
{
return Ok(());
}
dioxus_runtime.in_scope(dioxus_scope, || {
source_availability_signal.set(Some(initial_availability));
});
let context = AudioContext::new().map_err(audio_error_from_js)?;
pending.context = Some(context.clone());
settle_audio_promise(context.resume());
let analyser_node = context.create_analyser().map_err(audio_error_from_js)?;
analyser_node.set_fft_size(options.fft_size);
analyser_node.set_smoothing_time_constant(options.smoothing);
let source = context
.create_media_stream_source(
pending
.stream
.as_ref()
.expect("pending capture owns its stream"),
)
.map_err(audio_error_from_js)?;
source
.connect_with_audio_node(&analyser_node)
.map_err(audio_error_from_js)?;
let recorder_options = MediaRecorderOptions::new();
if let Some(mime_type) = options
.mime_types
.iter()
.find(|mime_type| MediaRecorder::is_type_supported(mime_type))
{
recorder_options.set_mime_type(mime_type);
}
if let Some(bits_per_second) = options.audio_bits_per_second {
recorder_options.set_audio_bits_per_second(bits_per_second);
}
let recorder = MediaRecorder::new_with_media_stream_and_media_recorder_options(
pending
.stream
.as_ref()
.expect("pending capture owns its stream"),
&recorder_options,
)
.map_err(audio_error_from_js)?;
let track = pending.track().clone();
let settings = pending.settings();
let chunks = Rc::new(RefCell::new(Vec::<Blob>::new()));
let chunk_delivery = options.chunk_delivery.clone().map(|delivery| {
ChunkDelivery::new(
delivery,
dioxus_runtime.clone(),
dioxus_scope,
chunk_delivery_failure,
)
});
let start_resolver = Rc::new(RefCell::new(None::<js_sys::Function>));
let resolver_for_promise = start_resolver.clone();
let start_promise = js_sys::Promise::new(&mut move |resolve, _| {
resolver_for_promise.borrow_mut().replace(resolve);
});
let start_succeeded = Rc::new(Cell::new(false));
let resolver_for_start = start_resolver.clone();
let succeeded_for_start = start_succeeded.clone();
let on_start = Closure::wrap(Box::new(move || {
succeeded_for_start.set(true);
settle_start(&resolver_for_start);
}) as Box<dyn FnMut()>);
recorder.set_onstart(Some(on_start.as_ref().unchecked_ref()));
let chunks_for_data = chunks.clone();
let delivery_for_data = chunk_delivery.clone();
let runtime_for_data = Rc::downgrade(runtime);
let recorder_for_data = recorder.clone();
let on_data = Closure::wrap(Box::new(move |event: BlobEvent| {
if let Some(blob) = event.data()
&& blob.size() > 0.0
{
chunks_for_data.borrow_mut().push(blob.clone());
let sequence = runtime_for_data.upgrade().and_then(|runtime| {
runtime
.borrow_mut()
.lifecycle
.next_chunk_sequence(recording_id)
});
if let (Some(delivery), Some(sequence)) = (&delivery_for_data, sequence) {
delivery.enqueue(PendingChunk {
recording_id,
sequence,
blob,
media_type: recorder_for_data.mime_type(),
});
}
}
}) as Box<dyn FnMut(BlobEvent)>);
recorder.set_ondataavailable(Some(on_data.as_ref().unchecked_ref()));
let runtime_for_stop = Rc::downgrade(runtime);
let recorder_for_stop = recorder.clone();
let dioxus_runtime_for_stop = dioxus_runtime.clone();
let resolver_for_stop = start_resolver.clone();
let on_stop = Closure::wrap(Box::new(move || {
settle_start(&resolver_for_stop);
dioxus_runtime_for_stop.in_scope(dioxus_scope, || {
let Some(runtime) = runtime_for_stop.upgrade() else {
return;
};
let (
disposition,
completion_cause,
terminal_error,
chunks,
chunk_delivery,
duration,
peaks,
selected_device,
mime_type,
) = {
let mut runtime = runtime.borrow_mut();
if matches!(
runtime.lifecycle.status(),
RecorderStatus::Recording | RecorderStatus::Paused
) {
runtime.accumulate_elapsed();
}
let disposition = runtime.lifecycle.begin_finalize(recording_id);
let completion_cause = runtime.lifecycle.completion_cause(recording_id);
(
disposition,
completion_cause,
runtime.terminal_error.take(),
runtime
.recording
.as_ref()
.map(|recording| recording.chunks.borrow().clone())
.unwrap_or_default(),
runtime
.recording
.as_ref()
.and_then(|recording| recording.chunk_delivery.clone()),
duration_from_ms(runtime.elapsed_ms),
runtime.peaks.clone(),
runtime.selected_device.clone(),
recorder_for_stop.mime_type(),
)
};
let Some(disposition) = disposition else {
return;
};
if disposition == CompletionDisposition::Save {
debug_assert!(completion_cause.is_some());
}
source_availability_signal.set(None);
analyser_signal.set(None);
publish_status(&runtime, &mut status, &mut microphone);
spawn(async move {
gloo_timers::future::TimeoutFuture::new(0).await;
if terminal_error.is_none()
&& disposition == CompletionDisposition::Save
&& let Some(delivery) = &chunk_delivery
{
delivery.wait_until_drained().await;
}
if !runtime.borrow().mounted {
return;
}
runtime.borrow_mut().recording.take();
if let Some(error) = terminal_error {
let mut signals = RecorderTerminalSignals {
status,
analyser: analyser_signal,
source_availability: source_availability_signal,
microphone,
outcome: outcome_signal,
};
publish_recording_failure(recording_id, error, &runtime, &mut signals);
return;
}
if disposition == CompletionDisposition::Discard {
let mut runtime_mut = runtime.borrow_mut();
if !runtime_mut.lifecycle.complete_finalize(recording_id)
|| !runtime_mut.mounted
{
return;
}
runtime_mut.muted = false;
drop(runtime_mut);
outcome_signal.set(Some(RecordingOutcome::Discarded(recording_id)));
publish_status(&runtime, &mut status, &mut microphone);
return;
}
match collect_audio(chunks, mime_type).await {
Ok(audio) => {
let mut runtime_mut = runtime.borrow_mut();
if !runtime_mut.lifecycle.complete_finalize(recording_id)
|| !runtime_mut.mounted
{
return;
}
runtime_mut.muted = false;
drop(runtime_mut);
completed.set(Some(RecordedAudio {
recording_id,
audio,
duration,
peaks,
input_device: selected_device.clone(),
}));
outcome_signal.set(Some(RecordingOutcome::Completed {
recording_id,
cause: completion_cause
.expect("saved Recording finalization has a completion cause"),
}));
publish_status(&runtime, &mut status, &mut microphone);
}
Err(error) => {
let mut signals = RecorderTerminalSignals {
status,
analyser: analyser_signal,
source_availability: source_availability_signal,
microphone,
outcome: outcome_signal,
};
publish_recording_failure(recording_id, error, &runtime, &mut signals);
}
}
});
});
}) as Box<dyn FnMut()>);
recorder.set_onstop(Some(on_stop.as_ref().unchecked_ref()));
let runtime_for_error = Rc::downgrade(runtime);
let dioxus_runtime_for_error = dioxus_runtime.clone();
let resolver_for_error = start_resolver.clone();
let on_error = Closure::wrap(Box::new(move || {
settle_start(&resolver_for_error);
dioxus_runtime_for_error.in_scope(dioxus_scope, || {
let Some(runtime) = runtime_for_error.upgrade() else {
return;
};
let should_finalize = {
let mut runtime = runtime.borrow_mut();
let active = runtime.lifecycle.active_recording == Some(recording_id);
let status = runtime.lifecycle.status().clone();
if active
&& (matches!(status, RecorderStatus::Recording | RecorderStatus::Paused)
|| (matches!(status, RecorderStatus::Stopping)
&& runtime.lifecycle.completion_cause(recording_id).is_some()))
{
if let Some(delivery) = runtime
.recording
.as_ref()
.and_then(|recording| recording.chunk_delivery.as_ref())
{
delivery.invalidate();
}
runtime.terminal_error = Some(AudioError::new(
AudioErrorKind::RecorderFailure,
"media recorder failed",
));
if matches!(status, RecorderStatus::Recording | RecorderStatus::Paused) {
runtime.accumulate_elapsed();
runtime.lifecycle.stop().is_ok()
} else {
false
}
} else {
false
}
};
if should_finalize {
analyser_signal.set(None);
publish_status(&runtime, &mut status, &mut microphone);
}
});
}) as Box<dyn FnMut()>);
recorder.set_onerror(Some(on_error.as_ref().unchecked_ref()));
let runtime_for_mute = Rc::downgrade(runtime);
let dioxus_runtime_for_mute = dioxus_runtime.clone();
let mut source_availability_for_mute = source_availability_signal;
let on_mute = Closure::wrap(Box::new(move || {
dioxus_runtime_for_mute.in_scope(dioxus_scope, || {
if let Some(runtime) = runtime_for_mute.upgrade() {
let mut runtime_mut = runtime.borrow_mut();
if matches!(
runtime_mut.lifecycle.status(),
RecorderStatus::Preparing | RecorderStatus::Recording | RecorderStatus::Paused
) {
if runtime_mut.acquires_device {
runtime_mut.muted = true;
}
drop(runtime_mut);
source_availability_for_mute
.set(Some(RecordingSourceAvailability::Interrupted));
publish_status(&runtime, &mut status, &mut microphone);
}
}
});
}) as Box<dyn FnMut()>);
let _ = track.add_event_listener_with_callback("mute", on_mute.as_ref().unchecked_ref());
let runtime_for_unmute = Rc::downgrade(runtime);
let dioxus_runtime_for_unmute = dioxus_runtime.clone();
let mut source_availability_for_unmute = source_availability_signal;
let on_unmute = Closure::wrap(Box::new(move || {
dioxus_runtime_for_unmute.in_scope(dioxus_scope, || {
if let Some(runtime) = runtime_for_unmute.upgrade() {
let mut runtime_mut = runtime.borrow_mut();
if matches!(
runtime_mut.lifecycle.status(),
RecorderStatus::Preparing | RecorderStatus::Recording | RecorderStatus::Paused
) {
if runtime_mut.acquires_device {
runtime_mut.muted = false;
}
drop(runtime_mut);
source_availability_for_unmute.set(Some(RecordingSourceAvailability::Live));
publish_status(&runtime, &mut status, &mut microphone);
}
}
});
}) as Box<dyn FnMut()>);
let _ = track.add_event_listener_with_callback("unmute", on_unmute.as_ref().unchecked_ref());
let runtime_for_ended = Rc::downgrade(runtime);
let dioxus_runtime_for_ended = dioxus_runtime.clone();
let recorder_for_ended = recorder.clone();
let mut source_availability_for_ended = source_availability_signal;
let on_ended = Closure::wrap(Box::new(move || {
dioxus_runtime_for_ended.in_scope(dioxus_scope, || {
if let Some(runtime) = runtime_for_ended.upgrade() {
let should_stop = {
let mut runtime = runtime.borrow_mut();
if matches!(
runtime.lifecycle.status(),
RecorderStatus::Recording | RecorderStatus::Paused
) {
runtime.accumulate_elapsed();
let accepted = runtime.lifecycle.source_ended();
if accepted && runtime.acquires_device {
runtime.muted = true;
}
accepted
} else {
false
}
};
if should_stop {
source_availability_for_ended.set(None);
}
publish_status(&runtime, &mut status, &mut microphone);
if should_stop {
let _ = recorder_for_ended.stop();
}
}
});
}) as Box<dyn FnMut()>);
let _ = track.add_event_listener_with_callback("ended", on_ended.as_ref().unchecked_ref());
let initially_muted = track.muted();
let (stream, track, track_cleanup, context) = pending.into_parts();
let (analyser_control, analyser) =
AudioAnalyserControl::new(analyser_node, context.sample_rate());
analyser_control.set_available(true);
let recording = WebRecording {
recorder: recorder.clone(),
start_resolver: start_resolver.clone(),
_stream: stream,
context,
_source: source,
analyser: analyser_control,
chunks,
chunk_delivery,
_on_data: on_data,
_on_stop: on_stop,
_on_error: on_error,
track,
track_cleanup,
on_mute,
on_unmute,
on_ended,
};
let start_result = if let Some(delivery) = options.chunk_delivery.as_ref() {
recording
.recorder
.start_with_time_slice(delivery.time_slice_millis())
} else {
recording.recorder.start()
};
if let Err(error) = start_result {
drop(recording);
return Err(audio_error_from_js(error));
}
{
let mut runtime_mut = runtime.borrow_mut();
if runtime_mut.lifecycle.active_recording != Some(recording_id) || !runtime_mut.mounted {
drop(runtime_mut);
drop(recording);
return Ok(());
}
runtime_mut.recording = Some(recording);
}
let _ = wasm_bindgen_futures::JsFuture::from(start_promise).await;
recorder.set_onstart(None);
drop(on_start);
if !start_succeeded.get() {
return Err(AudioError::new(
AudioErrorKind::RecorderFailure,
"media recorder failed before starting",
));
}
let media_type = recorder.mime_type();
let mut runtime_mut = runtime.borrow_mut();
if !runtime_mut.lifecycle.started(recording_id) || !runtime_mut.mounted {
let recording = runtime_mut.recording.take();
drop(runtime_mut);
drop(recording);
return Ok(());
}
runtime_mut.segment_started_at = Some(now_ms());
runtime_mut.last_peak_at = now_ms();
runtime_mut.muted = runtime_mut.acquires_device && initially_muted;
if runtime_mut.acquires_device {
runtime_mut.microphone_permission = MicrophonePermission::Granted;
}
drop(runtime_mut);
dioxus_runtime.in_scope(dioxus_scope, || {
analyser_signal.set(Some(analyser));
elapsed.set(Duration::ZERO);
settings_signal.set(settings);
media_type_signal.set(Some(media_type));
publish_status(runtime, &mut status, &mut microphone);
});
Ok(())
}
#[derive(Clone, Copy)]
enum SourceTrackError {
AudioTrackCount,
Ended,
InvalidTrack,
}
impl SourceTrackError {
fn command_message(self) -> &'static str {
match self {
Self::AudioTrackCount => "Recording Source must contain exactly one audio track",
Self::Ended => "Recording Source audio track must be live",
Self::InvalidTrack => "Recording Source contains an invalid audio track",
}
}
}
fn prepare_supplied_source(
source: &RecordingSource,
) -> Result<PreparedSource, RecorderCommandError> {
let track = single_live_audio_track(&source.stream)
.map_err(|error| command_error(error.command_message()))?;
let stream = audio_only_stream(&track)
.map_err(|_| command_error("browser rejected the Recording Source"))?;
Ok(PreparedSource {
stream,
track,
track_cleanup: match source.shutdown {
RecordingSourceShutdown::PreserveTracks => TrackCleanup::Preserve,
RecordingSourceShutdown::StopAudioTracks => TrackCleanup::Stop,
},
settings: None,
})
}
async fn prepare_acquired_source(
input_device: Option<&AudioInputId>,
requested: &RecordingConstraints,
) -> Result<PreparedSource, AudioError> {
let acquired = acquire_stream(input_device, requested).await?;
let track = match single_live_audio_track(&acquired) {
Ok(track) => track,
Err(error) => {
stop_stream(&acquired);
let kind = match error {
SourceTrackError::AudioTrackCount => AudioErrorKind::DeviceNotFound,
SourceTrackError::Ended => AudioErrorKind::DeviceUnavailable,
SourceTrackError::InvalidTrack => AudioErrorKind::Backend,
};
return Err(AudioError::new(kind, error.command_message()));
}
};
let stream = match audio_only_stream(&track) {
Ok(stream) => stream,
Err(error) => {
stop_stream(&acquired);
return Err(audio_error_from_js(error));
}
};
let settings = Some(read_settings(&track));
Ok(PreparedSource {
stream,
track,
track_cleanup: TrackCleanup::Stop,
settings,
})
}
fn single_live_audio_track(stream: &MediaStream) -> Result<MediaStreamTrack, SourceTrackError> {
let tracks = stream.get_audio_tracks();
if tracks.length() != 1 {
return Err(SourceTrackError::AudioTrackCount);
}
let track = tracks
.get(0)
.dyn_into::<MediaStreamTrack>()
.map_err(|_| SourceTrackError::InvalidTrack)?;
if track.ready_state() != MediaStreamTrackState::Live {
return Err(SourceTrackError::Ended);
}
Ok(track)
}
fn audio_only_stream(track: &MediaStreamTrack) -> Result<MediaStream, JsValue> {
let tracks = Array::new();
tracks.push(track.as_ref());
MediaStream::new_with_tracks(tracks.as_ref())
}
async fn acquire_stream(
input_device: Option<&AudioInputId>,
requested: &RecordingConstraints,
) -> Result<MediaStream, AudioError> {
let constraints = web_sys::MediaStreamConstraints::new();
let audio = web_sys::MediaTrackConstraints::new();
if let Some(input_device) = input_device {
let exact = web_sys::ConstrainDomStringParameters::new();
exact.set_exact_str(input_device.as_str());
audio.set_device_id_constrain_dom_string_parameters(&exact);
}
set_constraint(
&audio,
"channelCount",
requested.channel_count.as_ref(),
|value| JsValue::from_f64(*value as f64),
)?;
set_constraint(
&audio,
"sampleRate",
requested.sample_rate.as_ref(),
|value| JsValue::from_f64(*value as f64),
)?;
set_constraint(
&audio,
"echoCancellation",
requested.echo_cancellation.as_ref(),
|value| JsValue::from_bool(*value),
)?;
set_constraint(
&audio,
"noiseSuppression",
requested.noise_suppression.as_ref(),
|value| JsValue::from_bool(*value),
)?;
set_constraint(&audio, "latency", requested.latency.as_ref(), |value| {
JsValue::from_f64(value.as_secs_f64())
})?;
constraints.set_audio_media_track_constraints(&audio);
let value = wasm_bindgen_futures::JsFuture::from(
media_devices()?
.get_user_media_with_constraints(&constraints)
.map_err(audio_error_from_js)?,
)
.await
.map_err(audio_error_from_js)?;
value.dyn_into::<MediaStream>().map_err(audio_error_from_js)
}
fn read_constraint_capabilities() -> Option<RecorderConstraintCapabilities> {
let supported = media_devices().ok()?.get_supported_constraints();
Some(RecorderConstraintCapabilities {
channel_count: supported.get_channel_count().unwrap_or(false),
sample_rate: supported.get_sample_rate().unwrap_or(false),
echo_cancellation: supported.get_echo_cancellation().unwrap_or(false),
noise_suppression: supported.get_noise_suppression().unwrap_or(false),
latency: supported.get_latency().unwrap_or(false),
})
}
fn read_settings(track: &MediaStreamTrack) -> RecordingSourceSettings {
let settings = track.get_settings();
let settings = settings.as_ref();
RecordingSourceSettings {
channel_count: read_u32(settings, "channelCount"),
sample_rate: read_u32(settings, "sampleRate"),
echo_cancellation: read_bool(settings, "echoCancellation"),
noise_suppression: read_bool(settings, "noiseSuppression"),
latency: read_number(settings, "latency").and_then(|seconds| {
if seconds.is_finite() && seconds >= 0.0 {
Duration::try_from_secs_f64(seconds).ok()
} else {
None
}
}),
}
}
fn read_u32(value: &JsValue, field: &str) -> Option<u32> {
let number = read_number(value, field)?;
if number.is_finite() && number >= 0.0 && number <= u32::MAX as f64 && number.fract() == 0.0 {
Some(number as u32)
} else {
None
}
}
fn read_number(value: &JsValue, field: &str) -> Option<f64> {
Reflect::get(value, &JsValue::from_str(field))
.ok()?
.as_f64()
}
fn read_bool(value: &JsValue, field: &str) -> Option<bool> {
Reflect::get(value, &JsValue::from_str(field))
.ok()?
.as_bool()
}
fn set_constraint<T>(
target: &web_sys::MediaTrackConstraints,
field: &str,
constraint: Option<&RecordingConstraint<T>>,
into_js: impl FnOnce(&T) -> JsValue,
) -> Result<(), AudioError> {
if let Some(constraint) = constraint {
let (kind, value) = match constraint {
RecordingConstraint::Ideal(value) => ("ideal", value),
RecordingConstraint::Exact(value) => ("exact", value),
};
set_js_constraint(target, field, kind, into_js(value))?;
}
Ok(())
}
fn set_js_constraint(
target: &web_sys::MediaTrackConstraints,
field: &str,
kind: &str,
value: JsValue,
) -> Result<(), AudioError> {
let constraint = js_sys::Object::new();
Reflect::set(&constraint, &JsValue::from_str(kind), &value).map_err(audio_error_from_js)?;
Reflect::set(
target.as_ref(),
&JsValue::from_str(field),
constraint.as_ref(),
)
.map_err(audio_error_from_js)?;
Ok(())
}
fn settle_start(resolver: &Rc<RefCell<Option<js_sys::Function>>>) {
if let Some(resolve) = resolver.borrow_mut().take() {
let _ = resolve.call0(&JsValue::UNDEFINED);
}
}
async fn run_timer(
recording_id: RecordingId,
peak_interval: Duration,
runtime: Rc<RefCell<Runtime>>,
mut elapsed: Signal<Duration>,
) {
loop {
gloo_timers::future::TimeoutFuture::new(30).await;
let mut runtime = runtime.borrow_mut();
if runtime.lifecycle.active_recording != Some(recording_id) {
break;
}
if !matches!(runtime.lifecycle.status(), RecorderStatus::Recording) {
continue;
}
let now = now_ms();
let current_ms = runtime.elapsed_ms
+ runtime
.segment_started_at
.map(|start| (now - start).max(0.0))
.unwrap_or(0.0);
elapsed.set(duration_from_ms(current_ms));
if now - runtime.last_peak_at >= peak_interval.as_secs_f64() * 1000.0 {
runtime.last_peak_at = now;
if let Some(recording) = runtime.recording.as_ref() {
let mut samples = vec![0_u8; recording.analyser.node().fft_size() as usize];
recording
.analyser
.node()
.get_byte_time_domain_data(&mut samples);
let peak = peak_amplitude(&samples);
runtime.peaks.push(peak);
}
}
}
}
fn stop_or_cancel(
cancel: bool,
runtime: &Rc<RefCell<Runtime>>,
elapsed: &mut Signal<Duration>,
signals: &mut RecorderTerminalSignals,
) -> Result<(), RecorderCommandError> {
let mut runtime_mut = runtime.borrow_mut();
let cancelled_recording = cancel
.then_some(runtime_mut.lifecycle.active_recording)
.flatten();
if cancel {
runtime_mut.lifecycle.cancel()?;
if let Some(delivery) = runtime_mut
.recording
.as_ref()
.and_then(|recording| recording.chunk_delivery.as_ref())
{
delivery.invalidate();
}
} else {
runtime_mut.lifecycle.stop()?;
}
if matches!(runtime_mut.lifecycle.status(), RecorderStatus::Idle) {
let recording = runtime_mut.recording.take();
signals.source_availability.set(None);
if let Some(recording_id) = cancelled_recording {
signals
.outcome
.set(Some(RecordingOutcome::Discarded(recording_id)));
}
signals.status.set(RecorderStatus::Idle);
runtime_mut.microphone_permission = MicrophonePermission::Unknown;
signals.microphone.set(MicrophoneStatus {
permission: runtime_mut.microphone_permission,
recorder: RecorderStatus::Idle,
input_device: runtime_mut.selected_device.clone(),
muted: false,
});
drop(runtime_mut);
drop(recording);
return Ok(());
}
runtime_mut.accumulate_elapsed();
elapsed.set(duration_from_ms(runtime_mut.elapsed_ms));
let recorder = runtime_mut
.recording
.as_ref()
.map(|recording| recording.recorder.clone())
.ok_or_else(|| command_error("no active recorder"))?;
drop(runtime_mut);
publish_status(runtime, &mut signals.status, &mut signals.microphone);
if recorder.stop().is_err() {
let error = AudioError::new(
AudioErrorKind::RecorderFailure,
"browser rejected recording stop",
);
let recording_id = runtime.borrow().lifecycle.active_recording;
if let Some(recording_id) = recording_id {
publish_recording_failure(recording_id, error.clone(), runtime, signals);
}
return Err(command_error("browser rejected stop"));
}
Ok(())
}
fn fail_start(
recording_id: RecordingId,
error: AudioError,
runtime: &Rc<RefCell<Runtime>>,
signals: &mut RecorderTerminalSignals,
) {
{
let mut runtime = runtime.borrow_mut();
if !runtime.mounted || runtime.lifecycle.active_recording != Some(recording_id) {
return;
}
runtime.microphone_permission =
if runtime.acquires_device && error.kind() == AudioErrorKind::PermissionDenied {
MicrophonePermission::Denied
} else {
MicrophonePermission::Unknown
};
}
publish_recording_failure(recording_id, error, runtime, signals);
}
fn publish_recording_failure(
recording_id: RecordingId,
error: AudioError,
runtime: &Rc<RefCell<Runtime>>,
signals: &mut RecorderTerminalSignals,
) -> bool {
let (input_device, permission) = {
let mut runtime = runtime.borrow_mut();
if !runtime.mounted || !runtime.lifecycle.failed(recording_id, error.clone()) {
return false;
}
runtime.muted = false;
runtime.recording.take();
(
runtime.selected_device.clone(),
runtime.microphone_permission,
)
};
signals.outcome.set(Some(RecordingOutcome::Failed {
recording_id,
error: error.clone(),
}));
signals.analyser.set(None);
signals.source_availability.set(None);
signals.status.set(RecorderStatus::Failed(error.clone()));
signals.microphone.set(MicrophoneStatus {
permission,
recorder: RecorderStatus::Failed(error),
input_device,
muted: false,
});
true
}
fn publish_status(
runtime: &Rc<RefCell<Runtime>>,
status: &mut Signal<RecorderStatus>,
microphone: &mut Signal<MicrophoneStatus>,
) {
let runtime = runtime.borrow();
let recorder = runtime.lifecycle.status().clone();
status.set(recorder.clone());
microphone.set(MicrophoneStatus {
permission: runtime.microphone_permission,
recorder,
input_device: runtime.selected_device.clone(),
muted: runtime.muted,
});
}
async fn collect_audio(chunks: Vec<Blob>, mime_type: String) -> Result<AudioData, AudioError> {
if chunks.is_empty() {
return Err(AudioError::new(
AudioErrorKind::RecorderFailure,
"recording produced no audio data",
));
}
let parts = Array::new();
for chunk in chunks {
parts.push(&chunk);
}
let properties = BlobPropertyBag::new();
properties.set_type(&mime_type);
let blob = Blob::new_with_blob_sequence_and_options(&parts, &properties)
.map_err(audio_error_from_js)?;
let buffer = wasm_bindgen_futures::JsFuture::from(blob.array_buffer())
.await
.map_err(audio_error_from_js)?;
let bytes = Uint8Array::new(&buffer).to_vec();
if bytes.is_empty() {
return Err(AudioError::new(
AudioErrorKind::RecorderFailure,
"recording produced empty audio data",
));
}
Ok(AudioData::new(bytes, mime_type))
}
fn now_ms() -> f64 {
web_sys::window()
.and_then(|window| window.performance())
.map(|performance| performance.now())
.unwrap_or_else(js_sys::Date::now)
}
fn duration_from_ms(milliseconds: f64) -> Duration {
Duration::from_secs_f64((milliseconds / 1000.0).max(0.0))
}
fn settle_audio_promise(promise: Result<js_sys::Promise, JsValue>) {
if let Ok(promise) = promise {
wasm_bindgen_futures::spawn_local(async move {
let _ = wasm_bindgen_futures::JsFuture::from(promise).await;
});
}
}