pub mod audio;
pub mod capture;
pub mod display;
pub mod gpu;
pub mod options;
pub mod video;
use std::
{
sync::
{
Arc,
RwLock,
atomic::{ AtomicBool, Ordering },
},
};
use tokio::
{
task,
io::AsyncWriteExt,
net::tcp::OwnedWriteHalf,
sync::
{
Mutex,
mpsc::
{
self,
Sender,
Receiver,
UnboundedSender,
},
},
};
use winit::event_loop::EventLoopProxy;
use crate::
{
crypto,
options as chat_options,
network::
{
self,
client::{ self, ClientEvent },
codes::PacketCode,
screen::
{
self,
consts,
ScreenPacketCode,
client::audio::AudioFrame,
},
},
};
pub struct ScreenShareRequest
{
pub rx: Receiver<Vec<u8>>,
pub running: Arc<AtomicBool>,
pub deattach: UnboundedSender<()>, }
pub enum UserEvent {
NewSession(ScreenShareRequest),
NewFrame,
}
pub static SCREEN_SHARE_PROXY: RwLock<Option<EventLoopProxy<UserEvent>>> = RwLock::new(None);
pub static SCREEN_FRAME_SINK: RwLock<Option<UnboundedSender<Vec<u8>>>> = RwLock::new(None);
pub async fn screen(token: [u8; 32], events: Sender<ClientEvent>)
{
let (_read_stream, mut write_stream) = client::connect(chat_options::get_server_address()).await
.expect("Screen upload connection failed");
screen::cap_socket_buffers(write_stream.as_ref());
write_stream.write_all(&token).await.unwrap();
let (tx, mut rx) = mpsc::channel(consts::MULTIPLEX_CHANNEL_BOUND);
let (audio_tx, mut audio_rx) = mpsc::channel(consts::MULTIPLEX_CHANNEL_BOUND);
let running = Arc::new(AtomicBool::new(true));
let running_capture = running.clone();
let running_audio = running.clone();
let capture = task::spawn_blocking(move || capture::capture_loop(tx, running_capture, consts::TARGET_FPS));
tokio::spawn(audio::spawn_audio_capture(audio_tx, running_audio));
let mut seq = 0usize;
let mut rex_stream = crypto::init_rex_stream(chat_options::get_keys().as_ref().unwrap(), &token).unwrap();
loop
{
if !options::get_use_screen() { break; }
tokio::select!
{
msg = rx.recv() =>
{
let compressed_frame = match msg
{
Some(f) => f,
None => break,
};
screen::send_frame(&mut write_stream,
ScreenPacketCode::Video { data: compressed_frame }, &mut rex_stream, Some(&mut seq)).await;
},
msg = audio_rx.recv() =>
{
let audio_frame = match msg
{
Some(f) => f,
None => break,
};
screen::send_frame(&mut write_stream,
ScreenPacketCode::Audio { data: audio_frame.data }, &mut rex_stream, Some(&mut seq)).await;
}
}
}
running.store(false, Ordering::Relaxed);
let reason = match capture.await
{
Ok(Err(reason)) => Some(reason),
Err(_) => Some("screen capture crashed".to_owned()), Ok(Ok(())) => None,
};
if let Some(reason) = reason
{
events.send(ClientEvent::ScreenFailed(reason)).await.ok();
}
}
pub async fn attach(token: [u8; 32], main_stream: Arc<Mutex<OwnedWriteHalf>>)
{
let (mut read_stream, mut write_stream) = client::connect(chat_options::get_server_address()).await
.expect("Screen download connection failed");
screen::cap_socket_buffers(read_stream.as_ref());
write_stream.write_all(&token).await.unwrap();
let (tx, rx) = mpsc::channel(consts::MULTIPLEX_CHANNEL_BOUND);
let (audio_tx, audio_rx) = mpsc::channel(consts::NETWORK_CHANNEL_BOUND);
let running = Arc::new(AtomicBool::new(true));
let running_audio = running.clone();
tokio::spawn(audio::spawn_audio_playback(audio_rx, running_audio));
let mut rex_stream = crypto::init_rex_stream(chat_options::get_keys().as_ref().unwrap(), &token).unwrap();
let (deattach_tx, mut deattach_rx) = mpsc::unbounded_channel::<()>();
tokio::spawn(async move
{
while deattach_rx.recv().await.is_some()
{
network::send(&mut *main_stream.lock().await,
PacketCode::Deattach { username: None }, chat_options::get_keys().as_ref()).await;
}
});
let running_net = running.clone();
tokio::spawn(async move
{
let mut streams = (&mut read_stream, Arc::new(Mutex::new(write_stream)));
let mut seq = 0usize;
while running_net.load(Ordering::Relaxed)
{
if !options::get_attach_screen()
{
running_net.store(false, Ordering::Relaxed);
return;
}
let read = match screen::receive_frame(&mut streams, &mut rex_stream, &mut seq).await
{
Some(r) => r,
None =>
{
running_net.store(false, Ordering::Relaxed);
return;
}
};
match read
{
ScreenPacketCode::Video { data } =>
{
let sink = SCREEN_FRAME_SINK.read().unwrap().clone();
match sink
{
Some(sink) => { sink.send(data).ok(); },
None =>
{
tx.send(data).await.ok();
if let Some(proxy) = SCREEN_SHARE_PROXY.read().unwrap().as_ref()
{
proxy.send_event(UserEvent::NewFrame).ok();
}
},
}
},
ScreenPacketCode::Audio { data } =>
{
audio_tx.try_send(AudioFrame { data }).ok();
},
}
}
});
if SCREEN_FRAME_SINK.read().unwrap().is_some() { return }
if let Some(proxy) = SCREEN_SHARE_PROXY.read().unwrap().as_ref()
{
proxy.send_event(UserEvent::NewSession(ScreenShareRequest
{
rx,
running,
deattach: deattach_tx,
})).ok();
}
}