From 47fa48d787893d68bed0508e7ad757ea198c2c48 Mon Sep 17 00:00:00 2001 From: Lucas Gomide Date: Wed, 19 Aug 2026 03:07:33 -0300 Subject: [PATCH] fix: close the agent scope on every failed attempt (#6997) * fix: close the agent scope on every failed attempt `_check_execution_error` only emitted `AgentExecutionErrorEvent` once the retries were exhausted, but each retry re-enters `execute_task` and opens a new `agent_execution_started` scope. The scopes left open were then popped by the next ending event, so `task_failed` closed an agent scope instead of `task_started` and the task never got its own terminal pairing. Passthrough exceptions keep bubbling untouched, since a HITL pause must leave its scope open for the resume. * fix: return the retried result instead of finalizing it twice A retry reenters `execute_task`, whose own `_finalize_task_execution` already emitted `AgentExecutionCompletedEvent`, and the outer frame then finalized the same result again. The duplicate used to be absorbed by the `agent_execution_started` scope that a failed attempt left open, so closing every attempt exposed it: the extra completed event popped `task_started`, and the task and crew ends paired with the wrong scopes. --------- Co-authored-by: Vidit Ostwal <110953813+Vidit-Ostwal@users.noreply.github.com> --- lib/crewai/src/crewai/agent/core.py | 36 ++++---- lib/crewai/tests/utilities/test_events.py | 105 ++++++++++++++++++++++ 2 files changed, 121 insertions(+), 20 deletions(-) diff --git a/lib/crewai/src/crewai/agent/core.py b/lib/crewai/src/crewai/agent/core.py index 70e090c6f..156e576c0 100644 --- a/lib/crewai/src/crewai/agent/core.py +++ b/lib/crewai/src/crewai/agent/core.py @@ -728,28 +728,22 @@ class Agent(BaseAgent): Raises: Exception: If the error is from litellm, a passthrough, or retries are exhausted. """ - if e.__class__.__module__.startswith("litellm"): - crewai_event_bus.emit( - self, - event=AgentExecutionErrorEvent( - agent=self, - task=task, - error=str(e), - ), - ) - raise e if isinstance(e, _passthrough_exceptions): raise + # A retry re-enters execute_task, which opens a new agent_execution_started + # scope, so every failed attempt has to close its own. + crewai_event_bus.emit( + self, + event=AgentExecutionErrorEvent( + agent=self, + task=task, + error=str(e), + ), + ) + if e.__class__.__module__.startswith("litellm"): + raise e self._times_executed += 1 if self._times_executed > self.max_retry_limit: - crewai_event_bus.emit( - self, - event=AgentExecutionErrorEvent( - agent=self, - task=task, - error=str(e), - ), - ) raise e def _handle_execution_error( @@ -886,7 +880,8 @@ class Agent(BaseAgent): ) raise e except Exception as e: - result = self._handle_execution_error(e, task, context, tools) + # The retry runs a whole execute_task of its own, result already finalized. + return self._handle_execution_error(e, task, context, tools) return self._finalize_task_execution(task, result) @@ -1020,7 +1015,8 @@ class Agent(BaseAgent): ) raise e except Exception as e: - result = await self._handle_execution_error_async(e, task, context, tools) + # The retry runs a whole aexecute_task of its own, result already finalized. + return await self._handle_execution_error_async(e, task, context, tools) return self._finalize_task_execution(task, result) diff --git a/lib/crewai/tests/utilities/test_events.py b/lib/crewai/tests/utilities/test_events.py index f4bf43100..77a4a6578 100644 --- a/lib/crewai/tests/utilities/test_events.py +++ b/lib/crewai/tests/utilities/test_events.py @@ -373,6 +373,111 @@ def test_agent_emits_execution_error_event(base_agent, base_task): assert received_events[0].type == "agent_execution_error" +def test_agent_retries_close_the_scope_of_every_attempt(): + agent_started = [] + agent_errored = [] + task_started = [] + task_failed = [] + + @crewai_event_bus.on(AgentExecutionStartedEvent) + def handle_agent_started(source, event): + agent_started.append(event) + + @crewai_event_bus.on(AgentExecutionErrorEvent) + def handle_agent_error(source, event): + agent_errored.append(event) + + @crewai_event_bus.on(TaskStartedEvent) + def handle_task_started(source, event): + task_started.append(event) + + @crewai_event_bus.on(TaskFailedEvent) + def handle_task_failed(source, event): + task_failed.append(event) + + from crewai.experimental.agent_executor import AgentExecutor + + agent = Agent( + role="retrying_agent", + llm="gpt-4o-mini", + goal="Just say hi", + backstory="You are a helpful assistant that just says hi", + max_retry_limit=1, + ) + task = Task(description="Just say hi", expected_output="hi", agent=agent) + crew = Crew(agents=[agent], tasks=[task]) + + with patch.object(AgentExecutor, "invoke", side_effect=Exception("boom")): + with pytest.raises(Exception): # noqa: B017 + crew.kickoff() + + wait_for_event_handlers() + + assert len(agent_started) == 2 + assert {event.started_event_id for event in agent_errored} == { + event.event_id for event in agent_started + } + assert len(task_failed) == 1 + assert task_failed[0].started_event_id == task_started[0].event_id + + +def test_agent_retry_that_succeeds_closes_one_scope_per_attempt(): + agent_started = [] + agent_errored = [] + agent_completed = [] + task_started = [] + task_completed = [] + + @crewai_event_bus.on(AgentExecutionStartedEvent) + def handle_agent_started(source, event): + agent_started.append(event) + + @crewai_event_bus.on(AgentExecutionErrorEvent) + def handle_agent_error(source, event): + agent_errored.append(event) + + @crewai_event_bus.on(AgentExecutionCompletedEvent) + def handle_agent_completed(source, event): + agent_completed.append(event) + + @crewai_event_bus.on(TaskStartedEvent) + def handle_task_started(source, event): + task_started.append(event) + + @crewai_event_bus.on(TaskCompletedEvent) + def handle_task_completed(source, event): + task_completed.append(event) + + from crewai.experimental.agent_executor import AgentExecutor + + agent = Agent( + role="retrying_agent", + llm="gpt-4o-mini", + goal="Just say hi", + backstory="You are a helpful assistant that just says hi", + max_retry_limit=2, + ) + task = Task(description="Just say hi", expected_output="hi", agent=agent) + crew = Crew(agents=[agent], tasks=[task]) + + with patch.object( + AgentExecutor, "invoke", side_effect=[Exception("boom"), {"output": "hi"}] + ): + crew.kickoff() + + wait_for_event_handlers() + + assert len(agent_started) == 2 + assert len(agent_errored) == 1 + assert len(agent_completed) == 1 + assert { + agent_errored[0].started_event_id, + agent_completed[0].started_event_id, + } == {event.event_id for event in agent_started} + assert len(task_completed) == 1 + assert task_completed[0].started_event_id == task_started[0].event_id + + class SayHiTool(BaseTool): name: str = Field(default="say_hi", description="The name of the tool") description: str = Field(