#[cfg(test)]
mod agent_test;
use shared::error::*;
use std::collections::{HashMap, VecDeque};
use std::time::Instant;
use crate::message::*;
#[derive(Default)]
pub struct Agent {
transactions: HashMap<TransactionId, AgentTransaction>,
closed: bool,
events_queue: VecDeque<Event>,
}
#[derive(Debug)] pub struct Event {
pub id: TransactionId,
pub evt: StunEvent,
}
#[derive(Debug)] pub enum StunEvent {
AgentClosed,
TransactionStopped,
TransactionTimeOut,
Message(Message),
}
pub(crate) struct AgentTransaction {
id: TransactionId,
deadline: Instant,
}
const AGENT_COLLECT_CAP: usize = 100;
#[derive(Debug)]
pub enum ClientAgent {
Process(Message),
Collect(Instant),
Start(TransactionId, Instant),
Stop(TransactionId),
Close,
}
impl Agent {
pub fn new() -> Self {
Agent {
transactions: HashMap::new(),
closed: false,
events_queue: VecDeque::new(),
}
}
pub fn handle_event(&mut self, client_agent: ClientAgent) -> Result<()> {
match client_agent {
ClientAgent::Process(message) => self.process(message),
ClientAgent::Collect(deadline) => self.collect(deadline),
ClientAgent::Start(tid, deadline) => self.start(tid, deadline),
ClientAgent::Stop(tid) => self.stop(tid),
ClientAgent::Close => self.close(),
}
}
pub fn poll_timeout(&mut self) -> Option<Instant> {
let mut deadline = None;
for transaction in self.transactions.values() {
if deadline.is_none() || transaction.deadline < *deadline.as_ref().unwrap() {
deadline = Some(transaction.deadline);
}
}
deadline
}
pub fn poll_event(&mut self) -> Option<Event> {
self.events_queue.pop_front()
}
fn process(&mut self, message: Message) -> Result<()> {
if self.closed {
return Err(Error::ErrAgentClosed);
}
self.transactions.remove(&message.transaction_id);
self.events_queue.push_back(Event {
id: message.transaction_id,
evt: StunEvent::Message(message),
});
Ok(())
}
fn close(&mut self) -> Result<()> {
if self.closed {
return Err(Error::ErrAgentClosed);
}
for id in self.transactions.keys() {
self.events_queue.push_back(Event {
id: *id,
evt: StunEvent::AgentClosed,
});
}
self.transactions.clear();
self.closed = true;
Ok(())
}
fn start(&mut self, id: TransactionId, deadline: Instant) -> Result<()> {
if self.closed {
return Err(Error::ErrAgentClosed);
}
if self.transactions.contains_key(&id) {
return Err(Error::ErrTransactionExists);
}
self.transactions
.insert(id, AgentTransaction { id, deadline });
Ok(())
}
fn stop(&mut self, id: TransactionId) -> Result<()> {
if self.closed {
return Err(Error::ErrAgentClosed);
}
let v = self.transactions.remove(&id);
if let Some(t) = v {
self.events_queue.push_back(Event {
id: t.id,
evt: StunEvent::TransactionStopped,
});
Ok(())
} else {
Err(Error::ErrTransactionNotExists)
}
}
fn collect(&mut self, deadline: Instant) -> Result<()> {
if self.closed {
return Err(Error::ErrAgentClosed);
}
let mut to_remove: Vec<TransactionId> = Vec::with_capacity(AGENT_COLLECT_CAP);
for (id, t) in &self.transactions {
if t.deadline < deadline {
to_remove.push(*id);
}
}
for id in &to_remove {
self.transactions.remove(id);
}
for id in to_remove {
self.events_queue.push_back(Event {
id,
evt: StunEvent::TransactionTimeOut,
});
}
Ok(())
}
}