diff --git a/backend/.sqlx/query-8f0d98661f2a32659f37a123320d14288ce88c9d6e6279bdc36334410023f5e1.json b/backend/.sqlx/query-0694f4be1a7966126f20afbd9ee094088db8ce98aefdbe96144c42e38e5bd813.json similarity index 70% rename from backend/.sqlx/query-8f0d98661f2a32659f37a123320d14288ce88c9d6e6279bdc36334410023f5e1.json rename to backend/.sqlx/query-0694f4be1a7966126f20afbd9ee094088db8ce98aefdbe96144c42e38e5bd813.json index dba901932e..e4cfd00b92 100644 --- a/backend/.sqlx/query-8f0d98661f2a32659f37a123320d14288ce88c9d6e6279bdc36334410023f5e1.json +++ b/backend/.sqlx/query-0694f4be1a7966126f20afbd9ee094088db8ce98aefdbe96144c42e38e5bd813.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "select tag, concurrent_limit, concurrency_time_window_s from script where hash = $1 AND workspace_id = $2", + "query": "select tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where hash = $1 AND workspace_id = $2", "describe": { "columns": [ { @@ -17,6 +17,11 @@ "ordinal": 2, "name": "concurrency_time_window_s", "type_info": "Int4" + }, + { + "ordinal": 3, + "name": "cache_ttl", + "type_info": "Int4" } ], "parameters": { @@ -26,10 +31,11 @@ ] }, "nullable": [ + true, true, true, true ] }, - "hash": "8f0d98661f2a32659f37a123320d14288ce88c9d6e6279bdc36334410023f5e1" + "hash": "0694f4be1a7966126f20afbd9ee094088db8ce98aefdbe96144c42e38e5bd813" } diff --git a/backend/.sqlx/query-9289709b2727425be0e7f1e341d6acd46e1e79d5018f6a18cb41d50ca3f4ba11.json b/backend/.sqlx/query-35d2c5708645e60dd5ab56d1e6cfcd197254a534fed178a46de4a2d0f0f5d348.json similarity index 60% rename from backend/.sqlx/query-9289709b2727425be0e7f1e341d6acd46e1e79d5018f6a18cb41d50ca3f4ba11.json rename to backend/.sqlx/query-35d2c5708645e60dd5ab56d1e6cfcd197254a534fed178a46de4a2d0f0f5d348.json index 6f298c678b..76b8ee6cdd 100644 --- a/backend/.sqlx/query-9289709b2727425be0e7f1e341d6acd46e1e79d5018f6a18cb41d50ca3f4ba11.json +++ b/backend/.sqlx/query-35d2c5708645e60dd5ab56d1e6cfcd197254a534fed178a46de4a2d0f0f5d348.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "select hash, tag, concurrent_limit, concurrency_time_window_s from script where path = $1 AND workspace_id = $2 AND\n created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND\n deleted = false AND archived = false)", + "query": "select hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where path = $1 AND workspace_id = $2 AND\n created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND\n deleted = false AND archived = false AND lock IS not NULL AND lock_error_logs IS NULL)", "describe": { "columns": [ { @@ -22,6 +22,11 @@ "ordinal": 3, "name": "concurrency_time_window_s", "type_info": "Int4" + }, + { + "ordinal": 4, + "name": "cache_ttl", + "type_info": "Int4" } ], "parameters": { @@ -34,8 +39,9 @@ false, true, true, + true, true ] }, - "hash": "9289709b2727425be0e7f1e341d6acd46e1e79d5018f6a18cb41d50ca3f4ba11" + "hash": "35d2c5708645e60dd5ab56d1e6cfcd197254a534fed178a46de4a2d0f0f5d348" } diff --git a/backend/.sqlx/query-9ee1423945740c17dedc2fc7ea7e7183380c43d0feeb615c651aa486eef3bd53.json b/backend/.sqlx/query-52dbe77fad5729e9e121609eeb6a3573e6d0e1d1c16e29b15c7ad810536e249c.json similarity index 84% rename from backend/.sqlx/query-9ee1423945740c17dedc2fc7ea7e7183380c43d0feeb615c651aa486eef3bd53.json rename to backend/.sqlx/query-52dbe77fad5729e9e121609eeb6a3573e6d0e1d1c16e29b15c7ad810536e249c.json index cd3ae03d9f..ac9f9288ea 100644 --- a/backend/.sqlx/query-9ee1423945740c17dedc2fc7ea7e7183380c43d0feeb615c651aa486eef3bd53.json +++ b/backend/.sqlx/query-52dbe77fad5729e9e121609eeb6a3573e6d0e1d1c16e29b15c7ad810536e249c.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "INSERT INTO script (workspace_id, hash, path, parent_hashes, summary, description, content, created_by, schema, is_template, extra_perms, lock, language, kind, tag, draft_only, envs, concurrent_limit, concurrency_time_window_s) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9::text::json, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19)", + "query": "INSERT INTO script (workspace_id, hash, path, parent_hashes, summary, description, content, created_by, schema, is_template, extra_perms, lock, language, kind, tag, draft_only, envs, concurrent_limit, concurrency_time_window_s, cache_ttl) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9::text::json, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20)", "describe": { "columns": [], "parameters": { @@ -57,10 +57,11 @@ "Bool", "VarcharArray", "Int4", + "Int4", "Int4" ] }, "nullable": [] }, - "hash": "9ee1423945740c17dedc2fc7ea7e7183380c43d0feeb615c651aa486eef3bd53" + "hash": "52dbe77fad5729e9e121609eeb6a3573e6d0e1d1c16e29b15c7ad810536e249c" } diff --git a/backend/.sqlx/query-01d5e4e371a507602fca01190965a4f20ff806315a9ed02b372977e2d9479338.json b/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json similarity index 84% rename from backend/.sqlx/query-01d5e4e371a507602fca01190965a4f20ff806315a9ed02b372977e2d9479338.json rename to backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json index 14bba19013..6308bf3bb2 100644 --- a/backend/.sqlx/query-01d5e4e371a507602fca01190965a4f20ff806315a9ed02b372977e2d9479338.json +++ b/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "INSERT INTO queue\n (workspace_id, id, running, parent_job, created_by, permissioned_as, scheduled_for, \n script_hash, script_path, raw_code, raw_lock, args, job_kind, schedule_path, raw_flow, flow_status, is_flow_step, language, started_at, same_worker, pre_run_error, email, visible_to_owner, root_job, tag, concurrent_limit, concurrency_time_window_s, timeout, flow_step_id)\n VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, now()), $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, CASE WHEN $3 THEN now() END, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28) RETURNING id", + "query": "INSERT INTO queue\n (workspace_id, id, running, parent_job, created_by, permissioned_as, scheduled_for, \n script_hash, script_path, raw_code, raw_lock, args, job_kind, schedule_path, raw_flow, flow_status, is_flow_step, language, started_at, same_worker, pre_run_error, email, visible_to_owner, root_job, tag, concurrent_limit, concurrency_time_window_s, timeout, flow_step_id, cache_ttl)\n VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, now()), $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, CASE WHEN $3 THEN now() END, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29) RETURNING id", "describe": { "columns": [ { @@ -80,12 +80,13 @@ "Int4", "Int4", "Int4", - "Varchar" + "Varchar", + "Int4" ] }, "nullable": [ false ] }, - "hash": "01d5e4e371a507602fca01190965a4f20ff806315a9ed02b372977e2d9479338" + "hash": "9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b" } diff --git a/backend/.sqlx/query-7435be9a71a6b30980d180047486a47d69e36132d40ef0f53448e96a9866842b.json b/backend/.sqlx/query-cf21b61aa5b5e47744235c860e2d9b90d7654a4b5d0ec32085b178399f701bfc.json similarity index 63% rename from backend/.sqlx/query-7435be9a71a6b30980d180047486a47d69e36132d40ef0f53448e96a9866842b.json rename to backend/.sqlx/query-cf21b61aa5b5e47744235c860e2d9b90d7654a4b5d0ec32085b178399f701bfc.json index ad92a17eb4..58c600f2ca 100644 --- a/backend/.sqlx/query-7435be9a71a6b30980d180047486a47d69e36132d40ef0f53448e96a9866842b.json +++ b/backend/.sqlx/query-cf21b61aa5b5e47744235c860e2d9b90d7654a4b5d0ec32085b178399f701bfc.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "select hash, tag, concurrent_limit, concurrency_time_window_s from script where path = $1 AND workspace_id = $2 AND\n created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND\n deleted = false AND archived = false AND lock IS not NULL AND lock_error_logs IS NULL)", + "query": "select hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where path = $1 AND workspace_id = $2 AND\n created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND\n deleted = false AND archived = false)", "describe": { "columns": [ { @@ -22,6 +22,11 @@ "ordinal": 3, "name": "concurrency_time_window_s", "type_info": "Int4" + }, + { + "ordinal": 4, + "name": "cache_ttl", + "type_info": "Int4" } ], "parameters": { @@ -34,8 +39,9 @@ false, true, true, + true, true ] }, - "hash": "7435be9a71a6b30980d180047486a47d69e36132d40ef0f53448e96a9866842b" + "hash": "cf21b61aa5b5e47744235c860e2d9b90d7654a4b5d0ec32085b178399f701bfc" } diff --git a/backend/.sqlx/query-d164caec6637ab53d8c2e0faad653a458f31a976bb23c3a43d35b8476b066954.json b/backend/.sqlx/query-ebac391c83dae67831502fbe3f1034210afce7902a69f7c42bb68160e0464d71.json similarity index 72% rename from backend/.sqlx/query-d164caec6637ab53d8c2e0faad653a458f31a976bb23c3a43d35b8476b066954.json rename to backend/.sqlx/query-ebac391c83dae67831502fbe3f1034210afce7902a69f7c42bb68160e0464d71.json index ce1e44088d..ed90669699 100644 --- a/backend/.sqlx/query-d164caec6637ab53d8c2e0faad653a458f31a976bb23c3a43d35b8476b066954.json +++ b/backend/.sqlx/query-ebac391c83dae67831502fbe3f1034210afce7902a69f7c42bb68160e0464d71.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "select path, tag, concurrent_limit, concurrency_time_window_s from script where hash = $1 AND workspace_id = $2", + "query": "select path, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where hash = $1 AND workspace_id = $2", "describe": { "columns": [ { @@ -22,6 +22,11 @@ "ordinal": 3, "name": "concurrency_time_window_s", "type_info": "Int4" + }, + { + "ordinal": 4, + "name": "cache_ttl", + "type_info": "Int4" } ], "parameters": { @@ -34,8 +39,9 @@ false, true, true, + true, true ] }, - "hash": "d164caec6637ab53d8c2e0faad653a458f31a976bb23c3a43d35b8476b066954" + "hash": "ebac391c83dae67831502fbe3f1034210afce7902a69f7c42bb68160e0464d71" } diff --git a/backend/migrations/20230828202124_add_cache_ttl.down.sql b/backend/migrations/20230828202124_add_cache_ttl.down.sql new file mode 100644 index 0000000000..d2f607c5b8 --- /dev/null +++ b/backend/migrations/20230828202124_add_cache_ttl.down.sql @@ -0,0 +1 @@ +-- Add down migration script here diff --git a/backend/migrations/20230828202124_add_cache_ttl.up.sql b/backend/migrations/20230828202124_add_cache_ttl.up.sql new file mode 100644 index 0000000000..bec8c5f44f --- /dev/null +++ b/backend/migrations/20230828202124_add_cache_ttl.up.sql @@ -0,0 +1,3 @@ +-- Add up migration script here +ALTER TABLE script ADD COLUMN cache_ttl INTEGER; +ALTER TABLE queue ADD COLUMN cache_ttl INTEGER; diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index fbf228cda8..b3fd54db8d 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1619,6 +1619,7 @@ func main(derp string) (string, error) { language: ScriptLang::Go, concurrent_limit: None, concurrency_time_window_s: None, + cache_ttl: None })) .arg("derp", json!("world")) .run_until_complete(&db, port) @@ -1648,6 +1649,7 @@ echo "hello $msg" language: ScriptLang::Bash, concurrent_limit: None, concurrency_time_window_s: None, + cache_ttl: None })) .arg("msg", json!("world")) .run_until_complete(&db, port) @@ -1675,6 +1677,7 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, + cache_ttl: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) @@ -1707,6 +1710,7 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, + cache_ttl: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) @@ -1738,6 +1742,7 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, + cache_ttl: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 7def7adfc4..27b05fa726 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -5831,7 +5831,8 @@ components: type: array items: type: integer - + cache_ttl: + type: number required: - hash - path @@ -5905,6 +5906,8 @@ components: type: array items: type: integer + cache_ttl: + type: number required: - path - summary diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index bf2d7ef403..7e4f6158c4 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -210,7 +210,7 @@ async fn create_flow( w_id, nf.path, nf.summary, - nf.description, + nf.description.unwrap_or_else(String::new), nf.value, &authed.username, nf.schema.and_then(|x| serde_json::to_string(&x.0).ok()), @@ -352,7 +352,7 @@ async fn update_flow( edited_at = now(), schema = $6::text::json, dependency_job = NULL, draft_only = NULL, tag = $9 WHERE path = $7 AND workspace_id = $8", nf.path, nf.summary, - nf.description, + nf.description.unwrap_or_else(String::new), nf.value, &authed.username, schema.and_then(|x| serde_json::to_string(&x).ok()), @@ -752,6 +752,7 @@ mod tests { concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }; let expect = serde_json::json!({ "modules": [ diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index b846e9ec9e..4d75b78971 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -305,13 +305,19 @@ pub async fn get_path_for_hash<'c>( Ok(path) } -pub async fn get_path_tag_and_limits_for_hash( +pub async fn get_path_tag_limits_cache_for_hash( db: &DB, w_id: &str, hash: i64, -) -> error::Result<(String, Option, Option, Option)> { +) -> error::Result<( + String, + Option, + Option, + Option, + Option, +)> { let script = sqlx::query!( - "select path, tag, concurrent_limit, concurrency_time_window_s from script where hash = $1 AND workspace_id = $2", + "select path, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where hash = $1 AND workspace_id = $2", hash, w_id ) @@ -327,6 +333,7 @@ pub async fn get_path_tag_and_limits_for_hash( script.tag, script.concurrent_limit, script.concurrency_time_window_s, + script.cache_ttl, )) } @@ -1340,6 +1347,7 @@ impl From for Job { concurrency_time_window_s: uj.concurrency_time_window_s, timeout: None, flow_step_id: None, + cache_ttl: None, }), t => panic!("job type {} not valid", t), } @@ -2027,8 +2035,8 @@ pub async fn run_wait_result_script_by_hash( check_queue_too_long(&db, run_query.queue_limit).await?; let hash = script_hash.0; - let (path, tag, concurrent_limit, concurrency_time_window_s) = - get_path_tag_and_limits_for_hash(&db, &w_id, hash).await?; + let (path, tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = + get_path_tag_limits_cache_for_hash(&db, &w_id, hash).await?; check_scopes(&authed, || format!("run:script/{path}"))?; let args = run_query.add_include_headers(headers, args.unwrap_or_default()); @@ -2045,6 +2053,7 @@ pub async fn run_wait_result_script_by_hash( path: path, concurrent_limit: concurrent_limit, concurrency_time_window_s: concurrency_time_window_s, + cache_ttl, }, args, &authed.username, @@ -2218,6 +2227,7 @@ async fn run_preview_job( lock: None, concurrent_limit: None, // TODO(gbouv): once I find out how to store limits in the content of a script, should be easy to plug limits here concurrency_time_window_s: None, // TODO(gbouv): same as above + cache_ttl: None, }), }, args, @@ -2350,8 +2360,8 @@ pub async fn run_job_by_hash( JsonOrForm(args, raw_string): JsonOrForm, ) -> error::Result<(StatusCode, String)> { let hash = script_hash.0; - let (path, tag, concurrent_limit, concurrency_time_window_s) = - get_path_tag_and_limits_for_hash(&db, &w_id, hash).await?; + let (path, tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = + get_path_tag_limits_cache_for_hash(&db, &w_id, hash).await?; check_scopes(&authed, || format!("run:script/{path}"))?; let scheduled_for = run_query.get_scheduled_for(&db).await?; @@ -2369,6 +2379,7 @@ pub async fn run_job_by_hash( path: path, concurrent_limit: concurrent_limit, concurrency_time_window_s: concurrency_time_window_s, + cache_ttl, }, args, &authed.username, diff --git a/backend/windmill-api/src/oauth2.rs b/backend/windmill-api/src/oauth2.rs index 065603aebd..0e86bb4d74 100644 --- a/backend/windmill-api/src/oauth2.rs +++ b/backend/windmill-api/src/oauth2.rs @@ -799,7 +799,7 @@ async fn slack_command( (JobPayload::Flow(path.to_string()), None) } else { let path = path.strip_prefix("script/").unwrap_or_else(|| path); - let (script_hash, tag, concurrent_limit, concurrency_time_window_s) = + let (script_hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = windmill_common::get_latest_deployed_hash_for_path( &db, &settings.workspace_id, @@ -812,6 +812,7 @@ async fn slack_command( path: path.to_owned(), concurrent_limit, concurrency_time_window_s, + cache_ttl, }, tag, ) diff --git a/backend/windmill-api/src/scripts.rs b/backend/windmill-api/src/scripts.rs index 919daed481..b4f05d3093 100644 --- a/backend/windmill-api/src/scripts.rs +++ b/backend/windmill-api/src/scripts.rs @@ -69,6 +69,8 @@ pub struct ScriptWDraft { pub concurrent_limit: Option, #[serde(skip_serializing_if = "Option::is_none")] pub concurrency_time_window_s: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cache_ttl: Option, } pub fn global_service() -> Router { @@ -370,8 +372,8 @@ async fn create_script( sqlx::query!( "INSERT INTO script (workspace_id, hash, path, parent_hashes, summary, description, \ content, created_by, schema, is_template, extra_perms, lock, language, kind, tag, \ - draft_only, envs, concurrent_limit, concurrency_time_window_s) \ - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9::text::json, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19)", + draft_only, envs, concurrent_limit, concurrency_time_window_s, cache_ttl) \ + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9::text::json, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20)", &w_id, &hash.0, ns.path, @@ -390,7 +392,8 @@ async fn create_script( ns.draft_only, envs, ns.concurrent_limit, - ns.concurrency_time_window_s + ns.concurrency_time_window_s, + ns.cache_ttl, ) .execute(&mut tx) .await?; @@ -587,7 +590,7 @@ async fn get_script_by_path_w_draft( let mut tx = user_db.begin(&authed).await?; let script_o = sqlx::query_as::<_, ScriptWDraft>( - "SELECT hash, script.path, summary, description, content, language, kind, tag, schema, draft_only, envs, concurrent_limit, concurrency_time_window_s, draft.value as draft FROM script LEFT JOIN draft ON + "SELECT hash, script.path, summary, description, content, language, kind, tag, schema, draft_only, envs, concurrent_limit, concurrency_time_window_s, cache_ttl, draft.value as draft FROM script LEFT JOIN draft ON script.path = draft.path AND script.workspace_id = draft.workspace_id AND draft.typ = 'script' WHERE script.path = $1 AND script.workspace_id = $2 \ AND script.created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND \ diff --git a/backend/windmill-common/src/flows.rs b/backend/windmill-common/src/flows.rs index 33e663ec86..a06f54129a 100644 --- a/backend/windmill-common/src/flows.rs +++ b/backend/windmill-common/src/flows.rs @@ -60,7 +60,7 @@ pub struct ListableFlow { pub struct NewFlow { pub path: String, pub summary: String, - pub description: String, + pub description: Option, pub value: serde_json::Value, pub schema: Option, pub draft_only: Option, @@ -83,6 +83,8 @@ pub struct FlowValue { pub concurrency_time_window_s: Option, #[serde(skip_serializing_if = "Option::is_none")] pub skip_expr: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cache_ttl: Option, } #[derive(Deserialize, Serialize, Debug, Clone)] diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index 3206e241c0..5e6e772165 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -90,6 +90,8 @@ pub struct QueuedJob { pub timeout: Option, #[serde(skip_serializing_if = "Option::is_none")] pub flow_step_id: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cache_ttl: Option, } impl QueuedJob { @@ -99,6 +101,9 @@ impl QueuedJob { .map(String::as_str) .unwrap_or("tmp/main") } + pub fn is_flow(&self) -> bool { + matches!(self.job_kind, JobKind::Flow | JobKind::FlowPreview) + } } impl QueuedJob { @@ -156,6 +161,7 @@ impl Default for QueuedJob { concurrency_time_window_s: None, timeout: None, flow_step_id: None, + cache_ttl: None, } } } @@ -170,6 +176,7 @@ pub enum JobPayload { path: String, concurrent_limit: Option, concurrency_time_window_s: Option, + cache_ttl: Option, }, Code(RawCode), Dependencies { @@ -202,6 +209,7 @@ pub struct RawCode { pub lock: Option, pub concurrent_limit: Option, pub concurrency_time_window_s: Option, + pub cache_ttl: Option, } type Tag = String; @@ -216,7 +224,7 @@ pub async fn script_path_to_payload( let (job_payload, tag) = if script_path.starts_with("hub/") { (JobPayload::ScriptHub { path: script_path.to_owned() }, None) } else { - let (script_hash, tag, concurrent_limit, concurrency_time_window_s) = + let (script_hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = get_latest_deployed_hash_for_path(db, w_id, script_path).await?; ( JobPayload::ScriptHash { @@ -224,6 +232,7 @@ pub async fn script_path_to_payload( path: script_path.to_owned(), concurrent_limit, concurrency_time_window_s, + cache_ttl: cache_ttl, }, tag, ) @@ -235,9 +244,9 @@ pub async fn script_hash_to_tag_and_limits<'c>( script_hash: &ScriptHash, db: &mut Transaction<'c, Postgres>, w_id: &String, -) -> error::Result<(Option, Option, Option)> { +) -> error::Result<(Option, Option, Option, Option)> { let script = sqlx::query!( - "select tag, concurrent_limit, concurrency_time_window_s from script where hash = $1 AND workspace_id = $2", + "select tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where hash = $1 AND workspace_id = $2", script_hash.0, w_id ) @@ -252,6 +261,7 @@ pub async fn script_hash_to_tag_and_limits<'c>( script.tag, script.concurrent_limit, script.concurrency_time_window_s, + script.cache_ttl, )) } diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 7347dc8c5d..ab03430814 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -163,9 +163,15 @@ pub async fn get_latest_deployed_hash_for_path( db: &DB, w_id: &str, script_path: &str, -) -> error::Result<(scripts::ScriptHash, Option, Option, Option)> { +) -> error::Result<( + scripts::ScriptHash, + Option, + Option, + Option, + Option, +)> { let r_o = sqlx::query!( - "select hash, tag, concurrent_limit, concurrency_time_window_s from script where path = $1 AND workspace_id = $2 AND + "select hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where path = $1 AND workspace_id = $2 AND created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND deleted = false AND archived = false AND lock IS not NULL AND lock_error_logs IS NULL)", script_path, @@ -181,6 +187,7 @@ pub async fn get_latest_deployed_hash_for_path( script.tag, script.concurrent_limit, script.concurrency_time_window_s, + script.cache_ttl, )) } @@ -188,9 +195,15 @@ pub async fn get_latest_hash_for_path<'c>( db: &mut sqlx::Transaction<'c, sqlx::Postgres>, w_id: &str, script_path: &str, -) -> error::Result<(scripts::ScriptHash, Option, Option, Option)> { +) -> error::Result<( + scripts::ScriptHash, + Option, + Option, + Option, + Option, +)> { let r_o = sqlx::query!( - "select hash, tag, concurrent_limit, concurrency_time_window_s from script where path = $1 AND workspace_id = $2 AND + "select hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl from script where path = $1 AND workspace_id = $2 AND created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND deleted = false AND archived = false)", script_path, @@ -206,5 +219,6 @@ pub async fn get_latest_hash_for_path<'c>( script.tag, script.concurrent_limit, script.concurrency_time_window_s, + script.cache_ttl, )) } diff --git a/backend/windmill-common/src/scripts.rs b/backend/windmill-common/src/scripts.rs index 31a313bce8..dacb538d6a 100644 --- a/backend/windmill-common/src/scripts.rs +++ b/backend/windmill-common/src/scripts.rs @@ -214,6 +214,7 @@ pub struct NewScript { pub envs: Option>, pub concurrent_limit: Option, pub concurrency_time_window_s: Option, + pub cache_ttl: Option, } #[derive(Deserialize)] diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 53c71ffacc..63b0b98912 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1409,8 +1409,15 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( language, concurrent_limit, concurrency_time_window_s, + cache_ttl, ) = match job_payload { - JobPayload::ScriptHash { hash, path, concurrent_limit, concurrency_time_window_s } => { + JobPayload::ScriptHash { + hash, + path, + concurrent_limit, + concurrency_time_window_s, + cache_ttl, + } => { let language = fetch_scalar_isolated!(sqlx::query_scalar!( "SELECT language as \"language: ScriptLang\" FROM script WHERE hash = $1 AND workspace_id = $2", hash.0, @@ -1431,6 +1438,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( Some(language), concurrent_limit, concurrency_time_window_s, + cache_ttl, ) } JobPayload::ScriptHub { path } => { @@ -1444,6 +1452,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, None, None, + None, ) } JobPayload::Code(RawCode { @@ -1453,6 +1462,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( lock, concurrent_limit, concurrency_time_window_s, + cache_ttl, }) => ( None, path, @@ -1462,6 +1472,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( Some(language), concurrent_limit, concurrency_time_window_s, + cache_ttl, ), JobPayload::Dependencies { hash, dependencies, language, path } => ( Some(hash.0), @@ -1472,6 +1483,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( Some(language), None, None, + None, ), JobPayload::FlowDependencies { path } => { let value_json = fetch_scalar_isolated!( @@ -1497,6 +1509,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, None, None, + None, ) } JobPayload::AppDependencies { path, version } => ( @@ -1508,6 +1521,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, None, None, + None, ), JobPayload::RawFlow { value, path } => ( None, @@ -1518,6 +1532,7 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, value.concurrent_limit.clone(), value.concurrency_time_window_s, + value.cache_ttl.map(|x| x as i32), ), JobPayload::Flow(flow) => { let value_json = fetch_scalar_isolated!( @@ -1543,10 +1558,31 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( None, value.concurrent_limit.clone(), value.concurrency_time_window_s, + value.cache_ttl.map(|x| x as i32), ) } - JobPayload::Identity => (None, None, None, JobKind::Identity, None, None, None, None), - JobPayload::Noop => (None, None, None, JobKind::Noop, None, None, None, None), + JobPayload::Identity => ( + None, + None, + None, + JobKind::Identity, + None, + None, + None, + None, + None, + ), + JobPayload::Noop => ( + None, + None, + None, + JobKind::Noop, + None, + None, + None, + None, + None, + ), }; let is_running = same_worker; @@ -1657,8 +1693,8 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( "INSERT INTO queue (workspace_id, id, running, parent_job, created_by, permissioned_as, scheduled_for, script_hash, script_path, raw_code, raw_lock, args, job_kind, schedule_path, raw_flow, \ - flow_status, is_flow_step, language, started_at, same_worker, pre_run_error, email, visible_to_owner, root_job, tag, concurrent_limit, concurrency_time_window_s, timeout, flow_step_id) - VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, now()), $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, CASE WHEN $3 THEN now() END, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28) \ + flow_status, is_flow_step, language, started_at, same_worker, pre_run_error, email, visible_to_owner, root_job, tag, concurrent_limit, concurrency_time_window_s, timeout, flow_step_id, cache_ttl) + VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, now()), $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, CASE WHEN $3 THEN now() END, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29) \ RETURNING id", workspace_id, job_id, @@ -1687,7 +1723,8 @@ pub async fn push<'c, R: rsmq_async::RsmqConnection + Send + 'c>( concurrent_limit, concurrency_time_window_s, custom_timeout, - flow_step_id + flow_step_id, + cache_ttl ) .fetch_one(&mut tx) .await diff --git a/backend/windmill-queue/src/schedule.rs b/backend/windmill-queue/src/schedule.rs index 843a5f4476..460e2dc603 100644 --- a/backend/windmill-queue/src/schedule.rs +++ b/backend/windmill-queue/src/schedule.rs @@ -74,7 +74,7 @@ pub async fn push_scheduled_job<'c, R: rsmq_async::RsmqConnection + Send + 'c>( let (payload, tag) = if schedule.is_flow { (JobPayload::Flow(schedule.script_path), None) } else { - let (hash, tag, concurrent_limit, concurrency_time_window_s) = + let (hash, tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = windmill_common::get_latest_hash_for_path( tx.transaction_mut(), &schedule.workspace_id, @@ -87,6 +87,7 @@ pub async fn push_scheduled_job<'c, R: rsmq_async::RsmqConnection + Send + 'c>( path: schedule.script_path, concurrent_limit: concurrent_limit, concurrency_time_window_s: concurrency_time_window_s, + cache_ttl: cache_ttl, }, tag, ) diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 01b0fa5575..29498b7673 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -2,6 +2,7 @@ use async_recursion::async_recursion; use serde_json::{json, Value}; use sqlx::{Pool, Postgres}; use tokio::{fs::File, io::AsyncReadExt}; +use windmill_api_client::{types::CreateResource, Client}; use windmill_common::{ error::{self, Error}, jobs::QueuedJob, @@ -10,7 +11,12 @@ use windmill_queue::CLOUD_HOSTED; use anyhow::Result; use std::{ - borrow::Borrow, collections::HashMap, io, os::unix::process::ExitStatusExt, panic, + borrow::Borrow, + collections::{hash_map::DefaultHasher, HashMap}, + hash::{Hash, Hasher}, + io, + os::unix::process::ExitStatusExt, + panic, time::Duration, }; @@ -566,6 +572,32 @@ fn append_with_limit(dst: &mut String, src: &str, limit: &mut usize) { } } +pub fn hash_args(v: &serde_json::Value) -> i64 { + let mut dh = DefaultHasher::new(); + serde_json::to_string(v).unwrap().hash(&mut dh); + dh.finish() as i64 +} + +pub async fn save_in_cache(client: &AuthedClient, job: &QueuedJob, cached_path: String, r: &Value) { + let client: &Client = client.get_client(); + let expire = chrono::Utc::now().timestamp() + job.cache_ttl.unwrap() as i64; + let cr = &CreateResource { + path: cached_path, + description: None, + resource_type: "cache".to_string(), + value: serde_json::json!({ + "value": r, + "expire": expire + }), + }; + if let Err(e) = client + .create_resource(&job.workspace_id, Some(true), cr) + .await + { + tracing::error!("Error creating cache resource {e}") + } +} + /* TODO retry this? */ #[tracing::instrument(level = "trace", skip_all)] async fn append_logs(job_id: uuid::Uuid, logs: impl AsRef, db: impl Borrow>) { diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 1c184b0f3f..59ec4762c6 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -14,9 +14,7 @@ use sqlx::{Pool, Postgres}; use windmill_api_client::Client; use std::{ collections::HashMap, time::Duration, - sync::{Arc, atomic::Ordering}, - collections::hash_map::DefaultHasher, - hash::{Hasher, Hash}, + sync::{Arc, atomic::Ordering} }; use tracing::{trace_span, Instrument}; @@ -26,7 +24,7 @@ use windmill_common::{ flows::{FlowModuleValue, FlowValue, FlowModule}, scripts::{ScriptHash, ScriptLang, get_full_hub_script_by_path}, utils::{rd_string, StripPath}, - users::SUPERADMIN_SECRET_EMAIL, jobs::{JobKind, QueuedJob, Metrics}, METRICS_ENABLED, IS_READY, + users::SUPERADMIN_SECRET_EMAIL, jobs::{JobKind, QueuedJob, Metrics}, METRICS_ENABLED, IS_READY, DB, }; use windmill_queue::{canceled_job_to_result, get_queued_job, pull, CLOUD_HOSTED, HTTP_CLIENT, ACCEPTED_TAGS, IS_WORKER_TAGS_DEFINED}; @@ -43,7 +41,6 @@ use tokio::{ use futures::future::FutureExt; use async_recursion::async_recursion; -use windmill_api_client::types::CreateResource; #[cfg(feature = "enterprise")] use rand::Rng; @@ -55,8 +52,8 @@ use windmill_queue::{add_completed_job, add_completed_job_error,IDLE_WORKERS}; use crate::{ worker_flow::{ - handle_flow, update_flow_status_after_job_completion, update_flow_status_in_progress, - }, python_executor::{create_dependencies_dir, pip_compile, handle_python_job, handle_python_reqs}, common::{read_result, set_logs, write_file, transform_json_value}, go_executor::{handle_go_job, install_go_dependencies}, js_eval::{transpile_ts, eval_fetch_timeout}, pg_executor::do_postgresql, mysql_executor::do_mysql, graphql_executor::do_graphql, bun_executor::{handle_bun_job, gen_lockfile}, bash_executor::{ANSI_ESCAPE_RE, handle_powershell_job, handle_bash_job}, deno_executor::{handle_deno_job, generate_deno_lock}, + handle_flow, update_flow_status_after_job_completion, update_flow_status_in_progress, + }, python_executor::{create_dependencies_dir, pip_compile, handle_python_job, handle_python_reqs}, common::{read_result, set_logs, write_file, transform_json_value, save_in_cache, hash_args}, go_executor::{handle_go_job, install_go_dependencies}, js_eval::{transpile_ts, eval_fetch_timeout}, pg_executor::do_postgresql, mysql_executor::do_mysql, graphql_executor::do_graphql, bun_executor::{handle_bun_job, gen_lockfile}, bash_executor::{ANSI_ESCAPE_RE, handle_powershell_job, handle_bash_job}, deno_executor::{handle_deno_job, generate_deno_lock}, }; #[cfg(feature = "enterprise")] @@ -908,12 +905,6 @@ pub struct JobCompleted { pub success: bool } -fn hash_args(v: &serde_json::Value) -> i64 { - let mut dh = DefaultHasher::new(); - serde_json::to_string(v).unwrap().hash(&mut dh); - dh.finish() as i64 -} - pub async fn get_content(job: &QueuedJob, db: &Pool) -> Result { @@ -960,7 +951,7 @@ async fn do_nativets(job: QueuedJob, logs: String, client: &AuthedClient, code: #[tracing::instrument(level = "trace", skip_all)] async fn handle_queued_job( job: QueuedJob, - db: &sqlx::Pool, + db: &DB, client: &AuthedClientBackgroundTask, worker_name: &str, worker_dir: &str, @@ -978,20 +969,63 @@ async fn handle_queued_job( return Err(Error::ExecutionErr(e)); } - let (cache_ttl, step) = if job.is_flow_step { - update_flow_status_in_progress( + let step = if job.is_flow_step { + let r = update_flow_status_in_progress( db, &job.workspace_id, job.parent_job .ok_or_else(|| Error::InternalErr(format!("expected parent job")))?, job.id, ) - .await? + .await?; + + r } else { - (None, None) + None }; - + let cached_res_path = if job.cache_ttl.is_some() { + let args_hash = hash_args(&job.args.clone().unwrap_or_else(|| json!({}))); + let permissioned_as = &job.permissioned_as; + if job.is_flow_step { + let flow_path = sqlx::query_scalar!( + "SELECT script_path FROM queue WHERE id = $1", + &job.parent_job.unwrap() + ) + .fetch_one(db) + .await + .map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))? + .ok_or_else(|| Error::InternalErr(format!("Expected script_path")))?; + let step = step.unwrap_or(-1); + Some(format!("{permissioned_as}/cache/{flow_path}/{step}/{args_hash}")) + } else if let Some(script_path) = &job.script_path { + let is_flow = if job.is_flow() { "flow/" } else { "" }; + Some(format!("{permissioned_as}/cache/{is_flow}{script_path}/{args_hash}")) + } else { + None + } + + } else { + None + }; + + if let Some(cached_res_path) = cached_res_path.clone() { + let authed_client = client.get_authed().await; + let client: &Client = authed_client.get_client(); + let resource = client.get_resource_value(&job.workspace_id, &cached_res_path).await; + if let Ok(resource) = resource { + let v = resource.into_inner(); + if let Some(o) = v.as_object() { + let expire = o.get("expire"); + if expire.is_some() && expire.unwrap().as_i64().map(|x| x > chrono::Utc::now().timestamp()).unwrap_or(false) { + let result = v.get("value").map(|x| x.to_owned()).unwrap_or_else(|| json!({})); + let logs = "Job skipped because args & path found in cache and not expired".to_string(); + process_result(authed_client, job, Ok(result), None, db, worker_dir, job_dir, metrics, same_worker_tx, base_internal_url, rsmq, job_completed_tx, logs).await?; + return Ok(()) + } + } + } + }; match job.job_kind { JobKind::FlowPreview | JobKind::Flow => { let args = job.args.clone().unwrap_or(Value::Null); @@ -1019,25 +1053,6 @@ async fn handle_queued_job( set_logs(&logs, &job.id, db).await; - - - let cached_res_path = if cache_ttl.is_some() { - let flow_path = sqlx::query_scalar!( - "SELECT script_path FROM queue WHERE id = $1", - &job.parent_job.unwrap() - ) - .fetch_one(db) - .await - .map_err(|e| Error::InternalErr(format!("fetching step flow status: {e}")))? - .ok_or_else(|| Error::InternalErr(format!("Expected script_path")))?; - let step = step.unwrap_or(-1); - let args_hash = hash_args(&job.args.clone().unwrap_or_else(|| json!({}))); - let permissioned_as = &job.permissioned_as; - Some(format!("{permissioned_as}/cache/{flow_path}/{step}/{args_hash}")) - } else { - None - }; - tracing::debug!( worker = %worker_name, job_id = %job.id, @@ -1046,32 +1061,7 @@ async fn handle_queued_job( job.id ); - let cached_res = if let Some(cached_res_path) = cached_res_path.clone() { - let authed_client = client.get_authed().await; - let client: &Client = authed_client.get_client(); - let resource = client.get_resource_value(&job.workspace_id, &cached_res_path).await; - resource.ok() - .and_then(|x| { - let v = x.into_inner(); - if let Some(o) = v.as_object() { - let expire = o.get("expire"); - if expire.is_some() && expire.unwrap().as_i64().map(|x| x > chrono::Utc::now().timestamp()).unwrap_or(false) { - v.get("value").map(|x| x.to_owned()) - } else { - None - } - } else { - None - } - }) - } else { - None - }; - - let result = if let Some(cached_res) = cached_res { - Ok(cached_res) - } else { - match job.job_kind { + let result = match job.job_kind { JobKind::Dependencies => { handle_dependency_job(&job, &mut logs, job_dir, db, worker_name, worker_dir, base_internal_url, &client.get_token().await).await } @@ -1105,119 +1095,127 @@ async fn handle_queued_job( worker_name ) .await } - } - }; + }; //it's a test job, no need to update the db if job.workspace_id == "" { return Ok(()); } - let client = &client.get_authed().await; - match result { - Ok(r) => { - // println!("bef completed job{:?}", SystemTime::now()); - if let Some(cached_path) = cached_res_path { - let client: &Client = client.get_client(); - let expire = chrono::Utc::now().timestamp() + cache_ttl.unwrap() as i64; - let cr = &CreateResource { - path: cached_path, - description: None, - resource_type: "cache".to_string(), - value: serde_json::json!({ - "value": r, - "expire": expire - }) - }; - if let Err(e) = client.create_resource(&job.workspace_id, Some(true), cr).await { - tracing::error!("Error creating cache resource {e}") - } - } - if job.is_flow_step { - add_completed_job(db, &job, true, false, r.clone(), logs, rsmq.clone()).await?; - if let Some(parent_job) = job.parent_job { - update_flow_status_after_job_completion( - db, - client, - parent_job, - &job.id, - &job.workspace_id, - true, - r, - metrics.clone(), - false, - same_worker_tx.clone(), - worker_dir, - None, - base_internal_url, - rsmq.clone() - ) - .await?; - } - } else { - // in the happy path and if job not a flow step, we can delegate updating the completed job in the background - job_completed_tx.send(JobCompleted{job,result:r,logs:logs, success: true}).await.expect("send job completed"); - - } - } - Err(e) => { - let error_value = match e { - Error::ExitStatus(i) => { - let res = read_result(job_dir).await.ok(); - - if res.is_some() && res.clone().unwrap().is_object() { - res.unwrap() - } else { - let last_10_log_lines = logs - .lines() - .skip(logs.lines().count().max(13) - 13) - .join("\n") - .to_string() - .replace("\n\n", "\n"); - - let log_lines = last_10_log_lines - .split("CODE EXECUTION ---") - .last() - .unwrap_or(&logs); - - extract_error_value(log_lines, i) - } - } - err @ _ => { - json!({"message": format!("error during execution of the script:\n{}", err), "name": "ExecutionErr"}) - } - }; - - let result = - add_completed_job_error(db, &job, logs, error_value, metrics.clone(), rsmq.clone()) - .await?; - if job.is_flow_step { - if let Some(parent_job) = job.parent_job { - update_flow_status_after_job_completion( - db, - client, - parent_job, - &job.id, - &job.workspace_id, - false, - result, - metrics, - false, - same_worker_tx, - worker_dir, - None, - base_internal_url, - rsmq - ) - .await?; - } - } - } - }; + let client = client.get_authed().await; + process_result(client, job, result, cached_res_path, db, worker_dir, job_dir, metrics, same_worker_tx, base_internal_url, rsmq, job_completed_tx, logs).await?; } } Ok(()) } +async fn process_result( + client: AuthedClient, + job: QueuedJob, + result: error::Result, + cached_res_path: Option, + db: &DB, + worker_dir: &str, + job_dir: &str, + metrics: Option, + same_worker_tx: Sender, + base_internal_url: &str, + rsmq: Option, + job_completed_tx: Sender, + logs: String, +) -> error::Result<()> { + match result { + Ok(r) => { + // println!("bef completed job{:?}", SystemTime::now()); + if let Some(cached_path) = cached_res_path { + save_in_cache(&client, &job, cached_path, &r).await; + } + if job.is_flow_step { + + add_completed_job(db, &job, true, false, r.clone(), logs, rsmq.clone()).await?; + if let Some(parent_job) = job.parent_job { + update_flow_status_after_job_completion( + db, + &client, + parent_job, + &job.id, + &job.workspace_id, + true, + r, + metrics.clone(), + false, + same_worker_tx.clone(), + worker_dir, + None, + base_internal_url, + rsmq.clone() + ) + .await?; + } + } else { + // in the happy path and if job not a flow step, we can delegate updating the completed job in the background + job_completed_tx.send(JobCompleted{job,result:r,logs:logs, success: true}).await.expect("send job completed"); + + } + } + Err(e) => { + let error_value = match e { + Error::ExitStatus(i) => { + let res = read_result(job_dir).await.ok(); + + if res.is_some() && res.clone().unwrap().is_object() { + res.unwrap() + } else { + let last_10_log_lines = logs + .lines() + .skip(logs.lines().count().max(13) - 13) + .join("\n") + .to_string() + .replace("\n\n", "\n"); + + let log_lines = last_10_log_lines + .split("CODE EXECUTION ---") + .last() + .unwrap_or(&logs); + + extract_error_value(log_lines, i) + } + } + err @ _ => { + json!({"message": format!("error during execution of the script:\n{}", err), "name": "ExecutionErr"}) + } + }; + + let result = + add_completed_job_error(db, &job, logs, error_value, metrics.clone(), rsmq.clone()) + .await?; + if job.is_flow_step { + if let Some(parent_job) = job.parent_job { + update_flow_status_after_job_completion( + db, + &client, + parent_job, + &job.id, + &job.workspace_id, + false, + result, + metrics, + false, + same_worker_tx, + worker_dir, + None, + base_internal_url, + rsmq + ) + .await?; + } + } + } + }; + Ok(()) +} + + + #[tracing::instrument(level = "trace", skip_all)] async fn handle_code_execution_job( job: &QueuedJob, diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index cd1da706c9..5e3eb341cd 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -9,6 +9,7 @@ use std::collections::HashMap; use std::time::Duration; +use crate::common::{hash_args, save_in_cache}; use crate::js_eval::{eval_timeout, IdContext}; use crate::{AuthedClient, KEEP_JOB_DIR}; use anyhow::Context; @@ -580,6 +581,16 @@ pub async fn update_flow_status_after_job_completion_internal< ) .await?; } else { + if flow_job.cache_ttl.is_some() { + let cached_res_path = { + let args_hash = hash_args(&flow_job.args.clone().unwrap_or_else(|| json!({}))); + let permissioned_as = &flow_job.permissioned_as; + let flow_path = flow_job.script_path(); + format!("{permissioned_as}/cache/flow/{flow_path}/{args_hash}") + }; + + save_in_cache(&client, &flow_job, cached_res_path, &nresult).await; + } add_completed_job( db, &flow_job, @@ -745,28 +756,26 @@ async fn compute_bool_from_expr( } } -type CacheAndStep = (Option, Option); pub async fn update_flow_status_in_progress( db: &DB, w_id: &str, flow: Uuid, job_in_progress: Uuid, -) -> error::Result { +) -> error::Result> { let step = get_step_of_flow_status(db, flow).await?; - let cache_ttl = if let Step::Step(step) = step { - let ttl = sqlx::query_scalar(&format!( + if let Step::Step(step) = step { + sqlx::query(&format!( "UPDATE queue SET flow_status = jsonb_set(jsonb_set(flow_status, '{{modules, {step}, job}}', $1), '{{modules, {step}, type}}', $2) - WHERE id = $3 AND workspace_id = $4 - RETURNING (raw_flow->'modules'->{step}->>'cache_ttl')::int as cache_ttl", + WHERE id = $3 AND workspace_id = $4", )) .bind(json!(job_in_progress.to_string())) .bind(json!("InProgress")) .bind(flow) .bind(w_id) - .fetch_one(db) + .execute(db) .await?; - (ttl, Some(step)) + Ok(Some(step)) } else { sqlx::query(&format!( "UPDATE queue @@ -779,9 +788,8 @@ pub async fn update_flow_status_in_progress( .bind(w_id) .execute(db) .await?; - (None, None) - }; - Ok(cache_ttl) + Ok(None) + } } pub enum Step { @@ -1410,6 +1418,7 @@ async fn push_next_flow_job concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }, path: Some(format!("{}/forloop", flow_job.script_path())), }, @@ -1777,7 +1786,7 @@ async fn compute_next_flow_transform( } else { let hash = script_hash.clone().unwrap(); let mut tx: sqlx::Transaction<'_, sqlx::Postgres> = db.begin().await?; - let (tag, concurrent_limit, concurrency_time_window_s) = + let (tag, concurrent_limit, concurrency_time_window_s, cache_ttl) = script_hash_to_tag_and_limits(&hash, &mut tx, &flow_job.workspace_id).await?; ( JobPayload::ScriptHash { @@ -1785,6 +1794,7 @@ async fn compute_next_flow_transform( path: script_path.to_owned(), concurrent_limit, concurrency_time_window_s, + cache_ttl: module.cache_ttl.map(|x| x as i32).ok_or(cache_ttl).ok(), }, tag, ) @@ -1816,6 +1826,7 @@ async fn compute_next_flow_transform( lock: lock.clone(), concurrent_limit: *concurrent_limit, concurrency_time_window_s: *concurrency_time_window_s, + cache_ttl: module.cache_ttl.map(|x| x as i32), }), tag: tag.clone(), }), @@ -1928,6 +1939,7 @@ async fn compute_next_flow_transform( concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }, path: Some(format!("{}/loop-{}", flow_job.script_path(), ns.index)), }, @@ -2002,6 +2014,7 @@ async fn compute_next_flow_transform( concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }, path: Some(format!( "{}/branchone-{}", @@ -2045,6 +2058,7 @@ async fn compute_next_flow_transform( concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }, path: Some(format!( "{}/branchall-{}", @@ -2119,6 +2133,7 @@ async fn compute_next_flow_transform( concurrent_limit: None, concurrency_time_window_s: None, skip_expr: None, + cache_ttl: None, }, path: Some(format!( "{}/branchall-{}", diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index f9106a3658..e7c6062735 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -10,7 +10,7 @@ import ScriptEditor from './ScriptEditor.svelte' import ScriptSchema from './ScriptSchema.svelte' import { dirtyStore } from './common/confirmationModal/dirtyStore' - import { Alert, Badge, Button, Drawer, Kbd } from './common' + import { Alert, Badge, Button, Drawer, Kbd, SecondsInput } from './common' import { faPlus, faSave } from '@fortawesome/free-solid-svg-icons' import LanguageIcon from './common/languageIcons/LanguageIcon.svelte' import type { SupportedLanguage } from '$lib/common' @@ -27,6 +27,7 @@ import { Icon } from 'svelte-awesome' import { fade } from 'svelte/transition' import Popover from './Popover.svelte' + import Toggle from './Toggle.svelte' export let script: NewScript export let initialPath: string = '' @@ -168,7 +169,8 @@ tag: script.tag, envs: script.envs, concurrent_limit: script.concurrent_limit, - concurrency_time_window_s: script.concurrency_time_window_s + concurrency_time_window_s: script.concurrency_time_window_s, + cache_ttl: script.cache_ttl } }) history.replaceState(history.state, '', `/scripts/edit/${script.path}`) @@ -210,7 +212,8 @@ draft_only: true, envs: script.envs, concurrent_limit: script.concurrent_limit, - concurrency_time_window_s: script.concurrency_time_window_s + concurrency_time_window_s: script.concurrency_time_window_s, + cache_ttl: script.cache_ttl } }) } @@ -443,6 +446,34 @@ {:else} {/if} + +

Cache { + if (script.cache_ttl && script.cache_ttl != undefined) { + script.cache_ttl = undefined + } else { + script.cache_ttl = 300 + } + }} + options={{ + right: 'Cache the results for each possible inputs' + }} + />

+ +
+
How long to the keep cache valid
+ + {#if script.cache_ttl} + + {:else} + + {/if} +
+ {#if !isCloudHosted()}

Custom env variables diff --git a/frontend/src/lib/components/flows/content/FlowModuleCache.svelte b/frontend/src/lib/components/flows/content/FlowModuleCache.svelte index 3740b220aa..8ac578c8b1 100644 --- a/frontend/src/lib/components/flows/content/FlowModuleCache.svelte +++ b/frontend/src/lib/components/flows/content/FlowModuleCache.svelte @@ -24,7 +24,7 @@ if (isCacheEnabled && flowModule.cache_ttl != undefined) { flowModule.cache_ttl = undefined } else { - flowModule.cache_ttl = 60 * 60 * 24 * 2 + flowModule.cache_ttl = 600 } }} options={{ diff --git a/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte b/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte index 694a04b4ac..34b777e7f7 100644 --- a/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte +++ b/frontend/src/lib/components/flows/content/FlowModuleComponent.svelte @@ -295,8 +295,8 @@ {#if !$selectedId.includes('failure')} Cache Concurrency - Early Stop/Break - Suspend/Approval + Early Stop + Suspend Sleep Mock Shared Directory diff --git a/frontend/src/lib/components/flows/content/FlowSettings.svelte b/frontend/src/lib/components/flows/content/FlowSettings.svelte index 5561b9034f..7279daecfa 100644 --- a/frontend/src/lib/components/flows/content/FlowSettings.svelte +++ b/frontend/src/lib/components/flows/content/FlowSettings.svelte @@ -61,6 +61,7 @@ Early Stop Worker Group Concurrency + Cache @@ -238,6 +239,36 @@ }} /> + +

Cache { + if ($flowStore.value.cache_ttl && $flowStore.value.cache_ttl != undefined) { + $flowStore.value.cache_ttl = undefined + } else { + $flowStore.value.cache_ttl = 300 + } + }} + options={{ + right: 'Cache the results for each possible inputs' + }} + />

+ +
+
How long to keep the cache valid
+
+ {#if $flowStore.value.cache_ttl} + + {:else} + + {/if} +
+
+
+ When a worker group is defined at the flow level, any steps inside the flow will run diff --git a/openflow.openapi.yaml b/openflow.openapi.yaml index b7f9efd3c9..9cada14180 100644 --- a/openflow.openapi.yaml +++ b/openflow.openapi.yaml @@ -52,6 +52,8 @@ components: type: number skip_expr: type: string + cache_ttl: + type: number required: - modules