homelab-codex-ws/services/kb-query/tests/test_search.py
oskar e7625cd322 feat(kb): aktywny fallback embeddingów SOLARIA→PIHA dla kb-query (faza 4 Krok 2)
Ostatni krok fazy 4 KB (plan §2 Decyzja 2, §5): kb-query przestaje być martwe
przez ~16 h/dobę, gdy SOLARIA (GPU) śpi — zapytania embeduje wtedy lokalna
Ollama CPU na PIHA (wolniej: ~790 ms+ vs ~207 ms na GPU, ale działa).

Nowy serwis services/ollama-piha (GitOps, owner_node: piha):
- ollama/ollama:latest (arm64 natywnie), OLLAMA_KEEP_ALIVE=0 — model zwalnia
  RAM natychmiast po każdym wywołaniu (spike, nie rezydent; PIHA dzieli 8 GB z HA)
- bind wyłącznie 127.0.0.1 + LAN_BIND_IP (192.168.31.5), nigdy 0.0.0.0/Tailscale
- named volume ollama_piha_models (NVMe data-root) zamiast bind-mounta — obraz
  biega jako root w kontenerze i bind łamałby wzorzec uid PIHA (oskar=1004,
  kontenery uid 1000, setgid pi)
- override hosts/piha/runtime/ollama-piha: mem_limit 2560m (wartość startowa
  z planu, do potwierdzenia kalibracją na żywo), świadomie bez mem_reservation
- pull bge-m3 to jawny, ręczny krok deployu (README) — obraz nie ma modeli

kb-query — maszyna stanów fallbacku (app/embed_router.py):
- health-check SOLARII (GET /api/tags, timeout 1.5 s) z cache 30 s — zero
  sondowania per request; po powrocie SOLARII ruch wraca na GPU w ≤30 s
- primary up → embed na SOLARII z twardym timeoutem 3 s; błąd W TRAKCIE
  zapytania = jednorazowe przełączenie (krok 3b planu): status down na 30 s
  i TO SAMO zapytanie leci na fallback — user nie widzi błędu SOLARII
- primary down → embed prosto na ollama-piha (bez twardego timeoutu: CPU +
  zimny load modelu to legalnie pojedyncze sekundy)
- 503 tylko gdy oba backendy padłe (lub fallback nieskonfigurowany)
- inwariant modelu, druga połowa: każdy backend weryfikowany raz, leniwie przy
  pierwszym użyciu, że /api/tags zawiera EMBED_MODEL (bge-m3 — ta sama wartość
  co startowy check przeciw document_chunk.model/document_summary.embedding_model);
  niezgodność = ERROR log + 500, nigdy ciche liczenie dystansów między
  różnymi przestrzeniami embeddingów; leniwie, bo śpiąca SOLARIA nie może
  blokować startu serwisu
- odpowiedź /search: nowe pole embed_backend ("solaria"|"piha") + sol_status
  wg realnego świata routera (UI już renderuje down jako "offline (fallback
  embed)"); log INFO backend=... elapsed_ms=... per zapytanie
- /healthz: sol_status przez cache routera (spójny widok z routingiem) +
  fallback_status (żywa, tania sonda /api/tags)

Konfiguracja spójnie przez env (compose + env.example + service.yaml + README):
EMBED_PRIMARY_URL (zastępuje OLLAMA_URL), EMBED_FALLBACK_URL (pusty = brak
fallbacku, zachowanie sprzed kroku 2), EMBED_{PRIMARY,FALLBACK}_NAME,
EMBED_HEALTH_TTL_S/EMBED_HEALTH_TIMEOUT_S/EMBED_PRIMARY_TIMEOUT_S.

Testy: 39 pass (14 nowych w test_embed_router.py: cache TTL, failover w trakcie
zapytania, powrót po TTL, oba padłe, mismatch modelu na primary i fallbacku,
tag "bge-m3:latest" vs "bge-m3"); docker build + smoke (importy + uvicorn do
guardu KB_DSN) OK; compose config OK dla obu stacków.

Deploy (Oskar, na PIHA z mastera po merge):
  cd ~/homelab-codex-ws && git pull
  # 1. ollama-piha
  cp services/ollama-piha/env.example services/ollama-piha/.env
  docker compose -f services/ollama-piha/docker-compose.yml \
    -f hosts/piha/runtime/ollama-piha/docker-compose.override.yml \
    --env-file services/ollama-piha/.env up -d
  docker exec ollama-piha ollama pull bge-m3     # ręczny krok, obowiązkowy
  services/ollama-piha/healthcheck.sh
  # 2. kb-query (dopisać fallback do istniejącego .env)
  echo 'EMBED_FALLBACK_URL=http://192.168.31.5:11434' >> services/kb-query/.env
  docker compose -f services/kb-query/docker-compose.yml \
    -f hosts/piha/runtime/kb-query/docker-compose.override.yml up -d --build
  services/kb-query/healthcheck.sh
  # (deploy-node.sh też podniesie oba serwisy z hosts/piha/services.yaml,
  #  ale pull bge-m3 i .env pozostają ręczne)
Weryfikacja: testy A/B/C w services/kb-query/README.md (backend=solaria przy
SOLARII online; backend=piha przy symulacji offline; powrót na GPU w ≤30 s).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-29 19:01:29 +02:00

249 lines
11 KiB
Python

"""Unit tests for /search's core logic (app.search.run_search) -- no real DB, no real Ollama.
Same mocking style as packages/kb-retrieval/tests/test_retrieval.py, extended with an
`envelope` table fixture for the join app/db.py adds on top of kb_retrieval."""
from __future__ import annotations
import pathlib
import sys
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[1]))
from app.search import run_search # noqa: E402
class _FakeConn:
"""summaries: [(envelope_id, dist), ...] -- cascade stage-1 pre-filter. chunks_by_envelope:
envelope_id -> [(chunk_index, text, dist), ...]. envelopes: envelope_id -> {"source": ...,
"entities": [...]}. summary_texts: envelope_id -> {"summary": ..., "tags": [...]} -- the
document_summary row fetched for the result header (app/db.py fetch_summaries)."""
def __init__(
self, summaries=None, chunks_by_envelope=None, envelopes=None, summary_texts=None,
mail_chunks_by_source=None,
):
self._summaries = list(summaries or [])
self._chunks_by_envelope = chunks_by_envelope or {}
self._envelopes = envelopes or {}
self._summary_texts = summary_texts or {}
self._mail_chunks_by_source = mail_chunks_by_source or {}
async def fetch(self, query, *params):
if "FROM document_summary" in query and "= ANY" in query:
(envelope_ids, _model) = params
return [
{"envelope_id": eid, "summary": s["summary"], "tags": s["tags"]}
for eid, s in self._summary_texts.items()
if eid in envelope_ids
]
if "FROM document_summary" in query:
_, _model, limit = params
return [{"envelope_id": eid, "dist": dist} for eid, dist in self._summaries[:limit]]
if "JOIN envelope" in query: # hybrid's direct summaryless-source chunk scan
_, sources, limit = params
rows = [
{"envelope_id": eid, "chunk_index": idx, "text": text, "dist": dist}
for source in sources
for eid, idx, text, dist in self._mail_chunks_by_source.get(source, [])
]
rows.sort(key=lambda r: r["dist"])
return rows[:limit]
if "FROM document_chunk" in query and "= ANY" in query:
_, envelope_ids, limit = params
rows = [
{"envelope_id": eid, "chunk_index": idx, "text": text, "dist": dist}
for eid in envelope_ids
for idx, text, dist in self._chunks_by_envelope.get(eid, [])
]
rows.sort(key=lambda r: r["dist"])
return rows[:limit]
if "FROM document_chunk" in query: # flat path
_, limit = params
rows = [
{"envelope_id": eid, "chunk_index": idx, "text": text, "dist": dist}
for eid, chunk_list in self._chunks_by_envelope.items()
for idx, text, dist in chunk_list
]
rows.sort(key=lambda r: r["dist"])
return rows[:limit]
if "FROM envelope" in query:
(envelope_ids,) = params
return [
{"id": eid, "source": self._envelopes[eid]["source"], "entities": self._envelopes[eid]["entities"]}
for eid in envelope_ids
if eid in self._envelopes
]
raise AssertionError(f"unexpected query: {query}")
class _FakeSession:
"""run_search no longer talks HTTP itself -- embedding goes through the router
(below), so the session is just passed through untouched."""
class _FakeRouter:
"""Stands in for app.embed_router.EmbedRouter: returns a fixed embedding and the name
of the backend that 'served' it, mirroring EmbedRouter.embed's contract."""
def __init__(self, backend="solaria", primary_name="solaria"):
self._backend = backend
self.primary = type("B", (), {"name": primary_name})()
self.embed_calls: list[str] = []
async def embed(self, session, text):
self.embed_calls.append(text)
return [0.01] * 1024, self._backend
class TestRunSearchHappyPath:
async def test_cascade_hit_joins_envelope_and_shapes_paperless_link(self):
conn = _FakeConn(
summaries=[("paperless:119", 0.1)],
chunks_by_envelope={"paperless:119": [(2, "hit text", 0.34)]},
envelopes={"paperless:119": {"source": "paperless", "entities": []}},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "polisa PZU", "cascade", "claude-haiku-4-5"
)
assert result["query"] == "polisa PZU"
assert result["mode"] == "cascade"
assert result["sol_status"] == "up"
assert result["embed_backend"] == "solaria"
assert len(result["results"]) == 1
hit = result["results"][0]
assert hit["envelope_id"] == "paperless:119"
assert hit["source"] == "paperless"
assert hit["dist"] == 0.34
assert hit["chunk_index"] == 2
assert hit["text"] == "hit text"
assert hit["link"] == "https://paper.kapala.org/documents/119/details"
async def test_flat_mode_skips_cascade_stage1(self):
conn = _FakeConn(
chunks_by_envelope={"paperless:1": [(0, "a", 0.2)]},
envelopes={"paperless:1": {"source": "paperless", "entities": []}},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "q", "flat", "claude-haiku-4-5"
)
assert result["mode"] == "flat"
assert len(result["results"]) == 1
async def test_summary_attached_when_document_summary_row_exists(self):
conn = _FakeConn(
summaries=[("paperless:119", 0.1)],
chunks_by_envelope={"paperless:119": [(2, "hit text", 0.34)]},
envelopes={"paperless:119": {"source": "paperless", "entities": []}},
summary_texts={"paperless:119": {"summary": "Polisa OC 2024", "tags": ["ubezpieczenia"]}},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "polisa PZU", "cascade", "claude-haiku-4-5"
)
hit = result["results"][0]
assert hit["summary"] == "Polisa OC 2024"
assert hit["summary_tags"] == ["ubezpieczenia"]
async def test_summary_defaults_to_none_when_no_document_summary_row(self):
conn = _FakeConn(
summaries=[("paperless:119", 0.1)],
chunks_by_envelope={"paperless:119": [(2, "hit text", 0.34)]},
envelopes={"paperless:119": {"source": "paperless", "entities": []}},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "polisa PZU", "cascade", "claude-haiku-4-5"
)
hit = result["results"][0]
assert hit["summary"] is None
assert hit["summary_tags"] == []
async def test_hybrid_mode_merges_cascade_and_mail_branches(self):
conn = _FakeConn(
summaries=[("paperless:1", 0.3)],
chunks_by_envelope={"paperless:1": [(0, "doc text", 0.3)]},
envelopes={
"paperless:1": {"source": "paperless", "entities": []},
"<msgid@example.com>": {
"source": "gmail",
"entities": [{"type": "headers", "from": None, "subject": "s", "date_raw": "d"}],
},
},
mail_chunks_by_source={"gmail": [("<msgid@example.com>", 0, "mail text", 0.2)]},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "q", "hybrid", "claude-haiku-4-5"
)
assert result["mode"] == "hybrid"
envelope_ids = [r["envelope_id"] for r in result["results"]]
assert envelope_ids == ["<msgid@example.com>", "paperless:1"]
async def test_gmail_hit_carries_header_metadata_not_a_link(self):
conn = _FakeConn(
summaries=[("<msgid@example.com>", 0.1)],
chunks_by_envelope={"<msgid@example.com>": [(0, "body text", 0.4)]},
envelopes={
"<msgid@example.com>": {
"source": "gmail",
"entities": [
{"type": "headers", "from": {"name": "A", "address": "a@b.com"}, "subject": "s", "date_raw": "d"}
],
}
},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "q", "cascade", "claude-haiku-4-5"
)
hit = result["results"][0]
assert hit["source"] == "gmail"
assert hit["subject"] == "s"
assert hit["link"] is None
assert hit["mail_ui_url"] is None
class TestRunSearchEmbedBackendMarking:
async def test_fallback_embed_marks_backend_piha_and_sol_status_down(self):
# Task spec: the response must say WHICH backend embedded the query (quality
# debugging), and sol_status must reflect the router's world view -- the frontend
# renders "down" as "offline (fallback embed)".
conn = _FakeConn(
summaries=[("paperless:119", 0.1)],
chunks_by_envelope={"paperless:119": [(2, "hit text", 0.34)]},
envelopes={"paperless:119": {"source": "paperless", "entities": []}},
)
result = await run_search(
conn, _FakeSession(), _FakeRouter(backend="piha"), "q", "cascade", "claude-haiku-4-5"
)
assert result["embed_backend"] == "piha"
assert result["sol_status"] == "down"
assert len(result["results"]) == 1
class TestRunSearchNoGoodResults:
async def test_results_above_no_answer_threshold_are_still_returned_unfiltered(self):
# Plan §7: the 0.55 "no answer" colour threshold is a frontend concern -- the API
# must not silently drop/hide a poor match, only report its true dist so the caller
# (UI or eval harness) can apply that policy itself.
conn = _FakeConn(
summaries=[("paperless:1", 0.6)],
chunks_by_envelope={"paperless:1": [(0, "unrelated text", 0.62)]},
envelopes={"paperless:1": {"source": "paperless", "entities": []}},
)
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "unrelated query", "cascade", "claude-haiku-4-5"
)
assert len(result["results"]) == 1
assert result["results"][0]["dist"] == 0.62
async def test_no_summaries_yields_empty_results_not_an_error(self):
conn = _FakeConn(summaries=[], chunks_by_envelope={}, envelopes={})
session = _FakeSession()
result = await run_search(
conn, session, _FakeRouter(), "nothing matches", "cascade", "claude-haiku-4-5"
)
assert result["results"] == []