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 <agent@opencode.local>
This commit is contained in:
Sergey 2026-07-26 19:19:55 +03:00 committed by GitHub
parent 4b377f3e4d
commit 25cf10baa6
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 595 additions and 32 deletions

View file

@ -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

View file

@ -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 не изменён.

View file

@ -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).

View file

@ -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

View file

@ -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

View file

@ -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}")

View file

@ -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]]

View file

@ -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)

View file

@ -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"