Skip to content

Commit 818a5f6

Browse files
test: prove Python local activity retry against Server
1 parent 26d611f commit 818a5f6

1 file changed

Lines changed: 106 additions & 1 deletion

File tree

‎tests/integration/test_local_activity.py‎

Lines changed: 106 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
"""Qualify Python local activity recording and cold replay against Server."""
2+
23
from __future__ import annotations
34

45
import uuid
@@ -29,6 +30,29 @@ async def local_greet(name: str) -> str:
2930
return f"hello, {name}"
3031

3132

33+
@workflow.defn(name="tests.python-local-retry")
34+
class LocalRetryWorkflow:
35+
def run(self, ctx: Any, name: str) -> Any:
36+
result = yield ctx.local_activity(
37+
"tests.python-local-retry-greet",
38+
[name],
39+
retry_policy={"max_attempts": 2, "backoff_seconds": [0]},
40+
)
41+
return {"greeting": result}
42+
43+
44+
_retry_executions = 0
45+
46+
47+
@activity.defn(name="tests.python-local-retry-greet")
48+
async def local_retry_greet(name: str) -> str:
49+
global _retry_executions
50+
_retry_executions += 1
51+
if _retry_executions == 1:
52+
raise RuntimeError("transient")
53+
return f"hello, {name}"
54+
55+
3256
@pytest.mark.asyncio
3357
async def test_local_activity_completion_survives_cold_replay(
3458
server_url: str,
@@ -82,7 +106,8 @@ async def test_local_activity_completion_survives_cold_replay(
82106
commands = await worker._run_workflow_task(task)
83107
assert commands is not None
84108
assert [command["type"] for command in commands] == [
85-
"record_local_activity", "complete_workflow",
109+
"record_local_activity",
110+
"complete_workflow",
86111
]
87112
assert _executions == 1
88113
assert await handle.result(timeout=10.0) == {"greeting": "hello, Ada"}
@@ -107,3 +132,83 @@ async def test_local_activity_completion_survives_cold_replay(
107132
assert _executions == 1
108133
finally:
109134
await client.deregister_worker_registration(worker.worker_id)
135+
136+
137+
@pytest.mark.asyncio
138+
async def test_local_activity_retries_commit_one_terminal_record_and_replay(
139+
server_url: str,
140+
server_token: str,
141+
) -> None:
142+
global _retry_executions
143+
_retry_executions = 0
144+
suffix = uuid.uuid4().hex[:8]
145+
queue = f"py-local-retry-{suffix}"
146+
workflow_id = f"py-local-retry-{suffix}"
147+
manifest = {
148+
**PORTABLE_WORKER_AFFINITY_CAPABILITY_MANIFEST,
149+
"local_activities": {
150+
"supported": True,
151+
"minimum_protocol_version": "1.18",
152+
"implementation": "record_local_activity",
153+
},
154+
}
155+
156+
async with Client(server_url, token=server_token, namespace="default") as client:
157+
worker = Worker(
158+
client,
159+
task_queue=queue,
160+
workflows=[LocalRetryWorkflow],
161+
activities=[local_retry_greet],
162+
worker_id=f"py-local-retry-worker-{suffix}",
163+
)
164+
await client.register_worker(
165+
worker_id=worker.worker_id,
166+
task_queue=queue,
167+
supported_workflow_types=list(worker.workflows),
168+
supported_activity_types=list(worker.activities),
169+
workflow_definition_fingerprints=worker.workflow_definition_fingerprints,
170+
workflow_command_contracts=worker.workflow_command_contracts,
171+
capabilities=["local_activities"],
172+
capability_manifest=manifest,
173+
)
174+
try:
175+
handle = await client.start_workflow(
176+
workflow_type="tests.python-local-retry",
177+
task_queue=queue,
178+
workflow_id=workflow_id,
179+
input=["Ada"],
180+
)
181+
task = await client.poll_workflow_task(
182+
worker_id=worker.worker_id,
183+
task_queue=queue,
184+
timeout=10.0,
185+
)
186+
assert task is not None
187+
commands = await worker._run_workflow_task(task)
188+
assert commands is not None
189+
assert [command["type"] for command in commands] == [
190+
"record_local_activity",
191+
"complete_workflow",
192+
]
193+
assert [attempt["outcome"] for attempt in commands[0]["attempts"]] == ["failed", "completed"]
194+
assert commands[0]["attempts"][0]["retry_reason"] == "failure"
195+
assert _retry_executions == 2
196+
assert await handle.result(timeout=10.0) == {"greeting": "hello, Ada"}
197+
198+
history = await handle.get_history()
199+
events = history.get("events", history.get("history_events", []))
200+
event_types = [event["event_type"] for event in events]
201+
assert event_types.count("ActivityCompleted") == 1
202+
assert event_types.count("WorkflowCompleted") == 1
203+
204+
outcome = replay(
205+
LocalRetryWorkflow,
206+
events,
207+
["Ada"],
208+
workflow_id=workflow_id,
209+
run_id=handle.run_id or "",
210+
)
211+
assert [command.__class__.__name__ for command in outcome.commands] == ["CompleteWorkflow"]
212+
assert _retry_executions == 2
213+
finally:
214+
await client.deregister_worker_registration(worker.worker_id)

0 commit comments

Comments
 (0)