use std::io::{self, BufReader};
use std::io::prelude::*;
use std::sync::mpsc::{channel, Receiver};
use std::{thread, result};
use std::{time, fmt};
use crate::errors::*; pub use regex::Regex;
#[derive(Debug)]
enum PipeError {
IO(io::Error),
}
#[derive(Debug)]
enum PipedChar {
Char(u8),
EOF,
}
pub enum ReadUntil {
String(String),
Regex(Regex),
EOF,
NBytes(usize),
Any(Vec<ReadUntil>),
}
impl fmt::Display for ReadUntil {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let printable = match self {
&ReadUntil::String(ref s) if s == "\n" => "\\n (newline)".to_string(),
&ReadUntil::String(ref s) if s == "\r" => "\\r (carriage return)".to_string(),
&ReadUntil::String(ref s) => format!("\"{}\"", s),
&ReadUntil::Regex(ref r) => format!("Regex: \"{}\"", r),
&ReadUntil::EOF => "EOF (End of File)".to_string(),
&ReadUntil::NBytes(n) => format!("reading {} bytes", n),
&ReadUntil::Any(ref v) => {
let mut res = Vec::new();
for r in v {
res.push(r.to_string());
}
res.join(", ")
}
};
write!(f, "{}", printable)
}
}
pub fn find(needle: &ReadUntil, buffer: &str, eof: bool) -> Option<(usize, usize)> {
match needle {
&ReadUntil::String(ref s) => buffer.find(s).and_then(|pos| Some((pos, pos + s.len()))),
&ReadUntil::Regex(ref pattern) => {
if let Some(mat) = pattern.find(buffer) {
Some((mat.start(), mat.end()))
} else {
None
}
}
&ReadUntil::EOF => if eof { Some((0, buffer.len())) } else { None },
&ReadUntil::NBytes(n) => {
if n <= buffer.len() {
Some((0, n))
} else if eof && buffer.len() > 0 {
Some((0, buffer.len()))
} else {
None
}
}
&ReadUntil::Any(ref any) => {
for read_until in any {
if let Some(pos_tuple) = find(&read_until, buffer, eof) {
return Some(pos_tuple);
}
}
None
}
}
}
pub struct NBReader {
reader: Receiver<result::Result<PipedChar, PipeError>>,
buffer: String,
eof: bool,
timeout: Option<time::Duration>,
}
impl NBReader {
pub fn new<R: Read + Send + 'static>(f: R, timeout: Option<u64>) -> NBReader {
let (tx, rx) = channel();
thread::spawn(move || {
let _ = || -> Result<()> {
let mut reader = BufReader::new(f);
let mut byte = [0u8];
loop {
match reader.read(&mut byte) {
Ok(0) => {
let _ = tx.send(Ok(PipedChar::EOF)).chain_err(|| "cannot send")?;
break;
}
Ok(_) => {
tx.send(Ok(PipedChar::Char(byte[0])))
.chain_err(|| "cannot send")?;
}
Err(error) => {
tx.send(Err(PipeError::IO(error)))
.chain_err(|| "cannot send")?;
}
}
}
Ok(())
}();
});
NBReader {
reader: rx,
buffer: String::with_capacity(1024),
eof: false,
timeout: timeout.and_then(|millis| Some(time::Duration::from_millis(millis))),
}
}
fn read_into_buffer(&mut self) -> Result<()> {
if self.eof {
return Ok(());
}
while let Ok(from_channel) = self.reader.try_recv() {
match from_channel {
Ok(PipedChar::Char(c)) => self.buffer.push(c as char),
Ok(PipedChar::EOF) => self.eof = true,
Err(PipeError::IO(ref err)) if err.kind() == io::ErrorKind::Other => {
self.eof = true
}
Err(_) => {}
}
}
Ok(())
}
pub fn read_until(&mut self, needle: &ReadUntil) -> Result<(String, String)> {
let start = time::Instant::now();
loop {
self.read_into_buffer()?;
if let Some(tuple_pos) = find(needle, &self.buffer, self.eof) {
let first = self.buffer.drain(..tuple_pos.0).collect();
let second = self.buffer.drain(..tuple_pos.1 - tuple_pos.0).collect();
return Ok((first, second));
}
if self.eof {
return Err(ErrorKind::EOF(needle.to_string(), self.buffer.clone(), None).into());
}
if let Some(timeout) = self.timeout {
if start.elapsed() > timeout {
return Err(ErrorKind::Timeout(needle.to_string(),
self.buffer.clone()
.replace("\n", "`\\n`\n")
.replace("\r", "`\\r`")
.replace('\u{1b}', "`^`"),
timeout)
.into());
}
}
thread::sleep(time::Duration::from_millis(100));
}
}
pub fn try_read(&mut self) -> Option<char> {
let _ = self.read_into_buffer();
if self.buffer.len() > 0 {
self.buffer.drain(..1).last()
} else {
None
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_expect_melon() {
let f = io::Cursor::new("a melon\r\n");
let mut r = NBReader::new(f, None);
assert_eq!(("a melon".to_string(), "\r\n".to_string()),
r.read_until(&ReadUntil::String("\r\n".to_string()))
.expect("cannot read line"));
match r.read_until(&ReadUntil::NBytes(10)) {
Ok(_) => assert!(false),
Err(Error(ErrorKind::EOF(_, _, _), _)) => {}
Err(Error(_, _)) => assert!(false),
}
}
#[test]
fn test_regex() {
let f = io::Cursor::new("2014-03-15");
let mut r = NBReader::new(f, None);
let re = Regex::new(r"^\d{4}-\d{2}-\d{2}$").unwrap();
r.read_until(&ReadUntil::Regex(re))
.expect("regex doesn't match");
}
#[test]
fn test_regex2() {
let f = io::Cursor::new("2014-03-15");
let mut r = NBReader::new(f, None);
let re = Regex::new(r"-\d{2}-").unwrap();
assert_eq!(("2014".to_string(), "-03-".to_string()),
r.read_until(&ReadUntil::Regex(re))
.expect("regex doesn't match"));
}
#[test]
fn test_nbytes() {
let f = io::Cursor::new("abcdef");
let mut r = NBReader::new(f, None);
assert_eq!(("".to_string(), "ab".to_string()),
r.read_until(&ReadUntil::NBytes(2)).expect("2 bytes"));
assert_eq!(("".to_string(), "cde".to_string()),
r.read_until(&ReadUntil::NBytes(3)).expect("3 bytes"));
assert_eq!(("".to_string(), "f".to_string()),
r.read_until(&ReadUntil::NBytes(4)).expect("4 bytes"));
}
#[test]
fn test_eof() {
let f = io::Cursor::new("lorem ipsum dolor sit amet");
let mut r = NBReader::new(f, None);
r.read_until(&ReadUntil::NBytes(2)).expect("2 bytes");
assert_eq!(("".to_string(), "rem ipsum dolor sit amet".to_string()),
r.read_until(&ReadUntil::EOF).expect("reading until EOF"));
}
#[test]
fn test_try_read() {
let f = io::Cursor::new("lorem");
let mut r = NBReader::new(f, None);
r.read_until(&ReadUntil::NBytes(4)).expect("4 bytes");
assert_eq!(Some('m'), r.try_read());
assert_eq!(None, r.try_read());
assert_eq!(None, r.try_read());
assert_eq!(None, r.try_read());
assert_eq!(None, r.try_read());
}
}