From aae70ac6cbe33de16718e89e6b4a9621598cd524 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 12 May 2024 04:16:04 +0200 Subject: [PATCH] feat: logs can be downloaded directly from server/frontend if using shared volume --- backend/windmill-api/src/jobs.rs | 109 ++++++++++-------- backend/windmill-common/src/worker.rs | 2 + backend/windmill-worker/src/common.rs | 8 +- backend/windmill-worker/src/worker.rs | 7 +- docker-compose.yml | 7 +- frontend/src/lib/components/LogViewer.svelte | 45 +++++++- .../src/lib/components/TestJobLoader.svelte | 21 +++- 7 files changed, 133 insertions(+), 66 deletions(-) diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 62fbd4caf2..fd1f6c8428 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -12,6 +12,7 @@ use serde_json::value::RawValue; use std::collections::HashMap; #[cfg(feature = "prometheus")] use std::sync::atomic::Ordering; +use tokio::io::AsyncReadExt; #[cfg(feature = "prometheus")] use tokio::time::Instant; use windmill_common::flow_status::{JobResult, RestartedFrom}; @@ -19,6 +20,7 @@ use windmill_common::jobs::{ format_completed_job_result, format_result, CompletedJobWithFormattedResult, FormattedResult, ENTRYPOINT_OVERRIDE, }; +use windmill_common::worker::TMP_DIR; #[cfg(all(feature = "enterprise", feature = "parquet"))] use windmill_common::scripts::PREVIEW_IS_CODEBASE_HASH; @@ -664,17 +666,17 @@ async fn get_job_internal( async fn get_logs_from_store( log_offset: i32, logs: &str, - log_file_index: Option>, + log_file_index: &Option>, ) -> Option> { if log_offset > 0 { - if let Some(file_index) = log_file_index { + if let Some(file_index) = log_file_index.clone() { tracing::debug!("Getting logs from store: {file_index:?}"); if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() { tracing::debug!("object store client present, streaming from there"); let logs = logs.to_string(); let stream = async_stream::stream! { - for file_p in file_index { + for file_p in file_index.clone() { let file_p_2 = file_p.clone(); let file = os.get(&object_store::path::Path::from(file_p)).await; if let Ok(file) = file { @@ -696,6 +698,40 @@ async fn get_logs_from_store( } return None; } + +async fn get_logs_from_disk( + log_offset: i32, + logs: &str, + log_file_index: &Option>, +) -> Option> { + if log_offset > 0 { + if let Some(file_index) = log_file_index.clone() { + for file_p in &file_index { + if !tokio::fs::metadata(format!("{TMP_DIR}/{file_p}")) + .await + .is_ok() + { + return None; + } + } + + let logs = logs.to_string(); + let stream = async_stream::stream! { + for file_p in file_index.clone() { + let mut file = tokio::fs::File::open(format!("{TMP_DIR}/{file_p}")).await.map_err(to_anyhow)?; + let mut buffer = Vec::new(); + file.read_to_end(&mut buffer).await.map_err(to_anyhow)?; + yield Ok(bytes::Bytes::from(buffer)) as anyhow::Result; + } + + yield Ok(bytes::Bytes::from(logs)) + }; + return Some(Ok(Body::from_stream(stream))); + } + } + return None; +} + async fn get_job_logs( Extension(db): Extension, Path((w_id, id)): Path<(String, Uuid)>, @@ -714,7 +750,11 @@ async fn get_job_logs( if let Some(record) = record { let logs = record.logs.unwrap_or_default(); #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(r) = get_logs_from_store(record.log_offset, &logs, record.log_file_index).await + if let Some(r) = get_logs_from_store(record.log_offset, &logs, &record.log_file_index).await + { + return r; + } + if let Some(r) = get_logs_from_disk(record.log_offset, &logs, &record.log_file_index).await { return r; } @@ -734,7 +774,10 @@ async fn get_job_logs( let logs = text.logs.unwrap_or_default(); #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(r) = get_logs_from_store(text.log_offset, &logs, text.log_file_index).await { + if let Some(r) = get_logs_from_store(text.log_offset, &logs, &text.log_file_index).await { + return r; + } + if let Some(r) = get_logs_from_disk(text.log_offset, &logs, &text.log_file_index).await { return r; } Ok(Body::from(logs)) @@ -3649,45 +3692,20 @@ pub struct JobUpdate { pub flow_status: Option, } -// #[cfg(all(feature = "enterprise", feature = "parquet"))] -// async fn get_logs_from_store( -// log_offset: i32, -// logs: &str, -// log_file_index: Option>, -// ) -> Option> { -// if log_offset > 0 { -// if let Some(file_index) = log_file_index { -// tracing::debug!("Getting logs from store: {file_index:?}"); -// if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() { -// tracing::debug!("object store client present, streaming from there"); - -// let logs = logs.to_string(); -// let stream = async_stream::stream! { -// for file_p in file_index { -// let file_p_2 = file_p.clone(); -// let file = os.get(&object_store::path::Path::from(file_p)).await; -// if let Ok(file) = file { -// if let Ok(bytes) = file.bytes().await { -// yield Ok(bytes::Bytes::from(bytes)) as object_store::Result; -// } -// } else { -// tracing::debug!("error getting file from store: {file_p_2}: {}", file.err().unwrap()); -// } -// } - -// yield Ok(bytes::Bytes::from(logs)) -// }; -// return Some(Ok(Body::from_stream(stream))); -// } else { -// tracing::debug!("object store client not present, cannot stream logs from store"); -// } -// } -// } -// return None; -// } - -#[cfg(all(feature = "enterprise", feature = "parquet"))] async fn get_log_file(Path((_w_id, file_p)): Path<(String, String)>) -> error::Result { + let local_file = format!("{TMP_DIR}/logs/{file_p}"); + if tokio::fs::metadata(&local_file).await.is_ok() { + let mut file = tokio::fs::File::open(local_file).await.map_err(to_anyhow)?; + let mut buffer = Vec::new(); + file.read_to_end(&mut buffer).await.map_err(to_anyhow)?; + let res = Response::builder() + .header(http::header::CONTENT_TYPE, "text/plain") + .body(Body::from(bytes::Bytes::from(buffer))) + .unwrap(); + return Ok(res); + } + + #[cfg(all(feature = "enterprise", feature = "parquet"))] if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() { let file = os .get(&object_store::path::Path::from(format!("logs/{file_p}"))) @@ -3717,12 +3735,9 @@ async fn get_log_file(Path((_w_id, file_p)): Path<(String, String)>) -> error::R "Object store client not present, cannot stream logs from store".to_string(), )); } -} -#[cfg(not(all(feature = "enterprise", feature = "parquet")))] -async fn get_log_file(Path((_w_id, file_p)): Path<(String, String)>) -> error::Result { return Err(error::Error::NotFound(format!( - "Get log file is an EE feature: {}", + "File not found on server logs volume /tmp/windmill/logs and no distributed logs s3 storage for {}", file_p ))); } diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index f9db742a49..1b3a2c46ab 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -70,6 +70,8 @@ lazy_static::lazy_static! { } +pub const TMP_DIR: &str = "/tmp/windmill"; + pub async fn reload_custom_tags_setting(db: &DB) -> error::Result<()> { let q = sqlx::query!( "SELECT value FROM global_settings WHERE name = $1", diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 1c9a5c7860..1c55f7e697 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -26,7 +26,7 @@ use windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS; use windmill_common::s3_helpers::{ get_etag_or_empty, LargeFileStorage, ObjectStoreResource, S3Object, }; -use windmill_common::worker::{CLOUD_HOSTED, WORKER_CONFIG}; +use windmill_common::worker::{CLOUD_HOSTED, TMP_DIR, WORKER_CONFIG}; use windmill_common::{ error::{self, Error}, jobs::QueuedJob, @@ -70,7 +70,7 @@ use futures::{ use crate::{ AuthedClient, AuthedClientBackgroundTask, JOB_DEFAULT_TIMEOUT, MAX_RESULT_SIZE, - MAX_TIMEOUT_DURATION, MAX_WAIT_FOR_SIGINT, MAX_WAIT_FOR_SIGTERM, ROOT_CACHE_DIR, TMP_DIR, + MAX_TIMEOUT_DURATION, MAX_WAIT_FOR_SIGINT, MAX_WAIT_FOR_SIGTERM, ROOT_CACHE_DIR, }; pub async fn build_args_map<'a>( @@ -717,9 +717,9 @@ async fn compact_logs( ); let mut new_current_logs = match compact_kind { - CompactLogs::NoS3 => format!("[windmill] worker {worker_name}: Logs length has exceeded a threshold\n[windmill] Previous logs have been saved to disk at {path}, add object storage in the instance settings to save it on distributed storage and allow direct download from Windmill\n"), + CompactLogs::NoS3 => format!("[windmill] No object storage set in instance settings. Previous logs have been saved to disk at {path}\n"), CompactLogs::S3 => format!("[windmill] Previous logs have been saved to object storage at {path}\n"), - CompactLogs::NotEE => format!("[windmill] worker {worker_name}: Logs length has exceeded a threshold\n[windmill] Previous logs have been saved to disk at {path}\n[windmill] Upgrade to EE and add object storage to save it persistentely on distributed storage and allow direct download from Windmill\n"), + CompactLogs::NotEE => format!("[windmill] Previous logs have been saved to disk at {path}\n"), }; new_current_logs.push_str(¤t_logs); diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index d6acedbb92..d55694a3b4 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -6,6 +6,8 @@ * LICENSE-AGPL for a copy of the license. */ +use windmill_common::worker::TMP_DIR; + use anyhow::Result; use const_format::concatcp; #[cfg(feature = "prometheus")] @@ -185,10 +187,9 @@ pub async fn create_token_for_owner( Ok(token) } -pub const TMP_DIR: &str = "/tmp/windmill"; -pub const TMP_LOGS_DIR: &str = "/tmp/windmill/logs"; +pub const TMP_LOGS_DIR: &str = concatcp!(TMP_DIR, "/logs"); -pub const ROOT_CACHE_DIR: &str = "/tmp/windmill/cache/"; +pub const ROOT_CACHE_DIR: &str = concatcp!(TMP_DIR, "/cache/"); pub const LOCK_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "lock"); pub const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip"); pub const TAR_PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "tar/pip"); diff --git a/docker-compose.yml b/docker-compose.yml index 7b7d9d58cf..d5053ecd54 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -37,6 +37,8 @@ services: depends_on: db: condition: service_healthy + volumes: + - worker_logs:/tmp/windmill/logs windmill_worker: image: ${WM_IMAGE} @@ -60,6 +62,7 @@ services: # mount the docker socket to allow to run docker containers from within the workers - /var/run/docker.sock:/var/run/docker.sock - worker_dependency_cache:/tmp/windmill/cache + - worker_logs:/tmp/windmill/logs ## This worker is specialized for "native" jobs. Native jobs run in-process and thus are much more lightweight than other jobs windmill_worker_native: @@ -80,7 +83,8 @@ services: depends_on: db: condition: service_healthy - + volumes: + - worker_logs:/tmp/windmill/logs ## This worker is specialized for reports or scraping jobs. It is assigned the "reports" worker group which has an init script that installs chromium and can be targeted by using the "chromium" worker tag. # windmill_worker_reports: # image: ${WM_IMAGE} @@ -142,4 +146,5 @@ services: volumes: db_data: null worker_dependency_cache: null + worker_logs: null lsp_cache: null diff --git a/frontend/src/lib/components/LogViewer.svelte b/frontend/src/lib/components/LogViewer.svelte index 9fc5f0941b..ec6718fb83 100644 --- a/frontend/src/lib/components/LogViewer.svelte +++ b/frontend/src/lib/components/LogViewer.svelte @@ -1,5 +1,9 @@