use crate::common::Kcp2KMode;
use crate::kcp2k_channel::Kcp2KChannel;
use crate::kcp2k_config::Kcp2KConfig;
use crate::kcp2k_state::Kcp2KPeerState;
use bytes::{BufMut, Bytes, BytesMut};
use kcp::{Kcp, KCP_OVERHEAD};
use socket2::{SockAddr, Socket};
use std::io;
use std::io::Write;
use std::sync::{Arc, RwLock};
use std::time::{Duration, Instant};
use tklog::error;
#[derive(Debug)]
pub struct Kcp2KPeer {
pub cookie: Arc<Bytes>, pub state: RwLock<Kcp2KPeerState>, pub kcp: RwLock<Kcp<UdpOutput>>, pub watch: Instant,
pub timeout_duration: Duration, pub last_recv_time: RwLock<Duration>, pub last_send_ping_time: RwLock<Duration>, }
impl Kcp2KPeer {
pub fn new(
kcp2k_mode: Arc<Kcp2KMode>,
config: Arc<Kcp2KConfig>,
cookie: Arc<Bytes>,
socket: Arc<Socket>,
client_sock_addr: Arc<SockAddr>,
) -> Self {
let udp_output = UdpOutput::new(
kcp2k_mode,
Arc::clone(&cookie),
Arc::clone(&socket),
Arc::clone(&client_sock_addr),
);
let mut kcp = Kcp::new(0, udp_output);
kcp.set_nodelay(
if config.no_delay { true } else { false },
config.interval,
config.fast_resend,
!config.congestion_window,
);
kcp.set_wndsize(config.send_window_size, config.receive_window_size);
kcp.set_mtu(config.mtu - Kcp2KConfig::METADATA_SIZE_RELIABLE)
.expect("set_mtu failed");
kcp.set_maximum_resend_times(config.max_retransmits);
Self {
kcp: RwLock::new(kcp),
cookie,
state: RwLock::new(Kcp2KPeerState::Connected),
timeout_duration: Duration::from_millis(config.timeout),
watch: Instant::now(),
last_recv_time: RwLock::new(Duration::from_secs(0)),
last_send_ping_time: RwLock::new(Duration::from_secs(0)),
}
}
pub fn reliable_max_message_size_unconstrained(mtu: u32, rcv_wnd: u32) -> usize {
((mtu - KCP_OVERHEAD as u32 - 5) * (rcv_wnd - 1) - 1) as usize
}
pub fn reliable_max_message_size(mtu: u32, rcv_wnd: u32) -> usize {
Self::reliable_max_message_size_unconstrained(mtu, rcv_wnd.min(255))
}
pub fn unreliable_max_message_size(mtu: u32) -> usize {
(mtu - KCP_OVERHEAD as u32 - 1) as usize
}
}
#[derive(Debug)]
pub struct UdpOutput {
kcp2k_mode: Arc<Kcp2KMode>, cookie: Arc<Bytes>, socket: Arc<Socket>, client_sock_addr: Arc<SockAddr>, }
impl UdpOutput {
pub fn new(
kcp2k_mode: Arc<Kcp2KMode>,
cookie: Arc<Bytes>,
socket: Arc<Socket>,
client_sock_addr: Arc<SockAddr>,
) -> UdpOutput {
UdpOutput {
kcp2k_mode,
cookie,
socket,
client_sock_addr,
}
}
}
impl Write for UdpOutput {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let mut buffer = BytesMut::new();
buffer.put_u8(Kcp2KChannel::Reliable.to_u8());
buffer.put_slice(&self.cookie);
buffer.put_slice(buf);
match match *self.kcp2k_mode {
Kcp2KMode::Client => self.socket.send(&buffer),
Kcp2KMode::Server => self.socket.send_to(&buffer, &self.client_sock_addr),
} {
Ok(_) => Ok(buf.len()),
Err(err) => {
error!(format!("UdpOutput write error: {:?}", err));
Err(err)
}
}
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}