feat: implement memory and evidence administration with guided repairs
Publish documentation / publish (push) Successful in 1m27s

Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation.

Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
This commit is contained in:
Codex
2026-09-10 10:31:34 +02:00
parent 8fe526dd6e
commit 82e2c91f42
168 changed files with 11914 additions and 1772 deletions
+82
View File
@@ -0,0 +1,82 @@
"""Apply only removals confirmed by a successful Catalog physical synchronization."""
import json
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import text
from .models import MemoryConflict
from .review import digest
class RemovedColumn(BaseModel):
model_config = ConfigDict(extra="forbid")
table: str = Field(min_length=1, max_length=200)
column: str = Field(min_length=1, max_length=200)
class CleanupRequest(BaseModel):
model_config = ConfigDict(extra="forbid")
sync_id: str = Field(min_length=1, max_length=200)
database: str = Field(min_length=1, max_length=200)
schema_name: str = Field(min_length=1, max_length=200)
removed_tables: list[str] = Field(default_factory=list, max_length=10000)
removed_columns: list[RemovedColumn] = Field(default_factory=list, max_length=100000)
def cleanup(service, request: CleanupRequest):
service._admin()
fingerprint = digest(request.model_dump(mode="json"))
with service.repository.operation() as repo:
with repo.transaction() as connection:
params = {"w": repo.workspace_id, "sync": request.sync_id}
receipt = (
connection.execute(
text(
"SELECT request_hash,card_ids "
"FROM thoth_memory.cleanup_receipts WHERE workspace_id=:w AND sync_id=:sync"
),
params,
)
.mappings()
.first()
)
if receipt:
if receipt["request_hash"] != fingerprint:
raise MemoryConflict(
"Physical cleanup identity was reused with different removals"
)
identities = receipt["card_ids"]
else:
identities = list(
connection.execute(
text(
"SELECT DISTINCT card_id "
"FROM thoth_memory.dependencies d WHERE workspace_id=:w "
"AND database_id=:db AND schema_name=:schema AND (table_name=ANY(:tables) "
"OR EXISTS (SELECT 1 FROM jsonb_to_recordset(CAST(:columns AS jsonb)) "
'AS removed("table" text, "column" text) WHERE '
'removed."table"=d.table_name AND removed."column"=d.column_name))'
),
{
**params,
"db": request.database,
"schema": request.schema_name,
"tables": request.removed_tables,
"columns": json.dumps(
[c.model_dump() for c in request.removed_columns]
),
},
).scalars()
)
for identity in identities:
repo.delete(identity)
connection.execute(
text(
"INSERT INTO thoth_memory.cleanup_receipts "
"VALUES (:w,:sync,:hash,CAST(:ids AS jsonb))"
),
{**params, "hash": fingerprint, "ids": json.dumps(identities)},
)
results = [service._propagate(repo, identity) for identity in identities]
return {"deleted": len(identities), "indexed": all(r["indexed"] for r in results)}
+1 -1
View File
@@ -25,7 +25,7 @@ class MemoryRecord(BaseModel):
concepts: list[str] = []
_MEM_ID_RE = re.compile(r"\bmem-\d{4,}\b")
_MEM_ID_RE = re.compile(r"\bmem-(?:[0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}|\d{4,})\b")
def decided_memory_ids(decisions: list[DecisionRecord]) -> set[str]:
+74
View File
@@ -0,0 +1,74 @@
"""Versioned Memory migration pack, run only by installation preparation."""
import hashlib
import os
from functools import lru_cache
from importlib.resources import files
from pathlib import Path
from sqlalchemy import URL, create_engine, text
from sqlalchemy.pool import NullPool
def installation_url(*, migrator: bool = False) -> str:
role = "MIGRATOR" if migrator else "RUNTIME"
direct = os.environ.get(f"THT_CATALOG_{role}_DATABASE_URL")
if not migrator:
direct = direct or os.environ.get("THT_CATALOG_DATABASE_URL")
if direct:
return direct.replace("postgresql://", "postgresql+psycopg2://", 1)
prefix = "THT_CATALOG_"
try:
password = Path(os.environ[prefix + role + "_PASSWORD_FILE"]).read_text().strip()
url = URL.create(
"postgresql+psycopg2", host=os.environ[prefix + "DB_HOST"],
port=int(os.environ.get(prefix + "DB_PORT", "5432")),
database=os.environ[prefix + "DB_NAME"],
username=os.environ[prefix + role + "_USER"], password=password,
)
return url.render_as_string(hide_password=False)
except (KeyError, OSError, ValueError):
raise ValueError("Memory PostgreSQL installation configuration is unavailable") from None
@lru_cache(maxsize=1)
def expected_migrations() -> dict[str, str]:
return {p.name: hashlib.sha256(p.read_text().encode()).hexdigest()
for p in files("tht").joinpath("migrations/memory").iterdir()
if p.name.endswith(".sql")}
def migrate(database_url: str) -> None:
engine = create_engine(database_url, poolclass=NullPool)
try:
with engine.begin() as connection:
connection.execute(text("SELECT pg_advisory_xact_lock(792114203)"))
connection.execute(text("CREATE SCHEMA IF NOT EXISTS thoth_memory"))
connection.execute(text("CREATE TABLE IF NOT EXISTS thoth_memory.migrations "
"(version text PRIMARY KEY, checksum text NOT NULL)"))
applied = dict(connection.execute(text(
"SELECT version, checksum FROM thoth_memory.migrations"
)).all())
pack = sorted(files("tht").joinpath("migrations/memory").iterdir(), key=lambda p: p.name)
known = {p.name for p in pack if p.name.endswith(".sql")}
if set(applied) - known:
raise ValueError("Memory schema is newer than this application")
for path in pack:
if path.name not in known:
continue
sql = path.read_text()
digest = hashlib.sha256(sql.encode()).hexdigest()
if path.name in applied:
if applied[path.name] != digest:
raise ValueError("Memory migration checksum mismatch")
continue
connection.execute(text(sql))
connection.execute(text("INSERT INTO thoth_memory.migrations VALUES (:v, :c)"),
{"v": path.name, "c": digest})
finally:
engine.dispose()
if __name__ == "__main__":
migrate(installation_url(migrator=True))
print("Memory migrations: ready")
+112
View File
@@ -0,0 +1,112 @@
"""Authoritative Memory contracts, independent of workflow decision kinds."""
from datetime import datetime
from typing import Literal, Self
from pydantic import BaseModel, ConfigDict, Field, model_validator
Family = Literal["domain_clarification", "sql_rule", "solved_question", "explained_error"]
class Dependency(BaseModel):
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
database: str = Field(min_length=1, max_length=200)
schema_name: str = Field(default="", max_length=200)
table: str = Field(default="", max_length=200)
column: str = Field(default="", max_length=200)
@model_validator(mode="after")
def structured(self) -> Self:
if self.column and not self.table:
raise ValueError("A column dependency requires a table")
if self.table and not self.schema_name:
raise ValueError("A table dependency requires a schema")
return self
class LinkInput(BaseModel):
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
target_id: str = Field(min_length=1, max_length=100)
meaning: str = Field(min_length=1, max_length=1000)
class CardInput(BaseModel):
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
family: Family
subject: str = Field(min_length=1, max_length=1000)
detail: str = Field(default="", max_length=50000)
scope: str = Field(min_length=1, max_length=10000)
rationale: str = Field(default="", max_length=10000)
question: str = Field(default="", max_length=10000)
sql: str = Field(default="", max_length=100000)
concepts: list[str] = Field(default_factory=list, max_length=100)
dependencies: list[Dependency] = Field(default_factory=list, max_length=200)
links: list[LinkInput] = Field(default_factory=list, max_length=200)
@model_validator(mode="after")
def valid_family(self) -> Self:
if self.family == "solved_question" and (not self.question or not self.sql):
raise ValueError("A solved question requires its question and approved SQL")
if self.family == "explained_error" and (not self.detail or not self.rationale):
raise ValueError("An explained error requires a correction and rationale")
if self.family != "solved_question" and self.sql:
raise ValueError("Only solved questions carry exemplar SQL")
if any(not c.strip() or len(c) > 200 for c in self.concepts):
raise ValueError("Concepts must contain between 1 and 200 characters")
if len({link.target_id for link in self.links}) != len(self.links):
raise ValueError("Each linked destination must be unique")
return self
class Card(CardInput):
id: str
workspace_id: str
origin: Literal["manual", "workflow"]
session_id: str | None = None
decision_seq: int | None = None
created_at: datetime
updated_at: datetime
revision: str
indexed: bool = False
class CardQuery(BaseModel):
model_config = ConfigDict(extra="forbid")
q: str = Field(default="", max_length=1000)
family: Family | None = None
concept: str = Field(default="", max_length=200)
database: str = Field(default="", max_length=200)
table: str = Field(default="", max_length=200)
column: str = Field(default="", max_length=200)
origin: Literal["manual", "workflow"] | None = None
updated_after: datetime | None = None
updated_before: datetime | None = None
page: int = Field(default=1, ge=1)
page_size: int = Field(default=25, ge=1, le=100)
sort: Literal["updated_at", "created_at", "subject", "family"] = "updated_at"
direction: Literal["asc", "desc"] = "desc"
class MemoryError(Exception):
code = "memory_operation_failed"
status = 500
class MemoryUnavailable(MemoryError):
code = "memory_unavailable"
status = 503
class MemoryNotFound(MemoryError):
code = "memory_not_found"
status = 404
class MemoryForbidden(MemoryError):
code = "memory_forbidden"
status = 403
class MemoryConflict(MemoryError):
code = "memory_conflict"
status = 409
+245
View File
@@ -0,0 +1,245 @@
"""Workspace-scoped PostgreSQL persistence and durable projection work."""
import json
from contextlib import contextmanager
from uuid import uuid4
from sqlalchemy import create_engine, text
from sqlalchemy.exc import IntegrityError, SQLAlchemyError
from sqlalchemy.pool import NullPool
from .migrate import expected_migrations
from .models import Card, CardInput, CardQuery, MemoryConflict, MemoryNotFound, MemoryUnavailable
class MemoryRepository:
def __init__(self, database_url: str, workspace_id: str, *, engine=None, connection=None):
self.workspace_id = workspace_id
self.engine = engine or create_engine(
database_url, poolclass=NullPool, connect_args={"connect_timeout": 5},
)
self.connection = connection
def close(self):
self.engine.dispose()
@contextmanager
def transaction(self):
connection = self.connection
try:
connection = connection or self.engine.connect()
with (connection.begin_nested() if connection.in_transaction() else connection.begin()):
connection.execute(text("SET LOCAL ROLE thoth_memory_runtime"))
connection.execute(text("SELECT set_config('thoth.memory_workspace', :w, true)"),
{"w": self.workspace_id})
connection.execute(text("SET LOCAL statement_timeout = '15s'"))
installed = dict(connection.execute(text(
"SELECT version, checksum FROM thoth_memory.migrations"
)).all())
if installed != expected_migrations():
raise MemoryUnavailable("Memory schema is incompatible; run installation migrations")
yield connection
except IntegrityError:
raise MemoryConflict("Memory references conflict with the current archive") from None
except SQLAlchemyError:
raise MemoryUnavailable("Memory archive is unavailable; check its migrations and access") \
from None
finally:
if self.connection is None and connection is not None:
connection.close()
@contextmanager
def operation(self):
"""Serialize each workspace across SQL commits and the bounded vector call."""
try:
connection = self.engine.connect()
connection.execute(text("SET statement_timeout = '15s'"))
connection.execute(text("SELECT pg_advisory_lock(hashtextextended(:w, 792114204))"),
{"w": self.workspace_id})
connection.commit()
except SQLAlchemyError:
if 'connection' in locals():
connection.close()
raise MemoryUnavailable("Memory archive is busy or unavailable") from None
try:
yield MemoryRepository("", self.workspace_id, engine=self.engine, connection=connection)
finally:
# NullPool closes the physical connection, releasing the session advisory lock.
connection.close()
def _card(self, connection, row) -> Card:
params = {"w": self.workspace_id, "id": row["id"]}
links = connection.execute(text(
"SELECT target_id, meaning FROM thoth_memory.links "
"WHERE workspace_id=:w AND source_id=:id ORDER BY target_id"
), params).mappings().all()
dependencies = connection.execute(text(
'SELECT database_id AS database, schema_name, table_name AS "table", '
'column_name AS "column" FROM thoth_memory.dependencies '
"WHERE workspace_id=:w AND card_id=:id "
"ORDER BY database_id, schema_name, table_name, column_name"
), params).mappings().all()
return Card.model_validate({
**row["data"], "id": row["id"], "workspace_id": self.workspace_id,
"family": row["family"], "subject": row["subject"], "origin": row["origin"],
"created_at": row["created_at"], "updated_at": row["updated_at"],
"revision": row["revision"], "indexed": row["indexed"],
"links": [dict(v) for v in links], "dependencies": [dict(v) for v in dependencies],
})
@staticmethod
def _selection():
return ("SELECT c.*, (p.revision=c.revision AND NOT p.pending "
"AND p.action='upsert' AND p.format=2) AS indexed FROM thoth_memory.cards c "
"JOIN thoth_memory.projections p ON p.workspace_id=c.workspace_id "
"AND p.card_id=c.id ")
def get(self, card_id: str) -> Card:
with self.transaction() as c:
row = c.execute(text(self._selection()+"WHERE c.workspace_id=:w AND c.id=:id"),
{"w": self.workspace_id, "id": card_id}).mappings().first()
if row is None:
raise MemoryNotFound("Memory card was not found in this workspace")
return self._card(c, row)
def list(self, query: CardQuery) -> dict:
conditions = ["c.workspace_id=:w"]
params = {"w": self.workspace_id}
if query.q:
conditions.append("(c.id ILIKE :q OR c.subject ILIKE :q OR c.data::text ILIKE :q)")
params["q"] = "%" + query.q.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%"
for key in ("family", "origin"):
if value := getattr(query, key):
conditions.append(f"c.{key}=:{key}")
params[key] = value
if query.concept:
conditions.append("c.data->'concepts' @> CAST(:concept AS jsonb)")
params["concept"] = json.dumps([query.concept])
refs = []
for key, column in [("database", "database_id"), ("table", "table_name"),
("column", "column_name")]:
if value := getattr(query, key):
refs.append(f"d.{column}=:{key}")
params[key] = value
if refs:
conditions.append("EXISTS (SELECT 1 FROM thoth_memory.dependencies d WHERE "
"d.workspace_id=c.workspace_id AND d.card_id=c.id AND "
+ " AND ".join(refs) + ")")
for key, comparison in [("updated_after", ">="), ("updated_before", "<=")]:
if value := getattr(query, key):
conditions.append(f"c.updated_at {comparison} :{key}")
params[key] = value
where = " WHERE " + " AND ".join(conditions)
with self.transaction() as c:
total = c.execute(text("SELECT count(*) FROM thoth_memory.cards c" + where),
params).scalar_one()
rows = c.execute(text(self._selection() + where
+ f" ORDER BY c.{query.sort} {query.direction}, c.id ASC LIMIT :limit OFFSET :offset"),
{**params, "limit": query.page_size, "offset": (query.page - 1) * query.page_size},
).mappings().all()
return {"items": [self._card(c, row).model_dump(mode="json") for row in rows],
"total": total, "page": query.page, "page_size": query.page_size}
def exact_match(self, value: CardInput) -> Card | None:
"""Match authored content only; provenance and index state do not create new knowledge."""
data = value.model_dump(mode="json", exclude={"dependencies", "links"})
with self.transaction() as c:
rows = c.execute(text(self._selection() +
"WHERE c.workspace_id=:w AND c.family=:family AND c.subject=:subject "
"AND c.data - 'session_id' - 'decision_seq'=CAST(:data AS jsonb) ORDER BY c.id"),
{"w": self.workspace_id, "family": value.family, "subject": value.subject,
"data": json.dumps(data)}).mappings()
for row in rows:
candidate = self._card(c, row)
def ordered(values):
return sorted(json.dumps(v.model_dump(), sort_keys=True) for v in values)
if (ordered(candidate.dependencies) == ordered(value.dependencies)
and ordered(candidate.links) == ordered(value.links)):
return candidate
return None
def source(self, source_key: str):
with self.transaction() as c:
row = c.execute(text("SELECT card_id, action FROM thoth_memory.projections "
"WHERE workspace_id=:w AND source_key=:s"),
{"w": self.workspace_id, "s": source_key}).mappings().first()
return dict(row) if row else None
def save(self, value: CardInput, *, card_id: str | None = None, source_key: str | None = None,
session_id: str | None = None, decision_seq: int | None = None,
new_id: str | None = None) -> str:
creating = card_id is None
card_id = card_id or new_id or "mem-" + str(uuid4())
revision = str(uuid4())
data = value.model_dump(mode="json", exclude={"links", "dependencies"})
with self.transaction() as c:
if not creating:
old = c.execute(text("SELECT data FROM thoth_memory.cards "
"WHERE workspace_id=:w AND id=:id"),
{"w": self.workspace_id, "id": card_id}).scalar_one_or_none()
if old is None:
raise MemoryNotFound("Memory card was not found in this workspace")
data.update({k: old.get(k) for k in ("session_id", "decision_seq")})
else:
if c.execute(text("SELECT 1 FROM thoth_memory.projections "
"WHERE workspace_id=:w AND card_id=:id"),
{"w": self.workspace_id, "id": card_id}).first():
raise MemoryConflict("A proposed card identity was already used")
data.update(session_id=session_id, decision_seq=decision_seq)
params = {"w": self.workspace_id, "id": card_id, "r": revision,
"data": json.dumps(data), "family": value.family, "subject": value.subject,
"origin": "workflow" if source_key else "manual", "source": source_key}
c.execute(text("INSERT INTO thoth_memory.cards "
"(workspace_id,id,family,subject,origin,data,revision) "
"VALUES (:w,:id,:family,:subject,:origin,CAST(:data AS jsonb),:r) "
"ON CONFLICT (workspace_id,id) DO UPDATE SET family=EXCLUDED.family, "
"subject=EXCLUDED.subject,data=EXCLUDED.data,revision=EXCLUDED.revision, "
"updated_at=clock_timestamp()"), params)
c.execute(text("DELETE FROM thoth_memory.links WHERE workspace_id=:w AND source_id=:id"), params)
for link in value.links:
c.execute(text("INSERT INTO thoth_memory.links VALUES (:w,:id,:target,:meaning)"),
{**params, "target": link.target_id, "meaning": link.meaning})
c.execute(text("DELETE FROM thoth_memory.dependencies "
"WHERE workspace_id=:w AND card_id=:id"), params)
for dep in {tuple(d.model_dump().values()) for d in value.dependencies}:
c.execute(text("INSERT INTO thoth_memory.dependencies VALUES "
"(:w,:id,:database,:schema,:table,:column)"),
{**params, **dict(zip(("database", "schema", "table", "column"), dep))})
c.execute(text("INSERT INTO thoth_memory.projections "
"(workspace_id,card_id,revision,action,source_key) VALUES (:w,:id,:r,'upsert',:source) "
"ON CONFLICT (workspace_id,card_id) DO UPDATE SET revision=EXCLUDED.revision, "
"action='upsert',pending=true,error=NULL,updated_at=clock_timestamp()"), params)
return card_id
def delete(self, card_id: str):
with self.transaction() as c:
params = {"w": self.workspace_id, "id": card_id, "r": str(uuid4())}
deleted = c.execute(text("DELETE FROM thoth_memory.cards "
"WHERE workspace_id=:w AND id=:id"), params).rowcount
if not deleted:
raise MemoryNotFound("Memory card was not found in this workspace")
c.execute(text("UPDATE thoth_memory.projections SET action='delete',revision=:r,"
"pending=true,error=NULL,updated_at=clock_timestamp() "
"WHERE workspace_id=:w AND card_id=:id"), params)
def projections(self, *, pending: bool = True):
with self.transaction() as c:
needs_update = "(pending OR (action='upsert' AND format<>2))"
rows = c.execute(text("SELECT card_id,revision,action,"+needs_update+" AS pending,"
"error,updated_at "
"FROM thoth_memory.projections WHERE workspace_id=:w "
+ ("AND "+needs_update+" " if pending else "") + "ORDER BY updated_at,card_id"),
{"w": self.workspace_id}).mappings().all()
return [dict(row) for row in rows]
def projection_result(self, card_id: str, revision: str, error: str | None):
with self.transaction() as c:
c.execute(text("UPDATE thoth_memory.projections SET pending=:p,error=:error,format=2,"
"updated_at=clock_timestamp() WHERE workspace_id=:w AND card_id=:id AND revision=:r"),
{"p": error is not None, "error": error, "w": self.workspace_id,
"id": card_id, "r": revision})
def invalidate_all(self):
with self.transaction() as c:
c.execute(text("UPDATE thoth_memory.projections SET pending=true,error=NULL "
"WHERE workspace_id=:w"), {"w": self.workspace_id})
+131
View File
@@ -0,0 +1,131 @@
"""Bounded graph recall over current Memory cards; no model or approval side effects."""
from dataclasses import dataclass
from pydantic import BaseModel, ConfigDict, Field, model_validator
from .models import Card, Family, MemoryNotFound
PROJECTION_FORMAT = 2
MAX_SEEDS = 100
MAX_DEPTH = 2
MAX_LINKS_PER_CARD = 20
MAX_VISITED = 200
MAX_EDGES = 400
RRF_K = 60
class RecallScope(BaseModel):
"""Exact business scope/concepts and hierarchical physical context.
Cards without dependencies are workspace-wide. A dependency applies to its
database and every descendant of the schema/table/column it names.
All physical fields must match the SAME dependency, never separate entries.
"""
model_config = ConfigDict(extra="forbid", str_strip_whitespace=True)
scope: str = Field(default="", max_length=10000)
database: str = Field(default="", max_length=200)
schema_name: str = Field(default="", max_length=200)
table: str = Field(default="", max_length=200)
column: str = Field(default="", max_length=200)
concepts: list[str] = Field(default_factory=list, max_length=100)
@model_validator(mode="after")
def validate_context(self):
if ((self.schema_name and not self.database) or (self.table and not self.schema_name)
or (self.column and not self.table)):
raise ValueError("Physical recall scope requires its database/schema/table ancestors")
if any(not value.strip() or len(value) > 200 for value in self.concepts):
raise ValueError("Recall concepts must contain between 1 and 200 characters")
return self
def matches(self, card: Card) -> bool:
if self.scope and self.scope != card.scope:
return False
if not set(self.concepts) <= set(card.concepts):
return False
if not self.database or not card.dependencies:
return True
return any(all(not getattr(self, key) or getattr(dep, key) in ("", getattr(self, key))
for key in ("database", "schema_name", "table", "column"))
for dep in card.dependencies)
def vector_filter(self, family: Family | None) -> dict:
return {"memory": {**self.model_dump(), "family": family,
"format": PROJECTION_FORMAT}}
@dataclass(frozen=True)
class RecalledCard:
card: Card
score: float
path: tuple[str, ...]
def expand_and_rank(repo, hits, *, scope: RecallScope, family: Family | None,
excluded: set[str], top: int) -> list[RecalledCard]:
"""RRF direct rank + strongest link path, decayed by 0.5 per outgoing hop.
Roots, nodes, fan-out and depth are all bounded. Repeated paths do not add
votes: cycles and highly connected cards cannot amplify their own relevance.
The caller holds the workspace operation lock while resolving authority.
"""
cache: dict[str, Card | None] = {}
def current(identity):
if identity not in cache:
if len(cache) >= MAX_VISITED:
return None
try:
card = repo.get(identity)
except MemoryNotFound:
card = None
if card is not None and (not card.indexed or card.id in excluded
or (family and card.family != family) or not scope.matches(card)):
card = None
cache[identity] = card
return cache[identity]
direct: dict[str, float] = {}
graph: dict[str, tuple[float, tuple[str, ...]]] = {}
seeds = []
for rank, hit in enumerate(hits[:MAX_SEEDS], 1):
card = current(hit.ref)
if (card is None or hit.metadata.get("memory_revision") != card.revision
or hit.metadata.get("memory_format") != PROJECTION_FORMAT
or card.id in direct):
continue
direct[card.id] = 1 / (RRF_K + rank)
seeds.append(card)
traversed = 0
for seed in seeds:
frontier = [(seed, (seed.id,))]
visited = {seed.id}
for depth in range(1, MAX_DEPTH + 1):
next_frontier = []
for source, path in frontier:
for link in sorted(source.links, key=lambda link: link.target_id)[:MAX_LINKS_PER_CARD]:
if traversed >= MAX_EDGES:
break
traversed += 1
if link.target_id in visited:
continue
visited.add(link.target_id)
target = current(link.target_id)
if target is None:
continue
target_path = (*path, target.id)
score = direct[seed.id] * 0.5 ** depth
previous = graph.get(target.id)
if previous is None or (-score, target_path) < (-previous[0], previous[1]):
graph[target.id] = (score, target_path)
next_frontier.append((target, target_path))
frontier = next_frontier
ranked = []
for identity in direct.keys() | graph.keys():
graph_score, path = graph.get(identity, (0, (identity,)))
ranked.append(RecalledCard(cache[identity], direct.get(identity, 0) + graph_score, path))
return sorted(ranked, key=lambda result: (-result.score, result.card.id))[:top]
+313
View File
@@ -0,0 +1,313 @@
"""Reviewer-edited additions/updates, grounded in effective approved session decisions."""
import hashlib
import json
from uuid import NAMESPACE_URL, uuid5
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import text
from tht.phase import current_phase, effective_decisions
from .models import CardInput, MemoryConflict
from .solved import _build_solved_snapshot
APPROVED_SOURCES = {
"concept_clarified",
"join_modified",
"column_corrected",
"concept_formula_approved",
"cte_corrected",
"cte_approved",
"sql_approved",
}
class Proposal(BaseModel):
model_config = ConfigDict(extra="forbid")
id: str = Field(pattern=r"^[a-zA-Z0-9_-]{1,80}$")
source_seqs: list[int] = Field(min_length=1, max_length=20)
card: CardInput
target_id: str | None = None
target_revision: str | None = None
reason: str = Field(min_length=1, max_length=10000)
class Selection(BaseModel):
model_config = ConfigDict(extra="forbid")
id: str
card: CardInput
class ReviewResponse(BaseModel):
model_config = ConfigDict(extra="forbid")
summary_id: str
items: list[Selection] = Field(max_length=20)
def digest(value):
return hashlib.sha256(
json.dumps(value, sort_keys=True, ensure_ascii=False).encode()
).hexdigest()
def context_hash(snapshot):
return digest(
{
"decisions": [
d.model_dump(mode="json")
for d in effective_decisions(snapshot)
if d.type != "memory_summary_reviewed"
],
"proposals": snapshot.artifacts.get("memory_proposals"),
"sql": snapshot.artifacts.get("sql_final"),
}
)
def solved_snapshot(snapshot):
linking = json.loads(snapshot.artifacts.get("schema_linking", "{}"))
tables = {
c["name"]
for c in linking.get("candidates", [])
if c.get("kind") == "table" and c.get("decision") == "promoted"
}
return _build_solved_snapshot(snapshot, tables)
def validate_proposals(snapshot, raw):
if not isinstance(raw, list) or len(raw) > 20:
raise ValueError("Memory proposals must be a list of at most 20 cards")
proposals = [Proposal.model_validate(p) for p in raw]
effective = {d.seq: d for d in effective_decisions(snapshot)}
if len({p.id for p in proposals}) != len(proposals):
raise ValueError("Memory proposal identities must be unique")
for p in proposals:
sources = [effective.get(seq) for seq in p.source_seqs]
if any(d is None or d.type not in APPROVED_SOURCES for d in sources):
raise MemoryConflict("Memory proposals require effective approved source decisions")
if p.card.family == "explained_error" and not any(d.rationale.strip() for d in sources):
raise MemoryConflict("Explained errors require an approved explanation")
if p.card.family == "solved_question":
solved = solved_snapshot(snapshot)
if p.card.sql != solved.metadata["sql"]:
raise MemoryConflict("Exemplar SQL must match the current approved solution")
if bool(p.target_id) != bool(p.target_revision):
raise ValueError("Updates require the identity and revision of the card being replaced")
return proposals
def prepare(service, snapshot):
service._session(snapshot)
if current_phase(snapshot) != 8 or snapshot.manifest.status in {"finalized", "archived"}:
raise MemoryConflict("The Memory summary is reviewed at the end of F8")
with service.repository.transaction() as connection:
receipt = (
connection.execute(
text(
"SELECT summary_id,result FROM thoth_memory.reviews "
"WHERE workspace_id=:w AND session_id=:s AND result->>'context_hash'=:h "
"ORDER BY created_at DESC LIMIT 1"
),
{
"w": service.repository.workspace_id,
"s": snapshot.manifest.id,
"h": context_hash(snapshot),
},
)
.mappings()
.first()
)
if receipt:
return {
"reviewed": True,
"summary_id": receipt["summary_id"],
"saved": len(receipt["result"]["saved"]),
}
raw = json.loads(snapshot.artifacts.get("memory_proposals", "[]"))
proposals = validate_proposals(snapshot, raw)
covered = {seq for p in proposals for seq in p.source_seqs}
for d in service.promotions(snapshot):
if d["decision_seq"] not in covered:
proposals.append(
Proposal(
id=f"decision-{d['decision_seq']}",
source_seqs=[d["decision_seq"]],
reason="Reusable domain clarification",
card=CardInput(
family="domain_clarification",
subject=d["subject"],
detail=d["detail"],
rationale=d["rationale"],
question=d["question_context"],
scope=service.repository.workspace_id,
concepts=[d["subject"]],
),
)
)
if not any(p.card.family == "solved_question" for p in proposals):
solved = solved_snapshot(snapshot)
approved = [d for d in effective_decisions(snapshot) if d.type == "sql_approved"][-1]
proposals.append(
Proposal(
id="solved-question",
source_seqs=[approved.seq],
reason="Approved solution, for consultation in future questions",
card=CardInput(
family="solved_question",
subject=solved.title,
question=solved.content,
sql=solved.metadata["sql"],
scope=service.repository.workspace_id,
dependencies=[
{
"database": snapshot.manifest.database,
"schema_name": snapshot.manifest.db_schema,
"table": table,
}
for table in solved.metadata["tables"]
],
),
)
)
if len(proposals) > 20:
raise MemoryConflict("Reduce the final Memory summary to at most 20 cards")
items = []
targets = set()
content = {}
aliases = {}
for p in proposals:
before = None
if p.target_id:
old = service.repository.get(p.target_id)
if old.revision != p.target_revision:
raise MemoryConflict(
"A Memory card changed; refresh the proposed update before review"
)
if p.target_id in targets:
raise MemoryConflict("Propose only one update to each Memory card")
targets.add(p.target_id)
before = old.model_dump(mode="json")
key = digest(p.card.model_dump(mode="json"))
duplicate = service.repository.exact_match(p.card)
if duplicate and (not p.target_id or duplicate.id == p.target_id):
aliases["proposal:" + p.id] = duplicate.id
continue
if not p.target_id and key in content:
aliases["proposal:" + p.id] = "proposal:" + content[key]
continue
content[key] = p.id
items.append({**p.model_dump(mode="json"), "before": before})
for item in items:
for link in item["card"]["links"]:
link["target_id"] = aliases.get(link["target_id"], link["target_id"])
return {
"summary_id": digest(
{
"items": items,
"decisions": [d.model_dump(mode="json") for d in effective_decisions(snapshot)],
}
),
"items": items,
}
def apply(service, snapshot, response: ReviewResponse):
service._session(snapshot)
request_hash = digest(response.model_dump(mode="json"))
session_id = snapshot.manifest.id
with service.repository.operation() as repo:
with repo.transaction() as connection:
params = {"w": repo.workspace_id, "s": session_id, "id": response.summary_id}
receipt = (
connection.execute(
text(
"SELECT request_hash,result FROM thoth_memory.reviews "
"WHERE workspace_id=:w AND session_id=:s AND summary_id=:id"
),
params,
)
.mappings()
.first()
)
if receipt:
if receipt["request_hash"] != request_hash:
raise MemoryConflict(
"This summary was already reviewed with different selections"
)
saved = receipt["result"]["saved"]
else:
# Resolve the same locked repository for preview and optimistic update checks.
original = service.repository
service.repository = repo
try:
summary = prepare(service, snapshot)
finally:
service.repository = original
if summary.get("reviewed") or summary["summary_id"] != response.summary_id:
raise MemoryConflict("The Memory summary changed; review it again")
choices = {choice.id: choice for choice in response.items}
candidates = {p["id"]: p for p in summary["items"]}
if len(choices) != len(response.items) or choices.keys() - candidates.keys():
raise ValueError("Memory review contains duplicate or unknown choices")
selected = validate_proposals(
snapshot,
[
{
**{k: v for k, v in candidates[key].items() if k != "before"},
"card": choice.card.model_dump(mode="json"),
}
for key, choice in choices.items()
],
)
identities = {
p.id: p.target_id
or "mem-"
+ str(uuid5(NAMESPACE_URL, f"thothii:{repo.workspace_id}:{session_id}:{p.id}"))
for p in selected
}
saved = []
for p in selected:
# Create/update every selected card before inserting links among new cards.
identity = repo.save(
p.card.model_copy(update={"links": []}),
card_id=p.target_id,
new_id=identities[p.id],
source_key=f"review:{session_id}:{p.id}",
session_id=session_id,
decision_seq=p.source_seqs[0],
)
saved.append({"id": identity, "proposal_id": p.id})
for p in selected:
value = p.card.model_dump(mode="json")
for link in value["links"]:
if link["target_id"].startswith("proposal:"):
target = link["target_id"].removeprefix("proposal:")
if target not in identities:
raise ValueError("Select the linked card or remove its link")
link["target_id"] = identities[target]
repo.save(CardInput.model_validate(value), card_id=identities[p.id])
connection.execute(
text(
"INSERT INTO thoth_memory.reviews "
"(workspace_id,session_id,summary_id,request_hash,result) "
"VALUES (:w,:s,:id,:hash,CAST(:result AS jsonb))"
),
{
**params,
"hash": request_hash,
"result": json.dumps(
{
"saved": saved,
"context_hash": context_hash(snapshot),
}
),
},
)
results = [service._propagate(repo, item["id"]) for item in saved]
return {
"saved": len(saved),
"declined": len(summary["items"]) - len(saved) if not receipt else None,
"indexed": all(r["indexed"] for r in results),
"results": results,
}
+63
View File
@@ -0,0 +1,63 @@
"""Installation binding for Memory; vector dependencies are opened only when needed."""
import os
from tht.session.models import PrincipalContext
from .migrate import installation_url
from .models import MemoryForbidden, MemoryUnavailable
from .repository import MemoryRepository
from .service import MemoryService
def _principal():
issuer = os.environ.get("THT_PRINCIPAL_ISSUER", "").strip()
subject = os.environ.get("THT_PRINCIPAL_SUBJECT", "").strip()
if not issuer or not subject:
raise MemoryForbidden("A trusted runtime principal is required for Memory")
return PrincipalContext(
issuer=issuer, subject=subject,
is_admin=os.environ.get("THT_PRINCIPAL_IS_ADMIN", "").lower() in {"1", "true"},
)
def _repository(workspace_id):
try:
url = installation_url()
except ValueError:
raise MemoryUnavailable("Memory PostgreSQL installation configuration is unavailable") \
from None
return MemoryRepository(url, workspace_id)
def memory_service(cfg):
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder
principal = _principal()
return MemoryService(_repository(cfg._workspace_id), principal,
language=cfg.language,
store_factory=lambda: build_vector_store(cfg, require_write=True),
embedder_factory=lambda: make_embedder(cfg.embeddings))
def admin_service(workspace_id, runtime):
"""Admin access needs no DWH binding, active session or Evidence materialization."""
from tht.adapters.vector.qdrant import QdrantVectorStore
from tht.config import EmbeddingsConfig
from tht.vectorstore.embeddings import OllamaEmbeddings
principal = _principal()
if not principal.is_admin:
raise MemoryForbidden("Memory administration requires an administrator")
return MemoryService(_repository(workspace_id), principal,
language=runtime.get("memoryLanguage", "en"),
store_factory=lambda: QdrantVectorStore(
base_url=runtime["internalQdrantUrl"], workspace_id=workspace_id,
collections={"reference": workspace_id+"-reference", "memory": workspace_id+"-memory"},
expected_dimension=runtime["internalEmbeddingDimensions"],
),
embedder_factory=lambda: OllamaEmbeddings(EmbeddingsConfig(
base_url=runtime["internalEmbeddingUrl"], model=runtime["internalEmbeddingModel"],
dimensions=runtime["internalEmbeddingDimensions"], timeout=30,
)))
+240
View File
@@ -0,0 +1,240 @@
"""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()
self.repository.list(CardQuery(page_size=1))
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"])