diff --git a/docs/sessions/2026-07-27-kb-f4-fallback.md b/docs/sessions/2026-07-27-kb-f4-fallback.md new file mode 100644 index 0000000..8c27f0f --- /dev/null +++ b/docs/sessions/2026-07-27-kb-f4-fallback.md @@ -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. diff --git a/hosts/piha/runtime/ollama-piha/docker-compose.override.yml b/hosts/piha/runtime/ollama-piha/docker-compose.override.yml new file mode 100644 index 0000000..dce83ed --- /dev/null +++ b/hosts/piha/runtime/ollama-piha/docker-compose.override.yml @@ -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 diff --git a/hosts/piha/services.yaml b/hosts/piha/services.yaml index 825a589..2e34684 100644 --- a/hosts/piha/services.yaml +++ b/hosts/piha/services.yaml @@ -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. diff --git a/jobs/documents-ingest/README.md b/jobs/documents-ingest/README.md index bf3fce1..624a05f 100644 --- a/jobs/documents-ingest/README.md +++ b/jobs/documents-ingest/README.md @@ -517,6 +517,12 @@ python eval/retrieval_eval.py --dsn postgresql://kb:@piha:5433/kb \ --ollama-url http://solaria:11434 --n-sweep 5,10,20 ``` +`--transport http --base-url http://: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, diff --git a/jobs/documents-ingest/eval/retrieval_eval.py b/jobs/documents-ingest/eval/retrieval_eval.py index 416eab3..f0aa063 100644 --- a/jobs/documents-ingest/eval/retrieval_eval.py +++ b/jobs/documents-ingest/eval/retrieval_eval.py @@ -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:@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) diff --git a/packages/kb-retrieval/src/kb_retrieval/embed.py b/packages/kb-retrieval/src/kb_retrieval/embed.py index a2aad28..20ca0aa 100644 --- a/packages/kb-retrieval/src/kb_retrieval/embed.py +++ b/packages/kb-retrieval/src/kb_retrieval/embed.py @@ -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 diff --git a/services/kb-query/README.md b/services/kb-query/README.md index 0982597..5e469f9 100644 --- a/services/kb-query/README.md +++ b/services/kb-query/README.md @@ -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. diff --git a/services/kb-query/app/fallback.py b/services/kb-query/app/fallback.py new file mode 100644 index 0000000..5b44e09 --- /dev/null +++ b/services/kb-query/app/fallback.py @@ -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" diff --git a/services/kb-query/app/main.py b/services/kb-query/app/main.py index 1bf7d3d..f854bf0 100644 --- a/services/kb-query/app/main.py +++ b/services/kb-query/app/main.py @@ -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 diff --git a/services/kb-query/app/search.py b/services/kb-query/app/search.py index be1fee6..46c8b81 100644 --- a/services/kb-query/app/search.py +++ b/services/kb-query/app/search.py @@ -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, } diff --git a/services/kb-query/docker-compose.yml b/services/kb-query/docker-compose.yml index dc524dd..e14083d 100644 --- a/services/kb-query/docker-compose.yml +++ b/services/kb-query/docker-compose.yml @@ -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. diff --git a/services/kb-query/env.example b/services/kb-query/env.example index 4b23f5d..170c75c 100644 --- a/services/kb-query/env.example +++ b/services/kb-query/env.example @@ -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. diff --git a/services/kb-query/service.yaml b/services/kb-query/service.yaml index 37c7bb4..bd1a18f 100644 --- a/services/kb-query/service.yaml +++ b/services/kb-query/service.yaml @@ -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 diff --git a/services/kb-query/tests/test_fallback.py b/services/kb-query/tests/test_fallback.py new file mode 100644 index 0000000..96f1793 --- /dev/null +++ b/services/kb-query/tests/test_fallback.py @@ -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 diff --git a/services/kb-query/tests/test_search.py b/services/kb-query/tests/test_search.py index cdcb97a..85637b3 100644 --- a/services/kb-query/tests/test_search.py +++ b/services/kb-query/tests/test_search.py @@ -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"] == [] diff --git a/services/ollama-piha/README.md b/services/ollama-piha/README.md new file mode 100644 index 0000000..db27689 --- /dev/null +++ b/services/ollama-piha/README.md @@ -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. diff --git a/services/ollama-piha/docker-compose.yml b/services/ollama-piha/docker-compose.yml new file mode 100644 index 0000000..5a9803b --- /dev/null +++ b/services/ollama-piha/docker-compose.yml @@ -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"). diff --git a/services/ollama-piha/env.example b/services/ollama-piha/env.example new file mode 100644 index 0000000..b2c3f6a --- /dev/null +++ b/services/ollama-piha/env.example @@ -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 diff --git a/services/ollama-piha/healthcheck.sh b/services/ollama-piha/healthcheck.sh new file mode 100755 index 0000000..2f81528 --- /dev/null +++ b/services/ollama-piha/healthcheck.sh @@ -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 diff --git a/services/ollama-piha/service.yaml b/services/ollama-piha/service.yaml new file mode 100644 index 0000000..1c0ea82 --- /dev/null +++ b/services/ollama-piha/service.yaml @@ -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.