use crate::core::context_data::ContextData;
use crate::core::control::{PipelineControl, PipelineResult};
use crate::error::OrkaError; use crate::pipeline::definition::Pipeline; use tracing::{event, instrument, span, Level};
impl<TData, Err> Pipeline<TData, Err>
where
TData: 'static + Send + Sync,
Err: std::error::Error + From<OrkaError> + Send + Sync + 'static,
{
#[instrument(
name = "Pipeline::run",
skip_all,
fields(
pipeline_context_data_type = %std::any::type_name::<TData>(),
pipeline_error_type = %std::any::type_name::<Err>(),
num_steps = self.steps.len(),
),
err(Display)
)]
pub async fn run(&self, ctx_data: ContextData<TData>) -> Result<PipelineResult, Err>
{
event!(Level::DEBUG, "Pipeline execution starting.");
for (step_idx, step_def) in self.steps.iter().enumerate() {
let step_name_str = step_def.name.as_str();
let step_span = span!(
Level::INFO,
"pipeline_step_execution",
step_name = step_name_str,
step_index = step_idx,
optional = step_def.optional
);
let _step_span_guard = step_span.enter();
event!(Level::DEBUG, "Processing step.");
if let Some(skip_cond_fn) = &step_def.skip_if {
if skip_cond_fn(ctx_data.clone()) {
event!(Level::INFO, "Step skipped due to 'skip_if' condition.");
continue;
}
}
let has_before_handlers = self.before.get(step_name_str).is_some_and(|v| !v.is_empty());
let has_on_handlers = self.on.get(step_name_str).is_some_and(|v| !v.is_empty());
let has_after_handlers = self.after.get(step_name_str).is_some_and(|v| !v.is_empty());
if !has_before_handlers && !has_on_handlers && !has_after_handlers {
if step_def.optional {
event!(Level::DEBUG, "Optional step has no handlers, skipping.");
continue;
} else {
event!(Level::ERROR, "Non-optional step has no handlers.");
return Err(Err::from(OrkaError::HandlerMissing {
step_name: step_def.name.clone(),
}));
}
}
if let Some(handlers) = self.before.get(step_name_str) {
if !handlers.is_empty() {
event!(Level::TRACE, "Executing 'before' handlers.");
for (handler_idx, handler_fn) in handlers.iter().enumerate() {
let handler_span = span!(Level::DEBUG, "before_handler", handler_index = handler_idx);
let _handler_span_guard = handler_span.enter();
match handler_fn(ctx_data.clone()).await {
Ok(PipelineControl::Continue) => {}
Ok(PipelineControl::Stop) => {
event!(Level::INFO, "Pipeline stopped by a 'before' handler.");
return Ok(PipelineResult::Stopped);
}
Err(e) => {
event!(Level::ERROR, error = %e, "'before' handler failed.");
return Err(e);
}
}
}
}
}
if let Some(handlers) = self.on.get(step_name_str) {
if !handlers.is_empty() {
event!(Level::TRACE, "Executing 'on' handlers.");
for (handler_idx, handler_fn) in handlers.iter().enumerate() {
let handler_span = span!(Level::DEBUG, "on_handler", handler_index = handler_idx);
let _handler_span_guard = handler_span.enter();
match handler_fn(ctx_data.clone()).await {
Ok(PipelineControl::Continue) => {}
Ok(PipelineControl::Stop) => {
event!(Level::INFO, "Pipeline stopped by an 'on' handler.");
return Ok(PipelineResult::Stopped);
}
Err(e) => {
event!(Level::ERROR, error = %e, "'on' handler failed.");
return Err(e);
}
}
}
}
}
if let Some(handlers) = self.after.get(step_name_str) {
if !handlers.is_empty() {
event!(Level::TRACE, "Executing 'after' handlers.");
for (handler_idx, handler_fn) in handlers.iter().enumerate() {
let handler_span = span!(Level::DEBUG, "after_handler", handler_index = handler_idx);
let _handler_span_guard = handler_span.enter();
match handler_fn(ctx_data.clone()).await {
Ok(PipelineControl::Continue) => {}
Ok(PipelineControl::Stop) => {
event!(Level::INFO, "Pipeline stopped by an 'after' handler.");
return Ok(PipelineResult::Stopped);
}
Err(e) => {
event!(Level::ERROR, error = %e, "'after' handler failed.");
return Err(e);
}
}
}
}
}
event!(Level::DEBUG, "Step processing finished successfully.");
}
event!(Level::DEBUG, "Pipeline execution completed successfully.");
Ok(PipelineResult::Completed)
}
}