use std::marker::PhantomData;
use serde::Serialize;
use serde::de::DeserializeOwned;
use crate::codec::Codec;
use crate::{BatchSubscriber, Broker, Connected, Subscriber, SubscriptionSource};
use crate::runtime::batch::{BatchDef, SliceHandler, TypedBatch, batch_metadata, typed_batch};
use crate::runtime::batch_publishing::{BatchPublishingDef, batch_publishing_metadata};
use crate::runtime::dispatch::Workers;
use crate::runtime::failure::FailurePolicies;
use crate::runtime::handler::Handler;
use crate::runtime::input::DecodeWith;
use crate::runtime::metadata::HandlerMetadata;
use crate::runtime::middleware::{BlanketLayer, Identity, Stack};
use crate::runtime::publish::{PublishPipeline, PublishTransform, TypedPublisher};
use crate::runtime::publishing::{PublishingDef, publishing_metadata};
use crate::runtime::subscriber_def::{SubscriberDef, subscriber_metadata};
use crate::runtime::typed::Typed;
use super::routes::{
BatchPublishingRoute, BatchRoute, HandleRoute, MountRoute, PublishingRoute, RouteMeta,
RouterDef, RouterHandlers, SubscribeRoute,
};
use super::sink::RouterSink;
use super::{
BatchPublishingRouter, IncludedBatchRouter, IncludedRouter, MergedRouter, PublishingRouter,
SourceMessage, SubscribedBatchRouter,
};
pub struct Router<B, Routes = (), C = (), Layers = Identity> {
pub(super) routes: Routes,
pub(super) codec: C,
pub(super) layers: Layers,
pub(super) _broker: PhantomData<fn() -> B>,
}
impl<B: Broker + 'static> Default for Router<B, (), (), Identity> {
fn default() -> Self {
Self {
routes: (),
codec: (),
layers: Identity,
_broker: PhantomData,
}
}
}
impl<B, Routes, C, Layers> std::fmt::Debug for Router<B, Routes, C, Layers> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Router").finish_non_exhaustive()
}
}
impl<B: Broker + 'static> Router<B, ()> {
#[must_use]
pub fn new() -> Self {
Self::default()
}
}
impl<B: Broker + 'static, Routes, RouteCodec, RouteLayers>
Router<B, Routes, RouteCodec, RouteLayers>
{
#[must_use]
pub fn with_codec<C>(self, codec: C) -> Router<B, Routes, C, RouteLayers> {
Router {
routes: self.routes,
codec,
layers: self.layers,
_broker: PhantomData,
}
}
#[must_use]
pub fn layer<N>(self, layer: N) -> Router<B, Routes, RouteCodec, Stack<N, RouteLayers>> {
Router {
routes: self.routes,
codec: self.codec,
layers: Stack::new(layer, self.layers),
_broker: PhantomData,
}
}
#[must_use]
pub fn merge<R2, C2, L2>(
self,
other: Router<B, R2, C2, L2>,
) -> MergedRouter<B, R2, C2, L2, RouteCodec, RouteLayers, Routes>
where
L2: BlanketLayer,
{
Router {
routes: (other, self.routes),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub fn handle<S, H>(
self,
subscriber: S,
handler: H,
meta: HandlerMetadata,
) -> Router<B, (HandleRoute<S, H>, Routes), RouteCodec, RouteLayers>
where
S: Subscriber + Send + 'static,
H: Handler<S::Message> + 'static,
{
Router {
routes: (
HandleRoute {
subscriber,
handler,
meta,
policies: FailurePolicies::default(),
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub fn subscribe<S, H>(
self,
source: S,
handler: H,
meta: HandlerMetadata,
) -> Router<B, (SubscribeRoute<S, H>, Routes), RouteCodec, RouteLayers>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: Send + 'static,
H: Handler<SourceMessage<B, S>> + 'static,
{
Router {
routes: (
SubscribeRoute {
source,
handler,
meta,
policies: FailurePolicies::default(),
workers: Workers::sequential(),
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub(super) fn push_batch_route<S, T, C, H>(
self,
source: S,
handler: H,
codec: C,
meta: HandlerMetadata,
) -> SubscribedBatchRouter<B, S, T, C, H, RouteCodec, RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
T: DeserializeOwned + Send + Sync + 'static,
C: Codec + 'static,
H: SliceHandler<T> + 'static,
{
Router {
routes: (
BatchRoute {
source,
handler: typed_batch(codec, handler),
meta,
policies: FailurePolicies::default(),
workers: Workers::sequential(),
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub(super) fn mount_subscriber<Source, Def, DecodeCodec>(
self,
source: Source,
def: Def,
codec: DecodeCodec,
) -> IncludedRouter<B, Source, Def, DecodeCodec, RouteCodec, RouteLayers, Routes>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: SubscriberDef,
Def::Input: DecodeWith<DecodeCodec>,
Def::Handler: 'static,
DecodeCodec: Send + Sync + 'static,
{
let meta = subscriber_metadata(source.name().to_owned(), &def);
let policies = def.failure_policies();
let workers = def.workers();
let handler = Typed::over(codec, def.into_handler()).on_decode_failure(policies.decode);
Router {
routes: (
SubscribeRoute {
source,
handler,
meta,
policies,
workers,
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub(super) fn mount_batch<Source, Def, DecodeCodec>(
self,
source: Source,
def: Def,
codec: DecodeCodec,
) -> IncludedBatchRouter<B, Source, Def, DecodeCodec, RouteCodec, RouteLayers, Routes>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchDef,
Def::Input: DecodeWith<DecodeCodec>,
Def::Handler: 'static,
DecodeCodec: Send + Sync + 'static,
{
let meta = batch_metadata(source.name().to_owned(), &def);
let policies = def.failure_policies();
let workers = def.workers();
let handler = TypedBatch::<_, Def::Input, _, _>::over(codec, def.into_handler())
.with_decode(policies.decode);
Router {
routes: (
BatchRoute {
source,
handler,
meta,
policies,
workers,
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub(super) fn mount_batch_publishing<Source, Def, DecodeCodec, ReplySource>(
self,
source: Source,
def: Def,
codec: DecodeCodec,
publisher: ReplySource,
) -> BatchPublishingRouter<
B,
Source,
Def,
DecodeCodec,
ReplySource,
RouteCodec,
RouteLayers,
Routes,
>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchPublishingDef + 'static,
Def::Input: DecodeWith<DecodeCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
DecodeCodec: Send + Sync + 'static,
ReplySource: 'static,
{
let meta = batch_publishing_metadata(source.name().to_owned(), &def);
let policies = def.failure_policies();
let workers = def.workers();
Router {
routes: (
BatchPublishingRoute {
source,
def,
codec,
publisher,
meta,
policies,
workers,
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
pub(super) fn mount_publishing<Source, Def, DecodeCodec, Leaf, ReplyCodec, Transforms>(
self,
source: Source,
def: Def,
codec: DecodeCodec,
publisher: TypedPublisher<Leaf, ReplyCodec, Transforms>,
) -> PublishingRouter<
B,
Source,
Def,
DecodeCodec,
Leaf,
ReplyCodec,
Transforms,
RouteCodec,
RouteLayers,
Routes,
>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: PublishingDef + 'static,
Def::Input: DecodeWith<DecodeCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
DecodeCodec: Codec + 'static,
Leaf: 'static,
ReplyCodec: Codec + 'static,
Transforms: PublishTransform<Def::Context> + 'static,
{
let meta = publishing_metadata(source.name().to_owned(), &def);
let policies = def.failure_policies();
let workers = def.workers();
Router {
routes: (
PublishingRoute {
source,
def,
codec,
publisher,
meta,
policies,
workers,
},
self.routes,
),
codec: self.codec,
layers: self.layers,
_broker: PhantomData,
}
}
}
impl<B, S, H, Routes, RouteCodec, RouteLayers>
Router<B, (SubscribeRoute<S, H>, Routes), RouteCodec, RouteLayers>
{
#[must_use]
pub fn workers(mut self, workers: Workers) -> Self {
self.routes.0.workers = workers;
self
}
}
impl<B, S, H, Routes, RouteCodec, RouteLayers>
Router<B, (BatchRoute<S, H>, Routes), RouteCodec, RouteLayers>
{
#[must_use]
pub fn workers(mut self, workers: Workers) -> Self {
self.routes.0.workers = workers;
self
}
}
impl<B: Broker + 'static, Routes: RouterHandlers, C, Layers> Router<B, Routes, C, Layers> {
#[must_use]
pub fn handlers(&self) -> Vec<HandlerMetadata> {
let mut out = Vec::new();
self.routes.collect_handlers(&mut out);
out
}
}
#[derive(Clone)]
struct ComposedBlanket<Outer, Inner> {
outer: Outer,
inner: Inner,
}
impl<Outer: BlanketLayer, Inner: BlanketLayer> BlanketLayer for ComposedBlanket<Outer, Inner> {
fn apply<M, C, S, H>(&self, handler: H) -> impl Handler<M, C, S> + 'static
where
M: Send + Sync + 'static,
C: Send + 'static,
S: Send + Sync + 'static,
H: Handler<M, C, S> + 'static,
{
self.outer
.apply::<M, C, S, _>(self.inner.apply::<M, C, S, _>(handler))
}
}
impl<B, Routes, C, Layers, State> RouterDef<B, State> for Router<B, Routes, C, Layers>
where
B: Broker + 'static,
Routes: RouterDef<B, State>,
Layers: BlanketLayer + Clone + Send + Sync + 'static,
{
fn mount<G, PP>(self, global: &G, pipeline: &PP, sink: &mut RouterSink<B, State>)
where
G: BlanketLayer + Clone + Send + Sync + 'static,
PP: PublishPipeline + Clone + Send + 'static,
{
let composed = ComposedBlanket {
outer: global.clone(),
inner: self.layers,
};
self.routes.mount(&composed, pipeline, sink);
}
}
impl<B, Routes, C, Layers> RouterHandlers for Router<B, Routes, C, Layers>
where
Routes: RouterHandlers,
{
fn collect_handlers(&self, out: &mut Vec<HandlerMetadata>) {
self.routes.collect_handlers(out);
}
}
impl<B, Routes, C, Layers, State> MountRoute<B, State> for Router<B, Routes, C, Layers>
where
B: Broker + 'static,
Routes: RouterDef<B, State>,
Layers: BlanketLayer + Clone + Send + Sync + 'static,
{
fn mount_one<G, PP>(self, global: &G, pipeline: &PP, sink: &mut RouterSink<B, State>)
where
G: BlanketLayer + Clone + Send + Sync + 'static,
PP: PublishPipeline + Clone + Send + 'static,
{
RouterDef::mount(self, global, pipeline, sink);
}
}
impl<B, Routes, C, Layers> RouteMeta for Router<B, Routes, C, Layers>
where
Routes: RouterHandlers,
{
fn collect(&self, out: &mut Vec<HandlerMetadata>) {
self.routes.collect_handlers(out);
}
}