mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-09 08:03:50 +00:00
fix: add settable poll delay for sse streams
This commit is contained in:
@@ -1751,6 +1751,7 @@ pub struct RunJobQuery {
|
||||
pub timeout: Option<i32>,
|
||||
pub cache_ttl: Option<i32>,
|
||||
pub skip_preprocessor: Option<bool>,
|
||||
pub poll_delay_ms: Option<u64>,
|
||||
}
|
||||
|
||||
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<bool>,
|
||||
pub fast: Option<bool>,
|
||||
pub is_flow: Option<bool>,
|
||||
pub poll_delay_ms: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Serialize, Debug)]
|
||||
@@ -6696,6 +6700,7 @@ async fn get_job_update_sse(
|
||||
only_result,
|
||||
fast,
|
||||
is_flow,
|
||||
poll_delay_ms,
|
||||
}): Query<JobUpdateQuery>,
|
||||
) -> error::Result<Response> {
|
||||
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<bool>,
|
||||
is_flow: Option<bool>,
|
||||
tx: tokio::sync::mpsc::Sender<JobUpdateSSEStream>,
|
||||
poll_delay_ms: Option<u64>,
|
||||
) -> () {
|
||||
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}");
|
||||
|
||||
Reference in New Issue
Block a user