use tokio::sync::{broadcast, oneshot};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use rskit_errors::AppResult;
use crate::event::Event;
pub struct TaskHandle<O: Clone + Send + 'static> {
pub id: Uuid,
events_rx: broadcast::Receiver<Event<O>>,
result_rx: oneshot::Receiver<AppResult<O>>,
cancel: CancellationToken,
}
impl<O: Clone + Send + 'static> TaskHandle<O> {
pub(crate) fn new(
id: Uuid,
events_rx: broadcast::Receiver<Event<O>>,
result_rx: oneshot::Receiver<AppResult<O>>,
cancel: CancellationToken,
) -> Self {
Self {
id,
events_rx,
result_rx,
cancel,
}
}
pub async fn result(self) -> AppResult<O> {
match self.result_rx.await {
Ok(r) => r,
Err(_) => Err(rskit_errors::AppError::new(
rskit_errors::ErrorCode::Internal,
"worker task dropped before completing",
)),
}
}
pub fn events(&self) -> broadcast::Receiver<Event<O>> {
self.events_rx.resubscribe()
}
pub fn cancel(&self) {
self.cancel.cancel();
}
pub fn cancel_token(&self) -> CancellationToken {
self.cancel.clone()
}
}