mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-08-23 00:00:33 +00:00
optimize skip failure query
This commit is contained in:
@@ -292,6 +292,16 @@ pub struct FlowModuleValueWithSkipFailures {
|
||||
pub parallelism: Option<u16>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct BranchWithSkipFailures {
|
||||
pub skip_failure: Option<bool>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct FlowModuleWithBranches {
|
||||
pub branches: Vec<BranchWithSkipFailures>,
|
||||
}
|
||||
|
||||
impl FlowModule {
|
||||
pub fn id_append(&mut self, s: &str) {
|
||||
self.id = format!("{}-{}", self.id, s);
|
||||
@@ -305,6 +315,11 @@ impl FlowModule {
|
||||
.map_err(crate::error::to_anyhow)
|
||||
}
|
||||
|
||||
pub fn get_branches_skip_failures(&self) -> anyhow::Result<FlowModuleWithBranches> {
|
||||
serde_json::from_str::<FlowModuleWithBranches>(self.value.get())
|
||||
.map_err(crate::error::to_anyhow)
|
||||
}
|
||||
|
||||
pub fn is_flow(&self) -> bool {
|
||||
self.get_type().is_ok_and(|x| x == "flow")
|
||||
}
|
||||
|
||||
@@ -369,12 +369,11 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
parallel,
|
||||
..
|
||||
} => compute_skip_branchall_failure(
|
||||
flow,
|
||||
job_id_for_status,
|
||||
old_status.step,
|
||||
*branch,
|
||||
*parallel,
|
||||
db,
|
||||
current_module.as_ref(),
|
||||
)
|
||||
.await?
|
||||
.unwrap_or(false),
|
||||
@@ -1205,12 +1204,11 @@ fn get_module(flow_job: &QueuedJob, module_step: &Step) -> Option<FlowModule> {
|
||||
}
|
||||
|
||||
async fn compute_skip_branchall_failure<'c>(
|
||||
flow: Uuid,
|
||||
job: &Uuid,
|
||||
step: i32,
|
||||
branch: usize,
|
||||
parallel: bool,
|
||||
db: &DB,
|
||||
flow_module: Option<&FlowModule>,
|
||||
) -> Result<Option<bool>, Error> {
|
||||
let branch = if parallel {
|
||||
sqlx::query_scalar!("SELECT script_path FROM completed_job WHERE id = $1", job)
|
||||
@@ -1234,22 +1232,13 @@ async fn compute_skip_branchall_failure<'c>(
|
||||
} else {
|
||||
branch as i32
|
||||
};
|
||||
sqlx::query_as(
|
||||
"SELECT (raw_flow->'modules'->$1->'value'->'branches'->$2->>'skip_failure')::bool
|
||||
FROM queue
|
||||
WHERE id = $3",
|
||||
)
|
||||
.bind(step)
|
||||
.bind(branch)
|
||||
.bind(flow)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.map(|(v,)| v)
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"error during retrieval of skip_loop_failures: {e:#}"
|
||||
))
|
||||
})
|
||||
Ok(flow_module
|
||||
.and_then(|x| x.get_branches_skip_failures().ok())
|
||||
.and_then(|x| {
|
||||
x.branches
|
||||
.get(branch as usize)
|
||||
.map(|x| x.skip_failure.unwrap_or(false))
|
||||
}))
|
||||
}
|
||||
|
||||
async fn has_failure_module<'c>(flow: Uuid, db: &DB) -> Result<bool, Error> {
|
||||
|
||||
Reference in New Issue
Block a user