From 869d0bc85466b71a367755148e61c4ebf606d4bc Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 22 Sep 2025 10:57:33 +0000 Subject: [PATCH] fix: add settable poll delay for sse streams --- backend/windmill-api/src/jobs.rs | 23 ++++++++++++++++++++++- 1 file changed, 22 insertions(+), 1 deletion(-) diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 8089e39fda..d1815a4994 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1751,6 +1751,7 @@ pub struct RunJobQuery { pub timeout: Option, pub cache_ttl: Option, pub skip_preprocessor: Option, + pub poll_delay_ms: Option, } impl RunJobQuery { @@ -5294,6 +5295,7 @@ pub async fn stream_job( .to_args_from_runnable(&db, &w_id, runnable_id.clone(), run_query.skip_preprocessor) .await?; + let poll_delay_ms = run_query.poll_delay_ms; let uuid = match runnable_id { RunnableId::ScriptId(ScriptId::ScriptPath(script_path)) | RunnableId::HubScript(script_path) => { @@ -5363,6 +5365,7 @@ pub async fn stream_job( None, None, tx, + poll_delay_ms, ); let body = axum::body::Body::from_stream(stream.map(Result::<_, std::convert::Infallible>::Ok)); @@ -6535,6 +6538,7 @@ pub struct JobUpdateQuery { pub only_result: Option, pub fast: Option, pub is_flow: Option, + pub poll_delay_ms: Option, } #[derive(Serialize, Debug)] @@ -6696,6 +6700,7 @@ async fn get_job_update_sse( only_result, fast, is_flow, + poll_delay_ms, }): Query, ) -> error::Result { let (tx, rx) = tokio::sync::mpsc::channel(32); @@ -6715,6 +6720,7 @@ async fn get_job_update_sse( no_logs, is_flow, tx, + poll_delay_ms, ); let stream = tokio_stream::wrappers::ReceiverStream::new(rx).map(|x| { @@ -6760,6 +6766,7 @@ fn start_job_update_sse_stream( no_logs: Option, is_flow: Option, tx: tokio::sync::mpsc::Sender, + poll_delay_ms: Option, ) -> () { tokio::spawn(async move { let mut log_offset = initial_log_offset; @@ -6842,13 +6849,27 @@ fn start_job_update_sse_stream( loop { i += 1; - let ms_duration = if i > 100 || !fast.unwrap_or(false) { + let mut ms_duration = if i > 100 || !fast.unwrap_or(false) { 3000 } else if i > 10 { 500 } else { 100 }; + + if let Some(poll_delay_ms) = poll_delay_ms { + #[cfg(feature = "enterprise")] + if poll_delay_ms < 50 { + tracing::warn!("Poll delay ms is less than 50, setting it to 50"); + ms_duration = 50; + } else { + ms_duration = poll_delay_ms; + } + + #[cfg(not(feature = "enterprise"))] + tracing::warn!("Settable poll delay requires EE"); + } + if last_ping.elapsed().as_secs() > 5 { if tx.send(JobUpdateSSEStream::Ping).await.is_err() { tracing::warn!("Failed to send job ping for job {job_id}");