diff --git a/CHANGELOG.md b/CHANGELOG.md index b26f945..7853757 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,9 @@ ## Unreleased +- Group curation sessions by routed topic pages instead of workspace folder names. + The nightly curator asks the model which topic each container-folder session is + about before building the backlog. - Label pages and links written by the local curator as `local-curator`. ## 0.13.0 - 2026-09-29 diff --git a/README.md b/README.md index d7d42a8..e890d88 100644 --- a/README.md +++ b/README.md @@ -227,6 +227,10 @@ wikibricks curate --prune-archived-sessions-after-days 90 ## Nightly curation (optional) +Sessions captured under container folders such as `~/code` do not identify their +topic. Before building the backlog, the curator asks the model to route those +sessions to an existing or new page, subject to `curation.generic_workspaces`. + `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 diff --git a/src/wikibricks/config/__init__.py b/src/wikibricks/config/__init__.py index 9180e00..b9956b4 100644 --- a/src/wikibricks/config/__init__.py +++ b/src/wikibricks/config/__init__.py @@ -16,6 +16,7 @@ "database": {"path": None, "url": None}, "search": {"default_results": None, "maximum_results": None}, "maintenance": {"prune_archived_sessions_after_days": None}, + "curation": {"generic_workspaces": None}, "automation": { "enabled": None, "poll_seconds": None, @@ -41,6 +42,7 @@ class WikiBricksConfig: search_default_results: int search_maximum_results: int prune_archived_sessions_after_days: int | None + curation_generic_workspaces: tuple[str, ...] automation_enabled: bool automation_poll_seconds: int automation_local_maintenance_hours: int @@ -115,6 +117,21 @@ def _string(value: Any, path: str, *, optional: bool = False) -> str | None: return value.strip() +def _string_list(value: Any, path: str) -> list[str]: + if not isinstance(value, list) or not all(isinstance(item, str) for item in value): + raise ValueError(f"{path} must be a list of strings") + if any(not item.strip() for item in value): + raise ValueError(f"{path} must not contain empty strings") + return [item.strip() for item in value] + + +def _comma_separated_strings(value: str) -> list[str]: + values = [item.strip() for item in value.split(",")] + if any(not item for item in values): + raise ValueError("must be a comma-separated list of non-empty strings") + return values + + def _environment_overlay(environ: Mapping[str, str]) -> dict[str, Any]: result: dict[str, Any] = {} mappings = { @@ -127,6 +144,11 @@ def _environment_overlay(environ: Mapping[str, str]) -> dict[str, Any]: "prune_archived_sessions_after_days", int, ), + "WIKIBRICKS_CURATION_GENERIC_WORKSPACES": ( + "curation", + "generic_workspaces", + _comma_separated_strings, + ), "WIKIBRICKS_SYNC_BATCH_SIZE": ("sync", "batch_size", int), "WIKIBRICKS_SYNC_APPLY_POLICY": ("sync", "apply_policy", str), "WIKIBRICKS_AUTOMATION_ENABLED": ("automation", "enabled", str), @@ -215,6 +237,10 @@ def load_config( minimum=1, maximum=36500, ) + generic_workspaces = _string_list( + value["curation"]["generic_workspaces"], + "curation.generic_workspaces", + ) automation_enabled = _boolean( value["automation"]["enabled"], "automation.enabled", @@ -260,6 +286,7 @@ def load_config( search_default_results=default_results, search_maximum_results=maximum_results, prune_archived_sessions_after_days=retention, + curation_generic_workspaces=tuple(generic_workspaces), automation_enabled=automation_enabled, automation_poll_seconds=poll_seconds, automation_local_maintenance_hours=local_maintenance_hours, diff --git a/src/wikibricks/config/defaults.yml b/src/wikibricks/config/defaults.yml index 0809dcd..1ce7ee1 100644 --- a/src/wikibricks/config/defaults.yml +++ b/src/wikibricks/config/defaults.yml @@ -7,6 +7,18 @@ search: maximum_results: 20 maintenance: prune_archived_sessions_after_days: null +curation: + generic_workspaces: + - code + - emails + - work + - projects + - repos + - src + - documents + - desktop + - downloads + - tmp automation: enabled: true poll_seconds: 300 diff --git a/src/wikibricks/curation/backlog.py b/src/wikibricks/curation/backlog.py index e044064..7dbd3fb 100644 --- a/src/wikibricks/curation/backlog.py +++ b/src/wikibricks/curation/backlog.py @@ -5,11 +5,23 @@ from pathlib import Path from typing import Any +from wikibricks.config import load_config + def _normalize(value: str) -> str: return value.strip().lower().replace(" ", "-").replace("_", "-") +def _normalized_name(value: str) -> str: + return _normalize(Path(value).name) + + +def _unrouted_workspace(workspace: str | None, *, generic: set[str], home: Path) -> bool: + if workspace is None or not workspace or Path(workspace) == home: + return True + return _normalized_name(workspace) in generic + + def _timestamp(value: str | datetime) -> datetime: if isinstance(value, datetime): parsed = value @@ -21,56 +33,70 @@ def _timestamp(value: str | datetime) -> datetime: def curation_backlog( - sessions: Iterable[tuple[str | None, str | datetime]], + sessions: Iterable[tuple[str, str | None, str | datetime]], pages: Iterable[tuple[str, str, str | datetime]], *, now: datetime, days: int = 7, limit: int = 5, - home: Path | None = None, + targets: dict[str, str | None], ) -> list[dict[str, Any]]: """Return recent sessions that are newer than their covering pages.""" now = _timestamp(now) - home = home or Path.home() cutoff = now - timedelta(days=days) + page_updates = { + path: _timestamp(updated_at) for path, _, updated_at in pages if not path.startswith("_meta/") + } + page_titles = {path: title for path, title, _ in pages} projects: dict[str, dict[str, Any]] = {} - for workspace, updated_at in sessions: - if workspace is None or not workspace or Path(workspace) == home: + for session_id, workspace, updated_at in sessions: + target = targets.get(session_id) + if target is None: continue timestamp = _timestamp(updated_at) if timestamp < cutoff or timestamp > now: continue - project = _normalize(Path(workspace).name) + project = Path(target).name item = projects.setdefault( - project, + target, { "project": project, + "living_page": target, "workspace": workspace, + "workspaces": set(), "session_timestamps": [], }, ) + if workspace is not None and workspace: + item["workspaces"].add(workspace) + if ( + item["workspace"] is None + or not item["workspace"] + or workspace < item["workspace"] + ): + item["workspace"] = workspace item["session_timestamps"].append(timestamp) - covering: dict[str, list[tuple[str, datetime]]] = {} - for path, title, updated_at in pages: - if path.startswith("_meta/"): - continue - normalized_path = _normalize(path) - normalized_title = _normalize(title) - timestamp = _timestamp(updated_at) - for project, item in projects.items(): - if project in normalized_path or project in normalized_title: - covering.setdefault(project, []).append((path, timestamp)) - result: list[dict[str, Any]] = [] - for project, item in projects.items(): - project_pages = sorted(covering.get(project, [])) + for target, item in projects.items(): + project = item["project"] + living_page = item["living_page"] + last_page_update = page_updates.get(living_page) + project_pages: list[str] = [] + if last_page_update is not None: + project_pages.append(living_page) + project_pages.extend( + path + for path in sorted(page_updates) + if path != living_page + and path.startswith("topics/") + and ( + project in _normalize(path) + or project in _normalize(page_titles.get(path, "")) + ) + ) last_session_at = max(item["session_timestamps"]) - if project_pages: - last_page_update = max(timestamp for _, timestamp in project_pages) - else: - last_page_update = None new_sessions = sum( last_page_update is None or timestamp > last_page_update for timestamp in item["session_timestamps"] @@ -81,8 +107,10 @@ def curation_backlog( result.append( { **item, + "workspaces": sorted(item["workspaces"]), + "workspace": item["workspace"], "new_sessions": new_sessions, - "pages": [path for path, _ in project_pages[:3]], + "pages": project_pages[:3], "last_page_update": ( last_page_update.isoformat() if last_page_update else None ), @@ -95,16 +123,109 @@ def curation_backlog( )[:limit] +def session_targets( + conn: Any, + *, + generic: Iterable[str] | None = None, + home: Path | None = None, +) -> dict[str, str | None]: + """Route every session to the durable page that covers its work.""" + home = home or Path.home() + normalized_generic = { + _normalize(value) for value in ( + load_config().curation_generic_workspaces + if generic is None + else generic + ) + } + session_rows = conn.execute( + "SELECT s.session_id, s.workspace, t.session_id, t.page_path FROM sessions s " + "LEFT JOIN session_topics t ON t.session_id = s.session_id" + ).fetchall() + page_rows = conn.execute( + "SELECT p.path, v.title FROM pages p " + "JOIN page_versions v ON v.version_id = p.current_version_id " + "WHERE p.status = 'active' AND p.path LIKE 'topics/%' ORDER BY p.path" + ).fetchall() + pages = [(row[0], row[1]) for row in page_rows] + targets: dict[str, str | None] = {} + for session_id, workspace, topic_session_id, routed_page in session_rows: + if topic_session_id is not None: + targets[session_id] = routed_page + continue + if _unrouted_workspace(workspace, generic=normalized_generic, home=home): + targets[session_id] = None + continue + key = _normalized_name(workspace) + target = next( + ( + path + for path, title in pages + if key in _normalize(path) or key in _normalize(title) + ), + f"topics/{key}", + ) + targets[session_id] = target + return targets + + +def unrouted_sessions( + conn: Any, + *, + since: str | datetime, + generic: Iterable[str] | None = None, + home: Path | None = None, +) -> list[str]: + """Return container-folder sessions still waiting for topic routing.""" + home = home or Path.home() + normalized_generic = { + _normalize(value) for value in ( + load_config().curation_generic_workspaces + if generic is None + else generic + ) + } + cutoff = _timestamp(since) + rows = conn.execute( + "SELECT s.session_id, s.workspace, " + "COALESCE(s.source_updated_at, s.updated_at) FROM sessions s " + "LEFT JOIN session_topics t ON t.session_id = s.session_id " + "WHERE t.session_id IS NULL ORDER BY s.session_id" + ).fetchall() + return [ + session_id + for session_id, workspace, updated_at in rows + if _timestamp(updated_at) >= cutoff + and _unrouted_workspace(workspace, generic=normalized_generic, home=home) + ] + + def load_curation_backlog(conn: Any, **kwargs: Any) -> list[dict[str, Any]]: """Load sessions and active pages, then calculate the curation backlog.""" # Count sessions by when the work happened, not when a backfill imported them. sessions = conn.execute( - "SELECT workspace, COALESCE(source_updated_at, updated_at) FROM sessions" + "SELECT s.session_id, s.workspace, COALESCE(s.source_updated_at, s.updated_at), " + "t.created_at " + "FROM sessions s LEFT JOIN session_topics t ON t.session_id = s.session_id" ).fetchall() + sessions = [ + ( + row[0], + row[1], + _timestamp(row[2]) if row[3] is None else max(_timestamp(row[2]), _timestamp(row[3])), + ) + for row in sessions + ] pages = conn.execute( "SELECT p.path, v.title, p.updated_at FROM pages p " "JOIN page_versions v ON v.version_id = p.current_version_id " "WHERE p.status = 'active'" ).fetchall() + home = kwargs.pop("home", None) + generic = kwargs.pop("generic", None) + if generic is None: + generic = load_config().curation_generic_workspaces + targets = session_targets(conn, generic=generic, home=home) kwargs.setdefault("now", datetime.now(timezone.utc)) + kwargs["targets"] = targets return curation_backlog(sessions, pages, **kwargs) diff --git a/src/wikibricks/sql/migrations/0007_session_topics.sql b/src/wikibricks/sql/migrations/0007_session_topics.sql new file mode 100644 index 0000000..80daa18 --- /dev/null +++ b/src/wikibricks/sql/migrations/0007_session_topics.sql @@ -0,0 +1,6 @@ +CREATE TABLE session_topics ( + session_id uuid PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE, + page_path text, + origin text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() +); diff --git a/src/wikibricks/sql/sqlite/0004_session_topics.sql b/src/wikibricks/sql/sqlite/0004_session_topics.sql new file mode 100644 index 0000000..b11f0f9 --- /dev/null +++ b/src/wikibricks/sql/sqlite/0004_session_topics.sql @@ -0,0 +1,6 @@ +CREATE TABLE IF NOT EXISTS session_topics ( + session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE, + page_path TEXT, + origin TEXT NOT NULL, + created_at TEXT NOT NULL +); diff --git a/src/wikibricks/storage/store.py b/src/wikibricks/storage/store.py index 3a2109c..50987ca 100644 --- a/src/wikibricks/storage/store.py +++ b/src/wikibricks/storage/store.py @@ -76,8 +76,8 @@ def migrate(self) -> None: def clear_all(self) -> None: tables = ( - "remote_maintenance_runs, curation_conflicts, curation_receipts, " - "page_aliases, curation_patches, " + "session_topics, remote_maintenance_runs, curation_conflicts, " + "curation_receipts, page_aliases, curation_patches, " "curation_runs, sync_replicas, archive_batch_events, archive_events, " "archive_batches, curated_pages, archive_pages, sync_state, sync_outbox, " "session_search_chunks, session_event_versions, session_events, sessions, " diff --git a/src/wikibricks_curator/curator.py b/src/wikibricks_curator/curator.py index 558000f..9514c29 100644 --- a/src/wikibricks_curator/curator.py +++ b/src/wikibricks_curator/curator.py @@ -4,6 +4,7 @@ import hashlib from dataclasses import replace +from datetime import datetime, timedelta, timezone from importlib.resources import files from pathlib import Path from typing import Any, Callable @@ -19,19 +20,13 @@ from wikibricks.curation.backlog import load_curation_backlog from wikibricks.storage.sqlite_store import SQLiteStore from wikibricks_curator.evidence import build_request +from wikibricks_curator.router import route_sessions 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]] -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 [ { @@ -141,6 +136,12 @@ def run_curator( ) schema = load_schema() policy = replace(load_policy(), allowed_operations=("create_page", "update_page", "add_link")) + routing = route_sessions( + store, + chat, + since=datetime.now(timezone.utc) - timedelta(days=7), + dry_run=dry_run, + ) with store.connection() as conn: backlog = load_curation_backlog(conn, limit=projects) replica_id = get_or_create_replica_id(store) @@ -149,7 +150,7 @@ def run_curator( for item in backlog: project = item["project"] - living_page = _living_page(item) + living_page = item["living_page"] result: dict[str, Any] = { "project": project, "living_page": living_page, @@ -242,4 +243,4 @@ def run_curator( ) results.append(result) - return {"projects": results, "errors": errors} + return {"projects": results, "errors": errors, "routing": routing} diff --git a/src/wikibricks_curator/evidence.py b/src/wikibricks_curator/evidence.py index c2e3510..e0635eb 100644 --- a/src/wikibricks_curator/evidence.py +++ b/src/wikibricks_curator/evidence.py @@ -8,13 +8,10 @@ from pathlib import Path from typing import Any +from wikibricks.curation.backlog import session_targets 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 "" @@ -91,7 +88,7 @@ def _moment(value: Any) -> datetime: def _evidence( conn: sqlite3.Connection, *, - project: str, + living_page: str, cursor: str | None, last_page_update: str | None, max_events: int, @@ -100,13 +97,20 @@ def _evidence( ) -> 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] + targets = session_targets(conn) + session_rows = conn.execute( + "SELECT s.session_id, s.workspace, COALESCE(s.source_updated_at, s.updated_at), " + "t.created_at " + "FROM sessions s LEFT JOIN session_topics t ON t.session_id = s.session_id" + ).fetchall() 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) + ( + _moment(row[2]) if row[3] is None else max(_moment(row[2]), _moment(row[3])), + row[0], + ) + for row in session_rows + if targets.get(row[0]) == living_page + and all(_moment(row[2] if row[3] is None else max(_moment(row[2]), _moment(row[3]))) > bound for bound in after) ] if not sessions: return [], "" @@ -166,10 +170,7 @@ def build_request( """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}", - ) + living_page = item["living_page"] paths = set(requested_paths) | {living_page} pages = _active_pages( conn, @@ -179,7 +180,7 @@ def build_request( ) evidence, newest_evidence_at = _evidence( conn, - project=project, + living_page=living_page, cursor=cursor, last_page_update=item.get("last_page_update"), max_events=max_events, @@ -192,7 +193,7 @@ def build_request( page_ids = {page["evidence_id"] for page in pages} return { "request": { - "project": project, + "project": item["project"], "living_page": living_page, "current_pages": pages, "evidence": evidence, diff --git a/src/wikibricks_curator/resources/router.md b/src/wikibricks_curator/resources/router.md new file mode 100644 index 0000000..d99b598 --- /dev/null +++ b/src/wikibricks_curator/resources/router.md @@ -0,0 +1,17 @@ +# Route sessions to wiki topics + +Each session in `sessions` ran in a container folder such as the home directory, `code`, +or `emails`. The folder says nothing about the subject. Assign each session to the wiki +page it is about. + +For each session return one route with `session_id`, `page_path`, and a short `reason`: + +1. Use a `candidates` path when the session's main subject is that page's topic. +2. Use a new path `topics/` when the session is about a durable subject that no + candidate covers. The slug names the subject (a customer, product, project, or + concept) in lowercase letters, digits, and hyphens. Never name it after the folder. +3. Use `null` only when the session holds nothing worth remembering: a one-off factual + question, a test of a tool, or small talk. A work product (script, guide, deck, email, + analysis) or a solved problem is durable: route it with rule 1 or 2. + +Return `{"routes": [...]}` with one entry per session. diff --git a/src/wikibricks_curator/resources/router.schema.json b/src/wikibricks_curator/resources/router.schema.json new file mode 100644 index 0000000..ad393e5 --- /dev/null +++ b/src/wikibricks_curator/resources/router.schema.json @@ -0,0 +1,18 @@ +{ + "type": "object", + "properties": { + "routes": { + "type": "array", + "items": { + "type": "object", + "properties": { + "session_id": {"type": "string"}, + "page_path": {"type": ["string", "null"]}, + "reason": {"type": "string"} + }, + "required": ["session_id", "page_path", "reason"] + } + } + }, + "required": ["routes"] +} diff --git a/src/wikibricks_curator/router.py b/src/wikibricks_curator/router.py new file mode 100644 index 0000000..c5b013e --- /dev/null +++ b/src/wikibricks_curator/router.py @@ -0,0 +1,190 @@ +"""Route container-folder sessions to topic pages with one model call.""" + +from __future__ import annotations + +import json +import re +from collections.abc import Callable +from datetime import datetime, timezone +from importlib.resources import files +from typing import Any + +from wikibricks.config import load_config +from wikibricks.curation.backlog import unrouted_sessions +from wikibricks.storage.sqlite_store import SQLiteStore + +Chat = Callable[[str, dict[str, Any], dict[str, Any]], dict[str, Any]] +_NEW_TOPIC = re.compile(r"^topics/[a-z0-9]+(-[a-z0-9]+)*$") +_ROUTER_PROMPT = files("wikibricks_curator").joinpath("resources", "router.md") +_ROUTER_SCHEMA = files("wikibricks_curator").joinpath("resources", "router.schema.json") + + +def _normalize(value: str) -> str: + return value.strip().lower().replace(" ", "-").replace("_", "-") + + +def _sessions(conn: Any, ids: list[str], max_chars: int) -> list[dict[str, str]]: + if not ids: + return [] + placeholders = ",".join("?" for _ in ids) + rows = conn.execute( + "SELECT s.session_id, s.workspace, s.title FROM sessions s " + f"WHERE s.session_id IN ({placeholders})", + ids, + ).fetchall() + event_rows = conn.execute( + "SELECT e.session_id, v.content FROM session_events e " + "JOIN session_event_versions v ON v.version_id = e.current_version_id " + "WHERE e.active = 1 AND v.kind = 'user' " + f"AND e.session_id IN ({placeholders}) ORDER BY e.session_id, e.position", + ids, + ).fetchall() + texts: dict[str, list[str]] = {} + for session_id, content in event_rows: + texts.setdefault(session_id, []).append(content) + return [ + { + "session_id": row["session_id"], + "title": row["title"], + "workspace": row["workspace"], + "text": "\n".join(texts.get(row["session_id"], []))[:max_chars], + } + for row in sorted(rows, key=lambda row: ids.index(row["session_id"])) + ] + + +def _candidates(conn: Any, generic: set[str]) -> list[dict[str, str]]: + # Pages named after a container folder (e.g. an old `topics/emails`) are not topics. + rows = conn.execute( + "SELECT p.path, v.title, v.content FROM pages p " + "JOIN page_versions v ON v.version_id = p.current_version_id " + "WHERE p.status = 'active' " + "AND (p.path LIKE 'topics/%' OR p.path LIKE 'projects/%') ORDER BY p.path" + ).fetchall() + return [ + { + "path": row["path"], + "title": row["title"], + "summary": str(json.loads(row["content"]).get("summary", ""))[:200], + } + for row in rows + if _normalize(row["path"].rsplit("/", 1)[-1]) not in generic | {"home"} + ] + + +def _new_path(path: str, *, candidates: set[str], generic: set[str]) -> bool: + if not isinstance(path, str) or path in candidates: + return False + last_segment = path.rsplit("/", 1)[-1] + return _NEW_TOPIC.fullmatch(path) is not None and _normalize(last_segment) not in generic | {"home"} + + +def _valid_path(path: Any, *, candidates: set[str], generic: set[str]) -> bool: + if not isinstance(path, str): + return False + return path in candidates or _new_path(path, candidates=candidates, generic=generic) + + +def _schema() -> dict[str, Any]: + return json.loads(_ROUTER_SCHEMA.read_text(encoding="utf-8")) + + +def route_sessions( + store: SQLiteStore, + chat: Chat, + *, + since: str | datetime, + limit: int = 20, + max_chars: int = 600, + dry_run: bool = False, +) -> dict[str, Any]: + """Ask the model for topic routes and store the accepted answers.""" + generic = {_normalize(value) for value in load_config().curation_generic_workspaces} + with store.connection() as conn: + ids = unrouted_sessions(conn, since=since)[:limit] + if not ids: + return { + "sessions": 0, + "routed": 0, + "to_none": 0, + "new_topics": [], + "invalid": 0, + "error": None, + } + sessions = _sessions(conn, ids, max_chars) + candidates = _candidates(conn, generic) + candidate_paths = {candidate["path"] for candidate in candidates} + prompt = _ROUTER_PROMPT.read_text(encoding="utf-8") + schema = _schema() + empty = {session["session_id"] for session in sessions if not session["text"]} + model_ids = [session_id for session_id in ids if session_id not in empty] + model_id_set = set(model_ids) + request = { + "sessions": [session for session in sessions if session["session_id"] in model_id_set], + "candidates": candidates, + } + raw: Any = None + error: str | None = None + for _attempt in (1, 2): + try: + raw = chat(prompt, request, schema) if model_ids else {"routes": []} + error = None + break + except Exception as exc: + error = str(exc)[:300] + if error is not None: + return { + "sessions": len(ids), + "routed": 0, + "to_none": 0, + "new_topics": [], + "invalid": 0, + "error": error, + } + + routes = raw.get("routes") if isinstance(raw, dict) else None + invalid_routes = not isinstance(routes, list) + routes = routes if isinstance(routes, list) else [] + accepted: dict[str, str | None] = {} + invalid = 0 + for route in routes: + session_id = route.get("session_id") if isinstance(route, dict) else None + page_path = route.get("page_path") if isinstance(route, dict) else None + if ( + not isinstance(session_id, str) + or session_id not in model_id_set + or session_id in accepted + or not ( + page_path is None + or _valid_path(page_path, candidates=candidate_paths, generic=generic) + ) + ): + invalid += 1 + continue + accepted[session_id] = page_path + for session_id in empty: + accepted[session_id] = None + if invalid_routes: + invalid += len(model_ids) + new_topics = sorted( + {path for path in accepted.values() if path is not None and path not in candidate_paths} + ) + if not dry_run: + timestamp = datetime.now(timezone.utc).isoformat() + with store.connection(write=True) as conn: + conn.executemany( + "INSERT OR IGNORE INTO session_topics" + "(session_id, page_path, origin, created_at) VALUES (?, ?, 'local-curator-router', ?)", + [ + (session_id, page_path, timestamp) + for session_id, page_path in accepted.items() + ], + ) + return { + "sessions": len(ids), + "routed": len(accepted), + "to_none": sum(path is None for path in accepted.values()), + "new_topics": new_topics, + "invalid": invalid, + "error": None, + } diff --git a/tests/test_curation_backlog.py b/tests/test_curation_backlog.py index a1ee986..79142a3 100644 --- a/tests/test_curation_backlog.py +++ b/tests/test_curation_backlog.py @@ -23,17 +23,24 @@ def test_curation_backlog_counts_recent_sessions_without_a_page(): from wikibricks.curation.backlog import curation_backlog sessions = [ - ("/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), - ("/Users/u/work/slide-hub", _timestamp(2, reference=NOW)), - ("/Users/u/work/slide-hub", _timestamp(3, reference=NOW)), + ("one", "/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), + ("two", "/Users/u/work/slide-hub", _timestamp(2, reference=NOW)), + ("three", "/Users/u/work/slide-hub", _timestamp(3, reference=NOW)), ] + targets = { + "one": "topics/slide-hub", + "two": "topics/slide-hub", + "three": "topics/slide-hub", + } - result = curation_backlog(sessions, [], now=NOW) + result = curation_backlog(sessions, [], now=NOW, targets=targets) assert result == [ { "project": "slide-hub", + "living_page": "topics/slide-hub", "workspace": "/Users/u/work/slide-hub", + "workspaces": ["/Users/u/work/slide-hub"], "new_sessions": 3, "last_session_at": _timestamp(1, reference=NOW), "pages": [], @@ -46,27 +53,35 @@ def test_curation_backlog_ignores_sessions_covered_by_a_newer_page(): from wikibricks.curation.backlog import curation_backlog sessions = [ - ("/Users/u/work/agent-compliance-cockpit", _timestamp(3, reference=NOW)), - ("/Users/u/work/agent-compliance-cockpit", _timestamp(4, reference=NOW)), - ("/Users/u/work/agent-compliance-cockpit", _timestamp(5, reference=NOW)), + ("one", "/Users/u/work/agent-compliance-cockpit", _timestamp(3, reference=NOW)), + ("two", "/Users/u/work/agent-compliance-cockpit", _timestamp(4, reference=NOW)), + ("three", "/Users/u/work/agent-compliance-cockpit", _timestamp(5, reference=NOW)), ] pages = [ ("topics/agent-atlas", "Agent Atlas (agent-compliance-cockpit)", _timestamp(2, reference=NOW)) ] - assert curation_backlog(sessions, pages, now=NOW) == [] + assert curation_backlog(sessions, pages, now=NOW, targets={ + "one": "topics/agent-atlas", + "two": "topics/agent-atlas", + "three": "topics/agent-atlas", + }) == [] def test_curation_backlog_counts_only_sessions_newer_than_the_page(): from wikibricks.curation.backlog import curation_backlog sessions = [ - ("/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), - ("/Users/u/work/slide-hub", _timestamp(2, reference=NOW)), - ("/Users/u/work/slide-hub", _timestamp(5, reference=NOW)), + ("one", "/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), + ("two", "/Users/u/work/slide-hub", _timestamp(2, reference=NOW)), + ("three", "/Users/u/work/slide-hub", _timestamp(5, reference=NOW)), ] pages = [("topics/slide-hub", "Slide Hub", _timestamp(3, reference=NOW))] - result = curation_backlog(sessions, pages, now=NOW) + result = curation_backlog(sessions, pages, now=NOW, targets={ + "one": "topics/slide-hub", + "two": "topics/slide-hub", + "three": "topics/slide-hub", + }) assert result[0]["new_sessions"] == 2 assert result[0]["pages"] == ["topics/slide-hub"] @@ -77,21 +92,32 @@ def test_curation_backlog_filters_expired_or_invalid_workspaces_and_orders_resul from wikibricks.curation.backlog import curation_backlog sessions = [ - ("/Users/u/work/slide-hub", _timestamp(24 * 8, reference=NOW)), - ("/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), - (str(Path.home()), _timestamp(1, reference=NOW)), - ("", _timestamp(1, reference=NOW)), - (None, _timestamp(1, reference=NOW)), - ("/Users/u/work/zeta", _timestamp(2, reference=NOW)), - ("/Users/u/work/zeta", _timestamp(3, reference=NOW)), + ("old-slide", "/Users/u/work/slide-hub", _timestamp(24 * 8, reference=NOW)), + ("slide", "/Users/u/work/slide-hub", _timestamp(1, reference=NOW)), + ("home", str(Path.home()), _timestamp(1, reference=NOW)), + ("empty", "", _timestamp(1, reference=NOW)), + ("missing", None, _timestamp(1, reference=NOW)), + ("zeta-one", "/Users/u/work/zeta", _timestamp(2, reference=NOW)), + ("zeta-two", "/Users/u/work/zeta", _timestamp(3, reference=NOW)), ] + targets = { + "old-slide": "topics/slide-hub", + "slide": "topics/slide-hub", + "home": None, + "empty": None, + "missing": None, + "zeta-one": "topics/zeta", + "zeta-two": "topics/zeta", + } assert curation_backlog( - sessions, [], now=NOW, days=7, home=Path.home() + sessions, [], now=NOW, days=7, targets=targets ) == [ { "project": "zeta", + "living_page": "topics/zeta", "workspace": "/Users/u/work/zeta", + "workspaces": ["/Users/u/work/zeta"], "new_sessions": 2, "last_session_at": _timestamp(2, reference=NOW), "pages": [], @@ -99,7 +125,9 @@ def test_curation_backlog_filters_expired_or_invalid_workspaces_and_orders_resul }, { "project": "slide-hub", + "living_page": "topics/slide-hub", "workspace": "/Users/u/work/slide-hub", + "workspaces": ["/Users/u/work/slide-hub"], "new_sessions": 1, "last_session_at": _timestamp(1, reference=NOW), "pages": [], diff --git a/tests/test_local_curator.py b/tests/test_local_curator.py index 7c986c6..b67cbc2 100644 --- a/tests/test_local_curator.py +++ b/tests/test_local_curator.py @@ -94,10 +94,28 @@ def _populate(tmp_path: Path) -> tuple[SQLiteStore, dict]: "beta", "/Users/u/work/beta", NEW_SESSION, - [SessionEvent("user-beta", "user", "beta user")], + [SessionEvent("user-beta", "user", "beta user")], ) + with store.connection(write=True) as conn: + session_ids = { + row["external_id"]: row["session_id"] + for row in conn.execute( + "SELECT external_id, session_id FROM sessions WHERE external_id IN " + "('old-alpha', 'middle-alpha', 'new-alpha')" + ).fetchall() + } + conn.executemany( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/alpha', 'manual', ?)", + [ + (session_ids["old-alpha"], OLD_SESSION), + (session_ids["middle-alpha"], MIDDLE_SESSION), + (session_ids["new-alpha"], NEW_SESSION), + ], + ) item = { "project": "alpha-one", + "living_page": "topics/alpha", "workspace": "/Users/u/work/Alpha One", "new_sessions": 2, "last_session_at": NEW_SESSION, @@ -212,6 +230,7 @@ def test_build_request_defaults_living_page_for_a_project_without_coverage(tmp_p conn, { "project": "beta", + "living_page": "topics/beta", "workspace": "/Users/u/work/beta", "new_sessions": 1, "last_session_at": NEW_SESSION, @@ -431,7 +450,12 @@ def test_build_request_compares_timestamps_across_utc_offsets(tmp_path: Path): [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"} + item = { + "project": "offsets", + "living_page": "topics/offsets", + "pages": [], + "last_page_update": "2026-01-04T11:00:00+00:00", + } with store.connection() as conn: built = build_request(conn, item) diff --git a/tests/test_migration_immutability.py b/tests/test_migration_immutability.py index 73a8613..1b2dbcb 100644 --- a/tests/test_migration_immutability.py +++ b/tests/test_migration_immutability.py @@ -20,11 +20,17 @@ "ba5ad411cff8fe7572900dcd1493efca4442e22818e4ef80a4192f46345cff5c" ), f"{LOCAL_MIGRATIONS}/0006_add_link.sql": "90de2f20cb319896227329a4de24f57cdac50f2c752ca1188781952bf4a5d4d8", + f"{LOCAL_MIGRATIONS}/0007_session_topics.sql": ( + "b24e80325ed2b46b2a472b790da84755da3dc0a2feb4f5feb4d836139fb4e6f6" + ), f"{SQLITE_MIGRATIONS}/0001_core.sql": "622f8257f4a8dd9d5db3419ae30f3484f77566094dd416f33f13974d3ca18e8e", f"{SQLITE_MIGRATIONS}/0002_sync.sql": "70e60dd7f65b11df52feb20e1a411cd88485a1d2a45de6daae74cb6781eb17cc", f"{SQLITE_MIGRATIONS}/0003_backfill_late_tables.sql": ( "cb5b62c81aa4c8d7c8aa22e11c558c13305ef16806a258567ea0c98cf4c55fbf" ), + f"{SQLITE_MIGRATIONS}/0004_session_topics.sql": ( + "8fad04bb0f94813aea4817acb82d10f465cc8b35a481121a1c9e9de912f8df6e" + ), f"{REMOTE_MIGRATIONS}/0001_lakebase_search.sql": "b9b2da1ebcf902ed6a10e48dd7bdb4f3d3a9adcb28752bf8dc2e4a1120b17144", } diff --git a/tests/test_postgres_store.py b/tests/test_postgres_store.py index 354341a..2ff0b68 100644 --- a/tests/test_postgres_store.py +++ b/tests/test_postgres_store.py @@ -43,7 +43,7 @@ def test_migrations_are_repeatable_and_create_required_indexes(postgres_url: str assert {"page_search_chunks_vector_idx", "session_search_chunks_vector_idx"} <= indexes assert {"pages_path_trgm_idx", "page_versions_title_trgm_idx"} <= indexes assert {"pages_active_path_idx", "curation_conflicts_pending_idx"} <= indexes - assert migration_count == 6 + assert migration_count == 7 def test_store_exposes_focused_repositories(store: PostgresStore): diff --git a/tests/test_topic_routing.py b/tests/test_topic_routing.py new file mode 100644 index 0000000..3bf325d --- /dev/null +++ b/tests/test_topic_routing.py @@ -0,0 +1,556 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any + +import pytest + +from wikibricks.config import load_config +from wikibricks.curation.backlog import ( + curation_backlog, + load_curation_backlog, + session_targets, + unrouted_sessions, +) +from wikibricks.models import SessionEvent, SessionRecord +from wikibricks.storage.sqlite_store import SQLiteStore +from wikibricks_curator.curator import run_curator +from wikibricks_curator.evidence import build_request +from wikibricks_curator.router import route_sessions + +NOW = datetime(2026, 9, 29, 12, 0, tzinfo=timezone.utc) +HOME = Path("/Users/u") + + +def _timestamp(offset_hours: int) -> str: + return (NOW - timedelta(hours=offset_hours)).isoformat() + + +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 = _timestamp(1), +) -> str: + store.ingest_session( + SessionRecord( + harness="test-harness", + external_id=external_id, + user_id="user", + workspace=workspace, + updated_at=updated_at, + events=[SessionEvent("0", "user", f"{external_id} content")], + ) + ) + with store.connection(write=True) as conn: + return conn.execute( + "SELECT session_id FROM sessions " + "WHERE harness = 'test-harness' AND external_id = ?", + (external_id,), + ).fetchone()[0] + + +def test_session_targets_route_by_rows_and_pages(tmp_path: Path): + store = _store(tmp_path) + store.write_page( + "topics/agent-atlas", + "Agent Atlas (agent-compliance-cockpit)", + {"summary": "atlas", "body": "page"}, + ) + routed = _ingest(store, "routed", "/Users/u/work/covered") + routed_none = _ingest(store, "routed-none", "/Users/u/work/covered") + generic = _ingest(store, "generic", str(HOME / "code")) + home = _ingest(store, "home", str(HOME)) + covered = _ingest(store, "covered", str(HOME / "agent-compliance-cockpit")) + uncovered = _ingest(store, "uncovered", str(HOME / "uncovered-project")) + + with store.connection(write=True) as conn: + conn.executemany( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, ?, 'manual', ?)", + [ + (routed, "topics/agent-atlas", _timestamp(2)), + (routed_none, None, _timestamp(2)), + ], + ) + targets = session_targets(conn, home=HOME) + + assert targets[routed] == "topics/agent-atlas" + assert targets[routed_none] is None + assert targets[generic] is None + assert targets[home] is None + assert targets[covered] == "topics/agent-atlas" + assert targets[uncovered] == "topics/uncovered-project" + + +def test_session_targets_honours_a_generic_override(tmp_path: Path): + store = _store(tmp_path) + overridden = _ingest(store, "overridden", str(HOME / "custom")) + project = _ingest(store, "project", str(HOME / "custom" / "project")) + + with store.connection(write=True) as conn: + targets = session_targets(conn, generic=("custom",), home=HOME) + + assert targets[overridden] is None + assert targets[project] == "topics/project" + + +def test_unrouted_sessions_returns_only_recent_container_folder_sessions(tmp_path: Path): + store = _store(tmp_path) + generic = _ingest(store, "generic", str(HOME / "code"), _timestamp(1)) + home = _ingest(store, "home", str(HOME), _timestamp(2)) + old = _ingest(store, "old", str(HOME / "code"), _timestamp(10)) + routed = _ingest(store, "routed", str(HOME / "code"), _timestamp(3)) + project = _ingest(store, "project", str(HOME / "project"), _timestamp(1)) + with store.connection(write=True) as conn: + conn.execute( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/agent-atlas', 'manual', ?)", + (routed, _timestamp(4)), + ) + + with store.connection() as conn: + result = unrouted_sessions(conn, since=_timestamp(5), home=HOME) + + assert set(result) == {generic, home} + assert old not in result + assert routed not in result + assert project not in result + + +def test_curation_backlog_groups_by_routed_pages(tmp_path: Path): + store = _store(tmp_path) + store.write_page( + "topics/agent-atlas", + "Agent Atlas (agent-compliance-cockpit)", + {"summary": "atlas", "body": "page"}, + ) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = 'topics/agent-atlas'", + (_timestamp(3),), + ) + project_one = _ingest(store, "one", str(HOME / "agent-compliance-cockpit")) + project_two = _ingest(store, "two", str(HOME / "agent-compliance-cockpit")) + routed = _ingest(store, "routed", str(HOME / "code")) + unrouted = _ingest(store, "unrouted", str(HOME / "code")) + routed_none = _ingest(store, "none", str(HOME / "agent-compliance-cockpit")) + with store.connection(write=True) as conn: + conn.executemany( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, ?, 'manual', ?)", + [ + (routed, "topics/agent-atlas", _timestamp(2)), + (routed_none, None, _timestamp(2)), + ], + ) + targets = session_targets(conn, home=HOME) + + sessions = [ + (project_one, str(HOME / "agent-compliance-cockpit"), _timestamp(1)), + (project_two, str(HOME / "agent-compliance-cockpit"), _timestamp(1)), + (routed, str(HOME / "code"), _timestamp(1)), + (unrouted, str(HOME / "code"), _timestamp(1)), + (routed_none, str(HOME / "agent-compliance-cockpit"), _timestamp(1)), + ] + pages = [ + ( + "topics/agent-atlas", + "Agent Atlas (agent-compliance-cockpit)", + _timestamp(3), + ) + ] + result = curation_backlog(sessions, pages, now=NOW, targets=targets) + + assert targets[unrouted] is None + assert len(result) == 1 + item = result[0] + assert item["project"] == "agent-atlas" + assert item["living_page"] == "topics/agent-atlas" + assert item["workspaces"] == [ + str(HOME / "agent-compliance-cockpit"), + str(HOME / "code"), + ] + assert item["workspace"] == str(HOME / "agent-compliance-cockpit") + assert item["new_sessions"] == 3 + assert item["pages"] == ["topics/agent-atlas"] + + +def test_build_request_selects_routed_sessions(tmp_path: Path): + store = _store(tmp_path) + store.write_page("topics/agent-atlas", "Agent Atlas", {"summary": "a", "body": "b"}) + routed = _ingest(store, "routed", str(HOME / "code")) + unrouted = _ingest(store, "unrouted", str(HOME / "code")) + with store.connection(write=True) as conn: + conn.execute( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/agent-atlas', 'manual', ?)", + (routed, _timestamp(2)), + ) + item = { + "project": "agent-atlas", + "living_page": "topics/agent-atlas", + "new_sessions": 1, + "last_session_at": _timestamp(1), + "pages": ["topics/agent-atlas"], + "last_page_update": None, + } + + with store.connection(write=True) as conn: + result = build_request(conn, item, related=0) + + texts = [entry["text"] for entry in result["request"]["evidence"]] + assert texts == ["routed content"] + assert result["request"]["living_page"] == "topics/agent-atlas" + with store.connection() as conn: + assert session_targets(conn, home=HOME)[unrouted] is None + + +def test_config_defaults_environment_and_validation(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + defaults = load_config(home=tmp_path / "home", environ={}) + assert defaults.curation_generic_workspaces == ( + "code", + "emails", + "work", + "projects", + "repos", + "src", + "documents", + "desktop", + "downloads", + "tmp", + ) + + monkeypatch.setenv("WIKIBRICKS_CURATION_GENERIC_WORKSPACES", "code, Custom Workspace") + overridden = load_config(home=tmp_path / "home") + assert overridden.curation_generic_workspaces == ("code", "Custom Workspace") + + with pytest.raises(ValueError, match="WIKIBRICKS_CURATION_GENERIC_WORKSPACES"): + load_config( + home=tmp_path / "home", + environ={"WIKIBRICKS_CURATION_GENERIC_WORKSPACES": "code,"}, + ) + + +def test_migrations_apply_and_delete_session_topics_with_sessions(tmp_path: Path): + store = _store(tmp_path) + session_id = _ingest(store, "routed", str(HOME / "code")) + with store.connection(write=True) as conn: + conn.execute( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/agent-atlas', 'manual', ?)", + (session_id, _timestamp(2)), + ) + count = conn.execute("SELECT count(*) FROM session_topics").fetchone()[0] + store.migrate() + with store.connection(write=True) as conn: + conn.execute("DROP TABLE session_topics") + conn.execute( + "DELETE FROM schema_migrations WHERE name = '0004_session_topics.sql'" + ) + store.migrate() + with store.connection() as conn: + assert count == 1 + conn.execute("DELETE FROM sessions WHERE session_id = ?", (session_id,)) + remaining = conn.execute("SELECT count(*) FROM session_topics").fetchone()[0] + + assert remaining == 0 + + +def _topic_routes(store: SQLiteStore, external_ids: tuple[str, ...]) -> list[str]: + return [_ingest(store, external_id, str(HOME / "code")) for external_id in external_ids] + + +def test_route_sessions_writes_model_routes_and_builds_backlog(tmp_path: Path): + store = _store(tmp_path) + store.write_page( + "topics/agent-atlas", + "Agent Atlas", + {"summary": "Atlas summary", "body": "Atlas body"}, + ) + store.write_page("projects/unity", "Unity", {"summary": "Unity", "body": "project"}) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = 'topics/agent-atlas'", + (_timestamp(4),), + ) + atlas, unity, quick = _topic_routes(store, ("atlas", "unity", "quick")) + calls: list[tuple[str, dict[str, Any], dict[str, Any]]] = [] + + def chat(prompt: str, request: dict[str, Any], schema: dict[str, Any]) -> dict[str, Any]: + calls.append((prompt, request, schema)) + return { + "routes": [ + {"session_id": atlas, "page_path": "topics/agent-atlas", "reason": "atlas"}, + {"session_id": unity, "page_path": "topics/unity-gateway", "reason": "new"}, + {"session_id": quick, "page_path": None, "reason": "quick"}, + ] + } + + result = route_sessions(store, chat, since=_timestamp(3)) + + assert result == { + "sessions": 3, + "routed": 3, + "to_none": 1, + "new_topics": ["topics/unity-gateway"], + "invalid": 0, + "error": None, + } + assert calls[0][0].startswith("# Route sessions to wiki topics") + assert [item["path"] for item in calls[0][1]["candidates"]] == [ + "projects/unity", + "topics/agent-atlas", + ] + assert sorted(item["session_id"] for item in calls[0][1]["sessions"]) == sorted([atlas, unity, quick]) + assert calls[0][2]["properties"]["routes"]["items"]["properties"]["page_path"]["type"] == [ + "string", + "null", + ] + with store.connection() as conn: + rows = { + row[0]: (row[1], row[2]) + for row in conn.execute( + "SELECT session_id, page_path, origin FROM session_topics" + ) + } + targets = session_targets(conn, home=HOME) + assert rows == { + atlas: ("topics/agent-atlas", "local-curator-router"), + unity: ("topics/unity-gateway", "local-curator-router"), + quick: (None, "local-curator-router"), + } + assert targets == { + atlas: "topics/agent-atlas", + unity: "topics/unity-gateway", + quick: None, + } + with store.connection() as conn: + backlog = load_curation_backlog(conn, home=HOME, now=NOW) + assert {item["living_page"] for item in backlog} == { + "topics/unity-gateway", + "topics/agent-atlas", + } + assert all(item["new_sessions"] == 1 for item in backlog) + + +def test_route_sessions_skips_invalid_routes(tmp_path: Path): + store = _store(tmp_path) + store.write_page("topics/valid", "Valid", {"summary": "s", "body": "b"}) + (folder, email, home, project, bad_slug) = _topic_routes( + store, + ("folder", "email", "home", "project", "bad-slug"), + ) + + def chat(_prompt: str, _request: dict[str, Any], _schema: dict[str, Any]) -> dict[str, Any]: + return { + "routes": [ + {"session_id": "missing", "page_path": "topics/valid", "reason": "unknown"}, + {"session_id": folder, "page_path": "topics/code", "reason": "folder"}, + {"session_id": email, "page_path": "topics/emails", "reason": "folder"}, + {"session_id": home, "page_path": "topics/home", "reason": "folder"}, + {"session_id": project, "page_path": "projects/x", "reason": "not candidate"}, + {"session_id": bad_slug, "page_path": "Topics/Bad Slug", "reason": "bad"}, + ] + } + + with store.connection() as conn: + before = set(unrouted_sessions(conn, since=_timestamp(2), home=HOME)) + result = route_sessions(store, chat, since=_timestamp(2)) + with store.connection() as conn: + after = set(unrouted_sessions(conn, since=_timestamp(2), home=HOME)) + + assert result["sessions"] == 5 + assert result["routed"] == 0 + assert result["to_none"] == 0 + assert result["new_topics"] == [] + assert result["invalid"] == 6 + assert result["error"] is None + assert before == after + + +def test_route_sessions_routes_empty_session_without_chat(tmp_path: Path): + store = _store(tmp_path) + store.ingest_session( + SessionRecord( + harness="test-harness", + external_id="empty", + user_id="user", + workspace=str(HOME / "code"), + updated_at=_timestamp(1), + events=[], + ) + ) + with store.connection(write=True) as conn: + session_id = conn.execute("SELECT session_id FROM sessions").fetchone()[0] + + def chat(*_args: Any) -> dict[str, Any]: + raise AssertionError("chat should not be called") + + result = route_sessions(store, chat, since=_timestamp(2)) + + assert result == { + "sessions": 1, + "routed": 1, + "to_none": 1, + "new_topics": [], + "invalid": 0, + "error": None, + } + with store.connection() as conn: + row = conn.execute( + "SELECT page_path, origin FROM session_topics WHERE session_id = ?", + (session_id,), + ).fetchone() + assert tuple(row) == (None, "local-curator-router") + + +def test_route_sessions_retries_once_and_reports_failure(tmp_path: Path): + store = _store(tmp_path) + _topic_routes(store, ("one", "two")) + attempts = 0 + + def succeeding(_prompt: str, _request: dict[str, Any], _schema: dict[str, Any]) -> dict[str, Any]: + nonlocal attempts + attempts += 1 + assert attempts == 2 + return {"routes": []} + + result = route_sessions(store, succeeding, since=_timestamp(2)) + assert result["error"] is None + + def failing(_prompt: str, _request: dict[str, Any], _schema: dict[str, Any]) -> dict[str, Any]: + raise RuntimeError("model unavailable") + + failure = route_sessions(store, failing, since=_timestamp(2)) + assert failure["error"] == "model unavailable" + with store.connection() as conn: + assert conn.execute( + "SELECT count(*) FROM session_topics WHERE origin = 'local-curator-router'" + ).fetchone()[0] == 0 + + +def test_route_sessions_dry_run_writes_nothing(tmp_path: Path): + store = _store(tmp_path) + (session_id,) = _topic_routes(store, ("dry",)) + + def chat(_prompt: str, _request: dict[str, Any], _schema: dict[str, Any]) -> dict[str, Any]: + return {"routes": [{"session_id": session_id, "page_path": "topics/dry", "reason": "ok"}]} + + result = route_sessions(store, chat, since=_timestamp(2), dry_run=True) + + assert result["routed"] == 1 + assert result["new_topics"] == ["topics/dry"] + with store.connection() as conn: + assert conn.execute("SELECT count(*) FROM session_topics").fetchone()[0] == 0 + + +def test_routed_time_counts_as_new_evidence(tmp_path: Path): + store = _store(tmp_path) + store.write_page("topics/agent-atlas", "Agent Atlas", {"summary": "a", "body": "b"}) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = 'topics/agent-atlas'", + (NOW.isoformat(),), + ) + session_id = _ingest( + store, + "routed-at", + str(HOME / "code"), + (NOW - timedelta(hours=1)).isoformat(), + ) + with store.connection(write=True) as conn: + conn.execute( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/agent-atlas', 'local-curator-router', ?)", + (session_id, (NOW + timedelta(hours=1)).isoformat()), + ) + targets = session_targets(conn, home=HOME) + assert targets[session_id] == "topics/agent-atlas" + with store.connection() as conn: + result = load_curation_backlog( + conn, + now=NOW + timedelta(hours=2), + home=HOME, + ) + assert result[0]["new_sessions"] == 1 + + item = { + "project": "agent-atlas", + "living_page": "topics/agent-atlas", + "new_sessions": 1, + "last_session_at": (NOW + timedelta(hours=1)).isoformat(), + "pages": ["topics/agent-atlas"], + "last_page_update": NOW.isoformat(), + } + with store.connection(write=True) as conn: + built = build_request(conn, item) + assert built is not None + assert built["request"]["evidence"][0]["text"] == "routed-at content" + assert built["newest_evidence_at"] == (NOW + timedelta(hours=1)).isoformat() + + +def test_run_curator_routes_before_backlog_and_hides_routing_errors(tmp_path: Path): + store = _store(tmp_path) + store.write_page("topics/agent-atlas", "Agent Atlas", {"summary": "a", "body": "b"}) + with store.connection(write=True) as conn: + conn.execute( + "UPDATE pages SET updated_at = ? WHERE path = 'topics/agent-atlas'", + (_timestamp(5),), + ) + routed = _ingest(store, "routed", str(HOME / "code"), _timestamp(4)) + _ingest(store, "router", str(HOME / "code")) + with store.connection(write=True) as conn: + conn.execute( + "INSERT INTO session_topics(session_id, page_path, origin, created_at) " + "VALUES (?, 'topics/agent-atlas', 'manual', ?)", + (routed, _timestamp(3)), + ) + calls = 0 + + def chat(prompt: str, request: dict[str, Any], schema: dict[str, Any]) -> dict[str, Any]: + nonlocal calls + calls += 1 + if prompt.startswith("# Route sessions"): + raise RuntimeError("router temporarily unavailable") + return {"proposals": []} + + result = run_curator(tmp_path / "wikibricks.db", chat=chat) + + assert calls == 3 + assert result["errors"] == 0 + assert result["routing"]["error"] == "router temporarily unavailable" + assert result["routing"]["sessions"] == 1 + assert result["routing"]["invalid"] == 0 + assert result["projects"][0]["status"] == "no_changes" + with store.connection() as conn: + assert conn.execute( + "SELECT count(*) FROM session_topics WHERE origin = 'local-curator-router'" + ).fetchone()[0] == 0 + + +def test_pages_named_after_a_container_folder_are_not_candidates(tmp_path: Path): + # The live store already has a `topics/emails` page from before routing existed. + store = _store(tmp_path) + store.write_page("topics/emails", "Emails", {"summary": "mixed", "body": "Unrelated work."}) + store.write_page("topics/agent-atlas", "Agent Atlas", {"summary": "s", "body": "Page."}) + email = _ingest(store, "email-session", "/Users/u/work/50-communications/emails", _timestamp(1)) + seen: dict[str, Any] = {} + + def chat(_prompt: str, request: dict[str, Any], _schema: dict[str, Any]) -> dict[str, Any]: + seen["candidates"] = [candidate["path"] for candidate in request["candidates"]] + return {"routes": [{"session_id": email, "page_path": "topics/emails", "reason": "folder"}]} + + result = route_sessions(store, chat, since=NOW - timedelta(days=7)) + + assert seen["candidates"] == ["topics/agent-atlas"] + assert result["invalid"] == 1 + with store.connection() as conn: + assert email in unrouted_sessions(conn, since=NOW - timedelta(days=7))