use crate::conditional::provider::{FunctionalPipelineProvider, PipelineProvider, StaticPipelineProvider};
use crate::conditional::scope::{AnyConditionalScope, ConditionalScope};
use crate::core::context::{ExtractorFn, Handler, MergeFn};
use crate::core::context_data::ContextData;
use crate::core::control::PipelineControl;
use crate::error::OrkaError;
use crate::pipeline::Pipeline;
use std::future::Future;
use std::marker::PhantomData;
use std::sync::Arc;
use tracing::{event, instrument, Level};
#[must_use = "conditional scopes are only applied when you call .finalize_conditional_step()"]
pub struct ConditionalScopeBuilder<'pipeline, TData, Err>
where
TData: 'static + Send + Sync,
Err: std::error::Error + From<OrkaError> + Send + Sync + 'static,
{
pipeline: &'pipeline mut Pipeline<TData, Err>,
step_name: String,
collected_scopes: Vec<Arc<dyn AnyConditionalScope<TData, Err>>>,
on_no_match_behavior: PipelineControl,
}
impl<'pipeline, TData, Err> ConditionalScopeBuilder<'pipeline, TData, Err>
where
TData: 'static + Send + Sync,
Err: std::error::Error + From<OrkaError> + Send + Sync + 'static,
{
pub(crate) fn new(pipeline: &'pipeline mut Pipeline<TData, Err>, step_name: String) -> Self {
if !pipeline.steps.iter().any(|s| s.name == step_name) {
pipeline.steps.push(crate::core::step::StepDef {
name: step_name.clone(),
optional: false,
skip_if: None,
});
}
pipeline.pending_conditional.insert(step_name.clone());
Self {
pipeline,
step_name,
collected_scopes: Vec::new(),
on_no_match_behavior: PipelineControl::Continue,
}
}
pub fn add_static_scope<SData>(
self,
static_pipeline: Arc<Pipeline<SData, Err>>, extractor_fn: impl Fn(ContextData<TData>) -> Result<ContextData<SData>, OrkaError> + Send + Sync + 'static,
) -> ConditionalScopeConfigurator<'pipeline, TData, SData, Err, StaticPipelineProvider<SData, Err>>
where
SData: 'static + Send + Sync,
{
ConditionalScopeConfigurator {
builder: self,
provider: Arc::new(StaticPipelineProvider::new(static_pipeline)),
extractor: Arc::new(extractor_fn),
merge: None,
_phantom_sdata: PhantomData,
}
}
pub fn add_dynamic_scope<SData, F, Fut>(
self,
pipeline_factory: F,
extractor_fn: impl Fn(ContextData<TData>) -> Result<ContextData<SData>, OrkaError> + Send + Sync + 'static,
) -> ConditionalScopeConfigurator<'pipeline, TData, SData, Err, FunctionalPipelineProvider<TData, SData, Err, F, Fut>>
where
SData: 'static + Send + Sync,
F: Fn(ContextData<TData>) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<Arc<Pipeline<SData, Err>>, OrkaError>> + Send + 'static,
{
ConditionalScopeConfigurator {
builder: self,
provider: Arc::new(FunctionalPipelineProvider::new(pipeline_factory)),
extractor: Arc::new(extractor_fn),
merge: None,
_phantom_sdata: PhantomData,
}
}
pub fn if_no_scope_matches(mut self, behavior: PipelineControl) -> Self {
self.on_no_match_behavior = behavior;
self
}
#[instrument(
name = "ConditionalScopeBuilder::finalize_conditional_step",
skip_all,
fields(step_name = %self.step_name, num_scopes = self.collected_scopes.len())
)]
pub fn finalize_conditional_step(self, optional_for_main_step: bool) {
let step_name_captured = self.step_name.clone();
let scopes_for_closure_capture = Arc::new(self.collected_scopes); let on_no_match_behavior_captured = self.on_no_match_behavior;
let master_handler: Handler<TData, Err> = Box::new(move |main_ctx_data: ContextData<TData>| {
let scopes_to_check = scopes_for_closure_capture.clone();
let step_name_log_ctx = step_name_captured.clone();
let current_main_ctx_data = main_ctx_data.clone();
let is_step_optional_captured = optional_for_main_step;
Box::pin(async move {
for scope_candidate in scopes_to_check.iter() {
if scope_candidate.is_condition_met(current_main_ctx_data.clone()) {
event!(Level::DEBUG, step_name = %step_name_log_ctx, "Conditional scope matched. Executing.");
match scope_candidate
.execute_scoped_pipeline(current_main_ctx_data.clone())
.await
{
Ok(control) => return Ok(control),
Err(e) => {
event!(Level::ERROR, step_name = %step_name_log_ctx, error = %e, "Error during conditional scope execution.");
if is_step_optional_captured {
event!(Level::WARN, step_name = %step_name_log_ctx, "Conditional step is optional, swallowing error and continuing main pipeline.");
return Ok(PipelineControl::Continue); } else {
return Err(e); }
}
}
}
}
event!(Level::DEBUG, step_name = %step_name_log_ctx, "No conditional scope matched. Defaulting to {:?}.", on_no_match_behavior_captured);
Ok(on_no_match_behavior_captured)
})
});
if let Some(step_def) = self.pipeline.steps.iter_mut().find(|s| s.name == self.step_name) {
step_def.optional = optional_for_main_step;
} else {
event!(Level::WARN, step_name = %self.step_name, "Step definition not found during finalize_conditional_step. This may indicate an internal issue.");
}
self.pipeline.pending_conditional.remove(&self.step_name);
self
.pipeline
.on
.entry(self.step_name.clone())
.or_default()
.push(master_handler);
event!(Level::INFO, step_name = %self.step_name, "Conditional scopes finalized and master handler registered.");
}
}
#[must_use = "this scope is only registered when you call .on_condition(), and applied when you call .finalize_conditional_step()"]
pub struct ConditionalScopeConfigurator<
'pipeline,
TData: 'static + Send + Sync,
SData: 'static + Send + Sync,
Err: std::error::Error + From<OrkaError> + Send + Sync + 'static,
P: PipelineProvider<TData, SData, Err> + 'static, > {
builder: ConditionalScopeBuilder<'pipeline, TData, Err>,
provider: Arc<P>,
extractor: ExtractorFn<TData, SData>,
merge: Option<MergeFn<TData, SData>>,
_phantom_sdata: PhantomData<SData>, }
impl<'pipeline, TData, SData, Err, P> ConditionalScopeConfigurator<'pipeline, TData, SData, Err, P>
where
TData: 'static + Send + Sync,
SData: 'static + Send + Sync,
Err: std::error::Error + From<OrkaError> + Send + Sync + 'static,
P: PipelineProvider<TData, SData, Err> + 'static,
{
pub fn with_merge(mut self, merge_fn: impl Fn(&mut TData, &SData) + Send + Sync + 'static) -> Self {
self.merge = Some(Arc::new(merge_fn));
self
}
#[instrument(
name = "ConditionalScopeConfigurator::on_condition",
skip_all,
fields(builder_step_name = %self.builder.step_name)
)]
pub fn on_condition(
mut self,
condition_fn: impl Fn(ContextData<TData>) -> bool + Send + Sync + 'static,
) -> ConditionalScopeBuilder<'pipeline, TData, Err> {
let final_scope_definition = ConditionalScope::<TData, SData, Err> {
pipeline_provider: self.provider, extractor: self.extractor,
condition: Arc::new(condition_fn),
merge: self.merge,
_phantom_main_err: PhantomData, };
event!(Level::DEBUG, "Conditional scope configured with condition.");
self.builder.collected_scopes.push(Arc::new(final_scope_definition));
self.builder
}
}