Skip to content
13 changes: 12 additions & 1 deletion dapr/ext/workflow/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ dapr/ext/workflow/
├── workflow_context.py # WorkflowContext ABC
├── workflow_activity_context.py # WorkflowActivityContext wrapper
├── workflow_state.py # WorkflowState, WorkflowStatus enum
├── workflow_management.py # WorkflowHistoryEvent(Type), WorkflowInstanceIdPage
├── retry_policy.py # RetryPolicy wrapper
├── util.py # gRPC address resolution
├── logger/options.py # LoggerOptions
Expand All @@ -24,6 +25,7 @@ tests/ext/workflow/
├── test_dapr_workflow_context.py # Context method proxying
├── test_workflow_activity_context.py # Activity context properties
├── test_workflow_client.py # Sync client (mock gRPC)
├── test_workflow_management.py # list/history/rerun on both clients
├── test_workflow_client_aio.py # Async client (IsolatedAsyncioTestCase)
├── test_workflow_runtime.py # Registration, decorators, worker readiness
├── test_workflow_util.py # Address resolution
Expand Down Expand Up @@ -86,6 +88,11 @@ from dapr.ext.workflow import (
when_any, # Race combinator — wait for first task
alternate_name, # Decorator to set a custom registration name
RetryPolicy, # Retry config for activities/child workflows
WorkflowHistoryEvent, # One event from an instance's execution history
WorkflowHistoryEventType, # Enum of history event kinds; unknown ones map to UNKNOWN
WorkflowInstanceIdPage, # One page of instance IDs plus the continuation token
FailureDetails, # Error carried by WorkflowHistoryEvent / TaskFailedError
UNSET, # Sentinel default for rerun_workflow_from_event's input
)

# Async client:
Expand Down Expand Up @@ -139,9 +146,13 @@ Client for workflow lifecycle management:
- `terminate_workflow(instance_id, *, output, recursive)`
- `pause_workflow(instance_id)` / `resume_workflow(instance_id)`
- `purge_workflow(instance_id, *, recursive)`
- `list_workflow_instance_ids(*, page_size, continuation_token)` → `WorkflowInstanceIdPage`. No `app_id`: `ListInstanceIDsRequest` and `GetInstanceHistoryRequest` carry no `router`, and the runtime scopes both to the calling app. Only rerun can be routed cross-app.
- `iter_workflow_instance_ids(*, page_size=1024)` → iterator over instance IDs, paging internally (`async for` on the async client)
- `get_workflow_history(instance_id)` → `list[WorkflowHistoryEvent]`
- `rerun_workflow_from_event(instance_id, event_id, *, new_instance_id, input, new_child_workflow_instance_id, app_id)` → new `instance_id`. Omitting `input` keeps the original; passing `None` clears it. That pair is an `input` + `overwriteInput` pair on the wire, which the engine layer takes as-is; the public method collapses it into one argument via the exported `UNSET` sentinel.
- `close()` — close gRPC connection

Converts gRPC "no such instance exists" errors to `None` returns. The async variant in `aio/` has the same API with `async` methods.
`get_workflow_state` converts gRPC "no such instance exists" errors to a `None` return. The other methods let the error propagate, including `get_workflow_history`, which raises NOT_FOUND for a missing or purged instance. The async variant in `aio/` has the same API with `async` methods.

### DaprWorkflowContext (`dapr_workflow_context.py`)

Expand Down
13 changes: 12 additions & 1 deletion dapr/ext/workflow/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
"""

# Import your main classes here
from dapr.ext.workflow._durabletask.task import TaskFailedError
from dapr.ext.workflow._durabletask.task import FailureDetails, TaskFailedError
from dapr.ext.workflow.dapr_workflow_client import DaprWorkflowClient
from dapr.ext.workflow.dapr_workflow_context import DaprWorkflowContext, when_all, when_any
from dapr.ext.workflow.mcp import DaprMCPClient, MCPToolDef
Expand All @@ -28,6 +28,12 @@
)
from dapr.ext.workflow.retry_policy import RetryPolicy
from dapr.ext.workflow.workflow_activity_context import WorkflowActivityContext
from dapr.ext.workflow.workflow_management import (
UNSET,
WorkflowHistoryEvent,
WorkflowHistoryEventType,
WorkflowInstanceIdPage,
)
from dapr.ext.workflow.workflow_runtime import WorkflowRuntime, alternate_name
from dapr.ext.workflow.workflow_state import WorkflowState, WorkflowStatus

Expand All @@ -38,11 +44,16 @@
'WorkflowActivityContext',
'WorkflowState',
'WorkflowStatus',
'WorkflowHistoryEvent',
'WorkflowHistoryEventType',
'WorkflowInstanceIdPage',
'UNSET',
'when_all',
'when_any',
'alternate_name',
'RetryPolicy',
'TaskFailedError',
'FailureDetails',
'PropagationScope',
Comment on lines 31 to 57

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

can you pls confirm that we must expose all of these? which of these can be internal only?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I kept most of them, just moved to private this one RERUNNABLE_EVENT_TYPES

'PropagatedHistory',
'PropagationNotFoundError',
Expand Down
37 changes: 37 additions & 0 deletions dapr/ext/workflow/_durabletask/aio/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,8 @@
TOutput,
WorkflowIdReusePolicy,
WorkflowState,
_new_list_request,
_new_rerun_request,
_TransientTimeout,
new_orchestration_state,
new_task_router,
Expand Down Expand Up @@ -374,3 +376,38 @@ async def purge_orchestration(
)
self._logger.info(f"Purging instance '{instance_id}'.")
await self._get_stub().PurgeInstances(req)

async def list_instance_ids(
self, *, page_size: Optional[int] = None, continuation_token: Optional[str] = None
) -> pb.ListInstanceIDsResponse:
req = _new_list_request(page_size, continuation_token)
return await self._get_stub().ListInstanceIDs(req)

async def get_instance_history(self, instance_id: str) -> list[pb.HistoryEvent]:
req = pb.GetInstanceHistoryRequest(instanceId=instance_id)
res: pb.GetInstanceHistoryResponse = await self._get_stub().GetInstanceHistory(req)
return list(res.events)

async def rerun_orchestration_from_event(
self,
instance_id: str,
event_id: int,
*,
new_instance_id: Optional[str] = None,
input: Optional[Any] = None,
overwrite_input: bool = False,
new_child_instance_id: Optional[str] = None,
app_id: Optional[str] = None,
) -> str:
req = _new_rerun_request(
instance_id,
event_id,
new_instance_id=new_instance_id,
input=input,
overwrite_input=overwrite_input,
new_child_instance_id=new_child_instance_id,
app_id=app_id,
)
self._logger.info(f"Rerunning instance '{instance_id}' from event {event_id}.")
res: pb.RerunWorkflowFromEventResponse = await self._get_stub().RerunWorkflowFromEvent(req)
return res.newInstanceID
121 changes: 121 additions & 0 deletions dapr/ext/workflow/_durabletask/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,92 @@ def new_orchestration_state(
)


# eventID and pageSize are both uint32 on the wire; protobuf rejects anything
# outside this with a message naming neither the argument nor the reason.
_MAX_UINT32 = 2**32 - 1


def _new_list_request(
page_size: Optional[int], continuation_token: Optional[str]
) -> pb.ListInstanceIDsRequest:
"""Build a ListInstanceIDs request, rejecting page sizes the stores mishandle.

Raises:
ValueError: If page_size is given and is not between 1 and the uint32 the
wire allows. Zero is not "no limit": state stores disagree on it, and at
least one loops forever while another indexes out of range on it.
"""
if page_size is not None and not 1 <= page_size <= _MAX_UINT32:
raise ValueError(
f'page_size must be between 1 and {_MAX_UINT32}, got {page_size}. '
'Pass None rather than 0 to leave the limit to the runtime; a 0 reaches '
'the state store, which may return an endless page or fail on it.'
)

return pb.ListInstanceIDsRequest(pageSize=page_size, continuationToken=continuation_token)


def _new_rerun_request(
instance_id: str,
event_id: int,
*,
new_instance_id: Optional[str],
input: Optional[Any],
overwrite_input: bool,
new_child_instance_id: Optional[str],
app_id: Optional[str],
) -> pb.RerunWorkflowFromEventRequest:
"""Build a RerunWorkflowFromEvent request.

``input`` and ``overwrite_input`` mirror the wire, which needs both: the
replacement input rides in a non-optional ``StringValue``, so "leave the
input alone" and "replace it with null" are only distinguishable through the
flag. Callers of the public client express that as one argument; see
:data:`dapr.ext.workflow.UNSET`.

Raises:
ValueError: If event_id falls outside the uint32 range the wire allows,
or if either instance ID is an empty string. The runtime validates
neither: protobuf rejects an out-of-range event_id with a message naming
neither the argument nor the reason, and an empty ID is taken literally.
"""
if not 0 <= event_id <= _MAX_UINT32:
negative_hint = (
' The runtime reports -1 for history events it assigns no ID to, and those '
'cannot be rerun from.'
if event_id < 0
else ''
)
raise ValueError(
f'event_id must be between 0 and {_MAX_UINT32}, got {event_id}.{negative_hint}'
)

# Both ID fields have explicit presence: None leaves them unset and the runtime
# generates an ID, while '' sets them to empty and the runtime takes it literally,
# producing an instance it cannot schedule reminders for.
for name, value in (
('new_instance_id', new_instance_id),
('new_child_instance_id', new_child_instance_id),
):
if value == '':
raise ValueError(
f'{name} must be a non-empty ID or None, got an empty string. None asks '
'the runtime to generate one; an empty string is used as the ID itself.'
)

return pb.RerunWorkflowFromEventRequest(
sourceInstanceID=instance_id,
eventID=event_id,
newInstanceID=new_instance_id,
input=wrappers_pb2.StringValue(value=shared.to_json(input))
if overwrite_input and input is not None
else None,
overwriteInput=overwrite_input,
newChildWorkflowInstanceID=new_child_instance_id,
router=new_task_router(app_id),
)


class TaskHubGrpcClient:
def __init__(
self,
Expand Down Expand Up @@ -506,3 +592,38 @@ def purge_orchestration(
)
self._logger.info(f"Purging instance '{instance_id}'.")
self._stub.PurgeInstances(req)

def list_instance_ids(
self, *, page_size: Optional[int] = None, continuation_token: Optional[str] = None
) -> pb.ListInstanceIDsResponse:
req = _new_list_request(page_size, continuation_token)
return self._stub.ListInstanceIDs(req)

def get_instance_history(self, instance_id: str) -> list[pb.HistoryEvent]:
req = pb.GetInstanceHistoryRequest(instanceId=instance_id)
res: pb.GetInstanceHistoryResponse = self._stub.GetInstanceHistory(req)
return list(res.events)

def rerun_orchestration_from_event(
self,
instance_id: str,
event_id: int,
*,
new_instance_id: Optional[str] = None,
input: Optional[Any] = None,
overwrite_input: bool = False,
new_child_instance_id: Optional[str] = None,
app_id: Optional[str] = None,
) -> str:
req = _new_rerun_request(
instance_id,
event_id,
new_instance_id=new_instance_id,
input=input,
overwrite_input=overwrite_input,
new_child_instance_id=new_child_instance_id,
app_id=app_id,
)
self._logger.info(f"Rerunning instance '{instance_id}' from event {event_id}.")
res: pb.RerunWorkflowFromEventResponse = self._stub.RerunWorkflowFromEvent(req)
return res.newInstanceID
Loading
Loading