feat: large log disk and distributed storage compaction

This commit is contained in:
Ruben Fiszel
2024-03-23 10:54:44 +01:00
parent 461243a7a5
commit 75e9e67d7a
38 changed files with 665 additions and 296 deletions
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT running, substr(concat(coalesce(queue.logs, ''), job_logs.logs), $1) as logs, mem_peak, \n CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status \n FROM queue\n LEFT JOIN job_logs ON job_logs.job_id = queue.id \n WHERE queue.workspace_id = $2 AND queue.id = $3",
"query": "SELECT running, substr(concat(coalesce(queue.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) as logs, mem_peak, \n CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status,\n job_logs.log_offset + char_length(job_logs.logs) + 1 as log_offset\n FROM queue\n LEFT JOIN job_logs ON job_logs.job_id = queue.id \n WHERE queue.workspace_id = $2 AND queue.id = $3",
"describe": {
"columns": [
{
@@ -22,6 +22,11 @@
"ordinal": 3,
"name": "flow_status",
"type_info": "Jsonb"
},
{
"ordinal": 4,
"name": "log_offset",
"type_info": "Int4"
}
],
"parameters": {
@@ -35,8 +40,9 @@
false,
null,
true,
null,
null
]
},
"hash": "04d9cb2edf6933a3b8efbe274e10227a37475ff4ec351577bfdade19a15596d4"
"hash": "04be51a152d7c9644f11173da2cc386a71e178685364e7da4b910d1648ea55ba"
}
@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT CONCAT(coalesce(completed_job.logs, ''), coalesce(job_logs.logs, '')) \n FROM completed_job \n LEFT JOIN job_logs ON job_logs.job_id = completed_job.id \n WHERE completed_job.id = $1 AND completed_job.workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "concat",
"type_info": "Text"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
null
]
},
"hash": "3dfa0bf34c38b5529f2bbad405d23bda84db69757e132859d7e04df38e503f3e"
}
@@ -0,0 +1,18 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE job_logs SET logs = $1, log_offset = $2, \n log_file_index = array_append(coalesce(log_file_index, array[]::text[]), $3) \n WHERE workspace_id = $4 AND job_id = $5",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Int4",
"Text",
"Text",
"Uuid"
]
},
"nullable": []
},
"hash": "528cdbb75f1c5135170a58fce3fda464be138272487639d0ffbbbe6961ec5c37"
}
@@ -1,36 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT substr(concat(coalesce(completed_job.logs, ''), job_logs.logs), $1) as logs, mem_peak, \n CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status \n FROM completed_job \n LEFT JOIN job_logs ON job_logs.job_id = completed_job.id \n WHERE completed_job.workspace_id = $2 AND id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "logs",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "mem_peak",
"type_info": "Int4"
},
{
"ordinal": 2,
"name": "flow_status",
"type_info": "Jsonb"
}
],
"parameters": {
"Left": [
"Int4",
"Text",
"Uuid"
]
},
"nullable": [
null,
true,
null
]
},
"hash": "5cb644f89a94b6e6a1d7a84155bafc460c1749946056ded8ca712dacecbf427d"
}
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "SELECT char_length(logs) FROM job_logs WHERE job_id = $1 AND workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "char_length",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
null
]
},
"hash": "9a5c7e8b60a260085b438bd300972ebf948ea26f313ee5a73b85574becdd7dc7"
}
@@ -0,0 +1,42 @@
{
"db_name": "PostgreSQL",
"query": "SELECT substr(concat(coalesce(completed_job.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) as logs, mem_peak, \n CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status,\n job_logs.log_offset + char_length(job_logs.logs) + 1 as log_offset\n FROM completed_job \n LEFT JOIN job_logs ON job_logs.job_id = completed_job.id \n WHERE completed_job.workspace_id = $2 AND id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "logs",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "mem_peak",
"type_info": "Int4"
},
{
"ordinal": 2,
"name": "flow_status",
"type_info": "Jsonb"
},
{
"ordinal": 3,
"name": "log_offset",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Int4",
"Text",
"Uuid"
]
},
"nullable": [
null,
true,
null,
null
]
},
"hash": "a73b57229602d68cc25a8d963753271619e74df9fd7fc6cb05fc0614b10e6001"
}
@@ -0,0 +1,35 @@
{
"db_name": "PostgreSQL",
"query": "SELECT CONCAT(coalesce(queue.logs, ''), coalesce(job_logs.logs, '')) as logs, job_logs.log_offset, job_logs.log_file_index\n FROM queue \n LEFT JOIN job_logs ON job_logs.job_id = queue.id \n WHERE queue.id = $1 AND queue.workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "logs",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "log_offset",
"type_info": "Int4"
},
{
"ordinal": 2,
"name": "log_file_index",
"type_info": "TextArray"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
null,
false,
true
]
},
"hash": "bd213fca18a04d9e34405fd753c02398b72afb17b1ce650dbac947490634009d"
}
@@ -0,0 +1,35 @@
{
"db_name": "PostgreSQL",
"query": "SELECT CONCAT(coalesce(completed_job.logs, ''), coalesce(job_logs.logs, '')) as logs, job_logs.log_offset, job_logs.log_file_index\n FROM completed_job \n LEFT JOIN job_logs ON job_logs.job_id = completed_job.id \n WHERE completed_job.id = $1 AND completed_job.workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "logs",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "log_offset",
"type_info": "Int4"
},
{
"ordinal": 2,
"name": "log_file_index",
"type_info": "TextArray"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
null,
false,
true
]
},
"hash": "d0ce6a8ff7b89a0e99902989e57e7cae736e346c067261de34100f4fb97ca6db"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT right(logs, 300) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1",
"query": "SELECT right(logs, 600) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1",
"describe": {
"columns": [
{
@@ -19,5 +19,5 @@
null
]
},
"hash": "ab3c979c20a8ba9f6d906193c3bc7bc1941f0ecbe7ee2d8a62d2e02faeaae3f0"
"hash": "eb610cc508048cdcc030369784c274c828d67ee46ae7260b3c2fb04901b782f2"
}
@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT CONCAT(coalesce(queue.logs, ''), coalesce(job_logs.logs, '')) \n FROM queue \n LEFT JOIN job_logs ON job_logs.job_id = queue.id \n WHERE queue.id = $1 AND queue.workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "concat",
"type_info": "Text"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
null
]
},
"hash": "efc22aae27c30894f3d8ad9484f48439711707dba1f3112b495b1b18cbd31b31"
}
+1
View File
@@ -9664,6 +9664,7 @@ dependencies = [
"argon2",
"async-oauth2",
"async-recursion",
"async-stream",
"async-stripe",
"async_zip",
"axum",
+2 -1
View File
@@ -238,4 +238,5 @@ aws-sdk-sts = "^1"
crc = "^3"
tar = "^0"
http = "^1"
http = "^1"
async-stream = "^0"
@@ -0,0 +1,3 @@
-- Add down migration script here
ALTER TABLE job_logs DROP COLUMN log_offset
ALTER TABLE job_logs DROP COLUMN log_file_index;
@@ -0,0 +1,4 @@
-- Add up migration script here
ALTER TABLE job_logs ADD COLUMN log_offset int NOT NULL DEFAULT 0;
ALTER TABLE job_logs ADD COLUMN log_file_index text[];
@@ -142,8 +142,6 @@ pub fn otyp_to_string(otyp: Option<String>) -> String {
#[cfg(test)]
mod tests {
use windmill_parser::{Arg, MainArgSignature, ObjectProperty, Typ};
use super::*;
#[test]
@@ -1,4 +1,4 @@
use serde_json::{self, json};
use serde_json::json;
use wasm_bindgen::prelude::*;
use windmill_parser::MainArgSignature;
use windmill_parser_ts::{parse_expr_for_ids, parse_expr_for_imports};
+8 -11
View File
@@ -24,8 +24,7 @@ use windmill_common::{
JOB_DEFAULT_TIMEOUT_SECS_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING,
NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, PIP_INDEX_URL_SETTING,
REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING,
RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING,
SCIM_TOKEN_SETTING,
RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING, SCIM_TOKEN_SETTING,
},
stats_ee::schedule_stats,
utils::{rd_string, Mode},
@@ -40,10 +39,9 @@ use windmill_common::METRICS_ADDR;
use windmill_common::global_settings::OBJECT_STORE_CACHE_CONFIG_SETTING;
use windmill_worker::{
BUN_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM,
GO_BIN_CACHE_DIR,
GO_CACHE_DIR, HUB_CACHE_DIR, LOCK_CACHE_DIR,
PIP_CACHE_DIR, TAR_PIP_CACHE_DIR, POWERSHELL_CACHE_DIR,
BUN_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, GO_BIN_CACHE_DIR,
GO_CACHE_DIR, HUB_CACHE_DIR, LOCK_CACHE_DIR, PIP_CACHE_DIR, POWERSHELL_CACHE_DIR,
TAR_PIP_CACHE_DIR, TMP_LOGS_DIR,
};
use crate::monitor::{
@@ -51,8 +49,8 @@ use crate::monitor::{
monitor_db, monitor_pool, reload_base_url_setting, reload_bunfig_install_scopes_setting,
reload_extra_pip_index_url_setting, reload_job_default_timeout_setting, reload_license_key,
reload_npm_config_registry_setting, reload_pip_index_url_setting,
reload_retention_period_setting, reload_scim_token_setting,
reload_server_config, reload_worker_config,
reload_retention_period_setting, reload_scim_token_setting, reload_server_config,
reload_worker_config,
};
#[cfg(feature = "parquet")]
@@ -666,10 +664,9 @@ pub async fn run_workers<R: rsmq_async::RsmqConnection + Send + Sync + Clone + '
let mut handles = Vec::with_capacity(num_workers as usize);
for x in [
LOCK_CACHE_DIR,
TMP_LOGS_DIR,
PIP_CACHE_DIR,
TAR_PIP_CACHE_DIR,
DENO_CACHE_DIR,
@@ -679,7 +676,7 @@ pub async fn run_workers<R: rsmq_async::RsmqConnection + Send + Sync + Clone + '
GO_CACHE_DIR,
GO_BIN_CACHE_DIR,
HUB_CACHE_DIR,
POWERSHELL_CACHE_DIR
POWERSHELL_CACHE_DIR,
] {
DirBuilder::new()
.recursive(true)
+53 -51
View File
@@ -19,18 +19,25 @@ use windmill_api::{
DEFAULT_BODY_LIMIT, IS_SECURE, OAUTH_CLIENTS, REQUEST_SIZE_LIMIT, SAML_METADATA, SCIM_TOKEN,
};
use windmill_common::{
error, flow_status::FlowStatusModule, global_settings::{
error,
flow_status::FlowStatusModule,
global_settings::{
BASE_URL_SETTING, BUNFIG_INSTALL_SCOPES_SETTING, DEFAULT_TAGS_PER_WORKSPACE_SETTING,
EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING,
JOB_DEFAULT_TIMEOUT_SECS_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING,
NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, PIP_INDEX_URL_SETTING,
REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING,
RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING,
SCIM_TOKEN_SETTING,
}, jobs::QueuedJob, oauth2::REQUIRE_PREEXISTING_USER_FOR_OAUTH, server::load_server_config, users::truncate_token, worker::{
RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING, SCIM_TOKEN_SETTING,
},
jobs::QueuedJob,
oauth2::REQUIRE_PREEXISTING_USER_FOR_OAUTH,
server::load_server_config,
users::truncate_token,
worker::{
load_worker_config, reload_custom_tags_setting, DEFAULT_TAGS_PER_WORKSPACE, SERVER_CONFIG,
WORKER_CONFIG,
}, BASE_URL, DB, METRICS_DEBUG_ENABLED, METRICS_ENABLED
},
BASE_URL, DB, METRICS_DEBUG_ENABLED, METRICS_ENABLED,
};
use windmill_queue::cancel_job;
use windmill_worker::{
@@ -40,7 +47,10 @@ use windmill_worker::{
};
#[cfg(feature = "parquet")]
use windmill_common::s3_helpers::{build_object_store_from_settings, build_s3_client_from_settings, OBJECT_STORE_CACHE_SETTINGS, S3Settings};
use windmill_common::s3_helpers::{
build_object_store_from_settings, build_s3_client_from_settings, S3Settings,
OBJECT_STORE_CACHE_SETTINGS,
};
#[cfg(feature = "parquet")]
use windmill_common::global_settings::OBJECT_STORE_CACHE_CONFIG_SETTING;
@@ -159,7 +169,8 @@ pub async fn load_metrics_enabled(db: &DB) -> error::Result<()> {
}
pub async fn load_tag_per_workspace_enabled(db: &DB) -> error::Result<()> {
let metrics_enabled = load_value_from_global_settings(db, DEFAULT_TAGS_PER_WORKSPACE_SETTING).await;
let metrics_enabled =
load_value_from_global_settings(db, DEFAULT_TAGS_PER_WORKSPACE_SETTING).await;
match metrics_enabled {
Ok(Some(serde_json::Value::Bool(t))) => {
@@ -171,10 +182,7 @@ pub async fn load_tag_per_workspace_enabled(db: &DB) -> error::Result<()> {
}
pub async fn load_metrics_debug_enabled(db: &DB) -> error::Result<()> {
let metrics_enabled = load_value_from_global_settings(db,
EXPOSE_DEBUG_METRICS_SETTING
)
.await;
let metrics_enabled = load_value_from_global_settings(db, EXPOSE_DEBUG_METRICS_SETTING).await;
match metrics_enabled {
Ok(Some(serde_json::Value::Bool(t))) => METRICS_DEBUG_ENABLED.store(t, Ordering::Relaxed),
_ => (),
@@ -183,10 +191,7 @@ pub async fn load_metrics_debug_enabled(db: &DB) -> error::Result<()> {
}
pub async fn load_keep_job_dir(db: &DB) {
let value = load_value_from_global_settings(db,
KEEP_JOB_DIR_SETTING
)
.await;
let value = load_value_from_global_settings(db, KEEP_JOB_DIR_SETTING).await;
match value {
Ok(Some(serde_json::Value::Bool(t))) => KEEP_JOB_DIR.store(t, Ordering::Relaxed),
Err(e) => {
@@ -197,9 +202,8 @@ pub async fn load_keep_job_dir(db: &DB) {
}
pub async fn load_require_preexisting_user(db: &DB) {
let value = load_value_from_global_settings(db,
REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING
).await;
let value =
load_value_from_global_settings(db, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING).await;
match value {
Ok(Some(serde_json::Value::Bool(t))) => {
REQUIRE_PREEXISTING_USER_FOR_OAUTH.store(t, Ordering::Relaxed)
@@ -394,16 +398,22 @@ pub async fn reload_retention_period_setting(db: &DB) {
}
}
#[cfg(feature = "parquet")]
pub async fn reload_s3_cache_setting(db: &DB) {
use windmill_common::s3_helpers::ObjectSettings;
use windmill_common::{
ee::{get_license_plan, LicensePlan},
s3_helpers::ObjectSettings,
};
let s3_config = load_value_from_global_settings(db, OBJECT_STORE_CACHE_CONFIG_SETTING).await;
if let Err(e) = s3_config {
tracing::error!("Error reloading s3 cache config: {:?}", e)
} else {
if let Some(v) = s3_config.unwrap() {
if matches!(get_license_plan().await, LicensePlan::Pro) {
tracing::error!("S3 cache is not available for pro plan");
return;
}
let mut s3_cache_settings = OBJECT_STORE_CACHE_SETTINGS.write().await;
let setting = serde_json::from_value::<ObjectSettings>(v);
if let Err(e) = setting {
@@ -419,15 +429,21 @@ pub async fn reload_s3_cache_setting(db: &DB) {
} else {
let mut s3_cache_settings = OBJECT_STORE_CACHE_SETTINGS.write().await;
if std::env::var("S3_CACHE_BUCKET").is_ok() {
if matches!(get_license_plan().await, LicensePlan::Pro) {
tracing::error!("S3 cache is not available for pro plan");
return;
}
*s3_cache_settings = build_s3_client_from_settings(S3Settings {
bucket: None,
region: None,
bucket: None,
region: None,
access_key: None,
secret_key: None,
endpoint: None,
store_logs: None,
allow_http: None
}).await.ok();
allow_http: None,
})
.await
.ok();
} else {
*s3_cache_settings = None;
}
@@ -461,9 +477,7 @@ pub async fn reload_request_size(db: &DB) {
}
pub async fn reload_license_key(db: &DB) -> error::Result<()> {
let q = load_value_from_global_settings(db,
LICENSE_KEY_SETTING
).await?;
let q = load_value_from_global_settings(db, LICENSE_KEY_SETTING).await?;
let mut value = std::env::var("LICENSE_KEY")
.ok()
@@ -498,13 +512,17 @@ pub async fn reload_option_setting_with_tracing<T: FromStr + DeserializeOwned>(
}
}
async fn load_value_from_global_settings(db: &DB, setting_name: &str) -> error::Result<Option<serde_json::Value>> {
async fn load_value_from_global_settings(
db: &DB,
setting_name: &str,
) -> error::Result<Option<serde_json::Value>> {
let r = sqlx::query!(
"SELECT value FROM global_settings WHERE name = $1",
setting_name
)
.fetch_optional(db)
.await?.map(|x| x.value);
.await?
.map(|x| x.value);
Ok(r)
}
pub async fn reload_option_setting<T: FromStr + DeserializeOwned>(
@@ -521,10 +539,7 @@ pub async fn reload_option_setting<T: FromStr + DeserializeOwned>(
if let Some(q) = q {
if let Ok(v) = serde_json::from_value::<T>(q.clone()) {
tracing::info!(
"Loaded setting {setting_name} from db config: {:#?}",
&q
);
tracing::info!("Loaded setting {setting_name} from db config: {:#?}", &q);
value = Some(v)
} else {
tracing::error!("Could not parse {setting_name} found: {:#?}", &q);
@@ -559,10 +574,7 @@ pub async fn reload_setting<T: FromStr + DeserializeOwned + Display>(
if let Some(q) = q {
if let Ok(v) = serde_json::from_value::<T>(q.clone()) {
tracing::info!(
"Loaded setting {setting_name} from db config: {:#?}",
&q
);
tracing::info!("Loaded setting {setting_name} from db config: {:#?}", &q);
value = transformer(v);
} else {
tracing::error!("Could not parse {setting_name} found: {:#?}", &q);
@@ -738,9 +750,7 @@ pub async fn reload_worker_config(
}
pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> {
let q_base_url = load_value_from_global_settings(db,
BASE_URL_SETTING
).await?;
let q_base_url = load_value_from_global_settings(db, BASE_URL_SETTING).await?;
let std_base_url = std::env::var("BASE_URL")
.ok()
@@ -763,21 +773,13 @@ pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> {
std_base_url
};
let q_oauth = load_value_from_global_settings(db,
OAUTH_SETTING
)
.await?;
let q_oauth = load_value_from_global_settings(db, OAUTH_SETTING).await?;
let oauths = if let Some(q) = q_oauth {
if let Ok(v) =
serde_json::from_value::<Option<HashMap<String, OAuthClient>>>(q.clone())
{
if let Ok(v) = serde_json::from_value::<Option<HashMap<String, OAuthClient>>>(q.clone()) {
v
} else {
tracing::error!(
"Could not parse oauth setting as a json, found: {:#?}",
&q
);
tracing::error!("Could not parse oauth setting as a json, found: {:#?}", &q);
None
}
} else {
+1 -4
View File
@@ -171,10 +171,8 @@ fn find_module_in_vec(modules: Vec<FlowStatusModule>, id: &str) -> Option<FlowSt
mod suspend_resume {
use futures::{Stream, StreamExt};
use serde_json::json;
use sqlx::{query_scalar, types::Uuid};
use windmill_common::{flows::FlowValue, jobs::JobPayload};
use sqlx::query_scalar;
use super::*;
@@ -437,7 +435,6 @@ mod suspend_resume {
mod retry {
use serde_json::json;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use windmill_common::flows::FlowValue;
use super::*;
+1
View File
@@ -91,3 +91,4 @@ openidconnect = { workspace = true, optional = true}
pin-project.workspace = true
crc.workspace = true
http.workspace = true
async-stream.workspace = true
+2
View File
@@ -5467,6 +5467,8 @@ paths:
type: boolean
new_logs:
type: string
log_offset:
type: integer
mem_peak:
type: integer
flow_status:
+1 -6
View File
@@ -27,7 +27,6 @@ use hyper::StatusCode;
use serde::{Deserialize, Serialize};
use serde_json::json;
use sql_builder::prelude::*;
use sql_builder::SqlBuilder;
use sqlx::{FromRow, Postgres, Transaction};
use windmill_audit::audit_ee::audit_log;
use windmill_audit::ActionKind;
@@ -173,11 +172,7 @@ async fn list_hub_flows(Extension(db): Extension<DB>) -> impl IntoResponse {
&db,
)
.await?;
Ok::<_, Error>((
status_code,
headers,
response
))
Ok::<_, Error>((status_code, headers, response))
}
async fn list_paths(
-1
View File
@@ -1,7 +1,6 @@
use crate::db::DB;
use axum::{extract::Path, routing::post, Extension, Json, Router};
use hyper::http;
use serde::{Deserialize, Serialize};
use tower_http::cors::{Any, CorsLayer};
use uuid::Uuid;
+65 -17
View File
@@ -6,6 +6,7 @@
* LICENSE-AGPL for a copy of the license.
*/
use axum::body::Body;
use axum::http::HeaderValue;
use serde_json::value::RawValue;
use std::collections::HashMap;
@@ -37,9 +38,9 @@ use axum::{
use base64::Engine;
use chrono::Utc;
use hmac::Mac;
use hyper::{http, Request, StatusCode};
use hyper::{Request, StatusCode};
use serde::{de::DeserializeOwned, Deserialize, Serialize};
use sql_builder::{prelude::*, quote, SqlBuilder};
use sql_builder::prelude::*;
use sqlx::types::JsonRawValue;
use sqlx::{types::Uuid, FromRow, Postgres, Transaction};
use tower_http::cors::{Any, CorsLayer};
@@ -59,6 +60,8 @@ use windmill_common::{
utils::{not_found_if_none, now_from_db, paginate, require_admin, Pagination, StripPath},
};
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS;
#[cfg(feature = "prometheus")]
use windmill_common::{METRICS_DEBUG_ENABLED, METRICS_ENABLED};
@@ -591,12 +594,40 @@ async fn get_job_internal(db: &DB, workspace_id: &str, job_id: Uuid) -> error::R
}
}
#[cfg(all(feature = "enterprise", feature = "parquet"))]
async fn get_logs_from_store(
log_offset: i32,
logs: &str,
log_file_index: Option<Vec<String>>,
) -> Option<error::Result<Body>> {
if log_offset > 0 {
if let Some(file_index) = log_file_index {
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
let logs = logs.to_string();
let stream = async_stream::stream! {
for file in file_index {
let file = os.get(&object_store::path::Path::from(file)).await;
if let Ok(file) = file {
if let Ok(bytes) = file.bytes().await {
yield Ok(bytes::Bytes::from(bytes)) as object_store::Result<bytes::Bytes>;
}
}
}
yield Ok(bytes::Bytes::from(logs))
};
return Some(Ok(Body::from_stream(stream)));
}
}
}
return None;
}
async fn get_job_logs(
Extension(db): Extension<DB>,
Path((w_id, id)): Path<(String, Uuid)>,
) -> error::Result<String> {
let text = sqlx::query_scalar!(
"SELECT CONCAT(coalesce(completed_job.logs, ''), coalesce(job_logs.logs, ''))
) -> error::Result<Body> {
let record = sqlx::query!(
"SELECT CONCAT(coalesce(completed_job.logs, ''), coalesce(job_logs.logs, '')) as logs, job_logs.log_offset, job_logs.log_file_index
FROM completed_job
LEFT JOIN job_logs ON job_logs.job_id = completed_job.id
WHERE completed_job.id = $1 AND completed_job.workspace_id = $2",
@@ -604,23 +635,35 @@ async fn get_job_logs(
w_id
)
.fetch_optional(&db)
.await?
.flatten();
if let Some(text) = text {
Ok(text)
.await?;
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
{
return r;
}
Ok(Body::from(logs))
} else {
let text = sqlx::query_scalar!(
"SELECT CONCAT(coalesce(queue.logs, ''), coalesce(job_logs.logs, ''))
let text = sqlx::query!(
"SELECT CONCAT(coalesce(queue.logs, ''), coalesce(job_logs.logs, '')) as logs, job_logs.log_offset, job_logs.log_file_index
FROM queue
LEFT JOIN job_logs ON job_logs.job_id = queue.id
WHERE queue.id = $1 AND queue.workspace_id = $2",
id,
w_id
)
.fetch_one(&db)
.fetch_optional(&db)
.await?;
let text = not_found_if_none(text, "Job Logs", id.to_string())?;
Ok(text)
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 {
return r;
}
Ok(Body::from(logs))
}
}
@@ -3301,6 +3344,7 @@ pub struct JobUpdate {
pub running: Option<bool>,
pub completed: Option<bool>,
pub new_logs: Option<String>,
pub log_offset: Option<i32>,
pub mem_peak: Option<i32>,
pub flow_status: Option<serde_json::Value>,
}
@@ -3311,8 +3355,9 @@ async fn get_job_update(
Query(JobUpdateQuery { running, log_offset }): Query<JobUpdateQuery>,
) -> error::JsonResult<JobUpdate> {
let record = sqlx::query!(
"SELECT running, substr(concat(coalesce(queue.logs, ''), job_logs.logs), $1) as logs, mem_peak,
CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status
"SELECT running, substr(concat(coalesce(queue.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) as logs, mem_peak,
CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status,
job_logs.log_offset + char_length(job_logs.logs) + 1 as log_offset
FROM queue
LEFT JOIN job_logs ON job_logs.job_id = queue.id
WHERE queue.workspace_id = $2 AND queue.id = $3",
@@ -3330,6 +3375,7 @@ async fn get_job_update(
} else {
None
},
log_offset: record.log_offset,
completed: None,
new_logs: record.logs,
mem_peak: record.mem_peak,
@@ -3337,8 +3383,9 @@ async fn get_job_update(
}))
} else {
let record = sqlx::query!(
"SELECT substr(concat(coalesce(completed_job.logs, ''), job_logs.logs), $1) as logs, mem_peak,
CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status
"SELECT substr(concat(coalesce(completed_job.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) as logs, mem_peak,
CASE WHEN is_flow_step is true then NULL else flow_status END as flow_status,
job_logs.log_offset + char_length(job_logs.logs) + 1 as log_offset
FROM completed_job
LEFT JOIN job_logs ON job_logs.job_id = completed_job.id
WHERE completed_job.workspace_id = $2 AND id = $3",
@@ -3352,6 +3399,7 @@ async fn get_job_update(
Ok(Json(JobUpdate {
running: Some(false),
completed: Some(true),
log_offset: record.log_offset,
new_logs: record.logs,
mem_peak: record.mem_peak,
flow_status: record.flow_status,
+8 -11
View File
@@ -23,7 +23,6 @@ use axum::extract::DefaultBodyLimit;
use axum::{middleware::from_extractor, routing::get, Extension, Router};
use db::DB;
use git_version::git_version;
use hyper::http;
use reqwest::Client;
use std::collections::HashMap;
use std::{net::SocketAddr, sync::Arc};
@@ -175,31 +174,28 @@ pub async fn run_server(
let embeddings_db = if server_mode {
#[cfg(feature = "embedding")]
{
Some(load_embeddings_db(&db))
Some(load_embeddings_db(&db))
}
#[cfg(not(feature = "embedding"))]
{
Some(())
Some(())
}
} else {
None
};
let job_helpers_service = {
#[cfg(feature = "parquet")]
{
job_helpers_ee::workspaced_service()
job_helpers_ee::workspaced_service()
}
#[cfg(not(feature = "parquet"))]
{
Router::new()
Router::new()
}
};
// build our application with a route
let app = Router::new()
.nest(
@@ -309,14 +305,15 @@ pub async fn run_server(
let instance_name = rd_string(5);
let listener = tokio::net::TcpListener::bind(addr).await.unwrap();
let port = listener.local_addr().map(|x| x.port()).unwrap_or(8000);
let ip = listener.local_addr().map(|x| x.ip().to_string()).unwrap_or("localhost".to_string());
let ip = listener
.local_addr()
.map(|x| x.ip().to_string())
.unwrap_or("localhost".to_string());
let server = axum::serve(listener, app.into_make_service());
tracing::info!(
instance = %instance_name,
"server started on port={} and addr={}",
+1 -1
View File
@@ -31,7 +31,7 @@ use windmill_common::{
utils::{not_found_if_none, paginate, Pagination, StripPath},
};
use windmill_git_sync::{handle_deployment_metadata, DeployedObject};
use windmill_queue::{self, schedule::push_scheduled_job, QueueTransaction};
use windmill_queue::{schedule::push_scheduled_job, QueueTransaction};
pub fn workspaced_service() -> Router {
Router::new()
+2 -7
View File
@@ -23,7 +23,6 @@ use hyper::StatusCode;
use serde::{Deserialize, Serialize};
use serde_json::json;
use sql_builder::prelude::*;
use sql_builder::SqlBuilder;
use sqlx::{FromRow, Postgres, Transaction};
use std::{
collections::{hash_map::DefaultHasher, HashMap},
@@ -48,7 +47,7 @@ use windmill_common::{
};
use windmill_git_sync::{handle_deployment_metadata, DeployedObject};
use windmill_parser_ts::remove_pinned_imports;
use windmill_queue::{self, schedule::push_scheduled_job, PushIsolationLevel, QueueTransaction};
use windmill_queue::{schedule::push_scheduled_job, PushIsolationLevel, QueueTransaction};
const MAX_HASH_HISTORY_LENGTH_STORED: usize = 20;
@@ -294,11 +293,7 @@ async fn get_top_hub_scripts(
&db,
)
.await?;
Ok::<_, Error>((
status_code,
headers,
response,
))
Ok::<_, Error>((status_code, headers, response))
}
fn hash_script(ns: &NewScript) -> i64 {
+29 -12
View File
@@ -31,7 +31,6 @@ use windmill_common::{
};
pub fn global_service() -> Router {
#[warn(unused_mut)]
let r = Router::new()
.route("/envs", get(get_local_settings))
@@ -42,7 +41,7 @@ pub fn global_service() -> Router {
.route("/test_smtp", post(test_email))
.route("/test_license_key", post(test_license_key))
.route("/send_stats", post(send_stats));
#[cfg(feature = "parquet")]
{
return r.route("/test_s3_config", post(test_s3_bucket));
@@ -50,9 +49,8 @@ pub fn global_service() -> Router {
#[cfg(not(feature = "parquet"))]
{
return r
return r;
}
}
#[derive(Deserialize)]
@@ -106,8 +104,6 @@ use windmill_common::s3_helpers::ObjectSettings;
#[cfg(feature = "parquet")]
use windmill_common::s3_helpers::build_object_store_from_settings;
#[cfg(feature = "parquet")]
pub async fn test_s3_bucket(
Extension(db): Extension<DB>,
@@ -115,16 +111,37 @@ pub async fn test_s3_bucket(
Json(test_s3_bucket): Json<ObjectSettings>,
) -> error::Result<String> {
use bytes::Bytes;
use windmill_common::ee::{get_license_plan, LicensePlan};
if matches!(get_license_plan().await, LicensePlan::Pro) {
return Err(error::Error::InternalErr(
"This feature is only available in Enterprise, not Pro".to_string(),
));
}
require_super_admin(&db, &authed.email).await?;
let client = build_object_store_from_settings(test_s3_bucket).await?;
let path = object_store::path::Path::from(format!("/test-s3-bucket-{uuid}", uuid = uuid::Uuid::new_v4()));
let path = object_store::path::Path::from(format!(
"/test-s3-bucket-{uuid}",
uuid = uuid::Uuid::new_v4()
));
tracing::info!("Testing s3 bucket at path: {path}");
client.put(&path, Bytes::from_static(b"hello")).await.map_err(to_anyhow)?;
let content = client.get(&path).await.map_err(to_anyhow)?.bytes().await.map_err(to_anyhow)?;
client
.put(&path, Bytes::from_static(b"hello"))
.await
.map_err(to_anyhow)?;
let content = client
.get(&path)
.await
.map_err(to_anyhow)?
.bytes()
.await
.map_err(to_anyhow)?;
if content != Bytes::from_static(b"hello") {
return Err(error::Error::InternalErr("Failed to read back from s3".to_string()));
return Err(error::Error::InternalErr(
"Failed to read back from s3".to_string(),
));
}
client.delete(&path).await.map_err(to_anyhow)?;
Ok("Tested bucket successfully".to_string())
@@ -162,7 +179,7 @@ pub async fn get_local_settings(
#[derive(serde::Deserialize)]
pub struct Value {
pub value: serde_json::Value,
pub value: Option<serde_json::Value>,
}
pub async fn delete_global_setting(db: &DB, key: &str) -> error::Result<()> {
@@ -179,7 +196,7 @@ pub async fn set_global_setting(
Json(value): Json<Value>,
) -> error::Result<()> {
require_super_admin(&db, &authed.email).await?;
set_global_setting_internal(&db, key, value.value).await
set_global_setting_internal(&db, key, value.value.unwrap_or(serde_json::Value::Null)).await
}
pub async fn set_global_setting_internal(
+1 -1
View File
@@ -25,7 +25,7 @@ use argon2::{password_hash::SaltString, Argon2, PasswordHash, PasswordHasher, Pa
use axum::{
async_trait,
extract::{Extension, FromRequestParts, OriginalUri, Path, Query},
http::{self, request::Parts},
http::request::Parts,
response::{IntoResponse, Response},
routing::{delete, get, post},
Json, Router,
+4 -7
View File
@@ -13,7 +13,7 @@ use std::{
};
use rand::Rng;
use serde::{self, Deserialize, Serialize, Serializer};
use serde::{Deserialize, Serialize, Serializer};
use crate::{
more_serde::{
@@ -22,8 +22,7 @@ use crate::{
scripts::{Schema, ScriptHash, ScriptLang},
};
#[derive(Serialize)]
#[derive(sqlx::FromRow)]
#[derive(Serialize, sqlx::FromRow)]
pub struct Flow {
pub workspace_id: String,
pub path: String,
@@ -47,8 +46,7 @@ pub struct Flow {
pub timeout: Option<i32>,
}
#[derive(Serialize)]
#[derive(sqlx::FromRow)]
#[derive(Serialize, sqlx::FromRow)]
pub struct ListableFlow {
pub workspace_id: String,
pub path: String,
@@ -66,8 +64,7 @@ pub struct ListableFlow {
pub ws_error_handler_muted: Option<bool>,
}
#[derive(Deserialize)]
#[derive(sqlx::FromRow)]
#[derive(Deserialize, sqlx::FromRow)]
pub struct NewFlow {
pub path: String,
pub summary: String,
+2 -2
View File
@@ -227,7 +227,7 @@ pub async fn cancel_job<'c: 'async_recursion>(
#[tracing::instrument(level = "trace", skip_all)]
pub async fn append_logs(
job_id: uuid::Uuid,
workspace: String,
workspace: impl AsRef<str>,
logs: impl AsRef<str>,
db: impl Borrow<Pool<Postgres>>,
) {
@@ -243,7 +243,7 @@ pub async fn append_logs(
"INSERT INTO job_logs (logs, job_id, workspace_id) VALUES ($1, $2, $3) ON CONFLICT (job_id) DO UPDATE SET logs = concat(job_logs.logs, $1::text)",
logs.as_ref(),
job_id,
workspace,
workspace.as_ref(),
)
.execute(db.borrow())
.await
+220 -2
View File
@@ -1,4 +1,5 @@
use async_recursion::async_recursion;
use deno_ast::swc::parser::lexer::util::CharExt;
use futures::Future;
use itertools::Itertools;
@@ -7,6 +8,8 @@ use nix::sys::signal::{self, Signal};
#[cfg(any(target_os = "linux", target_os = "macos"))]
use nix::unistd::Pid;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use object_store::path::Path;
use regex::Regex;
use serde::{Deserialize, Serialize};
use serde_json::value::RawValue;
@@ -17,6 +20,8 @@ use tokio::process::Command;
use tokio::{fs::File, io::AsyncReadExt};
use windmill_common::error::to_anyhow;
use windmill_common::jobs::ENTRYPOINT_OVERRIDE;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS;
#[cfg(feature = "parquet")]
use windmill_common::s3_helpers::{
get_etag_or_empty, LargeFileStorage, ObjectStoreResource, S3Object,
@@ -34,6 +39,8 @@ use windmill_queue::{append_logs, CanceledBy};
#[cfg(any(target_os = "linux", target_os = "macos"))]
use std::os::unix::process::ExitStatusExt;
use std::sync::atomic::AtomicU32;
use std::sync::Arc;
use std::{
collections::{hash_map::DefaultHasher, HashMap},
hash::{Hash, Hasher},
@@ -62,7 +69,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,
MAX_TIMEOUT_DURATION, MAX_WAIT_FOR_SIGINT, MAX_WAIT_FOR_SIGTERM, ROOT_CACHE_DIR, TMP_DIR,
};
pub async fn build_args_map<'a>(
@@ -588,6 +595,202 @@ pub async fn update_job_poller<F, Fut>(
tracing::info!("job {job_id} finished");
}
pub enum CompactLogs {
NotEE,
NoS3,
S3,
}
async fn compact_logs(
job_id: Uuid,
w_id: &str,
db: &DB,
nlogs: String,
total_size: Arc<AtomicU32>,
compact_kind: CompactLogs,
worker_name: &str,
) -> error::Result<(String, String)> {
let size = sqlx::query_scalar!(
"SELECT char_length(logs) FROM job_logs WHERE job_id = $1 AND workspace_id = $2",
job_id,
w_id
)
.fetch_optional(db)
.await?
.flatten()
.unwrap_or(0);
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 nlogs_len = nlogs.char_indices().count();
let modulo = nlogs_len % LARGE_LOG_THRESHOLD_SIZE;
let extra_split = modulo < nlogs_len;
let excess_size_modulo = if extra_split { nlogs_len - modulo } else { 0 };
let excess_size = excess_size_modulo
+ nlogs[excess_size_modulo..]
.chars()
.find_position(|x| x.is_line_break())
.map(|(i, _)| i + 1)
.unwrap_or(0);
let (excess_prev_logs, current_logs) = if extra_split {
let (excess_prev_logs, current_logs) = nlogs.split_at(excess_size as usize);
(excess_prev_logs, current_logs.to_string())
} else {
("", nlogs.to_string())
};
let new_size_with_excess = size + excess_size 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!("[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::S3 => format!("[windmill] worker {worker_name}: Logs length has exceeded a threshold\n[windmill] Previous logs have been saved to object storage at {path}\n[windmill] Download logs in expanded drawer to get full logs."),
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"),
};
new_current_logs.push_str(&current_logs);
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(&excess_prev_logs);
return Ok((prev_logs, path));
}
async fn default_disk_log_storage(
job_id: Uuid,
w_id: &str,
db: &DB,
nlogs: String,
total_size: Arc<AtomicU32>,
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_dir = format!("{}/{}", TMP_DIR, path);
tokio::fs::create_dir_all(&path_dir)
.await
.map_err(|e| {
tracing::error!("Could not create logs directory: {e:?}",);
e
})
.ok();
tokio::fs::write(&path, prev_logs)
.await
.map_err(|e| {
tracing::error!("Could not save logs to disk: {e:?}",);
e
})
.ok();
}
}
}
async fn append_job_logs(
job_id: Uuid,
w_id: String,
logs: String,
db: DB,
must_compact_logs: bool,
total_size: Arc<AtomicU32>,
worker_name: String,
) -> () {
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 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 saved to object storage at {path2}");
}
}
} else {
default_disk_log_storage(
job_id,
&w_id,
&db,
logs,
total_size,
CompactLogs::NoS3,
&worker_name,
)
.await;
}
#[cfg(not(all(feature = "enterprise", feature = "parquet")))]
{
default_disk_log_storage(
job_id,
&w_id,
&db,
logs,
total_size,
CompactLogs::NotEE,
&worker_name,
)
.await;
}
} else {
append_logs(job_id, w_id, logs, db).await;
}
}
pub const LARGE_LOG_THRESHOLD_SIZE: usize = 5000;
/// - wait until child exits and return with exit status
/// - read lines from stdout and stderr and append them to the "queue"."logs"
/// quitting early if output exceedes MAX_LOG_SIZE characters (not bytes)
@@ -759,6 +962,9 @@ pub async fn handle_child(
* It's useful to know if the task completed. */
let (mut do_write, mut write_result) = tokio::spawn(ready(())).remote_handle();
let mut log_total_size: u64 = 0;
let pg_log_total_size = Arc::new(AtomicU32::new(0));
while let Some(line) = output.by_ref().next().await {
let do_write_ = do_write.shared();
@@ -830,7 +1036,19 @@ pub async fn handle_child(
panic::resume_unwind(p);
}
(do_write, write_result) = tokio::spawn(append_logs(job_id, w_id.to_string(), joined, db.clone())).remote_handle();
let joined_len = joined.len() as u64;
log_total_size += joined_len;
let compact_logs = log_total_size > LARGE_LOG_THRESHOLD_SIZE as u64;
if compact_logs {
log_total_size = 0;
}
let worker_name = worker_name.to_string();
let w_id2 = w_id.to_string();
(do_write, write_result) = tokio::spawn(append_job_logs(job_id, w_id2, joined, db.clone(), compact_logs, pg_log_total_size.clone(), worker_name)).remote_handle();
if let Err(err) = result {
tracing::error!(%job_id, %err, "error reading output for job {job_id}: {err}");
+6 -6
View File
@@ -1,5 +1,5 @@
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use crate::{ROOT_CACHE_DIR, PIP_CACHE_DIR};
use crate::{PIP_CACHE_DIR, ROOT_CACHE_DIR};
// #[cfg(feature = "enterprise")]
// use rand::Rng;
@@ -7,7 +7,7 @@ use crate::{ROOT_CACHE_DIR, PIP_CACHE_DIR};
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use tokio::time::Instant;
#[cfg(feature = "parquet")]
#[cfg(all(feature = "enterprise", feature = "parquet"))]
use object_store::ObjectStore;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
@@ -17,7 +17,10 @@ use windmill_common::error;
use std::sync::Arc;
#[cfg(all(feature = "enterprise", feature = "parquet"))]
pub async fn build_tar_and_push(s3_client: Arc<dyn ObjectStore>, folder: String) -> error::Result<()> {
pub async fn build_tar_and_push(
s3_client: Arc<dyn ObjectStore>,
folder: String,
) -> error::Result<()> {
use bytes::Bytes;
use object_store::path::Path;
@@ -38,7 +41,6 @@ pub async fn build_tar_and_push(s3_client: Arc<dyn ObjectStore>, folder: String)
)));
}
// let s3_settings = S3_CACHE_SETTINGS.read().await;
// let s3_client = s3_settings.as_ref().ok_or_else(|| {
// error::Error::ExecutionErr("Failed to read s3 cache settings".to_string())
@@ -69,10 +71,8 @@ pub async fn build_tar_and_push(s3_client: Arc<dyn ObjectStore>, folder: String)
Ok(())
}
#[cfg(all(feature = "enterprise", feature = "parquet"))]
pub async fn pull_from_tar(client: Arc<dyn ObjectStore>, folder: String) -> error::Result<()> {
use object_store::path::Path;
use tokio::fs::metadata;
let folder_name = folder.split("/").last().unwrap();
+64 -53
View File
@@ -152,7 +152,13 @@ pub async fn pip_compile(
write_file(job_dir, file, &requirements).await?;
let mut args = vec!["-q", "--no-header", file, "--resolver=backtracking", "--strip-extras"];
let mut args = vec![
"-q",
"--no-header",
file,
"--resolver=backtracking",
"--strip-extras",
];
let mut pip_args = vec![];
let pip_extra_index_url = PIP_EXTRA_INDEX_URL
.read()
@@ -776,7 +782,6 @@ pub async fn handle_python_reqs(
.await?;
};
let mut req_with_penv: Vec<(String, String)> = vec![];
for req in requirements {
@@ -800,67 +805,73 @@ pub async fn handle_python_reqs(
#[cfg(all(feature = "enterprise", feature = "parquet"))]
if req_with_penv.len() > 0 {
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
if matches!(get_license_plan().await, LicensePlan::Pro) {
append_logs(job_id.clone(), w_id.to_string(), format!("s3 cache not available in Pro Plan"), db).await;
tracing::warn!("S3 cache not available in the pro plan");
} else {
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel(1);
let job_id_2 = job_id.clone();
let db_2 = db.clone();
tokio::spawn(async move {
loop {
tokio::select! {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => {
if let Err(e) = sqlx::query_scalar!("UPDATE queue SET last_ping = now() WHERE id = $1", &job_id_2)
.execute(&db_2)
.await {
tracing::error!("failed to update last_ping: {}", e);
}
}
_ = done_rx.recv() => {
break;
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel(1);
let job_id_2 = job_id.clone();
let db_2 = db.clone();
tokio::spawn(async move {
loop {
tokio::select! {
_ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => {
if let Err(e) = sqlx::query_scalar!("UPDATE queue SET last_ping = now() WHERE id = $1", &job_id_2)
.execute(&db_2)
.await {
tracing::error!("failed to update last_ping: {}", e);
}
}
}
});
let start = std::time::Instant::now();
let futures = req_with_penv.clone().into_iter().map(|(req, venv_p)| {
let os = os.clone();
async move {
if pull_from_tar(os, venv_p.clone()).await.is_ok() {
PullFromTar::Pulled(venv_p.to_string())
} else {
PullFromTar::NotPulled(req.to_string(), venv_p.to_string())
}
}}).collect::<Vec<_>>();
let results = futures::future::join_all(futures).await;
req_with_penv.clear();
done_tx.send(()).await.expect("failed to send done");
let mut pulled = vec![];
for result in results {
match result {
PullFromTar::Pulled(venv_p) => {
pulled.push(venv_p.split("/").last().unwrap_or_default().to_string());
req_paths.push(venv_p);
}
PullFromTar::NotPulled(req, venv_p) => {
req_with_penv.push((req, venv_p));
_ = done_rx.recv() => {
break;
}
}
}
if pulled.len() > 0 {
append_logs(job_id.clone(), w_id.to_string(), format!("pulled {} from s3 cache in {}ms", pulled.join(", "), start.elapsed().as_millis()), db).await;
});
let start = std::time::Instant::now();
let futures = req_with_penv
.clone()
.into_iter()
.map(|(req, venv_p)| {
let os = os.clone();
async move {
if pull_from_tar(os, venv_p.clone()).await.is_ok() {
PullFromTar::Pulled(venv_p.to_string())
} else {
PullFromTar::NotPulled(req.to_string(), venv_p.to_string())
}
}
})
.collect::<Vec<_>>();
let results = futures::future::join_all(futures).await;
req_with_penv.clear();
done_tx.send(()).await.expect("failed to send done");
let mut pulled = vec![];
for result in results {
match result {
PullFromTar::Pulled(venv_p) => {
pulled.push(venv_p.split("/").last().unwrap_or_default().to_string());
req_paths.push(venv_p);
}
PullFromTar::NotPulled(req, venv_p) => {
req_with_penv.push((req, venv_p));
}
}
}
}
if pulled.len() > 0 {
append_logs(
job_id.clone(),
w_id.to_string(),
format!(
"pulled {} from s3 cache in {}ms",
pulled.join(", "),
start.elapsed().as_millis()
),
db,
)
.await;
}
}
}
for (req, venv_p) in req_with_penv {
let mut logs1 = String::new();
logs1.push_str("\n\n--- PIP INSTALL ---\n");
logs1.push_str(&format!("\n{req} is being installed for the first time.\n It will be cached for all ulterior uses."));
@@ -2,7 +2,6 @@ use base64::{engine, Engine as _};
use core::fmt::Write;
use futures::TryFutureExt;
use jsonwebtoken::{encode, Algorithm, EncodingKey, Header};
use pem;
use serde_json::{json, value::RawValue, Value};
use sha2::{Digest, Sha256};
use windmill_common::error::to_anyhow;
+3 -1
View File
@@ -185,6 +185,8 @@ pub async fn create_token_for_owner(
}
pub const TMP_DIR: &str = "/tmp/windmill";
pub const TMP_LOGS_DIR: &str = "/tmp/windmill/logs";
pub const ROOT_CACHE_DIR: &str = "/tmp/windmill/cache/";
pub const LOCK_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "lock");
pub const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip");
@@ -2871,7 +2873,7 @@ async fn process_result(
res.unwrap()
} else {
let last_10_log_lines = sqlx::query_scalar!(
"SELECT right(logs, 300) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1",
"SELECT right(logs, 600) FROM job_logs WHERE job_id = $1 AND workspace_id = $2 ORDER BY created_at DESC LIMIT 1",
&job.id,
&job.workspace_id
).fetch_one(db).await.ok().flatten().unwrap_or("".to_string());
+11 -11
View File
@@ -2980,18 +2980,18 @@ async fn get_transform_context(
Ok(IdContext { flow_job: flow_job.id, steps_results, previous_id: previous_id.to_string() })
}
trait IntoArray: Sized {
fn into_array(self) -> Result<Vec<Value>, Self>;
}
// trait IntoArray: Sized {
// fn into_array(self) -> Result<Vec<Value>, Self>;
// }
impl IntoArray for Value {
fn into_array(self) -> Result<Vec<Value>, Self> {
match self {
Value::Array(array) => Ok(array),
not_array => Err(not_array),
}
}
}
// impl IntoArray for Value {
// fn into_array(self) -> Result<Vec<Value>, Self> {
// match self {
// Value::Array(array) => Ok(array),
// not_array => Err(not_array),
// }
// }
// }
fn from_now(duration: Duration) -> chrono::DateTime<chrono::Utc> {
// "This function errors when original duration is larger than
@@ -21,6 +21,8 @@
let syncIteration: number = 0
let errorIteration = 0
let logOffset = 0
let ITERATIONS_BEFORE_SLOW_REFRESH = 10
let ITERATIONS_BEFORE_SUPER_SLOW_REFRESH = 100
@@ -135,6 +137,7 @@
}
export async function watchJob(testId: string) {
logOffset = 0
syncIteration = 0
errorIteration = 0
currentId = testId
@@ -156,8 +159,13 @@
workspace: workspace!,
id,
running: job.running,
logOffset: job.logs?.length ? job.logs?.length + 1 : 0
logOffset: logOffset == 0 ? job?.logs?.length + 1 ?? 0 : logOffset
})
console.log(logOffset, previewJobUpdates.log_offset, previewJobUpdates.new_logs)
if (previewJobUpdates.log_offset) {
logOffset = previewJobUpdates.log_offset ?? 0
}
if (previewJobUpdates.new_logs) {
job.logs = (job?.logs ?? '').concat(previewJobUpdates.new_logs)
}