Source code for pipecat.evals.matcher

#
# Copyright (c) 2024-2026, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#

"""Matching a scenario's expectations against the bot's events.

:class:`ExpectationMatcher` waits on the run's
:class:`~pipecat.evals.events.EvalEventStream` for the event an expectation
names and verifies it: a payload check, a judge verdict, the aggregation of a
reply across its segments, the any-order matching of a turn's function calls,
and the inverted ``absent:`` check.
"""

import time

from loguru import logger

from pipecat.evals.events import EvalEventStream
from pipecat.evals.judge import EvalJudge
from pipecat.evals.results import EvalAssertionFailure, EvalTrace
from pipecat.evals.script import FUNCTION_CALL_EVENTS, EvalExpectation

# How long a judge's "no" waits for more of the reply once the bot has stopped
# speaking: the transcription of its last sentence lands after the stop.
JUDGE_NO_GRACE_S = 2.0


[docs] class ExpectationMatcher: """Matches one expectation at a time against the event stream. Expected events must appear in order, with unmatched events allowed in between. A reply with a content check aggregates its segments and re-checks on each, so an interim "Let me check" is rolled past rather than taken for the answer. A turn's function calls match by name in any order. """
[docs] def __init__(self, *, stream: EvalEventStream, judge: EvalJudge | None, trace: EvalTrace): """Initialize the matcher. Args: stream: The bot's events, consumed as expectations match. judge: The judge for ``eval:`` assertions, or ``None`` when the scenario has none. trace: The run's trace, for the matcher's progress. """ self._stream = stream self._judge = judge self._trace = trace # function_call events popped while matching another expectation, held so # the turn's calls can be matched by name in any order (reset per turn). self._pending_function_calls: list[dict] = [] # Text content of the most recently matched event (the bot's response, or # a user transcript), surfaced to verbose progress. Empty for events with # no text (llm_started, function_call, speaking events). self.last_match_text: str = "" # When the most recent expectation matched, so an absence check can tell # a reply that began after it from the rest of the one it matched. self._last_match_at: float = 0.0
[docs] def reset_turn(self) -> None: """Forget the previous turn's unclaimed function calls.""" self._pending_function_calls = []
[docs] async def match( self, expectation: EvalExpectation, anchor: float, budget_ms: int, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Wait for the expected event and verify it. Stale bot output is not filtered here: the stream drops it on an interruption and the driver before each send. Args: expectation: The expectation to match. anchor: Monotonic time the turn's budget is measured from. budget_ms: The expectation's latency budget. turn_idx: Index of the turn, for the failure record. exp_idx: Index of the expectation within the turn, for the record. Returns: The failure, or ``None`` when the expectation is satisfied. Raises: TimeoutError: When no matching event arrives at all (so the caller can report "no matching event arrived"). A response that arrives but never satisfies the content check returns a failure instead. """ deadline = anchor + (budget_ms / 1000.0) self.last_match_text = "" if expectation.absent: return await self._match_absent(expectation, deadline, budget_ms, turn_idx, exp_idx) if expectation.aggregates: return await self._match_aggregating( expectation, deadline, budget_ms, turn_idx, exp_idx ) if expectation.event in FUNCTION_CALL_EVENTS: # A call expectation holds the set of calls the turn should make; it # completes only when all are found, in any order (a response # arriving doesn't short-circuit it). return await self._match_function_calls(expectation, deadline, turn_idx, exp_idx) return await self._match_one(expectation, deadline, turn_idx, exp_idx)
async def _match_one( self, expectation: EvalExpectation, deadline: float, turn_idx: int, exp_idx: int ) -> EvalAssertionFailure | None: """Match a single event and check its payload and judge assertion.""" self._trace.log(f"match: waiting for {expectation.event!r}") event = await self._stream.next_event(expectation.event, deadline) payload_failure = self._check_payload(event, expectation, turn_idx, exp_idx) if payload_failure: return payload_failure judge_failure = await self._check_judge(event, expectation, turn_idx, exp_idx) if judge_failure is None: self.last_match_text = self._match_summary(event) self._last_match_at = time.monotonic() return judge_failure async def _match_aggregating( self, expectation: EvalExpectation, deadline: float, budget_ms: int, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Accumulate reply segments until the content check passes or fails.""" if expectation.eval is not None and self._judge is None: return self._failure( expectation, turn_idx, exp_idx, "scenario uses 'eval:' but no judge could be built", "no_judge", ) check = "+".join( name for name, val in ( ("text_contains", expectation.text_contains), ("eval", expectation.eval), ) if val is not None ) self._trace.log(f"match: waiting for {expectation.event!r} ({check})") aggregate = "" last_reason = "" seen_any = False pending: dict | None = None while True: if pending is not None: event, pending = pending, None else: try: event = await self._stream.next_event(expectation.event, deadline) except TimeoutError: if not seen_any: raise # no response at all: the caller reports the missing event return self._unsatisfied(expectation, turn_idx, exp_idx, budget_ms, last_reason) seen_any = True delta = self._event_text(event) aggregate += delta excluded = self._text_excluded(aggregate, expectation, turn_idx, exp_idx) if excluded is not None: return excluded # Feed each segment to the judge as its own assistant message, so it # judges the bot's reply in the conversation's context (the cumulative # `aggregate` is kept only for text_contains and the match summary). if expectation.eval is not None and self._judge is not None: self._judge.add_assistant_message(delta) status, reason = await self._evaluate_aggregate(aggregate, expectation) self._trace.log(f"eval: {status} (aggregate={aggregate.strip()!r}) {reason}") if status == "pass": self.last_match_text = aggregate self._last_match_at = time.monotonic() return None if status == "fail": # Only the judge can affirmatively fail an aggregate, and only # once the reply is complete: a "no" on a reply still being # spoken, or whose last sentence is still being transcribed, is # a "continue" that the next segment may turn into a "yes". pending = await self._rest_of_reply(expectation.event, deadline) if pending is None: return self._failure(expectation, turn_idx, exp_idx, reason, "judge_no") self._trace.log("eval: no, but the reply goes on: judging the rest") # "continue": wait for the next segment, separated by a space so # sentences don't run together (e.g. "...that. The weather..."). aggregate += " " last_reason = reason async def _rest_of_reply(self, event_type: str, deadline: float) -> dict | None: """The next segment of a reply the judge rejected, or ``None`` when the reply is over. While the bot is speaking the next segment is awaited within the turn's budget; once it has stopped, only for the grace its last sentence's transcription needs. """ if self._stream.bot_speaking: until = deadline else: until = min(deadline, time.monotonic() + JUDGE_NO_GRACE_S) try: return await self._stream.next_event(event_type, until) except TimeoutError: return None def _unsatisfied( self, expectation: EvalExpectation, turn_idx: int, exp_idx: int, budget_ms: int, last_reason: str, ) -> EvalAssertionFailure: """The failure of a reply that never satisfied its check within the budget. With ``eval:`` the judge kept saying ``continue``; without it the only way to be unsatisfied is a missing substring. """ self._trace.log(f"eval: timeout, not satisfied: {last_reason}") return self._failure( expectation, turn_idx, exp_idx, f"not satisfied within {budget_ms}ms: {last_reason}", "judge_continue" if expectation.eval is not None else "text_mismatch", ) async def _match_absent( self, expectation: EvalExpectation, deadline: float, budget_ms: int, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Pass when no event of this type arrives before the deadline; an arriving one fails at once, with its content. A reply reaches the stream in segments, one per pause in the bot's speech, so a ``response`` that continues the reply an earlier expectation matched is not a new one: only a response the bot began after that match counts. """ self._trace.log(f"match: expecting NO {expectation.event!r} for {budget_ms}ms") while True: try: event = await self._stream.next_event(expectation.event, deadline) except TimeoutError: # The quiet window held: absence confirmed. self.last_match_text = f"no {expectation.event!r} for {budget_ms}ms" return None if expectation.event == "response" and not self._reply_began_after_last_match(): self._trace.log( f"absent: the matched reply goes on, not a new one: " f"{self._match_summary(event)!r}" ) continue break return self._failure( expectation, turn_idx, exp_idx, f"expected no {expectation.event!r} within {budget_ms}ms, " f"but one arrived: {self._match_summary(event)}", "unexpected_event", ) def _reply_began_after_last_match(self) -> bool: """Whether the bot started a reply since the most recent match, by its LLM or its speech.""" times = self._stream.latest_event_times began = max(times.get("llm_started", 0.0), times.get("bot_started_speaking", 0.0)) return began > self._last_match_at async def _match_function_calls( self, expectation: EvalExpectation, deadline: float, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Match every call in the expectation, in any order, within the budget; else a failure naming the call that was missing or whose args did not match. With ``eval:``, each matched call is also put to the judge, and the first one it rejects fails the expectation. """ if expectation.eval is not None and self._judge is None: return self._failure( expectation, turn_idx, exp_idx, "scenario uses 'eval:' but no judge could be built", "no_judge", ) matched: list[str] = [] for spec in expectation.calls or []: want = spec.args or None self._trace.log(f"match: waiting for {expectation.event!r} ({spec.signature})") try: event = await self._next_function_call(spec.name, deadline, want, expectation.event) except TimeoutError: # A call of the right name with the wrong arguments is a different # failure from the call never being made, and the bot's arguments # are what the reader needs to see. near = [ ev.get("args") for ev in self._pending_function_calls if ev.get("type") == expectation.event and (spec.name is None or ev.get("name") == spec.name) ] if want is not None and near: return self._failure( expectation, turn_idx, exp_idx, f"no {spec.name!r} call had args {want!r} (saw {near!r})", "function_args_mismatch", ) missing = spec.name or "any function" seen = ", ".join(matched) if matched else "none" return self._failure( expectation, turn_idx, exp_idx, f"function call {missing!r} not seen (matched: {seen})", "missing_function_call", ) judge_failure = await self._check_call_judge(event, expectation, turn_idx, exp_idx) if judge_failure: return judge_failure matched.append(str(event.get("name"))) self.last_match_text = ", ".join(matched) or "function call" self._last_match_at = time.monotonic() return None async def _check_call_judge( self, event: dict, expectation: EvalExpectation, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Put a matched call to the judge if ``eval:`` was set on the expectation. The judge is asked about the call by name and arguments, in the light of the spoken conversation so far. A call is not a partial reply, so a ``continue`` fails it like a ``no``. """ if expectation.eval is None: return None # _match_function_calls fails before matching when there is no judge. assert self._judge is not None name = str(event.get("name") or "?") args = event.get("args") or {} with logger.contextualize(eval_pipeline="judge"): verdict = await self._judge.evaluate_call(name, args, expectation.eval) self._trace.log(f"eval: {verdict.verdict} ({self._match_summary(event)}) {verdict.reason}") if verdict.passed: return None return self._failure( expectation, turn_idx, exp_idx, f"eval {expectation.eval!r} on {self._match_summary(event)}: " f"judge said {verdict.verdict} — {verdict.reason}", "judge_no", ) async def _next_function_call( self, name: str | None, deadline: float, args: dict | None = None, event_type: str = "function_call", ) -> dict: """The next call event matching ``name`` (``None`` for any) and ``args``. Calls seen but not yet claimed are buffered, so a turn's calls can arrive in any order and a call the LLM corrects and repeats still satisfies it. Raises TimeoutError at ``deadline``. """ def matches(ev: dict) -> bool: if ev.get("type") != event_type: return False if name is not None and ev.get("name") != name: return False if args is None: return True actual = ev.get("args") or {} return all(actual.get(k) == v for k, v in args.items()) for i, ev in enumerate(self._pending_function_calls): if matches(ev): return self._pending_function_calls.pop(i) while True: event = await self._stream.next_any(deadline) if event.get("type") not in FUNCTION_CALL_EVENTS: continue if matches(event): return event self._pending_function_calls.append(event) async def _evaluate_aggregate( self, aggregate: str, expectation: EvalExpectation ) -> tuple[str, str]: """Check the accumulated reply text: ``pass``, ``fail``, or ``continue`` for more text. A missing substring is ``continue``; only the judge can ``fail``. """ if expectation.text_contains is not None and not self._text_contains( aggregate, expectation.text_contains ): return ("continue", f"does not contain {expectation.text_contains!r}") if expectation.eval is not None: if not aggregate.strip(): return ("continue", "no response text yet") # match() guarantees a judge exists before aggregating eval:. assert self._judge is not None with logger.contextualize(eval_pipeline="judge"): # The reply segments were added to the judge's conversation in the # aggregation loop; the judge evaluates that context, not `aggregate`. verdict = await self._judge.evaluate(expectation.eval) if verdict.verdict == "no": return ("fail", f"judge said no: {verdict.reason}") if verdict.verdict == "continue": return ("continue", f"judge said continue: {verdict.reason}") return ("pass", f"judge said yes: {verdict.reason}") return ("pass", "") def _check_payload( self, event: dict, expectation: EvalExpectation, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Apply payload-level checks to a matched event. Returns the first failure or None.""" if expectation.text_contains is not None: content = self._event_text(event) if not self._text_contains(content, expectation.text_contains): return self._failure( expectation, turn_idx, exp_idx, f"text {content!r} does not contain {expectation.text_contains!r}", "text_mismatch", ) excluded = self._text_excluded(self._event_text(event), expectation, turn_idx, exp_idx) if excluded is not None: return excluded if expectation.marker is not None: kind = event.get("kind") wanted = ( ("short", "long") if expectation.marker == "incomplete" else (expectation.marker,) ) if kind not in wanted: return self._failure( expectation, turn_idx, exp_idx, f"marker {self._event_text(event)!r} is {kind or 'of no known kind'}, " f"expected {expectation.marker}", "marker_mismatch", ) problem = self._check_marker_format(event, expectation) if problem is not None: return self._failure(expectation, turn_idx, exp_idx, problem, "marker_format") return None def _check_marker_format(self, event: dict, expectation: EvalExpectation) -> str | None: """What is wrong with the shape of the response's raw text, if anything. The checks read the raw text the marker event carries against every marker the bot recognizes, so they see what the LLM wrote before the bot held any of it back. """ if ( expectation.marker_first is None and expectation.markers is None and expectation.text_after is None ): return None raw: str = event.get("raw") or "" known: list[str] = event.get("markers") or [] found = sorted((raw.find(m), m) for m in known if m in raw) count = sum(raw.count(m) for m in known) first_at, first = found[0] if found else (-1, None) excerpt = repr(raw[:80]) if expectation.markers is not None and count != expectation.markers: return f"raw text holds {count} marker(s), expected {expectation.markers}: {excerpt}" if expectation.marker_first is not None: is_first = first is not None and not raw[:first_at].strip() if is_first != expectation.marker_first: return ( f"raw text {'starts' if is_first else 'does not start'} with a marker, " f"expected it {'to' if expectation.marker_first else 'not to'}: {excerpt}" ) if expectation.text_after is not None: if first is None: return f"raw text holds no marker to check text after: {excerpt}" has_after = bool(raw[first_at + len(first) :].strip()) if has_after != expectation.text_after: return ( f"raw text has {'text' if has_after else 'nothing'} after the marker, " f"expected {'text' if expectation.text_after else 'nothing'}: {excerpt}" ) return None async def _check_judge( self, event: dict, expectation: EvalExpectation, turn_idx: int, exp_idx: int, ) -> EvalAssertionFailure | None: """Run the judge assertion if ``eval:`` was set on this expectation.""" if expectation.eval is None: return None if self._judge is None: return self._failure( expectation, turn_idx, exp_idx, "scenario uses 'eval:' but no judge could be built", "no_judge", ) content = event.get("text") or event.get("transcript") if not content: return self._failure( expectation, turn_idx, exp_idx, f"event has no text/transcript to judge: {event!r}", "no_content", ) self._judge.add_assistant_message(content) verdict = await self._judge.evaluate(expectation.eval) if not verdict.passed: return self._failure( expectation, turn_idx, exp_idx, f"eval {expectation.eval!r}: judge said no — {verdict.reason}", "judge_no", ) return None def _failure( self, expectation: EvalExpectation, turn_idx: int, exp_idx: int, reason: str, kind: str ) -> EvalAssertionFailure: """Build the failure record for ``expectation``.""" return EvalAssertionFailure(turn_idx, exp_idx, expectation.event, reason, kind) def _match_summary(self, event: dict) -> str: """A short label for a matched event: the call signature, or the event's text.""" if event.get("type") == "function_call": args = event.get("args") or {} sig = ", ".join(f"{k}={v}" for k, v in args.items()) return f"{event.get('name') or '?'}({sig})" return self._event_text(event) def _event_text(self, event: dict) -> str: """The text an event carries: reply events use ``text``, ``user_transcription`` ``transcript``.""" return event.get("text") or event.get("transcript") or "" def _text_excluded( self, content: str, expectation: EvalExpectation, turn_idx: int, exp_idx: int ) -> EvalAssertionFailure | None: """The failure for ``text_excludes`` when ``content`` holds it, else ``None``.""" if expectation.text_excludes is None or not self._text_contains( content, expectation.text_excludes ): return None return self._failure( expectation, turn_idx, exp_idx, f"text {content.strip()!r} contains {expectation.text_excludes!r}", "text_present", ) def _text_contains(self, content: str, needle: str) -> bool: """Whether ``needle`` occurs in ``content``, ignoring spacing.""" return " ".join(needle.split()) in " ".join(content.split())