Skip to content

Commit fa91ca8

Browse files
authored
Merge pull request #69 from durable-workflow/feat/python-local-activity-execution
Execute Python local activities with durable replay
2 parents efeede8 + 8e58dc7 commit fa91ca8

18 files changed

Lines changed: 2335 additions & 72 deletions

‎docker-compose.test.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ services:
2929
CACHE_STORE: redis
3030
WORKFLOW_SERVER_AUTH_DRIVER: token
3131
WORKFLOW_SERVER_AUTH_TOKEN: "test-token"
32+
DW_WORKFLOW_TASK_TIMEOUT: 10
3233
depends_on:
3334
mysql:
3435
condition: service_healthy

‎docs/sdk-reference.md‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -185,6 +185,72 @@ receipt = yield ctx.start_child_workflow(
185185
)
186186
```
187187

188+
## Local activities
189+
190+
Use a local activity for short work that should run in the workflow worker
191+
process without a separate activity task. Register it as an activity on the
192+
same `Worker`, then yield `ctx.local_activity(...)` from the workflow:
193+
194+
```python
195+
import asyncio
196+
import os
197+
from uuid import uuid4
198+
199+
from durable_workflow import Client, Worker, activity, workflow
200+
201+
202+
@activity.defn(name="orders.normalize")
203+
def normalize(order_id: str) -> str:
204+
return order_id.strip().upper()
205+
206+
207+
@workflow.defn(name="orders.prepare")
208+
class PrepareOrder:
209+
def run(self, ctx, order_id):
210+
normalized = yield ctx.local_activity(
211+
"orders.normalize",
212+
[order_id],
213+
retry_policy={"max_attempts": 2, "backoff_seconds": [1]},
214+
start_to_close_timeout=5,
215+
)
216+
return {"order_id": normalized}
217+
218+
219+
async def main():
220+
async with Client(
221+
os.getenv("DURABLE_WORKFLOW_RUNTIME_URL", "http://127.0.0.1:8080"),
222+
token=os.getenv("DURABLE_WORKFLOW_TOKEN", "dev-token-123"),
223+
namespace="default",
224+
) as client:
225+
worker = Worker(
226+
client,
227+
task_queue="orders",
228+
workflows=[PrepareOrder],
229+
activities=[normalize],
230+
)
231+
handle = await client.start_workflow(
232+
workflow_type="orders.prepare",
233+
workflow_id=f"order-{uuid4().hex}",
234+
task_queue="orders",
235+
input=[" a-123 "],
236+
)
237+
await worker.run_until(workflow_id=handle.workflow_id, timeout=30.0)
238+
print(await handle.result(timeout=10.0))
239+
240+
241+
asyncio.run(main())
242+
```
243+
244+
Use the Server setup from the [quickstart](index.md#first-workflow), or replace
245+
the URL and credentials with those of a provisioned Cloud namespace. The worker
246+
records the local result with its workflow task. Once Server accepts that
247+
record, replay reads the result instead of running the handler again. If the
248+
worker loses its lease or crashes before the record is committed, the handler
249+
may run again on a replacement worker. Make side effects idempotent and keep
250+
local attempts short; use `ctx.schedule_activity(...)` for separately queued
251+
work. Local activities are yielded one at a time; they are not members of a
252+
parallel workflow group.
253+
188254
## Deterministic parallel groups
189255

190256
Yield a list to schedule one durable parallel barrier. Lists can nest and mix

‎pyproject.toml‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
44

55
[project]
66
name = "durable-workflow"
7-
version = "2.1.2"
7+
version = "2.2.0"
88
description = "Python client and worker SDK for Durable Workflow Cloud and self-hosted Server"
99
readme = "README.md"
1010
requires-python = ">=3.10"
@@ -71,8 +71,8 @@ durable-workflow-replay-conformance = "durable_workflow.replay_conformance:main"
7171
durable-workflow-workflow-updates-conformance = "durable_workflow.workflow_updates_conformance:main"
7272

7373
[tool.durable-workflow]
74-
product-train = "2.1.2"
75-
registry-version = "2.1.2"
74+
product-train = "2.2.0"
75+
registry-version = "2.2.0"
7676
supported-server-versions = "2.4.0"
7777
worker-protocol-version = "1.19"
7878
control-plane-version = "2"

‎src/durable_workflow/activity.py‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
from __future__ import annotations
33

44
import contextvars
5-
from collections.abc import Callable
5+
from collections.abc import Awaitable, Callable
66
from dataclasses import dataclass
77
from typing import TYPE_CHECKING, Any
88

@@ -38,9 +38,11 @@ def __init__(
3838
*,
3939
info: ActivityInfo,
4040
client: Client,
41+
heartbeat_callback: Callable[[dict[str, Any] | None], Awaitable[None]] | None = None,
4142
) -> None:
4243
self._info = info
4344
self._client = client
45+
self._heartbeat_callback = heartbeat_callback
4446
self._cancel_requested = False
4547

4648
@property
@@ -65,6 +67,13 @@ async def heartbeat(self, details: dict[str, Any] | None = None) -> None:
6567
owning workflow has requested cancellation, so the activity can exit
6668
cleanly at its next natural break point.
6769
"""
70+
if self._heartbeat_callback is not None:
71+
try:
72+
await self._heartbeat_callback(details)
73+
except ActivityCancelled:
74+
self._cancel_requested = True
75+
raise
76+
return
6877
resp = await self._client.heartbeat_activity_task(
6978
task_id=self._info.task_id,
7079
activity_attempt_id=self._info.activity_attempt_id,

‎src/durable_workflow/client.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -66,9 +66,9 @@
6666
CONTROL_PLANE_VERSION = "2"
6767
PORTABLE_WORKER_AFFINITY_CAPABILITY_MANIFEST: dict[str, dict[str, str | bool]] = {
6868
"local_activities": {
69-
"supported": False,
69+
"supported": True,
7070
"minimum_protocol_version": "1.18",
71-
"reason": "python_worker_does_not_execute_record_local_activity",
71+
"implementation": "record_local_activity",
7272
},
7373
"worker_sessions": {
7474
"supported": False,

‎src/durable_workflow/external_storage.py‎

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,12 @@
1212
from typing import Any, Protocol
1313
from urllib.parse import quote, unquote, urlparse
1414

15-
from .errors import ExternalPayloadIntegrityMismatch, ExternalPayloadUnsupported
15+
from .errors import (
16+
ExternalPayloadError,
17+
ExternalPayloadIntegrityMismatch,
18+
ExternalPayloadUnavailable,
19+
ExternalPayloadUnsupported,
20+
)
1621

1722
EXTERNAL_PAYLOAD_REFERENCE_SCHEMA = "durable-workflow.v2.external-payload-reference.v1"
1823
RUNTIME_EXTERNAL_PAYLOAD_REFERENCE_SCHEMA = (
@@ -556,7 +561,12 @@ def store_external_payload(
556561
if expires_at is not None:
557562
_validate_rfc3339(expires_at)
558563
sha256 = hashlib.sha256(data).hexdigest()
559-
uri = driver.put(data, sha256=sha256, codec=codec)
564+
try:
565+
uri = driver.put(data, sha256=sha256, codec=codec)
566+
except ExternalPayloadError:
567+
raise
568+
except Exception as exc:
569+
raise ExternalPayloadUnavailable("external payload storage upload failed") from exc
560570
return ExternalPayloadReference(
561571
uri=uri,
562572
sha256=sha256,

0 commit comments

Comments
 (0)