Zero copy IMAP passing tests
This commit is contained in:
@@ -128,11 +128,11 @@ impl Filter {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn contains(field: impl Into<u8>, value: Vec<u8>) -> Self {
|
||||
pub fn contains(field: impl Into<u8>, value: &str) -> Self {
|
||||
Filter::MatchValue {
|
||||
field: field.into(),
|
||||
op: Operator::Contains,
|
||||
value,
|
||||
value: value.to_lowercase().into_bytes(),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -14,6 +14,7 @@ use std::{
|
||||
|
||||
use nlp::tokenizers::word::WordTokenizer;
|
||||
use rand::Rng;
|
||||
use rkyv::util::AlignedVec;
|
||||
use roaring::RoaringBitmap;
|
||||
use utils::BlobHash;
|
||||
|
||||
@@ -29,10 +30,12 @@ pub mod key;
|
||||
pub mod log;
|
||||
pub mod serialize;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub(crate) const ARCHIVE_ALIGNMENT: usize = 16;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum Archive {
|
||||
Raw(Vec<u8>),
|
||||
Uncompressed(Vec<u8>),
|
||||
Aligned(AlignedVec<ARCHIVE_ALIGNMENT>),
|
||||
Vec(Vec<u8>),
|
||||
}
|
||||
|
||||
#[repr(transparent)]
|
||||
|
||||
@@ -4,31 +4,27 @@
|
||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||
*/
|
||||
|
||||
use crate::{Deserialize, Serialize, SerializeInfallible, Value};
|
||||
use rkyv::util::AlignedVec;
|
||||
|
||||
use super::{Archive, Archiver, LegacyBincode, assert::HashedValue};
|
||||
use crate::{Deserialize, Serialize, SerializeInfallible, U32_LEN, Value};
|
||||
|
||||
use super::{ARCHIVE_ALIGNMENT, Archive, Archiver, LegacyBincode, assert::HashedValue};
|
||||
|
||||
const MAGIC_MARKER: u8 = 1 << 7;
|
||||
const LZ4_COMPRESSES: u8 = 1 << 6;
|
||||
const ARCHIVE_UNCOMPRESSED: u8 = MAGIC_MARKER;
|
||||
const ARCHIVE_LZ4_COMPRESSED: u8 = MAGIC_MARKER | LZ4_COMPRESSES;
|
||||
const COMPRESS_WATERMARK: usize = 8192;
|
||||
const COMPRESS_DATA_OFFSET: usize = std::mem::size_of::<u32>() + 1;
|
||||
|
||||
impl Deserialize for Archive {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
match bytes.first().copied() {
|
||||
Some(ARCHIVE_UNCOMPRESSED) => Ok(Archive::Raw(bytes.to_vec())),
|
||||
Some(ARCHIVE_LZ4_COMPRESSED) => {
|
||||
lz4_flex::decompress_size_prepended(bytes.get(1..).unwrap_or_default())
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DecompressError
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
.map(Archive::Uncompressed)
|
||||
match bytes.split_last() {
|
||||
Some((&ARCHIVE_UNCOMPRESSED, archive)) => {
|
||||
let mut bytes = AlignedVec::with_capacity(archive.len());
|
||||
bytes.extend_from_slice(archive);
|
||||
Ok(Archive::Aligned(bytes))
|
||||
}
|
||||
Some((&ARCHIVE_LZ4_COMPRESSED, archive)) => aligned_lz4_deflate(archive),
|
||||
_ => Err(trc::StoreEvent::DataCorruption
|
||||
.into_err()
|
||||
.details("Invalid archive marker.")
|
||||
@@ -37,18 +33,20 @@ impl Deserialize for Archive {
|
||||
}
|
||||
}
|
||||
|
||||
fn deserialize_owned(bytes: Vec<u8>) -> trc::Result<Self> {
|
||||
match bytes.first().copied() {
|
||||
Some(ARCHIVE_UNCOMPRESSED) => Ok(Archive::Raw(bytes)),
|
||||
Some(ARCHIVE_LZ4_COMPRESSED) => {
|
||||
lz4_flex::decompress_size_prepended(bytes.get(1..).unwrap_or_default())
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DecompressError
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
.map(Archive::Uncompressed)
|
||||
fn deserialize_owned(mut bytes: Vec<u8>) -> trc::Result<Self> {
|
||||
match bytes.last() {
|
||||
Some(&ARCHIVE_UNCOMPRESSED) => {
|
||||
bytes.pop();
|
||||
if bytes.as_ptr().addr() & (ARCHIVE_ALIGNMENT - 1) == 0 {
|
||||
Ok(Archive::Vec(bytes))
|
||||
} else {
|
||||
let mut aligned = AlignedVec::with_capacity(bytes.len());
|
||||
aligned.extend_from_slice(&bytes);
|
||||
Ok(Archive::Aligned(aligned))
|
||||
}
|
||||
}
|
||||
Some(&ARCHIVE_LZ4_COMPRESSED) => {
|
||||
aligned_lz4_deflate(bytes.get(..bytes.len() - 1).unwrap_or_default())
|
||||
}
|
||||
_ => Err(trc::StoreEvent::DataCorruption
|
||||
.into_err()
|
||||
@@ -59,6 +57,22 @@ impl Deserialize for Archive {
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
fn aligned_lz4_deflate(archive: &[u8]) -> trc::Result<Archive> {
|
||||
lz4_flex::block::uncompressed_size(archive)
|
||||
.and_then(|(uncompressed_size, archive)| {
|
||||
let mut bytes = AlignedVec::with_capacity(uncompressed_size);
|
||||
lz4_flex::decompress_into(archive, &mut bytes)?;
|
||||
Ok(Archive::Aligned(bytes))
|
||||
})
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DecompressError
|
||||
.ctx(trc::Key::Value, archive)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
|
||||
impl<T> Serialize for Archiver<T>
|
||||
where
|
||||
T: rkyv::Archive
|
||||
@@ -81,28 +95,27 @@ where
|
||||
let input = input.as_ref();
|
||||
let input_len = input.len();
|
||||
if input_len > COMPRESS_WATERMARK {
|
||||
let mut bytes = vec![
|
||||
ARCHIVE_LZ4_COMPRESSED;
|
||||
lz4_flex::block::get_maximum_output_size(input_len)
|
||||
+ COMPRESS_DATA_OFFSET
|
||||
];
|
||||
bytes[1..COMPRESS_DATA_OFFSET]
|
||||
.copy_from_slice(&(input_len as u32).to_le_bytes());
|
||||
let bytes_len =
|
||||
lz4_flex::compress_into(input, &mut bytes[COMPRESS_DATA_OFFSET..]).unwrap()
|
||||
+ COMPRESS_DATA_OFFSET;
|
||||
let mut bytes =
|
||||
vec![
|
||||
ARCHIVE_LZ4_COMPRESSED;
|
||||
lz4_flex::block::get_maximum_output_size(input_len) + U32_LEN + 1
|
||||
];
|
||||
bytes[0..U32_LEN].copy_from_slice(&(input_len as u32).to_le_bytes());
|
||||
let bytes_len = lz4_flex::compress_into(input, &mut bytes[U32_LEN..]).unwrap()
|
||||
+ U32_LEN
|
||||
+ 1;
|
||||
if bytes_len < input_len {
|
||||
bytes.truncate(bytes_len);
|
||||
} else {
|
||||
bytes.clear();
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes.extend_from_slice(input);
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
}
|
||||
bytes
|
||||
} else {
|
||||
let mut bytes = Vec::with_capacity(input_len + 1);
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes.extend_from_slice(input);
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes
|
||||
}
|
||||
})
|
||||
@@ -110,62 +123,51 @@ where
|
||||
}
|
||||
|
||||
impl Archive {
|
||||
pub fn unarchive<T>(&self) -> trc::Result<&T>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + Sync
|
||||
+ Send,
|
||||
{
|
||||
#[inline]
|
||||
pub fn as_bytes(&self) -> &[u8] {
|
||||
match self {
|
||||
Archive::Raw(bytes) => rkyv::access::<T, rkyv::rancor::Error>(bytes.get(1..).unwrap())
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.caused_by(trc::location!())
|
||||
.ctx(trc::Key::Value, bytes.as_slice())
|
||||
.reason(err)
|
||||
}),
|
||||
Archive::Uncompressed(bytes) => Ok(unsafe { rkyv::access_unchecked::<T>(bytes) }),
|
||||
Archive::Vec(bytes) => bytes.as_slice(),
|
||||
Archive::Aligned(bytes) => bytes.as_slice(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn deserialize<T, V>(&self) -> trc::Result<V>
|
||||
pub fn unarchive<T>(&self) -> trc::Result<&<T as rkyv::Archive>::Archived>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<
|
||||
T: rkyv::Archive,
|
||||
T::Archived: for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + Sync
|
||||
+ Send
|
||||
+ rkyv::Deserialize<V, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
> + rkyv::Deserialize<T, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
self.unarchive::<T>().and_then(|value| {
|
||||
rkyv::deserialize::<V, rkyv::rancor::Error>(value).map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.ctx(
|
||||
trc::Key::Value,
|
||||
match self {
|
||||
Archive::Raw(bytes) => bytes,
|
||||
Archive::Uncompressed(bytes) => bytes,
|
||||
}
|
||||
.as_slice(),
|
||||
)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
rkyv::access::<T::Archived, rkyv::rancor::Error>(self.as_bytes()).map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.caused_by(trc::location!())
|
||||
.ctx(trc::Key::Value, self.as_bytes())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn deserialize<T>(&self) -> trc::Result<T>
|
||||
where
|
||||
T: rkyv::Archive,
|
||||
T::Archived: for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + rkyv::Deserialize<T, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
rkyv::from_bytes(self.as_bytes()).map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.ctx(trc::Key::Value, self.as_bytes())
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn into_inner(self) -> Vec<u8> {
|
||||
match self {
|
||||
Archive::Raw(bytes) => bytes,
|
||||
Archive::Uncompressed(bytes) => {
|
||||
let mut result = Vec::with_capacity(bytes.len() + 1);
|
||||
result.push(ARCHIVE_UNCOMPRESSED);
|
||||
result.extend_from_slice(&bytes);
|
||||
result
|
||||
}
|
||||
}
|
||||
let mut bytes = match self {
|
||||
Archive::Vec(bytes) => bytes,
|
||||
Archive::Aligned(bytes) => bytes.to_vec(),
|
||||
};
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes
|
||||
}
|
||||
}
|
||||
|
||||
@@ -190,30 +192,27 @@ where
|
||||
}
|
||||
|
||||
impl HashedValue<Archive> {
|
||||
pub fn to_unarchived<T>(&self) -> trc::Result<HashedValue<&T>>
|
||||
pub fn to_unarchived<T>(&self) -> trc::Result<HashedValue<&<T as rkyv::Archive>::Archived>>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<
|
||||
T: rkyv::Archive,
|
||||
T::Archived: for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + Sync
|
||||
+ Send,
|
||||
> + rkyv::Deserialize<T, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
self.inner.unarchive().map(|inner| HashedValue {
|
||||
self.inner.unarchive::<T>().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn into_deserialized<T, V>(self) -> trc::Result<HashedValue<V>>
|
||||
pub fn into_deserialized<T>(&self) -> trc::Result<HashedValue<T>>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<
|
||||
T: rkyv::Archive,
|
||||
T::Archived: for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + Sync
|
||||
+ Send
|
||||
+ rkyv::Deserialize<V, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
> + rkyv::Deserialize<T, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
self.inner.deserialize::<T, V>().map(|inner| HashedValue {
|
||||
self.inner.deserialize::<T>().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
@@ -236,14 +235,13 @@ where
|
||||
})
|
||||
}
|
||||
|
||||
pub fn rkyv_unarchive<T>(input: &[u8]) -> trc::Result<&T>
|
||||
pub fn rkyv_unarchive<T>(input: &[u8]) -> trc::Result<&<T as rkyv::Archive>::Archived>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>>
|
||||
+ Sync
|
||||
+ Send,
|
||||
T: rkyv::Archive,
|
||||
T::Archived: for<'a> rkyv::bytecheck::CheckBytes<rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>>
|
||||
+ rkyv::Deserialize<T, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
rkyv::access::<T, rkyv::rancor::Error>(input).map_err(|err| {
|
||||
rkyv::access::<T::Archived, rkyv::rancor::Error>(input).map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.caused_by(trc::location!())
|
||||
.ctx(trc::Key::Value, input)
|
||||
|
||||
Reference in New Issue
Block a user