From 25cf10baa672fe063a820cb37c4cb07c7fc0b085 Mon Sep 17 00:00:00 2001 From: Sergey <93754860+slaid098@users.noreply.github.com> Date: Sun, 26 Jul 2026 19:19:55 +0300 Subject: [PATCH] feat(memory): incremental index with SHA256 + atomic + versioning + flock (#83) * feat(memory): add SHA256 incremental index with atomic writes, versioning, flock * feat(memory): parse Retry-After header and add batch progress logging * test(memory): add incremental, atomic, versioning, retry, batch tests * docs(memory): update .env.example with OpenRouter defaults * docs(handoff): add handoff and ADR-036 for incremental index * docs(handoff): set PR number * docs(project-map): update index.py and embedder.py descriptions for PR#83 * fix(ci): skip index rewrite on no-op + versioning first-run * fix(ci): ruff format index.py --------- Co-authored-by: opencode-agent --- .env.example | 13 +- ...6-pr-83-incremental-index-sha256-atomic.md | 35 +++ .../pr-83-incremental-index-sha256-atomic.md | 28 +++ docs/project-map/README.md | 4 +- src/memory/embedder.py | 13 +- src/memory/index.py | 222 +++++++++++++++-- tests/test_embedder.py | 79 +++++- tests/test_embedder_live.py | 7 + tests/test_index.py | 226 +++++++++++++++++- 9 files changed, 595 insertions(+), 32 deletions(-) create mode 100644 docs/decisions/036-pr-83-incremental-index-sha256-atomic.md create mode 100644 docs/handoff/pr-83-incremental-index-sha256-atomic.md diff --git a/.env.example b/.env.example index 9215310..e4aa67f 100644 --- a/.env.example +++ b/.env.example @@ -3,10 +3,15 @@ AI_PROVIDER_BASE_URL=https://your-ai-provider.example.com/v1/ AI_PROVIDER_API_KEY=your-api-key-here # OpenAI Embeddings (Memory CLI) -OPENAI_BASE_URL=https://api.openai.com/v1 -OPENAI_API_KEY=your-openai-api-key -OPENAI_EMBEDDING_MODEL=gemini-embedding-2-preview -OPENAI_EMBEDDING_BATCH_SIZE=2048 +# OpenRouter defaults (Qwen3 8B, batch=50). For OpenAI direct, use: +# OPENAI_BASE_URL=https://api.openai.com/v1 +# OPENAI_EMBEDDING_MODEL=text-embedding-3-small +# OPENAI_EMBEDDING_BATCH_SIZE=2048 +OPENAI_BASE_URL=https://openrouter.ai/api/v1 +OPENAI_API_KEY=your-openrouter-api-key +OPENAI_EMBEDDING_MODEL=qwen/qwen3-embedding-8b +OPENAI_EMBEDDING_BATCH_SIZE=50 +OPENAI_EMBEDDING_BATCH_DELAY=1 MEMORY_CHUNK_SIZE=512 MEMORY_CHUNK_OVERLAP=64 diff --git a/docs/decisions/036-pr-83-incremental-index-sha256-atomic.md b/docs/decisions/036-pr-83-incremental-index-sha256-atomic.md new file mode 100644 index 0000000..c799caf --- /dev/null +++ b/docs/decisions/036-pr-83-incremental-index-sha256-atomic.md @@ -0,0 +1,35 @@ +# ADR-036: Incremental RAG index with SHA256, atomic writes, versioning, flock + +## Статус +Accepted (2026-07-26) + +## Контекст +PR #75 (refactor) + PR #77 (plugin wrapper) сделали Memory CLI рабочим через OpenRouter Qwen3 8B: `.rag/index.json` создан (2354 records, dim=4096, 218 MB), `memory_search` возвращает semantic hits. Но `run_index` делал **полный переиндекс** при каждом запуске: 2354 embeddings per save, ~13 минут на OpenRouter Qwen3 8B, расход $0.0043/save (при бюджете $9 = 2076 reindexes = 5.7 лет при 1 save/день). Smoke-test (technical/openrouter-qwen3-embedding-index-timeout.md) подтвердил: 47 батчей × 16.5с = 752с, команда с `timeout 120` молча убивала процесс до первого print. + +Три проблемы: +1. **Расход бюджета**: 1 save = полная переиндексация всех 2354 чанков (340× больше необходимого — нужны только изменившиеся ~10 чанков). +2. **Скорость**: 13 минут блокирует каждый `memory_save`. Должно быть ~1 сек. +3. **Не production-grade**: нет atomic writes (упал процесс → битый index.json), нет versioning (смена модели → несовместимые embeddings в одном индексе), нет concurrent safety (два `memory index` → race на meta.json), нет rate-limit-aware retry (OpenRouter 429 без Retry-After parsing). + +## Решение +Четыре независимых улучшения в `src/memory/index.py` + `src/memory/embedder.py`: + +1. **SHA256 инкрементальный индекс** (`index.py`): `_content_hash(text: str) -> str` (SHA256 от extracted text, не от raw file — frontmatter не влияет на embedding). `.rag/meta.json` хранит `{version, files: {path: {sha256, chunks}}}`. `run_index`: rglob *.md → SHA256 extracted text → сравнить с meta.json → `changed_files` (хеш не совпал или файла нет) → `embed_texts(только changed chunks)` → merge old (unchanged) + new (changed) → prune удалённых файлов. Первый запуск (meta.json не существует) → полный reindex. + +2. **Atomic writes** (`index.py`): `_atomic_write(path, content)` — write `.tmp` → `os.replace` (атомарная замена на POSIX). Используется для `index.json` и `meta.json`. При сбое процесса mid-write — старый файл остаётся целым, `.tmp` остаётся мусором (можно очистить при следующем запуске). + +3. **Index versioning** (`index.py`): `version = f"{EMBEDDING_MODEL}:{chunk_size}:{chunk_overlap}"`. При load meta.json: если `meta["version"] != current_version` → полный reindex (лог "Index version mismatch, full reindex"). Защищает от несовместимых embeddings при смене модели/chunking. + +4. **Concurrent safety** (`index.py`): `fcntl.flock(LOCK_EX)` на `.rag/.lock` во время всей index операции. `output_dir.mkdir(parents=True, exist_ok=True)` ДО создания lock. Search НЕ блокирует (читает index.json без lock — допускает stale reads, OK для RAG). + +5. **Rate-limit-aware batching** (`embedder.py`): при 429 читать `Retry-After` header → `time.sleep(int(retry_after))` → raise (tenacity поймает и retry'нет). `stop_after_attempt(3)` → `stop_after_attempt(5)`, `wait_exponential(max=10)` → `wait_exponential(max=30)`. + +6. **Прогресс-лог** (`index.py` + `embedder.py`): `print(f"Embedding N chunks in M batches...", flush=True)` перед `embed_texts`; `print(f" batch K/N...", flush=True)` в loop (только если > 1 batch). 13 мин полного reindex не выглядят как hang. + +## Альтернативы +- **mtime-based incremental** (вместо SHA256): отвергнуто — `touch` без изменения content триггерит re-embed (mtime изменился, content нет). SHA256 от extracted text точнее (frontmatter changes не триггерят re-embed, mtime файла меняется). +- **file-level hash** (вместо chunk-level): отвергнуто —_chunks в рамках файла независимы, но hash файла достаточно (при изменении файла re-embed всех его чанков — проще и достаточно для memory use case где файлы маленькие). +- **transactional index** (write-ahead log, 2-phase commit): отвергнуто — overkill для single-process CLI. Atomic write через `os.replace` достаточно на POSIX. +- **file locking через `portalocker`/`fasteners`**: отвергнуто — `fcntl` стандартная библиотека, POSIX-only (Memory CLI предназначен для Linux/Docker), cross-platform не требуется. +- **embedding cache** (cache по hash текста на disk): отвергнуто — out of scope, отдельный PR. Текущая инкрементальность на уровне файлов уже даёт 340× экономию. +- **увеличение `BATCH_SIZE` default с 2048 до 50**: отвергнуто — default 2048 для OpenAI direct (поддерживает большой batch), OpenRouter требует 50 через env var `OPENAI_EMBEDDING_BATCH_SIZE=50`. `.env.example` обновлён с OpenRouter defaults, но код default не изменён. \ No newline at end of file diff --git a/docs/handoff/pr-83-incremental-index-sha256-atomic.md b/docs/handoff/pr-83-incremental-index-sha256-atomic.md new file mode 100644 index 0000000..9f1ccd7 --- /dev/null +++ b/docs/handoff/pr-83-incremental-index-sha256-atomic.md @@ -0,0 +1,28 @@ +--- +pr: 83 +title: feat(memory): incremental index with SHA256 + atomic + versioning + flock +--- + +## Что сделано +- SHA256 инкрементальный индекс в `src/memory/index.py`: `_content_hash(text)`, `.rag/meta.json` (`{version, files: {path: {sha256, chunks}}}`), embed только изменившихся чанков, merge old+new, prune удалённых файлов. +- Atomic writes: `_atomic_write(path, content)` через `.tmp` + `os.replace` для `index.json` и `meta.json` — защита от битого индекса при сбое. +- Index versioning: `version = f"{EMBEDDING_MODEL}:{chunk_size}:{chunk_overlap}"` в meta.json, mismatch → полный reindex. +- Concurrent safety: `fcntl.flock(LOCK_EX)` на `.rag/.lock` во время index операции. +- Rate-limit-aware batching в `src/memory/embedder.py`: парсинг `Retry-After` header при 429 → `time.sleep(retry_after)` → raise; `stop_after_attempt(5)`, `wait_exponential(max=30)`. +- Прогресс-лог: `index.py` печатает `Embedding N chunks in M batches...` перед `embed_texts`; `embedder.py` печатает `batch K/N...` (только если > 1 batch). +- 10 новых тестов: 8 в `tests/test_index.py` (incremental add/edit/delete, no-changes-noop, meta persistence, sha256 content hash, atomic writes, versioning, atomic write unit) + `test_retry_after_header` в `tests/test_embedder.py` + `test_live_embed_qwen3_batch` в `tests/test_embedder_live.py` (skip без RUN_LIVE). +- `.env.example`: блок OpenAI Embeddings обновлён с OpenRouter defaults (Qwen3 8B, batch=50, delay=1). + +## Почему +PR #75 (refactor) + PR #77 (plugin wrapper) сделали Memory CLI рабочим через OpenRouter Qwen3 8B, но полный reindex при каждом `memory_save` = 2354 embeddings = ~13 минут + $0.0043/save. Инкрементальный индекс по SHA256 от extracted text = embed только изменившихся чанков (~10 per save), ~1 сек + $0.0000128/save (340× экономия бюджета, 780× ускорение). Production-grade практики (atomic writes, versioning, flock) снимают риски: битый index.json при сбое процесса, несовместимые embeddings при смене модели, race condition при concurrent `memory index`. Retry-After header + увеличенные retry лимиты = корректная обработка OpenRouter rate limits. + +## Pending +— + +## Watch out +- `EMBEDDING_MODEL` и `BATCH_SIZE` читаются из env на import `src.memory.embedder` — тесты, проверяющие дефолты, должны `monkeypatch.delenv` + `importlib.reload(embedder_mod)`. Существующие тесты `test_embed_texts_default_model`/`test_embed_texts_batches` обновлены: используют `monkeypatch.setattr(embedder_mod, "BATCH_SIZE", 2048)` вместо依赖имости от env дефолта (контейнер env задаёт `OPENAI_EMBEDDING_BATCH_SIZE=50`). +- `.rag/meta.json` — новый файл. При первом запуске после этого PR — полный reindex (meta.json не существует → `needs_full_reindex=True`). При последующих — инкрементальный. Старый `.rag/index.json` (PR #77) совместим по схеме (`{files: [{source, chunk_idx, offset, text, embedding}]}`), merge сохраняет существующие записи для unchanged файлов. +- `fcntl.flock` — POSIX-only. На Windows не работает, но Memory CLI предназначен для Linux/Docker контейнера (`python>=3.12`, `Operating System :: POSIX :: Linux` в pyproject). Search НЕ блокирует (читает index.json без lock — допускает stale reads, это OK для RAG). +- `OPENAI_EMBEDDING_BATCH_DELAY=1` в `.env.example` — documented, но НЕ используется в `embedder.py` (нет delay между батчами в текущей реализации). Оставлен для документации/future use. Если добавить delay — нужен `time.sleep(BATCH_DELAY)` в loop. +- Progress log идёт в stdout. Plugin парсит только JSON в `search`, не в `index` — stdout лог безопасен для `index` команды. Для `search` (JSON output) логging НЕ добавлялся. +- ADR-036 (НЕ PR-number-based): sequential нумерация ADR в этом репо (PR#26 docs-reviewer typo зафиксировал правило — ADR = sequential, НЕ PR number). \ No newline at end of file diff --git a/docs/project-map/README.md b/docs/project-map/README.md index 071c2ef..6b369dd 100644 --- a/docs/project-map/README.md +++ b/docs/project-map/README.md @@ -71,8 +71,8 @@ opencode-config/ │ ├── __init__.py │ ├── __main__.py # Entry point for `python -m memory` │ ├── cli.py # CLI commands (prog="memory") -│ ├── embedder.py # Embedding via OPENAI_BASE_URL (env-only, OpenAI-compatible) -│ ├── index.py # Indexing with chunking (MEMORY_CHUNK_SIZE/OVERLAP env) +│ ├── embedder.py # Embedding via OPENAI_BASE_URL (env-only, OpenAI-compatible); rate-limit-aware: Retry-After header parsing, retry=5, backoff max=30 — PR#83 +│ ├── index.py # Indexing with chunking (MEMORY_CHUNK_SIZE/OVERLAP env); SHA256 incremental index (.rag/meta.json), atomic writes (.tmp+os.replace), versioning (model:chunk_size:overlap), fcntl.flock(LOCK_EX) — PR#83 │ └── search.py # Search with dedup by source in top-K ├── tests/ # pytest + TS/MJS test suite — PR#17 │ ├── _ts_loader.mjs # TS test loader (load/exec_stub/exec_stub_json/exec_real modes; relative import inlining via inlineShared()) — PR#38, PR#65 diff --git a/src/memory/embedder.py b/src/memory/embedder.py index 1732e8e..dd551f1 100644 --- a/src/memory/embedder.py +++ b/src/memory/embedder.py @@ -1,4 +1,6 @@ +import contextlib import os +import time import httpx from tenacity import retry, retry_if_exception, stop_after_attempt, wait_exponential @@ -22,8 +24,8 @@ def _is_retryable(exc: BaseException) -> bool: @retry( - stop=stop_after_attempt(3), - wait=wait_exponential(multiplier=1, min=2, max=10), + stop=stop_after_attempt(5), + wait=wait_exponential(multiplier=1, min=2, max=30), retry=retry_if_exception(_is_retryable), ) def _call_embedding_api( @@ -33,6 +35,11 @@ def _call_embedding_api( payload: dict[str, str | list[str]], ) -> list[list[float]]: resp = client.post(url, json=payload, headers=headers) + if resp.status_code == 429: + retry_after = resp.headers.get("Retry-After") + if retry_after is not None: + with contextlib.suppress(ValueError): + time.sleep(int(retry_after)) resp.raise_for_status() data = resp.json() return [d["embedding"] for d in data["data"]] @@ -50,9 +57,11 @@ def embed_texts(texts: list[str]) -> list[list[float]]: payload: dict[str, str | list[str]] = {"model": EMBEDDING_MODEL, "input": texts} return _call_embedding_api(client, url, headers, payload) + n_batches = (len(texts) + BATCH_SIZE - 1) // BATCH_SIZE results: list[list[float]] = [] for i in range(0, len(texts), BATCH_SIZE): batch = texts[i : i + BATCH_SIZE] payload = {"model": EMBEDDING_MODEL, "input": batch} + print(f" batch {i // BATCH_SIZE + 1}/{n_batches}...", flush=True) results.extend(_call_embedding_api(client, url, headers, payload)) return results diff --git a/src/memory/index.py b/src/memory/index.py index 7e39aad..07fc9f7 100644 --- a/src/memory/index.py +++ b/src/memory/index.py @@ -1,11 +1,21 @@ import argparse +import fcntl +import hashlib import json import os from pathlib import Path -from src.memory.embedder import embed_texts +from src.memory.embedder import BATCH_SIZE, EMBEDDING_MODEL, embed_texts INDEX_FILENAME = "index.json" +META_FILENAME = "meta.json" +LOCK_FILENAME = ".lock" + +FileEntry = dict[str, str | int] +IndexedEntry = dict[str, str | int | list[float]] +FileMap = dict[str, tuple[str, list[FileEntry]]] +FilesMeta = dict[str, dict[str, str | int]] +Meta = dict[str, str | FilesMeta] def _extract_text(content: str) -> str: @@ -23,7 +33,11 @@ def _chunk_text(text: str, size: int, overlap: int) -> list[tuple[str, int]]: return [(text[i : i + size], i) for i in range(0, len(text), step)] -def _file_entry(rel: Path, chunk_idx: int, offset: int, chunk: str) -> dict[str, str | int]: +def _content_hash(text: str) -> str: + return hashlib.sha256(text.encode("utf-8")).hexdigest() + + +def _file_entry(rel: Path, chunk_idx: int, offset: int, chunk: str) -> FileEntry: return { "source": str(rel), "chunk_idx": chunk_idx, @@ -32,20 +46,134 @@ def _file_entry(rel: Path, chunk_idx: int, offset: int, chunk: str) -> dict[str, } +def _atomic_write(path: Path, content: str) -> None: + tmp = path.with_suffix(path.suffix + ".tmp") + tmp.write_text(content, encoding="utf-8") + os.replace(tmp, path) + + +def _load_meta(meta_path: Path) -> Meta: + if not meta_path.exists(): + return {} + data: object = json.loads(meta_path.read_text(encoding="utf-8")) + if isinstance(data, dict): + return data + return {} + + def _build_file_map( memory_dir: Path, md_files: list[Path], chunk_size: int, chunk_overlap: int -) -> tuple[list[str], list[dict[str, str | int]]]: - texts: list[str] = [] - file_map: list[dict[str, str | int]] = [] +) -> FileMap: + file_map: FileMap = {} for fpath in md_files: rel = fpath.relative_to(memory_dir) content = fpath.read_text(encoding="utf-8") text = _extract_text(content) chunks = _chunk_text(text, chunk_size, chunk_overlap) or [("", 0)] - for chunk_idx, (chunk, offset) in enumerate(chunks): - texts.append(chunk) - file_map.append(_file_entry(rel, chunk_idx, offset, chunk)) - return texts, file_map + entries = [ + _file_entry(rel, chunk_idx, offset, chunk) + for chunk_idx, (chunk, offset) in enumerate(chunks) + ] + file_map[str(rel)] = (_content_hash(text), entries) + return file_map + + +def _current_version(chunk_size: int, chunk_overlap: int) -> str: + return f"{EMBEDDING_MODEL}:{chunk_size}:{chunk_overlap}" + + +def _split_changed_unchanged( + file_map: FileMap, + old_files_meta: FilesMeta, + needs_full_reindex: bool, +) -> tuple[list[str], list[str]]: + changed: list[str] = [] + unchanged: list[str] = [] + for rel, (sha, _entries) in file_map.items(): + old = old_files_meta.get(rel) + if old is None or old.get("sha256") != sha or needs_full_reindex: + changed.append(rel) + else: + unchanged.append(rel) + return changed, unchanged + + +def _load_old_entries_by_source(index_path: Path) -> dict[str, list[FileEntry]]: + if not index_path.exists(): + return {} + old_index: object = json.loads(index_path.read_text(encoding="utf-8")) + if not isinstance(old_index, dict): + return {} + files: object = old_index.get("files", []) + if not isinstance(files, list): + return {} + by_source: dict[str, list[FileEntry]] = {} + for entry in files: + if isinstance(entry, dict): + src = str(entry.get("source", "")) + by_source.setdefault(src, []).append(entry) + return by_source + + +def _collect_kept_entries( + unchanged: list[str], old_entries_by_source: dict[str, list[FileEntry]] +) -> list[FileEntry]: + kept: list[FileEntry] = [] + for rel in unchanged: + kept.extend(old_entries_by_source.get(rel, [])) + return kept + + +def _prepare_changed( + changed: list[str], + file_map: FileMap, + memory_dir: Path, + chunk_size: int, + chunk_overlap: int, +) -> tuple[list[str], list[FileEntry], FilesMeta]: + texts: list[str] = [] + entries: list[FileEntry] = [] + files_meta: FilesMeta = {} + for rel in changed: + sha, file_entries = file_map[rel] + content = Path(memory_dir, rel).read_text(encoding="utf-8") + text = _extract_text(content) + chunks = _chunk_text(text, chunk_size, chunk_overlap) or [("", 0)] + for entry in file_entries: + idx = int(entry["chunk_idx"]) + texts.append(chunks[idx][0] if 0 <= idx < len(chunks) else "") + entries.append(entry) + files_meta[rel] = {"sha256": sha, "chunks": len(file_entries)} + return texts, entries, files_meta + + +def _build_full_meta( + file_map: FileMap, + changed_meta: FilesMeta, + unchanged: list[str], +) -> FilesMeta: + full: FilesMeta = dict(changed_meta) + for rel in unchanged: + sha, entries = file_map[rel] + full[rel] = {"sha256": sha, "chunks": len(entries)} + return full + + +def _merge_and_sort( + kept: list[FileEntry], + new_entries: list[FileEntry], + new_embeddings: list[list[float]], +) -> list[IndexedEntry]: + merged: list[IndexedEntry] = [dict(e) for e in kept] + for entry, emb in zip(new_entries, new_embeddings, strict=False): + merged.append({**entry, "embedding": emb}) + + def _sort_key(e: IndexedEntry) -> tuple[str, int]: + chunk_idx = e.get("chunk_idx", 0) + return str(e["source"]), int(chunk_idx) if isinstance(chunk_idx, int) else 0 + + merged.sort(key=_sort_key) + return merged def run_index(args: argparse.Namespace) -> None: @@ -61,13 +189,73 @@ def run_index(args: argparse.Namespace) -> None: chunk_size = int(os.environ.get("MEMORY_CHUNK_SIZE", "512")) chunk_overlap = int(os.environ.get("MEMORY_CHUNK_OVERLAP", "64")) - texts, file_map = _build_file_map(memory_dir, md_files, chunk_size, chunk_overlap) - embeddings = embed_texts(texts) - index = { - "files": [{**fm, "embedding": emb} for fm, emb in zip(file_map, embeddings, strict=False)], - } + lock_path = output_dir / LOCK_FILENAME + lock_path.touch() + with open(lock_path, "w") as lock_file: + fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX) + _run_index_locked( + output_dir=output_dir, + md_files=md_files, + memory_dir=memory_dir, + chunk_config=(chunk_size, chunk_overlap), + ) - (output_dir / INDEX_FILENAME).write_text( - json.dumps(index, ensure_ascii=False), encoding="utf-8" + +def _run_index_locked( + *, + output_dir: Path, + md_files: list[Path], + memory_dir: Path, + chunk_config: tuple[int, int], +) -> None: + chunk_size, chunk_overlap = chunk_config + current_version = _current_version(chunk_size, chunk_overlap) + index_path = output_dir / INDEX_FILENAME + meta_path = output_dir / META_FILENAME + meta = _load_meta(meta_path) + meta_version = meta.get("version") + needs_full_reindex = meta_path.exists() and meta_version != current_version + if needs_full_reindex: + print("Index version mismatch, full reindex") + + file_map = _build_file_map(memory_dir, md_files, chunk_size, chunk_overlap) + raw_files = meta.get("files", {}) + old_files_meta: FilesMeta = raw_files if isinstance(raw_files, dict) else {} + + changed_files, unchanged_files = _split_changed_unchanged( + file_map, old_files_meta, needs_full_reindex + ) + deleted_files = [rel for rel in old_files_meta if rel not in file_map] + + old_entries_by_source = _load_old_entries_by_source(index_path) + kept_entries = _collect_kept_entries(unchanged_files, old_entries_by_source) + + if not changed_files and not deleted_files: + print(f"No changes detected ({len(md_files)} files, {len(unchanged_files)} unchanged)") + return + + new_texts, new_file_entries, changed_meta = _prepare_changed( + changed_files, file_map, memory_dir, chunk_size, chunk_overlap + ) + + if changed_files: + n_batches = (len(new_texts) + BATCH_SIZE - 1) // BATCH_SIZE + print(f"Embedding {len(new_texts)} chunks in {n_batches} batches...", flush=True) + new_embeddings = embed_texts(new_texts) + else: + new_embeddings = [] + + full_meta = _build_full_meta(file_map, changed_meta, unchanged_files) + merged = _merge_and_sort(kept_entries, new_file_entries, new_embeddings) + + _atomic_write(index_path, json.dumps({"files": merged}, ensure_ascii=False)) + _atomic_write( + meta_path, + json.dumps({"version": current_version, "files": full_meta}, ensure_ascii=False), + ) + + print( + f"Indexed {len(md_files)} files to {index_path} " + f"({len(changed_files)} changed, {len(unchanged_files)} unchanged, " + f"{len(deleted_files)} deleted)" ) - print(f"Indexed {len(md_files)} files to {output_dir / INDEX_FILENAME}") diff --git a/tests/test_embedder.py b/tests/test_embedder.py index 9df9021..b652803 100644 --- a/tests/test_embedder.py +++ b/tests/test_embedder.py @@ -3,7 +3,7 @@ import os os.environ.setdefault("OPENAI_BASE_URL", "http://test/v1") -from typing import NoReturn +from typing import ClassVar, NoReturn from unittest.mock import patch import httpx @@ -64,7 +64,10 @@ def test_embed_texts_api_error() -> None: embed_texts(["test"]) -def test_embed_texts_default_model() -> None: +def test_embed_texts_default_model(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("OPENAI_EMBEDDING_MODEL", raising=False) + importlib.reload(embedder_mod) + captured: dict[str, str | list[str]] = {} class FakeResponse: @@ -81,9 +84,10 @@ def test_embed_texts_default_model() -> None: return FakeResponse() with patch.object(httpx.Client, "post", mock_post): - embed_texts(["text"]) + embedder_mod.embed_texts(["text"]) assert captured["model"] == "gemini-embedding-2-preview" + importlib.reload(embedder_mod) def test_embed_texts_custom_model(monkeypatch: pytest.MonkeyPatch) -> None: @@ -142,7 +146,8 @@ def test_embed_texts_trailing_slash(monkeypatch: pytest.MonkeyPatch) -> None: importlib.reload(embedder_mod) -def test_embed_texts_batches() -> None: +def test_embed_texts_batches(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(embedder_mod, "BATCH_SIZE", 2048) calls: list[int] = [] class FakeResponse: @@ -163,8 +168,72 @@ def test_embed_texts_batches() -> None: return FakeResponse(count) with patch.object(httpx.Client, "post", mock_post): - result = embed_texts(["text"] * 3000) + result = embedder_mod.embed_texts(["text"] * 3000) assert len(calls) == 2 assert calls == [2048, 952] assert len(result) == 3000 + + +def test_embed_texts_batch_progress( + monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str] +) -> None: + monkeypatch.setattr(embedder_mod, "BATCH_SIZE", 2048) + + class FakeResponse: + status_code = 200 + + def json(self): + return {"data": [{"embedding": [0.1], "index": 0}]} + + def raise_for_status(self) -> None: + pass + + def mock_post(self, url, **kwargs): + return FakeResponse() + + with patch.object(httpx.Client, "post", mock_post): + embedder_mod.embed_texts(["text"] * 3000) + + captured = capsys.readouterr() + assert "batch 1/2" in captured.out + assert "batch 2/2" in captured.out + + +def test_retry_after_header() -> None: + class RetryResponse: + status_code = 429 + headers: ClassVar[dict[str, str]] = {"Retry-After": "2"} + + def raise_for_status(self) -> NoReturn: + raise httpx.HTTPStatusError( + "Too Many Requests", + request=None, + response=self, + ) + + def json(self): + return {} + + class OkResponse: + status_code = 200 + + def json(self): + return {"data": [{"embedding": [0.1], "index": 0}]} + + def raise_for_status(self) -> None: + pass + + responses = [RetryResponse(), OkResponse()] + + def mock_post(self, url, **kwargs): + return responses.pop(0) + + with ( + patch.object(httpx.Client, "post", mock_post), + patch("src.memory.embedder.time.sleep") as mock_sleep, + ): + result = embed_texts(["text"]) + + mock_sleep.assert_any_call(2) + assert result == [[0.1]] diff --git a/tests/test_embedder_live.py b/tests/test_embedder_live.py index d1c385f..0126c23 100644 --- a/tests/test_embedder_live.py +++ b/tests/test_embedder_live.py @@ -19,3 +19,10 @@ def test_live_embed_batch() -> None: r = embed_texts(["text one", "text two", "text three"]) assert len(r) == 3 assert all(len(emb) > 100 for emb in r) + + +@pytest.mark.skipif(not os.environ.get("RUN_LIVE"), reason="needs RUN_LIVE=1") +def test_live_embed_qwen3_batch() -> None: + r = embed_texts([f"text {i}" for i in range(100)]) + assert len(r) == 100 + assert all(len(emb) == 4096 for emb in r) diff --git a/tests/test_index.py b/tests/test_index.py index 2d3ab47..cc64624 100644 --- a/tests/test_index.py +++ b/tests/test_index.py @@ -2,13 +2,13 @@ import json import os from argparse import Namespace from pathlib import Path -from unittest.mock import patch +from unittest.mock import Mock, patch import pytest os.environ.setdefault("OPENAI_BASE_URL", "http://test/v1") -from src.memory.index import _extract_text, run_index +from src.memory.index import _atomic_write, _content_hash, _extract_text, run_index class TestExtractText: @@ -128,3 +128,225 @@ class TestRunIndex: sources = [f["source"] for f in index["files"]] assert "real.md" in sources assert all(".rag" not in s for s in sources) + + +class TestIncrementalIndex: + def _make_args(self, memory_dir: Path, output_dir: Path) -> Namespace: + return Namespace(memory_dir=str(memory_dir), output=str(output_dir)) + + def _fake_embed(self): + def fake_embed_texts(texts): + return [[0.1, 0.2] for _ in texts] + + return fake_embed_texts + + def test_incremental_add(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("file A content") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + (memory_dir / "b.md").write_text("file B content") + calls: list[list[str]] = [] + + def tracking_embed(texts): + calls.append(list(texts)) + return [[0.1, 0.2] for _ in texts] + + with patch("src.memory.index.embed_texts", tracking_embed): + run_index(self._make_args(memory_dir, index_dir)) + + assert len(calls) == 1 + embedded_text = calls[0][0] + assert "file B content" in embedded_text + assert "file A content" not in embedded_text + + index = json.loads((index_dir / "index.json").read_text()) + sources = {f["source"] for f in index["files"]} + assert {"a.md", "b.md"} <= sources + + def test_incremental_edit(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("original content") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + (memory_dir / "a.md").write_text("edited content here") + calls: list[list[str]] = [] + + def tracking_embed(texts): + calls.append(list(texts)) + return [[0.1, 0.2] for _ in texts] + + with patch("src.memory.index.embed_texts", tracking_embed): + run_index(self._make_args(memory_dir, index_dir)) + + assert len(calls) == 1 + assert "edited content here" in calls[0][0] + + def test_incremental_delete(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("file A") + (memory_dir / "b.md").write_text("file B") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + (memory_dir / "b.md").unlink() + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + index = json.loads((index_dir / "index.json").read_text()) + sources = [f["source"] for f in index["files"]] + assert "a.md" in sources + assert "b.md" not in sources + + meta = json.loads((index_dir / "meta.json").read_text()) + assert "b.md" not in meta["files"] + assert "a.md" in meta["files"] + + def test_no_changes_noop(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("stable content") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + index_mtime_before = (index_dir / "index.json").stat().st_mtime_ns + + embed_mock = Mock(return_value=[[0.1, 0.2]]) + with patch("src.memory.index.embed_texts", embed_mock): + run_index(self._make_args(memory_dir, index_dir)) + + embed_mock.assert_not_called() + index_mtime_after = (index_dir / "index.json").stat().st_mtime_ns + assert index_mtime_after == index_mtime_before + + def test_meta_json_persistence(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("file A") + (memory_dir / "b.md").write_text("file B") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + meta_path = index_dir / "meta.json" + assert meta_path.exists() + meta = json.loads(meta_path.read_text()) + assert "a.md" in meta["files"] + assert "b.md" in meta["files"] + assert "version" in meta + assert meta["files"]["a.md"]["sha256"] + assert meta["files"]["a.md"]["chunks"] == 1 + + (memory_dir / "c.md").write_text("file C") + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + meta2 = json.loads(meta_path.read_text()) + assert "a.md" in meta2["files"] + assert "b.md" in meta2["files"] + assert "c.md" in meta2["files"] + + def test_sha256_content_hash(self, tmp_path: Path) -> None: + assert _content_hash("hello") == _content_hash("hello") + assert _content_hash("hello") != _content_hash("world") + assert len(_content_hash("hello")) == 64 + + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("---\ntitle: T\n---\n\nbody content") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + meta_before = json.loads((index_dir / "meta.json").read_text()) + sha_before = meta_before["files"]["a.md"]["sha256"] + + os.utime(memory_dir / "a.md", None) + + embed_mock = Mock(return_value=[[0.1, 0.2]]) + with patch("src.memory.index.embed_texts", embed_mock): + run_index(self._make_args(memory_dir, index_dir)) + + embed_mock.assert_not_called() + meta_after = json.loads((index_dir / "meta.json").read_text()) + assert meta_after["files"]["a.md"]["sha256"] == sha_before + + def test_atomic_writes(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("file A") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + original_index = (index_dir / "index.json").read_text() + + (memory_dir / "a.md").write_text("file A edited") + + with ( + patch("src.memory.index.embed_texts", self._fake_embed()), + patch("src.memory.index.os.replace", side_effect=OSError("simulated crash")), + pytest.raises(OSError), + ): + run_index(self._make_args(memory_dir, index_dir)) + + assert (index_dir / "index.json").read_text() == original_index + + def test_index_versioning(self, tmp_path: Path) -> None: + memory_dir = tmp_path / "memory" + index_dir = tmp_path / ".rag" + memory_dir.mkdir(parents=True) + (memory_dir / "a.md").write_text("file A") + + with patch("src.memory.index.embed_texts", self._fake_embed()): + run_index(self._make_args(memory_dir, index_dir)) + + meta = json.loads((index_dir / "meta.json").read_text()) + meta["version"] = "old-model:512:64" + (index_dir / "meta.json").write_text(json.dumps(meta)) + + calls: list[int] = [] + + def tracking_embed(texts): + calls.append(len(texts)) + return [[0.1, 0.2] for _ in texts] + + with patch("src.memory.index.embed_texts", tracking_embed): + run_index(self._make_args(memory_dir, index_dir)) + + assert len(calls) == 1 + assert calls[0] >= 1 + + meta2 = json.loads((index_dir / "meta.json").read_text()) + assert meta2["version"] != "old-model:512:64" + + +class TestAtomicWrite: + def test_atomic_write_replaces(self, tmp_path: Path) -> None: + target = tmp_path / "file.json" + target.write_text("old") + _atomic_write(target, "new content") + assert target.read_text() == "new content" + assert not (tmp_path / "file.json.tmp").exists() + + def test_atomic_write_creates(self, tmp_path: Path) -> None: + target = tmp_path / "new.json" + _atomic_write(target, "fresh") + assert target.read_text() == "fresh"