Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand All @@ -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"
Expand Down
2 changes: 1 addition & 1 deletion scripts/ci/test-publish-rust-sdk.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
67 changes: 66 additions & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")]
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<AtomicUsize>) -> Worker {
Expand Down
Original file line number Diff line number Diff line change
@@ -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"}
]
}
}
23 changes: 22 additions & 1 deletion tests/replay_regression_corpus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)),
Expand Down Expand Up @@ -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
Expand All @@ -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"))
Expand Down
Loading