use crate::tx_thread::TxThread; use crate::udp_connection::OpusUdpConnection; use crate::util::{SAMPLE_RATE, SAMPLE_RATE_RAW, STEREO_FRAME_SIZE, ThreadMessage}; use anyhow::Error; use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; use cpal::{BufferSize, FromSample, Host, Sample, SampleRate, StreamConfig}; use log::{error, info}; use opus::{Application, Bitrate, Channels, Encoder, Signal}; use std::io::{Read, stdin}; use std::net::SocketAddr; use std::sync::mpsc::{Sender, channel}; use std::thread; pub fn tap( device: &str, src_addr: SocketAddr, dest_addr: SocketAddr, host: Host, ) -> Result<(), Error> { // Set up the input device and stream with the default input config. let device = if device == "default" { host.default_input_device() } else { let id = &device.parse().expect("failed to parse input device id"); host.device_by_id(id) } .expect("failed to find input device"); for dev in host.input_devices()? { info!("ID: {}", dev.id()?); } info!("Input device: {}", device.id()?); let config = device .default_input_config() .expect("Failed to get default input config"); let mut stream_config = StreamConfig::from(config); let count = stream_config.channels; stream_config.sample_rate = SAMPLE_RATE_RAW as SampleRate; stream_config.buffer_size = BufferSize::Fixed(STEREO_FRAME_SIZE as u32 / count as u32); for cfg in device.supported_input_configs()? { info!("{:?}", cfg) } info!("Stream config: {stream_config:?}"); let udp_connection = OpusUdpConnection::new(src_addr, dest_addr)?; let mut encoder = Encoder::new(SAMPLE_RATE, Channels::Stereo, Application::Audio)?; encoder.set_bitrate(Bitrate::Max)?; encoder.set_complexity(10)?; encoder.set_signal(Signal::Music)?; let (sender, receiver) = channel(); let tx_thread = TxThread::new(receiver, encoder, udp_connection); thread::spawn(move || tx_thread.thread()); let err_fn = move |err| { error!("an error occurred on stream: {err}"); }; let stream = match config.sample_format() { cpal::SampleFormat::I8 => device.build_input_stream( stream_config, move |data, _: &_| write_input_data::(data, &sender), err_fn, None, )?, cpal::SampleFormat::I16 => device.build_input_stream( stream_config, move |data, _: &_| write_input_data::(data, &sender), err_fn, None, )?, cpal::SampleFormat::I32 => device.build_input_stream( stream_config, move |data, _: &_| write_input_data::(data, &sender), err_fn, None, )?, cpal::SampleFormat::F32 => device.build_input_stream( stream_config, move |data, _: &_| write_input_data::(data, &sender), err_fn, None, )?, sample_format => { return Err(anyhow::Error::msg(format!( "Unsupported sample format '{sample_format}'" ))); } }; info!("Running tap..."); stream.play()?; let mut stdin = stdin().lock(); let mut exit = [0u8; 1]; stdin.read_exact(&mut exit)?; info!("Exiting..."); Ok(()) } fn write_input_data(input: &[T], ctx: &Sender) where T: Sample, f32: FromSample, { let sample: Vec = input.iter().map(|s| f32::from_sample(*s)).collect(); if let Err(err) = ctx.send(ThreadMessage::Data(sample)) { error!("Got error when sending data: {err}") } }