#![allow(unused_imports)]
#![allow(dead_code)]
#![allow(unused_variables)]
extern crate core;
use std::sync::Arc;
use std::ptr;
use std::rt::heap::allocate;
use std::rt::heap::deallocate;
use std::mem::size_of;
use std::sync::Mutex;
use std::intrinsics::copy_memory;
use std::intrinsics::transmute;
use std::io::timer::sleep;
use std::time::duration::Duration;
use time::get_time;
use time::Timespec;
use rawmessage::RawMessage;
use endpoint::Endpoint;
use message::Message;
use message::MessagePayload;
use tcp;
use tcp::TcpBridgeListener;
use tcp::TcpBridgeConnector;
pub type ID = u64;
pub const UNUSED_ID: ID = !0u64;
struct Internal {
endpoints: Vec<Endpoint>,
hueid: ID, }
pub struct Net {
i: Arc<Mutex<Internal>>,
sid: ID,
}
impl Clone for Net {
fn clone(&self) -> Net {
Net {
i: self.i.clone(),
sid: self.sid,
}
}
}
const NET_MAXLATENCY: i64 = 100000; const NET_MINLATENCY: i64 = 100;
impl Net {
fn sleeperthread(net: Net) {
let mut latency: i64 = NET_MAXLATENCY;
let mut wokesomeone: bool;
loop {
sleep(Duration::microseconds(latency));
let ctime: Timespec = get_time();
{
let mut i = net.i.lock();
wokesomeone = false;
for ep in i.endpoints.iter_mut() {
if ctime > ep.getwaketime() {
if ep.wakeonewaiter() {
wokesomeone = true;
}
}
}
}
if wokesomeone {
latency = latency / 2;
if latency < NET_MINLATENCY {
latency = NET_MINLATENCY;
}
} else {
latency = latency * 2;
if latency > NET_MAXLATENCY {
latency = NET_MAXLATENCY;
}
}
}
}
pub fn getepcount(&self) -> uint {
self.i.lock().endpoints.len()
}
pub fn new(sid: ID) -> Net {
let net = Net {
i: Arc::new(Mutex::new(Internal {
endpoints: Vec::new(),
hueid: 0x10000,
})),
sid: sid,
};
let netclone = net.clone();
spawn(move || { Net::sleeperthread(netclone) });
net
}
pub fn tcplisten(&self, addr: String) -> TcpBridgeListener {
tcp::listener::TcpBridgeListener::new(self, addr)
}
pub fn tcpconnect(&self, addr: String) -> TcpBridgeConnector {
tcp::connector::TcpBridgeConnector::new(self, addr)
}
pub fn sendas(&self, msg: &mut Message, fromsid: ID, fromeid: ID) -> uint {
msg.srcsid = fromsid;
msg.srceid = fromeid;
self.send(msg)
}
pub fn send(&self, msg: &Message) -> uint {
if msg.is_sync() {
panic!("You must send a sync message with `sendsync` or `sendsyncas`!");
}
if msg.is_raw() {
self.send_internal(&(msg.dup()))
} else {
self.send_internal(msg)
}
}
pub fn sendsyncas(&self, mut msg: Message, frmsid: ID, frmeid: ID) -> uint {
msg.srcsid = frmsid;
msg.srceid = frmeid;
self.sendsync(msg)
}
pub fn sendsync(&self, msg: Message) -> uint {
self.send_internal(&msg)
}
fn send_internal(&self, msg: &Message) -> uint {
let mut ocnt = 0u;
let mut i = self.i.lock();
for ep in i.endpoints.iter_mut() {
if ep.give(msg) {
ocnt += 1;
}
}
ocnt
}
pub fn get_neweid(&mut self) -> ID {
let mut i = self.i.lock();
let eid = i.hueid;
i.hueid += 1;
eid
}
pub fn new_endpoint_withid(&mut self, eid: ID) -> Endpoint {
let mut i = self.i.lock();
if eid > i.hueid {
i.hueid = eid + 1;
}
let ep = Endpoint::new(self.sid, eid, self.clone());
i.endpoints.push(ep.clone());
ep
}
pub fn add_endpoint(&mut self, ep: Endpoint) {
self.i.lock().endpoints.push(ep);
}
pub fn new_endpoint(&mut self) -> Endpoint {
let mut i = self.i.lock();
let ep = Endpoint::new(self.sid, i.hueid, self.clone());
i.hueid += 1;
let epclone = ep.clone();
i.endpoints.push(epclone);
ep
}
pub fn getserveraddr(&self) -> ID {
self.sid
}
}