From 75e9e67d7a8fa07d1dee6c1194261cbeacb69c8d Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 23 Mar 2024 10:54:44 +0100 Subject: [PATCH] feat: large log disk and distributed storage compaction --- ...386a71e178685364e7da4b910d1648ea55ba.json} | 10 +- ...23bda84db69757e132859d7e04df38e503f3e.json | 23 -- ...da464be138272487639d0ffbbbe6961ec5c37.json | 18 ++ ...afc460c1749946056ded8ca712dacecbf427d.json | 36 --- ...72ebf948ea26f313ee5a73b85574becdd7dc7.json | 23 ++ ...3271619e74df9fd7fc6cb05fc0614b10e6001.json | 42 ++++ ...02398b72afb17b1ce650dbac947490634009d.json | 35 +++ ...e7cae736e346c067261de34100f4fb97ca6db.json | 35 +++ ...74c828d67ee46ae7260b3c2fb04901b782f2.json} | 4 +- ...48439711707dba1f3112b495b1b18cbd31b31.json | 23 -- backend/Cargo.lock | 1 + backend/Cargo.toml | 3 +- ...20240322130200_add_offset_to_logs.down.sql | 3 + .../20240322130200_add_offset_to_logs.up.sql | 4 + backend/parsers/windmill-parser-go/src/lib.rs | 2 - .../parsers/windmill-parser-wasm/src/lib.rs | 2 +- backend/src/main.rs | 19 +- backend/src/monitor.rs | 104 ++++---- backend/tests/worker.rs | 5 +- backend/windmill-api/Cargo.toml | 1 + backend/windmill-api/openapi.yaml | 2 + backend/windmill-api/src/flows.rs | 7 +- backend/windmill-api/src/job_metrics.rs | 1 - backend/windmill-api/src/jobs.rs | 82 +++++-- backend/windmill-api/src/lib.rs | 19 +- backend/windmill-api/src/schedule.rs | 2 +- backend/windmill-api/src/scripts.rs | 9 +- backend/windmill-api/src/settings.rs | 41 +++- backend/windmill-api/src/users.rs | 2 +- backend/windmill-common/src/flows.rs | 11 +- backend/windmill-queue/src/jobs.rs | 4 +- backend/windmill-worker/src/common.rs | 222 +++++++++++++++++- backend/windmill-worker/src/global_cache.rs | 12 +- .../windmill-worker/src/python_executor.rs | 117 ++++----- .../windmill-worker/src/snowflake_executor.rs | 1 - backend/windmill-worker/src/worker.rs | 4 +- backend/windmill-worker/src/worker_flow.rs | 22 +- .../src/lib/components/TestJobLoader.svelte | 10 +- 38 files changed, 665 insertions(+), 296 deletions(-) rename backend/.sqlx/{query-04d9cb2edf6933a3b8efbe274e10227a37475ff4ec351577bfdade19a15596d4.json => query-04be51a152d7c9644f11173da2cc386a71e178685364e7da4b910d1648ea55ba.json} (54%) delete mode 100644 backend/.sqlx/query-3dfa0bf34c38b5529f2bbad405d23bda84db69757e132859d7e04df38e503f3e.json create mode 100644 backend/.sqlx/query-528cdbb75f1c5135170a58fce3fda464be138272487639d0ffbbbe6961ec5c37.json delete mode 100644 backend/.sqlx/query-5cb644f89a94b6e6a1d7a84155bafc460c1749946056ded8ca712dacecbf427d.json create mode 100644 backend/.sqlx/query-9a5c7e8b60a260085b438bd300972ebf948ea26f313ee5a73b85574becdd7dc7.json create mode 100644 backend/.sqlx/query-a73b57229602d68cc25a8d963753271619e74df9fd7fc6cb05fc0614b10e6001.json create mode 100644 backend/.sqlx/query-bd213fca18a04d9e34405fd753c02398b72afb17b1ce650dbac947490634009d.json create mode 100644 backend/.sqlx/query-d0ce6a8ff7b89a0e99902989e57e7cae736e346c067261de34100f4fb97ca6db.json rename backend/.sqlx/{query-ab3c979c20a8ba9f6d906193c3bc7bc1941f0ecbe7ee2d8a62d2e02faeaae3f0.json => query-eb610cc508048cdcc030369784c274c828d67ee46ae7260b3c2fb04901b782f2.json} (70%) delete mode 100644 backend/.sqlx/query-efc22aae27c30894f3d8ad9484f48439711707dba1f3112b495b1b18cbd31b31.json create mode 100644 backend/migrations/20240322130200_add_offset_to_logs.down.sql create mode 100644 backend/migrations/20240322130200_add_offset_to_logs.up.sql diff --git a/backend/.sqlx/query-04d9cb2edf6933a3b8efbe274e10227a37475ff4ec351577bfdade19a15596d4.json b/backend/.sqlx/query-04be51a152d7c9644f11173da2cc386a71e178685364e7da4b910d1648ea55ba.json similarity index 54% rename from backend/.sqlx/query-04d9cb2edf6933a3b8efbe274e10227a37475ff4ec351577bfdade19a15596d4.json rename to backend/.sqlx/query-04be51a152d7c9644f11173da2cc386a71e178685364e7da4b910d1648ea55ba.json index 5085b2b037..976674b566 100644 --- a/backend/.sqlx/query-04d9cb2edf6933a3b8efbe274e10227a37475ff4ec351577bfdade19a15596d4.json +++ b/backend/.sqlx/query-04be51a152d7c9644f11173da2cc386a71e178685364e7da4b910d1648ea55ba.json @@ -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" } diff --git a/backend/.sqlx/query-3dfa0bf34c38b5529f2bbad405d23bda84db69757e132859d7e04df38e503f3e.json b/backend/.sqlx/query-3dfa0bf34c38b5529f2bbad405d23bda84db69757e132859d7e04df38e503f3e.json deleted file mode 100644 index cd60dffec3..0000000000 --- a/backend/.sqlx/query-3dfa0bf34c38b5529f2bbad405d23bda84db69757e132859d7e04df38e503f3e.json +++ /dev/null @@ -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" -} diff --git a/backend/.sqlx/query-528cdbb75f1c5135170a58fce3fda464be138272487639d0ffbbbe6961ec5c37.json b/backend/.sqlx/query-528cdbb75f1c5135170a58fce3fda464be138272487639d0ffbbbe6961ec5c37.json new file mode 100644 index 0000000000..e998bea512 --- /dev/null +++ b/backend/.sqlx/query-528cdbb75f1c5135170a58fce3fda464be138272487639d0ffbbbe6961ec5c37.json @@ -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" +} diff --git a/backend/.sqlx/query-5cb644f89a94b6e6a1d7a84155bafc460c1749946056ded8ca712dacecbf427d.json b/backend/.sqlx/query-5cb644f89a94b6e6a1d7a84155bafc460c1749946056ded8ca712dacecbf427d.json deleted file mode 100644 index dddedac81f..0000000000 --- a/backend/.sqlx/query-5cb644f89a94b6e6a1d7a84155bafc460c1749946056ded8ca712dacecbf427d.json +++ /dev/null @@ -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" -} diff --git a/backend/.sqlx/query-9a5c7e8b60a260085b438bd300972ebf948ea26f313ee5a73b85574becdd7dc7.json b/backend/.sqlx/query-9a5c7e8b60a260085b438bd300972ebf948ea26f313ee5a73b85574becdd7dc7.json new file mode 100644 index 0000000000..ebb54c644c --- /dev/null +++ b/backend/.sqlx/query-9a5c7e8b60a260085b438bd300972ebf948ea26f313ee5a73b85574becdd7dc7.json @@ -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" +} diff --git a/backend/.sqlx/query-a73b57229602d68cc25a8d963753271619e74df9fd7fc6cb05fc0614b10e6001.json b/backend/.sqlx/query-a73b57229602d68cc25a8d963753271619e74df9fd7fc6cb05fc0614b10e6001.json new file mode 100644 index 0000000000..4083089a04 --- /dev/null +++ b/backend/.sqlx/query-a73b57229602d68cc25a8d963753271619e74df9fd7fc6cb05fc0614b10e6001.json @@ -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" +} diff --git a/backend/.sqlx/query-bd213fca18a04d9e34405fd753c02398b72afb17b1ce650dbac947490634009d.json b/backend/.sqlx/query-bd213fca18a04d9e34405fd753c02398b72afb17b1ce650dbac947490634009d.json new file mode 100644 index 0000000000..643920205b --- /dev/null +++ b/backend/.sqlx/query-bd213fca18a04d9e34405fd753c02398b72afb17b1ce650dbac947490634009d.json @@ -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" +} diff --git a/backend/.sqlx/query-d0ce6a8ff7b89a0e99902989e57e7cae736e346c067261de34100f4fb97ca6db.json b/backend/.sqlx/query-d0ce6a8ff7b89a0e99902989e57e7cae736e346c067261de34100f4fb97ca6db.json new file mode 100644 index 0000000000..200d86c4c6 --- /dev/null +++ b/backend/.sqlx/query-d0ce6a8ff7b89a0e99902989e57e7cae736e346c067261de34100f4fb97ca6db.json @@ -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" +} diff --git a/backend/.sqlx/query-ab3c979c20a8ba9f6d906193c3bc7bc1941f0ecbe7ee2d8a62d2e02faeaae3f0.json b/backend/.sqlx/query-eb610cc508048cdcc030369784c274c828d67ee46ae7260b3c2fb04901b782f2.json similarity index 70% rename from backend/.sqlx/query-ab3c979c20a8ba9f6d906193c3bc7bc1941f0ecbe7ee2d8a62d2e02faeaae3f0.json rename to backend/.sqlx/query-eb610cc508048cdcc030369784c274c828d67ee46ae7260b3c2fb04901b782f2.json index fc37929689..8f61a30dd7 100644 --- a/backend/.sqlx/query-ab3c979c20a8ba9f6d906193c3bc7bc1941f0ecbe7ee2d8a62d2e02faeaae3f0.json +++ b/backend/.sqlx/query-eb610cc508048cdcc030369784c274c828d67ee46ae7260b3c2fb04901b782f2.json @@ -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" } diff --git a/backend/.sqlx/query-efc22aae27c30894f3d8ad9484f48439711707dba1f3112b495b1b18cbd31b31.json b/backend/.sqlx/query-efc22aae27c30894f3d8ad9484f48439711707dba1f3112b495b1b18cbd31b31.json deleted file mode 100644 index 024d164a7e..0000000000 --- a/backend/.sqlx/query-efc22aae27c30894f3d8ad9484f48439711707dba1f3112b495b1b18cbd31b31.json +++ /dev/null @@ -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" -} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index d38b462f99..c7d3780564 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -9664,6 +9664,7 @@ dependencies = [ "argon2", "async-oauth2", "async-recursion", + "async-stream", "async-stripe", "async_zip", "axum", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 10405dfd2a..da34b6ec2c 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -238,4 +238,5 @@ aws-sdk-sts = "^1" crc = "^3" tar = "^0" -http = "^1" \ No newline at end of file +http = "^1" +async-stream = "^0" \ No newline at end of file diff --git a/backend/migrations/20240322130200_add_offset_to_logs.down.sql b/backend/migrations/20240322130200_add_offset_to_logs.down.sql new file mode 100644 index 0000000000..2eb2d1b570 --- /dev/null +++ b/backend/migrations/20240322130200_add_offset_to_logs.down.sql @@ -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; diff --git a/backend/migrations/20240322130200_add_offset_to_logs.up.sql b/backend/migrations/20240322130200_add_offset_to_logs.up.sql new file mode 100644 index 0000000000..a4a658a90d --- /dev/null +++ b/backend/migrations/20240322130200_add_offset_to_logs.up.sql @@ -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[]; + diff --git a/backend/parsers/windmill-parser-go/src/lib.rs b/backend/parsers/windmill-parser-go/src/lib.rs index 01012291f1..cfb20b203c 100644 --- a/backend/parsers/windmill-parser-go/src/lib.rs +++ b/backend/parsers/windmill-parser-go/src/lib.rs @@ -142,8 +142,6 @@ pub fn otyp_to_string(otyp: Option) -> String { #[cfg(test)] mod tests { - use windmill_parser::{Arg, MainArgSignature, ObjectProperty, Typ}; - use super::*; #[test] diff --git a/backend/parsers/windmill-parser-wasm/src/lib.rs b/backend/parsers/windmill-parser-wasm/src/lib.rs index e9c491f850..5619e0142d 100644 --- a/backend/parsers/windmill-parser-wasm/src/lib.rs +++ b/backend/parsers/windmill-parser-wasm/src/lib.rs @@ -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}; diff --git a/backend/src/main.rs b/backend/src/main.rs index 5f4ca64478..0380c3c2f6 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -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 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::(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( } } -async fn load_value_from_global_settings(db: &DB, setting_name: &str) -> error::Result> { +async fn load_value_from_global_settings( + db: &DB, + setting_name: &str, +) -> error::Result> { 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( @@ -521,10 +539,7 @@ pub async fn reload_option_setting( if let Some(q) = q { if let Ok(v) = serde_json::from_value::(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( if let Some(q) = q { if let Ok(v) = serde_json::from_value::(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::>>(q.clone()) - { + if let Ok(v) = serde_json::from_value::>>(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 { diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index f500e155e0..4296e2f023 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -171,10 +171,8 @@ fn find_module_in_vec(modules: Vec, id: &str) -> Option) -> impl IntoResponse { &db, ) .await?; - Ok::<_, Error>(( - status_code, - headers, - response - )) + Ok::<_, Error>((status_code, headers, response)) } async fn list_paths( diff --git a/backend/windmill-api/src/job_metrics.rs b/backend/windmill-api/src/job_metrics.rs index 6f95f11fb8..4a57ed28d4 100644 --- a/backend/windmill-api/src/job_metrics.rs +++ b/backend/windmill-api/src/job_metrics.rs @@ -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; diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index a991c88bc9..eec01b6f1c 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -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>, +) -> Option> { + 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; + } + } + } + + 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)>, -) -> error::Result { - let text = sqlx::query_scalar!( - "SELECT CONCAT(coalesce(completed_job.logs, ''), coalesce(job_logs.logs, '')) +) -> error::Result { + 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, pub completed: Option, pub new_logs: Option, + pub log_offset: Option, pub mem_peak: Option, pub flow_status: Option, } @@ -3311,8 +3355,9 @@ async fn get_job_update( Query(JobUpdateQuery { running, log_offset }): Query, ) -> error::JsonResult { 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, diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 2595238dcb..8604464fb8 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -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={}", diff --git a/backend/windmill-api/src/schedule.rs b/backend/windmill-api/src/schedule.rs index 667d94951e..17e0ad5f84 100644 --- a/backend/windmill-api/src/schedule.rs +++ b/backend/windmill-api/src/schedule.rs @@ -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() diff --git a/backend/windmill-api/src/scripts.rs b/backend/windmill-api/src/scripts.rs index 488d216962..8d5367df2f 100644 --- a/backend/windmill-api/src/scripts.rs +++ b/backend/windmill-api/src/scripts.rs @@ -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 { diff --git a/backend/windmill-api/src/settings.rs b/backend/windmill-api/src/settings.rs index 3d70ee8515..43c144066b 100644 --- a/backend/windmill-api/src/settings.rs +++ b/backend/windmill-api/src/settings.rs @@ -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, @@ -115,16 +111,37 @@ pub async fn test_s3_bucket( Json(test_s3_bucket): Json, ) -> error::Result { 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, } 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, ) -> 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( diff --git a/backend/windmill-api/src/users.rs b/backend/windmill-api/src/users.rs index 457c837f53..37c9909cda 100644 --- a/backend/windmill-api/src/users.rs +++ b/backend/windmill-api/src/users.rs @@ -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, diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index c99e16d814..6d30fb5cc1 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -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, } -#[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, } -#[derive(Deserialize)] -#[derive(sqlx::FromRow)] +#[derive(Deserialize, sqlx::FromRow)] pub struct NewFlow { pub path: String, pub summary: String, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 8f4be9b8fe..7a14b87cb8 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -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, logs: impl AsRef, db: impl Borrow>, ) { @@ -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 diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 07dbd5defd..2eb40d34c4 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -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( 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, + 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(¤t_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, + 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, + 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}"); diff --git a/backend/windmill-worker/src/global_cache.rs b/backend/windmill-worker/src/global_cache.rs index 062b7e876a..56aa4c8eae 100644 --- a/backend/windmill-worker/src/global_cache.rs +++ b/backend/windmill-worker/src/global_cache.rs @@ -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, folder: String) -> error::Result<()> { +pub async fn build_tar_and_push( + s3_client: Arc, + 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, 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, folder: String) Ok(()) } - #[cfg(all(feature = "enterprise", feature = "parquet"))] pub async fn pull_from_tar(client: Arc, folder: String) -> error::Result<()> { - use object_store::path::Path; use tokio::fs::metadata; let folder_name = folder.split("/").last().unwrap(); diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 2800880de1..58f24ea64d 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -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::>(); - 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::>(); + 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.")); diff --git a/backend/windmill-worker/src/snowflake_executor.rs b/backend/windmill-worker/src/snowflake_executor.rs index c3d8cb4fb4..8f5afadd2d 100644 --- a/backend/windmill-worker/src/snowflake_executor.rs +++ b/backend/windmill-worker/src/snowflake_executor.rs @@ -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; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 1f3ecba0b7..547aa7c411 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -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()); diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 3bb864c2c1..e8e38ba314 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -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, Self>; -} +// trait IntoArray: Sized { +// fn into_array(self) -> Result, Self>; +// } -impl IntoArray for Value { - fn into_array(self) -> Result, Self> { - match self { - Value::Array(array) => Ok(array), - not_array => Err(not_array), - } - } -} +// impl IntoArray for Value { +// fn into_array(self) -> Result, Self> { +// match self { +// Value::Array(array) => Ok(array), +// not_array => Err(not_array), +// } +// } +// } fn from_now(duration: Duration) -> chrono::DateTime { // "This function errors when original duration is larger than diff --git a/frontend/src/lib/components/TestJobLoader.svelte b/frontend/src/lib/components/TestJobLoader.svelte index 7c637bfe81..3456b50977 100644 --- a/frontend/src/lib/components/TestJobLoader.svelte +++ b/frontend/src/lib/components/TestJobLoader.svelte @@ -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) }