mirror of
https://github.com/crewAIInc/crewAI.git
synced 2026-08-10 08:21:54 +00:00
Compare commits
1 Commits
docs/agent
...
fix/task-c
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
328e0e36e9 |
@@ -1065,6 +1065,7 @@ class Crew(FlowTrackable, BaseModel):
|
|||||||
self,
|
self,
|
||||||
CrewKickoffFailedEvent(
|
CrewKickoffFailedEvent(
|
||||||
error=str(e),
|
error=str(e),
|
||||||
|
error_type=type(e).__name__,
|
||||||
crew_name=self.name,
|
crew_name=self.name,
|
||||||
started_event_id=self._kickoff_event_id,
|
started_event_id=self._kickoff_event_id,
|
||||||
),
|
),
|
||||||
@@ -1279,6 +1280,7 @@ class Crew(FlowTrackable, BaseModel):
|
|||||||
self,
|
self,
|
||||||
CrewKickoffFailedEvent(
|
CrewKickoffFailedEvent(
|
||||||
error=str(e),
|
error=str(e),
|
||||||
|
error_type=type(e).__name__,
|
||||||
crew_name=self.name,
|
crew_name=self.name,
|
||||||
started_event_id=self._kickoff_event_id,
|
started_event_id=self._kickoff_event_id,
|
||||||
),
|
),
|
||||||
|
|||||||
@@ -198,6 +198,11 @@ class EventListener(BaseEventListener):
|
|||||||
|
|
||||||
@crewai_event_bus.on(CrewKickoffFailedEvent)
|
@crewai_event_bus.on(CrewKickoffFailedEvent)
|
||||||
def on_crew_failed(source: Any, event: CrewKickoffFailedEvent) -> None:
|
def on_crew_failed(source: Any, event: CrewKickoffFailedEvent) -> None:
|
||||||
|
# Previously this handler never touched telemetry, so a crew that
|
||||||
|
# raised left its execution span open: never ended, never exported,
|
||||||
|
# and the failure invisible downstream.
|
||||||
|
self._telemetry.crew_failed(source, event.error_type)
|
||||||
|
|
||||||
self.formatter.handle_crew_status(
|
self.formatter.handle_crew_status(
|
||||||
event.crew_name or "Crew",
|
event.crew_name or "Crew",
|
||||||
source.id,
|
source.id,
|
||||||
@@ -261,8 +266,10 @@ class EventListener(BaseEventListener):
|
|||||||
def on_task_failed(source: Any, event: TaskFailedEvent) -> None:
|
def on_task_failed(source: Any, event: TaskFailedEvent) -> None:
|
||||||
span = self.execution_spans.pop(source, None)
|
span = self.execution_spans.pop(source, None)
|
||||||
if span:
|
if span:
|
||||||
if source.agent and source.agent.crew:
|
# Closed unconditionally: previously the span was only ended
|
||||||
self._telemetry.task_ended(span, source, source.agent.crew)
|
# when source.agent.crew was present, so any task that failed
|
||||||
|
# without one leaked its span and was never exported.
|
||||||
|
self._telemetry.task_failed(span, source, event.error_type)
|
||||||
|
|
||||||
task_name = get_task_name(source)
|
task_name = get_task_name(source)
|
||||||
self.formatter.handle_task_status(
|
self.formatter.handle_task_status(
|
||||||
|
|||||||
@@ -52,6 +52,13 @@ class CrewKickoffFailedEvent(CrewBaseEvent):
|
|||||||
"""Event emitted when a crew fails to complete execution"""
|
"""Event emitted when a crew fails to complete execution"""
|
||||||
|
|
||||||
error: str
|
error: str
|
||||||
|
error_type: str | None = None
|
||||||
|
"""Exception class name (e.g. "ValidationError").
|
||||||
|
|
||||||
|
Kept separate from ``error`` so telemetry can record what kind of failure
|
||||||
|
occurred without ever touching the message, which routinely contains
|
||||||
|
prompts, model output, or credentials.
|
||||||
|
"""
|
||||||
type: Literal["crew_kickoff_failed"] = "crew_kickoff_failed"
|
type: Literal["crew_kickoff_failed"] = "crew_kickoff_failed"
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -49,6 +49,13 @@ class TaskFailedEvent(BaseEvent):
|
|||||||
"""Event emitted when a task fails"""
|
"""Event emitted when a task fails"""
|
||||||
|
|
||||||
error: str
|
error: str
|
||||||
|
error_type: str | None = None
|
||||||
|
"""Exception class name (e.g. "ValidationError").
|
||||||
|
|
||||||
|
Kept separate from ``error`` so telemetry can record what kind of failure
|
||||||
|
occurred without ever touching the message, which routinely contains
|
||||||
|
prompts, model output, or credentials.
|
||||||
|
"""
|
||||||
type: Literal["task_failed"] = "task_failed"
|
type: Literal["task_failed"] = "task_failed"
|
||||||
task: Any | None = None
|
task: Any | None = None
|
||||||
|
|
||||||
|
|||||||
@@ -797,7 +797,10 @@ class Task(BaseModel):
|
|||||||
return task_output
|
return task_output
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
self.end_time = datetime.datetime.now()
|
self.end_time = datetime.datetime.now()
|
||||||
crewai_event_bus.emit(self, TaskFailedEvent(error=str(e), task=self))
|
crewai_event_bus.emit(
|
||||||
|
self,
|
||||||
|
TaskFailedEvent(error=str(e), error_type=type(e).__name__, task=self),
|
||||||
|
)
|
||||||
raise e
|
raise e
|
||||||
finally:
|
finally:
|
||||||
clear_task_files(self.id)
|
clear_task_files(self.id)
|
||||||
@@ -953,7 +956,10 @@ class Task(BaseModel):
|
|||||||
return task_output
|
return task_output
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
self.end_time = datetime.datetime.now()
|
self.end_time = datetime.datetime.now()
|
||||||
crewai_event_bus.emit(self, TaskFailedEvent(error=str(e), task=self))
|
crewai_event_bus.emit(
|
||||||
|
self,
|
||||||
|
TaskFailedEvent(error=str(e), error_type=type(e).__name__, task=self),
|
||||||
|
)
|
||||||
raise e
|
raise e
|
||||||
finally:
|
finally:
|
||||||
clear_task_files(self.id)
|
clear_task_files(self.id)
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ from crewai.telemetry.utils import (
|
|||||||
add_crew_and_task_attributes,
|
add_crew_and_task_attributes,
|
||||||
add_crew_attributes,
|
add_crew_attributes,
|
||||||
close_span,
|
close_span,
|
||||||
|
close_span_with_error,
|
||||||
)
|
)
|
||||||
from crewai.utilities.i18n import I18N_DEFAULT
|
from crewai.utilities.i18n import I18N_DEFAULT
|
||||||
from crewai.utilities.logger_utils import suppress_warnings
|
from crewai.utilities.logger_utils import suppress_warnings
|
||||||
@@ -571,6 +572,30 @@ class Telemetry:
|
|||||||
|
|
||||||
self._safe_telemetry_operation(_operation)
|
self._safe_telemetry_operation(_operation)
|
||||||
|
|
||||||
|
def task_failed(
|
||||||
|
self, span: Span, task: Task, error_type: str | None = None
|
||||||
|
) -> None:
|
||||||
|
"""Records that a task execution failed and closes its span with ERROR.
|
||||||
|
|
||||||
|
Previously failures were routed through task_ended, which closes every
|
||||||
|
span as OK - making failed and successful tasks indistinguishable
|
||||||
|
downstream and leaving error counts permanently at zero.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
span: The OpenTelemetry span tracking the task execution.
|
||||||
|
task: The task that failed.
|
||||||
|
error_type: Exception class name. The error message is never
|
||||||
|
recorded - it routinely contains prompts and model output.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _operation() -> None:
|
||||||
|
if hasattr(task, "fingerprint") and task.fingerprint:
|
||||||
|
self._add_attribute(span, "task_fingerprint", task.fingerprint.uuid_str)
|
||||||
|
|
||||||
|
close_span_with_error(span, error_type)
|
||||||
|
|
||||||
|
self._safe_telemetry_operation(_operation)
|
||||||
|
|
||||||
def tool_repeated_usage(self, llm: Any, tool_name: str, attempts: int) -> None:
|
def tool_repeated_usage(self, llm: Any, tool_name: str, attempts: int) -> None:
|
||||||
"""Records when a tool is used repeatedly, which might indicate an issue.
|
"""Records when a tool is used repeatedly, which might indicate an issue.
|
||||||
|
|
||||||
@@ -922,6 +947,28 @@ class Telemetry:
|
|||||||
if crew.share_crew:
|
if crew.share_crew:
|
||||||
self._safe_telemetry_operation(_operation)
|
self._safe_telemetry_operation(_operation)
|
||||||
|
|
||||||
|
def crew_failed(self, crew: Any, error_type: str | None = None) -> None:
|
||||||
|
"""Records that a crew execution failed and closes its span.
|
||||||
|
|
||||||
|
Without this, a crew that raises leaves its execution span open: it is
|
||||||
|
never ended, never exported, and the failure is invisible downstream.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
crew: The crew whose execution failed.
|
||||||
|
error_type: Exception class name. The error message is never
|
||||||
|
recorded - it routinely contains prompts and model output.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def _operation() -> None:
|
||||||
|
span = getattr(crew, "_execution_span", None)
|
||||||
|
if span is None:
|
||||||
|
return
|
||||||
|
self._add_attribute(span, "crewai_version", version("crewai"))
|
||||||
|
close_span_with_error(span, error_type)
|
||||||
|
crew._execution_span = None
|
||||||
|
|
||||||
|
self._safe_telemetry_operation(_operation)
|
||||||
|
|
||||||
def _add_attribute(self, span: Span, key: str, value: Any) -> None:
|
def _add_attribute(self, span: Span, key: str, value: Any) -> None:
|
||||||
"""Add an attribute to a span.
|
"""Add an attribute to a span.
|
||||||
|
|
||||||
|
|||||||
@@ -111,3 +111,23 @@ def close_span(span: Span) -> None:
|
|||||||
"""
|
"""
|
||||||
span.set_status(Status(StatusCode.OK))
|
span.set_status(Status(StatusCode.OK))
|
||||||
span.end()
|
span.end()
|
||||||
|
|
||||||
|
|
||||||
|
def close_span_with_error(span: Span, error_type: str | None = None) -> None:
|
||||||
|
"""Set span status to ERROR and end it.
|
||||||
|
|
||||||
|
Used for spans representing work that failed, so failures are
|
||||||
|
distinguishable from successes downstream. Only the exception's *type* is
|
||||||
|
recorded - never the message, which routinely contains prompts, model
|
||||||
|
output, or credentials.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
span: The span to close.
|
||||||
|
error_type: Exception class name (e.g. "ValidationError"). Anything
|
||||||
|
that is not a plain identifier is discarded rather than recorded,
|
||||||
|
so a message can never be passed in by mistake.
|
||||||
|
"""
|
||||||
|
span.set_status(Status(StatusCode.ERROR))
|
||||||
|
if error_type and error_type.isidentifier():
|
||||||
|
span.set_attribute("error_type", error_type)
|
||||||
|
span.end()
|
||||||
|
|||||||
210
lib/crewai/tests/telemetry/test_failure_instrumentation.py
Normal file
210
lib/crewai/tests/telemetry/test_failure_instrumentation.py
Normal file
@@ -0,0 +1,210 @@
|
|||||||
|
"""Tests that failed executions are recorded as failures, not successes.
|
||||||
|
|
||||||
|
Regression coverage for telemetry that reported every task as OK, leaving
|
||||||
|
downstream error counts permanently at zero.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from unittest.mock import Mock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from opentelemetry.sdk.trace import TracerProvider
|
||||||
|
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
|
||||||
|
from opentelemetry.sdk.trace.export.in_memory_span_exporter import (
|
||||||
|
InMemorySpanExporter,
|
||||||
|
)
|
||||||
|
from opentelemetry.trace import StatusCode
|
||||||
|
|
||||||
|
from crewai.telemetry.utils import close_span, close_span_with_error
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def enable_otel_sdk(monkeypatch):
|
||||||
|
"""Ensure the OTel SDK is active for these tests.
|
||||||
|
|
||||||
|
The suite runs with OTEL_SDK_DISABLED=true, which makes TracerProvider hand
|
||||||
|
out non-recording spans that are never exported. Set explicitly rather than
|
||||||
|
relying on the root conftest teardown, which pops the variable and would
|
||||||
|
otherwise leave only the first test in a session running against a
|
||||||
|
disabled SDK.
|
||||||
|
"""
|
||||||
|
monkeypatch.delenv("OTEL_SDK_DISABLED", raising=False)
|
||||||
|
monkeypatch.delenv("CREWAI_DISABLE_TELEMETRY", raising=False)
|
||||||
|
monkeypatch.delenv("CREWAI_DISABLE_TRACKING", raising=False)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def exporter():
|
||||||
|
exp = InMemorySpanExporter()
|
||||||
|
provider = TracerProvider()
|
||||||
|
provider.add_span_processor(SimpleSpanProcessor(exp))
|
||||||
|
# yield rather than return: the generator frame keeps `provider` alive for
|
||||||
|
# the test. If it is collected, its processor shuts down and spans are lost.
|
||||||
|
yield exp, provider.get_tracer("test")
|
||||||
|
|
||||||
|
|
||||||
|
def test_close_span_with_error_sets_error_status(exporter):
|
||||||
|
exp, tracer = exporter
|
||||||
|
|
||||||
|
close_span_with_error(tracer.start_span("Task Execution"), "ValidationError")
|
||||||
|
|
||||||
|
span = exp.get_finished_spans()[0]
|
||||||
|
assert span.status.status_code is StatusCode.ERROR
|
||||||
|
assert span.attributes["error_type"] == "ValidationError"
|
||||||
|
|
||||||
|
|
||||||
|
def test_successful_and_failed_spans_are_distinguishable(exporter):
|
||||||
|
"""The whole point: a downstream count of failures must be possible."""
|
||||||
|
exp, tracer = exporter
|
||||||
|
|
||||||
|
close_span(tracer.start_span("Task Execution"))
|
||||||
|
close_span_with_error(tracer.start_span("Task Execution"), "TimeoutError")
|
||||||
|
close_span(tracer.start_span("Task Execution"))
|
||||||
|
|
||||||
|
spans = exp.get_finished_spans()
|
||||||
|
failed = [s for s in spans if s.status.status_code is StatusCode.ERROR]
|
||||||
|
assert len(spans) == 3
|
||||||
|
assert len(failed) == 1
|
||||||
|
assert failed[0].attributes["error_type"] == "TimeoutError"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"not_an_identifier",
|
||||||
|
[
|
||||||
|
"Rate limit exceeded for gpt-4o",
|
||||||
|
"API key sk-live-1234 is invalid",
|
||||||
|
"connection to db://user:pass@host failed",
|
||||||
|
"",
|
||||||
|
" ",
|
||||||
|
"429",
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_error_message_can_never_be_recorded(exporter, not_an_identifier):
|
||||||
|
"""PII guard: only identifier-shaped values survive.
|
||||||
|
|
||||||
|
Error messages routinely contain prompts, model output, and credentials.
|
||||||
|
Passing one where an exception class name belongs must record nothing.
|
||||||
|
"""
|
||||||
|
exp, tracer = exporter
|
||||||
|
|
||||||
|
close_span_with_error(tracer.start_span("Task Execution"), not_an_identifier)
|
||||||
|
|
||||||
|
span = exp.get_finished_spans()[0]
|
||||||
|
assert span.status.status_code is StatusCode.ERROR
|
||||||
|
assert "error_type" not in (span.attributes or {})
|
||||||
|
|
||||||
|
|
||||||
|
def test_error_type_is_optional(exporter):
|
||||||
|
exp, tracer = exporter
|
||||||
|
|
||||||
|
close_span_with_error(tracer.start_span("Task Execution"))
|
||||||
|
|
||||||
|
span = exp.get_finished_spans()[0]
|
||||||
|
assert span.status.status_code is StatusCode.ERROR
|
||||||
|
assert "error_type" not in (span.attributes or {})
|
||||||
|
|
||||||
|
|
||||||
|
def test_real_exception_class_names_are_accepted(exporter):
|
||||||
|
"""Every builtin exception name is a valid identifier, so none are dropped."""
|
||||||
|
exp, tracer = exporter
|
||||||
|
|
||||||
|
for exc in (ValueError, TimeoutError, KeyError, RuntimeError, ConnectionError):
|
||||||
|
close_span_with_error(tracer.start_span("Task Execution"), exc.__name__)
|
||||||
|
|
||||||
|
recorded = [s.attributes["error_type"] for s in exp.get_finished_spans()]
|
||||||
|
assert recorded == [
|
||||||
|
"ValueError",
|
||||||
|
"TimeoutError",
|
||||||
|
"KeyError",
|
||||||
|
"RuntimeError",
|
||||||
|
"ConnectionError",
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_failed_closes_span_with_error():
|
||||||
|
from crewai.telemetry.telemetry import Telemetry
|
||||||
|
|
||||||
|
exp = InMemorySpanExporter()
|
||||||
|
provider = TracerProvider()
|
||||||
|
provider.add_span_processor(SimpleSpanProcessor(exp))
|
||||||
|
|
||||||
|
telemetry = Telemetry()
|
||||||
|
telemetry.ready = True
|
||||||
|
span = provider.get_tracer("test").start_span("Task Execution")
|
||||||
|
|
||||||
|
telemetry.task_failed(span, Mock(fingerprint=None), "ValueError")
|
||||||
|
|
||||||
|
finished = exp.get_finished_spans()[0]
|
||||||
|
assert finished.status.status_code is StatusCode.ERROR
|
||||||
|
assert finished.attributes["error_type"] == "ValueError"
|
||||||
|
|
||||||
|
|
||||||
|
def test_crew_failed_closes_leaked_execution_span():
|
||||||
|
"""A crew that raises must not leave its span open and unexported."""
|
||||||
|
from crewai.telemetry.telemetry import Telemetry
|
||||||
|
|
||||||
|
exp = InMemorySpanExporter()
|
||||||
|
provider = TracerProvider()
|
||||||
|
provider.add_span_processor(SimpleSpanProcessor(exp))
|
||||||
|
|
||||||
|
telemetry = Telemetry()
|
||||||
|
telemetry.ready = True
|
||||||
|
|
||||||
|
crew = Mock()
|
||||||
|
crew._execution_span = provider.get_tracer("test").start_span("Crew Execution")
|
||||||
|
|
||||||
|
telemetry.crew_failed(crew, "RuntimeError")
|
||||||
|
|
||||||
|
finished = exp.get_finished_spans()
|
||||||
|
assert len(finished) == 1, "span was never ended - it would never be exported"
|
||||||
|
assert finished[0].status.status_code is StatusCode.ERROR
|
||||||
|
assert finished[0].attributes["error_type"] == "RuntimeError"
|
||||||
|
assert crew._execution_span is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_crew_failed_is_safe_when_no_span_exists():
|
||||||
|
"""share_crew=False crews have no execution span; this must not raise."""
|
||||||
|
from crewai.telemetry.telemetry import Telemetry
|
||||||
|
|
||||||
|
telemetry = Telemetry()
|
||||||
|
telemetry.ready = True
|
||||||
|
|
||||||
|
crew = Mock()
|
||||||
|
crew._execution_span = None
|
||||||
|
|
||||||
|
telemetry.crew_failed(crew, "RuntimeError")
|
||||||
|
|
||||||
|
|
||||||
|
def test_task_failed_event_carries_error_type():
|
||||||
|
"""The exception class must reach the event without the message."""
|
||||||
|
from crewai.events.types.task_events import TaskFailedEvent
|
||||||
|
|
||||||
|
try:
|
||||||
|
raise TimeoutError("request to gpt-4o timed out after 60s")
|
||||||
|
except TimeoutError as e:
|
||||||
|
event = TaskFailedEvent(error=str(e), error_type=type(e).__name__, task=None)
|
||||||
|
|
||||||
|
assert event.error_type == "TimeoutError"
|
||||||
|
assert "gpt-4o" not in event.error_type
|
||||||
|
|
||||||
|
|
||||||
|
def test_crew_kickoff_failed_event_carries_error_type():
|
||||||
|
from crewai.events.types.crew_events import CrewKickoffFailedEvent
|
||||||
|
|
||||||
|
try:
|
||||||
|
raise ValueError("bad input: {'api_key': 'sk-live-1234'}")
|
||||||
|
except ValueError as e:
|
||||||
|
event = CrewKickoffFailedEvent(
|
||||||
|
error=str(e), error_type=type(e).__name__, crew_name="TestCrew"
|
||||||
|
)
|
||||||
|
|
||||||
|
assert event.error_type == "ValueError"
|
||||||
|
assert "sk-live" not in event.error_type
|
||||||
|
|
||||||
|
|
||||||
|
def test_error_type_defaults_to_none_for_backwards_compatibility():
|
||||||
|
"""Existing callers that omit error_type must keep working."""
|
||||||
|
from crewai.events.types.crew_events import CrewKickoffFailedEvent
|
||||||
|
from crewai.events.types.task_events import TaskFailedEvent
|
||||||
|
|
||||||
|
assert TaskFailedEvent(error="boom", task=None).error_type is None
|
||||||
|
assert CrewKickoffFailedEvent(error="boom", crew_name="C").error_type is None
|
||||||
Reference in New Issue
Block a user