Files
treg/tests/test_embed.py
SToneX 7fc6ee3f83 feat(find): semantic channel
v2's recall gains its second channel. infra/embed.py is an
OpenAI-compatible /embeddings client that never raises and caches a
query's vector in-process by model and folded text. application/
find_index.py builds one card matrix per catalog in the background on
the first v2 find: vectors are read from the archive's object store
under find-vectors/<model>/<card sha256>, only missing cards are
embedded, and those are written back. Without a store they live in the
process; without the API the channel stays off and the build is retried
later. The judged event and the searchlog row record the query's
embedding time and error.

The object store gains named objects for these vectors only, each body
carrying its own card hash and size. numpy joins the server extra for the
matrix product. Settings find_embed_api_key (falls back to treg's
OpenRouter key on the OpenRouter URL), find_embed_model, find_embed_url,
find_embed_timeout_s.

The bench builds the vectors before scoring when a key is set, and its
judge cache key now includes the job question's wording.
2026-09-30 18:19:49 +08:00

71 lines
2.8 KiB
Python

"""The embeddings client (infra.embed): answers in input order, abstains instead of raising, and caches
a query's vector by its folded text."""
from __future__ import annotations
import json
import httpx
from treg.infra import embed as embed_infra
KW = dict(api_key="k", model="voyageai/voyage-4-lite", url="https://embed.test/v1/embeddings", timeout_s=1.0)
def _transport(handler):
return httpx.MockTransport(handler)
def _answer(vectors, tokens=7):
return httpx.Response(200, json={"data": [{"index": i, "embedding": v} for i, v in reversed(list(enumerate(vectors)))],
"usage": {"prompt_tokens": tokens}})
async def test_vectors_come_back_in_input_order_with_usage():
seen = []
def handler(request):
seen.append(json.loads(request.content))
assert request.headers["authorization"] == "Bearer k"
return _answer([[1.0, 0.0], [0.0, 1.0]])
e = await embed_infra.embed(["a", "b"], transport=_transport(handler), **KW)
assert e.vectors == [[1.0, 0.0], [0.0, 1.0]] and e.tokens == 7 and e.error is None
assert seen == [{"model": "voyageai/voyage-4-lite", "input": ["a", "b"]}]
async def test_it_abstains_on_timeout_http_error_bad_body_and_a_wrong_size():
def timeout(request):
raise httpx.ReadTimeout("slow", request=request)
cases = [
(timeout, None, "timeout"),
(lambda r: httpx.Response(502, text="bad gateway"), None, "http_502"),
(lambda r: httpx.Response(200, json={"data": []}), None, "bad_body"),
(lambda r: httpx.Response(200, text="not json"), None, "bad_body"),
(lambda r: _answer([[1.0, 0.0, 0.0]]), 2, "dim_3"),
]
for handler, dim, error in cases:
e = await embed_infra.embed(["a"], transport=_transport(handler), dim=dim, **KW)
assert e.vectors is None and e.error == error
async def test_a_query_is_embedded_once_per_folded_text():
embed_infra.clear_cache()
calls = []
def handler(request):
calls.append(json.loads(request.content)["input"])
return _answer([[0.6, 0.8]])
first = await embed_infra.embed_query("Café Reviews", transport=_transport(handler), **KW)
again = await embed_infra.embed_query("cafe reviews", transport=_transport(handler), **KW)
assert first.vector == again.vector == [0.6, 0.8] and again.cached and not first.cached
assert calls == [["cafe reviews"]]
# a failed answer is not cached
await embed_infra.embed_query("other", transport=_transport(lambda r: httpx.Response(500)), **KW)
await embed_infra.embed_query("other", transport=_transport(handler), **KW)
assert calls[-1] == ["other"]
# another model is another vector
await embed_infra.embed_query("cafe reviews", transport=_transport(handler), **{**KW, "model": "m2"})
assert len(calls) == 3