use crate::transport::{
error::{TransportError, TransportResult},
memory_mock::memory_server,
protocol::{TransportCommand, TransportEvent},
transport_trait::Transport,
TransportId, TransportIdRef,
};
use lib3h_protocol::DidWork;
use std::collections::{HashMap, HashSet, VecDeque};
pub struct TransportMemory {
cmd_inbox: VecDeque<TransportCommand>,
my_servers: HashSet<String>,
connections: HashMap<TransportId, String>,
n_id: u32,
}
impl TransportMemory {
pub fn new() -> Self {
TransportMemory {
cmd_inbox: VecDeque::new(),
my_servers: HashSet::new(),
connections: HashMap::new(),
n_id: 0,
}
}
}
impl Transport for TransportMemory {
fn transport_id_list(&self) -> TransportResult<Vec<TransportId>> {
Ok(self.connections.keys().map(|id| id.to_string()).collect())
}
fn connect(&mut self, uri: &str) -> TransportResult<TransportId> {
let server_map = memory_server::MEMORY_SERVER_MAP.read().unwrap();
let maybe_server = server_map.get(uri);
if let None = maybe_server {
return Err(TransportError::new(format!(
"No Memory server at this url address: {}",
uri
)));
}
let mut server = maybe_server.unwrap().lock().unwrap();
self.n_id += 1;
let id = format!("{}__{}", uri, self.n_id);
server.connect(&id)?;
self.connections.insert(id.clone(), uri.to_string());
Ok(id)
}
fn close(&mut self, id: &TransportIdRef) -> TransportResult<()> {
let maybe_url = self.connections.get(id);
if let None = maybe_url {
return Err(TransportError::new(format!(
"No known connection for TransportId {}",
id
)));
}
let url = maybe_url.unwrap();
let server_map = memory_server::MEMORY_SERVER_MAP.read().unwrap();
let maybe_server = server_map.get(url);
if let None = maybe_server {
return Err(TransportError::new(format!(
"No Memory server at this url: {}",
url,
)));
}
let mut server = maybe_server.unwrap().lock().unwrap();
server.close(&id)?;
self.connections.remove(id);
Ok(())
}
fn close_all(&mut self) -> TransportResult<()> {
let id_list = self.transport_id_list()?;
for id in id_list {
self.close(&id)?;
}
Ok(())
}
fn send(&mut self, id_list: &[&TransportIdRef], payload: &[u8]) -> TransportResult<()> {
for id in id_list {
let maybe_url = self.connections.get(*id);
if let None = maybe_url {
println!("[w] No known connection for TransportId {}", id);
continue;
}
let url = maybe_url.unwrap();
let server_map = memory_server::MEMORY_SERVER_MAP.read().unwrap();
let maybe_server = server_map.get(url);
if let None = maybe_server {
return Err(TransportError::new(format!(
"No Memory server at this url address: {}",
url
)));
}
let mut server = maybe_server.unwrap().lock().unwrap();
server
.post(id, payload)
.expect("Post on memory server should work");
}
Ok(())
}
fn send_all(&mut self, payload: &[u8]) -> TransportResult<()> {
let id_list = self.transport_id_list()?;
for id in id_list {
self.send(&[id.as_str()], payload)?;
}
Ok(())
}
fn post(&mut self, command: TransportCommand) -> TransportResult<()> {
self.cmd_inbox.push_back(command);
Ok(())
}
fn bind(&mut self, url: &str) -> TransportResult<String> {
memory_server::set_server(url)?;
self.my_servers.insert(url.to_string());
Ok(url.to_string())
}
fn process(&mut self) -> TransportResult<(DidWork, Vec<TransportEvent>)> {
let mut outbox = Vec::new();
let mut did_work = false;
loop {
let cmd = match self.cmd_inbox.pop_front() {
None => break,
Some(msg) => msg,
};
let res = self.serve_TransportCommand(&cmd);
if let Ok(mut output) = res {
did_work = true;
outbox.append(&mut output);
}
}
for server_url in &self.my_servers {
let server_map = memory_server::MEMORY_SERVER_MAP.read().unwrap();
let server = server_map.get(server_url).expect("My server should exist.");
let (success, mut output) = server.lock().unwrap().process()?;
if success {
did_work = true;
outbox.append(&mut output);
}
}
Ok((did_work, outbox))
}
}
impl TransportMemory {
#[allow(non_snake_case)]
fn serve_TransportCommand(
&mut self,
cmd: &TransportCommand,
) -> TransportResult<Vec<TransportEvent>> {
println!("(log.d) >>> '(TransportMemory)' recv cmd: {:?}", cmd);
match cmd {
TransportCommand::Connect(url) => {
let id = self.connect(url)?;
let evt = TransportEvent::ConnectResult(id);
Ok(vec![evt])
}
TransportCommand::Send(id_list, payload) => {
let mut id_ref_list = Vec::with_capacity(id_list.len());
for id in id_list {
id_ref_list.push(id.as_str());
}
let _id = self.send(&id_ref_list, payload)?;
Ok(vec![])
}
TransportCommand::SendAll(payload) => {
let _id = self.send_all(payload)?;
Ok(vec![])
}
TransportCommand::Close(id) => {
self.close(id)?;
let evt = TransportEvent::Closed(id.to_string());
Ok(vec![evt])
}
TransportCommand::CloseAll => {
self.close_all()?;
let mut outbox = Vec::new();
for (id, _url) in &self.connections {
let evt = TransportEvent::Closed(id.to_string());
outbox.push(evt);
}
Ok(outbox)
}
TransportCommand::Bind(url) => {
self.bind(url)?;
Ok(vec![])
}
}
}
}