Skip to content
Merged
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -345,6 +345,10 @@ Frame shapes:
{"type": "error", "message": "..."}
```

For `/reindex` the `upserting` frames carry `current` as a monotonically increasing count of files
resolved — indexed, skipped as unchanged, or dropped by a fetch/parse failure — so it always ends at
`total`. Files are indexed concurrently, so frames are emitted per batch rather than per file.

For `/reindex-history` the `phase` value is `discovery|embedding|upserting` and the `done` result is
`{"new": int, "skipped": int, "diff_updated": int}`.

Expand Down
33 changes: 21 additions & 12 deletions docs/docs/ingestion.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ Ingestion is managed by `IndexPipeline` (`server/indexer/pipeline.py`). For each

1. Discovers all indexable files from GitHub
2. Skips files whose content hasn't changed since the last index run
3. Downloads changed file content, parses it into `CodeSymbol` entries, generates dense and sparse embeddings, and upserts them into Qdrant
3. Downloads changed file content concurrently, parses it into `CodeSymbol` entries, generates dense and sparse embeddings in batches that span multiple files, and upserts them into Qdrant
4. Removes index entries for files that have been deleted from the repository

The pipeline is triggered via the `/reindex` HTTP endpoint (streaming NDJSON progress) or the `index_all` MCP admin tool.
Expand Down Expand Up @@ -76,10 +76,14 @@ content = await fetch_blob_content(

Fetching by blob SHA is more efficient than path-based fetching during indexing: the SHA is already known from the tree response, and the blob API is a direct content lookup with no ref resolution overhead.

Fetches run concurrently, bounded by a semaphore (`_FETCH_CONCURRENCY`, default 10 — the same bound `github_source.py` uses for tree walks and commit diffs). Each fetched file is parsed and handed to a bounded queue (`_PARSED_QUEUE_SIZE`), so downloading continues while the previous group of symbols is being embedded. Every request goes through `_gh_get`, which retries rate limits, 5xx, and transport errors. A fetch that still fails logs an error and drops only that file; its existing index entries are preserved.

### 4. Parsing

`parse_file(content, stored_path)` dispatches to the language-specific parser via the registry. The result is a `list[CodeSymbol]` — one entry per indexable symbol (class, method, function, interface, etc.).

Parsing runs on the event loop rather than in a thread pool, deliberately: `registry.py` builds one parser instance per language and shares it across every file, and tree-sitter `Parser` objects are not safe for concurrent use. Moving `parse_file` to a worker thread would be a data race.

If a file produces no symbols (empty, unsupported format, or parse failure), any existing index entries for that file are cleaned up and the file is skipped.

### 5. Embedding Text Construction
Expand Down Expand Up @@ -121,25 +125,32 @@ This text is then pre-processed by `split_code_identifiers` (see [sparse-vectors

### 6. Embedding

Both embedding calls are made sequentially per file batch:
Symbols are accumulated **across files** until they fill the provider's batch size (`EmbeddingProvider.batch_size` — 128 for Voyage/OpenAI/Jina's hosted API, 32 for self-hosted Jina TEI and Ollama), or until the batch reaches `_MAX_BATCH_CHARS`. A single file usually yields only a handful of symbols, so batching across files is what keeps requests full instead of sending one tiny request per file.

The dense and sparse embeds for a batch run concurrently:

```python
dense_vectors = await self._embedder.embed_batch(texts_dense)
sparse_vectors = await self._sparse_embedder.embed_batch(texts_sparse)
dense, sparse = await asyncio.gather(
self._embedder.embed_batch(dense_texts),
self._sparse_embedder.embed_batch(sparse_texts),
)
```

If either call raises an exception, the file is skipped and existing index entries are preserved until the next successful run.
Dense is network-bound and sparse runs in a thread executor, so the two overlap for free.

If a batch call raises, the pipeline retries that batch **file by file**, so a single unembeddable file costs one extra round-trip rather than dropping every file batched alongside it. A file that still fails is skipped, and its existing index entries are preserved until the next successful run.

### 7. Upsert

Before inserting new vectors, all existing entries for the file are removed:
Writes stay scoped to one file at a time, even though embedding is batched. New vectors are upserted *before* stale ones are deleted, so the file never has zero indexed symbols:

```python
await self._store.delete_by_file(svc.name, stored_path)
await self._store.upsert_chunks(payloads, dense_vectors, sparse_vectors)
previous_ids = await self._store.get_point_ids_by_file(service_name, stored_path)
new_ids = await self._store.upsert_chunks(payloads, dense, sparse)
await self._store.delete_by_ids(list(previous_ids - set(new_ids)))
```

This ensures clean replacement when symbols are added, removed, or renamed within a file. Each point's ID is a deterministic `uuid5` derived from `service:file_path:symbol_name:start_line`, so symbols moving to a new line produce new IDs (handled correctly by the delete-first approach).
Each point's ID is a deterministic `uuid5` derived from `service:file_path:symbol_name:start_line`, so symbols moving to a new line produce new IDs and the old ones fall out as stale. Because the ID includes `file_path`, two different files can never produce colliding IDs — which is what makes batching across files safe.

Each point carries a payload with 20+ fields (see Data Model below).

Expand Down Expand Up @@ -197,13 +208,11 @@ All `CodeSymbol` fields are stored verbatim, plus:

## Observations

**Sequential embedding calls** — `embed_batch` for dense and `embed_batch` for sparse are awaited sequentially. They are independent operations targeting different providers; wrapping them in `asyncio.gather` would reduce per-file embedding latency by ~50%.

**No embedding retry** — a transient API error on either embedding call causes the file to be silently skipped, leaving its existing index stale indefinitely. There is no exponential backoff or retry queue. Reindexing requires either a force reindex or waiting for the file's content to change.

**BM25 text still omits some dense-only metadata** — `_build_bm25_text` folds in name, package, annotations, and HTTP method/route, but the dense preamble's service name, language, and symbol-type phrasing (e.g. "Java method") are still dense-only. A BM25 query for "Python method" will not match unless the word "Python" or "method" appears elsewhere in the folded-in fields or the source code itself.

**delete-before-upsert gap** — The pipeline deletes all entries for a file before upserting the new ones. If the process is interrupted between delete and upsert, the file has no index entries. The next incremental run will redownload and reindex the file correctly — but until then, queries miss the file entirely.
**GitHub retries are bounded** — `_gh_get` retries rate limits (403/429, waiting for the window named by `Retry-After` / `X-RateLimit-Reset`, capped at 120s), 5xx, and transport errors, for `_GH_ATTEMPTS` attempts total. A failure that outlives those attempts surfaces as a per-file fetch error, leaving that file un-reindexed until the next run rather than failing the whole service.

**GitHub Trees truncation** — Very large repositories may have their tree response silently truncated by the GitHub API. The pipeline logs a warning but does not retry or paginate to recover the missing entries.

Expand Down
14 changes: 14 additions & 0 deletions server/embeddings/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,26 @@

from typing import Protocol, runtime_checkable

# Conservative ceiling for providers that do not declare their own. Matches the
# smallest batch size any bundled provider uses (jina TEI, ollama).
_DEFAULT_BATCH_SIZE = 32


@runtime_checkable
class EmbeddingProvider(Protocol):
@property
def dimensions(self) -> int: ...

@property
def batch_size(self) -> int:
"""Max texts the provider accepts per request.

Concrete default rather than `...` — every bundled provider subclasses
this Protocol, so an ellipsis body would silently return None for any
provider that forgot to override it.
"""
return _DEFAULT_BATCH_SIZE

async def embed_batch(self, texts: list[str]) -> list[list[float]]: ...

async def embed_query(self, text: str) -> list[float]: ...
4 changes: 4 additions & 0 deletions server/embeddings/jina.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,10 @@ def __init__(self) -> None:
def dimensions(self) -> int:
return self._dims

@property
def batch_size(self) -> int:
return _BATCH_SIZE

@staticmethod
def _extract(data) -> list[list[float]]:
# TEI returns a list of vectors directly
Expand Down
4 changes: 4 additions & 0 deletions server/embeddings/jina_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,10 @@ def __init__(self) -> None:
def dimensions(self) -> int:
return self._dims

@property
def batch_size(self) -> int:
return _BATCH_SIZE

async def embed_batch(self, texts: list[str]) -> list[list[float]]:
return await self._embed(texts, task="retrieval.passage")

Expand Down
4 changes: 4 additions & 0 deletions server/embeddings/ollama.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,10 @@ def __init__(self) -> None:
def dimensions(self) -> int:
return self._dims

@property
def batch_size(self) -> int:
return _BATCH_SIZE

async def embed_batch(self, texts: list[str]) -> list[list[float]]:
return await embed_in_batches(
texts,
Expand Down
4 changes: 4 additions & 0 deletions server/embeddings/openai.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,10 @@ def __init__(self) -> None:
def dimensions(self) -> int:
return self._dims

@property
def batch_size(self) -> int:
return _BATCH_SIZE

def _make_body(self, inputs: list[str]) -> dict:
body: dict = {"model": self._model, "input": inputs}
if self._dims_override is not None:
Expand Down
4 changes: 4 additions & 0 deletions server/embeddings/voyage.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,10 @@ def __init__(self) -> None:
def dimensions(self) -> int:
return self._dims

@property
def batch_size(self) -> int:
return _BATCH_SIZE

async def embed_batch(self, texts: list[str]) -> list[list[float]]:
return await self._embed(texts, input_type="document")

Expand Down
86 changes: 70 additions & 16 deletions server/indexer/github_source.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,15 @@
_DIFF_CONCURRENCY = 10
_TREE_WALK_CONCURRENCY = 10

# Attempts per GitHub GET, shared by rate limits, 5xx and transport errors.
_GH_ATTEMPTS = 3
# Backoff for 5xx and transport errors. Rate limits ignore this and wait for the
# window named by Retry-After / X-RateLimit-Reset instead.
_GH_BACKOFF_DELAYS = (1.0, 5.0)
# Transient network failures worth another attempt. httpx.TimeoutException is a
# subclass of TransportError, but naming it keeps the intent obvious.
_RETRYABLE_EXCEPTIONS = (httpx.TransportError, httpx.TimeoutException)

logger = logging.getLogger(__name__)


Expand Down Expand Up @@ -83,23 +92,68 @@ async def _gh_get(
params: dict | None = None,
timeout: float = 30.0,
) -> Any:
"""GET a GitHub API URL, retrying up to 3 times on rate-limit responses (403/429)."""
"""GET a GitHub API URL, retrying rate limits (403/429), 5xx and transport errors.

Rate limits wait for the window GitHub names via ``Retry-After`` /
``X-RateLimit-Reset`` (capped at 120s); 5xx and transport errors use a short
fixed backoff, since they carry no such hint and usually clear immediately.
"""
headers = _auth_headers(token)
for _ in range(3):
r = await client.get(url, headers=headers, params=params, timeout=timeout)
if r.status_code not in (403, 429):
r.raise_for_status()
return r.json()
reset_ts = float(r.headers.get("X-RateLimit-Reset", 0))
retry_after = float(r.headers.get("Retry-After", 60))
now = time.time()
wait = min(max(retry_after, reset_ts - now if reset_ts > now else 0.0), 120.0)
logger.warning(
"GitHub rate-limited (HTTP %d) — retrying in %.0fs", r.status_code, wait
)
await asyncio.sleep(wait)
r.raise_for_status() # raise final rate-limit error after exhausting retries
return r.json() # unreachable
r: httpx.Response | None = None

for attempt in range(_GH_ATTEMPTS):
last_attempt = attempt == _GH_ATTEMPTS - 1
try:
r = await client.get(url, headers=headers, params=params, timeout=timeout)
except _RETRYABLE_EXCEPTIONS as exc:
if last_attempt:
raise
wait = _GH_BACKOFF_DELAYS[min(attempt, len(_GH_BACKOFF_DELAYS) - 1)]
logger.warning(
"GitHub request failed (%s: %s) — retrying in %.0fs (attempt %d/%d)",
type(exc).__name__,
exc,
wait,
attempt + 1,
_GH_ATTEMPTS,
)
await asyncio.sleep(wait)
continue

if r.status_code in (403, 429):
if last_attempt:
break
reset_ts = float(r.headers.get("X-RateLimit-Reset", 0))
retry_after = float(r.headers.get("Retry-After", 60))
now = time.time()
wait = min(
max(retry_after, reset_ts - now if reset_ts > now else 0.0), 120.0
)
logger.warning(
"GitHub rate-limited (HTTP %d) — retrying in %.0fs", r.status_code, wait
)
await asyncio.sleep(wait)
continue

if r.status_code >= 500:
if last_attempt:
break
wait = _GH_BACKOFF_DELAYS[min(attempt, len(_GH_BACKOFF_DELAYS) - 1)]
logger.warning(
"GitHub server error (%d) — retrying in %.0fs (attempt %d/%d)",
r.status_code,
wait,
attempt + 1,
_GH_ATTEMPTS,
)
await asyncio.sleep(wait)
continue

break

assert r is not None # a transport error on the last attempt re-raises
r.raise_for_status() # surfaces the final rate-limit / 5xx error
return r.json()


def _filter_tree_blobs(
Expand Down
Loading
Loading