From a2a76cd5ee9235b13d64152b38fd05911b0faf8c Mon Sep 17 00:00:00 2001 From: Philipp Date: Tue, 29 Sep 2026 09:36:32 +0200 Subject: [PATCH 1/3] feat(curator): evidence pack and gateway client for a local curator New optional package `wikibricks_curator` (stdlib only, no psycopg or Databricks SDK; the core library does not import it). - `build_request` turns one curation-backlog item into a bounded request: the project's living page (first covering `topics/` page, else `topics/`), covering and related pages in the shape `build_patches` expects, and only user/assistant messages from the project's sessions after the page update or cursor, newest kept first, capped at 40 events and 30,000 characters. - `chat_json` calls an OpenAI-compatible chat-completions endpoint with the remote curator's message layout and JSON parsing; `resolve_token` reads WIKIBRICKS_GATEWAY_TOKEN or `databricks auth token --profile`. Implemented by GLM 5.3 Flash; timestamp comparison across UTC offsets and per-project event loading fixed in review. Co-authored-by: Isaac --- pyproject.toml | 2 +- src/wikibricks_curator/__init__.py | 3 + src/wikibricks_curator/evidence.py | 204 ++++++++++++++ src/wikibricks_curator/gateway.py | 109 ++++++++ tests/test_local_curator.py | 435 +++++++++++++++++++++++++++++ 5 files changed, 752 insertions(+), 1 deletion(-) create mode 100644 src/wikibricks_curator/__init__.py create mode 100644 src/wikibricks_curator/evidence.py create mode 100644 src/wikibricks_curator/gateway.py create mode 100644 tests/test_local_curator.py diff --git a/pyproject.toml b/pyproject.toml index 4535451..ac31809 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -33,7 +33,7 @@ requires = ["hatchling"] build-backend = "hatchling.build" [tool.hatch.build.targets.wheel] -packages = ["src/wikibricks", "src/wikibricks_remote"] +packages = ["src/wikibricks", "src/wikibricks_remote", "src/wikibricks_curator"] [tool.ruff] target-version = "py314" diff --git a/src/wikibricks_curator/__init__.py b/src/wikibricks_curator/__init__.py new file mode 100644 index 0000000..e63643f --- /dev/null +++ b/src/wikibricks_curator/__init__.py @@ -0,0 +1,3 @@ +"""Local curator request and gateway clients.""" + +__all__: list[str] = [] diff --git a/src/wikibricks_curator/evidence.py b/src/wikibricks_curator/evidence.py new file mode 100644 index 0000000..c2e3510 --- /dev/null +++ b/src/wikibricks_curator/evidence.py @@ -0,0 +1,204 @@ +"""Build bounded model requests from local curation evidence.""" + +from __future__ import annotations + +import json +import sqlite3 +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +from wikibricks.storage.sqlite_store import SQLiteStore + + +def _normalize(value: str) -> str: + return Path(value).name.strip().lower().replace(" ", "-").replace("_", "-") + + +def _database_path(conn: sqlite3.Connection) -> Path | None: + row = conn.execute("PRAGMA database_list").fetchone() + path = str(row[2]) if row else "" + return Path(path) if path else None + + +def _active_pages( + conn: sqlite3.Connection, + *, + paths: set[str], + project: str, + related: int, +) -> list[dict[str, Any]]: + candidates: list[str] = [] + if paths: + placeholders = ",".join("?" for _ in paths) + rows = conn.execute( + "SELECT p.path FROM pages p WHERE p.status = 'active' " + f"AND p.path IN ({placeholders})", + tuple(paths), + ).fetchall() + candidates.extend(row[0] for row in rows) + + if related > 0 and len(candidates) < related + len(paths): + database_path = _database_path(conn) + if database_path is not None: + hits = SQLiteStore(database_path).search( + project.replace("-", " "), + num_results=related + len(candidates) + len(paths), + ) + related_limit = len(candidates) + related + for hit in hits: + path = hit["path"] + if ( + len(candidates) >= related_limit + or path.startswith("_meta/") + or path in candidates + ): + continue + candidates.append(path) + + if not candidates: + return [] + placeholders = ",".join("?" for _ in candidates) + rows = conn.execute( + "SELECT v.version_id, p.path, v.title, v.page_type, v.content, v.tags, " + "v.source_ids, v.content_hash FROM pages p " + "JOIN page_versions v ON v.version_id = p.current_version_id " + "WHERE p.status = 'active' " + f"AND p.path IN ({placeholders}) ORDER BY p.path", + candidates, + ).fetchall() + return [ + { + "evidence_id": f"page-version:{row['version_id']}", + "path": row["path"], + "title": row["title"], + "page_type": row["page_type"], + "content": json.loads(row["content"]), + "tags": json.loads(row["tags"]), + "source_ids": json.loads(row["source_ids"]) if row["source_ids"] else [], + "base_version_id": row["version_id"], + "base_content_hash": row["content_hash"], + } + for row in rows + ] + + +def _moment(value: Any) -> datetime: + parsed = value if isinstance(value, datetime) else datetime.fromisoformat(str(value)) + return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc) + + +def _evidence( + conn: sqlite3.Connection, + *, + project: str, + cursor: str | None, + last_page_update: str | None, + max_events: int, + max_chars: int, + max_event_chars: int, +) -> tuple[list[dict[str, Any]], str]: + # Compare parsed datetimes: stored ISO strings carry different UTC offsets. + after = [_moment(value) for value in (cursor, last_page_update) if value] + sessions = [ + (_moment(row[2]), row[0]) + for row in conn.execute( + "SELECT session_id, workspace, COALESCE(source_updated_at, updated_at) " + "FROM sessions WHERE workspace IS NOT NULL" + ).fetchall() + if _normalize(row[1]) == project and all(_moment(row[2]) > bound for bound in after) + ] + if not sessions: + return [], "" + sessions.sort() + order = {session_id: index for index, (_, session_id) in enumerate(sessions)} + moments = {session_id: moment for moment, session_id in sessions} + placeholders = ",".join("?" for _ in sessions) + rows = conn.execute( + "SELECT s.session_id, s.page_path, e.position, " + "v.version_id, v.kind, v.content, v.source_created_at " + "FROM sessions s JOIN session_events e ON e.session_id = s.session_id " + "JOIN session_event_versions v ON v.version_id = e.current_version_id " + "WHERE e.active = 1 AND v.kind IN ('user', 'assistant') " + f"AND s.session_id IN ({placeholders})", + list(order), + ).fetchall() + candidates = sorted(rows, key=lambda row: (order[row["session_id"]], row["position"])) + + selected: list[sqlite3.Row] = [] + total_chars = 0 + for row in candidates: + length = min(len(row["content"]), max_event_chars) + if len(selected) == max_events: + total_chars -= min(len(selected[0]["content"]), max_event_chars) + selected.pop(0) + while selected and total_chars + length > max_chars: + total_chars -= min(len(selected[0]["content"]), max_event_chars) + selected.pop(0) + if total_chars + length > max_chars: + continue + total_chars += length + selected.append(row) + evidence = [ + { + "evidence_id": f"session-event:{row['version_id']}", + "session": row["page_path"], + "kind": row["kind"], + "created_at": row["source_created_at"], + "text": row["content"][:max_event_chars], + } + for row in selected + ] + newest = max(moments[row["session_id"]] for row in selected) + return evidence, newest.isoformat() + + +def build_request( + conn: sqlite3.Connection, + item: dict[str, Any], + *, + cursor: str | None = None, + max_events: int = 40, + max_chars: int = 30000, + max_event_chars: int = 2000, + related: int = 5, +) -> dict[str, Any] | None: + """Build one bounded curator request for a backlog item.""" + project = item["project"] + requested_paths = list(dict.fromkeys(item["pages"])) + living_page = next( + (path for path in requested_paths if path.startswith("topics/")), + f"topics/{project}", + ) + paths = set(requested_paths) | {living_page} + pages = _active_pages( + conn, + paths=paths, + project=project, + related=related, + ) + evidence, newest_evidence_at = _evidence( + conn, + project=project, + cursor=cursor, + last_page_update=item.get("last_page_update"), + max_events=max_events, + max_chars=max_chars, + max_event_chars=max_event_chars, + ) + if not evidence: + return None + evidence_ids = {entry["evidence_id"] for entry in evidence} + page_ids = {page["evidence_id"] for page in pages} + return { + "request": { + "project": project, + "living_page": living_page, + "current_pages": pages, + "evidence": evidence, + "similarity_candidates": [], + }, + "pages": pages, + "evidence_ids": evidence_ids | page_ids, + "newest_evidence_at": newest_evidence_at, + } diff --git a/src/wikibricks_curator/gateway.py b/src/wikibricks_curator/gateway.py new file mode 100644 index 0000000..195e99f --- /dev/null +++ b/src/wikibricks_curator/gateway.py @@ -0,0 +1,109 @@ +"""Databricks AI Gateway chat-completions client.""" + +from __future__ import annotations + +import json +import os +import subprocess +from typing import Any +from urllib.error import HTTPError +from urllib.request import Request, urlopen + + +def resolve_token(profile: str | None) -> str: + """Resolve a gateway token from the environment or the Databricks CLI.""" + token = os.environ.get("WIKIBRICKS_GATEWAY_TOKEN") + if token: + return token + command = ["databricks", "auth", "token"] + if profile is not None: + command.extend(["--profile", profile]) + try: + result = subprocess.run( + command, + check=True, + capture_output=True, + text=True, + ) + except subprocess.CalledProcessError as error: + raise ValueError("Databricks auth token command failed") from error + try: + payload = json.loads(result.stdout) + except json.JSONDecodeError as error: + raise ValueError("Databricks auth token response was not JSON") from error + token = payload.get("access_token") if isinstance(payload, dict) else None + if not token: + raise ValueError("Databricks auth token response did not contain access_token") + return str(token) + + +def _json_content(response: Any) -> dict[str, Any]: + content = response.read() + text = content.decode("utf-8") if isinstance(content, bytes) else content + try: + payload = json.loads(text) + except json.JSONDecodeError as error: + raise RuntimeError("curation model response could not be decoded") from error + if not isinstance(payload, dict) or not payload.get("choices"): + raise RuntimeError("curation model returned no content") + message = payload["choices"][0].get("message") + content = message.get("content") if isinstance(message, dict) else None + if not isinstance(content, str) or not content.strip(): + raise RuntimeError("curation model returned empty content") + text = content.strip() + if text.startswith("```"): + lines = text.splitlines() + text = "\n".join(lines[1:-1]) + try: + value, _ = json.JSONDecoder().raw_decode(text) + except json.JSONDecodeError as error: + raise RuntimeError("curation model output could not be decoded") from error + if not isinstance(value, dict): + raise RuntimeError("curation model output must be a JSON object") + return value + + +def chat_json( + system_prompt: str, + request: dict[str, Any], + schema: dict[str, Any], + *, + base_url: str, + token: str, + model: str, + temperature: float = 0.0, + max_tokens: int = 8192, + timeout: float = 180, + opener=urlopen, +) -> dict[str, Any]: + """Call an OpenAI-compatible chat endpoint and return its JSON object.""" + messages = [ + { + "role": "system", + "content": ( + f"{system_prompt}\n\nReturn JSON matching this schema:\n" + f"{json.dumps(schema, separators=(',', ':'))}" + ), + }, + {"role": "user", "content": json.dumps(request, ensure_ascii=False)}, + ] + payload = { + "model": model, + "messages": messages, + "temperature": temperature, + "max_tokens": max_tokens, + } + http_request = Request( + f"{base_url.rstrip('/')}/chat/completions", + data=json.dumps(payload, ensure_ascii=False).encode("utf-8"), + headers={ + "Authorization": f"Bearer {token}", + "Content-Type": "application/json", + }, + method="POST", + ) + try: + response = opener(http_request, timeout) + except HTTPError as error: + raise RuntimeError(f"curation model request failed: HTTP {error.code}") from error + return _json_content(response) diff --git a/tests/test_local_curator.py b/tests/test_local_curator.py new file mode 100644 index 0000000..ce2984e --- /dev/null +++ b/tests/test_local_curator.py @@ -0,0 +1,435 @@ +from __future__ import annotations + +import json +import subprocess +from pathlib import Path +from urllib.error import HTTPError +from urllib.request import Request +from uuid import uuid4 + +import pytest + +from wikibricks.models import SessionEvent, SessionRecord +from wikibricks.storage.sqlite_store import SQLiteStore +from wikibricks_remote.resources import load_policy + +TEN = "2026-01-03T10:00:00+00:00" +OLD_SESSION = "2026-01-03T09:00:00+00:00" +MIDDLE_SESSION = "2026-01-04T09:00:00+00:00" +NEW_SESSION = "2026-01-04T11:00:00+00:00" + + +def _store(tmp_path: Path) -> SQLiteStore: + store = SQLiteStore(tmp_path / "wikibricks.db") + store.migrate() + return store + + +def _ingest( + store: SQLiteStore, + external_id: str, + workspace: str, + updated_at: str, + events: list[SessionEvent], +) -> None: + store.ingest_session( + SessionRecord( + harness="test-harness", + external_id=external_id, + user_id="user", + workspace=workspace, + updated_at=updated_at, + events=events, + ) + ) + + +def _populate(tmp_path: Path) -> tuple[SQLiteStore, dict]: + store = _store(tmp_path) + store.write_page("topics/alpha", "Alpha", {"summary": "old", "body": "alpha page"}) + store.write_page( + "topics/alpha-notes", "Alpha Notes", {"summary": "related", "body": "alpha notes"} + ) + store.write_page("topics/beta", "Beta", {"summary": "other", "body": "beta page"}) + _ingest( + store, + "old-alpha", + "/Users/u/work/Alpha One", + OLD_SESSION, + [ + SessionEvent("user-1", "user", "old user"), + SessionEvent("tool-1", "tool_call", "old tool"), + ], + ) + _ingest( + store, + "middle-alpha", + "/Users/u/work/alpha_one", + MIDDLE_SESSION, + [ + SessionEvent("user-2", "user", "middle user"), + SessionEvent("assistant-2", "assistant", "middle assistant"), + SessionEvent("tool-2", "tool_result", "middle tool"), + ], + ) + _ingest( + store, + "new-alpha", + "/Users/u/work/Alpha One", + NEW_SESSION, + [ + SessionEvent("user-3", "user", "new user"), + SessionEvent("assistant-3", "assistant", "new assistant"), + SessionEvent("tool-3", "tool_call", "new tool"), + ], + ) + _ingest( + store, + "beta", + "/Users/u/work/beta", + NEW_SESSION, + [SessionEvent("user-beta", "user", "beta user")], + ) + item = { + "project": "alpha-one", + "workspace": "/Users/u/work/Alpha One", + "new_sessions": 2, + "last_session_at": NEW_SESSION, + "pages": ["topics/alpha"], + "last_page_update": TEN, + } + return store, item + + +def test_build_request_selects_active_pages_and_user_assistant_evidence(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store, item = _populate(tmp_path) + with store.connection() as conn: + result = build_request(conn, item, related=1) + + assert result["request"]["living_page"] == "topics/alpha" + assert [page["path"] for page in result["pages"]] == ["topics/alpha", "topics/alpha-notes"] + with store.connection() as conn: + version = conn.execute( + "SELECT version_id, content_hash, content, tags, source_ids, page_type, title " + "FROM page_versions v JOIN pages p ON p.page_id = v.page_id " + "WHERE p.path = ?", + ("topics/alpha",), + ).fetchone() + page = result["pages"][0] + assert page["evidence_id"] == f"page-version:{version['version_id']}" + assert page["base_version_id"] == version["version_id"] + assert page["base_content_hash"] == version["content_hash"] + assert page["content"] == json.loads(version["content"]) + + evidence = result["request"]["evidence"] + assert [entry["kind"] for entry in evidence] == ["user", "assistant", "user", "assistant"] + assert [entry["text"] for entry in evidence] == [ + "middle user", + "middle assistant", + "new user", + "new assistant", + ] + session_paths = {entry["session"] for entry in evidence} + assert all( + path.endswith(("middle-alpha", "new-alpha")) for path in session_paths + ) + assert result["evidence_ids"] == {entry["evidence_id"] for entry in evidence} | { + page["evidence_id"] for page in result["pages"] + } + assert result["newest_evidence_at"] == NEW_SESSION + assert result["request"]["current_pages"] == result["pages"] + assert result["request"]["evidence"] == evidence + assert result["request"]["similarity_candidates"] == [] + + +def test_build_request_trims_to_the_newest_evidence(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store, item = _populate(tmp_path) + with store.connection() as conn: + limited = build_request(conn, item, max_events=3) + assert [entry["text"] for entry in limited["request"]["evidence"]] == [ + "middle assistant", + "new user", + "new assistant", + ] + + with store.connection() as conn: + bounded = build_request(conn, item, max_chars=24, max_event_chars=20) + assert [entry["text"] for entry in bounded["request"]["evidence"]] == [ + "new user", + "new assistant", + ] + assert bounded["newest_evidence_at"] == NEW_SESSION + + +def test_build_request_cursor_excludes_older_sessions(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store, item = _populate(tmp_path) + with store.connection() as conn: + result = build_request(conn, item, cursor="2026-01-04T10:00:00+00:00") + + assert [entry["text"] for entry in result["request"]["evidence"]] == [ + "new user", + "new assistant", + ] + assert result["newest_evidence_at"] == NEW_SESSION + + +def test_build_request_truncates_each_event_and_returns_none_without_evidence(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store, item = _populate(tmp_path) + with store.connection() as conn: + result = build_request(conn, item, max_event_chars=6) + assert [entry["text"] for entry in result["request"]["evidence"]] == [ + "middle", + "middle", + "new us", + "new as", + ] + + old_item = {**item, "last_page_update": NEW_SESSION} + with store.connection() as conn: + assert build_request(conn, old_item) is None + + +def test_build_request_defaults_living_page_for_a_project_without_coverage(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store, _ = _populate(tmp_path) + with store.connection() as conn: + result = build_request( + conn, + { + "project": "beta", + "workspace": "/Users/u/work/beta", + "new_sessions": 1, + "last_session_at": NEW_SESSION, + "pages": [], + "last_page_update": None, + }, + related=0, + ) + + assert result["request"]["living_page"] == "topics/beta" + assert [page["path"] for page in result["pages"]] == ["topics/beta"] + assert [entry["text"] for entry in result["request"]["evidence"]] == ["beta user"] + + +def test_evidence_pages_support_a_cited_update_page_patch(tmp_path: Path): + from wikibricks_curator.evidence import build_request + from wikibricks_remote.proposals import build_patches + + store, item = _populate(tmp_path) + with store.connection() as conn: + result = build_request(conn, item) + evidence_id = result["request"]["evidence"][0]["evidence_id"] + raw_result = { + "proposals": [ + { + "group": "update-alpha", + "operation": "update_page", + "path": "topics/alpha", + "title": "Updated Alpha", + "page_type": "concept", + "summary": "Updated summary", + "body": "Derived from session evidence.", + "tags": ["alpha"], + "source_ids": [], + "target_path": None, + "evidence_ids": [evidence_id], + "reason": "Session evidence changed the project state.", + "risk_class": "low", + } + ] + } + + patches = build_patches( + raw_result, + run_id=uuid4(), + pages=result["pages"], + evidence_ids=result["evidence_ids"], + policy=load_policy(), + ) + + assert len(patches) == 1 + patch = patches[0] + assert patch["operation"] == "update_page" + assert patch["path"] == "topics/alpha" + assert patch["proposal"]["content"] == { + "summary": "Updated summary", + "body": "Derived from session evidence.", + } + assert patch["base_version_id"] == result["pages"][0]["base_version_id"] + assert patch["base_content_hash"] == result["pages"][0]["base_content_hash"] + assert patch["evidence_ids"] == [evidence_id] + + +class _Response: + def __init__(self, content: str) -> None: + self.content = content + + def read(self) -> bytes: + return self.content.encode("utf-8") + + +def _opener(content: str): + calls: list[Request] = [] + + def fake_opener(request: Request, timeout: float) -> _Response: + calls.append(request) + if content is None: + response = "" + else: + response = json.dumps({"choices": [{"message": {"content": content}}]}) + return _Response(response) + + return calls, fake_opener + + +def _payload(request: Request) -> dict: + return json.loads(request.data) + + +def test_chat_json_posts_model_messages_and_parses_plain_json(): + from wikibricks_curator.gateway import chat_json + + calls, opener = _opener('{"proposals": []}') + schema = {"type": "object"} + request = {"project": "alpha"} + result = chat_json( + "Curate pages", + request, + schema, + base_url="https://gateway.example/", + token="secret", + model="databricks-model", + opener=opener, + ) + + assert result == {"proposals": []} + assert len(calls) == 1 + assert calls[0].get_full_url() == "https://gateway.example/chat/completions" + assert calls[0].get_header("Authorization") == "Bearer secret" + payload = _payload(calls[0]) + assert payload["model"] == "databricks-model" + assert payload["temperature"] == 0.0 + assert payload["max_tokens"] == 8192 + assert payload["messages"][0]["role"] == "system" + assert payload["messages"][0]["content"].startswith("Curate pages") + assert 'Return JSON matching this schema:' in payload["messages"][0]["content"] + assert json.loads(payload["messages"][0]["content"].split(":", 1)[1]) == schema + assert payload["messages"][1] == {"role": "user", "content": json.dumps(request, ensure_ascii=False)} + + +def test_chat_json_parses_fenced_and_trailing_json(): + from wikibricks_curator.gateway import chat_json + + _, fenced_opener = _opener('```json\n{"value": 1}\n```') + assert chat_json( + "s", + {}, + {}, + base_url="https://x.example", + token="t", + model="m", + opener=fenced_opener, + ) == {"value": 1} + + _, trailing_opener = _opener('{"value": 2}\n\nModel commentary.') + assert chat_json( + "s", + {}, + {}, + base_url="https://x.example", + token="t", + model="m", + opener=trailing_opener, + ) == {"value": 2} + + +@pytest.mark.parametrize( + ("content", "message"), + [ + (None, "empty content"), + ("[]", "JSON object"), + ("not json", "decode"), + ], +) +def test_chat_json_rejects_invalid_model_output(content: str | None, message: str): + from wikibricks_curator.gateway import chat_json + + _, opener = _opener(content or "") + with pytest.raises(RuntimeError, match=message): + chat_json( + "s", + {}, + {}, + base_url="https://x.example", + token="t", + model="m", + opener=opener, + ) + + +def test_chat_json_wraps_http_errors(): + from wikibricks_curator.gateway import chat_json + + def opener(request: Request, timeout: float): + raise HTTPError("url", 500, "Server Error", None, None) + + with pytest.raises(RuntimeError, match="HTTP 500"): + chat_json( + "s", + {}, + {}, + base_url="https://x.example", + token="t", + model="m", + opener=opener, + ) + + +def test_resolve_token_prefers_env(monkeypatch: pytest.MonkeyPatch): + from wikibricks_curator.gateway import resolve_token + + monkeypatch.setenv("WIKIBRICKS_GATEWAY_TOKEN", "env-token") + assert resolve_token("profile") == "env-token" + + +def test_resolve_token_calls_databricks_profile(monkeypatch: pytest.MonkeyPatch): + from wikibricks_curator.gateway import resolve_token + + calls: list[list[str]] = [] + + def fake_run(command: list[str], **_kwargs): + calls.append(command) + return subprocess.CompletedProcess(command, 0, stdout='{"access_token": "cli-token"}', stderr="") + + monkeypatch.delenv("WIKIBRICKS_GATEWAY_TOKEN", raising=False) + monkeypatch.setattr("subprocess.run", fake_run) + assert resolve_token("staging") == "cli-token" + assert calls == [["databricks", "auth", "token", "--profile", "staging"]] + + +def test_build_request_compares_timestamps_across_utc_offsets(tmp_path: Path): + from wikibricks_curator.evidence import build_request + + store = _store(tmp_path) + # 12:30+02:00 is 10:30 UTC: before the 11:00 UTC page update, although it sorts after as text. + _ingest(store, "before", "/Users/u/work/offsets", "2026-01-04T12:30:00+02:00", + [SessionEvent("0", "user", "before the page update")]) + _ingest(store, "after", "/Users/u/work/offsets", "2026-01-04T11:30:00+00:00", + [SessionEvent("0", "user", "after the page update")]) + item = {"project": "offsets", "pages": [], "last_page_update": "2026-01-04T11:00:00+00:00"} + + with store.connection() as conn: + built = build_request(conn, item) + later = build_request(conn, item, cursor="2026-01-04T13:15:00+02:00") + + assert [entry["text"] for entry in built["request"]["evidence"]] == ["after the page update"] + assert [entry["text"] for entry in later["request"]["evidence"]] == ["after the page update"] From d0b545f456a81d2473add7f022c267f726275d64 Mon Sep 17 00:00:00 2001 From: Philipp Date: Tue, 29 Sep 2026 09:49:59 +0200 Subject: [PATCH 2/3] feat(curator): nightly local curation with safe auto-apply `wikibricks-curator propose --base-url URL --profile P` sends each of the top backlog projects' new session text to one model (default GLM 5.3 Flash) and turns the reply into a local curation run. - Keeps one living page per project (`topics/` or the existing covering topics page) through a local prompt addendum on top of the remote curator prompt; only create_page, update_page and add_link. - Auto-applies low-risk groups with the existing `safe` policy. Updates that would drop more than half of a page's text are forced to high risk and wait for review; `wiki_index` lists pending review runs. - `store_manifest` stores a run locally; `pull_manifests` now uses it. - A cursor per project stops the same sessions from being proposed twice. Implemented by GLM 5.3 Flash. Review fixed the gateway call (positional timeout sent as the request body; `reasoning_effort: low`, without which GLM 5.3 Flash spends all 8,192 output tokens on reasoning), a dry run that advanced the cursor, and a guard that crashed on pages without a summary. Co-authored-by: Isaac --- AGENTS.md | 2 + CHANGELOG.md | 7 + README.md | 23 + pyproject.toml | 1 + src/wikibricks/curation/__init__.py | 6 + src/wikibricks/curation/repository.py | 35 +- src/wikibricks/mcp_server.py | 18 + src/wikibricks_curator/cli.py | 54 +++ src/wikibricks_curator/curator.py | 191 ++++++++ src/wikibricks_curator/gateway.py | 6 +- .../resources/local-curator.md | 19 + tests/test_local_curator.py | 444 ++++++++++++++++++ 12 files changed, 798 insertions(+), 8 deletions(-) create mode 100644 src/wikibricks_curator/cli.py create mode 100644 src/wikibricks_curator/curator.py create mode 100644 src/wikibricks_curator/resources/local-curator.md diff --git a/AGENTS.md b/AGENTS.md index 2fcd180..e8e7478 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -19,6 +19,8 @@ WikiBricks is shared memory for multiple agent harnesses. when Lakebase, Databricks credentials, and the network are absent. - `import omnigent-server` is an explicit, opt-in network command like `sync lakebase`. +- `wikibricks-curator` is an optional, explicit network command that calls one + model. The core library still never calls a model. The library does not call a language model. The active agent makes semantic decisions about page content. Deterministic local maintenance repairs indexes diff --git a/CHANGELOG.md b/CHANGELOG.md index baefaa1..a9cf422 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,12 @@ # Changelog +## Unreleased + +- Added optional `wikibricks-curator propose` for local backlog curation. It + stores one model-generated run, auto-applies low-risk groups under the safe + policy, leaves other groups for review, and lists pending review runs in + `wiki_index`. + ## 0.12.3 - 2026-09-29 - Local curation now reports the newest projects whose sessions are newer than diff --git a/README.md b/README.md index 27e5275..abb80b7 100644 --- a/README.md +++ b/README.md @@ -225,6 +225,29 @@ only after every immutable event version has a committed archive receipt: wikibricks curate --prune-archived-sessions-after-days 90 ``` +## Nightly curation (optional) + +`wikibricks-curator propose` asks one Databricks model for updates to the top +local backlog projects. It writes those proposals as a curation run, applies +low-risk groups with the safe policy, and leaves high-risk or conflicting +groups for review. Run it only when you intend to call the model: + +```bash +wikibricks-curator propose \ + --database-path ~/.wikibricks/wikibricks.db \ + --base-url https://YOUR-DATABRICKS-WORKSPACE/serving-endpoints \ + --profile PROFILE +``` + +Use `--dry-run` to print proposals without writing. Use `--no-apply` to store +every run for inspection. `wiki_index` lists pending runs as +`_meta/curation-review`. After reading them, run: + +```bash +wikibricks sync plan RUN_ID +wikibricks sync apply RUN_ID --policy all +``` + ## Optional Lakebase curation Local memory does not need Lakebase. Configure it only when you want a remote diff --git a/pyproject.toml b/pyproject.toml index ac31809..e07a838 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -24,6 +24,7 @@ postgres-migration = [ [project.scripts] wikibricks = "wikibricks.cli:main" +wikibricks-curator = "wikibricks_curator.cli:main" wiki-init = "wikibricks.cli:init_main" wikibricks-mcp = "wikibricks.mcp_server:main" wikibricks-remote-maintenance = "wikibricks_remote.main:main" diff --git a/src/wikibricks/curation/__init__.py b/src/wikibricks/curation/__init__.py index 8878296..6302c2f 100644 --- a/src/wikibricks/curation/__init__.py +++ b/src/wikibricks/curation/__init__.py @@ -4,25 +4,31 @@ from wikibricks.curation.planning import plan_run from wikibricks.curation.protocol import ( build_manifest, + canonical_json, create_patch, validate_manifest, ) from wikibricks.curation.repository import ( get_or_create_replica_id, list_conflicts, + pending_review_runs, publish_manifest, pull_manifests, + store_manifest, ) __all__ = [ "apply_run", "build_manifest", + "canonical_json", "create_patch", "get_or_create_replica_id", "list_conflicts", + "pending_review_runs", "plan_run", "publish_manifest", "pull_manifests", "resolve_conflict", + "store_manifest", "validate_manifest", ] diff --git a/src/wikibricks/curation/repository.py b/src/wikibricks/curation/repository.py index 95ea035..047e146 100644 --- a/src/wikibricks/curation/repository.py +++ b/src/wikibricks/curation/repository.py @@ -136,6 +136,18 @@ def publish_manifest( return _insert_manifest(conn, checked, received=False) +def store_manifest( + store: PostgresStore | SQLiteStore, + manifest: dict[str, Any], +) -> bool: + checked = validate_manifest(manifest) + if isinstance(store, SQLiteStore): + with store.connection(write=True) as conn: + return _insert_manifest(conn, checked, received=True) + with store.connection() as conn, conn.transaction(): + return _insert_manifest(conn, checked, received=True) + + def pull_manifests( local: PostgresStore | SQLiteStore, remote: PostgresStore, @@ -153,13 +165,8 @@ def pull_manifests( received_runs = 0 received_patches = 0 for row in rows: - manifest = validate_manifest(dict(row[0])) - if isinstance(local, SQLiteStore): - with local.connection(write=True) as conn: - inserted = _insert_manifest(conn, manifest, received=True) - else: - with local.connection() as conn, conn.transaction(): - inserted = _insert_manifest(conn, manifest, received=True) + manifest = dict(row[0]) + inserted = store_manifest(local, manifest) if inserted: received_runs += 1 received_patches += len(manifest["patches"]) @@ -417,3 +424,17 @@ def list_conflicts(store: PostgresStore | SQLiteStore) -> list[dict[str, Any]]: } for row in rows ] + + +def pending_review_runs(store: PostgresStore | SQLiteStore) -> list[dict[str, Any]]: + store.migrate() + with store.connection() as conn: + rows = conn.execute( + "SELECT p.run_id, count(*) FROM curation_patches p " + "LEFT JOIN curation_receipts r ON r.patch_id = p.patch_id " + "WHERE r.patch_id IS NULL GROUP BY p.run_id ORDER BY p.run_id" + ).fetchall() + return [ + {"run_id": str(row[0]), "pending_patches": int(row[1])} + for row in rows + ] diff --git a/src/wikibricks/mcp_server.py b/src/wikibricks/mcp_server.py index af124d1..abd61c3 100644 --- a/src/wikibricks/mcp_server.py +++ b/src/wikibricks/mcp_server.py @@ -20,6 +20,7 @@ def _build_tools() -> dict[str, Any]: write_tools = make_agent_tools(database_path=str(client.database_path)) def wiki_index(prefix=None): from wikibricks.curation.backlog import load_curation_backlog + from wikibricks.curation.repository import pending_review_runs from wikibricks.maintenance import capture_status pages = [ @@ -56,6 +57,23 @@ def wiki_index(prefix=None): "items": backlog, } ) + pending_runs = pending_review_runs(client.store) + if pending_runs: + pending = sum(run["pending_patches"] for run in pending_runs) + pages.append( + { + "path": "_meta/curation-review", + "page_type": "review", + "title": f"Curation review: {pending} proposed changes wait for review", + "summary": ( + f"Tell the user about these proposed changes, or for example run " + f"`wikibricks sync plan {pending_runs[0]['run_id']}` and " + f"`wikibricks sync apply {pending_runs[0]['run_id']} --policy all` " + "after checking." + ), + "items": pending_runs, + } + ) return pages return { diff --git a/src/wikibricks_curator/cli.py b/src/wikibricks_curator/cli.py new file mode 100644 index 0000000..7850e5e --- /dev/null +++ b/src/wikibricks_curator/cli.py @@ -0,0 +1,54 @@ +"""Command-line entry point for the optional local curator.""" + +from __future__ import annotations + +import argparse +import json +from pathlib import Path +from typing import Any + +from wikibricks_curator.gateway import chat_json, resolve_token + + +def build_parser() -> argparse.ArgumentParser: + from wikibricks.config import load_config + + config = load_config() + parser = argparse.ArgumentParser(prog="wikibricks-curator") + commands = parser.add_subparsers(dest="command", required=True) + propose = commands.add_parser("propose", help="Propose local curation updates") + propose.add_argument("--database-path", type=Path, default=config.database_path) + propose.add_argument("--base-url", required=True) + propose.add_argument("--profile", default=config.sync_profile) + propose.add_argument("--model", default="system.ai.glm-5-3-flash") + propose.add_argument("--projects", type=int, default=3) + propose.add_argument("--no-apply", action="store_true") + propose.add_argument("--dry-run", action="store_true") + return parser + + +def main(argv: list[str] | None = None) -> int: + args = build_parser().parse_args(argv) + from wikibricks_curator.curator import run_curator + + token = resolve_token(args.profile) + + def chat(system_prompt: str, request: dict[str, Any], schema: dict[str, Any]) -> dict[str, Any]: + return chat_json( + system_prompt, + request, + schema, + base_url=args.base_url, + token=token, + model=args.model, + ) + + result = run_curator( + args.database_path, + chat=chat, + projects=args.projects, + apply=not args.no_apply, + dry_run=args.dry_run, + ) + print(json.dumps(result, default=str, indent=2)) + return 1 if result["errors"] else 0 diff --git a/src/wikibricks_curator/curator.py b/src/wikibricks_curator/curator.py new file mode 100644 index 0000000..942af06 --- /dev/null +++ b/src/wikibricks_curator/curator.py @@ -0,0 +1,191 @@ +"""Nightly local curation orchestration.""" + +from __future__ import annotations + +import hashlib +from dataclasses import replace +from importlib.resources import files +from pathlib import Path +from typing import Any, Callable +from uuid import uuid5 + +from wikibricks.curation import ( + apply_run, + build_manifest, + canonical_json, + get_or_create_replica_id, + store_manifest, +) +from wikibricks.curation.backlog import load_curation_backlog +from wikibricks.storage.sqlite_store import SQLiteStore +from wikibricks_curator.evidence import build_request +from wikibricks_remote.proposals import build_patches +from wikibricks_remote.resources import load_policy, load_prompt, load_schema + +Chat = Callable[[str, dict[str, Any], dict[str, Any]], dict[str, Any]] + + +def _living_page(item: dict[str, Any]) -> str: + return next( + (path for path in item.get("pages", []) if path.startswith("topics/")), + f"topics/{item['project']}", + ) + + +def _proposal_result(raw: dict[str, Any]) -> list[dict[str, Any]]: + return [ + { + "operation": proposal.get("operation"), + "path": proposal.get("path"), + "title": proposal.get("title"), + "risk_class": proposal.get("risk_class"), + "reason": proposal.get("reason"), + } + for proposal in raw.get("proposals", []) + if isinstance(proposal, dict) + ] + + +def _guard_shrinking_updates( + raw: dict[str, Any], + pages: list[dict[str, Any]], +) -> int: + proposals = raw.get("proposals") + if not isinstance(proposals, list): + return 0 + current = {page["path"]: page for page in pages} + guarded = 0 + for proposal in proposals: + if not isinstance(proposal, dict): + continue + if proposal.get("operation") != "update_page": + continue + path = proposal.get("path") + page = current.get(path) if isinstance(path, str) else None + if page is None: + continue + content = page.get("content") or {} + current_length = len(str(content.get("summary", ""))) + len( + str(content.get("body", "")) + ) + proposed_length = len(str(proposal.get("summary", ""))) + len( + str(proposal.get("body", "")) + ) + if proposed_length < current_length / 2: + proposal["risk_class"] = "high" + guarded += 1 + return guarded + + +def run_curator( + database_path: str | Path, + *, + chat: Chat, + projects: int = 3, + apply: bool = True, + dry_run: bool = False, +) -> dict[str, Any]: + """Propose curation for local backlog projects and optionally apply safe groups.""" + store = SQLiteStore(database_path) + store.migrate() + prompt = load_prompt() + "\n\n" + ( + files("wikibricks_curator") + .joinpath("resources", "local-curator.md") + .read_text(encoding="utf-8") + ) + schema = load_schema() + policy = replace(load_policy(), allowed_operations=("create_page", "update_page", "add_link")) + with store.connection() as conn: + backlog = load_curation_backlog(conn, limit=projects) + replica_id = get_or_create_replica_id(store) + results: list[dict[str, Any]] = [] + errors = 0 + + for item in backlog: + project = item["project"] + living_page = _living_page(item) + result: dict[str, Any] = { + "project": project, + "living_page": living_page, + "events": 0, + "status": "no_evidence", + "proposals": [], + "guarded": 0, + "run_id": None, + "counts": {}, + "error": None, + } + with store.connection() as conn: + cursor = store.get_sync_cursor(f"curator:{project}").get("newest_evidence_at") + built = build_request(conn, item, cursor=cursor) + if built is None: + results.append(result) + continue + result["events"] = len(built["request"]["evidence"]) + try: + raw = chat(prompt, built["request"], schema) + except Exception as exc: + result["status"] = "error" + result["error"] = str(exc)[:300] + results.append(result) + errors += 1 + continue + result["guarded"] = _guard_shrinking_updates(raw, built["pages"]) + result["proposals"] = _proposal_result(raw) + digest = hashlib.sha256( + canonical_json(built["request"]).encode("utf-8") + ).hexdigest() + run_id = uuid5(replica_id, f"wikibricks:local-curator:{project}:{digest}") + try: + patches = build_patches( + raw, + run_id=run_id, + pages=built["pages"], + evidence_ids=built["evidence_ids"], + policy=policy, + ) + except ValueError as exc: + result["status"] = "error" + result["error"] = str(exc)[:300] + results.append(result) + errors += 1 + continue + if not patches: + result["status"] = "no_changes" + results.append(result) + if not dry_run: + store.set_sync_cursor( + f"curator:{project}", + {"newest_evidence_at": built["newest_evidence_at"]}, + ) + continue + if dry_run: + result["status"] = "dry_run" + result["run_id"] = str(run_id) + results.append(result) + continue + manifest = build_manifest( + replica_id=replica_id, + input_watermark=0, + patches=patches, + run_id=run_id, + ) + store_manifest(store, manifest) + result["run_id"] = str(run_id) + if apply: + outcome = apply_run(store, run_id, policy="safe") + result["counts"] = outcome["counts"] + result["status"] = ( + "applied" + if outcome["counts"].get("applied") == len(outcome["groups"]) + else "review_required" + ) + else: + result["status"] = "review_required" + store.set_sync_cursor( + f"curator:{project}", + {"newest_evidence_at": built["newest_evidence_at"]}, + ) + results.append(result) + + return {"projects": results, "errors": errors} diff --git a/src/wikibricks_curator/gateway.py b/src/wikibricks_curator/gateway.py index 195e99f..d4a2d73 100644 --- a/src/wikibricks_curator/gateway.py +++ b/src/wikibricks_curator/gateway.py @@ -74,6 +74,7 @@ def chat_json( temperature: float = 0.0, max_tokens: int = 8192, timeout: float = 180, + reasoning_effort: str | None = "low", opener=urlopen, ) -> dict[str, Any]: """Call an OpenAI-compatible chat endpoint and return its JSON object.""" @@ -93,6 +94,9 @@ def chat_json( "temperature": temperature, "max_tokens": max_tokens, } + if reasoning_effort is not None: + # GLM 5.3 Flash otherwise spends the whole output budget on reasoning. + payload["reasoning_effort"] = reasoning_effort http_request = Request( f"{base_url.rstrip('/')}/chat/completions", data=json.dumps(payload, ensure_ascii=False).encode("utf-8"), @@ -103,7 +107,7 @@ def chat_json( method="POST", ) try: - response = opener(http_request, timeout) + response = opener(http_request, timeout=timeout) except HTTPError as error: raise RuntimeError(f"curation model request failed: HTTP {error.code}") from error return _json_content(response) diff --git a/src/wikibricks_curator/resources/local-curator.md b/src/wikibricks_curator/resources/local-curator.md new file mode 100644 index 0000000..8869403 --- /dev/null +++ b/src/wikibricks_curator/resources/local-curator.md @@ -0,0 +1,19 @@ +# Local nightly curation for one project + +This run covers one project. `request.project` names it and `request.living_page` is the +path of its living page. + +1. Keep one living page per project at `request.living_page`. If that path is in + `current_pages`, propose `update_page` for it; otherwise propose `create_page` with + `page_type` `entity`. +2. Write the living page body as Markdown with the sections `## Current state`, + `## Decisions`, `## Open issues`, and `## Key facts`. Keep it under 6000 characters. + Integrate the new evidence into the existing text: keep facts that are still true and + replace facts the evidence supersedes. +3. Propose `update_page` for another current page only when the evidence directly changes + a fact on it. +4. Propose `add_link` with `related` from the living page to current pages it builds on. +5. Use `risk_class` `low` for `create_page`, `update_page`, and `add_link`. Never propose + `retarget_links`, `add_alias`, or `supersede_page`. +6. Write dates as absolute dates. Never copy secrets, tokens, passwords, or credentials. +7. Return `{"proposals": []}` when the evidence holds nothing durable. diff --git a/tests/test_local_curator.py b/tests/test_local_curator.py index ce2984e..8c19d17 100644 --- a/tests/test_local_curator.py +++ b/tests/test_local_curator.py @@ -2,6 +2,8 @@ import json import subprocess +from datetime import datetime, timedelta, timezone +from importlib.resources import files from pathlib import Path from urllib.error import HTTPError from urllib.request import Request @@ -19,6 +21,10 @@ NEW_SESSION = "2026-01-04T11:00:00+00:00" +def _recent_timestamp(offset_hours: int = 1) -> str: + return (datetime.now(timezone.utc) - timedelta(hours=offset_hours)).isoformat() + + def _store(tmp_path: Path) -> SQLiteStore: store = SQLiteStore(tmp_path / "wikibricks.db") store.migrate() @@ -433,3 +439,441 @@ def test_build_request_compares_timestamps_across_utc_offsets(tmp_path: Path): assert [entry["text"] for entry in built["request"]["evidence"]] == ["after the page update"] assert [entry["text"] for entry in later["request"]["evidence"]] == ["after the page update"] + + +def _backlog_project( + store: SQLiteStore, + *, + project: str, + page: bool = False, +) -> dict: + from wikibricks.curation.backlog import load_curation_backlog + + if page: + store.write_page( + f"topics/{project}", + project.title(), + {"summary": "current summary", "body": "The project page contains durable facts."}, + ) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = ?", + (_recent_timestamp(3), f"topics/{project}"), + ) + store.ingest_session( + SessionRecord( + harness="test-harness", + external_id=f"{project}-new", + user_id="user", + workspace=f"/Users/u/work/{project}", + updated_at=_recent_timestamp(1), + events=[SessionEvent("user-1", "user", f"{project} learned one durable fact")], + ) + ) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE sessions SET updated_at = ? WHERE external_id = ?", + (_recent_timestamp(1), f"{project}-new"), + ) + with store.connection() as conn: + backlog = load_curation_backlog(conn, limit=3) + return next(item for item in backlog if item["project"] == project) + + +def _proposal( + *, + operation: str, + path: str, + evidence_id: str, + body: str = "Derived from session evidence.", + summary: str = "Updated summary", +) -> dict: + return { + "group": "main", + "operation": operation, + "path": path, + "title": path.rsplit("/", 1)[-1].title(), + "page_type": "entity", + "summary": summary, + "body": body, + "tags": [], + "source_ids": [], + "target_path": None, + "evidence_ids": [evidence_id], + "reason": "The session evidence changed the project state.", + "risk_class": "low", + } + + +def test_run_curator_creates_page_then_stops_at_cursor(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="create-project") + + def chat(system_prompt: str, request: dict, schema: dict) -> dict: + assert system_prompt + assert schema + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + )]} + + result = run_curator(store.database_path, chat=chat, projects=1) + project = result["projects"][0] + assert project["status"] == "applied" + assert project["guarded"] == 0 + assert project["counts"] == {"applied": 1} + assert project["proposals"][0]["operation"] == "create_page" + assert store.read_page("topics/create-project")["version"] == 1 + assert store.get_sync_cursor("curator:create-project")["newest_evidence_at"] + + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = 'topics/create-project'", + (_recent_timestamp(3),), + ) + + second = run_curator(store.database_path, chat=chat, projects=1) + assert second["projects"][0]["status"] == "no_evidence" + assert second["projects"][0]["run_id"] is None + + +def test_run_curator_updates_covered_page(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + item = _backlog_project(store, project="update-project", page=True) + assert item["pages"] == ["topics/update-project"] + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + return {"proposals": [_proposal( + operation="update_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + summary="New current summary", + body="New body is at least half the current page length.", + )]} + + result = run_curator(store.database_path, chat=chat, projects=1) + assert result["projects"][0]["status"] == "applied" + page = store.read_page("topics/update-project") + assert page["version"] == 2 + assert page["content"] == { + "summary": "New current summary", + "body": "New body is at least half the current page length.", + } + + +def test_shrink_guard_leaves_page_for_review( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +): + from uuid import UUID + + from wikibricks import mcp_server + from wikibricks.curation import apply_run + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="guard-project", page=True) + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + return {"proposals": [_proposal( + operation="update_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + summary="tiny", + body="tiny", + )]} + + result = run_curator(store.database_path, chat=chat, projects=1) + project = result["projects"][0] + assert project["status"] == "review_required" + assert project["guarded"] == 1 + assert store.read_page("topics/guard-project")["version"] == 1 + + monkeypatch.setenv("WIKIBRICKS_DATABASE_PATH", str(store.database_path)) + indexed = mcp_server.dispatch_tool( + "wiki_index", {}, tools=mcp_server._build_tools() + ) + review = [page for page in indexed if page["path"] == "_meta/curation-review"] + assert len(review) == 1 + assert review[0]["items"] == [{ + "run_id": project["run_id"], + "pending_patches": 1, + }] + + applied = apply_run(store, UUID(project["run_id"]), policy="all") + assert applied["counts"] == {"applied": 1} + indexed = mcp_server.dispatch_tool( + "wiki_index", {}, tools=mcp_server._build_tools() + ) + assert not [page for page in indexed if page["path"] == "_meta/curation-review"] + + +def test_run_curator_reports_errors_and_processes_other_projects(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="failed-project") + _backlog_project(store, project="passed-project") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + if request["project"] == "failed-project": + raise RuntimeError("model failed") + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + )]} + + result = run_curator(store.database_path, chat=chat, projects=2) + statuses = {item["project"]: item["status"] for item in result["projects"]} + assert statuses == {"failed-project": "error", "passed-project": "applied"} + assert result["errors"] == 1 + assert store.get_sync_cursor("curator:failed-project") == {} + assert store.get_sync_cursor("curator:passed-project") + + +def test_run_curator_rejects_unknown_evidence(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="unknown-evidence") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id="session-event:not-real", + )]} + + result = run_curator(store.database_path, chat=chat, projects=1) + assert result["errors"] == 1 + assert "unknown evidence" in result["projects"][0]["error"] + assert store.read_page("topics/unknown-evidence") is None + + +def test_dry_run_reports_proposals_without_writes(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="dry-run") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + )]} + + result = run_curator(store.database_path, chat=chat, projects=1, dry_run=True) + project = result["projects"][0] + assert project["status"] == "dry_run" + assert project["proposals"] == [{ + "operation": "create_page", + "path": "topics/dry-run", + "title": "Dry-Run", + "risk_class": "low", + "reason": "The session evidence changed the project state.", + }] + assert store.read_page("topics/dry-run") is None + with store.connection() as conn: + assert conn.execute("SELECT count(*) FROM curation_runs").fetchone()[0] == 0 + assert store.get_sync_cursor("curator:dry-run") == {} + + +def test_prompt_schema_and_policy_reject_forbidden_operations(tmp_path: Path): + from wikibricks_curator.curator import run_curator + from wikibricks_remote.resources import load_prompt, load_schema + + store = _store(tmp_path) + _backlog_project(store, project="policy-project", page=True) + calls: list[tuple[str, dict, dict]] = [] + + def chat(system_prompt: str, request: dict, schema: dict) -> dict: + calls.append((system_prompt, request, schema)) + return {"proposals": [{ + **_proposal( + operation="supersede_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + ), + "target_path": "topics/other", + }]} + + result = run_curator(store.database_path, chat=chat, projects=1) + prompt, _request, schema = calls[0] + assert prompt.startswith(load_prompt() + "\n\n") + assert schema == load_schema() + assert "# Local nightly curation for one project" in prompt + assert result["errors"] == 1 + assert "disabled by remote policy" in result["projects"][0]["error"] + + +def test_cli_propose_binds_gateway_and_reports_errors( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + capsys: pytest.CaptureFixture[str], +): + from wikibricks_curator import cli + + database_path = tmp_path / "wikibricks.db" + store = SQLiteStore(database_path) + store.migrate() + _backlog_project(store, project="cli-project") + tokens: list[str | None] = [] + + def resolve_token(profile: str | None) -> str: + tokens.append(profile) + return "test-token" + + def chat_json(_prompt, request, _schema, *, base_url, token, model): + assert base_url == "https://gateway.example" + assert token == "test-token" + assert model == "system.ai.glm-5-3-flash" + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + )]} + + monkeypatch.setattr(cli, "resolve_token", resolve_token) + monkeypatch.setattr(cli, "chat_json", chat_json) + code = cli.main([ + "propose", + "--database-path", + str(database_path), + "--base-url", + "https://gateway.example", + "--projects", + "1", + ]) + output = json.loads(capsys.readouterr().out) + assert code == 0 + assert tokens == [None] + assert output["projects"][0]["project"] == "cli-project" + assert store.read_page("topics/cli-project") + + _backlog_project(store, project="cli-error") + + def failing_chat_json(*_args, **_kwargs): + raise RuntimeError("model failed") + + monkeypatch.setattr(cli, "chat_json", failing_chat_json) + assert cli.main([ + "propose", + "--database-path", + str(database_path), + "--base-url", + "https://gateway.example", + "--projects", + "2", + "--profile", + "staging", + "--no-apply", + ]) == 1 + output = json.loads(capsys.readouterr().out) + assert output["errors"] == 1 + + +def test_store_manifest_is_idempotent(tmp_path: Path): + from uuid import uuid4 + + from wikibricks.curation import build_manifest, store_manifest + from wikibricks_remote.proposals import build_patches + from wikibricks_remote.resources import load_policy + + store = _store(tmp_path) + patches = build_patches( + {"proposals": [_proposal( + operation="create_page", + path="topics/new", + evidence_id="session-event:test", + )]}, + run_id=uuid4(), + pages=[], + evidence_ids={"session-event:test"}, + policy=load_policy(), + ) + manifest = build_manifest( + replica_id=uuid4(), + input_watermark=0, + patches=patches, + ) + assert store_manifest(store, manifest) is True + assert store_manifest(store, manifest) is False + + +def test_run_curator_uses_a_deterministic_run_id_for_the_same_input(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="stable-run") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + return {"proposals": [_proposal( + operation="create_page", + path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"], + )]} + + first = run_curator(store.database_path, chat=chat, projects=1, dry_run=True) + second = run_curator(store.database_path, chat=chat, projects=1, dry_run=True) + assert first["projects"][0]["run_id"] == second["projects"][0]["run_id"] + + +def test_local_curator_resource_is_packaged(): + resource = files("wikibricks_curator").joinpath("resources", "local-curator.md") + assert resource.read_text(encoding="utf-8").startswith( + "# Local nightly curation for one project" + ) + + +def test_dry_run_without_changes_does_not_advance_the_cursor(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="quiet") + + result = run_curator( + store.database_path, chat=lambda *_: {"proposals": []}, projects=1, dry_run=True + ) + + assert result["projects"][0]["status"] == "no_changes" + assert store.get_sync_cursor("curator:quiet") == {} + + +def test_shrink_guard_handles_pages_without_a_summary(): + from wikibricks_curator.curator import _guard_shrinking_updates + + raw = {"proposals": [{"operation": "update_page", "path": "topics/x", "summary": "", "body": "x"}]} + pages = [{"path": "topics/x", "content": {"body": "long body " * 20}}] + + assert _guard_shrinking_updates(raw, pages) == 1 + assert raw["proposals"][0]["risk_class"] == "high" + + +def test_chat_json_passes_timeout_by_keyword_and_asks_for_low_reasoning(): + # urllib.request.urlopen(url, data, timeout): a positional timeout becomes the body. + from wikibricks_curator.gateway import chat_json + + seen: dict = {} + + def keyword_only_opener(request: Request, *, timeout: float) -> _Response: + seen["timeout"] = timeout + seen["payload"] = _payload(request) + return _Response(json.dumps({"choices": [{"message": {"content": '{"proposals": []}'}}]})) + + result = chat_json( + "prompt", {"a": 1}, {"type": "object"}, + base_url="https://gateway.example/v1", token="t", model="m", timeout=42, + opener=keyword_only_opener, + ) + + assert result == {"proposals": []} + assert seen["timeout"] == 42 + # GLM 5.3 Flash spends the whole output budget on reasoning without this. + assert seen["payload"]["reasoning_effort"] == "low" From 9281ed7310c237708e29bb83f708ee87b606f463 Mon Sep 17 00:00:00 2001 From: Philipp Date: Tue, 29 Sep 2026 10:01:25 +0200 Subject: [PATCH 3/3] fix(curator): tolerate model omissions found in acceptance Acceptance runs against GLM 5.3 Flash on a copy of the live store showed three failure modes; each lost a whole project run: - A link from the living page created in the same reply: build_patches only knows existing pages. Such links are now dropped and counted; `curate` adds mention links once the page exists. - Omitted fields with no content value (risk_class, tags, source_ids, target_path, group; title/summary/body on links) get neutral defaults, and unknown keys are dropped. Page proposals without content still fail. - An undecodable or invalid reply is retried once. After the fixes: 8 of 8 projects applied, 0 errors, 100 s for 6 projects. Co-authored-by: Isaac --- src/wikibricks_curator/curator.py | 91 ++++++++++++++++++++++++------- tests/test_local_curator.py | 86 +++++++++++++++++++++++++++++ 2 files changed, 156 insertions(+), 21 deletions(-) diff --git a/src/wikibricks_curator/curator.py b/src/wikibricks_curator/curator.py index 942af06..2759177 100644 --- a/src/wikibricks_curator/curator.py +++ b/src/wikibricks_curator/curator.py @@ -19,7 +19,7 @@ from wikibricks.curation.backlog import load_curation_backlog from wikibricks.storage.sqlite_store import SQLiteStore from wikibricks_curator.evidence import build_request -from wikibricks_remote.proposals import build_patches +from wikibricks_remote.proposals import _PROPOSAL_FIELDS, build_patches from wikibricks_remote.resources import load_policy, load_prompt, load_schema Chat = Callable[[str, dict[str, Any], dict[str, Any]], dict[str, Any]] @@ -77,6 +77,52 @@ def _guard_shrinking_updates( return guarded +_LINK_ONLY_FIELDS = ("title", "page_type", "summary", "body") + + +def _fill_neutral_defaults(raw: dict[str, Any]) -> None: + """Fill fields a model may omit whose value carries no content; drop unknown keys. + + Never invents content, evidence, or reasons: page proposals without a title or + body still fail validation. + """ + proposals = raw.get("proposals") if isinstance(raw, dict) else None + if not isinstance(proposals, list): + return + for proposal in proposals: + if not isinstance(proposal, dict): + continue + for key in set(proposal) - _PROPOSAL_FIELDS: + del proposal[key] + proposal.setdefault("group", "main") + proposal.setdefault("tags", []) + proposal.setdefault("source_ids", []) + proposal.setdefault("target_path", None) + proposal.setdefault("risk_class", "low") + if proposal.get("operation") == "add_link": + for key in _LINK_ONLY_FIELDS: + proposal.setdefault(key, "") + + +def _drop_unresolvable_links(raw: dict[str, Any], pages: list[dict[str, Any]]) -> int: + """Drop links whose ends are not existing pages; `curate` links new pages later.""" + proposals = raw.get("proposals") + if not isinstance(proposals, list): + return 0 + existing = {page["path"] for page in pages} + kept = [ + proposal + for proposal in proposals + if not ( + isinstance(proposal, dict) + and proposal.get("operation") == "add_link" + and not {proposal.get("path"), proposal.get("target_path")} <= existing + ) + ] + raw["proposals"] = kept + return len(proposals) - len(kept) + + def run_curator( database_path: str | Path, *, @@ -111,6 +157,8 @@ def run_curator( "status": "no_evidence", "proposals": [], "guarded": 0, + "dropped_links": 0, + "attempts": 0, "run_id": None, "counts": {}, "error": None, @@ -122,31 +170,32 @@ def run_curator( results.append(result) continue result["events"] = len(built["request"]["evidence"]) - try: - raw = chat(prompt, built["request"], schema) - except Exception as exc: - result["status"] = "error" - result["error"] = str(exc)[:300] - results.append(result) - errors += 1 - continue - result["guarded"] = _guard_shrinking_updates(raw, built["pages"]) - result["proposals"] = _proposal_result(raw) digest = hashlib.sha256( canonical_json(built["request"]).encode("utf-8") ).hexdigest() run_id = uuid5(replica_id, f"wikibricks:local-curator:{project}:{digest}") - try: - patches = build_patches( - raw, - run_id=run_id, - pages=built["pages"], - evidence_ids=built["evidence_ids"], - policy=policy, - ) - except ValueError as exc: + # Model output varies between calls; one retry absorbs an invalid reply. + for attempt in (1, 2): + result["attempts"] = attempt + try: + raw = chat(prompt, built["request"], schema) + _fill_neutral_defaults(raw) + result["dropped_links"] = _drop_unresolvable_links(raw, built["pages"]) + result["guarded"] = _guard_shrinking_updates(raw, built["pages"]) + result["proposals"] = _proposal_result(raw) + patches = build_patches( + raw, + run_id=run_id, + pages=built["pages"], + evidence_ids=built["evidence_ids"], + policy=policy, + ) + result["error"] = None + break + except Exception as exc: + result["error"] = str(exc)[:300] + if result["error"] is not None: result["status"] = "error" - result["error"] = str(exc)[:300] results.append(result) errors += 1 continue diff --git a/tests/test_local_curator.py b/tests/test_local_curator.py index 8c19d17..f6d6aac 100644 --- a/tests/test_local_curator.py +++ b/tests/test_local_curator.py @@ -877,3 +877,89 @@ def keyword_only_opener(request: Request, *, timeout: float) -> _Response: assert seen["timeout"] == 42 # GLM 5.3 Flash spends the whole output budget on reasoning without this. assert seen["payload"]["reasoning_effort"] == "low" + + +def test_links_from_a_page_created_in_the_same_run_are_dropped_not_fatal(tmp_path: Path): + # Real GLM output: create the living page and link it in one reply. build_patches only + # knows existing pages, so the link must not cost the page. + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + store.write_page("topics/other", "Other", {"summary": "s", "body": "Existing page."}) + _backlog_project(store, project="fresh") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + evidence_id = request["evidence"][0]["evidence_id"] + link = _proposal(operation="add_link", path=request["living_page"], evidence_id=evidence_id) + link.update({"target_path": "topics/other", "link_type": "related", "title": "", + "page_type": "entity", "summary": "", "body": ""}) + return {"proposals": [ + _proposal(operation="create_page", path=request["living_page"], evidence_id=evidence_id), + link, + ]} + + result = run_curator(store.database_path, chat=chat, projects=1) + project = result["projects"][0] + + assert project["status"] == "applied", project["error"] + assert project["dropped_links"] == 1 + assert store.read_page("topics/fresh") is not None + + +def test_neutral_defaults_fill_fields_models_omit(tmp_path: Path): + # Real GLM output omitted risk_class, and links rarely carry title/body/tags. + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="omits", page=True) + store.write_page("topics/other", "Other", {"summary": "s", "body": "Existing page."}) + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + evidence_id = request["evidence"][0]["evidence_id"] + update = _proposal(operation="update_page", path=request["living_page"], evidence_id=evidence_id, + body="The project page contains durable facts and one more fact.") + del update["risk_class"], update["tags"], update["group"] + update["confidence"] = 0.9 + link = {"operation": "add_link", "path": request["living_page"], "target_path": "topics/other", + "link_type": "related", "evidence_ids": [evidence_id], "reason": "Builds on it."} + return {"proposals": [update, link]} + + project = run_curator(store.database_path, chat=chat, projects=1)["projects"][0] + + assert project["status"] == "applied", project["error"] + assert "one more fact" in store.read_page("topics/omits")["content_text"] + + +def test_page_proposals_without_content_still_fail(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="empty") + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + proposal = _proposal(operation="create_page", path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"]) + del proposal["body"] + return {"proposals": [proposal]} + + assert run_curator(store.database_path, chat=chat, projects=1)["projects"][0]["status"] == "error" + + +def test_one_retry_after_an_invalid_model_reply(tmp_path: Path): + from wikibricks_curator.curator import run_curator + + store = _store(tmp_path) + _backlog_project(store, project="flaky") + replies = [RuntimeError("curation model output could not be decoded"), None] + + def chat(_prompt: str, request: dict, _schema: dict) -> dict: + reply = replies.pop(0) + if isinstance(reply, Exception): + raise reply + return {"proposals": [_proposal(operation="create_page", path=request["living_page"], + evidence_id=request["evidence"][0]["evidence_id"])]} + + project = run_curator(store.database_path, chat=chat, projects=1)["projects"][0] + + assert project["status"] == "applied", project["error"] + assert project["attempts"] == 2