diff --git a/src/durable_workflow/client.py b/src/durable_workflow/client.py index e5565b3..a293f2c 100644 --- a/src/durable_workflow/client.py +++ b/src/durable_workflow/client.py @@ -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`.""" @@ -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: @@ -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})" diff --git a/tests/test_client.py b/tests/test_client.py index 6b2f9a3..12c13e6 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -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") @@ -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") diff --git a/tests/test_workflow_result_timeout.py b/tests/test_workflow_result_timeout.py index 2f98d47..1b17c55 100644 --- a/tests/test_workflow_result_timeout.py +++ b/tests/test_workflow_result_timeout.py @@ -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()