mirror of
https://github.com/superdesigndev/treg.git
synced 2026-10-02 03:24:35 +08:00
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.
71 lines
2.8 KiB
Python
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
|