Commit ·
70b1f2f
1
Parent(s): 71ad734
Fused retrieval from directive and løn
Browse files- README.md +84 -1
- data/processed/index/chunks.jsonl +0 -0
- scripts/refresh_metadata.py +75 -0
- src/backend/README.md +22 -15
- src/backend/indexing/ingest.py +18 -0
- src/backend/indexing/store.py +30 -0
- src/backend/rag.py +38 -9
- src/backend/schemas.py +5 -0
- src/frontend/README.md +2 -2
README.md
CHANGED
|
@@ -12,4 +12,87 @@ license: apache-2.0
|
|
| 12 |
short_description: Multilingual AI assistant to guarantee fair pay
|
| 13 |
---
|
| 14 |
|
| 15 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 12 |
short_description: Multilingual AI assistant to guarantee fair pay
|
| 13 |
---
|
| 14 |
|
| 15 |
+
# Pay Equity for EU
|
| 16 |
+
|
| 17 |
+
Multilingual Gradio Space for Danish users: **expected pay** (IDA lønstatistik)
|
| 18 |
+
and **pay-transparency rights** (EU Directive 2023/970), grounded with citations.
|
| 19 |
+
|
| 20 |
+
## Quick start
|
| 21 |
+
|
| 22 |
+
```bash
|
| 23 |
+
uv sync
|
| 24 |
+
uv run python main.py
|
| 25 |
+
```
|
| 26 |
+
|
| 27 |
+
Requires a built vector index at `data/processed/index/` (see below). Backend
|
| 28 |
+
details: [`src/backend/README.md`](src/backend/README.md).
|
| 29 |
+
|
| 30 |
+
## Build the vector index
|
| 31 |
+
|
| 32 |
+
Embed the corpus once (directive articles + IDA salary rows), then commit the
|
| 33 |
+
output:
|
| 34 |
+
|
| 35 |
+
```bash
|
| 36 |
+
uv run python scripts/build_index.py # GPU / ZeroGPU Space admin button
|
| 37 |
+
```
|
| 38 |
+
|
| 39 |
+
Produces `data/processed/index/index.faiss` + `chunks.jsonl`.
|
| 40 |
+
|
| 41 |
+
### Refresh metadata only (no re-embed)
|
| 42 |
+
|
| 43 |
+
After changing `Chunk.metadata` mapping in `src/backend/indexing/ingest.py`:
|
| 44 |
+
|
| 45 |
+
```bash
|
| 46 |
+
uv run python scripts/refresh_metadata.py
|
| 47 |
+
```
|
| 48 |
+
|
| 49 |
+
Rewrites `chunks.jsonl` only; vectors stay unchanged. Re-embed if chunk ids or
|
| 50 |
+
order changed.
|
| 51 |
+
|
| 52 |
+
## Run the RAG benchmark
|
| 53 |
+
|
| 54 |
+
Offline evaluation of retrieval (and lightweight generation checks) against
|
| 55 |
+
15 ground-truth questions in `scripts/eval/benchmark_questions.json`.
|
| 56 |
+
|
| 57 |
+
**Prerequisites:** committed `chunks.jsonl` plus a local `index.faiss`. If
|
| 58 |
+
`index.faiss` is missing or a Git LFS pointer, rebuild it from the committed
|
| 59 |
+
chunks (CPU-only, ~few minutes):
|
| 60 |
+
|
| 61 |
+
```bash
|
| 62 |
+
uv run python scripts/eval/build_local_index.py
|
| 63 |
+
```
|
| 64 |
+
|
| 65 |
+
**Run the benchmark:**
|
| 66 |
+
|
| 67 |
+
```bash
|
| 68 |
+
# Full run: retrieval metrics + generation fact/citation checks (default k=5)
|
| 69 |
+
uv run python scripts/eval/run_benchmark.py
|
| 70 |
+
|
| 71 |
+
# Retrieval only (faster)
|
| 72 |
+
uv run python scripts/eval/run_benchmark.py --no-generation
|
| 73 |
+
|
| 74 |
+
# Custom top-k
|
| 75 |
+
uv run python scripts/eval/run_benchmark.py --k 10
|
| 76 |
+
```
|
| 77 |
+
|
| 78 |
+
**Output:** human-readable summary on stdout + JSON report at
|
| 79 |
+
`scripts/eval/results/benchmark_YYYYMMDD_HHMMSS.json`.
|
| 80 |
+
|
| 81 |
+
**Metrics:** Recall@k, MRR, NDCG@k (retrieval); fact coverage and citation
|
| 82 |
+
presence (generation). See [`scripts/eval/BENCHMARK_FINDINGS.md`](scripts/eval/BENCHMARK_FINDINGS.md)
|
| 83 |
+
for the latest scored run and failure analysis.
|
| 84 |
+
|
| 85 |
+
> **Note:** The benchmark script uses a single global FAISS top-k search, not
|
| 86 |
+
> the production per-source + RRF fusion in `rag.py`. Use it to track embedding
|
| 87 |
+
> quality; for end-to-end production behaviour, test via the Gradio UI or
|
| 88 |
+
> `/chat` API.
|
| 89 |
+
|
| 90 |
+
## Data
|
| 91 |
+
|
| 92 |
+
| Path | Contents |
|
| 93 |
+
|---|---|
|
| 94 |
+
| `data/processed/directive/` | Parsed EU Directive 2023/970 |
|
| 95 |
+
| `data/processed/lonstatistik/` | IDA + Djøf salary statistics ([README](data/processed/lonstatistik/README.md)) |
|
| 96 |
+
| `data/processed/index/` | FAISS index + aligned `chunks.jsonl` |
|
| 97 |
+
|
| 98 |
+
Raw PDFs live under `data/raw/` (gitignored).
|
data/processed/index/chunks.jsonl
CHANGED
|
The diff for this file is too large to render.
See raw diff
|
|
|
scripts/refresh_metadata.py
ADDED
|
@@ -0,0 +1,75 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""Regenerate chunks.jsonl metadata without re-embedding (CPU-only, offline).
|
| 2 |
+
|
| 3 |
+
The dual-source retrieval work enriches each ``Chunk`` with structured
|
| 4 |
+
``metadata`` (IDA: sector/category/experience/measures; directive: article
|
| 5 |
+
title). That metadata is **not embedded**, so the vectors and chunk order are
|
| 6 |
+
unchanged — meaning the committed ``index.faiss`` stays byte-identical and we
|
| 7 |
+
only need to rewrite ``chunks.jsonl``.
|
| 8 |
+
|
| 9 |
+
This script loads the existing index + old chunks, rebuilds the corpus with the
|
| 10 |
+
new metadata, asserts the rebuilt chunks line up 1:1 with the live vectors
|
| 11 |
+
(same ids, same order, same count), then re-persists. No GPU, no embedding —
|
| 12 |
+
runs in well under a second.
|
| 13 |
+
|
| 14 |
+
Usage::
|
| 15 |
+
|
| 16 |
+
uv run python scripts/refresh_metadata.py
|
| 17 |
+
"""
|
| 18 |
+
|
| 19 |
+
from __future__ import annotations
|
| 20 |
+
|
| 21 |
+
import sys
|
| 22 |
+
from pathlib import Path
|
| 23 |
+
|
| 24 |
+
# Allow running as a plain script (`python scripts/refresh_metadata.py`).
|
| 25 |
+
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
| 26 |
+
|
| 27 |
+
from src.backend.indexing.ingest import ( # noqa: E402
|
| 28 |
+
build_corpus,
|
| 29 |
+
load_directive_chunks,
|
| 30 |
+
load_ida_chunks,
|
| 31 |
+
)
|
| 32 |
+
from src.backend.indexing.store import DEFAULT_INDEX_DIR, VectorStore # noqa: E402
|
| 33 |
+
|
| 34 |
+
|
| 35 |
+
def main() -> None:
|
| 36 |
+
store = VectorStore()
|
| 37 |
+
if not store.load():
|
| 38 |
+
raise SystemExit(
|
| 39 |
+
f"No index at {DEFAULT_INDEX_DIR}. Build it first (scripts/build_index.py "
|
| 40 |
+
"or scripts/build_index_modal.py) before refreshing metadata."
|
| 41 |
+
)
|
| 42 |
+
|
| 43 |
+
old_ids = [c.id for c in store._chunks]
|
| 44 |
+
ntotal = store._index.ntotal
|
| 45 |
+
|
| 46 |
+
n_dir = len(load_directive_chunks())
|
| 47 |
+
n_ida = len(load_ida_chunks())
|
| 48 |
+
new_chunks = build_corpus()
|
| 49 |
+
new_ids = [c.id for c in new_chunks]
|
| 50 |
+
|
| 51 |
+
# Alignment guard: the rebuilt chunks must line up exactly with the stored
|
| 52 |
+
# vectors, else chunks.jsonl would point at the wrong embeddings.
|
| 53 |
+
if len(new_chunks) != ntotal:
|
| 54 |
+
raise SystemExit(
|
| 55 |
+
f"Count mismatch: rebuilt {len(new_chunks)} chunks but index has "
|
| 56 |
+
f"{ntotal} vectors. Re-embed (build_index) instead of refreshing."
|
| 57 |
+
)
|
| 58 |
+
if new_ids != old_ids:
|
| 59 |
+
raise SystemExit(
|
| 60 |
+
"Order/id mismatch between rebuilt chunks and the existing index. "
|
| 61 |
+
"The corpus changed shape — re-embed (build_index) instead."
|
| 62 |
+
)
|
| 63 |
+
|
| 64 |
+
store._chunks = new_chunks
|
| 65 |
+
store.persist() # rewrites chunks.jsonl (+ identical index.faiss)
|
| 66 |
+
|
| 67 |
+
n_meta = sum(1 for c in new_chunks if c.metadata)
|
| 68 |
+
print(
|
| 69 |
+
f"refreshed: directive={n_dir} ida={n_ida} total={len(new_chunks)} "
|
| 70 |
+
f"ntotal={ntotal} with_metadata={n_meta} → {DEFAULT_INDEX_DIR}"
|
| 71 |
+
)
|
| 72 |
+
|
| 73 |
+
|
| 74 |
+
if __name__ == "__main__":
|
| 75 |
+
main()
|
src/backend/README.md
CHANGED
|
@@ -14,33 +14,38 @@ From the repo root:
|
|
| 14 |
uv run python main.py
|
| 15 |
```
|
| 16 |
|
| 17 |
-
`main.py`
|
| 18 |
-
(default `http://127.0.0.1:7860`). Open it for the
|
| 19 |
-
|
| 20 |
|
| 21 |
Before first run, build the vector index (see **RAG pipeline** below).
|
| 22 |
|
| 23 |
## RAG pipeline
|
| 24 |
|
| 25 |
```
|
| 26 |
-
query ─► embed (Nemotron) ─►
|
| 27 |
```
|
| 28 |
|
| 29 |
- **Ingestion** (`indexing/ingest.py`): the Directive `document.json` is grouped
|
| 30 |
-
into one `Chunk` per article (`Art. N` citations
|
| 31 |
-
(`union == "IDA"`) become chunks from their
|
| 32 |
-
source page)
|
|
|
|
|
|
|
|
|
|
| 33 |
- **Offline index build** (`scripts/build_index.py`): embeds the whole corpus
|
| 34 |
on ZeroGPU (the embedding call is wrapped in `@spaces.GPU`) and persists
|
| 35 |
-
`data/processed/index/{index.faiss, chunks.jsonl}`. Run
|
| 36 |
-
Space (or any GPU machine) and **commit the output** —
|
| 37 |
*loads* the index and never re-embeds the corpus at runtime.
|
|
|
|
|
|
|
|
|
|
| 38 |
- **Runtime** (`rag.py`): `answer_stream()` is the single grounding path shared
|
| 39 |
-
by the UI and (via `answer()`) the `/chat` API.
|
| 40 |
-
|
| 41 |
-
|
| 42 |
-
|
| 43 |
-
consumes the same stream and returns the final response.
|
| 44 |
|
| 45 |
Models are pinned by revision in `configs/models.yaml`. The LLM backend is
|
| 46 |
selected by `llm.<role>.backend` (`tiny_aya` | `minicpm` | `llama_cpp` |
|
|
@@ -87,5 +92,7 @@ indexing/embedder.py Embedder (config-driven, sentence-transformers | transforme
|
|
| 87 |
indexing/store.py VectorStore (FAISS IndexFlatIP + persist/load/search)
|
| 88 |
parsing/ directive.py (→ ingest), lonstatistik.py (placeholder)
|
| 89 |
api/endpoints.py register(app): the @app.api endpoints
|
| 90 |
-
scripts/build_index.py
|
|
|
|
|
|
|
| 91 |
```
|
|
|
|
| 14 |
uv run python main.py
|
| 15 |
```
|
| 16 |
|
| 17 |
+
`main.py` builds `demo = build_blocks()` and calls `demo.launch(...)`. The
|
| 18 |
+
server prints a local URL (default `http://127.0.0.1:7860`). Open it for the
|
| 19 |
+
chat UI.
|
| 20 |
|
| 21 |
Before first run, build the vector index (see **RAG pipeline** below).
|
| 22 |
|
| 23 |
## RAG pipeline
|
| 24 |
|
| 25 |
```
|
| 26 |
+
query ─► embed (Nemotron) ─► per-source top-k ─► RRF fuse ─► grounded prompt ─► Tiny Aya ─► cited answer
|
| 27 |
```
|
| 28 |
|
| 29 |
- **Ingestion** (`indexing/ingest.py`): the Directive `document.json` is grouped
|
| 30 |
+
into one `Chunk` per article (`Art. N` citations + `metadata.article_number` /
|
| 31 |
+
`article_title`); IDA rows (`union == "IDA"`) become chunks from their
|
| 32 |
+
ready-made `rag_text` (cited by source page) with structured salary
|
| 33 |
+
`metadata` (sector, category, experience, measures, …). Djøf is excluded for
|
| 34 |
+
now. Only `Chunk.text` is embedded; `metadata` is persisted in `chunks.jsonl`
|
| 35 |
+
for downstream use.
|
| 36 |
- **Offline index build** (`scripts/build_index.py`): embeds the whole corpus
|
| 37 |
on ZeroGPU (the embedding call is wrapped in `@spaces.GPU`) and persists
|
| 38 |
+
`data/processed/index/{index.faiss, chunks.jsonl}`. Run once on the ZeroGPU
|
| 39 |
+
Space (or any GPU machine) and **commit the output** — the live app only
|
| 40 |
*loads* the index and never re-embeds the corpus at runtime.
|
| 41 |
+
- **Metadata refresh** (`scripts/refresh_metadata.py`): rewrites `chunks.jsonl`
|
| 42 |
+
from updated ingest logic **without** re-embedding. Safe only when chunk ids and
|
| 43 |
+
order are unchanged; otherwise re-run `build_index.py`.
|
| 44 |
- **Runtime** (`rag.py`): `answer_stream()` is the single grounding path shared
|
| 45 |
+
by the UI and (via `answer()`) the `/chat` API. Retrieval runs **per source**
|
| 46 |
+
(`directive`, `lonstatistik`) — top-5 each — then **Reciprocal Rank Fusion**
|
| 47 |
+
into the final context (default 6 chunks). Query embedding + generation run
|
| 48 |
+
inside one `@spaces.GPU` function (`_retrieve_and_stream`).
|
|
|
|
| 49 |
|
| 50 |
Models are pinned by revision in `configs/models.yaml`. The LLM backend is
|
| 51 |
selected by `llm.<role>.backend` (`tiny_aya` | `minicpm` | `llama_cpp` |
|
|
|
|
| 92 |
indexing/store.py VectorStore (FAISS IndexFlatIP + persist/load/search)
|
| 93 |
parsing/ directive.py (→ ingest), lonstatistik.py (placeholder)
|
| 94 |
api/endpoints.py register(app): the @app.api endpoints
|
| 95 |
+
scripts/build_index.py offline corpus embedding → data/processed/index/
|
| 96 |
+
scripts/refresh_metadata.py rewrite chunks.jsonl metadata without re-embedding
|
| 97 |
+
scripts/eval/run_benchmark.py offline retrieval benchmark (see root README.md)
|
| 98 |
```
|
src/backend/indexing/ingest.py
CHANGED
|
@@ -82,6 +82,10 @@ def load_directive_chunks(
|
|
| 82 |
article=f"Art. {cur_num}",
|
| 83 |
page=cur_page,
|
| 84 |
citation=f"Art. {cur_num}",
|
|
|
|
|
|
|
|
|
|
|
|
|
| 85 |
)
|
| 86 |
)
|
| 87 |
cur_body = []
|
|
@@ -134,6 +138,20 @@ def load_ida_chunks(path: str | Path = SALARY_JSONL) -> list[Chunk]:
|
|
| 134 |
source="lonstatistik",
|
| 135 |
page=page,
|
| 136 |
citation=f"{rec['source_doc']}, p.{page}",
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 137 |
)
|
| 138 |
)
|
| 139 |
return chunks
|
|
|
|
| 82 |
article=f"Art. {cur_num}",
|
| 83 |
page=cur_page,
|
| 84 |
citation=f"Art. {cur_num}",
|
| 85 |
+
metadata={
|
| 86 |
+
"article_number": cur_num,
|
| 87 |
+
"article_title": cur_title,
|
| 88 |
+
},
|
| 89 |
)
|
| 90 |
)
|
| 91 |
cur_body = []
|
|
|
|
| 138 |
source="lonstatistik",
|
| 139 |
page=page,
|
| 140 |
citation=f"{rec['source_doc']}, p.{page}",
|
| 141 |
+
metadata={
|
| 142 |
+
"sector": rec.get("sector"),
|
| 143 |
+
"category": rec.get("table_title"),
|
| 144 |
+
"row_label": rec.get("row_label"),
|
| 145 |
+
"segment_label": rec.get("segment_label"),
|
| 146 |
+
"experience_years_min": rec.get("experience_years_min"),
|
| 147 |
+
"experience_years_max": rec.get("experience_years_max"),
|
| 148 |
+
"graduation_year_start": rec.get("graduation_year_start"),
|
| 149 |
+
"graduation_year_end": rec.get("graduation_year_end"),
|
| 150 |
+
"period": rec.get("period"),
|
| 151 |
+
"pay_concept": rec.get("pay_concept"),
|
| 152 |
+
"currency": rec.get("currency"),
|
| 153 |
+
"measure": rec.get("measure"),
|
| 154 |
+
},
|
| 155 |
)
|
| 156 |
)
|
| 157 |
return chunks
|
src/backend/indexing/store.py
CHANGED
|
@@ -13,6 +13,7 @@ On-disk layout (``data/processed/index/``)::
|
|
| 13 |
|
| 14 |
from __future__ import annotations
|
| 15 |
|
|
|
|
| 16 |
from pathlib import Path
|
| 17 |
|
| 18 |
import numpy as np
|
|
@@ -33,6 +34,9 @@ class VectorStore:
|
|
| 33 |
self.backend = backend
|
| 34 |
self._index = None # faiss.IndexFlatIP
|
| 35 |
self._chunks: list[Chunk] = []
|
|
|
|
|
|
|
|
|
|
| 36 |
|
| 37 |
def add(self, chunks: list[Chunk], vectors) -> int:
|
| 38 |
"""Add chunks and their (N, D) vectors; returns the number added."""
|
|
@@ -76,8 +80,34 @@ class VectorStore:
|
|
| 76 |
for line in chunks_path.read_text(encoding="utf-8").splitlines()
|
| 77 |
if line.strip()
|
| 78 |
]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 79 |
return True
|
| 80 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 81 |
def search(self, query_vector, k: int = 5) -> list[tuple[Chunk, float]]:
|
| 82 |
"""Return up to ``k`` nearest (chunk, score) pairs, best first."""
|
| 83 |
if self._index is None or not self._chunks:
|
|
|
|
| 13 |
|
| 14 |
from __future__ import annotations
|
| 15 |
|
| 16 |
+
from collections import defaultdict
|
| 17 |
from pathlib import Path
|
| 18 |
|
| 19 |
import numpy as np
|
|
|
|
| 34 |
self.backend = backend
|
| 35 |
self._index = None # faiss.IndexFlatIP
|
| 36 |
self._chunks: list[Chunk] = []
|
| 37 |
+
# Populated on load() for source-masked retrieval (see search_by_source).
|
| 38 |
+
self._vectors: np.ndarray | None = None # (N, D) float32, row-aligned to _chunks
|
| 39 |
+
self._rows_by_source: dict[str, list[int]] = {}
|
| 40 |
|
| 41 |
def add(self, chunks: list[Chunk], vectors) -> int:
|
| 42 |
"""Add chunks and their (N, D) vectors; returns the number added."""
|
|
|
|
| 80 |
for line in chunks_path.read_text(encoding="utf-8").splitlines()
|
| 81 |
if line.strip()
|
| 82 |
]
|
| 83 |
+
# Reconstruct the stored vectors once (flat index → exact) so we can run
|
| 84 |
+
# cheap per-source masked search in numpy. The corpus is ~2k vectors, so
|
| 85 |
+
# this is trivial in both time and memory (~16 MB).
|
| 86 |
+
n = self._index.ntotal
|
| 87 |
+
self._vectors = self._index.reconstruct_n(0, n) if n else None
|
| 88 |
+
rows: dict[str, list[int]] = defaultdict(list)
|
| 89 |
+
for i, c in enumerate(self._chunks):
|
| 90 |
+
rows[c.source].append(i)
|
| 91 |
+
self._rows_by_source = dict(rows)
|
| 92 |
return True
|
| 93 |
|
| 94 |
+
def search_by_source(
|
| 95 |
+
self, query_vector, source: str, k: int = 5
|
| 96 |
+
) -> list[tuple[Chunk, float]]:
|
| 97 |
+
"""Return up to ``k`` nearest (chunk, score) pairs from one ``source`` only.
|
| 98 |
+
|
| 99 |
+
Searches just that source's rows, so each source contributes its own
|
| 100 |
+
top-k regardless of the other's score distribution. Scores are cosine
|
| 101 |
+
(vectors are L2-normalized at build time). Requires ``load()`` first.
|
| 102 |
+
"""
|
| 103 |
+
rows = self._rows_by_source.get(source, [])
|
| 104 |
+
if not rows or self._vectors is None:
|
| 105 |
+
return []
|
| 106 |
+
q = np.asarray(query_vector, dtype="float32")
|
| 107 |
+
scores = self._vectors[rows] @ q
|
| 108 |
+
top = np.argsort(-scores)[: min(k, len(rows))]
|
| 109 |
+
return [(self._chunks[rows[i]], float(scores[i])) for i in top]
|
| 110 |
+
|
| 111 |
def search(self, query_vector, k: int = 5) -> list[tuple[Chunk, float]]:
|
| 112 |
"""Return up to ``k`` nearest (chunk, score) pairs, best first."""
|
| 113 |
if self._index is None or not self._chunks:
|
src/backend/rag.py
CHANGED
|
@@ -23,6 +23,33 @@ from .schemas import ChatMessage, ChatResponse, Chunk
|
|
| 23 |
_store: VectorStore | None = None
|
| 24 |
_llm: LLMClient | None = None
|
| 25 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 26 |
|
| 27 |
def _get_store() -> VectorStore:
|
| 28 |
"""Load the persisted FAISS index once (lazy singleton).
|
|
@@ -113,16 +140,18 @@ def build_index() -> str:
|
|
| 113 |
|
| 114 |
@spaces.GPU(duration=180)
|
| 115 |
def _retrieve_and_stream(query: str, lang: str, k: int):
|
| 116 |
-
"""Embed query, retrieve
|
| 117 |
|
| 118 |
-
|
| 119 |
-
|
| 120 |
-
|
| 121 |
-
|
|
|
|
| 122 |
"""
|
| 123 |
qvec = Embedder.get().embed_query(query)
|
| 124 |
-
|
| 125 |
-
|
|
|
|
| 126 |
prompt = build_grounded_prompt(query, chunks, lang)
|
| 127 |
acc = ""
|
| 128 |
for piece in _get_llm().stream(prompt):
|
|
@@ -130,7 +159,7 @@ def _retrieve_and_stream(query: str, lang: str, k: int):
|
|
| 130 |
yield acc, chunks
|
| 131 |
|
| 132 |
|
| 133 |
-
def answer_stream(messages: list[ChatMessage], lang: str = "en", k: int =
|
| 134 |
"""Stream grounded answers as ``ChatResponse`` snapshots (growing reply)."""
|
| 135 |
query = next((m.content for m in reversed(messages) if m.role == "user"), "")
|
| 136 |
if not query.strip():
|
|
@@ -140,7 +169,7 @@ def answer_stream(messages: list[ChatMessage], lang: str = "en", k: int = 10):
|
|
| 140 |
yield ChatResponse(reply=acc, citations=_citations(chunks))
|
| 141 |
|
| 142 |
|
| 143 |
-
def answer(messages: list[ChatMessage], lang: str = "en", k: int =
|
| 144 |
"""Non-streaming answer (used by the structured ``/chat`` API)."""
|
| 145 |
resp = ChatResponse(reply="", citations=[])
|
| 146 |
for resp in answer_stream(messages, lang=lang, k=k):
|
|
|
|
| 23 |
_store: VectorStore | None = None
|
| 24 |
_llm: LLMClient | None = None
|
| 25 |
|
| 26 |
+
# Retrieval is run per knowledge source then fused, so each source contributes
|
| 27 |
+
# regardless of the other's similarity-score scale (directive = long English
|
| 28 |
+
# legal text; IDA = short salary summaries). Sources to retrieve from.
|
| 29 |
+
_SOURCES = ("directive", "lonstatistik")
|
| 30 |
+
K_PER_SOURCE = 5 # top-k pulled from each source before fusion
|
| 31 |
+
K_FINAL = 6 # chunks kept after fusion (passed to the LLM)
|
| 32 |
+
RRF_C = 60 # Reciprocal Rank Fusion constant (standard default)
|
| 33 |
+
|
| 34 |
+
|
| 35 |
+
def _rrf_fuse(
|
| 36 |
+
ranked_lists: list[list[tuple[Chunk, float]]], k_final: int, c: int = RRF_C
|
| 37 |
+
) -> list[Chunk]:
|
| 38 |
+
"""Fuse per-source ranked lists with Reciprocal Rank Fusion.
|
| 39 |
+
|
| 40 |
+
Rank-based and score-scale-agnostic: each chunk scores ``sum 1/(c + rank)``
|
| 41 |
+
across the lists it appears in. Our sources are disjoint, so this cleanly
|
| 42 |
+
interleaves the lists by rank. Returns the top ``k_final`` chunks.
|
| 43 |
+
"""
|
| 44 |
+
score: dict[str, float] = {}
|
| 45 |
+
keep: dict[str, Chunk] = {}
|
| 46 |
+
for lst in ranked_lists:
|
| 47 |
+
for rank, (chunk, _s) in enumerate(lst):
|
| 48 |
+
score[chunk.id] = score.get(chunk.id, 0.0) + 1.0 / (c + rank + 1)
|
| 49 |
+
keep[chunk.id] = chunk
|
| 50 |
+
ordered = sorted(keep, key=lambda cid: score[cid], reverse=True)
|
| 51 |
+
return [keep[cid] for cid in ordered[:k_final]]
|
| 52 |
+
|
| 53 |
|
| 54 |
def _get_store() -> VectorStore:
|
| 55 |
"""Load the persisted FAISS index once (lazy singleton).
|
|
|
|
| 140 |
|
| 141 |
@spaces.GPU(duration=180)
|
| 142 |
def _retrieve_and_stream(query: str, lang: str, k: int):
|
| 143 |
+
"""Embed query, retrieve per source + fuse, and stream generation.
|
| 144 |
|
| 145 |
+
One GPU allocation. Retrieves the top ``K_PER_SOURCE`` from each knowledge
|
| 146 |
+
source independently, fuses with Reciprocal Rank Fusion, and keeps the top
|
| 147 |
+
``k`` fused chunks — so both the Directive and IDA statistics can surface
|
| 148 |
+
even though their similarity scores live on different scales. Yields
|
| 149 |
+
``(partial_reply, chunks)`` as tokens arrive so the UI renders incrementally.
|
| 150 |
"""
|
| 151 |
qvec = Embedder.get().embed_query(query)
|
| 152 |
+
store = _get_store()
|
| 153 |
+
ranked = [store.search_by_source(qvec, s, K_PER_SOURCE) for s in _SOURCES]
|
| 154 |
+
chunks = _rrf_fuse(ranked, k_final=k)
|
| 155 |
prompt = build_grounded_prompt(query, chunks, lang)
|
| 156 |
acc = ""
|
| 157 |
for piece in _get_llm().stream(prompt):
|
|
|
|
| 159 |
yield acc, chunks
|
| 160 |
|
| 161 |
|
| 162 |
+
def answer_stream(messages: list[ChatMessage], lang: str = "en", k: int = K_FINAL):
|
| 163 |
"""Stream grounded answers as ``ChatResponse`` snapshots (growing reply)."""
|
| 164 |
query = next((m.content for m in reversed(messages) if m.role == "user"), "")
|
| 165 |
if not query.strip():
|
|
|
|
| 169 |
yield ChatResponse(reply=acc, citations=_citations(chunks))
|
| 170 |
|
| 171 |
|
| 172 |
+
def answer(messages: list[ChatMessage], lang: str = "en", k: int = K_FINAL) -> ChatResponse:
|
| 173 |
"""Non-streaming answer (used by the structured ``/chat`` API)."""
|
| 174 |
resp = ChatResponse(reply="", citations=[])
|
| 175 |
for resp in answer_stream(messages, lang=lang, k=k):
|
src/backend/schemas.py
CHANGED
|
@@ -27,6 +27,11 @@ class Chunk(BaseModel):
|
|
| 27 |
article: str = Field(default="", description="Directive citation label, e.g. 'Art. 7'.")
|
| 28 |
page: int | None = Field(default=None, description="Source page (lønstatistik).")
|
| 29 |
citation: str = Field(default="", description="Human-readable citation label.")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 30 |
|
| 31 |
|
| 32 |
class ParseResult(BaseModel):
|
|
|
|
| 27 |
article: str = Field(default="", description="Directive citation label, e.g. 'Art. 7'.")
|
| 28 |
page: int | None = Field(default=None, description="Source page (lønstatistik).")
|
| 29 |
citation: str = Field(default="", description="Human-readable citation label.")
|
| 30 |
+
metadata: dict = Field(
|
| 31 |
+
default_factory=dict,
|
| 32 |
+
description="Structured, non-embedded source metadata for filtering/display "
|
| 33 |
+
"(e.g. IDA sector/category/measures, directive article title).",
|
| 34 |
+
)
|
| 35 |
|
| 36 |
|
| 37 |
class ParseResult(BaseModel):
|
src/frontend/README.md
CHANGED
|
@@ -8,8 +8,8 @@ selector (en/da/fr) and inline source citations.
|
|
| 8 |
The Scandinavian/sage style is implemented through Gradio's own theming engine
|
| 9 |
(`gr.themes.Soft(...).set(...)` + a small CSS block for the Bricolage Grotesque
|
| 10 |
/ Inter fonts and message-bubble colors) — **not** custom HTML. Theme and CSS
|
| 11 |
-
are passed to `
|
| 12 |
-
constructor).
|
| 13 |
|
| 14 |
## Start it
|
| 15 |
|
|
|
|
| 8 |
The Scandinavian/sage style is implemented through Gradio's own theming engine
|
| 9 |
(`gr.themes.Soft(...).set(...)` + a small CSS block for the Bricolage Grotesque
|
| 10 |
/ Inter fonts and message-bubble colors) — **not** custom HTML. Theme and CSS
|
| 11 |
+
are passed to `demo.launch(...)` in `main.py` (Gradio 6 moved them off the
|
| 12 |
+
`Blocks` constructor).
|
| 13 |
|
| 14 |
## Start it
|
| 15 |
|