use std::sync::Arc;
use futures::StreamExt;
use crate::byte_range::ByteRangeIterator;
use crate::{
AsyncListableStorageTraits, AsyncReadableStorageTraits, AsyncWritableStorageTraits, Bytes,
ListableStorageTraits, MaybeBytesIterator, MaybeSend, MaybeSync, OffsetBytesIterator,
ReadableStorageTraits, StorageError, StoreKey, StoreKeys, StoreKeysPrefixes, StorePrefix,
WritableStorageTraits,
};
pub trait AsyncToSyncBlockOn: MaybeSend + MaybeSync {
fn block_on<F: core::future::Future>(&self, future: F) -> F::Output;
}
pub struct AsyncToSyncStorageAdapter<TStorage: ?Sized, TBlockOn: AsyncToSyncBlockOn> {
storage: Arc<TStorage>,
block_on: TBlockOn,
}
impl<TStorage: ?Sized, TBlockOn: AsyncToSyncBlockOn> AsyncToSyncStorageAdapter<TStorage, TBlockOn> {
#[must_use]
pub fn new(storage: Arc<TStorage>, block_on: TBlockOn) -> Self {
Self { storage, block_on }
}
fn block_on<F: core::future::Future>(&self, future: F) -> F::Output {
self.block_on.block_on(future)
}
}
impl<TStorage: ?Sized + AsyncReadableStorageTraits, TBlockOn: AsyncToSyncBlockOn>
ReadableStorageTraits for AsyncToSyncStorageAdapter<TStorage, TBlockOn>
{
fn get_partial_many<'a>(
&'a self,
key: &StoreKey,
byte_ranges: ByteRangeIterator<'a>,
) -> Result<MaybeBytesIterator<'a>, StorageError> {
let results = self.block_on(self.storage.get_partial_many(key, byte_ranges))?;
if let Some(results) = results {
let results = self.block_on(results.collect::<Vec<_>>());
Ok(Some(Box::new(results.into_iter())))
} else {
Ok(None)
}
}
fn size_key(&self, key: &StoreKey) -> Result<Option<u64>, StorageError> {
self.block_on(self.storage.size_key(key))
}
fn supports_get_partial(&self) -> bool {
self.storage.supports_get_partial()
}
}
impl<TStorage: ?Sized + AsyncListableStorageTraits, TBlockOn: AsyncToSyncBlockOn>
ListableStorageTraits for AsyncToSyncStorageAdapter<TStorage, TBlockOn>
{
fn list(&self) -> Result<StoreKeys, StorageError> {
self.block_on(self.storage.list())
}
fn list_prefix(&self, prefix: &StorePrefix) -> Result<StoreKeys, StorageError> {
self.block_on(self.storage.list_prefix(prefix))
}
fn list_dir(&self, prefix: &StorePrefix) -> Result<StoreKeysPrefixes, StorageError> {
self.block_on(self.storage.list_dir(prefix))
}
fn size_prefix(&self, prefix: &StorePrefix) -> Result<u64, StorageError> {
self.block_on(self.storage.size_prefix(prefix))
}
}
impl<TStorage: ?Sized + AsyncWritableStorageTraits, TBlockOn: AsyncToSyncBlockOn>
WritableStorageTraits for AsyncToSyncStorageAdapter<TStorage, TBlockOn>
{
fn set(&self, key: &StoreKey, value: Bytes) -> Result<(), StorageError> {
self.block_on(self.storage.set(key, value))
}
fn set_partial_many(
&self,
key: &StoreKey,
offset_values: OffsetBytesIterator,
) -> Result<(), StorageError> {
self.block_on(self.storage.set_partial_many(key, offset_values))
}
fn erase(&self, key: &StoreKey) -> Result<(), StorageError> {
self.block_on(self.storage.erase(key))
}
fn erase_many(&self, keys: &[StoreKey]) -> Result<(), StorageError> {
self.block_on(self.storage.erase_many(keys))
}
fn erase_prefix(&self, prefix: &StorePrefix) -> Result<(), StorageError> {
self.block_on(self.storage.erase_prefix(prefix))
}
fn supports_set_partial(&self) -> bool {
self.storage.supports_set_partial()
}
}