use crate::stages::common::handlers::source::traits::{FiniteSourceHandler, InfiniteSourceHandler};
use crate::stages::common::handlers::{SinkHandler, TransformHandler};
use obzenflow_core::{SccId, StageId};
use serde::{Deserialize, Serialize};
use std::collections::HashSet;
use std::sync::Arc;
use super::MaxIterations;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CycleGuardConfig {
pub max_iterations: MaxIterations,
pub scc_id: SccId,
pub external_upstreams: HashSet<StageId>,
pub internal_upstreams: HashSet<StageId>,
pub is_entry_point: bool,
pub scc_internal_edges: Vec<(StageId, StageId)>,
}
pub enum StageHandlerType {
FiniteSource(Box<dyn FiniteSourceHandler>),
InfiniteSource(Box<dyn InfiniteSourceHandler>),
Transform(Box<dyn TransformHandler>),
Sink(Box<dyn SinkHandler>),
}
pub struct StageConfig {
pub stage_id: StageId,
pub name: String,
pub flow_name: String,
pub cycle_guard: Option<CycleGuardConfig>,
pub lineage: obzenflow_core::config::LineagePolicy,
pub effective_config: Arc<crate::runtime_config::FlowEffectiveConfig>,
}
pub struct ObserverConfig {
pub name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MiddlewareStackConfig {
pub stack: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub circuit_breaker: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rate_limiter: Option<serde_json::Value>,
}
impl MiddlewareStackConfig {
pub fn names_only(stack: Vec<String>) -> Self {
Self {
stack,
circuit_breaker: None,
rate_limiter: None,
}
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(crate) enum SourceContractStrictMode {
#[default]
Abort,
Warn,
}
impl SourceContractStrictMode {
pub(crate) fn from_token(token: &str) -> Self {
match token {
"warn" => SourceContractStrictMode::Warn,
_ => SourceContractStrictMode::Abort,
}
}
}