"""Memory operations: SQL authority, explicit projection recovery, verified recall.""" import hashlib from tht.phase import effective_decisions from tht.ports.vector import VectorWriteRecord from tht.session.models import PrincipalContext from tht.vectorstore.records import VectorRecord from .core import decided_memory_ids, declined_promotion_seqs, question_context from .models import ( CardInput, CardQuery, Dependency, MemoryConflict, MemoryForbidden, MemoryNotFound, ) from .repository import MemoryRepository from .retrieval import MAX_SEEDS, PROJECTION_FORMAT, RecallScope, expand_and_rank from .solved import _build_solved_snapshot class MemoryService: def __init__(self, repository: MemoryRepository, principal: PrincipalContext, *, store_factory, embedder_factory, language: str = "en"): self.repository = repository self.principal = principal self.store_factory = store_factory self.embedder_factory = embedder_factory if language not in {"en", "it"}: raise ValueError("Memory language must be en or it") self.query_language = {"en": "english", "it": "italian"}[language] def close(self): self.repository.close() def _admin(self): if not self.principal.is_admin: raise MemoryForbidden("Memory administration requires an administrator") def list(self, query: CardQuery): self._admin() return self.repository.list(query) def get(self, card_id: str): self._admin() return self.repository.get(card_id).model_dump(mode="json") def pending(self): self._admin() return self.repository.projections() def save(self, value: CardInput, card_id: str | None = None): self._admin() with self.repository.operation() as repo: card_id = repo.save(value, card_id=card_id) return self._propagate(repo, card_id) def delete(self, card_id: str): self._admin() with self.repository.operation() as repo: repo.delete(card_id) return self._propagate(repo, card_id) def retry(self, card_id: str): self._admin() with self.repository.operation() as repo: return self._propagate(repo, card_id) def _propagate(self, repo, card_id): operation = next((p for p in repo.projections(pending=False) if p["card_id"] == card_id), None) if operation is None: raise MemoryNotFound("Memory operation was not found in this workspace") error = None if operation["pending"]: try: store = self.store_factory() if operation["action"] == "delete": store.delete_memory_records([f"card:{card_id}"]) else: card = repo.get(card_id) kind = "solved_question" if card.family == "solved_question" else "memory" content = "\n".join(filter(None, [card.subject, card.detail, card.scope, card.rationale, card.question, card.sql, " ".join(card.concepts), "\n".join(".".join(filter(None, [d.database, d.schema_name, d.table, d.column])) for d in card.dependencies)])) record = VectorRecord( id=f"card:{card.id}", kind=kind, ref=card.id, title=card.subject, content=content, metadata={"memory_revision": card.revision, "memory_format": PROJECTION_FORMAT, "memory_family": card.family, "memory_scope": card.scope, "memory_concepts": card.concepts, "memory_dependencies": [d.model_dump() for d in card.dependencies]}, ) embedding = self.embedder_factory().embed_documents([content])[0] # A family change may change the vector kind (and point identity). store.delete_memory_records([f"card:{card_id}"]) store.upsert("memory", [VectorWriteRecord( record=record, embedding=embedding, content_hash=hashlib.sha256(card.revision.encode()).hexdigest(), sparse_text=content, sparse_language=self.query_language, )]) except Exception: # noqa: BLE001 - durable pending state covers adapter/factory failures. error = "Memory change is saved; index update is incomplete. Retry the index update." repo.projection_result(card_id, operation["revision"], error) result = {"id": card_id, "saved": True, "indexed": error is None, "action": operation["action"], "error": error} if operation["action"] != "delete": result["card"] = repo.get(card_id).model_dump(mode="json") return result def rebuild(self): self._admin() with self.repository.operation() as repo: # Persist invalidation before deleting anything. A crash remains recoverable. repo.invalidate_all() store = self.store_factory() store.prepare_memory_index() store.delete_kinds("memory", ["memory", "solved_question"]) return [self._propagate(repo, p["card_id"]) for p in repo.projections()] def retrieve(self, question: str, *, searcher, embedder, top: int = 5, family=None, scope: RecallScope | None = None, excluded=()): if type(top) is not int or not 1 <= top <= 100: raise ValueError("Recall limit must be between 1 and 100") if not question.strip(): raise ValueError("Recall question must not be empty") scope = scope or RecallScope() # PostgreSQL is authoritative: an empty archive needs no search projection. # Still read it first so archive failures cannot masquerade as zero results. if self.repository.list(CardQuery(page_size=1))["total"] == 0: return [] kinds = (["solved_question"] if family == "solved_question" else ["memory"] if family else ["memory", "solved_question"]) hits = searcher.search(embedder.embed_query(question), top_n=min(MAX_SEEDS, max(20, top * 3)), kinds=kinds, query_text=question, query_language=self.query_language, metadata_filter=scope.vector_filter(family)) # Mutations use the same lock: links, eligibility and payload are resolved # together against current authority, after the potentially slow vector call. with self.repository.operation() as repo: return expand_and_rank(repo, hits, scope=scope, family=family, excluded=set(excluded), top=top) def recall(self, question: str, *, searcher, embedder, top: int = 5, solved: bool = False, decisions=(), scope: RecallScope | None = None): family = "solved_question" if solved else "domain_clarification" candidates = self.retrieve(question, searcher=searcher, embedder=embedder, top=top, family=family, scope=scope, excluded=decided_memory_ids(list(decisions))) result = [] for candidate in candidates: card = candidate.card common = {"id": card.id, "session_id": card.session_id, "revision": card.revision, "family": card.family, "scope": card.scope, "dependencies": [d.model_dump() for d in card.dependencies], "tables": sorted({d.table for d in card.dependencies if d.table}), "score": round(candidate.score, 6), "retrieval": {"path": list(candidate.path), "method": "hybrid_links"}} if solved: result.append({**common, "question": card.question, "sql": card.sql}) else: # New workflow category consumption belongs to M3. result.append({**common, "type": "concept_clarified", "subject": card.subject, "detail": card.detail, "rationale": card.rationale, "question_context": card.question, "scope": card.scope, "concepts": card.concepts}) return result def _session(self, snapshot): if (snapshot.manifest.workspace_id != self.repository.workspace_id or (not self.principal.is_admin and snapshot.manifest.author != self.principal.subject)): raise MemoryForbidden("Memory source session is outside the authorized context") def promotions(self, snapshot): self._session(snapshot) decisions = effective_decisions(snapshot) declined = declined_promotion_seqs(decisions) result = [] seen = set() for d in decisions: if d.type != "concept_clarified" or d.seq in declined: continue source_key = f"decision:{snapshot.manifest.id}:{d.seq}" if self.repository.source(source_key) is not None: continue content = (d.subject, d.detail, d.rationale) if content in seen: continue seen.add(content) result.append({"decision_seq": d.seq, "type": d.type, "subject": d.subject, "detail": d.detail, "rationale": d.rationale, "question_context": question_context(decisions, snapshot.manifest)}) return result def promote(self, snapshot, seqs): self._session(snapshot) decisions = effective_decisions(snapshot) selected = [d for d in decisions if d.seq in seqs and d.type == "concept_clarified"] results = [] with self.repository.operation() as repo: for d in selected: key = f"decision:{snapshot.manifest.id}:{d.seq}" existing = repo.source(key) if existing and existing["action"] == "delete": continue card_id = existing["card_id"] if existing else repo.save(CardInput( family="domain_clarification", subject=d.subject, detail=d.detail, rationale=d.rationale, scope=self.repository.workspace_id, question=question_context(decisions, snapshot.manifest), concepts=[d.subject], ), source_key=key, session_id=snapshot.manifest.id, decision_seq=d.seq) results.append(self._propagate(repo, card_id)) return results def save_solved(self, snapshot, promoted_tables=None): self._session(snapshot) if snapshot.manifest.status != "finalized": raise MemoryConflict("Only a finalized session can produce a solved-question card") record = _build_solved_snapshot(snapshot, promoted_tables) with self.repository.operation() as repo: existing = repo.source(f"solved:{snapshot.manifest.id}") if existing and existing["action"] == "delete": return {"saved": True, "indexed": True, "action": "delete", "error": None} card_id = existing["card_id"] if existing else repo.save(CardInput( family="solved_question", subject=record.title, scope=self.repository.workspace_id, question=record.content, sql=record.metadata["sql"], dependencies=[Dependency(database=snapshot.manifest.database, schema_name=snapshot.manifest.db_schema, table=table) for table in record.metadata["tables"]], ), source_key=f"solved:{snapshot.manifest.id}", session_id=snapshot.manifest.id) return self._propagate(repo, card_id) def retry_solved(self, snapshot): self._session(snapshot) with self.repository.operation() as repo: existing = repo.source(f"solved:{snapshot.manifest.id}") if existing is None: raise MemoryNotFound("No authoritative exemplar exists for this session") return self._propagate(repo, existing["card_id"])