extern crate rand;
extern crate xdr_codec;
extern crate zmq;
use std::cmp;
use std::mem;
use std::process;
use self::xdr_codec::{Pack, Unpack};
use self::zmq::Message;
use std::io::{Cursor, Read};
use rand::{thread_rng, Rng};
use crate::quality::{qserverclient::*, quality::*};
const DEBUG: u64 = 0;
pub struct QualityClient {
_server: String,
_port: u64,
socket: zmq::Socket,
_filename: String,
}
impl QualityClient {
pub fn new(server: String, port: u64, filename: String) -> QualityClient {
let context = zmq::Context::new();
let socket = context.socket(zmq::REQ).unwrap();
let mut connection_str = String::from("tcp://");
connection_str.push_str(&server);
connection_str.push(':');
connection_str.push_str(&port.to_string());
log::info!(" connecting to {} ... ", connection_str);
let resconn = socket.connect(&connection_str);
assert!(resconn.is_ok());
match resconn {
Ok(_) => {
println!(" connect Ok");
}
Err(e) => {
println!(" got an error in connect : {}", e);
process::exit(1);
}
}
QualityClient {
_server: server,
_port: port,
socket: socket,
_filename: filename,
}
}
fn get_request_handle(&self) -> Qhandle {
thread_rng().gen::<u64>() as Qhandle
}
pub fn get_quality_sequence(&self, numseq: u64) -> Result<QSequenceRaw, StatusCode> {
log::debug!("get_quality requesting qsequence {}", numseq);
let len = mem::size_of::<RequestId>() + 8 + 4;
let vec: Vec<u8> = Vec::with_capacity(len);
let mut buffer = Cursor::new(vec);
let handle = self.get_request_handle();
let req_code = RequestCode::GetQRead;
log::debug!("get_quality request handle {}", handle);
if handle.pack(&mut buffer).is_err() {
return Err(StatusCode::ErrXdr);
}
if (req_code as u64).pack(&mut buffer).is_err() {
return Err(StatusCode::ErrXdr);
}
if numseq.pack(&mut buffer).is_err() {
return Err(StatusCode::ErrXdr);
}
let rawbuf = buffer.get_ref().as_slice();
let resmsg = Message::try_from(rawbuf);
if resmsg.is_err() {
println!(" construction of msg from slice failed ... ");
}
let msg = resmsg.unwrap();
log::debug!(" sending message ");
let _res = self.socket.send(msg, 0).unwrap();
if DEBUG > 0 {
println!(" waiting for server ... ");
}
if let Ok(v) = self.socket.recv_bytes(0) {
log::info!(" client receiving a response .. nb_bytes : {} ", v.len());
let mut qualv: Vec<u8>;
let mut cursor = Cursor::new(v);
let (handle_got, _) = u64::unpack(&mut cursor).unwrap();
assert_eq!(handle_got, handle);
let (ret_code, _) = u64::unpack(&mut cursor).unwrap();
if ret_code != StatusCode::Ok as u64 {
log::info!("get_quality got an error status from server");
return Err(StatusCode::ErrGen);
}
if let Ok((nb_qual, _)) = u64::unpack(&mut cursor) {
log::info!("got a qseq size: {} ", nb_qual);
qualv = Vec::with_capacity(nb_qual as usize);
unsafe { qualv.set_len(nb_qual as usize) };
cursor.read_exact(qualv.as_mut_slice()).unwrap();
if DEBUG > 0 {
for i in 0..cmp::min(10, qualv.len()) {
println!("client got qual {} {}", i, qualv[i]);
}
}
log::info!(
"returning a QSequenceRaw with size : {:?} ",
cursor.position()
);
return Ok(QSequenceRaw {
read_num: numseq as usize,
qseq: qualv,
});
} else {
return Err(StatusCode::ErrGen);
}
} else {
return Err(StatusCode::ErrGen);
}
} }