pub struct Connection<D: literustlib::packet::PacketData> {
pub(crate) id: i64,
pub(crate) pool: tokio::sync::Mutex<literustlib::serdes::PacketPool<D>>,
pub(crate) addr: core::net::SocketAddr,
pub(crate) socket: std::sync::Arc<tokio::net::UdpSocket>,
pub(crate) reliable_unordered: Channel,
pub(crate) sequenced: Channel,
pub(crate) reliable_ordered: Channel,
pub(crate) state_channel: Channel,
pub(crate) auth_channel: Channel, pub(crate) ping_channel: Channel,
pub(crate) last_ping_time: tokio::sync::Mutex<Option<chrono::DateTime<chrono::Utc>>>,
pub(crate) round_trip: std::sync::atomic::AtomicU64, pub(crate) is_connected: std::sync::atomic::AtomicBool,
pub(crate) is_certified: std::sync::atomic::AtomicBool,
pub(crate) last_seen: tokio::sync::Mutex<chrono::DateTime<chrono::Utc>>,
}
impl <D: literustlib::packet::PacketData> Connection<D> {
pub fn id(&self) -> i64 {
self.id
}
pub fn disconnect(&self) -> bool {
self.is_connected.swap(false, literustlib::serdes::ATOMIC_ORDERING)
}
pub async fn goodbye(&self, sender: &super::DataSender<D>) -> bool {
let packet = literustlib::packet::Packet::without_data(
literustlib::packet::Header::with_prop(literustlib::packet::Property::Disconnect)
);
if let Err(e) = sender.raw_send_to(&packet, self).await {
log::error!("Failed to send disconnect packet for connection {}: {}", self.id, e);
}
self.disconnect()
}
pub fn is_connected(&self) -> bool {
self.is_connected.load(literustlib::serdes::ATOMIC_ORDERING)
}
pub fn certify(&self) {
self.is_certified.store(true, literustlib::serdes::ATOMIC_ORDERING);
}
pub fn is_certified(&self) -> bool {
self.is_certified.load(literustlib::serdes::ATOMIC_ORDERING)
}
pub(crate) fn channel_by_prop(&self, property: literustlib::packet::Property) -> Option<&'_ Channel> {
match property {
literustlib::packet::Property::Reliable => Some(&self.reliable_unordered),
literustlib::packet::Property::Sequenced => Some(&self.sequenced),
literustlib::packet::Property::ReliableOrdered => Some(&self.reliable_ordered),
literustlib::packet::Property::StateUpdate => Some(&self.state_channel),
literustlib::packet::Property::Auth => Some(&self.auth_channel),
_ => None,
}
}
pub(crate) fn channel_by_ack_prop(&self, property: literustlib::packet::Property) -> &'_ Channel {
match property {
literustlib::packet::Property::AckReliable => &self.reliable_unordered,
literustlib::packet::Property::AckReliableOrdered => &self.reliable_ordered,
literustlib::packet::Property::AckAuth => &self.auth_channel,
_ => unreachable!("There are only three ack property variants"),
}
}
pub(crate) async fn last_seen_now(&self) {
*self.last_seen.lock().await = chrono::Utc::now();
}
pub(crate) async fn last_seen_delta(&self) -> chrono::TimeDelta {
chrono::Utc::now().signed_duration_since(&*self.last_seen.lock().await)
}
pub(crate) fn new(connect_id: i64, address: core::net::SocketAddr, socket: std::sync::Arc<tokio::net::UdpSocket>, window: u16) -> Self {
Self {
id: connect_id,
pool: tokio::sync::Mutex::new(literustlib::serdes::PacketPool::new()),
addr: address,
socket: socket,
reliable_unordered: super::Channel::new(window, true),
sequenced: super::Channel::new(window, false),
reliable_ordered: super::Channel::new(window, true),
state_channel: super::Channel::new(window, false),
auth_channel: super::Channel::new(window, true),
ping_channel: super::Channel::new(0, false),
last_ping_time: tokio::sync::Mutex::new(None),
round_trip: std::sync::atomic::AtomicU64::new(0),
is_connected: std::sync::atomic::AtomicBool::new(true),
is_certified: std::sync::atomic::AtomicBool::new(false),
last_seen: tokio::sync::Mutex::new(chrono::Utc::now()),
}
}
}
pub(crate) const SEEN_WINDOW: usize = 16;
pub struct Channel {
pub(crate) local_sequence: std::sync::atomic::AtomicU16,
pub(crate) remote_sequence: std::sync::atomic::AtomicU16,
pub(crate) pending: tokio::sync::Mutex<std::collections::VecDeque<literustlib::packet::Packet>>,
pub(crate) window_start: std::sync::atomic::AtomicU16,
pub(crate) remote_window_start: std::sync::atomic::AtomicU16,
pub(crate) window_size: u16,
pub(crate) reliable: bool,
pub(crate) seen: tokio::sync::Mutex<[u8; SEEN_WINDOW]>,
}
impl Channel {
pub fn new(window: u16, is_reliable: bool) -> Self {
Self {
local_sequence: std::sync::atomic::AtomicU16::new(0),
remote_sequence: std::sync::atomic::AtomicU16::new(0),
pending: tokio::sync::Mutex::new(std::collections::VecDeque::with_capacity(window as usize * 2)),
window_start: std::sync::atomic::AtomicU16::new(0),
remote_window_start: std::sync::atomic::AtomicU16::new(0),
window_size: window,
reliable: is_reliable,
seen: tokio::sync::Mutex::new([0u8; SEEN_WINDOW]),
}
}
pub(crate) fn is_window_full(&self, pending: &std::collections::VecDeque<literustlib::packet::Packet>) -> bool {
pending.len() >= (self.window_size as usize)
}
pub fn is_reliable(&self) -> bool {
self.reliable
}
pub(crate) fn move_seen_window(seen: &mut [u8; SEEN_WINDOW]) -> usize {
let mut iterations = 0;
while (seen[iterations / 8] >> (7 - iterations % 8)) & 1 == 1 {
iterations += 1;
}
let mut new_seen = [0u8; SEEN_WINDOW];
let bytes_to_move = iterations / 8;
for i in bytes_to_move..SEEN_WINDOW {
new_seen[i - bytes_to_move] = seen[i];
}
*seen = new_seen;
bytes_to_move * 8
}
}