use bevy::prelude::*;
use matchbox_socket::PeerId;
use serde::{Serialize, de::DeserializeOwned};
use std::collections::VecDeque;
#[derive(Message, Debug, Clone)]
pub struct Broadcast<T: Serialize + DeserializeOwned + Send + Sync + 'static> {
pub payload: T,
pub channel: ChannelKind,
}
#[derive(Debug, Clone)]
pub struct NetworkReceived<T: Serialize + DeserializeOwned + Send + Sync + 'static> {
pub payload: T,
pub sender: PeerId,
pub channel: ChannelKind,
pub(crate) packet_size: usize,
}
const MAX_NETWORK_QUEUE_LEN: usize = 4096;
const MAX_NETWORK_QUEUE_BYTES: usize = 64 * 1024 * 1024;
#[derive(Resource, Debug)]
pub struct NetworkQueue<T: Serialize + DeserializeOwned + Send + Sync + 'static> {
incoming: VecDeque<NetworkReceived<T>>,
total_bytes: usize,
warned_full: bool,
}
impl<T: Serialize + DeserializeOwned + Send + Sync + 'static> Default for NetworkQueue<T> {
fn default() -> Self {
Self {
incoming: VecDeque::new(),
total_bytes: 0,
warned_full: false,
}
}
}
impl<T: Serialize + DeserializeOwned + Send + Sync + 'static> NetworkQueue<T> {
pub(crate) fn push(&mut self, msg: NetworkReceived<T>) {
if self.incoming.len() >= MAX_NETWORK_QUEUE_LEN {
return;
}
if self.total_bytes.saturating_add(msg.packet_size) > MAX_NETWORK_QUEUE_BYTES {
return;
}
self.total_bytes += msg.packet_size;
self.incoming.push_back(msg);
}
pub fn drain(&mut self) -> impl Iterator<Item = NetworkReceived<T>> {
self.total_bytes = 0;
self.warned_full = false;
std::mem::take(&mut self.incoming).into_iter()
}
pub fn is_empty(&self) -> bool {
self.incoming.is_empty()
}
pub fn len(&self) -> usize {
self.incoming.len()
}
pub(crate) fn would_drop(&mut self, packet_size: usize) -> bool {
if self.incoming.len() >= MAX_NETWORK_QUEUE_LEN {
if !self.warned_full {
bevy::log::warn!(
"NetworkQueue full ({MAX_NETWORK_QUEUE_LEN} messages), dropping incoming messages"
);
self.warned_full = true;
}
return true;
}
if self.total_bytes.saturating_add(packet_size) > MAX_NETWORK_QUEUE_BYTES {
if !self.warned_full {
bevy::log::warn!(
total_bytes = self.total_bytes,
msg_bytes = packet_size,
"NetworkQueue byte budget exceeded ({MAX_NETWORK_QUEUE_BYTES} bytes), dropping incoming messages"
);
self.warned_full = true;
}
return true;
}
false
}
}
#[derive(Debug, Clone)]
pub struct PeerStateChanged {
pub peer: PeerId,
pub state: PeerConnectionState,
}
#[derive(Resource, Debug)]
pub struct PeerStateQueue<T: Send + Sync + 'static> {
events: VecDeque<PeerStateChanged>,
warned_full: bool,
_marker: std::marker::PhantomData<T>,
}
impl<T: Send + Sync + 'static> Default for PeerStateQueue<T> {
fn default() -> Self {
Self {
events: VecDeque::new(),
warned_full: false,
_marker: std::marker::PhantomData,
}
}
}
impl<T: Send + Sync + 'static> PeerStateQueue<T> {
pub(crate) fn push(&mut self, event: PeerStateChanged) {
if self.events.len() >= MAX_NETWORK_QUEUE_LEN {
if !self.warned_full {
bevy::log::warn!(
"PeerStateQueue full ({MAX_NETWORK_QUEUE_LEN} events), dropping incoming event"
);
self.warned_full = true;
}
return;
}
self.events.push_back(event);
}
pub fn drain(&mut self) -> impl Iterator<Item = PeerStateChanged> {
self.warned_full = false;
std::mem::take(&mut self.events).into_iter()
}
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
pub fn len(&self) -> usize {
self.events.len()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PeerConnectionState {
Connected,
Disconnected,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ChannelKind {
#[default]
Reliable,
Unreliable,
}
impl ChannelKind {
pub(crate) fn index(self) -> usize {
match self {
Self::Reliable => 0,
Self::Unreliable => 1,
}
}
}