mirror of
https://github.com/crewAIInc/crewAI.git
synced 2026-09-20 18:13:49 +00:00
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>
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user