use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
#[cfg(test)]
mod catalog_visibility_tests;
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;
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(crate) use checkpoint::CreateCheckpointReceipt;
pub use context::SessionContext;
pub(crate) use context::{SessionBranch, load_default_branch_id_from_index};
pub use create_branch::{CreateBranchOptions, CreateBranchReceipt};
pub use execute::{
CoherentReadBatch, CommitReceipt, CommitSpan, ExecuteBatchResult, ExecuteBatchStatement,
ExecuteOptions, ExecuteResult, ResultRowRef, Row, 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,
};
pub use observe::ObserveEvent;
pub(crate) use observe::ObserveEvents as SessionObserveEvents;
pub use switch_branch::{SwitchBranchOptions, SwitchBranchReceipt};
pub use transaction::SessionTransaction;
#[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 acquire_session(
&self,
) -> impl Future<Output = Result<crate::storage::StorageSessionToken, StorageError>> + Send
{
self.0.acquire_session()
}
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)
}
}
}
pub(crate) use media_upload::{export_recoverable_uploads, has_recoverable_uploads};
pub(crate) use execute::{
discover_read_fulfillment, prepare_partial_candidate_read_scope,
seed_foreground_filesystem_interest,
};
pub(crate) use merge::{
MergeAnalysis, analyze_incoming_rows, stage_merge_native_heads, stage_native_change_application,
};