use std::fmt::Debug;
use crate::{
entry::Entrylike,
prelude::*,
storage::{
AppendToPayloadPrefixError, GetVerifiableStreamError, InternalOrNoSuchEntryError,
PayloadPrefixStore,
},
};
use bab_rs::CHUNK_SIZE;
use futures_lite::future::zip;
use ufotofu::{
ExpectedFinalError, PipeError, ProduceAtLeastError,
channels::{new_sssr, sssr::Receiver},
consumer::compat::fn_mut::*,
pipe,
prelude::*,
queues::{UnboundedElastic, new_unbounded_elastic},
};
mod decode;
pub use decode::*;
mod encode;
pub use encode::*;
mod slice_metadata;
pub use slice_metadata::{DropSliceMetadata, SliceMetadataError};
pub async fn export_drop<S, C, P>(
store: &mut S,
entry_producer: &mut P,
consumer: &mut C,
) -> Result<(), ExportDropError<S::InternalError, C::Error, P::Error>>
where
C: Consumer<Item = (DropSliceMetadata, Receiver<UnboundedElastic<u8>, ()>), Final = ()>,
S: PayloadPrefixStore + Clone,
P: Producer<Item = AuthorisedEntry>,
{
let mut s = store.clone();
let f =
async |value: Either<AuthorisedEntry, _>| -> Result<(), ExportDropError<S::InternalError, C::Error, P::Error>> {
match value {
Left(authorised_entry) => {
let byte_count = match s
.length_of_payload_prefix(
&authorised_entry.namespace_id(),
&authorised_entry,
Some(authorised_entry.payload_digest().clone()),
)
.await
{
Ok(byte_count) => byte_count,
Err(err) => match err {
InternalOrNoSuchEntryError::NoSuchEntry => {
return Ok(());
}
InternalOrNoSuchEntryError::StoreError(err) => {
return Err(ExportDropError::StoreError(err));
}
},
};
let chunk_count = byte_count.div_ceil(CHUNK_SIZE as u64);
let digest = authorised_entry.payload_digest().clone();
let metadata =
DropSliceMetadata::new(
authorised_entry.clone(),
0, chunk_count,
0, false ).expect("metadata generated from a known entry in the store should be valid");
let (mut sender, receiver) = new_sssr(new_unbounded_elastic());
let stream_options = metadata.streaming_options();
let slice_stream_generator = async {
match s
.get_verifiable_stream(
&authorised_entry.namespace_id(),
&authorised_entry,
Some(digest),
0,
byte_count,
stream_options,
&mut sender,
)
.await
{
Ok(_) => {
sender.consume_final(()).await.expect("sender is infallible");
Ok(())
},
Err(err) => match err {
GetVerifiableStreamError::ConsumerError(_) => unreachable!("sender is infallible"),
GetVerifiableStreamError::StoreError(err) => {
Err(ExportDropError::StoreError(err))
}
GetVerifiableStreamError::NoSuchEntry => {
Err(ExportDropError::EntryDeleted)
}
},
}
};
let (generator_result, consumer_result) = zip(
slice_stream_generator,
consumer.consume_item((metadata, receiver)),
)
.await;
generator_result?;
consumer_result.map_err(ExportDropError::ConsumerError)
}
Right(_) => consumer.consume_final(()).await.map_err(ExportDropError::ConsumerError)
}
};
pipe(entry_producer, &mut ClosureConsumer::new(f))
.await
.map_err(|err| match err {
PipeError::Producer(err) => ExportDropError::ProducerError(err),
PipeError::Consumer(err) => err,
})
}
pub async fn import_drop<S, P, PP>(
store: &mut S,
producer: &mut P,
skip_incompatible_slices: bool,
) -> Result<(), ImportDropError<P::Error, PP::Error, S::InternalError>>
where
S: PayloadPrefixStore + Clone,
P: Producer<Item = (DropSliceMetadata, PP)>,
PP: BulkProducer<Item = u8>,
{
while let Left((metadata, mut slice)) = producer
.produce()
.await
.map_err(ImportDropError::ImportProducerError)?
{
let entry = metadata.entry();
if !store
.insert_entry(entry.clone())
.await
.map_err(ImportDropError::StoreError)?
{
skip_slice(&metadata, slice).await?;
continue;
}
let Some(resumption_info) = store
.prefix_stream_resumption_info(entry.namespace_id(), entry, None)
.await
.map_err(|err| match err {
InternalOrNoSuchEntryError::NoSuchEntry => ImportDropError::EntryDeleted,
InternalOrNoSuchEntryError::StoreError(err) => ImportDropError::StoreError(err),
})?
else {
skip_slice(&metadata, slice).await?;
continue;
};
if metadata.resumption_info() == resumption_info && metadata.chunk_count() > 0 {
store
.append_to_payload_prefix(
entry.namespace_id(),
entry,
&mut slice,
metadata.streaming_options(),
)
.await?;
slice.produce_final().await?;
continue;
}
skip_slice(&metadata, slice).await?;
if metadata.first_chunk() > 0
&& resumption_info.start_chunk == 0
&& skip_incompatible_slices
{
return Err(ImportDropError::PartialPayloadSliceEncountered);
}
}
Ok(())
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub enum ExportDropError<StoreError, ConsumerError, ProducerError> {
StoreError(StoreError),
ConsumerError(ConsumerError),
ProducerError(ProducerError),
EntryDeleted,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ImportDropError<ImportProducerError, SliceProducerError, StoreError> {
SliceStreamBytesMismatch,
ImportProducerError(ImportProducerError),
SliceProducerError(SliceProducerError),
StoreError(StoreError),
VerificationError,
EntryDeleted,
ArchitectureTooSmall,
PartialPayloadSliceEncountered,
}
async fn skip_slice<P: BulkProducer<Item = u8>, DropProducerError, StoreError>(
metadata: &DropSliceMetadata,
mut slice: P,
) -> Result<(), ImportDropError<DropProducerError, P::Error, StoreError>> {
slice
.skip(
metadata
.expected_bytes()
.try_into()
.map_err(|_| ImportDropError::ArchitectureTooSmall)?,
)
.await?;
slice.produce_final().await?;
Ok(())
}
impl<Fin, ImportProducerError, SliceProducerError, StoreError>
From<ProduceAtLeastError<Fin, SliceProducerError>>
for ImportDropError<ImportProducerError, SliceProducerError, StoreError>
{
fn from(
ProduceAtLeastError { count: _, reason }: ProduceAtLeastError<Fin, SliceProducerError>,
) -> Self {
match reason {
Ok(_) => ImportDropError::SliceStreamBytesMismatch,
Err(err) => ImportDropError::SliceProducerError(err),
}
}
}
impl<ImportProducerError, SliceProducerError, StoreError>
From<AppendToPayloadPrefixError<SliceProducerError, StoreError>>
for ImportDropError<ImportProducerError, SliceProducerError, StoreError>
{
fn from(value: AppendToPayloadPrefixError<SliceProducerError, StoreError>) -> Self {
match value {
AppendToPayloadPrefixError::ProducerError(err) => Self::SliceProducerError(err),
AppendToPayloadPrefixError::UnexpectedEndOfStream => Self::SliceStreamBytesMismatch,
AppendToPayloadPrefixError::VerificationError => Self::VerificationError,
AppendToPayloadPrefixError::NoSuchEntry => Self::EntryDeleted,
AppendToPayloadPrefixError::StoreError(err) => Self::StoreError(err),
}
}
}
impl<ImportProducerError, SliceProducerError, StoreError>
From<ExpectedFinalError<u8, SliceProducerError>>
for ImportDropError<ImportProducerError, SliceProducerError, StoreError>
{
fn from(value: ExpectedFinalError<u8, SliceProducerError>) -> Self {
match value {
ExpectedFinalError::Item(_) => Self::SliceStreamBytesMismatch,
ExpectedFinalError::Error(err) => Self::SliceProducerError(err),
}
}
}