diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 1a95bfe721..32671f9c64 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -5938b3f66b070ba3578aeeda89e6cc17c4cbda2d \ No newline at end of file +053cf0a539c64efbe8eee0e2a330b72114a43193 \ No newline at end of file diff --git a/backend/windmill-worker/src/job_logger.rs b/backend/windmill-worker/src/job_logger.rs index 05b9e5812b..8a6d824e6e 100644 --- a/backend/windmill-worker/src/job_logger.rs +++ b/backend/windmill-worker/src/job_logger.rs @@ -1,15 +1,6 @@ -use itertools::Itertools; -use swc_ecma_parser::lexer::util::CharExt; - -#[cfg(all(feature = "enterprise", feature = "parquet"))] -use object_store::path::Path; use regex::Regex; -#[cfg(all(feature = "enterprise", feature = "parquet"))] -use windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS; - -use windmill_common::error::{self}; -use windmill_common::worker::{CLOUD_HOSTED, TMP_DIR}; +use windmill_common::worker::CLOUD_HOSTED; use windmill_queue::append_logs; @@ -19,6 +10,12 @@ use std::sync::Arc; use uuid::Uuid; use windmill_common::DB; +use crate::job_logger_ee::default_disk_log_storage; + +#[cfg(all(feature = "enterprise", feature = "parquet"))] +use crate::job_logger_ee::s3_storage; + + pub enum CompactLogs { #[cfg(not(all(feature = "enterprise", feature = "parquet")))] NotEE, @@ -28,150 +25,7 @@ pub enum CompactLogs { S3, } -async fn compact_logs( - job_id: Uuid, - w_id: &str, - db: &DB, - nlogs: String, - total_size: Arc, - compact_kind: CompactLogs, - _worker_name: &str, -) -> error::Result<(String, String)> { - let mut prev_logs = sqlx::query_scalar!( - "SELECT logs FROM job_logs WHERE job_id = $1 AND workspace_id = $2", - job_id, - w_id - ) - .fetch_optional(db) - .await? - .flatten() - .unwrap_or_default(); - let size = prev_logs.char_indices().count() as i32; - let nlogs_len = nlogs.char_indices().count(); - let to_keep_in_db = usize::max( - usize::min(nlogs_len, 3000), - nlogs_len % LARGE_LOG_THRESHOLD_SIZE, - ); - let extra_split = to_keep_in_db < nlogs_len; - let stored_in_storage_len = if extra_split { - nlogs_len - to_keep_in_db - } else { - 0 - }; - let extra_to_newline = nlogs - .chars() - .skip(stored_in_storage_len) - .find_position(|x| x.is_line_break()) - .map(|(i, _)| i) - .unwrap_or(to_keep_in_db); - let stored_in_storage_to_newline = stored_in_storage_len + extra_to_newline; - let (append_to_storage, stored_in_db) = if extra_split { - if stored_in_storage_to_newline == nlogs.len() { - (nlogs.as_ref(), "".to_string()) - } else { - let split_idx = nlogs - .char_indices() - .nth(stored_in_storage_to_newline) - .map(|(i, _)| i) - .unwrap_or(0); - let (append_to_storage, stored_in_db) = nlogs.split_at(split_idx); - // tracing::error!("{append_to_storage} ||||| {stored_in_db}"); - // tracing::error!( - // "{:?} {:?} {} {}", - // excess_prev_logs.lines().last(), - // current_logs.lines().next(), - // split_idx, - // excess_size_modulo - // ); - (append_to_storage, stored_in_db.to_string()) - } - } else { - // tracing::error!("{:?}", nlogs.lines().last()); - ("", nlogs.to_string()) - }; - - let new_size_with_excess = size + stored_in_storage_to_newline as i32; - - let new_size = total_size.fetch_add( - new_size_with_excess as u32, - std::sync::atomic::Ordering::SeqCst, - ) + new_size_with_excess as u32; - - let path = format!( - "logs/{job_id}/{}_{new_size}.txt", - chrono::Utc::now().timestamp_millis() - ); - - let mut new_current_logs = match compact_kind { - CompactLogs::NoS3 => format!("\n[windmill] No object storage set in instance settings. Previous logs have been saved to disk at {path}"), - CompactLogs::S3 => format!("\n[windmill] Previous logs have been saved to object storage at {path}"), - #[cfg(not(all(feature = "enterprise", feature = "parquet")))] - CompactLogs::NotEE => format!("\n[windmill] Previous logs have been saved to disk at {path}"), - }; - new_current_logs.push_str(&stored_in_db); - - sqlx::query!( - "UPDATE job_logs SET logs = $1, log_offset = $2, - log_file_index = array_append(coalesce(log_file_index, array[]::text[]), $3) - WHERE workspace_id = $4 AND job_id = $5", - new_current_logs, - new_size as i32, - path, - w_id, - job_id - ) - .execute(db) - .await?; - prev_logs.push_str(&append_to_storage); - - return Ok((prev_logs, path)); -} - -async fn default_disk_log_storage( - job_id: Uuid, - w_id: &str, - db: &DB, - nlogs: String, - total_size: Arc, - compact_kind: CompactLogs, - worker_name: &str, -) { - match compact_logs( - job_id, - &w_id, - &db, - nlogs, - total_size, - compact_kind, - worker_name, - ) - .await - { - Err(e) => tracing::error!("Could not compact logs for job {job_id}: {e:?}",), - Ok((prev_logs, path)) => { - let path = format!("{}/{}", TMP_DIR, path); - let splitted = &path.split("/").collect_vec(); - tokio::fs::create_dir_all(splitted.into_iter().take(splitted.len() - 1).join("/")) - .await - .map_err(|e| { - tracing::error!("Could not create logs directory: {e:?}",); - e - }) - .ok(); - let created = tokio::fs::File::create(&path).await; - if let Err(e) = created { - tracing::error!("Could not create logs file {path}: {e:?}",); - return; - } - if let Err(e) = tokio::fs::write(&path, prev_logs).await { - tracing::error!("Could not write to logs file {path}: {e:?}"); - } else { - tracing::info!("Logs length of {job_id} has exceeded a threshold. Previous logs have been saved to disk at {path}"); - } - } - } -} pub(crate) async fn append_job_logs( job_id: Uuid, @@ -184,43 +38,7 @@ pub(crate) async fn append_job_logs( ) -> () { if must_compact_logs { #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() { - match compact_logs( - job_id, - &w_id, - &db, - logs, - total_size, - CompactLogs::S3, - &worker_name, - ) - .await - { - Err(e) => tracing::error!("Could not compact logs for job {job_id}: {e:?}",), - Ok((prev_logs, path)) => { - tracing::info!("Logs length of {job_id} has exceeded a threshold. Previous logs have been saved to object storage at {path}"); - let path2 = path.clone(); - if let Err(e) = os - .put(&Path::from(path), prev_logs.to_string().into_bytes().into()) - .await - { - tracing::error!("Could not save logs to s3: {e:?}"); - } - tracing::info!("Logs of {job_id} saved to object storage at {path2}"); - } - } - } else { - default_disk_log_storage( - job_id, - &w_id, - &db, - logs, - total_size, - CompactLogs::NoS3, - &worker_name, - ) - .await; - } + s3_storage(job_id, &w_id, &db, &logs, &total_size, &worker_name).await; #[cfg(not(all(feature = "enterprise", feature = "parquet")))] { diff --git a/backend/windmill-worker/src/job_logger_ee.rs b/backend/windmill-worker/src/job_logger_ee.rs new file mode 100644 index 0000000000..0274c6c777 --- /dev/null +++ b/backend/windmill-worker/src/job_logger_ee.rs @@ -0,0 +1,26 @@ +use std::sync::atomic::AtomicU32; +use std::sync::Arc; + +use uuid::Uuid; +use windmill_common::DB; + +use crate::job_logger::CompactLogs; + +#[cfg(all(feature = "enterprise", feature = "parquet"))] +pub(crate) async fn s3_storage(_job_id: Uuid, _w_id: &String, _db: &sqlx::Pool, _logs: &String, _total_size: &Arc, _worker_name: &String) { + tracing::info!("Logs length of {job_id} has exceeded a threshold. Implementation to store excess on s3 in not OSS"); +} + +pub(crate) async fn default_disk_log_storage( + job_id: Uuid, + _w_id: &str, + _db: &DB, + _nlogs: String, + _total_size: Arc, + _compact_kind: CompactLogs, + _worker_name: &str, +) { + tracing::info!("Logs length of {job_id} has exceeded a threshold. Implementation to store excess on disk in not OSS"); +} + + diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index b790eb3faf..680f3ca178 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -29,7 +29,7 @@ mod rust_executor; mod worker; mod worker_flow; mod worker_lockfiles; - +mod job_logger_ee; pub use worker::*; pub use result_processor::handle_job_error;