use crate::errors::FlowError;
use crate::pipeline::FlowHandle;
use crate::stages::common::stage_handle::StageHandle;
use crate::supervised_base::handle::ExecutionCancellation;
use std::sync::Arc;
pub struct ExecutionGuard {
supervisor: Option<ExecutionCancellation>,
stages: Vec<Arc<dyn StageHandle>>,
metrics: Arc<crate::pipeline::resources::MetricsOwner>,
}
impl ExecutionGuard {
pub(crate) fn new(
supervisor: ExecutionCancellation,
stages: Vec<Arc<dyn StageHandle>>,
metrics: Arc<crate::pipeline::resources::MetricsOwner>,
) -> Self {
Self {
supervisor: Some(supervisor),
stages,
metrics,
}
}
pub fn disarm(mut self) {
self.supervisor = None;
}
}
impl Drop for ExecutionGuard {
fn drop(&mut self) {
if let Some(supervisor) = &self.supervisor {
supervisor.abort();
self.metrics.request_abort();
for stage in &self.stages {
stage.request_abort();
}
}
}
}
pub fn guard_execution(flow: &FlowHandle) -> ExecutionGuard {
flow.execution_guard()
}
pub async fn wait(flow: &FlowHandle) -> Result<(), FlowError> {
flow.wait_for_resources().await
}