#![feature(slicing_syntax)]
#![allow(deprecated)]
extern crate time;
extern crate water;
use water::Net;
use water::Endpoint;
use water::RawMessage;
use water::NoPointers;
use water::Message;
use std::thread::JoinGuard;
use std::thread::Thread;
use time::Timespec;
use std::io::timer::sleep;
use std::time::duration::Duration;
struct SafeStructure {
a: u64,
b: u32,
c: u8,
}
impl NoPointers for SafeStructure {}
const THREADCNT: uint = 3;
fn funnyworker(mut net: Net, dbgid: uint) {
let ep: Endpoint = net.new_endpoint();
while net.getepcount() < THREADCNT + 1 { }
let limit = 1u;
let mut sentmsgcnt: uint = 0u;
let mut recvmsgcnt: uint = 0u;
let mut msgtosend = Message::new_raw(64);
msgtosend.dstsid = 0;
msgtosend.dsteid = 0;
let mut got: Vec<Vec<u8>> = Vec::new();
for i in range(0u, THREADCNT) {
got.push(Vec::new());
for _ in range(0u, limit) {
got[i].push(0u8);
}
}
while recvmsgcnt < limit * (THREADCNT - 1) || sentmsgcnt < limit {
loop {
let result = ep.recv();
if result.is_err() {
break;
}
let safestruct: SafeStructure = result.ok().get_raw().readstruct(0);
if safestruct.c != 0x10 {
if got[safestruct.a as uint][safestruct.b as uint] != 0u8 {
panic!("got {} twice from thread {}", safestruct.b, safestruct.a);
}
got[safestruct.a as uint][safestruct.b as uint] = 1u8;
recvmsgcnt += 1;
assert!(safestruct.c == 0x12u8);
}
}
if sentmsgcnt < limit {
let safestruct = SafeStructure {
a: dbgid as u64,
b: sentmsgcnt as u32,
c: 0x12,
};
msgtosend.get_rawmutref().writestructref(0, &safestruct);
let rby = ep.send(msgtosend.clone());
if rby < THREADCNT {
panic!("send {} less than {}", rby, THREADCNT);
}
sentmsgcnt += 1;
}
}
let safestruct = SafeStructure {
a: dbgid as u64,
b: 0x00,
c: 0x10,
};
msgtosend.get_rawmutref().writestructref(0, &safestruct);
ep.send(msgtosend);
println!("thread[{}]: exiting", dbgid);
}
#[test]
fn rawmessage() {
let m = RawMessage::new_fromstr("ABCDE");
assert!(m.readu8(0) == 65);
assert!(m.readu8(1) == 66);
assert!(m.readu8(2) == 67);
assert!(m.readu8(3) == 68);
assert!(m.readu8(4) == 69);
assert!(m.len() == 5);
}
#[test]
fn rawmsgstress() {
let mut v: Vec<RawMessage> = Vec::new();
for _ in range(0u, 10000u) {
let rm = RawMessage::new(32);
v.push(rm.dup());
v.push(rm);
}
}
#[test]
fn basicio() {
for _ in range(0u, 100u) {
_basicio();
}
}
fn _basicio() {
let mut net: Net = Net::new(234);
let ep = net.new_endpoint();
let mut completedcnt: u32 = 0u32;
let mut threadterm = [0u;THREADCNT];
let mut threads: Vec<JoinGuard<()>> = Vec::new();
for i in range(0, THREADCNT) {
let netclone = net.clone();
threads.push(Thread::spawn(move || { funnyworker(netclone, i); }));
}
let mut sectowait = 6i64;
loop {
let result = ep.recvorblock(Timespec { sec: sectowait, nsec: 0 });
sectowait = 6i64;
if result.is_err() {
panic!("timed out waiting for messages likely..");
}
let raw = result.ok().get_raw();
let safestruct: SafeStructure = raw.readstruct(0);
if safestruct.c == 0x10 {
if threadterm[safestruct.a as uint] != 0 {
panic!("got termination message from same thread twice!");
}
threadterm[safestruct.a as uint] = 1;
completedcnt += 1;
if completedcnt > 2 {
break;
}
}
println!("got msg");
}
}