pub use tcp::listener::TcpBridgeListener;
pub use tcp::connector::TcpBridgeConnector;
use std::io::TcpStream;
use endpoint::Endpoint;
use message::Message;
use rawmessage::RawMessage;
use net::ID;
use net::UNUSED_ID;
use time::Timespec;
pub mod listener;
pub mod connector;
pub struct TerminateMessage;
impl Copy for TerminateMessage { }
impl Clone for TerminateMessage {
fn clone(&self) -> TerminateMessage { TerminateMessage }
}
fn getok<T, E>(result: Result<T, E>) -> T {
match result {
Ok(r) => r,
Err(e) => panic!("I/O error"),
}
}
pub enum Which<T, V> {
Listener(T),
Connector(V),
}
pub fn thread_rx(mut which: Which<TcpBridgeListener, TcpBridgeConnector>, mut ep: Endpoint, mut stream: TcpStream) {
let rsid: u64 = getok(stream.read_be_u64());
match which {
Which::Listener(ref mut bridge) => bridge.negcountinc(),
Which::Connector(ref mut bridge) => bridge.setconnected(true),
}
ep.setsid(rsid);
loop {
let ioresult = stream.read_be_u64();
let mut msgsize: u64;
msgsize = match ioresult {
Ok(r) => r,
Err(e) => {
return;
}
};
let msg_type: u8 = getok(stream.read_u8());
let msg_srcsid: u64 = getok(stream.read_be_u64());
let msg_srceid: u64 = getok(stream.read_be_u64());
let msg_dstsid: u64 = getok(stream.read_be_u64());
let msg_dsteid: u64 = getok(stream.read_be_u64());
msgsize -= 1 + 8 * 4;
let mut vbuf: Vec<u8> = Vec::from_elem(msgsize as uint, 0u8);
stream.read_at_least(msgsize as uint, vbuf.as_mut_slice());
if msg_type != 1u8 {
panic!("got message type {} instead of 1", msg_type);
}
let mut rmsg = RawMessage::new(msgsize as uint);
rmsg.write_from_slice(0, vbuf.as_slice());
let mut msg = Message::new_fromraw(rmsg);
msg.dstsid = msg_dstsid;
msg.dsteid = msg_dsteid;
msg.srcsid = msg_srcsid;
msg.srceid = msg_srceid;
ep.sendraw(&msg);
}
}
pub fn thread_tx(mut which: Which<TcpBridgeListener, TcpBridgeConnector>, mut ep: Endpoint, mut stream: TcpStream, sid: ID) {
stream.write_be_u64(sid);
loop {
let result = ep.recvorblock(Timespec { sec: 900i64, nsec: 0i32 });
if result.is_err() {
continue;
}
let msg = result.ok();
if msg.is_type::<TerminateMessage>() {
stream.close_read();
stream.close_write();
return;
}
if !msg.is_raw() {
continue;
}
let srcsid = msg.srcsid;
let srceid = msg.srceid;
let dstsid = msg.dstsid;
let dsteid = msg.dsteid;
let rmsg = msg.get_raw();
stream.write_be_u64((1 + 8 * 4 + rmsg.len()) as u64);
stream.write_u8(1u8);
stream.write_be_u64(srcsid);
stream.write_be_u64(srceid);
stream.write_be_u64(dstsid);
stream.write_be_u64(dsteid);
stream.write(rmsg.as_slice());
}
}