mod builder;
mod controller;
mod driver;
mod runtime;
pub use builder::{BuildError, ChainCtx, Pipeline, PipelineError, SinkOptions};
pub use runtime::{PipelineRuntime, RuntimeOptions, ShutdownHandle, StartError, metrics_settings};
use crate::error::FatalError;
use crate::record::PartitionId;
use crate::sink::ShardQueues;
use crate::source::{DrainBarrier, LaneId};
use std::time::Instant;
pub(crate) enum ThreadControl<L> {
AddLane(L),
StopLanes {
lanes: Vec<LaneId>,
barrier: DrainBarrier,
deadline: Instant,
},
FlushNow,
DropLanes { lanes: Vec<LaneId> },
Shutdown {
barrier: DrainBarrier,
deadline: Instant,
},
}
#[derive(Debug)]
pub(crate) enum DriverEvent {
PauseLanes { lanes: Vec<LaneId> },
ResumeLanes { lanes: Vec<LaneId> },
Fatal { thread: usize, error: FatalError },
}
pub use crate::sink::DrainReport;
pub struct SinkRuntime {
pub queues: Vec<ShardQueues>,
pub drain: SinkDrainFn,
pub probe: Option<SinkProbeFn>,
}
impl std::fmt::Debug for SinkRuntime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SinkRuntime")
.field("queues", &self.queues)
.finish_non_exhaustive()
}
}
pub use crate::sink::{SinkDrainFn, SinkProbeFn};
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum ExitState {
Completed,
Failed(FatalErrorReport),
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct FatalErrorReport {
pub component: String,
pub reason: String,
}
impl std::fmt::Display for FatalErrorReport {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "pipeline failed in {}: {}", self.component, self.reason)
}
}
impl std::error::Error for FatalErrorReport {}
#[derive(Debug)]
#[non_exhaustive]
pub struct ExitReport {
pub state: ExitState,
pub sink_drain: Option<DrainReport>,
pub final_watermarks: Vec<(PartitionId, i64)>,
}
impl ExitReport {
pub fn log(&self) {
match &self.state {
ExitState::Completed => tracing::info!(
state = ?self.state,
drain = ?self.sink_drain,
watermarks = ?self.final_watermarks,
"pipeline finished"
),
ExitState::Failed(failure) => tracing::error!(
component = %failure.component,
reason = %failure.reason,
drain = ?self.sink_drain,
watermarks = ?self.final_watermarks,
"pipeline failed"
),
}
}
#[must_use]
pub fn exit_code(&self) -> i32 {
match self.state {
ExitState::Completed => 0,
ExitState::Failed(_) => 1,
}
}
pub fn ok(self) -> Result<ExitReport, FatalErrorReport> {
match &self.state {
ExitState::Completed => Ok(self),
ExitState::Failed(failure) => Err(failure.clone()),
}
}
}
#[cfg(all(test, not(loom)))]
pub(crate) mod fakes;
#[cfg(all(test, not(loom)))]
mod tests;