#[cfg(feature = "std")]
use std::sync::atomic::{AtomicI32};
#[cfg(feature = "std")]
use std::{sync::{Arc, mpsc}, time::{Duration, Instant}};
use super::*;
#[cfg(not(feature = "std"))]
macro_rules! println {
() => {
};
($($arg:tt)*) => {{
}};
}
#[test]
fn simple_test()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
let buf0_res = bufs.allocate();
assert_eq!(buf0_res.is_ok(), true, "{:?}", buf0_res.err().unwrap());
let buf0 = buf0_res.unwrap();
let buf0_w = buf0.write();
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0_w);
let buf0_r = buf0.read();
assert_eq!(buf0_r.is_ok(), true, "{:?}", buf0_r.err().unwrap());
assert_eq!(buf0.write(), Err(RwBufferError::WriteTryAgianLater));
let buf0_1 = buf0.clone();
assert_eq!(buf0_1.write(), Err(RwBufferError::WriteTryAgianLater));
let flags0 = buf0.get_flags();
let flags0_1 = buf0_1.get_flags();
assert_eq!(flags0, flags0_1);
assert_eq!(flags0.base, 3);
assert_eq!(flags0.read, 1);
assert_eq!(flags0.write, false);
}
#[test]
fn simple_test_dopped_in_place()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
let buf0_res = bufs.allocate();
assert_eq!(buf0_res.is_ok(), true, "{:?}", buf0_res.err().unwrap());
let buf0 = buf0_res.unwrap();
println!("{:?}", buf0.get_flags());
let buf0_w = buf0.write();
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0);
let buf0_flags = bufs.get_flags_by_index(0);
assert_eq!(buf0_flags.is_some(), true, "no flags");
let buf0_flags = buf0_flags.unwrap();
println!("{:?}", buf0_flags);
assert_eq!(buf0_flags.base, 1);
assert_eq!(buf0_flags.read, 0);
assert_eq!(buf0_flags.write, true);
drop(buf0_w.unwrap());
let buf0_flags = bufs.get_flags_by_index(0);
assert_eq!(buf0_flags.is_some(), true, "no flags");
let buf0_flags = buf0_flags.unwrap();
println!("{:?}", buf0_flags);
assert_eq!(buf0_flags.base, 1);
assert_eq!(buf0_flags.read, 0);
assert_eq!(buf0_flags.write, false);
}
#[test]
fn simple_test_dropped_in_place_downgrade()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
let buf0_res = bufs.allocate();
assert_eq!(buf0_res.is_ok(), true, "{:?}", buf0_res.err().unwrap());
let buf0 = buf0_res.unwrap();
println!("{:?}", buf0.get_flags());
let buf0_w = buf0.write();
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0);
let buf0_rd = buf0_w.unwrap().downgrade();
assert_eq!(buf0_rd.is_ok(), true, "{:?}", buf0_rd.err().unwrap());
let buf0_flags = bufs.get_flags_by_index(0);
assert_eq!(buf0_flags.is_some(), true, "no flags");
let buf0_flags = buf0_flags.unwrap();
println!("{:?}", buf0_flags);
assert_eq!(buf0_flags.base, 1);
assert_eq!(buf0_flags.read, 1);
assert_eq!(buf0_flags.write, false);
}
#[test]
fn simple_test_drop_in_place_downgrade()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
let buf0_w =
{
let buf0 = bufs.allocate_in_place();
println!("1: {:?}", buf0.get_flags());
let buf0_w = buf0.write();
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0);
buf0_w
};
let buf0_rd = buf0_w.unwrap().downgrade();
assert_eq!(buf0_rd.is_ok(), true, "{:?}", buf0_rd.err().unwrap());
let buf0_flags = bufs.get_flags_by_index(0);
assert_eq!(buf0_flags.is_some(), false, "flags");
let buf0_rd = buf0_rd.unwrap();
let buf0_flags = buf0_rd.get_flags();
println!("2: {:?}", buf0_flags);
assert_eq!(buf0_flags.base, 0);
assert_eq!(buf0_flags.read, 1);
assert_eq!(buf0_flags.write, false);
}
#[test]
fn timing_test()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
for _ in 0..10
{
#[cfg(feature = "std")]
let inst = Instant::now();
let buf0_res = bufs.allocate_in_place();
#[cfg(feature = "std")]
let end = inst.elapsed();
#[cfg(feature = "std")]
println!("alloc: {:?}", end);
drop(buf0_res);
}
let buf0_res = bufs.allocate();
assert_eq!(buf0_res.is_ok(), true, "{:?}", buf0_res.err().unwrap());
let buf0 = buf0_res.unwrap();
for _ in 0..10
{
#[cfg(feature = "std")]
let inst = Instant::now();
let buf0_w = buf0.write();
#[cfg(feature = "std")]
let end = inst.elapsed();
#[cfg(feature = "std")]
println!("write: {:?}", end);
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0_w);
}
for _ in 0..10
{
#[cfg(feature = "std")]
let inst = Instant::now();
let buf0_r = buf0.read();
#[cfg(feature = "std")]
let end = inst.elapsed();
#[cfg(feature = "std")]
println!("read: {:?}", end);
assert_eq!(buf0_r.is_ok(), true, "{:?}", buf0_r.err().unwrap());
assert_eq!(buf0.write(), Err(RwBufferError::WriteTryAgianLater));
drop(buf0_r);
}
}
#[cfg(feature = "std")]
#[test]
fn simple_test_mth()
{
let mut bufs = RwBuffers::new(4096, 1, 3).unwrap();
let buf0 = bufs.allocate().unwrap();
let buf0_rd = buf0.write().unwrap().downgrade().unwrap();
let join1=
std::thread::spawn(move ||
{
println!("{:?}", buf0_rd);
std::thread::sleep(Duration::from_secs(2));
return;
}
);
let buf1_rd = buf0.read().unwrap();
let join2=
std::thread::spawn(move ||
{
println!("{:?}", buf1_rd);
std::thread::sleep(Duration::from_secs(2));
return;
}
);
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 2);
assert_eq!(flags.write, false);
let _ = join1.join();
let _ = join2.join();
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 0);
assert_eq!(flags.write, false);
}
#[cfg(feature = "std")]
#[test]
fn high_load_concurrent()
{
let w_flag = Arc::new(AtomicI32::new(0));
let (s1, r1) = mpsc::channel::<()>();
let mut bufs = RwBuffers::new(4096, 1, 3).unwrap();
let buf0 = bufs.allocate().unwrap();
let cbuf0 = buf0.clone();
let cw_flag = w_flag.clone();
let join1=
std::thread::spawn(move ||
{
std::thread::park();
for i in 0..100
{
let w = cbuf0.write();
if let Ok(ww) = w
{
cw_flag.store(i, Ordering::Relaxed);
let flags = cbuf0.get_flags();
assert_eq!(flags.read, 0);
assert_eq!(flags.write, true);
let Ok(_) = r1.recv() else { return };
drop(ww);
}
std::thread::sleep(Duration::from_nanos(20));
}
return;
}
);
let mut prev_cwf = -1;
join1.thread().unpark();
let mut cntr = 0;
loop
{
if join1.is_finished() == true
{
break;
}
let re = buf0.read();
if let Err(_e) = re
{
let cwf = w_flag.load(Ordering::Relaxed);
if cwf != prev_cwf
{
prev_cwf = cwf;
let flags = buf0.get_flags();
println!("w->{} cwf->{} {:?}", cntr, cwf, flags);
assert_eq!(flags.read, 0);
assert_eq!(flags.write, true);
s1.send(()).unwrap();
}
}
else
{
let flags = buf0.get_flags();
assert_eq!(flags.read, 1);
assert_eq!(flags.write, false);
cntr += 1;
}
}
}
#[cfg(feature = "std")]
#[test]
fn simple_test_concurent()
{
let mut bufs = RwBuffers::new(4096, 1, 3).unwrap();
let buf0 = bufs.allocate().unwrap();
let buf0_w = buf0.write().unwrap();
let join1=
std::thread::spawn(move ||
{
std::thread::park();
println!("{:?}", buf0_w);
std::thread::sleep(Duration::from_secs(3));
return;
}
);
join1.thread().unpark();
let s = Instant::now();
let buf1_rd =
loop
{
match buf0.read()
{
Ok(r) => break r,
Err(e) =>
{
assert_eq!(e, RwBufferError::ReadTryAgianLater);
continue;
}
}
};
let e = s.elapsed();
println!("read await {:?} {}", e, e.as_millis());
assert!(2990 <= e.as_millis() && e.as_millis() <= 3010);
let _ = join1.join();
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 1);
assert_eq!(flags.write, false);
drop(buf1_rd);
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 0);
assert_eq!(flags.write, false);
}
#[test]
fn test_try_into_read()
{
let mut bufs = RwBuffers::new(4096, 1, 2).unwrap();
let buf0 = bufs.allocate_in_place();
println!("{:?}", buf0.get_flags());
let buf0_w = buf0.write();
assert_eq!(buf0_w.is_ok(), true, "{:?}", buf0_w.err().unwrap());
assert_eq!(buf0.read(), Err(RwBufferError::ReadTryAgianLater));
drop(buf0);
let buf0_rd = buf0_w.unwrap().downgrade();
assert_eq!(buf0_rd.is_ok(), true, "{:?}", buf0_rd.err().unwrap());
let buf0_flags = bufs.get_flags_by_index(0);
assert_eq!(buf0_flags.is_some(), false, "flags");
let buf0_rd = buf0_rd.unwrap();
let buf0_flags = buf0_rd.get_flags();
println!("{:?}", buf0_flags);
assert_eq!(buf0_flags.base, 0);
assert_eq!(buf0_flags.read, 1);
assert_eq!(buf0_flags.write, false);
#[cfg(feature = "std")]
let inst = Instant::now();
let ve = buf0_rd.try_inner();
#[cfg(feature = "std")]
let end = inst.elapsed();
println!("try inner: {:?}", end);
assert_eq!(ve.is_ok(), true);
}
#[cfg(all(feature = "enable_async", feature = "std"))]
#[test]
fn test_multithreading_async()
{
use futures::{executor::{self, ThreadPool}, task::SpawnExt};
let pool = ThreadPool::new().expect("Failed to build pool");
let mut bufs = RwBuffers::new(4096, 1, 3).unwrap();
let buf0 = bufs.allocate().unwrap();
let cbuf0 = buf0.clone();
let hndl =
pool
.spawn_with_handle(async move
{
let buf0_w = cbuf0.write_async().await.unwrap();
println!("{:?}", buf0_w);
std::thread::sleep(Duration::from_secs(3));
async_drop(buf0_w).await;
return;
}
).unwrap();
pool
.spawn_ok(async move
{
let s = std::time::Instant::now();
let buf1_rd = buf0.read_async().await.unwrap();
let e = s.elapsed();
println!("read await {:?} {}", e, e.as_millis());
assert!(2998 <= e.as_millis() && e.as_millis() <= 3005);
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 1);
assert_eq!(flags.write, false);
async_drop(buf1_rd).await;
let flags = buf0.get_flags();
assert_eq!(flags.base, 2);
assert_eq!(flags.read, 0);
assert_eq!(flags.write, false);
}
);
executor::block_on(hndl);
return;
}
#[cfg(feature = "enable_async")]
#[test]
fn test_async_clone()
{
use futures::executor::ThreadPool;
let pool = ThreadPool::new().expect("Failed to build pool");
pool
.spawn_ok(async move
{
let mut bufs = RwBuffers::new(4096, 1, 3).unwrap();
let buf0 = bufs.allocate().unwrap();
let buf0_r = buf0.read_async().await.unwrap();
let c_buf0_r = buf0_r.async_clone().await;
println!("{:?} {:?}", buf0_r, c_buf0_r);
assert_eq!(buf0_r, c_buf0_r);
}
);
}