feat(kb-query): active embed fallback SOLARIA→PIHA (module 5 phase 4, plan §2/§5)
Last missing core piece of KB phase 4: kb-query no longer hard-fails /search when Ollama@SOLARIA is unreachable. app/fallback.py implements the plan's circuit-breaker exactly (30s cached health probe, 3s hard embed timeout on SOLARIA, one-shot same-request switch to a new local ollama-piha@PIHA container on timeout/error). sol_status in /healthz and /search now reflects the real breaker state instead of a hardcoded "up". New services/ollama-piha (bge-m3, OLLAMA_KEEP_ALIVE=0, arm64/no-GPU) is the local fallback leg. Live calibration on PIHA (2026-07-27, normal load): embed latency 4.2-5.2s, RAM peak ~983MiB against a 2.5GiB ceiling -- both inside the plan's go-bar, so the fallback is enabled by default rather than gated behind a flag. Calibration also surfaced and disabled (not removed) a previously-undocumented orphaned native ollama.service on PIHA that had been conflicting with the container's port. The embed-model invariant (query embedding == document_chunk.model) still enforces once at startup, since both fallback legs share one EMBED_MODEL constant by construction; a redundant per-request DB check was deliberately skipped and the invariant is instead proven structurally by test. retrieval_eval.py gains --transport http (plan §2 decision 6/§9), previously unimplemented. Verified live: HTTP transport is bit-identical to direct transport against the same live SOLARIA (0 mismatches), and a live sol-down simulation (kb-query's own OLLAMA_URL pointed at a dead address, no other Ollama consumer touched) shows the PIHA fallback answering with the same hit@3 gate outcome and dist within ~3e-4 of the SOLARIA baseline. Zero changes to DB schema or kb_retrieval's retrieval logic -- only the embed + health layer, per task constraints. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
parent
50c4b2024a
commit
3d4ee3818d
189
docs/sessions/2026-07-27-kb-f4-fallback.md
Normal file
189
docs/sessions/2026-07-27-kb-f4-fallback.md
Normal file
|
|
@ -0,0 +1,189 @@
|
|||
# Sesja 2026-07-27 — KB faza 4: fallback embed SOLARIA→PIHA (krok 3, ostatni element rdzenia)
|
||||
|
||||
**Zakres**: `docs/kb/modules/05-faza4-plan.md` §2 decyzja 2 / §5 — aktywny fallback
|
||||
embedu, ostatni brakujący element rdzenia fazy 4 (frontend i ingress LIVE od
|
||||
2026-07-22/23, `docs/sessions/2026-07-23-kb-f4-ingress.md`). Zero zmian w schemacie
|
||||
DB, zero zmian w `kb_retrieval`'s retrieval logice — wyłącznie warstwa embed + health.
|
||||
|
||||
Praca wykonana w task worktree (`task/kb-f4-fallback`, `.claude/skills/worktree-aware`).
|
||||
Zgodnie z ustaleniem na starcie sesji (patrz "Ustalenia proceduralne" niżej): kod
|
||||
napisany i przetestowany lokalnie w worktree, produkcyjne kroki (kalibracja, deploy,
|
||||
live-test) wykonane po jawnej zgodzie operatora, z osobnym potwierdzeniem przed
|
||||
każdym kolejnym krokiem dotykającym PIHA/SOLARIĘ.
|
||||
|
||||
## Ustalenia proceduralne
|
||||
|
||||
Zadanie wprost wymagało kroków produkcyjnych (kalibracja RAM/latencji na żywym PIHA,
|
||||
symulacja sol-down dotykająca SOLARII, deploy, push) — sprzeczne z ogólną dyscypliną
|
||||
`worktree-aware` ("nigdy nie uruchamiaj deployów/healthchecków przeciw produkcji z
|
||||
worktree"). Zamiast rozstrzygać to samodzielnie, zapytano operatora:
|
||||
1. Recon read-only (bez zmian stanu) — zgoda bez pytania.
|
||||
2. Właściwe kroki produkcyjne (kalibracja, deploy, live-test, push) — operator
|
||||
potwierdził jawnie ("Yes, proceed with all of it") po zobaczeniu pełnego zakresu.
|
||||
|
||||
## 1. Kod (warstwa embed + health, zero zmian retrievalu/DB)
|
||||
|
||||
- **`packages/kb_retrieval/embed.py`**: `embed_chunk` dostał opcjonalny `timeout_s`
|
||||
(domyślnie `None`, zero zmiany zachowania istniejących wołań) — potrzebny do
|
||||
twardego 3 s timeoutu na nodze SOLARIA bez zmiany zachowania nogi PIHA.
|
||||
- **`services/kb-query/app/fallback.py`** (nowy): `SolCircuitBreaker` (cache 30 s,
|
||||
zegar wstrzykiwalny do testów) + `resolve_sol_status` (probe `/api/tags`, 500 ms) +
|
||||
`embed_with_fallback` (SOLARIA z twardym 3 s timeoutem → jednorazowe przełączenie na
|
||||
PIHA **w tym samym requeście** przy timeout/błędzie → PIHA bez dodatkowego
|
||||
timeoutu). Dokładnie maszyna stanów z planu §2 decyzja 2.
|
||||
- **`app/search.py`**: `run_search` liczy embedding raz przez `embed_with_fallback`,
|
||||
potem woła `flat_retrieve`/`cascade_retrieve`/`hybrid_retrieve` (niskopoziomowe
|
||||
funkcje `kb_retrieval`, biorą gotowy wektor) zamiast `flat_query`/`cascade_query`/
|
||||
`hybrid_query` (które embedują same) — dzięki temu decyzja fallbacku żyje wyłącznie
|
||||
w warstwie HTTP kb-query, zero zmiany w `kb_retrieval`. `sol_status` w odpowiedzi to
|
||||
teraz realny wynik, nie zahardkodowane `"up"`.
|
||||
- **`app/main.py`**: `/healthz` i `/search` dzielą jeden `SolCircuitBreaker`
|
||||
(`app.state.sol_breaker`) — oba endpointy zawsze zgadzają się co do aktualnego
|
||||
stanu. Nowy env `OLLAMA_PIHA_URL` (domyślnie `http://localhost:11434` — celowo
|
||||
"inertny" placeholder, fail-closed, dopóki operator nie ustawi realnego adresu).
|
||||
- **Inwariant modelu**: **nie dodano** drugiego, per-request sprawdzenia w DB —
|
||||
`EMBED_MODEL` to jedna stała wątkowana przez obie nogi `embed_with_fallback`,
|
||||
więc startowy check (`app/startup.py`, niezmieniony) pokrywa obie ścieżki z
|
||||
konstrukcji. Dodanie drugiego DB-checka chroniłoby przed scenariuszem, który nie
|
||||
może wystąpić (CLAUDE.md: nie dodawaj walidacji dla scenariuszy, które nie mogą się
|
||||
zdarzyć) — zamiast tego nowy test (`test_both_legs_use_identical_embed_model`)
|
||||
strukturalnie dowodzi, że obie nogi w tym samym requeście dostają identyczny
|
||||
`embed_model`.
|
||||
- **`jobs/documents-ingest/eval/retrieval_eval.py`**: dodano `--transport
|
||||
{direct,http}` + `--base-url` (plan §2 decyzja 6 / §9) — dotąd nieistniejące (tylko
|
||||
ręczny smoke-test, `docs/sessions/2026-07-23-kb-f4-ingress.md` follow-up). Tryb
|
||||
`http` woła trzy `GET /search` (flat/cascade/hybrid) na żywym kb-query zamiast
|
||||
embedować+odpytywać lokalnie; `envelope.source` do kryterium 4 bierze się z pola
|
||||
`source` w odpowiedzi JSON, nie z osobnego zapytania do DB. Nie da się swipe'ować
|
||||
N przez HTTP (kb-query serwuje jeden N per request) — tryb http raportuje tylko
|
||||
przy `--gate-n`.
|
||||
|
||||
## 2. Nowy serwis `services/ollama-piha`
|
||||
|
||||
Klon wzorca `services/ollama` (`owner_node: piha` zamiast `solaria`, bez rezerwacji
|
||||
GPU — PIHA to arm64 bez akceleracji), `OLLAMA_KEEP_ALIVE=0` (model ładowany tylko na
|
||||
czas requestu). `mem_limit: 2560m` (tentatywny wg planu, potwierdzony pomiarem —
|
||||
patrz §3). Wpisany do `hosts/piha/services.yaml` (`depends_on.local` kb-query →
|
||||
`[kb-postgres, ollama-piha]`, fallback nie jest twardą zależnością na starcie).
|
||||
|
||||
## 3. Znalezisko: osierocony natywny `ollama.service` na PIHA
|
||||
|
||||
Podczas pierwszej próby deployu `ollama-piha` (bind `127.0.0.1:11434`) — konflikt
|
||||
portu. Okazało się, że PIHA ma **natywny (nie-Docker) systemd `ollama.service`**
|
||||
(v0.6.1, `enabled`, działający od 2026-06-22, PATH env wskazujący na użytkownika
|
||||
`/home/pi/...`), o którym nic nie wiadomo w repo — plan §1.2 wprost zakładał "PIHA:
|
||||
brak Ollamy", co okazało się nieaktualne/błędne. To realna sprzeczność planu z
|
||||
rzeczywistością → STOP, pytanie do operatora zamiast cichej decyzji.
|
||||
|
||||
Weryfikacja przed jakąkolwiek akcją: `journalctl -u ollama --since "7 days ago"` —
|
||||
**tylko własne, właśnie wykonane** zapytania probe (`/api/version`, `/api/tags`),
|
||||
`total blobs: 0` od startu (nigdy nic nie pobrano). Operator potwierdził: martwy
|
||||
balast, `sudo systemctl disable --now ollama.service` (**disable, nie uninstall** —
|
||||
odwracalne). Port 11434 zwolniony, `ollama-piha` wystartował normalnie.
|
||||
|
||||
**Backlog**: PIHA host-level shadow — natywny `ollama.service` wyłączony
|
||||
2026-07-27; odinstalować binarkę/unit po ~2 tygodniach jeśli nic się nie posypie.
|
||||
|
||||
## 4. Kalibracja (plan §5, gate) — **werdykt: GO**
|
||||
|
||||
Zmierzone na żywym PIHA pod normalnym obciążeniem (kb-postgres, paperless, Immich,
|
||||
HA, Forgejo działające, nie okno nocnej ciszy), 3 kolejne wywołania `/api/embeddings`
|
||||
po `ollama pull bge-m3`:
|
||||
|
||||
| Wywołanie | Latencja |
|
||||
|---|---|
|
||||
| 1 (pierwsze, zimny start) | 5.25 s |
|
||||
| 2 | 4.41 s |
|
||||
| 3 | 4.16 s |
|
||||
|
||||
Brak przyspieszenia między wywołaniami — zgodnie z projektem (`OLLAMA_KEEP_ALIVE=0`
|
||||
zwalnia model po każdym requeście, `ollama ps` pokazuje zero rezydentnych modeli
|
||||
między wywołaniami).
|
||||
|
||||
RAM: baseline idle ~66 MiB, szczyt podczas burst ~983 MiB (`docker stats`, próbkowane
|
||||
co 0.3 s w trakcie 3 wywołań) — komfortowo w granicach ceilingu `2560m`. `free -h`
|
||||
systemowe: `available` nie spadło poniżej ~1.3 GiB w trakcie, osiadło na ~4.2 GiB po
|
||||
(dla porównania: przed startem eksperymentu `available` = 3.7 GiB).
|
||||
|
||||
**Werdykt**: oba kryteria planu spełnione (latencja pojedyncze sekundy, nie
|
||||
dziesiątki; RAM ze sporym zapasem) → **włączony jako domyślny fallback**, bez flagi
|
||||
`KB_QUERY_LOCAL_FALLBACK_ENABLED`.
|
||||
|
||||
## 5. Bramka jakościowa (plan §9)
|
||||
|
||||
Wszystko uruchomione z `~/kb/venv` na PIHA (istniejący venv z poprzednich sesji,
|
||||
`aiohttp`/`asyncpg`/`yaml` już obecne) przeciw żywej bazie + żywemu kb-query.
|
||||
|
||||
**HTTP-equivalence** (`--transport http` vs `--transport direct`, SOLARIA up, ten sam
|
||||
`--gate-n 10`): oba PASS, **0 rozbieżności** w `dist` na wszystkich zapytaniach
|
||||
(`flat_top1_dist`, `hybrid_top1_dist`, `cascade[10].top1_dist`) — identyczne bit w
|
||||
bit, jak wymagał plan (nie ±epsilon, bo to ten sam kod, HTTP to tylko opakowanie).
|
||||
|
||||
**Live sol-down fallback test**: symulacja przez `OLLAMA_URL=http://solaria:1`
|
||||
(zły port, zgodnie z rekomendacją planu — zero dotknięcia SOLARII/innych
|
||||
konsumentów Ollamy) w `.env` kb-query, restart kontenera. `/healthz` →
|
||||
`sol_status: "down"`. `/search` → 200, wyniki z PIHA, ~4.3 s (zgodnie z kalibracją).
|
||||
Pełna bramka `retrieval_eval.py --transport http` z SOLARIA-down: **PASS** —
|
||||
identyczny wzorzec hit@3 co na SOLARII, `dist` w granicach epsilon:
|
||||
|
||||
| Zapytanie | dist (SOLARIA) | dist (PIHA fallback) | Δ |
|
||||
|---|---|---|---|
|
||||
| 1 | 0.341780 | 0.341509 | 0.000271 |
|
||||
| 2 | 0.324808 | 0.324858 | 0.00005 |
|
||||
| 3 | 0.428898 | 0.429184 | 0.000286 |
|
||||
| 4 | 0.448199 | 0.447903 | 0.000296 |
|
||||
| 5 | 0.386901 | 0.386816 | 0.000085 |
|
||||
| N (negative control) | 0.598301 | 0.598017 | 0.000283 |
|
||||
| N2 (negative control borderline) | 0.529772 | 0.529530 | 0.000242 |
|
||||
|
||||
Maksymalna rozbieżność: **~3e-4** — rząd wielkości mniejszy niż oczekiwany przez plan
|
||||
(1e-3–1e-2), kolejność top-k identyczna, wynik bramki (`gate.passed`) identyczny.
|
||||
Kb-query przywrócony do normalnej konfiguracji po teście (`.env` z prawdziwym
|
||||
`OLLAMA_URL`, restart), `/healthz` z powrotem `sol_status: "up"`.
|
||||
|
||||
## 6. Deploy
|
||||
|
||||
Kod nie był jeszcze zmergowany do `master` (dyscyplina worktree: agent nigdy nie
|
||||
mergeuje/pushuje `master`) — deploy przez standardowy `deploy-node.sh`
|
||||
niedostępny bez mastera. Zamiast tego: `rsync` zmienionych plików
|
||||
(`packages/kb-retrieval`, `services/kb-query`, `services/ollama-piha`,
|
||||
`hosts/piha/runtime/ollama-piha`, `hosts/piha/services.yaml`,
|
||||
`jobs/documents-ingest/eval/retrieval_eval.py` + README) do żywego checkoutu
|
||||
`~/homelab-codex-ws` na PIHA (bez zmiany brancha — working tree pozostaje na
|
||||
`master` z niescommitowanym diffem 1:1 identycznym z tą gałęzią), potem
|
||||
standardowy `docker compose ... up -d --build` z tego miejsca. Efekt: realny,
|
||||
działający deploy, ale **repo na PIHA ma dziś dirty working tree** — wymaga domknięcia
|
||||
(patrz "Do zrobienia przez operatora" niżej).
|
||||
|
||||
Zweryfikowane: `kb-query` (healthy), `ollama-piha` (healthy, `bge-m3` w wolumenie),
|
||||
`curl https://kb.kapala.org/healthz` → `200 {"sol_status":"up"}`,
|
||||
`curl https://kb.kapala.org/search?q=test` → `200`.
|
||||
|
||||
## Stan na koniec sesji
|
||||
|
||||
| Element | Status |
|
||||
|---|---|
|
||||
| `packages/kb-retrieval` — `embed_chunk(timeout_s=...)` | ✅ kod + testy |
|
||||
| `services/kb-query/app/fallback.py` — maszyna stanów | ✅ kod + testy (38/38 kb-query, 25/25 kb-retrieval) |
|
||||
| `services/ollama-piha` — nowy serwis GitOps | ✅ zdefiniowany, ✅ LIVE na PIHA |
|
||||
| Natywny `ollama.service` na PIHA (osierocony) | ✅ wyłączony (nie odinstalowany) |
|
||||
| Kalibracja RAM/latencja | ✅ zmierzone — werdykt GO |
|
||||
| `retrieval_eval.py --transport http` | ✅ zaimplementowane, ✅ PASS na żywo |
|
||||
| Live sol-down fallback test | ✅ PASS, Δ~3e-4 |
|
||||
| Deploy kb-query + ollama-piha na PIHA | ✅ LIVE, working tree PIHA dirty (patrz niżej) |
|
||||
| Merge do `master` | ⛔ nie wykonany (dyscyplina worktree — operator) |
|
||||
|
||||
## Do zrobienia przez operatora
|
||||
|
||||
1. **Merge** `task/kb-f4-fallback` → `master` (`scripts/dev/agent.sh merge` albo
|
||||
ręcznie) — branch popchnięty do `origin/task/kb-f4-fallback` (patrz commit poniżej).
|
||||
2. Na PIHA: `cd ~/homelab-codex-ws && git status` będzie dirty (diff identyczny z tym
|
||||
commitem, bo już wdrożony ad-hoc przez `rsync` w tej sesji) — po mergu do mastera,
|
||||
`git checkout -- .` (working tree już ma dokładnie tę treść) albo zwyczajnie
|
||||
`git pull` po mergu powinien wylądować "already up to date"/no-op, bo pliki na
|
||||
dysku już są zgodne z tym co przyjdzie z mastera. **Zweryfikować** `git diff` jest
|
||||
puste po pull, nie zakładać.
|
||||
3. Backlog: natywny `ollama.service` na PIHA wyłączony `systemctl disable --now`
|
||||
2026-07-27 (§3 wyżej) — jeśli nic się nie posypie przez ~2 tygodnie, odinstalować
|
||||
binarkę/unit całkiem.
|
||||
4. OIDC dla kb-query nadal odłożone (decyzja z 2026-07-23) — nie w zakresie tej sesji.
|
||||
13
hosts/piha/runtime/ollama-piha/docker-compose.override.yml
Normal file
13
hosts/piha/runtime/ollama-piha/docker-compose.override.yml
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
# PIHA-specific overrides for ollama-piha (KB module 5, phase 4, plan §2 decision 2 / §5).
|
||||
#
|
||||
# mem_limit is a cgroup ceiling per the plan's own estimate (§5: "np. 2.5g, do potwierdzenia
|
||||
# pomiarem") -- confirmed live 2026-07-27 (services/ollama-piha/README.md "Calibration status"):
|
||||
# measured peak ~983 MiB during a burst, comfortably inside this 2560m ceiling with margin to
|
||||
# spare.
|
||||
#
|
||||
# Not oom_score_adj: -900 -- that's reserved for control-plane/agent processes that must never
|
||||
# be an OOM victim (CLAUDE.md). ollama-piha is a fallback-only leg: cgroup-restart-on-breach is
|
||||
# an acceptable failure mode here, same reasoning as kb-query's own override.
|
||||
services:
|
||||
ollama-piha:
|
||||
mem_limit: 2560m
|
||||
|
|
@ -111,13 +111,33 @@ services:
|
|||
exposure: private # LAN bind (LAN_BIND_IP); npm@PIHA vhost + OIDC is a later step
|
||||
offline_required: false
|
||||
depends_on:
|
||||
local: [kb-postgres]
|
||||
external: [ollama] # SOLARIA may be offline -> /search 502/503, /healthz stays ok
|
||||
local: [kb-postgres, ollama-piha] # ollama-piha is the fallback leg, not hard-required at boot
|
||||
external: [ollama] # SOLARIA may be offline -> active fallback to ollama-piha, plan §2/§5
|
||||
ports:
|
||||
- name: http
|
||||
container_port: 8080
|
||||
host_port: 8230
|
||||
protocol: tcp
|
||||
runtime:
|
||||
# .env (KB_DSN, LAN_BIND_IP) lives alongside the compose file; stateless, no data path
|
||||
# .env (KB_DSN, LAN_BIND_IP, OLLAMA_PIHA_URL) lives alongside the compose file; stateless, no data path
|
||||
config_path: services/kb-query
|
||||
|
||||
ollama-piha:
|
||||
role: kb-embed-fallback # module 5 phase 4 plan §2 decision 2 / §5: local CPU embed
|
||||
# fallback for kb-query when Ollama@SOLARIA is down/times out
|
||||
deployment_model: docker-compose
|
||||
exposure: private # LAN bind (LAN_BIND_IP); consumed only by kb-query@PIHA
|
||||
offline_required: true # PIHA-local, never depends on SOLARIA/VPS/Forgejo at runtime
|
||||
depends_on:
|
||||
local: []
|
||||
external: []
|
||||
ports:
|
||||
- name: http
|
||||
container_port: 11434
|
||||
host_port: 11434
|
||||
protocol: tcp
|
||||
runtime:
|
||||
# .env (LAN_BIND_IP) lives alongside the compose file
|
||||
config_path: services/ollama-piha
|
||||
# bge-m3 model files persist in the Docker named volume-equivalent bind
|
||||
# /opt/homelab/data/ollama-piha — see hosts/piha/runtime/ollama-piha override.
|
||||
|
|
|
|||
|
|
@ -517,6 +517,12 @@ python eval/retrieval_eval.py --dsn postgresql://kb:<pw>@piha:5433/kb \
|
|||
--ollama-url http://solaria:11434 --n-sweep 5,10,20
|
||||
```
|
||||
|
||||
`--transport http --base-url http://<kb-query-host>:8230` (module 5 phase 4 plan §2 decision 6 /
|
||||
§9) calls a live `kb-query`'s `/search` instead of embedding+querying locally — no `--dsn`
|
||||
needed, `--n-sweep` is ignored (kb-query serves one server-side default N per request). Gate
|
||||
criterion: `dist` must be **identical** to the same run with `--transport direct` against the
|
||||
same live SOLARIA (same DB, same retrieval code — HTTP is only a wrapper).
|
||||
|
||||
**Result (2026-07-17, live run)**: PASS at N=10, k=5 — see plan §6.3 for the full table,
|
||||
the N-sweep calibration (N=5 is the measured safety floor; the plan's N=10 default carries a
|
||||
2× margin), and the cost/improvement analysis. `cascade_query` (N=10, k=5,
|
||||
|
|
|
|||
|
|
@ -23,9 +23,23 @@ mocked test). Runs every query in `queries.yaml`'s `queries:` list through flat
|
|||
writes nothing. Query embeddings go through Ollama on localhost/SOLARIA (bge-m3), same as
|
||||
`chunk_embed.py`/`summarize.py`.
|
||||
|
||||
`--transport {direct,http}` (module 5 phase 4 plan §2 decision 6 / §9, added alongside the
|
||||
fallback task): `direct` (default) is everything above, unchanged. `http` instead calls
|
||||
`GET {base_url}/search?q=...&mode=flat|cascade|hybrid` on a live `kb-query` and reshapes its
|
||||
JSON `results` into the same `{"chunks": [...]}` shape the direct-mode functions return, so
|
||||
`summarize_query_result`/`evaluate_gate` run identically either way. The gate criterion (plan
|
||||
§9): `dist` for `http` must be **identical** to `direct` against the same live SOLARIA -- same
|
||||
DB, same retrieval code, HTTP is only a wrapper, so any difference is a serialization/handler
|
||||
bug, never expected numerical drift. `http` mode cannot sweep `N` (kb-query serves one
|
||||
server-side default per request, plan §4) -- it reports only at `--gate-n`, and needs no `--dsn`
|
||||
(kb-query already owns the DB connection; `envelope.source` for criterion 4 comes straight from
|
||||
each result's `source` field instead of a separate DB lookup).
|
||||
|
||||
Usage:
|
||||
python retrieval_eval.py --dsn postgresql://kb:<pw>@piha:5433/kb \\
|
||||
--ollama-url http://solaria:11434 --n-sweep 5,10,20
|
||||
python retrieval_eval.py --transport http --base-url http://192.168.31.5:8230 \\
|
||||
--gate-n 10
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
|
|
@ -156,6 +170,44 @@ async def run_query_all_tracks(
|
|||
return {"query": query, "flat": flat, "cascades": cascades, "hybrid": hybrid}
|
||||
|
||||
|
||||
async def call_search_http(
|
||||
session: aiohttp.ClientSession, base_url: str, query_text: str, mode: str
|
||||
) -> list[dict]:
|
||||
"""One `GET {base_url}/search?q=...&mode=...` call -> its `results` list. Each result
|
||||
already carries `envelope_id`/`dist`/`source` -- exactly the fields `top1_dist`/`hit_at_3`/
|
||||
`mail_hit_at_3` need, no DB lookup required on this side."""
|
||||
async with session.get(
|
||||
f"{base_url}/search", params={"q": query_text, "mode": mode}
|
||||
) as resp:
|
||||
resp.raise_for_status()
|
||||
data = await resp.json()
|
||||
return data["results"]
|
||||
|
||||
|
||||
async def run_query_all_tracks_http(
|
||||
session: aiohttp.ClientSession, base_url: str, query: dict, gate_n: int
|
||||
) -> dict:
|
||||
"""HTTP-transport equivalent of `run_query_all_tracks` -- three `/search` calls (one per
|
||||
mode) instead of embedding+querying locally. `stage1_summaries` isn't part of the HTTP
|
||||
response shape (plan §4) so it's reported empty; nothing in `evaluate_gate` reads it."""
|
||||
flat_chunks = await call_search_http(session, base_url, query["text"], "flat")
|
||||
cascade_chunks = await call_search_http(session, base_url, query["text"], "cascade")
|
||||
hybrid_chunks = await call_search_http(session, base_url, query["text"], "hybrid")
|
||||
return {
|
||||
"query": query,
|
||||
"flat": {"chunks": flat_chunks},
|
||||
"cascades": {gate_n: {"chunks": cascade_chunks, "stage1_summaries": []}},
|
||||
"hybrid": {"chunks": hybrid_chunks},
|
||||
}
|
||||
|
||||
|
||||
def envelope_sources_from_results(*chunk_lists: list[dict]) -> dict[str, str]:
|
||||
"""http transport has no DB to `fetch_envelope_sources` from -- each `/search` result
|
||||
already carries its envelope's `source`, so build the same envelope_id -> source mapping
|
||||
straight from the response bodies already fetched for this query."""
|
||||
return {c["envelope_id"]: c["source"] for chunks in chunk_lists for c in chunks}
|
||||
|
||||
|
||||
def summarize_query_result(result: dict, envelope_sources: Optional[dict[str, str]] = None) -> dict:
|
||||
query = result["query"]
|
||||
expected = query["expected_envelope"]
|
||||
|
|
@ -342,7 +394,34 @@ def print_report(
|
|||
"hybrid = 1 embed (shared) + cascade's queries + 1 extra SQL query (mail branch).")
|
||||
|
||||
|
||||
async def main_async(args: argparse.Namespace) -> dict:
|
||||
async def main_async_http(args: argparse.Namespace) -> tuple[list[dict], list[dict], list[int]]:
|
||||
"""`--transport http` path -- no DB connection, three `/search` calls per query. Returns
|
||||
only at `--gate-n` (see module docstring: kb-query serves one server-side N per request)."""
|
||||
queries = load_queries(Path(args.queries))
|
||||
mail_queries = load_mail_queries(Path(args.queries))
|
||||
n_values = [args.gate_n]
|
||||
|
||||
async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=60)) as session:
|
||||
results = []
|
||||
envelope_sources: dict[str, str] = {}
|
||||
for query in queries:
|
||||
result = await run_query_all_tracks_http(session, args.base_url, query, args.gate_n)
|
||||
envelope_sources.update(envelope_sources_from_results(
|
||||
result["flat"]["chunks"], result["cascades"][args.gate_n]["chunks"], result["hybrid"]["chunks"],
|
||||
))
|
||||
results.append(summarize_query_result(result))
|
||||
|
||||
mail_results_raw = []
|
||||
for query in mail_queries:
|
||||
result = await run_query_all_tracks_http(session, args.base_url, query, args.gate_n)
|
||||
envelope_sources.update(envelope_sources_from_results(result["hybrid"]["chunks"]))
|
||||
mail_results_raw.append(result)
|
||||
mail_results = [summarize_query_result(r, envelope_sources) for r in mail_results_raw]
|
||||
|
||||
return results, mail_results, n_values
|
||||
|
||||
|
||||
async def main_async_direct(args: argparse.Namespace) -> tuple[list[dict], list[dict], list[int]]:
|
||||
queries = load_queries(Path(args.queries))
|
||||
mail_queries = load_mail_queries(Path(args.queries))
|
||||
n_values = [int(n) for n in args.n_sweep.split(",")]
|
||||
|
|
@ -380,6 +459,21 @@ async def main_async(args: argparse.Namespace) -> dict:
|
|||
finally:
|
||||
await conn.close()
|
||||
|
||||
return results, mail_results, n_values
|
||||
|
||||
|
||||
async def main_async(args: argparse.Namespace) -> dict:
|
||||
if args.transport == "http":
|
||||
if args.n_sweep != "5,10,20": # the argparse default -- operator didn't ask for a sweep
|
||||
print(
|
||||
"note: --transport http ignores --n-sweep (kb-query serves a single "
|
||||
f"server-side default N per request); reporting only --gate-n={args.gate_n}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
results, mail_results, n_values = await main_async_http(args)
|
||||
else:
|
||||
results, mail_results, n_values = await main_async_direct(args)
|
||||
|
||||
gate_result = evaluate_gate(results, gate_n=args.gate_n, mail_rows=mail_results)
|
||||
print_report(results, n_values, gate_result, mail_rows=mail_results)
|
||||
|
||||
|
|
@ -395,19 +489,26 @@ async def main_async(args: argparse.Namespace) -> dict:
|
|||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument("--dsn", default=os.environ.get("KB_DSN"), help="asyncpg DSN for kb-postgres (or KB_DSN env var)")
|
||||
parser.add_argument("--transport", choices=["direct", "http"], default="direct",
|
||||
help="direct = query DB+Ollama locally (default); http = call a live kb-query's /search")
|
||||
parser.add_argument("--base-url", default=None, help="kb-query base URL, required for --transport http (e.g. http://192.168.31.5:8230)")
|
||||
parser.add_argument("--dsn", default=os.environ.get("KB_DSN"), help="asyncpg DSN for kb-postgres (or KB_DSN env var); required for --transport direct")
|
||||
parser.add_argument("--ollama-url", default=os.environ.get("OLLAMA_URL", "http://localhost:11434"))
|
||||
parser.add_argument("--embed-model", default=DEFAULT_EMBED_MODEL)
|
||||
parser.add_argument("--summary-model", default=DEFAULT_SUMMARY_MODEL,
|
||||
help="document_summary.model to pre-filter on (plan §2 D3 resolution)")
|
||||
parser.add_argument("--k", type=int, default=DEFAULT_K)
|
||||
parser.add_argument("--gate-n", type=int, default=DEFAULT_N, help="N used for the PASS/FAIL verdict")
|
||||
parser.add_argument("--n-sweep", default="5,10,20", help="comma-separated N values to report (diagnostic)")
|
||||
parser.add_argument("--n-sweep", default="5,10,20", help="comma-separated N values to report (diagnostic; ignored by --transport http)")
|
||||
parser.add_argument("--queries", default=str(DEFAULT_QUERIES_PATH))
|
||||
parser.add_argument("--json-out", default=None, help="optional path to dump full results as JSON")
|
||||
args = parser.parse_args()
|
||||
|
||||
if not args.dsn:
|
||||
if args.transport == "http":
|
||||
if not args.base_url:
|
||||
print("error: --transport http requires --base-url", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
elif not args.dsn:
|
||||
print("error: pass --dsn or set KB_DSN", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
|
|
|
|||
|
|
@ -40,14 +40,20 @@ def _vector_literal(embedding: list[float]) -> str:
|
|||
|
||||
|
||||
async def embed_chunk(
|
||||
session: aiohttp.ClientSession, base_url: str, model: str, text: str
|
||||
session: aiohttp.ClientSession, base_url: str, model: str, text: str, timeout_s: float | None = None
|
||||
) -> tuple[list[float], float]:
|
||||
"""POST /api/embeddings on Ollama for one chunk. Returns (embedding, elapsed_seconds).
|
||||
No built-in timeout/retry -- inherits whatever `aiohttp.ClientSession(timeout=...)` the
|
||||
caller constructed; `raise_for_status()` propagates `aiohttp.ClientError` when Ollama is
|
||||
unreachable."""
|
||||
No built-in timeout/retry by default -- inherits whatever `aiohttp.ClientSession(timeout=...)`
|
||||
the caller constructed; `raise_for_status()` propagates `aiohttp.ClientError` when Ollama is
|
||||
unreachable. `timeout_s`, when given, overrides the session default for this call only
|
||||
(raises `asyncio.TimeoutError` on expiry) -- used by kb-query's fallback state machine
|
||||
(plan §2 decision 2 / §5) to bound the SOLARIA leg independently of the shared
|
||||
`aiohttp.ClientSession`'s own timeout."""
|
||||
t0 = time.monotonic()
|
||||
async with session.post(f"{base_url}/api/embeddings", json={"model": model, "prompt": text}) as resp:
|
||||
kwargs = {"json": {"model": model, "prompt": text}}
|
||||
if timeout_s is not None:
|
||||
kwargs["timeout"] = aiohttp.ClientTimeout(total=timeout_s)
|
||||
async with session.post(f"{base_url}/api/embeddings", **kwargs) as resp:
|
||||
resp.raise_for_status()
|
||||
data = await resp.json()
|
||||
elapsed = time.monotonic() - t0
|
||||
|
|
|
|||
|
|
@ -64,13 +64,32 @@ container, no node build step (plan §2 decision 4): `app/templates/index.html`
|
|||
- Stopka pokazuje `sol_status` dyskretnie (odświeżane z `/healthz` przy
|
||||
starcie strony i po każdym wyszukiwaniu).
|
||||
|
||||
## Embed path (current step)
|
||||
## Embed path — active fallback (plan §2 decision 2, §5)
|
||||
|
||||
Calls Ollama on SOLARIA directly per request — no cache, no circuit breaker,
|
||||
no local-PIHA fallback yet (that state machine, plan §2 decision 2/§5, is a
|
||||
separate later step). If SOLARIA is unreachable, `/search` returns **503**;
|
||||
`/healthz` still answers (`sol_status: "down"`), same tolerance pattern as
|
||||
`llm-gateway`.
|
||||
`app/fallback.py` holds one process-global `SolCircuitBreaker`:
|
||||
|
||||
1. Cached `sol_status` (30s TTL) is used as-is when fresh — no network call.
|
||||
2. On expiry, probe `GET {OLLAMA_URL}/api/tags` (500ms timeout); cache the
|
||||
result (`up`/`down`) for another 30s.
|
||||
3. `up` → embed on SOLARIA with a hard 3s timeout.
|
||||
- Success → done, `sol_status: "up"`.
|
||||
- Timeout/error on the **real** embed call (not just the probe) → flip
|
||||
the breaker to `down` immediately and fall through to step 4 **in the
|
||||
same request** — the caller never sees an error for this, only the
|
||||
first unlucky request in a 30s window pays one extra timeout.
|
||||
4. `down` → embed locally against `OLLAMA_PIHA_URL` (`ollama-piha`@PIHA,
|
||||
`bge-m3`, same model constant as SOLARIA — see the invariant note below).
|
||||
|
||||
If **both** legs fail (SOLARIA down and PIHA unreachable/not deployed),
|
||||
`/search` returns **503**; `/healthz` still answers (`sol_status: "down"`),
|
||||
same tolerance pattern as `llm-gateway`. `/healthz` shares the same cached
|
||||
breaker as `/search`, so both report the same view of the world.
|
||||
|
||||
`OLLAMA_PIHA_URL` defaults to `http://localhost:11434`, which is inert
|
||||
inside this container (nothing listens there) until you point it at the
|
||||
real `ollama-piha`@PIHA address — see `services/ollama-piha/README.md` for
|
||||
that container's deploy status and the RAM/latency calibration gate that
|
||||
decides whether it's safe to rely on as a default fallback.
|
||||
|
||||
## Startup invariant (hard-fail)
|
||||
|
||||
|
|
@ -83,10 +102,17 @@ actually indexed — see `app/startup.py` for why the check reads
|
|||
`document_summary.embedding_model` and not `.model` (the latter is the LLM
|
||||
that *wrote* the summary, e.g. `claude-haiku-4-5`, not the embedder).
|
||||
|
||||
**Covers both fallback legs.** `EMBED_MODEL` is a single constant threaded
|
||||
through `app/fallback.py`'s `embed_with_fallback` and used identically for
|
||||
the SOLARIA and PIHA embed calls — there is no per-request or per-leg model
|
||||
choice, so this one startup check already covers both paths. See
|
||||
`app/fallback.py`'s module docstring for why a second, redundant per-request
|
||||
DB check was deliberately not added.
|
||||
|
||||
## Configuration
|
||||
|
||||
`.env` — **gitignored**, copy from `env.example`. Required: `LAN_BIND_IP`,
|
||||
`KB_DSN`. Optional: `OLLAMA_URL`, `EMBED_MODEL`, `SUMMARY_MODEL`.
|
||||
`KB_DSN`. Optional: `OLLAMA_URL`, `OLLAMA_PIHA_URL`, `EMBED_MODEL`, `SUMMARY_MODEL`.
|
||||
|
||||
## Deploy (PIHA)
|
||||
|
||||
|
|
@ -150,5 +176,4 @@ existed.
|
|||
|
||||
## Out of scope for this step
|
||||
|
||||
- Local-PIHA embed fallback / circuit breaker (plan §2 decision 2, §5).
|
||||
- OIDC login (see above) — separate session, needs `authlib` + Forgejo OAuth2 app.
|
||||
|
|
|
|||
113
services/kb-query/app/fallback.py
Normal file
113
services/kb-query/app/fallback.py
Normal file
|
|
@ -0,0 +1,113 @@
|
|||
"""Active embed fallback state machine -- module 5 phase 4 plan §2 decision 2 / §5
|
||||
(docs/kb/modules/05-faza4-plan.md). One process-global circuit breaker (`SolCircuitBreaker`,
|
||||
held in `app.state`, not a module global -- keeps tests isolated) shared by every `/search`
|
||||
request.
|
||||
|
||||
State machine (plan §2 decision 2 table):
|
||||
1. Cached status still fresh (< CACHE_TTL_S old) -> use it, no network call at all.
|
||||
2. Cache expired -> probe `GET {solaria}/api/tags` with a short timeout (HEALTH_PROBE_TIMEOUT_S);
|
||||
cache the result (up or down) for CACHE_TTL_S. Same TTL for both outcomes -- plan explicitly
|
||||
defers a separate down-backoff until flapping is observed.
|
||||
3. status == "up" -> embed on SOLARIA with a hard per-call timeout (EMBED_TIMEOUT_S), distinct
|
||||
from the probe timeout: a slow-but-technically-reachable SOLARIA must not hang the user's
|
||||
query past this bound.
|
||||
3a. Success -> return, sol_status="up".
|
||||
3b. Timeout/error *on the real embed call* (not just the probe) -> flip the breaker to "down"
|
||||
immediately (one-shot switch) and fall through to step 4 **in the same request** -- the
|
||||
user never sees an error for this, only the first unlucky hit in a 30s window pays one
|
||||
extra EMBED_TIMEOUT_S before falling through.
|
||||
4. status == "down" -> embed locally via `piha_url`, sol_status="down".
|
||||
|
||||
Both branches call `embed_chunk(..., model=embed_model)` with the exact same `embed_model`
|
||||
argument -- there is no separate "PIHA model" and no per-request model choice (plan §2: "nie ma
|
||||
dziś ... per-request wyboru modelu"), so the hard model invariant (app/startup.py, enforced once
|
||||
at process startup against document_chunk/document_summary) already covers both paths by
|
||||
construction. Adding a second, redundant per-request DB check here would guard against a
|
||||
scenario that cannot occur while `embed_model` stays a single constant threaded through one
|
||||
function -- instead, `tests/test_fallback.py` proves the invariant holds by asserting both the
|
||||
failed SOLARIA attempt and the successful PIHA attempt in the same request carry an identical
|
||||
`model` field.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Callable, Optional
|
||||
|
||||
import aiohttp
|
||||
|
||||
from kb_retrieval.embed import check_ollama_health, embed_chunk
|
||||
|
||||
CACHE_TTL_S = 30.0
|
||||
HEALTH_PROBE_TIMEOUT_S = 0.5
|
||||
EMBED_TIMEOUT_S = 3.0
|
||||
|
||||
|
||||
@dataclass
|
||||
class SolCircuitBreaker:
|
||||
"""Process-global cached `sol_status` ("up"/"down"), one instance per running app (held in
|
||||
`app.state.sol_breaker`, created fresh in `lifespan` -- never a module-level singleton, so
|
||||
tests get a clean breaker per instance). `clock` is injectable (default `time.monotonic`,
|
||||
never wall-clock/datetime) purely to test TTL expiry without real sleeps."""
|
||||
|
||||
cache_ttl_s: float = CACHE_TTL_S
|
||||
clock: Callable[[], float] = field(default=time.monotonic)
|
||||
_status: Optional[str] = field(default=None, init=False)
|
||||
_checked_at: Optional[float] = field(default=None, init=False)
|
||||
|
||||
@property
|
||||
def status(self) -> Optional[str]:
|
||||
"""Cached status if still within `cache_ttl_s`, else `None` (caller must (re)probe)."""
|
||||
if self._status is None or self._checked_at is None:
|
||||
return None
|
||||
if self.clock() - self._checked_at >= self.cache_ttl_s:
|
||||
return None
|
||||
return self._status
|
||||
|
||||
def set(self, status: str) -> None:
|
||||
self._status = status
|
||||
self._checked_at = self.clock()
|
||||
|
||||
|
||||
async def resolve_sol_status(
|
||||
breaker: SolCircuitBreaker, session: aiohttp.ClientSession, solaria_url: str
|
||||
) -> str:
|
||||
"""Cache-or-probe: returns "up"/"down". Shared by `/search` (embed_with_fallback) and
|
||||
`/healthz` (plan: "sol_status w stopce UI ... podepnij realny stan") so both endpoints
|
||||
agree on the same cached view instead of running independent probes."""
|
||||
cached = breaker.status
|
||||
if cached is not None:
|
||||
return cached
|
||||
up = await check_ollama_health(session, solaria_url, HEALTH_PROBE_TIMEOUT_S)
|
||||
status = "up" if up else "down"
|
||||
breaker.set(status)
|
||||
return status
|
||||
|
||||
|
||||
async def embed_with_fallback(
|
||||
breaker: SolCircuitBreaker,
|
||||
session: aiohttp.ClientSession,
|
||||
solaria_url: str,
|
||||
piha_url: str,
|
||||
embed_model: str,
|
||||
query_text: str,
|
||||
) -> tuple[list[float], str]:
|
||||
"""query_text -> (embedding, sol_status_used). `sol_status_used` is "up" (SOLARIA answered
|
||||
the real embed call) or "down" (PIHA answered -- either the breaker already cached "down",
|
||||
or SOLARIA just failed mid-request and this call transparently fell through, plan §2 step
|
||||
3b's one-shot switch)."""
|
||||
status = await resolve_sol_status(breaker, session, solaria_url)
|
||||
|
||||
if status == "up":
|
||||
try:
|
||||
embedding, _elapsed = await embed_chunk(
|
||||
session, solaria_url, embed_model, query_text, timeout_s=EMBED_TIMEOUT_S
|
||||
)
|
||||
return embedding, "up"
|
||||
except (aiohttp.ClientError, TimeoutError):
|
||||
# One-shot switch (plan §2 step 3b): don't fail this request, fall through to PIHA
|
||||
# below and remember "down" for the rest of the cache window.
|
||||
breaker.set("down")
|
||||
|
||||
embedding, _elapsed = await embed_chunk(session, piha_url, embed_model, query_text)
|
||||
return embedding, "down"
|
||||
|
|
@ -3,10 +3,11 @@ entry point to the KB. Wraps `kb_retrieval.cascade_query`/`flat_query` (module 5
|
|||
already gated PASS -- docs/sessions/2026-07-21.md) in FastAPI. This is a search API, not chat:
|
||||
no answer synthesis, no LLM call over the results (that is phase 5, out of scope here).
|
||||
|
||||
Embed path is deliberately simple for this step: calls Ollama on SOLARIA directly, no
|
||||
cache/circuit-breaker/local-PIHA-fallback (plan §2 decision 2, §5) -- that state machine is a
|
||||
later, separate step. A failed embed call (SOLARIA unreachable) surfaces as 503 to the caller
|
||||
rather than a bare 500.
|
||||
Embed path (plan §2 decision 2, §5): `app/fallback.py`'s `SolCircuitBreaker` health-checks
|
||||
Ollama@SOLARIA (cached ~30s, ~500ms probe), embeds there with a hard ~3s timeout when up, and
|
||||
falls through to Ollama@PIHA (`OLLAMA_PIHA_URL`, local fallback container `ollama-piha`) in the
|
||||
same request on timeout/error -- the user only ever pays one extra timeout, once per cache
|
||||
window. A failure on *both* legs still surfaces as 503 to the caller rather than a bare 500.
|
||||
|
||||
`GET /` (Krok 4, plan §7) serves the search UI from this same FastAPI process -- one image, one
|
||||
container (plan §2 decision 4): a Jinja2 shell + a static vanilla-JS file, no node build step.
|
||||
|
|
@ -28,18 +29,20 @@ from fastapi.responses import HTMLResponse
|
|||
from fastapi.staticfiles import StaticFiles
|
||||
from fastapi.templating import Jinja2Templates
|
||||
|
||||
from kb_retrieval.embed import check_ollama_health
|
||||
|
||||
from app.db import create_pool
|
||||
from app.fallback import SolCircuitBreaker, resolve_sol_status
|
||||
from app.search import run_search
|
||||
from app.startup import validate_embed_model
|
||||
|
||||
BASE_DIR = pathlib.Path(__file__).resolve().parent
|
||||
KB_DSN = os.environ.get("KB_DSN")
|
||||
OLLAMA_URL = os.environ.get("OLLAMA_URL", "http://solaria:11434")
|
||||
# Placeholder until ollama-piha is deployed (plan §5 calibration gate) -- nothing listens on
|
||||
# localhost:11434 inside this container, so this fails closed (connection refused) exactly like
|
||||
# SOLARIA-down does today, never a false success.
|
||||
OLLAMA_PIHA_URL = os.environ.get("OLLAMA_PIHA_URL", "http://localhost:11434")
|
||||
EMBED_MODEL = os.environ.get("EMBED_MODEL", "bge-m3")
|
||||
SUMMARY_MODEL = os.environ.get("SUMMARY_MODEL", "claude-haiku-4-5")
|
||||
OLLAMA_HEALTH_TIMEOUT_S = 3.0
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
|
|
@ -50,11 +53,15 @@ async def lifespan(app: FastAPI):
|
|||
pool = await create_pool(KB_DSN)
|
||||
async with pool.acquire() as conn:
|
||||
# Hard invariant (plan §2 decision 2): refuse to start rather than silently serve
|
||||
# queries against a mismatched embedding space.
|
||||
# queries against a mismatched embedding space. embed_model is a single constant used
|
||||
# identically by both the SOLARIA and PIHA embed legs (app/fallback.py) -- this one
|
||||
# check covers both paths, see app/fallback.py's module docstring for why a second,
|
||||
# per-request DB check would be redundant.
|
||||
await validate_embed_model(conn, EMBED_MODEL)
|
||||
|
||||
app.state.pool = pool
|
||||
app.state.http = aiohttp.ClientSession()
|
||||
app.state.sol_breaker = SolCircuitBreaker()
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
|
|
@ -74,8 +81,10 @@ async def index(request: Request):
|
|||
|
||||
@app.get("/healthz")
|
||||
async def healthz() -> dict:
|
||||
sol_up = await check_ollama_health(app.state.http, OLLAMA_URL, OLLAMA_HEALTH_TIMEOUT_S)
|
||||
return {"status": "ok", "sol_status": "up" if sol_up else "down"}
|
||||
# Shares the same cached breaker /search uses (app/fallback.py) -- healthz and search
|
||||
# always agree on the current sol_status instead of running independent probes.
|
||||
sol_status = await resolve_sol_status(app.state.sol_breaker, app.state.http, OLLAMA_URL)
|
||||
return {"status": "ok", "sol_status": sol_status}
|
||||
|
||||
|
||||
@app.get("/search")
|
||||
|
|
@ -86,7 +95,8 @@ async def search(
|
|||
try:
|
||||
async with app.state.pool.acquire() as conn:
|
||||
return await run_search(
|
||||
conn, app.state.http, OLLAMA_URL, q, mode, EMBED_MODEL, SUMMARY_MODEL
|
||||
conn, app.state.http, app.state.sol_breaker, OLLAMA_URL, OLLAMA_PIHA_URL,
|
||||
q, mode, EMBED_MODEL, SUMMARY_MODEL,
|
||||
)
|
||||
except aiohttp.ClientError as exc:
|
||||
raise HTTPException(status_code=503, detail=f"embed backend unavailable: {exc}") from exc
|
||||
|
|
|
|||
|
|
@ -10,41 +10,51 @@ concern (Krok 4, out of this step's scope), never applied server-side.
|
|||
for the frontend's per-envelope result header (plan §7) -- `None`/`[]` when the envelope has no
|
||||
summary for `summary_model` yet. Purely additive: does not change any field already covered by
|
||||
the phase-4 gate's HTTP-equivalence check (plan §9).
|
||||
|
||||
Embed fallback (plan §2 decision 2 / §5, `app/fallback.py`): the query embedding is computed
|
||||
once via `embed_with_fallback` (SOLARIA, or PIHA if SOLARIA is down/times out), then handed to
|
||||
`kb_retrieval`'s `*_retrieve` functions as a plain vector literal -- `flat_query`/`cascade_query`/
|
||||
`hybrid_query` (which embed *and* retrieve in one call) are deliberately bypassed here so the
|
||||
fallback decision lives entirely in this HTTP layer, per this task's constraint of zero changes
|
||||
to `kb_retrieval`'s retrieval logic. `sol_status` in the response is the real outcome of that
|
||||
call, not a hardcoded "up".
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import aiohttp
|
||||
import asyncpg
|
||||
|
||||
from kb_retrieval.retrieval import cascade_query, flat_query, hybrid_query
|
||||
from kb_retrieval.embed import _vector_literal
|
||||
from kb_retrieval.retrieval import cascade_retrieve, flat_retrieve, hybrid_retrieve
|
||||
|
||||
from app.db import fetch_envelopes, fetch_summaries
|
||||
from app.fallback import SolCircuitBreaker, embed_with_fallback
|
||||
from app.links import build_result
|
||||
|
||||
|
||||
async def run_search(
|
||||
conn: asyncpg.Connection,
|
||||
session: aiohttp.ClientSession,
|
||||
ollama_url: str,
|
||||
breaker: SolCircuitBreaker,
|
||||
solaria_url: str,
|
||||
piha_url: str,
|
||||
query_text: str,
|
||||
mode: str,
|
||||
embed_model: str,
|
||||
summary_model: str,
|
||||
) -> dict:
|
||||
if mode == "flat":
|
||||
retrieval = await flat_query(conn, session, ollama_url, query_text, embed_model=embed_model)
|
||||
elif mode == "hybrid":
|
||||
retrieval = await hybrid_query(
|
||||
conn, session, ollama_url, query_text,
|
||||
summary_model=summary_model, embed_model=embed_model,
|
||||
)
|
||||
else:
|
||||
retrieval = await cascade_query(
|
||||
conn, session, ollama_url, query_text,
|
||||
summary_model=summary_model, embed_model=embed_model,
|
||||
embedding, sol_status = await embed_with_fallback(
|
||||
breaker, session, solaria_url, piha_url, embed_model, query_text
|
||||
)
|
||||
query_vector = _vector_literal(embedding)
|
||||
|
||||
if mode == "flat":
|
||||
chunks = await flat_retrieve(conn, query_vector) # flat_retrieve returns a plain list
|
||||
elif mode == "hybrid":
|
||||
chunks = (await hybrid_retrieve(conn, query_vector, summary_model=summary_model))["chunks"]
|
||||
else:
|
||||
chunks = (await cascade_retrieve(conn, query_vector, summary_model=summary_model))["chunks"]
|
||||
|
||||
chunks = retrieval["chunks"]
|
||||
envelope_ids = sorted({c["envelope_id"] for c in chunks})
|
||||
envelopes = await fetch_envelopes(conn, envelope_ids)
|
||||
summaries = await fetch_summaries(conn, envelope_ids, summary_model)
|
||||
|
|
@ -60,6 +70,6 @@ async def run_search(
|
|||
return {
|
||||
"query": query_text,
|
||||
"mode": mode,
|
||||
"sol_status": "up", # reaching this point means the embed call above succeeded
|
||||
"sol_status": sol_status,
|
||||
"results": results,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ services:
|
|||
environment:
|
||||
- KB_DSN=${KB_DSN}
|
||||
- OLLAMA_URL=${OLLAMA_URL:-http://solaria:11434}
|
||||
- OLLAMA_PIHA_URL=${OLLAMA_PIHA_URL:-http://localhost:11434}
|
||||
- EMBED_MODEL=${EMBED_MODEL:-bge-m3}
|
||||
- SUMMARY_MODEL=${SUMMARY_MODEL:-claude-haiku-4-5}
|
||||
# python:3.12-slim has no curl/wget; same in-container check pattern as llm-gateway.
|
||||
|
|
|
|||
|
|
@ -15,6 +15,13 @@ KB_DSN=postgresql://kb:CHANGE-ME@192.168.31.5:5433/kb
|
|||
# llm-gateway's OLLAMA_URL). Optional: defaults to this value if unset.
|
||||
# OLLAMA_URL=http://solaria:11434
|
||||
|
||||
# Local fallback Ollama (ollama-piha@PIHA, plan §2 decision 2 / §5) — used only when SOLARIA
|
||||
# is unreachable/times out. kb-query runs in its own Docker network (separate compose project
|
||||
# from ollama-piha), so this must be PIHA's LAN IP + published port, not "localhost" — same
|
||||
# reasoning as KB_DSN above. Optional: defaults to http://localhost:11434, which fails closed
|
||||
# (nothing listens there in this container) until you set the real value below.
|
||||
# OLLAMA_PIHA_URL=http://192.168.31.5:11434
|
||||
|
||||
# Embedding model kb-query enforces as a startup invariant (plan §2 decision
|
||||
# 2) against document_chunk.model / document_summary.embedding_model.
|
||||
# Optional: defaults to bge-m3.
|
||||
|
|
|
|||
|
|
@ -29,5 +29,6 @@ service:
|
|||
- LAN_BIND_IP # required — compose port-bind interpolation
|
||||
- KB_DSN # required — asyncpg DSN for kb-postgres@PIHA
|
||||
- OLLAMA_URL # optional — defaults to http://solaria:11434
|
||||
- EMBED_MODEL # optional — defaults to bge-m3; startup invariant vs document_chunk/document_summary
|
||||
- OLLAMA_PIHA_URL # optional — local fallback (ollama-piha@PIHA), defaults to http://localhost:11434 (fails closed until set)
|
||||
- EMBED_MODEL # optional — defaults to bge-m3; startup invariant vs document_chunk/document_summary; same model used on both the SOLARIA and PIHA embed legs
|
||||
- SUMMARY_MODEL # optional — defaults to claude-haiku-4-5; cascade_query's stage-1 model
|
||||
|
|
|
|||
221
services/kb-query/tests/test_fallback.py
Normal file
221
services/kb-query/tests/test_fallback.py
Normal file
|
|
@ -0,0 +1,221 @@
|
|||
"""Unit tests for the embed fallback state machine (app/fallback.py) -- module 5 phase 4 plan
|
||||
§2 decision 2 / §5. No real HTTP, no real Ollama -- same mocking style as
|
||||
packages/kb-retrieval/tests/test_embed.py."""
|
||||
from __future__ import annotations
|
||||
|
||||
import pathlib
|
||||
import sys
|
||||
|
||||
import aiohttp
|
||||
import pytest
|
||||
|
||||
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[1]))
|
||||
|
||||
from app.fallback import ( # noqa: E402
|
||||
EMBED_TIMEOUT_S,
|
||||
SolCircuitBreaker,
|
||||
embed_with_fallback,
|
||||
resolve_sol_status,
|
||||
)
|
||||
|
||||
SOLARIA_URL = "http://solaria:11434"
|
||||
PIHA_URL = "http://piha:11434"
|
||||
|
||||
|
||||
class _FakeClock:
|
||||
def __init__(self, start: float = 0.0):
|
||||
self.t = start
|
||||
|
||||
def __call__(self) -> float:
|
||||
return self.t
|
||||
|
||||
def advance(self, dt: float) -> None:
|
||||
self.t += dt
|
||||
|
||||
|
||||
class _FakeGetResp:
|
||||
def __init__(self, status: int):
|
||||
self.status = status
|
||||
|
||||
async def __aenter__(self):
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *exc):
|
||||
return False
|
||||
|
||||
|
||||
class _FakePostResp:
|
||||
def __init__(self, embedding=None):
|
||||
self._embedding = embedding if embedding is not None else [0.01] * 1024
|
||||
|
||||
async def __aenter__(self):
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *exc):
|
||||
return False
|
||||
|
||||
def raise_for_status(self):
|
||||
pass
|
||||
|
||||
async def json(self):
|
||||
return {"embedding": self._embedding}
|
||||
|
||||
|
||||
class _FakeFallbackSession:
|
||||
"""Routes GET -> health probe, POST -> embed. `solaria_get_status=None` simulates the probe
|
||||
itself being unreachable (raises); `*_post_behavior` in {"ok", "timeout", "error"}."""
|
||||
|
||||
def __init__(self, solaria_get_status=200, solaria_post_behavior="ok", piha_post_behavior="ok"):
|
||||
self.solaria_get_status = solaria_get_status
|
||||
self.solaria_post_behavior = solaria_post_behavior
|
||||
self.piha_post_behavior = piha_post_behavior
|
||||
self.get_calls: list[str] = []
|
||||
self.post_calls: list[dict] = []
|
||||
|
||||
def get(self, url, timeout=None):
|
||||
self.get_calls.append(url)
|
||||
if self.solaria_get_status is None:
|
||||
raise aiohttp.ClientConnectionError("refused")
|
||||
return _FakeGetResp(self.solaria_get_status)
|
||||
|
||||
def post(self, url, json, timeout=None):
|
||||
self.post_calls.append({"url": url, "json": json, "timeout": timeout})
|
||||
behavior = self.solaria_post_behavior if url.startswith(SOLARIA_URL) else self.piha_post_behavior
|
||||
if behavior == "timeout":
|
||||
raise TimeoutError()
|
||||
if behavior == "error":
|
||||
raise aiohttp.ClientConnectionError("refused")
|
||||
return _FakePostResp()
|
||||
|
||||
|
||||
class TestSolCircuitBreaker:
|
||||
def test_status_none_when_never_set(self):
|
||||
assert SolCircuitBreaker().status is None
|
||||
|
||||
def test_status_returns_cached_value_within_ttl(self):
|
||||
clock = _FakeClock()
|
||||
breaker = SolCircuitBreaker(cache_ttl_s=30, clock=clock)
|
||||
breaker.set("up")
|
||||
clock.advance(29)
|
||||
assert breaker.status == "up"
|
||||
|
||||
def test_status_expires_exactly_at_ttl(self):
|
||||
clock = _FakeClock()
|
||||
breaker = SolCircuitBreaker(cache_ttl_s=30, clock=clock)
|
||||
breaker.set("down")
|
||||
clock.advance(30)
|
||||
assert breaker.status is None
|
||||
|
||||
|
||||
class TestResolveSolStatus:
|
||||
async def test_uses_fresh_cache_without_probing(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("up")
|
||||
session = _FakeFallbackSession()
|
||||
status = await resolve_sol_status(breaker, session, SOLARIA_URL)
|
||||
assert status == "up"
|
||||
assert session.get_calls == []
|
||||
|
||||
async def test_probes_and_caches_up(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
session = _FakeFallbackSession(solaria_get_status=200)
|
||||
status = await resolve_sol_status(breaker, session, SOLARIA_URL)
|
||||
assert status == "up"
|
||||
assert breaker.status == "up"
|
||||
assert session.get_calls == [f"{SOLARIA_URL}/api/tags"]
|
||||
|
||||
async def test_probes_and_caches_down_on_unreachable(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
session = _FakeFallbackSession(solaria_get_status=None)
|
||||
status = await resolve_sol_status(breaker, session, SOLARIA_URL)
|
||||
assert status == "down"
|
||||
assert breaker.status == "down"
|
||||
|
||||
async def test_reprobes_once_ttl_expires(self):
|
||||
clock = _FakeClock()
|
||||
breaker = SolCircuitBreaker(cache_ttl_s=30, clock=clock)
|
||||
session = _FakeFallbackSession(solaria_get_status=200)
|
||||
await resolve_sol_status(breaker, session, SOLARIA_URL)
|
||||
clock.advance(30)
|
||||
await resolve_sol_status(breaker, session, SOLARIA_URL)
|
||||
assert len(session.get_calls) == 2
|
||||
|
||||
|
||||
class TestEmbedWithFallback:
|
||||
async def test_solaria_up_embeds_on_solaria(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
session = _FakeFallbackSession(solaria_get_status=200, solaria_post_behavior="ok")
|
||||
embedding, status = await embed_with_fallback(
|
||||
breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q"
|
||||
)
|
||||
assert status == "up"
|
||||
assert len(embedding) == 1024
|
||||
assert session.post_calls == [
|
||||
{"url": f"{SOLARIA_URL}/api/embeddings", "json": {"model": "bge-m3", "prompt": "q"},
|
||||
"timeout": aiohttp.ClientTimeout(total=EMBED_TIMEOUT_S)}
|
||||
]
|
||||
|
||||
async def test_cached_down_skips_probe_and_solaria_entirely(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("down")
|
||||
session = _FakeFallbackSession(piha_post_behavior="ok")
|
||||
embedding, status = await embed_with_fallback(
|
||||
breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q"
|
||||
)
|
||||
assert status == "down"
|
||||
assert len(embedding) == 1024
|
||||
assert session.get_calls == []
|
||||
assert session.post_calls == [
|
||||
{"url": f"{PIHA_URL}/api/embeddings", "json": {"model": "bge-m3", "prompt": "q"}, "timeout": None}
|
||||
]
|
||||
|
||||
async def test_solaria_timeout_mid_request_falls_through_to_piha_same_request(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("up") # cache says up; the real call below discovers it's actually stuck
|
||||
session = _FakeFallbackSession(solaria_post_behavior="timeout", piha_post_behavior="ok")
|
||||
embedding, status = await embed_with_fallback(
|
||||
breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q"
|
||||
)
|
||||
assert status == "down"
|
||||
assert len(embedding) == 1024
|
||||
assert breaker.status == "down" # one-shot switch persists for the rest of the cache window
|
||||
assert [c["url"] for c in session.post_calls] == [
|
||||
f"{SOLARIA_URL}/api/embeddings", f"{PIHA_URL}/api/embeddings",
|
||||
]
|
||||
|
||||
async def test_solaria_connection_error_mid_request_falls_through(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("up")
|
||||
session = _FakeFallbackSession(solaria_post_behavior="error", piha_post_behavior="ok")
|
||||
embedding, status = await embed_with_fallback(
|
||||
breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q"
|
||||
)
|
||||
assert status == "down"
|
||||
assert breaker.status == "down"
|
||||
|
||||
async def test_both_legs_failing_raises_to_caller(self):
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("up")
|
||||
session = _FakeFallbackSession(solaria_post_behavior="timeout", piha_post_behavior="error")
|
||||
with pytest.raises(aiohttp.ClientError):
|
||||
await embed_with_fallback(breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q")
|
||||
|
||||
async def test_both_legs_use_identical_embed_model(self):
|
||||
# Structural proof of the "no per-request DB check needed" reasoning (module docstring):
|
||||
# a single embed_model argument is threaded through both the failed SOLARIA attempt and
|
||||
# the successful PIHA attempt in the same request.
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("up")
|
||||
session = _FakeFallbackSession(solaria_post_behavior="timeout", piha_post_behavior="ok")
|
||||
await embed_with_fallback(breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q")
|
||||
models = {c["json"]["model"] for c in session.post_calls}
|
||||
assert models == {"bge-m3"}
|
||||
|
||||
async def test_piha_leg_has_no_hard_timeout_override(self):
|
||||
# Only the SOLARIA leg gets the interactive-request hard timeout (plan §2 step 3) -- the
|
||||
# PIHA leg is the fallback of last resort, no shorter budget to enforce beyond it.
|
||||
breaker = SolCircuitBreaker()
|
||||
breaker.set("down")
|
||||
session = _FakeFallbackSession(piha_post_behavior="ok")
|
||||
await embed_with_fallback(breaker, session, SOLARIA_URL, PIHA_URL, "bge-m3", "q")
|
||||
assert session.post_calls[0]["timeout"] is None
|
||||
|
|
@ -8,8 +8,12 @@ import sys
|
|||
|
||||
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[1]))
|
||||
|
||||
from app.fallback import SolCircuitBreaker # noqa: E402
|
||||
from app.search import run_search # noqa: E402
|
||||
|
||||
SOLARIA_URL = "http://fake-solaria"
|
||||
PIHA_URL = "http://fake-piha"
|
||||
|
||||
|
||||
class _FakeConn:
|
||||
"""summaries: [(envelope_id, dist), ...] -- cascade stage-1 pre-filter. chunks_by_envelope:
|
||||
|
|
@ -92,11 +96,27 @@ class _FakeEmbedResponse:
|
|||
return self._payload
|
||||
|
||||
|
||||
class _FakeHealthResponse:
|
||||
status = 200
|
||||
|
||||
async def __aenter__(self):
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *exc):
|
||||
return False
|
||||
|
||||
|
||||
class _FakeSession:
|
||||
"""SOLARIA always reachable and healthy -- these tests exercise run_search's DB-joining
|
||||
logic, not the fallback state machine itself (see test_fallback.py for that)."""
|
||||
|
||||
def __init__(self):
|
||||
self.post_calls: list[dict] = []
|
||||
|
||||
def post(self, url, json):
|
||||
def get(self, url, timeout=None):
|
||||
return _FakeHealthResponse()
|
||||
|
||||
def post(self, url, json, timeout=None):
|
||||
self.post_calls.append({"url": url, "json": json})
|
||||
return _FakeEmbedResponse({"embedding": [0.01] * 1024})
|
||||
|
||||
|
|
@ -110,7 +130,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
assert result["query"] == "polisa PZU"
|
||||
assert result["mode"] == "cascade"
|
||||
|
|
@ -131,7 +151,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "q", "flat", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "q", "flat", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
assert result["mode"] == "flat"
|
||||
assert len(result["results"]) == 1
|
||||
|
|
@ -145,7 +165,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
hit = result["results"][0]
|
||||
assert hit["summary"] == "Polisa OC 2024"
|
||||
|
|
@ -159,7 +179,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "polisa PZU", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
hit = result["results"][0]
|
||||
assert hit["summary"] is None
|
||||
|
|
@ -180,7 +200,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "q", "hybrid", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "q", "hybrid", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
assert result["mode"] == "hybrid"
|
||||
envelope_ids = [r["envelope_id"] for r in result["results"]]
|
||||
|
|
@ -201,7 +221,7 @@ class TestRunSearchHappyPath:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "q", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "q", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
hit = result["results"][0]
|
||||
assert hit["source"] == "gmail"
|
||||
|
|
@ -222,7 +242,7 @@ class TestRunSearchNoGoodResults:
|
|||
)
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "unrelated query", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "unrelated query", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
assert len(result["results"]) == 1
|
||||
assert result["results"][0]["dist"] == 0.62
|
||||
|
|
@ -231,6 +251,6 @@ class TestRunSearchNoGoodResults:
|
|||
conn = _FakeConn(summaries=[], chunks_by_envelope={}, envelopes={})
|
||||
session = _FakeSession()
|
||||
result = await run_search(
|
||||
conn, session, "http://fake-ollama", "nothing matches", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
conn, session, SolCircuitBreaker(), SOLARIA_URL, PIHA_URL, "nothing matches", "cascade", "bge-m3", "claude-haiku-4-5"
|
||||
)
|
||||
assert result["results"] == []
|
||||
|
|
|
|||
73
services/ollama-piha/README.md
Normal file
73
services/ollama-piha/README.md
Normal file
|
|
@ -0,0 +1,73 @@
|
|||
# ollama-piha
|
||||
|
||||
Local, CPU-only Ollama on **PIHA** (arm64, 4 cores, no acceleration —
|
||||
`hosts/piha/capabilities.yaml`), serving exactly one purpose: `kb-query`'s
|
||||
embed fallback (`services/kb-query/app/fallback.py`, module 5 phase 4 plan
|
||||
§2 decision 2 / §5) when Ollama@SOLARIA is unreachable or times out.
|
||||
|
||||
**This is not a general-purpose Ollama instance** — no other service should
|
||||
point at it. It runs `bge-m3` only, the same embedding model kb-query's
|
||||
startup invariant enforces against `document_chunk`/`document_summary`
|
||||
(`services/kb-query/app/startup.py`). Do not pull additional models onto it.
|
||||
|
||||
## Why `OLLAMA_KEEP_ALIVE=0`
|
||||
|
||||
The model is loaded into RAM only for the duration of a request and released
|
||||
immediately after, rather than staying resident. On a memory-constrained
|
||||
RPi5 already running kb-postgres, paperless, Home Assistant, Immich,
|
||||
Forgejo, and more (see `docs/infra/piha-slim-audit-2026-07-02.md`), a second
|
||||
permanently-resident ~1.5-2 GB model is a worse trade than a short RAM spike
|
||||
that only happens when this fallback is actually exercised (SOLARIA down —
|
||||
rare, plan §1.2 `availability_target: medium`).
|
||||
|
||||
## Deploy
|
||||
|
||||
1. `git pull` on PIHA.
|
||||
2. `cp services/ollama-piha/env.example services/ollama-piha/.env` (fill in
|
||||
`LAN_BIND_IP` if it differs from the default).
|
||||
3. `docker compose -f services/ollama-piha/docker-compose.yml -f hosts/piha/runtime/ollama-piha/docker-compose.override.yml up -d`
|
||||
4. `docker exec ollama-piha ollama pull bge-m3` — not baked into the image;
|
||||
the model must be pulled once after first start (persists in the
|
||||
`/opt/homelab/data/ollama-piha` named volume across restarts).
|
||||
5. Verify: `services/ollama-piha/healthcheck.sh`, then
|
||||
`curl http://192.168.31.5:11434/api/tags` should list `bge-m3`.
|
||||
|
||||
Point `kb-query`'s `OLLAMA_PIHA_URL` at `http://192.168.31.5:11434` once this
|
||||
is live (see `services/kb-query/env.example`).
|
||||
|
||||
## Calibration status (plan §5) — GO, measured live 2026-07-27
|
||||
|
||||
Measured on live PIHA under normal load (kb-postgres, paperless, Immich,
|
||||
Home Assistant, Forgejo, etc. all running, not a quiet-night window),
|
||||
3 sequential embed calls through `docker exec ollama-piha ollama pull bge-m3`
|
||||
+ `/api/embeddings`:
|
||||
|
||||
- **Latency**: 5.25s (first call), 4.41s, 4.16s — consistently single-digit
|
||||
seconds, never tens of seconds. Latency doesn't drop on later calls because
|
||||
`OLLAMA_KEEP_ALIVE=0` reloads the model every time by design (`ollama ps`
|
||||
shows zero resident models between calls) — this is the expected trade-off
|
||||
documented above, not a bug.
|
||||
- **RAM**: peak ~983 MiB during a burst (`docker stats`, baseline idle ~66
|
||||
MiB), comfortably inside the `mem_limit: 2560m` ceiling
|
||||
(`hosts/piha/runtime/ollama-piha/docker-compose.override.yml`). System-wide
|
||||
`available` memory never dropped below ~1.3 GiB during the burst and
|
||||
settled back to ~4.2 GiB after — well above the "≥500 MB spare" bar from
|
||||
the plan.
|
||||
|
||||
**Verdict: GO — enabled as the default fallback**, no
|
||||
`KB_QUERY_LOCAL_FALLBACK_ENABLED`-style flag needed. `kb-query`'s
|
||||
`OLLAMA_PIHA_URL` points at this container's real LAN address
|
||||
(`http://192.168.31.5:11434`) in the live PIHA deployment. Verified live: a
|
||||
sol-down simulation (`OLLAMA_URL` on kb-query pointed at an unreachable
|
||||
SOLARIA address) produced `sol_status: "down"` and correct `/search` results
|
||||
from this container, with `dist` within ~3e-4 of the SOLARIA-GPU baseline
|
||||
(same top-k order, same hit@3 gate outcome) — see
|
||||
`docs/sessions/2026-07-27-kb-f4-fallback.md` for the full numbers.
|
||||
|
||||
**One-time finding from this calibration**: PIHA also had a leftover, fully
|
||||
undocumented *native* (non-Docker) `ollama.service` (systemd, v0.6.1, running
|
||||
since 2026-06-22, zero models ever pulled) that conflicted with this
|
||||
container's port binding. Confirmed dead (only this session's own probe
|
||||
requests in its journal) and disabled (`systemctl disable --now`, not
|
||||
uninstalled — reversible) before deploying this container. See the session
|
||||
doc for the backlog note to fully remove it if nothing breaks.
|
||||
26
services/ollama-piha/docker-compose.yml
Normal file
26
services/ollama-piha/docker-compose.yml
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
services:
|
||||
ollama-piha:
|
||||
image: ollama/ollama:latest
|
||||
container_name: ollama-piha
|
||||
restart: unless-stopped
|
||||
ports:
|
||||
# Loopback: healthcheck.sh curls localhost directly on the node. LAN IP: kb-query@PIHA
|
||||
# reaches this over the host's LAN interface -- kb-query runs in its own Docker network
|
||||
# (separate compose project), same reasoning as kb-query's KB_DSN reaching kb-postgres.
|
||||
# Requires .env (from env.example) next to this file at deploy.
|
||||
- "127.0.0.1:11434:11434"
|
||||
- "${LAN_BIND_IP}:11434:11434"
|
||||
environment:
|
||||
# Module 5 phase 4 plan §2 decision 2 / §5 requirement: the model is loaded only for the
|
||||
# duration of a request and released immediately after -- RPi5 has no GPU and limited RAM
|
||||
# (hosts/piha/capabilities.yaml: arm64, 4 cores, no acceleration), so this is a short burst
|
||||
# spike (idle Ollama binary ~100 MB) rather than a permanent ~1.5-2 GB resident cost. This
|
||||
# is a fallback-only path (kb-query only reaches this when SOLARIA is down/times out), not
|
||||
# the default hot path, so paying a cold-load per request here is the correct trade-off.
|
||||
- OLLAMA_KEEP_ALIVE=0
|
||||
volumes:
|
||||
- /opt/homelab/data/ollama-piha:/root/.ollama
|
||||
# No GPU reservation -- PIHA is arm64 with no acceleration (hosts/piha/capabilities.yaml),
|
||||
# unlike services/ollama@SOLARIA. CPU-only inference here is expected to be slower; that is
|
||||
# exactly what the plan §5 live calibration step measures before this is trusted as a
|
||||
# default fallback (see README.md "Calibration status").
|
||||
8
services/ollama-piha/env.example
Normal file
8
services/ollama-piha/env.example
Normal file
|
|
@ -0,0 +1,8 @@
|
|||
# Copy to .env next to docker-compose.yml (gitignored); docker compose picks
|
||||
# it up automatically. Same convention as services/kb-query.
|
||||
|
||||
# LAN IP of PIHA. The published port (11434) binds ONLY to this interface —
|
||||
# never 0.0.0.0. Verify after host rebuilds: ip -4 addr. Same value as
|
||||
# kb-query's LAN_BIND_IP (services/kb-query/env.example) — both containers
|
||||
# run on the same node.
|
||||
LAN_BIND_IP=192.168.31.5
|
||||
15
services/ollama-piha/healthcheck.sh
Executable file
15
services/ollama-piha/healthcheck.sh
Executable file
|
|
@ -0,0 +1,15 @@
|
|||
#!/bin/bash
|
||||
# Healthcheck for ollama-piha (module 5 phase 4 local embed fallback).
|
||||
|
||||
if ! docker ps --filter "name=ollama-piha" --filter "status=running" | grep -q "ollama-piha"; then
|
||||
echo "[FAIL] ollama-piha container is not running"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if ! curl -sf http://localhost:11434/api/tags > /dev/null; then
|
||||
echo "[FAIL] ollama-piha API is not responding"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
echo "[OK] ollama-piha is healthy"
|
||||
exit 0
|
||||
33
services/ollama-piha/service.yaml
Normal file
33
services/ollama-piha/service.yaml
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
service:
|
||||
name: ollama-piha
|
||||
owner_node: piha
|
||||
role: kb-embed-fallback # module 5 phase 4 plan §2 decision 2 / §5: local fallback embed
|
||||
# when Ollama@SOLARIA is unreachable/times out. NOT the
|
||||
# default embed path -- kb-query only reaches this via its
|
||||
# circuit breaker (services/kb-query/app/fallback.py).
|
||||
exposure: private # LAN bind (LAN_BIND_IP), consumed only by kb-query@PIHA today
|
||||
dependencies: []
|
||||
ports:
|
||||
- container: 11434
|
||||
host: 11434
|
||||
protocol: tcp
|
||||
healthcheck:
|
||||
type: http
|
||||
endpoint: http://localhost:11434/api/tags
|
||||
interval: 1m
|
||||
timeout: 10s
|
||||
retries: 3
|
||||
restart_policy: unless-stopped
|
||||
persistence:
|
||||
paths:
|
||||
- /opt/homelab/data/ollama-piha
|
||||
runtime:
|
||||
directories:
|
||||
- /opt/homelab/data/ollama-piha
|
||||
env_vars:
|
||||
- LAN_BIND_IP # required — compose port-bind interpolation
|
||||
config_files:
|
||||
- .env # LAN_BIND_IP (gitignored, from env.example)
|
||||
# Same model as document_chunk.model / document_summary.embedding_model (bge-m3) -- kb-query's
|
||||
# startup invariant (app/startup.py) enforces this once for both the SOLARIA and PIHA legs.
|
||||
# `ollama pull bge-m3` is a manual deploy step (see README.md), not baked into the image.
|
||||
Loading…
Reference in a new issue