diff --git a/src/aipass/memory/apps/handlers/monitor/memory_watcher.py b/src/aipass/memory/apps/handlers/monitor/memory_watcher.py index 18905626..2ee993cf 100644 --- a/src/aipass/memory/apps/handlers/monitor/memory_watcher.py +++ b/src/aipass/memory/apps/handlers/monitor/memory_watcher.py @@ -43,11 +43,11 @@ except ImportError: FileSystemEventHandler = object # type: ignore[assignment,misc] logger.info("Optional dependency 'watchdog' not available") -# Handler imports (relative within package) -from aipass.memory.apps.handlers.tracking.line_counter import update_line_count -from aipass.memory.apps.handlers.monitor.detector import check_single_file -from aipass.prax.apps.modules.logger import get_system_logger -from aipass.memory.apps.handlers.json import json_handler +# Handler imports (relative within package — after conditional watchdog block) +from aipass.memory.apps.handlers.tracking.line_counter import update_line_count # noqa: E402 +from aipass.memory.apps.handlers.monitor.detector import check_single_file # noqa: E402 +from aipass.prax.apps.modules.logger import get_system_logger # noqa: E402 +from aipass.memory.apps.handlers.json import json_handler # noqa: E402 logger = get_system_logger() @@ -200,8 +200,6 @@ def check_and_rollover() -> Dict[str, Any]: return results # Check each branch for memory files over limit - # Also sync current_lines metadata to keep it accurate - lines_synced = 0 for branch_path in branch_paths: branch = Path(branch_path) # Find memory files in .trinity/ subdirectory @@ -213,33 +211,23 @@ def check_and_rollover() -> Dict[str, Any]: results["files_checked"] += 1 try: - line_count = len(memory_file.read_text(encoding="utf-8").splitlines()) + # Auto-heal: reconcile file against template (strips orphan keys) + from aipass.memory.apps.handlers.schema.normalize import normalize_memory_file - # Sync current_lines metadata if stale - try: - import json as _json - - _data = _json.loads(memory_file.read_text(encoding="utf-8")) - meta_lines = _data.get("document_metadata", {}).get("status", {}).get("current_lines") - if meta_lines != line_count: - sync_result = update_line_count(memory_file) - if sync_result.get("success"): - lines_synced += 1 - except Exception as e: - logger.warning(f"[memory_watcher] Non-critical metadata sync failed for {memory_file}: {e}") + normalize_memory_file(memory_file) # Use detector for trigger decision (handles both v1 line-based and v2 entry-count) from aipass.memory.apps.handlers.monitor.detector import _should_rollover - triggered, _, _, _, _ = _should_rollover(memory_file) + triggered, current_lines, _, _, _ = _should_rollover(memory_file) if triggered: results["files_over_limit"].append( - {"file": str(memory_file), "lines": line_count, "threshold": 0} + {"file": str(memory_file), "lines": current_lines, "threshold": 0} ) except Exception as e: logger.warning(f"[memory_watcher] Failed to read memory file {memory_file}: {e}") - results["lines_synced"] = lines_synced + results["lines_synced"] = 0 # Trigger rollover if any files are over limit if results["files_over_limit"]: @@ -537,19 +525,23 @@ class MemoryFileWatcher(FileSystemEventHandler): # type: ignore[misc] logger.info(f"[memory_watcher] Detected modification: {file_path.name}") - # Step 1: Update line count metadata + # Step 1: Auto-heal schema drift (strips orphan keys) + from aipass.memory.apps.handlers.schema.normalize import normalize_memory_file + + norm_result = normalize_memory_file(file_path) + if norm_result.get("changes"): + self._recent_modifications.add(file_key) + + # Step 2: Update health check metadata update_result = update_line_count(file_path) if not update_result["success"]: logger.error( - f"[memory_watcher] Failed to update line count for {file_path.name}: {update_result.get('error')}" + f"[memory_watcher] Failed to update metadata for {file_path.name}: {update_result.get('error')}" ) return - current_lines = update_result.get("lines", 0) - logger.info(f"[memory_watcher] Updated {file_path.name}: {current_lines} lines") - - # Step 2: Check if rollover needed + # Step 3: Check if rollover needed check_result = check_single_file(file_path) if not check_result["success"]: diff --git a/src/aipass/memory/apps/handlers/schema/normalize.py b/src/aipass/memory/apps/handlers/schema/normalize.py index 897c0d9d..bd39f6dc 100644 --- a/src/aipass/memory/apps/handlers/schema/normalize.py +++ b/src/aipass/memory/apps/handlers/schema/normalize.py @@ -1,24 +1,17 @@ # =================== AIPass ==================== # Name: normalize.py # Description: Memory File Schema Normalizer -# Version: 0.2.0 +# Version: 0.3.0 # Created: 2026-01-22 -# Modified: 2026-03-06 +# Modified: 2026-06-08 # ============================================= """ Memory File Schema Normalizer -Fixes inconsistent schema in memory JSON files: -1. Moves root-level 'limits' into document_metadata.limits -2. Removes redundant root-level 'status' -3. Removes auto_compress_at (redundant with max_lines) -4. Ensures document_metadata.status has current_lines - -Supports two schema versions: - v1 (schema_version <2.0.0): { "limits": { "max_lines": N } } - v2 (schema_version >=2.0.0): { "limits": { "max_sessions": N, "max_key_learnings": N, - "session_summary_max_chars": N, "learning_value_max_chars": N } } +Reconciles memory files against their canonical template schema. +Strips any key not present in the template at every level (root, +document_metadata, limits, status). Template = the whole truth. """ import json @@ -31,6 +24,36 @@ from aipass.memory.apps.handlers.json import json_handler logger = get_system_logger() +_MEMORY_ROOT = Path(__file__).parents[3] # normalize.py -> schema/ -> handlers/ -> apps/ -> memory/ + + +def _load_template(file_path: Path) -> Dict[str, Any] | None: + """Load the matching template for a memory file (local or observations).""" + templates_dir = _MEMORY_ROOT / "templates" + name = file_path.name.lower() + + if "local" in name: + tmpl_path = templates_dir / "LOCAL.template.json" + elif "observation" in name: + tmpl_path = templates_dir / "OBSERVATIONS.template.json" + else: + return None + + try: + with open(tmpl_path, "r", encoding="utf-8") as f: + return json.load(f) + except Exception as e: + logger.warning(f"[normalize] Failed to load template {tmpl_path}: {e}") + return None + + +def _strip_orphan_keys(data: Dict, allowed: set, level_name: str, changes: list) -> None: + """Remove keys from data that aren't in the allowed set.""" + orphans = set(data.keys()) - allowed + for key in orphans: + del data[key] + changes.append(f"Stripped orphan '{key}' from {level_name}") + def _find_repo_root() -> Path: """Walk up from this file to find repo root (contains AIPASS_REGISTRY.json).""" @@ -43,7 +66,7 @@ def _find_repo_root() -> Path: def normalize_memory_file(file_path: Path, dry_run: bool = False) -> Dict[str, Any]: """ - Normalize schema for a single memory file. + Normalize a memory file against its canonical template. Args: file_path: Path to memory JSON file @@ -71,59 +94,53 @@ def normalize_memory_file(file_path: Path, dry_run: bool = False) -> Dict[str, A metadata = data["document_metadata"] - # 1. Move root 'limits' into document_metadata.limits + # Legacy fix: move root 'limits' into document_metadata.limits if "limits" in data and "limits" not in metadata: metadata["limits"] = data.pop("limits") changes.append("Moved root 'limits' into document_metadata") elif "limits" in data and "limits" in metadata: - # Both exist - merge, preferring document_metadata values root_limits = data.pop("limits") for key, val in root_limits.items(): if key not in metadata["limits"]: metadata["limits"][key] = val changes.append("Merged root 'limits' into document_metadata.limits") - # 2. Remove root 'status' (redundant) + # Legacy fix: move root 'status' into document_metadata.status if "status" in data: - root_status = data.pop("status") - # If document_metadata.status doesn't have current_lines, copy it + data.pop("status") if "status" not in metadata: metadata["status"] = {} - if "current_lines" not in metadata["status"] and "current_lines" in root_status: - metadata["status"]["current_lines"] = root_status["current_lines"] changes.append("Removed redundant root 'status'") - # 3. Remove auto_compress_at from document_metadata.status (redundant with max_lines) - if "status" in metadata and "auto_compress_at" in metadata["status"]: - del metadata["status"]["auto_compress_at"] - changes.append("Removed redundant 'auto_compress_at'") - - # 4. Remove unused limits fields (max_word_count, max_token_count - no code uses these) - # Preserve v2 fields: max_sessions, max_key_learnings, session_summary_max_chars, learning_value_max_chars, - # max_observations, max_lines, note - if "limits" in metadata: - for unused_field in ["max_word_count", "max_token_count"]: - if unused_field in metadata["limits"]: - del metadata["limits"][unused_field] - changes.append(f"Removed unused '{unused_field}'") - - # 4. Ensure status has required fields + # Ensure status has required fields if "status" not in metadata: metadata["status"] = {} - if "current_lines" not in metadata["status"]: - # Count actual lines - try: - with open(file_path, "r", encoding="utf-8") as f: - metadata["status"]["current_lines"] = len(f.readlines()) - changes.append("Added current_lines count") - except Exception as e: - logger.warning(f"[normalize] Failed to count lines in {file_path}: {e}") - if "last_health_check" not in metadata["status"]: metadata["status"]["last_health_check"] = datetime.now().strftime("%Y-%m-%d") changes.append("Added last_health_check") + # Template-conformance: strip orphan keys at every level + template = _load_template(file_path) + if template is not None: + tmpl_meta = template.get("document_metadata", {}) + + # Root level + _strip_orphan_keys(data, set(template.keys()), "root", changes) + + # document_metadata level + _strip_orphan_keys(metadata, set(tmpl_meta.keys()), "document_metadata", changes) + + # limits level + tmpl_limits = tmpl_meta.get("limits", {}) + if "limits" in metadata: + _strip_orphan_keys(metadata["limits"], set(tmpl_limits.keys()), "limits", changes) + + # status level + tmpl_status = tmpl_meta.get("status", {}) + if "status" in metadata: + _strip_orphan_keys(metadata["status"], set(tmpl_status.keys()), "status", changes) + # Write if changes made and not dry run if changes and not dry_run: try: diff --git a/src/aipass/memory/apps/handlers/tracking/line_counter.py b/src/aipass/memory/apps/handlers/tracking/line_counter.py index 1ae44339..a6464126 100644 --- a/src/aipass/memory/apps/handlers/tracking/line_counter.py +++ b/src/aipass/memory/apps/handlers/tracking/line_counter.py @@ -62,24 +62,20 @@ def _count_physical_lines(file_path: Path) -> int: def update_line_count(file_path: Path) -> Dict[str, Any]: """ - Update current_lines in document_metadata.status - - Reads file, counts lines, updates metadata field using safe json_handler. + Update health check metadata after file modification. Args: file_path: Path to memory JSON file Returns: - Dict with success status and updated line count + Dict with success status """ if not file_path.exists(): return {"success": False, "error": f"File not found: {file_path}"} - # Count lines line_count = _count_physical_lines(file_path) - # Update metadata using safe handler (atomic write) - result = update_metadata(file_path, current_lines=line_count, last_health_check=datetime.now().strftime("%Y-%m-%d")) + result = update_metadata(file_path, last_health_check=datetime.now().strftime("%Y-%m-%d")) if not result["success"]: return {"success": False, "error": f"Failed to update metadata: {result['error']}"} diff --git a/src/aipass/memory/templates/LOCAL.template.json b/src/aipass/memory/templates/LOCAL.template.json index c050ba85..3ffb46d2 100644 --- a/src/aipass/memory/templates/LOCAL.template.json +++ b/src/aipass/memory/templates/LOCAL.template.json @@ -23,8 +23,7 @@ }, "status": { "health": "healthy", - "last_health_check": "{{DATE}}", - "current_lines": 0 + "last_health_check": "{{DATE}}" } }, "key_learnings": {}, diff --git a/src/aipass/memory/templates/OBSERVATIONS.template.json b/src/aipass/memory/templates/OBSERVATIONS.template.json index 15d36d99..31a7afb1 100644 --- a/src/aipass/memory/templates/OBSERVATIONS.template.json +++ b/src/aipass/memory/templates/OBSERVATIONS.template.json @@ -13,12 +13,11 @@ "{{BRANCHNAME}}" ], "limits": { - "max_lines": 600, - "note": "DO NOT trim, prune, or delete entries. Auto-rollover to @memory when max_lines exceeded." + "max_observations": 25, + "note": "DO NOT trim, prune, or delete entries. Auto-rollover to @memory when max_observations exceeded." }, "status": { "health": "healthy", - "current_lines": 0, "last_health_check": "{{DATE}}" } }, diff --git a/src/aipass/memory/tests/test_handlers.py b/src/aipass/memory/tests/test_handlers.py index 7f8cd6da..ba7e0750 100644 --- a/src/aipass/memory/tests/test_handlers.py +++ b/src/aipass/memory/tests/test_handlers.py @@ -350,7 +350,7 @@ class TestCountPhysicalLines: def test_counts_lines_correctly(self, monkeypatch, tmp_path): lc, _ = _import_line_counter(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" f.write_text("line1\nline2\nline3\n", encoding="utf-8") assert lc._count_physical_lines(f) == 3 @@ -376,7 +376,7 @@ class TestUpdateLineCount: def test_updates_line_count_successfully(self, monkeypatch, tmp_path): lc, mocks = _import_line_counter(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" f.write_text('{\n "a": 1\n}\n', encoding="utf-8") result = lc.update_line_count(f) @@ -387,7 +387,7 @@ class TestUpdateLineCount: def test_reports_failure_when_metadata_update_fails(self, monkeypatch, tmp_path): lc, mocks = _import_line_counter(monkeypatch) mocks["memory_files"].update_metadata.return_value = {"success": False, "error": "write error"} - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" f.write_text("{}\n", encoding="utf-8") result = lc.update_line_count(f) @@ -413,12 +413,12 @@ class TestNormalizeMemoryFile: def test_moves_root_limits_into_metadata(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { - "document_metadata": {"status": {"current_lines": 10}}, - "limits": {"max_lines": 600}, + "document_metadata": {"status": {}}, + "limits": {"max_sessions": 20}, "sessions": [], }, ) @@ -427,19 +427,19 @@ class TestNormalizeMemoryFile: data = json.loads(f.read_text(encoding="utf-8")) assert "limits" not in {k for k in data if k != "document_metadata"} - assert data["document_metadata"]["limits"]["max_lines"] == 600 + assert data["document_metadata"]["limits"]["max_sessions"] == 20 def test_merges_root_limits_preserving_metadata_values(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { "document_metadata": { - "limits": {"max_lines": 500}, - "status": {"current_lines": 10}, + "limits": {"max_sessions": 20}, + "status": {}, }, - "limits": {"max_lines": 600, "extra_field": 42}, + "limits": {"max_sessions": 30, "max_key_learnings": 25}, "sessions": [], }, ) @@ -447,19 +447,19 @@ class TestNormalizeMemoryFile: assert result["success"] is True data = json.loads(f.read_text(encoding="utf-8")) - # metadata value (500) wins over root value (600) - assert data["document_metadata"]["limits"]["max_lines"] == 500 - # extra_field from root gets merged in - assert data["document_metadata"]["limits"]["extra_field"] == 42 + # metadata value (20) wins over root value (30) + assert data["document_metadata"]["limits"]["max_sessions"] == 20 + # valid key from root gets merged in + assert data["document_metadata"]["limits"]["max_key_learnings"] == 25 def test_removes_root_status(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { - "document_metadata": {"status": {"current_lines": 10}}, - "status": {"health": "ok", "current_lines": 5}, + "document_metadata": {"status": {"last_health_check": "2026-01-01"}}, + "status": {"health": "ok"}, "sessions": [], }, ) @@ -472,12 +472,12 @@ class TestNormalizeMemoryFile: def test_removes_auto_compress_at(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { "document_metadata": { - "status": {"current_lines": 10, "auto_compress_at": 500}, + "status": {"auto_compress_at": 500, "last_health_check": "2026-01-01"}, }, "sessions": [], }, @@ -490,7 +490,7 @@ class TestNormalizeMemoryFile: def test_dry_run_does_not_write(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" original = { "document_metadata": {}, "limits": {"max_lines": 600}, @@ -508,13 +508,13 @@ class TestNormalizeMemoryFile: def test_no_changes_when_already_normalized(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { "document_metadata": { "limits": {"max_sessions": 20}, - "status": {"current_lines": 10, "last_health_check": "2026-03-31"}, + "status": {"last_health_check": "2026-03-31"}, }, "sessions": [], }, @@ -525,13 +525,13 @@ class TestNormalizeMemoryFile: def test_removes_unused_limit_fields(self, monkeypatch, tmp_path): norm, _ = _import_normalize(monkeypatch) - f = tmp_path / "test.json" + f = tmp_path / "test.local.json" self._write_json( f, { "document_metadata": { - "limits": {"max_lines": 600, "max_word_count": 9999, "max_token_count": 5000}, - "status": {"current_lines": 10, "last_health_check": "2026-03-31"}, + "limits": {"max_sessions": 20, "max_lines": 600, "max_word_count": 9999, "max_token_count": 5000}, + "status": {"last_health_check": "2026-03-31"}, }, "sessions": [], }, @@ -542,7 +542,8 @@ class TestNormalizeMemoryFile: data = json.loads(f.read_text(encoding="utf-8")) assert "max_word_count" not in data["document_metadata"]["limits"] assert "max_token_count" not in data["document_metadata"]["limits"] - assert data["document_metadata"]["limits"]["max_lines"] == 600 + assert "max_lines" not in data["document_metadata"]["limits"] + assert data["document_metadata"]["limits"]["max_sessions"] == 20 class TestTodosOperational: diff --git a/src/aipass/memory/tests/test_watcher.py b/src/aipass/memory/tests/test_watcher.py index 1fe7bd4b..35096477 100644 --- a/src/aipass/memory/tests/test_watcher.py +++ b/src/aipass/memory/tests/test_watcher.py @@ -71,6 +71,21 @@ def _prepare_watcher_mocks(monkeypatch): monkeypatch.setitem(sys.modules, "watchdog.observers", mock_watchdog_observers) monkeypatch.setitem(sys.modules, "watchdog.events", mock_watchdog_events) + # Mock normalize (lazy import inside on_modified and check_and_rollover) + mock_normalize_memory_file = MagicMock(return_value={"success": True, "changes": []}) + mock_normalize = MagicMock() + mock_normalize.normalize_memory_file = mock_normalize_memory_file + monkeypatch.setitem( + sys.modules, + "aipass.memory.apps.handlers.schema", + MagicMock(), + ) + monkeypatch.setitem( + sys.modules, + "aipass.memory.apps.handlers.schema.normalize", + mock_normalize, + ) + # Mock rollover orchestrator (lazy import inside on_modified) mock_execute_rollover = MagicMock(return_value={"success": True}) mock_orchestrator = MagicMock() @@ -92,6 +107,7 @@ def _prepare_watcher_mocks(monkeypatch): "observer_instance": mock_observer_instance, "observer_cls": mock_observer_cls, "execute_rollover": mock_execute_rollover, + "normalize_memory_file": mock_normalize_memory_file, } @@ -495,3 +511,61 @@ class TestMemoryFileWatcherOnModified: watcher.on_modified(event) mocks["execute_rollover"].assert_called_once() + + def test_normalize_called_on_modification(self, monkeypatch): + """on_modified calls normalize_memory_file before update_line_count.""" + mod, mocks = _import_watcher(monkeypatch) + + watcher = mod.MemoryFileWatcher() + event = MagicMock() + event.is_directory = False + event.src_path = "/some/branch/.trinity/local.json" + + watcher.on_modified(event) + + mocks["normalize_memory_file"].assert_called_once() + mocks["update_line_count"].assert_called_once() + + def test_normalize_changes_guard_write_loop(self, monkeypatch): + """When normalize makes changes, file_key is added to _recent_modifications to prevent write-loop.""" + mod, mocks = _import_watcher(monkeypatch) + + mocks["normalize_memory_file"].return_value = { + "success": True, + "changes": ["Stripped orphan 'stale_key' from root"], + } + + watcher = mod.MemoryFileWatcher() + event = MagicMock() + event.is_directory = False + event.src_path = "/some/branch/.trinity/local.json" + + watcher.on_modified(event) + + from pathlib import Path as _P + + file_key = str(_P("/some/branch/.trinity/local.json")) + assert file_key in watcher._recent_modifications + + def test_normalize_no_changes_no_guard(self, monkeypatch): + """When normalize makes no changes, _recent_modifications is not populated by normalize step.""" + mod, mocks = _import_watcher(monkeypatch) + + mocks["normalize_memory_file"].return_value = { + "success": True, + "changes": [], + } + + watcher = mod.MemoryFileWatcher() + event = MagicMock() + event.is_directory = False + event.src_path = "/some/branch/.trinity/local.json" + + watcher.on_modified(event) + + from pathlib import Path as _P + + file_key = str(_P("/some/branch/.trinity/local.json")) + # file_key should NOT be in _recent_modifications from normalize (may be from rollover) + mocks["check_single_file"].return_value = {"success": True, "should_rollover": False} + assert file_key not in watcher._recent_modifications