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
4 changes: 2 additions & 2 deletions playground/templates/rust/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

50 changes: 49 additions & 1 deletion polyglot/python_worker/scripts/sdk_updates.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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"),
Expand Down
46 changes: 46 additions & 0 deletions polyglot/python_worker/tests/test_sdk_updates.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
2 changes: 1 addition & 1 deletion polyglot/qualified-artifact-tuple.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
4 changes: 2 additions & 2 deletions polyglot/rust_worker/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion polyglot/rust_worker/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
22 changes: 20 additions & 2 deletions polyglot/rust_worker/src/updates.rs
Original file line number Diff line number Diff line change
@@ -1,20 +1,38 @@
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)}))
});
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<()> {
Expand Down
12 changes: 9 additions & 3 deletions polyglot/updates/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
3 changes: 2 additions & 1 deletion scripts/sdk-updates.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)"
Expand Down
Loading