use std::net;
use std::{net::ToSocketAddrs, process, io::Write, io::Read, ops,
thread, time::Duration, clone, fs, sync, collections::HashSet,
env, path, process::Command, process::Stdio};
use portman_client;
use nscldaq_ringbuffer::ringbuffer;
use url::Url;
use local_ip_address;
use std::os::fd::IntoRawFd;
use std::os::fd::FromRawFd;
const PORTMAN_PORT: u16 =30000;
const RINGMASTER_SERVICE : &str = "RingMaster";
const RINGBUFFER_DIRECTORY : &str = "/dev/shm";
const DEFAULT_RING_SIZE : u32 =8*1024*1024;
fn parse_tcl_list(input: &str) -> Vec<String> {
let mut result = Vec::new();
let mut chars = input.chars().peekable();
while let Some(c) = chars.peek() {
match c {
'{' => {
let mut brace_count = 0;
let mut current_item = String::new();
chars.next(); while let Some(inner_c) = chars.next() {
if inner_c == '{' {
brace_count += 1;
} else if inner_c == '}' {
if brace_count == 0 {
break; } else {
brace_count -= 1;
}
}
current_item.push(inner_c);
}
result.push(current_item.trim().to_string()); }
'"' => {
let mut current_item = String::new();
chars.next(); while let Some(inner_c) = chars.next() {
if inner_c == '"' {
break; }
current_item.push(inner_c);
}
result.push(current_item);
}
_ if c.is_whitespace() => {
chars.next();
}
_ => {
let mut current_item = String::new();
while let Some(inner_c) = chars.peek() {
if inner_c.is_whitespace() || *inner_c == '{' || *inner_c == '"' {
break;
}
current_item.push(chars.next().unwrap());
}
if !current_item.is_empty() {
result.push(current_item);
}
}
}
}
result
}
fn stdin_to_ring_path() -> Result<String, String> {
let program_file = "stdintoring";
let bindir_env = "DAQBIN";
match env::var(bindir_env) {
Ok(bindir) => {
let full_path = format!("{}/{}", bindir, program_file);
let path = path::Path::new(&full_path);
match path.try_exists() {
Ok(yesno) => {
if yesno {
return Ok(String::from(path.to_str().unwrap()));
} else {
return Err(format!("{} can't be found", path.to_str().unwrap()));
}
},
Err(reason) => {
return Err(format!("Could not look up {} in the filesystem: {}", path.to_str().unwrap(), reason));
}
}
},
Err(reason) =>
return Err(format!("DAQBIN must be defined: {} ", reason))
}
}
fn proxy_ring_name(host: &str, remote_ring: &str) -> String {
format!("{}.{}", host, remote_ring)
}
fn proxy_ring_size() -> u32 {
let mut result = DEFAULT_RING_SIZE;
match env::var("NSCLDAQ_DEFAULT_PROXYMB") {
Ok(strval) => {
let intval = strval.parse::<u32>();
if let Ok(mb) = intval {
result = mb*1024*1024;
}
},
Err(_) => {
}
};
return result;
}
pub fn ring_path(ring: &str) -> String {
format!("{}/{}", RINGBUFFER_DIRECTORY, ring)
}
#[derive(Debug)]
pub struct Client {
host : String, socket : Option<net::TcpStream>, }
impl ops::Drop for Client {
#[doc = r"Shutdown the socket if it exists -- don't care if the shutdown fails."]
fn drop(&mut self) {
if let Some(s) = &self.socket {
let _ = s.shutdown(net::Shutdown::Both); }
}
}
impl clone::Clone for Client {
#[doc = r"Provide clone... if can't clone socket it'll be None."]
fn clone(&self) -> Self {
Client {
host: self.host.clone(),
socket: match &self.socket {
Some(s) => {
match s.try_clone() {
Ok(s) => Some(s),
Err(_) => None
}
},
None => None
}
}
}
}
#[derive(Debug)]
pub struct RingInformation {
pub name : String,
pub size : u32,
pub free : u32,
pub maxconsumers : u32,
pub producer_pid : i32, pub max_get : u32,
pub min_get : u32,
pub consumers : Vec<ringbuffer::ConsumerUsage>,
}
impl Client {
fn get_port(&self) -> Result<u16, portman_client::Error> {
let mut portman = portman_client::Client::new(PORTMAN_PORT);
let matches = portman.find_by_service(RINGMASTER_SERVICE)?;
if matches.len() > 0 {
Ok(matches[0].port)
} else {
Err(portman_client::Error::Unimplemented)
}
}
fn connect(&self) -> std::io::Result<net::TcpStream> {
let port = self.get_port();
if let Ok(num) = port {
net::TcpStream::connect(format!("{}:{}", self.host, num))
} else {
Err(std::io::Error::other("Failed to get ringmaster port"))
}
}
fn connect_persistent(&mut self) -> Result<net::TcpStream, std::io::Error> {
let sock_result = self.connect();
match sock_result {
Ok(socket) => {
self.socket = Some(socket.try_clone().unwrap());
return Ok(socket);
},
Err(e) => Err(e)
}
}
fn read_line(sock: &mut net::TcpStream) -> Result<String, std::io::Error> {
let mut reply = String::from(""); let mut buf : [u8; 1] = [0]; loop {
let n = sock.read(&mut buf)?;
if n == 0 {
break; }
let c = buf[0] as char;
if c == '\n' {
break; }
if c != '\r' {
reply.push(buf[0] as char);
}
}
Ok(reply)
}
fn transaction(&mut self, sock : &mut net::TcpStream, request : &str) -> Result<String, std::io::Error> {
sock.write_all(request.as_bytes())?;
sock.flush()?;
Client::read_line(sock)
}
fn get_socket_or_create(&mut self) -> Result<net::TcpStream, String> {
if let None = self.socket {
if let Err(result) = self.connect_persistent() {
return Err(format!("Failed to open socket to ringmaster {}", result));
}
}
return Ok(self.socket()
.expect("Ring master client expected to have socket but did not"));
}
fn transact_and_analyze(&mut self, request : &str) -> Result<(), String> {
let mut socket = self.get_socket_or_create()?;
let response = match self.transaction(&mut socket, &request) {
Ok(reply) => Ok(reply),
Err(reason) => Err(format!("Ring master transaction failed {}", reason))
}?;
if response == "OK" {
return Ok(());
} else if response == "OK BINARY FOLLOWS"
{
return Ok(());
} else {
let _ = socket.shutdown(net::Shutdown::Both);
self.socket = None;
return Err(response);
}
}
fn parse_else<T: std::str::FromStr>(string : &str) -> Result<T, String> {
let parsed = string.parse::<T>();
match parsed {
Ok(result) => Ok(result),
Err(_) => Err(format!("Numeric parse failed for {} ", string))
}
}
fn analyze_consumer(consumer: &str, ring_size: u32) -> Result<ringbuffer::ConsumerUsage, String> {
let consumer_list = parse_tcl_list(consumer);
if consumer_list.len() != 2 {
return Err(format!("Expected consumer info to have {} entries, got {}", 2, consumer_list.len()));
}
let consumer_pid : u32 = Self::parse_else(&consumer_list[0])?;
let consumer_maxget : u32 = Self::parse_else(&consumer_list[1])?;
let result = ringbuffer::ConsumerUsage {
pid: consumer_pid,
free : (ring_size - consumer_maxget) as usize,
available : consumer_maxget as usize
};
Ok(result)
}
fn analyze_ring(ring : &str) -> Result<RingInformation, String> {
let ring_list = parse_tcl_list(ring);
if ring_list.len() != 2 {
return Err(format!("Ring info was not correct expected {} got {}", 2, ring_list.len()));
}
let name = ring_list[0].clone();
let info_list = parse_tcl_list(&ring_list[1]);
if info_list.len() != 7 {
return Err(format!("Size/consumer list incorrect size expected {} got {}",7, info_list.len()));
}
let ring_size : u32 = Self::parse_else(&info_list[0])?;
let ring_free : u32 = Self::parse_else(&info_list[1])?;
let max_consumer : u32 = Self::parse_else(&info_list[2])?;
let producer_pid :i32 = Self::parse_else(&info_list[3])?;
let max_get : u32 = Self::parse_else(&info_list[4])?;
let min_get : u32 = Self::parse_else(&info_list[5])?;
let consumers = parse_tcl_list(&info_list[6]);
let mut result = RingInformation {
name : name.clone(),
size : ring_size,
free : ring_free,
maxconsumers : max_consumer,
producer_pid: producer_pid,
max_get : max_get,
min_get: min_get,
consumers : vec![]
};
for consumer in consumers {
result.consumers.push(Self::analyze_consumer(&consumer, ring_size)?);
}
Ok(result)
}
fn analyze_ring_listing(tcl_list : &str) -> Result<Vec<RingInformation>, String> {
let mut result = Vec::<RingInformation>::new();
let ring_list = parse_tcl_list(tcl_list); for ring in ring_list {
result.push(Client::analyze_ring(&ring)?);
}
Ok(result)
}
pub fn new(host : &str) -> Client {
Client {
host: host.to_string(),
socket: None
}
}
pub fn host(&self) -> String {
self.host.clone()
}
pub fn socket(&self) -> Option<net::TcpStream> {
match &self.socket {
None => None,
Some(s) => Some(s.try_clone().expect("Ringmaster client unable to clone the socket"))
}
}
pub fn register_ring(&mut self, ringname : &str) -> Result<(), String> {
let request = format!("REGISTER {} \n", ringname);
self.transact_and_analyze(&request)
}
pub fn connect_as_producer(&mut self, ringname : &str, comment : &str) -> Result<(), String> {
let pid = process::id();
let request = format!("CONNECT {{{}}} producer {} \"{}\" \n", ringname, pid, comment);
self.transact_and_analyze(&request)
}
pub fn disconnect_producer(&mut self, ringname: &str) -> Result<(), String> {
let pid = process::id();
let request = format!("DISCONNECT {{{}}} producer {}\n", ringname, pid);
self.transact_and_analyze(&request)
}
pub fn connect_as_consumer(&mut self, ring: &str, slot : u32, comment: &str) -> Result<(), String> {
let pid = process::id();
let request = format!("CONNECT {{{}}} consumer.{} {} \"{}\" \n", ring, slot, pid, comment);
self.transact_and_analyze(&request)
}
pub fn disconnect_consumer(&mut self, ring : &str, slot : u32) -> Result<(), String> {
let pid = process::id();
let request = format!("DISCONNECT {{{}}} consumer.{} {} \n", ring, slot, pid);
self.transact_and_analyze(&request)
}
pub fn list_rings(&mut self) -> Result<Vec<RingInformation>, String> {
self.transact_and_analyze("LIST\n")?;
let mut socket = self.socket()
.expect("If list_rings got here there shoulid still be a socket!! but there wasn't");
let tcl_list = Client::read_line(&mut socket);
if let Err(ioerr) = tcl_list {
return Err(format!("Error reading listing: {}", ioerr));
}
let tcl_list = tcl_list.unwrap();
Client::analyze_ring_listing(&tcl_list)
}
pub fn unregister_ring(&mut self, name :&str) -> Result<(), String> {
let request = format!("UNREGISTER {}\n", name);
self.transact_and_analyze(&request)
}
pub fn get_data(&mut self, ring: &str) -> Result<net::TcpStream, String> {
let request = format!("REMOTE {} \n", ring);
if let Err(s) = self.transact_and_analyze(&request) {
Err(s)
} else {
thread::sleep(Duration::from_secs(2));
let result = self.socket().unwrap().try_clone().unwrap();
self.socket = None;
Ok(result)
}
}
}
pub struct RingBufferProducer {
ringmaster : Client,
name : String,
pub ring : ringbuffer::producer::Producer }
impl ops::Drop for RingBufferProducer {
fn drop(&mut self) {
let _ = self.ringmaster.disconnect_producer(&self.name);
}
}
impl RingBufferProducer {
fn produce(name : &str, path: &str) -> Result<RingBufferProducer, String> {
let map = ringbuffer::RingBufferMap::new(path)?;
let ring = ringbuffer::ThreadSafeRingBuffer::new(sync::Mutex::new(map));
if let Ok(producer) = ringbuffer::producer::Producer::attach(&ring) {
let mut c = Client::new("localhost"); c.connect_as_producer(name, "Rust Ring producer")?;
return Ok(RingBufferProducer {
ringmaster: c,
name: String::from(name),
ring: producer
});
} else {
return Err(String::from("Could not create a producer object for the ring"));
}
}
pub fn make_and_register(path: &str, name: &str) -> Result<(), String> {
let mut c = Client::new("localhost");
let ring_size : u32 = DEFAULT_RING_SIZE;
ringbuffer::RingBufferMap::create(path, ring_size)?;
c.register_ring(name)?;
Ok(())
}
pub fn attach(name : &str) -> Result<RingBufferProducer, String> {
let path = ring_path(name);
match fs::exists(&path) {
Ok(exists) => { if !exists {
return Err(format!("Ring buffer file for {} does not exist", name));
} else {
return Self::produce(name, &path);
}},
Err(_) => {
return Err(String::from("Could not check existence of ring file"));
}
};
}
pub fn create_and_attach(name: &str) -> Result<RingBufferProducer, String> {
let path = ring_path(name);
if let Ok(exists) = fs::exists(&path) {
if !exists {
Self::make_and_register(&path, name)?;
}
return Self::produce(name, &path);
} else {
return Err(String::from("Unable to check ringbuffer existence"));
}
}
}
pub fn kill_ring(name : &str) -> Result<(), String> {
let path = ring_path(name);
let mut c = Client::new("localhost");
c.unregister_ring(name)?;
ringbuffer::RingBufferMap::delete(&path)?;
Ok(())
}
pub struct RingBufferConsumer {
ringmaster : Client,
name : String,
pub consumer : ringbuffer::consumer::Consumer
}
impl Drop for RingBufferConsumer {
fn drop(&mut self) {
let idx = self.consumer.get_index();
let _ignore = self.ringmaster
.disconnect_consumer(&self.name, idx);
}
}
impl RingBufferConsumer {
fn is_local(host : &str) -> Result<bool, String> {
let host_and_port = String::from(host) + ":30000";
let ips = host_and_port.to_socket_addrs();
if let Err(reason) = ips {
return Err(format!("Could not resolve URL host to an ip {}", reason));
}
let ips = ips.unwrap();
if let Ok(my_addresses) = local_ip_address::list_afinet_netifas() {
let mut my_address_list = HashSet::<net::IpAddr>::new();
for (_, ip) in my_addresses.iter() {
my_address_list.insert(ip.clone());
}
for host_ip in ips {
if my_address_list.contains(&host_ip.ip()) {
return Ok(true);
}
}
Ok(false)
} else {
Err(String::from("Could not get our local addresses to check if the URI is local"))
}
}
fn no_such(local_ring: &str) -> bool {
let path =ring_path(local_ring);
if let Ok(e) = fs::exists(&path) {
if !e {
return true;
}
let mut c = Client::new("localhost");
let listing = c.list_rings();
if let Err(_e) = listing {
false } else {
let list = listing.unwrap();
let mut ring_list = HashSet::new();
for ring in list {
ring_list.insert(ring.name);
}
!ring_list.contains(local_ring)
}
} else { false
}
}
fn start_stdintoring(path : &str, ring: &str, stdin : &mut net::TcpStream) -> Result<String, String> {
let command = Command::new("sh")
.stdin(unsafe {Stdio::from_raw_fd(stdin.try_clone().unwrap().into_raw_fd())})
.arg("nohup")
.arg(path)
.arg(ring)
.arg(">")
.arg("/dev/null")
.arg("2>&1").spawn();
if let Err(reason) = command {
Err(format!("Could not spawn stdintoring: {}", reason))
} else {
Ok(String::from(ring)) }
}
fn start_hoister(host : &str, remote_ring: &str) -> Result<String, String> {
let stdintoring = stdin_to_ring_path();
if let Err(reason) = stdintoring {
return Err(format!("Could not find stdintoring: {}", reason));
}
let stdintoring = stdintoring.unwrap();
let mut c = Client::new(host); let hoist_socket = c.get_data(remote_ring);
if let Err(reason) = hoist_socket {
return Err(format!("Failed to set up remote host hoist: {}", reason));
}
let mut hoist_socket = hoist_socket.unwrap();
let ring_name = proxy_ring_name(host, remote_ring);
let ring_path = ring_path(&ring_name);
let ring_size = proxy_ring_size();
if let Err(reason) = ringbuffer::RingBufferMap::create(&ring_path, ring_size) {
return Err(format!("Could not create proxy ring: {}", reason));
}
let mut local_master = Client::new("localhost");
if let Err(reason) = local_master.register_ring(&ring_name) {
return Err(format!("Unable to register proxy ring: {}", reason));
}
Self::start_stdintoring(&stdintoring, &ring_name, &mut hoist_socket)
}
pub fn attach(uri : &str) -> Result<RingBufferConsumer, String> {
let parsed_uri = Url::parse(uri);
if let Err(e) = parsed_uri {
return Err(format!("Failed to parse the ring uri '{}' : {}", uri, e));
}
let parsed_uri = parsed_uri.unwrap();
let scheme = parsed_uri.scheme(); let host = parsed_uri.host_str();
let mut ring = String::from(parsed_uri.path());
let _ = ring.remove(0);
if scheme != "tcp" {
return Err(String::from("Invalid scheme/protocol for ringbuffer must be tcp:"));
}
if host.is_none() {
return Err(String::from("Ringbuffer URIs' must supply a host."));
}
let host = String::from(host.unwrap());
if ! Self::is_local(&host)? {
let hoist = Self::start_hoister(&host, &ring);
if let Err(s) = hoist {
return Err(format!("Failed to start hoister for {} : {}", uri, s));
}
ring = hoist.unwrap();
}
if Self::no_such(&ring) {
return Err(format!("There is no such local ring: {}", ring));
}
let path = ring_path(&ring); let map = ringbuffer::RingBufferMap::new(&path)?;
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(sync::Mutex::new(map));
let consumer = ringbuffer::consumer::Consumer::attach(&ringbuffer);
if let Err(e) = consumer {
return Err(format!("Failed to attach as consumer to {} : {:?}", uri, e));
}
let consumer = consumer.unwrap();
let mut c = Client::new("localhost");
c.connect_as_consumer(&ring, consumer.get_index(), "")?;
Ok(RingBufferConsumer {
ringmaster: c,
name : ring.clone(),
consumer,
})
}
}
#[cfg(test)]
mod list_parse_tests {
use super::*;
#[test]
fn simple_list() {
let list = "a b cd";
let parse = parse_tcl_list(list);
assert_eq!(3, parse.len());
assert_eq!(vec!["a", "b", "cd"], parse);
}
#[test]
fn empty_list() {
let list="";
let parse = parse_tcl_list(list);
assert_eq!(0, parse.len());
}
#[test]
fn nested_list() {
let list="a b {cd ef}";
let parse = parse_tcl_list(list);
assert_eq!(3, parse.len());
assert_eq!(vec!["a", "b", "cd ef"], parse);
let sublist = parse_tcl_list(&parse[2]);
assert_eq!(2, sublist.len());
assert_eq!(vec!["cd", "ef"], sublist);
}
#[test]
fn nested_nested() {
let list = "a b {c d {e f}}";
let outer = parse_tcl_list(list);
assert_eq!(3, outer.len());
assert_eq!(vec!["a", "b", "c d {e f}"], outer);
let inner=parse_tcl_list(&outer[2]);
assert_eq!(3, inner.len());
assert_eq!(vec!["c", "d", "e f"], inner);
let innermost = parse_tcl_list(&inner[2]);
assert_eq!(2, innermost.len());
assert_eq!(vec!["e", "f"], innermost);
}
}
#[cfg(test)]
mod client_tests {
use super::*;
use nscldaq_ringbuffer::ringbuffer::RingBufferMap;
use std::collections::HashSet;
use std::sync::Mutex;
use std::time::Duration;
fn ring_name(base_name : &str) -> String {
format!("{}/{}", RINGBUFFER_DIRECTORY, base_name)
}
fn create_and_register(path : &str, name: &str) {
let mut c = Client::new("localhost");
let ring_size : u32 = 1024*1024;
RingBufferMap::create(path, ring_size).expect("Failed to create ring");
c.register_ring(name).expect("Failed to register the ring");
}
fn kill_ring_pre_unregister(path : &str, name: &str) {
let mut c = Client::new("localhost");
let unreg_req = format!("UNREGISTER {} \n", name);
c.transact_and_analyze(&unreg_req).expect("Unable to unregister a ring manually");
RingBufferMap::delete(path).expect("Unable to delete ring");
}
fn kill_ring(path : &str, name: &str) {
let mut c = Client::new("localhost");
c.unregister_ring(name).expect("Unable to unregister ring");
RingBufferMap::delete(path).expect("Unable to delete ring buffer file");
}
#[test]
fn new() {
let c = Client::new("localhost");
assert_eq!("localhost".to_string(), c.host);
assert!(c.socket.is_none());
}
#[test]
fn host() {
let c = Client::new("localhost");
assert_eq!("localhost".to_string(), c.host());
}
#[test]
fn socket_none() {
let c = Client::new("localhost");
assert!(c.socket().is_none());
}
#[test]
fn get_port() {
let c = Client::new("localhost");
let port = c.get_port();
assert!(port.is_ok());
}
#[test]
fn register_1() {
let mut c = Client::new("localhost");
let result = c.register_ring("no-such-ring");
assert!(result.is_err());
}
#[test]
fn register_2() {
let name = "register2_ring";
let path = ring_name(name);
RingBufferMap::create(&path, 1024*1024).expect("register_1 could not make ring");
let mut c = Client::new("localhost");
let result = c.register_ring(name);
assert!(result.is_ok());
kill_ring_pre_unregister(&path, name);
}
#[test]
fn register_3() {
let name = "register3_ring";
let path = ring_name(name);
RingBufferMap::create(&path, 1024*1024).expect("register_1 could not make ring");
let mut c = Client::new("localhost");
c.register_ring(name).expect("Failed initial ring registration.");
assert!(c.register_ring(name).is_ok());
kill_ring_pre_unregister(&path, name);
}
#[test]
fn list_1() {
let mut c = Client::new("localhost");
let info = c.list_rings();
assert!(info.is_ok());
let listing = info.unwrap();
assert_eq!(0, listing.len());
}
#[test]
fn list_2() {
let ring_size : u32 = 1024*1024;
let name = "list_2";
let ring_path = ring_name(name);
println!("Formatted : '{}'", ring_path);
RingBufferMap::create(&ring_path, ring_size).expect("list_2 failed to make ringbuffer");
let mut c = Client::new("localhost");
c.register_ring(name).expect("list_2 failed to register ring");
let list_result = c.list_rings();
assert!(list_result.is_ok());
let listing = list_result.unwrap();
assert_eq!(1, listing.len());
let ring_info = &listing[0];
assert_eq!(String::from(name), ring_info.name);
assert_eq!(ring_size, ring_info.size);
assert_eq!(ring_size, ring_info.free); assert_eq!(100, ring_info.maxconsumers);
assert_eq!(-1, ring_info.producer_pid);
assert_eq!(0, ring_info.max_get);
assert_eq!(0, ring_info.min_get);
assert_eq!(0, ring_info.consumers.len());
kill_ring_pre_unregister(&ring_path, name);
}
#[test]
fn list_3() {
let ring1_name ="list3_1";
let ring2_name = "list3_2";
let ring1 = ring_name(ring1_name);
let ring2 = ring_name(ring2_name);
create_and_register(&ring1, &ring1_name);
create_and_register(&ring2, &ring2_name);
let mut c = Client::new("localhost");
let listing_reply = c.list_rings();
assert!(listing_reply.is_ok());
let listing = listing_reply.unwrap();
assert_eq!(2, listing.len());
let mut set = HashSet::new(); set.insert(listing[0].name.clone());
set.insert(listing[1].name.clone());
assert!(set.contains(&String::from(ring1_name)));
assert!(set.contains(&String::from(ring2_name)));
kill_ring_pre_unregister(&ring1, ring1_name);
kill_ring_pre_unregister(&ring2, ring2_name);
}
#[test]
fn unregister_1() {
let mut c = Client::new("localhost");
let response = c.unregister_ring("does_not_exist");
assert!(response.is_ok());
}
#[test]
fn unregister_2() {
let name ="unregister_1";
let path = ring_name(name);
create_and_register(&path, name);
let mut c = Client::new("localhost");
let response = c.unregister_ring(name);
assert!(response.is_ok());
assert_eq!(0, c.list_rings().expect("Unable to list rings").len());
RingBufferMap::delete(&path).expect("Unable to delete ring buffer file");
}
#[test]
fn pconnect_1() {
let mut c = Client::new("localhost");
let response = c.connect_as_producer("aring", "Why not");
assert!(response.is_err());
}
#[test]
fn pconnect_2() {
let name = "pconnect_2";
let path = ring_name(name);
create_and_register(&path,name);
let mut ring = RingBufferMap::new(&path)
.expect("Could not map the ring we made");
let pid = process::id();
let setp = ring.set_producer(pid);
assert!(setp.is_ok());
{
let mut c = Client::new("localhost");
let connect_response = c.connect_as_producer(&name, "It-is-me");
assert!(connect_response.is_ok());
let list = c.list_rings().expect("Could not list the rings");
assert_eq!(1, list.len());
assert_eq!(pid as i32, list[0].producer_pid);
}
let mut c = Client::new("localhost");
let list = c.list_rings().expect("Unable to get ring list");
assert_eq!(1, list.len());
assert_eq!(-1, list[0].producer_pid);
assert_eq!(ringbuffer::UNUSED_ENTRY, ring.producer().get_pid());
kill_ring(&path, name);
}
#[test]
fn pconnect_3() {
let name = "pconnect_3";
let path = ring_name(name);
create_and_register(&path,name);
let mut ring = RingBufferMap::new(&path)
.expect("Could not map the ring we made");
let pid = process::id();
let setp = ring.set_producer(pid);
assert!(setp.is_ok());
{
let mut c = Client::new("localhost");
let connect_response = c.connect_as_producer(&name, "It-is-me");
assert!(connect_response.is_ok());
let request = format!("CONNECT {{{}}} producer {} \"Testing\" \n", name, pid+1);
let response = c.transact_and_analyze(&request);
assert!(response.is_err());
}
let mut c = Client::new("localhost");
let list = c.list_rings().expect("Unable to get ring list");
assert_eq!(1, list.len());
assert_eq!(-1, list[0].producer_pid);
assert_eq!(ringbuffer::UNUSED_ENTRY, ring.producer().get_pid());
kill_ring(&path, name);
}
#[test]
fn pdisconnect_1() {
let name = "pdisconnect_1";
let mut c = Client::new("localhost");
assert!(c.disconnect_producer(name).is_err());
}
#[test]
fn pdisconnect_2() {
let name = "pdisconnect_2";
let path = ring_name(name);
create_and_register(&path, name);
let mut c = Client::new("localhost");
assert!(c.disconnect_producer(name).is_err());
kill_ring(&path, name);
}
#[test]
fn pdisconnect_3() {
let name = "pdisconnect_3";
let path = ring_name(name);
create_and_register(&path, name);
let mut map =
RingBufferMap::new(&path).expect("failed to map ring");
map.set_producer(process::id()).expect("failed to set ring producer pid");
let mut c = Client::new("localhost");
c.connect_as_producer(name, "A comment").expect("Failed to connect as producer");
let status = c.disconnect_producer(name);
assert!(status.is_ok());
assert!(map.free_producer(process::id()).is_ok());
let list = c.list_rings().expect("Failed to list ringgs");
assert_eq!(1, list.len());
assert_eq!(-1, list[0].producer_pid);
kill_ring(&path, name);
}
#[test]
fn cconnect_1() {
let mut c = Client::new("localhost");
assert!(c.connect_as_consumer("aring", 1, "Comment").is_err());
}
#[test]
fn cconnect_2() {
let ring = "cconnect_2";
let path = ring_name(ring);
create_and_register(&path, ring);
let map = RingBufferMap::new(&path).expect("failed to map ring");
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(Mutex::new(map));
let consumer_stat = ringbuffer::consumer::Consumer::attach(&ringbuffer);
assert!(consumer_stat.is_ok());
let consumer = consumer_stat.unwrap();
let idx = consumer.get_index();
let mut c = Client::new("localhost");
let status = c.connect_as_consumer(ring, idx, "WAWA");
assert!(status.is_ok());
let list = c.list_rings().expect("Failed to list rings");
assert_eq!(1, list.len());
let consumers = &(list[0].consumers);
assert_eq!(1, consumers.len()); assert_eq!(process::id(), consumers[0].pid as u32);
drop(consumer);
kill_ring(&path, ring);
}
#[test]
fn cconnect_3() {
let ring = "cconnect_3";
let path = ring_name(ring);
create_and_register(&path, ring);
let map = RingBufferMap::new(&path).expect("failed to map ring");
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(Mutex::new(map));
let consumer = ringbuffer::consumer::Consumer::attach(&ringbuffer)
.expect("Failed to get a consumer for the ring");
{
let mut c = Client::new("localhost");
c
.connect_as_consumer(ring, consumer.get_index(), "Some comment")
.expect("Failed to register consumer");
}
let mut c = Client::new("localhost");
let listing = c.list_rings().expect("Could not list rings");
assert_eq!(1, listing.len());
let clients = &listing[0].consumers;
assert_eq!(0, clients.len());
drop(consumer); kill_ring(&path, ring);
}
#[test]
fn cdisconnect_1() {
let mut c = Client::new("localhost");
assert!(c.disconnect_consumer("Nosuch", 1).is_err());
}
#[test]
fn cdisconnect_2() {
let ring = "cdisconnect_2";
let path = ring_name(ring);
create_and_register(&path, ring);
let mut c = Client::new("localhost");
assert!(c.disconnect_consumer(ring, 1).is_err());
kill_ring(&path, ring);
}
#[test]
fn cdisconnect_3() {
let ring = "cdisconnect_3";
let path = ring_name(ring);
create_and_register(&path, ring);
let map = RingBufferMap::new(&path).expect("failed to map ring");
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(Mutex::new(map));
let consumer = ringbuffer::consumer::Consumer::attach(&ringbuffer)
.expect("Failed to get a consumer for the ring");
let mut c = Client::new("localhost");
c.connect_as_consumer(ring, consumer.get_index(), "Junk")
.expect("Failed to register as conumser");
let result = c.disconnect_consumer(ring, consumer.get_index());
assert!(result.is_ok());
let listing = c.list_rings().expect("could not list rings");
assert_eq!(1, listing.len());
let clist = &listing[0].consumers;
assert_eq!(0, clist.len());
drop(consumer);
kill_ring(&path, ring);
}
#[test]
fn cdisconnect_4() {
let ring = "cdisconnect_4";
let path = ring_name(ring);
create_and_register(&path, ring);
let map = RingBufferMap::new(&path).expect("failed to map ring");
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(Mutex::new(map));
let consumer1 = ringbuffer::consumer::Consumer::attach(&ringbuffer)
.expect("Failed to make consumer1 for the ring");
let consumer2 = ringbuffer::consumer::Consumer::attach(&ringbuffer)
.expect("Failed to make consumer 2 for the ring");
let mut c = Client::new("localhost");
c.connect_as_consumer(ring, consumer1.get_index(), "Consumer1")
.expect("Failed to connect consumer 1 in ringmaster");
c.connect_as_consumer(ring, consumer2.get_index(), "Consumer 2")
.expect("Failed to register consumer2");
c.disconnect_consumer(ring, consumer1.get_index()).expect("Could not unregister consumer1");
drop(consumer1);
let list = c.list_rings().expect("Failed to list rings");
assert_eq!(1, list.len());
let consumers = &list[0].consumers;
assert_eq!(1, consumers.len());
assert_eq!(process::id(), consumers[0].pid);
c.disconnect_consumer(ring, consumer2.get_index()).expect("Failed to unregister consumer 2");
drop(consumer2);
let list = c.list_rings().expect("Failed to list rings");
assert_eq!(1, list.len());
let consumers = &list[0].consumers;
assert_eq!(0, consumers.len());
kill_ring(&path, ring);
}
#[test]
fn ring_path_1() {
let ring = "ring_path_1";
assert_eq!(ring_name(ring), ring_path(ring));
}
#[test]
fn remote_1() {
let mut c = Client::new("localhost");
assert!(c.get_data("nosuch").is_err());
}
#[test]
fn remote_2() {
let ring = "remote_2";
let path = ring_path(ring);
create_and_register(&path, ring);
let mut c = Client::new("localhost");
let status = c.get_data(ring);
assert!(status.is_ok());
if let Ok(socket) = status {
let _ = socket.shutdown(net::Shutdown::Both); drop(socket); }
kill_ring(&path, ring);
}
#[test]
fn remote_3() {
let ring = "remote_3";
let path = ring_path(ring);
create_and_register(&path, ring);
let map = RingBufferMap::new(&path).expect("failed to map ring");
let ringbuffer = ringbuffer::ThreadSafeRingBuffer::new(Mutex::new(map));
let mut producer = ringbuffer::producer::Producer::attach(&ringbuffer)
.expect("Could not attach as producer");
let mut c = Client::new("localhost");
c.connect_as_producer(ring, "Producer")
.expect("Could not connect as a producer");
let mut c_data = Client::new("localhost");
let mut sock = c_data.get_data(ring).expect("Filed to set up socket");
let data : [u8; 20] = [
'a' as u8 ; 20
];
producer.write_all(&data).expect("Could not put test data in ring");
sock.set_read_timeout(Some(Duration::from_secs(5))).expect("could not set read timout for socket");
let mut received_data : [u8;18] = [0;18];
let result = sock.read(&mut received_data);
assert!(result.is_ok());
assert_eq!(18, result.unwrap());
for (i,b) in received_data.into_iter().enumerate() {
assert_eq!(data[i], b);
}
let _ = sock.shutdown(net::Shutdown::Both);
c.disconnect_producer(ring).expect("Failed to disconnect producer");
drop(producer);
kill_ring(&path, ring);
}
}
#[cfg(test)]
mod producer_tests {
use super::*;
use std::collections::HashSet;
#[test]
fn attach_1() {
assert!(RingBufferProducer::attach("nosuch").is_err());
}
#[test]
fn attach_2() {
let name ="attach_2";
let path = ring_path(name);
RingBufferProducer::make_and_register(&path, name).expect("Failed to make/register the ring");
let producer_stat = RingBufferProducer::attach(name);
assert!(producer_stat.is_ok());
let mut c = Client::new("localhost");
let l = c.list_rings().expect("failed to list rings");
let mut names =HashSet::new();
for ring in l {
names.insert(ring.name);
}
assert!(names.contains(name));
let mut producer = producer_stat.unwrap();
let data : [u8; 10] = ['b' as u8 ;10];
assert!(producer.ring.blocking_put(&data).is_ok());
kill_ring(name).expect("could not kill off ringbuffer");
}
}
#[cfg(test)]
mod consumer_tests {
use super::*;
use nscldaq_ringbuffer::ringbuffer::RingBufferMap;
fn create_and_register(path : &str, name: &str) {
let mut c = Client::new("localhost");
let ring_size : u32 = 1024*1024;
RingBufferMap::create(path, ring_size).expect("Failed to create ring");
c.register_ring(name).expect("Failed to register the ring");
}
fn kill_ring(path : &str, name: &str) {
let mut c = Client::new("localhost");
c.unregister_ring(name).expect("Unable to unregister ring");
RingBufferMap::delete(path).expect("Unable to delete ring buffer file");
}
#[test]
fn islocal_1() {
assert!(RingBufferConsumer::is_local("localhost").unwrap());
assert!(RingBufferConsumer::is_local("127.0.0.1").unwrap());
assert!(RingBufferConsumer::is_local("::1").unwrap());
}
#[test]
fn islocal_2() {
assert!(!RingBufferConsumer::is_local("www.google.com").unwrap()); }
#[test]
fn nosuch_1() {
assert!(RingBufferConsumer::no_such("there_is_no_such_ring"));
}
#[test]
fn nosuch_2() {
let ring_name = "nosuch_2";
let ring_path = ring_path(ring_name);
create_and_register(&ring_path, ring_name);
assert!(!RingBufferConsumer::no_such(ring_name));
kill_ring(&ring_path, ring_name);
}
#[test]
fn attach_1() {
let status = RingBufferConsumer::attach("tcp://localhost/no_such_ring");
assert!(status.is_err());
}
#[test]
fn attach_2() {
let ring = "attach_2";
let ring_path = ring_path(ring);
create_and_register(&ring_path, ring);
let ring_url = format!("tcp://localhost/{}", ring);
let status = RingBufferConsumer::attach(&ring_url);
assert!(status.is_ok());
drop(status);
kill_ring(&ring_path, ring);
}
}