From 9d7001ddc0b02274c5bb34d00a62d859df0179c0 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 16 Sep 2026 00:18:30 +0200 Subject: [PATCH] fix: key a cached inline task on a fingerprint of its code and its arguments --- backend/tests/wac_task_cache_key.rs | 144 ++++++++++++++++++ backend/windmill-common/src/wac.rs | 10 ++ backend/windmill-worker/src/bun_executor.rs | 6 +- backend/windmill-worker/src/common.rs | 58 +++++-- backend/windmill-worker/src/wac_executor.rs | 4 + cli/src/guidance/skills.gen.ts | 30 ++-- python-client/wmill/tests/test_workflow.py | 16 ++ python-client/wmill/wmill/client.py | 31 +++- system_prompts/auto-generated/prompts.ts | 30 ++-- system_prompts/auto-generated/script.md | 10 +- system_prompts/auto-generated/sdks/python.md | 10 +- .../auto-generated/sdks/wac-python.md | 10 +- .../auto-generated/sdks/wac-typescript.md | 10 +- .../skills/write-script-python3/SKILL.md | 10 +- .../skills/write-workflow-as-code/SKILL.md | 20 +-- typescript-client/client.ts | 29 +++- .../tests/workflow_task_identity.test.ts | 93 +++++++++++ 17 files changed, 412 insertions(+), 109 deletions(-) create mode 100644 backend/tests/wac_task_cache_key.rs create mode 100644 typescript-client/tests/workflow_task_identity.test.ts diff --git a/backend/tests/wac_task_cache_key.rs b/backend/tests/wac_task_cache_key.rs new file mode 100644 index 0000000000..0cf04816d2 --- /dev/null +++ b/backend/tests/wac_task_cache_key.rs @@ -0,0 +1,144 @@ +//! A workflow-as-code task child runs its parent's code with its parent's +//! arguments, so its result-cache key comes from what the dispatch seeded in its +//! checkpoint: the SDK's fingerprint of the task and the arguments the task was +//! called with, or, from an SDK that sends no fingerprint, its step key and the +//! parent's arguments. + +use serde_json::{json, value::RawValue, Value}; +use sqlx::{Pool, Postgres}; +use std::collections::HashMap; +use uuid::Uuid; +use windmill_common::client::AuthedClient; +use windmill_common::jobs::JobKind; +use windmill_common::scripts::{ScriptHash, ScriptLang}; +use windmill_queue::MiniPulledJob; +use windmill_worker::common::cached_result_path; + +const W_ID: &str = "test-workspace"; + +async fn insert_job(db: &Pool, id: Uuid, parent: Option) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, created_by, created_at, permissioned_as, \ + permissioned_as_email, kind, script_lang, runnable_path, tag, visible_to_owner, parent_job) \ + VALUES ($1, $2, 'test-user', now(), 'u/test-user', 'test@windmill.dev', \ + 'flowscript', 'bun', 'f/system/wac/a', 'bun', true, $3)", + ) + .bind(id) + .bind(W_ID) + .bind(parent) + .execute(db) + .await?; + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag) \ + VALUES ($1, $2, now(), true, 'bun')", + ) + .bind(id) + .bind(W_ID) + .execute(db) + .await?; + Ok(()) +} + +/// The cache path of a task child carrying `parent_args`, whose checkpoint was +/// seeded with step key `key`, fingerprint `fn_id` and call arguments `{n}`. +async fn child_cache_path( + db: &Pool, + parent_args: Value, + key: &str, + fn_id: Option<&str>, + n: i64, +) -> anyhow::Result { + let parent = Uuid::new_v4(); + let child = Uuid::new_v4(); + insert_job(db, parent, None).await?; + insert_job(db, child, Some(parent)).await?; + let mut checkpoint = json!({ + "completed_steps": {}, + "_executing_key": key, + "_executing_args": { "n": n }, + }); + if let Some(fn_id) = fn_id { + checkpoint["_executing_fn"] = json!(fn_id); + } + sqlx::query("INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, $2)") + .bind(child) + .bind(json!({ "_checkpoint": checkpoint })) + .execute(db) + .await?; + + let args: HashMap> = serde_json::from_str(&parent_args.to_string())?; + let mut job = MiniPulledJob::new_inline( + W_ID.to_string(), + Some(args), + "test-user".to_string(), + "u/test-user".to_string(), + "test@windmill.dev".to_string(), + Some("f/system/wac/a".to_string()), + JobKind::FlowScript, + Some(ScriptHash(42)), + "bun".to_string(), + Some(ScriptLang::Bun), + ); + job.id = child; + job.parent_job = Some(parent); + job.cache_ttl = Some(60); + let client = AuthedClient::new( + "http://localhost".to_string(), + W_ID.to_string(), + "tok".to_string(), + None, + ); + Ok(cached_result_path(db, &client, &job, None).await?) +} + +#[sqlx::test(fixtures("base"))] +async fn a_task_child_is_cached_under_its_fingerprint_or_else_its_step_key( + db: Pool, +) -> anyhow::Result<()> { + let parent_args = || json!({ "x": 1 }); + + let fingerprinted = child_cache_path(&db, parent_args(), "step", Some("f1"), 1).await?; + assert_ne!( + fingerprinted, + child_cache_path(&db, parent_args(), "step", Some("f2"), 1).await?, + "two tasks called at one position: the fingerprint tells them apart" + ); + assert_eq!( + fingerprinted, + child_cache_path(&db, parent_args(), "step_2", Some("f1"), 1).await?, + "one task called at two positions keeps one identity" + ); + assert_eq!( + fingerprinted, + child_cache_path(&db, json!({ "x": 2 }), "step", Some("f1"), 1).await?, + "with a fingerprint, the parent's arguments are not in the key" + ); + assert_ne!( + fingerprinted, + child_cache_path(&db, parent_args(), "step", Some("f1"), 2).await?, + "and the task's own arguments are" + ); + + let unfingerprinted = child_cache_path(&db, parent_args(), "step", None, 1).await?; + assert_ne!( + unfingerprinted, + child_cache_path(&db, parent_args(), "step_2", None, 1).await?, + "without a fingerprint, the step key is the identity" + ); + assert_ne!( + unfingerprinted, + child_cache_path(&db, json!({ "x": 2 }), "step", None, 1).await?, + "and the parent's arguments are in the key" + ); + assert_eq!( + unfingerprinted, + child_cache_path(&db, parent_args(), "step", None, 2).await?, + "and the task's own arguments are not" + ); + assert_ne!( + child_cache_path(&db, parent_args(), "f1", None, 1).await?, + fingerprinted, + "a step key never reads as a fingerprint" + ); + Ok(()) +} diff --git a/backend/windmill-common/src/wac.rs b/backend/windmill-common/src/wac.rs index 05c343b427..5e06e650e7 100644 --- a/backend/windmill-common/src/wac.rs +++ b/backend/windmill-common/src/wac.rs @@ -40,6 +40,16 @@ pub struct WacCheckpoint { #[serde(skip_serializing_if = "Option::is_none")] #[serde(default)] pub _executing_key: Option, + /// With `_executing_key`: the SDK's fingerprint of the task's code and the + /// arguments the task was called with, what its cached result is keyed on. + /// A child whose SDK sent no fingerprint is keyed on its step key and its + /// parent's arguments instead (`cached_result_path`). + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(default)] + pub _executing_fn: Option, + #[serde(skip_serializing_if = "Option::is_none")] + #[serde(default)] + pub _executing_args: Option>, /// `resume_job.id` values already consumed by earlier approval steps (the /// row primary key, not the distinct integer `resume_id` column). Rows are /// never deleted, so a workflow with several sequential wait_for_approval() diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 117fbda086..4ea15d7b03 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -3234,10 +3234,14 @@ pub async fn handle_wac_v2_output( // _executing_key to know which step to run). External // scripts/flows don't need a WAC checkpoint. if !is_external { - let child_checkpoint_json = serde_json::json!({ + let mut child_checkpoint_json = serde_json::json!({ "completed_steps": &checkpoint.completed_steps, "_executing_key": &step.key, + "_executing_args": &step.args, }); + if let Some(fn_id) = &step.fn_id { + child_checkpoint_json["_executing_fn"] = serde_json::json!(fn_id); + } sqlx::query( "INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, jsonb_build_object('_checkpoint', $2::jsonb)) diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 9c012fd117..9d01cf8b07 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -1574,18 +1574,28 @@ pub async fn cached_result_path( _ => {} } } - // A workflow-as-code task child runs its parent's code with the parent's - // arguments; the step it executes is what tells its result from the parent's - // and from its siblings'. - if let Some(step_key) = wac_executing_key(db, job).await? { - hasher.update(b"wac_step:"); - hasher.update(step_key.as_bytes()); - } + // A task child runs its parent's code with its parent's arguments. With the SDK's + // fingerprint it is keyed on that and its own call arguments. Without one it is keyed + // on its step key and the parent's arguments: a step key names a position, not a + // task, and only the parent's arguments tell apart tasks a branch puts at one. + let args = match wac_task_identity(db, job).await? { + Some(WacTaskIdentity::Fingerprint { fn_id, args }) => { + hasher.update(b"wac_fn:"); + hasher.update(fn_id.as_bytes()); + Some(Json(args)) + } + Some(WacTaskIdentity::StepKey(key)) => { + hasher.update(b"wac_step:"); + hasher.update(key.as_bytes()); + job.args.clone() + } + None => job.args.clone(), + }; hash_args( db, client, &job.workspace_id, - &job.args, + &args, &mut hasher, &job.id, job.cache_ignore_s3_path.unwrap_or(false), @@ -1594,23 +1604,41 @@ pub async fn cached_result_path( Ok(format!("g/results/{:064x}", hasher.finalize())) } -/// The checkpoint step key a workflow-as-code parent seeded for this child at push -/// time; `None` for any job that is not such a child. -async fn wac_executing_key( +/// What a workflow-as-code parent seeded in a task child's checkpoint at push time +/// to key the child's cached result on. +enum WacTaskIdentity { + Fingerprint { fn_id: String, args: HashMap> }, + StepKey(String), +} + +/// `None` for any job that is not a workflow-as-code task child. +async fn wac_task_identity( db: &DB, job: &MiniPulledJob, -) -> windmill_common::error::Result> { +) -> windmill_common::error::Result> { if job.parent_job.is_none() || job.flow_step_id.is_some() { return Ok(None); } - let key: Option> = sqlx::query_scalar( - "SELECT workflow_as_code_status->'_checkpoint'->>'_executing_key' \ + let identity: Option<( + Option, + Option, + Option>>>, + )> = sqlx::query_as( + "SELECT workflow_as_code_status->'_checkpoint'->>'_executing_fn', \ + workflow_as_code_status->'_checkpoint'->>'_executing_key', \ + workflow_as_code_status->'_checkpoint'->'_executing_args' \ FROM v2_job_status WHERE id = $1", ) .bind(job.id) .fetch_optional(db) .await?; - Ok(key.flatten()) + Ok(match identity { + Some((Some(fn_id), _, Some(Json(args)))) => { + Some(WacTaskIdentity::Fingerprint { fn_id, args }) + } + Some((_, Some(key), _)) => Some(WacTaskIdentity::StepKey(key)), + _ => None, + }) } #[cfg(feature = "parquet")] diff --git a/backend/windmill-worker/src/wac_executor.rs b/backend/windmill-worker/src/wac_executor.rs index c188792a5e..ab86c0e2fa 100644 --- a/backend/windmill-worker/src/wac_executor.rs +++ b/backend/windmill-worker/src/wac_executor.rs @@ -85,6 +85,10 @@ pub struct WacStepDispatch { pub concurrency_key: Option, #[serde(default)] pub concurrency_time_window_s: Option, + /// The SDK's fingerprint of the task's code, the identity its cached result + /// is keyed on; absent from an SDK that predates it. + #[serde(default)] + pub fn_id: Option, } fn default_dispatch_type() -> String { diff --git a/cli/src/guidance/skills.gen.ts b/cli/src/guidance/skills.gen.ts index 767f6a23e0..1dc4bea8c1 100644 --- a/cli/src/guidance/skills.gen.ts +++ b/cli/src/guidance/skills.gen.ts @@ -4609,12 +4609,10 @@ def parse_sql_client_name(name: str) -> tuple[str, Optional[str]] # no \`\`delay\`\` all go out in a single round. # # \`\`cache_ttl\`\` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A \`\`task_script\`\` target is keyed on the -# arguments it is called with. It has no effect on a \`\`task_flow\`\` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# \`\`task_flow\`\` target, which keeps its flow's own cache policy. # # Usage:: # @@ -6750,12 +6748,10 @@ export interface TaskOptions { timeout?: number; tag?: string; /** Seconds during which a previous result of this task is served instead of - * running it again. A task written inline in the workflow is keyed on its - * step key (its name and call order) and the workflow's input, not on the - * arguments it is called with, so cache one only when whether it runs, and - * what it receives, follow from the workflow's input alone. A \`taskScript\` - * target is keyed on the arguments it is called with. It has no effect on a - * \`taskFlow\` target, which keeps its flow's own cache policy. */ + * running it again. The result is keyed on the task and the arguments it is + * called with, so anything a cached task reads from its closure must be + * passed in as an argument. It has no effect on a \`taskFlow\` target, which + * keeps its flow's own cache policy. */ cache_ttl?: number; priority?: number; concurrency_limit?: number; @@ -6948,12 +6944,10 @@ def get_resume_urls(approver: str = None, flow_level: bool = None) -> dict # no \`\`delay\`\` all go out in a single round. # # \`\`cache_ttl\`\` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A \`\`task_script\`\` target is keyed on the -# arguments it is called with. It has no effect on a \`\`task_flow\`\` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# \`\`task_flow\`\` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/python-client/wmill/tests/test_workflow.py b/python-client/wmill/tests/test_workflow.py index ceed39bfeb..3502cb9d47 100644 --- a/python-client/wmill/tests/test_workflow.py +++ b/python-client/wmill/tests/test_workflow.py @@ -1847,3 +1847,19 @@ class TestApprovalKeys: await wait_for_approval() assert _run_workflow(wf, {"completed_steps": {"approval": {}}}, {})["key"] == "approval_2" + + +class TestTaskFingerprint: + """The worker keys a task child's cached result on the ``fn_id`` its dispatch + carries, so it has to follow the task, not the position it is called at.""" + + def test_one_task_keeps_its_fingerprint_and_another_gets_its_own(self): + first = _run_workflow(double_parallel_wf, {}, {}) + assert [s["key"] for s in first["steps"]] == ["double", "double_2"] + assert first["steps"][0]["fn_id"] == first["steps"][1]["fn_id"] + + second = _run_workflow( + double_parallel_wf, {"completed_steps": {"double": 2, "double_2": 4}}, {} + ) + assert [s["key"] for s in second["steps"]] == ["add_one", "add_one_2"] + assert second["steps"][0]["fn_id"] != first["steps"][0]["fn_id"] diff --git a/python-client/wmill/wmill/client.py b/python-client/wmill/wmill/client.py index dbc9d9b193..a5ee1bb151 100644 --- a/python-client/wmill/wmill/client.py +++ b/python-client/wmill/wmill/client.py @@ -3026,7 +3026,7 @@ class WorkflowCtx: print(f"\n--- WAC: {key} ---") info = {"name": name or key, "script": script or key, "args": kwargs, "key": key, "dispatch_type": dispatch_type} if _task_options: - for opt_key in ("timeout", "tag", "cache_ttl", "priority", "concurrent_limit", "concurrency_key", "concurrency_time_window_s"): + for opt_key in ("timeout", "tag", "cache_ttl", "priority", "concurrent_limit", "concurrency_key", "concurrency_time_window_s", "fn_id"): if opt_key in _task_options and _task_options[opt_key] is not None: info[opt_key] = _task_options[opt_key] self._pending.append(info) @@ -3285,6 +3285,22 @@ class WorkflowCtx: }) +def _fn_fingerprint(func) -> str: + """A stable identity for a task's code, what its cached result is keyed on: a + name is shared by any two tasks called the same, and a step key by any two + tasks called at the same position, so neither can tell them apart.""" + import hashlib + import inspect + import marshal + + try: + src = inspect.getsource(func).encode() + except (OSError, TypeError): + # No source on disk: the whole code object, constants and names included. + src = marshal.dumps(func.__code__) + return hashlib.sha256(src).hexdigest() + + def task( _func=None, *, @@ -3328,12 +3344,10 @@ def task( no ``delay`` all go out in a single round. ``cache_ttl`` serves a previous result of the task for that many seconds - instead of running it again. A task is keyed on its step key (its name and - call order) and the workflow's input, not on the arguments it is called - with, so cache one only when whether it runs, and what it receives, follow - from the workflow's input alone. A ``task_script`` target is keyed on the - arguments it is called with. It has no effect on a ``task_flow`` target, - which keeps its flow's own cache policy. + instead of running it again. The result is keyed on the task and the + arguments it is called with, so anything a cached task reads from its + closure must be passed in as an argument. It has no effect on a + ``task_flow`` target, which keeps its flow's own cache policy. Usage:: @@ -3364,6 +3378,7 @@ def task( def decorator(func) -> Callable[..., Any]: task_path = path task_name = func.__name__ + _fn_opts = {**(_task_opts or {}), "fn_id": _fn_fingerprint(func)} _params_list = list(_sig(func).parameters) @@ -3389,7 +3404,7 @@ def task( if ctx is not None: script = task_path if task_path else task_name merged = _merge_args(args, kwargs) - return ctx._next_step(task_name, script, func, _task_options=_task_opts, **merged) + return ctx._next_step(task_name, script, func, _task_options=_fn_opts, **merged) # WAC v1: running inside a Windmill job but not in a @workflow if ( diff --git a/system_prompts/auto-generated/prompts.ts b/system_prompts/auto-generated/prompts.ts index a59d4b37fa..e2cb1c93cb 100644 --- a/system_prompts/auto-generated/prompts.ts +++ b/system_prompts/auto-generated/prompts.ts @@ -2588,12 +2588,10 @@ def parse_sql_client_name(name: str) -> tuple[str, Optional[str]] # no \`\`delay\`\` all go out in a single round. # # \`\`cache_ttl\`\` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A \`\`task_script\`\` target is keyed on the -# arguments it is called with. It has no effect on a \`\`task_flow\`\` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# \`\`task_flow\`\` target, which keeps its flow's own cache policy. # # Usage:: # @@ -2746,12 +2744,10 @@ export interface TaskOptions { timeout?: number; tag?: string; /** Seconds during which a previous result of this task is served instead of - * running it again. A task written inline in the workflow is keyed on its - * step key (its name and call order) and the workflow's input, not on the - * arguments it is called with, so cache one only when whether it runs, and - * what it receives, follow from the workflow's input alone. A \`taskScript\` - * target is keyed on the arguments it is called with. It has no effect on a - * \`taskFlow\` target, which keeps its flow's own cache policy. */ + * running it again. The result is keyed on the task and the arguments it is + * called with, so anything a cached task reads from its closure must be + * passed in as an argument. It has no effect on a \`taskFlow\` target, which + * keeps its flow's own cache policy. */ cache_ttl?: number; priority?: number; concurrency_limit?: number; @@ -2944,12 +2940,10 @@ def get_resume_urls(approver: str = None, flow_level: bool = None) -> dict # no \`\`delay\`\` all go out in a single round. # # \`\`cache_ttl\`\` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A \`\`task_script\`\` target is keyed on the -# arguments it is called with. It has no effect on a \`\`task_flow\`\` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# \`\`task_flow\`\` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/system_prompts/auto-generated/script.md b/system_prompts/auto-generated/script.md index f278a67fba..e4171a7150 100644 --- a/system_prompts/auto-generated/script.md +++ b/system_prompts/auto-generated/script.md @@ -2729,12 +2729,10 @@ def parse_sql_client_name(name: str) -> tuple[str, Optional[str]] # no ``delay`` all go out in a single round. # # ``cache_ttl`` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A ``task_script`` target is keyed on the -# arguments it is called with. It has no effect on a ``task_flow`` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# ``task_flow`` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/system_prompts/auto-generated/sdks/python.md b/system_prompts/auto-generated/sdks/python.md index 4f727c7d45..c3b34b69db 100644 --- a/system_prompts/auto-generated/sdks/python.md +++ b/system_prompts/auto-generated/sdks/python.md @@ -673,12 +673,10 @@ def parse_sql_client_name(name: str) -> tuple[str, Optional[str]] # no ``delay`` all go out in a single round. # # ``cache_ttl`` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A ``task_script`` target is keyed on the -# arguments it is called with. It has no effect on a ``task_flow`` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# ``task_flow`` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/system_prompts/auto-generated/sdks/wac-python.md b/system_prompts/auto-generated/sdks/wac-python.md index a8e98a9d4f..2e54c5bd06 100644 --- a/system_prompts/auto-generated/sdks/wac-python.md +++ b/system_prompts/auto-generated/sdks/wac-python.md @@ -59,12 +59,10 @@ def get_resume_urls(approver: str = None, flow_level: bool = None) -> dict # no ``delay`` all go out in a single round. # # ``cache_ttl`` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A ``task_script`` target is keyed on the -# arguments it is called with. It has no effect on a ``task_flow`` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# ``task_flow`` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/system_prompts/auto-generated/sdks/wac-typescript.md b/system_prompts/auto-generated/sdks/wac-typescript.md index 608d75d18c..cc59e74954 100644 --- a/system_prompts/auto-generated/sdks/wac-typescript.md +++ b/system_prompts/auto-generated/sdks/wac-typescript.md @@ -35,12 +35,10 @@ export interface TaskOptions { timeout?: number; tag?: string; /** Seconds during which a previous result of this task is served instead of - * running it again. A task written inline in the workflow is keyed on its - * step key (its name and call order) and the workflow's input, not on the - * arguments it is called with, so cache one only when whether it runs, and - * what it receives, follow from the workflow's input alone. A `taskScript` - * target is keyed on the arguments it is called with. It has no effect on a - * `taskFlow` target, which keeps its flow's own cache policy. */ + * running it again. The result is keyed on the task and the arguments it is + * called with, so anything a cached task reads from its closure must be + * passed in as an argument. It has no effect on a `taskFlow` target, which + * keeps its flow's own cache policy. */ cache_ttl?: number; priority?: number; concurrency_limit?: number; diff --git a/system_prompts/auto-generated/skills/write-script-python3/SKILL.md b/system_prompts/auto-generated/skills/write-script-python3/SKILL.md index 1ccb376a9f..f6a2bbe70e 100644 --- a/system_prompts/auto-generated/skills/write-script-python3/SKILL.md +++ b/system_prompts/auto-generated/skills/write-script-python3/SKILL.md @@ -858,12 +858,10 @@ def parse_sql_client_name(name: str) -> tuple[str, Optional[str]] # no ``delay`` all go out in a single round. # # ``cache_ttl`` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A ``task_script`` target is keyed on the -# arguments it is called with. It has no effect on a ``task_flow`` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# ``task_flow`` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/system_prompts/auto-generated/skills/write-workflow-as-code/SKILL.md b/system_prompts/auto-generated/skills/write-workflow-as-code/SKILL.md index 8e5ab1f06c..f30acc0e47 100644 --- a/system_prompts/auto-generated/skills/write-workflow-as-code/SKILL.md +++ b/system_prompts/auto-generated/skills/write-workflow-as-code/SKILL.md @@ -278,12 +278,10 @@ export interface TaskOptions { timeout?: number; tag?: string; /** Seconds during which a previous result of this task is served instead of - * running it again. A task written inline in the workflow is keyed on its - * step key (its name and call order) and the workflow's input, not on the - * arguments it is called with, so cache one only when whether it runs, and - * what it receives, follow from the workflow's input alone. A `taskScript` - * target is keyed on the arguments it is called with. It has no effect on a - * `taskFlow` target, which keeps its flow's own cache policy. */ + * running it again. The result is keyed on the task and the arguments it is + * called with, so anything a cached task reads from its closure must be + * passed in as an argument. It has no effect on a `taskFlow` target, which + * keeps its flow's own cache policy. */ cache_ttl?: number; priority?: number; concurrency_limit?: number; @@ -476,12 +474,10 @@ def get_resume_urls(approver: str = None, flow_level: bool = None) -> dict # no ``delay`` all go out in a single round. # # ``cache_ttl`` serves a previous result of the task for that many seconds -# instead of running it again. A task is keyed on its step key (its name and -# call order) and the workflow's input, not on the arguments it is called -# with, so cache one only when whether it runs, and what it receives, follow -# from the workflow's input alone. A ``task_script`` target is keyed on the -# arguments it is called with. It has no effect on a ``task_flow`` target, -# which keeps its flow's own cache policy. +# instead of running it again. The result is keyed on the task and the +# arguments it is called with, so anything a cached task reads from its +# closure must be passed in as an argument. It has no effect on a +# ``task_flow`` target, which keeps its flow's own cache policy. # # Usage:: # diff --git a/typescript-client/client.ts b/typescript-client/client.ts index dd2123e746..faf82e8e42 100644 --- a/typescript-client/client.ts +++ b/typescript-client/client.ts @@ -1713,12 +1713,10 @@ export interface TaskOptions { timeout?: number; tag?: string; /** Seconds during which a previous result of this task is served instead of - * running it again. A task written inline in the workflow is keyed on its - * step key (its name and call order) and the workflow's input, not on the - * arguments it is called with, so cache one only when whether it runs, and - * what it receives, follow from the workflow's input alone. A `taskScript` - * target is keyed on the arguments it is called with. It has no effect on a - * `taskFlow` target, which keeps its flow's own cache policy. */ + * running it again. The result is keyed on the task and the arguments it is + * called with, so anything a cached task reads from its closure must be + * passed in as an argument. It has no effect on a `taskFlow` target, which + * keeps its flow's own cache policy. */ cache_ttl?: number; priority?: number; concurrency_limit?: number; @@ -1856,6 +1854,7 @@ export class WorkflowCtx { args: Record = {}, dispatch_type: string = "inline", options?: TaskOptions, + fnId?: string, ): PromiseLike { this._rethrowSwallowed(); const stepName = name || script || "step"; @@ -1909,6 +1908,7 @@ export class WorkflowCtx { } const stepInfo: any = { name: name || key, script: script || key, args, key, dispatch_type }; + if (fnId) stepInfo.fn_id = fnId; if (options) { if (options.timeout !== undefined) stepInfo.timeout = options.timeout; if (options.tag !== undefined) stepInfo.tag = options.tag; @@ -2285,6 +2285,20 @@ export async function step( return jsonRoundTrip(await fn()); } +// A stable identity for a task's code, what its cached result is keyed on: a +// name is shared by any two tasks called the same, and a step key by any two +// tasks called at the same position, so neither can tell them apart. The name +// stays in: every bound function stringifies to the same native stub, and only +// its name (`bound orders`) tells two of them apart. +function fnFingerprint(fn: Function): string { + const src = JSON.stringify([fn.name, fn.toString()]); + let h = 0xcbf29ce484222325n; + for (let i = 0; i < src.length; i++) { + h = ((h ^ BigInt(src.charCodeAt(i))) * 0x100000001b3n) & 0xffffffffffffffffn; + } + return h.toString(16); +} + /** * Wrap an async function as a workflow task. * @@ -2322,6 +2336,7 @@ export function task Promise>( assertUsableRetry(taskOptions?.retry); const taskName = fn.name || taskPath || ""; + const fnId = fnFingerprint(fn); // NOT async — in workflow context we return the thenable directly so that // unawaited task calls leave the step in ctx.pending (for _flushPending). @@ -2342,7 +2357,7 @@ export function task Promise>( kwargs[`arg${i}`] = args[i]; } } - const stepResult = ctx._nextStep(taskName, script, kwargs, "inline", taskOptions); + const stepResult = ctx._nextStep(taskName, script, kwargs, "inline", taskOptions, fnId); // If this step should execute directly (child job mode), run the inner function // and throw StepSuspend with mode "step_complete" to signal that we're done if ((stepResult as any)?._execute_directly) { diff --git a/typescript-client/tests/workflow_task_identity.test.ts b/typescript-client/tests/workflow_task_identity.test.ts new file mode 100644 index 0000000000..5f4f644eba --- /dev/null +++ b/typescript-client/tests/workflow_task_identity.test.ts @@ -0,0 +1,93 @@ +/** + * The fingerprint a dispatched task carries, against the real client. + * + * Run with: bun test typescript-client/tests/workflow_task_identity.test.ts + * + * The worker keys a task child's cached result on this `fn_id`, so two tasks + * that share a step key (anonymous tasks in exclusive branches) or whose source + * reads the same (bound functions) must not share one. Imports client.ts itself, + * with the two generated modules stubbed so the import works without ./build.sh. + */ +import { expect, test, describe, mock, beforeAll, afterAll } from "bun:test"; + +mock.module("../services.gen", () => ({ + ResourceService: {}, + VariableService: {}, + JobService: {}, + HelpersService: {}, + AppService: {}, + MetricsService: {}, + OidcService: {}, + UserService: {}, + KafkaTriggerService: {}, +})); +mock.module("../core/OpenAPI", () => ({ + OpenAPI: { BASE: "http://localhost:8000/api", TOKEN: "tok" }, +})); + +const { WorkflowCtx, task, setWorkflowCtx, StepSuspend } = await import("../client.ts"); +import { isSuspendSignal } from "../wacError"; + +/** The steps a first round of `body` dispatches. */ +async function dispatched(body: () => Promise): Promise { + const ctx = new WorkflowCtx({ completed_steps: {} } as any); + setWorkflowCtx(ctx); + try { + await body(); + } catch (e: any) { + if (isSuspendSignal(e, StepSuspend)) return e.dispatchInfo.steps; + throw e; + } finally { + setWorkflowCtx(null); + } + throw new Error("the body completed without dispatching"); +} + +// These assert the suspend a worker acts on, the legacy inline path; the v2 fast +// path is on by default, so pin it off rather than depend on `WM_JOB_ID` being +// absent. +const priorFastPath = process.env.WM_WAC_INLINE_FAST_PATH; +beforeAll(() => { + process.env.WM_WAC_INLINE_FAST_PATH = "0"; +}); +afterAll(() => { + if (priorFastPath === undefined) delete process.env.WM_WAC_INLINE_FAST_PATH; + else process.env.WM_WAC_INLINE_FAST_PATH = priorFastPath; +}); + +describe("task fingerprint", () => { + test("two anonymous tasks dispatched at the same position carry different fingerprints", async () => { + const tasks = { + a: task(async () => "A", { cache_ttl: 60 }), + b: task(async () => "B", { cache_ttl: 60 }), + }; + const [a] = await dispatched(() => tasks.a() as Promise); + const [b] = await dispatched(() => tasks.b() as Promise); + expect(a.key).toBe(b.key); + expect(typeof a.fn_id).toBe("string"); + expect(a.fn_id).not.toBe(b.fn_id); + }); + + test("one task keeps its fingerprint across positions", async () => { + const double = task(async (n: number) => n * 2, { cache_ttl: 60 }); + const steps = await dispatched(async () => { + await Promise.all([double(1), double(2)]); + }); + expect(steps.map((s) => s.key)).toEqual(["step", "step_2"]); + expect(steps[0].fn_id).toBe(steps[1].fn_id); + }); + + test("two bound functions carry different fingerprints", async () => { + const api = { + async orders() { + return "orders"; + }, + async users() { + return "users"; + }, + }; + const [orders] = await dispatched(() => task(api.orders.bind(api))() as Promise); + const [users] = await dispatched(() => task(api.users.bind(api))() as Promise); + expect(orders.fn_id).not.toBe(users.fn_id); + }); +});