#![forbid(unsafe_code)]
mod broker;
mod buffered;
mod capability;
mod error;
mod field;
mod headers;
mod message;
mod publisher;
mod schema;
mod subscriber;
mod subscription;
pub mod testing;
#[cfg(feature = "testing")]
#[doc(hidden)]
pub use inventory;
pub use broker::{Broker, Connected, ConnectedBroker};
pub use buffered::{Buffered, BufferedSubscriber};
pub use capability::{
ApiKeyLocation, BatchSubscriber, DescribeServer, HttpApiKeyLocation, OwnedTransactions,
Partitioned, Positioned, RequestReply, SecurityScheme, Seekable, Seeker, ServerSpec, Subscribe,
Transaction, TransactionalPublisher,
};
pub use error::AckError;
pub use field::{BuildContext, ContextField, Field, FieldMut};
pub use headers::Headers;
pub use message::{IncomingMessage, OutgoingMessage, RawMessage};
pub use publisher::{DefaultPublish, PairError, PublishPolicy, Publisher};
pub use schema::Message;
pub use subscriber::Subscriber;
pub use subscription::{
Name, SeekerPendingError, SeekerToken, StartAt, SubscriptionSource, WithSeeker,
};
pub mod codec;
#[cfg(feature = "memory")]
pub mod memory;
pub mod runtime;
pub use runtime::RustStream;
#[cfg(feature = "macros")]
pub use ruststream_macros::subscriber;
#[cfg(feature = "macros")]
pub use ruststream_macros::app;
#[cfg(feature = "macros")]
pub use ruststream_macros::Message;
#[cfg(feature = "macros")]
pub use ruststream_macros::FromRef;
#[cfg(feature = "conformance")]
pub mod conformance;
#[cfg(feature = "asyncapi")]
pub mod asyncapi;
#[cfg(feature = "asyncapi")]
pub use schemars;
#[cfg(feature = "metrics")]
pub mod metrics;
#[cfg(feature = "logging")]
pub mod logging;
#[cfg(feature = "otel")]
pub mod otel;
#[doc(hidden)]
pub mod __private {
use core::marker::PhantomData;
#[derive(Debug)]
pub struct Probe<T>(pub PhantomData<T>);
impl<T> Probe<T> {
#[must_use]
pub const fn new() -> Self {
Self(PhantomData)
}
}
impl<T> Default for Probe<T> {
fn default() -> Self {
Self::new()
}
}
pub trait NoSchemaProbe {
fn schema_json(&self) -> Option<String>;
}
impl<T> NoSchemaProbe for Probe<T> {
fn schema_json(&self) -> Option<String> {
None
}
}
#[cfg(feature = "asyncapi")]
impl<T: schemars::JsonSchema> Probe<T> {
#[must_use]
pub fn schema_json(&self) -> Option<String> {
serde_json::to_string(&schemars::schema_for!(T)).ok()
}
}
pub trait NoMessageProbe {
fn message_name(&self) -> Option<&'static str>;
fn message_description(&self) -> Option<&'static str>;
}
impl<T> NoMessageProbe for Probe<T> {
fn message_name(&self) -> Option<&'static str> {
None
}
fn message_description(&self) -> Option<&'static str> {
None
}
}
impl<T: crate::Message> Probe<T> {
#[must_use]
pub fn message_name(&self) -> Option<&'static str> {
Some(T::NAME)
}
#[must_use]
pub fn message_description(&self) -> Option<&'static str> {
T::DESCRIPTION
}
}
}
#[macro_export]
macro_rules! nonzero {
($value:expr) => {
const {
match ::core::num::NonZero::new($value) {
::core::option::Option::Some(value) => value,
::core::option::Option::None => panic!("nonzero!(..) requires a non-zero value"),
}
}
};
}