#![allow(unused_imports)]
#![allow(dead_code)]
#![allow(unused_variables)]
use std;
use std::sync::Arc;
use std::sync::atomic;
use std::sync::atomic::AtomicUint;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::Ordering;
use std::ptr;
use std::rt::heap::allocate;
use std::mem::size_of;
use std::mem::align_of;
use std::rt::heap::deallocate;
use std::sync::Mutex;
use std::sync::MutexGuard;
use std::sync::Condvar;
use std::time::duration::Duration;
use std::mem::uninitialized;
use std::intrinsics::copy_memory;
use std::mem::transmute;
use std::mem::transmute_copy;
use std::thread::Thread;
use Queue;
use time::Timespec;
use time::get_time;
use timespec;
use net::Net;
use net::ID;
use net::UNUSED_ID;
use rawmessage::RawMessage;
use message::Message;
use message::MessagePayload;
pub enum IoErrorCode {
TimedOut,
NoMessages,
}
impl Copy for IoErrorCode { }
pub struct IoError {
pub code: IoErrorCode,
}
impl Copy for IoError { }
pub enum IoResult<T> {
Err(IoError),
Ok(T),
}
impl<T> IoResult<T> {
pub fn ok(self) -> T {
match self {
IoResult::Ok(v) => v,
IoResult::Err(e) => panic!("not `IoResult::Ok`!"),
}
}
pub fn err(self) -> IoError {
match self {
IoResult::Ok(v) => panic!("not `IoResult::Err`!"),
IoResult::Err(e) => e,
}
}
pub fn is_ok(&self) -> bool {
match *self {
IoResult::Ok(_) => true,
IoResult::Err(_) => false,
}
}
pub fn is_err(&self) -> bool {
match *self {
IoResult::Ok(_) => false,
IoResult::Err(_) => true,
}
}
}
struct SleepToken {
condvar: Condvar,
mutex: Mutex<()>,
}
impl SleepToken {
pub fn new() -> SleepToken {
SleepToken {
condvar: Condvar::new(),
mutex: Mutex::new(()),
}
}
pub fn wait(&self) {
let lock = self.mutex.lock().unwrap();
self.condvar.wait(lock).unwrap();
}
pub fn notify_all(&self) {
self.condvar.notify_all();
}
}
struct Internal {
stoken: SleepToken,
messages: Queue<Message>,
wakeupat: Timespec,
memoryused: AtomicUint,
net: Net,
refcnt: AtomicUint,
slpcnt: AtomicUint,
limitpending: uint,
limitmemory: uint,
eid: u64,
sid: u64,
gid: u64,
}
unsafe impl Send for Internal { }
pub struct Endpoint {
i: Arc<Internal>,
ui: *mut Internal,
}
unsafe impl Send for Endpoint { }
impl Clone for Endpoint {
fn clone(&self) -> Endpoint {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.refcnt.fetch_add(1, Ordering::SeqCst);
Endpoint {
i: self.i.clone(),
ui: self.ui,
}
}
}
impl Drop for Endpoint {
fn drop(&mut self) {
let i: &mut Internal = unsafe { transmute (self.ui) };
let refcnt = i.refcnt.fetch_sub(1, Ordering::SeqCst) - 1;
if refcnt == 1 {
let mut net = i.net.clone();
net.drop_endpoint(self);
}
}
}
impl Internal {
fn neverwakeme(&mut self) {
self.wakeupat = Timespec { sec: 0x7fffffffffffffffi64, nsec: 0i32 };
}
fn recv(&mut self) -> IoResult<Message> {
loop {
let result = self.messages.get();
if result.is_none() {
return IoResult::Err(IoError { code: IoErrorCode::NoMessages });
}
let msg = result.unwrap().dup_ifok();
self.memoryused.fetch_sub(msg.cap(), Ordering::SeqCst);
match msg.payload {
MessagePayload::Raw(_) => { return IoResult::Ok(msg.dup()); },
MessagePayload::Sync(_) => {
if msg.get_syncref().takeasvalid() {
return IoResult::Ok(msg);
}
},
MessagePayload::Clone(_) => { return IoResult::Ok(msg); },
}
}
}
}
impl Endpoint {
pub fn id(&self) -> uint {
self.ui as uint
}
pub fn new(sid: u64, eid: u64, net: Net) -> Endpoint {
let mut ep = Endpoint {
i: Arc::new(Internal {
messages: Queue::new(),
stoken: SleepToken::new(),
wakeupat: Timespec { nsec: 0i32, sec: 0x7fffffffffffffffi64 },
limitpending: 0,
limitmemory: 0,
memoryused: AtomicUint::new(0),
sid: sid,
eid: eid,
net: net,
gid: UNUSED_ID,
refcnt: AtomicUint::new(1),
slpcnt: AtomicUint::new(0),
}),
ui: 0 as *mut Internal,
};
ep.ui = unsafe { transmute(&*(ep.i)) };
ep
}
pub fn getwaketime(&self) -> Timespec {
unsafe { (*(self.ui)).wakeupat }
}
pub fn getpeercount(&self) -> uint {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.net.getepcount()
}
pub fn give(&mut self, msg: &Message) -> bool {
let i: &mut Internal = unsafe { transmute(self.ui) };
if !msg.canloop && msg.srceid == i.eid && msg.srcsid == i.sid {
return false;
}
if msg.dstsid != 0 {
if msg.dstsid != 1 {
if msg.dstsid != i.sid {
return false;
}
} else {
if i.sid != i.net.getserveraddr() {
return false;
}
}
}
if msg.dsteid != 0 && msg.dsteid != i.eid {
return false;
}
if i.limitpending > 0 && i.messages.len() >= i.limitpending {
return false;
}
if i.limitmemory > 0 && i.memoryused.load(Ordering::SeqCst) >= i.limitmemory {
return false;
}
let cloned;
if msg.is_sync() {
cloned = (*msg).internal_clone(0x879);
} else {
cloned = (*msg).clone();
}
i.messages.put(cloned);
i.memoryused.fetch_add(msg.cap(), Ordering::SeqCst);
self.wakeonewaiter();
true
}
pub fn sleepercount(&self) -> uint {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.slpcnt.load(Ordering::Relaxed)
}
pub fn hasmessages(&self) -> bool {
let i: &mut Internal = unsafe { transmute (self.ui) };
if i.messages.len() > 0 {
true
} else {
false
}
}
pub fn wakeonewaiter(&self) {
let i: &mut Internal = unsafe { transmute (self.ui) };
}
pub fn setlimitpending(&mut self, limit: uint) {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.limitpending = limit;
}
pub fn setlimitmemory(&mut self, limit: uint) {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.limitmemory = limit;
}
pub fn getsid(&self) -> ID {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.sid
}
pub fn geteid(&self) -> ID {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.eid
}
pub fn getgid(&self) -> ID {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.gid
}
pub fn setgid(&mut self, id: ID) {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.gid = id;
}
pub fn setsid(&mut self, id: ID) {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.sid = id;
}
pub fn seteid(&mut self, id: ID) {
let i: &mut Internal = unsafe { transmute (self.ui) };
i.eid = id;
}
pub fn sendx(&self, msg: Message) -> uint {
let i: &mut Internal = unsafe { transmute (self.ui) };
let net = i.net.clone();
net.send(msg)
}
pub fn send(&self, msg: Message) -> uint {
let i: &mut Internal = unsafe { transmute (self.ui) };
let net;
let sid;
let eid;
net = i.net.clone();
sid = i.sid;
eid = i.eid;
net.sendas(msg, sid, eid)
}
pub fn sendsynctype<T: Send>(&self, t: T) -> uint {
let mut msg = Message::new_sync(t);
msg.dstsid = 1; msg.dsteid = 0; self.send(msg)
}
pub fn sendclonetype<T: Send + Clone>(&self, t: T) -> uint {
let mut msg = Message::new_clone(t);
msg.dstsid = 1; msg.dsteid = 0; self.send(msg)
}
pub fn recvorblockforever(&self) -> IoResult<Message> {
let ui: &mut Internal = unsafe { transmute(self.ui) };
loop {
let r = ui.recv();
if r.is_ok() {
return r;
}
Thread::yield_now();
}
}
pub fn recvorblock(&self, duration: Timespec) -> IoResult<Message> {
let ui: &mut Internal = unsafe { transmute(self.ui) };
let mut when: Timespec = get_time();
when = timespec::add(when, duration);
while ui.messages.len() < 1 {
let ctime: Timespec = get_time();
if ctime > when && ui.messages.len() < 1 {
ui.neverwakeme();
return IoResult::Err(IoError { code: IoErrorCode::TimedOut });
}
}
ui.neverwakeme();
ui.recv()
}
pub fn recv(&self) -> IoResult<Message> {
let ui: &mut Internal = unsafe { transmute(self.ui) };
ui.recv()
}
}