Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
104 changes: 70 additions & 34 deletions src/durable_workflow/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -1271,11 +1271,21 @@ async def describe(self) -> WorkflowExecution:
"""Return the server's current view of this workflow. See :meth:`Client.describe_workflow`."""
return await self._client.describe_workflow(self.workflow_id)

async def get_history(self) -> Any:
"""Fetch this run's durable history. See :meth:`Client.get_history`."""
async def get_history(
self,
*,
page_size: int | None = None,
next_page_token: str | None = None,
) -> Any:
"""Fetch one page of this run's history. See :meth:`Client.get_history`."""
if self.run_id is None:
raise ValueError("run_id is required to fetch workflow history from a handle")
return await self._client.get_history(self.workflow_id, self.run_id)
options: dict[str, Any] = {}
if page_size is not None:
options["page_size"] = page_size
if next_page_token is not None:
options["next_page_token"] = next_page_token
return await self._client.get_history(self.workflow_id, self.run_id, **options)

async def export_history(self) -> Any:
"""Export this run's history as a replay bundle. See :meth:`Client.export_history`."""
Expand Down Expand Up @@ -3462,10 +3472,26 @@ async def list_workflows(
next_page_token=data.get("next_page_token"),
)

async def get_history(self, workflow_id: str, run_id: str) -> Any:
"""Fetch the full durable history for one specific run of a workflow."""
async def get_history(
self,
workflow_id: str,
run_id: str,
*,
page_size: int | None = None,
next_page_token: str | None = None,
) -> Any:
"""Fetch one history page; pass its opaque next_page_token for the next page."""
params: dict[str, str] = {}
if page_size is not None:
params["page_size"] = str(page_size)
if next_page_token is not None:
params["next_page_token"] = next_page_token
query = urlencode(params)
path = f"/workflows/{workflow_id}/runs/{run_id}/history"
if query:
path += f"?{query}"
return await self._request(
"GET", f"/workflows/{workflow_id}/runs/{run_id}/history", context=workflow_id
"GET", path, context=workflow_id
)

async def export_history(self, workflow_id: str, run_id: str) -> Any:
Expand Down Expand Up @@ -4240,34 +4266,44 @@ async def get_result(
run_id = handle.run_id or desc.run_id
if run_id is None:
raise WorkflowFailed("no run_id available to fetch history")
history = await self.get_history(handle.workflow_id, run_id)
events = history.get("events", [])
for ev in reversed(events):
etype = ev.get("event_type")
payload = ev.get("payload") or {}
if etype == "WorkflowCompleted":
return serializer.decode_envelope(
payload.get("output"),
codec=payload.get("payload_codec") or desc.payload_codec,
external_storage=self.external_storage,
external_storage_cache=self.external_storage_cache,
)
if etype == "WorkflowFailed":
raise WorkflowFailed(
payload.get("message", "workflow failed"),
payload.get("exception_class"),
)
if etype == "WorkflowTerminated":
raise WorkflowTerminated(
payload.get("reason", "workflow was terminated")
)
if etype == "WorkflowCancelled":
raise WorkflowCancelled(
payload.get("reason", "workflow was cancelled")
)
if etype == "WorkflowTimedOut":
raise WorkflowTimedOut()
return None
page_token: str | None = None
seen_tokens: set[str] = set()
while True:
history = await self.get_history(
handle.workflow_id, run_id, page_size=1000, next_page_token=page_token
)
events = history.get("events", [])
for ev in reversed(events):
etype = ev.get("event_type")
payload = ev.get("payload") or {}
if etype == "WorkflowCompleted":
return serializer.decode_envelope(
payload.get("output"),
codec=payload.get("payload_codec") or desc.payload_codec,
external_storage=self.external_storage,
external_storage_cache=self.external_storage_cache,
)
if etype == "WorkflowFailed":
raise WorkflowFailed(
payload.get("message", "workflow failed"),
payload.get("exception_class"),
)
if etype == "WorkflowTerminated":
raise WorkflowTerminated(
payload.get("reason", "workflow was terminated")
)
if etype == "WorkflowCancelled":
raise WorkflowCancelled(
payload.get("reason", "workflow was cancelled")
)
if etype == "WorkflowTimedOut":
raise WorkflowTimedOut()
page_token = history.get("next_page_token")
if not page_token:
return None
if page_token in seen_tokens:
raise RuntimeError("history pagination repeated a page token")
seen_tokens.add(page_token)
if asyncio.get_running_loop().time() > deadline:
raise TimeoutError(
f"workflow {handle.workflow_id} not terminal after {timeout}s (status={status})"
Expand Down
47 changes: 47 additions & 0 deletions tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -791,6 +791,17 @@ async def test_describe_run_request_matches_polyglot_fixture(self, client: Clien


class TestWorkflowHandleControlPlane:
@pytest.mark.asyncio
async def test_history_page_token_is_encoded(self, client: Client) -> None:
response = _mock_response(200, {"events": [], "next_page_token": None})
with patch.object(client._http, "request", new_callable=AsyncMock, return_value=response) as request:
page = await client.get_history("wf-1", "run-1", page_size=1000, next_page_token="opaque+/=")

assert page["events"] == []
assert request.call_args.args[1] == (
"/api/workflows/wf-1/runs/run-1/history?page_size=1000&next_page_token=opaque%2B%2F%3D"
)

@pytest.mark.asyncio
async def test_run_visibility_delegates_to_client(self, client: Client) -> None:
handle = WorkflowHandle(client, workflow_id="wf-1", run_id="r1", workflow_type="greeter")
Expand Down Expand Up @@ -3181,6 +3192,42 @@ async def test_update_request_matches_polyglot_fixture(self, client: Client) ->


class TestGetResult:
@pytest.mark.asyncio
async def test_completed_result_on_later_history_page(self, client: Client) -> None:
handle = WorkflowHandle(client, workflow_id="wf-1", run_id="run-1", workflow_type="greeter")
client.describe_workflow = AsyncMock(
return_value=WorkflowExecution(
workflow_id="wf-1", run_id="run-1", workflow_type="greeter", status="completed",
)
)
client.get_history = AsyncMock(side_effect=[
{"events": [{"event_type": "WorkflowStarted"}], "next_page_token": "page two"},
{"events": [{"event_type": "WorkflowCompleted", "payload": {
"output": serializer.encode(42, codec="avro"), "payload_codec": "avro",
}}], "next_page_token": None},
])

assert await client.get_result(handle) == 42
assert client.get_history.await_args_list[0].args == ("wf-1", "run-1")
assert client.get_history.await_args_list[0].kwargs == {"page_size": 1000, "next_page_token": None}
assert client.get_history.await_args_list[1].kwargs == {
"page_size": 1000, "next_page_token": "page two",
}

@pytest.mark.asyncio
async def test_repeated_history_page_token_fails(self, client: Client) -> None:
handle = WorkflowHandle(client, workflow_id="wf-1", run_id="run-1", workflow_type="greeter")
client.describe_workflow = AsyncMock(
return_value=WorkflowExecution(
workflow_id="wf-1", run_id="run-1", workflow_type="greeter", status="completed",
)
)
client.get_history = AsyncMock(return_value={"events": [], "next_page_token": "same"})

with pytest.raises(RuntimeError, match="repeated a page token"):
await client.get_result(handle)
assert client.get_history.await_count == 2

@pytest.mark.asyncio
async def test_completed_result_uses_event_payload_codec(self, client: Client) -> None:
handle = WorkflowHandle(client, workflow_id="wf-1", run_id="run-1", workflow_type="greeter")
Expand Down
4 changes: 3 additions & 1 deletion tests/test_workflow_result_timeout.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,9 @@ async def test_persisted_deadline_raises_typed_timeout_for_selected_history(kind
try:
with pytest.raises(WorkflowTimedOut, match="workflow execution timed out"):
await client.get_result(handle, timeout=0)
client.get_history.assert_awaited_once_with("order", "selected-run")
client.get_history.assert_awaited_once_with(
"order", "selected-run", page_size=1000, next_page_token=None
)
finally:
await client.aclose()

Expand Down
Loading