Files
crewAI/lib/crewai/tests/telemetry/test_flow_telemetry.py
Joao Moura bda1bd4026 feat(flow): tag flow origin and report resumed runs
Two gaps found while testing the pause/resume path end to end.

Resumed runs were invisible. There is no resume event: a restored run re-enters
through kickoff(), so it looked identical to a fresh start. flow:resumed is
derived from _is_execution_resuming at flow start, which makes
flow:paused - flow:resumed the abandonment rate.

Flow counts are dominated by CrewAI's own AgentExecutor, which is itself a Flow
and runs once per agent execution - it is the top flow in the warehouse by a
wide margin. Nothing distinguished it from a user's flows except guessing at the
name. Both Flow Execution and Flow Completed now carry origin: "internal" when
the flow class is defined under crewai.*, "user" otherwise. Tagging only the new
span would have left the existing daily count unsplittable.

Both span methods take origin with a default, so their signatures stay
backward compatible.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ASfWmW3RGy4qAQm6s8U9jH
2026-08-11 15:51:18 -07:00

418 lines
13 KiB
Python

"""Flow outcome and human-in-the-loop signals must reach telemetry.
Driven through real ``Flow`` executions rather than by emitting events directly,
so these fail if the event bus, the listener wiring, or the emitting call site
changes - not just if the listener body does.
Before this, a flow reported only that it *started*: ``FlowFinishedEvent``,
``FlowFailedEvent``, ``MethodExecutionFailedEvent``, ``MethodExecutionPausedEvent``
and ``FlowPausedEvent`` all reached the console formatter and stopped there, and
the input and conversation-failure events had no listener at all.
"""
from __future__ import annotations
import contextlib
import time
import pytest
from crewai.flow.async_feedback import HumanFeedbackPending, PendingFeedbackContext
from crewai.flow.flow import Flow, listen, start
from crewai.flow.human_feedback import human_feedback
from crewai.flow.input_provider import InputResponse
def _reregister_listener() -> None:
"""Re-subscribe the global listener to the event bus.
The repo-wide ``cleanup_event_handlers`` fixture clears every handler after
each test, so anything relying on the shared listener sees an empty bus
unless it happens to run first.
"""
from crewai.events import event_listener as listener_module
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.flow_events import FlowStartedEvent
# Only when the bus is empty: subscribing a second time registers a fresh
# set of closures, and every handler then fires twice.
if crewai_event_bus._sync_handlers.get(FlowStartedEvent):
return
listener_module.event_listener.setup_listeners(crewai_event_bus)
@pytest.fixture
def flow_spans(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, str]]:
"""Record (flow_name, origin) for every Flow Execution span."""
from crewai.events import event_listener as listener_module
_reregister_listener()
recorded: list[tuple[str, str]] = []
monkeypatch.setattr(
listener_module.event_listener._telemetry,
"flow_execution_span",
lambda flow_name, node_names, origin="user": recorded.append(
(flow_name, origin)
),
)
return recorded
@pytest.fixture
def durations(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, float, str]]:
"""Record every (flow_name, duration_ms, outcome) the listener reports."""
from crewai.events import event_listener as listener_module
_reregister_listener()
recorded: list[tuple[str, float, str]] = []
monkeypatch.setattr(
listener_module.event_listener._telemetry,
"flow_completed_span",
lambda flow_name, duration_ms, outcome, origin="user": recorded.append(
(flow_name, duration_ms, outcome)
),
)
return recorded
@pytest.fixture
def features(monkeypatch: pytest.MonkeyPatch) -> list[str]:
"""Record every feature the listener reports for a real flow run.
Observes the telemetry boundary rather than exported spans: the suite builds
the Telemetry singleton with collection disabled, so it has no provider to
export through, and replacing that singleton mid-session leaves the event
bus without its handlers. That the recorded features become spans is covered
by ``test_tracer_isolation``.
"""
from crewai.events import event_listener as listener_module
_reregister_listener()
recorded: list[str] = []
monkeypatch.setattr(
listener_module.event_listener._telemetry,
"feature_usage_span",
recorded.append,
)
return recorded
def test_completed_flow_reports_its_outcome(features: list[str]) -> None:
class OkFlow(Flow):
@start()
def go(self) -> str:
return "ok"
OkFlow().kickoff()
assert "flow:completed" in features
def test_failed_flow_reports_the_failure_and_the_method(features: list[str]) -> None:
class BoomFlow(Flow):
@start()
def go(self) -> str:
raise RuntimeError("boom")
with pytest.raises(RuntimeError, match="boom"):
BoomFlow().kickoff()
emitted = features
assert "flow:failed" in emitted
assert "flow:method_failed" in emitted
assert "flow:completed" not in emitted
def test_a_failed_flow_is_still_counted_as_an_execution(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""The start-time span must survive, or aborted runs vanish from counts.
``flow_executions_daily_target`` counts ``Flow Execution`` spans, emitted
when the flow starts. Holding that span open until completion to measure
duration - the obvious way to add duration - would drop every run that never
finishes, so the outcome signals are reported separately instead.
"""
from crewai.events import event_listener as listener_module
_reregister_listener()
started: list[str] = []
monkeypatch.setattr(
listener_module.event_listener._telemetry,
"flow_execution_span",
lambda flow_name, node_names, origin="user": started.append(flow_name),
)
class BoomFlow(Flow):
@start()
def go(self) -> str:
raise RuntimeError("boom")
with pytest.raises(RuntimeError, match="boom"):
BoomFlow().kickoff()
assert "BoomFlow" in started
def test_requesting_input_reports_both_sides(features: list[str]) -> None:
class StubProvider:
def request_input(self, message: str, flow: Flow, metadata=None):
return InputResponse(value="typed answer")
class AskFlow(Flow):
@start()
def go(self) -> str:
return self.ask("What topic?")
AskFlow(input_provider=StubProvider()).kickoff()
emitted = features
assert "flow:input_requested" in emitted
assert "flow:input_received" in emitted
def test_paused_flow_reports_the_pause(features: list[str]) -> None:
"""An async feedback provider pauses the flow; both signals must land."""
class AsyncProvider:
def request_feedback(self, context: PendingFeedbackContext, flow: Flow) -> str:
raise HumanFeedbackPending(context=context)
class PausingFlow(Flow):
@start()
@human_feedback(message="Review:", provider=AsyncProvider())
def generate(self) -> str:
return "content"
@listen(generate)
def process(self, result) -> str:
return f"processed: {result.feedback}"
# Whether the pause surfaces as an exception depends on the persistence
# backend in use; the signals must land either way.
with contextlib.suppress(BaseException):
PausingFlow().kickoff()
emitted = features
assert "flow:hitl_paused" in emitted
assert "flow:paused" in emitted
def test_failed_conversation_turn_is_reported(features: list[str]) -> None:
"""Only completed turns were tracked, so failure rate was unknowable."""
class FailingChat(Flow):
conversational = True
@start()
def begin(self) -> str:
raise RuntimeError("turn exploded")
with pytest.raises(RuntimeError, match="turn exploded"):
FailingChat().handle_turn("hello")
assert "flow:conversation_turn_failed" in features
def test_no_user_authored_strings_are_recorded(features: list[str]) -> None:
"""Method names, flow names and error text must not reach telemetry."""
class SecretNamedFlow(Flow):
@start()
def my_secret_method_name(self) -> str:
raise RuntimeError("secret error detail")
with pytest.raises(RuntimeError, match="secret error detail"):
SecretNamedFlow().kickoff()
emitted = features
assert emitted
for feature in emitted:
assert "my_secret_method_name" not in feature
assert "secret error detail" not in feature
assert "SecretNamedFlow" not in feature
def test_completed_flow_reports_a_real_duration(
durations: list[tuple[str, float, str]],
) -> None:
"""Elapsed time must be measured, not merely present."""
class SlowFlow(Flow):
@start()
def go(self) -> str:
time.sleep(0.05)
return "ok"
SlowFlow().kickoff()
assert len(durations) == 1
flow_name, duration_ms, outcome = durations[0]
assert flow_name == "SlowFlow"
assert outcome == "completed"
assert duration_ms >= 50
def test_failed_flow_reports_its_duration_and_outcome(
durations: list[tuple[str, float, str]],
) -> None:
class SlowBoomFlow(Flow):
@start()
def go(self) -> str:
time.sleep(0.05)
raise RuntimeError("boom")
with pytest.raises(RuntimeError, match="boom"):
SlowBoomFlow().kickoff()
assert len(durations) == 1
flow_name, duration_ms, outcome = durations[0]
assert flow_name == "SlowBoomFlow"
assert outcome == "failed"
assert duration_ms >= 50
def test_no_duration_is_reported_without_a_recorded_start(
durations: list[tuple[str, float, str]],
) -> None:
"""A completion with no observed start reports nothing, and does not raise.
A conversational turn can re-emit completion for a restored run, so the
stamp is genuinely absent rather than impossible.
"""
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.flow_events import FlowFinishedEvent
class NeverStartedFlow(Flow):
@start()
def go(self) -> str:
return "ok"
flow = NeverStartedFlow()
crewai_event_bus.emit(
flow,
FlowFinishedEvent(flow_name="NeverStartedFlow", result="ok", state={}),
)
assert durations == []
def test_duration_is_reported_once_per_run(
durations: list[tuple[str, float, str]],
) -> None:
"""The stamp is cleared on use, so a repeated completion cannot double-count."""
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.flow_events import FlowFinishedEvent
class OkFlow(Flow):
@start()
def go(self) -> str:
return "ok"
flow = OkFlow()
flow.kickoff()
crewai_event_bus.emit(
flow, FlowFinishedEvent(flow_name="OkFlow", result="ok", state={})
)
assert len(durations) == 1
def test_user_authored_flows_are_tagged_as_user(flow_spans) -> None:
class MyOwnFlow(Flow):
@start()
def go(self) -> str:
return "ok"
MyOwnFlow().kickoff()
assert ("MyOwnFlow", "user") in flow_spans
def test_crewais_own_agent_executor_is_tagged_internal(flow_spans) -> None:
"""The agent executor is a Flow and runs once per agent execution.
Without an origin tag it is indistinguishable from a user's flows in the
daily counts, and it dominates them.
"""
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()
origins = {name: origin for name, origin in flow_spans}
assert origins.get("AgentExecutor") == "internal"
def test_resumed_flow_is_reported(tmp_path, features: list[str]) -> None:
"""A restored run is only visible here - there is no resume event.
Without it, a paused flow that was abandoned cannot be told apart from one
the user came back to.
"""
from pydantic import BaseModel
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.flow_events import FlowPausedEvent
from crewai.flow.persistence.sqlite import SQLiteFlowPersistence
persistence = SQLiteFlowPersistence(str(tmp_path / "flows.db"))
class State(BaseModel):
id: str = "resume-test-1"
class AsyncProvider:
def request_feedback(self, context: PendingFeedbackContext, flow: Flow) -> str:
raise HumanFeedbackPending(context=context)
class ReviewFlow(Flow[State]):
@start()
@human_feedback(message="Review:", provider=AsyncProvider())
def draft(self) -> str:
return "draft"
@listen(draft)
def finish(self, result) -> str:
return f"final: {result.feedback}"
paused: dict[str, str] = {}
@crewai_event_bus.on(FlowPausedEvent)
def _capture(source, event) -> None:
paused["flow_id"] = event.flow_id
with contextlib.suppress(BaseException):
ReviewFlow(persistence=persistence).kickoff()
assert "flow:paused" in features
assert "flow:resumed" not in features
flow = ReviewFlow.from_pending(paused["flow_id"], persistence)
flow.resume("looks good")
assert "flow:resumed" in features
assert "flow:completed" in features