From deaa9c92339f450133ccc188e99250e09a9afcd8 Mon Sep 17 00:00:00 2001 From: Durable Workflow Date: Sat, 3 Oct 2026 06:47:56 +0000 Subject: [PATCH] Use ordinary task identity for Rust workflow streams --- CHANGELOG.md | 7 ++ Cargo.toml | 4 +- scripts/ci/test-publish-rust-sdk.py | 2 +- src/lib.rs | 67 ++++++++++++++++++- .../workflow-stream-task-identity.json | 41 ++++++++++++ tests/replay_regression_corpus.rs | 23 ++++++- 6 files changed, 139 insertions(+), 5 deletions(-) create mode 100644 tests/fixtures/replay-regressions/workflow-stream-task-identity.json diff --git a/CHANGELOG.md b/CHANGELOG.md index 363dc77..3a09945 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,12 @@ # Changelog +## 2.1.5 + +- Allow the ordinary Worker to append and close workflow streams when Server + supplies a durable task ID without a separate workflow command ID. Prefer an + explicit command ID when available. Recorded stream effects still replay + without producing duplicate output. + ## 2.1.4 - Keep fresh installations compatible with Rust 1.86 by selecting uuid 1.26.1. diff --git a/Cargo.toml b/Cargo.toml index f6d60b6..f9e633c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "durable-workflow" -version = "2.1.4" +version = "2.1.5" edition = "2021" description = "Rust client and worker SDK for Durable Workflow Cloud and self-hosted Server" license = "MIT" @@ -14,7 +14,7 @@ categories = ["api-bindings", "asynchronous"] include = ["/src/**", "/schema/**", "/examples/**", "/Cargo.toml", "/README.md", "/CHANGELOG.md", "/LICENSE"] [package.metadata.durable-workflow] -product-train = "2.1.4" +product-train = "2.1.5" compatibility-authority = "protocol-manifests" supported-server-versions = "2.0.0" qualified-server-version = "2.4.0" diff --git a/scripts/ci/test-publish-rust-sdk.py b/scripts/ci/test-publish-rust-sdk.py index 78e675a..8f0f832 100644 --- a/scripts/ci/test-publish-rust-sdk.py +++ b/scripts/ci/test-publish-rust-sdk.py @@ -18,7 +18,7 @@ PUBLISH = ROOT / "scripts" / "ci" / "publish-rust-sdk.sh" RELEASE_WORKFLOW = ROOT / ".github" / "workflows" / "release.yml" RELEASE_TOOLING_INSTALLER = ROOT / "scripts" / "ci" / "install-release-tooling.sh" -PACKAGE_VERSION = "2.1.4" +PACKAGE_VERSION = "2.1.5" PRODUCT_TRAIN = PACKAGE_VERSION SERVER_VERSIONS = "2.0.0" QUALIFIED_SERVER_VERSION = "2.4.0" diff --git a/src/lib.rs b/src/lib.rs index 30f8aeb..09450f2 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -213,7 +213,7 @@ pub enum Error { #[error("workflow future yielded without emitting a durable command")] WorkflowYieldedWithoutCommand, #[error( - "workflow_stream_command_identity_missing: workflow stream authoring requires a non-empty server-provided workflow_command_id" + "workflow_stream_command_identity_missing: workflow stream authoring requires a non-empty server-provided workflow_command_id or task_id" )] MissingWorkflowCommandIdentity, #[error("workflow state lock is poisoned")] @@ -7622,6 +7622,7 @@ impl Worker { .workflow_command_id .clone() .filter(|identity| !identity.is_empty()) + .or_else(|| (!task.task_id.is_empty()).then(|| task.task_id.clone())) .unwrap_or_default(); let mut workflow_state = WorkflowState::new_with_identity( task.history_events, @@ -17270,6 +17271,70 @@ mod tests { assert!(context.take_commands().expect("commands").is_empty()); } + #[test] + fn worker_stream_authoring_uses_task_identity_and_replays_without_output() { + let client = Client::builder("http://localhost:8080").build().unwrap(); + let mut worker = Worker::new(client, "rust-workers"); + worker.register_workflow("streams.worker", |ctx, _| async move { + ctx.append_workflow_stream( + "output", + &[WorkflowStreamAppendItem::new(json!("hello"))?], + None, + )?; + ctx.close_workflow_stream("output", None)?; + Ok(json!("done")) + }); + for command_id in [None, Some(String::new()), Some("command-42".to_string())] { + let mut task = workflow_task("streams.worker", Vec::new(), DEFAULT_CODEC); + task.workflow_command_id = command_id.clone(); + let expected_identity = command_id + .as_deref() + .filter(|id| !id.is_empty()) + .unwrap_or(&task.task_id) + .to_string(); + let commands = worker + .execute_workflow_task(task) + .expect("ordinary worker stream output"); + assert_eq!(commands.len(), 3); + assert_eq!( + commands[0]["workflow_stream"]["command_identity"], + expected_identity + ); + assert_eq!( + commands[0]["workflow_stream"]["items"][0]["idempotency_key"], + format!("dw-stream:{expected_identity}:0:0") + ); + assert_eq!( + commands[1]["workflow_stream"]["command_identity"], + expected_identity + ); + assert_eq!(commands[1]["workflow_stream"]["operation"], "close"); + let history = commands[..2] + .iter() + .enumerate() + .map(|(index, command)| { + history_event( + "SideEffectRecorded", + json!({"sequence":index + 1,"result":command["result"]}), + ) + }) + .collect(); + let mut replay = workflow_task("streams.worker", history, DEFAULT_CODEC); + replay.task_id = "replacement-task".to_string(); + let replayed = worker + .execute_workflow_task(replay) + .expect("replacement worker consumes recorded stream effects"); + assert_eq!(replayed.len(), 1); + assert_eq!(replayed[0]["type"], "complete_workflow"); + } + let mut missing = workflow_task("streams.worker", Vec::new(), DEFAULT_CODEC); + missing.task_id.clear(); + assert!(matches!( + worker.execute_workflow_task(missing), + Err(Error::MissingWorkflowCommandIdentity) + )); + } + #[test] fn cold_worker_replay_does_not_repeat_committed_side_effects_or_markers() { fn worker(calls: Arc) -> Worker { diff --git a/tests/fixtures/replay-regressions/workflow-stream-task-identity.json b/tests/fixtures/replay-regressions/workflow-stream-task-identity.json new file mode 100644 index 0000000..85339ed --- /dev/null +++ b/tests/fixtures/replay-regressions/workflow-stream-task-identity.json @@ -0,0 +1,41 @@ +{ + "$schema": "https://raw.githubusercontent.com/durable-workflow/.github/main/regression-corpus/evidence-schema.json", + "fixture_schema": "durable-workflow.replay-regression/v1", + "id": "rust-workflow-stream-task-identity", + "protocol_version": "1.19", + "bindings": ["rust"], + "workflow": {"type": "corpus.workflow-stream", "input": [], "payload_codec": "avro"}, + "worker_task": {"payload_codec": "avro", "task_id": "01JTASK0000000000000000000", "workflow_command_id": null}, + "command_sequence": [ + { + "type": "record_side_effect", + "workflow_stream": { + "operation": "append", "stream_name": "tokens", "command_identity": "01JTASK0000000000000000000", "command_ordinal": 0, + "items": [ + {"payload": {"token": "hello"}, "payload_codec": "avro", "idempotency_key": "dw-stream:01JTASK0000000000000000000:0:0"}, + {"payload_reference": "s3://payloads/token-2", "idempotency_key": "dw-stream:01JTASK0000000000000000000:0:1"} + ] + }, + "result": null + }, + {"type": "record_side_effect", "workflow_stream": {"operation": "close", "stream_name": "tokens", "command_identity": "01JTASK0000000000000000000", "command_ordinal": 1}, "result": null}, + {"type": "complete_workflow", "result": "done"} + ], + "expected": { + "command_sequence": [ + { + "type": "record_side_effect", + "workflow_stream": { + "operation": "append", "stream_name": "tokens", "command_identity": "01JTASK0000000000000000000", "command_ordinal": 0, + "items": [ + {"payload": {"token": "hello"}, "payload_codec": "avro", "idempotency_key": "dw-stream:01JTASK0000000000000000000:0:0"}, + {"payload_reference": "s3://payloads/token-2", "idempotency_key": "dw-stream:01JTASK0000000000000000000:0:1"} + ] + }, + "result": null + }, + {"type": "record_side_effect", "workflow_stream": {"operation": "close", "stream_name": "tokens", "command_identity": "01JTASK0000000000000000000", "command_ordinal": 1}, "result": null}, + {"type": "complete_workflow", "result": "done"} + ] + } +} diff --git a/tests/replay_regression_corpus.rs b/tests/replay_regression_corpus.rs index 8a7a17d..8a0de8d 100644 --- a/tests/replay_regression_corpus.rs +++ b/tests/replay_regression_corpus.rs @@ -439,7 +439,10 @@ async fn execute_fixture_delivery(fixture: &Value, delivery_id: &str) -> Result< return Err(format!("{fixture_id}.history must be an array")); } - let task_id = format!("regression-corpus-{fixture_id}-{delivery_id}"); + let task_id = fixture["worker_task"]["task_id"] + .as_str() + .map(str::to_string) + .unwrap_or_else(|| format!("regression-corpus-{fixture_id}-{delivery_id}")); let task_payload_codec = match fixture.get("worker_task") { Some(worker_task) => worker_task.get("payload_codec").cloned(), None => Some(json!(payload_codec)), @@ -988,6 +991,7 @@ async fn workflow_stream_command_without_durable_identity_fails_closed() { .as_object_mut() .expect("workflow stream worker task") .remove("workflow_command_id"); + fixture["worker_task"]["task_id"] = json!(""); let error = execute_fixture_delivery(&fixture, "missing-identity") .await @@ -1000,6 +1004,23 @@ async fn workflow_stream_command_without_durable_identity_fails_closed() { assert!(error.contains("workflow_command_id"), "{error}"); } +#[tokio::test] +async fn workflow_stream_task_identity_is_stable_across_cold_worker_redelivery() { + let fixture_path = Path::new(env!("CARGO_MANIFEST_DIR")) + .join("tests/fixtures/replay-regressions/workflow-stream-task-identity.json"); + let fixture: Value = serde_json::from_str( + &fs::read_to_string(fixture_path).expect("read ordinary-task stream fixture"), + ) + .expect("parse ordinary-task stream fixture"); + let first = execute_fixture_delivery(&fixture, "delivery-a") + .await + .expect("first ordinary task delivery"); + let restarted = execute_fixture_delivery(&fixture, "delivery-b") + .await + .expect("replacement ordinary task delivery"); + assert_eq!(first, restarted); +} + #[tokio::test] async fn adjacent_condition_wait_occurrences_are_deterministic_across_cold_workers() { let fixture_path = Path::new(env!("CARGO_MANIFEST_DIR"))