#![doc = include_str!("../README.md")]
pub mod backpressure;
pub mod bootstrap;
pub mod contracts;
pub mod control_plane;
pub mod effects;
pub mod errors;
pub mod execution;
pub mod feed_plan;
pub mod id_conversions;
pub mod message_bus;
pub mod replay;
pub mod run_context;
pub mod runtime_config;
pub mod runtime_resource_limits;
pub mod supervised_base;
pub mod typing;
pub mod messaging;
pub mod metrics;
pub mod pipeline;
pub mod stages;
#[doc(hidden)]
pub mod __private {
pub mod lifecycle;
pub use crate::stages::common::handlers::join::{
ErasedJoinInvocation, TypedJoinHandlerAdapter, UnifiedJoinHandler,
};
pub use crate::stages::common::handlers::sink::SinkWriterAdapter;
pub use crate::stages::common::handlers::source::typed::{
TypedAsyncFiniteSourceHandlerAdapter, TypedAsyncInfiniteSourceHandlerAdapter,
TypedFiniteSourceHandlerAdapter, TypedInfiniteSourceHandlerAdapter,
};
pub use crate::stages::common::handlers::source::{
ErasedSourceCompletion, ErasedSourceInvocation, ErasedSourceOutcome,
UnifiedAsyncFiniteSourceHandler, UnifiedAsyncInfiniteSourceHandler,
UnifiedFiniteSourceHandler, UnifiedInfiniteSourceHandler,
};
pub use obzenflow_core::event::schema::{EmptySet, WithMember};
}
#[cfg(any(test, feature = "test-support"))]
pub mod testing;
pub mod prelude {
pub use crate::errors::{FlowError, MessageBusError, PipelineSupervisorError, RuntimeResult};
pub use crate::pipeline::{
FlowHandle, FlowStopMode, ObserverConfig, PipelineBuilder, PipelineControl,
PipelineStageConfig, PipelineState,
};
pub use crate::message_bus::{FsmMessageBus, StageCommand};
pub use crate::stages::{
EffectfulStatefulHandler, EffectfulTransformHandler, HostedIngressSource, InferenceHandler,
IngressDecodeError, IngressDecoder, InlineSink, ObserverHandler, ResourceManaged,
SinkConnector, SinkDescription, SinkWriter, SourceError, SourceObservationSink,
TypedAsyncFiniteSourceHandler, TypedAsyncInfiniteSourceHandler, TypedFiniteSourceHandler,
TypedInfiniteSourceHandler, TypedJoinHandler, TypedStatefulHandler, TypedTransformHandler,
};
pub use crate::typing::{SourceTyping, TransformTyping};
pub use crate::effects::{
DomainFacts, Effect, EffectCommitHandle, EffectContext, EffectDeclaration, EffectError,
EffectOutcomeKind, EffectOutcomePayload, EffectSafety, Effects, IdempotencyKey,
RecordedReply, SinkRedeliverySafety, TransactionalEffectPort,
};
pub use crate::messaging::UpstreamSubscription;
}