diff --git a/backend/.sqlx/query-b0ee9e084ae867700f1964b3fb56e9e55d0d15c90fb9244066ac4d409f1b7e1b.json b/backend/.sqlx/query-b0ee9e084ae867700f1964b3fb56e9e55d0d15c90fb9244066ac4d409f1b7e1b.json new file mode 100644 index 0000000000..40c6eea515 --- /dev/null +++ b/backend/.sqlx/query-b0ee9e084ae867700f1964b3fb56e9e55d0d15c90fb9244066ac4d409f1b7e1b.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "\nDELETE FROM flow_iterator_data\nWHERE job_id NOT IN (SELECT id FROM v2_job_queue)\nRETURNING job_id\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "job_id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false + ] + }, + "hash": "b0ee9e084ae867700f1964b3fb56e9e55d0d15c90fb9244066ac4d409f1b7e1b" +} diff --git a/backend/migrations/20251218220001_remove_flow_iterator_data_fk.down.sql b/backend/migrations/20251218220001_remove_flow_iterator_data_fk.down.sql new file mode 100644 index 0000000000..4f1c17cb01 --- /dev/null +++ b/backend/migrations/20251218220001_remove_flow_iterator_data_fk.down.sql @@ -0,0 +1,8 @@ +-- Re-add created_at column +ALTER TABLE flow_iterator_data +ADD COLUMN created_at TIMESTAMP WITH TIME ZONE DEFAULT now(); + +-- Re-add foreign key constraint +ALTER TABLE flow_iterator_data +ADD CONSTRAINT flow_iterator_data_job_id_fkey +FOREIGN KEY (job_id) REFERENCES v2_job_queue (id) ON DELETE CASCADE; diff --git a/backend/migrations/20251218220001_remove_flow_iterator_data_fk.up.sql b/backend/migrations/20251218220001_remove_flow_iterator_data_fk.up.sql new file mode 100644 index 0000000000..d1709f27e7 --- /dev/null +++ b/backend/migrations/20251218220001_remove_flow_iterator_data_fk.up.sql @@ -0,0 +1,7 @@ +-- Remove foreign key constraint - we'll handle cleanup manually in monitor.rs +ALTER TABLE flow_iterator_data +DROP CONSTRAINT IF EXISTS flow_iterator_data_job_id_fkey; + +-- Remove created_at column as it's not needed +ALTER TABLE flow_iterator_data +DROP COLUMN IF EXISTS created_at; diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 1a53418d28..51ab25fee9 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1620,6 +1620,16 @@ pub async fn monitor_db( } }; + let cleanup_flow_iterator_data_f = async { + if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(10) { + if let Some(db) = conn.as_sql() { + if let Err(e) = cleanup_flow_iterator_data_orphaned_jobs(&db).await { + tracing::error!("Error cleaning up flow_iterator_data: {:?}", e); + } + } + } + }; + // run every hour (60 minutes / 30 seconds = 120) let cleanup_worker_group_stats_f = async { if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(120) { @@ -1741,6 +1751,7 @@ pub async fn monitor_db( cleanup_concurrency_counters_empty_keys_f, cleanup_debounce_keys_f, cleanup_debounce_keys_completed_f, + cleanup_flow_iterator_data_f, cleanup_worker_group_stats_f, ); } @@ -2839,3 +2850,23 @@ RETURNING key,job_id } Ok(()) } + +async fn cleanup_flow_iterator_data_orphaned_jobs(db: &DB) -> error::Result<()> { + let result = sqlx::query!( + " +DELETE FROM flow_iterator_data +WHERE job_id NOT IN (SELECT id FROM v2_job_queue) +RETURNING job_id + ", + ) + .fetch_all(db) + .await?; + + if result.len() > 0 { + tracing::info!( + "Cleaned up {} orphaned flow_iterator_data rows", + result.len() + ); + } + Ok(()) +} diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 07b9cc3315..c15c5068b4 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -601,10 +601,7 @@ pub async fn update_flow_status_after_job_completion_internal( .. } if *parallel => { let (nindex, len) = match (iterator, branchall) { - (Some(FlowIterator { itered, .. }), _) => { - // Read itered from separate table or fallback to JSONB - let itered = read_itered_from_db(db, flow, itered).await?; - + (Some(FlowIterator { itered, itered_len, .. }), _) => { let position = if flow_jobs_success.is_some() { find_flow_job_index(jobs, job_id_for_status) } else { @@ -663,12 +660,18 @@ pub async fn update_flow_status_after_job_completion_internal( // tracing::error!("status_for_debug: {:?}", status_for_debug.flow_status); + let itered_len = if let Some(itered_len) = itered_len { + *itered_len + } else { + itered.as_ref().map(|itered| itered.len()).unwrap_or(0) + }; + tracing::info!( "parallel iteration {job_id_for_status} of flow {flow} update nindex: {nindex} len: {len}", nindex = nindex, - len = itered.len() + len = itered_len ); - (nindex, itered.len() as i32) + (nindex, itered_len as i32) } (_, Some(BranchAllStatus { len, .. })) => { let position = if flow_jobs_success.is_some() { @@ -3241,8 +3244,11 @@ async fn push_next_flow_job( simple_input_transforms, } => { if let Ok(args) = args.as_ref() { - // Read itered from separate table or fallback to JSONB - let itered = read_itered_from_db(db, flow_job.id, itered).await?; + let Some(itered) = itered else { + return Err(Error::ExecutionErr(format!( + "iterator itered should never be none for parallel for loop jobs" + ))); + }; let mut hm = HashMap::new(); for (k, v) in args.iter() { @@ -3582,12 +3588,8 @@ async fn push_next_flow_job( } } NextStatus::AllFlowJobs { iterator, branchall, .. } => { - // Conditionally write itered to separate table if all workers support it - // write_itered_to_db returns None if written to table, Some(itered) if should be in JSONB let iterator_for_status = if let Some(mut iter) = iterator { - if let Some(ref itered) = iter.itered { - iter.itered = write_itered_to_db(db, flow_job.id, itered).await?; - } + iter.itered = None; // in parallel for loops, itered is only useful when starting the jobs, no need to store it in status/or the dedicated table Some(iter) } else { None @@ -3864,7 +3866,7 @@ struct ForloopNextIteration { } enum ForLoopStatus { - ParallelIteration { itered: Vec> }, + ParallelIteration { itered: Vec>, itered_len: usize }, NextIteration(ForloopNextIteration), EmptyIterator, } @@ -4225,7 +4227,7 @@ async fn compute_next_flow_transform( ) .await } - ForLoopStatus::ParallelIteration { itered, .. } => { + ForLoopStatus::ParallelIteration { itered, itered_len } => { // let inner_path = Some(format!("{}/loop-parallel", flow_job.script_path(),)); // let value = &modules[0].get_value()?; @@ -4237,7 +4239,7 @@ async fn compute_next_flow_transform( // ContinuePayload::ForloopJobs { n: itered.len(), payload: payload } // } else { - let payloads = (0..itered.len()) + let payloads = (0..itered_len) .into_iter() .filter_map(|i| { let Some(payload) = payload_from_modules( @@ -4269,7 +4271,7 @@ async fn compute_next_flow_transform( branchall: None, iterator: Some(FlowIterator { index: 0, - itered_len: Some(itered.len()), + itered_len: Some(itered_len), itered: Some(itered), }), // we removed the is_simple_case for simple_input_transforms @@ -4619,7 +4621,7 @@ async fn next_forloop_status( if itered.is_empty() { ForLoopStatus::EmptyIterator } else if *parallel { - ForLoopStatus::ParallelIteration { itered } + ForLoopStatus::ParallelIteration { itered_len: itered.len(), itered } } else if let Some(first) = itered.first() { let iter = Iter { index: 0 as i32, value: first.to_owned() }; ForLoopStatus::NextIteration(ForloopNextIteration {