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(