use std::thread;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::{Duration, Instant};
use connection::{Event, Writer};
use message::Prefix;
#[derive(Clone)]
enum MonitorStatus {
Activity(Instant),
Ping(Instant),
}
#[derive(Clone)]
enum ConnectionStatus {
Connected(MonitorStatus),
Disconnected,
Quit,
}
#[derive(Clone)]
struct State {
status: Arc<Mutex<ConnectionStatus>>,
server: Arc<Mutex<Option<String>>>,
}
impl State {
fn new(ts: Instant) -> State {
let conn_status = ConnectionStatus::Connected(MonitorStatus::Activity(ts));
State {
status: Arc::new(Mutex::new(conn_status)),
server: Arc::new(Mutex::new(None)),
}
}
fn set_activity(&self, ts: Instant) {
*self.status.lock().unwrap() = ConnectionStatus::Connected(MonitorStatus::Activity(ts));
}
fn set_disconnected(&self) {
*self.status.lock().unwrap() = ConnectionStatus::Disconnected;
}
fn quit(&self) {
*self.status.lock().unwrap() = ConnectionStatus::Quit;
}
fn connection_status<'a>(&self) -> MutexGuard<ConnectionStatus> {
self.status.lock().unwrap()
}
fn has_server(&self) -> bool {
self.server.lock().unwrap().is_some()
}
fn get_server(&self) -> MutexGuard<Option<String>> {
self.server.lock().unwrap()
}
fn set_server(&self, name: String) {
*self.server.lock().unwrap() = Some(name);
}
fn unset_server(&self) {
*self.server.lock().unwrap() = None;
}
}
fn periodic_checker(state: State, handle: Writer, settings: MonitorSettings) {
loop {
let mut conn_status = state.connection_status();
match conn_status.clone() {
ConnectionStatus::Connected(ref mon_status) => {
match *mon_status {
MonitorStatus::Activity(activity_ts) => {
let diff = activity_ts.elapsed();
if diff > settings.activity_timeout {
match *state.get_server() {
Some(ref server) => {
*conn_status = ConnectionStatus::Connected(MonitorStatus::Ping(Instant::now()));
let _ = handle.raw(format!("PING {}\n", server));
}
None => {
panic!("Server is None! This scenario is highly unlikely, please report this issue!");
}
}
}
}
MonitorStatus::Ping(ping_ts) => {
let diff = ping_ts.elapsed();
if diff > settings.ping_timeout {
let _ = handle.disconnect();
}
},
}
},
ConnectionStatus::Disconnected => {},
ConnectionStatus::Quit => break,
}
drop(conn_status);
thread::sleep(Duration::from_secs(1));
}
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub struct MonitorSettings {
pub activity_timeout: Duration,
pub ping_timeout: Duration,
}
impl Default for MonitorSettings {
fn default() -> MonitorSettings {
MonitorSettings {
activity_timeout: Duration::from_secs(60),
ping_timeout: Duration::from_secs(15),
}
}
}
pub struct ActivityMonitor {
state: State,
}
impl ActivityMonitor {
pub fn new(handle: &Writer, settings: MonitorSettings) -> ActivityMonitor {
let state = State::new(Instant::now());
let state_clone = state.clone();
let handle_clone = handle.clone();
thread::spawn(move || {
periodic_checker(state_clone, handle_clone, settings);
});
ActivityMonitor {
state: state,
}
}
pub fn feed(&self, event: &Event) {
match *event {
Event::Closed(_) => {
self.state.quit();
}
Event::Disconnected => {
self.state.set_disconnected();
self.state.unset_server();
}
Event::Reconnected => {
self.state.set_activity(Instant::now());
}
Event::Message(ref msg) => {
self.state.set_activity(Instant::now());
if let Some(ref prefix) = msg.prefix {
match *prefix {
Prefix::Server(ref name) => {
if !self.state.has_server() {
self.state.set_server(name.clone());
}
}
_ => {}
}
}
}
_ => {}
}
}
}
impl Drop for ActivityMonitor {
fn drop(&mut self) {
self.state.quit();
}
}