diff --git a/playground/templates/rust/Cargo.lock b/playground/templates/rust/Cargo.lock index d6367ff..75a93e8 100644 --- a/playground/templates/rust/Cargo.lock +++ b/playground/templates/rust/Cargo.lock @@ -268,9 +268,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.3.0" +version = "3.3.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c5f5cc23a542ec5249798a396ac3e5745c0819ec94f682372d453e338d88278c" +checksum = "10f10a62d9691094f654557ae8c45a2abd0b608e389860b0683a559b9936912f" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/python_worker/scripts/sdk_updates.py b/polyglot/python_worker/scripts/sdk_updates.py index 77158c4..f09d59e 100644 --- a/polyglot/python_worker/scripts/sdk_updates.py +++ b/polyglot/python_worker/scripts/sdk_updates.py @@ -171,6 +171,54 @@ async def replacement(client): emit(scenario="replacement-and-duplicate", runtime="rust", update_id=update_id, run_id=execution.run_id, result=expected, events=related) await matrix(client) + await snapshot(client, "replacement") + + +async def snapshot(client, stage="initial"): + request_id = f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-snapshot-{stage}" + signal = {"request_id": f"{required('DURABLE_WORKFLOW_UPDATE_ID')}-touch", "delta": 7} + signal_arguments = [[signal], [[1, 2]], []] + if stage == "initial": + for arguments in signal_arguments: + await client.get_workflow_handle(workflow_id("rust")).signal("updates-touch", args=arguments) + execution, history = await observe(client, "rust") + deliveries = [event for event in history["events"] if event["event_type"] == "SignalReceived" + and event["payload"].get("signal_name") == "updates-touch"] + if (len(deliveries) != len(signal_arguments) or not same_result( + [serializer.decode_envelope(event["payload"]["arguments"], codec="avro") for event in deliveries], + signal_arguments)): + raise RuntimeError("The original snapshot signals were not durably recorded once each.") + expected = {"workflow_id": execution.workflow_id, "run_id": execution.run_id, + "workflow_input": [workflow_id("rust")], "signals": signal_arguments} + query = await client.query_workflow(execution.workflow_id, "snapshot") + response = await client.update_workflow(execution.workflow_id, "snapshot", + args=[request("python", request_id)], wait_for="completed", request_id=request_id) + result = serializer.decode_envelope(response["result_envelope"], codec="avro") + execution, history = await observe(client, "rust") + update_id, related = update_events(history, request_id) + completed = [event for event in related if event["event_type"] == "UpdateCompleted"] + if (response.get("update_status") != "completed" or response.get("update_id") != update_id + or len(completed) != 1 or completed[0]["payload"].get("failure_id") + or not same_result(serializer.decode_envelope(completed[0]["payload"]["result"], codec="avro"), result)): + fields = ("signal_id", "signal_name", "workflow_command_id", "update_id", "sequence", "failure_id", "message") + emit(scenario="rust-update-snapshot-incomplete", stage=stage, expected=expected, + query=query.get("result"), update=result, + response={key: response.get(key) for key in ("update_id", "update_status", "ordering_state", + "queued_behind_command_id", "queued_behind_command_type")}, + history=[{"event_type": event["event_type"], "sequence": event.get("sequence"), + "payload": {key: event["payload"][key] for key in fields if key in event["payload"]}} + for event in history["events"]]) + raise RuntimeError("Snapshot update did not retain its one original completion.") + applied = [event for event in history["events"] if event["event_type"] == "SignalApplied" + and event["payload"].get("signal_name") == "updates-touch"] + signal_ids = {event["payload"].get("signal_id") for event in deliveries} + if (len(applied) != len(deliveries) or None in signal_ids or len(signal_ids) != len(deliveries) + or {event["payload"].get("signal_id") for event in applied} != signal_ids): + raise RuntimeError("Snapshot signals were not applied once each.") + emit(scenario="rust-update-snapshot", stage=stage, expected=expected, + query=query.get("result"), update=result, run_id=execution.run_id, update_id=update_id) + if not same_result(query.get("result"), expected) or not same_result(result, expected): + raise RuntimeError("Query and update must expose the original workflow input and committed signal snapshot.") async def failure(client): @@ -232,7 +280,7 @@ async def finish(client): async def main(): parser = argparse.ArgumentParser() - parser.add_argument("phase", choices=["start", "call", "matrix", "queued", "replacement", "failure", "finish"]) + parser.add_argument("phase", choices=["start", "call", "matrix", "queued", "replacement", "snapshot", "failure", "finish"]) parser.add_argument("arguments", nargs="*") args = parser.parse_args() async with Client(required("DURABLE_WORKFLOW_SERVER_URL"), token=required("DURABLE_WORKFLOW_AUTH_TOKEN"), diff --git a/polyglot/python_worker/tests/test_sdk_updates.py b/polyglot/python_worker/tests/test_sdk_updates.py index 7b73674..d946178 100644 --- a/polyglot/python_worker/tests/test_sdk_updates.py +++ b/polyglot/python_worker/tests/test_sdk_updates.py @@ -158,6 +158,52 @@ async def test_sdk_failure_rejects_wrong_diagnostics(self): await updates.failure(client) +class SnapshotObservationTest(unittest.IsolatedAsyncioTestCase): + async def test_update_and_query_read_the_same_committed_snapshot(self): + signal = {"request_id": "example-touch", "delta": 7} + expected = {"workflow_id": "example-rust", "run_id": "original", + "workflow_input": ["example-rust"], "signals": [[signal], [[1, 2]], []]} + execution = types.SimpleNamespace(workflow_id="example-rust", run_id="original") + history = {"events": [ + {"event_type": "SignalReceived", "payload": {"signal_name": "updates-touch", + "signal_id": "map", "arguments": serializer.envelope([signal])}}, + {"event_type": "SignalReceived", "payload": {"signal_name": "updates-touch", + "signal_id": "nested", "arguments": serializer.envelope([[1, 2]])}}, + {"event_type": "SignalReceived", "payload": {"signal_name": "updates-touch", + "signal_id": "empty", "arguments": serializer.envelope([])}}, + *[{"event_type": "SignalApplied", "payload": {"signal_name": "updates-touch", "signal_id": identity}} + for identity in ("map", "nested", "empty")], + {"event_type": "UpdateAccepted", "payload": {"update_id": "snapshot-update", + "arguments": serializer.envelope([updates.request("python", "example-snapshot-initial")])}}, + {"event_type": "UpdateCompleted", "payload": {"update_id": "snapshot-update", + "result": serializer.envelope(expected)}}, + ]} + client = types.SimpleNamespace( + get_workflow_handle=lambda _id: types.SimpleNamespace(signal=AsyncMock()), + query_workflow=AsyncMock(return_value={"result": expected}), + update_workflow=AsyncMock(return_value={"update_status": "completed", + "update_id": "snapshot-update", "result_envelope": serializer.envelope(expected)}), + ) + with (patch.dict(os.environ, {"DURABLE_WORKFLOW_UPDATE_ID": "example"}), + patch.object(updates, "observe", AsyncMock(return_value=(execution, history))), + patch.object(updates, "emit")): + await updates.snapshot(client) + for change in ({"workflow_input": None}, {"signals": []}, {"run_id": "replacement"}): + with self.subTest(change=change): + result = {**expected, **change} + client.update_workflow.return_value["result_envelope"] = serializer.envelope(result) + history["events"][-1]["payload"]["result"] = serializer.envelope(result) + with self.assertRaisesRegex(RuntimeError, "committed signal snapshot"): + await updates.snapshot(client) + client.update_workflow.return_value["result_envelope"] = serializer.envelope(expected) + with self.assertRaisesRegex(RuntimeError, "original completion"): + await updates.snapshot(client) + history["events"][-1]["payload"]["result"] = serializer.envelope(expected) + history["events"][3]["payload"]["signal_id"] = "nested" + with self.assertRaisesRegex(RuntimeError, "applied once each"): + await updates.snapshot(client) + + class HistoryPaginationTest(unittest.IsolatedAsyncioTestCase): async def test_observer_reads_later_pages(self): client = types.SimpleNamespace( diff --git a/polyglot/qualified-artifact-tuple.json b/polyglot/qualified-artifact-tuple.json index 58aa071..20a10ce 100644 --- a/polyglot/qualified-artifact-tuple.json +++ b/polyglot/qualified-artifact-tuple.json @@ -6,7 +6,7 @@ "cli": "2.1.3", "sdk-php": "2.2.3", "sdk-python": "2.4.3", - "sdk-rust": "3.3.0", + "sdk-rust": "3.3.3", "server": "2.5.9", "waterline": "2.3.1", "workflow": "2.5.3" diff --git a/polyglot/rust_worker/Cargo.lock b/polyglot/rust_worker/Cargo.lock index 9bf7bcd..8d3968a 100644 --- a/polyglot/rust_worker/Cargo.lock +++ b/polyglot/rust_worker/Cargo.lock @@ -244,9 +244,9 @@ dependencies = [ [[package]] name = "durable-workflow" -version = "3.3.0" +version = "3.3.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c5f5cc23a542ec5249798a396ac3e5745c0819ec94f682372d453e338d88278c" +checksum = "10f10a62d9691094f654557ae8c45a2abd0b608e389860b0683a559b9936912f" dependencies = [ "apache-avro", "base64", diff --git a/polyglot/rust_worker/Cargo.toml b/polyglot/rust_worker/Cargo.toml index a99a433..546cd3e 100644 --- a/polyglot/rust_worker/Cargo.toml +++ b/polyglot/rust_worker/Cargo.toml @@ -7,7 +7,7 @@ publish = false [dependencies] apache-avro = "=0.21.0" base64 = "=0.22.1" -durable-workflow = "3.3.0" +durable-workflow = "3.3.3" tokio = { version = "1.47.1", features = ["macros", "rt-multi-thread", "signal"] } [dev-dependencies] diff --git a/polyglot/rust_worker/src/updates.rs b/polyglot/rust_worker/src/updates.rs index 9df71ff..13d8566 100644 --- a/polyglot/rust_worker/src/updates.rs +++ b/polyglot/rust_worker/src/updates.rs @@ -1,13 +1,25 @@ -use durable_workflow::{json, Client, Error, Result, Worker}; +use durable_workflow::{json, Client, Error, QueryContext, Result, Value, Worker}; + +fn snapshot(ctx: QueryContext) -> Value { + json!({ + "workflow_id": ctx.workflow_id, + "run_id": ctx.run_id, + "workflow_input": ctx.workflow_input(), + "signals": ctx.signals("updates-touch"), + }) +} pub fn register(worker: &mut Worker) { worker.register_workflow("polyglot.rust.updates", |ctx, input| async move { let request = super::first_argument(&input); + for _ in 0..3 { + ctx.wait_signal("updates-touch").await?; + } ctx.wait_signal("updates-finish").await?; Ok(json!({"workflow_runtime": "rust", "request": request})) }); worker - .declare_workflow_signals("polyglot.rust.updates", &["updates-finish"]) + .declare_workflow_signals("polyglot.rust.updates", &["updates-finish", "updates-touch"]) .expect("declare the update workflow's completion signal"); worker.register_update("polyglot.rust.updates", "echo", |_ctx, input| async move { Ok(json!({"handler_runtime": "rust", "request": super::first_argument(&input)})) @@ -15,6 +27,12 @@ pub fn register(worker: &mut Worker) { worker.register_update("polyglot.rust.updates", "fail", |_ctx, _input| async move { Err(Error::Codec("update-probe-failure".into())) }); + worker.register_query("polyglot.rust.updates", "snapshot", |ctx, _input| async move { + Ok(snapshot(ctx)) + }); + worker.register_update("polyglot.rust.updates", "snapshot", |ctx, _input| async move { + Ok(snapshot(ctx)) + }); } pub async fn call(client: Client) -> Result<()> { diff --git a/polyglot/updates/README.md b/polyglot/updates/README.md index bb0cff1..0306e15 100644 --- a/polyglot/updates/README.md +++ b/polyglot/updates/README.md @@ -17,9 +17,8 @@ each of the three handler languages. These are client-to-handler directions. The Server's existing update experiment separately covers its embedded probe, PHP process boundary, Python SDK surface fixtures, validators and diagnostics. -Rust update workers require SDK 3.3.0 or later. Their registration includes the -handler argument contracts and declared completion signal Server records when -starting a run. +The Rust worker's registration includes handler argument contracts and declared +signals that Server records when starting a run. Each call checks the named handler's result through the real SDK client and persisted `UpdateAccepted`/`UpdateCompleted` history. It then kills the Rust @@ -33,6 +32,13 @@ The Python client must raise `UpdateFailed` with the handler's message and matching workflow, run, update and failure IDs. Repeating that failed request must return the same error identities with one durable failed completion. +A Rust query and update inspect the same original workflow input and committed +signals before and after worker replacement, including a map argument, one nested +array argument and no arguments. The update result must also match +its persisted completion. This exercises the immutable state snapshot that a +stateful handler uses to reconstruct its input and prior signal deliveries. +The Rust workflow consumes these three signals before waiting for completion. + Rust does not support synchronous pre-accept update validators. The installed crate must return `UnsupportedUpdateValidators` for a contract claiming one. This check does not claim validator execution, process loss during an external diff --git a/scripts/sdk-updates.sh b/scripts/sdk-updates.sh index 2485a06..e1c22a1 100755 --- a/scripts/sdk-updates.sh +++ b/scripts/sdk-updates.sh @@ -3,7 +3,7 @@ set -euo pipefail if [[ "${1:-}" == --help ]]; then printf '%s\n' 'Usage: scripts/sdk-updates.sh' \ - 'Runs nine PHP/Python/Rust client/update-handler directions, Rust worker replacement, duplicate requests, handler failure and validator refusal.' \ + 'Runs nine PHP/Python/Rust client/update-handler directions, committed snapshots, Rust worker replacement, duplicate requests, handler failure and validator refusal.' \ 'Requires Docker Compose and exact assignments from scripts/resolve-current-artifacts.sh.' \ 'SDK_UPDATES_COMPOSE_PROJECT_NAME selects an isolated project. All project resources are removed on exit.' exit 0 @@ -75,6 +75,7 @@ for direction in php:php php:python php:rust python:php python:python python:rus DURABLE_WORKFLOW_UPDATE_RESULTS+="${result}"$'\n' done observer matrix +observer snapshot "${compose[@]}" kill --signal SIGKILL rust-workflow-worker DURABLE_WORKFLOW_UPDATE_QUEUED="$(observer queued)"