From ba16f12804d2e61d0903715e538cdb01b7fab981 Mon Sep 17 00:00:00 2001 From: lipluscodex <268560960+lipluscodex@users.noreply.github.com> Date: Sat, 22 Aug 2026 23:59:39 +0900 Subject: [PATCH 1/2] feat(decisions): make SQLite the canonical judgment graph MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 判断の stable identity、revision、lifecycle、typed relation を SQLite に保存し、原子的な domain write API と MCP surface を追加する。論理的忘却と制限付き物理削除を分離する。 決定論的 import/export、SQLite backup/restore、integrity 検証、および lifecycle・atomicity・round-trip・既存 contract 回帰テストを追加する。 --- docs/Decision-Structure.md | 4 + docs/canonical-sqlite-judgment-graph.md | 27 ++ docs/requirements.md | 2 + src/neuron_graph_rag/__init__.py | 3 + src/neuron_graph_rag/engine.py | 2 + src/neuron_graph_rag/judgments.py | 321 ++++++++++++++++++++++++ src/neuron_graph_rag/storage.py | 53 +++- src/neuron_graph_rag_mcp/server.py | 86 ++++++- tests/test_judgments.py | 100 ++++++++ tests/test_mcp_adapter.py | 51 +++- tools/judgment_graph.py | 100 ++++++++ 11 files changed, 741 insertions(+), 8 deletions(-) create mode 100644 docs/canonical-sqlite-judgment-graph.md create mode 100644 src/neuron_graph_rag/judgments.py create mode 100644 tests/test_judgments.py create mode 100644 tools/judgment_graph.py diff --git a/docs/Decision-Structure.md b/docs/Decision-Structure.md index c708cd6..cd8df6a 100644 --- a/docs/Decision-Structure.md +++ b/docs/Decision-Structure.md @@ -21,6 +21,7 @@ | [repository-native-controlled-corpus](https://github.com/Liplus-Project/neuron-graph-rag/wiki/repository-native-controlled-corpus) | active | repository-native controlled corpus v2 は、固定 SHA の公開 documentation と本文中の明示的な同一 directory 相対 link だけから、node、doc path、source URL、credited edge identity が相互に分離した development / holdout の各 3-edge path を導出する。v1 は provenance として保持する。これは controlled benchmark であり、外部 corpus への一般化、評価 query、gold、result、既定値変更を含まない。 | | [soft-start-feedback-reinforcement](https://github.com/Liplus-Project/neuron-graph-rag/wiki/soft-start-feedback-reinforcement) | active | 最初の credited `used` に通常 bounded updateの小さなprovisional fractionを適用し、最初の独立`confirmed`がremainder、後続confirmationがgeometric decayを適用する。v1 snapshot評価の不支持を保持し、baseline-aware successorはfresh initial evidenceからq3 first mutationを導出する。v2 freeze-only PRとsquash後のobserved registrationを分離し、development全gate通過時だけholdoutを一度開く。source database、live config、defaultを変更しない。 | | [github-rag-mcp-replacement-compatibility](https://github.com/Liplus-Project/neuron-graph-rag/wiki/github-rag-mcp-replacement-compatibility) | active | public GitHub repository一つのread-only snapshotをNGR local indexへ接続する。github-rag-mcp `search` の保存済み raw capture と source URL、根拠を比較する。共有 source identity を確認しても最小 doc 検索 path の候補に限り、production github-rag-mcp、MCP authentication / transport、remote deployment、default変更は含まない。 | +| [sqlite-canonical-judgment-graph](https://github.com/Liplus-Project/neuron-graph-rag/wiki/sqlite-canonical-judgment-graph) | active | NGR 自身の判断構造は SQLite の stable identity、revision、lifecycle、typed relation を machine-native 正本とし、raw SQL でなく atomic domain API で変更する。Wiki は移行 fixture の検証後に optional generated view へ下げる。 | ## Entry format @@ -50,6 +51,9 @@ edge target は Decision Structure node slug とする。外部資料は `Edges` ## Source-of-truth boundary +- 新規に domain API から登録された judgment graph の machine-readable 正本は SQLite である。詳細契約は [Canonical SQLite judgment graph](canonical-sqlite-judgment-graph.md) に置く。 +- 既存 Wiki entry は検証済み import が行われるまで従来の正本境界を維持する。本変更だけで本番 Wiki entry を自動移行または削除しない。 + - `docs/Decision-Structure.md` は main repository における索引、format、vocabulary、所有境界、lifecycle の正本である。これは docs-owned であり、GitHub Wiki の `Decision-Structure.md` へ mirror する。 - lowercase kebab-case の Decision Structure entry と `_Sidebar.md` は Wiki-only である。docs-to-Wiki synchronization は、`docs/` に対応物がないことを理由にこれらを create、overwrite、delete しない。 - 個別 Wiki entry は、その判断の current state の正本である。GitHub issue、pull request、commit、test output、fixture、gold、manifest、gate、result artifact は、それぞれの所有境界に従う根拠または契約であり、リンクしただけで Decision Structure entry にはならない。 diff --git a/docs/canonical-sqlite-judgment-graph.md b/docs/canonical-sqlite-judgment-graph.md new file mode 100644 index 0000000..18c75b1 --- /dev/null +++ b/docs/canonical-sqlite-judgment-graph.md @@ -0,0 +1,27 @@ +# Canonical SQLite judgment graph + +## 正本境界 + +NGR の新規 judgment は SQLite の `judgments`、`judgment_revisions`、`judgment_relations` を machine-readable 正本とする。`nodes` と `edges` は既存 retrieval contract を保つ検索 projection であり、判断の lifecycle や provenance の正本ではない。既存 Wiki entry の本番移行は fixture による import / export 検証後の別操作とし、この実装は自動移行しない。 + +## Domain API + +`NeuronGraphRAG.judgments` は add、update、supersede、archive、restore、hard delete を transaction 単位で提供する。update、supersede、archive、restore、hard delete は `expected_revision` による楽観的 concurrency check を要求する。relation target 不在、部分更新、stale revision、同一 predecessor の再 supersede は transaction 全体を失敗させる。 + +supersede は新しい stable identity を作り、旧判断を archive し、successor から predecessor への `supersedes` relation と predecessor の `superseded_by` を保持する。superseded judgment は restore できない。 + +archive は通常 retrieval から外す論理的忘却であり、revision、provenance、relation は監査 API から取得できる。hard delete は明示操作で、archived、revision 一致、inbound relation と successor history がない候補にだけ許可する。 + +MCP の `write_judgment` は同じ domain API へ写像し、model に raw SQL を公開しない。既存 `search`、`record_source_use`、`record_outcome` の contract と既定値は変更しない。 + +## Portability and recovery + +`tools/judgment_graph.py` は次を提供する。 + +- `export SOURCE_DB OUTPUT_JSON`: key と judgment / relation 順序を固定した UTF-8 LF JSON を出力する。 +- `import INPUT_JSON DESTINATION_DB`: 全 identity と relation target を事前検証し、一つの transaction で再構築する。 +- `backup SOURCE_DB BACKUP_DB`: SQLite backup API で transaction-consistent copy を作る。既存出力は上書きしない。 +- `restore BACKUP_DB DESTINATION_DB`: integrity check 後に新規 destination へ復元する。既存 database は上書きしない。 +- `integrity DATABASE`: SQLite integrity、foreign key、dangling relation、二重 successor を fail closed に検査する。 + +backup は export より多くの revision history と retrieval / feedback state を保持するため、完全復旧の正本である。export は current judgment graph の決定論的 portability surface である。 diff --git a/docs/requirements.md b/docs/requirements.md index 7c58c26..b8ab00f 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -104,6 +104,7 @@ 85. baseline-aware soft-start snapshot evaluation は、v1のprotocol、gate、observed result、private snapshotを変更、再実行、再集計、入力再利用せず、fresh snapshot、新規namespace、新規output、v1 observed developmentと異なるcredited edge identityを使う。各relation caseのinitial weight、reinforced count、evidence count、confirmation countを結果前に登録し、q3/s1のfirst mutationを`max(1, quorum - initial evidence count)`で導出する。導出不能、baseline不一致、event budget内のquorum capacity不足はregistered resultを作らずfailure reportで停止する。v2 protocolはregistered output不在のfreeze-only PRで固定し、そのsquash merge後の別Issueでdevelopmentを一度だけ実行する。全8 hard gate通過時だけholdoutを一度開き、支持結果もlocal cutover候補に限定してsource database、live config、library defaultを変更しない。 86. outcome-driven feedback deactivation はsoft-startと同時にだけ有効化できるdefault-off candidateとする。provisional / confirmationごとにcredited加算と同時発生したsame-source sibling normalization減算を一つのsigned mutation journalへ永続化し、因果帰属できる`corrected` / `rolled_back`だけが未反転contributionを基礎weight未満へ下げずexact reversalする。`superseded`はedge、evidence、trace、outcomeを削除せずrelation edgeをdormantにして通常activationから除外し、同じ保存済みcredited pathの後続`confirmed`で再活性化する。duplicate、retry、restart、transaction failure、lexical、zero-hop、別source、uncredited edge、因果帰属不能outcomeは二重減算または局所外mutationを行わない。 87. outcome-driven deactivation evaluation はcontrol / candidate、`corrected` / `rolled_back` / `superseded`、exact credited / sibling inverse、baseline floor、dormancy / reactivation、rank / locality、source isolation、exclusive outputを結果観測前に固定する。protocolはregistered output不在のfreeze-only PRで固定し、そのsquash merge後のsuccessor Issueでdevelopmentを一度だけ実行する。全hard gate通過時だけholdoutを一度開き、観測前後にquery、case、schedule、metric、gate、default、live configを変更しない。 +88. NGR 自身の新規 judgment graph は SQLite の stable identity、revision、lifecycle、provenance、typed relation を machine-readable 正本とする。add / update / supersede / archive / restore / hard-delete candidate は raw SQL でなく atomic domain API を通し、stale revision、dangling relation、部分更新、二重 successor を fail closed にする。archive は通常 retrieval から外す論理的忘却、hard delete は履歴参照のない archived candidate だけに許す物理削除として分離する。current graph の deterministic export / atomic import と SQLite backup / integrity-checked restore を維持し、既存 Wiki entry の本番移行は fixture 検証後に分離する。 75. v3 implementation、prompt、manifest、query override、schema、集約、path audit、hash規則、gate、stop rule、testsをresult-free commitでpushした後、development stage / 4 case packet / 12 responses / resultを各一度だけ生成する。 76. development全12 gate通過時だけholdout stageを一度生成し、異なるfresh 12 judgesで同じgateを評価する。packet、response、resultの上書き、観測後の規則変更、実LLM品質値のCI再生成を拒否する。 @@ -143,3 +144,4 @@ - [Soft-start snapshot evaluation](soft-start-snapshot-evaluation.md) が transaction-consistent private snapshot、固定4 arm、result-free freeze、privacy、one-time development、conditional holdout、local cutover境界を定義する。 - [Outcome-driven feedback deactivation](outcome-driven-feedback-deactivation.md) がsigned contribution journal、exact reversal、dormancy / reactivation、default-off境界、result-free freezeを定義する。 - [Baseline-aware soft-start snapshot evaluation](baseline-aware-soft-start-snapshot-evaluation.md) がfresh baseline stateからのq3 boundary導出、v1 evidence isolation、capacity preflight、新規one-time result境界を定義する。 +- [Canonical SQLite judgment graph](canonical-sqlite-judgment-graph.md) が judgment source-of-truth、domain write API、logical forgetting、hard-delete boundary、deterministic portability、backup / restore を定義する。 diff --git a/src/neuron_graph_rag/__init__.py b/src/neuron_graph_rag/__init__.py index 71fcd2d..5575349 100644 --- a/src/neuron_graph_rag/__init__.py +++ b/src/neuron_graph_rag/__init__.py @@ -1,5 +1,6 @@ from .engine import EngineConfig, NeuronGraphRAG from .feedback import FeedbackLedger +from .judgments import JudgmentContractError, JudgmentGraph from .models import ( ActivationPath, ConfirmedEdge, @@ -39,6 +40,8 @@ "FeedbackEvidence", "FeedbackLedger", "FeedbackReceipt", + "JudgmentContractError", + "JudgmentGraph", "NeuronGraphRAG", "NormalizedSiblingEdge", "OutcomeReceipt", diff --git a/src/neuron_graph_rag/engine.py b/src/neuron_graph_rag/engine.py index df302b6..ac8dc3f 100644 --- a/src/neuron_graph_rag/engine.py +++ b/src/neuron_graph_rag/engine.py @@ -23,6 +23,7 @@ from .dynamics import DynamicsSettings, propagate from .retrieval import BM25Retriever, DenseEncoder, DenseRetriever, normalize_scores from .storage import SQLiteStore +from .judgments import JudgmentGraph @dataclass(frozen=True, slots=True) @@ -134,6 +135,7 @@ def __init__( ) -> None: self.config = config or EngineConfig() self.store = SQLiteStore(database) + self.judgments = JudgmentGraph(self.store) self.sparse_retriever = BM25Retriever() self.dense_retriever = DenseRetriever(dense_encoder) diff --git a/src/neuron_graph_rag/judgments.py b/src/neuron_graph_rag/judgments.py new file mode 100644 index 0000000..0ceaaee --- /dev/null +++ b/src/neuron_graph_rag/judgments.py @@ -0,0 +1,321 @@ +from __future__ import annotations + +import json +import re +from datetime import UTC, datetime +from typing import Any, Iterable + +from .models import DocumentNode, TypedEdge +from .storage import SQLiteStore + + +_ID = re.compile(r"^[a-z0-9][a-z0-9._:-]{0,127}$") +_RELATION = re.compile(r"^[a-z][a-z0-9_-]{0,63}$") + + +class JudgmentContractError(ValueError): + """A fail-closed domain contract violation.""" + + +class JudgmentGraph: + """Atomic domain API for the canonical SQLite judgment graph.""" + + def __init__(self, store: SQLiteStore) -> None: + self.store = store + + @staticmethod + def _now() -> str: + return datetime.now(UTC).isoformat(timespec="microseconds").replace("+00:00", "Z") + + @staticmethod + def _identity(value: str) -> str: + if not isinstance(value, str) or _ID.fullmatch(value) is None: + raise JudgmentContractError("invalid judgment identity") + return value + + @staticmethod + def _text(value: str, name: str) -> str: + if not isinstance(value, str) or not value.strip(): + raise JudgmentContractError(f"{name} is required") + return value.strip() + + @staticmethod + def _provenance(value: dict[str, Any]) -> dict[str, Any]: + if not isinstance(value, dict) or not value: + raise JudgmentContractError("provenance must be a non-empty object") + try: + json.dumps(value, sort_keys=True, ensure_ascii=False, allow_nan=False) + except (TypeError, ValueError) as error: + raise JudgmentContractError("provenance must be deterministic JSON") from error + return value + + @staticmethod + def _relations( + relations: Iterable[dict[str, str]], *, source_id: str + ) -> tuple[tuple[str, str], ...]: + normalized: list[tuple[str, str]] = [] + for relation in relations: + if not isinstance(relation, dict) or set(relation) != {"target_id", "relation_type"}: + raise JudgmentContractError("relation requires target_id and relation_type") + target = JudgmentGraph._identity(relation["target_id"]) + kind = relation["relation_type"] + if not isinstance(kind, str) or _RELATION.fullmatch(kind) is None: + raise JudgmentContractError("invalid relation type") + if target == source_id: + raise JudgmentContractError("self relations are not allowed") + normalized.append((target, kind)) + if len(set(normalized)) != len(normalized): + raise JudgmentContractError("duplicate relation") + return tuple(sorted(normalized)) + + def add( + self, + judgment_id: str, + statement: str, + rationale: str, + provenance: dict[str, Any], + *, + relations: Iterable[dict[str, str]] = (), + ) -> dict[str, Any]: + judgment_id = self._identity(judgment_id) + statement = self._text(statement, "statement") + rationale = self._text(rationale, "rationale") + provenance = self._provenance(provenance) + normalized = self._relations(relations, source_id=judgment_id) + now = self._now() + with self.store.transaction() as connection: + if connection.execute( + "SELECT 1 FROM judgments WHERE judgment_id = ?", (judgment_id,) + ).fetchone(): + raise JudgmentContractError("judgment already exists") + self._require_targets(connection, normalized) + connection.execute( + "INSERT INTO judgments VALUES (?, 1, 'active', NULL, ?, ?)", + (judgment_id, now, now), + ) + self._insert_revision(connection, judgment_id, 1, statement, rationale, provenance, now) + self._replace_relations(connection, judgment_id, normalized, now) + self._sync_retrieval(connection, judgment_id, statement, rationale, provenance, normalized) + return self.get(judgment_id) + + def update( + self, + judgment_id: str, + statement: str, + rationale: str, + provenance: dict[str, Any], + *, + expected_revision: int, + relations: Iterable[dict[str, str]] = (), + ) -> dict[str, Any]: + judgment_id = self._identity(judgment_id) + statement = self._text(statement, "statement") + rationale = self._text(rationale, "rationale") + provenance = self._provenance(provenance) + normalized = self._relations(relations, source_id=judgment_id) + now = self._now() + with self.store.transaction() as connection: + row = self._require(connection, judgment_id) + if row["lifecycle"] != "active" or row["superseded_by"] is not None: + raise JudgmentContractError("only current active judgments can be updated") + if expected_revision != row["current_revision"]: + raise JudgmentContractError("stale expected_revision") + self._require_targets(connection, normalized) + revision = expected_revision + 1 + self._insert_revision(connection, judgment_id, revision, statement, rationale, provenance, now) + connection.execute( + "UPDATE judgments SET current_revision = ?, updated_at = ? WHERE judgment_id = ?", + (revision, now, judgment_id), + ) + self._replace_relations(connection, judgment_id, normalized, now) + self._sync_retrieval(connection, judgment_id, statement, rationale, provenance, normalized) + return self.get(judgment_id) + + def supersede( + self, + predecessor_id: str, + successor_id: str, + statement: str, + rationale: str, + provenance: dict[str, Any], + *, + expected_revision: int, + relations: Iterable[dict[str, str]] = (), + ) -> dict[str, Any]: + predecessor_id = self._identity(predecessor_id) + successor_id = self._identity(successor_id) + if predecessor_id == successor_id: + raise JudgmentContractError("successor must have a new identity") + statement = self._text(statement, "statement") + rationale = self._text(rationale, "rationale") + provenance = self._provenance(provenance) + normalized = self._relations(relations, source_id=successor_id) + normalized = tuple(sorted((*normalized, (predecessor_id, "supersedes")))) + now = self._now() + with self.store.transaction() as connection: + old = self._require(connection, predecessor_id) + if old["lifecycle"] != "active" or old["superseded_by"] is not None: + raise JudgmentContractError("predecessor already inactive or superseded") + if old["current_revision"] != expected_revision: + raise JudgmentContractError("stale expected_revision") + if connection.execute("SELECT 1 FROM judgments WHERE judgment_id = ?", (successor_id,)).fetchone(): + raise JudgmentContractError("successor already exists") + self._require_targets(connection, tuple(item for item in normalized if item[0] != predecessor_id)) + connection.execute( + "INSERT INTO judgments VALUES (?, 1, 'active', NULL, ?, ?)", + (successor_id, now, now), + ) + self._insert_revision(connection, successor_id, 1, statement, rationale, provenance, now) + self._replace_relations(connection, successor_id, normalized, now) + self._sync_retrieval(connection, successor_id, statement, rationale, provenance, normalized) + connection.execute( + "UPDATE judgments SET lifecycle = 'archived', superseded_by = ?, updated_at = ? WHERE judgment_id = ?", + (successor_id, now, predecessor_id), + ) + return self.get(successor_id) + + def archive(self, judgment_id: str, *, expected_revision: int) -> dict[str, Any]: + return self._lifecycle(judgment_id, expected_revision, "archived") + + def restore(self, judgment_id: str, *, expected_revision: int) -> dict[str, Any]: + return self._lifecycle(judgment_id, expected_revision, "active") + + def _lifecycle(self, judgment_id: str, expected_revision: int, lifecycle: str) -> dict[str, Any]: + judgment_id = self._identity(judgment_id) + with self.store.transaction() as connection: + row = self._require(connection, judgment_id) + if row["current_revision"] != expected_revision: + raise JudgmentContractError("stale expected_revision") + if lifecycle == "active" and row["superseded_by"] is not None: + raise JudgmentContractError("superseded judgments cannot be restored") + connection.execute( + "UPDATE judgments SET lifecycle = ?, updated_at = ? WHERE judgment_id = ?", + (lifecycle, self._now(), judgment_id), + ) + return self.get(judgment_id) + + def hard_delete(self, judgment_id: str, *, expected_revision: int) -> None: + judgment_id = self._identity(judgment_id) + with self.store.transaction() as connection: + row = self._require(connection, judgment_id) + if row["current_revision"] != expected_revision or row["lifecycle"] != "archived": + raise JudgmentContractError("hard delete requires an archived current revision") + inbound = connection.execute( + "SELECT 1 FROM judgment_relations WHERE target_id = ? LIMIT 1", (judgment_id,) + ).fetchone() + if inbound or row["superseded_by"] is not None: + raise JudgmentContractError("hard delete candidate has retained graph history") + connection.execute("DELETE FROM nodes WHERE node_id = ?", (judgment_id,)) + connection.execute("DELETE FROM judgments WHERE judgment_id = ?", (judgment_id,)) + + def get(self, judgment_id: str) -> dict[str, Any]: + judgment_id = self._identity(judgment_id) + row = self.store.connection.execute( + """ + SELECT j.*, r.statement, r.rationale, r.provenance_json + FROM judgments j JOIN judgment_revisions r + ON r.judgment_id = j.judgment_id AND r.revision = j.current_revision + WHERE j.judgment_id = ? + """, (judgment_id,), + ).fetchone() + if row is None: + raise KeyError(f"Unknown judgment: {judgment_id}") + relations = self.store.connection.execute( + "SELECT target_id, relation_type FROM judgment_relations WHERE source_id = ? ORDER BY target_id, relation_type", + (judgment_id,), + ).fetchall() + return { + "judgment_id": row["judgment_id"], "revision": row["current_revision"], + "statement": row["statement"], "rationale": row["rationale"], + "provenance": json.loads(row["provenance_json"]), "lifecycle": row["lifecycle"], + "superseded_by": row["superseded_by"], + "relations": [dict(item) for item in relations], + } + + def export(self) -> dict[str, Any]: + ids = [row[0] for row in self.store.connection.execute("SELECT judgment_id FROM judgments ORDER BY judgment_id")] + return {"format": "ngr-judgment-graph/v1", "judgments": [self.get(item) for item in ids]} + + def import_graph(self, payload: dict[str, Any]) -> None: + if not isinstance(payload, dict) or payload.get("format") != "ngr-judgment-graph/v1": + raise JudgmentContractError("unsupported judgment graph format") + records = payload.get("judgments") + if not isinstance(records, list): + raise JudgmentContractError("judgments must be a list") + ids = {record.get("judgment_id") for record in records if isinstance(record, dict)} + if len(ids) != len(records) or None in ids: + raise JudgmentContractError("judgment identities must be unique") + for record in records: + for relation in record.get("relations", []): + if relation.get("target_id") not in ids: + raise JudgmentContractError("dangling relation") + now = self._now() + with self.store.transaction() as connection: + for record in sorted(records, key=lambda item: item["judgment_id"]): + judgment_id = self._identity(record["judgment_id"]) + revision = record.get("revision") + if isinstance(revision, bool) or not isinstance(revision, int) or revision < 1: + raise JudgmentContractError("invalid revision") + statement = self._text(record["statement"], "statement") + rationale = self._text(record["rationale"], "rationale") + provenance = self._provenance(record["provenance"]) + lifecycle = record.get("lifecycle", "active") + if lifecycle not in {"active", "archived"}: + raise JudgmentContractError("invalid lifecycle") + connection.execute( + "INSERT INTO judgments VALUES (?, ?, ?, NULL, ?, ?)", + (judgment_id, revision, lifecycle, now, now), + ) + self._insert_revision(connection, judgment_id, revision, statement, rationale, provenance, now) + self._sync_retrieval(connection, judgment_id, statement, rationale, provenance, ()) + for record in sorted(records, key=lambda item: item["judgment_id"]): + relations = self._relations(record.get("relations", []), source_id=record["judgment_id"]) + self._require_targets(connection, relations) + self._replace_relations(connection, record["judgment_id"], relations, now) + self._sync_retrieval( + connection, record["judgment_id"], record["statement"], record["rationale"], + record["provenance"], relations, + ) + for record in sorted(records, key=lambda item: item["judgment_id"]): + successor = record.get("superseded_by") + if successor is not None and successor not in ids: + raise JudgmentContractError("unknown successor") + connection.execute( + "UPDATE judgments SET superseded_by = ? WHERE judgment_id = ?", + (successor, record["judgment_id"]), + ) + + @staticmethod + def _require(connection: Any, judgment_id: str) -> Any: + row = connection.execute("SELECT * FROM judgments WHERE judgment_id = ?", (judgment_id,)).fetchone() + if row is None: + raise KeyError(f"Unknown judgment: {judgment_id}") + return row + + @staticmethod + def _require_targets(connection: Any, relations: tuple[tuple[str, str], ...]) -> None: + for target, _ in relations: + if connection.execute("SELECT 1 FROM judgments WHERE judgment_id = ?", (target,)).fetchone() is None: + raise JudgmentContractError(f"unknown relation target: {target}") + + @staticmethod + def _insert_revision(connection: Any, judgment_id: str, revision: int, statement: str, rationale: str, provenance: dict[str, Any], now: str) -> None: + connection.execute( + "INSERT INTO judgment_revisions VALUES (?, ?, ?, ?, ?, ?)", + (judgment_id, revision, statement, rationale, json.dumps(provenance, sort_keys=True, ensure_ascii=False, separators=(",", ":")), now), + ) + + @staticmethod + def _replace_relations(connection: Any, source_id: str, relations: tuple[tuple[str, str], ...], now: str) -> None: + connection.execute("DELETE FROM judgment_relations WHERE source_id = ?", (source_id,)) + connection.executemany("INSERT INTO judgment_relations VALUES (?, ?, ?, ?)", [(source_id, target, kind, now) for target, kind in relations]) + + @staticmethod + def _sync_retrieval(connection: Any, judgment_id: str, statement: str, rationale: str, provenance: dict[str, Any], relations: tuple[tuple[str, str], ...]) -> None: + metadata = {"kind": "judgment", "provenance": provenance} + connection.execute( + "INSERT INTO nodes VALUES (?, ?, ?, 1.0) ON CONFLICT(node_id) DO UPDATE SET text=excluded.text, metadata_json=excluded.metadata_json, confidence=1.0", + (judgment_id, f"{statement}\n\n{rationale}", json.dumps(metadata, sort_keys=True, ensure_ascii=False)), + ) + connection.execute("DELETE FROM edges WHERE source_id = ?", (judgment_id,)) + connection.executemany("INSERT INTO edges(source_id,target_id,edge_type,weight,factuality) VALUES (?,?,?,1.0,1.0)", [(judgment_id, target, kind) for target, kind in relations]) diff --git a/src/neuron_graph_rag/storage.py b/src/neuron_graph_rag/storage.py index a44c6e2..88dcbd1 100644 --- a/src/neuron_graph_rag/storage.py +++ b/src/neuron_graph_rag/storage.py @@ -279,6 +279,34 @@ def _create_schema(self) -> None: FOREIGN KEY (source_id, target_id, edge_type) REFERENCES edges(source_id, target_id, edge_type) ON DELETE CASCADE ); + + CREATE TABLE IF NOT EXISTS judgments ( + judgment_id TEXT PRIMARY KEY, + current_revision INTEGER NOT NULL CHECK(current_revision >= 1), + lifecycle TEXT NOT NULL CHECK(lifecycle IN ('active', 'archived')), + superseded_by TEXT UNIQUE REFERENCES judgments(judgment_id), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + ); + + CREATE TABLE IF NOT EXISTS judgment_revisions ( + judgment_id TEXT NOT NULL REFERENCES judgments(judgment_id) ON DELETE CASCADE, + revision INTEGER NOT NULL CHECK(revision >= 1), + statement TEXT NOT NULL, + rationale TEXT NOT NULL, + provenance_json TEXT NOT NULL, + created_at TEXT NOT NULL, + PRIMARY KEY (judgment_id, revision) + ); + + CREATE TABLE IF NOT EXISTS judgment_relations ( + source_id TEXT NOT NULL REFERENCES judgments(judgment_id) ON DELETE CASCADE, + target_id TEXT NOT NULL REFERENCES judgments(judgment_id) ON DELETE RESTRICT, + relation_type TEXT NOT NULL, + created_at TEXT NOT NULL, + PRIMARY KEY (source_id, target_id, relation_type), + CHECK(source_id <> target_id) + ); """ ) self.connection.commit() @@ -332,12 +360,28 @@ def get_node(self, node_id: str) -> DocumentNode: return self._node_from_row(row) def list_nodes(self) -> list[DocumentNode]: - rows = self.connection.execute("SELECT * FROM nodes ORDER BY node_id").fetchall() + rows = self.connection.execute( + """ + SELECT nodes.* FROM nodes + LEFT JOIN judgments ON judgments.judgment_id = nodes.node_id + WHERE judgments.judgment_id IS NULL OR judgments.lifecycle = 'active' + ORDER BY nodes.node_id + """ + ).fetchall() return [self._node_from_row(row) for row in rows] def list_edges(self) -> list[TypedEdge]: rows = self.connection.execute( - "SELECT * FROM edges ORDER BY source_id, target_id, edge_type" + """ + SELECT edges.* FROM edges + LEFT JOIN judgments AS source_judgment + ON source_judgment.judgment_id = edges.source_id + LEFT JOIN judgments AS target_judgment + ON target_judgment.judgment_id = edges.target_id + WHERE (source_judgment.judgment_id IS NULL OR source_judgment.lifecycle = 'active') + AND (target_judgment.judgment_id IS NULL OR target_judgment.lifecycle = 'active') + ORDER BY edges.source_id, edges.target_id, edges.edge_type + """ ).fetchall() return [self._edge_from_row(row) for row in rows] @@ -350,6 +394,11 @@ def outgoing_edges(self, node_id: str) -> list[TypedEdge]: AND dormancy.target_id = edges.target_id AND dormancy.edge_type = edges.edge_type WHERE edges.source_id = ? AND COALESCE(dormancy.dormant, 0) = 0 + AND NOT EXISTS ( + SELECT 1 FROM judgments + WHERE judgments.judgment_id = edges.target_id + AND judgments.lifecycle <> 'active' + ) ORDER BY edges.target_id, edges.edge_type """, (node_id,), diff --git a/src/neuron_graph_rag_mcp/server.py b/src/neuron_graph_rag_mcp/server.py index de40e53..ed9cd70 100644 --- a/src/neuron_graph_rag_mcp/server.py +++ b/src/neuron_graph_rag_mcp/server.py @@ -152,6 +152,28 @@ def _object(properties: dict[str, Any], required: list[str]) -> dict[str, Any]: }, ["contract_version", "idempotency_key", "trace_id", "node_ids", "outcome", "summary"], ) +JUDGMENT_WRITE_INPUT = { + "type": "object", + "properties": { + "contract_version": {"type": "string", "const": CONTRACT_VERSION}, + "action": {"type": "string", "enum": ["add", "update", "supersede", "archive", "restore", "hard_delete"]}, + "judgment_id": {"type": "string", "minLength": 1, "maxLength": 128}, + "successor_id": {"type": "string", "minLength": 1, "maxLength": 128}, + "statement": {"type": "string", "minLength": 1}, + "rationale": {"type": "string", "minLength": 1}, + "provenance": {"type": "object"}, + "expected_revision": {"type": "integer", "minimum": 1}, + "relations": { + "type": "array", + "items": _object( + {"target_id": {"type": "string"}, "relation_type": {"type": "string"}}, + ["target_id", "relation_type"], + ), + }, + }, + "required": ["contract_version", "action", "judgment_id"], + "additionalProperties": False, +} _STEP_OUTPUT = _object( { @@ -412,6 +434,27 @@ def _annotations(*, idempotent: bool) -> types.ToolAnnotations: output_schema=OUTCOME_OUTPUT, annotations=_annotations(idempotent=True), ), + types.Tool( + name="write_judgment", + description=( + "Atomically add, update, supersede, archive, restore, or explicitly hard-delete " + "a canonical SQLite judgment. Never use raw SQL. Updates and lifecycle changes " + "require expected_revision; hard delete is restricted to safe archived candidates." + ), + input_schema=JUDGMENT_WRITE_INPUT, + output_schema=_object( + { + "contract_version": {"type": "string"}, + "action": {"type": "string"}, + "judgment": {"type": ["object", "null"]}, + }, + ["contract_version", "action", "judgment"], + ), + annotations=types.ToolAnnotations( + read_only_hint=False, destructive_hint=True, idempotent_hint=False, + open_world_hint=False, + ), + ), ) @@ -453,6 +496,7 @@ def _tools( output_schema=OUTCOME_OUTPUT, annotations=_annotations(idempotent=True), ), + TOOLS[3], ) @@ -490,8 +534,10 @@ async def call_tool(self, _context: object, params: types.CallToolRequestParams) output = self._search(arguments) elif params.name == "record_source_use": output = self._record_source_use(arguments) - else: + elif params.name == "record_outcome": output = self._record_outcome(arguments) + else: + output = self._write_judgment(arguments) return self._success(output) except FeedbackContractError as error: return self._error(error.code, str(error), error.retryable) @@ -760,6 +806,44 @@ def _record_outcome(self, data: dict[str, Any]) -> dict[str, Any]: ], } + def _write_judgment(self, data: dict[str, Any]) -> dict[str, Any]: + allowed = { + "contract_version", "action", "judgment_id", "successor_id", + "statement", "rationale", "provenance", "expected_revision", "relations", + } + self._keys(data, allowed, {"contract_version", "action", "judgment_id"}) + self._version(data) + action = data["action"] + identifier = data["judgment_id"] + graph = self.engine.judgments + if action == "add": + required = {"statement", "rationale", "provenance"} + if not required <= data.keys(): + raise ValueError("add requires statement, rationale, and provenance") + judgment = graph.add(identifier, data["statement"], data["rationale"], data["provenance"], relations=data.get("relations", [])) + elif action == "update": + required = {"statement", "rationale", "provenance", "expected_revision"} + if not required <= data.keys(): + raise ValueError("update requires content, provenance, and expected_revision") + judgment = graph.update(identifier, data["statement"], data["rationale"], data["provenance"], expected_revision=data["expected_revision"], relations=data.get("relations", [])) + elif action == "supersede": + required = {"successor_id", "statement", "rationale", "provenance", "expected_revision"} + if not required <= data.keys(): + raise ValueError("supersede requires successor content and expected_revision") + judgment = graph.supersede(identifier, data["successor_id"], data["statement"], data["rationale"], data["provenance"], expected_revision=data["expected_revision"], relations=data.get("relations", [])) + elif action in {"archive", "restore"}: + if "expected_revision" not in data: + raise ValueError("lifecycle change requires expected_revision") + judgment = getattr(graph, action)(identifier, expected_revision=data["expected_revision"]) + elif action == "hard_delete": + if "expected_revision" not in data: + raise ValueError("hard_delete requires expected_revision") + graph.hard_delete(identifier, expected_revision=data["expected_revision"]) + judgment = None + else: + raise ValueError("unsupported judgment action") + return {"contract_version": CONTRACT_VERSION, "action": action, "judgment": judgment} + @staticmethod def _success(output: dict[str, Any]) -> types.CallToolResult: text = json.dumps(output, ensure_ascii=False, sort_keys=True, separators=(",", ":")) diff --git a/tests/test_judgments.py b/tests/test_judgments.py new file mode 100644 index 0000000..5fcdae0 --- /dev/null +++ b/tests/test_judgments.py @@ -0,0 +1,100 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +from pathlib import Path + +from neuron_graph_rag import JudgmentContractError, NeuronGraphRAG +from tools.judgment_graph import backup, integrity, restore + + +PROVENANCE = {"source": "issue:116", "actor": "test"} + + +class JudgmentGraphTests(unittest.TestCase): + def test_lifecycle_supersede_and_audit_history(self) -> None: + with NeuronGraphRAG() as engine: + engine.judgments.add("base", "Use the old path", "Initial decision", PROVENANCE) + updated = engine.judgments.update( + "base", "Use the old path", "Clarified decision", PROVENANCE, + expected_revision=1, + ) + self.assertEqual(updated["revision"], 2) + successor = engine.judgments.supersede( + "base", "successor", "Use the new path", "Evidence changed", PROVENANCE, + expected_revision=2, + ) + self.assertEqual( + successor["relations"], + [{"target_id": "base", "relation_type": "supersedes"}], + ) + self.assertEqual(engine.judgments.get("base")["lifecycle"], "archived") + self.assertEqual(engine.judgments.get("base")["superseded_by"], "successor") + self.assertEqual([node.node_id for node in engine.store.list_nodes()], ["successor"]) + with self.assertRaises(JudgmentContractError): + engine.judgments.supersede( + "base", "other", "Other", "No", PROVENANCE, expected_revision=2 + ) + with self.assertRaises(JudgmentContractError): + engine.judgments.restore("base", expected_revision=2) + + def test_dangling_relation_and_stale_update_roll_back_atomically(self) -> None: + with NeuronGraphRAG() as engine: + engine.judgments.add("one", "One", "Reason", PROVENANCE) + with self.assertRaises(JudgmentContractError): + engine.judgments.update( + "one", "Changed", "Reason", PROVENANCE, expected_revision=1, + relations=[{"target_id": "missing", "relation_type": "supports"}], + ) + self.assertEqual(engine.judgments.get("one")["revision"], 1) + self.assertEqual(engine.judgments.get("one")["statement"], "One") + with self.assertRaises(JudgmentContractError): + engine.judgments.update( + "one", "Changed", "Reason", PROVENANCE, expected_revision=2 + ) + + def test_archive_is_logical_and_hard_delete_is_explicitly_restricted(self) -> None: + with NeuronGraphRAG() as engine: + engine.judgments.add("candidate", "Candidate", "Reason", PROVENANCE) + engine.judgments.archive("candidate", expected_revision=1) + self.assertEqual(engine.store.list_nodes(), []) + self.assertEqual(engine.judgments.get("candidate")["statement"], "Candidate") + engine.judgments.restore("candidate", expected_revision=1) + self.assertEqual([node.node_id for node in engine.store.list_nodes()], ["candidate"]) + engine.judgments.archive("candidate", expected_revision=1) + engine.judgments.hard_delete("candidate", expected_revision=1) + with self.assertRaises(KeyError): + engine.judgments.get("candidate") + + def test_deterministic_round_trip_and_backup_restore(self) -> None: + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + database = root / "source.sqlite" + exported = root / "graph.json" + imported = root / "imported.sqlite" + backup_path = root / "backup.sqlite" + restored = root / "restored.sqlite" + with NeuronGraphRAG(database) as engine: + engine.judgments.add("a", "Alpha", "Reason A", PROVENANCE) + engine.judgments.add( + "b", "Beta", "Reason B", PROVENANCE, + relations=[{"target_id": "a", "relation_type": "supports"}], + ) + expected = engine.judgments.export() + exported.write_text( + json.dumps(expected, ensure_ascii=False, sort_keys=True, indent=2) + "\n", + encoding="utf-8", newline="\n", + ) + with NeuronGraphRAG(imported) as engine: + engine.judgments.import_graph(json.loads(exported.read_text(encoding="utf-8"))) + self.assertEqual(engine.judgments.export(), expected) + backup(database, backup_path) + restore(backup_path, restored) + self.assertEqual(integrity(restored)["sqlite_integrity"], "ok") + with NeuronGraphRAG(restored) as engine: + self.assertEqual(engine.judgments.export(), expected) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_mcp_adapter.py b/tests/test_mcp_adapter.py index 4072ebf..bca568f 100644 --- a/tests/test_mcp_adapter.py +++ b/tests/test_mcp_adapter.py @@ -50,21 +50,24 @@ def tearDown(self) -> None: async def test_tools_list_contract_and_structured_search_result(self) -> None: listed = await self.adapter.list_tools() - self.assertEqual([tool.name for tool in listed.tools], ["search", "record_source_use", "record_outcome"]) + self.assertEqual([tool.name for tool in listed.tools], ["search", "record_source_use", "record_outcome", "write_judgment"]) self.assertEqual( - [tool.description for tool in listed.tools], + [tool.description for tool in listed.tools[:3]], [SEARCH_DESCRIPTION, SOURCE_USE_DESCRIPTION, OUTCOME_DESCRIPTION], ) self.assertEqual( [tool.annotations.idempotent_hint for tool in listed.tools], - [False, True, True], + [False, True, True, False], ) - for tool in listed.tools: + for tool in listed.tools[:3]: self.assertFalse(tool.annotations.read_only_hint) self.assertFalse(tool.annotations.destructive_hint) self.assertFalse(tool.annotations.open_world_hint) self.assertFalse(tool.input_schema["additionalProperties"]) self.assertFalse(tool.output_schema["additionalProperties"]) + self.assertTrue(listed.tools[3].annotations.destructive_hint) + self.assertFalse(listed.tools[3].input_schema["additionalProperties"]) + self.assertFalse(listed.tools[3].output_schema["additionalProperties"]) result = await self.adapter.call_tool( None, @@ -147,6 +150,44 @@ async def test_feedback_loop_and_safe_tool_errors(self) -> None: self.assertNotIn("secret query text", invalid.content[0].text) self.assertNotIn("private value", invalid.content[0].text) + async def test_judgment_write_uses_atomic_domain_surface(self) -> None: + added = await self.adapter.call_tool( + None, + types.CallToolRequestParams( + name="write_judgment", + arguments={ + "contract_version": CONTRACT_VERSION, + "action": "add", + "judgment_id": "mcp-judgment", + "statement": "Use the domain API", + "rationale": "Raw SQL is outside the model-facing contract", + "provenance": {"source": "test"}, + }, + ), + ) + self.assertFalse(added.is_error) + self.assertEqual(added.structured_content["judgment"]["revision"], 1) + stale = await self.adapter.call_tool( + None, + types.CallToolRequestParams( + name="write_judgment", + arguments={ + "contract_version": CONTRACT_VERSION, + "action": "update", + "judgment_id": "mcp-judgment", + "statement": "Changed", + "rationale": "Stale writes fail closed", + "provenance": {"source": "test"}, + "expected_revision": 2, + }, + ), + ) + self.assertTrue(stale.is_error) + self.assertEqual( + self.adapter.engine.judgments.get("mcp-judgment")["statement"], + "Use the domain API", + ) + async def test_feedback_receipt_exposes_quorum_evidence_and_replays(self) -> None: self.adapter.close() self.adapter = FeedbackMCPAdapter( @@ -481,7 +522,7 @@ async def test_stdio_protocol_smoke(self) -> None: ): await session.initialize() listed = await session.list_tools() - self.assertEqual([tool.name for tool in listed.tools], ["search", "record_source_use", "record_outcome"]) + self.assertEqual([tool.name for tool in listed.tools], ["search", "record_source_use", "record_outcome", "write_judgment"]) result = await session.call_tool( "search", {"contract_version": CONTRACT_VERSION, "query": "cache"}, diff --git a/tools/judgment_graph.py b/tools/judgment_graph.py new file mode 100644 index 0000000..5de9d1a --- /dev/null +++ b/tools/judgment_graph.py @@ -0,0 +1,100 @@ +from __future__ import annotations + +import argparse +import json +import sqlite3 +from contextlib import closing +from pathlib import Path + +from neuron_graph_rag import NeuronGraphRAG + + +def _write_json(path: Path, payload: object) -> None: + path.write_text( + json.dumps(payload, ensure_ascii=False, sort_keys=True, indent=2) + "\n", + encoding="utf-8", + newline="\n", + ) + + +def export_graph(database: Path, output: Path) -> None: + with NeuronGraphRAG(database) as engine: + _write_json(output, engine.judgments.export()) + + +def import_graph(database: Path, source: Path) -> None: + payload = json.loads(source.read_text(encoding="utf-8")) + with NeuronGraphRAG(database) as engine: + engine.judgments.import_graph(payload) + + +def backup(source: Path, destination: Path) -> None: + if destination.exists(): + raise FileExistsError(f"refusing to overwrite backup: {destination}") + with closing(sqlite3.connect(source)) as source_connection, closing(sqlite3.connect(destination)) as destination_connection: + source_connection.backup(destination_connection) + + +def restore(source: Path, destination: Path) -> None: + if destination.exists(): + raise FileExistsError(f"refusing to overwrite database: {destination}") + with closing(sqlite3.connect(source)) as source_connection, closing(sqlite3.connect(destination)) as destination_connection: + if source_connection.execute("PRAGMA integrity_check").fetchone()[0] != "ok": + raise RuntimeError("backup integrity check failed") + source_connection.backup(destination_connection) + + +def integrity(database: Path) -> dict[str, object]: + with closing(sqlite3.connect(database)) as connection: + connection.row_factory = sqlite3.Row + sqlite_status = connection.execute("PRAGMA integrity_check").fetchone()[0] + foreign_keys = [dict(row) for row in connection.execute("PRAGMA foreign_key_check")] + dangling = connection.execute( + """ + SELECT source_id, target_id, relation_type FROM judgment_relations + WHERE source_id NOT IN (SELECT judgment_id FROM judgments) + OR target_id NOT IN (SELECT judgment_id FROM judgments) + ORDER BY source_id, target_id, relation_type + """ + ).fetchall() + duplicate_successors = connection.execute( + """ + SELECT superseded_by, COUNT(*) AS predecessor_count FROM judgments + WHERE superseded_by IS NOT NULL GROUP BY superseded_by HAVING COUNT(*) > 1 + """ + ).fetchall() + result = { + "sqlite_integrity": sqlite_status, + "foreign_key_violations": foreign_keys, + "dangling_relations": [dict(row) for row in dangling], + "duplicate_successors": [dict(row) for row in duplicate_successors], + } + if sqlite_status != "ok" or foreign_keys or dangling or duplicate_successors: + raise RuntimeError(json.dumps(result, sort_keys=True)) + return result + + +def main() -> None: + parser = argparse.ArgumentParser(description="Operate the canonical SQLite judgment graph") + subparsers = parser.add_subparsers(dest="command", required=True) + for name in ("export", "import", "backup", "restore"): + command = subparsers.add_parser(name) + command.add_argument("source", type=Path) + command.add_argument("destination", type=Path) + check = subparsers.add_parser("integrity") + check.add_argument("database", type=Path) + arguments = parser.parse_args() + if arguments.command == "export": + export_graph(arguments.source, arguments.destination) + elif arguments.command == "import": + import_graph(arguments.destination, arguments.source) + elif arguments.command == "backup": + backup(arguments.source, arguments.destination) + elif arguments.command == "restore": + restore(arguments.source, arguments.destination) + else: + print(json.dumps(integrity(arguments.database), sort_keys=True)) + + +if __name__ == "__main__": + main() From d9a57c2ec5511ccbd1bafeef04fb85b9b7570aeb Mon Sep 17 00:00:00 2001 From: lipluscodex <268560960+lipluscodex@users.noreply.github.com> Date: Sun, 23 Aug 2026 00:06:27 +0900 Subject: [PATCH 2/2] fix(decisions): enforce judgment failure boundaries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Finding 1 accepted: hard delete の active、inbound relation、supersede history 拒否は Issue #116 の物理削除境界を直接証明するため、各失敗前後の graph 一致をテストする。 Finding 2 accepted: lifecycle no-op の成功は同一 revision で updated_at を書き換えるため、再 archive / 再 restore を fail closed にし timestamp 不変を検証する。 Finding 3 accepted: non-dangling だけでは supersession の意味整合を保証できないため、superseded_by、archived lifecycle、明示 supersedes relation の双方向整合を integrity check と corruption tests に追加する。 --- docs/canonical-sqlite-judgment-graph.md | 4 +- src/neuron_graph_rag/judgments.py | 2 + tests/test_judgments.py | 85 +++++++++++++++++++++++++ tools/judgment_graph.py | 39 +++++++++++- 4 files changed, 128 insertions(+), 2 deletions(-) diff --git a/docs/canonical-sqlite-judgment-graph.md b/docs/canonical-sqlite-judgment-graph.md index 18c75b1..d542009 100644 --- a/docs/canonical-sqlite-judgment-graph.md +++ b/docs/canonical-sqlite-judgment-graph.md @@ -8,6 +8,8 @@ NGR の新規 judgment は SQLite の `judgments`、`judgment_revisions`、`judg `NeuronGraphRAG.judgments` は add、update、supersede、archive、restore、hard delete を transaction 単位で提供する。update、supersede、archive、restore、hard delete は `expected_revision` による楽観的 concurrency check を要求する。relation target 不在、部分更新、stale revision、同一 predecessor の再 supersede は transaction 全体を失敗させる。 +archive 済み judgment の再 archive と active judgment の再 restore は no-op として成功させず、fail closed にする。同一 revision の反復操作によって `updated_at` を暗黙更新しない。 + supersede は新しい stable identity を作り、旧判断を archive し、successor から predecessor への `supersedes` relation と predecessor の `superseded_by` を保持する。superseded judgment は restore できない。 archive は通常 retrieval から外す論理的忘却であり、revision、provenance、relation は監査 API から取得できる。hard delete は明示操作で、archived、revision 一致、inbound relation と successor history がない候補にだけ許可する。 @@ -22,6 +24,6 @@ MCP の `write_judgment` は同じ domain API へ写像し、model に raw SQL - `import INPUT_JSON DESTINATION_DB`: 全 identity と relation target を事前検証し、一つの transaction で再構築する。 - `backup SOURCE_DB BACKUP_DB`: SQLite backup API で transaction-consistent copy を作る。既存出力は上書きしない。 - `restore BACKUP_DB DESTINATION_DB`: integrity check 後に新規 destination へ復元する。既存 database は上書きしない。 -- `integrity DATABASE`: SQLite integrity、foreign key、dangling relation、二重 successor を fail closed に検査する。 +- `integrity DATABASE`: SQLite integrity、foreign key、dangling relation、二重 successor、および `superseded_by` と archived lifecycle / successor の明示的 `supersedes` relation の双方向整合を fail closed に検査する。 backup は export より多くの revision history と retrieval / feedback state を保持するため、完全復旧の正本である。export は current judgment graph の決定論的 portability surface である。 diff --git a/src/neuron_graph_rag/judgments.py b/src/neuron_graph_rag/judgments.py index 0ceaaee..10f4015 100644 --- a/src/neuron_graph_rag/judgments.py +++ b/src/neuron_graph_rag/judgments.py @@ -186,6 +186,8 @@ def _lifecycle(self, judgment_id: str, expected_revision: int, lifecycle: str) - row = self._require(connection, judgment_id) if row["current_revision"] != expected_revision: raise JudgmentContractError("stale expected_revision") + if row["lifecycle"] == lifecycle: + raise JudgmentContractError(f"judgment is already {lifecycle}") if lifecycle == "active" and row["superseded_by"] is not None: raise JudgmentContractError("superseded judgments cannot be restored") connection.execute( diff --git a/tests/test_judgments.py b/tests/test_judgments.py index 5fcdae0..3fbf930 100644 --- a/tests/test_judgments.py +++ b/tests/test_judgments.py @@ -67,6 +67,91 @@ def test_archive_is_logical_and_hard_delete_is_explicitly_restricted(self) -> No with self.assertRaises(KeyError): engine.judgments.get("candidate") + def test_hard_delete_failure_boundaries_preserve_graph_atomically(self) -> None: + with NeuronGraphRAG() as engine: + engine.judgments.add("active", "Active", "Reason", PROVENANCE) + before = engine.judgments.export() + with self.assertRaises(JudgmentContractError): + engine.judgments.hard_delete("active", expected_revision=1) + self.assertEqual(engine.judgments.export(), before) + + engine.judgments.add("referenced", "Referenced", "Reason", PROVENANCE) + engine.judgments.add( + "source", "Source", "Reason", PROVENANCE, + relations=[{"target_id": "referenced", "relation_type": "supports"}], + ) + engine.judgments.archive("referenced", expected_revision=1) + before = engine.judgments.export() + with self.assertRaises(JudgmentContractError): + engine.judgments.hard_delete("referenced", expected_revision=1) + self.assertEqual(engine.judgments.export(), before) + + engine.judgments.add("predecessor", "Old", "Reason", PROVENANCE) + engine.judgments.supersede( + "predecessor", "successor", "New", "Reason", PROVENANCE, + expected_revision=1, + ) + before = engine.judgments.export() + with self.assertRaises(JudgmentContractError): + engine.judgments.hard_delete("predecessor", expected_revision=1) + self.assertEqual(engine.judgments.export(), before) + + def test_lifecycle_no_op_fails_without_touching_timestamp(self) -> None: + with NeuronGraphRAG() as engine: + engine.judgments.add("state", "State", "Reason", PROVENANCE) + active_timestamp = engine.store.connection.execute( + "SELECT updated_at FROM judgments WHERE judgment_id = 'state'" + ).fetchone()[0] + with self.assertRaisesRegex(JudgmentContractError, "already active"): + engine.judgments.restore("state", expected_revision=1) + self.assertEqual( + engine.store.connection.execute( + "SELECT updated_at FROM judgments WHERE judgment_id = 'state'" + ).fetchone()[0], + active_timestamp, + ) + engine.judgments.archive("state", expected_revision=1) + archived_timestamp = engine.store.connection.execute( + "SELECT updated_at FROM judgments WHERE judgment_id = 'state'" + ).fetchone()[0] + with self.assertRaisesRegex(JudgmentContractError, "already archived"): + engine.judgments.archive("state", expected_revision=1) + self.assertEqual( + engine.store.connection.execute( + "SELECT updated_at FROM judgments WHERE judgment_id = 'state'" + ).fetchone()[0], + archived_timestamp, + ) + + def test_integrity_rejects_inconsistent_supersession_state(self) -> None: + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + for mutation in ("active_predecessor", "missing_relation", "wrong_successor"): + with self.subTest(mutation=mutation): + database = root / f"{mutation}.sqlite" + with NeuronGraphRAG(database) as engine: + engine.judgments.add("old", "Old", "Reason", PROVENANCE) + engine.judgments.supersede( + "old", "new", "New", "Reason", PROVENANCE, + expected_revision=1, + ) + if mutation == "active_predecessor": + engine.store.connection.execute( + "UPDATE judgments SET lifecycle = 'active' WHERE judgment_id = 'old'" + ) + elif mutation == "missing_relation": + engine.store.connection.execute( + "DELETE FROM judgment_relations WHERE source_id = 'new' AND target_id = 'old' AND relation_type = 'supersedes'" + ) + else: + engine.judgments.add("other", "Other", "Reason", PROVENANCE) + engine.store.connection.execute( + "UPDATE judgments SET superseded_by = 'other' WHERE judgment_id = 'old'" + ) + engine.store.connection.commit() + with self.assertRaisesRegex(RuntimeError, "supersession_inconsistencies"): + integrity(database) + def test_deterministic_round_trip_and_backup_restore(self) -> None: with tempfile.TemporaryDirectory() as directory: root = Path(directory) diff --git a/tools/judgment_graph.py b/tools/judgment_graph.py index 5de9d1a..cd5ce7b 100644 --- a/tools/judgment_graph.py +++ b/tools/judgment_graph.py @@ -63,13 +63,50 @@ def integrity(database: Path) -> dict[str, object]: WHERE superseded_by IS NOT NULL GROUP BY superseded_by HAVING COUNT(*) > 1 """ ).fetchall() + supersession_inconsistencies = connection.execute( + """ + SELECT predecessor.judgment_id, predecessor.lifecycle, + predecessor.superseded_by + FROM judgments AS predecessor + WHERE predecessor.superseded_by IS NOT NULL + AND ( + predecessor.lifecycle <> 'archived' + OR NOT EXISTS ( + SELECT 1 FROM judgment_relations AS relation + WHERE relation.source_id = predecessor.superseded_by + AND relation.target_id = predecessor.judgment_id + AND relation.relation_type = 'supersedes' + ) + ) + UNION ALL + SELECT relation.target_id, target.lifecycle, relation.source_id + FROM judgment_relations AS relation + JOIN judgments AS target ON target.judgment_id = relation.target_id + WHERE relation.relation_type = 'supersedes' + AND ( + target.lifecycle <> 'archived' + OR target.superseded_by IS NULL + OR target.superseded_by <> relation.source_id + ) + ORDER BY 1, 3 + """ + ).fetchall() result = { "sqlite_integrity": sqlite_status, "foreign_key_violations": foreign_keys, "dangling_relations": [dict(row) for row in dangling], "duplicate_successors": [dict(row) for row in duplicate_successors], + "supersession_inconsistencies": [ + dict(row) for row in supersession_inconsistencies + ], } - if sqlite_status != "ok" or foreign_keys or dangling or duplicate_successors: + if ( + sqlite_status != "ok" + or foreign_keys + or dangling + or duplicate_successors + or supersession_inconsistencies + ): raise RuntimeError(json.dumps(result, sort_keys=True)) return result