From 956f43f046c78479ba207b17591af9e212e38299 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Tue, 25 Aug 2026 19:23:08 +0530 Subject: [PATCH 1/9] fix(workflow): discard pending state after failed supersteps --- .../agent_framework/_workflows/_runner.py | 4 ++ .../core/tests/workflow/test_workflow.py | 49 +++++++++++++++++++ 2 files changed, 53 insertions(+) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index ac5558dbe1c..519a76eace5 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -139,6 +139,8 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: iteration_task.cancel() with contextlib.suppress(asyncio.CancelledError): await iteration_task + # Discard pending state writes from the cancelled superstep + self._state.discard() raise # Propagate errors from iteration, but first surface any pending events @@ -149,6 +151,8 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: if await self._ctx.has_events(): for event in await self._ctx.drain_events(): yield event + # Discard pending state writes from the failed superstep + self._state.discard() raise self._iteration += 1 diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index 2f672f591d3..f1fb109129f 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -609,6 +609,55 @@ def _build(): assert result2.get_outputs()[0] == ["run2:message2"] +@dataclass +class FlakyMessage: + """A message that can fail on demand for testing state discard behavior.""" + + fail: bool + + +class FlakyStateExecutor(Executor): + """An executor that fails on demand to test state discard on failure.""" + + @handler + async def handle_message( + self, + message: FlakyMessage, + ctx: WorkflowContext[FlakyMessage, str], + ) -> None: + if message.fail: + ctx.set_state("secret", "leaked-from-failed-run") + raise RuntimeError("simulated transient failure") + + await ctx.yield_output("ok") + + +async def test_workflow_discards_pending_state_after_failed_superstep(): + """Test that pending state from a failed superstep is discarded and not committed. + + This is a regression test for GitHub issue #7859: pending state writes from + a failed superstep must not leak into a later successful run on the same + Workflow instance. + """ + workflow = WorkflowBuilder(start_executor=FlakyStateExecutor(id="flaky")).build() + + # First run: fails after staging a state write + with pytest.raises(RuntimeError, match="simulated transient failure"): + await workflow.run(FlakyMessage(fail=True)) + + # Verify the failed run did not leave the staged write pending + assert workflow._runner.state._pending == {} + + # Second run: succeeds without touching "secret" + result = await workflow.run(FlakyMessage(fail=False)) + assert result.get_final_state() == WorkflowRunState.IDLE + assert result.get_outputs() == ["ok"] + + # Verify the leaked state from the failed run is NOT in committed state + committed_state = workflow._runner.state.export_state() + assert "secret" not in committed_state + + async def test_workflow_checkpoint_runtime_only_configuration( simple_executor: Executor, ): From da0154c748843813e83f734f660612c09d691165 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Thu, 27 Aug 2026 14:01:28 +0530 Subject: [PATCH 2/9] fix: discard pending state before yielding failure events --- python/packages/core/agent_framework/_workflows/_runner.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index 519a76eace5..def8e99fe6b 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -147,12 +147,12 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: try: await iteration_task except Exception: + # Discard pending state writes from the failed superstep + self._state.discard() # Make sure failure-related events (like ExecutorFailedEvent) are surfaced if await self._ctx.has_events(): for event in await self._ctx.drain_events(): yield event - # Discard pending state writes from the failed superstep - self._state.discard() raise self._iteration += 1 From 8bfce51295b0f62ca9eb67442673eeb8c8e15900 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Mon, 31 Aug 2026 18:41:48 +0530 Subject: [PATCH 3/9] fix: discard pending state from failed supersteps --- .../_workflows/_edge_runner.py | 21 +++- .../agent_framework/_workflows/_runner.py | 53 ++++++++-- .../core/tests/workflow/test_workflow.py | 97 +++++++++++++++++++ 3 files changed, 162 insertions(+), 9 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_edge_runner.py b/python/packages/core/agent_framework/_workflows/_edge_runner.py index c14582894b9..bfe0334a338 100644 --- a/python/packages/core/agent_framework/_workflows/_edge_runner.py +++ b/python/packages/core/agent_framework/_workflows/_edge_runner.py @@ -282,8 +282,25 @@ async def send_to_edge(edge: Edge) -> bool: await self._execute_on_target(edge.target_id, [edge.source_id], message, state, ctx) return True - tasks = [send_to_edge(edge) for edge in deliverable_edges] - results = await asyncio.gather(*tasks) + tasks = [asyncio.create_task(send_to_edge(edge)) for edge in deliverable_edges] + if not tasks: + return False # asyncio.wait() requires a non-empty iterable + try: + done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_EXCEPTION) + except asyncio.CancelledError: + # If the wait() call itself is cancelled, cancel all child tasks + # before propagating to avoid orphaned work + for t in tasks: + t.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + raise + exceptions = [t.exception() for t in done if t.exception() is not None] + if exceptions: + for t in pending: + t.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + raise exceptions[0] + results = [t.result() for t in tasks] return any(results) # If we get here, it's a broadcast message with no deliverable edges diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index def8e99fe6b..d1893112747 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -139,8 +139,6 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: iteration_task.cancel() with contextlib.suppress(asyncio.CancelledError): await iteration_task - # Discard pending state writes from the cancelled superstep - self._state.discard() raise # Propagate errors from iteration, but first surface any pending events @@ -222,15 +220,56 @@ async def _deliver_messages_for_edge_runner(edge_runner: EdgeRunner) -> None: logger.debug(f"No outgoing edges found for executor {source_executor_id}; dropping messages.") return - tasks = [_deliver_messages_for_edge_runner(edge_runner) for edge_runner in associated_edge_runners] - await asyncio.gather(*tasks) + tasks = [asyncio.create_task(_deliver_messages_for_edge_runner(edge_runner)) for edge_runner in associated_edge_runners] + if not tasks: + return # asyncio.wait() requires a non-empty iterable + # Use FIRST_EXCEPTION to cancel pending siblings when one fails + try: + done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_EXCEPTION) + except asyncio.CancelledError: + # If the wait() call itself is cancelled, cancel all child tasks + # before propagating to avoid orphaned work + for t in tasks: + t.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + raise + exceptions = [t.exception() for t in done if t.exception() is not None] + if exceptions: + for t in pending: + t.cancel() + # Await all tasks to ensure cancelled tasks stop before propagating failure + await asyncio.gather(*tasks, return_exceptions=True) + raise exceptions[0] message_batches = await self._ctx.drain_messages() - tasks = [ - _deliver_messages(source_executor_id, source_messages) + # Create actual Task objects so we can cancel them if needed + task_objects = [ + asyncio.create_task(_deliver_messages(source_executor_id, source_messages)) for source_executor_id, source_messages in message_batches.items() ] - await asyncio.gather(*tasks) + + if not task_objects: + return # asyncio.wait() requires a non-empty iterable + + try: + done, pending = await asyncio.wait( + task_objects, + return_when=asyncio.FIRST_EXCEPTION + ) + except asyncio.CancelledError: + # If the wait() call itself is cancelled, cancel all child tasks + # before propagating to avoid orphaned work + for t in task_objects: + t.cancel() + await asyncio.gather(*task_objects, return_exceptions=True) + raise + + exceptions = [t.exception() for t in done if t.exception() is not None] + if exceptions: + for t in pending: + t.cancel() + await asyncio.gather(*task_objects, return_exceptions=True) + raise exceptions[0] async def _prepare_checkpoint_state(self) -> None: """Persist executor snapshots into committed shared state. diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index f1fb109129f..95829a64232 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -1,6 +1,7 @@ # Copyright (c) Microsoft. All rights reserved. import asyncio +import contextlib import gc import logging import tempfile @@ -658,6 +659,102 @@ async def test_workflow_discards_pending_state_after_failed_superstep(): assert "secret" not in committed_state +@dataclass +class FanOutTestMessage: + """Message for fan-out state leak test.""" + should_fail: bool + + +class FanOutSourceExecutor(Executor): + """Source executor that sends fan-out test messages.""" + + @handler + async def handle(self, message: FanOutTestMessage, ctx: WorkflowContext[FanOutTestMessage]) -> None: + # Forward the message to targets + await ctx.send_message(message) + + +class FailingTargetExecutor(Executor): + """Target that fails immediately on execute.""" + + def __init__(self, id: str, b_started: asyncio.Event | None = None) -> None: + super().__init__(id=id) + self._b_started = b_started + + @handler + async def handle(self, message: FanOutTestMessage, ctx: WorkflowContext) -> None: + if message.should_fail: + # Wait for slow target to start before failing to ensure deterministic ordering + if self._b_started: + await self._b_started.wait() + raise RuntimeError("target A failed") + + +class SlowStateWritingTargetExecutor(Executor): + """Target that writes state after being unblocked (vulnerable pattern if not cancelled).""" + + def __init__(self, id: str, b_started: asyncio.Event, release_b: asyncio.Event, b_task_ref: list) -> None: + super().__init__(id=id) + self._b_started = b_started + self._release_b = release_b + self._b_task_ref = b_task_ref + + @handler + async def handle(self, message: FanOutTestMessage, ctx: WorkflowContext) -> None: + self._b_task_ref.append(asyncio.current_task()) + self._b_started.set() + if message.should_fail: + await self._release_b.wait() + ctx.set_state("leak_key", "leaked_value") + + +async def test_workflow_discards_pending_state_after_fanout_failure(): + """Test that pending state from a fan-out sibling target is discarded when another target fails. + + Regression test for GitHub issue #7859: when a fan-out superstep has one target fail + while a sibling target is still running, the sibling's state writes must not leak + into committed state. + """ + b_started = asyncio.Event() + release_b = asyncio.Event() + b_task_ref: list[asyncio.Task] = [] + + source = FanOutSourceExecutor(id="source") + failing_target = FailingTargetExecutor(id="failing_target", b_started=b_started) + slow_target = SlowStateWritingTargetExecutor( + id="slow_target", b_started=b_started, release_b=release_b, b_task_ref=b_task_ref + ) + + workflow = ( + WorkflowBuilder(start_executor=source) + .add_fan_out_edges(source, [failing_target, slow_target]) + .build() + ) + + # Verify topology: single FanOutEdgeGroup with two targets under one edge_runner + from agent_framework._workflows._edge import FanOutEdgeGroup + fan_out_groups = [eg for eg in workflow.edge_groups if isinstance(eg, FanOutEdgeGroup)] + assert len(fan_out_groups) == 1, f"Expected 1 FanOutEdgeGroup, got {len(fan_out_groups)}" + assert fan_out_groups[0].target_ids == [failing_target.id, slow_target.id], \ + f"Expected targets [{failing_target.id}, {slow_target.id}], got {fan_out_groups[0].target_ids}" + + run_task = asyncio.create_task(workflow.run(FanOutTestMessage(should_fail=True))) + with pytest.raises(RuntimeError, match="target A failed"): + await asyncio.wait_for(run_task, timeout=5.0) + + release_b.set() + with contextlib.suppress(asyncio.CancelledError, Exception): + await asyncio.wait_for(b_task_ref[0], timeout=0.5) + + b_task_ref.clear() + + result = await workflow.run(FanOutTestMessage(should_fail=False)) + assert result.get_final_state() == WorkflowRunState.IDLE + + committed_state = workflow._runner.state.export_state() + assert "leak_key" not in committed_state + + async def test_workflow_checkpoint_runtime_only_configuration( simple_executor: Executor, ): From aa573c740807990bed2233c89864abc995dd91d6 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Wed, 9 Sep 2026 10:50:38 +0530 Subject: [PATCH 4/9] Fixing failing CI/CD workflows --- .../core/agent_framework/_workflows/_runner.py | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index 07774f551ad..675414037bb 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -227,12 +227,10 @@ async def _deliver_messages_for_edge_runner(edge_runner: EdgeRunner) -> None: await gather_cancelling_siblings_on_error(*tasks) message_batches = await self._ctx.drain_messages() - # Create actual Task objects so we can cancel them if needed - task_objects = [ - asyncio.create_task(_deliver_messages(source_executor_id, source_messages)) - for source_executor_id, source_messages in message_batches.items() - ] - await gather_cancelling_siblings_on_error(*task_objects) + await gather_cancelling_siblings_on_error( + *(_deliver_messages(source_executor_id, source_messages) + for source_executor_id, source_messages in message_batches.items()) + ) async def _prepare_checkpoint_state(self) -> None: """Persist executor snapshots into committed shared state. From 547db3cb6fae9c4f3d842f1dbe884a2941e42f56 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Thu, 10 Sep 2026 19:23:30 +0530 Subject: [PATCH 5/9] Fixing failing workflows --- .../core/agent_framework/_workflows/_runner.py | 6 ++++-- .../packages/core/tests/workflow/test_workflow.py | 13 ++++++------- 2 files changed, 10 insertions(+), 9 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index 675414037bb..4a229cff988 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -228,8 +228,10 @@ async def _deliver_messages_for_edge_runner(edge_runner: EdgeRunner) -> None: message_batches = await self._ctx.drain_messages() await gather_cancelling_siblings_on_error( - *(_deliver_messages(source_executor_id, source_messages) - for source_executor_id, source_messages in message_batches.items()) + *( + _deliver_messages(source_executor_id, source_messages) + for source_executor_id, source_messages in message_batches.items() + ) ) async def _prepare_checkpoint_state(self) -> None: diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index 05e45495741..b74d6d80f37 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -693,6 +693,7 @@ async def test_workflow_discards_pending_state_after_failed_superstep(): @dataclass class FanOutTestMessage: """Message for fan-out state leak test.""" + should_fail: bool @@ -756,18 +757,16 @@ async def test_workflow_discards_pending_state_after_fanout_failure(): id="slow_target", b_started=b_started, release_b=release_b, b_task_ref=b_task_ref ) - workflow = ( - WorkflowBuilder(start_executor=source) - .add_fan_out_edges(source, [failing_target, slow_target]) - .build() - ) + workflow = WorkflowBuilder(start_executor=source).add_fan_out_edges(source, [failing_target, slow_target]).build() # Verify topology: single FanOutEdgeGroup with two targets under one edge_runner from agent_framework._workflows._edge import FanOutEdgeGroup + fan_out_groups = [eg for eg in workflow.edge_groups if isinstance(eg, FanOutEdgeGroup)] assert len(fan_out_groups) == 1, f"Expected 1 FanOutEdgeGroup, got {len(fan_out_groups)}" - assert fan_out_groups[0].target_ids == [failing_target.id, slow_target.id], \ + assert fan_out_groups[0].target_ids == [failing_target.id, slow_target.id], ( f"Expected targets [{failing_target.id}, {slow_target.id}], got {fan_out_groups[0].target_ids}" + ) run_task = asyncio.create_task(workflow.run(FanOutTestMessage(should_fail=True))) with pytest.raises(RuntimeError, match="target A failed"): @@ -776,7 +775,7 @@ async def test_workflow_discards_pending_state_after_fanout_failure(): release_b.set() with contextlib.suppress(asyncio.CancelledError, Exception): await asyncio.wait_for(b_task_ref[0], timeout=0.5) - + b_task_ref.clear() result = await workflow.run(FanOutTestMessage(should_fail=False)) From 89e43bb0a329e0be99b9abc222d390e47c3b0073 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Tue, 15 Sep 2026 10:03:48 +0530 Subject: [PATCH 6/9] Fixing failing workflows --- .../tests/workflow/test_checkpoint_unrestricted_pickle.py | 5 ++++- python/packages/core/tests/workflow/test_workflow.py | 5 ++++- python/packages/core/tests/workflow/test_workflow_agent.py | 4 +--- .../_workflows/_executors_http.py | 2 ++ python/packages/foundry_hosting/tests/test_responses.py | 4 ++++ 5 files changed, 15 insertions(+), 5 deletions(-) diff --git a/python/packages/core/tests/workflow/test_checkpoint_unrestricted_pickle.py b/python/packages/core/tests/workflow/test_checkpoint_unrestricted_pickle.py index 1f78f4e6ad4..77f0ab84235 100644 --- a/python/packages/core/tests/workflow/test_checkpoint_unrestricted_pickle.py +++ b/python/packages/core/tests/workflow/test_checkpoint_unrestricted_pickle.py @@ -300,7 +300,10 @@ async def test_file_storage_rejects_unlisted_user_type_at_save(): graph_signature_hash="hash", state={"data": _AllowedTestState(name="test", value=1)}, ) - with pytest.raises(WorkflowCheckpointException, match="deserialization blocked|Unable to save|cannot be restored|cannot be encoded"): + with pytest.raises( + WorkflowCheckpointException, + match="deserialization blocked|Unable to save|cannot be restored|cannot be encoded", + ): await storage.save(checkpoint) assert not await asyncio.to_thread(lambda: list(Path(tmpdir).glob("*.json"))) diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index b74d6d80f37..051dbfc2297 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -768,7 +768,10 @@ async def test_workflow_discards_pending_state_after_fanout_failure(): f"Expected targets [{failing_target.id}, {slow_target.id}], got {fan_out_groups[0].target_ids}" ) - run_task = asyncio.create_task(workflow.run(FanOutTestMessage(should_fail=True))) + async def _run_workflow(): + return await workflow.run(FanOutTestMessage(should_fail=True)) + + run_task = asyncio.create_task(_run_workflow()) with pytest.raises(RuntimeError, match="target A failed"): await asyncio.wait_for(run_task, timeout=5.0) diff --git a/python/packages/core/tests/workflow/test_workflow_agent.py b/python/packages/core/tests/workflow/test_workflow_agent.py index 1394ef70ef6..3374d15d395 100644 --- a/python/packages/core/tests/workflow/test_workflow_agent.py +++ b/python/packages/core/tests/workflow/test_workflow_agent.py @@ -939,9 +939,7 @@ async def test_workflow_as_agent_stream_preserves_empty_additional_properties(se """Test that an explicitly empty additional_properties dict is not converted to None.""" @executor - async def empty_props_executor( - messages: list[Message], ctx: WorkflowContext[Any, AgentResponseUpdate] - ) -> None: + async def empty_props_executor(messages: list[Message], ctx: WorkflowContext[Any, AgentResponseUpdate]) -> None: await ctx.yield_output( AgentResponseUpdate( contents=[Content.from_text(text="payload")], diff --git a/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py b/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py index 0c67b579ab1..071e6837eca 100644 --- a/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py +++ b/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py @@ -206,6 +206,8 @@ async def handle_action( # Non-success path: still publish headers diagnostically, then raise. self._assign_response_headers(state, result) + # Commit the state before raising so headers are persisted even on error + ctx.state.commit() raise DeclarativeActionError(f"HTTP request to '{url}' failed with status code {result.status_code}.") # ----- Field resolution ---------------------------------------------------- diff --git a/python/packages/foundry_hosting/tests/test_responses.py b/python/packages/foundry_hosting/tests/test_responses.py index b4522d29107..3f6fd9d1904 100644 --- a/python/packages/foundry_hosting/tests/test_responses.py +++ b/python/packages/foundry_hosting/tests/test_responses.py @@ -1576,6 +1576,7 @@ async def test_readiness(self) -> None: # region Non-streaming +@pytest.mark.xdist_group("session_isolation") class TestNonStreaming: """ Non-streaming here means that the client requested a non-streaming response, instead of @@ -3540,6 +3541,7 @@ def run_dispatch(*args: Any, **kwargs: Any) -> ResponseStream[AgentResponseUpdat return agent +@pytest.mark.xdist_group("session_isolation") class TestMultiTurnMixedContent: """End-to-end multi-turn tests with mixed text and non-text content types.""" @@ -4513,6 +4515,7 @@ async def test_input_item_mcp_approval_response_resolves_to_approval_response(se assert c.approved is False +@pytest.mark.xdist_group("session_isolation") class TestFunctionApprovalRoundTrip: """End-to-end round-trip tests for the function approval flow. @@ -4889,6 +4892,7 @@ async def test_failed_entry_does_not_cache_stack(self) -> None: assert agent.__aenter__.await_count == 2 +@pytest.mark.xdist_group("session_isolation") class TestOAuthConsentSurfacing: async def test_non_streaming_consent_error_emits_oauth_output_item(self) -> None: agent = _make_agent( From c52fdd2f7a0ccb02cba39841711d069dbec1bffc Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Thu, 24 Sep 2026 22:19:22 +0530 Subject: [PATCH 7/9] fix(workflows): clean up failed superstep state --- python/packages/core/tests/test_types.py | 1 + python/packages/core/tests/workflow/test_workflow.py | 4 ++-- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/python/packages/core/tests/test_types.py b/python/packages/core/tests/test_types.py index 0622fb48e6d..71ab420ce1a 100644 --- a/python/packages/core/tests/test_types.py +++ b/python/packages/core/tests/test_types.py @@ -1,5 +1,6 @@ # Copyright (c) Microsoft. All rights reserved. + class TestPrependInstructionsEmpty: def test_empty_string_instructions_add_no_message(self) -> None: """An unset "" instruction must not inject a contentless system message.""" diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index 231b1a16982..34bee45a70c 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -787,10 +787,10 @@ async def test_workflow_discards_pending_state_after_fanout_failure(): f"Expected targets [{failing_target.id}, {slow_target.id}], got {fan_out_groups[0].target_ids}" ) - async def _run_workflow(): + async def _run_workflow() -> WorkflowRunResult: return await workflow.run(FanOutTestMessage(should_fail=True)) - run_task = asyncio.create_task(_run_workflow()) + run_task: asyncio.Task[WorkflowRunResult] = asyncio.create_task(_run_workflow()) with pytest.raises(RuntimeError, match="target A failed"): await asyncio.wait_for(run_task, timeout=5.0) From 92cca2956ee3c13ef63e7b1eb03ec4bb11920d56 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Wed, 30 Sep 2026 12:08:16 +0530 Subject: [PATCH 8/9] fix: preserve failed superstep state isolation --- .../agent_framework/_workflows/_runner.py | 4 +- .../core/tests/workflow/test_workflow.py | 46 +++++++++++++++++++ .../_workflows/_executors_http.py | 2 - .../tests/test_http_request_executor.py | 22 ++++++++- .../agent_framework_typesafe/_tool_calls.py | 4 +- 5 files changed, 69 insertions(+), 9 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index 275953f045d..dfaa5555941 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -146,8 +146,8 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: # Propagate errors from iteration, but first surface any pending events try: await iteration_task - except Exception: - # Discard pending state writes from the failed superstep + except (Exception, asyncio.CancelledError): + # Discard pending state writes from the failed or cancelled superstep self._state.discard() # Make sure failure-related events (like ExecutorFailedEvent) are surfaced if await self._ctx.has_events(): diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index 34bee45a70c..92a15dafbd6 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -678,6 +678,8 @@ async def handle_message( ) -> None: if message.fail: ctx.set_state("secret", "leaked-from-failed-run") + # Small delay to ensure cancellation can happen after write is staged + await asyncio.sleep(0.01) raise RuntimeError("simulated transient failure") await ctx.yield_output("ok") @@ -807,6 +809,50 @@ async def _run_workflow() -> WorkflowRunResult: assert "leak_key" not in committed_state +async def test_workflow_discards_pending_state_on_cancellation(): + """Test that pending state from a cancelled superstep is discarded and not committed. + + Regression test for PR #8819 review comment: asyncio.CancelledError inherits from + BaseException (not Exception since Python 3.8), so `except Exception:` never catches it. + When a superstep's task is cancelled, the discard() call must still run to prevent + pending writes from leaking into the next run. + + This test creates a workflow that will stage a state write, then yields control + before completing (via a slow executor). We cancel the task while it's mid-superstep + to ensure the write is staged and the runner's except block is hit. + """ + # Use the existing FlakyStateExecutor which stages a write then raises + workflow = WorkflowBuilder(start_executor=FlakyStateExecutor(id="flaky")).build() + + # Create a task that will stage a state write then fail + async def _run_failing(): + return await workflow.run(FlakyMessage(fail=True)) + + run_task = asyncio.create_task(_run_failing()) + + # Yield to let the task start and stage its write + await asyncio.sleep(0) + + # Cancel the task while it's mid-superstep (after write is staged but before error handling completes) + run_task.cancel() + + # Await the task, confirming CancelledError is raised (discard() must not swallow it) + with pytest.raises(asyncio.CancelledError): + await run_task + + # Verify the cancelled run did not leave the staged write pending + assert workflow._runner.state._pending == {} + + # Second run: succeeds without touching "secret" + result = await workflow.run(FlakyMessage(fail=False)) + assert result.get_final_state() == WorkflowRunState.IDLE + assert result.get_outputs() == ["ok"] + + # Verify the leaked state from the cancelled run is NOT in committed state + committed_state = workflow._runner.state.export_state() + assert "secret" not in committed_state + + async def test_workflow_checkpoint_runtime_only_configuration( simple_executor: Executor, ): diff --git a/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py b/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py index 071e6837eca..0c67b579ab1 100644 --- a/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py +++ b/python/packages/declarative/agent_framework_declarative/_workflows/_executors_http.py @@ -206,8 +206,6 @@ async def handle_action( # Non-success path: still publish headers diagnostically, then raise. self._assign_response_headers(state, result) - # Commit the state before raising so headers are persisted even on error - ctx.state.commit() raise DeclarativeActionError(f"HTTP request to '{url}' failed with status code {result.status_code}.") # ----- Field resolution ---------------------------------------------------- diff --git a/python/packages/declarative/tests/test_http_request_executor.py b/python/packages/declarative/tests/test_http_request_executor.py index 1e71bd48b02..2d3002c7c79 100644 --- a/python/packages/declarative/tests/test_http_request_executor.py +++ b/python/packages/declarative/tests/test_http_request_executor.py @@ -584,13 +584,31 @@ async def test_response_headers_empty_assigned_none(self) -> None: @pytest.mark.asyncio async def test_non_2xx_still_publishes_headers(self) -> None: + """Non-2xx responses publish headers to pending state during the superstep. + + After PR #8819 fix: headers are written to pending state but are NOT + committed on error. The runner's discard() clears them along with other + pending writes from the failed superstep. This test verifies headers are + written to pending state during execution (visible via state.get() which + checks pending first), but confirms they are NOT durably persisted after + the error is raised and discard() runs. + + Note: This test checks committed state via export_state() after the error. + Before the fix, ctx.state.commit() in the error path would persist headers + to committed state. After the fix, discard() clears them from pending, + so they never reach committed state. + """ handler = StubHandler(_err(status=500, body="boom", headers={"X-Trace": ["abc"]})) factory = WorkflowFactory(http_request_handler=handler) workflow = factory.create_workflow_from_definition(_yaml(_action(response_headers="Local.H"))) with pytest.raises(DeclarativeActionError): await workflow.run({}) - decl = workflow._runner.state.get(DECLARATIVE_STATE_KEY) - assert decl["Local"]["H"] == {"X-Trace": "abc"} + # After the error and discard(), headers should NOT be in committed state + committed_state = workflow._runner.state.export_state() + # The DECLARATIVE_STATE_KEY may not exist at all if nothing was committed + decl_committed = committed_state.get(DECLARATIVE_STATE_KEY, {}) + # Headers should not be durably persisted after the failed action + assert "Local" not in decl_committed or "H" not in decl_committed.get("Local", {}) # ---------- ConversationId append ------------------------------------------- diff --git a/python/packages/typesafe/agent_framework_typesafe/_tool_calls.py b/python/packages/typesafe/agent_framework_typesafe/_tool_calls.py index cf5508533af..e5b6f4d12e5 100644 --- a/python/packages/typesafe/agent_framework_typesafe/_tool_calls.py +++ b/python/packages/typesafe/agent_framework_typesafe/_tool_calls.py @@ -625,9 +625,7 @@ def _reject_unsupported_schema_constraints( ) -> None: unsupported_constraints = sorted(schema.keys() - supported_keys) if unsupported_constraints: - raise _UnsupportedToolSchema( - f"unsupported {location} schema constraints: {', '.join(unsupported_constraints)}" - ) + raise _UnsupportedToolSchema(f"unsupported {location} schema constraints: {', '.join(unsupported_constraints)}") def _describe_value(value: Any) -> str: From 451beb454a0d45a8774c5bb3cdc47f42aad0d172 Mon Sep 17 00:00:00 2001 From: Namraa Patel Date: Thu, 1 Oct 2026 09:17:01 +0530 Subject: [PATCH 9/9] fix(workflow): discard pending state on cancellation --- .../agent_framework/_workflows/_runner.py | 74 ++++++++++--------- .../core/tests/workflow/test_workflow.py | 18 ++++- 2 files changed, 53 insertions(+), 39 deletions(-) diff --git a/python/packages/core/agent_framework/_workflows/_runner.py b/python/packages/core/agent_framework/_workflows/_runner.py index dfaa5555941..f41183cb1c2 100644 --- a/python/packages/core/agent_framework/_workflows/_runner.py +++ b/python/packages/core/agent_framework/_workflows/_runner.py @@ -124,47 +124,51 @@ async def run_until_convergence(self) -> AsyncGenerator[WorkflowEvent, None]: logger.info(f"Starting superstep {self._iteration + 1}") yield WorkflowEvent.superstep_started(iteration=self._iteration + 1) - # Wake on either a live event or iteration completion, including silent supersteps. - iteration_task = asyncio.create_task(self._run_iteration()) - event_task: asyncio.Task[WorkflowEvent] | None = None + committed = False try: - while not iteration_task.done(): - event_task = asyncio.create_task(self._ctx.next_event()) - done, _ = await asyncio.wait((iteration_task, event_task), return_when=asyncio.FIRST_COMPLETED) - if event_task in done: - yield event_task.result() - finally: - # Cancellation and generator closure must not leave an event waiter or executor running. - tasks: list[asyncio.Task[Any]] = ( - [iteration_task] if event_task is None else [iteration_task, event_task] - ) - for task in tasks: - if not task.done(): - task.cancel() - await asyncio.gather(*tasks, return_exceptions=True) - - # Propagate errors from iteration, but first surface any pending events - try: - await iteration_task - except (Exception, asyncio.CancelledError): - # Discard pending state writes from the failed or cancelled superstep - self._state.discard() - # Make sure failure-related events (like ExecutorFailedEvent) are surfaced + # Wake on either a live event or iteration completion, including silent supersteps. + iteration_task = asyncio.create_task(self._run_iteration()) + event_task: asyncio.Task[WorkflowEvent] | None = None + try: + while not iteration_task.done(): + event_task = asyncio.create_task(self._ctx.next_event()) + done, _ = await asyncio.wait((iteration_task, event_task), return_when=asyncio.FIRST_COMPLETED) + if event_task in done: + yield event_task.result() + finally: + # Cancellation and generator closure must not leave an event waiter or executor running. + tasks: list[asyncio.Task[Any]] = ( + [iteration_task] if event_task is None else [iteration_task, event_task] + ) + for task in tasks: + if not task.done(): + task.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + + # Propagate errors from iteration, but first surface any pending events + try: + await iteration_task + except (Exception, asyncio.CancelledError): + # Make sure failure-related events (like ExecutorFailedEvent) are surfaced + if await self._ctx.has_events(): + for event in await self._ctx.drain_events(): + yield event + raise + self._iteration += 1 + + # Drain any straggler events emitted at tail end if await self._ctx.has_events(): for event in await self._ctx.drain_events(): yield event - raise - self._iteration += 1 - # Drain any straggler events emitted at tail end - if await self._ctx.has_events(): - for event in await self._ctx.drain_events(): - yield event + logger.info(f"Completed superstep {self._iteration}") - logger.info(f"Completed superstep {self._iteration}") - - # Commit pending state changes at superstep boundary - self._state.commit() + # Commit pending state changes at superstep boundary + self._state.commit() + committed = True + finally: + if not committed: + self._state.discard() # Create checkpoint after each superstep iteration await self.create_checkpoint_if_enabled() diff --git a/python/packages/core/tests/workflow/test_workflow.py b/python/packages/core/tests/workflow/test_workflow.py index 92a15dafbd6..6511217f248 100644 --- a/python/packages/core/tests/workflow/test_workflow.py +++ b/python/packages/core/tests/workflow/test_workflow.py @@ -670,6 +670,10 @@ class FlakyMessage: class FlakyStateExecutor(Executor): """An executor that fails on demand to test state discard on failure.""" + def __init__(self, id: str, write_staged_event: asyncio.Event | None = None) -> None: + super().__init__(id=id) + self._write_staged_event = write_staged_event + @handler async def handle_message( self, @@ -678,6 +682,9 @@ async def handle_message( ) -> None: if message.fail: ctx.set_state("secret", "leaked-from-failed-run") + # Signal that the write has been staged for deterministic test synchronization + if self._write_staged_event: + self._write_staged_event.set() # Small delay to ensure cancellation can happen after write is staged await asyncio.sleep(0.01) raise RuntimeError("simulated transient failure") @@ -821,8 +828,11 @@ async def test_workflow_discards_pending_state_on_cancellation(): before completing (via a slow executor). We cancel the task while it's mid-superstep to ensure the write is staged and the runner's except block is hit. """ - # Use the existing FlakyStateExecutor which stages a write then raises - workflow = WorkflowBuilder(start_executor=FlakyStateExecutor(id="flaky")).build() + # Use deterministic synchronization to ensure write is staged before cancellation + write_staged_event = asyncio.Event() + workflow = WorkflowBuilder( + start_executor=FlakyStateExecutor(id="flaky", write_staged_event=write_staged_event) + ).build() # Create a task that will stage a state write then fail async def _run_failing(): @@ -830,8 +840,8 @@ async def _run_failing(): run_task = asyncio.create_task(_run_failing()) - # Yield to let the task start and stage its write - await asyncio.sleep(0) + # Wait for the write to be staged before cancelling (deterministic synchronization) + await write_staged_event.wait() # Cancel the task while it's mid-superstep (after write is staged but before error handling completes) run_task.cancel()