use std::time::Instant;
use std::sync::Arc;
use std::collections::HashMap;
use std::fmt;
use serde_json::Value as Json;
use serde::ser::{Serialize, Serializer, SerializeTuple};
use futures::sync::mpsc::{UnboundedSender as Sender};
use config;
use intern::{Topic, SessionId, SessionPoolName, Lattice as Namespace};
use intern::LatticeKey;
use chat::{Cid, ConnectionSender, CloseReason};
use chat::message::{Meta, MetaWithExtra};
use chat::error::MessageError;
mod main;
mod pool;
mod public;
mod session;
mod heap;
mod try_iter; mod connection;
mod lattice;
pub use self::public::{Processor, ProcessorPool};
pub use self::lattice::Delta;
#[derive(Debug)]
pub struct Event {
pool: SessionPoolName,
timestamp: Instant,
action: Action,
}
#[derive(Debug)]
pub enum ConnectionMessage {
Publish(Topic, Arc<Json>),
Hello(SessionId, Arc<Json>),
Result(Arc<Meta>, Json),
Lattice(Namespace, Arc<HashMap<LatticeKey, lattice::Values>>),
Error(Arc<Meta>, MessageError),
StopSocket(CloseReason),
}
#[derive(Debug)]
pub enum PoolMessage {
InactiveSession {
session_id: SessionId,
connections_active: usize,
metadata: Arc<Json>,
},
}
pub enum Action {
NewSessionPool {
config: Arc<config::SessionPool>,
channel: Sender<PoolMessage>,
},
StopSessionPool,
NewConnection {
conn_id: Cid,
channel: ConnectionSender,
},
Associate {
conn_id: Cid,
session_id: SessionId,
metadata: Arc<Json>
},
UpdateActivity {
conn_id: Cid,
timestamp: Instant,
},
Disconnect {
conn_id: Cid,
},
Subscribe {
conn_id: Cid,
topic: Topic,
},
Unsubscribe {
conn_id: Cid,
topic: Topic,
},
Publish {
topic: Topic,
data: Arc<Json>,
},
Attach {
namespace: Namespace,
conn_id: Cid,
},
Lattice {
namespace: Namespace,
delta: Delta,
},
Detach {
namespace: Namespace,
conn_id: Cid,
},
}
impl Serialize for ConnectionMessage {
fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error>
{
use self::ConnectionMessage::*;
let mut tup = serializer.serialize_tuple(3)?;
match *self {
Publish(ref topic, ref json) => {
#[derive(Serialize)]
struct Meta<'a> {
topic: &'a Topic,
}
tup.serialize_element("message")?;
tup.serialize_element(&Meta { topic: topic })?;
tup.serialize_element(json)?;
}
Hello(_, ref json) => {
tup.serialize_element("hello")?;
tup.serialize_element(&json!({}))?;
tup.serialize_element(json)?;
}
Lattice(ref namespace, ref json) => {
#[derive(Serialize)]
struct Meta<'a> {
namespace: &'a Namespace,
}
tup.serialize_element("lattice")?;
tup.serialize_element(&Meta { namespace: namespace })?;
tup.serialize_element(json)?;
}
Result(ref meta, ref json) => {
tup.serialize_element("result")?;
tup.serialize_element(&meta)?;
tup.serialize_element(json)?;
}
Error(ref meta, ref err) => {
tup.serialize_element("error")?;
let extra = match err {
&MessageError::HttpError(ref status, _) => {
json!({
"error_kind": "http_error",
"http_error": status.code(),
})
}
&MessageError::Utf8Error(_) => {
json!({"error_kind": "data_error"})
}
&MessageError::JsonError(_) => {
json!({"error_kind": "data_error"})
}
&MessageError::ValidationError(_) => {
json!({"error_kind": "validation_error"})
}
_ => {
json!({"error_kind": "internal_error"})
}
};
tup.serialize_element(&MetaWithExtra {
meta: meta, extra: extra
})?;
tup.serialize_element(&err)?;
}
StopSocket(ref reason) => {
tup.serialize_element("stop")?;
tup.serialize_element(&format!("{:?}", reason))?;
tup.serialize_element(&json!(null))?;
}
}
tup.end()
}
}
impl fmt::Debug for Action {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
use self::Action::*;
match self {
&NewSessionPool {..} => {
write!(f, "Action::NewSessionPool")
}
&StopSessionPool => {
write!(f, "Action::StopSessionPool")
}
&NewConnection { ref conn_id, .. } => {
write!(f, "Action::NewConnection({:?})", conn_id)
}
&Associate { ref conn_id, ref session_id, .. } => {
write!(f, "Action::Associate({:?}, {:?})", conn_id, session_id)
}
&UpdateActivity { ref conn_id, .. } => {
write!(f, "Action::UpdateActivity({:?})", conn_id)
}
&Disconnect { ref conn_id } => {
write!(f, "Action::Disconnect({:?})", conn_id)
}
&Subscribe { ref conn_id, ref topic } => {
write!(f, "Action::Subscribe({:?}, {:?})", conn_id, topic)
}
&Unsubscribe { ref conn_id, ref topic } => {
write!(f, "Action::Unsubscribe({:?}, {:?})", conn_id, topic)
}
&Publish { ref topic, .. } => {
write!(f, "Action::Publish({:?})", topic)
}
&Attach { ref conn_id, ref namespace } => {
write!(f, "Action::Attach({:?}, {:?})", conn_id, namespace)
}
&Lattice { ref namespace, .. } => {
write!(f, "Action::Lattice({:?})", namespace)
}
&Detach { ref conn_id, ref namespace } => {
write!(f, "Action::Detach({:?}, {:?})", conn_id, namespace)
}
}
}
}