feat: add trigger, tag, run-as and digest claims to job OIDC tokens (#11481)

* feat: add trigger, tag, run-as and digest claims to job OIDC tokens

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: parse digest floats exactly and derive the codebase digest like push

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: give flows that reference flow nodes no digest and keep --json clean

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: list flow step digests from module positions only

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* test: skip the pulled module-script digest check on windows

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* chore: update ee-repo-ref to 71b8c1042fd8188c2d2882476f12b671cb5ba421

This commit updates the EE repository reference after PR #840 was merged in windmill-ee-private.

Previous ee-repo-ref: ba1870536789a43c5b6ec18522a93edb9cb07f3d

New ee-repo-ref: 71b8c1042fd8188c2d2882476f12b671cb5ba421

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
Ruben Fiszel
2026-10-02 19:02:17 +02:00
committed by GitHub
co-authored by Claude Opus 5.5 windmill-internal-app[bot]
parent 4bc7e0d7ba
commit 4beb1f9420
20 changed files with 933 additions and 154 deletions
File diff suppressed because one or more lines are too long
@@ -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"
}
+1
View File
@@ -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",
+2 -1
View File
@@ -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 = [
+1 -1
View File
@@ -1 +1 @@
ff4d04f17721d84fb9a1f655af50ff7e4e7dcbc7
71b8c1042fd8188c2d2882476f12b671cb5ba421
+101
View File
@@ -150,3 +150,104 @@ async fn flow_steps_are_latest_only_under_a_current_unrestarted_flow(db: Pool<Po
assert_eq!(p.latest, latest, "job {job}");
}
}
/// `wmill digest` hashes a sync checkout of the same items to the same values
/// (`cli/test/item_digest_unit.test.ts`), so a trust policy can compare the claim to git.
#[sqlx::test(fixtures("base"))]
async fn digests_match_the_cli_vectors(db: Pool<Postgres>) {
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);
}
+1
View File
@@ -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
+132
View File
@@ -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 `{<relative path>: {"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::<FlowValue>(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}"
);
}
}
+91 -5
View File
@@ -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<String>,
pub root_kind: JobKind,
pub root_trigger_kind: Option<String>,
/// What triggered the root job, as recorded in `v2_job.trigger`.
pub root_trigger: Option<String>,
pub root_version: Option<String>,
/// See [`item_digest`]. Only scripts, flows and flow steps' inline scripts have one.
pub digest: Option<String>,
pub root_digest: Option<String>,
/// The job's worker tag.
pub tag: String,
/// The group of the worker running the job, if it is running.
pub worker_group: Option<String>,
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<String>,
runnable_id: Option<i64>,
trigger_kind: Option<String>,
trigger: Option<String>,
tag: String,
permissioned_as: String,
origin_verified: bool,
current_version: bool,
@@ -66,6 +101,10 @@ struct LineageJob {
restarted_from: Option<String>,
app_stamped: bool,
args_modules: bool,
worker_group: Option<String>,
run_as_type: Option<String>,
/// What the digest hashes, for the job and the root only.
digest_source: Option<Value>,
}
/// 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<Option
LineageJob,
r#"WITH RECURSIVE lineage AS (
SELECT id, parent_job, kind, runnable_path, runnable_id, trigger_kind, trigger,
permissioned_as, args, 0 AS depth
tag, permissioned_as, args, 0 AS depth
FROM v2_job WHERE id = $1 AND workspace_id = $2
UNION ALL
SELECT p.id, p.parent_job, p.kind, p.runnable_path, p.runnable_id, p.trigger_kind,
p.trigger, p.permissioned_as, p.args, l.depth + 1
p.trigger, p.tag, p.permissioned_as, p.args, l.depth + 1
FROM v2_job p JOIN lineage l ON p.id = l.parent_job
WHERE p.workspace_id = $2 AND l.depth < 100
)
SELECT id AS "id!", parent_job, kind AS "kind!: JobKind", runnable_path, runnable_id,
trigger_kind::text AS trigger_kind, permissioned_as AS "permissioned_as!",
trigger_kind::text AS trigger_kind, trigger, tag AS "tag!",
permissioned_as AS "permissioned_as!",
-- A restart takes its flow version from the request, so a trusted flow path
-- can carry another flow's code: the version must belong to that path.
CASE kind
@@ -121,7 +161,30 @@ pub async fn job_provenance(db: &DB, job_id: &Uuid, w_id: &str) -> Result<Option
AS "app_stamped!",
-- A preview's modules come from its args, which its parent may have taken from
-- the caller.
COALESCE(jsonb_typeof(args->'_MODULES') = 'object', false) AS "args_modules!"
COALESCE(jsonb_typeof(args->'_MODULES') = 'object', false) AS "args_modules!",
CASE WHEN depth = 0 THEN (SELECT wp.worker_group FROM v2_job_queue q
JOIN worker_ping wp ON wp.worker = q.worker WHERE q.id = lineage.id)
END AS worker_group,
CASE WHEN depth = 0 THEN CASE
WHEN starts_with(permissioned_as, 'g/') THEN 'group'
ELSE (SELECT CASE WHEN u.is_service_account THEN 'service_account' ELSE 'user' END
FROM usr u WHERE starts_with(permissioned_as, 'u/')
AND u.username = substr(permissioned_as, 3) AND u.workspace_id = $2)
END END AS run_as_type,
-- The canonical item `item_digest` hashes, null members included.
CASE WHEN depth = 0 OR depth = max(depth) OVER () THEN CASE kind
WHEN 'script' THEN (SELECT jsonb_build_object('content', s.content,
'lock', s.lock, 'codebase', s.codebase,
'modules', (SELECT jsonb_object_agg(m.key, m.value - 'language')
FROM jsonb_each(s.modules) m))
FROM script s WHERE s.hash = runnable_id AND s.workspace_id = $2
AND NOT s.deleted)
WHEN 'flow' THEN (SELECT fv.value FROM flow_version fv
WHERE fv.id = runnable_id AND fv.workspace_id = $2)
WHEN 'flowscript' THEN (SELECT jsonb_build_object('content', n.code,
'lock', n.lock)
FROM flow_node n WHERE n.id = runnable_id AND n.workspace_id = $2)
END END AS digest_source
FROM lineage ORDER BY depth"#,
job_id,
w_id
@@ -174,6 +237,15 @@ pub async fn job_provenance(db: &DB, job_id: &Uuid, w_id: &str) -> Result<Option
}
}
}
let job = &lineage[0];
let digest = item_digest_of(job);
let root_digest = item_digest_of(root);
let run_as_type = match job.run_as_type.as_deref() {
Some("group") => RunAsType::Group,
Some("service_account") => RunAsType::ServiceAccount,
Some("user") => RunAsType::User,
_ => RunAsType::NonMember,
};
Ok(Some(JobProvenance {
deployed,
latest: deployed && all_current,
@@ -186,7 +258,13 @@ pub async fn job_provenance(db: &DB, job_id: &Uuid, w_id: &str) -> Result<Option
root_path: root.runnable_path.clone(),
root_kind: root.kind,
root_trigger_kind: root.trigger_kind.clone(),
root_trigger: root.trigger.clone(),
root_version: version(root),
digest,
root_digest,
tag: job.tag.clone(),
worker_group: job.worker_group.clone(),
run_as_type,
}))
}
@@ -220,6 +298,14 @@ async fn restart_origin_is_current(db: &DB, origin: &Uuid, w_id: &str) -> Result
Ok(current)
}
fn item_digest_of(job: &LineageJob) -> Option<String> {
let source = job.digest_source.as_ref()?;
if job.kind == JobKind::Flow && item_digest::references_flow_nodes(source) {
return None;
}
Some(item_digest::digest(source))
}
fn version(job: &LineageJob) -> Option<String> {
match job.kind {
JobKind::Script => job.runnable_id.map(|h| ScriptHash(h).to_string()),
+1
View File
@@ -90,6 +90,7 @@ pub mod workspace_dependencies;
#[cfg(feature = "private")]
pub mod git_sync_ee;
pub mod git_sync_oss;
pub mod item_digest;
pub mod job_provenance;
pub mod jobs;
pub mod jwt;
+6
View File
@@ -70,9 +70,15 @@ pub struct JobClaim {
pub root_path: Option<String>,
pub root_job_kind: String,
pub root_trigger_kind: Option<String>,
pub root_trigger: Option<String>,
pub latest: bool,
pub version: Option<String>,
pub root_version: Option<String>,
pub digest: Option<String>,
pub root_digest: Option<String>,
pub tag: String,
pub worker_group: Option<String>,
pub run_as_type: String,
pub jti: String,
}
+219
View File
@@ -0,0 +1,219 @@
import { createHash } from "node:crypto";
import { stat } from "node:fs/promises";
import { sep as SEP } from "node:path";
import { Command } from "@cliffy/command";
import * as log from "../../core/log.ts";
import { mergeConfigWithConfigFile } from "../../core/conf.ts";
import { GlobalOptions } from "../../types.ts";
import { readLocalFlow } from "../flow/flow.ts";
import {
findContentFile,
readModulesFromDisk,
} from "../script/script.ts";
import { replaceLock } from "../../utils/metadata.ts";
import { listSyncCodebases, SyncCodebase } from "../../utils/codebase.ts";
import { findCodebase } from "../sync/sync.ts";
import {
extractFolderPath,
getModuleFolderSuffix,
isFlowPath,
isMissingDbtDescriptor,
isModuleEntryPoint,
scriptPathToRemotePath,
} from "../../utils/resource_folders.ts";
import { inferContentTypeFromFilePath } from "../../utils/script_common.ts";
import { readTextFile } from "../../utils/utils.ts";
import { yamlParseFile } from "../../utils/yaml.ts";
// The server computes the same digests for the `digest` and `root_digest` claims of job
// OIDC tokens (backend/windmill-common/src/item_digest.rs): both must hash the same bytes.
/** RFC 8785 canonical JSON, with null and undefined object members left out. */
export function canonicalJson(value: unknown): string {
if (value === null || value === undefined) return "null";
if (Array.isArray(value)) {
return "[" + value.map(canonicalJson).join(",") + "]";
}
if (typeof value === "object") {
// Sorting strings in JS compares UTF-16 code units, as RFC 8785 requires.
const members = Object.entries(value as Record<string, unknown>)
.filter(([, v]) => v !== null && v !== undefined)
.sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0));
return (
"{" +
members
.map(([k, v]) => JSON.stringify(k) + ":" + canonicalJson(v))
.join(",") +
"}"
);
}
return JSON.stringify(value);
}
export function digestOf(value: unknown): string {
return createHash("sha256").update(canonicalJson(value), "utf8").digest("hex");
}
/** The inline scripts of a flow's steps, by step id. Only module positions are walked: a
* static input shaped like a step is data. */
function inlineSteps(modules: any[] | undefined, steps: Record<string, string>) {
for (const m of modules ?? []) {
stepValue(m?.id, m?.value, steps);
}
}
function stepValue(id: unknown, v: any, steps: Record<string, string>) {
switch (v?.type) {
case "rawscript":
if (typeof id === "string") {
steps[id] = digestOf({ content: v.content, lock: v.lock });
}
break;
case "forloopflow":
case "whileloopflow":
inlineSteps(v.modules, steps);
break;
case "branchone":
inlineSteps(v.default, steps);
for (const b of v.branches ?? []) inlineSteps(b?.modules, steps);
break;
case "branchall":
for (const b of v.branches ?? []) inlineSteps(b?.modules, steps);
break;
case "aiagent":
for (const t of v.tools ?? []) {
if (t?.value?.tool_type === "flowmodule") stepValue(t.id, t.value, steps);
}
break;
}
}
export interface ItemDigest {
path: string;
kind: "script" | "flow";
digest: string;
/** For a flow, the digest of each step's inline script, by step id. */
steps?: Record<string, string>;
}
export async function flowDigest(folder: string): Promise<ItemDigest> {
const { flow, missingFiles } = await readLocalFlow(folder);
if (missingFiles.length > 0) {
throw new Error(`Missing inline script file(s): ${missingFiles.join(", ")}`);
}
const steps: Record<string, string> = {};
const value = flow.value as any;
inlineSteps(
[...(value.modules ?? []), value.failure_module, value.preprocessor_module].filter(Boolean),
steps
);
const path = folder
.replaceAll(SEP, "/")
.replace(/\/$/, "")
.replace(/(\.flow|__flow)$/, "");
return { path, kind: "flow", digest: digestOf(flow.value), steps };
}
export async function scriptDigest(
contentFile: string,
defaultTs?: "bun" | "deno",
codebases: SyncCodebase[] = []
): Promise<ItemDigest> {
const moduleEntry = isModuleEntryPoint(contentFile);
const remotePath = scriptPathToRemotePath(contentFile);
const base = remotePath.replaceAll("/", SEP);
const language = inferContentTypeFromFilePath(contentFile, defaultTs);
const metaBase = moduleEntry ? base + getModuleFolderSuffix() + SEP + "script" : base + ".script";
let meta: any = undefined;
for (const ext of [".yaml", ".json"]) {
if (await stat(metaBase + ext).catch(() => undefined)) {
meta = await yamlParseFile(metaBase + ext);
break;
}
}
replaceLock(meta);
let content = "";
try {
content = await readTextFile(contentFile);
} catch (e) {
if (!isMissingDbtDescriptor(contentFile, e)) throw e;
}
// What `wmill sync push` sends: the codebase is derived from wmill.yaml, never read
// from the metadata file. Only bundling tells whether push adds `.tar`, so a pulled
// value that is this digest plus `.tar` is taken as is.
const codebase = language == "bun" ? findCodebase(contentFile, codebases) : undefined;
let codebaseDigest = codebase ? await codebase.getDigest() : undefined;
if (codebaseDigest && meta?.codebase === codebaseDigest + ".tar") {
codebaseDigest = meta.codebase;
}
const isDbt = language === "dbt";
const modules = await readModulesFromDisk(
base + getModuleFolderSuffix(language),
defaultTs,
moduleEntry,
isDbt
);
const digest = digestOf({
content,
lock: meta?.lock,
codebase: codebaseDigest,
modules: modules
? Object.fromEntries(
Object.entries(modules).map(([p, m]) => [p, { content: m.content, lock: m.lock }])
)
: undefined,
});
return { path: remotePath, kind: "script", digest };
}
/** The item at a local path of a sync checkout: a flow folder, a script file, or a script's
* metadata file. Paths resolve against the current directory, which must be the
* checkout's root, as for `wmill sync push`. */
export async function itemDigest(
localPath: string,
defaultTs?: "bun" | "deno",
codebases: SyncCodebase[] = []
): Promise<ItemDigest> {
const folderPath = localPath.endsWith(SEP) ? localPath : localPath + SEP;
if (isFlowPath(folderPath)) {
return flowDigest(extractFolderPath(folderPath, "flow")!);
}
if (/\.script\.(yaml|json)$/.test(localPath) || /script\.(yaml|json)$/.test(localPath)) {
return scriptDigest(await findContentFile(localPath), defaultTs, codebases);
}
return scriptDigest(localPath, defaultTs, codebases);
}
async function digest(
opts: GlobalOptions & { json?: boolean },
...paths: string[]
) {
if (opts.json) log.setSilent(true);
const merged = await mergeConfigWithConfigFile(opts);
const codebases = listSyncCodebases(merged);
const results: ItemDigest[] = [];
for (const p of paths) {
results.push(await itemDigest(p, merged.defaultTs, codebases));
}
if (opts.json) {
console.log(JSON.stringify(results, null, 2));
return;
}
for (const r of results) {
console.log(`${r.digest} ${r.path}`);
for (const [id, d] of Object.entries(r.steps ?? {})) {
console.log(`${d} ${r.path}/${id}`);
}
}
}
const command = new Command()
.description(
"Compute the content digest of local scripts and flows, as the `digest` and `root_digest` claims of job OIDC tokens carry it. Run from the root of a sync checkout, pulled after the deployment's dependency jobs completed: they write the lockfiles the digest covers. A flow also lists the digest of each inline step, as `<flow path>/<step id>`."
)
.arguments("<paths...:string>")
.option("--json", "Output the digests as JSON")
.action(digest as any);
export default command;
+30 -22
View File
@@ -170,28 +170,10 @@ function collectStepPaths(flowValue: any): string[] {
return paths;
}
export async function pushFlow(
workspace: string,
remotePath: string,
localPath: string,
message?: string,
permissionedAsContext?: PermissionedAsContext
): Promise<void> {
if (alreadySynced.includes(localPath)) {
return;
}
alreadySynced.push(localPath);
remotePath = remotePath.replaceAll(SEP, "/");
let flow: Flow | undefined = undefined;
try {
flow = await wmill.getFlowByPath({
workspace: workspace,
path: remotePath,
});
} catch {
// flow doesn't exist
}
/** The flow in the folder `localPath`, with its inline scripts read back in. */
export async function readLocalFlow(
localPath: string
): Promise<{ flow: FlowFile; missingFiles: string[] }> {
if (!localPath.endsWith(SEP)) {
localPath += SEP;
}
@@ -214,6 +196,32 @@ export async function pushFlow(
if (localFlow.value.preprocessor_module) {
await replaceInlineScripts([localFlow.value.preprocessor_module], fileReader, log, localPath, SEP, undefined, missingFiles);
}
return { flow: localFlow, missingFiles };
}
export async function pushFlow(
workspace: string,
remotePath: string,
localPath: string,
message?: string,
permissionedAsContext?: PermissionedAsContext
): Promise<void> {
if (alreadySynced.includes(localPath)) {
return;
}
alreadySynced.push(localPath);
remotePath = remotePath.replaceAll(SEP, "/");
let flow: Flow | undefined = undefined;
try {
flow = await wmill.getFlowByPath({
workspace: workspace,
path: remotePath,
});
} catch {
// flow doesn't exist
}
const { flow: localFlow, missingFiles } = await readLocalFlow(localPath);
if (missingFiles.length > 0) {
// Hard-fail rather than push the literal `!inline path` text as
// rawscript.content. That string would be persisted in flow_version.value
+9
View File
@@ -7952,6 +7952,15 @@ Watch local file changes and live-reload the dev page for preview. Does NOT depl
- \`--path <path:string>\` - Watch a specific windmill path (e.g., u/admin/my_script or f/my_flow)
- \`--no-open\` - Do not open the browser automatically
### digest
Compute the content digest of local scripts and flows, as the \`digest\` and \`root_digest\` claims of job OIDC tokens carry it. Run from the root of a sync checkout, pulled after the deployment's dependency jobs completed: they write the lockfiles the digest covers. A flow also lists the digest of each inline step, as \`<flow path>/<step id>\`.
**Arguments:** \`<paths...:string>\`
**Options:**
- \`--json\` - Output the digests as JSON
### docs
Search Windmill documentation.
+2
View File
@@ -25,6 +25,7 @@ import protectionRules from "./commands/protection-rules/protection-rules.ts";
import instance from "./commands/instance/instance.ts";
import workerGroups from "./commands/worker-groups/worker-groups.ts";
import lint from "./commands/lint/lint.ts";
import digest from "./commands/digest/digest.ts";
import dev from "./commands/dev/dev.ts";
import { GlobalOptions } from "./types.ts";
@@ -203,6 +204,7 @@ const command = new Command()
.command("dev", dev)
.command("sync", sync)
.command("lint", lint)
.command("digest", digest)
.command("gitsync-settings", gitsyncSettings)
.command("protection-rules", protectionRules)
.command("instance", instance)
+67
View File
@@ -0,0 +1,67 @@
{
"script": {
"path": "f/digest/script",
"language": "bun",
"content": "import { helper } from \"./helper.ts\";\r\n\r\nexport async function main(name = \"wörld 🌍\") {\n return helper(name);\n}\n",
"lock": "{\n \"dependencies\": {}\n}\n//bun.lock\n<empty>",
"modules": {
"helper.ts": {
"content": "export const helper = (s: string) => `hi ${s}\\t`;\n",
"language": "bun"
}
},
"digest": "fcf85b60c3d9aef073e6911f417fdd19bedeb6d9e0f36b2589df07900997f29d"
},
"flow": {
"path": "f/digest/flow",
"value": {
"modules": [
{
"id": "a",
"value": {
"type": "rawscript",
"content": "export async function main(x: number) {\n return x * 1.5;\n}\n",
"language": "bun",
"lock": "{\n \"dependencies\": {}\n}\n//bun.lock\n<empty>",
"input_transforms": {
"x": { "type": "static", "value": 1.0 },
"y": { "type": "static", "value": [92.46585493017267, 15624.632993405901, 1e-7] },
"z": { "type": "static", "value": [{ "type": "flowscript", "id": 1, "modules_node": 2 }, { "id": "a", "value": { "type": "rawscript", "content": "data" } }] }
}
},
"summary": "Multiply",
"retry": { "constant": { "attempts": 2, "seconds": 5 } }
},
{
"id": "b",
"value": {
"type": "forloopflow",
"iterator": { "type": "javascript", "expr": "[1, 2, 3]" },
"skip_failures": false,
"parallel": true,
"modules": [
{
"id": "c",
"value": {
"type": "rawscript",
"content": "def main(x: int):\n return {\"é\": x}\n",
"language": "python3",
"input_transforms": {
"x": { "type": "javascript", "expr": "flow_input.iter.value" }
},
"tag": null
}
}
]
}
}
],
"same_worker": false
},
"digest": "0d8ae0dd8eb7272264ee802481a5abfe67501a759268fb2aa7e8277f0a2ab44c",
"steps": {
"a": "abf4646174c8d6984ea9f6d9c6746b61c3d95cd0cb91a4ca89b884b0f61600ae",
"c": "ee70e41e8647bab22bc4f7ac943c76355e2157015e71538a8f3adbb492c76a45"
}
}
}
+88
View File
@@ -0,0 +1,88 @@
import { afterAll, beforeAll, expect, test } from "bun:test";
import JSZip from "jszip";
import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { dirname, join } from "node:path";
import {
readDirRecursiveWithIgnore,
ZipFSElement,
} from "../src/commands/sync/sync.ts";
import { itemDigest } from "../src/commands/digest/digest.ts";
import vectors from "./fixtures/item_digest_vectors.json";
// The server hashes the same vectors in backend/tests/job_provenance.rs, so a checkout
// `wmill sync pull` writes for them must digest to the values the job OIDC claims carry.
const originalCwd = process.cwd();
beforeAll(() => {
process.chdir(mkdtempSync(join(tmpdir(), "wmill-digest-")));
});
afterAll(() => {
process.chdir(originalCwd);
});
/** The workspace export of the vectors, as `GET .../workspaces/tarball` writes it. */
function exportZip(): JSZip {
const zip = new JSZip();
const file = (path: string, content: string) =>
zip.file(path, content, { createFolders: false });
const { script, flow } = vectors;
file(`${script.path}.ts`, script.content);
file(
`${script.path}.script.json`,
JSON.stringify({
summary: "",
description: "",
schema: {},
lock: script.lock,
kind: "script",
modules: script.modules,
})
);
file(
`${flow.path}.flow.json`,
JSON.stringify({ summary: "", description: "", value: flow.value, schema: {} })
);
return zip;
}
async function writeTree(el: any) {
for await (const e of readDirRecursiveWithIgnore(() => false, el)) {
if (!e.isDirectory) {
mkdirSync(dirname(e.path), { recursive: true });
writeFileSync(e.path, await e.getContentText());
}
}
}
beforeAll(async () => {
await writeTree(ZipFSElement(exportZip(), true, "bun", {}, {}, false, true));
});
// On Windows the pull leaves a multi-file script's entry point outside its `__mod/` folder
// (`ZipFSElement` matches module bases against `/`-separated zip names), a layout neither
// push nor this command reads back.
test.skipIf(process.platform === "win32")(
"a pulled script digests to the server's digest",
async () => {
const script = await itemDigest("f/digest/script__mod/script.ts", "bun");
expect(script).toEqual({
path: vectors.script.path,
kind: "script",
digest: vectors.script.digest,
});
expect(
(await itemDigest("f/digest/script__mod/script.yaml", "bun")).digest
).toBe(vectors.script.digest);
}
);
test("a pulled flow digests to the server's digests", async () => {
const flow = await itemDigest(join("f", "digest", "flow.flow"), "bun");
expect(flow).toEqual({
path: vectors.flow.path,
kind: "flow",
digest: vectors.flow.digest,
steps: vectors.flow.steps,
});
});
@@ -124,6 +124,15 @@ Watch local file changes and live-reload the dev page for preview. Does NOT depl
- `--path <path:string>` - Watch a specific windmill path (e.g., u/admin/my_script or f/my_flow)
- `--no-open` - Do not open the browser automatically
### digest
Compute the content digest of local scripts and flows, as the `digest` and `root_digest` claims of job OIDC tokens carry it. Run from the root of a sync checkout, pulled after the deployment's dependency jobs completed: they write the lockfiles the digest covers. A flow also lists the digest of each inline step, as `<flow path>/<step id>`.
**Arguments:** `<paths...:string>`
**Options:**
- `--json` - Output the digests as JSON
### docs
Search Windmill documentation.
+9
View File
@@ -3470,6 +3470,15 @@ Watch local file changes and live-reload the dev page for preview. Does NOT depl
- \`--path <path:string>\` - Watch a specific windmill path (e.g., u/admin/my_script or f/my_flow)
- \`--no-open\` - Do not open the browser automatically
### digest
Compute the content digest of local scripts and flows, as the \`digest\` and \`root_digest\` claims of job OIDC tokens carry it. Run from the root of a sync checkout, pulled after the deployment's dependency jobs completed: they write the lockfiles the digest covers. A flow also lists the digest of each inline step, as \`<flow path>/<step id>\`.
**Arguments:** \`<paths...:string>\`
**Options:**
- \`--json\` - Output the digests as JSON
### docs
Search Windmill documentation.
@@ -129,6 +129,15 @@ Watch local file changes and live-reload the dev page for preview. Does NOT depl
- `--path <path:string>` - Watch a specific windmill path (e.g., u/admin/my_script or f/my_flow)
- `--no-open` - Do not open the browser automatically
### digest
Compute the content digest of local scripts and flows, as the `digest` and `root_digest` claims of job OIDC tokens carry it. Run from the root of a sync checkout, pulled after the deployment's dependency jobs completed: they write the lockfiles the digest covers. A flow also lists the digest of each inline step, as `<flow path>/<step id>`.
**Arguments:** `<paths...:string>`
**Options:**
- `--json` - Output the digests as JSON
### docs
Search Windmill documentation.