Zero copy everything
This commit is contained in:
@@ -31,7 +31,7 @@ use super::{Field, postings::Postings};
|
||||
pub const TERM_INDEX_VERSION: u8 = 1;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct Text<'x, T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
pub(crate) struct Text<'x, T: Into<u8> + Display + std::fmt::Debug> {
|
||||
pub field: Field<T>,
|
||||
pub text: Cow<'x, str>,
|
||||
pub typ: Type,
|
||||
@@ -45,7 +45,7 @@ pub(crate) enum Type {
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct FtsDocument<'x, T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
pub struct FtsDocument<'x, T: Into<u8> + Display + std::fmt::Debug> {
|
||||
pub(crate) parts: Vec<Text<'x, T>>,
|
||||
pub(crate) default_language: Language,
|
||||
pub(crate) account_id: u32,
|
||||
@@ -53,7 +53,7 @@ pub struct FtsDocument<'x, T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
pub(crate) document_id: u32,
|
||||
}
|
||||
|
||||
impl<'x, T: Into<u8> + Display + Clone + std::fmt::Debug> FtsDocument<'x, T> {
|
||||
impl<'x, T: Into<u8> + Display + std::fmt::Debug> FtsDocument<'x, T> {
|
||||
pub fn with_default_language(default_language: Language) -> FtsDocument<'x, T> {
|
||||
FtsDocument {
|
||||
parts: vec![],
|
||||
@@ -107,7 +107,7 @@ impl<'x, T: Into<u8> + Display + Clone + std::fmt::Debug> FtsDocument<'x, T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Into<u8> + Display + Clone + std::fmt::Debug> From<Field<T>> for u8 {
|
||||
impl<T: Into<u8> + Display + std::fmt::Debug> From<Field<T>> for u8 {
|
||||
fn from(value: Field<T>) -> Self {
|
||||
match value {
|
||||
Field::Body => 0,
|
||||
@@ -119,7 +119,7 @@ impl<T: Into<u8> + Display + Clone + std::fmt::Debug> From<Field<T>> for u8 {
|
||||
}
|
||||
|
||||
impl Store {
|
||||
pub async fn fts_index<T: Into<u8> + Display + Clone + std::fmt::Debug>(
|
||||
pub async fn fts_index<T: Into<u8> + Display + std::fmt::Debug>(
|
||||
&self,
|
||||
document: FtsDocument<'_, T>,
|
||||
) -> trc::Result<()> {
|
||||
|
||||
@@ -13,7 +13,7 @@ pub mod postings;
|
||||
pub mod query;
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub enum Field<T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
pub enum Field<T: Into<u8> + Display + std::fmt::Debug> {
|
||||
Header(T),
|
||||
Body,
|
||||
Attachment,
|
||||
@@ -21,7 +21,7 @@ pub enum Field<T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub enum FtsFilter<T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
pub enum FtsFilter<T: Into<u8> + Display + std::fmt::Debug> {
|
||||
Exact {
|
||||
field: Field<T>,
|
||||
text: String,
|
||||
@@ -42,7 +42,7 @@ pub enum FtsFilter<T: Into<u8> + Display + Clone + std::fmt::Debug> {
|
||||
End,
|
||||
}
|
||||
|
||||
impl<T: Into<u8> + Display + Clone + std::fmt::Debug> FtsFilter<T> {
|
||||
impl<T: Into<u8> + Display + std::fmt::Debug> FtsFilter<T> {
|
||||
pub fn has_text_detect(
|
||||
field: Field<T>,
|
||||
text: impl Into<String>,
|
||||
|
||||
@@ -9,22 +9,15 @@ use std::{
|
||||
collections::HashSet,
|
||||
fmt::{self, Formatter},
|
||||
hash::Hash,
|
||||
slice::Iter,
|
||||
time::{Duration, SystemTime},
|
||||
};
|
||||
|
||||
use assert::HashedValue;
|
||||
use nlp::tokenizers::word::WordTokenizer;
|
||||
use rand::Rng;
|
||||
use roaring::RoaringBitmap;
|
||||
use utils::{
|
||||
BlobHash,
|
||||
codec::leb128::{Leb128Iterator, Leb128Vec},
|
||||
};
|
||||
use utils::BlobHash;
|
||||
|
||||
use crate::{
|
||||
BlobClass, Deserialize, Serialize, SerializeInfallible, Value, backend::MAX_TOKEN_LENGTH,
|
||||
};
|
||||
use crate::{BlobClass, SerializeInfallible, backend::MAX_TOKEN_LENGTH};
|
||||
|
||||
use self::assert::AssertValue;
|
||||
|
||||
@@ -34,6 +27,30 @@ pub mod blob;
|
||||
pub mod hash;
|
||||
pub mod key;
|
||||
pub mod log;
|
||||
pub mod serialize;
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum Archive {
|
||||
Raw(Vec<u8>),
|
||||
Uncompressed(Vec<u8>),
|
||||
}
|
||||
|
||||
#[repr(transparent)]
|
||||
pub struct Archiver<T>(pub T)
|
||||
where
|
||||
T: rkyv::Archive
|
||||
+ for<'a> rkyv::Serialize<
|
||||
rkyv::api::high::HighSerializer<
|
||||
rkyv::util::AlignedVec,
|
||||
rkyv::ser::allocator::ArenaHandle<'a>,
|
||||
rkyv::rancor::Error,
|
||||
>,
|
||||
>;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct LegacyBincode<T: serde::Serialize + serde::de::DeserializeOwned> {
|
||||
pub inner: T,
|
||||
}
|
||||
|
||||
pub trait SerializeWithId: Send + Sync {
|
||||
fn serialize_with_id(&self, ids: &AssignedIds) -> trc::Result<Vec<u8>>;
|
||||
@@ -303,177 +320,6 @@ impl<T> From<()> for TagValue<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u32 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for i64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u16 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for f64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for &str {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.as_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for String {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(String::from_utf8_lossy(bytes).into_owned())
|
||||
}
|
||||
|
||||
fn deserialize_owned(bytes: Vec<u8>) -> trc::Result<Self> {
|
||||
Ok(String::from_utf8(bytes)
|
||||
.unwrap_or_else(|err| String::from_utf8_lossy(err.as_bytes()).into_owned()))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for u64 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(u64::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for i64 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(i64::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for u32 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(u32::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
pub trait SerializeInto {
|
||||
fn serialize_into(&self, buf: &mut Vec<u8>);
|
||||
}
|
||||
|
||||
pub trait DeserializeFrom: Sized {
|
||||
fn deserialize_from(bytes: &mut Iter<'_, u8>) -> Option<Self>;
|
||||
}
|
||||
|
||||
pub struct ArchivedValue<T> {
|
||||
inner: Vec<u8>,
|
||||
_phantom: std::marker::PhantomData<T>,
|
||||
}
|
||||
|
||||
impl<T: SerializeInto> Serialize for Vec<T> {
|
||||
fn serialize(&self) -> trc::Result<Vec<u8>> {
|
||||
let mut bytes = Vec::with_capacity(self.len() * 4);
|
||||
bytes.push_leb128(self.len());
|
||||
for item in self.iter() {
|
||||
item.serialize_into(&mut bytes);
|
||||
}
|
||||
Ok(bytes)
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInto for String {
|
||||
fn serialize_into(&self, buf: &mut Vec<u8>) {
|
||||
buf.push_leb128(self.len());
|
||||
if !self.is_empty() {
|
||||
buf.extend_from_slice(self.as_bytes());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInto for Vec<u8> {
|
||||
fn serialize_into(&self, buf: &mut Vec<u8>) {
|
||||
buf.push_leb128(self.len());
|
||||
if !self.is_empty() {
|
||||
buf.extend_from_slice(self.as_slice());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInto for u32 {
|
||||
fn serialize_into(&self, buf: &mut Vec<u8>) {
|
||||
buf.push_leb128(*self);
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInto for u64 {
|
||||
fn serialize_into(&self, buf: &mut Vec<u8>) {
|
||||
buf.push_leb128(*self);
|
||||
}
|
||||
}
|
||||
|
||||
impl DeserializeFrom for u32 {
|
||||
fn deserialize_from(bytes: &mut Iter<'_, u8>) -> Option<Self> {
|
||||
bytes.next_leb128()
|
||||
}
|
||||
}
|
||||
|
||||
impl DeserializeFrom for u64 {
|
||||
fn deserialize_from(bytes: &mut Iter<'_, u8>) -> Option<Self> {
|
||||
bytes.next_leb128()
|
||||
}
|
||||
}
|
||||
|
||||
impl DeserializeFrom for String {
|
||||
fn deserialize_from(bytes: &mut Iter<'_, u8>) -> Option<Self> {
|
||||
<Vec<u8>>::deserialize_from(bytes).and_then(|s| String::from_utf8(s).ok())
|
||||
}
|
||||
}
|
||||
|
||||
impl DeserializeFrom for Vec<u8> {
|
||||
fn deserialize_from(bytes: &mut Iter<'_, u8>) -> Option<Self> {
|
||||
let len: usize = bytes.next_leb128()?;
|
||||
let mut buf = Vec::with_capacity(len);
|
||||
for _ in 0..len {
|
||||
buf.push(*bytes.next()?);
|
||||
}
|
||||
buf.into()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: DeserializeFrom + Sync + Send> Deserialize for Vec<T> {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
let mut bytes = bytes.iter();
|
||||
let len: usize = bytes
|
||||
.next_leb128()
|
||||
.ok_or_else(|| trc::StoreEvent::DataCorruption.caused_by(trc::location!()))?;
|
||||
let mut list = Vec::with_capacity(len);
|
||||
for _ in 0..len {
|
||||
list.push(
|
||||
T::deserialize_from(&mut bytes)
|
||||
.ok_or_else(|| trc::StoreEvent::DataCorruption.caused_by(trc::location!()))?,
|
||||
);
|
||||
}
|
||||
Ok(list)
|
||||
}
|
||||
}
|
||||
|
||||
pub trait TokenizeText {
|
||||
fn tokenize_into(&self, tokens: &mut HashSet<String>);
|
||||
fn to_tokens(&self) -> HashSet<String>;
|
||||
@@ -493,18 +339,6 @@ impl TokenizeText for &str {
|
||||
}
|
||||
}
|
||||
|
||||
impl Serialize for () {
|
||||
fn serialize(&self) -> trc::Result<Vec<u8>> {
|
||||
Ok(Vec::with_capacity(0))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for () {
|
||||
fn deserialize(_bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub trait IntoOperations {
|
||||
fn build(self, batch: &mut BatchBuilder) -> trc::Result<()>;
|
||||
}
|
||||
@@ -574,124 +408,6 @@ impl BlobClass {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Bincode<T: serde::Serialize + serde::de::DeserializeOwned> {
|
||||
pub inner: T,
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> Bincode<T> {
|
||||
pub fn new(inner: T) -> Self {
|
||||
Self { inner }
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> From<Value<'static>> for Bincode<T> {
|
||||
fn from(_: Value<'static>) -> Self {
|
||||
unreachable!("From Value called on Bincode<T>")
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> Serialize for Bincode<T> {
|
||||
fn serialize(&self) -> trc::Result<Vec<u8>> {
|
||||
bincode::serialize(&self.inner)
|
||||
.map(|bytes| lz4_flex::compress_prepend_size(&bytes))
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned + Sized + Sync + Send> Deserialize
|
||||
for Bincode<T>
|
||||
{
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
lz4_flex::decompress_size_prepended(bytes)
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DecompressError
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
.and_then(|result| {
|
||||
bincode::deserialize(&result).map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
})
|
||||
.map(|inner| Self { inner })
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: Sync + Send> Deserialize for ArchivedValue<T> {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(ArchivedValue {
|
||||
inner: bytes.to_vec(),
|
||||
_phantom: std::marker::PhantomData,
|
||||
})
|
||||
}
|
||||
|
||||
fn deserialize_owned(bytes: Vec<u8>) -> trc::Result<Self> {
|
||||
Ok(ArchivedValue {
|
||||
inner: bytes,
|
||||
_phantom: std::marker::PhantomData,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> ArchivedValue<T>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>>
|
||||
+ Sync
|
||||
+ Send,
|
||||
{
|
||||
pub fn unarchive(&self) -> trc::Result<&T> {
|
||||
rkyv::access::<T, rkyv::rancor::Error>(&self.inner).map_err(Into::into)
|
||||
}
|
||||
|
||||
pub fn unarchive_unsafe(&self) -> &T {
|
||||
unsafe { rkyv::access_unchecked::<T>(&self.inner) }
|
||||
}
|
||||
|
||||
pub fn deserialize<V>(&self) -> trc::Result<V>
|
||||
where
|
||||
T: rkyv::Deserialize<V, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
rkyv::access::<T, rkyv::rancor::Error>(&self.inner)
|
||||
.and_then(|value| rkyv::deserialize::<V, rkyv::rancor::Error>(value))
|
||||
.map_err(Into::into)
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> HashedValue<ArchivedValue<T>>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>>
|
||||
+ Sync
|
||||
+ Send,
|
||||
{
|
||||
pub fn to_unarchived(&self) -> trc::Result<HashedValue<&T>> {
|
||||
self.inner.unarchive().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn into_deserialized<V>(self) -> trc::Result<HashedValue<V>>
|
||||
where
|
||||
T: rkyv::Deserialize<V, rkyv::api::high::HighDeserializer<rkyv::rancor::Error>>,
|
||||
{
|
||||
self.inner.deserialize().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl AssignedIds {
|
||||
pub fn push_document_id(&mut self, id: u32) {
|
||||
self.document_ids.push(id);
|
||||
|
||||
382
crates/store/src/write/serialize.rs
Normal file
382
crates/store/src/write/serialize.rs
Normal file
@@ -0,0 +1,382 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs Ltd <hello@stalw.art>
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||
*/
|
||||
|
||||
use crate::{Deserialize, Serialize, SerializeInfallible, Value};
|
||||
|
||||
use super::{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)
|
||||
}
|
||||
_ => Err(trc::StoreEvent::DataCorruption
|
||||
.into_err()
|
||||
.details("Invalid archive marker.")
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())),
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
_ => Err(trc::StoreEvent::DataCorruption
|
||||
.into_err()
|
||||
.details("Invalid archive marker.")
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Serialize for Archiver<T>
|
||||
where
|
||||
T: rkyv::Archive
|
||||
+ for<'a> rkyv::Serialize<
|
||||
rkyv::api::high::HighSerializer<
|
||||
rkyv::util::AlignedVec,
|
||||
rkyv::ser::allocator::ArenaHandle<'a>,
|
||||
rkyv::rancor::Error,
|
||||
>,
|
||||
>,
|
||||
{
|
||||
fn serialize(&self) -> trc::Result<Vec<u8>> {
|
||||
rkyv::to_bytes::<rkyv::rancor::Error>(&self.0)
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
.map(|input| {
|
||||
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;
|
||||
if bytes_len < input_len {
|
||||
bytes.truncate(bytes_len);
|
||||
} else {
|
||||
bytes.clear();
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes.extend_from_slice(input);
|
||||
}
|
||||
bytes
|
||||
} else {
|
||||
let mut bytes = Vec::with_capacity(input_len + 1);
|
||||
bytes.push(ARCHIVE_UNCOMPRESSED);
|
||||
bytes.extend_from_slice(input);
|
||||
bytes
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
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,
|
||||
{
|
||||
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) }),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn deserialize<T, V>(&self) -> trc::Result<V>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ 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>>,
|
||||
{
|
||||
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)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Archiver<T>
|
||||
where
|
||||
T: rkyv::Archive
|
||||
+ for<'a> rkyv::Serialize<
|
||||
rkyv::api::high::HighSerializer<
|
||||
rkyv::util::AlignedVec,
|
||||
rkyv::ser::allocator::ArenaHandle<'a>,
|
||||
rkyv::rancor::Error,
|
||||
>,
|
||||
>,
|
||||
{
|
||||
pub fn new(inner: T) -> Self {
|
||||
Self(inner)
|
||||
}
|
||||
|
||||
pub fn into_inner(self) -> T {
|
||||
self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl HashedValue<Archive> {
|
||||
pub fn to_unarchived<T>(&self) -> trc::Result<HashedValue<&T>>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<
|
||||
rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>,
|
||||
> + Sync
|
||||
+ Send,
|
||||
{
|
||||
self.inner.unarchive().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn into_deserialized<T, V>(self) -> trc::Result<HashedValue<V>>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ 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>>,
|
||||
{
|
||||
self.inner.deserialize::<T, V>().map(|inner| HashedValue {
|
||||
hash: self.hash,
|
||||
inner,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[inline]
|
||||
pub fn rkyv_deserialize<T, V>(input: &T) -> trc::Result<V>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ 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::<V, rkyv::rancor::Error>(input).map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn rkyv_unarchive<T>(input: &[u8]) -> trc::Result<&T>
|
||||
where
|
||||
T: rkyv::Portable
|
||||
+ for<'a> rkyv::bytecheck::CheckBytes<rkyv::api::high::HighValidator<'a, rkyv::rancor::Error>>
|
||||
+ Sync
|
||||
+ Send,
|
||||
{
|
||||
rkyv::access::<T, rkyv::rancor::Error>(input).map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.caused_by(trc::location!())
|
||||
.ctx(trc::Key::Value, input)
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u32 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for i64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for u16 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for f64 {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.to_be_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl SerializeInfallible for &str {
|
||||
fn serialize(&self) -> Vec<u8> {
|
||||
self.as_bytes().to_vec()
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for String {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(String::from_utf8_lossy(bytes).into_owned())
|
||||
}
|
||||
|
||||
fn deserialize_owned(bytes: Vec<u8>) -> trc::Result<Self> {
|
||||
Ok(String::from_utf8(bytes)
|
||||
.unwrap_or_else(|err| String::from_utf8_lossy(err.as_bytes()).into_owned()))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for u64 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(u64::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for i64 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(i64::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for u32 {
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(u32::from_be_bytes(bytes.try_into().map_err(|_| {
|
||||
trc::StoreEvent::DataCorruption.caused_by(trc::location!())
|
||||
})?))
|
||||
}
|
||||
}
|
||||
|
||||
impl Deserialize for () {
|
||||
fn deserialize(_bytes: &[u8]) -> trc::Result<Self> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> LegacyBincode<T> {
|
||||
pub fn new(inner: T) -> Self {
|
||||
Self { inner }
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> Serialize for LegacyBincode<T> {
|
||||
fn serialize(&self) -> trc::Result<Vec<u8>> {
|
||||
bincode::serialize(&self.inner)
|
||||
.map(|bytes| lz4_flex::compress_prepend_size(&bytes))
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DeserializeError
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned + Sized + Sync + Send> Deserialize
|
||||
for LegacyBincode<T>
|
||||
{
|
||||
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
|
||||
lz4_flex::decompress_size_prepended(bytes)
|
||||
.map_err(|err| {
|
||||
trc::StoreEvent::DecompressError
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
.and_then(|result| {
|
||||
bincode::deserialize(&result).map_err(|err| {
|
||||
trc::StoreEvent::DataCorruption
|
||||
.ctx(trc::Key::Value, bytes)
|
||||
.caused_by(trc::location!())
|
||||
.reason(err)
|
||||
})
|
||||
})
|
||||
.map(|inner| Self { inner })
|
||||
}
|
||||
}
|
||||
|
||||
impl From<Value<'static>> for Archive {
|
||||
fn from(_: Value<'static>) -> Self {
|
||||
unimplemented!()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: serde::Serialize + serde::de::DeserializeOwned> From<Value<'static>> for LegacyBincode<T> {
|
||||
fn from(_: Value<'static>) -> Self {
|
||||
unimplemented!()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user