Merge pull request #3231 from pipecat-ai/aleix/improve-evals-assert-on-exit

evals: use EndFrame reason field to provide eval result
This commit is contained in:
Aleix Conchillo Flaqué
2025-12-15 13:23:29 -08:00
committed by GitHub
7 changed files with 40 additions and 21 deletions

View File

@@ -0,0 +1,3 @@
- Changed the `reason` field in `EndFrame`, `CancelFrame`, `EndTaskFrame`, and
`CancelTaskFrame` from `str` to `Any` to indicate that it can hold values
other than strings.

View File

@@ -17,7 +17,6 @@ from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.aggregators.llm_response import LLMAssistantAggregatorParams
from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair
from pipecat.processors.transcript_processor import TranscriptProcessor from pipecat.processors.transcript_processor import TranscriptProcessor
from pipecat.runner.types import RunnerArguments from pipecat.runner.types import RunnerArguments

View File

@@ -20,7 +20,6 @@ from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.aggregators.llm_response import LLMAssistantAggregatorParams
from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair
from pipecat.runner.types import RunnerArguments from pipecat.runner.types import RunnerArguments
from pipecat.runner.utils import create_transport from pipecat.runner.utils import create_transport

View File

@@ -18,7 +18,6 @@ from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.aggregators.llm_response import LLMAssistantAggregatorParams
from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair
from pipecat.runner.types import RunnerArguments from pipecat.runner.types import RunnerArguments
from pipecat.runner.utils import ( from pipecat.runner.utils import (

View File

@@ -31,7 +31,13 @@ from pipecat.adapters.schemas.function_schema import FunctionSchema
from pipecat.adapters.schemas.tools_schema import ToolsSchema from pipecat.adapters.schemas.tools_schema import ToolsSchema
from pipecat.audio.vad.silero import SileroVADAnalyzer from pipecat.audio.vad.silero import SileroVADAnalyzer
from pipecat.audio.vad.vad_analyzer import VADParams from pipecat.audio.vad.vad_analyzer import VADParams
from pipecat.frames.frames import EndTaskFrame, LLMRunFrame, OutputImageRawFrame from pipecat.frames.frames import (
CancelFrame,
EndFrame,
EndTaskFrame,
LLMRunFrame,
OutputImageRawFrame,
)
from pipecat.pipeline.pipeline import Pipeline from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
@@ -50,6 +56,7 @@ SCRIPT_DIR = Path(__file__).resolve().parent
PIPELINE_IDLE_TIMEOUT_SECS = 60 PIPELINE_IDLE_TIMEOUT_SECS = 60
EVAL_TIMEOUT_SECS = 120 EVAL_TIMEOUT_SECS = 120
EVAL_RESULT_TIMEOUT_SECS = 10
EvalPrompt = str | Tuple[str, ImageFile] EvalPrompt = str | Tuple[str, ImageFile]
@@ -78,7 +85,7 @@ class EvalRunner:
self._log_level = log_level self._log_level = log_level
self._total_success = 0 self._total_success = 0
self._tests: List[EvalResult] = [] self._tests: List[EvalResult] = []
self._queue = asyncio.Queue() self._result_future: Optional[asyncio.Future[bool]] = None
# We to save runner files. # We to save runner files.
name = name or f"{datetime.now().strftime('%Y%m%d_%H%M%S')}" name = name or f"{datetime.now().strftime('%Y%m%d_%H%M%S')}"
@@ -88,16 +95,16 @@ class EvalRunner:
os.makedirs(self._logs_dir, exist_ok=True) os.makedirs(self._logs_dir, exist_ok=True)
os.makedirs(self._recordings_dir, exist_ok=True) os.makedirs(self._recordings_dir, exist_ok=True)
async def assert_eval(self, params: FunctionCallParams): async def function_assert_eval(self, params: FunctionCallParams):
result = params.arguments["result"] result = params.arguments["result"]
reasoning = params.arguments["reasoning"] reasoning = params.arguments["reasoning"]
logger.debug(f"🧠 EVAL REASONING(result: {result}): {reasoning}") logger.debug(f"🧠 EVAL REASONING(result: {result}): {reasoning}")
await self._queue.put(result)
await params.result_callback(None) await params.result_callback(None)
await params.llm.push_frame(EndTaskFrame(), FrameDirection.UPSTREAM) await params.llm.push_frame(EndTaskFrame(reason=result), FrameDirection.UPSTREAM)
async def assert_eval_false(self): async def assert_eval(self, result: bool):
await self._queue.put(False) if self._result_future:
self._result_future.set_result(result)
async def run_eval( async def run_eval(
self, self,
@@ -117,6 +124,9 @@ class EvalRunner:
start_time = time.time() start_time = time.time()
# Create a future to store the eval result.
self._result_future = asyncio.get_running_loop().create_future()
try: try:
tasks = [ tasks = [
asyncio.create_task(run_example_pipeline(script_path, eval_config)), asyncio.create_task(run_example_pipeline(script_path, eval_config)),
@@ -136,8 +146,10 @@ class EvalRunner:
logger.error(f"ERROR: Unable to run {example_file}: {e}") logger.error(f"ERROR: Unable to run {example_file}: {e}")
try: try:
result = await asyncio.wait_for(self._queue.get(), timeout=1.0) # Wait for the future to resolve.
result = await asyncio.wait_for(self._result_future, timeout=EVAL_RESULT_TIMEOUT_SECS)
except asyncio.TimeoutError: except asyncio.TimeoutError:
logger.error(f"ERROR: Timeout waiting for eval result.")
result = False result = False
if result: if result:
@@ -244,7 +256,7 @@ async def run_eval_pipeline(
llm = OpenAILLMService(api_key=os.getenv("OPENAI_API_KEY")) llm = OpenAILLMService(api_key=os.getenv("OPENAI_API_KEY"))
llm.register_function("eval_function", eval_runner.assert_eval) llm.register_function("eval_function", eval_runner.function_assert_eval)
eval_function = FunctionSchema( eval_function = FunctionSchema(
name="eval_function", name="eval_function",
@@ -348,7 +360,14 @@ async def run_eval_pipeline(
@task.event_handler("on_idle_timeout") @task.event_handler("on_idle_timeout")
async def on_pipeline_idle_timeout(task): async def on_pipeline_idle_timeout(task):
await eval_runner.assert_eval_false() await eval_runner.assert_eval(False)
@task.event_handler("on_pipeline_finished")
async def on_pipeline_finished(task, frame):
if isinstance(frame, EndFrame):
await eval_runner.assert_eval(frame.reason)
elif isinstance(frame, CancelFrame):
await eval_runner.assert_eval(False)
# TODO(aleix): We should handle SIGINT and SIGTERM so we can cancel both the # TODO(aleix): We should handle SIGINT and SIGTERM so we can cancel both the
# eval and the example. # eval and the example.

View File

@@ -30,13 +30,13 @@ EVAL_SIMPLE_MATH = EvalConfig(
) )
EVAL_WEATHER = EvalConfig( EVAL_WEATHER = EvalConfig(
prompt="What's the weather in San Francisco (in farhenheit or celsius)?", prompt="What's the weather in San Francisco? Temperature should be in any unit, just pick one.",
eval="The user says something specific about the current weather in San Francisco, including the degrees (in farhenheit or celsius).", eval="The user talks about the weather in San Francisco, including the degrees.",
) )
EVAL_ONLINE_SEARCH = EvalConfig( EVAL_ONLINE_SEARCH = EvalConfig(
prompt="What's the date right now in London?", prompt="What's the current date in UTC?",
eval=f"The user says today is {datetime.now(timezone.utc).strftime('%B %d, %Y')} in London.", eval=f"Current date in UTC is {datetime.now(timezone.utc).strftime('%A, %B %d, %Y')}.",
) )
EVAL_SWITCH_LANGUAGE = EvalConfig( EVAL_SWITCH_LANGUAGE = EvalConfig(

View File

@@ -947,7 +947,7 @@ class CancelFrame(SystemFrame):
reason: Optional reason for pushing a cancel frame. reason: Optional reason for pushing a cancel frame.
""" """
reason: Optional[str] = None reason: Optional[Any] = None
def __str__(self): def __str__(self):
return f"{self.name}(reason: {self.reason})" return f"{self.name}(reason: {self.reason})"
@@ -1537,7 +1537,7 @@ class EndTaskFrame(TaskFrame):
reason: Optional reason for pushing an end frame. reason: Optional reason for pushing an end frame.
""" """
reason: Optional[str] = None reason: Optional[Any] = None
def __str__(self): def __str__(self):
return f"{self.name}(reason: {self.reason})" return f"{self.name}(reason: {self.reason})"
@@ -1555,7 +1555,7 @@ class CancelTaskFrame(TaskFrame):
reason: Optional reason for pushing a cancel frame. reason: Optional reason for pushing a cancel frame.
""" """
reason: Optional[str] = None reason: Optional[Any] = None
def __str__(self): def __str__(self):
return f"{self.name}(reason: {self.reason})" return f"{self.name}(reason: {self.reason})"
@@ -1634,7 +1634,7 @@ class EndFrame(ControlFrame):
reason: Optional reason for pushing an end frame. reason: Optional reason for pushing an end frame.
""" """
reason: Optional[str] = None reason: Optional[Any] = None
def __str__(self): def __str__(self):
return f"{self.name}(reason: {self.reason})" return f"{self.name}(reason: {self.reason})"