diff --git a/backend/.sqlx/query-76440a1c5b40a897078817cc960247270c9044b363b71496ba36e1357c63a2a4.json b/backend/.sqlx/query-76440a1c5b40a897078817cc960247270c9044b363b71496ba36e1357c63a2a4.json new file mode 100644 index 0000000000..8eecbfa7bc --- /dev/null +++ b/backend/.sqlx/query-76440a1c5b40a897078817cc960247270c9044b363b71496ba36e1357c63a2a4.json @@ -0,0 +1,155 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH RECURSIVE lineage AS (\n SELECT id, parent_job, kind, runnable_path, runnable_id, trigger_kind, trigger,\n tag, permissioned_as, args, 0 AS depth\n FROM v2_job WHERE id = $1 AND workspace_id = $2\n UNION ALL\n SELECT p.id, p.parent_job, p.kind, p.runnable_path, p.runnable_id, p.trigger_kind,\n p.trigger, p.tag, p.permissioned_as, p.args, l.depth + 1\n FROM v2_job p JOIN lineage l ON p.id = l.parent_job\n WHERE p.workspace_id = $2 AND l.depth < 100\n )\n SELECT id AS \"id!\", parent_job, kind AS \"kind!: JobKind\", runnable_path, runnable_id,\n trigger_kind::text AS trigger_kind, trigger, tag AS \"tag!\",\n permissioned_as AS \"permissioned_as!\",\n -- A restart takes its flow version from the request, so a trusted flow path\n -- can carry another flow's code: the version must belong to that path.\n CASE kind\n WHEN 'flow' THEN EXISTS (SELECT 1 FROM flow_version fv\n WHERE fv.id = runnable_id AND fv.path = runnable_path AND fv.workspace_id = $2)\n WHEN 'script' THEN EXISTS (SELECT 1 FROM script s\n WHERE s.hash = runnable_id AND s.path = runnable_path AND s.workspace_id = $2\n AND NOT s.deleted)\n ELSE true\n END AS \"origin_verified!\",\n -- What `run/f` and `run/p` resolve the path to now; for scripts the predicate of\n -- `get_latest_deployed_script_hash`.\n COALESCE(CASE kind\n WHEN 'flow' THEN runnable_id = (SELECT f.versions[array_upper(f.versions, 1)]\n FROM flow f WHERE f.path = runnable_path AND f.workspace_id = $2)\n WHEN 'script' THEN runnable_id = (SELECT s.hash FROM script s\n WHERE s.path = runnable_path AND s.workspace_id = $2 AND NOT s.deleted\n AND s.lock IS NOT NULL AND s.lock_error_logs IS NULL\n ORDER BY s.created_at DESC LIMIT 1)\n END, false) AS \"current_version!\",\n EXISTS (SELECT 1 FROM v2_job_status st WHERE st.id = lineage.id\n AND jsonb_typeof(st.flow_status->'restarted_from') = 'object')\n OR EXISTS (SELECT 1 FROM v2_job_completed c WHERE c.id = lineage.id\n AND jsonb_typeof(c.flow_status->'restarted_from') = 'object')\n AS \"restarted!\",\n COALESCE(\n (SELECT st.flow_status->'restarted_from'->>'flow_job_id' FROM v2_job_status st\n WHERE st.id = lineage.id),\n (SELECT c.flow_status->'restarted_from'->>'flow_job_id' FROM v2_job_completed c\n WHERE c.id = lineage.id)\n ) AS restarted_from,\n -- Only deployed-app runs are stamped with their app; an app editor preview\n -- runs app code at an app path it does not have to own.\n COALESCE(trigger_kind = 'app' AND starts_with(runnable_path, trigger || '/'), false)\n AS \"app_stamped!\",\n -- A preview's modules come from its args, which its parent may have taken from\n -- the caller.\n COALESCE(jsonb_typeof(args->'_MODULES') = 'object', false) AS \"args_modules!\",\n CASE WHEN depth = 0 THEN (SELECT wp.worker_group FROM v2_job_queue q\n JOIN worker_ping wp ON wp.worker = q.worker WHERE q.id = lineage.id)\n END AS worker_group,\n CASE WHEN depth = 0 THEN CASE\n WHEN starts_with(permissioned_as, 'g/') THEN 'group'\n ELSE (SELECT CASE WHEN u.is_service_account THEN 'service_account' ELSE 'user' END\n FROM usr u WHERE starts_with(permissioned_as, 'u/')\n AND u.username = substr(permissioned_as, 3) AND u.workspace_id = $2)\n END END AS run_as_type,\n -- The canonical item `item_digest` hashes, null members included.\n CASE WHEN depth = 0 OR depth = max(depth) OVER () THEN CASE kind\n WHEN 'script' THEN (SELECT jsonb_build_object('content', s.content,\n 'lock', s.lock, 'codebase', s.codebase,\n 'modules', (SELECT jsonb_object_agg(m.key, m.value - 'language')\n FROM jsonb_each(s.modules) m))\n FROM script s WHERE s.hash = runnable_id AND s.workspace_id = $2\n AND NOT s.deleted)\n WHEN 'flow' THEN (SELECT fv.value FROM flow_version fv\n WHERE fv.id = runnable_id AND fv.workspace_id = $2)\n WHEN 'flowscript' THEN (SELECT jsonb_build_object('content', n.code,\n 'lock', n.lock)\n FROM flow_node n WHERE n.id = runnable_id AND n.workspace_id = $2)\n END END AS digest_source\n FROM lineage ORDER BY depth", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "parent_job", + "type_info": "Uuid" + }, + { + "ordinal": 2, + "name": "kind!: JobKind", + "type_info": { + "Custom": { + "name": "job_kind", + "kind": { + "Enum": [ + "script", + "preview", + "flow", + "dependencies", + "flowpreview", + "script_hub", + "identity", + "flowdependencies", + "http", + "graphql", + "postgresql", + "noop", + "appdependencies", + "deploymentcallback", + "singlestepflow", + "flowscript", + "flownode", + "appscript", + "aiagent", + "unassigned_script", + "unassigned_flow", + "unassigned_singlestepflow" + ] + } + } + } + }, + { + "ordinal": 3, + "name": "runnable_path", + "type_info": "Varchar" + }, + { + "ordinal": 4, + "name": "runnable_id", + "type_info": "Int8" + }, + { + "ordinal": 5, + "name": "trigger_kind", + "type_info": "Text" + }, + { + "ordinal": 6, + "name": "trigger", + "type_info": "Varchar" + }, + { + "ordinal": 7, + "name": "tag!", + "type_info": "Varchar" + }, + { + "ordinal": 8, + "name": "permissioned_as!", + "type_info": "Varchar" + }, + { + "ordinal": 9, + "name": "origin_verified!", + "type_info": "Bool" + }, + { + "ordinal": 10, + "name": "current_version!", + "type_info": "Bool" + }, + { + "ordinal": 11, + "name": "restarted!", + "type_info": "Bool" + }, + { + "ordinal": 12, + "name": "restarted_from", + "type_info": "Text" + }, + { + "ordinal": 13, + "name": "app_stamped!", + "type_info": "Bool" + }, + { + "ordinal": 14, + "name": "args_modules!", + "type_info": "Bool" + }, + { + "ordinal": 15, + "name": "worker_group", + "type_info": "Varchar" + }, + { + "ordinal": 16, + "name": "run_as_type", + "type_info": "Text" + }, + { + "ordinal": 17, + "name": "digest_source", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null + ] + }, + "hash": "76440a1c5b40a897078817cc960247270c9044b363b71496ba36e1357c63a2a4" +} diff --git a/backend/.sqlx/query-ee92100fcd681515ba81f23970e64af7f33464879bef4d1eff682126dffeef64.json b/backend/.sqlx/query-ee92100fcd681515ba81f23970e64af7f33464879bef4d1eff682126dffeef64.json deleted file mode 100644 index 3c52f8ba96..0000000000 --- a/backend/.sqlx/query-ee92100fcd681515ba81f23970e64af7f33464879bef4d1eff682126dffeef64.json +++ /dev/null @@ -1,125 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "WITH RECURSIVE lineage AS (\n SELECT id, parent_job, kind, runnable_path, runnable_id, trigger_kind, trigger,\n permissioned_as, args, 0 AS depth\n FROM v2_job WHERE id = $1 AND workspace_id = $2\n UNION ALL\n SELECT p.id, p.parent_job, p.kind, p.runnable_path, p.runnable_id, p.trigger_kind,\n p.trigger, p.permissioned_as, p.args, l.depth + 1\n FROM v2_job p JOIN lineage l ON p.id = l.parent_job\n WHERE p.workspace_id = $2 AND l.depth < 100\n )\n SELECT id AS \"id!\", parent_job, kind AS \"kind!: JobKind\", runnable_path, runnable_id,\n trigger_kind::text AS trigger_kind, permissioned_as AS \"permissioned_as!\",\n -- A restart takes its flow version from the request, so a trusted flow path\n -- can carry another flow's code: the version must belong to that path.\n CASE kind\n WHEN 'flow' THEN EXISTS (SELECT 1 FROM flow_version fv\n WHERE fv.id = runnable_id AND fv.path = runnable_path AND fv.workspace_id = $2)\n WHEN 'script' THEN EXISTS (SELECT 1 FROM script s\n WHERE s.hash = runnable_id AND s.path = runnable_path AND s.workspace_id = $2\n AND NOT s.deleted)\n ELSE true\n END AS \"origin_verified!\",\n -- What `run/f` and `run/p` resolve the path to now; for scripts the predicate of\n -- `get_latest_deployed_script_hash`.\n COALESCE(CASE kind\n WHEN 'flow' THEN runnable_id = (SELECT f.versions[array_upper(f.versions, 1)]\n FROM flow f WHERE f.path = runnable_path AND f.workspace_id = $2)\n WHEN 'script' THEN runnable_id = (SELECT s.hash FROM script s\n WHERE s.path = runnable_path AND s.workspace_id = $2 AND NOT s.deleted\n AND s.lock IS NOT NULL AND s.lock_error_logs IS NULL\n ORDER BY s.created_at DESC LIMIT 1)\n END, false) AS \"current_version!\",\n EXISTS (SELECT 1 FROM v2_job_status st WHERE st.id = lineage.id\n AND jsonb_typeof(st.flow_status->'restarted_from') = 'object')\n OR EXISTS (SELECT 1 FROM v2_job_completed c WHERE c.id = lineage.id\n AND jsonb_typeof(c.flow_status->'restarted_from') = 'object')\n AS \"restarted!\",\n COALESCE(\n (SELECT st.flow_status->'restarted_from'->>'flow_job_id' FROM v2_job_status st\n WHERE st.id = lineage.id),\n (SELECT c.flow_status->'restarted_from'->>'flow_job_id' FROM v2_job_completed c\n WHERE c.id = lineage.id)\n ) AS restarted_from,\n -- Only deployed-app runs are stamped with their app; an app editor preview\n -- runs app code at an app path it does not have to own.\n COALESCE(trigger_kind = 'app' AND starts_with(runnable_path, trigger || '/'), false)\n AS \"app_stamped!\",\n -- A preview's modules come from its args, which its parent may have taken from\n -- the caller.\n COALESCE(jsonb_typeof(args->'_MODULES') = 'object', false) AS \"args_modules!\"\n FROM lineage ORDER BY depth", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "parent_job", - "type_info": "Uuid" - }, - { - "ordinal": 2, - "name": "kind!: JobKind", - "type_info": { - "Custom": { - "name": "job_kind", - "kind": { - "Enum": [ - "script", - "preview", - "flow", - "dependencies", - "flowpreview", - "script_hub", - "identity", - "flowdependencies", - "http", - "graphql", - "postgresql", - "noop", - "appdependencies", - "deploymentcallback", - "singlestepflow", - "flowscript", - "flownode", - "appscript", - "aiagent", - "unassigned_script", - "unassigned_flow", - "unassigned_singlestepflow" - ] - } - } - } - }, - { - "ordinal": 3, - "name": "runnable_path", - "type_info": "Varchar" - }, - { - "ordinal": 4, - "name": "runnable_id", - "type_info": "Int8" - }, - { - "ordinal": 5, - "name": "trigger_kind", - "type_info": "Text" - }, - { - "ordinal": 6, - "name": "permissioned_as!", - "type_info": "Varchar" - }, - { - "ordinal": 7, - "name": "origin_verified!", - "type_info": "Bool" - }, - { - "ordinal": 8, - "name": "current_version!", - "type_info": "Bool" - }, - { - "ordinal": 9, - "name": "restarted!", - "type_info": "Bool" - }, - { - "ordinal": 10, - "name": "restarted_from", - "type_info": "Text" - }, - { - "ordinal": 11, - "name": "app_stamped!", - "type_info": "Bool" - }, - { - "ordinal": 12, - "name": "args_modules!", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [ - null, - null, - null, - null, - null, - null, - null, - null, - null, - null, - null, - null, - null - ] - }, - "hash": "ee92100fcd681515ba81f23970e64af7f33464879bef4d1eff682126dffeef64" -} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 7f660b6b51..10f5a71687 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -15711,6 +15711,7 @@ dependencies = [ "rsa", "rustls 0.23.35", "rustls-native-certs 0.8.4", + "ryu-js", "schemars 0.8.22", "semver 1.0.28", "serde", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index a17197ee9f..263cb62907 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -458,7 +458,7 @@ tower-cookies = "^0.11" serde = "^1" # 1.0.151 introduced RawValue::from_string_unchecked, which the SQL executors use # to avoid re-parsing every collected row. -serde_json = { version = "^1.0.151", features = ["preserve_order", "raw_value"] } +serde_json = { version = "^1.0.151", features = ["preserve_order", "raw_value", "float_roundtrip"] } serde_yml = "0.0.12" uuid = { version = "^1", features = ["serde", "v4", "js"] } thiserror = "^2" @@ -571,6 +571,7 @@ base64 = "^0.22.1" base32 = "^0" hmac = "0.12.1" sha2 = "0.10.6" +ryu-js = "1.0.3" md-5 = "0.10.6" sha1 = "0.10.6" sqlx = { version = "0.8.0", features = [ diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index c2c253e8f0..0d0bb78414 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -ff4d04f17721d84fb9a1f655af50ff7e4e7dcbc7 +71b8c1042fd8188c2d2882476f12b671cb5ba421 diff --git a/backend/tests/job_provenance.rs b/backend/tests/job_provenance.rs index 42ccf829e6..4eda4e5e17 100644 --- a/backend/tests/job_provenance.rs +++ b/backend/tests/job_provenance.rs @@ -150,3 +150,104 @@ async fn flow_steps_are_latest_only_under_a_current_unrestarted_flow(db: Pool) { + use serde_json::{json, Value}; + use windmill_common::job_provenance::RunAsType; + + let v: Value = + serde_json::from_str(include_str!("../../cli/test/fixtures/item_digest_vectors.json")) + .unwrap(); + let (script, flow) = (&v["script"], &v["flow"]); + + sqlx::query( + "INSERT INTO script (workspace_id, hash, path, content, lock, modules, language, kind, created_by, schema, summary, description) + VALUES ('test-workspace', 777, $1, $2, $3, $4, 'bun', 'script', 'test-user', '{}', '', '')", + ) + .bind(script["path"].as_str()) + .bind(script["content"].as_str()) + .bind(script["lock"].as_str()) + .bind(&script["modules"]) + .execute(&db) + .await + .unwrap(); + + // Version 889 holds step `a` by reference, which no checkout can reproduce. + let step = &flow["value"]["modules"][0]["value"]; + let mut by_ref = flow["value"].clone(); + by_ref["modules"][0]["value"] = json!({ + "type": "flowscript", "id": 999, "language": step["language"], + "input_transforms": step["input_transforms"], + }); + sqlx::query( + "INSERT INTO flow (workspace_id, path, summary, description, value, edited_by, versions) + VALUES ('test-workspace', $1, '', '', $2, 'test-user', ARRAY[888::bigint, 889::bigint])", + ) + .bind(flow["path"].as_str()) + .bind(&flow["value"]) + .execute(&db) + .await + .unwrap(); + for (id, value) in [(888, &flow["value"]), (889, &by_ref)] { + sqlx::query( + "INSERT INTO flow_version (id, workspace_id, path, value, schema, created_by) + VALUES ($1, 'test-workspace', $2, $3, '{}', 'test-user')", + ) + .bind(id as i64) + .bind(flow["path"].as_str()) + .bind(value) + .execute(&db) + .await + .unwrap(); + } + sqlx::query( + "INSERT INTO flow_node (id, workspace_id, path, code, lock, hash_v2) + VALUES (999, 'test-workspace', $1, $2, $3, 'h')", + ) + .bind(flow["path"].as_str()) + .bind(step["content"].as_str()) + .bind(step["lock"].as_str()) + .execute(&db) + .await + .unwrap(); + + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, kind, runnable_path, runnable_id, parent_job, trigger, trigger_kind, tag, created_by, permissioned_as, permissioned_as_email) + VALUES + ('3bb0c0de-0000-4000-8000-000000000201', 'test-workspace', 'script', $1, 777, NULL, NULL, NULL, 'bun', 'test-user', 'g/all', 'group-all@windmill.dev'), + ('3bb0c0de-0000-4000-8000-000000000202', 'test-workspace', 'flow', $2, 888, NULL, 'f/digest/nightly', 'schedule', 'flow', 'test-user', 'u/test-user', 'test@windmill.dev'), + ('3bb0c0de-0000-4000-8000-000000000203', 'test-workspace', 'flowscript', $2 || '/a', 999, '3bb0c0de-0000-4000-8000-000000000202', NULL, NULL, 'gpu', 'test-user', 'u/test-user', 'test@windmill.dev'), + ('3bb0c0de-0000-4000-8000-000000000204', 'test-workspace', 'flow', $2, 889, NULL, NULL, NULL, 'flow', 'test-user', 'u/test-user', 'test@windmill.dev')", + ) + .bind(script["path"].as_str()) + .bind(flow["path"].as_str()) + .execute(&db) + .await + .unwrap(); + for q in [ + "INSERT INTO worker_ping (worker, worker_instance, worker_group) VALUES ('wk-gpu', 'wk', 'gpu-group')", + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, worker, tag) + VALUES ('3bb0c0de-0000-4000-8000-000000000203', 'test-workspace', now(), true, 'wk-gpu', 'gpu')", + ] { + sqlx::query(q).execute(&db).await.unwrap(); + } + + let s = provenance(&db, "3bb0c0de-0000-4000-8000-000000000201").await; + assert_eq!(s.digest.as_deref(), script["digest"].as_str()); + assert_eq!(s.root_digest, s.digest); + assert_eq!(s.run_as_type, RunAsType::Group); + assert_eq!(s.worker_group, None); + + let step = provenance(&db, "3bb0c0de-0000-4000-8000-000000000203").await; + assert_eq!(step.digest.as_deref(), flow["steps"]["a"].as_str()); + assert_eq!(step.root_digest.as_deref(), flow["digest"].as_str()); + assert_eq!(step.root_trigger.as_deref(), Some("f/digest/nightly")); + assert_eq!(step.tag, "gpu"); + assert_eq!(step.worker_group.as_deref(), Some("gpu-group")); + assert_eq!(step.run_as_type, RunAsType::User); + + assert_eq!(provenance(&db, "3bb0c0de-0000-4000-8000-000000000204").await.digest, None); +} diff --git a/backend/windmill-common/Cargo.toml b/backend/windmill-common/Cargo.toml index 96fa7ceb63..1a8027c149 100644 --- a/backend/windmill-common/Cargo.toml +++ b/backend/windmill-common/Cargo.toml @@ -33,6 +33,7 @@ path = "src/lib.rs" tar.workspace = true hmac.workspace = true sha2.workspace = true +ryu-js.workspace = true sha1.workspace = true thiserror.workspace = true anyhow.workspace = true diff --git a/backend/windmill-common/src/item_digest.rs b/backend/windmill-common/src/item_digest.rs new file mode 100644 index 0000000000..23110bcae7 --- /dev/null +++ b/backend/windmill-common/src/item_digest.rs @@ -0,0 +1,132 @@ +//! Content digests of deployed scripts and flows, reproducible with `wmill digest` +//! (`cli/src/commands/digest/digest.ts`) from a `wmill sync` checkout pulled after the +//! deployment's dependency jobs: they write the locks, and store loops and branches with +//! their defaults filled in. The two implementations must hash the same bytes: change one +//! only together with the other. Numbers hash as the doubles they parse to, which needs +//! serde_json's `float_roundtrip` to parse them as exactly as JavaScript does. +//! +//! A digest is the lowercase hex SHA-256 of the canonical JSON (RFC 8785) of the item, with +//! every object member whose value is null removed: +//! - a script: `{"content", "lock", "modules", "codebase"}`, where `modules` is the +//! multi-file script's `{: {"content", "lock"}}` and is left out when empty. +//! The language is not covered: a checkout only records it in the file extension; +//! - a flow step's inline script: `{"content", "lock"}`; +//! - a flow: its `value`. A flow whose value references code in `flow_node` has none, see +//! [`references_flow_nodes`]. + +use serde_json::Value; +use sha2::{Digest, Sha256}; + +use crate::flows::{Branch, FlowModule, FlowModuleValue, FlowValue, ToolValue}; + +pub fn digest(value: &Value) -> String { + let mut out = String::new(); + write_canonical(value, &mut out); + format!("{:x}", Sha256::digest(out.as_bytes())) +} + +pub fn canonical_json(value: &Value) -> String { + let mut out = String::new(); + write_canonical(value, &mut out); + out +} + +fn write_canonical(value: &Value, out: &mut String) { + match value { + Value::Null => out.push_str("null"), + Value::Bool(b) => out.push_str(if *b { "true" } else { "false" }), + // RFC 8785 serializes every number as an IEEE double, the way ECMAScript prints it. + Value::Number(n) => { + let f = n.as_f64().unwrap_or(0.0); + out.push_str(ryu_js::Buffer::new().format(f)) + } + Value::String(s) => out.push_str(&serde_json::to_string(s).unwrap_or_default()), + Value::Array(items) => { + out.push('['); + for (i, item) in items.iter().enumerate() { + if i > 0 { + out.push(','); + } + write_canonical(item, out); + } + out.push(']'); + } + Value::Object(members) => { + // RFC 8785 orders members by the UTF-16 code units of their names. + let mut members: Vec<_> = members.iter().filter(|(_, v)| !v.is_null()).collect(); + members.sort_by(|(a, _), (b, _)| a.encode_utf16().cmp(b.encode_utf16())); + out.push('{'); + for (i, (k, v)) in members.into_iter().enumerate() { + if i > 0 { + out.push(','); + } + out.push_str(&serde_json::to_string(k).unwrap_or_default()); + out.push(':'); + write_canonical(v, out); + } + out.push('}'); + } + } +} + +/// Whether a flow value runs code stored outside it, in `flow_node` rows: an inline script +/// step `{"type": "flowscript", "id"}`, or a loop's or branch's steps in a `*_node`. Its digest +/// would not cover that code, and a checkout cannot reproduce it either, since the export +/// carries the references. A value or step that does not parse runs nothing, so it counts as +/// none. +pub fn references_flow_nodes(value: &Value) -> bool { + let Ok(flow) = serde_json::from_value::(value.clone()) else { + return false; + }; + flow.modules + .iter() + .chain(flow.preprocessor_module.as_deref()) + .chain(flow.failure_module.as_deref()) + .any(module_references_flow_nodes) +} + +fn module_references_flow_nodes(module: &FlowModule) -> bool { + module + .get_value() + .is_ok_and(|v| module_value_references_flow_nodes(&v)) +} + +fn module_value_references_flow_nodes(value: &FlowModuleValue) -> bool { + let any = |modules: &[FlowModule]| modules.iter().any(module_references_flow_nodes); + let branches = |branches: &[Branch]| { + branches + .iter() + .any(|b| b.modules_node.is_some() || any(&b.modules)) + }; + match value { + FlowModuleValue::FlowScript { .. } => true, + FlowModuleValue::ForloopFlow { modules, modules_node, .. } + | FlowModuleValue::WhileloopFlow { modules, modules_node, .. } => { + modules_node.is_some() || any(modules) + } + FlowModuleValue::BranchOne { branches: b, default, default_node } => { + default_node.is_some() || any(default) || branches(b) + } + FlowModuleValue::BranchAll { branches: b, .. } => branches(b), + FlowModuleValue::AIAgent { tools, .. } => tools.iter().any(|tool| match &tool.value { + ToolValue::FlowModule(v) => module_value_references_flow_nodes(v), + _ => false, + }), + _ => false, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn canonical_form_follows_rfc_8785_without_nulls() { + let v = json!({"b": [1.0, 1e21, 0.1, -0.0, null], "a": null, "é": "\u{1f}\"\n", "\u{e000}": 1, "\u{1f600}": 2}); + assert_eq!( + canonical_json(&v), + "{\"b\":[1,1e+21,0.1,0,null],\"é\":\"\\u001f\\\"\\n\",\"\u{1f600}\":2,\"\u{e000}\":1}" + ); + } +} diff --git a/backend/windmill-common/src/job_provenance.rs b/backend/windmill-common/src/job_provenance.rs index 63075b0c37..96db3b8a87 100644 --- a/backend/windmill-common/src/job_provenance.rs +++ b/backend/windmill-common/src/job_provenance.rs @@ -1,6 +1,7 @@ +use serde_json::Value; use uuid::Uuid; -use crate::{db::DB, error::Result, jobs::JobKind, scripts::ScriptHash}; +use crate::{db::DB, error::Result, item_digest, jobs::JobKind, scripts::ScriptHash}; /// Where a job's code comes from, as far as its chain of parents can prove it. pub struct JobProvenance { @@ -28,7 +29,39 @@ pub struct JobProvenance { pub root_path: Option, pub root_kind: JobKind, pub root_trigger_kind: Option, + /// What triggered the root job, as recorded in `v2_job.trigger`. + pub root_trigger: Option, pub root_version: Option, + /// See [`item_digest`]. Only scripts, flows and flow steps' inline scripts have one. + pub digest: Option, + pub root_digest: Option, + /// The job's worker tag. + pub tag: String, + /// The group of the worker running the job, if it is running. + pub worker_group: Option, + pub run_as_type: RunAsType, +} + +/// What `permissioned_as` names. +#[derive(Debug, PartialEq, Eq, Clone, Copy)] +pub enum RunAsType { + User, + ServiceAccount, + Group, + /// No member of the workspace: a superadmin outside it, a built-in identity + /// (`superadmin_secret@windmill.dev`, ...) or a user since removed. + NonMember, +} + +impl RunAsType { + pub fn as_str(&self) -> &'static str { + match self { + RunAsType::User => "user", + RunAsType::ServiceAccount => "service_account", + RunAsType::Group => "group", + RunAsType::NonMember => "non_member", + } + } } #[derive(Debug, PartialEq, Eq)] @@ -58,6 +91,8 @@ struct LineageJob { runnable_path: Option, runnable_id: Option, trigger_kind: Option, + trigger: Option, + tag: String, permissioned_as: String, origin_verified: bool, current_version: bool, @@ -66,6 +101,10 @@ struct LineageJob { restarted_from: Option, app_stamped: bool, args_modules: bool, + worker_group: Option, + run_as_type: Option, + /// What the digest hashes, for the job and the root only. + digest_source: Option, } /// Reads any job of `w_id` regardless of the caller: authorize access to `job_id` first. @@ -74,16 +113,17 @@ pub async fn job_provenance(db: &DB, job_id: &Uuid, w_id: &str) -> Result