fix(flow): scope outcome and resume signals to user flows

Two findings from review, both confirmed against the code.

Outcome features counted CrewAI's own flows. The agent executor, memory
encoding and memory recall are all Flows and all set suppress_flow_events;
they run far more often than anything a user wrote, so flow:completed,
flow:failed and flow:method_failed were mostly bookkeeping. Those three are now
emitted only for flows the caller wrote. Internal outcomes are still recorded
on the Flow Completed span, which carries origin.

flow:resumed counted checkpoint restores. _is_execution_resuming is set both by
from_pending (a human pause) and by a checkpoint restore that never paused for
anyone, so resumes could exceed pauses and the abandonment rate was unusable.
Keyed off _pending_feedback_context instead, which only from_pending sets.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ASfWmW3RGy4qAQm6s8U9jH
This commit is contained in:
Joao Moura
2026-08-11 15:55:39 -07:00
parent bda1bd4026
commit e0d0bde9a6
2 changed files with 88 additions and 7 deletions

View File

@@ -308,6 +308,17 @@ class EventListener(BaseEventListener):
def on_flow_created(_: Any, event: FlowCreatedEvent) -> None:
self._telemetry.flow_creation_span(event.flow_name)
def _is_infrastructure_flow(source: Any) -> bool:
"""Flows CrewAI runs for its own bookkeeping.
The agent executor, memory encoding and memory recall are all Flows
and all set ``suppress_flow_events``. They run far more often than
anything a user wrote, so counting their outcomes in the same
feature as user flows makes that feature meaningless. Their outcome
is still recorded on the Flow Completed span, which carries origin.
"""
return bool(getattr(source, "suppress_flow_events", False))
def _flow_origin(source: Any) -> str:
"""Separate CrewAI's own flows from the caller's.
@@ -343,10 +354,12 @@ class EventListener(BaseEventListener):
event.flow_name, list(source._methods.keys()), _flow_origin(source)
)
source._telemetry_started_at = time.monotonic()
if getattr(source, "_is_execution_resuming", False):
if getattr(source, "_pending_feedback_context", None) is not None:
# No resume event exists, so a run restored from a pause is only
# visible here. Without it, paused flows can be counted but
# abandoned ones cannot be told apart from resumed ones.
# visible here. Keyed off the pending-feedback context rather
# than _is_execution_resuming, which is also set by checkpoint
# restores that never paused for a human - counting those would
# inflate resumes past pauses and break the abandonment rate.
self._telemetry.feature_usage_span("flow:resumed")
if not getattr(source, "suppress_flow_events", False):
self.formatter.handle_flow_created(event.flow_name, str(source.flow_id))
@@ -354,7 +367,8 @@ class EventListener(BaseEventListener):
@crewai_event_bus.on(FlowFinishedEvent)
def on_flow_finished(source: Any, event: FlowFinishedEvent) -> None:
self._telemetry.feature_usage_span("flow:completed")
if not _is_infrastructure_flow(source):
self._telemetry.feature_usage_span("flow:completed")
_report_flow_duration(source, event.flow_name, "completed")
if not getattr(source, "suppress_flow_events", False):
@@ -365,7 +379,8 @@ class EventListener(BaseEventListener):
@crewai_event_bus.on(FlowFailedEvent)
def on_flow_failed(source: Any, event: FlowFailedEvent) -> None:
self._telemetry.feature_usage_span("flow:failed")
if not _is_infrastructure_flow(source):
self._telemetry.feature_usage_span("flow:failed")
_report_flow_duration(source, event.flow_name, "failed")
if not getattr(source, "suppress_flow_events", False):
@@ -419,11 +434,12 @@ class EventListener(BaseEventListener):
@crewai_event_bus.on(MethodExecutionFailedEvent)
def on_method_execution_failed(
_: Any, event: MethodExecutionFailedEvent
source: Any, event: MethodExecutionFailedEvent
) -> None:
# The method name is not recorded: it is user-authored and would put
# arbitrary strings in telemetry.
self._telemetry.feature_usage_span("flow:method_failed")
if not _is_infrastructure_flow(source):
self._telemetry.feature_usage_span("flow:method_failed")
self.formatter.handle_method_status(
event.method_name,

View File

@@ -415,3 +415,68 @@ def test_resumed_flow_is_reported(tmp_path, features: list[str]) -> None:
assert "flow:resumed" in features
assert "flow:completed" in features
def test_infrastructure_flows_do_not_pollute_outcome_signals(
features: list[str], durations: list[tuple[str, float, str]]
) -> None:
"""CrewAI's own flows must not be counted as user flow outcomes.
The agent executor, memory encoding and memory recall are all Flows and run
far more often than anything a user wrote. Counting their outcomes in the
same feature would make ``flow:completed`` mostly bookkeeping. Their outcome
is still recorded on the Flow Completed span, which carries ``origin``.
"""
from crewai import Agent, Crew, Task
from crewai.llms.base_llm import BaseLLM
class StubLLM(BaseLLM):
def __init__(self) -> None:
super().__init__(model="stub-model")
def call(self, messages, **kwargs) -> str:
return "Final Answer: done"
def supports_function_calling(self) -> bool:
return False
def supports_stop_words(self) -> bool:
return False
def get_context_window_size(self) -> int:
return 8192
agent = Agent(role="R", goal="G", backstory="B", llm=StubLLM())
task = Task(description="Do it", expected_output="A result", agent=agent)
Crew(agents=[agent], tasks=[task]).kickoff()
assert "flow:completed" not in features
assert ("AgentExecutor", "completed") in [
(name, outcome) for name, _duration, outcome in durations
]
def test_a_checkpoint_restore_is_not_counted_as_a_resume(
features: list[str],
) -> None:
"""Only a run restored from a human pause counts as resumed.
``_is_execution_resuming`` is also set by checkpoint restores that never
paused for anyone. Counting those would push resumes above pauses and make
the abandonment rate meaningless.
"""
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.flow_events import FlowStartedEvent
class RestoredFlow(Flow):
@start()
def go(self) -> str:
return "ok"
flow = RestoredFlow()
flow._is_execution_resuming = True
assert flow._pending_feedback_context is None
crewai_event_bus.emit(flow, FlowStartedEvent(flow_name="RestoredFlow"))
assert "flow:resumed" not in features