#
# 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())