Merge remote-tracking branch 'origin/main' into joao/epd-181-conversational-flow-handle_turn-double-appends-the-reply

# Conflicts:
#	lib/crewai/src/crewai/experimental/conversational_mixin.py
#	lib/crewai/tests/test_flow_conversation.py
This commit is contained in:
Joao Moura
2026-07-10 18:07:07 -07:00
24 changed files with 2051 additions and 69 deletions

View File

@@ -568,12 +568,32 @@ FooterKey .footer-key--key {
self._default_inputs: dict[str, Any] | None = None
self._crew_result: Any = None
self._crew_json_path: Any = None
# Declarative-flow execution state. A flow renders per-method "STEPS"
# (built from flow method events) instead of the crew task list.
self._flow_inputs: dict[str, Any] | None = None
self._flow_method_types: dict[str, str] = {}
self._flow_steps: list[dict[str, Any]] = []
self._current_method: str | None = None
self._elapsed_frozen: float | None = None
self._want_deploy: bool = False
self._trace_url: str | None = None
self._consent_screen: TraceConsentScreen | None = None
self._telemetry: Telemetry | None = None
@property
def _is_flow_run(self) -> bool:
"""True for a non-conversational declarative flow (the STEPS view).
Gates every flow-specific rendering branch so crew and conversational
paths stay byte-identical.
"""
return self._flow is not None and not self._is_conversational
@property
def _run_noun(self) -> str:
"""User-facing noun for the run — 'flow' for a declarative flow, else 'crew'."""
return "flow" if self._is_flow_run else "crew"
# ── Layout ──────────────────────────────────────────────
def compose(self) -> ComposeResult:
@@ -602,6 +622,8 @@ FooterKey .footer-key--key {
self._tick_timer = self.set_interval(1 / 8, self._tick)
if self._is_conversational and self._flow:
self._start_conversational_session()
elif self._flow:
self._run_flow_worker()
elif self._crew:
self._run_crew_worker()
elif self._crew_json_path:
@@ -681,6 +703,49 @@ FooterKey .footer-key--key {
except Exception as e:
self.call_from_thread(self._on_crew_failed, str(e))
@work(thread=True, exclusive=True, group="flow")
def _run_flow_worker(self) -> None:
from crewai.events.listeners.tracing.utils import (
set_suppress_tracing_messages,
set_tui_mode,
)
set_tui_mode(True)
set_suppress_tracing_messages(True)
try:
# A declarative flow returns either a CrewOutput (has ``.raw``) or a
# bare value (str/dict/pydantic); _stringify_output handles both.
result = self._flow.kickoff(inputs=self._flow_inputs)
output = self._stringify_output(result)
with self._lock:
self._crew_result = result
self.call_from_thread(self._on_crew_done, output)
except Exception as e:
self.call_from_thread(self._on_crew_failed, str(e))
def _set_flow_step_status(self, name: str, status: str) -> None:
"""Update a flow method step's status. Caller must hold ``self._lock``."""
for step in self._flow_steps:
if step["name"] == name:
step["status"] = status
return
def _clear_current_method(self, finished_name: str) -> None:
"""Drop the header's active method once it ends. Caller holds the lock.
Falls back to another still-active step (methods can overlap) so the
header never keeps spinning a method the STEPS list already shows as
done or failed.
"""
if self._current_method != finished_name:
return
self._current_method = next(
(s["name"] for s in self._flow_steps if s["status"] == "active"), None
)
# The active method changed; drop its agent so the header doesn't show a
# stale agent until the next method's agent event arrives.
self._current_agent = ""
def _on_crew_done(self, output: str | None) -> None:
with self._lock:
self._status = "completed"
@@ -694,13 +759,18 @@ FooterKey .footer-key--key {
for k in self._task_statuses:
if self._task_statuses[k] == "active":
self._task_statuses[k] = "done"
for step in self._flow_steps:
if step["status"] == "active":
step["status"] = "done"
now = time.time()
for entry in self._log_entries:
if entry["status"] == "running":
if entry["tool_name"] == "memory_save":
continue
entry["status"] = "timeout"
entry["error"] = "No result received before crew completed"
entry["error"] = (
f"No result received before {self._run_noun} completed"
)
entry["duration"] = now - entry["start_time"]
try:
from crewai.events.listeners.tracing.trace_listener import (
@@ -739,13 +809,18 @@ FooterKey .footer-key--key {
self._is_streaming = False
self._current_step = None
self._elapsed_frozen = time.time() - self._start_time
for step in self._flow_steps:
if step["status"] == "active":
step["status"] = "failed"
now = time.time()
for entry in self._log_entries:
if entry["status"] == "running":
if entry["tool_name"] == "memory_save":
continue
entry["status"] = "error"
entry["error"] = "No result received before crew failed"
entry["error"] = (
f"No result received before {self._run_noun} failed"
)
entry["duration"] = now - entry["start_time"]
self._tick()
self.call_later(self._focus_activity_log)
@@ -1156,6 +1231,45 @@ FooterKey .footer-key--key {
widget.update(t)
return
if self._is_flow_run:
t.append(" STEPS\n", style=f"bold {_C_PRIMARY}")
t.append("\n")
if not self._flow_steps:
t.append(" ○ waiting…\n", style=_C_DIM)
for step in self._flow_steps:
name = step["name"]
max_name = sidebar_width - 6
if len(name) > max_name:
name = name[: max_name - 1] + ""
status = step.get("status", "pending")
if status == "done":
t.append("", style=_C_GREEN)
t.append(name, style=_C_DIM)
elif status == "active":
t.append(f" {self._spinner()} ", style=_C_PRIMARY)
t.append(name, style=f"bold {_C_TEXT}")
elif status == "failed":
t.append("", style=_C_RED)
t.append(name, style=_C_RED)
elif status == "paused":
t.append("", style=_C_TEAL)
t.append(name, style=_C_TEAL)
else:
t.append("", style=_C_DIM)
t.append(name, style=_C_DIM)
if step.get("call_type"):
t.append(f" ({step['call_type']})", style=_C_DIM)
t.append("\n")
t.append("\n")
t.append(" TOKENS\n", style=f"bold {_C_PRIMARY}")
t.append("\n")
out = self._output_tokens + self._live_out_tokens
t.append(f"{self._input_tokens:,}\n", style=_C_DIM)
t.append(f"{out:,}\n", style=_C_DIM)
widget.update(t)
return
t.append(" TASKS\n", style=f"bold {_C_PRIMARY}")
t.append("\n")
@@ -1225,6 +1339,55 @@ FooterKey .footer-key--key {
widget.update(t)
return
if self._is_flow_run:
if self._status == "completed":
elapsed = self._elapsed_frozen or (time.time() - self._start_time)
t.append("", style=f"bold {_C_GREEN}")
t.append("Flow complete", style=f"bold {_C_GREEN}")
t.append(f" {elapsed:.1f}s", style=_C_DIM)
out = self._output_tokens + self._live_out_tokens
parts = []
if self._input_tokens:
parts.append(f"{self._input_tokens:,}")
if out:
parts.append(f"{out:,}")
if parts:
t.append(f" {' '.join(parts)} tokens", style=_C_DIM)
elif self._status == "failed":
t.append("", style=f"bold {_C_RED}")
t.append("Failed", style=f"bold {_C_RED}")
if self._error:
t.append(f"\n{self._error[:120]}", style=_C_RED)
elif self._current_method:
paused = any(
s["name"] == self._current_method and s["status"] == "paused"
for s in self._flow_steps
)
if paused:
t.append("", style=_C_TEAL)
t.append(self._current_method, style=f"bold {_C_TEAL}")
else:
t.append(f"{self._spinner()} ", style=_C_PRIMARY)
t.append(self._current_method, style=f"bold {_C_PRIMARY}")
call_type = self._flow_method_types.get(self._current_method)
if call_type:
t.append(f" ({call_type})", style=_C_DIM)
if paused:
t.append(" waiting for feedback", style=_C_DIM)
elif self._current_agent:
t.append("\nAgent: ", style=_C_DIM)
t.append(self._current_agent, style=f"bold {_C_TEXT}")
else:
t.append(f"{self._spinner()} ", style=_C_PRIMARY)
# "Working…" once a step has run (between/after methods);
# "Starting flow…" only before the first method.
t.append(
"Working…" if self._flow_steps else "Starting flow…",
style=_C_DIM,
)
widget.update(t)
return
if self._status == "completed":
elapsed = self._elapsed_frozen or (time.time() - self._start_time)
t.append("", style=f"bold {_C_GREEN}")
@@ -1839,6 +2002,13 @@ FooterKey .footer-key--key {
def _subscribe(self) -> None:
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.crew_events import CrewKickoffStartedEvent
from crewai.events.types.flow_events import (
FlowStartedEvent,
MethodExecutionFailedEvent,
MethodExecutionFinishedEvent,
MethodExecutionPausedEvent,
MethodExecutionStartedEvent,
)
from crewai.events.types.llm_events import (
LLMCallCompletedEvent,
LLMCallStartedEvent,
@@ -1872,13 +2042,74 @@ FooterKey .footer-key--key {
@crewai_event_bus.on(CrewKickoffStartedEvent)
def on_crew_started(source: Any, event: CrewKickoffStartedEvent) -> None:
with self._lock:
if event.crew_name:
# In flow mode the app is named for the flow; a nested crew's
# kickoff (a `call: crew` step) must not rename it.
if event.crew_name and not self._is_flow_run:
self._crew_name = event.crew_name
self.title = f"CrewAI — {event.crew_name}"
self._status = "working"
self._register_handler(CrewKickoffStartedEvent, on_crew_started)
# ── Declarative-flow method events → STEPS panel ────────
@crewai_event_bus.on(FlowStartedEvent)
def on_flow_started(source: Any, event: FlowStartedEvent) -> None:
with self._lock:
self._status = "working"
self._register_handler(FlowStartedEvent, on_flow_started)
@crewai_event_bus.on(MethodExecutionStartedEvent)
def on_method_started(source: Any, event: MethodExecutionStartedEvent) -> None:
with self._lock:
name = event.method_name
self._current_method = name
# Agent is per-method; clear it so the header doesn't show the
# previous method's agent until a new agent event arrives.
self._current_agent = ""
for step in self._flow_steps:
if step["name"] == name:
step["status"] = "active"
break
else:
self._flow_steps.append(
{
"name": name,
"call_type": self._flow_method_types.get(name),
"status": "active",
}
)
self._register_handler(MethodExecutionStartedEvent, on_method_started)
@crewai_event_bus.on(MethodExecutionFinishedEvent)
def on_method_finished(
source: Any, event: MethodExecutionFinishedEvent
) -> None:
with self._lock:
self._set_flow_step_status(event.method_name, "done")
self._clear_current_method(event.method_name)
self._register_handler(MethodExecutionFinishedEvent, on_method_finished)
@crewai_event_bus.on(MethodExecutionFailedEvent)
def on_method_failed(source: Any, event: MethodExecutionFailedEvent) -> None:
with self._lock:
self._set_flow_step_status(event.method_name, "failed")
self._clear_current_method(event.method_name)
self._register_handler(MethodExecutionFailedEvent, on_method_failed)
@crewai_event_bus.on(MethodExecutionPausedEvent)
def on_method_paused(source: Any, event: MethodExecutionPausedEvent) -> None:
# A @human_feedback method paused; flow status panels are suppressed
# in TUI mode, so surface the wait in STEPS/header instead of leaving
# a spinner. _current_method stays pointed at it.
with self._lock:
self._set_flow_step_status(event.method_name, "paused")
self._register_handler(MethodExecutionPausedEvent, on_method_paused)
@crewai_event_bus.on(TaskStartedEvent)
def on_task_started(source: Any, event: TaskStartedEvent) -> None:
with self._lock:

View File

@@ -1,9 +1,10 @@
from __future__ import annotations
import json
import logging
from pathlib import Path
import subprocess
from typing import Any
from typing import TYPE_CHECKING, Any
import click
from crewai_core.project import ProjectDefinitionError, configured_project_definition
@@ -18,6 +19,13 @@ from crewai_cli.input_prompt import (
from crewai_cli.utils import build_env_with_all_tool_credentials
if TYPE_CHECKING:
from crewai.flow.flow import Flow
logger = logging.getLogger(__name__)
def run_declarative_flow_in_project_env(
definition: str | Path, inputs: str | None = None
) -> None:
@@ -66,17 +74,182 @@ def run_declarative_flow(definition: str | Path, inputs: str | None = None) -> N
flow = load_declarative_flow(definition)
resolved_inputs = _resolve_flow_inputs(flow, provided)
# The TUI is the interactive default. Headless contexts run directly on the
# terminal: deploy/CREWAI_DMN, piped output, CI — anything without an
# interactive TTY. is_interactive() already folds in the CREWAI_DMN check.
# Human-feedback flows also run on the terminal: their methods collect input
# via the flow runtime's blocking input()/Rich prompts (and async feedback
# returns a pending marker rather than completing), neither of which the
# Textual TUI can handle correctly.
if is_interactive() and not _flow_uses_human_feedback(flow):
_run_declarative_flow_tui(flow, resolved_inputs or None)
return
try:
result = flow.kickoff(inputs=resolved_inputs or None)
except Exception as exc:
click.echo(
f"An error occurred while running the declarative flow: {exc}", err=True
f"An error occurred while running the declarative flow: {exc}",
err=True,
)
raise SystemExit(1) from exc
click.echo(_format_result(result))
def _run_declarative_flow_tui(
flow: Flow[Any], resolved_inputs: dict[str, Any] | None
) -> Any:
"""Run a declarative flow on the CrewAI TUI (the interactive default).
Mirrors the declarative-crew TUI contract (``run_crew._run_json_crew``):
a failed flow exits non-zero, a user quit ends the process so in-flight LLM
work stops, and choosing Deploy chains into the deploy command.
"""
import os
import sys
from crewai.events.event_listener import EventListener
from crewai_cli.crew_run_tui import CrewRunApp
# The flow runtime (unlike a Crew constructor) doesn't create the event
# listener, and the TUI's trace/telemetry features depend on it.
EventListener()
# The STEPS panel and header are driven by flow method events. A flow may
# declare ``config.suppress_flow_events`` (a headless/production
# optimization) which would leave STEPS stuck on "waiting…" here — so force
# emission on for the interactive TUI run. The headless path never reaches
# this and keeps the flow's declared setting.
try:
flow.suppress_flow_events = False
except Exception:
logger.debug(
"Could not disable suppress_flow_events for the flow TUI", exc_info=True
)
app = CrewRunApp(crew_name=flow.name or type(flow).__name__)
app._flow = flow
app._flow_inputs = resolved_inputs
app._flow_method_types = _flow_method_types(flow)
app.run()
_print_flow_post_tui_summary(app)
if app._status == "failed":
raise SystemExit(1)
if app._status not in ("completed", "failed"):
# User quit mid-run. kickoff runs in a thread worker that cannot be
# force-cancelled, so end the process to stop in-flight LLM and tool
# work instead of letting it burn tokens in the background.
click.secho("\n Run cancelled.", fg="yellow")
sys.stdout.flush()
os._exit(130)
if getattr(app, "_want_deploy", False):
from crewai_cli.run_crew import _chain_deploy
_chain_deploy()
return app._crew_result
def _flow_uses_human_feedback(flow: Flow[Any]) -> bool:
"""True if any declarative method declares ``@human_feedback``.
Such flows need the flow runtime's interactive stdin / Rich prompts, which
don't compose with Textual — so they run on the terminal, not the TUI.
"""
try:
return any(
method.human_feedback is not None
for method in flow._definition.methods.values()
)
except Exception:
logger.debug("Could not inspect flow for human feedback", exc_info=True)
return False
def _flow_method_types(flow: Flow[Any]) -> dict[str, str]:
"""Map each declarative method name to its ``call`` type (crew/agent/…).
Best-effort: the STEPS panel shows this as a dim label. Method events don't
carry the call type, so it's read from the flow definition up front.
"""
method_types: dict[str, str] = {}
try:
for name, method_definition in flow._definition.methods.items():
method_types[name] = method_definition.do.call
except Exception:
logger.debug("Could not derive flow method types", exc_info=True)
return method_types
def _print_flow_post_tui_summary(app: Any) -> None:
"""Print a compact result panel after the flow TUI exits."""
import time
from rich.console import Console
from rich.markdown import Markdown
from rich.padding import Padding
from rich.panel import Panel
from rich.text import Text
console = Console()
elapsed = (app._elapsed_frozen or (time.time() - app._start_time)) or 0.0
out_tokens = app._output_tokens + app._live_out_tokens
token_parts = []
if app._input_tokens:
token_parts.append(f"{app._input_tokens:,}")
if out_tokens:
token_parts.append(f"{out_tokens:,}")
token_str = " ".join(token_parts)
if token_str:
token_str += " tokens"
crewai_red = "#FF5A50"
crewai_teal = "#1F7982"
if app._status == "completed":
summary = Text()
summary.append(" ✔ Flow complete", style=f"bold {crewai_teal}")
summary.append(f" in {elapsed:.1f}s", style="dim")
if token_str:
summary.append(f" {token_str}", style="dim")
console.print(
Panel(
summary,
title=f" {app._crew_name} ",
title_align="left",
border_style=crewai_teal,
padding=(0, 1),
)
)
if app._final_output:
console.print()
console.print(Text(" Final Result", style=f"bold {crewai_teal}"))
console.print()
console.print(Padding(Markdown(app._final_output), (0, 2)))
elif app._status == "failed":
content = Text()
content.append(" ✘ Failed", style=f"bold {crewai_red}")
content.append(f" after {elapsed:.1f}s\n", style="dim")
if app._error:
content.append(f"\n {app._error}\n", style=crewai_red)
console.print(
Panel(
content,
title=f" {app._crew_name} ",
title_align="left",
border_style=crewai_red,
padding=(0, 1),
)
)
def _resolve_flow_inputs(flow: Any, provided: dict[str, Any]) -> dict[str, Any]:
"""Resolve kickoff inputs from the flow's state schema.

View File

@@ -6,6 +6,14 @@ from unittest.mock import Mock
import pytest
from crewai.events.event_bus import crewai_event_bus
from crewai.events.types.crew_events import CrewKickoffStartedEvent
from crewai.events.types.flow_events import (
FlowStartedEvent,
MethodExecutionFailedEvent,
MethodExecutionFinishedEvent,
MethodExecutionPausedEvent,
MethodExecutionStartedEvent,
)
from crewai.events.types.memory_events import (
MemorySaveCompletedEvent,
MemorySaveFailedEvent,
@@ -959,6 +967,31 @@ async def test_crew_done_does_not_mark_unfinished_tool_successful() -> None:
assert app._plan_step_status == {1: "failed", 2: "done", 3: "done"}
@pytest.mark.asyncio
async def test_flow_done_uses_flow_wording_for_unfinished_tool() -> None:
# The shared completion handler reports "flow" (not "crew") in flow mode.
app = CrewRunApp(crew_name="Demo Flow")
app._flow = SimpleNamespace()
async with app.run_test(size=(100, 40)) as pilot:
app._log_entries = [
{
"tool_name": "search",
"status": "running",
"args": None,
"result": None,
"error": None,
"start_time": time.time() - 2,
"duration": None,
"task_idx": 1,
}
]
app._on_crew_done("final output")
await pilot.pause()
assert app._log_entries[0]["error"] == "No result received before flow completed"
@pytest.mark.asyncio
async def test_crew_done_does_not_timeout_memory_save() -> None:
app = _app_with_plan()
@@ -1481,3 +1514,210 @@ def test_overlapping_task_logs_keep_their_own_state() -> None:
assert any(step.get("summary") == "thinking" for step in entry2["steps"])
finally:
app._unsubscribe()
# ── Declarative-flow (non-conversational) TUI support ───────
def test_is_flow_run_gating() -> None:
"""The flow-render gate must be true only for a non-conversational flow."""
crew_app = CrewRunApp(total_tasks=1)
crew_app._crew = SimpleNamespace()
assert crew_app._is_flow_run is False
conv_app = CrewRunApp(conversational=True)
conv_app._flow = SimpleNamespace()
assert conv_app._is_flow_run is False
flow_app = CrewRunApp()
flow_app._flow = SimpleNamespace()
assert flow_app._is_flow_run is True
def test_flow_method_events_build_steps() -> None:
app = CrewRunApp(crew_name="Demo")
app._flow = SimpleNamespace()
app._flow_method_types = {"research": "crew", "summarize": "agent"}
app._subscribe()
try:
_emit_event(FlowStartedEvent(flow_name="Demo"))
assert app._status == "working"
_emit_event(
MethodExecutionStartedEvent(
flow_name="Demo", method_name="research", state={}
)
)
assert app._flow_steps == [
{"name": "research", "call_type": "crew", "status": "active"}
]
assert app._current_method == "research"
_emit_event(
MethodExecutionFinishedEvent(
flow_name="Demo", method_name="research", result="ok", state={}
)
)
_emit_event(
MethodExecutionStartedEvent(
flow_name="Demo", method_name="summarize", state={}
)
)
_emit_event(
MethodExecutionFailedEvent(
flow_name="Demo",
method_name="summarize",
error=RuntimeError("boom"),
)
)
finally:
app._unsubscribe()
assert app._flow_steps == [
{"name": "research", "call_type": "crew", "status": "done"},
{"name": "summarize", "call_type": "agent", "status": "failed"},
]
# The header must not keep spinning a method that already ended.
assert app._current_method is None
def test_current_method_clears_and_falls_back_across_overlap() -> None:
app = CrewRunApp(crew_name="Demo")
app._flow = SimpleNamespace()
app._subscribe()
try:
_emit_event(
MethodExecutionStartedEvent(flow_name="Demo", method_name="a", state={})
)
_emit_event(
MethodExecutionStartedEvent(flow_name="Demo", method_name="b", state={})
)
assert app._current_method == "b"
# 'a' finishes while 'b' is still active → header stays on 'b'.
_emit_event(
MethodExecutionFinishedEvent(
flow_name="Demo", method_name="a", result=None, state={}
)
)
assert app._current_method == "b"
# 'b' finishes → nothing active left → header clears.
_emit_event(
MethodExecutionFinishedEvent(
flow_name="Demo", method_name="b", result=None, state={}
)
)
assert app._current_method is None
finally:
app._unsubscribe()
def test_flow_method_transitions_clear_current_agent() -> None:
app = CrewRunApp(crew_name="Demo")
app._flow = SimpleNamespace()
app._subscribe()
try:
_emit_event(
MethodExecutionStartedEvent(flow_name="Demo", method_name="a", state={})
)
app._current_agent = "Researcher" # an agent ran during method 'a'
# Starting a new method clears the previous method's agent.
_emit_event(
MethodExecutionStartedEvent(flow_name="Demo", method_name="b", state={})
)
assert app._current_agent == ""
app._current_agent = "Writer"
# 'b' ending switches the active method ('a' still active) → agent clears.
_emit_event(
MethodExecutionFinishedEvent(
flow_name="Demo", method_name="b", result=None, state={}
)
)
assert app._current_agent == ""
finally:
app._unsubscribe()
def test_crew_kickoff_does_not_rename_flow_run() -> None:
# A `call: crew` step must not relabel the flow with the nested crew's name.
app = CrewRunApp(crew_name="My Flow")
app._flow = SimpleNamespace()
app._subscribe()
try:
_emit_event(CrewKickoffStartedEvent(crew_name="Nested Crew", inputs=None))
assert app._crew_name == "My Flow"
assert app._status == "working"
finally:
app._unsubscribe()
def test_crew_kickoff_renames_in_crew_mode() -> None:
# Regression: crew runs still adopt the crew name from the event.
app = CrewRunApp(crew_name="Crew")
app._crew = SimpleNamespace()
app._subscribe()
try:
_emit_event(CrewKickoffStartedEvent(crew_name="Real Crew", inputs=None))
assert app._crew_name == "Real Crew"
finally:
app._unsubscribe()
def test_method_paused_marks_step_paused() -> None:
app = CrewRunApp(crew_name="Demo")
app._flow = SimpleNamespace()
app._subscribe()
try:
_emit_event(
MethodExecutionStartedEvent(flow_name="Demo", method_name="ask", state={})
)
_emit_event(
MethodExecutionPausedEvent(
flow_name="Demo",
method_name="ask",
state={},
flow_id="flow-1",
message="Need your input",
)
)
assert app._flow_steps == [
{"name": "ask", "call_type": None, "status": "paused"}
]
assert app._current_method == "ask"
finally:
app._unsubscribe()
@pytest.mark.asyncio
async def test_declarative_flow_runs_on_tui() -> None:
"""End-to-end: on_mount dispatches _run_flow_worker → flow.kickoff →
_on_crew_done, and any still-active step is swept to done on completion."""
kicked: dict[str, object] = {}
class FakeFlow:
name = "Demo Flow"
def kickoff(self, inputs=None):
kicked["inputs"] = inputs
return "flow result"
app = CrewRunApp(crew_name="Demo Flow")
app._flow = FakeFlow()
app._flow_inputs = {"topic": "AI"}
# A step left active (no Finished event) must be swept to done by _on_crew_done.
app._flow_steps = [{"name": "compute", "call_type": "expression", "status": "active"}]
async with app.run_test() as pilot:
for _ in range(100):
await pilot.pause(0.05)
if app._status == "completed":
break
assert kicked["inputs"] == {"topic": "AI"}
assert app._status == "completed"
assert app._final_output == "flow result"
assert app._crew_result == "flow result"
assert app._flow_steps[0]["status"] == "done"

View File

@@ -2,6 +2,7 @@ from __future__ import annotations
import os
from pathlib import Path
from types import SimpleNamespace
import pytest
@@ -9,6 +10,17 @@ import crewai_cli.input_prompt as input_prompt_module
import crewai_cli.run_declarative_flow as run_declarative_flow_module
@pytest.fixture(autouse=True)
def _headless_by_default(monkeypatch: pytest.MonkeyPatch) -> None:
"""Default these tests to the headless/terminal path.
``run_declarative_flow`` now launches the TUI when interactive, which can't
run under pytest; tests here assert the terminal/headless contract. Tests
that exercise TUI routing override ``is_dmn_mode_enabled`` explicitly.
"""
monkeypatch.setenv("CREWAI_DMN", "true")
FLOW_YAML = """\
schema: crewai.flow/v1
name: TestFlow
@@ -400,3 +412,202 @@ def test_id_restore_still_drops_unknown_keys(
assert resolved == {"id": "run-123"} # id kept, typo dropped
assert "Ignoring unknown input 'prospect_emai'" in captured.err
assert "Ignoring unknown input 'id'" not in captured.err
# ── TUI vs terminal (headless/deploy) routing ──────────────────────
def _install_fake_flow_app(monkeypatch, *, status, want_deploy=False):
"""Replace CrewRunApp/EventListener/summary so _run_declarative_flow_tui is
driven by a controllable fake app."""
class FakeEventListener:
pass
class FakeApp:
def __init__(self, crew_name=""):
self._crew_name = crew_name
self._status = status
self._want_deploy = want_deploy
self._crew_result = "result"
def run(self):
pass
monkeypatch.setattr(
"crewai.events.event_listener.EventListener", FakeEventListener
)
monkeypatch.setattr("crewai_cli.crew_run_tui.CrewRunApp", FakeApp)
monkeypatch.setattr(
run_declarative_flow_module, "_print_flow_post_tui_summary", lambda app: None
)
def test_run_declarative_flow_dmn_uses_terminal(
tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setenv("CREWAI_DMN", "true")
monkeypatch.setattr(
run_declarative_flow_module,
"_run_declarative_flow_tui",
lambda *a, **k: pytest.fail("DMN/headless mode must not launch the TUI"),
)
path = _write(tmp_path, REQUIRED_FLOW_YAML)
run_declarative_flow_module.run_declarative_flow(
str(path), '{"prospect_email":"a@b.com"}'
)
assert capsys.readouterr().out == "a@b.com\n"
def test_run_declarative_flow_interactive_uses_tui(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr(run_declarative_flow_module, "is_interactive", lambda: True)
captured: dict[str, object] = {}
monkeypatch.setattr(
run_declarative_flow_module,
"_run_declarative_flow_tui",
lambda flow, resolved: captured.update(flow=flow, inputs=resolved),
)
path = _write(tmp_path, REQUIRED_FLOW_YAML)
run_declarative_flow_module.run_declarative_flow(
str(path), '{"prospect_email":"a@b.com"}'
)
assert captured["inputs"] == {"prospect_email": "a@b.com"}
assert captured["flow"] is not None
def test_run_declarative_flow_tui_failed_exits_nonzero(
monkeypatch: pytest.MonkeyPatch,
) -> None:
_install_fake_flow_app(monkeypatch, status="failed")
with pytest.raises(SystemExit) as exc_info:
run_declarative_flow_module._run_declarative_flow_tui(
SimpleNamespace(name="Flow"), None
)
assert exc_info.value.code == 1
def test_run_declarative_flow_tui_user_quit_exits_130(
monkeypatch: pytest.MonkeyPatch,
) -> None:
_install_fake_flow_app(monkeypatch, status="chatting")
exit_calls: list[int] = []
monkeypatch.setattr(os, "_exit", lambda code: exit_calls.append(code))
run_declarative_flow_module._run_declarative_flow_tui(
SimpleNamespace(name="Flow"), None
)
assert exit_calls == [130]
def test_run_declarative_flow_tui_chains_deploy(
monkeypatch: pytest.MonkeyPatch,
) -> None:
_install_fake_flow_app(monkeypatch, status="completed", want_deploy=True)
deploy_calls: list[bool] = []
monkeypatch.setattr(
"crewai_cli.run_crew._chain_deploy", lambda: deploy_calls.append(True)
)
run_declarative_flow_module._run_declarative_flow_tui(
SimpleNamespace(name="Flow"), None
)
assert deploy_calls == [True]
def test_run_declarative_flow_tui_no_deploy_when_not_requested(
monkeypatch: pytest.MonkeyPatch,
) -> None:
_install_fake_flow_app(monkeypatch, status="completed", want_deploy=False)
deploy_calls: list[bool] = []
monkeypatch.setattr(
"crewai_cli.run_crew._chain_deploy", lambda: deploy_calls.append(True)
)
run_declarative_flow_module._run_declarative_flow_tui(
SimpleNamespace(name="Flow"), None
)
assert deploy_calls == []
def test_run_declarative_flow_tui_enables_flow_events(
monkeypatch: pytest.MonkeyPatch,
) -> None:
# The STEPS panel depends on flow method events; a flow that declared
# suppress_flow_events must have it forced off for the interactive TUI run.
_install_fake_flow_app(monkeypatch, status="completed")
flow = SimpleNamespace(name="Flow", suppress_flow_events=True)
run_declarative_flow_module._run_declarative_flow_tui(flow, None)
assert flow.suppress_flow_events is False
def test_flow_uses_human_feedback_detection() -> None:
hf_flow = SimpleNamespace(
_definition=SimpleNamespace(
methods={
"ask": SimpleNamespace(human_feedback=SimpleNamespace(emit=None)),
"plain": SimpleNamespace(human_feedback=None),
}
)
)
assert run_declarative_flow_module._flow_uses_human_feedback(hf_flow) is True
no_hf = SimpleNamespace(
_definition=SimpleNamespace(
methods={"a": SimpleNamespace(human_feedback=None)}
)
)
assert run_declarative_flow_module._flow_uses_human_feedback(no_hf) is False
# No definition → False, no error.
assert run_declarative_flow_module._flow_uses_human_feedback(SimpleNamespace()) is False
def test_human_feedback_flow_uses_terminal_even_when_interactive(
tmp_path: Path, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
) -> None:
# A human-feedback flow must run on the terminal (blocking input / Rich
# prompts) even in an interactive session, never on the TUI.
monkeypatch.setattr(run_declarative_flow_module, "is_interactive", lambda: True)
monkeypatch.setattr(
run_declarative_flow_module, "_flow_uses_human_feedback", lambda flow: True
)
monkeypatch.setattr(
run_declarative_flow_module,
"_run_declarative_flow_tui",
lambda *a, **k: pytest.fail("human-feedback flow must run on the terminal"),
)
path = _write(tmp_path, FLOW_YAML)
run_declarative_flow_module.run_declarative_flow(str(path), '{"topic":"AI"}')
assert capsys.readouterr().out == "AI\n"
def test_flow_method_types_from_definition() -> None:
flow = SimpleNamespace(
_definition=SimpleNamespace(
methods={
"fetch": SimpleNamespace(do=SimpleNamespace(call="expression")),
"research": SimpleNamespace(do=SimpleNamespace(call="crew")),
}
)
)
assert run_declarative_flow_module._flow_method_types(flow) == {
"fetch": "expression",
"research": "crew",
}
# No definition → empty map, no error.
assert run_declarative_flow_module._flow_method_types(SimpleNamespace()) == {}

View File

@@ -86,6 +86,7 @@ from crewai.skills.models import Skill as SkillModel
from crewai.state.checkpoint_config import CheckpointConfig, apply_checkpoint
from crewai.tools.agent_tools.agent_tools import AgentTools
from crewai.types.callback import SerializableCallable
from crewai.types.usage_metrics import UsageMetrics
from crewai.utilities.agent_utils import (
get_tool_names,
is_inside_event_loop,
@@ -400,10 +401,29 @@ class Agent(BaseAgent):
return self.planning_config is not None or self.planning
def _setup_agent_executor(self) -> None:
"""Initialize the agent executor with a default cache handler."""
if not self.cache_handler:
self.cache_handler = CacheHandler()
self.set_cache_handler(self.cache_handler)
"""Initialize the agent's tools handler and optional tool cache.
Tool-result caching is opt-in: a standalone agent gets a cache only
when it was constructed with an explicit ``cache=True`` or a
``cache_handler``. Agents inside a crew additionally receive the
crew's shared handler when ``Crew(cache=True)``. Without an opt-in,
repeated tool calls with identical arguments always re-execute the
tool — the safe default for live-data and state-mutating tools.
"""
# Recorded before any crew can offer its shared handler at kickoff,
# so copy() can distinguish a construction-time opt-in from runtime
# crew wiring (which must not turn copies into cachers).
self._constructor_cache_opt_in = bool(
self.cache
and (self.cache_handler is not None or "cache" in self.model_fields_set)
)
opted_in = self.cache_handler is not None or (
"cache" in self.model_fields_set and self.cache
)
if opted_in:
if not self.cache_handler:
self.cache_handler = CacheHandler()
self.set_cache_handler(self.cache_handler)
def set_knowledge(self, crew_embedder: EmbedderConfig | None = None) -> None:
"""Initialize knowledge sources with the agent or crew embedder config."""
@@ -1582,9 +1602,18 @@ class Agent(BaseAgent):
crewai_event_bus.emit(self, event=started_event)
self._kickoff_event_id = started_event.event_id
output = self._execute_and_build_output(executor, inputs, response_format)
usage_baseline = self._current_usage_summary()
output = self._execute_and_build_output(
executor, inputs, response_format, usage_baseline
)
return self._finalize_kickoff(
output, executor, inputs, response_format, messages, agent_info
output,
executor,
inputs,
response_format,
messages,
agent_info,
usage_baseline,
)
except Exception as e:
@@ -1598,6 +1627,7 @@ class Agent(BaseAgent):
response_format: type[Any] | None,
messages: str | list[LLMMessage],
agent_info: dict[str, Any],
usage_baseline: UsageMetrics | None = None,
) -> LiteAgentOutput:
"""Apply guardrails, save to memory, and emit completion event.
@@ -1608,6 +1638,8 @@ class Agent(BaseAgent):
response_format: Optional response format.
messages: The original messages.
agent_info: Agent metadata for events.
usage_baseline: Usage snapshot taken at kickoff start, so retries
report per-call usage relative to it.
Returns:
The finalized output.
@@ -1618,6 +1650,7 @@ class Agent(BaseAgent):
executor=executor,
inputs=inputs,
response_format=response_format,
usage_baseline=usage_baseline,
)
self._save_kickoff_to_memory(messages, output.raw)
@@ -1669,11 +1702,24 @@ class Agent(BaseAgent):
except Exception as e:
self._logger.log("error", f"Failed to save kickoff result to memory: {e}")
def _current_usage_summary(self) -> UsageMetrics:
"""Snapshot the cumulative usage counters backing this agent's LLM.
The counters live on the LLM instance (or the agent's token process
for non-BaseLLM models) and grow for the object's lifetime — across
calls and across agents sharing the instance. Per-call usage is the
delta between two snapshots.
"""
if isinstance(self.llm, BaseLLM):
return self.llm.get_token_usage_summary()
return self._token_process.get_summary()
def _build_output_from_result(
self,
result: dict[str, Any],
executor: AgentExecutor,
response_format: type[Any] | None = None,
usage_baseline: UsageMetrics | None = None,
) -> LiteAgentOutput:
"""Build a LiteAgentOutput from an executor result dict.
@@ -1683,6 +1729,9 @@ class Agent(BaseAgent):
result: The result dictionary from executor.invoke / invoke_async.
executor: The executor instance.
response_format: Optional response format.
usage_baseline: Usage snapshot taken at kickoff start. When given,
the output carries only this call's usage (the delta) instead
of the LLM instance's cumulative lifetime counters.
Returns:
LiteAgentOutput with raw output, formatted result, and metrics.
@@ -1727,10 +1776,9 @@ class Agent(BaseAgent):
else:
raw_output = str(output) if not isinstance(output, str) else output
if isinstance(self.llm, BaseLLM):
usage_metrics = self.llm.get_token_usage_summary()
else:
usage_metrics = self._token_process.get_summary()
usage_metrics = self._current_usage_summary()
if usage_baseline is not None:
usage_metrics = usage_metrics.delta_since(usage_baseline)
raw_str = (
raw_output
@@ -1759,20 +1807,26 @@ class Agent(BaseAgent):
executor: AgentExecutor,
inputs: dict[str, str],
response_format: type[Any] | None = None,
usage_baseline: UsageMetrics | None = None,
) -> LiteAgentOutput:
"""Execute the agent synchronously and build the output object."""
result = cast(dict[str, Any], executor.invoke(inputs))
return self._build_output_from_result(result, executor, response_format)
return self._build_output_from_result(
result, executor, response_format, usage_baseline
)
async def _execute_and_build_output_async(
self,
executor: AgentExecutor,
inputs: dict[str, str],
response_format: type[Any] | None = None,
usage_baseline: UsageMetrics | None = None,
) -> LiteAgentOutput:
"""Execute the agent asynchronously and build the output object."""
result = await executor.invoke_async(inputs)
return self._build_output_from_result(result, executor, response_format)
return self._build_output_from_result(
result, executor, response_format, usage_baseline
)
def _process_kickoff_guardrail(
self,
@@ -1781,6 +1835,7 @@ class Agent(BaseAgent):
inputs: dict[str, str],
response_format: type[Any] | None = None,
retry_count: int = 0,
usage_baseline: UsageMetrics | None = None,
) -> LiteAgentOutput:
"""Process guardrail for kickoff execution with retry logic.
@@ -1790,6 +1845,9 @@ class Agent(BaseAgent):
inputs: Input dictionary for re-execution.
response_format: Optional response format.
retry_count: Current retry count.
usage_baseline: Usage snapshot taken at kickoff start, so a
retried output reports the whole call's usage, not just the
last attempt's.
Returns:
Validated/updated output.
@@ -1827,7 +1885,9 @@ class Agent(BaseAgent):
role="user",
)
output = self._execute_and_build_output(executor, inputs, response_format)
output = self._execute_and_build_output(
executor, inputs, response_format, usage_baseline
)
return self._process_kickoff_guardrail(
output=output,
@@ -1835,6 +1895,7 @@ class Agent(BaseAgent):
inputs=inputs,
response_format=response_format,
retry_count=retry_count + 1,
usage_baseline=usage_baseline,
)
if guardrail_result.result is not None:
@@ -1897,11 +1958,18 @@ class Agent(BaseAgent):
crewai_event_bus.emit(self, event=started_event)
self._kickoff_event_id = started_event.event_id
usage_baseline = self._current_usage_summary()
output = await self._execute_and_build_output_async(
executor, inputs, response_format
executor, inputs, response_format, usage_baseline
)
return self._finalize_kickoff(
output, executor, inputs, response_format, messages, agent_info
output,
executor,
inputs,
response_format,
messages,
agent_info,
usage_baseline,
)
except Exception as e:

View File

@@ -205,7 +205,11 @@ class BaseAgent(BaseModel, ABC, metaclass=AgentMeta):
role (str): Role of the agent.
goal (str): Objective of the agent.
backstory (str): Backstory of the agent.
cache (bool): Whether the agent should use a cache for tool usage.
cache (bool): Whether the agent participates in tool-result caching
when a cache is enabled. The default (True) only permits
participation — caching activates when the crew sets cache=True
or the agent explicitly opts in with cache=True or a
cache_handler; cache=False excludes the agent entirely.
config (dict[str, Any] | None): Configuration for the agent.
verbose (bool): Verbose mode for the Agent Execution.
max_rpm (int | None): Maximum number of requests per minute for the agent execution.
@@ -254,6 +258,7 @@ class BaseAgent(BaseModel, ABC, metaclass=AgentMeta):
_logger: Logger = PrivateAttr(default_factory=lambda: Logger(verbose=False))
_rpm_controller: RPMController | None = PrivateAttr(default=None)
_request_within_rpm_limit: SerializableCallable | None = PrivateAttr(default=None)
_constructor_cache_opt_in: bool = PrivateAttr(default=False)
_original_role: str | None = PrivateAttr(default=None)
_original_goal: str | None = PrivateAttr(default=None)
_original_backstory: str | None = PrivateAttr(default=None)
@@ -267,7 +272,14 @@ class BaseAgent(BaseModel, ABC, metaclass=AgentMeta):
description="Configuration for the agent", default=None, exclude=True
)
cache: bool = Field(
default=True, description="Whether the agent should use a cache for tool usage."
default=True,
description=(
"Whether the agent participates in tool-result caching when a "
"cache is enabled. Caching itself is opt-in: it activates only "
"when the crew sets cache=True or the agent explicitly opts in "
"(cache=True or a cache_handler at construction). Set False to "
"exclude this agent even when the crew enables caching."
),
)
verbose: bool = Field(
default=False, description="Verbose mode for the Agent Execution"
@@ -716,6 +728,19 @@ class BaseAgent(BaseModel, ABC, metaclass=AgentMeta):
copied_data = self.model_dump(exclude=exclude)
copied_data = {k: v for k, v in copied_data.items() if v is not None}
# Tool-result caching distinguishes "explicitly enabled" from the
# field default via model_fields_set; don't let the dump turn the
# default into an explicit opt-in on the copy. An agent that opted
# in at construction via an explicit cache_handler (excluded from
# the dump) must stay opted in — carry the consent as cache=True so
# the copy wires its own fresh handler. A handler merely offered by
# a crew at kickoff is runtime wiring, not consent, and must not
# opt the copy in; _constructor_cache_opt_in is recorded before any
# crew wiring can happen.
if "cache" not in self.model_fields_set:
copied_data.pop("cache", None)
if self._constructor_cache_opt_in:
copied_data["cache"] = True
return type(self)(
**copied_data,
llm=existing_llm,

View File

@@ -168,8 +168,11 @@ class Crew(FlowTrackable, BaseModel):
manager_agent: Custom agent that will be used as manager.
memory: Whether the crew should use memory to store memories of it's
execution.
cache: Whether the crew should use a cache to store the results of the
tools execution.
cache: Whether to cache tool results for the crew's agents. Off by
default; when enabled, repeated calls to the same tool with
identical arguments reuse the first result without re-executing —
avoid enabling for live-data or state-mutating tools unless they
gate writes with a cache_function.
function_calling_llm: The language model that will run the tool calling
for all the agents.
process: The process flow that the crew will follow (e.g., sequential,
@@ -216,7 +219,16 @@ class Crew(FlowTrackable, BaseModel):
_kickoff_event_id: str | None = PrivateAttr(default=None)
name: str | None = Field(default="crew")
cache: bool = Field(default=True)
cache: bool = Field(
default=False,
description=(
"Whether to cache tool results for the crew's agents. Opt-in: "
"when enabled, repeated calls to the same tool with identical "
"arguments return the first result without re-executing the "
"tool — do not enable for live-data or state-mutating tools "
"unless they set a cache_function that prevents caching."
),
)
tasks: list[Task] = Field(default_factory=list)
agents: Annotated[
list[BaseAgent],
@@ -1507,6 +1519,11 @@ class Crew(FlowTrackable, BaseModel):
)
self.manager_agent = manager
manager.crew = self
# The manager is created outside the agents loop that offers the
# crew's cache handler at validation time; offer it here so an
# opted-in crew (cache=True) also dedupes the manager's tool calls.
if self.cache:
manager.set_cache_handler(self._cache_handler)
def _get_execution_start_index(self, tasks: list[Task]) -> int | None:
if self.checkpoint_kickoff_event_id is None:

View File

@@ -24,9 +24,23 @@ class CrewOutput(BaseModel):
description="Output of each task", default_factory=list
)
token_usage: UsageMetrics = Field(
description="Processed token summary", default_factory=UsageMetrics
description=(
"Processed token summary; ``usage_metrics`` exposes the same "
"data as a plain dict"
),
default_factory=UsageMetrics,
)
@property
def usage_metrics(self) -> dict[str, Any]:
"""Token usage as a plain dict.
Same attribute name and shape as ``LiteAgentOutput.usage_metrics``
(the ``Agent.kickoff()`` result), so a usage accessor written for one
result type works on both.
"""
return self.token_usage.model_dump()
@property
def json(self) -> str | None: # type: ignore[override]
if self.tasks_output[-1].output_format != OutputFormat.JSON:

View File

@@ -211,6 +211,13 @@ To enable tracing, do any one of these:
"""Print a panel with consistent formatting if verbose is enabled."""
panel = self.create_panel(content, title, style)
if is_flow:
# A TUI (e.g. the CLI's CrewRunApp) owns the screen and renders flow
# progress in its own STEPS panel; emitting Rich panels here would
# interleave with and corrupt the TUI, so suppress them in TUI mode.
from crewai.events.listeners.tracing.utils import is_tui_mode
if is_tui_mode():
return
self.print(panel)
self.print()
else:

View File

@@ -133,6 +133,7 @@ class _ConversationalMixin:
_pending_user_message: str | dict[str, Any] | None
_pending_intents: Sequence[str] | None
_pending_intent_llm: str | BaseLLM | None
_turn_classified_intent: str | None
_assistant_reply_appended: bool
def _clear_or_listeners(self) -> None:
@@ -186,12 +187,22 @@ class _ConversationalMixin:
)
return configured_route
if state.last_intent:
turn_intent = self._turn_classified_intent
if turn_intent:
state.last_intent = turn_intent
self._emit_conversation_route_selected(
state.last_intent,
turn_intent,
previous_intent=previous_intent,
)
return state.last_intent
return turn_intent
if previous_intent:
logger.debug(
"route_turn() returned no route and no intent was classified "
"this turn; ignoring stale last_intent=%r from a previous turn "
"and falling back to built-in routing",
previous_intent,
)
if self.can_answer_from_history(context):
state.last_intent = "answer_from_history"
@@ -551,6 +562,11 @@ class _ConversationalMixin:
supply per-route descriptions, or change the default/fallback intent.
Override this method to bypass the LLM router entirely (e.g.,
permission gates before the LLM decision).
Returning a falsy value means "no routing decision": the turn falls
through to the built-in defaults (``answer_from_history`` when
configured, else ``converse``). It never replays a previous turn's
intent.
"""
config = self._conversation_config
if config is None:
@@ -726,6 +742,7 @@ class _ConversationalMixin:
context=self.conversation_messages,
)
state.last_intent = intent
object.__setattr__(self, "_turn_classified_intent", intent)
return intent
return text
@@ -792,6 +809,8 @@ class _ConversationalMixin:
object.__setattr__(self, "_pending_intent_llm", None)
if not hasattr(self, "_streaming_conversation_turn"):
object.__setattr__(self, "_streaming_conversation_turn", False)
if not hasattr(self, "_turn_classified_intent"):
object.__setattr__(self, "_turn_classified_intent", None)
if not hasattr(self, "_assistant_reply_appended"):
object.__setattr__(self, "_assistant_reply_appended", False)
@@ -858,6 +877,7 @@ class _ConversationalMixin:
self._method_call_counts.clear()
self._clear_or_listeners()
self._is_execution_resuming = False
object.__setattr__(self, "_turn_classified_intent", None)
def _apply_pending_conversational_turn(self) -> None:
"""Drain the stashed user message + classify if intents configured.
@@ -865,6 +885,7 @@ class _ConversationalMixin:
Called from ``Flow.kickoff_async`` AFTER persist state restore so
the appended message survives ``self.persistence.load_state(...)``.
"""
object.__setattr__(self, "_turn_classified_intent", None)
if self._pending_user_message is None:
return

View File

@@ -6,6 +6,7 @@ from typing import Any
from pydantic import BaseModel, Field
from crewai.types.usage_metrics import UsageMetrics
from crewai.utilities.planning_types import TodoItem
from crewai.utilities.types import LLMMessage
@@ -38,7 +39,13 @@ class LiteAgentOutput(BaseModel):
)
agent_role: str = Field(description="Role of the agent that produced this output")
usage_metrics: dict[str, Any] | None = Field(
description="Token usage metrics for this execution", default=None
description=(
"Token usage for this kickoff call only (guardrail retries "
"included), not the LLM instance's cumulative totals, as a "
"plain dict; ``token_usage`` exposes the same data as a "
"UsageMetrics object"
),
default=None,
)
messages: list[LLMMessage] = Field(
description="Messages of the agent", default_factory=list
@@ -86,6 +93,19 @@ class LiteAgentOutput(BaseModel):
return self.pydantic.model_dump()
return {}
@property
def token_usage(self) -> UsageMetrics:
"""Token usage as a ``UsageMetrics`` object.
Same attribute name and type as ``CrewOutput.token_usage``, so a
usage accessor written for one result type works on both. Returns
zeroed metrics when no usage was captured (``usage_metrics`` is
``None``).
"""
if not self.usage_metrics:
return UsageMetrics()
return UsageMetrics.model_validate(self.usage_metrics)
@property
def completed_todos(self) -> list[TodoExecutionResult]:
"""Get only the completed todos."""

View File

@@ -966,8 +966,15 @@ class BaseLLM(BaseModel, ABC):
def get_token_usage_summary(self) -> UsageMetrics:
"""Get summary of token usage for this LLM instance.
The counters are cumulative for the lifetime of this instance: they
grow across every call made through it, including calls issued by
different agents sharing the instance. For usage scoped to a single
call, snapshot before and after and use
``UsageMetrics.delta_since`` (agent kickoff results already report
per-call usage this way).
Returns:
Dictionary with token usage totals
UsageMetrics with this instance's lifetime token usage totals.
"""
return UsageMetrics(**self._token_usage)

View File

@@ -5,7 +5,6 @@ import asyncio
from collections.abc import Awaitable, Callable
import importlib
from inspect import Parameter, signature
import json
import threading
from typing import (
Any,
@@ -37,9 +36,9 @@ from crewai.tools.structured_tool import (
_infer_result_schema_from_callable,
_serialize_schema,
build_schema_hint,
format_description_for_llm,
)
from crewai.types.callback import SerializableCallable, _resolve_dotted_path
from crewai.utilities.pydantic_schema_utils import generate_model_description
from crewai.utilities.string_utils import sanitize_tool_name
@@ -479,15 +478,27 @@ class BaseTool(BaseModel, ABC):
f"{self.__class__.__name__}Schema", **fields
)
@property
def formatted_description(self) -> str:
"""LLM-facing composite of name, argument schema, and description.
Use this when rendering the tool into a prompt; ``description``
holds only the authored text.
"""
return format_description_for_llm(self.name, self.args_schema, self.description)
def _generate_description(self) -> None:
"""Generate the tool description with a JSON schema for arguments."""
schema = generate_model_description(self.args_schema)
args_json = json.dumps(schema["json_schema"]["schema"], indent=2)
self.description = (
f"Tool Name: {sanitize_tool_name(self.name)}\n"
f"Tool Arguments: {args_json}\n"
f"Tool Description: {self.description}"
)
"""Deprecated hook kept for backward compatibility; does nothing.
Historically this rewrote the public ``description`` field at
construction time into the LLM-facing composite (``Tool Name: …\\n
Tool Arguments: …\\nTool Description: <authored>``). The authored
``description`` is now preserved as written and the composite is
exposed separately via :attr:`formatted_description`.
``model_post_init`` still calls this so subclasses that override it
(e.g. adapters that customize the composite) keep working.
"""
_BASE_TOOL_CLS = BaseTool

View File

@@ -4,6 +4,7 @@ import asyncio
from collections.abc import Callable
import inspect
import json
import re
import textwrap
from typing import TYPE_CHECKING, Annotated, Any, cast, get_type_hints
import warnings
@@ -21,7 +22,10 @@ from pydantic import (
from typing_extensions import Self
from crewai.utilities.logger import Logger
from crewai.utilities.pydantic_schema_utils import create_model_from_schema
from crewai.utilities.pydantic_schema_utils import (
create_model_from_schema,
generate_model_description,
)
from crewai.utilities.string_utils import sanitize_tool_name
@@ -108,6 +112,70 @@ def build_schema_hint(args_schema: type[BaseModel]) -> str:
return ""
# Matches a description that IS a pre-composed LLM block (as written by
# older versions into the field, and by adapters that still bake it in).
# Anchored to the full three-line shape so authored prose that merely
# mentions "Tool Description:" is never mistaken for a composite. Greedy
# ``.*`` keeps only the text after the LAST marker, matching the historical
# split behavior for nested pre-baked blocks.
_COMPOSITE_DESCRIPTION_RE = re.compile(
r"^Tool Name:.*\nTool Arguments:.*\nTool Description:\s*",
re.DOTALL,
)
def strip_composite_description_prefix(description: str) -> str:
"""Return the authored text from a pre-composed LLM description block.
Descriptions that don't start with the composite shape are returned
unchanged.
"""
match = _COMPOSITE_DESCRIPTION_RE.match(description)
if match:
return description[match.end() :]
return description
def format_description_for_llm(
name: str,
args_schema: type[BaseModel] | None,
description: str,
) -> str:
"""Compose the LLM-facing tool description.
Combines the tool name, its argument JSON schema, and the authored
description into the prompt block agents see. The authored
``description`` field itself is never mutated — prompt rendering calls
this on demand.
Idempotent: if ``description`` already *is* a composed block (e.g. a
tool deserialized from a checkpoint written by an older version, or an
adapter that bakes the composite into the field), only the authored
text after the marker is used. The check is anchored to the composite
shape, so authored prose that merely mentions ``"Tool Description:"``
passes through untouched.
Args:
name: The tool name (sanitized for the prompt).
args_schema: The tool's argument schema, if any.
description: The authored tool description.
Returns:
The composed, LLM-facing description block.
"""
description = strip_composite_description_prefix(description)
if args_schema is not None:
schema = generate_model_description(args_schema)
args_json = json.dumps(schema["json_schema"]["schema"], indent=2)
else:
args_json = "{}"
return (
f"Tool Name: {sanitize_tool_name(name)}\n"
f"Tool Arguments: {args_json}\n"
f"Tool Description: {description}"
)
class ToolUsageLimitExceededError(Exception):
"""Exception raised when a tool has reached its maximum usage limit."""
@@ -141,6 +209,15 @@ class CrewStructuredTool(BaseModel):
_logger: Logger = PrivateAttr(default_factory=Logger)
_original_tool: Any = PrivateAttr(default=None)
@property
def formatted_description(self) -> str:
"""LLM-facing composite of name, argument schema, and description.
Use this when rendering the tool into a prompt; ``description``
holds only the authored text.
"""
return format_description_for_llm(self.name, self.args_schema, self.description)
@model_validator(mode="after")
def _validate_func(self) -> Self:
if self.func is not None:

View File

@@ -430,7 +430,7 @@ class ToolUsage:
).format(
error=e,
tool=sanitize_tool_name(tool.name),
tool_inputs=tool.description,
tool_inputs=tool.formatted_description,
)
result = ToolUsageError(
f"\n{error_message}.\nMoving on then. {I18N_DEFAULT.slice('format').format(tool_names=self.tools_names)}"
@@ -670,7 +670,7 @@ class ToolUsage:
).format(
error=e,
tool=sanitize_tool_name(tool.name),
tool_inputs=tool.description,
tool_inputs=tool.formatted_description,
)
result = ToolUsageError(
f"\n{error_message}.\nMoving on then. {I18N_DEFAULT.slice('format').format(tool_names=self.tools_names)}"
@@ -803,7 +803,7 @@ class ToolUsage:
def _render(self) -> str:
"""Render the tool name and description in plain text."""
descriptions = [tool.description for tool in self.tools]
descriptions = [tool.formatted_description for tool in self.tools]
return "\n--\n".join(descriptions)
def _function_calling(

View File

@@ -76,6 +76,38 @@ class UsageMetrics(BaseModel):
self.cache_creation_tokens += usage_metrics.cache_creation_tokens
self.successful_requests += usage_metrics.successful_requests
def delta_since(self, baseline: Self) -> Self:
"""Return the per-call usage accrued since ``baseline`` was captured.
Both objects must come from the same monotonically increasing
accumulator (e.g. an LLM instance's lifetime counters). Differences
are clamped at zero so a reset accumulator can't produce negative
usage.
Args:
baseline: A snapshot of the same accumulator taken earlier.
Returns:
A new UsageMetrics with the field-wise difference.
"""
return type(self)(
total_tokens=max(0, self.total_tokens - baseline.total_tokens),
prompt_tokens=max(0, self.prompt_tokens - baseline.prompt_tokens),
cached_prompt_tokens=max(
0, self.cached_prompt_tokens - baseline.cached_prompt_tokens
),
completion_tokens=max(
0, self.completion_tokens - baseline.completion_tokens
),
reasoning_tokens=max(0, self.reasoning_tokens - baseline.reasoning_tokens),
cache_creation_tokens=max(
0, self.cache_creation_tokens - baseline.cache_creation_tokens
),
successful_requests=max(
0, self.successful_requests - baseline.successful_requests
),
)
@classmethod
def from_provider_dict(cls, usage_data: dict[str, Any] | None) -> Self | None:
"""Normalize a provider's raw usage dict into a ``UsageMetrics``.

View File

@@ -27,7 +27,10 @@ from crewai.agents.parser import (
from crewai.llms.base_llm import BaseLLM, call_stop_override
from crewai.tools import BaseTool as CrewAITool
from crewai.tools.base_tool import BaseTool
from crewai.tools.structured_tool import CrewStructuredTool
from crewai.tools.structured_tool import (
CrewStructuredTool,
strip_composite_description_prefix,
)
from crewai.tools.tool_types import ToolResult
from crewai.utilities.errors import AgentRepositoryError
from crewai.utilities.exceptions.context_window_exceeding_exception import (
@@ -147,7 +150,14 @@ def render_text_description_and_args(
Returns:
Plain text description of tools.
"""
tool_strings = [tool.description for tool in tools]
# Fall back to the raw description for duck-typed tools (including test
# mocks) that don't provide a real formatted_description string.
tool_strings = [
formatted
if isinstance((formatted := getattr(tool, "formatted_description", None)), str)
else tool.description
for tool in tools
]
return "\n".join(tool_strings)
@@ -190,10 +200,10 @@ def convert_tools_to_openai_schema(
except Exception:
parameters = {}
# BaseTool formats description as "Tool Name: ...\nTool Arguments: ...\nTool Description: {original}"
description = tool.description
if "Tool Description:" in description:
description = description.split("Tool Description:")[-1].strip()
# Old checkpoints and some adapters bake the composed LLM block
# ("Tool Name: ...\nTool Arguments: ...\nTool Description: {authored}")
# into the description field; keep only the authored text here.
description = strip_composite_description_prefix(tool.description)
sanitized_name = sanitize_tool_name(tool.name)

View File

@@ -1138,3 +1138,160 @@ def test_lite_agent_memory_instance_recall_and_save_called():
mock_memory.remember_many.assert_called_once_with(
["Fact one.", "Fact two."], agent_role="Test"
)
class _FixedUsageLLM(BaseLLM):
"""Offline BaseLLM that records fixed usage (100/10 tokens) per call."""
def __init__(self):
super().__init__(model="fixed-usage-model")
def call(
self,
messages,
tools=None,
callbacks=None,
available_functions=None,
from_task=None,
from_agent=None,
response_model=None,
) -> str:
self._track_token_usage_internal(
{"prompt_tokens": 100, "completion_tokens": 10, "total_tokens": 110}
)
return "Thought: I know the answer.\nFinal Answer: fake answer"
def supports_function_calling(self) -> bool:
return False
def supports_stop_words(self) -> bool:
return False
def get_context_window_size(self) -> int:
return 4096
class TestKickoffUsageMetricsArePerCall:
"""Regression tests for EPD-177: kickoff results used to expose the LLM
instance's cumulative lifetime counters, so counts accumulated across
calls and pooled across agents sharing one LLM object.
"""
def _make_agent(self, role: str, llm: BaseLLM) -> Agent:
return Agent(
role=role,
goal="Answer questions.",
backstory="Test agent.",
llm=llm,
verbose=False,
)
def test_agents_sharing_one_llm_report_per_call_usage(self):
shared = _FixedUsageLLM()
r1 = self._make_agent("agent one", shared).kickoff("question one")
r2 = self._make_agent("agent two", shared).kickoff("question two")
assert r1.usage_metrics is not None
assert r1.usage_metrics["prompt_tokens"] > 0
# The second agent's call must not include the first agent's tokens.
assert r2.usage_metrics == r1.usage_metrics
# The shared LLM instance still exposes cumulative lifetime totals.
lifetime = shared.get_token_usage_summary()
assert lifetime.prompt_tokens == (
r1.usage_metrics["prompt_tokens"] + r2.usage_metrics["prompt_tokens"]
)
assert lifetime.successful_requests == (
r1.usage_metrics["successful_requests"]
+ r2.usage_metrics["successful_requests"]
)
def test_repeated_kickoffs_on_same_agent_report_per_call_usage(self):
agent = self._make_agent("agent", _FixedUsageLLM())
r1 = agent.kickoff("question one")
r2 = agent.kickoff("question two")
assert r1.usage_metrics is not None
assert r1.usage_metrics["prompt_tokens"] > 0
assert r2.usage_metrics == r1.usage_metrics
@pytest.mark.asyncio
async def test_async_kickoff_reports_per_call_usage(self):
shared = _FixedUsageLLM()
r1 = await self._make_agent("agent one", shared).kickoff_async("question one")
r2 = await self._make_agent("agent two", shared).kickoff_async("question two")
assert r1.usage_metrics is not None
assert r1.usage_metrics["prompt_tokens"] > 0
assert r2.usage_metrics == r1.usage_metrics
def test_guardrail_retry_usage_includes_all_attempts(self):
"""A guardrail retry re-invokes the LLM within the same kickoff, so
the result must report the whole call's usage — every attempt — not
just the last one."""
baseline = (
self._make_agent("baseline", _FixedUsageLLM())
.kickoff("question one")
.usage_metrics
)
attempts: list[str] = []
def flaky_guardrail(output):
attempts.append(output.raw)
if len(attempts) == 1:
return (False, "Please try again.")
return (True, output.raw)
agent = Agent(
role="agent",
goal="Answer questions.",
backstory="Test agent.",
llm=_FixedUsageLLM(),
guardrail=flaky_guardrail,
verbose=False,
)
result = agent.kickoff("question one")
assert len(attempts) == 2
assert result.usage_metrics["successful_requests"] == (
2 * baseline["successful_requests"]
)
assert result.usage_metrics["prompt_tokens"] == 2 * baseline["prompt_tokens"]
assert result.usage_metrics["total_tokens"] == 2 * baseline["total_tokens"]
class TestUsageMetricsDeltaSince:
def test_field_wise_difference(self):
baseline = UsageMetrics(
total_tokens=110,
prompt_tokens=100,
completion_tokens=10,
successful_requests=1,
)
current = UsageMetrics(
total_tokens=330,
prompt_tokens=300,
completion_tokens=30,
cached_prompt_tokens=5,
reasoning_tokens=7,
cache_creation_tokens=3,
successful_requests=3,
)
delta = current.delta_since(baseline)
assert delta == UsageMetrics(
total_tokens=220,
prompt_tokens=200,
completion_tokens=20,
cached_prompt_tokens=5,
reasoning_tokens=7,
cache_creation_tokens=3,
successful_requests=2,
)
def test_clamps_negative_differences_to_zero(self):
baseline = UsageMetrics(total_tokens=100, prompt_tokens=90, successful_requests=2)
delta = UsageMetrics().delta_since(baseline)
assert delta == UsageMetrics()

View File

@@ -859,6 +859,7 @@ def test_cache_hitting_between_agents(researcher, writer, ceo):
crew = Crew(
agents=[ceo, researcher],
tasks=tasks,
cache=True,
)
with patch.object(CacheHandler, "read") as read:
@@ -2246,7 +2247,9 @@ def test_tools_with_custom_caching():
agent=writer2,
)
crew = Crew(agents=[writer1, writer2], tasks=[task1, task2, task3, task4])
crew = Crew(
agents=[writer1, writer2], tasks=[task1, task2, task3, task4], cache=True
)
with patch.object(
CacheHandler, "add", wraps=crew._cache_handler.add

View File

@@ -1620,6 +1620,115 @@ class TestHandleTurnReplyFallback:
assert assistant_messages == ["computed reply"]
class TestFalsyRouteTurnFallback:
"""A falsy ``route_turn()`` must never replay a previous turn's intent.
Regression tests for EPD-176: an overridden ``route_turn()`` returning
``None`` on an unhandled input used to silently reuse the sticky
``state.last_intent`` from the *previous* turn, running the wrong handler
with no error or warning.
"""
def test_falsy_route_turn_does_not_replay_previous_turns_intent(self) -> None:
ran: list[str] = []
class Bot(ConversationalFlow):
def route_turn(self, context: dict[str, Any]) -> str | None:
message = context.get("current_user_message") or ""
if "hello" in message.lower():
return "GREETING"
return None # unhandled input -> falsy return
@listen("GREETING")
def greeting(self) -> str:
ran.append("GREETING")
reply = "Hi! I only do greetings."
self.append_assistant_message(reply)
return reply
@listen("WEATHER")
def weather(self) -> str:
ran.append("WEATHER")
reply = "It is sunny."
self.append_assistant_message(reply)
return reply
flow = Bot()
flow.handle_turn("hello there")
assert ran == ["GREETING"]
assert flow.state.last_intent == "GREETING"
flow.handle_turn("what is the meaning of life?")
assert ran == ["GREETING"], (
"an unhandled turn must not re-run the previous turn's handler"
)
# With no routing decision the turn falls through to the built-in
# 'converse' default instead of replaying the stale intent.
assert flow.state.last_intent == "converse"
assert flow.state.messages[-1].content != "Hi! I only do greetings."
def test_stale_intent_ignored_but_route_selected_event_still_emitted(
self,
) -> None:
class Bot(ConversationalFlow):
def route_turn(self, context: dict[str, Any]) -> str | None:
message = context.get("current_user_message") or ""
return "work" if "work" in message else None
@listen("work")
def do_work(self) -> str:
self.append_assistant_message("worked")
return "worked"
flow = Bot()
routes: list[ConversationRouteSelectedEvent] = []
with crewai_event_bus.scoped_handlers():
@crewai_event_bus.on(ConversationRouteSelectedEvent)
def capture(_: Any, event: ConversationRouteSelectedEvent) -> None:
routes.append(event)
flow.handle_turn("work please")
flow.handle_turn("something unrelated")
crewai_event_bus.flush()
assert [event.route for event in routes] == ["work", "converse"]
# The fallback decision still reports the prior intent for visibility.
assert routes[1].previous_intent == "work"
def test_fresh_intent_classified_this_turn_still_routes(self) -> None:
"""The legacy ``default_intents`` path classifies per turn and must
keep routing on the freshly classified intent — including when the
intent changes between turns."""
ran: list[str] = []
@ConversationConfig(
default_intents=["search", "weather"], intent_llm="gpt-4o-mini"
)
class LegacyFlow(ConversationalFlow):
@listen("search")
def handle_search(self) -> str:
ran.append("search")
self.append_assistant_message("searched")
return "searched"
@listen("weather")
def handle_weather(self) -> str:
ran.append("weather")
self.append_assistant_message("sunny")
return "sunny"
flow = LegacyFlow()
with patch.object(
flow, "_collapse_to_outcome", side_effect=["search", "weather"]
):
flow.handle_turn("look up crewai")
flow.handle_turn("how is the weather?")
assert ran == ["search", "weather"]
assert flow.state.last_intent == "weather"
class TestFlowTracingWhenSuppressed:
def test_flow_started_emitted_when_panel_events_suppressed(self) -> None:
class QuietFlow(Flow[ChatState]):

View File

@@ -0,0 +1,263 @@
# mypy: ignore-errors
"""Regression tests for EPD-180: tool-result caching used to be ON by default,
so an LLM calling the same tool with identical arguments twice in one run got
the first (possibly stale) result back without the tool executing — silently
wrong for live-data tools, and silently dropped actions for stateful tools.
Caching is now opt-in: ``Crew(cache=True)`` for crews, ``Agent(cache=True)``
(or an explicit ``cache_handler``) for standalone agents. The machinery —
including per-tool ``cache_function`` write gating — is unchanged once opted
in.
The end-to-end tests run fully offline: a fake OpenAI client scripts two
identical tool calls followed by a final answer, mirroring the EPD-180
clean-room repro.
"""
from openai.types.chat import ChatCompletion
from pydantic import BaseModel, Field
from crewai import LLM, Agent, Crew, Task
from crewai.agents.cache.cache_handler import CacheHandler
from crewai.tools import BaseTool
class LookupArgs(BaseModel):
city: str = Field(description="City to look up.")
def make_live_tool():
"""A tool returning a different value on every real execution."""
executions = []
class LiveLookupTool(BaseTool):
name: str = "live_lookup"
description: str = "Returns a live (time-varying) reading for a city."
args_schema: type[BaseModel] = LookupArgs
# cache_function deliberately NOT set — exercising the default.
def _run(self, city: str) -> str:
executions.append(city)
return f"reading #{len(executions)} for {city}"
return LiveLookupTool(), executions
def make_scripted_llm():
"""An offline LLM whose client scripts two identical tool calls."""
def tool_call_response(call_id: str):
return {
"index": 0,
"finish_reason": "tool_calls",
"message": {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": call_id,
"type": "function",
"function": {
"name": "live_lookup",
"arguments": '{"city": "paris"}',
},
}
],
},
}
scripted = [
tool_call_response("call_1"),
tool_call_response("call_2"), # identical name+args, new id
{
"index": 0,
"finish_reason": "stop",
"message": {"role": "assistant", "content": "Final answer: done."},
},
]
class FakeCompletions:
def __init__(self):
self.n = 0
def create(self, **params):
choice = scripted[min(self.n, len(scripted) - 1)]
self.n += 1
return ChatCompletion.model_validate(
{
"id": f"chatcmpl-fake-{self.n}",
"object": "chat.completion",
"created": 1,
"model": params.get("model", "gpt-4o"),
"choices": [choice],
"usage": {
"prompt_tokens": 10,
"completion_tokens": 5,
"total_tokens": 15,
},
}
)
class FakeClient:
def __init__(self):
self.chat = type("Chat", (), {"completions": FakeCompletions()})()
llm = LLM(model="openai/gpt-4o")
llm._client = FakeClient()
return llm
def run_crew(**crew_kwargs):
tool, executions = make_live_tool()
agent = Agent(
role="reader",
goal="Look things up.",
backstory="Test agent.",
llm=make_scripted_llm(),
tools=[tool],
verbose=False,
)
task = Task(
description="Look up paris twice and report.",
expected_output="A report.",
agent=agent,
)
crew = Crew(agents=[agent], tasks=[task], verbose=False, **crew_kwargs)
crew.kickoff()
return executions
class TestToolCachingIsOptIn:
def test_default_reexecutes_identical_tool_calls(self):
"""EPD-180: with no opt-in, both identical calls must really execute."""
executions = run_crew()
assert len(executions) == 2
def test_crew_cache_true_dedupes_identical_tool_calls(self):
"""Opting in via Crew(cache=True) restores the dedup behavior."""
executions = run_crew(cache=True)
assert len(executions) == 1
class TestAgentCacheWiring:
def _agent(self, **kwargs) -> Agent:
return Agent(
role="reader",
goal="Look things up.",
backstory="Test agent.",
**kwargs,
)
def test_standalone_agent_has_no_cache_by_default(self):
agent = self._agent()
assert agent.tools_handler.cache is None
assert agent.cache_handler is None
def test_standalone_agent_explicit_cache_true_opts_in(self):
agent = self._agent(cache=True)
assert agent.tools_handler.cache is not None
assert agent.cache_handler is not None
def test_standalone_agent_explicit_cache_handler_opts_in(self):
handler = CacheHandler()
agent = self._agent(cache_handler=handler)
assert agent.tools_handler.cache is handler
def test_explicit_cache_false_stays_off_even_with_handler(self):
agent = self._agent(cache=False, cache_handler=CacheHandler())
assert agent.tools_handler.cache is None
def test_agents_accept_a_crew_offered_handler_by_default(self):
"""``Crew(cache=True)`` offers its handler via set_cache_handler at
kickoff; agents that didn't explicitly opt out must accept it."""
agent = self._agent()
assert agent.tools_handler.cache is None
handler = CacheHandler()
agent.set_cache_handler(handler)
assert agent.tools_handler.cache is handler
def test_agents_that_opted_out_refuse_a_crew_offered_handler(self):
agent = self._agent(cache=False)
agent.set_cache_handler(CacheHandler())
assert agent.tools_handler.cache is None
def test_copy_of_default_agent_does_not_opt_in(self):
"""copy() rebuilds from model_dump(), which includes the field
default cache=True — that must not read as an explicit opt-in on
the copy (Bugbot review finding on the original PR)."""
copied = self._agent().copy()
assert copied.tools_handler.cache is None
assert copied.cache_handler is None
def test_copy_of_opted_in_agent_stays_opted_in(self):
copied = self._agent(cache=True).copy()
assert copied.tools_handler.cache is not None
def test_copy_of_handler_opted_in_agent_stays_opted_in(self):
"""An explicit cache_handler is an opt-in too; copy() excludes the
handler itself, but the consent must survive — the copy wires its
own fresh handler (Bugbot review finding on the original PR)."""
source = self._agent(cache_handler=CacheHandler())
copied = source.copy()
assert copied.tools_handler.cache is not None
assert copied.tools_handler.cache is not source.tools_handler.cache
def test_copy_of_explicit_cache_false_with_handler_stays_off(self):
copied = self._agent(cache=False, cache_handler=CacheHandler()).copy()
assert copied.tools_handler.cache is None
def test_copy_of_crew_wired_agent_does_not_opt_in(self):
"""A handler offered by a crew at kickoff (set_cache_handler) is
runtime wiring, not construction-time consent — copies of such
agents must not become standalone cachers (Bugbot review finding
on the original PR)."""
agent = self._agent()
agent.set_cache_handler(CacheHandler()) # what Crew(cache=True) does
assert agent.tools_handler.cache is not None
copied = agent.copy()
assert copied.tools_handler.cache is None
assert copied.cache_handler is None
class TestHierarchicalManagerCacheWiring:
"""The auto-created hierarchical manager is built outside the agents
loop that offers the crew's cache handler; an opted-in crew must wire
the manager too (Bugbot review finding on the original PR)."""
def _crew(self, **crew_kwargs) -> Crew:
from crewai.process import Process
agent = Agent(role="worker", goal="Do work.", backstory="Test agent.")
task = Task(description="Do the work.", expected_output="Done.")
return Crew(
agents=[agent],
tasks=[task],
process=Process.hierarchical,
manager_llm="gpt-4o",
**crew_kwargs,
)
def test_manager_gets_crew_handler_when_cache_enabled(self):
crew = self._crew(cache=True)
crew._create_manager_agent()
assert crew.manager_agent.tools_handler.cache is crew._cache_handler
def test_manager_has_no_cache_when_crew_did_not_opt_in(self):
crew = self._crew()
crew._create_manager_agent()
assert crew.manager_agent.tools_handler.cache is None
def test_user_provided_manager_with_cache_false_stays_excluded(self):
manager = Agent(
role="manager",
goal="Manage.",
backstory="Test manager.",
cache=False,
allow_delegation=True,
)
crew = self._crew(cache=True)
crew.manager_agent = manager
crew._create_manager_agent()
assert manager.tools_handler.cache is None

View File

@@ -0,0 +1,143 @@
# mypy: ignore-errors
"""Regression tests for EPD-178: token usage was exposed in different shapes
and attribute names per code path — ``Agent.kickoff()`` results carried a
plain dict at ``.usage_metrics`` (no ``token_usage`` attribute at all), while
``Crew.kickoff()`` results carried a ``UsageMetrics`` object at
``.token_usage`` (no ``usage_metrics`` attribute), so any single accessor
written for one path raised ``AttributeError`` on the other.
Both result types now expose both surfaces: ``.token_usage`` as a
``UsageMetrics`` object and ``.usage_metrics`` as a plain dict.
"""
from crewai import Agent, Crew, Task
from crewai.crews.crew_output import CrewOutput
from crewai.lite_agent_output import LiteAgentOutput
from crewai.llms.base_llm import BaseLLM
from crewai.types.usage_metrics import UsageMetrics
class _FixedUsageLLM(BaseLLM):
"""Offline BaseLLM that records fixed usage (100/10 tokens) per call."""
def __init__(self):
super().__init__(model="fixed-usage-model")
def call(
self,
messages,
tools=None,
callbacks=None,
available_functions=None,
from_task=None,
from_agent=None,
response_model=None,
) -> str:
self._track_token_usage_internal(
{"prompt_tokens": 100, "completion_tokens": 10, "total_tokens": 110}
)
return "Thought: I know the answer.\nFinal Answer: fake answer"
def supports_function_calling(self) -> bool:
return False
def supports_stop_words(self) -> bool:
return False
def get_context_window_size(self) -> int:
return 4096
class TestUsageShapeUnitParity:
def test_lite_agent_output_exposes_token_usage_object(self):
metrics = UsageMetrics(
total_tokens=110,
prompt_tokens=100,
completion_tokens=10,
successful_requests=1,
)
output = LiteAgentOutput(
agent_role="analyst", usage_metrics=metrics.model_dump()
)
assert output.token_usage == metrics
assert isinstance(output.token_usage, UsageMetrics)
def test_lite_agent_output_token_usage_zeroed_when_absent(self):
output = LiteAgentOutput(agent_role="analyst")
assert output.usage_metrics is None
assert output.token_usage == UsageMetrics()
def test_crew_output_exposes_usage_metrics_dict(self):
metrics = UsageMetrics(
total_tokens=110,
prompt_tokens=100,
completion_tokens=10,
successful_requests=1,
)
output = CrewOutput(token_usage=metrics)
assert output.usage_metrics == metrics.model_dump()
assert isinstance(output.usage_metrics, dict)
def test_both_shapes_carry_identical_keys(self):
"""The dict shape has exactly the UsageMetrics fields on both types."""
crew_dict = CrewOutput(token_usage=UsageMetrics()).usage_metrics
lite = LiteAgentOutput(
agent_role="analyst", usage_metrics=UsageMetrics().model_dump()
)
assert set(crew_dict) == set(UsageMetrics.model_fields)
assert set(lite.usage_metrics) == set(UsageMetrics.model_fields)
class TestUsageShapeEndToEnd:
"""Mirror of the EPD-178 clean-room repro, offline via a fake BaseLLM."""
@staticmethod
def _read_via_object(result) -> int:
"""Single accessor written against the CrewOutput shape."""
return result.token_usage.prompt_tokens
@staticmethod
def _read_via_dict(result) -> int:
"""Single accessor written against the LiteAgentOutput shape."""
return result.usage_metrics["prompt_tokens"]
def test_single_accessor_works_on_both_kickoff_paths(self):
agent_a = Agent(
role="analyst",
goal="Answer questions.",
backstory="Test agent.",
llm=_FixedUsageLLM(),
verbose=False,
)
result_agent = agent_a.kickoff("a question")
agent_b = Agent(
role="analyst",
goal="Answer questions.",
backstory="Test agent.",
llm=_FixedUsageLLM(),
verbose=False,
)
task = Task(
description="Answer: a question",
expected_output="A short answer.",
agent=agent_b,
)
crew = Crew(agents=[agent_b], tasks=[task], verbose=False)
result_crew = crew.kickoff()
assert isinstance(result_agent, LiteAgentOutput)
assert isinstance(result_crew, CrewOutput)
# Both accessors work on both result types and agree with each other.
for result in (result_agent, result_crew):
object_read = self._read_via_object(result)
dict_read = self._read_via_dict(result)
assert object_read == dict_read
assert object_read > 0
assert isinstance(result.token_usage, UsageMetrics)
assert isinstance(result.usage_metrics, dict)

View File

@@ -18,11 +18,16 @@ def test_creating_a_tool_using_annotation():
return question
assert my_tool.name == "Name of my tool"
assert "Tool Name: name_of_my_tool" in my_tool.description
assert "Tool Arguments:" in my_tool.description
assert '"question"' in my_tool.description
assert '"type": "string"' in my_tool.description
assert "Tool Description: Clear description for what this tool is useful for" in my_tool.description
# The authored description is preserved as written; the LLM-facing
# composite lives at formatted_description.
assert my_tool.description == (
"Clear description for what this tool is useful for, your agent will need this information to use it."
)
assert "Tool Name: name_of_my_tool" in my_tool.formatted_description
assert "Tool Arguments:" in my_tool.formatted_description
assert '"question"' in my_tool.formatted_description
assert '"type": "string"' in my_tool.formatted_description
assert "Tool Description: Clear description for what this tool is useful for" in my_tool.formatted_description
assert my_tool.args_schema.model_json_schema()["properties"] == {
"question": {"title": "Question", "type": "string"}
}
@@ -33,9 +38,10 @@ def test_creating_a_tool_using_annotation():
converted_tool = my_tool.to_structured_tool()
assert converted_tool.name == "Name of my tool"
assert "Tool Name: name_of_my_tool" in converted_tool.description
assert "Tool Arguments:" in converted_tool.description
assert '"question"' in converted_tool.description
assert converted_tool.description == my_tool.description
assert "Tool Name: name_of_my_tool" in converted_tool.formatted_description
assert "Tool Arguments:" in converted_tool.formatted_description
assert '"question"' in converted_tool.formatted_description
assert converted_tool.args_schema.model_json_schema()["properties"] == {
"question": {"title": "Question", "type": "string"}
}
@@ -56,11 +62,16 @@ def test_creating_a_tool_using_baseclass():
my_tool = MyCustomTool()
assert my_tool.name == "Name of my tool"
assert "Tool Name: name_of_my_tool" in my_tool.description
assert "Tool Arguments:" in my_tool.description
assert '"question"' in my_tool.description
assert '"type": "string"' in my_tool.description
assert "Tool Description: Clear description for what this tool is useful for" in my_tool.description
# The authored description is preserved as written; the LLM-facing
# composite lives at formatted_description.
assert my_tool.description == (
"Clear description for what this tool is useful for, your agent will need this information to use it."
)
assert "Tool Name: name_of_my_tool" in my_tool.formatted_description
assert "Tool Arguments:" in my_tool.formatted_description
assert '"question"' in my_tool.formatted_description
assert '"type": "string"' in my_tool.formatted_description
assert "Tool Description: Clear description for what this tool is useful for" in my_tool.formatted_description
assert my_tool.args_schema.model_json_schema()["properties"] == {
"question": {"title": "Question", "type": "string"}
}
@@ -69,9 +80,10 @@ def test_creating_a_tool_using_baseclass():
converted_tool = my_tool.to_structured_tool()
assert converted_tool.name == "Name of my tool"
assert "Tool Name: name_of_my_tool" in converted_tool.description
assert "Tool Arguments:" in converted_tool.description
assert '"question"' in converted_tool.description
assert converted_tool.description == my_tool.description
assert "Tool Name: name_of_my_tool" in converted_tool.formatted_description
assert "Tool Arguments:" in converted_tool.formatted_description
assert '"question"' in converted_tool.formatted_description
assert converted_tool.args_schema.model_json_schema()["properties"] == {
"question": {"title": "Question", "type": "string"}
}
@@ -695,3 +707,88 @@ class TestToolDecoratorArunValidation:
with pytest.raises(ValueError, match="validation failed"):
await async_execute.arun(wrong_arg="value")
class TestAuthoredDescriptionPreserved:
"""Regression tests for EPD-179: BaseTool.model_post_init silently
rewrote the authored ``description`` into the LLM-facing composite
(``Tool Name: …\\nTool Arguments: …\\nTool Description: <authored>``).
The authored field must survive construction as written, with the
composite exposed separately at ``formatted_description``.
"""
AUTHORED = "Returns the current temperature for a city."
def _make_tool(self) -> BaseTool:
class TempArgs(BaseModel):
city: str = Field(description="City name to look up.")
class TempTool(BaseTool):
name: str = "get_temperature"
description: str = TestAuthoredDescriptionPreserved.AUTHORED
args_schema: type[BaseModel] = TempArgs
def _run(self, city: str) -> str:
return f"22C in {city}"
return TempTool()
def test_description_equals_authored_text(self):
tool_instance = self._make_tool()
assert tool_instance.description == self.AUTHORED
def test_formatted_description_contains_composite(self):
tool_instance = self._make_tool()
formatted = tool_instance.formatted_description
assert "Tool Name: get_temperature" in formatted
assert "Tool Arguments:" in formatted
assert '"city"' in formatted
assert formatted.endswith(f"Tool Description: {self.AUTHORED}")
def test_formatted_description_tracks_later_description_edits(self):
tool_instance = self._make_tool()
tool_instance.description = "Edited description."
assert tool_instance.formatted_description.endswith(
"Tool Description: Edited description."
)
def test_prose_mentioning_the_marker_is_not_truncated(self):
"""Authored text that merely mentions "Tool Description:" must reach
the LLM untouched — only descriptions that ARE a pre-composed block
(anchored three-line shape) get stripped."""
tool_instance = self._make_tool()
prose = (
"Formats prompts. The output includes a line reading "
"'Tool Description:' followed by the tool's summary."
)
tool_instance.description = prose
assert tool_instance.formatted_description.endswith(
f"Tool Description: {prose}"
)
def test_composite_is_not_reapplied_to_prebaked_descriptions(self):
"""A description that already contains a composed block (old
checkpoints, adapters that bake the composite into the field) must
not be double-wrapped."""
tool_instance = self._make_tool()
tool_instance.description = (
"Tool Name: get_temperature\n"
'Tool Arguments: {"city": "str"}\n'
f"Tool Description: {self.AUTHORED}"
)
formatted = tool_instance.formatted_description
assert formatted.count("Tool Description:") == 1
assert formatted.endswith(f"Tool Description: {self.AUTHORED}")
def test_prompt_rendering_still_uses_composite(self):
from crewai.utilities.agent_utils import render_text_description_and_args
tool_instance = self._make_tool()
structured = tool_instance.to_structured_tool()
assert structured.description == self.AUTHORED
for candidate in (tool_instance, structured):
rendered = render_text_description_and_args([candidate])
assert "Tool Name: get_temperature" in rendered
assert "Tool Arguments:" in rendered
assert f"Tool Description: {self.AUTHORED}" in rendered

View File

@@ -0,0 +1,46 @@
"""Flow panels must be suppressed while a TUI owns the screen."""
from rich.text import Text
from crewai.events.listeners.tracing.utils import set_tui_mode
from crewai.events.utils.console_formatter import ConsoleFormatter
def _make_formatter(monkeypatch):
fmt = ConsoleFormatter(verbose=True)
calls: list[object] = []
monkeypatch.setattr(fmt, "print", lambda *a, **k: calls.append(a))
return fmt, calls
def test_flow_panel_suppressed_in_tui_mode(monkeypatch):
fmt, calls = _make_formatter(monkeypatch)
set_tui_mode(True)
try:
fmt.print_panel(Text("x"), "🌊 Flow Started", "blue", is_flow=True)
finally:
set_tui_mode(False)
assert calls == []
def test_flow_panel_prints_when_not_tui_mode(monkeypatch):
fmt, calls = _make_formatter(monkeypatch)
set_tui_mode(False)
fmt.print_panel(Text("x"), "🌊 Flow Started", "blue", is_flow=True)
# Panel + trailing blank line.
assert len(calls) == 2
def test_non_flow_panel_unaffected_by_tui_mode(monkeypatch):
# tui_mode only gates flow panels; regular panels still follow verbose.
fmt, calls = _make_formatter(monkeypatch)
set_tui_mode(True)
try:
fmt.print_panel(Text("x"), "Task", "blue", is_flow=False)
finally:
set_tui_mode(False)
assert len(calls) == 2