use std::borrow::Cow;
use std::marker::PhantomData;
use std::ops::Range;
use crate::op_set2::change::ChangeCollector;
use crate::storage::change::{OpReadState, Unverified, Verified};
use crate::storage::columns::compression;
use crate::storage::columns::{ColumnId, ColumnType};
use crate::storage::{parse, Header, RawColumns};
use crate::types::{ActorId, ChangeHash};
use crate::Change;
use super::{BundleChangeIter, BundleChangeIterUnverified, OpIter, OpIterUnverified, ParseError};
const ID_CTR_INVERSE_COL_ID: ColumnId = ColumnId::new(11);
const ID_COL_ID: ColumnId = ColumnId::new(2);
#[derive(Clone, Debug)]
pub(crate) struct BundleStorage<'a, OpReadState> {
pub(crate) bytes: Cow<'a, [u8]>,
pub(crate) compressed_bytes: Option<Cow<'a, [u8]>>,
pub(crate) header: Header,
pub(crate) deps: Vec<ChangeHash>,
pub(crate) actors: Vec<ActorId>,
pub(crate) ops_meta: RawColumns<compression::Uncompressed>,
pub(crate) ops_data: Range<usize>,
pub(crate) changes_meta: RawColumns<compression::Uncompressed>,
pub(crate) changes_data: Range<usize>,
pub(crate) id_ctr: Vec<i64>,
pub(crate) _phantom: PhantomData<OpReadState>,
}
impl<O: OpReadState> BundleStorage<'_, O> {
pub(crate) fn into_owned(self) -> BundleStorage<'static, O> {
BundleStorage {
bytes: Cow::Owned(self.bytes.into_owned()),
compressed_bytes: self.compressed_bytes.map(|c| Cow::Owned(c.into_owned())),
header: self.header,
deps: self.deps,
actors: self.actors,
ops_meta: self.ops_meta,
ops_data: self.ops_data,
changes_meta: self.changes_meta,
changes_data: self.changes_data,
id_ctr: self.id_ctr,
_phantom: self._phantom,
}
}
pub(crate) fn checksum_valid(&self) -> bool {
self.header.checksum_valid()
}
}
fn extract_id_ctr_values(
changes_meta: &RawColumns<compression::Uncompressed>,
changes_data: &[u8],
ops_meta: &RawColumns<compression::Uncompressed>,
ops_data: &[u8],
) -> Result<Vec<i64>, ParseError> {
let mut inverse_bytes: Option<&[u8]> = None;
let mut id_ctr_bytes: Option<&[u8]> = None;
for col in ops_meta.0.iter() {
let spec = col.spec();
if spec.col_type() != ColumnType::DeltaInteger {
continue;
}
match spec.id() {
id if id == ID_CTR_INVERSE_COL_ID => {
let d = col.data();
inverse_bytes = Some(&ops_data[d.start..d.end]);
}
id if id == ID_COL_ID => {
let d = col.data();
id_ctr_bytes = Some(&ops_data[d.start..d.end]);
}
_ => {}
}
}
if let Some(inverse_bytes) = inverse_bytes {
let inverse: Vec<i64> = decode_delta_int(inverse_bytes)?;
let mut change_meta: Vec<(usize, u64, u64, u64)> =
BundleChangeIterUnverified::try_new(changes_meta, changes_data)?
.map(|c| c.map(|c| (c.actor, c.seq, c.start_op, c.max_op)))
.collect::<Result<_, _>>()?;
change_meta.sort_unstable_by_key(|(actor, seq, _, _)| (*actor, *seq));
let mut counters = vec![0i64; inverse.len()];
let mut k = 0usize;
for (_actor, _seq, start_op, max_op) in &change_meta {
for ctr in *start_op..=*max_op {
if k >= inverse.len() {
return Err(ParseError::InverseLengthMismatch);
}
let doc_pos = inverse[k] as usize;
if doc_pos >= counters.len() {
return Err(ParseError::InverseDecode);
}
counters[doc_pos] = ctr as i64;
k += 1;
}
}
if k != inverse.len() {
return Err(ParseError::InverseLengthMismatch);
}
return Ok(counters);
}
if let Some(id_ctr_bytes) = id_ctr_bytes {
return decode_delta_int(id_ctr_bytes);
}
Ok(Vec::new())
}
fn decode_delta_int(bytes: &[u8]) -> Result<Vec<i64>, ParseError> {
hexane::DeltaDecoder::<Option<i64>>::new(bytes)
.map(|item| item.ok_or(ParseError::InverseDecode))
.collect()
}
impl<'a> BundleStorage<'a, Unverified> {
pub(crate) fn parse_following_header(
input: parse::Input<'a>,
header: Header,
) -> parse::ParseResult<'a, BundleStorage<'a, Unverified>, ParseError> {
let full_bytes = input.bytes();
let (i, prefix_r) = parse::range_of(
|i| -> parse::ParseResult<'_, _, ParseError> {
let (i, deps) = parse::length_prefixed(parse::change_hash)(i)?;
let (i, actors) = parse::length_prefixed(parse::actor_id)(i)?;
Ok((i, (deps, actors)))
},
input,
)?;
let (deps, actors) = prefix_r.value;
let prefix_end = prefix_r.range.end;
let (i, changes_meta_raw) = RawColumns::parse(i)?;
let (i, changes) =
parse::range_of(|i| parse::take_n(changes_meta_raw.total_column_len(), i), i)?;
let changes_data_range = changes.range.clone();
let (i, ops_meta_raw) = RawColumns::parse(i)?;
let (_, ops) = parse::range_of(|i| parse::take_n(ops_meta_raw.total_column_len(), i), i)?;
let ops_data_range = ops.range.clone();
if let (Some(changes_meta), Some(ops_meta)) =
(changes_meta_raw.uncompressed(), ops_meta_raw.uncompressed())
{
BundleChangeIterUnverified::try_new(&changes_meta, changes.value)
.map_err(|e| parse::ParseError::Error(ParseError::InvalidColumns(Box::new(e))))?;
let id_ctr = extract_id_ctr_values(&changes_meta, changes.value, &ops_meta, ops.value)
.map_err(parse::ParseError::Error)?;
OpIterUnverified::try_new(&ops_meta, ops.value, &id_ctr)
.map_err(|e| parse::ParseError::Error(ParseError::InvalidColumns(Box::new(e))))?;
return Ok((
parse::Input::empty(),
BundleStorage {
bytes: full_bytes.into(),
compressed_bytes: None,
header,
deps,
actors,
ops_meta,
ops_data: ops_data_range,
changes_meta,
changes_data: changes_data_range,
id_ctr,
_phantom: PhantomData,
},
));
}
let mut out = Vec::with_capacity(full_bytes.len());
out.extend_from_slice(&full_bytes[..prefix_end]);
let mut changes_data_buf = Vec::new();
let changes_meta = changes_meta_raw
.uncompress(
&full_bytes[changes_data_range.clone()],
&mut changes_data_buf,
)
.map_err(|_| parse::ParseError::Error(ParseError::CompressedChangeCols))?;
changes_meta.write(&mut out);
let new_changes_start = out.len();
out.extend_from_slice(&changes_data_buf);
let new_changes_end = out.len();
let mut ops_data_buf = Vec::new();
let ops_meta = ops_meta_raw
.uncompress(&full_bytes[ops_data_range.clone()], &mut ops_data_buf)
.map_err(|_| parse::ParseError::Error(ParseError::CompressedOpCols))?;
ops_meta.write(&mut out);
let new_ops_start = out.len();
out.extend_from_slice(&ops_data_buf);
let new_ops_end = out.len();
BundleChangeIterUnverified::try_new(
&changes_meta,
&out[new_changes_start..new_changes_end],
)
.map_err(|e| parse::ParseError::Error(ParseError::InvalidColumns(Box::new(e))))?;
let id_ctr = extract_id_ctr_values(
&changes_meta,
&out[new_changes_start..new_changes_end],
&ops_meta,
&out[new_ops_start..new_ops_end],
)
.map_err(parse::ParseError::Error)?;
OpIterUnverified::try_new(&ops_meta, &out[new_ops_start..new_ops_end], &id_ctr)
.map_err(|e| parse::ParseError::Error(ParseError::InvalidColumns(Box::new(e))))?;
Ok((
parse::Input::empty(),
BundleStorage {
bytes: Cow::Owned(out),
compressed_bytes: Some(full_bytes.into()),
header,
deps,
actors,
ops_meta,
ops_data: new_ops_start..new_ops_end,
changes_meta,
changes_data: new_changes_start..new_changes_end,
id_ctr,
_phantom: PhantomData,
},
))
}
pub(crate) fn verify(self) -> Result<BundleStorage<'a, Verified>, ParseError> {
for c in self.iter_change_meta() {
let _ = c?;
}
for o in self.iter_ops() {
let _ = o?;
}
Ok(BundleStorage {
bytes: self.bytes,
compressed_bytes: self.compressed_bytes,
header: self.header,
deps: self.deps,
actors: self.actors,
ops_meta: self.ops_meta,
ops_data: self.ops_data,
changes_meta: self.changes_meta,
changes_data: self.changes_data,
id_ctr: self.id_ctr,
_phantom: PhantomData,
})
}
pub(crate) fn iter_ops(&self) -> OpIterUnverified<'_> {
let bytes = &self.bytes[self.ops_data.clone()];
OpIterUnverified::new(&self.ops_meta, bytes, &self.id_ctr)
}
fn iter_change_meta(&self) -> BundleChangeIterUnverified<'_> {
let change_data = &self.bytes[self.changes_data.clone()];
BundleChangeIterUnverified::new(&self.changes_meta, change_data)
}
}
impl BundleStorage<'_, Verified> {
pub(crate) fn to_changes(&self) -> Result<Vec<Change>, ParseError> {
let change_meta = self.iter_change_meta().collect();
let mut collector = ChangeCollector::from_bundle_changes(change_meta, &self.actors);
for op in self.iter_ops() {
collector.add(op);
}
let bundle = collector
.unbundle(&self.actors, &self.deps)
.map_err(|e| ParseError::Unbundle(Box::new(e)))?;
Ok(bundle)
}
pub(crate) fn iter_ops(&self) -> OpIter<'_> {
let bytes = &self.bytes[self.ops_data.clone()];
OpIter::new(&self.ops_meta, bytes, &self.id_ctr)
}
pub(crate) fn iter_change_meta(&self) -> BundleChangeIter<'_> {
let change_data = &self.bytes[self.changes_data.clone()];
BundleChangeIter::new_from_verified(&self.changes_meta, change_data)
}
pub(crate) fn deps(&self) -> &[ChangeHash] {
&self.deps
}
}