mirror of
https://github.com/crewAIInc/crewAI.git
synced 2026-08-13 01:38:41 +00:00
* fix(telemetry): stop a failed turn from marking the next one failed
A conversational session that opts out of deferred finalization ends each
turn with its own FlowFailedEvent, emitted inside kickoff() before
handle_turn() emits ConversationTurnFailedEvent. The flag was therefore set
after the run that owned it had already cleared it, survived on the
instance, and reported the next healthy turn as failed.
Gate the flag on the run still having its start stamp: a deferring session
keeps it (no per-turn terminal event), so it still reports a failed turn at
session end.
Also aligns the Flow Lifecycle Signals privacy row with the rest of the
telemetry table, which qualifies every user-authored field it records with
"should not include personal info", and fixes a telemetry test that built
InputResponse with an unsupported `value` keyword - ask() swallowed the
TypeError, so the test asserted the signals while exercising the
provider-error path.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* test(telemetry): cover the streamed turn emitter of the failure flag
stream_turn() is the second emitter of ConversationTurnFailedEvent and
leaks the same flag as handle_turn(). Both regression tests fail on
77c68bd with ['failed', 'failed'].
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
764 lines
24 KiB
Python
764 lines
24 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.experimental.conversational import ConversationConfig
|
|
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
|
|
|
|
from ..utils import wait_for_event_handlers
|
|
|
|
|
|
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", resumed=False, conversational=False: recorded.append(
|
|
(flow_name, origin)
|
|
),
|
|
)
|
|
return recorded
|
|
|
|
|
|
@pytest.fixture
|
|
def starts(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, bool]]:
|
|
"""Record (flow_name, resumed) for every Flow Execution span."""
|
|
from crewai.events import event_listener as listener_module
|
|
|
|
_reregister_listener()
|
|
|
|
recorded: list[tuple[str, bool]] = []
|
|
monkeypatch.setattr(
|
|
listener_module.event_listener._telemetry,
|
|
"flow_execution_span",
|
|
lambda flow_name, node_names, origin="user", resumed=False, conversational=False: recorded.append(
|
|
(flow_name, resumed)
|
|
),
|
|
)
|
|
return recorded
|
|
|
|
|
|
@pytest.fixture
|
|
def conversational_marks(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, bool]]:
|
|
"""Record (flow_name, conversational) for every Flow Execution span."""
|
|
from crewai.events import event_listener as listener_module
|
|
|
|
_reregister_listener()
|
|
|
|
recorded: list[tuple[str, bool]] = []
|
|
monkeypatch.setattr(
|
|
listener_module.event_listener._telemetry,
|
|
"flow_execution_span",
|
|
lambda flow_name,
|
|
node_names,
|
|
origin="user",
|
|
resumed=False,
|
|
conversational=False: recorded.append((flow_name, conversational)),
|
|
)
|
|
return recorded
|
|
|
|
|
|
@pytest.fixture
|
|
def pauses(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, str]]:
|
|
"""Record (flow_name, origin) for every Flow Paused 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_paused_span",
|
|
lambda flow_name, origin="user": recorded.append((flow_name, origin)),
|
|
)
|
|
return recorded
|
|
|
|
|
|
@pytest.fixture
|
|
def method_failures(monkeypatch: pytest.MonkeyPatch) -> list[tuple[str, str]]:
|
|
"""Record (flow_name, origin) for every Flow Method Failed 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_method_failed_span",
|
|
lambda flow_name, 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", conversational=False: 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(
|
|
durations: list[tuple[str, float, str]],
|
|
) -> None:
|
|
"""Outcome is a lifecycle fact, so it belongs on a span, not a feature."""
|
|
|
|
class OkFlow(Flow):
|
|
@start()
|
|
def go(self) -> str:
|
|
return "ok"
|
|
|
|
OkFlow().kickoff()
|
|
|
|
assert [(n, o) for n, _d, o in durations] == [("OkFlow", "completed")]
|
|
|
|
|
|
def test_failed_flow_reports_the_failure_and_the_method(
|
|
durations: list[tuple[str, float, str]], method_failures: list[tuple[str, str]]
|
|
) -> None:
|
|
class BoomFlow(Flow):
|
|
@start()
|
|
def go(self) -> str:
|
|
raise RuntimeError("boom")
|
|
|
|
with pytest.raises(RuntimeError, match="boom"):
|
|
BoomFlow().kickoff()
|
|
|
|
assert [(n, o) for n, _d, o in durations] == [("BoomFlow", "failed")]
|
|
assert ("BoomFlow", "user") in method_failures
|
|
|
|
|
|
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", resumed=False, conversational=False: 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(text="typed answer")
|
|
|
|
class AskFlow(Flow):
|
|
@start()
|
|
def go(self) -> str:
|
|
return self.ask("What topic?")
|
|
|
|
# ask() swallows provider errors and returns None, so the answer is
|
|
# asserted too: a provider that raises would otherwise still emit both
|
|
# signals and pass this test.
|
|
assert AskFlow(input_provider=StubProvider()).kickoff() == "typed answer"
|
|
|
|
emitted = features
|
|
assert "flow:input_requested" in emitted
|
|
assert "flow:input_received" in emitted
|
|
|
|
|
|
def test_paused_flow_reports_the_pause(
|
|
features: list[str], pauses: list[tuple[str, 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()
|
|
|
|
# The pause itself is lifecycle and lands on a span; that a human-feedback
|
|
# method was what paused is genuine feature adoption.
|
|
assert ("PausingFlow", "user") in pauses
|
|
assert "flow:hitl_paused" in features
|
|
|
|
|
|
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_method_names_or_error_text_are_recorded(
|
|
method_failures: list[tuple[str, str]],
|
|
durations: list[tuple[str, float, str]],
|
|
features: list[str],
|
|
) -> None:
|
|
"""Method names and error text are user-authored and must not be sent.
|
|
|
|
The flow name is recorded, as it already is for flow creation and
|
|
execution, so it is deliberately not asserted against here.
|
|
"""
|
|
|
|
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()
|
|
|
|
assert method_failures, "the failure must still be reported"
|
|
recorded = [
|
|
str(value)
|
|
for row in (*method_failures, *durations)
|
|
for value in row
|
|
] + features
|
|
for value in recorded:
|
|
assert "my_secret_method_name" not in value
|
|
assert "secret error detail" not in value
|
|
|
|
|
|
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,
|
|
pauses: list[tuple[str, str]],
|
|
starts: list[tuple[str, bool]],
|
|
durations: list[tuple[str, float, 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 ("ReviewFlow", "user") in pauses
|
|
assert starts == [("ReviewFlow", False)]
|
|
|
|
flow = ReviewFlow.from_pending(paused["flow_id"], persistence)
|
|
flow.resume("looks good")
|
|
|
|
assert ("ReviewFlow", True) in starts
|
|
assert ("ReviewFlow", "completed") in [(n, o) for n, _d, o in durations]
|
|
|
|
|
|
def test_a_user_flow_that_suppresses_console_events_still_reports(
|
|
durations: list[tuple[str, float, str]],
|
|
) -> None:
|
|
"""``suppress_flow_events`` asks for console quiet, not for no telemetry."""
|
|
|
|
class QuietFlow(Flow):
|
|
suppress_flow_events: bool = True
|
|
|
|
@start()
|
|
def go(self) -> str:
|
|
return "ok"
|
|
|
|
QuietFlow().kickoff()
|
|
|
|
assert [(n, o) for n, _d, o in durations] == [("QuietFlow", "completed")]
|
|
|
|
|
|
def test_a_declarative_flow_is_not_treated_as_internal(
|
|
flow_spans: list[tuple[str, str]],
|
|
) -> None:
|
|
"""``Flow.from_declaration()`` yields a ``Flow``, defined inside crewai.
|
|
|
|
Deciding origin from the defining module would report a caller's
|
|
declarative flow as one of CrewAI's own.
|
|
"""
|
|
flow = Flow.from_declaration(contents={"name": "MyDeclarativeFlow"})
|
|
|
|
assert getattr(type(flow), "is_crewai_internal", False) is False
|
|
|
|
|
|
def test_a_failed_conversation_session_is_not_reported_completed(
|
|
features: list[str], durations: list[tuple[str, float, str]]
|
|
) -> None:
|
|
"""A conversational session closes with FlowFinishedEvent either way.
|
|
|
|
Reading that event at face value counted a failed session as a success,
|
|
alongside the turn-failure signal.
|
|
"""
|
|
|
|
class FailingChat(Flow):
|
|
conversational = True
|
|
|
|
@start()
|
|
def begin(self) -> str:
|
|
raise RuntimeError("turn exploded")
|
|
|
|
chat = FailingChat()
|
|
with pytest.raises(RuntimeError, match="turn exploded"):
|
|
chat.handle_turn("hello")
|
|
chat.finalize_session_traces()
|
|
|
|
assert "flow:conversation_turn_failed" in features
|
|
assert all(outcome != "completed" for _n, _d, outcome in durations)
|
|
|
|
|
|
def test_a_deferred_session_still_reports_a_failed_turn(
|
|
durations: list[tuple[str, float, str]],
|
|
) -> None:
|
|
"""A deferring session has no per-turn terminal event to carry the failure.
|
|
|
|
Its only outcome span is the one ``finalize_session_traces()`` triggers, so
|
|
the turn-failure flag is what makes that span say ``failed``. Deferral is
|
|
the default for a conversational flow, so this is the common path.
|
|
"""
|
|
|
|
class DeferringChat(Flow):
|
|
conversational = True
|
|
|
|
@start()
|
|
def begin(self) -> str:
|
|
raise RuntimeError("turn exploded")
|
|
|
|
chat = DeferringChat()
|
|
with pytest.raises(RuntimeError, match="turn exploded"):
|
|
chat.handle_turn("hello")
|
|
chat.finalize_session_traces()
|
|
# finalize_session_traces() emits without awaiting its handlers.
|
|
wait_for_event_handlers()
|
|
|
|
assert [outcome for _n, _d, outcome in durations] == ["failed"]
|
|
|
|
|
|
def test_a_failed_turn_does_not_mark_the_next_turn_failed(
|
|
durations: list[tuple[str, float, str]],
|
|
) -> None:
|
|
"""A session that opts out of deferral ends each turn with its own event.
|
|
|
|
That terminal event fires inside ``kickoff()``, before ``handle_turn()``
|
|
emits the turn-failure event, so the flag was set after the run that owned
|
|
it had already cleared it - and the next healthy turn read it as failed.
|
|
"""
|
|
|
|
turns: list[str] = []
|
|
|
|
@ConversationConfig(defer_trace_finalization=False)
|
|
class FlakyChat(Flow):
|
|
conversational = True
|
|
|
|
@start()
|
|
def begin(self) -> str:
|
|
turns.append("turn")
|
|
if len(turns) == 1:
|
|
raise RuntimeError("turn exploded")
|
|
return "second turn is fine"
|
|
|
|
chat = FlakyChat()
|
|
with pytest.raises(RuntimeError, match="turn exploded"):
|
|
chat.handle_turn("hello")
|
|
chat.handle_turn("again")
|
|
|
|
assert [outcome for _n, _d, outcome in durations] == ["failed", "completed"]
|
|
|
|
|
|
def test_a_failed_streamed_turn_does_not_mark_the_next_turn_failed(
|
|
durations: list[tuple[str, float, str]],
|
|
) -> None:
|
|
"""``stream_turn`` is the other emitter of the turn-failure event.
|
|
|
|
It emits from its own ``except`` block, after ``kickoff()`` has closed the
|
|
run out, so it leaks the same flag as the non-streamed path.
|
|
"""
|
|
|
|
turns: list[str] = []
|
|
|
|
@ConversationConfig(defer_trace_finalization=False)
|
|
class FlakyStreamingChat(Flow):
|
|
conversational = True
|
|
|
|
@start()
|
|
def begin(self) -> str:
|
|
turns.append("turn")
|
|
if len(turns) == 1:
|
|
raise RuntimeError("turn exploded")
|
|
return "second turn is fine"
|
|
|
|
chat = FlakyStreamingChat()
|
|
with pytest.raises(RuntimeError, match="turn exploded"):
|
|
list(chat.stream_turn("hello").events)
|
|
list(chat.stream_turn("again").events)
|
|
|
|
assert [outcome for _n, _d, outcome in durations] == ["failed", "completed"]
|
|
|
|
|
|
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()
|
|
|
|
# Internal outcomes are still recorded - on the span, tagged internal -
|
|
# they simply do not masquerade as a user's flow finishing.
|
|
assert ("AgentExecutor", "completed") in [
|
|
(name, outcome) for name, _duration, outcome in durations
|
|
]
|
|
assert "flow:completed" not in features
|
|
|
|
|
|
def test_a_checkpoint_restore_is_not_counted_as_a_resume(
|
|
starts: list[tuple[str, bool]],
|
|
) -> None:
|
|
"""Only a run restored from a human pause is marked 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"))
|
|
wait_for_event_handlers()
|
|
|
|
assert starts == [("RestoredFlow", False)]
|
|
|
|
|
|
def test_a_conversational_turn_is_marked(
|
|
conversational_marks: list[tuple[str, bool]],
|
|
) -> None:
|
|
"""Each turn is its own kickoff, but a session reports one completion.
|
|
|
|
Without the marker those spans run many-to-one against Flow Completed and
|
|
silently drag any completion rate computed across all flows.
|
|
"""
|
|
|
|
class Chatty(Flow):
|
|
conversational = True
|
|
|
|
@start()
|
|
def begin(self) -> str:
|
|
return "hi"
|
|
|
|
Chatty().handle_turn("hello")
|
|
|
|
assert ("Chatty", True) in conversational_marks
|
|
|
|
|
|
def test_an_ordinary_flow_is_not_marked_conversational(
|
|
conversational_marks: list[tuple[str, bool]],
|
|
) -> None:
|
|
class PlainFlow(Flow):
|
|
@start()
|
|
def go(self) -> str:
|
|
return "ok"
|
|
|
|
PlainFlow().kickoff()
|
|
|
|
assert ("PlainFlow", False) in conversational_marks
|