Add workflow management APIs: list, history and rerun - #1217
javier-aliaga wants to merge 8 commits into
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #1217 +/- ##
==========================================
+ Coverage 83.89% 84.20% +0.30%
==========================================
Files 123 124 +1
Lines 10265 10450 +185
==========================================
+ Hits 8612 8799 +187
+ Misses 1653 1651 -2 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
There was a problem hiding this comment.
🟡 Changes recommended
dapr/ext/workflow/AGENTS.md now contains a misleading statement about NOT_FOUND being converted to None, which conflicts with the new get_workflow_history() behavior.
Get a fresh assessment by requesting another Copilot review.
Pull request overview
This PR extends the Dapr Python SDK workflow extension by exposing three durabletask-backed workflow management capabilities (instance listing, history retrieval, and rerun-from-event) on both DaprWorkflowClient and dapr.ext.workflow.aio.DaprWorkflowClient. It also adds a runnable example plus unit and example tests to validate the new APIs end-to-end.
Changes:
- Added workflow management client APIs:
list_workflow_instances,iter_workflow_instances,get_workflow_history, andrerun_workflow_from_event(sync + async). - Introduced typed return models for instance pages and history events (
WorkflowInstanceIdPage,WorkflowHistoryEvent,WorkflowHistoryEventType), including “unknown event type” resilience. - Added comprehensive unit tests (client layer + engine client layer) and a new example validated by the examples test suite.
File summaries
| File | Description |
|---|---|
dapr/ext/workflow/dapr_workflow_client.py |
Adds the new management APIs to the sync workflow client. |
dapr/ext/workflow/aio/dapr_workflow_client.py |
Adds async equivalents of the new management APIs. |
dapr/ext/workflow/workflow_management.py |
New typed models and conversions for instance pages and history events. |
dapr/ext/workflow/_durabletask/client.py |
Adds engine-client RPC wrappers + sentinel handling for rerun input semantics. |
dapr/ext/workflow/_durabletask/aio/client.py |
Async engine-client wrappers for list/history/rerun using shared request builder. |
dapr/ext/workflow/__init__.py |
Exposes the new management types (and FailureDetails) at the extension top-level. |
dapr/ext/workflow/AGENTS.md |
Documents the new APIs (but needs a small correction to the NOT_FOUND/None note). |
tests/ext/workflow/test_workflow_management.py |
New unit tests covering sync + async workflow management APIs. |
tests/ext/workflow/durabletask/test_client_management_apis.py |
New engine-client tests validating request/response behavior for list/history/rerun. |
examples/workflow/workflow_management.py |
New example demonstrating list → history → rerun flow. |
examples/workflow/README.md |
Documents the new example and explains rerunnable event types and caveats. |
tests/examples/test_workflow.py |
Adds output-based validation for the new workflow management example. |
Review details
- Files reviewed: 12/12 changed files
- Comments generated: 1
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
There was a problem hiding this comment.
🟢 Approval recommended
The implementation matches the stated additive API goals and is backed by substantial unit and example test coverage, with only minor docstring formatting nits noted.
Review details
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
dapr/ext/workflow/aio/dapr_workflow_client.py:327
- The
page_sizeargument docstring continuation lines are mis-indented, which can cause doc rendering/readability issues (they appear as separate arguments). Indent continuation lines under thepage_size:description.
dapr/ext/workflow/dapr_workflow_client.py:325 - The
page_sizeargument docstring continuation lines are mis-indented, making them look like separate arguments when rendered (and harder to read in help()/Sphinx). Indent continuation lines under thepage_size:description.
- Files reviewed: 12/12 changed files
- Comments generated: 0 new
- Review effort level: Lite
What does the pass None to clear it actually mean? Does that essentially rerun the workflow from event with a new workflow id so it is essentially technically not a rerun? Can you pls help me understand that one? 🙏 |
sicoyle
left a comment
There was a problem hiding this comment.
thank youuuu!!
did you run checks to see if the changes in this PR are along the same lines in terms of behavior and api surface to that of the other sdks supporting this?
few comments in general from me pls. Also, from claude below to pls address:
Worth settling first
1. UNSET is a private sentinel sitting in a public signature.
dapr_workflow_client.py:382 defaults input to client.UNSET, where client is dapr.ext.workflow._durabletask.client. It isn't in dapr.ext.workflow.__all__, so anyone who wants to name the default has to import from a private module. That matters for wrapper code: a service layer or retry helper that forwards input through can't write input=UNSET, it has to build kwargs conditionally.
This PR already promoted FailureDetails to public for exactly this reason. UNSET deserves the same. Export it from dapr.ext.workflow and reference it in the docstring, or drop the sentinel for an explicit overwrite_input: bool = False pair that mirrors the wire. I prefer exporting the sentinel; the single-argument API reads better and the docstring explains it well.
2. list_workflow_instances returns IDs, not instances.
The return type is WorkflowInstanceIdPage, the iterator yields str, and the RPC is ListInstanceIDs. list_workflow_instance_ids / iter_workflow_instance_ids says what the caller gets, and leaves the unqualified name free if the runtime ever grows a filtered list that returns WorkflowState. I checked go-sdk and nothing there implements these yet, so whatever lands here is the precedent the other SDKs will copy, and renaming post-release is breaking.
On the question in the PR description
Letting NOT_FOUND propagate from get_workflow_history is the right call and I wouldn't invert it. get_workflow_state is the outlier in this client, not the rule: pause, resume, terminate and purge all propagate. And for history specifically, None versus [] would be a genuinely ambiguous return where a raised NOT_FOUND is not. The AGENTS.md wording after the second commit is accurate.
Smaller notes
- task_scheduled_id: getattr(payload, 'taskScheduledId', None) yields None only for payload types that don't define the field. For a type that defines it but leaves it unset you get 0, since it's a proto3 scalar with no presence. The runtime always sets it in practice, but the docstring sits right next to the event_id == -1 convention, which makes None read like the documented "absent" value. Either say so in the attribute docs or note that the field can't be distinguished from a real zero. Don't paper over it with or None.
- iter_workflow_instances loops forever if a sidecar keeps handing back the same non-empty continuation token. You've covered the empty-token case in tests. A one-line guard (stop when the token repeats) closes the remaining hole cheaply. Optional, since a server behaving that way is broken anyway.
- is_rerunnable encodes a server-side rule client-side. The docstring and README both hedge it correctly and test_every_event_type_the_proto_defines_is_mapped will fail loudly on a proto bump, which is the right safety net. I'd keep it. Just be aware that when the runtime adds a fourth rerunnable type, the example's find_charge_event_id pattern silently skips it until someone re-syncs the protos.
- The AGENTS.md "Public API" block at line 79 claims all public symbols are exported and then lists none of WorkflowHistoryEvent, WorkflowHistoryEventType, WorkflowInstanceIdPage or FailureDetails. That list was already stale before this PR, but since you're adding four symbols it's a cheap fix.
- Async side has no equivalent of test_fetches_lazily. The sync generator and the async generator have different laziness semantics in principle, so it's worth the eight lines.
None it meant the activity will receive NONE as input *** echo_act received 0 (type int) <- original run, fails For the second question a rerun is always a new workflow. |
sicoyle
left a comment
There was a problem hiding this comment.
final comments for you - thanks @javier-aliaga 🙌
| @@ -38,11 +44,16 @@ | |||
| 'WorkflowActivityContext', | |||
| 'WorkflowState', | |||
| 'WorkflowStatus', | |||
| 'WorkflowHistoryEvent', | |||
| 'WorkflowHistoryEventType', | |||
| 'WorkflowInstanceIdPage', | |||
| 'UNSET', | |||
| 'when_all', | |||
| 'when_any', | |||
| 'alternate_name', | |||
| 'RetryPolicy', | |||
| 'TaskFailedError', | |||
| 'FailureDetails', | |||
| 'PropagationScope', | |||
There was a problem hiding this comment.
can you pls confirm that we must expose all of these? which of these can be internal only?
There was a problem hiding this comment.
I kept most of them, just moved to private this one RERUNNABLE_EVENT_TYPES
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>
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>
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>
"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>
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>
7a217b4 to
fce303c
Compare
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>
CasperGN
left a comment
There was a problem hiding this comment.
Request changes — one small fix on page_size, the rest is non-blocking. Nice work on the rerun input semantics and the history mapping.
Blocking: validate page_size (>= 1, <= uint32) before it reaches the wire.
page_size=0 is sent as-is (sync engine, aio engine), and the stores don't agree on what 0 means. I ran each store's KeysLike directly (components-contrib main):
- in-memory returns
keys=[]with the same token ("0") every time, soiter_workflow_instance_ids(page_size=0)never ends (sync, aio). - sqlite panics:
index out of range [4294967295] with length 0(recs[*req.PageSize-1]wraps on uint32). I couldn't find a recover interceptor on the daprd gRPC path, so I expect this crashes the sidecar. I didn't run it against daprd.
A negative value or one above uint32 gets the same unclear protobuf ValueError you already guard against for event_id in _new_rerun_request. The same range check belongs in list_instance_ids. While you're there, the repeated-token guard sicoyle mentioned would have stopped the in-memory loop above, so I'd add it too. The store behaviour is worth a components-contrib issue too.
Non-blocking: async error translation has no test. Replacing raise NotImplementedError(advice) from error with a bare raise in the aio client still passes every test. ListingUnsupportedTest is sync-only, and SimulatedRpcError isn't an AioRpcError, so the async except never runs. Dropping app_id on the async rerun also passes, because test_rerun_forwards_every_argument sends app_id=None. Worth an async copy of the unsupported-store tests, plus a non-None app_id in that test.
Non-blocking: say what happens on older runtimes. ListInstanceIDs/GetInstanceHistory landed in dapr/dapr#9170 and are missing from v1.16.20. There they come back as a raw UNIMPLEMENTED RpcError. NotImplementedError is already this API's contract for "the sidecar can't list", so I'd map UNIMPLEMENTED to it as well, or state the minimum runtime in the docstring/README. The routed rerun is only on dapr master (v1.18.4 still forks on the local workflowActorType). The docstring's "older runtimes ignore app_id" covers that, but for rerun it means forking a local instance with the same ID if one exists. That deserves a sentence.
Non-blocking: _listing_unsupported_message and its two constants are duplicated in the sync and aio modules. Both strings match pkg/runtime/wfengine/state/list/list.go on master today. One copy (e.g. in workflow_management.py) keeps the two clients from drifting when that text changes.
Non-blocking: the rerun Raises: section (sync, aio) only lists negative event_id. It should also cover values above uint32 and an empty new_instance_id/new_child_workflow_instance_id.
What I checked: the full diff, plus the runtime side (dapr master actors.go/list.go, durabletask-go executor, components-contrib KeysLike for in-memory, sqlite, redis and postgres). ruff, ruff format and mypy are clean. The 71 new tests pass, and the full unit suite has 1721 passed, 35 deselected. I made 16 mutations to the new code, and 13 were caught by the tests. Two that weren't are the async ones above. The third, the engine dropping input when overwrite_input=False, can't happen through the public API. Public API changes are additive only, and there are no proto changes (vendored protos already had the RPCs). The sync and aio signatures match.
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>
CasperGN
left a comment
There was a problem hiding this comment.
Claude review · approved on Casper's behalf.
Thanks, the blocking item is fixed: _new_list_request rejects a page_size outside 1..uint32 on both engines, with a clear message. The rest of the earlier review is covered too:
- UNIMPLEMENTED now maps to
NotImplementedErrornaming 1.17, for listing and history, sync and aio. - The advice strings live once in
workflow_management.py. - The
Raises:sections cover uint32 and empty IDs. - There are async copies of the unsupported-store and old-runtime tests, and the rerun test now uses a non-None
app_id.
Minor
- The new sentence about older runtimes and rerun landed in
get_workflow_state'sapp_iddoc, whilererun_workflow_from_eventstill has the old "older runtimes ignore app_id" wording. Moving it would put it where it applies (sync and aio).
python-sdk/dapr/ext/workflow/dapr_workflow_client.py
Lines 174 to 177 in a5bdf2d
python-sdk/dapr/ext/workflow/dapr_workflow_client.py
Lines 551 to 552 in a5bdf2d
python-sdk/dapr/ext/workflow/aio/dapr_workflow_client.py
Lines 173 to 176 in a5bdf2d
iter_workflow_instance_idsstill has no repeated-token guard. Thepage_sizecheck closes the known in-memory loop, but a store that returns the same token again would still loop forever; stopping whenpage.continuation_token == continuation_tokenis cheap insurance.
python-sdk/dapr/ext/workflow/dapr_workflow_client.py
Lines 463 to 484 in a5bdf2d
Checked: 33a2161 against the earlier review. pytest tests/ext/workflow -m "not e2e" gives 537 passed, ruff check and ruff format --check are clean, and all 22 CI checks pass.
Reviewed at a5bdf2d.
Disagree, or want to talk it through? Reply /human and Casper will pick it up.
Description
The vendored durabletask protos already carry
ListInstanceIDs,GetInstanceHistoryandRerunWorkflowFromEvent, but neither client layer exposed them. This adds all three toDaprWorkflowClientand its async counterpart:list_workflow_instance_ids(*, page_size, continuation_token)WorkflowInstanceIdPage— one page plus the next tokeniter_workflow_instance_ids(*, page_size=1024)async foron the async client)get_workflow_history(instance_id)list[WorkflowHistoryEvent]rerun_workflow_from_event(instance_id, event_id, *, app_id, ...)All additive; no existing signature changes.
Three behaviours worth knowing:
rerun_workflow_from_eventtakes the replacement input as one argument. Omit it to keep the original input, passNoneto clear it, in which case the rerun activity receivesNone. Forwarding code can pass the exportedUNSETsentinel to mean "not supplied". The engine layer underneath takes the wire'sinput+overwrite_inputpair directly.WorkflowHistoryEventcarriesevent_id,timestamp,event_type,name,task_scheduled_idandfailure_details, so a rerun point can be chosen by activity name. Unrecognised event types map toUNKNOWNrather than raising.NotImplementedErrorthat names what is missing.Issue reference
Closes #1181
Part of dapr/dapr#9729