bevy_replicon_quinnet 0.10.0

Integration with bevy_quinnet for bevy_replicon
Documentation
use bevy::{
    app::{App, Plugin, PostUpdate, PreUpdate},
    ecs::{
        entity::Entity,
        system::{Commands, Query},
    },
    prelude::{EventReader, IntoSystemConfigs, IntoSystemSetConfigs, Local, Res, ResMut},
    time::Time,
};
use bevy_quinnet::{
    server::{QuinnetServer, QuinnetServerPlugin},
    shared::QuinnetSyncUpdate,
};
use bevy_replicon::{
    core::connected_client::{ClientId, ClientIdMap},
    prelude::{ConnectedClient, NetworkStats, RepliconServer},
    server::ServerSet,
};

use crate::BYTES_PER_SEC_PERIOD;

pub struct RepliconQuinnetServerPlugin;

impl Plugin for RepliconQuinnetServerPlugin {
    fn build(&self, app: &mut App) {
        app.add_plugins(QuinnetServerPlugin::default())
            .configure_sets(
                PreUpdate,
                ServerSet::ReceivePackets.after(QuinnetSyncUpdate),
            )
            .add_systems(
                PreUpdate,
                ((
                    Self::set_running.run_if(bevy_quinnet::server::server_just_opened),
                    Self::set_stopped.run_if(bevy_quinnet::server::server_just_closed),
                    (
                        Self::receive_packets,
                        Self::update_statistics,
                        Self::forward_server_events,
                    )
                        .run_if(bevy_quinnet::server::server_listening),
                )
                    .chain()
                    .in_set(ServerSet::ReceivePackets),),
            )
            .add_systems(
                PostUpdate,
                Self::send_packets
                    .in_set(ServerSet::SendPackets)
                    .run_if(bevy_quinnet::server::server_listening),
            );
    }
}

impl RepliconQuinnetServerPlugin {
    fn set_running(mut server: ResMut<RepliconServer>) {
        server.set_running(true);
    }

    fn set_stopped(mut server: ResMut<RepliconServer>) {
        server.set_running(false);
    }

    fn forward_server_events(
        mut commands: Commands,
        mut conn_events: EventReader<bevy_quinnet::server::ConnectionEvent>,
        mut conn_lost_events: EventReader<bevy_quinnet::server::ConnectionLostEvent>,
        client_map: Res<ClientIdMap>,
    ) {
        for event in conn_events.read() {
            let client_id = ClientId::new(event.id);
            const DEFAULT_INITIAL_MAX_DATAGRAM_SIZE: usize = 1200;
            commands.spawn(ConnectedClient::new(
                client_id,
                DEFAULT_INITIAL_MAX_DATAGRAM_SIZE,
            ));
        }
        for event in conn_lost_events.read() {
            let client_id = ClientId::new(event.id);
            let client_entity = *client_map
                .get(&client_id)
                .expect("clients should be connected before disconnection");
            commands.entity(client_entity).despawn();
        }
    }

    fn update_statistics(
        mut bps_timer: Local<f64>,
        mut clients: Query<(&mut ConnectedClient, &mut NetworkStats)>,
        mut quinnet_server: ResMut<QuinnetServer>,
        time: Res<Time>,
    ) {
        let Some(endpoint) = quinnet_server.get_endpoint_mut() else {
            return;
        };
        for (mut client, mut client_stats) in clients.iter_mut() {
            let Some(con) = endpoint.get_connection_mut(client.id().get()) else {
                return;
            };

            if let Some(max_size) = con.max_datagram_size() {
                client.max_size = max_size;
            }

            let stats = con.connection_stats();

            client_stats.rtt = stats.path.rtt.as_secs_f64();
            client_stats.packet_loss =
                100. * (stats.path.lost_packets as f64 / stats.path.sent_packets as f64);

            *bps_timer += time.delta_secs_f64();
            if *bps_timer >= BYTES_PER_SEC_PERIOD {
                *bps_timer = 0.;
                let received_bytes_count = con.clear_received_bytes_count() as f64;
                let sent_bytes_count = con.clear_sent_bytes_count() as f64;
                client_stats.received_bps = received_bytes_count / BYTES_PER_SEC_PERIOD;
                client_stats.sent_bps = sent_bytes_count / BYTES_PER_SEC_PERIOD;
            }
        }
    }

    fn receive_packets(
        mut quinnet_server: ResMut<QuinnetServer>,
        mut replicon_server: ResMut<RepliconServer>,
        mut clients: Query<(Entity, &ConnectedClient)>,
    ) {
        let Some(endpoint) = quinnet_server.get_endpoint_mut() else {
            return;
        };
        for (client_entity, client) in &mut clients {
            while let Some((channel_id, message)) =
                endpoint.try_receive_payload_from(client.id().get())
            {
                replicon_server.insert_received(client_entity, channel_id, message);
            }
        }
    }

    fn send_packets(
        mut quinnet_server: ResMut<QuinnetServer>,
        mut replicon_server: ResMut<RepliconServer>,
        clients: Query<&ConnectedClient>,
    ) {
        let Some(endpoint) = quinnet_server.get_endpoint_mut() else {
            return;
        };
        for (client_entity, channel_id, message) in replicon_server.drain_sent() {
            let client = clients
                .get(client_entity)
                .expect("messages should be sent only to connected clients");
            endpoint.try_send_payload_on(client.id().get(), channel_id, message);
        }
    }
}