Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c1c9d95147 | |||
| 9b2287a3f0 | |||
| e6520d99ec | |||
| 40cdec017c | |||
| a8fdbeb025 |
2 changed files with 64 additions and 24 deletions
|
|
@ -38,6 +38,7 @@ pub struct InternalConnection<C: CryptoInstance> {
|
|||
server_packet_counter: u16,
|
||||
client_packet_counter: u16,
|
||||
unacknowledged_packets: HashMap<u16, Arc<Vec<u8>>>,
|
||||
packet_buffer: Vec<u8>,
|
||||
packet_queue: HashMap<u16, (Instant, PRUDPV0Packet<Vec<u8>>)>,
|
||||
}
|
||||
pub struct Connection<C: CryptoInstance> {
|
||||
|
|
@ -96,7 +97,7 @@ impl<C: Crypto> Server<C> {
|
|||
.expect("packet malformed in creation"),
|
||||
);*/
|
||||
let mut inner = conn.inner.lock().await;
|
||||
let pieces = data.chunks(700);
|
||||
let pieces = data.chunks(962);
|
||||
let max_piece = pieces.len() - 1;
|
||||
let mut frag_num = 1;
|
||||
for (i, piece) in pieces.enumerate() {
|
||||
|
|
@ -140,9 +141,18 @@ impl<C: Crypto> Server<C> {
|
|||
.send_to(&data, conn.addr.regular_socket_addr)
|
||||
.await
|
||||
.ok();
|
||||
|
||||
break;
|
||||
sleep(Duration::from_millis(500)).await;
|
||||
}
|
||||
println!("connection exceeded max fail count, disconnecting");
|
||||
let Some(conn) = conn.upgrade() else {
|
||||
return;
|
||||
};
|
||||
let Some(this) = this.upgrade() else {
|
||||
return;
|
||||
};
|
||||
let mut conns = this.connections.write().await;
|
||||
conns.remove(&(conn.addr, conn.session_id));
|
||||
drop(conns);
|
||||
});
|
||||
frag_num += 1;
|
||||
}
|
||||
|
|
@ -282,6 +292,7 @@ impl<C: Crypto> Server<C> {
|
|||
server_packet_counter: 1,
|
||||
unacknowledged_packets: HashMap::new(),
|
||||
packet_queue: HashMap::new(),
|
||||
packet_buffer: vec![],
|
||||
}),
|
||||
});
|
||||
|
||||
|
|
@ -334,6 +345,13 @@ impl<C: Crypto> Server<C> {
|
|||
warn!("data packet on inactive connection from: {:?}", addr);
|
||||
return;
|
||||
};
|
||||
|
||||
if header.type_flags.get_flags() & ACK != 0 {
|
||||
let mut inner = res.inner.lock().await;
|
||||
inner.unacknowledged_packets.remove(&header.sequence_id);
|
||||
return;
|
||||
}
|
||||
|
||||
info!("frag: {}", frag_id);
|
||||
let mut conn = res.inner.lock().await;
|
||||
let ack = new_data_packet(
|
||||
|
|
@ -368,9 +386,16 @@ impl<C: Crypto> Server<C> {
|
|||
};
|
||||
|
||||
conn.crypto_instance.decrypt_incoming(payload);
|
||||
|
||||
res.target.send(payload.to_owned()).await;
|
||||
conn.packet_buffer.extend_from_slice(payload);
|
||||
conn.client_packet_counter += 1;
|
||||
if *packet.fragment_id().unwrap() != 0 {
|
||||
info!("handeling fragmented packet");
|
||||
continue;
|
||||
}
|
||||
|
||||
res.target
|
||||
.send(std::mem::take(&mut conn.packet_buffer))
|
||||
.await;
|
||||
}
|
||||
info!("finished handeling packets, dropping inner connection");
|
||||
drop(conn);
|
||||
|
|
@ -472,8 +497,8 @@ impl<C: Crypto> Server<C> {
|
|||
inner.last_action = Instant::now();
|
||||
drop(inner);
|
||||
};
|
||||
if header.type_flags.get_flags() & ACK != 0 {
|
||||
info!("got ack(acks are ignored for now)");
|
||||
if header.type_flags.get_flags() & ACK != 0 && header.type_flags.get_types() != DATA {
|
||||
info!("got ack(acks are ignored for now(unless they are data acks))");
|
||||
return;
|
||||
}
|
||||
println!("{:?}", header);
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ struct InternalConnection<E: CryptoHandlerConnectionInstance> {
|
|||
socket: Arc<UdpSocket>,
|
||||
packet_queue: HashMap<u16, PRUDPV1Packet>,
|
||||
last_packet_time: Instant,
|
||||
partial_packet: Vec<u8>,
|
||||
unacknowleged_packets: Vec<(Instant, PRUDPV1Packet)>,
|
||||
}
|
||||
|
||||
|
|
@ -431,6 +432,7 @@ impl<T: CryptoHandler> InternalSocket<T> {
|
|||
packet_queue: Default::default(),
|
||||
last_packet_time: Instant::now(),
|
||||
unacknowleged_packets: Vec::new(),
|
||||
partial_packet: Vec::new(),
|
||||
supported_function_version,
|
||||
};
|
||||
|
||||
|
|
@ -573,11 +575,24 @@ impl<T: CryptoHandler> InternalSocket<T> {
|
|||
while let Some(mut packet) = conn.packet_queue.remove(&counter) {
|
||||
conn.crypto_handler_instance
|
||||
.decrypt_incoming(packet.header.substream_id, &mut packet.payload[..]);
|
||||
|
||||
conn.data_sender.send(packet.payload).await.ok();
|
||||
|
||||
conn.partial_packet
|
||||
.extend_from_slice(&mut packet.payload[..]);
|
||||
conn.reliable_client_counter = conn.reliable_client_counter.overflowing_add(1).0;
|
||||
counter = conn.reliable_client_counter;
|
||||
if packet.options.iter().any(|v| {
|
||||
if let FragmentId(f) = v {
|
||||
*f != 0
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}) {
|
||||
println!("handeling fragmented packet");
|
||||
continue;
|
||||
}
|
||||
|
||||
let packet = std::mem::take(&mut conn.partial_packet);
|
||||
|
||||
conn.data_sender.send(packet).await.ok();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -690,19 +705,7 @@ impl<T: CryptoHandler> AnyInternalSocket for InternalSocket<T> {
|
|||
let conn = &**conn;
|
||||
let mut conn = conn.lock().await;
|
||||
|
||||
if conn.supported_function_version == 1 {
|
||||
let mut collected_ids: Vec<u16> = Vec::new();
|
||||
let mut cursor = Cursor::new(&packet.payload);
|
||||
|
||||
while let Ok(v) = read_u16(&mut cursor) {
|
||||
collected_ids.push(v);
|
||||
}
|
||||
|
||||
conn.unacknowleged_packets.retain_mut(|(_, up)| {
|
||||
!(collected_ids.iter().any(|id| up.header.sequence_id == *id)
|
||||
|| up.header.sequence_id <= packet.header.sequence_id)
|
||||
});
|
||||
} else {
|
||||
if packet.header.substream_id == 1 {
|
||||
let mut collected_ids: Vec<u16> = Vec::new();
|
||||
let mut cursor = Cursor::new(&packet.payload);
|
||||
|
||||
|
|
@ -729,10 +732,22 @@ impl<T: CryptoHandler> AnyInternalSocket for InternalSocket<T> {
|
|||
collected_ids.push(additional_sequence_id);
|
||||
}
|
||||
|
||||
conn.unacknowleged_packets.retain_mut(|(_, up)| {
|
||||
conn.unacknowleged_packets.retain(|(_, up)| {
|
||||
!(collected_ids.iter().any(|id| up.header.sequence_id == *id)
|
||||
|| up.header.sequence_id <= sequence_id)
|
||||
});
|
||||
} else {
|
||||
let mut collected_ids: Vec<u16> = Vec::new();
|
||||
let mut cursor = Cursor::new(&packet.payload);
|
||||
|
||||
while let Ok(v) = read_u16(&mut cursor) {
|
||||
collected_ids.push(v);
|
||||
}
|
||||
|
||||
conn.unacknowleged_packets.retain(|(_, up)| {
|
||||
!(collected_ids.iter().any(|id| up.header.sequence_id == *id)
|
||||
|| up.header.sequence_id <= packet.header.sequence_id)
|
||||
});
|
||||
}
|
||||
} else {
|
||||
error!("non connection acknowledgement packet on nonexistent connection...")
|
||||
|
|
|
|||
Loading…
Reference in a new issue