diff --git a/src/client_handler.rs b/src/client_handler.rs index b7fe1ba..bd184ab 100644 --- a/src/client_handler.rs +++ b/src/client_handler.rs @@ -1,6 +1,5 @@ use crate::protocol::Message; use crate::sound_scheduler::SoundScheduler; -use std::io::Read; use std::net::{TcpListener, TcpStream}; use std::thread; @@ -9,14 +8,19 @@ pub struct ClientHandler; impl ClientHandler { fn handle_client(mut stream: TcpStream) { println!("Handling client: {:?}", stream); - let mut buf: Vec = Vec::new(); - stream - .read_to_end(buf.as_mut()) - .expect("Could not read from stream"); - let message = bincode::deserialize::(&buf).expect("Could not deserialize message"); - match message { - Message::PlaySound(play_sound) => { - SoundScheduler::handle_scheduled_sound(play_sound); + loop { + match bincode::deserialize_from::<&mut TcpStream, Message>(&mut stream) { + Ok(message) => { + match message { + Message::PlaySound(play_sound) => { + SoundScheduler::handle_scheduled_sound(play_sound); + } + } + } + Err(e) => { + eprintln!("Error deserializing message: {}", e); + break; + } } } } diff --git a/src/dispatcher.rs b/src/dispatcher.rs index abcc54b..6094b06 100644 --- a/src/dispatcher.rs +++ b/src/dispatcher.rs @@ -1,41 +1,60 @@ -use std::io::{BufReader, Read}; +use std::io::BufReader; use std::net::TcpStream; -use std::time::SystemTime; +use std::time::{Duration, SystemTime}; use std::{fs::File, io::Write}; +use rodio::{Decoder, Source}; + use crate::protocol::{Message, PlaySound}; pub struct Dispatcher; impl Dispatcher { - fn load_sound(path: String) -> Vec { + fn load_sound(path: String) -> Decoder> { let file = BufReader::new(File::open(path).unwrap()); - let data: Vec = file.bytes().map(|byte| byte.unwrap()).collect(); - data + Decoder::new(file).unwrap() } pub fn handle_dispatch_sample(addrs: Vec, sound_path: String) { let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) - .unwrap() - .as_micros(); - let timestamp = now + 10000000; - let sound_data = Dispatcher::load_sound(sound_path); - let msg = Message::PlaySound(PlaySound { - timestamp: timestamp, - sound_data: sound_data, - }); + .unwrap(); + + let source = Dispatcher::load_sound(sound_path).buffered(); + + let channels = source.channels(); + let sample_rate = source.sample_rate(); + let duration = source.total_duration().unwrap(); + let chunk_len = Duration::from_secs(1); + + let mut offset = Duration::from_secs(0); + for addr in addrs { if addr.is_empty() { continue; - - } else { - let mut sock = TcpStream::connect(addr).expect("Failed to connect"); + } + + let mut sock = TcpStream::connect(addr).expect("Failed to connect"); + + while offset < duration { + let sound_data = source + .clone() + .skip_duration(offset) + .take_duration(chunk_len) + .collect::>(); + + offset += chunk_len; + + let msg = Message::PlaySound(PlaySound { + timestamp: (now + offset).as_micros(), + channels: channels, + sample_rate: sample_rate, + sound_data: sound_data, + }); + let buf = bincode::serialize(&msg).expect("Failed to serialize"); sock.write_all(&buf).unwrap(); } - } - } } diff --git a/src/protocol.rs b/src/protocol.rs index 3404033..84397d1 100644 --- a/src/protocol.rs +++ b/src/protocol.rs @@ -8,5 +8,7 @@ pub enum Message { #[derive(Serialize, Deserialize, Debug)] pub struct PlaySound { pub timestamp: u128, - pub sound_data: Vec, + pub channels: u16, + pub sample_rate: u32, + pub sound_data: Vec, } diff --git a/src/sound_player.rs b/src/sound_player.rs index d88e239..4bca269 100644 --- a/src/sound_player.rs +++ b/src/sound_player.rs @@ -1,12 +1,11 @@ -use rodio::{Decoder, OutputStream, Sink}; -use std::io::Cursor; +use rodio::{buffer::SamplesBuffer, OutputStream, Sink}; pub struct SoundPlayer; impl SoundPlayer { - pub fn play_sound(sound_data: Vec) { - let data = Cursor::new(sound_data); - let source = Decoder::new(data).unwrap(); + pub fn play_sound(channels: u16, sample_rate: u32, sound_data: Vec) { + // let data = Cursor::new(sound_data); + let source = SamplesBuffer::new(channels, sample_rate, sound_data); let (_stream, stream_handle) = OutputStream::try_default().unwrap(); let sink = Sink::try_new(&stream_handle).unwrap(); diff --git a/src/sound_scheduler.rs b/src/sound_scheduler.rs index d434359..05f29b0 100644 --- a/src/sound_scheduler.rs +++ b/src/sound_scheduler.rs @@ -9,6 +9,8 @@ impl SoundScheduler { pub fn handle_scheduled_sound(play_msg: PlaySound) { thread::spawn(move || { let PlaySound { + channels, + sample_rate, sound_data, timestamp, } = play_msg; @@ -28,7 +30,7 @@ impl SoundScheduler { } sleep(Duration::from_micros(wait as u64)); println!("Playing sound..."); - SoundPlayer::play_sound(sound_data); + SoundPlayer::play_sound(channels, sample_rate, sound_data); }); } }