"""M7: local vectors for imported passages, and what happens when there are none. The semantic half of retrieval. It uses the **existing** provider — the same `OpenAICompatibleProvider` the memory bank builds through `memorybank.embedding_provider` — and that is not a convenience. That path is where the endpoint allowlist is re-checked before every request, where the OS/private-CA trust store is unioned into verification, and where timeouts and error shapes are decided (ADR 011, `endpoints.py`, `tlstrust.py`). A second HTTP client here would be a second policy, and the one thing a local-only product cannot afford is two answers to "where may this connect". ## Failure is normal and must be visible Ollama is not running; the embedding model is not pulled; the LAN host is asleep. None of these may cost the reader their import. So: the source stays — content and classification are not derived from anything lexical retrieval keeps working — FTS5 is local SQLite and never touched the network the failure is recorded on the source — `embed_state`, `embed_detail` and on the campaign — `derived_status`, kind "knowledge" a retry fixes it — the next turn, or Reindex The campaign-level record reuses M6's `derived.py` rather than inventing a second status system. The per-source columns exist alongside it because "which file failed" is not a question a per-campaign row can answer, and it is the question a reader actually has. `derived.KNOWLEDGE` is its own kind rather than folded into `derived.EMBEDDING`. The memory bank's embeddings and the knowledge library's embeddings fail independently and are fixed by different actions, and M6's finding M6-F5 — reporting `ok` for work that never ran — is the same mistake as reporting one health for two subsystems. """ from __future__ import annotations import logging from sqlalchemy import select from sqlalchemy.orm import Session from .. import derived, memorybank, models, vectors from ..providers import ProviderError from . import fts log = logging.getLogger(__name__) #: Passages per embedding request. Matches the memory bank's batch size; the #: endpoint is the same one. MAX_BATCH = 32 #: How many passages one pass will embed. A first import of a large library #: would otherwise hold a turn's background task open for a long time; the #: remainder is picked up by the next pass, and `pending_count` says how many #: are left, so the state is legible rather than merely eventual. MAX_PER_RUN = 512 def model_name(settings: models.Settings) -> str: return (settings.embedding_model or "").strip() def enabled(settings: models.Settings) -> bool: """Whether semantic retrieval is configured at all. No embedding model is not a failure — it is a supported configuration in which retrieval is lexical. Reporting it as a failure would be M6-F5 again in the other direction: an alarm about a thing nobody asked for. """ return bool(model_name(settings)) def pending_chunks( db: Session, adventure_id: int, model: str, limit: int ) -> list[models.KnowledgeChunk]: """Passages of enabled, ready sources that have no current vector. "Current" means a vector from *this* embedding model at *this* parser and chunking version. A model change invalidates every vector, which is why the comparison is on the row's own metadata rather than on its presence. """ return list( db.execute( select(models.KnowledgeChunk) .join( models.KnowledgeSource, models.KnowledgeSource.id == models.KnowledgeChunk.source_id, ) .outerjoin( models.KnowledgeEmbedding, models.KnowledgeEmbedding.chunk_id == models.KnowledgeChunk.id, ) .where( models.KnowledgeChunk.adventure_id == adventure_id, models.KnowledgeSource.enabled.is_(True), models.KnowledgeSource.index_state == "ready", (models.KnowledgeEmbedding.id.is_(None)) | (models.KnowledgeEmbedding.model != model), ) .order_by(models.KnowledgeChunk.id) .limit(limit) ).scalars().all() ) def pending_count(db: Session, adventure_id: int, model: str) -> int: """How many passages are still waiting for a vector.""" return len(pending_chunks(db, adventure_id, model, MAX_PER_RUN + 1)) async def embed_pending( db: Session, adventure: models.Adventure, settings: models.Settings ) -> int: """Embeds what is missing. Returns how many vectors were written. Records its own outcome on every source it touched and on the campaign, and never raises: an embedding failure is not allowed to reach the turn that scheduled it. """ model = model_name(settings) if not model: derived.succeeded(db, adventure.id, derived.KNOWLEDGE, did_work=False) return 0 chunks = pending_chunks(db, adventure.id, model, MAX_PER_RUN) if not chunks: derived.succeeded(db, adventure.id, derived.KNOWLEDGE, did_work=False) _settle_sources(db, adventure.id, model) return 0 provider = memorybank.embedding_provider(settings) written = 0 try: for start in range(0, len(chunks), MAX_BATCH): batch = chunks[start:start + MAX_BATCH] payload = [fts.index_line(c.heading_path, c.text) for c in batch] produced = await provider.embed(payload) for chunk_row, vector in zip(batch, produced): _store(db, chunk_row, vector, model) written += 1 except ProviderError as exc: # Soft failure, loudly recorded. The chunks keep no vector, so the next # pass retries exactly them; the sources keep their content and their # lexical index, so the library still answers queries. derived.failed(db, adventure.id, derived.KNOWLEDGE, exc) _mark_sources(db, {c.source_id for c in chunks}, "failed", str(exc)) return written except Exception as exc: # pragma: no cover - defensive derived.failed(db, adventure.id, derived.KNOWLEDGE, exc) _mark_sources(db, {c.source_id for c in chunks}, "failed", str(exc)) return written derived.succeeded(db, adventure.id, derived.KNOWLEDGE, did_work=written > 0) _settle_sources(db, adventure.id, model) return written def _store( db: Session, chunk_row: models.KnowledgeChunk, vector: list[float], model: str ) -> None: """Writes or replaces one passage's vector, with the metadata to date it.""" row = db.execute( select(models.KnowledgeEmbedding).where( models.KnowledgeEmbedding.chunk_id == chunk_row.id ) ).scalars().first() if row is None: row = models.KnowledgeEmbedding( chunk_id=chunk_row.id, adventure_id=chunk_row.adventure_id ) db.add(row) row.vector = vectors.pack(vector) row.model = model row.dimensions = len(vector) row.parser_version = chunk_row.source.parser_version if chunk_row.source else 1 row.chunking_version = chunk_row.source.chunking_version if chunk_row.source else 1 row.created_at = models.utcnow() forget_cached(chunk_row.adventure_id) def _mark_sources(db: Session, source_ids: set[int], state: str, detail: str) -> None: if not source_ids: return db.query(models.KnowledgeSource).filter( models.KnowledgeSource.id.in_(source_ids) ).update( {"embed_state": state, "embed_detail": detail[:2000]}, synchronize_session=False, ) def _settle_sources(db: Session, adventure_id: int, model: str) -> None: """Marks each source `ok` or `pending` according to what it actually holds. Run after a successful pass so a source that was failing and has now been embedded stops saying so. A source with passages still waiting reports `pending` rather than `ok`, because `MAX_PER_RUN` can leave a large library part-way through and "ok" would be untrue. The flush is load-bearing. This session does not autoflush, so the rows `_store` just added are still pending in it, and the query below would not see them — every source would report `pending` immediately after being embedded, which is exactly the misleading status M6-F5 was about. """ db.flush() outstanding = { chunk.source_id for chunk in pending_chunks(db, adventure_id, model, MAX_PER_RUN + 1) } sources = db.execute( select(models.KnowledgeSource).where( models.KnowledgeSource.adventure_id == adventure_id ) ).scalars().all() for source in sources: if not source.enabled or source.index_state != "ready": continue if source.id in outstanding: source.embed_state = "pending" source.embed_detail = "" else: source.embed_state = "ok" source.embed_detail = "" def clear_vectors(db: Session, adventure_id: int) -> int: """Drops every vector in one campaign, so the next pass rebuilds them. This is the semantic half of Reindex. It touches no source, no passage, no story row, which is what `IMPORTED-KNOWLEDGE-DESIGN.md` §55 requires of a reindex — and it is the reason `KnowledgeEmbedding` is a table of its own. """ removed = db.query(models.KnowledgeEmbedding).filter( models.KnowledgeEmbedding.adventure_id == adventure_id ).delete(synchronize_session=False) db.query(models.KnowledgeSource).filter( models.KnowledgeSource.adventure_id == adventure_id ).update({"embed_state": "idle", "embed_detail": ""}, synchronize_session=False) forget_cached(adventure_id) return removed or 0 # ---------------------------------------------------------- the vector cache # # The same idea as the memory bank's, and for the same measured reason: turns # for one campaign arrive one after another, the library changes rarely between # them, and re-reading every vector on every turn is the largest read a turn # makes. `array("f")` holds four bytes a component, matching the column. # # Correctness rests on one rule: **every write to a vector calls # `forget_cached`.** There are three of them and they are all in this module. # Reads reconcile against the catalogue they were given, so a deletion needs no # invalidation at all — a chunk that is no longer listed is dropped from the # cache on the next read. _cache: dict[int, dict[int, object]] = {} CACHE_ADVENTURES = 8 def forget_cached(adventure_id: int) -> None: _cache.pop(adventure_id, None) def vectors_for( db: Session, adventure_id: int, chunk_ids: list[int] ) -> dict[int, object]: """The vectors for `chunk_ids`, reading only the ones not already held.""" held = _cache.get(adventure_id) if held is None: while len(_cache) >= CACHE_ADVENTURES: _cache.pop(next(iter(_cache))) held = _cache[adventure_id] = {} wanted = set(chunk_ids) for gone in set(held) - wanted: del held[gone] missing = [chunk_id for chunk_id in chunk_ids if chunk_id not in held] if missing: rows = db.execute( select(models.KnowledgeEmbedding.chunk_id, models.KnowledgeEmbedding.vector) .where(models.KnowledgeEmbedding.chunk_id.in_(missing)) ).all() for chunk_id, blob in rows: if blob: held[chunk_id] = vectors.unpack(blob) return held