use std::{fmt, future::Future, sync::Arc};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use crate::{BatchSubscriber, Broker, Connected, Subscriber, SubscriptionSource};
use crate::runtime::batch::BatchHandler;
use crate::runtime::dispatch::{
Delivery, Workers, spawn_batch_dispatch, spawn_dispatch, spawn_dispatch_workers,
};
use crate::runtime::failure::{DispatchFailure, ErrorShutdown, FailurePolicies};
use crate::runtime::handler::Handler;
use crate::runtime::lifecycle::{BoxError, BoxFuture};
use crate::runtime::metadata::HandlerMetadata;
use super::SourceMessage;
pub(crate) type BoundStarter<B, State> = Box<
dyn FnOnce(
Arc<Connected<B>>,
Arc<State>,
Arc<Delivery>,
ErrorShutdown,
CancellationToken,
) -> BoxFuture<'static, Result<JoinHandle<()>, BoxError>>
+ Send,
>;
pub struct RouterSink<B: Broker, State = ()> {
starters: Vec<BoundStarter<B, State>>,
handlers: Vec<HandlerMetadata>,
}
impl<B: Broker, State> fmt::Debug for RouterSink<B, State> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("RouterSink")
.field("handlers", &self.handlers.len())
.finish_non_exhaustive()
}
}
impl<B: Broker + 'static, State: Send + Sync + 'static> RouterSink<B, State> {
pub(crate) fn new() -> Self {
Self {
starters: Vec::new(),
handlers: Vec::new(),
}
}
pub(crate) fn push_handle<S, H, Cx>(
&mut self,
subscriber: S,
handler: H,
meta: HandlerMetadata,
policies: FailurePolicies,
) where
S: Subscriber + Send + 'static,
Cx: crate::BuildContext<S::Message> + Send + 'static,
H: Handler<S::Message, Cx, State> + 'static,
{
let handler = Arc::new(handler);
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |_connected, state, delivery, shutdown, token| {
Box::pin(async move {
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_dispatch(
subscriber, handler, token, name, state, delivery, failure,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn push_subscribe_batch<S, H>(
&mut self,
source: S,
handler: H,
meta: HandlerMetadata,
policies: FailurePolicies,
workers: Workers,
) where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
SourceMessage<B, S>: Send + 'static,
H: BatchHandler<SourceMessage<B, S>, State> + 'static,
{
let handler = Arc::new(handler);
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |connected: Arc<Connected<B>>, state, delivery, shutdown, token| {
Box::pin(async move {
let subscriber = source
.subscribe(connected.as_ref())
.await
.map_err(|e| Box::new(e) as BoxError)?;
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_batch_dispatch(
subscriber, handler, token, name, state, delivery, failure, workers,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn push_subscribe_workers<S, H, Cx>(
&mut self,
source: S,
handler: H,
meta: HandlerMetadata,
policies: FailurePolicies,
workers: Workers,
) where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: Send + 'static,
SourceMessage<B, S>: Send + Sync + 'static,
Cx: crate::BuildContext<SourceMessage<B, S>> + Send + 'static,
H: Handler<SourceMessage<B, S>, Cx, State> + 'static,
{
let handler = Arc::new(handler);
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |connected: Arc<Connected<B>>, state, delivery, shutdown, token| {
Box::pin(async move {
let subscriber = source
.subscribe(connected.as_ref())
.await
.map_err(|e| Box::new(e) as BoxError)?;
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_dispatch_workers(
subscriber, handler, token, name, state, delivery, failure, workers,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn push_subscribe<S, H, Cx>(
&mut self,
source: S,
handler: H,
meta: HandlerMetadata,
policies: FailurePolicies,
) where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: Send + 'static,
Cx: crate::BuildContext<SourceMessage<B, S>> + Send + 'static,
H: Handler<SourceMessage<B, S>, Cx, State> + 'static,
{
let handler = Arc::new(handler);
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |connected: Arc<Connected<B>>, state, delivery, shutdown, token| {
Box::pin(async move {
let subscriber = source
.subscribe(connected.as_ref())
.await
.map_err(|e| Box::new(e) as BoxError)?;
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_dispatch(
subscriber, handler, token, name, state, delivery, failure,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn push_raw(&mut self, starter: BoundStarter<B, State>, meta: HandlerMetadata) {
self.starters.push(starter);
self.handlers.push(meta);
}
pub(crate) fn push_injected_batch<Source, MakeHandler, HandlerFut, NewHandler>(
&mut self,
source: Source,
make_handler: MakeHandler,
meta: HandlerMetadata,
policies: FailurePolicies,
workers: Workers,
) where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: BatchSubscriber + Send + 'static,
SourceMessage<B, Source>: Send + 'static,
MakeHandler: FnOnce(Arc<Connected<B>>, Source::Subscriber) -> HandlerFut + Send + 'static,
HandlerFut: Future<Output = Result<(Source::Subscriber, NewHandler), BoxError>> + Send,
NewHandler: BatchHandler<SourceMessage<B, Source>, State> + 'static,
{
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |connected: Arc<Connected<B>>, state, delivery, shutdown, token| {
Box::pin(async move {
let subscriber = source
.subscribe(connected.as_ref())
.await
.map_err(|e| Box::new(e) as BoxError)?;
let (subscriber, handler) =
make_handler(Arc::clone(&connected), subscriber).await?;
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_batch_dispatch(
subscriber,
Arc::new(handler),
token,
name,
state,
delivery,
failure,
workers,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn push_injected_workers<Source, MakeHandler, HandlerFut, NewHandler, HandlerCx>(
&mut self,
source: Source,
make_handler: MakeHandler,
meta: HandlerMetadata,
policies: FailurePolicies,
workers: Workers,
) where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
SourceMessage<B, Source>: Send + Sync + 'static,
MakeHandler: FnOnce(Arc<Connected<B>>, Source::Subscriber) -> HandlerFut + Send + 'static,
HandlerFut: Future<Output = Result<(Source::Subscriber, NewHandler), BoxError>> + Send,
HandlerCx: crate::BuildContext<SourceMessage<B, Source>> + Send + 'static,
NewHandler: Handler<SourceMessage<B, Source>, HandlerCx, State> + 'static,
{
let name: Arc<str> = Arc::from(meta.name.as_ref());
self.starters.push(Box::new(
move |connected: Arc<Connected<B>>, state, delivery, shutdown, token| {
Box::pin(async move {
let subscriber = source
.subscribe(connected.as_ref())
.await
.map_err(|e| Box::new(e) as BoxError)?;
let (subscriber, handler) =
make_handler(Arc::clone(&connected), subscriber).await?;
let failure = DispatchFailure::new(policies, shutdown);
Ok(spawn_dispatch_workers(
subscriber,
Arc::new(handler),
token,
name,
state,
delivery,
failure,
workers,
))
})
},
));
self.handlers.push(meta);
}
pub(crate) fn into_parts(self) -> (Vec<BoundStarter<B, State>>, Vec<HandlerMetadata>) {
(self.starters, self.handlers)
}
}