test metawrite

This commit is contained in:
weisd
2024-12-16 17:45:28 +08:00
committed by weisd
parent adf4cfd76a
commit 99c229524d
8 changed files with 210 additions and 121 deletions
+154 -75
View File
@@ -1,9 +1,16 @@
use crate::disk::MetaCacheEntry;
use crate::error::Error;
use crate::error::Result;
use rmp::decode::RmpRead;
use rmp::encode::RmpWrite;
use rmp::Marker;
use std::io::Read;
use std::io::Write;
use std::str::from_utf8;
use tokio::io::AsyncRead;
use tokio::io::AsyncReadExt;
use tokio::io::AsyncWrite;
use tokio::io::AsyncWriteExt;
// use std::sync::Arc;
// use tokio::sync::mpsc;
// use tokio::sync::mpsc::Sender;
@@ -16,50 +23,64 @@ pub struct MetacacheWriter<W> {
wr: W,
created: bool,
// err: Option<Error>,
buf: Vec<u8>,
}
impl<W: Write> MetacacheWriter<W> {
impl<W: AsyncWrite + Unpin> MetacacheWriter<W> {
pub fn new(wr: W) -> Self {
Self {
wr,
created: false,
// err: None,
buf: Vec::new(),
}
}
pub fn init(&mut self) -> Result<()> {
pub async fn flush(&mut self) -> Result<()> {
self.wr.write_all(&self.buf).await?;
self.buf.clear();
Ok(())
}
pub async fn init(&mut self) -> Result<()> {
if !self.created {
rmp::encode::write_u8(&mut self.wr, METACACHE_STREAM_VERSION)?;
rmp::encode::write_u8(&mut self.buf, METACACHE_STREAM_VERSION).map_err(|e| Error::msg(format!("{:?}", e)))?;
self.flush().await?;
self.created = true;
}
Ok(())
}
pub fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> {
pub async fn write(&mut self, objs: &[MetaCacheEntry]) -> Result<()> {
if objs.is_empty() {
return Ok(());
}
self.init()?;
self.init().await?;
for obj in objs.iter() {
if obj.name.is_empty() {
return Err(Error::msg("metacacheWriter: no name"));
}
self.write_obj(obj)?;
self.write_obj(obj).await?;
}
Ok(())
}
pub fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> {
self.init()?;
rmp::encode::write_bool(&mut self.wr, true)?;
pub async fn write_obj(&mut self, obj: &MetaCacheEntry) -> Result<()> {
self.init().await?;
rmp::encode::write_str(&mut self.wr, &obj.name)?;
rmp::encode::write_bool(&mut self.buf, true).map_err(|e| Error::msg(format!("{:?}", e)))?;
rmp::encode::write_str(&mut self.buf, &obj.name).map_err(|e| Error::msg(format!("{:?}", e)))?;
rmp::encode::write_bin(&mut self.buf, &obj.metadata).map_err(|e| Error::msg(format!("{:?}", e)))?;
self.flush().await?;
rmp::encode::write_bin(&mut self.wr, &obj.metadata)?;
Ok(())
}
@@ -96,9 +117,9 @@ impl<W: Write> MetacacheWriter<W> {
// Ok(sender)
// }
pub fn close(&mut self) -> Result<()> {
rmp::encode::write_bool(&mut self.wr, false)?;
self.wr.flush()?;
pub async fn close(&mut self) -> Result<()> {
rmp::encode::write_bool(&mut self.buf, false).map_err(|e| Error::msg(format!("{:?}", e)))?;
self.flush().await?;
Ok(())
}
}
@@ -108,34 +129,60 @@ pub struct MetacacheReader<R> {
init: bool,
err: Option<Error>,
buf: Vec<u8>,
offset: usize,
}
impl<R: Read> MetacacheReader<R> {
impl<R: AsyncRead + Unpin> MetacacheReader<R> {
pub fn new(rd: R) -> Self {
Self {
rd,
init: false,
err: None,
buf: Vec::new(),
offset: 0,
}
}
fn check_init(&mut self) {
pub async fn read_more(&mut self, read_size: usize) -> Result<&[u8]> {
let ext_size = read_size + self.offset;
let extra = ext_size - self.offset;
if self.buf.capacity() >= ext_size {
// Extend the buffer if we have enough space.
self.buf.resize(ext_size, 0);
} else {
self.buf.extend(vec![0u8; extra]);
}
let pref = self.offset;
self.rd.read_exact(&mut self.buf[pref..ext_size]).await?;
self.offset += read_size;
let data = &self.buf[pref..ext_size];
println!("pref {} offset {},ext_size {}, data {:?}", pref, self.offset, ext_size, &data);
Ok(data)
}
fn reset(&mut self) {
self.buf.clear();
self.offset = 0;
}
async fn check_init(&mut self) -> Result<()> {
if !self.init {
// let mut buf = match self.read_buf(1).await {
// Ok(res) => res,
// Err(err) => {
// self.err = Some(Error::msg(err.to_string()));
// return;
// }
// };
let ver = match rmp::decode::read_u8(&mut self.rd) {
let ver = match rmp::decode::read_u8(&mut self.read_more(2).await?) {
Ok(res) => res,
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
self.err = Some(Error::msg(format!("{:?}", err)));
0
}
};
println!("ver {}", ver);
match ver {
1 | 2 => (),
_ => {
@@ -145,70 +192,99 @@ impl<R: Read> MetacacheReader<R> {
self.init = true;
}
Ok(())
}
pub fn peek(&mut self) -> Result<Option<MetaCacheEntry>> {
self.check_init();
async fn read_str_len(&mut self) -> Result<u32> {
let mark = match rmp::decode::read_marker(&mut self.read_more(1).await?) {
Ok(res) => res,
Err(err) => {
let serr = format!("{:?}", err);
self.err = Some(Error::msg(&serr));
return Err(Error::msg(&serr));
}
};
match mark {
Marker::FixStr(size) => Ok(u32::from(size)),
Marker::Str8 => Ok(u32::from(self.read_u8().await?)),
Marker::Str16 => Ok(u32::from(self.read_u16().await?)),
Marker::Str32 => Ok(self.read_u32().await?),
_ => Err(Error::msg("str marker err")),
}
}
async fn read_bin_len(&mut self) -> Result<u32> {
let mark = match rmp::decode::read_marker(&mut self.read_more(1).await?) {
Ok(res) => res,
Err(err) => {
let serr = format!("{:?}", err);
self.err = Some(Error::msg(&serr));
return Err(Error::msg(&serr));
}
};
match mark {
Marker::Bin8 => Ok(u32::from(self.read_u8().await?)),
Marker::Bin16 => Ok(u32::from(self.read_u16().await?)),
Marker::Bin32 => Ok(self.read_u32().await?),
_ => Err(Error::msg("bin marker err")),
}
}
async fn read_u8(&mut self) -> Result<u8> {
let a = self.read_more(1).await?;
Ok(a[0])
}
async fn read_u16(&mut self) -> Result<u16> {
rmp::decode::read_u16(&mut self.read_more(2).await?).map_err(|e| Error::msg(format!("{:?}", e)))
}
async fn read_u32(&mut self) -> Result<u32> {
rmp::decode::read_u32(&mut self.read_more(4).await?).map_err(|e| Error::msg(format!("{:?}", e)))
}
pub async fn peek(&mut self) -> Result<Option<MetaCacheEntry>> {
self.check_init().await;
if let Some(err) = &self.err {
return Err(err.clone());
}
match rmp::decode::read_bool(&mut self.rd) {
match rmp::decode::read_bool(&mut self.read_more(1).await?) {
Ok(res) => {
if !res {
return Ok(None);
}
}
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
return Err(Error::new(err));
let serr = format!("{:?}", err);
self.err = Some(Error::msg(&serr));
return Err(Error::msg(&serr));
}
};
let l = match rmp::decode::read_str_len(&mut self.rd) {
Ok(res) => res,
let l = self.read_str_len().await?;
let buf = self.read_more(l as usize).await?;
let name_buf = buf.to_vec();
let name = match from_utf8(&name_buf) {
Ok(decoded) => decoded.to_owned(),
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
return Err(Error::new(err));
return Err(Error::msg(err.to_string()));
}
};
self.buf.resize(l as usize, 0);
let name = match self.rd.read_exact(&mut self.buf) {
Ok(()) => {
let name_buf = self.buf.to_vec();
match from_utf8(&name_buf) {
Ok(decoded) => Ok(decoded.to_owned()),
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
Err(Error::msg(err.to_string()))
}
}
}
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
Err(Error::msg(err.to_string()))
}
}?;
let l = self.read_bin_len().await?;
let l = match rmp::decode::read_bin_len(&mut self.rd) {
Ok(res) => res,
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
return Err(Error::new(err));
}
};
self.buf.resize(l as usize, 0);
match self.rd.read_exact(&mut self.buf) {
Ok(res) => res,
Err(err) => {
self.err = Some(Error::msg(err.to_string()));
return Err(Error::new(err));
}
};
let buf = self.read_more(l as usize).await?;
let metadata = self.buf.clone();
let metadata = buf.to_vec();
self.reset();
Ok(Some(MetaCacheEntry {
name,
@@ -218,11 +294,11 @@ impl<R: Read> MetacacheReader<R> {
}))
}
pub fn read_all(&mut self) -> Result<Vec<MetaCacheEntry>> {
pub async fn read_all(&mut self) -> Result<Vec<MetaCacheEntry>> {
let mut ret = Vec::new();
loop {
if let Some(entry) = self.peek()? {
if let Some(entry) = self.peek().await? {
ret.push(entry);
continue;
}
@@ -236,13 +312,12 @@ impl<R: Read> MetacacheReader<R> {
#[tokio::test]
async fn test_writer() {
use crate::io::AsyncToSync;
use crate::io::VecAsyncReader;
use crate::io::VecAsyncWriter;
let mut f = VecAsyncWriter::new(Vec::new());
let mut w = MetacacheWriter::new(AsyncToSync::new_writer(&mut f));
let mut w = MetacacheWriter::new(&mut f);
let mut objs = Vec::new();
for i in 0..10 {
@@ -256,14 +331,18 @@ async fn test_writer() {
objs.push(info);
}
w.write(&objs).unwrap();
w.write(&objs).await.unwrap();
w.close().unwrap();
w.close().await.unwrap();
let nf = VecAsyncReader::new(f.get_buffer().to_vec());
let data = f.get_buffer().to_vec();
let mut r = MetacacheReader::new(AsyncToSync::new_reader(nf));
let nobjs = r.read_all().unwrap();
println!("data len {}", data.len());
let nf = VecAsyncReader::new(data);
let mut r = MetacacheReader::new(nf);
let nobjs = r.read_all().await.unwrap();
for info in nobjs.iter() {
println!("new {:?}", &info);