diff --git a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py index 760d6016721..6405baef943 100644 --- a/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py +++ b/python/packages/foundry_hosting/agent_framework_foundry_hosting/_responses.py @@ -14,6 +14,7 @@ from agent_framework import ( ChatOptions, + CheckpointID, CheckpointStorage, Content, ContextProvider, @@ -24,10 +25,11 @@ SessionStore, SupportsAgentRun, WorkflowAgent, + WorkflowCheckpoint, ) from agent_framework._telemetry import mark_feature_used from agent_framework.exceptions import AgentFrameworkException -from azure.ai.agentserver.core import get_request_context +from azure.ai.agentserver.core import FoundryAgentRequestContext, get_request_context from azure.ai.agentserver.responses import ( ResponseContext, ResponseProviderProtocol, @@ -77,6 +79,39 @@ _HOSTED_RESPONSES_HISTORY_SOURCE_ID = "_foundry_responses_history" +class _CapturingCheckpointStorage: + """Delegate storage while retaining the latest successful checkpoint save. + + The workflow runner does not expose the checkpoint it creates. Capturing its + ID lets the host copy the exact persisted response state to the conversation + without rescanning the full checkpoint history. + """ + + def __init__(self, storage: CheckpointStorage) -> None: + self._storage = storage + self.latest_checkpoint_id: CheckpointID | None = None + + async def save(self, checkpoint: WorkflowCheckpoint) -> CheckpointID: + checkpoint_id = await self._storage.save(checkpoint) + self.latest_checkpoint_id = checkpoint_id + return checkpoint_id + + async def load(self, checkpoint_id: CheckpointID) -> WorkflowCheckpoint: + return await self._storage.load(checkpoint_id) + + async def list_checkpoints(self, *, workflow_name: str) -> list[WorkflowCheckpoint]: + return await self._storage.list_checkpoints(workflow_name=workflow_name) + + async def delete(self, checkpoint_id: CheckpointID) -> bool: + return await self._storage.delete(checkpoint_id) + + async def get_latest(self, *, workflow_name: str) -> WorkflowCheckpoint | None: + return await self._storage.get_latest(workflow_name=workflow_name) + + async def list_checkpoint_ids(self, *, workflow_name: str) -> list[CheckpointID]: + return await self._storage.list_checkpoint_ids(workflow_name=workflow_name) + + def _validate_checkpoint_context_id(context_id: str) -> None: """Validate that a checkpoint context ID is a single safe path component in case file-based storage is used.""" if ( @@ -334,8 +369,9 @@ async def _handle_inner_agent( ) -> AsyncIterable[ResponseStreamEvent | dict[str, Any]]: """Handle a regular agent with Responses-managed MAF session continuity. - Conversation mode reads and writes one MAF session snapshot under - ``conversation_id``. Response chaining reads the snapshot under + Conversation mode reads the latest MAF session snapshot under + ``conversation_id`` and writes each turn under both its immutable + ``response_id`` and the conversation ID. Response chaining reads the snapshot under ``previous_response_id`` and writes the updated session under the current ``response_id``, allowing branches without changing the MAF session's own identifier. Hosted storage uses the request user as its isolation boundary. @@ -387,6 +423,8 @@ async def _handle_inner_agent( ) previous_response_id = request.get("previous_response_id") + if previous_response_id is not None and context.conversation_id is not None: + raise RuntimeError("Previous response ID cannot be used in conjunction with conversation ID.") session_load_id = context.conversation_id or previous_response_id session = await session_storage.get(session_load_id) if session_load_id is not None else None if session is None: @@ -395,7 +433,6 @@ async def _handle_inner_agent( f"Cannot find an existing agent session for previous_response_id={previous_response_id}." ) session = self._agent.create_session() - session_save_id = context.conversation_id or context.response_id except Exception as ex: logger.error("Failed to prepare state storage: %s", ex, exc_info=(type(ex), ex, ex.__traceback__)) for event in self._emit_failure(response_event_stream, None, ex): @@ -455,17 +492,37 @@ async def _handle_inner_agent( finally: if self._uses_hosted_responses_history: session.state.pop(_HOSTED_RESPONSES_HISTORY_SOURCE_ID, None) - try: - await session_storage.set(session_save_id, session) - except Exception as save_error: - save_failure = save_error - if request_interrupted: - message = "Failed to persist the Agent Framework session while unwinding an interrupted request" - elif request_failure is not None: - message = "Failed to persist the Agent Framework session after an agent failure" - else: - message = "Failed to persist the Agent Framework session after a successful request" - logger.error(message, exc_info=(type(save_error), save_error, save_error.__traceback__)) + if request_interrupted: + message = "Failed to persist the Agent Framework session while unwinding an interrupted request" + elif request_failure is not None: + message = "Failed to persist the Agent Framework session after an agent failure" + else: + message = "Failed to persist the Agent Framework session after a successful request" + + save_errors: list[tuple[str, Exception]] = [] + session_save_ids = [("response snapshot", context.response_id)] + if context.conversation_id is not None: + session_save_ids.append(("conversation", context.conversation_id)) + for save_target, session_save_id in session_save_ids: + try: + await session_storage.set(session_save_id, session) + except Exception as save_error: + save_errors.append((save_target, save_error)) + logger.error( + "%s (%s)", + message, + save_target, + exc_info=(type(save_error), save_error, save_error.__traceback__), + ) + + if len(save_errors) == 1: + save_failure = save_errors[0][1] + elif save_errors: + details = "; ".join( + f"{save_target}: {str(save_error) or type(save_error).__name__}" + for save_target, save_error in save_errors + ) + save_failure = RuntimeError(f"Multiple session persistence operations failed: {details}") if request_failure is not None and save_failure is not None: failure = RuntimeError( @@ -517,11 +574,12 @@ async def _handle_inner_workflow( # any future async resources owned by the workflow are entered here. await self._ensure_agent_ready() - checkpoint_save_id = context.conversation_id or context.response_id - _validate_checkpoint_context_id(checkpoint_save_id) + # Persist each turn under its immutable response ID first. For conversation turns, + # the latest successfully saved checkpoint is copied to the conversation ID below. + _validate_checkpoint_context_id(context.response_id) checkpoint_storage = self._checkpoint_storage_provider.get_store( config=self.config, - context_id=checkpoint_save_id, + context_id=context.response_id, platform_context=request_context, ) @@ -542,7 +600,7 @@ async def _handle_inner_workflow( restore_checkpoint_storage = checkpoint_storage if checkpoint_load_id is not None: _validate_checkpoint_context_id(checkpoint_load_id) - if checkpoint_load_id != checkpoint_save_id: + if checkpoint_load_id != context.response_id: restore_checkpoint_storage = self._checkpoint_storage_provider.get_store( config=self.config, context_id=checkpoint_load_id, @@ -554,57 +612,123 @@ async def _handle_inner_workflow( f"Cannot find an existing workflow checkpoint for previous_response_id={previous_response_id}." ) - # Multi-turn pattern: when we have a prior checkpoint, restore it - # first (drive the workflow back to idle with prior state intact), - # then make a separate call that delivers the new user input. This - # depends on Workflow.run preserving shared state across calls. The - # restore-only call may yield events from any pending in-flight - # work in the checkpoint; we consume those internally here so they - # don't surface to the response stream as duplicates. - # - # If the restored checkpoint had pending request_info events, the - # restore-only call replays them through - # ``WorkflowAgent._convert_workflow_event_to_agent_response_updates`` - # and populates ``self._agent.pending_requests``. That is the correct - # state: those requests are genuinely outstanding, and the next - # ``run(input_messages, ...)`` call may contain ``function_call_output`` - # items (carried as FunctionResult/FunctionApprovalResponse content) - # that fulfill them via :meth:`WorkflowAgent._process_pending_requests`. - if latest_checkpoint is not None: - async for _ in self._agent.run( + request_failure: Exception | None = None + save_failure: Exception | None = None + request_interrupted = False + write_checkpoint_storage = _CapturingCheckpointStorage(checkpoint_storage) + + try: + # Multi-turn pattern: when we have a prior checkpoint, restore it + # first (drive the workflow back to idle with prior state intact), + # then make a separate call that delivers the new user input. This + # depends on Workflow.run preserving shared state across calls. The + # restore-only call may yield events from any pending in-flight + # work in the checkpoint; we consume those internally here so they + # don't surface to the response stream as duplicates. + # + # If the restored checkpoint had pending request_info events, the + # restore-only call replays them through + # ``WorkflowAgent._convert_workflow_event_to_agent_response_updates`` + # and populates ``self._agent.pending_requests``. That is the correct + # state: those requests are genuinely outstanding, and the next + # ``run(input_messages, ...)`` call may contain ``function_call_output`` + # items (carried as FunctionResult/FunctionApprovalResponse content) + # that fulfill them via :meth:`WorkflowAgent._process_pending_requests`. + if latest_checkpoint is not None: + async for _ in self._agent.run( + stream=True, + checkpoint_id=latest_checkpoint.checkpoint_id, + checkpoint_storage=restore_checkpoint_storage, + ): + pass + + tracker = _OutputItemTracker(response_event_stream) + + # Run the workflow agent in streaming mode with the new user input. + async for update in self._agent.run( + input_messages, stream=True, - checkpoint_id=latest_checkpoint.checkpoint_id, - checkpoint_storage=restore_checkpoint_storage, + checkpoint_storage=write_checkpoint_storage, ): - pass - - tracker = _OutputItemTracker(response_event_stream) - - # Run the workflow agent in streaming mode with the new user input. - async for update in self._agent.run( - input_messages, - stream=True, - checkpoint_storage=checkpoint_storage, - ): - for content in update.contents: - for event in tracker.handle(content): - yield event - if tracker.needs_async: - async for item in _to_outputs( - response_event_stream, content, approval_storage=approval_storage - ): - yield item - tracker.needs_async = False + for content in update.contents: + for event in tracker.handle(content): + yield event + if tracker.needs_async: + async for item in _to_outputs( + response_event_stream, content, approval_storage=approval_storage + ): + yield item + tracker.needs_async = False + + # Close any remaining active builder + for event in tracker.close(): + yield event + except (asyncio.CancelledError, GeneratorExit): + request_interrupted = True + raise + except Exception as ex: + request_failure = ex + logger.error( + "Failed to produce response for workflow agent", + exc_info=(type(ex), ex, ex.__traceback__), + ) + finally: + try: + if ( + write_checkpoint_storage.latest_checkpoint_id is not None + and context.conversation_id is not None + ): + checkpoint = await write_checkpoint_storage.load(write_checkpoint_storage.latest_checkpoint_id) + await self._copy_workflow_checkpoint_to_conversation( + checkpoint, + conversation_id=context.conversation_id, + platform_context=request_context, + ) + except Exception as save_error: + save_failure = save_error + if request_interrupted: + message = "Failed to persist the workflow checkpoint while unwinding an interrupted request" + elif request_failure is not None: + message = "Failed to persist the workflow checkpoint after a workflow failure" + else: + message = "Failed to persist the workflow checkpoint after a successful request" + logger.error(message, exc_info=(type(save_error), save_error, save_error.__traceback__)) - # Close any remaining active builder - for event in tracker.close(): - yield event - yield response_event_stream.emit_completed() + if request_failure is not None and save_failure is not None: + failure = RuntimeError( + f"Workflow request failed: {str(request_failure) or type(request_failure).__name__}; " + f"checkpoint persistence also failed: {str(save_failure) or type(save_failure).__name__}" + ) + for event in self._emit_failure(response_event_stream, tracker, failure): + yield event + elif request_failure is not None: + for event in self._emit_failure(response_event_stream, tracker, request_failure): + yield event + elif save_failure is not None: + for event in self._emit_failure(response_event_stream, tracker, save_failure): + yield event + else: + yield response_event_stream.emit_completed() except Exception as ex: logger.exception("Failed to produce response for workflow agent") for event in self._emit_failure(response_event_stream, tracker, ex): yield event + async def _copy_workflow_checkpoint_to_conversation( + self, + checkpoint: WorkflowCheckpoint, + *, + conversation_id: str, + platform_context: FoundryAgentRequestContext, + ) -> None: + """Copy a response turn's latest workflow checkpoint to its conversation.""" + conversation_storage = self._checkpoint_storage_provider.get_store( + config=self.config, + context_id=conversation_id, + platform_context=platform_context, + ) + await conversation_storage.save(checkpoint) + @staticmethod def _emit_failure( response_event_stream: ResponseEventStream, diff --git a/python/packages/foundry_hosting/tests/test_responses.py b/python/packages/foundry_hosting/tests/test_responses.py index d36bf248eb2..0343dcf04f7 100644 --- a/python/packages/foundry_hosting/tests/test_responses.py +++ b/python/packages/foundry_hosting/tests/test_responses.py @@ -30,6 +30,7 @@ ChatMiddlewareLayer, ChatResponse, ChatResponseUpdate, + CheckpointStorage, Content, FunctionInvocationLayer, HistoryProvider, @@ -42,6 +43,7 @@ SupportsAgentRun, WorkflowAgent, WorkflowBuilder, + WorkflowCheckpoint, WorkflowContext, executor, tool, @@ -261,6 +263,22 @@ async def set(self, session_id: str, session: AgentSession) -> None: raise OSError("session storage is full") +class _FailingResponseSnapshotStore(SessionStore): + def __init__(self, *, fail_conversation: bool = False) -> None: + super().__init__() + self.fail_conversation = fail_conversation + self.set_attempts: list[str] = [] + + async def set(self, session_id: str, session: AgentSession) -> None: + self.set_attempts.append(session_id) + if session_id == "conversation-1": + if self.fail_conversation: + raise OSError("conversation storage is full") + await super().set(session_id, session) + return + raise OSError("response snapshot storage is full") + + _SESSION_STORE_UNSET = object() @@ -404,6 +422,23 @@ async def save_messages( with pytest.raises(RuntimeError, match="history provider"): ResponsesHostServer(agent) + async def test_previous_response_rejected_with_conversation(self) -> None: + agent = _make_agent() + server = _make_server(agent) + + response = await _post( + server, + previous_response_id="caresp_aaaaaaaaaaaaaaaa00" + "1" * 32, + conversation_id="conversation-1", + ) + + assert response.status_code == 200 + body = response.json() + assert body["status"] == "failed" + assert body["error"]["message"] == "Previous response ID cannot be used in conjunction with conversation ID." + agent.run.assert_not_called() + agent.create_session.assert_not_called() + async def test_previous_response_requires_existing_agent_session(self) -> None: agent = _make_agent() server = _make_server(agent, session_store=SessionStore()) @@ -430,6 +465,45 @@ async def test_previous_response_requires_existing_agent_session(self) -> None: class TestAgentSessionPersistence: + async def test_conversation_response_snapshots_support_branching(self) -> None: + seen_counts: list[int] = [] + + def run_with_state(*args: Any, **kwargs: Any) -> ResponseStream[AgentResponseUpdate, AgentResponse]: + del args + session = kwargs["session"] + assert isinstance(session, AgentSession) + count = int(session.state.get("turn_count", 0)) + 1 + session.state["turn_count"] = count + seen_counts.append(count) + + async def updates() -> AsyncIterator[AgentResponseUpdate]: + yield AgentResponseUpdate(contents=[Content.from_text(f"turn {count}")], role="assistant") + + return ResponseStream(updates(), finalizer=AgentResponse.from_updates) + + agent = _make_agent() + agent.run = MagicMock(side_effect=run_with_state) + store = SessionStore() + server = _make_server(agent, session_store=store) + + first = await _post(server, input_text="first", conversation_id="conversation-1") + await _post(server, input_text="second", conversation_id="conversation-1") + branch = await _post(server, input_text="branch", previous_response_id=first.json()["id"]) + + assert first.status_code == 200 + assert branch.status_code == 200 + assert seen_counts == [1, 2, 2] + + conversation_snapshot = await store.get("conversation-1") + first_snapshot = await store.get(first.json()["id"]) + branch_snapshot = await store.get(branch.json()["id"]) + assert conversation_snapshot is not None + assert first_snapshot is not None + assert branch_snapshot is not None + assert conversation_snapshot.state["turn_count"] == 2 + assert first_snapshot.state["turn_count"] == 1 + assert branch_snapshot.state["turn_count"] == 2 + async def test_previous_response_chain_restores_session_state(self) -> None: seen_counts: list[int] = [] seen_session_ids: list[str] = [] @@ -611,6 +685,44 @@ async def updates() -> AsyncIterator[AgentResponseUpdate]: assert "response.completed" not in event_types assert store.set_attempts == 1 + async def test_response_snapshot_failure_still_persists_conversation(self) -> None: + store = _FailingResponseSnapshotStore() + agent = _make_agent() + + def run(*_args: Any, **kwargs: Any) -> ResponseStream[AgentResponseUpdate, AgentResponse]: + session = kwargs["session"] + assert isinstance(session, AgentSession) + + async def updates() -> AsyncIterator[AgentResponseUpdate]: + session.state["run_complete"] = True + yield AgentResponseUpdate(contents=[Content.from_text("done")], role="assistant") + + return ResponseStream(updates(), finalizer=AgentResponse.from_updates) + + agent.run = MagicMock(side_effect=run) + server = _make_server(agent, session_store=store) + + response = await _post(server, conversation_id="conversation-1") + conversation = await store.get("conversation-1") + + assert response.json()["status"] == "failed" + assert response.json()["id"] == store.set_attempts[0] + assert store.set_attempts == [response.json()["id"], "conversation-1"] + assert conversation is not None + assert conversation.state["run_complete"] is True + + async def test_response_and_conversation_save_failures_are_both_reported(self) -> None: + store = _FailingResponseSnapshotStore(fail_conversation=True) + server = _make_server(_make_agent(), session_store=store) + + response = await _post(server, conversation_id="conversation-1") + error_message = response.json()["error"]["message"] + + assert response.json()["status"] == "failed" + assert "response snapshot: response snapshot storage is full" in error_message + assert "conversation: conversation storage is full" in error_message + assert store.set_attempts == [response.json()["id"], "conversation-1"] + async def test_run_and_save_failure_emit_one_combined_failure( self, caplog: pytest.LogCaptureFixture, @@ -4249,6 +4361,180 @@ async def test_basic_text_response_streaming(self) -> None: text_done = [e for e in events if e["event"] == "response.output_text.done"] assert any(e["data"]["text"] == "hello stream" for e in text_done) + @pytest.mark.parametrize("stream", [False, True]) + async def test_conversation_response_checkpoints_support_branching(self, stream: bool) -> None: + @executor + async def count_turns(messages: list[Message], ctx: WorkflowContext[Any, AgentResponse]) -> None: + del messages + turn_count = int(ctx.get_state("turn_count", 0)) + 1 + ctx.set_state("turn_count", turn_count) + await ctx.yield_output( + AgentResponse(messages=[Message("assistant", [Content.from_text(f"turn {turn_count}")])]) + ) + + workflow_agent = WorkflowAgent( + workflow=WorkflowBuilder(start_executor=count_turns).build(), + name="Counting Workflow Agent", + ) + server = _make_server(workflow_agent) + + def response_body(response: httpx.Response) -> dict[str, Any]: + if not stream: + return response.json() + return _parse_sse_events(response.text)[-1]["data"]["response"] + + first = await _post(server, input_text="first", conversation_id="conversation-1", stream=stream) + second = await _post(server, input_text="second", conversation_id="conversation-1", stream=stream) + first_body = response_body(first) + branch = await _post( + server, + input_text="branch", + previous_response_id=first_body["id"], + stream=stream, + ) + + assert first.status_code == 200 + assert second.status_code == 200 + assert branch.status_code == 200 + assert first_body["status"] == "completed" + assert response_body(second)["status"] == "completed" + assert response_body(branch)["status"] == "completed" + branch_text = [ + part["text"] + for item in response_body(branch)["output"] + if item["type"] == "message" + for part in item.get("content", []) + if part["type"] == "output_text" + ] + assert branch_text == ["turn 2"] + + @pytest.mark.parametrize( + ("termination", "save_new_checkpoint", "copy_failure"), + [ + pytest.param("failure", True, False, id="failure"), + pytest.param("failure", True, True, id="request-and-copy-failure"), + pytest.param("success", True, True, id="copy-failure"), + pytest.param("cancel", True, False, id="cancel"), + pytest.param("cancel", True, True, id="cancel-and-copy-failure"), + pytest.param("close", True, False, id="close"), + pytest.param("close", True, True, id="close-and-copy-failure"), + pytest.param("failure", False, False, id="failure-before-checkpoint"), + ], + ) + async def test_incomplete_conversation_workflow_snapshots_only_new_checkpoints( + self, + termination: str, + save_new_checkpoint: bool, + copy_failure: bool, + caplog: pytest.LogCaptureFixture, + ) -> None: + workflow_agent = _build_text_workflow_agent("ignored") + checkpoint = WorkflowCheckpoint( + workflow_name=workflow_agent.workflow.name, + graph_signature_hash="hash", + state={"nested": {"value": "saved"}}, + ) + server = _make_server(workflow_agent) + + if not save_new_checkpoint: + conversation_storage = server._checkpoint_storage_provider.get_store( # pyright: ignore[reportPrivateUsage] + config=server.config, + context_id="conversation-1", + platform_context=get_request_context(), + ) + await conversation_storage.save( + WorkflowCheckpoint( + workflow_name=workflow_agent.workflow.name, + graph_signature_hash="hash", + ) + ) + + async def updates(checkpoint_storage: CheckpointStorage) -> AsyncIterator[AgentResponseUpdate]: + if save_new_checkpoint: + await checkpoint_storage.save(checkpoint) + checkpoint.state["nested"]["value"] = "mutated" + if termination == "failure": + raise RuntimeError("workflow failed") + yield AgentResponseUpdate(contents=[Content.from_text("started")], role="assistant") + if termination in {"cancel", "close"}: + await asyncio.Event().wait() + + def run(*args: Any, **kwargs: Any) -> AsyncIterator[AgentResponseUpdate]: + del args + return updates(kwargs["checkpoint_storage"]) + + request = CreateResponse(model="m", input="hi", stream=True) + context = ResponseContext( + response_id="response-1", + conversation_id="conversation-1", + mode_flags=MagicMock(), + ) + copy_checkpoint = ( + AsyncMock(side_effect=RuntimeError("copy failed")) + if copy_failure + else AsyncMock(wraps=server._copy_workflow_checkpoint_to_conversation) # pyright: ignore[reportPrivateUsage] + ) + + with ( + patch.object(ResponseContext, "get_input_items", new=AsyncMock(return_value=[])), + patch.object(workflow_agent, "run", side_effect=run), + patch.object(server, "_copy_workflow_checkpoint_to_conversation", new=copy_checkpoint), + ): + if termination in {"failure", "success"}: + events = [ + event + async for event in server._handle_inner_workflow( # pyright: ignore[reportPrivateUsage] + request, + context, + ) + ] + assert events[-1].get("type") == "response.failed" + failed_event = cast(Mapping[str, Any], events[-1]) + response = cast(Mapping[str, Any], failed_event["response"]) + error = cast(Mapping[str, Any], response["error"]) + error_message = str(error["message"]) + if termination == "failure": + assert "workflow failed" in error_message + if copy_failure: + assert "copy failed" in error_message + else: + handler = cast( + AsyncGenerator[Any, None], + server._handle_inner_workflow(request, context), # pyright: ignore[reportPrivateUsage] + ) + await anext(handler) + await anext(handler) + await anext(handler) + if termination == "close": + await handler.aclose() + else: + with pytest.raises(asyncio.CancelledError): + await handler.athrow(asyncio.CancelledError()) + + response_storage = server._checkpoint_storage_provider.get_store( # pyright: ignore[reportPrivateUsage] + config=server.config, + context_id="response-1", + platform_context=get_request_context(), + ) + latest = await response_storage.get_latest(workflow_name=workflow_agent.workflow.name) + if save_new_checkpoint: + assert latest is not None + assert latest.checkpoint_id == checkpoint.checkpoint_id + assert latest.state["nested"]["value"] == "saved" + if not copy_failure: + conversation_storage = server._checkpoint_storage_provider.get_store( # pyright: ignore[reportPrivateUsage] + config=server.config, + context_id="conversation-1", + platform_context=get_request_context(), + ) + conversation_latest = await conversation_storage.get_latest(workflow_name=workflow_agent.workflow.name) + assert conversation_latest is not None + assert conversation_latest.to_dict() == latest.to_dict() + else: + assert latest is None + if termination in {"cancel", "close"} and copy_failure: + assert "while unwinding an interrupted request" in caplog.text + async def test_non_streaming_emits_mcp_approval_request_and_persists_to_storage(self) -> None: workflow_agent, mock_agent = _build_approval_workflow_agent(approval_request_id="apr_wf_ns") server = _make_server(workflow_agent)