Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions lifecycle/graph_enrich.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
33 changes: 22 additions & 11 deletions lifecycle/graph_miners.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
25 changes: 24 additions & 1 deletion shared/embeddings.py
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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
Expand All @@ -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

Expand Down
40 changes: 40 additions & 0 deletions tests/test_lifecycle/test_graph_enrich.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
24 changes: 24 additions & 0 deletions tests/test_lifecycle/test_graph_miners.py
Original file line number Diff line number Diff line change
Expand Up @@ -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, обе стороны)
Expand Down
113 changes: 113 additions & 0 deletions tests/test_shared/test_remote_embed_batching.py
Original file line number Diff line number Diff line change
@@ -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 == []
Loading