use std::{sync::Arc, time::Duration};
use compio::time::timeout;
use wresp::RespCommand;
use super::{
collection_item_broker::{
CollectionItemBroker, CollectionItemStore, CompioTaskSpawner, TaskSpawner,
},
collection_item_observer::{CollectionItemObserver, CollectionItemResult},
};
pub struct SharedItemBroker<S, Spawner = CompioTaskSpawner> {
inner: Arc<CollectionItemBroker<S, Spawner>>,
}
impl<S: CollectionItemStore + 'static> SharedItemBroker<S> {
pub fn new(broker: Arc<CollectionItemBroker<S>>) -> Self {
Self { inner: broker }
}
}
impl<S: CollectionItemStore + 'static, Spawner: TaskSpawner + 'static>
SharedItemBroker<S, Spawner>
{
#[inline]
pub fn start_wait(
&self,
command: RespCommand,
keys: Vec<Vec<u8>>,
session_id: usize,
cmd_args: Vec<Vec<u8>>,
) -> Arc<CollectionItemObserver> {
self.inner.start_wait(command, keys, session_id, cmd_args)
}
#[inline]
pub fn finish_wait(&self, observer: &Arc<CollectionItemObserver>) -> CollectionItemResult {
self.inner.finish_wait(observer)
}
#[inline]
pub fn handle_collection_update(&self, key: &[u8]) {
self.inner.handle_collection_update(key);
}
#[inline]
pub fn handle_session_disposed(&self, session_id: usize) {
self.inner.handle_session_disposed(session_id);
}
#[inline]
pub fn try_get_observer(&self, session_id: usize) -> Option<Arc<CollectionItemObserver>> {
self.inner.try_get_observer(session_id)
}
}
pub trait ItemBrokerFinisher: Send + Sync {
fn finish_wait(&self, observer: &Arc<CollectionItemObserver>) -> CollectionItemResult;
fn handle_session_disposed(&self, session_id: usize);
}
impl<S: CollectionItemStore + 'static, Spawner: TaskSpawner + 'static> ItemBrokerFinisher
for SharedItemBroker<S, Spawner>
{
fn finish_wait(&self, observer: &Arc<CollectionItemObserver>) -> CollectionItemResult {
self.finish_wait(observer)
}
fn handle_session_disposed(&self, session_id: usize) {
self.handle_session_disposed(session_id);
}
}
impl<T: ItemBrokerFinisher> ItemBrokerFinisher for Arc<T> {
fn finish_wait(&self, observer: &Arc<CollectionItemObserver>) -> CollectionItemResult {
(**self).finish_wait(observer)
}
fn handle_session_disposed(&self, session_id: usize) {
(**self).handle_session_disposed(session_id);
}
}
pub struct BlockedWait<B> {
broker: B,
observer: Arc<CollectionItemObserver>,
command: RespCommand,
timeout_secs: f64,
}
impl<B: ItemBrokerFinisher> BlockedWait<B> {
pub fn new(
broker: B,
observer: Arc<CollectionItemObserver>,
command: RespCommand,
timeout_secs: f64,
) -> Self {
Self {
broker,
observer,
command,
timeout_secs,
}
}
#[inline]
pub fn command(&self) -> RespCommand {
self.command
}
pub fn abort(&self) {
self
.broker
.handle_session_disposed(self.observer.session_id);
}
pub async fn resolve(self) -> (RespCommand, CollectionItemResult) {
if self.timeout_secs <= 0.0 {
self.observer.wait_result().await;
} else {
let _ = timeout(
Duration::from_secs_f64(self.timeout_secs),
self.observer.wait_result(),
)
.await;
}
let result = self.broker.finish_wait(&self.observer);
(self.command, result)
}
}