use serde::Serialize;
use serde::de::DeserializeOwned;
use crate::codec::Codec;
use crate::{BatchSubscriber, Broker, Connected, SubscriptionSource};
use crate::runtime::batch::{BatchDef, SliceHandler};
use crate::runtime::batch_publishing::BatchPublishingDef;
use crate::runtime::input::{DecodeWith, RawBytes};
use crate::runtime::metadata::HandlerMetadata;
use crate::runtime::publish::{PublishTransform, ReplyWiring, TypedPublisher};
use crate::runtime::publishing::PublishingDef;
use crate::runtime::subscriber_def::SubscriberDef;
use super::builder::Router;
use super::{
BatchPublishingRouter, IncludedBatchRouter, IncludedRouter, PublishingRouter,
SubscribedBatchRouter,
};
impl<B: Broker + 'static, Routes, RouteLayers> Router<B, Routes, (), RouteLayers> {
#[cfg(any(feature = "json", feature = "cbor", feature = "msgpack"))]
pub fn include<Def>(
self,
def: Def,
) -> IncludedRouter<B, Def::Source, Def, crate::codec::DefaultCodec, (), RouteLayers, Routes>
where
Def: SubscriberDef,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber: Send + 'static,
Def::Input: DecodeWith<crate::codec::DefaultCodec>,
Def::Handler: 'static,
{
let source = def.source();
self.mount_subscriber(source, def, crate::codec::DefaultCodec::default())
}
#[cfg(any(feature = "json", feature = "cbor", feature = "msgpack"))]
pub fn include_on<Source, Def>(
self,
source: Source,
def: Def,
) -> IncludedRouter<B, Source, Def, crate::codec::DefaultCodec, (), RouteLayers, Routes>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: SubscriberDef,
Def::Input: DecodeWith<crate::codec::DefaultCodec>,
Def::Handler: 'static,
{
self.mount_subscriber(source, def, crate::codec::DefaultCodec::default())
}
#[cfg(any(feature = "json", feature = "cbor", feature = "msgpack"))]
pub fn include_batch<Def>(
self,
def: Def,
) -> IncludedBatchRouter<B, Def::Source, Def, crate::codec::DefaultCodec, (), RouteLayers, Routes>
where
Def: BatchDef,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber:
BatchSubscriber + Send + 'static,
Def::Input: DecodeWith<crate::codec::DefaultCodec>,
Def::Handler: 'static,
{
let source = def.source();
self.mount_batch(source, def, crate::codec::DefaultCodec::default())
}
#[cfg(any(feature = "json", feature = "cbor", feature = "msgpack"))]
pub fn include_batch_on<S, Def>(
self,
source: S,
def: Def,
) -> IncludedBatchRouter<B, S, Def, crate::codec::DefaultCodec, (), RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchDef,
Def::Input: DecodeWith<crate::codec::DefaultCodec>,
Def::Handler: 'static,
{
self.mount_batch(source, def, crate::codec::DefaultCodec::default())
}
#[cfg(any(feature = "json", feature = "cbor", feature = "msgpack"))]
pub fn subscribe_batch<S, T, H>(
self,
source: S,
handler: H,
meta: HandlerMetadata,
) -> SubscribedBatchRouter<B, S, T, crate::codec::DefaultCodec, H, (), RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
T: DeserializeOwned + Send + Sync + 'static,
H: SliceHandler<T> + 'static,
{
self.push_batch_route(source, handler, crate::codec::DefaultCodec::default(), meta)
}
pub fn include_batch_publishing<Def, RP>(
self,
def: Def,
publisher: RP,
) -> BatchPublishingRouter<B, Def::Source, Def, RP::Codec, RP, (), RouteLayers, Routes>
where
Def: BatchPublishingDef + 'static,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber:
BatchSubscriber + Send + 'static,
Def::Input: DecodeWith<RP::Codec>,
Def::Reply: Serialize + Send + Sync + 'static,
RP: ReplyWiring + 'static,
{
let codec = publisher.decode_codec().clone();
let source = def.source();
self.mount_batch_publishing(source, def, codec, publisher)
}
pub fn include_batch_publishing_on<S, Def, RP>(
self,
source: S,
def: Def,
publisher: RP,
) -> BatchPublishingRouter<B, S, Def, RP::Codec, RP, (), RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchPublishingDef + 'static,
Def::Input: DecodeWith<RP::Codec>,
Def::Reply: Serialize + Send + Sync + 'static,
RP: ReplyWiring + 'static,
{
let codec = publisher.decode_codec().clone();
self.mount_batch_publishing(source, def, codec, publisher)
}
pub fn include_publishing<Def, Leaf, ReplyCodec, Transforms>(
self,
def: Def,
publisher: TypedPublisher<Leaf, ReplyCodec, Transforms>,
) -> PublishingRouter<
B,
Def::Source,
Def,
ReplyCodec,
Leaf,
ReplyCodec,
Transforms,
(),
RouteLayers,
Routes,
>
where
Def: PublishingDef + 'static,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber: Send + 'static,
Def::Input: DecodeWith<ReplyCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
Leaf: 'static,
ReplyCodec: Codec + Clone + 'static,
Transforms: PublishTransform<Def::Context> + 'static,
{
let codec = publisher.codec().clone();
let source = def.source();
self.mount_publishing(source, def, codec, publisher)
}
pub fn include_publishing_on<Source, Def, Leaf, ReplyCodec, Transforms>(
self,
source: Source,
def: Def,
publisher: TypedPublisher<Leaf, ReplyCodec, Transforms>,
) -> PublishingRouter<
B,
Source,
Def,
ReplyCodec,
Leaf,
ReplyCodec,
Transforms,
(),
RouteLayers,
Routes,
>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: PublishingDef + 'static,
Def::Input: DecodeWith<ReplyCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
Leaf: 'static,
ReplyCodec: Codec + Clone + 'static,
Transforms: PublishTransform<Def::Context> + 'static,
{
let codec = publisher.codec().clone();
self.mount_publishing(source, def, codec, publisher)
}
}
impl<B: Broker + 'static, Routes, RouteCodec, RouteLayers>
Router<B, Routes, RouteCodec, RouteLayers>
{
#[must_use]
pub fn include_raw<Def>(
self,
def: Def,
) -> IncludedRouter<B, Def::Source, Def, (), RouteCodec, RouteLayers, Routes>
where
Def: SubscriberDef<Input = RawBytes>,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber: Send + 'static,
Def::Handler: 'static,
{
let source = def.source();
self.mount_subscriber(source, def, ())
}
}
impl<B: Broker + 'static, Routes, RouteCodec: Codec + Clone + 'static, RouteLayers>
Router<B, Routes, RouteCodec, RouteLayers>
{
pub fn include<Def>(
self,
def: Def,
) -> IncludedRouter<B, Def::Source, Def, RouteCodec, RouteCodec, RouteLayers, Routes>
where
Def: SubscriberDef,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber: Send + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Handler: 'static,
{
let codec = self.codec.clone();
let source = def.source();
self.mount_subscriber(source, def, codec)
}
pub fn include_on<Source, Def>(
self,
source: Source,
def: Def,
) -> IncludedRouter<B, Source, Def, RouteCodec, RouteCodec, RouteLayers, Routes>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: SubscriberDef,
Def::Input: DecodeWith<RouteCodec>,
Def::Handler: 'static,
{
let codec = self.codec.clone();
self.mount_subscriber(source, def, codec)
}
pub fn include_batch<Def>(
self,
def: Def,
) -> IncludedBatchRouter<B, Def::Source, Def, RouteCodec, RouteCodec, RouteLayers, Routes>
where
Def: BatchDef,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber:
BatchSubscriber + Send + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Handler: 'static,
{
let codec = self.codec.clone();
let source = def.source();
self.mount_batch(source, def, codec)
}
pub fn include_batch_on<S, Def>(
self,
source: S,
def: Def,
) -> IncludedBatchRouter<B, S, Def, RouteCodec, RouteCodec, RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchDef,
Def::Input: DecodeWith<RouteCodec>,
Def::Handler: 'static,
{
let codec = self.codec.clone();
self.mount_batch(source, def, codec)
}
pub fn subscribe_batch<S, T, H>(
self,
source: S,
handler: H,
meta: HandlerMetadata,
) -> SubscribedBatchRouter<B, S, T, RouteCodec, H, RouteCodec, RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
T: DeserializeOwned + Send + Sync + 'static,
H: SliceHandler<T> + 'static,
{
let codec = self.codec.clone();
self.push_batch_route(source, handler, codec, meta)
}
pub fn include_batch_publishing<Def, RP>(
self,
def: Def,
publisher: RP,
) -> BatchPublishingRouter<B, Def::Source, Def, RouteCodec, RP, RouteCodec, RouteLayers, Routes>
where
Def: BatchPublishingDef + 'static,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber:
BatchSubscriber + Send + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
RP: 'static,
{
let codec = self.codec.clone();
let source = def.source();
self.mount_batch_publishing(source, def, codec, publisher)
}
pub fn include_batch_publishing_on<S, Def, RP>(
self,
source: S,
def: Def,
publisher: RP,
) -> BatchPublishingRouter<B, S, Def, RouteCodec, RP, RouteCodec, RouteLayers, Routes>
where
S: SubscriptionSource<Connected<B>> + Send + 'static,
S::Subscriber: BatchSubscriber + Send + 'static,
Def: BatchPublishingDef + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
RP: 'static,
{
let codec = self.codec.clone();
self.mount_batch_publishing(source, def, codec, publisher)
}
pub fn include_publishing<Def, Leaf, ReplyCodec, Transforms>(
self,
def: Def,
publisher: TypedPublisher<Leaf, ReplyCodec, Transforms>,
) -> PublishingRouter<
B,
Def::Source,
Def,
RouteCodec,
Leaf,
ReplyCodec,
Transforms,
RouteCodec,
RouteLayers,
Routes,
>
where
Def: PublishingDef + 'static,
Def::Source: SubscriptionSource<Connected<B>> + Send + 'static,
<Def::Source as SubscriptionSource<Connected<B>>>::Subscriber: Send + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
Leaf: 'static,
ReplyCodec: Codec + 'static,
Transforms: PublishTransform<Def::Context> + 'static,
{
let codec = self.codec.clone();
let source = def.source();
self.mount_publishing(source, def, codec, publisher)
}
pub fn include_publishing_on<Source, Def, Leaf, ReplyCodec, Transforms>(
self,
source: Source,
def: Def,
publisher: TypedPublisher<Leaf, ReplyCodec, Transforms>,
) -> PublishingRouter<
B,
Source,
Def,
RouteCodec,
Leaf,
ReplyCodec,
Transforms,
RouteCodec,
RouteLayers,
Routes,
>
where
Source: SubscriptionSource<Connected<B>> + Send + 'static,
Source::Subscriber: Send + 'static,
Def: PublishingDef + 'static,
Def::Input: DecodeWith<RouteCodec>,
Def::Reply: Serialize + Send + Sync + 'static,
Leaf: 'static,
ReplyCodec: Codec + 'static,
Transforms: PublishTransform<Def::Context> + 'static,
{
let codec = self.codec.clone();
self.mount_publishing(source, def, codec, publisher)
}
}