diff --git a/lifecycle/graph_enrich.py b/lifecycle/graph_enrich.py index b54868b..9c1a598 100644 --- a/lifecycle/graph_enrich.py +++ b/lifecycle/graph_enrich.py @@ -324,6 +324,15 @@ async def graph_enrich(layer: str = "user") -> dict[str, Any]: logger.warning("miner %s failed: %s", name, exc) miners[name] = {"edges": -1, "error": str(exc)[:200]} # -1 = failure (edge counts are non-negative) + # Per-miner edge counts, on one line. The phase timings above already say + # which miner ran; they cannot say that a miner ran and produced nothing. + # A miner returning a flat 0 used to be indistinguishable from a healthy + # miner with no work, and on 2026-10-05 that hid a six-day outage of + # `embedding`: `semantic_overlap` stayed frozen while every other relation + # kept growing. `-1` here means the miner failed outright and the reason is + # on the warning line above. + logger.info("miner edges: layer=%s %s", layer, " ".join(f"{name}={res.get('edges')}" for name, res in miners.items())) + # S17 addendum 11: HDBSCAN clusters of MIB vectors + cross-check against louvain — a # report section behind a flag (graph.embed_clusters, default off — the report shape stays stable). from config import config diff --git a/lifecycle/graph_miners.py b/lifecycle/graph_miners.py index 758a0a9..3ff81d6 100644 --- a/lifecycle/graph_miners.py +++ b/lifecycle/graph_miners.py @@ -1025,17 +1025,28 @@ async def miner_embedding(cm: AsyncConnectionManager, layer: str) -> dict[str, i crosscheck = bool(config.get("graph", "embedding_crosscheck", default=False)) - try: - # A-MEM rich embedding: f"{content} {tags}"; canonicalization (_canon from T2) - # applies to tags so that name/technology variants land in one meaning - # cache key. Cache key = raw content — reuses vectors seeded by the ingestor. - # anomaly:* tags (addendum 10) never enter the text — a flag does not - # change the node's vector. - vecs = await embed_texts([f"{c} {' '.join(sorted(t for t in tags.get(nid, []) if not t.startswith('anomaly:')))}" for nid, c in nodes]) - binary = [embed_to_binary(v, dim=len(v)) for v in vecs] - bits = [_bits_int(b) for b in binary] - except Exception: - return {"edges": 0} # embedding backend unavailable (no numpy/model) — miner skipped + # No try/except here, on purpose. This used to be + # `except Exception: return {"edges": 0}` with a comment claiming it only + # covered "no numpy/model" — two things that are both wrong. `embed_texts` + # does not raise when no model is available (it falls back to hash vectors, + # see `_compute_missing_embeddings`), and `embed_to_binary` has its own + # numpy-less path, so the only exceptions reaching here are real failures: + # a dead or slow embedding service above all. + # + # Swallowing them made a six-day outage of this one miner invisible: the + # graph lost `semantic_overlap` while the log stayed empty and this miner + # still reported a normal zero. `graph_enrich` already knows how to report + # a failed miner (`{"edges": -1}` plus a warning) and that branch is only + # reachable if the exception is allowed to escape. + # + # A-MEM rich embedding: f"{content} {tags}"; canonicalization (_canon from T2) + # applies to tags so that name/technology variants land in one meaning + # cache key. Cache key = raw content — reuses vectors seeded by the ingestor. + # anomaly:* tags (addendum 10) never enter the text — a flag does not + # change the node's vector. + vecs = await embed_texts([f"{c} {' '.join(sorted(t for t in tags.get(nid, []) if not t.startswith('anomaly:')))}" for nid, c in nodes]) + binary = [embed_to_binary(v, dim=len(v)) for v in vecs] + bits = [_bits_int(b) for b in binary] # S17 addendum 10: bit-degenerate vectors (0 bits — text without significant # tokens, junk from L3 dumps) → flagged `anomaly:junk_vector`, cleanup candidate. diff --git a/shared/embeddings.py b/shared/embeddings.py index bd4a9a1..56ea0dd 100644 --- a/shared/embeddings.py +++ b/shared/embeddings.py @@ -61,6 +61,17 @@ def _decode_int8(blob: bytes) -> list[float] | None: # limit is 32766 on modern builds but 999 on older ones, so stay well below both. _CACHE_LOOKUP_CHUNK = 500 +# Texts per HTTP request to the embedding service, and the per-request timeout. +# The whole list used to travel in ONE POST, bounded by a fixed 30 s timeout, +# and the measured batch time for a ~6k-text layer was 24.4 s — about 20 % of +# headroom. `miner_embedding` used to swallow the resulting timeout as a silent +# zero, which is how one relation stayed frozen for six days. Splitting keeps a +# single request near 2 s, so the timeout stops being a ceiling on layer size. +# 500 matches `_CACHE_LOOKUP_CHUNK`, already the convention for "how much goes +# into one round trip". +_REMOTE_BATCH_SIZE = 500 +_REMOTE_TIMEOUT_S = 30.0 + def _decode_blob(blob: bytes) -> list[float] | None: """Decode a stored cache blob; None for corrupt/truncated rows (→ cache miss). @@ -330,6 +341,18 @@ async def _get_results_from_cache(self, texts: list[str], cache_tag: str) -> tup return results, to_compute async def _remote_embed(self, texts: list[str]) -> list[list[float]]: + """Embed any number of texts, one bounded request per `_REMOTE_BATCH_SIZE`. + + Callers pass a whole layer; the service is asked in pieces so that one + slow response cannot exceed `_REMOTE_TIMEOUT_S` for the entire layer. + Vectors are returned in input order. + """ + out: list[list[float]] = [] + for batch in _chunked(texts, _REMOTE_BATCH_SIZE): + out.extend(await self._remote_embed_batch(batch)) + return out + + async def _remote_embed_batch(self, texts: list[str]) -> list[list[float]]: """POST /v1/embeddings to the configured embeddings.url. Response is model-tagged with the configured model name, so cache @@ -354,7 +377,7 @@ async def _remote_embed(self, texts: list[str]) -> list[list[float]]: def _post() -> dict[str, Any]: # Scheme is validated above (http/https only). - with urllib.request.urlopen(req, timeout=30) as resp: # noqa: S310 + with urllib.request.urlopen(req, timeout=_REMOTE_TIMEOUT_S) as resp: # noqa: S310 data: dict[str, Any] = json.loads(resp.read()) return data diff --git a/tests/test_lifecycle/test_graph_enrich.py b/tests/test_lifecycle/test_graph_enrich.py index 01cebdb..ad498bc 100644 --- a/tests/test_lifecycle/test_graph_enrich.py +++ b/tests/test_lifecycle/test_graph_enrich.py @@ -85,6 +85,46 @@ async def test_graph_enrich_miner_stubs_report_zero_edges(graph): assert all(v == {"edges": 0} for v in result["miners"].values()) +@pytest.mark.asyncio +async def test_graph_enrich_reports_a_failed_miner_as_minus_one(graph, monkeypatch): + """Сбой майнера — это `-1` и текст ошибки, а не честный ноль. + + 07.10.2026: молчаливый `except` внутри `miner_embedding` обесценивал эту + ветку — исключение до неё не доходило, и один тип связей пропал на шесть + дней без единой строки в логе. + """ + import lifecycle.graph_miners as gm + from lifecycle.graph_enrich import graph_enrich + + async def _boom(cm: Any, layer: str) -> dict[str, int]: + raise RuntimeError("embedding service is down") + + monkeypatch.setattr(gm, "MINERS", {"embedding": _boom}) + + result = await graph_enrich(layer="user") + + assert result["miners"]["embedding"]["edges"] == -1 + assert "embedding service is down" in result["miners"]["embedding"]["error"] + + +@pytest.mark.asyncio +async def test_graph_enrich_logs_edge_counts_per_miner(graph, monkeypatch, caplog): + """Число рёбер по каждому майнеру идёт в лог — «дал 0» видно сразу.""" + import lifecycle.graph_miners as gm + from lifecycle.graph_enrich import graph_enrich + + async def _seven(cm: Any, layer: str) -> dict[str, int]: + return {"edges": 7} + + monkeypatch.setattr(gm, "MINERS", {"embedding": _seven}) + + with caplog.at_level(logging.INFO, logger="lifecycle.graph_enrich"): + await graph_enrich(layer="user") + + lines = [r.getMessage() for r in caplog.records if r.name == "lifecycle.graph_enrich"] + assert any("miner edges:" in m and "embedding=7" in m for m in lines), lines + + @pytest.mark.asyncio async def test_graph_enrich_noop_layer_keeps_stats_shape(graph): from lifecycle.graph_enrich import graph_enrich diff --git a/tests/test_lifecycle/test_graph_miners.py b/tests/test_lifecycle/test_graph_miners.py index b949380..15ed56a 100644 --- a/tests/test_lifecycle/test_graph_miners.py +++ b/tests/test_lifecycle/test_graph_miners.py @@ -564,6 +564,30 @@ async def test_miner_embedding_skips_tool_junk(db): assert await _edges("semantic_overlap") == [] +@pytest.mark.asyncio +async def test_miner_embedding_lets_a_backend_failure_escape(db, monkeypatch): + """Сбой эмбеддингов обязан выйти наружу, а не стать тихим нулём. + + 07.10.2026: `except Exception: return {"edges": 0}` прятал сбой шесть дней — + `semantic_overlap` стоял на месте, а в логе было пусто. Договор теперь такой: + ноль означает «делать было нечего», сбой обязан дойти до `graph_enrich`, + который запишет `edges: -1` и предупреждение. + """ + await _node("первый текст про деплой сервиса", T) + await _node("второй текст про мониторинг", T) + + import shared.embeddings as emb + from lifecycle.graph_miners import miner_embedding + + async def _boom(texts: list[str], prefix: str = "") -> list[list[float]]: + raise RuntimeError("embeddings service is down") + + monkeypatch.setattr(emb, "embed_texts", _boom) + + with pytest.raises(RuntimeError, match="embeddings service is down"): + await miner_embedding(db, "user") + + @pytest.mark.asyncio async def test_miner_entities_synonym_canon_creates_co_mentions(db): # «Лили» и «Lily» — один канон-класс сущности (словарь rag.synonyms, обе стороны) diff --git a/tests/test_shared/test_remote_embed_batching.py b/tests/test_shared/test_remote_embed_batching.py new file mode 100644 index 0000000..2c74bfb --- /dev/null +++ b/tests/test_shared/test_remote_embed_batching.py @@ -0,0 +1,113 @@ +"""One bounded HTTP request per batch, instead of one request for the layer. + +Why these exist: `_remote_embed` used to POST the whole text list in a single +request bounded by a fixed 30 s timeout. A ~6k-text layer measured 24.4 s, so +the timeout was a ceiling on layer size with ~20 % of headroom — and +`miner_embedding` swallowed the resulting timeout as a silent `{"edges": 0}`, +which is how `semantic_overlap` stayed frozen for six days without a log line. + +These pin that the list is split, that nothing is lost or reordered by the +split, and that the ordinary small case still costs exactly one request. +""" + +import json +import urllib.request +from collections.abc import AsyncIterator +from typing import Self + +import pytest + +import shared.embeddings as emb +from shared.connection import connection_manager + +URL = "http://127.0.0.1:8710/v1/embeddings" + + +@pytest.fixture +async def cache(tmp_path, monkeypatch) -> AsyncIterator[emb.EmbeddingCache]: + """Своя база на tmp_path — и уборка за собой. + + `connection_manager._conns` кэширует соединение по имени файла, поэтому + фикстура, которая не чистит `_conns` на выходе, оставляет следующему тесту + соединение к уже удалённому каталогу. Проверено: без этой уборки соседний + `test_archived_memories` падал на `no such table: archived_memories` — он + не изолирован и берёт соединение из того же кэша. + """ + monkeypatch.setattr(connection_manager, "base_dir", tmp_path) + connection_manager._conns.clear() + c = emb.EmbeddingCache(cm=connection_manager) + await c.ensure() + yield c + connection_manager._conns.clear() + + +def _install_fake_service(monkeypatch, calls: list[dict]) -> None: + """Record every POST; answer with `index` reversed to exercise the sort.""" + + class _Resp: + def __init__(self, payload: dict) -> None: + self._payload = payload + + def read(self) -> bytes: + return json.dumps(self._payload).encode() + + def __enter__(self) -> Self: + return self + + def __exit__(self, *exc: object) -> bool: + return False + + def fake_urlopen(req, timeout=None): + texts = json.loads(req.data)["input"] + calls.append({"texts": texts, "timeout": timeout, "url": req.full_url}) + items = [{"index": i, "embedding": [float(t.lstrip("t") or 0), 0.5]} for i, t in enumerate(texts)] + return _Resp({"data": list(reversed(items))}) + + monkeypatch.setattr(emb, "_remote_url", lambda: URL) + monkeypatch.setattr(urllib.request, "urlopen", fake_urlopen) + + +async def test_remote_embed_splits_a_long_list_into_bounded_requests(cache, monkeypatch): + """25 текстов при батче 10 — три запроса, последний короткий.""" + monkeypatch.setattr(emb, "_REMOTE_BATCH_SIZE", 10) + calls: list[dict] = [] + _install_fake_service(monkeypatch, calls) + + out = await cache._remote_embed([f"t{i}" for i in range(25)]) + + assert len(out) == 25 + assert [len(c["texts"]) for c in calls] == [10, 10, 5], f"запросы: {[len(c['texts']) for c in calls]}" + assert all(c["timeout"] == emb._REMOTE_TIMEOUT_S for c in calls) + assert all(c["url"] == URL for c in calls) + + +async def test_remote_embed_keeps_input_order_across_batches(cache, monkeypatch): + """Склейка чанков обязана сохранить порядок входа, а не порядок ответов.""" + monkeypatch.setattr(emb, "_REMOTE_BATCH_SIZE", 5) + calls: list[dict] = [] + _install_fake_service(monkeypatch, calls) + + out = await cache._remote_embed([f"t{i}" for i in range(12)]) + + assert [v[0] for v in out] == [float(i) for i in range(12)] + assert len(calls) == 3 + + +async def test_remote_embed_uses_one_request_when_under_batch_size(cache, monkeypatch): + """Обычный случай не подорожал: список меньше батча — ровно один POST.""" + calls: list[dict] = [] + _install_fake_service(monkeypatch, calls) + + out = await cache._remote_embed(["t0", "t1", "t2"]) + + assert len(calls) == 1 + assert [v[0] for v in out] == [0.0, 1.0, 2.0] + + +async def test_remote_embed_empty_input_makes_no_request(cache, monkeypatch): + """Пустой список — ноль запросов, а не запрос с пустым телом.""" + calls: list[dict] = [] + _install_fake_service(monkeypatch, calls) + + assert await cache._remote_embed([]) == [] + assert calls == []