perf: cache workspace premium check

This commit is contained in:
Ruben Fiszel
2025-03-25 02:00:51 +01:00
parent 1d47c4ad08
commit 7e802bd611
27 changed files with 138 additions and 386 deletions
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM workspace_runnable_usage WHERE resource_path = $1 AND workspace_id = $2 and resource_kind = 'flow'",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": []
},
"hash": "1ea8c2ce29da8d1be19dfd3ecc7f84fef7557f3e3046e7188ee06f6c3a2b2add"
}
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE flow_workspace_runnables SET workspace_id = $1 WHERE workspace_id = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Text"
]
},
"nullable": []
},
"hash": "2bb2cf6accb18d3e37a63388cca52a6591e7593b1a7c3d7a6848587679a48187"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO flow_workspace_runnables (flow_path, runnable_path, runnable_is_flow, workspace_id) VALUES ($1, $2, TRUE, $3) ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar"
]
},
"nullable": []
},
"hash": "35795d27c4ca69d2f145b4dba08a6ed16c25aea4584103c1b9a3651eb31bfe53"
}
@@ -1,22 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT premium FROM workspace WHERE workspace.id = $1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "premium",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Text"
]
},
"nullable": [
false
]
},
"hash": "4c970f10d345bcdcf956dcbfa22b6e80888e511fb4787cb1a7976878abed1d30"
}
@@ -1,24 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT f.path\n FROM workspace_runnable_usage wru \n JOIN flow f\n ON wru.resource_path = f.path AND wru.workspace_id = f.workspace_id\n WHERE wru.runnable_path = $1 AND wru.runnable_is_flow = $2 AND wru.workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
false
]
},
"hash": "56a2fe44b73728c9aba10f3d1c7707bd1c48c58aaa5b2ae7560cbed3c8d46329"
}
@@ -1,17 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO workspace_runnable_usage (resource_path, runnable_path, script_hash, runnable_is_flow, workspace_id, resource_kind) VALUES ($1, $2, $3, FALSE, $4, 'flow') ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Int8",
"Varchar"
]
},
"nullable": []
},
"hash": "894e5080e83b098839a050bb472359d6d9fdfdbffb40a2464a45056660a3a2d4"
}
@@ -1,24 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "SELECT f.path\n FROM flow_workspace_runnables fwr \n JOIN flow f \n ON fwr.flow_path = f.path AND fwr.workspace_id = f.workspace_id\n WHERE fwr.runnable_path = $1 AND fwr.runnable_is_flow = $2 AND fwr.workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
false
]
},
"hash": "9aeee333b1dbe58ba819ba3b2713242b54d77b46b5f785ae66b8104f89f43219"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE workspace_runnable_usage SET resource_path = REGEXP_REPLACE(resource_path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE resource_path LIKE ('u/' || $2 || '/%') AND workspace_id = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text",
"Text"
]
},
"nullable": []
},
"hash": "9c956541e068154193243795d96cc62a0fc9d4ba56014f3e554cc2b50ce1da06"
}
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM workspace_runnable_usage WHERE resource_path = $1 AND workspace_id = $2 and resource_kind = 'app'",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": []
},
"hash": "a28f8697da75d7ac2c6e13dda394ee12e8cffe925bd7cc4e9d02dc90ffe200f3"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO workspace_runnable_usage (resource_path, runnable_path, runnable_is_flow, workspace_id, resource_kind) VALUES ($1, $2, TRUE, $3, 'flow') ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar"
]
},
"nullable": []
},
"hash": "ab9783a48e0f6cacd5dd86cfa21ba8354d4cd733335db3d76155b663f668b9d7"
}
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_workspace_runnables WHERE flow_path = $1 AND workspace_id = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": []
},
"hash": "acbea8740b28c26942c50edcf5618cd141e68cf83a7dcae7d3c1b8a7ba94425b"
}
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE workspace_runnable_usage SET workspace_id = $1 WHERE workspace_id = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Text"
]
},
"nullable": []
},
"hash": "af2f780da994befd09959d3b1b077eb9b71337a9aa2a64cd2db031e64e929251"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE flow_workspace_runnables SET flow_path = REGEXP_REPLACE(flow_path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE flow_path LIKE ('u/' || $2 || '/%') AND workspace_id = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text",
"Text"
]
},
"nullable": []
},
"hash": "b656927cd70b6667f3c72186ec04f0bf040da3af9e2eac3229264ec95b4755d8"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE workspace_runnable_usage SET runnable_path = REGEXP_REPLACE(runnable_path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE runnable_path LIKE ('u/' || $2 || '/%') AND workspace_id = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text",
"Text"
]
},
"nullable": []
},
"hash": "baf63525ca210c22d3ad8c0b197dc6b39dabab77ae7810f553f1773995db51f1"
}
@@ -1,17 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO flow_workspace_runnables (flow_path, runnable_path, script_hash, runnable_is_flow, workspace_id) VALUES ($1, $2, $3, FALSE, $4) ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Int8",
"Varchar"
]
},
"nullable": []
},
"hash": "c35f44f91b08fa57e29a2b4a685706f62e700695810f23108e975dfcd1fee7a3"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE flow_workspace_runnables SET runnable_path = REGEXP_REPLACE(runnable_path,'u/' || $2 || '/(.*)','u/' || $1 || '/\\1') WHERE runnable_path LIKE ('u/' || $2 || '/%') AND workspace_id = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text",
"Text"
]
},
"nullable": []
},
"hash": "dc58e5b4715601a93b3c01a2564a4420f232867a23cacb9a62b386f129a86a4b"
}
+1 -1
View File
@@ -58,7 +58,7 @@ parquet = ["windmill-api/parquet", "windmill-common/parquet", "windmill-worker/p
prometheus = ["windmill-common/prometheus", "windmill-api/prometheus", "windmill-worker/prometheus", "windmill-queue/prometheus", "dep:prometheus"]
flow_testing = ["windmill-worker/flow_testing"]
openidconnect = ["windmill-api/openidconnect"]
cloud = ["windmill-queue/cloud", "windmill-worker/cloud"]
cloud = ["windmill-queue/cloud", "windmill-worker/cloud", "windmill-common/cloud", "windmill-api/cloud"]
jemalloc = ["windmill-common/jemalloc", "dep:tikv-jemallocator", "dep:tikv-jemalloc-sys", "dep:tikv-jemalloc-ctl"]
tantivy = ["dep:windmill-indexer", "windmill-api/tantivy", "windmill-indexer/enterprise", "windmill-indexer/parquet", "windmill-common/tantivy", "enterprise", "parquet"]
sqlx = ["windmill-worker/sqlx"]
@@ -0,0 +1,3 @@
-- Add down migration script here
DROP TRIGGER workspace_premium_change_trigger ON workspace;
DROP FUNCTION notify_workspace_premium_change();
@@ -0,0 +1,13 @@
-- Add up migration script here
CREATE OR REPLACE FUNCTION notify_workspace_premium_change()
RETURNS TRIGGER AS $$
BEGIN
PERFORM pg_notify('notify_workspace_premium_change', NEW.id);
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
CREATE TRIGGER workspace_premium_change_trigger
AFTER UPDATE OF premium ON workspace
FOR EACH ROW
EXECUTE FUNCTION notify_workspace_premium_change();
+16 -9
View File
@@ -707,6 +707,11 @@ Windmill Community Edition {GIT_VERSION}
tracing::info!("Workspace envs change detected, invalidating workspace envs cache: {}", workspace_id);
windmill_common::variables::CUSTOM_ENVS_CACHE.remove(workspace_id);
},
"notify_workspace_premium_change" => {
let workspace_id = n.payload();
tracing::info!("Workspace premium change detected, invalidating workspace premium cache: {}", workspace_id);
windmill_common::workspaces::IS_PREMIUM_CACHE.remove(workspace_id);
},
"notify_global_setting_change" => {
tracing::info!("Global setting change detected: {}", n.payload());
match n.payload() {
@@ -980,15 +985,17 @@ async fn listen_pg(url: &str) -> Option<PgListener> {
}
};
if let Err(e) = listener
.listen_all(vec![
"notify_config_change",
"notify_global_setting_change",
"notify_webhook_change",
"notify_workspace_envs_change",
])
.await
{
#[allow(unused_mut)]
let mut channels = vec![
"notify_config_change",
"notify_global_setting_change",
"notify_webhook_change",
"notify_workspace_envs_change",
];
#[cfg(feature = "cloud")]
channels.push("notify_workspace_premium_change");
if let Err(e) = listener.listen_all(channels).await {
tracing::error!(error = %e, "Could not listen to database");
return None;
}
+1
View File
@@ -31,6 +31,7 @@ static_frontend = ["dep:rust-embed"]
postgres_trigger = ["dep:rust-postgres", "dep:pg_escape", "dep:byteorder", "dep:thiserror", "dep:rust_decimal", "dep:rust-postgres-native-tls"]
mqtt_trigger = ["dep:thiserror", "dep:rumqttc"]
sqs_trigger = ["dep:aws-sdk-sqs", "dep:thiserror", "dep:aws-config"]
cloud = ["windmill-common/cloud"]
[dependencies]
windmill-queue.workspace = true
+8 -12
View File
@@ -393,19 +393,15 @@ async fn list_pending_invites(
async fn is_premium(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Path(w_id): Path<String>,
Extension(_db): Extension<DB>,
Path(_w_id): Path<String>,
) -> JsonResult<bool> {
require_admin(authed.is_admin, &authed.username)?;
let mut tx = db.begin().await?;
let row = sqlx::query_scalar!(
"SELECT premium FROM workspace WHERE workspace.id = $1",
&w_id
)
.fetch_one(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(row))
#[cfg(feature = "cloud")]
let premium = windmill_common::workspaces::is_premium_workspace(&_db, &_w_id).await;
#[cfg(not(feature = "cloud"))]
let premium = false;
Ok(Json(premium))
}
async fn exists_workspace(
@@ -1362,7 +1358,7 @@ struct UsedTriggers {
pub nats_used: bool,
pub postgres_used: bool,
pub mqtt_used: bool,
pub sqs_used: bool
pub sqs_used: bool,
}
async fn get_used_triggers(
+1
View File
@@ -17,6 +17,7 @@ otel = ["dep:opentelemetry-semantic-conventions", "dep:opentelemetry-otlp", "dep
"dep:opentelemetry", "dep:tracing-opentelemetry", "dep:opentelemetry-appender-tracing", "dep:tonic"]
smtp = ["dep:mail-send"]
scoped_cache = []
cloud = []
[lib]
name = "windmill_common"
+20 -1
View File
@@ -1,3 +1,4 @@
use quick_cache::sync::Cache;
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize, Debug, Default)]
@@ -27,7 +28,7 @@ pub enum ObjectType {
ResourceType,
User,
Group,
Trigger
Trigger,
}
#[derive(Serialize, Deserialize, Debug)]
@@ -38,3 +39,21 @@ pub struct GitRepositorySettings {
pub group_by_folder: Option<bool>,
pub exclude_types_override: Option<Vec<ObjectType>>,
}
lazy_static::lazy_static! {
pub static ref IS_PREMIUM_CACHE: Cache<String, bool> = Cache::new(5000);
}
#[cfg(feature = "cloud")]
pub async fn is_premium_workspace(_db: &crate::DB, _w_id: &str) -> bool {
let cached = IS_PREMIUM_CACHE.get(_w_id);
if let Some(cached) = cached {
return cached;
}
let premium = sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", _w_id)
.fetch_one(_db)
.await
.unwrap_or(false);
IS_PREMIUM_CACHE.insert(_w_id.to_string(), premium);
premium
}
+2 -16
View File
@@ -928,13 +928,7 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
if *CLOUD_HOSTED && !queued_job.is_flow() && _duration > 1000 {
let additional_usage = _duration / 1000;
let w_id = &queued_job.workspace_id;
let premium_workspace =
sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", w_id)
.fetch_one(db)
.await
.map_err(|e| {
Error::internal_err(format!("fetching if {w_id} is premium: {e:#}"))
})?;
let premium_workspace = windmill_common::workspaces::is_premium_workspace(db, w_id).await;
let _ = sqlx::query!(
"INSERT INTO usage (id, is_workspace, month_, usage)
VALUES ($1, TRUE, EXTRACT(YEAR FROM current_date) * 12 + EXTRACT(MONTH FROM current_date), $2)
@@ -2996,15 +2990,7 @@ pub async fn push<'c, 'd>(
#[cfg(feature = "cloud")]
if *CLOUD_HOSTED {
let premium_workspace =
sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", workspace_id)
.fetch_one(_db)
.await
.map_err(|e| {
Error::internal_err(format!(
"fetching if {workspace_id} is premium and overquota: {e:#}"
))
})?;
windmill_common::workspaces::is_premium_workspace(_db, workspace_id).await;
// we track only non flow steps
let (workspace_usage, user_usage) = if !matches!(
job_payload,
+2 -8
View File
@@ -564,14 +564,8 @@ pub async fn resolve_job_timeout(
) -> (Duration, Option<String>, bool) {
let mut warn_msg: Option<String> = None;
#[cfg(feature = "cloud")]
let cloud_premium_workspace = *CLOUD_HOSTED
&& sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", _w_id)
.fetch_one(_db)
.await
.map_err(|e| {
tracing::error!(%e, "error getting premium workspace for job {_job_id}: {e:#}");
})
.unwrap_or(false);
let cloud_premium_workspace =
*CLOUD_HOSTED && windmill_common::workspaces::is_premium_workspace(_db, _w_id).await;
#[cfg(not(feature = "cloud"))]
let cloud_premium_workspace = false;
+71 -64
View File
@@ -1283,70 +1283,9 @@ pub async fn run_worker(
"received {} from same worker channel",
same_worker_job.job_id
);
let r = sqlx::query_as::<_, PulledJob>(
"WITH ping AS (
UPDATE v2_job_runtime SET ping = NOW() WHERE id = $1
),
started_at AS (
UPDATE v2_job_queue SET started_at = NOW() WHERE id = $1
)
SELECT
v2_job_queue.workspace_id,
v2_job_queue.id,
v2_job.args,
v2_job.parent_job,
v2_job.created_by,
v2_job_queue.started_at,
scheduled_for,
v2_job.runnable_path,
v2_job.kind,
v2_job.runnable_id,
v2_job_queue.canceled_reason,
v2_job_queue.canceled_by,
v2_job.permissioned_as,
v2_job.permissioned_as_email,
v2_job_status.flow_status,
v2_job.tag,
v2_job.script_lang,
v2_job.same_worker,
v2_job.pre_run_error,
v2_job.concurrent_limit,
v2_job.concurrency_time_window_s,
v2_job.flow_innermost_root_job,
v2_job.timeout,
v2_job.flow_step_id,
v2_job.cache_ttl,
v2_job_queue.priority,
v2_job.preprocessed,
v2_job.script_entrypoint_override,
v2_job.trigger,
v2_job.trigger_kind,
v2_job.visible_to_owner,
v2_job.raw_code,
v2_job.raw_lock,
v2_job.raw_flow,
pj.runnable_path as parent_runnable_path,
p.email as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin,
p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders
FROM v2_job_queue
INNER JOIN v2_job ON v2_job.id = v2_job_queue.id
LEFT JOIN v2_job_status ON v2_job_status.id = v2_job_queue.id
LEFT JOIN job_perms p ON p.job_id = v2_job.id
LEFT JOIN v2_job pj ON v2_job.parent_job = pj.id
WHERE v2_job_queue.id = $1
",
)
.bind(same_worker_job.job_id)
.fetch_optional(db)
.await
.map_err(|e| {
Error::internal_err(format!(
"Impossible to fetch same_worker job {}: {}",
same_worker_job.job_id, e
))
});
let job = get_same_worker_job(db, &same_worker_job).await;
// tracing::error!("r: {:?}", r);
if r.is_err() && !same_worker_job.recoverable {
if job.is_err() && !same_worker_job.recoverable {
tracing::error!(
worker = %worker_name, hostname = %hostname,
"failed to fetch same_worker job on a non recoverable job, exiting"
@@ -1358,7 +1297,7 @@ pub async fn run_worker(
.expect("send kill to job completed tx");
break;
} else {
r
job
}
} else if let Ok(_) = killpill_rx.try_recv() {
if !killed_but_draining_same_worker_jobs {
@@ -1862,6 +1801,74 @@ pub async fn run_worker(
tracing::info!(worker = %worker_name, hostname = %hostname, "number of jobs executed: {}", jobs_executed);
}
async fn get_same_worker_job(
db: &Pool<Postgres>,
same_worker_job: &SameWorkerPayload,
) -> windmill_common::error::Result<Option<PulledJob>> {
sqlx::query_as::<_, PulledJob>(
"WITH ping AS (
UPDATE v2_job_runtime SET ping = NOW() WHERE id = $1
),
started_at AS (
UPDATE v2_job_queue SET started_at = NOW() WHERE id = $1
)
SELECT
v2_job_queue.workspace_id,
v2_job_queue.id,
v2_job.args,
v2_job.parent_job,
v2_job.created_by,
v2_job_queue.started_at,
scheduled_for,
v2_job.runnable_path,
v2_job.kind,
v2_job.runnable_id,
v2_job_queue.canceled_reason,
v2_job_queue.canceled_by,
v2_job.permissioned_as,
v2_job.permissioned_as_email,
v2_job_status.flow_status,
v2_job.tag,
v2_job.script_lang,
v2_job.same_worker,
v2_job.pre_run_error,
v2_job.concurrent_limit,
v2_job.concurrency_time_window_s,
v2_job.flow_innermost_root_job,
v2_job.timeout,
v2_job.flow_step_id,
v2_job.cache_ttl,
v2_job_queue.priority,
v2_job.preprocessed,
v2_job.script_entrypoint_override,
v2_job.trigger,
v2_job.trigger_kind,
v2_job.visible_to_owner,
v2_job.raw_code,
v2_job.raw_lock,
v2_job.raw_flow,
pj.runnable_path as parent_runnable_path,
p.email as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin,
p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders
FROM v2_job_queue
INNER JOIN v2_job ON v2_job.id = v2_job_queue.id
LEFT JOIN v2_job_status ON v2_job_status.id = v2_job_queue.id
LEFT JOIN job_perms p ON p.job_id = v2_job.id
LEFT JOIN v2_job pj ON v2_job.parent_job = pj.id
WHERE v2_job_queue.id = $1
",
)
.bind(same_worker_job.job_id)
.fetch_optional(db)
.await
.map_err(|e| {
Error::internal_err(format!(
"Impossible to fetch same_worker job {}: {}",
same_worker_job.job_id, e
))
})
}
async fn queue_init_bash_maybe<'c>(
db: &Pool<Postgres>,
same_worker_tx: SameWorkerSender,