Partial: Dedicated Peer distribution system
All checks were successful
CI / Formatting (push) Successful in 1m7s
All checks were successful
CI / Formatting (push) Successful in 1m7s
This is an uncompleted commit to move the work over to my other machine. Should compile though.
This commit is contained in:
parent
4db82f328b
commit
1a5a628000
6 changed files with 161 additions and 53 deletions
|
@ -4,14 +4,14 @@ use crate::net::prelude::*;
|
||||||
|
|
||||||
pub fn handle_new_peer(new_peers: Query<&Peer, Added<Peer>>) {
|
pub fn handle_new_peer(new_peers: Query<&Peer, Added<Peer>>) {
|
||||||
for peer in new_peers {
|
for peer in new_peers {
|
||||||
info!("Peer {} was added", peer.uuid);
|
info!("Peer {} was added", peer.id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn handle_deleted_peer(mut old_peers: RemovedComponents<Peer>, peers: Query<&Peer>) -> Result {
|
pub fn handle_deleted_peer(mut old_peers: RemovedComponents<Peer>, peers: Query<&Peer>) -> Result {
|
||||||
for entity in old_peers.read() {
|
for entity in old_peers.read() {
|
||||||
if let Ok(peer) = peers.get(entity) {
|
if let Ok(peer) = peers.get(entity) {
|
||||||
info!("Peer {} was removed", peer.uuid);
|
info!("Peer {} was removed", peer.id);
|
||||||
} else {
|
} else {
|
||||||
info!("Peer {} was removed", entity);
|
info!("Peer {} was removed", entity);
|
||||||
}
|
}
|
||||||
|
|
|
@ -90,8 +90,8 @@ pub fn update_peer_ui_timings(
|
||||||
if let Ok((peer, recv, send)) = peers.get(id.0) {
|
if let Ok((peer, recv, send)) = peers.get(id.0) {
|
||||||
**row = format!(
|
**row = format!(
|
||||||
"{} {:.2} {:.2}",
|
"{} {:.2} {:.2}",
|
||||||
peer.uuid,
|
peer.id,
|
||||||
(time.elapsed() - recv.time().unwrap_or(default())).as_secs_f64(),
|
(time.elapsed() - recv.time()).as_secs_f64(),
|
||||||
(time.elapsed() - send.time()).as_secs_f64()
|
(time.elapsed() - send.time()).as_secs_f64()
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
|
@ -26,7 +26,7 @@ fn sender<T: Into<Vec<u8>> + Component + Clone>(
|
||||||
for entity in entities {
|
for entity in entities {
|
||||||
outbound.write(OutboundPacket(Packet::new(
|
outbound.write(OutboundPacket(Packet::new(
|
||||||
(*entity).clone().into(),
|
(*entity).clone().into(),
|
||||||
peer.uuid,
|
peer.id,
|
||||||
)));
|
)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
@ -3,16 +3,18 @@ use std::time::Duration;
|
||||||
use bevy::prelude::*;
|
use bevy::prelude::*;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
|
use crate::net::peer::PotentialPeer;
|
||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
packet::{InboundPacket, OutboundPacket, Packet},
|
packet::{InboundPacket, OutboundPacket, Packet},
|
||||||
peer::{Peer, PeerChangeEvent, PeerMap, PeerReceiveTiming, PeerSendTiming},
|
peer::{Peer, PeerChangeEvent, PeerData, PeerMap, PeerReceiveTiming, PeerSendTiming},
|
||||||
queues::{NetworkReceive, NetworkSend},
|
queues::{NetworkReceive, NetworkSend},
|
||||||
};
|
};
|
||||||
|
|
||||||
pub fn handle_network_input(
|
pub fn handle_network_input(
|
||||||
from_socket: Res<NetworkReceive>,
|
from_socket: Res<NetworkReceive>,
|
||||||
peer_map: Res<PeerMap>,
|
peer_map: Res<PeerMap>,
|
||||||
mut peers: Query<(&Peer, &mut PeerReceiveTiming)>,
|
mut peers: Query<(&PeerData, &mut PeerReceiveTiming)>,
|
||||||
mut to_app: EventWriter<InboundPacket>,
|
mut to_app: EventWriter<InboundPacket>,
|
||||||
time: Res<Time>,
|
time: Res<Time>,
|
||||||
mut change_peer: EventWriter<PeerChangeEvent>,
|
mut change_peer: EventWriter<PeerChangeEvent>,
|
||||||
|
@ -44,8 +46,7 @@ pub fn handle_network_input(
|
||||||
pub fn handle_network_output(
|
pub fn handle_network_output(
|
||||||
mut from_app: EventReader<OutboundPacket>,
|
mut from_app: EventReader<OutboundPacket>,
|
||||||
peer_map: Res<PeerMap>,
|
peer_map: Res<PeerMap>,
|
||||||
mut peers: Query<(&Peer, &mut PeerSendTiming)>,
|
mut peers: Query<(&PeerData, &mut PeerSendTiming)>,
|
||||||
config: Res<Config>,
|
|
||||||
to_socket: Res<NetworkSend>,
|
to_socket: Res<NetworkSend>,
|
||||||
time: Res<Time>,
|
time: Res<Time>,
|
||||||
) -> Result {
|
) -> Result {
|
||||||
|
@ -53,8 +54,8 @@ pub fn handle_network_output(
|
||||||
let peer_id = peer_map.try_get(&packet.0.peer)?;
|
let peer_id = peer_map.try_get(&packet.0.peer)?;
|
||||||
let (peer, mut last) = peers.get_mut(*peer_id)?;
|
let (peer, mut last) = peers.get_mut(*peer_id)?;
|
||||||
// Append our UUID for client identification
|
// Append our UUID for client identification
|
||||||
let message = [packet.0.message.as_slice(), config.id.as_bytes()].concat();
|
let message = [packet.0.message.as_slice(), peer.me.as_bytes()].concat();
|
||||||
to_socket.send(message, peer.addr)?;
|
to_socket.send(message, peer.addr.into())?;
|
||||||
last.update(&time);
|
last.update(&time);
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
|
@ -63,17 +64,24 @@ pub fn handle_network_output(
|
||||||
const TIMEOUT: Duration = Duration::from_secs(10);
|
const TIMEOUT: Duration = Duration::from_secs(10);
|
||||||
|
|
||||||
pub fn heartbeat(
|
pub fn heartbeat(
|
||||||
peers: Query<(&Peer, &PeerSendTiming)>,
|
peers: Query<(AnyOf<(&Peer, &PotentialPeer)>, &PeerSendTiming)>,
|
||||||
time: Res<Time>,
|
time: Res<Time>,
|
||||||
mut outbound: EventWriter<OutboundPacket>,
|
mut outbound: EventWriter<OutboundPacket>,
|
||||||
) {
|
) -> Result {
|
||||||
for (peer, last) in peers {
|
for (peer, last) in peers {
|
||||||
// Allow for 2 consecutive missed heartbeats without timing out
|
// Allow for 2 consecutive missed heartbeats without timing out
|
||||||
if last.time() + TIMEOUT / 3 > time.elapsed() {
|
if last.time() + TIMEOUT / 3 > time.elapsed() {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
outbound.write(OutboundPacket(Packet::new(Vec::new(), peer.uuid)));
|
let id = match peer {
|
||||||
|
(None, None) => return Err("No peer identification".into()),
|
||||||
|
(None, Some(potential)) => potential.id,
|
||||||
|
(Some(actual), None) => actual.id,
|
||||||
|
(Some(_), Some(_)) => return Err("Both a potential and actual peer".into()),
|
||||||
|
};
|
||||||
|
outbound.write(OutboundPacket(Packet::new(Vec::new(), id)));
|
||||||
}
|
}
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn timeout(
|
pub fn timeout(
|
||||||
|
@ -82,22 +90,9 @@ pub fn timeout(
|
||||||
mut delete: EventWriter<PeerChangeEvent>,
|
mut delete: EventWriter<PeerChangeEvent>,
|
||||||
) {
|
) {
|
||||||
for (peer, last) in peers {
|
for (peer, last) in peers {
|
||||||
if let Some(previous) = last.time() {
|
if last.time() + TIMEOUT < time.elapsed() {
|
||||||
if previous + TIMEOUT < time.elapsed() {
|
warn!("Peer {} timed out", peer.id);
|
||||||
warn!("Peer {} timed out", peer.uuid);
|
delete.write(PeerChangeEvent::new(peer.id, None));
|
||||||
delete.write(PeerChangeEvent::new(peer.uuid, None));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Resource)]
|
|
||||||
pub struct Config {
|
|
||||||
pub id: Uuid,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Config {
|
|
||||||
pub fn new() -> Self {
|
|
||||||
Self { id: Uuid::new_v4() }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
148
src/net/peer.rs
148
src/net/peer.rs
|
@ -1,4 +1,9 @@
|
||||||
use std::{collections::HashMap, net::SocketAddr, time::Duration};
|
use std::{
|
||||||
|
array::TryFromSliceError,
|
||||||
|
collections::HashMap,
|
||||||
|
net::{IpAddr, Ipv6Addr, SocketAddr},
|
||||||
|
time::Duration,
|
||||||
|
};
|
||||||
|
|
||||||
use bevy::prelude::*;
|
use bevy::prelude::*;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
@ -19,32 +24,122 @@ impl PeerSendTiming {
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Component, Debug, Default)]
|
#[derive(Component, Debug, Default)]
|
||||||
pub struct PeerReceiveTiming(Option<Duration>);
|
pub struct PeerReceiveTiming(Duration);
|
||||||
|
|
||||||
impl PeerReceiveTiming {
|
impl PeerReceiveTiming {
|
||||||
pub fn new(time: &Res<Time>) -> Self {
|
pub fn new(time: &Res<Time>) -> Self {
|
||||||
Self(Some(time.elapsed()))
|
Self(time.elapsed())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn update(&mut self, time: &Res<Time>) {
|
pub fn update(&mut self, time: &Res<Time>) {
|
||||||
self.0 = Some(time.elapsed())
|
self.0 = time.elapsed()
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn time(&self) -> Option<Duration> {
|
pub fn time(&self) -> Duration {
|
||||||
self.0
|
self.0
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Component, Debug)]
|
#[derive(Component, Debug)]
|
||||||
#[require(PeerSendTiming, PeerReceiveTiming)]
|
#[require(PeerData, PeerReceiveTiming)]
|
||||||
pub struct Peer {
|
pub struct Peer {
|
||||||
pub addr: SocketAddr,
|
pub id: Uuid,
|
||||||
pub uuid: Uuid,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Peer {
|
#[derive(Component, Debug)]
|
||||||
pub fn new(addr: SocketAddr, uuid: Uuid) -> Self {
|
#[require(PeerData)]
|
||||||
Self { addr, uuid }
|
pub struct PotentialPeer {
|
||||||
|
pub id: Uuid,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for PotentialPeer {
|
||||||
|
fn default() -> Self {
|
||||||
|
Self { id: Uuid::new_v4() }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Component, Copy, Debug)]
|
||||||
|
#[require(PeerSendTiming)]
|
||||||
|
pub struct PeerData {
|
||||||
|
pub addr: Address,
|
||||||
|
pub me: Uuid,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PeerData {
|
||||||
|
pub fn new(addr: SocketAddr, me: Uuid) -> Self {
|
||||||
|
Self {
|
||||||
|
addr: addr.into(),
|
||||||
|
me,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for PeerData {
|
||||||
|
fn default() -> Self {
|
||||||
|
Self {
|
||||||
|
addr: Default::default(),
|
||||||
|
me: Uuid::new_v4(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Copy, Debug)]
|
||||||
|
pub struct Address(SocketAddr);
|
||||||
|
|
||||||
|
impl From<Address> for SocketAddr {
|
||||||
|
fn from(value: Address) -> Self {
|
||||||
|
value.0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<SocketAddr> for Address {
|
||||||
|
fn from(value: SocketAddr) -> Self {
|
||||||
|
Self(value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PartialEq<SocketAddr> for Address {
|
||||||
|
fn eq(&self, other: &SocketAddr) -> bool {
|
||||||
|
self.0 == *other
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PartialEq<Address> for SocketAddr {
|
||||||
|
fn eq(&self, other: &Address) -> bool {
|
||||||
|
other == self
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl From<Address> for Vec<u8> {
|
||||||
|
fn from(value: Address) -> Self {
|
||||||
|
let mut bytes: Vec<u8> = match value.0.ip() {
|
||||||
|
IpAddr::V4(ipv4_addr) => ipv4_addr.octets().into(),
|
||||||
|
IpAddr::V6(ipv6_addr) => ipv6_addr.octets().into(),
|
||||||
|
};
|
||||||
|
bytes.extend(value.0.port().to_le_bytes());
|
||||||
|
bytes
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TryFrom<Vec<u8>> for Address {
|
||||||
|
type Error = TryFromSliceError;
|
||||||
|
|
||||||
|
fn try_from(value: Vec<u8>) -> Result<Self, Self::Error> {
|
||||||
|
let port = u16::from_le_bytes(TryInto::<[u8; 2]>::try_into(
|
||||||
|
value.clone().split_off(value.len() - 2).as_slice(),
|
||||||
|
)?);
|
||||||
|
let addr = if let Ok(bytes) = TryInto::<[u8; 4]>::try_into(value.as_slice()) {
|
||||||
|
SocketAddr::from((bytes, port))
|
||||||
|
} else {
|
||||||
|
SocketAddr::from((TryInto::<[u8; 16]>::try_into(value.as_slice())?, port))
|
||||||
|
};
|
||||||
|
Ok(Address(addr))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for Address {
|
||||||
|
fn default() -> Self {
|
||||||
|
Self(SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -90,7 +185,7 @@ impl PeerChangeEvent {
|
||||||
pub fn handle_peer_change(
|
pub fn handle_peer_change(
|
||||||
mut changes: EventReader<PeerChangeEvent>,
|
mut changes: EventReader<PeerChangeEvent>,
|
||||||
mut peer_map: ResMut<PeerMap>,
|
mut peer_map: ResMut<PeerMap>,
|
||||||
mut peers: Query<&mut Peer>,
|
mut peers: Query<&mut PeerData>,
|
||||||
mut commands: Commands,
|
mut commands: Commands,
|
||||||
time: Res<Time>,
|
time: Res<Time>,
|
||||||
) -> Result {
|
) -> Result {
|
||||||
|
@ -98,7 +193,7 @@ pub fn handle_peer_change(
|
||||||
if let Some(entity) = peer_map.get(&change.peer) {
|
if let Some(entity) = peer_map.get(&change.peer) {
|
||||||
if let Some(addr) = change.addr {
|
if let Some(addr) = change.addr {
|
||||||
if let Ok(mut peer) = peers.get_mut(*entity) {
|
if let Ok(mut peer) = peers.get_mut(*entity) {
|
||||||
peer.addr = addr;
|
peer.addr = Address(addr);
|
||||||
} else {
|
} else {
|
||||||
warn!("Peer {} doesn't exist (just added?)", change.peer);
|
warn!("Peer {} doesn't exist (just added?)", change.peer);
|
||||||
}
|
}
|
||||||
|
@ -111,7 +206,10 @@ pub fn handle_peer_change(
|
||||||
peer_map.insert(
|
peer_map.insert(
|
||||||
change.peer,
|
change.peer,
|
||||||
commands
|
commands
|
||||||
.spawn((Peer::new(addr, change.peer), PeerReceiveTiming::new(&time)))
|
.spawn((
|
||||||
|
PeerData::new(addr, Uuid::new_v4()),
|
||||||
|
PeerReceiveTiming::new(&time),
|
||||||
|
))
|
||||||
.id(),
|
.id(),
|
||||||
);
|
);
|
||||||
} else {
|
} else {
|
||||||
|
@ -126,11 +224,25 @@ pub fn handle_peer_change(
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[allow(clippy::type_complexity)]
|
||||||
pub fn handle_new_peer(
|
pub fn handle_new_peer(
|
||||||
new_peers: Query<&Peer, Added<Peer>>,
|
new_peers: Query<(Option<Ref<Peer>>, Option<&Peer>, &PeerData)>,
|
||||||
mut outbound: EventWriter<OutboundPacket>,
|
mut outbound: EventWriter<OutboundPacket>,
|
||||||
) {
|
) -> Result {
|
||||||
for peer in new_peers {
|
for (change, this, data) in new_peers {
|
||||||
outbound.write(OutboundPacket(Packet::new(Vec::new(), peer.uuid)));
|
if let Some(change) = change {
|
||||||
|
if change.is_added() {
|
||||||
|
if let Some(this) = this {
|
||||||
|
for (_, _, other) in new_peers {
|
||||||
|
if data.me != other.me {
|
||||||
|
outbound.write(OutboundPacket(Packet::new(other.addr.into(), this.id)));
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
return Err("Ref<Peer> without Peer".into());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
|
@ -4,10 +4,9 @@ use bevy::prelude::*;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use super::{
|
use super::{
|
||||||
distribution::Networked,
|
io::{handle_network_input, handle_network_output, heartbeat, timeout},
|
||||||
io::{Config, handle_network_input, handle_network_output, heartbeat, timeout},
|
|
||||||
packet::{InboundPacket, OutboundPacket},
|
packet::{InboundPacket, OutboundPacket},
|
||||||
peer::{Peer, PeerChangeEvent, PeerMap, handle_new_peer, handle_peer_change},
|
peer::{PeerChangeEvent, PeerData, PeerMap, handle_new_peer, handle_peer_change},
|
||||||
queues::{NetworkReceive, NetworkSend},
|
queues::{NetworkReceive, NetworkSend},
|
||||||
socket::bind_socket,
|
socket::bind_socket,
|
||||||
state::NetworkState,
|
state::NetworkState,
|
||||||
|
@ -41,7 +40,6 @@ impl Plugin for NetIOPlugin {
|
||||||
FixedPostUpdate,
|
FixedPostUpdate,
|
||||||
handle_network_output.run_if(in_state(NetworkState::MultiPlayer)),
|
handle_network_output.run_if(in_state(NetworkState::MultiPlayer)),
|
||||||
)
|
)
|
||||||
.insert_resource(Config::new())
|
|
||||||
.add_event::<PeerChangeEvent>()
|
.add_event::<PeerChangeEvent>()
|
||||||
.add_event::<InboundPacket>()
|
.add_event::<InboundPacket>()
|
||||||
.add_event::<OutboundPacket>();
|
.add_event::<OutboundPacket>();
|
||||||
|
@ -54,7 +52,10 @@ impl Plugin for NetIOPlugin {
|
||||||
|
|
||||||
let mut peer_map = PeerMap::default();
|
let mut peer_map = PeerMap::default();
|
||||||
if let Some(socket) = self.peer {
|
if let Some(socket) = self.peer {
|
||||||
let entity = app.world_mut().spawn(Peer::new(socket, Uuid::nil()));
|
let entity = app.world_mut().spawn(PeerData {
|
||||||
|
addr: socket.into(),
|
||||||
|
me: Uuid::nil(),
|
||||||
|
});
|
||||||
peer_map.insert(Uuid::nil(), entity.id());
|
peer_map.insert(Uuid::nil(), entity.id());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
Loading…
Add table
Add a link
Reference in a new issue