Repository navigation
Workflow: cross-app client operations - #1201
Conversation
Adds an optional app_id argument to every client-level workflow operation on DaprWorkflowClient and its async counterpart: schedule_new_workflow, get_workflow_state, wait_for_workflow_start, wait_for_workflow_completion, raise_workflow_event, terminate_workflow, pause_workflow, resume_workflow and purge_workflow. When set, the operation targets a workflow instance owned by another app in the same namespace, and the target app's WorkflowAccessPolicy decides whether it is permitted. When unset or equal to the local app, the behaviour is unchanged. The vendored durabletask client carries the value as a TaskRouter with targetAppID on each request, built by the new new_task_router helper. An older runtime ignores the field and applies the operation locally. Signed-off-by: joshvanl <me@joshvanl.dev>
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #1201 +/- ##
==========================================
+ Coverage 83.15% 83.70% +0.54%
==========================================
Files 123 123
Lines 10260 10266 +6
==========================================
+ Hits 8532 8593 +61
+ Misses 1728 1673 -55 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
There was a problem hiding this comment.
🟡 Changes recommended
Proto provenance is inconsistent with generated output, and cross-app purge semantics are documented incorrectly.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
Adds cross-app routing to synchronous and asynchronous workflow client operations using TaskRouter.
Changes:
- Adds optional
app_idforwarding across workflow lifecycle operations. - Regenerates durabletask protocol bindings and improves regeneration provenance.
- Adds unit and integration coverage for cross-app routing.
File summaries
| File | Description |
|---|---|
tools/regen_durabletask_protos.sh |
Supports local proto sources and records provenance. |
tests/integration/test_workflow_cross_app.py |
Tests cross-app workflow operations end to end. |
tests/integration/conftest.py |
Refines app-channel readiness checks. |
tests/integration/apps/workflow_host.py |
Provides the remote workflow host. |
tests/ext/workflow/test_workflow_client.py |
Tests sync app_id forwarding. |
tests/ext/workflow/test_workflow_client_aio.py |
Tests async app_id forwarding. |
tests/ext/workflow/durabletask/test_orchestration_executor.py |
Updates action-router assertions. |
tests/ext/workflow/durabletask/test_client_routing.py |
Verifies routing on generated requests. |
dapr/ext/workflow/dapr_workflow_client.py |
Exposes cross-app sync operations. |
dapr/ext/workflow/aio/dapr_workflow_client.py |
Exposes cross-app async operations. |
dapr/ext/workflow/_durabletask/internal/PROTO_SOURCE_COMMIT_HASH |
Updates recorded proto source revision. |
dapr/ext/workflow/_durabletask/internal/orchestrator_service_pb2.pyi |
Adds typed operation-router fields. |
dapr/ext/workflow/_durabletask/internal/orchestrator_service_pb2.py |
Updates generated operation protocol descriptors. |
dapr/ext/workflow/_durabletask/internal/orchestrator_actions_pb2.pyi |
Updates generated action types. |
dapr/ext/workflow/_durabletask/internal/orchestrator_actions_pb2.py |
Updates generated action descriptors. |
dapr/ext/workflow/_durabletask/internal/orchestration_pb2.pyi |
Adds generated workflow metadata types. |
dapr/ext/workflow/_durabletask/internal/orchestration_pb2.py |
Updates orchestration descriptors. |
dapr/ext/workflow/_durabletask/internal/history_events_pb2.pyi |
Adds child-retry metadata typing. |
dapr/ext/workflow/_durabletask/internal/history_events_pb2.py |
Updates history event descriptors. |
dapr/ext/workflow/_durabletask/internal/helpers.py |
Keeps routing on enclosing workflow actions. |
dapr/ext/workflow/_durabletask/internal/backend_service_pb2.pyi |
Adds unique-instance request typing. |
dapr/ext/workflow/_durabletask/internal/backend_service_pb2.py |
Updates backend protocol descriptors. |
dapr/ext/workflow/_durabletask/client.py |
Builds and attaches sync task routers. |
dapr/ext/workflow/_durabletask/aio/client.py |
Builds and attaches async task routers. |
Review details
Files not reviewed (5)
- dapr/ext/workflow/_durabletask/internal/backend_service_pb2.py: Generated file
- dapr/ext/workflow/_durabletask/internal/history_events_pb2.py: Generated file
- dapr/ext/workflow/_durabletask/internal/orchestration_pb2.py: Generated file
- dapr/ext/workflow/_durabletask/internal/orchestrator_actions_pb2.py: Generated file
- dapr/ext/workflow/_durabletask/internal/orchestrator_service_pb2.py: Generated file
Suppressed comments (1)
tests/ext/workflow/durabletask/test_orchestration_executor.py:770
- This is the
without_app_idcase, so the docstring still describes the opposite scenario.
"""Tests that the workflow action carries correct router fields when app_id is specified"""
- Files reviewed: 19/24 changed files
- Comments generated: 4
- Review effort level: Balanced
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Signed-off-by: joshvanl <me@joshvanl.dev>
sicoyle
left a comment
There was a problem hiding this comment.
this is great!! thanks 🚀 few comments so far pls
we already have examples/workflow/multi-app1.py and multi-app2.py showing the context-level cross-app calls. A client-level equivalent would go a long way for discoverability — right now there's no example anywhere showing you can pass app_id straight to schedule_new_workflow.
| schedule = pb.ScheduleTaskAction( | ||
| name=name, | ||
| input=get_string_value(encoded_input), | ||
| router=router, |
There was a problem hiding this comment.
Heads up — I think this one goes a bit beyond what the PR description covers. Dropping the router off the inner ScheduleTaskAction (and CreateChildWorkflowAction below) changes the wire format for cross-app activities and child workflows, which is separate from the client-level app_id feature this PR is about. Given our N-2 policy, I can't approve this unless every runtime in that window already reads WorkflowAction.router for these actions — if any of them still look at the inner field, cross-app calls break silently against a supported daprd, and the executor tests that would've caught it got relaxed in this PR rather than kept. Could you confirm the earliest runtime version that honors the outer router?
There was a problem hiding this comment.
the field is reserved in the pinned protos, so it cannot be set. No runtime in the N-2 window ever read it. Dapr 1.16.0 ships durabletask-go v0.9.0, 1.17.9 ships v0.11.5, 1.18.0 ships v0.12.1, 1.18.4 ships v0.12.5. In all four, applier.go builds the event with Router: action.Router, the outer field. Dapr itself never reads the inner field on any of those versions or master. The two dropped assertions referenced action.scheduleTask.router, which now raises AttributeError.
| @@ -322,8 +322,15 @@ class CreateInstanceRequest(_message.Message): | |||
| EXECUTIONID_FIELD_NUMBER: _builtins.int | |||
There was a problem hiding this comment.
can you pls confirm there are no manual edits to this file. seems commit 6732be7 includes a manual edit...
There was a problem hiding this comment.
commit 6732be7 is a clean regeneration, not a hand edit. Re-running the script against a clean checkout produces a byte-identical directory
| assert state.runtime_status.name == 'COMPLETED' | ||
|
|
||
| caller_client.purge_workflow(instance_id, app_id=HOST_APP_ID) | ||
| assert host_client.get_workflow_state(instance_id) is None |
There was a problem hiding this comment.
every other status check in file polls with a deadline. can you do the same here pls to avoid flake potential in future
| @@ -0,0 +1,170 @@ | |||
| # -*- coding: utf-8 -*- | |||
There was a problem hiding this comment.
every docstring here promises the target app's WorkflowAccessPolicy decides whether the operation is allowed, but nothing tests a rejected app_id. That's the security-relevant half of the feature — can we get coverage for the denial case pls?
There was a problem hiding this comment.
I removed the WAP tests bc they require mTLS which is difficult to setup in selfhosted
Signed-off-by: joshvanl <me@joshvanl.dev>
Signed-off-by: joshvanl <me@joshvanl.dev>
dapr#1201 added an optional app_id to every client-level operation while this branch was open. Of the three management APIs only rerun can follow: its request carries a router and the runtime reads it, routing to the target app's workflow actor type. Listing and history requests have no router field at all, and the runtime builds both against its own app ID, so an app_id argument there would be accepted and silently ignored. Rerun now takes app_id and builds the router through the same helper the other operations use. A test asserts the asymmetry against the protobuf descriptors, so if listing or history ever gains a router it fails and says to add the argument there too. Cross-app routing needs a 1.19 runtime, so this is covered by unit tests rather than against a live 1.18 sidecar, matching how dapr#1201 documents it. Signed-off-by: Javier Aliaga <javier@diagrid.io>
* Add workflow management APIs: list, history and rerun The three advanced workflow management operations from dapr/dapr#9729 had no Python surface: the vendored durabletask protos carried ListInstanceIDs, GetInstanceHistory and RerunWorkflowFromEvent, but neither client layer exposed them, so reaching them meant using the gRPC stub directly. DaprWorkflowClient and its async counterpart now expose: - list_workflow_instances(page_size, continuation_token) -> one page of instance IDs plus the token for the next, when the caller wants to hold the cursor themselves. - iter_workflow_instances(page_size) -> a lazy iterator that pages internally; an async generator on the async client. - get_workflow_history(instance_id) -> the instance's events as WorkflowHistoryEvent records. - rerun_workflow_from_event(instance_id, event_id, ...) -> the ID of a new instance that replays history up to the chosen event and resumes there. The rerun input is a single argument rather than a value plus a flag. The wire format pairs a non-optional StringValue with an overwriteInput bool precisely because StringValue cannot express absence, so the two are collapsed behind a sentinel default: omitting input keeps the original, passing None clears it. A falsy value such as 0 still overwrites. WorkflowHistoryEvent carries event_id, timestamp, event_type, name, task_scheduled_id and failure_details, which is enough to choose a rerun point by activity name instead of by raw event number. is_rerunnable reports this SDK's snapshot of which event types the runtime restarts from; the sidecar keeps the final say. Unrecognised event types map to UNKNOWN rather than raising, so a newer sidecar cannot break history reads. Verified end to end against runtime 1.18.0: a failed order is listed, its history read, and the failed charge rerun with a corrected input to completion. examples/workflow/workflow_management.py covers that flow and is asserted by tests/examples/test_workflow.py. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Javier Aliaga <javier@diagrid.io> * Narrow the AGENTS.md note on NOT_FOUND handling The blanket statement that the client converts "no such instance exists" to a None return only ever described get_workflow_state; the other methods propagate. Adding get_workflow_history, which raises NOT_FOUND for a missing or purged instance, made the sentence actively misleading. Also covers the input sentinel's repr, which is what help() and tracebacks show for the rerun default. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Javier Aliaga <javier@diagrid.io> * Address review feedback on the workflow management APIs Renames the listing methods to say what they return. They hand back instance IDs, not workflows, and a reviewer read them the other way. list_workflow_instance_ids and iter_workflow_instance_ids also line up with java-sdk#1798's listInstanceIds, and leave the unqualified name free if a filtered list returning WorkflowState ever lands. Moves the rerun input sentinel to the public API surface. The wire needs an input field plus an overwriteInput flag, so the engine layer now takes exactly that pair and knows nothing about sentinels. UNSET is defined in workflow_management.py and exported, which is what forwarding code needs: a wrapper passing an optional input through could not previously say "not supplied" without importing from a private module. Rejects a negative event_id before building the request. eventID is uint32 on the wire, so protobuf refused it with a message naming neither the argument nor the reason, and -1 is reachable precisely because it is what the runtime reports for history events it assigns no ID to. Documentation the review found misleading: - The continuation token is opaque and produced by the state store, not by Dapr, so it must not be parsed or expected to survive a component change. - task_scheduled_id has no presence on the wire, so a 0 is equally a real event ID or a field the runtime never set. Unlike event_id there is no sentinel to test for. - The AGENTS.md "Public API" block claimed to list every exported symbol and omitted the ones this PR adds. The example no longer sleeps after start(), which already waits for the worker's stream, and its rerun line now names both spellings of the same number: the runtime's error calls it "activity task dapr#2" where we call it "event dapr#2". Output recaptured from a live run rather than edited by hand. Also adds a laziness test for the async iterator, whose generator semantics differ from the sync one's, and makes the test fake build the real request so it rejects what the engine would reject. Signed-off-by: Javier Aliaga <javier@diagrid.io> * Say what clearing a rerun input does to the activity "Pass None to clear it" left a reviewer asking what clearing actually means. It means the activity being rerun receives None where it previously received its recorded input, so say that instead. Signed-off-by: Javier Aliaga <javier@diagrid.io> * Guard the remaining rerun inputs and explain an unlistable store Three things the runtime does not validate, all reaching the caller as errors that name neither the argument nor the reason. An empty new_instance_id or new_child_workflow_instance_id is taken literally rather than treated as absent. Both fields have explicit presence, so None leaves them unset and the runtime generates an ID, while '' becomes the ID itself and produces an instance the runtime cannot schedule reminders for. rerun.go forwards the field without checking it, so nothing downstream catches this. Guarding only event_id, as the previous commit did, left the asymmetry. event_id was checked at one end only. The field is uint32 on the wire, so 4294967296 fails exactly as -1 did, with the same protobuf message. It is now a range check, with tests on both bounds and on the largest value that must stay valid. Listing instances needs an actor state store that can list keys, and the runtime returns a bare error for both ways that can be missing, so they arrive as UNKNOWN with the reason buried in the details. Both are now translated into a NotImplementedError that says what is wrong and what it needs, keeping the original error as the cause. Verified against a sidecar with no actor state store; the unlistable-store branch is covered by unit test only, since every store that can run locally implements key listing. Signed-off-by: Javier Aliaga <javier@diagrid.io> * Let a rerun target another app, and say why listing cannot dapr#1201 added an optional app_id to every client-level operation while this branch was open. Of the three management APIs only rerun can follow: its request carries a router and the runtime reads it, routing to the target app's workflow actor type. Listing and history requests have no router field at all, and the runtime builds both against its own app ID, so an app_id argument there would be accepted and silently ignored. Rerun now takes app_id and builds the router through the same helper the other operations use. A test asserts the asymmetry against the protobuf descriptors, so if listing or history ever gains a router it fails and says to add the argument there too. Cross-app routing needs a 1.19 runtime, so this is covered by unit tests rather than against a live 1.18 sidecar, matching how dapr#1201 documents it. Signed-off-by: Javier Aliaga <javier@diagrid.io> * Reject page sizes the stores mishandle, and close two test gaps A page_size of 0 reached the wire as a literal 0, and the stores do not agree on what that means. Casper ran each one's KeysLike: in-memory returns an empty page with the same continuation token every time, so iter_workflow_instance_ids(page_size=0) never ends, and sqlite indexes recs[*req.PageSize-1], which wraps on uint32 and panics. A negative or over-range value failed the same opaque protobuf way event_id already guarded against. list_instance_ids now builds its request through a shared builder that rejects anything outside 1 to uint32, so both engines validate once. _MAX_EVENT_ID becomes _MAX_UINT32, since eventID and pageSize share the bound and two names for it invite drift. Two mutations passed the suite before this. Replacing the async client's error translation with a bare raise went unnoticed because SimulatedRpcError is not an AioRpcError, so the async except never ran at all; there is now a SimulatedAioRpcError and an async twin of the unsupported-store tests. Dropping app_id from the async rerun also passed, because the test forwarded None and could not tell a forwarded argument from a discarded one; it now forwards a real app ID. ListInstanceIDs and GetInstanceHistory landed in 1.17, so an older sidecar answers UNIMPLEMENTED with nothing to act on. Both now raise NotImplementedError naming the version, which is already this API's contract for "the sidecar cannot do this". Rerun is left alone: it has existed since 1.16, so the error is not reachable there. The app_id docstring said older runtimes "ignore" it. They drop the routing instead, which means rerunning a local instance with the same ID if one exists, so it now says that. _listing_unsupported_message and its two runtime strings were copied into both clients and now live once in workflow_management.py, and the rerun Raises section covers the over-range event_id and the empty instance IDs it already rejected. Signed-off-by: Javier Aliaga <javier@diagrid.io> --------- Signed-off-by: Javier Aliaga <javier@diagrid.io> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Adds an optional app_id argument to every client-level workflow operation on DaprWorkflowClient and its async counterpart: schedule_new_workflow, get_workflow_state, wait_for_workflow_start, wait_for_workflow_completion, raise_workflow_event, terminate_workflow, pause_workflow, resume_workflow and purge_workflow. When set, the operation targets a workflow instance owned by another app in the same namespace, and the target app's WorkflowAccessPolicy decides whether it is permitted. When unset or equal to the local app, the behaviour is unchanged.
The vendored durabletask client carries the value as a TaskRouter with targetAppID on each request, built by the new new_task_router helper. An older runtime ignores the field and applies the operation locally.