From e0d0bde9a66deefa72345f5f0d980ad24ad2ff32 Mon Sep 17 00:00:00 2001 From: Joao Moura Date: Tue, 11 Aug 2026 15:55:39 -0700 Subject: [PATCH] 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) Claude-Session: https://claude.ai/code/session_01ASfWmW3RGy4qAQm6s8U9jH --- .../src/crewai/events/event_listener.py | 30 +++++++-- .../tests/telemetry/test_flow_telemetry.py | 65 +++++++++++++++++++ 2 files changed, 88 insertions(+), 7 deletions(-) diff --git a/lib/crewai/src/crewai/events/event_listener.py b/lib/crewai/src/crewai/events/event_listener.py index 532bbe0314..1cb7b98dd7 100644 --- a/lib/crewai/src/crewai/events/event_listener.py +++ b/lib/crewai/src/crewai/events/event_listener.py @@ -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, diff --git a/lib/crewai/tests/telemetry/test_flow_telemetry.py b/lib/crewai/tests/telemetry/test_flow_telemetry.py index c404ad7d51..8fc78d573b 100644 --- a/lib/crewai/tests/telemetry/test_flow_telemetry.py +++ b/lib/crewai/tests/telemetry/test_flow_telemetry.py @@ -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