literustlib_server 0.4.0

Rust server for LiteNetLib
Documentation
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, // reliable
    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, // nanoseconds
    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
    }

    /// Disconnect connection silently
    pub fn disconnect(&self) -> bool {
        self.is_connected.swap(false, literustlib::serdes::ATOMIC_ORDERING)
    }

    /// Disconnect connection and tell them about it
    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)
    }

    /// Mark the connection as legitimate so DOS protection will be less strict
    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) fn channel_by_ack_prop_mut(&mut self, property: literustlib::packet::Property) -> &'_ mut Channel {
        match property {
            literustlib::packet::Property::AckReliable => &mut self.reliable_unordered,
            literustlib::packet::Property::AckReliableOrdered => &mut self.reliable_ordered,
            literustlib::packet::Property::AckAuth => &mut 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; // must be greater than window_size / 8 (ideally at least window_size / 4)

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 fn move_seen_window(&mut self) -> usize {
        let mut iterations = 0;
        while (self.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] = self.seen[i];
        }
        // TODO also bitshift
        self.seen = new_seen;
        bytes_to_move * 8
    }*/

    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];
        }
        // TODO also bitshift
        *seen = new_seen;
        bytes_to_move * 8
    }
}