mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
refactor(search): move the MCP search page into application/catalog_search
The orchestration behind MCP catalog_search (band, evidence rerank, hub merge, routed groups, the discovery experiment and the records) lived in mcp.py; the layering table says a router holds no query orchestration. It now lives in application/catalog_search.py and the MCP layer only resolves who is asking and shapes the rows. No behaviour change: both MCP surfaces answer byte-identically over a fixed query set with routed discovery on and off (verified locally with a snapshot, not committed because catalog ids change weekly). Fragment: search-experiment.md names the use case.
This commit is contained in:
@@ -2,6 +2,7 @@
|
||||
title: Discovery experiment — a relevance judge behind catalog search, measured on what the caller does next
|
||||
status: building
|
||||
sources:
|
||||
- src/treg/application/catalog_search.py
|
||||
- src/treg/application/search_experiment.py
|
||||
- src/treg/domain/catalog/interleave.py
|
||||
- src/treg/infra/judge.py
|
||||
@@ -41,10 +42,16 @@ Three layers, imports pointing inward:
|
||||
confidence and every option's probability); a candidate in `job_view`'s shape is asked about a
|
||||
whole job, under `job_criteria` ([find](find.md)). The experiment passes none of these, so its
|
||||
question and its cache key are unchanged while it runs.
|
||||
- **`application/search_experiment.py`** — the use case. Judges the candidates, builds the judged
|
||||
- **`application/catalog_search.py`** — the search use case behind MCP `catalog_search` on both
|
||||
surfaces: the shipped ranker's page (`store.rank_band` over a band wider than the page, the
|
||||
evidence rerank, listed hub tools merged by score, routed groups, cut to the page), then the
|
||||
experiment's say over it, then the records (`audit.record_search` while the experiment is on,
|
||||
`record_search_miss` whenever the LEXICAL page is empty). The hub read holds its own session and
|
||||
closes it before any judge request. The MCP layer resolves who is asking and shapes the rows.
|
||||
- **`application/search_experiment.py`** — the experiment. Judges the candidates, builds the judged
|
||||
page with the SAME finishing steps the baseline had (evidence rerank, routed grouping, cut to the
|
||||
page — the MCP layer passes that function in), deals the caller an arm, decides what is shown, and
|
||||
hands the MCP layer a row for `audit.record_search`.
|
||||
page — the use case passes that function in), deals the caller an arm, decides what is shown, and
|
||||
hands back a row for `audit.record_search`.
|
||||
|
||||
The judged page is **bucketed, not sorted by probability**: rows under `search_judge_keep` (0.4)
|
||||
are dropped, rows at or over `search_judge_high` (0.7) go first, and inside a bucket the lexical
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
"""Catalog search for agents - the use case behind MCP `catalog_search` on both surfaces.
|
||||
|
||||
The page is the shipped ranker: `store.rank_band` over a band wider than the page (so a routed group
|
||||
can collapse without starving it), the evidence rerank, the listed hub tools merged by score with no
|
||||
boost (docs/hub-listing-decisions.md), routed parents grouped with their children, cut to the
|
||||
page. Then the discovery experiment (`search_experiment`) may judge a wider recall and decide what
|
||||
this caller sees, and the search is recorded. Nothing here touches `store.search` or its scoring;
|
||||
the MCP layer only shapes the rows this returns.
|
||||
|
||||
Session discipline: the hub read opens, uses and closes its own session before any judge request,
|
||||
so no connection is held while the judge thinks.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Awaitable, Callable
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
from .. import audit
|
||||
from ..config import get_settings
|
||||
from ..domain.catalog import store as catalog_store
|
||||
from ..infra.db import session_maker
|
||||
from . import hub as hub_app
|
||||
from . import search_experiment
|
||||
|
||||
Rows = list[tuple[dict, float]]
|
||||
Observed = Callable[[list[str]], Awaitable[dict[str, dict]]]
|
||||
# (caller key, team id, email) for the experiment's log - resolved only while the experiment is on
|
||||
Identify = Callable[[], Awaitable[tuple[str | None, int | None, str | None]]]
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Caller:
|
||||
"""Who is searching, as far as the hub's lists need to know: the team slug and sign-in email
|
||||
that decide which listed hub tools this caller may see (None, None: an unknown reader, no hub)."""
|
||||
hub_slug: str | None = None
|
||||
hub_email: str | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class Page:
|
||||
rows: Rows # the page to serve, in order
|
||||
total: int # matches before the page cut (hub rows included)
|
||||
tie_truncated: bool # the tie group outran what the evidence sort weighs
|
||||
stats: dict[str, dict] # observed stats for every row on the page
|
||||
hidden: dict[str, int] # routed parent id -> children the page cut
|
||||
steering: bool # routed discovery on: parents lead their groups
|
||||
lexical_empty: bool = False # the shipped ranker admitted nothing (the miss log's signal)
|
||||
arm: str | None = None # the experiment's arm, when it ran
|
||||
caller_key: str | None = None # the experiment's handle on the caller, when it ran
|
||||
log: dict = field(default_factory=dict) # the experiment's SearchLog fields, when it ran
|
||||
|
||||
|
||||
def steering() -> bool:
|
||||
return str(get_settings().routed_discovery).strip().lower() not in ("off", "0", "false", "no")
|
||||
|
||||
|
||||
def _group(rows: Rows, limit: int) -> tuple[Rows, dict[str, int]]:
|
||||
"""A capability with a ROUTED row shows the parent first and its children right under it, so
|
||||
an agent sees "let treg choose" before the specific providers; the parent counts the children
|
||||
the page did not show."""
|
||||
grouped = catalog_store.group_routed(
|
||||
[{"ep": ep, "score": score, "capability": ep.get("capability"), "kind": ep.get("kind")} for ep, score in rows],
|
||||
max_children=catalog_store.MAX_ROUTED_CHILDREN)
|
||||
hidden = {r["ep"]["id"]: r["children_hidden"] for r in grouped if r.get("children_hidden")}
|
||||
return [(r["ep"], r["score"]) for r in grouped][:limit], hidden
|
||||
|
||||
|
||||
async def lexical(query: str, cat: catalog_store.Catalog, limit: int, *, observed: Observed,
|
||||
caller: Caller) -> Page:
|
||||
"""The shipped ranker's page. Score, then let the evidence break the ties: token scoring
|
||||
produces ties by the dozen, and with a page of 8 the rows an agent sees would otherwise be
|
||||
decided by file order. The band is widened only so routed groups can collapse without
|
||||
starving the page; with steering off there is no collapsing, so the page is the band."""
|
||||
steer = steering()
|
||||
ranked, total, tie_truncated = catalog_store.rank_band(query, cat, min(100, limit * 4) if steer else limit)
|
||||
stats = await observed([ep["id"] for ep, _ in ranked])
|
||||
ranked = catalog_store.rerank(ranked, stats, cat)
|
||||
async with session_maker() as db:
|
||||
hub_ranked, hub_stats = await hub_app.search_listed(db, query, cat, org_slug=caller.hub_slug,
|
||||
email=caller.hub_email)
|
||||
if hub_ranked:
|
||||
stats = {**stats, **hub_stats}
|
||||
ranked = catalog_store.merge_by_score(ranked, hub_ranked)
|
||||
total += len(hub_ranked)
|
||||
rows, hidden = _group(ranked, limit)
|
||||
return Page(rows, total, tie_truncated, stats, hidden, steer)
|
||||
|
||||
|
||||
async def search(query: str, limit: int, *, cat: catalog_store.Catalog, source: str, caller: Caller,
|
||||
observed: Observed, identify: Identify) -> Page:
|
||||
"""The page this caller sees, recorded. The lexical page first; while the experiment is on, the
|
||||
judge reads a wider recall and the arm decides what is shown - whatever it does, the result is
|
||||
a page. The miss log is judged by the LEXICAL page: it measures the shipped ranker's coverage,
|
||||
and a judged page that found something is the experiment's result, not a reason to stop
|
||||
recording the gap."""
|
||||
page = await lexical(query, cat, limit, observed=observed, caller=caller)
|
||||
page.lexical_empty = not page.rows and bool(query.strip())
|
||||
if query.strip() and search_experiment.mode() != "off":
|
||||
async def finish(rows: Rows) -> tuple[Rows, dict[str, dict]]:
|
||||
st = await observed([ep["id"] for ep, _ in rows])
|
||||
grouped, _ = _group(catalog_store.rerank(rows, st, cat), limit)
|
||||
return grouped, st
|
||||
key, org_id, email = await identify()
|
||||
exp = await search_experiment.run(query, cat, baseline=page.rows, baseline_total=page.total,
|
||||
limit=limit, caller=key, finish=finish)
|
||||
page.stats = {**exp.stats, **page.stats}
|
||||
page.rows, page.arm, page.log, page.caller_key = exp.shown, exp.arm, exp.log, key
|
||||
audit.record_search(query=query.strip(), source=source, org_id=org_id, user_email=email, **exp.log)
|
||||
if page.lexical_empty:
|
||||
audit.record_search_miss(query=query.strip(), source=source)
|
||||
return page
|
||||
+17
-61
@@ -54,6 +54,7 @@ from mcp.shared.exceptions import MCPError
|
||||
from mcp.types import AudioContent, CallToolResult, METHOD_NOT_FOUND, TextContent, ToolAnnotations
|
||||
|
||||
from . import analytics, audit, hints
|
||||
from .application import catalog_search as search_app
|
||||
from .application import search_experiment
|
||||
from .domain.catalog import store as catalog_store
|
||||
from .config import PUBLIC_HOST_ALIASES, get_settings
|
||||
@@ -715,69 +716,30 @@ async def _search_identity(ctx: Context | None) -> tuple[str | None, int | None,
|
||||
async def _catalog_search_impl(
|
||||
query: str, limit: int = 8, *, ctx: Context | None = None, surface: _SurfacePolicy
|
||||
) -> SearchOut:
|
||||
"""The use case is `application.catalog_search`: the shipped ranker's page, the discovery
|
||||
experiment's say over it, and the records. This layer resolves who is asking (the hub's lists
|
||||
and the experiment's log both need it) and shapes the rows for an agent."""
|
||||
cat = catalog_store.load()
|
||||
limit = max(1, min(limit, 25))
|
||||
# Score, then let the evidence break the ties. Token scoring produces ties by the dozen — every
|
||||
# one of the 24 "ad library" matches scores 6 — so with a default limit of 8 the rows an agent
|
||||
# actually sees were decided by file order. That handed back seven tikhub rows (one of them
|
||||
# uncallable) and hid the cheapest endpoint with a perfect measured record.
|
||||
# The band is widened only so routed groups can collapse below without starving the page; with
|
||||
# steering off there is no collapsing, so the original band is the right one.
|
||||
_steering = str(get_settings().routed_discovery).strip().lower() not in ("off", "0", "false", "no")
|
||||
ranked, total, tie_truncated = catalog_store.rank_band(
|
||||
query, cat, min(100, limit * 4) if _steering else limit)
|
||||
stats = await _observed_stats([ep["id"] for ep, _ in ranked])
|
||||
ranked = catalog_store.rerank(ranked, stats, cat)
|
||||
# Listed hub tools ride in by score, no boost (docs/hub-listing-decisions.md, decision 2).
|
||||
from .application import hub as hub_app
|
||||
from .infra.db import session_maker
|
||||
# While a list limits the hub, only a caller in it (by team or by email) sees hub rows.
|
||||
hub_slug = hub_email = None
|
||||
if get_settings().hub_enabled and get_settings().hub_limited:
|
||||
token = _bearer(ctx) if ctx is not None else ""
|
||||
if token:
|
||||
hub_slug, hub_email = await _hub_reader(token)
|
||||
async with session_maker() as _s:
|
||||
hub_ranked, hub_stats = await hub_app.search_listed(_s, query, cat, org_slug=hub_slug, email=hub_email)
|
||||
if hub_ranked:
|
||||
stats = {**stats, **hub_stats}
|
||||
ranked = catalog_store.merge_by_score(ranked, hub_ranked)
|
||||
total += len(hub_ranked)
|
||||
results = []
|
||||
# Same order the HTTP route serves: a capability with a ROUTED row shows the parent first and
|
||||
# its children right under it (catalog_store.group_routed), so an agent sees "let treg choose"
|
||||
# before the specific providers.
|
||||
grouped = catalog_store.group_routed(
|
||||
[{"ep": ep, "score": score, "capability": ep.get("capability"), "kind": ep.get("kind")} for ep, score in ranked],
|
||||
max_children=catalog_store.MAX_ROUTED_CHILDREN)
|
||||
hidden = {r["ep"]["id"]: r["children_hidden"] for r in grouped if r.get("children_hidden")}
|
||||
ranked = [(r["ep"], r["score"]) for r in grouped][:limit]
|
||||
baseline_page = ranked
|
||||
if query.strip() and search_experiment.mode() != "off":
|
||||
# The discovery experiment (application.search_experiment): a relevance judge over a wider
|
||||
# recall, compared with the page above on what the caller does next. `shadow` serves this
|
||||
# page unchanged and only logs; `interleave` may serve a merge. Whatever the judge does,
|
||||
# `ranked` stays a page — an abstaining judge leaves the baseline in place.
|
||||
async def _finish(rows):
|
||||
st = await _observed_stats([ep["id"] for ep, _ in rows])
|
||||
rows = catalog_store.rerank(rows, st, cat)
|
||||
g = catalog_store.group_routed(
|
||||
[{"ep": ep, "score": sc, "capability": ep.get("capability"), "kind": ep.get("kind")}
|
||||
for ep, sc in rows], max_children=catalog_store.MAX_ROUTED_CHILDREN)
|
||||
return [(r["ep"], r["score"]) for r in g][:limit], st
|
||||
key, org_id, email = await _search_identity(ctx)
|
||||
exp = await search_experiment.run(query, cat, baseline=baseline_page, baseline_total=total,
|
||||
limit=limit, caller=key, finish=_finish)
|
||||
stats = {**exp.stats, **stats}
|
||||
ranked = exp.shown
|
||||
audit.record_search(query=query.strip(), source=surface.event_source, org_id=org_id,
|
||||
user_email=email, **exp.log)
|
||||
analytics.capture(key or "anonymous", "catalog_search_judged", {
|
||||
"source": surface.event_source, "mode": exp.log["mode"], "arm": exp.arm,
|
||||
"baseline_total": total, "baseline_empty": not baseline_page,
|
||||
"differs": exp.log["differs"], "judge_ms": exp.log["judge_ms"],
|
||||
"judge_error": exp.log["judge_error"], "judge_tokens_in": exp.log["judge_tokens_in"],
|
||||
page = await search_app.search(
|
||||
query, limit, cat=cat, source=surface.event_source,
|
||||
caller=search_app.Caller(hub_slug=hub_slug, hub_email=hub_email),
|
||||
observed=_observed_stats, identify=lambda: _search_identity(ctx))
|
||||
ranked, stats, hidden, total, _steering = page.rows, page.stats, page.hidden, page.total, page.steering
|
||||
if page.arm is not None:
|
||||
analytics.capture(page.caller_key or "anonymous", "catalog_search_judged", {
|
||||
"source": surface.event_source, "mode": page.log["mode"], "arm": page.arm,
|
||||
"baseline_total": total, "baseline_empty": page.lexical_empty,
|
||||
"differs": page.log["differs"], "judge_ms": page.log["judge_ms"],
|
||||
"judge_error": page.log["judge_error"], "judge_tokens_in": page.log["judge_tokens_in"],
|
||||
"shown": len(ranked)})
|
||||
results = []
|
||||
for ep, score in ranked:
|
||||
obs = stats.get(ep["id"]) or {}
|
||||
cost = (ep.get("cost") if ep.get("kind") == "hub" else cat.cost_view(ep.get("cost"), ep.get("provider"))) or {}
|
||||
@@ -811,18 +773,12 @@ async def _catalog_search_impl(
|
||||
"samples": obs.get("samples") or 0,
|
||||
})
|
||||
out = {"query": query, "count": len(results), "total_matches": total, "results": results}
|
||||
if tie_truncated:
|
||||
if page.tie_truncated:
|
||||
# No silent caps: past this many equally-scoring rows the evidence sort never saw the rest,
|
||||
# so the tail is ordered by nothing in particular and must not read as a ranked answer.
|
||||
out["ranking_note"] = (f"{query!r} matches too broadly to rank on measured reliability past "
|
||||
f"the first {catalog_store.RERANK_BAND} equally-scoring rows — "
|
||||
f"add a word to narrow it")
|
||||
if not baseline_page and query.strip():
|
||||
# Same miss log as GET /catalog/search (see models.SearchMiss) — this tool reads the catalog
|
||||
# in-process, so the HTTP route's logging never sees an MCP agent's empty search. Judged by
|
||||
# the LEXICAL page: the miss log measures the shipped ranker's coverage, and a judged page
|
||||
# that found something is the experiment's result, not a reason to stop recording the gap.
|
||||
audit.record_search_miss(query=query.strip(), source=surface.event_source)
|
||||
if not results:
|
||||
# the zero-result answer carries the rows that JUST missed the gate and which words they
|
||||
# missed — the caller is an LLM, and told exactly what to drop it re-queries correctly
|
||||
|
||||
Reference in New Issue
Block a user