use std::collections::HashMap;
use std::net::{SocketAddrV4, SocketAddr};
use std::sync::Arc;
use std::time::{Duration, Instant};
use std::io;
use blowfish::Blowfish;
use rand::rngs::OsRng;
use rand::RngCore;
use super::bundle::{ElementReader, TopElementReader, BundleElement, BundleResult, BundleError, Bundle, ReplyElementReader, BundleElementWriter};
use super::socket::{BundleSocket, Bun};
use super::element::TopElement;
use super::packet::Packet;
pub mod login;
pub use login::{LoginInterface, LoginShared};
pub struct Interface<S: Shared> {
socket: BundleSocket,
shared: S,
bundle: Bundle,
top_callbacks: Box<[Option<TopCallback<S>>; 255]>,
request_manager: RequestManager<S>,
}
impl<S: Shared> Interface<S> {
const INIT_TOP_CALLBACK: Option<TopCallback<S>> = None;
pub fn new(addr: SocketAddrV4, shared: S) -> io::Result<Self> {
Ok(Self {
socket: BundleSocket::new(addr)?,
shared,
bundle: Bundle::new(),
top_callbacks: Box::new([Self::INIT_TOP_CALLBACK; 255]),
request_manager: RequestManager::new(),
})
}
#[inline]
pub fn shared(&self) -> &S {
&self.shared
}
#[inline]
pub fn shared_mut(&mut self) -> &mut S {
&mut self.shared
}
#[inline]
pub fn register_raw<U>(&mut self, id: u8, callback: U)
where
U: 'static + FnMut(&mut S, TopElementReader, Peer<S>) -> BundleResult<bool>,
{
assert_ne!(id, 0xFF, "id 0xFF reserved for reply elements");
self.top_callbacks[id as usize] = Some(Box::new(callback));
}
#[inline]
pub fn register_simple<E, U>(&mut self, id: u8, mut callback: U)
where
E: TopElement<Config = ()>,
U: 'static + FnMut(&mut S, BundleElement<E>, Peer<S>),
{
self.register_raw(id, move |state, reader, peer| {
callback(state, reader.read_simple::<E>()?, peer);
Ok(true) });
}
#[inline]
pub fn register<E, C, U, V>(&mut self, id: u8, mut callback: U, mut callback_config: V)
where
E: TopElement<Config = C>,
U: 'static + FnMut(&mut S, BundleElement<E>, Peer<S>),
V: 'static + FnMut(&mut S, SocketAddr) -> C,
{
self.register_raw(id, move |state, reader, peer| {
let config = callback_config(state, peer.addr);
callback(state, reader.read::<E>(&config)?, peer);
Ok(true) });
}
pub fn poll(&mut self, events: &mut Vec<Event>, timeout: Option<Duration>) -> Result<(), InterfaceError> {
self.socket.poll(events, timeout)?;
for event in events {
match &event.kind {
EventKind::Bundle(bundle) => {
let mut reader = bundle.element_reader();
while let Some(element) = reader.next_element() {
self.handle_element(event.addr, element)?;
}
}
EventKind::PacketError(packet, error) => {
self.shared.on_packet_error(&**packet, error);
}
}
}
Ok(())
}
fn handle_element(&mut self, addr: SocketAddr, element: ElementReader) -> BundleResult<bool> {
debug_assert!(self.bundle.is_empty());
match element {
ElementReader::Top(reader) => {
if let Some(mut callback) = self.top_callbacks[reader.id() as usize].take() {
callback(&mut self.shared, reader, Peer {
addr,
bundle: &mut self.bundle,
socket: &mut self.socket,
request_manager: &mut self.request_manager,
})
} else {
self.shared.on_element(reader)
}
}
ElementReader::Reply(reader) => {
if let Some((
mut callback,
_instant
)) = self.request_manager.callbacks.remove(&reader.request_id()) {
callback(&mut self.shared, reader, Peer {
addr,
bundle: &mut self.bundle,
socket: &mut self.socket,
request_manager: &mut self.request_manager,
})
} else {
self.shared.on_reply(reader)
}
}
}
}
pub fn peer(&mut self, addr: SocketAddr) -> Peer<S> {
debug_assert!(self.bundle.is_empty());
Peer {
addr,
bundle: &mut self.bundle,
socket: &mut self.socket,
request_manager: &mut self.request_manager,
}
}
}
type TopCallback<S> = Box<dyn FnMut(&mut S, TopElementReader, Peer<S>) -> BundleResult<bool>>;
type ReplyCallback<S> = Box<dyn FnMut(&mut S, ReplyElementReader, Peer<S>) -> BundleResult<bool>>;
pub trait Shared: 'static {
fn on_packet_error(&mut self, packet: &Packet, error: &BundleError) {
let _ = (packet, error);
}
fn on_element(&mut self, reader: TopElementReader) -> BundleResult<bool> {
let _ = reader;
Ok(false) }
fn on_reply(&mut self, reader: ReplyElementReader) -> BundleResult<bool> {
let _ = reader;
Ok(false) }
}
pub struct RequestManager<S> {
callbacks: HashMap<u32, (ReplyCallback<S>, Instant)>,
next_request_id: u32,
pending_callbacks: Vec<(u32, ReplyCallback<S>)>,
}
impl<S> RequestManager<S> {
fn new() -> Self {
Self {
callbacks: HashMap::new(),
next_request_id: OsRng.next_u32(),
pending_callbacks: Vec::new(),
}
}
fn flush(&mut self) {
if self.pending_callbacks.is_empty() {
let instant = Instant::now();
for (request_id, callback) in self.pending_callbacks.drain(..) {
self.callbacks.insert(request_id, (callback, instant));
}
}
}
fn abort(&mut self) {
self.pending_callbacks.clear();
}
fn add_pending_callback(&mut self, callback: ReplyCallback<S>) -> u32 {
let request_id = self.next_request_id;
self.next_request_id = request_id.wrapping_add(1);
self.pending_callbacks.push((request_id, callback));
request_id
}
#[inline]
pub fn register_raw<U>(&mut self, callback: U) -> u32
where
U: 'static + FnMut(&mut S, ReplyElementReader, Peer<S>) -> BundleResult<bool>,
{
self.add_pending_callback(Box::new(callback))
}
#[inline]
pub fn register_simple<E, U>(&mut self, mut callback: U) -> u32
where
E: TopElement<Config = ()>,
U: 'static + FnMut(&mut S, BundleElement<E>, Peer<S>),
{
self.register_raw(move |state, reader, peer| {
callback(state, reader.read_simple::<E>()?, peer);
Ok(true) })
}
#[inline]
pub fn register<E, C, U, V>(&mut self, mut callback: U, mut callback_config: V) -> u32
where
E: TopElement<Config = C>,
U: 'static + FnMut(&mut S, BundleElement<E>, Peer<S>),
V: 'static + FnMut(&mut S, SocketAddr) -> C,
{
self.register_raw(move |state, reader, peer| {
let config = callback_config(state, peer.addr);
callback(state, reader.read::<E>(&config)?, peer);
Ok(true) })
}
}
pub struct Peer<'a, S> {
addr: SocketAddr,
bundle: &'a mut Bundle,
socket: &'a mut BundleSocket,
request_manager: &'a mut RequestManager<S>,
}
impl<'a, S> Peer<'a, S> {
#[inline]
pub fn addr(&self) -> SocketAddr {
self.addr
}
#[inline]
pub fn set_channel(&mut self, blowfish: Arc<Blowfish>) {
self.socket.set_channel(self.addr, blowfish);
}
#[inline]
pub fn element_writer(&mut self) -> BundleElementWriter<'_> {
self.bundle.element_writer()
}
#[inline]
pub fn request_manager(&mut self) -> &mut RequestManager<S> {
self.request_manager
}
pub fn flush(&mut self) {
if !self.bundle.is_empty() {
self.socket.send(self.bundle, self.addr).expect("TODO: Change this");
self.bundle.clear();
self.request_manager.flush();
} else {
self.request_manager.abort();
}
}
pub fn abort(&mut self) {
self.bundle.clear();
self.request_manager.abort();
}
}
impl<'a, S> Drop for Peer<'a, S> {
fn drop(&mut self) {
self.flush();
}
}
#[derive(Debug, thiserror::Error)]
pub enum InterfaceError {
#[error("bundle error: {0}")]
Bundle(#[from] BundleError),
#[error("io error: {0}")]
Io(#[from] io::Error),
}
pub type InterfaceResult<T> = Result<T, InterfaceError>;