use std::sync::Arc;
use libsync::ReceiveResult;
use tokio::time::{Duration, Instant};
use crossbeam::queue::ArrayQueue;
use fastwebsockets::upgrade::{IncomingUpgrade, UpgradeFut};
use fastwebsockets::{FragmentCollectorRead, Frame, OpCode, WebSocketError};
use tokio::select;
use crate::{ConnectionStateId, ConnectionStateMessage};
use super::{OwnedFrame, ReadWebSocketActorInputMessage, ReadWebSocketActorOutputMessage, WebSocketActorOutputMessage, WebSocketServerActorInputMessage};
use super::WebSocketReader;
use super::WebSocketWriteHalf;
use act_rs::{impl_pre_run_async, impl_post_run_async, impl_pre_and_post_run_async, impl_mac_task_actor, impl_mac_task_actor_built_state};
use paste::paste;
use tokio::task::JoinHandle;
use tokio::time::timeout_at;
use libsync::crossbeam::mpmc::tokio::array_queue::{Sender, Receiver, channel, io_channels::{IOClient, IOServer, io_channels, io_channel_both}};
pub struct WebSocketServerActorState
{
read_web_socket_actor_io_client: IOClient<ConnectionStateMessage<ReadWebSocketActorInputMessage>, ConnectionStateMessage<ReadWebSocketActorOutputMessage>>, io_server: IOServer<ConnectionStateMessage<WebSocketServerActorInputMessage>, ConnectionStateMessage<WebSocketActorOutputMessage>>,
connection_state_id: ConnectionStateId
}
impl WebSocketServerActorState
{
pub fn new(io_server: IOServer<ConnectionStateMessage<WebSocketServerActorInputMessage>, ConnectionStateMessage<WebSocketActorOutputMessage>>) -> Self
{
let reader_actor_io_client = WebSocketServerReaderActorState::spawn(io_server.output_sender_ref());
Self
{
writer,
actor_io_receiver: actor_io_server.input_receiver_ref().clone()
}
}
impl_pre_and_post_run_async!();
async fn run_async(&mut self) -> bool
{
enum SelectResult
{
WriteFrame(ReceiveResult<OwnedFrame>),
FromSimpleWebSocketReaderActor(ReceiveResult<SimpleWebSocketActorInputMessage>)
}
let reader_actor_io_client_recv = self.reader_actor_io_client.output_receiver_ref().recv();
let write_frame_future = self.actor_io_receiver.recv();
let select_result;
select!
{
biased;
res = reader_actor_io_client_recv =>
{
select_result = SelectResult::FromSimpleWebSocketReaderActor(res);
}
res = write_frame_future =>
{
select_result = SelectResult::WriteFrame(res);
}
}
match select_result
{
SelectResult::WriteFrame(opt_owned_frame) =>
{
if let Ok(mut owned_frame) = opt_owned_frame
{
if owned_frame.opcode == OpCode::Close
{
let now = Instant::now();
let soon = now.checked_add(Duration::from_secs(10)).expect("Error: Instant problems");
let reader_actor_io_client_recv = self.reader_actor_io_client.output_receiver_ref().recv();
match timeout_at(soon, reader_actor_io_client_recv).await
{
Ok(opt_res) =>
{
if let Ok(res) = opt_res
{
match res
{
SimpleWebSocketActorInputMessage::Disconnect => {}
SimpleWebSocketActorInputMessage::WriteFrame(mut owned_frame) =>
{
let mut opcode = owned_frame.opcode;
let frame = owned_frame.new_frame_to_be_written();
if let Err(err) = self.writer.write_frame(frame).await
{
print!("{}", err);
}
loop
{
if opcode != OpCode::Close
{
if let ReceiveResult::Ok(message) = self.reader_actor_io_client.output_receiver_ref().try_recv()
{
match message
{
SimpleWebSocketActorInputMessage::Disconnect =>
{
break;
}
SimpleWebSocketActorInputMessage::WriteFrame(mut owned_frame) =>
{
let frame = owned_frame.new_frame_to_be_written();
if let Err(err) = self.writer.write_frame(frame).await
{
print!("{}", err);
break;
}
opcode = owned_frame.opcode;
}
}
}
}
else
{
break;
}
}
}
}
}
}
Err(err) =>
{
print!("{}", err);
}
}
return false;
}
let frame = owned_frame.new_frame_to_be_written();
if let Err(err) = self.writer.write_frame(frame).await
{
print!("{}", err);
}
else
{
return true;
}
}
}
SelectResult::FromSimpleWebSocketReaderActor(opt_simple_web_socket_actor_input_message) =>
{
if let Ok(input_message) = opt_simple_web_socket_actor_input_message
{
match input_message
{
SimpleWebSocketActorInputMessage::Disconnect => {}
SimpleWebSocketActorInputMessage::WriteFrame(mut owned_frame) =>
{
let frame = owned_frame.new_frame_to_be_written();
if let Err(err) = self.writer.write_frame(frame).await
{
print!("{}", err);
}
else
{
let opcode = owned_frame.opcode;
if opcode == OpCode::Close
{
let _ = self.reader_actor_io_client.input_sender_ref().send(());
return false;
}
return true;
}
}
}
}
}
}
false
}
}
impl_mac_task_actor!(WebSocketServerActor);
pub struct WebSocketServerReaderActorState
{
write_frame_processor_actor_io_sender: Sender<ConnectionStateMessage<WebSocketActorOutputMessage>>, connection_state_id: ConnectionStateId,
obligated_send_frame_holder: Arc<ArrayQueue<OwnedFrame>>,
io_server: IOServer<ConnectionStateMessage<ReadWebSocketActorInputMessage>, ConnectionStateMessage<ReadWebSocketActorOutputMessage>>
}
impl WebSocketServerReaderActorState
{
pub fn new(reader: WebSocketReader, actor_io_sender: Sender<OwnedFrame>) -> (Self, IOClient<(), SimpleWebSocketActorInputMessage>) {
let (io_client, io_server) = io_channel_both(1);
(Self
{
reader,
io_server,
obligated_send_frame_holder: Arc::new(ArrayQueue::new(1)),
actor_io_sender
},
io_client)
}
pub fn spawn(reader: WebSocketReader, actor_io_sender: Sender<OwnedFrame>) -> IOClient<(), SimpleWebSocketActorInputMessage> {
let (state, actor_io_client) = SimpleWebSocketReaderActorState::new(reader, actor_io_sender);
SimpleWebSocketReaderActor::spawn(state);
actor_io_client
}
impl_pre_and_post_run_async!();
async fn run_async(&mut self) -> bool
{
enum SelectResult<'f>
{
ReadFrame(Result<Frame<'f>, WebSocketError>),
Connection(ReceiveResult<()>)
}
let obligated_send_frame_holder = self.obligated_send_frame_holder.clone();
let mut send_fn = |obligated_send_frame: Frame|
{
let mut of = OwnedFrame::new();
of.copy_all_from_read_frame(&obligated_send_frame);
let _ = obligated_send_frame_holder.push(of).expect("Error: The obligated_send_frame_holder should've been checked.");
async
{
Result::<(), WebSocketError>::Ok(())
}
};
loop
{
let should_exit = self.io_server.input_receiver_ref().recv();
let read_frame_future = self.reader.read_frame(&mut send_fn);
let select_result;
select!
{
biased;
res = should_exit =>
{
select_result = SelectResult::Connection(res);
}
res = read_frame_future =>
{
select_result = SelectResult::ReadFrame(res);
}
}
if let Some(mut of) = obligated_send_frame_holder.pop()
{
match of.opcode
{
OpCode::Close =>
{
of.clear_payload();
if let Err(_err) = self.io_server.output_sender_ref().send(SimpleWebSocketActorInputMessage::WriteFrame(of)).await
{
break;
}
}
OpCode::Ping =>
{
of.pong_setup();
if let Err(_err) = self.io_server.output_sender_ref().send(SimpleWebSocketActorInputMessage::WriteFrame(of)).await
{
break;
}
}
OpCode::Continuation | OpCode::Text | OpCode::Binary | OpCode::Pong =>
{
if let Err(_err) = self.io_server.output_sender_ref().send(SimpleWebSocketActorInputMessage::WriteFrame(of)).await
{
break;
}
}
}
}
match select_result
{
SelectResult::ReadFrame(frame_res) =>
{
match frame_res
{
Ok(frame) =>
{
let mut of = OwnedFrame::new();
of.copy_all_from_read_frame(&frame);
if let Err(_err) = self.actor_io_sender.send(of).await
{
break;
}
}
Err(err) =>
{
print!("{}", err);
break;
}
}
}
SelectResult::Connection(_input_opt) =>
{
break;
}
}
}
false
}
}
impl_mac_task_actor!(WebSocketServerReaderActor);