From 2d4fadb590590837412d638192fbd62bdc9331e8 Mon Sep 17 00:00:00 2001 From: hugocasa Date: Tue, 21 Apr 2026 19:06:42 +0200 Subject: [PATCH] feat: add s3 stream progress logs to other DB executors (#8898) Extract the MSSQL s3 ingest+upload logging pattern into a reusable `s3_stream_and_upload_with_logs` helper and apply it to the PostgreSQL, MySQL, OracleDB, BigQuery, and Snowflake executors. Each s3 streamed query now emits periodic progress lines, an ingest-done line, and an upload+transcode-done line to the job output, matching MSSQL. `convert_json_line_stream` now returns `BoxStream<'static, _>` so the output stream can be forwarded to `s3.upload` from inside the generic helper without lifetime gymnastics; the two existing callers already boxed the result, so behavior is unchanged. Co-authored-by: Claude Opus 4.7 (1M context) --- backend/windmill-object-store/src/lib.rs | 9 +-- .../windmill-worker/src/bigquery_executor.rs | 24 +++++-- backend/windmill-worker/src/common.rs | 71 +++++++++++++++++++ backend/windmill-worker/src/mssql_executor.rs | 59 +++------------ backend/windmill-worker/src/mysql_executor.rs | 24 +++++-- .../windmill-worker/src/oracledb_executor.rs | 23 ++++-- backend/windmill-worker/src/pg_executor.rs | 23 ++++-- .../windmill-worker/src/snowflake_executor.rs | 24 +++++-- 8 files changed, 176 insertions(+), 81 deletions(-) diff --git a/backend/windmill-object-store/src/lib.rs b/backend/windmill-object-store/src/lib.rs index 78a0e22c92..1a25513e54 100644 --- a/backend/windmill-object-store/src/lib.rs +++ b/backend/windmill-object-store/src/lib.rs @@ -995,17 +995,18 @@ pub async fn convert_json_line_stream( _output_format: S3ModeFormat, _progress: Option>, ) -> anyhow::Result<( - impl futures::TryStreamExt>, + futures::stream::BoxStream<'static, anyhow::Result>, IngestStats, )> where V: serde::Serialize, E: Into, { + use futures::StreamExt; let stream = async_stream::stream! { yield Err(anyhow::anyhow!("Parquet feature is not enabled. Cannot convert JSON line stream.")); }; - Ok((stream, IngestStats::default())) + Ok((stream.boxed(), IngestStats::default())) } #[cfg(feature = "parquet")] @@ -1025,7 +1026,7 @@ pub async fn convert_json_line_stream( output_format: S3ModeFormat, progress: Option>, ) -> anyhow::Result<( - impl TryStreamExt>, + futures::stream::BoxStream<'static, anyhow::Result>, IngestStats, )> where @@ -1204,7 +1205,7 @@ where }); Ok(( - tokio_stream::wrappers::ReceiverStream::new(rx), + tokio_stream::wrappers::ReceiverStream::new(rx).boxed(), ingest_stats, )) } diff --git a/backend/windmill-worker/src/bigquery_executor.rs b/backend/windmill-worker/src/bigquery_executor.rs index c9e803deaf..4ab2a4039a 100644 --- a/backend/windmill-worker/src/bigquery_executor.rs +++ b/backend/windmill-worker/src/bigquery_executor.rs @@ -4,11 +4,11 @@ use futures::future::BoxFuture; use futures::{FutureExt, StreamExt}; use reqwest::Client; use serde_json::{json, value::RawValue, Value}; +use uuid::Uuid; use windmill_common::client::AuthedClient; use windmill_common::error::to_anyhow; use windmill_common::worker::{Connection, SqlResultCollectionStrategy}; use windmill_common::{error::Error, worker::to_raw_value}; -use windmill_object_store::convert_json_line_stream; use windmill_parser_sql::{ parse_bigquery_sig, parse_db_resource, parse_s3_mode, parse_sql_blocks, parse_sql_statement_named_params, @@ -19,8 +19,8 @@ use serde::Deserialize; use crate::common::{build_args_values, resolve_job_timeout}; use crate::common::{ - build_http_client, get_reserved_variables, s3_mode_args_to_worker_data, OccupancyMetrics, - S3ModeWorkerData, + build_http_client, get_reserved_variables, s3_mode_args_to_worker_data, + s3_stream_and_upload_with_logs, OccupancyMetrics, S3ModeWorkerData, }; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; @@ -89,6 +89,9 @@ fn do_bigquery_inner<'a>( first_row_only: bool, http_client: &'a Client, s3: Option, + job_id: Uuid, + workspace_id: &'a str, + log_conn: &'a Connection, ) -> windmill_common::error::Result>>>> { let param_names = parse_sql_statement_named_params(query, '@'); @@ -203,9 +206,15 @@ fn do_bigquery_inner<'a>( } }; - let (stream, _) = - convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; - s3.upload(stream.boxed()).await?; + s3_stream_and_upload_with_logs( + "BigQuery", + rows_stream.boxed(), + &s3, + job_id, + workspace_id, + log_conn, + ) + .await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); } @@ -455,6 +464,9 @@ pub async fn do_bigquery( collection_strategy.collect_first_row_only(), &http_client, s3.clone(), + job.id, + &job.workspace_id, + conn, )? .await?; results.push(result); diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 964e136f81..be5da520bd 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -1570,6 +1570,77 @@ impl S3ModeWorkerData { } } +/// Stream rows to S3 (via `convert_json_line_stream` + `s3.upload`) while appending +/// periodic progress, ingest-done, and upload-done logs to the job output. +/// +/// `db_name` is a human-readable prefix for the log lines (e.g. "MSSQL", "PostgreSQL"). +pub async fn s3_stream_and_upload_with_logs( + db_name: &str, + rows_stream: S, + s3: &S3ModeWorkerData, + job_id: Uuid, + workspace_id: &str, + conn: &Connection, +) -> anyhow::Result<()> +where + S: futures::TryStreamExt> + Unpin, + V: serde::Serialize, + E: Into, +{ + let s3_format = s3.format; + let (progress_tx, mut progress_rx) = + tokio::sync::mpsc::channel::(8); + let progress_conn = conn.clone(); + let progress_workspace = workspace_id.to_string(); + let progress_db_name = db_name.to_string(); + let progress_task = tokio::spawn(async move { + while let Some(stats) = progress_rx.recv().await { + let logs = format!( + "\n{} s3 stream progress: {} rows, {:.1} MB | elapsed {:.1}s (db-fetch {:.1}s, write {:.1}s)", + progress_db_name, + stats.rows, + stats.bytes as f64 / 1_048_576.0, + stats.elapsed.as_secs_f64(), + stats.fetch_wait.as_secs_f64(), + stats.write_time.as_secs_f64(), + ); + windmill_queue::append_logs(&job_id, &progress_workspace, logs, &progress_conn).await; + } + }); + + let (stream, ingest_stats) = + windmill_object_store::convert_json_line_stream(rows_stream, s3_format, Some(progress_tx)) + .await?; + let _ = progress_task.await; + + let ingest_log = format!( + "\n{} s3 stream ingest done ({:?}): {} rows, {:.1} MB in {:.2}s (db-fetch {:.2}s, write {:.2}s, first row after {:.2}s)", + db_name, + s3_format, + ingest_stats.rows, + ingest_stats.bytes as f64 / 1_048_576.0, + ingest_stats.elapsed.as_secs_f64(), + ingest_stats.fetch_wait.as_secs_f64(), + ingest_stats.write_time.as_secs_f64(), + ingest_stats.first_row_latency.unwrap_or_default().as_secs_f64(), + ); + windmill_queue::append_logs(&job_id, workspace_id, ingest_log, conn).await; + + let upload_start = std::time::Instant::now(); + s3.upload(stream).await?; + let upload_elapsed = upload_start.elapsed(); + + let final_log = format!( + "\n{} s3 stream upload+transcode done: {:.2}s | total {:.2}s", + db_name, + upload_elapsed.as_secs_f64(), + (ingest_stats.elapsed + upload_elapsed).as_secs_f64(), + ); + windmill_queue::append_logs(&job_id, workspace_id, final_log, conn).await; + + Ok(()) +} + pub fn s3_mode_args_to_worker_data( s3: S3ModeArgs, client: AuthedClient, diff --git a/backend/windmill-worker/src/mssql_executor.rs b/backend/windmill-worker/src/mssql_executor.rs index d7be6b0d1c..f1e8d93171 100644 --- a/backend/windmill-worker/src/mssql_executor.rs +++ b/backend/windmill-worker/src/mssql_executor.rs @@ -18,13 +18,13 @@ use windmill_common::{ utils::empty_as_none, worker::{to_raw_value, Connection}, }; -use windmill_object_store::convert_json_line_stream; use windmill_parser_sql::{parse_db_resource, parse_mssql_sig, parse_s3_mode}; use windmill_queue::MiniPulledJob; use windmill_queue::{append_logs, CanceledBy}; use crate::common::{ - build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, OccupancyMetrics, + build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, + s3_stream_and_upload_with_logs, OccupancyMetrics, }; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; @@ -244,7 +244,6 @@ pub async fn do_mssql( // fetching data in an asynchronous manner, if needed. if let Some(s3) = s3 { - let s3_format = s3.format; let rows_stream = async_stream::stream! { let mut stream = prepared_query.query(&mut client).await.map_err(to_anyhow)?.into_row_stream().map(|row| { row_to_json(row.map_err(to_anyhow)?).map_err(to_anyhow) @@ -254,51 +253,15 @@ pub async fn do_mssql( } }; - let (progress_tx, mut progress_rx) = - tokio::sync::mpsc::channel::(8); - let progress_conn = conn.clone(); - let progress_job_id = job.id; - let progress_workspace = job.workspace_id.clone(); - let progress_task = tokio::spawn(async move { - while let Some(stats) = progress_rx.recv().await { - let logs = format!( - "\nMSSQL s3 stream progress: {} rows, {:.1} MB | elapsed {:.1}s (db-fetch {:.1}s, write {:.1}s)", - stats.rows, - stats.bytes as f64 / 1_048_576.0, - stats.elapsed.as_secs_f64(), - stats.fetch_wait.as_secs_f64(), - stats.write_time.as_secs_f64(), - ); - append_logs(&progress_job_id, &progress_workspace, logs, &progress_conn).await; - } - }); - - let (stream, ingest_stats) = - convert_json_line_stream(rows_stream.boxed(), s3_format, Some(progress_tx)).await?; - let _ = progress_task.await; - - let ingest_log = format!( - "\nMSSQL s3 stream ingest done ({:?}): {} rows, {:.1} MB in {:.2}s (db-fetch {:.2}s, write {:.2}s, first row after {:.2}s)", - s3_format, - ingest_stats.rows, - ingest_stats.bytes as f64 / 1_048_576.0, - ingest_stats.elapsed.as_secs_f64(), - ingest_stats.fetch_wait.as_secs_f64(), - ingest_stats.write_time.as_secs_f64(), - ingest_stats.first_row_latency.unwrap_or_default().as_secs_f64(), - ); - append_logs(&job.id, &job.workspace_id, ingest_log, conn).await; - - let upload_start = std::time::Instant::now(); - s3.upload(stream.boxed()).await?; - let upload_elapsed = upload_start.elapsed(); - - let final_log = format!( - "\nMSSQL s3 stream upload+transcode done: {:.2}s | total {:.2}s", - upload_elapsed.as_secs_f64(), - (ingest_stats.elapsed + upload_elapsed).as_secs_f64(), - ); - append_logs(&job.id, &job.workspace_id, final_log, conn).await; + s3_stream_and_upload_with_logs( + "MSSQL", + rows_stream.boxed(), + &s3, + job.id, + &job.workspace_id, + conn, + ) + .await?; Ok(to_raw_value(&s3.to_return_s3_obj())) } else { diff --git a/backend/windmill-worker/src/mysql_executor.rs b/backend/windmill-worker/src/mysql_executor.rs index aee7de16f8..b2ec9b074d 100644 --- a/backend/windmill-worker/src/mysql_executor.rs +++ b/backend/windmill-worker/src/mysql_executor.rs @@ -12,12 +12,12 @@ use serde::{Deserialize, Serialize}; use serde_json::{json, value::RawValue, Value}; use std::str::FromStr; use tokio::sync::Mutex; +use uuid::Uuid; use windmill_common::{ client::AuthedClient, error::{to_anyhow, Error}, worker::{to_raw_value, Connection, SqlResultCollectionStrategy}, }; -use windmill_object_store::convert_json_line_stream; use windmill_parser_sql::{ parse_db_resource, parse_mysql_sig, parse_s3_mode, parse_sql_blocks, parse_sql_statement_named_params, RE_ARG_MYSQL_NAMED, @@ -27,8 +27,8 @@ use windmill_queue::MiniPulledJob; use crate::{ common::{ - build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, OccupancyMetrics, - S3ModeWorkerData, + build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, + s3_stream_and_upload_with_logs, OccupancyMetrics, S3ModeWorkerData, }, handle_child::run_future_with_polling_update_job_poller, sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args, @@ -52,6 +52,9 @@ fn do_mysql_inner<'a>( skip_collect: bool, first_row_only: bool, s3: Option, + job_id: Uuid, + workspace_id: &'a str, + log_conn: &'a Connection, ) -> windmill_common::error::Result>>>> { let param_names = parse_sql_statement_named_params(query, ':') @@ -107,9 +110,15 @@ fn do_mysql_inner<'a>( } }; - let (stream, _) = - convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; - s3.upload(stream.boxed()).await?; + s3_stream_and_upload_with_logs( + "MySQL", + rows_stream.boxed(), + s3, + job_id, + workspace_id, + log_conn, + ) + .await?; Ok(vec![to_raw_value(&s3.to_return_s3_obj())]) } else { @@ -321,6 +330,9 @@ pub async fn do_mysql( && i < queries.len() - 1, collection_strategy.collect_first_row_only(), s3.clone(), + job.id, + &job.workspace_id, + conn, )? .await?; results.push(result); diff --git a/backend/windmill-worker/src/oracledb_executor.rs b/backend/windmill-worker/src/oracledb_executor.rs index 55ffac467b..3381ff2326 100644 --- a/backend/windmill-worker/src/oracledb_executor.rs +++ b/backend/windmill-worker/src/oracledb_executor.rs @@ -8,11 +8,11 @@ use itertools::Itertools; use oracle::sql_type::{InnerValue, OracleType, ToSql}; use serde::{Deserialize, Serialize}; use serde_json::{json, value::RawValue, Value}; +use uuid::Uuid; use windmill_common::{ error::{to_anyhow, Error}, worker::{to_raw_value, Connection, SqlResultCollectionStrategy}, }; -use windmill_object_store::convert_json_line_stream; use windmill_queue::MiniPulledJob; use windmill_parser_sql::{ @@ -24,7 +24,8 @@ use windmill_queue::CanceledBy; use crate::{ common::{ build_args_values, check_executor_binary_exists, get_reserved_variables, - s3_mode_args_to_worker_data, OccupancyMetrics, S3ModeWorkerData, + s3_mode_args_to_worker_data, s3_stream_and_upload_with_logs, OccupancyMetrics, + S3ModeWorkerData, }, handle_child::run_future_with_polling_update_job_poller, sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args, @@ -50,6 +51,9 @@ pub fn do_oracledb_inner<'a>( skip_collect: bool, first_row_only: bool, s3: Option, + job_id: Uuid, + workspace_id: &'a str, + log_conn: &'a Connection, ) -> windmill_common::error::Result>>>> { let qw = query.trim_end_matches(';').to_string(); @@ -162,9 +166,15 @@ pub fn do_oracledb_inner<'a>( } if let Some(s3) = s3 { - let (stream, _) = - convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; - s3.upload(stream.boxed()).await?; + s3_stream_and_upload_with_logs( + "OracleDB", + rows_stream.boxed(), + &s3, + job_id, + workspace_id, + log_conn, + ) + .await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); } else { let rows: Vec<_> = rows_stream.collect().await; @@ -444,6 +454,9 @@ pub async fn do_oracledb( && i < queries.len() - 1, collection_strategy.collect_first_row_only(), s3.clone(), + job.id, + &job.workspace_id, + conn, )? .await?; results.push(result); diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 2e577da6a8..f9b673216c 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -29,7 +29,6 @@ use windmill_common::worker::{ }; use windmill_common::workspaces::get_datatable_resource_from_db_unchecked; use windmill_common::{PgDatabase, PrepareQueryColumnInfo, PrepareQueryResult, DB}; -use windmill_object_store::convert_json_line_stream; use windmill_parser::{Arg, Typ}; use windmill_parser_sql::{ parse_db_resource, parse_pg_statement_arg_indices, parse_pgsql_sig_with_typed_schema, @@ -39,8 +38,8 @@ use windmill_queue::{CanceledBy, MiniPulledJob}; use crate::agent_workers::get_datatable_resource_from_agent_http; use crate::common::{ - build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, sizeof_val, - OccupancyMetrics, S3ModeWorkerData, + build_args_values, get_reserved_variables, s3_mode_args_to_worker_data, + s3_stream_and_upload_with_logs, sizeof_val, OccupancyMetrics, S3ModeWorkerData, }; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; @@ -139,6 +138,9 @@ fn do_postgresql_inner<'a>( first_row_only: bool, s3: Option, typed_schema: bool, + job_id: Uuid, + workspace_id: &'a str, + log_conn: &'a Connection, ) -> error::Result>>>> { let mut query_params = vec![]; let mut param_types = vec![]; @@ -203,9 +205,15 @@ fn do_postgresql_inner<'a>( row_result.and_then(|row| postgres_row_to_json_value(row).map_err(to_anyhow)) }); - let (stream, _) = - convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; - s3.upload(stream.boxed()).await?; + s3_stream_and_upload_with_logs( + "PostgreSQL", + rows_stream.boxed(), + s3, + job_id, + workspace_id, + log_conn, + ) + .await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); } else { @@ -456,6 +464,9 @@ pub async fn do_postgresql( collection_strategy.collect_first_row_only(), s3.clone(), typed_schema, + job.id, + &job.workspace_id, + conn, )? .await?; results.push(result); diff --git a/backend/windmill-worker/src/snowflake_executor.rs b/backend/windmill-worker/src/snowflake_executor.rs index 0873624f54..95757112bc 100644 --- a/backend/windmill-worker/src/snowflake_executor.rs +++ b/backend/windmill-worker/src/snowflake_executor.rs @@ -8,9 +8,9 @@ use reqwest::{Client, Response}; use serde_json::{json, value::RawValue, Value}; use sha2::{Digest, Sha256}; use std::collections::HashMap; +use uuid::Uuid; use windmill_common::error::to_anyhow; use windmill_common::worker::{Connection, SqlResultCollectionStrategy}; -use windmill_object_store::convert_json_line_stream; use windmill_common::{error::Error, worker::to_raw_value}; use windmill_parser_sql::{ @@ -22,8 +22,8 @@ use serde::{Deserialize, Serialize}; use crate::common::{build_args_values, get_reserved_variables}; use crate::common::{ - build_http_client, resolve_job_timeout, s3_mode_args_to_worker_data, OccupancyMetrics, - S3ModeWorkerData, + build_http_client, resolve_job_timeout, s3_mode_args_to_worker_data, + s3_stream_and_upload_with_logs, OccupancyMetrics, S3ModeWorkerData, }; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; @@ -194,6 +194,9 @@ fn do_snowflake_inner<'a>( s3: Option, reserved_variables: &HashMap, deadline: std::time::Instant, + job_id: Uuid, + workspace_id: &'a str, + log_conn: &'a Connection, ) -> windmill_common::error::Result>>>> { let sig = parse_snowflake_sig(&query) @@ -501,9 +504,15 @@ fn do_snowflake_inner<'a>( if let Some(s3) = s3 { let rows_stream = rows_stream.map(|r| serde_json::value::to_value(&r?).map_err(to_anyhow)); - let (stream, _) = - convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; - s3.upload(stream.boxed()).await?; + s3_stream_and_upload_with_logs( + "Snowflake", + rows_stream.boxed(), + &s3, + job_id, + workspace_id, + log_conn, + ) + .await?; Ok(vec![to_raw_value(&s3.to_return_s3_obj())]) } else { let rows = rows_stream @@ -688,6 +697,9 @@ pub async fn do_snowflake( s3.clone(), &reserved_variables, deadline, + job.id, + &job.workspace_id, + conn, )? .await?; results.push(result);