diff --git a/pyproject.toml b/pyproject.toml index 3dac39be..68bdf29a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,7 +16,7 @@ dependencies = [ "sentence-transformers>=2.2.0", # MCP server - "mcp>=1.0.0", + "mcp>=1.0.0,<2", # CLI "typer>=0.9.0", diff --git a/scripts/retag_t3_app_provenance.py b/scripts/retag_t3_app_provenance.py new file mode 100644 index 00000000..2e68cceb --- /dev/null +++ b/scripts/retag_t3_app_provenance.py @@ -0,0 +1,239 @@ +#!/usr/bin/env python3 +"""Re-tag existing Codex chunks whose session is explicitly linked by T3.""" + +from __future__ import annotations + +import argparse +import json +import os +import sqlite3 +import sys +import tempfile +from pathlib import Path +from typing import Any + +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) + +from brainlayer.provenance import PROVENANCE_RANK +from brainlayer.t3_provenance import T3_APP_SESSION, codex_session_id_from_source, t3_app_codex_session_ids + + +def _candidates(connection: sqlite3.Connection, linked_session_ids: set[str]) -> list[tuple[str, str | None]]: + rows = connection.execute( + """ + SELECT id, source_file, provenance_class + FROM chunks + WHERE source_file LIKE '%/.codex/sessions/%' + ORDER BY id + """ + ) + return [ + (chunk_id, provenance_class) + for chunk_id, source_file, provenance_class in rows + if codex_session_id_from_source(source_file) in linked_session_ids + and provenance_class == "codex-session" + and provenance_class not in PROVENANCE_RANK + ] + + +def _atomic_write_jsonl(path: Path, rows: list[tuple[str, str | None]]) -> None: + """Durably replace a JSONL artifact without exposing a partial file.""" + path.parent.mkdir(parents=True, exist_ok=True) + temp_fd, temp_name = tempfile.mkstemp(prefix=f".{path.name}.", suffix=".tmp", dir=path.parent) + temp_path = Path(temp_name) + try: + with os.fdopen(temp_fd, "w", encoding="utf-8") as artifact: + for chunk_id, provenance_class in rows: + artifact.write(json.dumps({"id": chunk_id, "provenance_class": provenance_class}) + "\n") + artifact.flush() + os.fsync(artifact.fileno()) + os.replace(temp_path, path) + directory_fd = os.open(path.parent, os.O_RDONLY) + try: + os.fsync(directory_fd) + finally: + os.close(directory_fd) + finally: + if temp_path.exists(): + temp_path.unlink() + + +def _write_rollback_artifact(path: Path, candidates: list[tuple[str, str | None]]) -> None: + existing = ( + { + row["id"]: row["provenance_class"] + for line in path.read_text(encoding="utf-8").splitlines() + if line + for row in [json.loads(line)] + } + if path.exists() + else {} + ) + existing.update(dict(candidates)) + _atomic_write_jsonl(path, sorted(existing.items())) + + +def retag_t3_app_chunks( + *, + db_path: str | Path, + state_db: str | Path, + apply: bool = False, + rollback_artifact: str | Path | None = None, + batch_size: int = 5_000, +) -> dict[str, int]: + """Report or apply the deterministic T3 provenance re-tagging operation.""" + if batch_size <= 0: + raise ValueError("batch_size must be positive") + if apply and rollback_artifact is None: + raise ValueError("--rollback-artifact is required with --apply") + + linked_session_ids = t3_app_codex_session_ids(state_db) + path = Path(db_path).expanduser() + connection = ( + sqlite3.connect(path, timeout=1.0) + if apply + else sqlite3.connect(f"{path.absolute().as_uri()}?mode=ro&immutable=0", uri=True, timeout=1.0) + ) + try: + if apply: + connection.execute("PRAGMA busy_timeout = 30000") + candidates = _candidates(connection, linked_session_ids) + report = { + "linked_sessions": len(linked_session_ids), + "candidate_chunks": len(candidates), + "matched_chunks": 0, + "retagged_chunks": 0, + } + if not apply: + return report + + artifact_path = Path(rollback_artifact).expanduser() + _write_rollback_artifact(artifact_path, candidates) + connection.execute("PRAGMA wal_checkpoint(FULL)") + for batch_start in range(0, len(candidates), batch_size): + batch = candidates[batch_start : batch_start + batch_size] + cursor = connection.executemany( + "UPDATE chunks SET provenance_class = ? WHERE id = ? AND provenance_class = ?", + [(T3_APP_SESSION, chunk_id, "codex-session") for chunk_id, _ in batch], + ) + report["matched_chunks"] += cursor.rowcount + connection.commit() + if ((batch_start // batch_size) + 1) % 3 == 0: + connection.execute("PRAGMA wal_checkpoint(FULL)") + connection.execute("PRAGMA wal_checkpoint(FULL)") + report["retagged_chunks"] = report["matched_chunks"] + return report + finally: + connection.close() + + +def rollback_t3_app_chunks( + *, + db_path: str | Path, + rollback_artifact: str | Path, + only_null_prior_values: bool = False, + pre_restore_artifact: str | Path | None = None, + batch_size: int = 5_000, +) -> dict[str, int]: + """Restore provenance values recorded by a prior T3 re-tag artifact.""" + if batch_size <= 0: + raise ValueError("batch_size must be positive") + + artifact_path = Path(rollback_artifact).expanduser() + rows = [json.loads(line) for line in artifact_path.read_text(encoding="utf-8").splitlines() if line] + candidates = [(row["id"], row["provenance_class"]) for row in rows] + if only_null_prior_values: + candidates = [ + (chunk_id, provenance_class) for chunk_id, provenance_class in candidates if provenance_class is None + ] + if pre_restore_artifact is not None: + pre_restore_path = Path(pre_restore_artifact).expanduser() + if pre_restore_path.exists(): + raise ValueError(f"pre-restore artifact already exists: {pre_restore_path}") + else: + pre_restore_path = None + connection = sqlite3.connect(Path(db_path).expanduser(), timeout=1.0) + try: + connection.execute("PRAGMA busy_timeout = 30000") + current_values = ( + { + chunk_id: provenance_class + for chunk_id, provenance_class in connection.execute( + "SELECT id, provenance_class FROM chunks WHERE id IN ({})".format( + ", ".join("?" for _ in candidates) + ), + [chunk_id for chunk_id, _ in candidates], + ) + } + if candidates + else {} + ) + missing_chunks = sum(chunk_id not in current_values for chunk_id, _ in candidates) + matched_candidates = [ + (chunk_id, provenance_class) + for chunk_id, provenance_class in candidates + if current_values.get(chunk_id) == T3_APP_SESSION + ] + skipped_non_t3_chunks = len(candidates) - missing_chunks - len(matched_candidates) + if pre_restore_path is not None and matched_candidates: + _atomic_write_jsonl( + pre_restore_path, + [(chunk_id, current_values[chunk_id]) for chunk_id, _ in matched_candidates], + ) + restored_chunks = 0 + for batch_start in range(0, len(matched_candidates), batch_size): + batch = matched_candidates[batch_start : batch_start + batch_size] + cursor = connection.executemany( + "UPDATE chunks SET provenance_class = ? WHERE id = ? AND provenance_class = ?", + [(provenance_class, chunk_id, T3_APP_SESSION) for chunk_id, provenance_class in batch], + ) + restored_chunks += cursor.rowcount + connection.commit() + if ((batch_start // batch_size) + 1) % 3 == 0: + connection.execute("PRAGMA wal_checkpoint(FULL)") + connection.execute("PRAGMA wal_checkpoint(FULL)") + return { + "requested_chunks": len(candidates), + "matched_chunks": len(matched_candidates), + "restored_chunks": restored_chunks, + "missing_chunks": missing_chunks, + "skipped_non_t3_chunks": skipped_non_t3_chunks, + } + finally: + connection.close() + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--db-path", type=Path, default=Path.home() / ".local/share/brainlayer/brainlayer.db") + parser.add_argument("--state-db", type=Path, default=Path.home() / ".t3/userdata/state.sqlite") + parser.add_argument("--apply", action="store_true", help="Perform writes; default is read-only dry run") + parser.add_argument("--rollback", action="store_true", help="Restore values from --rollback-artifact") + parser.add_argument("--only-null-prior-values", action="store_true") + parser.add_argument("--pre-restore-artifact", type=Path) + parser.add_argument("--rollback-artifact", type=Path, help="JSONL (id, provenance_class) captured before writes") + parser.add_argument("--batch-size", type=int, default=5_000) + args = parser.parse_args() + if args.rollback: + if args.rollback_artifact is None: + parser.error("--rollback-artifact is required with --rollback") + report: dict[str, Any] = rollback_t3_app_chunks( + db_path=args.db_path, + rollback_artifact=args.rollback_artifact, + only_null_prior_values=args.only_null_prior_values, + pre_restore_artifact=args.pre_restore_artifact, + batch_size=args.batch_size, + ) + else: + report = retag_t3_app_chunks( + db_path=args.db_path, + state_db=args.state_db, + apply=args.apply, + rollback_artifact=args.rollback_artifact, + batch_size=args.batch_size, + ) + print(json.dumps(report, sort_keys=True)) + + +if __name__ == "__main__": + main() diff --git a/src/brainlayer/agent_provenance.py b/src/brainlayer/agent_provenance.py index 72bd7ac0..b7604d35 100644 --- a/src/brainlayer/agent_provenance.py +++ b/src/brainlayer/agent_provenance.py @@ -16,6 +16,7 @@ from .content_class import normalize_content_class from .ingest_denylist import is_denylisted +from .t3_provenance import T3_APP_SESSION, is_t3_app_initiated_codex_session SearchPolicy = Literal["KEEP", "ISOLATE", "OUT"] EffectiveVisibility = Literal["default", "operational", "cold"] @@ -138,11 +139,23 @@ def classify_provenance( content_class: str | None = None, *, content: str | None = None, + t3_state_db: str | Path | None = None, + t3_linked_session_ids: set[str] | None = None, ) -> ProvenanceDecision: """Classify a source path into an auditable provenance search policy.""" del content_class path = _abspath(source_file) + is_t3_app_session = _under_provider_sessions(path, ".codex") and ( + is_t3_app_initiated_codex_session( + source_file, + state_db=t3_state_db, + linked_session_ids=t3_linked_session_ids, + ) + ) + if is_t3_app_session: + return ProvenanceDecision(T3_APP_SESSION, "KEEP", "T3 runtime cursor links Codex session") + if has_recon_agent_signature(content) or _has_recon_path_signature(path): return ProvenanceDecision("recon-agent", "OUT", "recon Agent-tool signature wins precedence") diff --git a/src/brainlayer/t3_provenance.py b/src/brainlayer/t3_provenance.py new file mode 100644 index 00000000..495424f4 --- /dev/null +++ b/src/brainlayer/t3_provenance.py @@ -0,0 +1,97 @@ +"""Read-only discriminator for Codex sessions initiated by the T3 Code app.""" + +from __future__ import annotations + +import json +import os +import re +import sqlite3 +from pathlib import Path + +from .alarm import raise_alarm + +T3_APP_SESSION = "t3-app-session" +DEFAULT_T3_STATE_DB = Path.home() / ".t3" / "userdata" / "state.sqlite" +_CODEX_SESSION_ID_RE = re.compile(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", re.IGNORECASE) +_REQUIRED_RUNTIME_COLUMNS = frozenset({"thread_id", "provider_name", "resume_cursor_json"}) + + +def codex_session_id_from_source(source_file: str | Path) -> str | None: + """Return the final UUID in a Codex JSONL filename, if present.""" + matches = _CODEX_SESSION_ID_RE.findall(Path(source_file).stem) + return matches[-1].lower() if matches else None + + +def is_t3_app_initiated_codex_session( + source_file: str | Path, + *, + state_db: str | Path | None = None, + linked_session_ids: set[str] | None = None, +) -> bool: + """Return whether a Codex transcript is explicitly linked by T3 runtime state. + + A missing T3 database means there is no local T3 app installation to link + against. An existing database with a changed or unreadable schema is fatal: + silently treating those sessions as ordinary Codex would reintroduce the + provenance collision this module prevents. + """ + session_id = codex_session_id_from_source(source_file) + if session_id is None: + return False + + if linked_session_ids is not None: + return session_id in linked_session_ids + + path = Path(state_db or os.environ.get("BRAINLAYER_T3_STATE_DB", DEFAULT_T3_STATE_DB)).expanduser() + if not path.exists(): + return False + return session_id in t3_app_codex_session_ids(path) + + +def t3_app_codex_session_ids(state_db: str | Path = DEFAULT_T3_STATE_DB) -> set[str]: + """Return Codex session IDs explicitly linked by T3 runtime cursors.""" + path = Path(state_db).expanduser() + + try: + connection = sqlite3.connect(f"{path.absolute().as_uri()}?mode=ro&immutable=0", uri=True, timeout=1.0) + except sqlite3.Error as exc: + raise_alarm( + "t3_runtime_unavailable", + "could not open T3 runtime state read-only", + {"path": str(path), "error": str(exc)}, + ) + + try: + columns = {row[1] for row in connection.execute("PRAGMA table_info(provider_session_runtime)")} + missing = sorted(_REQUIRED_RUNTIME_COLUMNS - columns) + if missing: + raise_alarm( + "t3_runtime_schema_drift", + "provider_session_runtime no longer exposes the required T3 linkage columns", + {"path": str(path), "missing_columns": missing}, + ) + + session_ids: set[str] = set() + for provider_name, resume_cursor_json in connection.execute( + "SELECT provider_name, resume_cursor_json FROM provider_session_runtime WHERE provider_name = ?", ("codex",) + ): + try: + resume_cursor = json.loads(resume_cursor_json) + except (TypeError, json.JSONDecodeError) as exc: + raise_alarm( + "t3_runtime_linkage_invalid", + "provider_session_runtime.resume_cursor_json is not valid JSON", + {"path": str(path), "provider_name": provider_name, "error": str(exc)}, + ) + thread_id = resume_cursor.get("threadId") if isinstance(resume_cursor, dict) else None + if isinstance(thread_id, str) and _CODEX_SESSION_ID_RE.fullmatch(thread_id): + session_ids.add(thread_id.lower()) + return session_ids + except sqlite3.Error as exc: + raise_alarm( + "t3_runtime_schema_drift", + "could not query provider_session_runtime for T3 provenance", + {"path": str(path), "error": str(exc)}, + ) + finally: + connection.close() diff --git a/src/brainlayer/watcher_bridge.py b/src/brainlayer/watcher_bridge.py index 982a40da..7eb83fff 100644 --- a/src/brainlayer/watcher_bridge.py +++ b/src/brainlayer/watcher_bridge.py @@ -34,6 +34,7 @@ from .pipeline.correction_detection import build_correction_tags from .pipeline.secret_scrub import scrub_secrets from .queue_io import enqueue_watcher_chunk +from .t3_provenance import DEFAULT_T3_STATE_DB, t3_app_codex_session_ids from .vector_store import VectorStore logger = logging.getLogger(__name__) @@ -288,6 +289,8 @@ def flush_to_db(entries: list[dict[str, Any]]) -> FlushWatermarks: import time as _time flush_start = _time.monotonic() + t3_state_db = Path(os.environ.get("BRAINLAYER_T3_STATE_DB", DEFAULT_T3_STATE_DB)).expanduser() + linked_t3_session_ids = t3_app_codex_session_ids(t3_state_db) if t3_state_db.exists() else set() cursor = None if store is None else store.conn.cursor() inserted = 0 skipped = 0 @@ -419,6 +422,7 @@ def enqueue_chunk( source_file, base_content_class, content=clean_content, + t3_linked_session_ids=linked_t3_session_ids, ) visibility = effective_visibility(provenance_decision, base_content_class) content_class = _content_class_for_visibility(base_content_class, visibility) diff --git a/tests/conftest.py b/tests/conftest.py index 0d0cbaa4..806bc070 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -87,6 +87,12 @@ def isolate_writer_telemetry(monkeypatch, tmp_path): monkeypatch.setenv("BRAINLAYER_WRITER_HEARTBEAT_DIR", str(tmp_path / "writer-heartbeats")) +@pytest.fixture(autouse=True) +def isolate_t3_runtime_state(monkeypatch, tmp_path): + """Prevent unit provenance tests from reading the developer's live T3 app DB.""" + monkeypatch.setenv("BRAINLAYER_T3_STATE_DB", str(tmp_path / "missing-t3-state.sqlite")) + + @pytest.fixture def test_user() -> str: """Username for path-based tests. diff --git a/tests/test_retag_t3_app_provenance.py b/tests/test_retag_t3_app_provenance.py new file mode 100644 index 00000000..62edd77c --- /dev/null +++ b/tests/test_retag_t3_app_provenance.py @@ -0,0 +1,232 @@ +import json +import sqlite3 +from pathlib import Path + +import pytest + +import scripts.retag_t3_app_provenance as retag_module +from scripts.retag_t3_app_provenance import retag_t3_app_chunks, rollback_t3_app_chunks + + +def _create_t3_state_db(path: Path, session_id: str) -> None: + with sqlite3.connect(path) as conn: + conn.execute( + """ + CREATE TABLE provider_session_runtime ( + thread_id TEXT PRIMARY KEY, + provider_name TEXT NOT NULL, + resume_cursor_json TEXT NOT NULL + ) + """ + ) + conn.execute( + "INSERT INTO provider_session_runtime VALUES (?, ?, ?)", + ("t3-thread", "codex", json.dumps({"threadId": session_id})), + ) + + +def _create_brain_db(path: Path, app_session_id: str, plain_session_id: str) -> None: + with sqlite3.connect(path) as conn: + conn.execute("CREATE TABLE chunks (id TEXT PRIMARY KEY, source_file TEXT, provenance_class TEXT)") + conn.executemany( + "INSERT INTO chunks VALUES (?, ?, ?)", + [ + ("app-1", f"/home/etan/.codex/sessions/rollout-{app_session_id}.jsonl", "codex-session"), + ("app-2", f"/home/etan/.codex/sessions/rollout-{app_session_id}.jsonl", "RAW-ETAN-DIRECT"), + ("app-none", f"/home/etan/.codex/sessions/rollout-{app_session_id}.jsonl", None), + ("app-recon", f"/home/etan/.codex/sessions/rollout-{app_session_id}.jsonl", "recon-agent"), + ("plain", f"/home/etan/.codex/sessions/rollout-{plain_session_id}.jsonl", "codex-session"), + ], + ) + + +def test_retag_only_explicitly_t3_linked_sessions_and_write_rollback_artifact(tmp_path: Path) -> None: + app_session_id = "11111111-1111-4111-8111-111111111111" + plain_session_id = "22222222-2222-4222-8222-222222222222" + state_db = tmp_path / "state.sqlite" + brain_db = tmp_path / "brainlayer.db" + rollback_artifact = tmp_path / "rollback.jsonl" + _create_t3_state_db(state_db, app_session_id) + _create_brain_db(brain_db, app_session_id, plain_session_id) + + report = retag_t3_app_chunks( + db_path=brain_db, + state_db=state_db, + apply=True, + rollback_artifact=rollback_artifact, + batch_size=1, + ) + + assert report == {"linked_sessions": 1, "candidate_chunks": 1, "matched_chunks": 1, "retagged_chunks": 1} + assert [json.loads(line) for line in rollback_artifact.read_text().splitlines()] == [ + {"id": "app-1", "provenance_class": "codex-session"}, + ] + with sqlite3.connect(brain_db) as conn: + assert dict(conn.execute("SELECT id, provenance_class FROM chunks")) == { + "app-1": "t3-app-session", + "app-2": "RAW-ETAN-DIRECT", + "app-none": None, + "app-recon": "recon-agent", + "plain": "codex-session", + } + assert rollback_t3_app_chunks(db_path=brain_db, rollback_artifact=rollback_artifact, batch_size=1) == { + "requested_chunks": 1, + "matched_chunks": 1, + "restored_chunks": 1, + "missing_chunks": 0, + "skipped_non_t3_chunks": 0, + } + with sqlite3.connect(brain_db) as conn: + assert dict(conn.execute("SELECT id, provenance_class FROM chunks")) == { + "app-1": "codex-session", + "app-2": "RAW-ETAN-DIRECT", + "app-none": None, + "app-recon": "recon-agent", + "plain": "codex-session", + } + + +def test_rollback_can_restore_only_null_prior_values_and_snapshot_current_rows(tmp_path: Path) -> None: + brain_db = tmp_path / "brainlayer.db" + rollback_artifact = tmp_path / "rollback.jsonl" + pre_restore_artifact = tmp_path / "pre-restore.jsonl" + with sqlite3.connect(brain_db) as conn: + conn.execute("CREATE TABLE chunks (id TEXT PRIMARY KEY, source_file TEXT, provenance_class TEXT)") + conn.executemany( + "INSERT INTO chunks VALUES (?, ?, ?)", + [ + ("prior-null", "/home/etan/.codex/sessions/a.jsonl", "t3-app-session"), + ("prior-codex", "/home/etan/.codex/sessions/b.jsonl", "t3-app-session"), + ], + ) + rollback_artifact.write_text( + "\n".join( + [ + json.dumps({"id": "prior-null", "provenance_class": None}), + json.dumps({"id": "prior-codex", "provenance_class": "codex-session"}), + ] + ) + + "\n" + ) + + assert rollback_t3_app_chunks( + db_path=brain_db, + rollback_artifact=rollback_artifact, + only_null_prior_values=True, + pre_restore_artifact=pre_restore_artifact, + batch_size=1, + ) == { + "requested_chunks": 1, + "matched_chunks": 1, + "restored_chunks": 1, + "missing_chunks": 0, + "skipped_non_t3_chunks": 0, + } + assert [json.loads(line) for line in pre_restore_artifact.read_text().splitlines()] == [ + {"id": "prior-null", "provenance_class": "t3-app-session"} + ] + with sqlite3.connect(brain_db) as conn: + assert dict(conn.execute("SELECT id, provenance_class FROM chunks")) == { + "prior-null": None, + "prior-codex": "t3-app-session", + } + + +def test_retag_does_not_overwrite_a_value_changed_after_the_artifact( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + app_session_id = "11111111-1111-4111-8111-111111111111" + state_db = tmp_path / "state.sqlite" + brain_db = tmp_path / "brainlayer.db" + rollback_artifact = tmp_path / "rollback.jsonl" + _create_t3_state_db(state_db, app_session_id) + _create_brain_db(brain_db, app_session_id, "22222222-2222-4222-8222-222222222222") + original_write = retag_module._write_rollback_artifact + + def change_value_after_snapshot(path: Path, candidates: list[tuple[str, str | None]]) -> None: + original_write(path, candidates) + with sqlite3.connect(brain_db) as conn: + conn.execute("UPDATE chunks SET provenance_class = 'RAW-ETAN-DIRECT' WHERE id = 'app-1'") + + monkeypatch.setattr(retag_module, "_write_rollback_artifact", change_value_after_snapshot) + + report = retag_t3_app_chunks( + db_path=brain_db, + state_db=state_db, + apply=True, + rollback_artifact=rollback_artifact, + ) + + assert report["candidate_chunks"] == 1 + assert report["matched_chunks"] == 0 + assert report["retagged_chunks"] == 0 + with sqlite3.connect(brain_db) as conn: + assert conn.execute("SELECT provenance_class FROM chunks WHERE id = 'app-1'").fetchone()[0] == "RAW-ETAN-DIRECT" + + +def test_rollback_does_not_overwrite_a_value_changed_after_the_snapshot( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + brain_db = tmp_path / "brainlayer.db" + rollback_artifact = tmp_path / "rollback.jsonl" + pre_restore_artifact = tmp_path / "pre-restore.jsonl" + with sqlite3.connect(brain_db) as conn: + conn.execute("CREATE TABLE chunks (id TEXT PRIMARY KEY, source_file TEXT, provenance_class TEXT)") + conn.execute("INSERT INTO chunks VALUES ('app-1', 'source', 't3-app-session')") + rollback_artifact.write_text(json.dumps({"id": "app-1", "provenance_class": "codex-session"}) + "\n") + original_write = retag_module._atomic_write_jsonl + + def change_value_after_snapshot(path: Path, rows: list[tuple[str, str | None]], **kwargs: object) -> None: + original_write(path, rows, **kwargs) + with sqlite3.connect(brain_db) as conn: + conn.execute("UPDATE chunks SET provenance_class = 'RAW-ETAN-DIRECT' WHERE id = 'app-1'") + + monkeypatch.setattr(retag_module, "_atomic_write_jsonl", change_value_after_snapshot) + + report = rollback_t3_app_chunks( + db_path=brain_db, + rollback_artifact=rollback_artifact, + pre_restore_artifact=pre_restore_artifact, + ) + + assert report["matched_chunks"] == 1 + assert report["restored_chunks"] == 0 + with sqlite3.connect(brain_db) as conn: + assert conn.execute("SELECT provenance_class FROM chunks WHERE id = 'app-1'").fetchone()[0] == "RAW-ETAN-DIRECT" + + +def test_rollback_skips_missing_ids_without_creating_an_empty_pre_restore_artifact(tmp_path: Path) -> None: + brain_db = tmp_path / "brainlayer.db" + rollback_artifact = tmp_path / "rollback.jsonl" + pre_restore_artifact = tmp_path / "pre-restore.jsonl" + with sqlite3.connect(brain_db) as conn: + conn.execute("CREATE TABLE chunks (id TEXT PRIMARY KEY, source_file TEXT, provenance_class TEXT)") + rollback_artifact.write_text(json.dumps({"id": "deleted", "provenance_class": "codex-session"}) + "\n") + + report = rollback_t3_app_chunks( + db_path=brain_db, + rollback_artifact=rollback_artifact, + pre_restore_artifact=pre_restore_artifact, + ) + + assert report["missing_chunks"] == 1 + assert report["matched_chunks"] == 0 + assert not pre_restore_artifact.exists() + + +def test_rollback_artifact_replacement_preserves_the_old_file_on_replace_failure( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + artifact = tmp_path / "rollback.jsonl" + original = json.dumps({"id": "old", "provenance_class": "codex-session"}) + "\n" + artifact.write_text(original) + + def fail_replace(source: str, destination: str) -> None: + raise OSError("simulated replace failure") + + monkeypatch.setattr(retag_module.os, "replace", fail_replace) + + with pytest.raises(OSError, match="simulated replace failure"): + retag_module._write_rollback_artifact(artifact, [("new", "codex-session")]) + + assert artifact.read_text() == original diff --git a/tests/test_t3_app_provenance.py b/tests/test_t3_app_provenance.py new file mode 100644 index 00000000..6cf4c109 --- /dev/null +++ b/tests/test_t3_app_provenance.py @@ -0,0 +1,177 @@ +import json +import sqlite3 +from pathlib import Path + +import pytest + +import brainlayer.watcher_bridge as watcher_bridge +from brainlayer.agent_provenance import classify_provenance +from brainlayer.alarm import BrainLayerAlarm +from brainlayer.watcher_bridge import create_flush_callback + + +def _state_db(tmp_path: Path, runtime_rows: list[tuple[str, str, str]]) -> Path: + path = tmp_path / "state.sqlite" + with sqlite3.connect(path) as conn: + conn.execute( + """ + CREATE TABLE provider_session_runtime ( + thread_id TEXT PRIMARY KEY, + provider_name TEXT NOT NULL, + resume_cursor_json TEXT NOT NULL, + runtime_payload_json TEXT + ) + """ + ) + conn.executemany( + """ + INSERT INTO provider_session_runtime + (thread_id, provider_name, resume_cursor_json, runtime_payload_json) + VALUES (?, ?, ?, ?) + """, + [ + (thread_id, provider_name, resume_cursor_json, "{}") + for thread_id, provider_name, resume_cursor_json in runtime_rows + ], + ) + return path + + +def _codex_source(tmp_path: Path, session_id: str) -> Path: + return tmp_path / "home" / ".codex" / "sessions" / "2026" / "08" / f"rollout-2026-08-01-{session_id}.jsonl" + + +def test_plain_codex_working_on_t3layer_is_not_tagged_as_t3_app(tmp_path: Path) -> None: + state_db = _state_db( + tmp_path, + [("another-t3-thread", "codex", json.dumps({"threadId": "aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa"}))], + ) + + decision = classify_provenance( + str(_codex_source(tmp_path, "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb")), + t3_state_db=state_db, + ) + + assert (decision.provenance_tag, decision.search_policy) == ("codex-session", "KEEP") + + +def test_t3_app_initiated_codex_session_is_tagged_distinctly(tmp_path: Path) -> None: + session_id = "019fb03f-0650-7f53-850b-921246951edc" + state_db = _state_db( + tmp_path, + [("5ca576df-592c-406f-a8f1-7db9c56d36c9", "codex", json.dumps({"threadId": session_id}))], + ) + + decision = classify_provenance(str(_codex_source(tmp_path, session_id)), t3_state_db=state_db) + + assert (decision.provenance_tag, decision.search_policy) == ("t3-app-session", "KEEP") + + +def test_t3_app_linkage_outranks_recon_content_signature(tmp_path: Path) -> None: + session_id = "019fb03f-0650-7f53-850b-921246951edc" + state_db = _state_db( + tmp_path, + [("5ca576df-592c-406f-a8f1-7db9c56d36c9", "codex", json.dumps({"threadId": session_id}))], + ) + + decision = classify_provenance( + str(_codex_source(tmp_path, session_id)), + content="Task for brain-worker: mine the transcript.", + t3_state_db=state_db, + ) + + assert decision.provenance_tag == "t3-app-session" + + +def test_new_t3_app_codex_ingestion_gets_distinct_provenance(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + session_id = "cccccccc-cccc-4ccc-8ccc-cccccccccccc" + state_db = _state_db( + tmp_path, + [("5ca576df-592c-406f-a8f1-7db9c56d36c9", "codex", json.dumps({"threadId": session_id}))], + ) + monkeypatch.setenv("BRAINLAYER_T3_STATE_DB", str(state_db)) + source_file = _codex_source(tmp_path, session_id) + flush = create_flush_callback(db_path=tmp_path / "brainlayer.db", arbitrated=False) + + result = flush( + [ + { + "type": "user", + "message": { + "content": [{"type": "text", "text": "T3 app sessions must retain their explicit origin."}] + }, + "timestamp": "2026-08-01T12:00:00Z", + "_source_file": str(source_file), + "_line_end_offset": 100, + } + ] + ) + + assert result.inserted == 1 + with sqlite3.connect(tmp_path / "brainlayer.db") as conn: + assert conn.execute("SELECT provenance_class FROM chunks").fetchone()[0] == "t3-app-session" + + +def test_watcher_loads_t3_session_ids_once_per_flush(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + session_id = "cccccccc-cccc-4ccc-8ccc-cccccccccccc" + state_db = _state_db( + tmp_path, + [("5ca576df-592c-406f-a8f1-7db9c56d36c9", "codex", json.dumps({"threadId": session_id}))], + ) + calls: list[Path] = [] + + def load_once(path: Path) -> set[str]: + calls.append(path) + return {session_id} + + monkeypatch.setenv("BRAINLAYER_T3_STATE_DB", str(state_db)) + monkeypatch.setattr(watcher_bridge, "t3_app_codex_session_ids", load_once) + source_file = _codex_source(tmp_path, session_id) + flush = create_flush_callback(db_path=tmp_path / "brainlayer.db", arbitrated=False) + + result = flush( + [ + { + "type": "user", + "message": {"content": [{"type": "text", "text": "First T3 watcher message."}]}, + "timestamp": "2026-08-01T12:00:00Z", + "_source_file": str(source_file), + "_line_end_offset": 100, + }, + { + "type": "user", + "message": {"content": [{"type": "text", "text": "Second T3 watcher message."}]}, + "timestamp": "2026-08-01T12:00:01Z", + "_source_file": str(source_file), + "_line_end_offset": 200, + }, + ] + ) + + assert result.inserted == 2 + assert calls == [state_db] + + +def test_canonical_systems_codex_session_uses_runtime_cursor_linkage(tmp_path: Path) -> None: + session_id = "019fb03f-0650-7f53-850b-921246951edc" + state_db = _state_db( + tmp_path, + [("5ca576df-592c-406f-a8f1-7db9c56d36c9", "codex", json.dumps({"threadId": session_id}))], + ) + + decision = classify_provenance(str(_codex_source(tmp_path, session_id)), t3_state_db=state_db) + + assert decision.provenance_tag == "t3-app-session" + assert "runtime cursor" in decision.reason + + +def test_missing_runtime_schema_raises_loud_alarm(tmp_path: Path) -> None: + state_db = tmp_path / "state.sqlite" + with sqlite3.connect(state_db) as conn: + conn.execute("CREATE TABLE projection_thread_sessions (provider_session_id TEXT)") + + with pytest.raises(BrainLayerAlarm, match="t3_runtime_schema_drift"): + classify_provenance( + str(_codex_source(tmp_path, "bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb")), + t3_state_db=state_db, + )