diff --git a/backend/.sqlx/query-54fef88cc6b9e8db7c07fccbaa845edfefa339153a32e468faad9f063008863d.json b/backend/.sqlx/query-177661a6487cfef198c2ceca4c81e546230c2fd7756d13197a7c19a78eb4b3d6.json similarity index 54% rename from backend/.sqlx/query-54fef88cc6b9e8db7c07fccbaa845edfefa339153a32e468faad9f063008863d.json rename to backend/.sqlx/query-177661a6487cfef198c2ceca4c81e546230c2fd7756d13197a7c19a78eb4b3d6.json index 3c529dbd84..4a1be7e7b4 100644 --- a/backend/.sqlx/query-54fef88cc6b9e8db7c07fccbaa845edfefa339153a32e468faad9f063008863d.json +++ b/backend/.sqlx/query-177661a6487cfef198c2ceca4c81e546230c2fd7756d13197a7c19a78eb4b3d6.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2, occupancy_rate = $3, current_job_id = NULL, current_job_workspace_id = NULL WHERE worker = $4", + "query": "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2, occupancy_rate = $3, memory_usage = $4, current_job_id = NULL, current_job_workspace_id = NULL WHERE worker = $5", "describe": { "columns": [], "parameters": { @@ -8,10 +8,11 @@ "Int4", "TextArray", "Float4", + "Int8", "Text" ] }, "nullable": [] }, - "hash": "54fef88cc6b9e8db7c07fccbaa845edfefa339153a32e468faad9f063008863d" + "hash": "177661a6487cfef198c2ceca4c81e546230c2fd7756d13197a7c19a78eb4b3d6" } diff --git a/backend/.sqlx/query-e00171cc3fc8f32922562d4fc4d3db56955fbc6ee9872d4b2e99157ebed131e0.json b/backend/.sqlx/query-99d79e2a3792296b0582bbaf8aa1d39855d2109578d0a62d7f2cf1444e44550e.json similarity index 73% rename from backend/.sqlx/query-e00171cc3fc8f32922562d4fc4d3db56955fbc6ee9872d4b2e99157ebed131e0.json rename to backend/.sqlx/query-99d79e2a3792296b0582bbaf8aa1d39855d2109578d0a62d7f2cf1444e44550e.json index 67aec442e1..553cf6494e 100644 --- a/backend/.sqlx/query-e00171cc3fc8f32922562d4fc4d3db56955fbc6ee9872d4b2e99157ebed131e0.json +++ b/backend/.sqlx/query-99d79e2a3792296b0582bbaf8aa1d39855d2109578d0a62d7f2cf1444e44550e.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, CASE WHEN $4 IS TRUE THEN current_job_id ELSE NULL END as current_job_id, CASE WHEN $4 IS TRUE THEN current_job_workspace_id ELSE NULL END as current_job_workspace_id, custom_tags, worker_group, wm_version, occupancy_rate\n FROM worker_ping\n WHERE ($1::integer IS NULL AND ping_at > now() - interval '5 minute') OR (ping_at > now() - ($1 || ' seconds')::interval)\n ORDER BY ping_at desc LIMIT $2 OFFSET $3", + "query": "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, CASE WHEN $4 IS TRUE THEN current_job_id ELSE NULL END as current_job_id, CASE WHEN $4 IS TRUE THEN current_job_workspace_id ELSE NULL END as current_job_workspace_id, custom_tags, worker_group, wm_version, occupancy_rate, memory, vcpus, memory_usage\n FROM worker_ping\n WHERE ($1::integer IS NULL AND ping_at > now() - interval '5 minute') OR (ping_at > now() - ($1 || ' seconds')::interval)\n ORDER BY ping_at desc LIMIT $2 OFFSET $3", "describe": { "columns": [ { @@ -62,6 +62,21 @@ "ordinal": 11, "name": "occupancy_rate", "type_info": "Float4" + }, + { + "ordinal": 12, + "name": "memory", + "type_info": "Int8" + }, + { + "ordinal": 13, + "name": "vcpus", + "type_info": "Int8" + }, + { + "ordinal": 14, + "name": "memory_usage", + "type_info": "Int8" } ], "parameters": { @@ -84,8 +99,11 @@ true, false, false, + true, + true, + true, true ] }, - "hash": "e00171cc3fc8f32922562d4fc4d3db56955fbc6ee9872d4b2e99157ebed131e0" + "hash": "99d79e2a3792296b0582bbaf8aa1d39855d2109578d0a62d7f2cf1444e44550e" } diff --git a/backend/migrations/20240527151853_add_worker_memory_usage.down.sql b/backend/migrations/20240527151853_add_worker_memory_usage.down.sql new file mode 100644 index 0000000000..d2f607c5b8 --- /dev/null +++ b/backend/migrations/20240527151853_add_worker_memory_usage.down.sql @@ -0,0 +1 @@ +-- Add down migration script here diff --git a/backend/migrations/20240527151853_add_worker_memory_usage.up.sql b/backend/migrations/20240527151853_add_worker_memory_usage.up.sql new file mode 100644 index 0000000000..3678d9ecd7 --- /dev/null +++ b/backend/migrations/20240527151853_add_worker_memory_usage.up.sql @@ -0,0 +1,2 @@ +-- Add up migration script here +ALTER TABLE worker_ping ADD COLUMN memory_usage BIGINT; \ No newline at end of file diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index f68eeddb8c..01a934f4f4 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -9907,6 +9907,12 @@ components: type: string occupancy_rate: type: number + memory: + type: number + vcpus: + type: number + memory_usage: + type: number required: - worker - worker_instance diff --git a/backend/windmill-api/src/workers.rs b/backend/windmill-api/src/workers.rs index 8bb56e78fd..71c5be4dfd 100644 --- a/backend/windmill-api/src/workers.rs +++ b/backend/windmill-api/src/workers.rs @@ -51,7 +51,14 @@ struct WorkerPing { custom_tags: Option>, worker_group: String, wm_version: String, + #[serde(skip_serializing_if = "Option::is_none")] occupancy_rate: Option, + #[serde(skip_serializing_if = "Option::is_none")] + memory: Option, + #[serde(skip_serializing_if = "Option::is_none")] + vcpus: Option, + #[serde(skip_serializing_if = "Option::is_none")] + memory_usage: Option, } #[derive(Serialize, Deserialize)] @@ -79,7 +86,7 @@ async fn list_worker_pings( let rows = sqlx::query_as!( WorkerPing, - "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, CASE WHEN $4 IS TRUE THEN current_job_id ELSE NULL END as current_job_id, CASE WHEN $4 IS TRUE THEN current_job_workspace_id ELSE NULL END as current_job_workspace_id, custom_tags, worker_group, wm_version, occupancy_rate + "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, CASE WHEN $4 IS TRUE THEN current_job_id ELSE NULL END as current_job_id, CASE WHEN $4 IS TRUE THEN current_job_workspace_id ELSE NULL END as current_job_workspace_id, custom_tags, worker_group, wm_version, occupancy_rate, memory, vcpus, memory_usage FROM worker_ping WHERE ($1::integer IS NULL AND ping_at > now() - interval '5 minute') OR (ping_at > now() - ($1 || ' seconds')::interval) ORDER BY ping_at desc LIMIT $2 OFFSET $3", diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index e89661e5fd..bee9fb83ec 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -132,6 +132,24 @@ fn process_custom_tags(tags: Vec) -> (Vec, HashMap Option { + let mut memory = std::process::Command::new("cat") + .args(["/sys/fs/cgroup/memory.current"]) + .output() + .ok() + .map(|o| String::from_utf8_lossy(&o.stdout).to_string()); + + if memory.is_none() { + memory = std::process::Command::new("cat") + .args(["/sys/fs/cgroup/memory/memory.usage_in_bytes"]) + .output() + .ok() + .map(|o| String::from_utf8_lossy(&o.stdout).to_string()) + } + + memory.map(|x| x.parse::().ok()).flatten() +} + pub async fn update_ping(worker_instance: &str, worker_name: &str, ip: &str, db: &DB) { let (tags, dw) = { let wc = WORKER_CONFIG.read().await.clone(); diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 68870ac893..0e0697add9 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -6,7 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ -use windmill_common::worker::TMP_DIR; +use windmill_common::worker::{get_worker_memory_usage, TMP_DIR}; use anyhow::Result; use const_format::concatcp; @@ -1426,11 +1426,14 @@ pub async fn run_worker NUM_SECS_PING { let tags = WORKER_CONFIG.read().await.worker_tags.clone(); + let memory_usage = get_worker_memory_usage(); + if let Err(e) = sqlx::query!( - "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2, occupancy_rate = $3, current_job_id = NULL, current_job_workspace_id = NULL WHERE worker = $4", + "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2, occupancy_rate = $3, memory_usage = $4, current_job_id = NULL, current_job_workspace_id = NULL WHERE worker = $5", jobs_executed, tags.as_slice(), worker_code_execution_metric / start_time.elapsed().as_secs_f32(), + memory_usage, &worker_name ).execute(db).await { tracing::error!("failed to update worker ping, exiting: {}", e); diff --git a/frontend/src/routes/(root)/(logged)/workers/+page.svelte b/frontend/src/routes/(root)/(logged)/workers/+page.svelte index b8bc77f96e..a4b11fa018 100644 --- a/frontend/src/routes/(root)/(logged)/workers/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/workers/+page.svelte @@ -345,6 +345,8 @@ Current job Occupancy rate {/if} + Memory usage + vCPUs/memory limits Version Liveness @@ -354,7 +356,7 @@ @@ -367,7 +369,7 @@ {#if workers} - {#each workers as { worker, custom_tags, last_ping, started_at, jobs_executed, current_job_id, current_job_workspace_id, occupancy_rate, wm_version }} + {#each workers as { worker, custom_tags, last_ping, started_at, jobs_executed, current_job_id, current_job_workspace_id, occupancy_rate, wm_version, vcpus, memory, memory_usage }} {worker} @@ -394,11 +396,16 @@ {Math.ceil(occupancy_rate ?? 0 * 100)}% {/if} -
{wm_version.split('-')[0]}{wm_version}
+ {memory_usage ? Math.round(memory_usage / 1000000) + 'MB' : '--'} + + {vcpus ? (vcpus / 100000).toFixed(1) + ' vCPUs' : '--'} + / {memory ? Math.round(memory / 1000000) + 'MB' : '--'} + + +
+ {wm_version.split('-')[0]}{wm_version} +
+