use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
mod checkpoint;
mod context;
mod create_branch;
mod execute;
mod gc;
pub(crate) mod idempotency;
pub(crate) mod media_upload;
mod merge;
pub(crate) mod observe;
mod switch_branch;
mod transaction;
mod undo_redo;
pub(crate) use media_upload::stage_reclaimable_upload_receipts;
#[cfg(feature = "storage-benches")]
pub(crate) use merge::{MergeCommitsForBench, analyze_merge_for_bench};
pub(crate) use media_upload::{UPLOAD_MANIFEST_LEAF_SPACE, UPLOAD_STATE_SPACE};
pub(crate) use crate::common::ExecuteStatementMetadata;
#[cfg(feature = "server-protocol")]
pub(crate) use crate::common::VerifiedRequestBlob;
pub use checkpoint::CreateCheckpointReceipt;
pub(crate) use context::SessionBranch;
pub use context::SessionContext;
pub use create_branch::{CreateBranchOptions, CreateBranchReceipt};
pub use execute::{
CoherentReadBatch, ExecuteBatchStatement, ExecuteOptions, ExecuteResult, Row, RowRef,
TryFromValue,
};
pub(crate) use execute::{ExecutionDisposition, FileRead};
pub(crate) use idempotency::ExecuteIdempotency;
pub(crate) use idempotency::{
EXECUTE_IDEMPOTENCY_RECEIPT_SPACE, ExecuteIdempotencyReceipt, encode_receipt,
};
pub(crate) use media_upload::FileUploadProgress;
pub use merge::{
MergeBranchOptions, MergeBranchOutcome, MergeBranchPreview, MergeBranchPreviewOptions,
MergeBranchReceipt, MergeChangeStats, MergeConflict, MergeConflictChangeKind,
MergeConflictKind, MergeConflictSide,
};
pub use observe::{ObserveEvent, ObserveEvents};
pub use switch_branch::{SwitchBranchOptions, SwitchBranchReceipt};
pub use transaction::SessionTransaction;
pub use undo_redo::{RedoReceipt, UndoReceipt};
#[repr(transparent)]
pub(crate) struct AssumeSendFuture<F>(F);
impl<F> AssumeSendFuture<F> {
pub(crate) unsafe fn new(future: F) -> Self {
Self(future)
}
}
unsafe impl<F> Send for AssumeSendFuture<F>
where
F: Future,
F::Output: Send,
{
}
impl<F> Future for AssumeSendFuture<F>
where
F: Future,
{
type Output = F::Output;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
unsafe { self.map_unchecked_mut(|wrapped| &mut wrapped.0) }.poll(context)
}
}
#[cfg(test)]
pub(crate) mod borrowing_proof_storage {
use std::future::Future;
use crate::storage_adapter::{
Memory, MemoryRead, MemoryWrite, Storage, StorageBeginScanOptions, StorageError,
StorageGetManyRequest, StorageGetManyResult, StorageKeyRange, StorageRead,
StorageReadOptions, StorageScanCursor, StorageSpace, StorageWriteOptions,
};
pub(crate) struct BorrowingRead<'a> {
inner: MemoryRead,
_borrow: &'a Memory,
}
impl StorageRead for BorrowingRead<'_> {
fn snapshot_cache_key(&self) -> Option<u128> {
self.inner.snapshot_cache_key()
}
fn get_many(
&self,
requests: &[StorageGetManyRequest<'_>],
) -> impl Future<Output = Result<StorageGetManyResult, StorageError>> + Send {
self.inner.get_many(requests)
}
fn begin_scan(
&self,
space: StorageSpace,
range: StorageKeyRange,
opts: StorageBeginScanOptions,
) -> impl Future<Output = Result<StorageScanCursor<'_>, StorageError>> + Send {
self.inner.begin_scan(space, range, opts)
}
}
#[derive(Clone, Default)]
pub(crate) struct BorrowingStorage(Memory);
impl Storage for BorrowingStorage {
type Read<'a>
= BorrowingRead<'a>
where
Self: 'a;
type Write<'a>
= MemoryWrite
where
Self: 'a;
fn begin_read(
&self,
opts: StorageReadOptions,
) -> impl Future<Output = Result<Self::Read<'_>, StorageError>> + Send {
async move {
Ok(BorrowingRead {
inner: self.0.begin_read(opts).await?,
_borrow: &self.0,
})
}
}
fn begin_write(
&self,
opts: StorageWriteOptions,
) -> impl Future<Output = Result<Self::Write<'_>, StorageError>> + Send {
self.0.begin_write(opts)
}
}
}