use obzenflow_core::StageId;
use obzenflow_fsm::StateVariant;
use std::time::Duration;
#[derive(Clone, Debug)]
pub enum FlowStopMode {
Cancel,
Graceful { timeout: Duration },
}
#[derive(Clone, Debug, PartialEq)]
pub enum PipelineState {
Created,
Materializing,
Materialized,
ReadyForRun,
Running,
SourceCompleted,
AbortRequested {
reason: obzenflow_core::event::types::ViolationCause,
upstream: Option<StageId>,
},
Draining,
Drained,
Failed {
reason: String,
failure_cause: Option<obzenflow_core::event::types::ViolationCause>,
},
}
impl StateVariant for PipelineState {
fn variant_name(&self) -> &str {
match self {
PipelineState::Created => "Created",
PipelineState::Materializing => "Materializing",
PipelineState::Materialized => "Materialized",
PipelineState::ReadyForRun => "ReadyForRun",
PipelineState::Running => "Running",
PipelineState::SourceCompleted => "SourceCompleted",
PipelineState::AbortRequested { .. } => "AbortRequested",
PipelineState::Draining => "Draining",
PipelineState::Drained => "Drained",
PipelineState::Failed { .. } => "Failed",
}
}
}
impl PipelineState {
pub fn is_terminal(&self) -> bool {
matches!(self, PipelineState::Drained | PipelineState::Failed { .. })
}
}
#[non_exhaustive]
#[derive(Clone, Debug)]
pub enum PipelineControl {
Start,
Stop { mode: FlowStopMode },
Abort { reason: String },
}
#[derive(Debug, Clone, PartialEq)]
pub enum FlowStartControlOutcome {
Submitted { observed_state: PipelineState },
AlreadyRunning { state: PipelineState },
Rejected {
state: PipelineState,
reason: &'static str,
},
}