diff --git a/src/aipass/flow/apps/handlers/mbank/process.py b/src/aipass/flow/apps/handlers/mbank/process.py index c512a8b1..7fd427e6 100644 --- a/src/aipass/flow/apps/handlers/mbank/process.py +++ b/src/aipass/flow/apps/handlers/mbank/process.py @@ -19,13 +19,16 @@ Key Functions: - verify_and_heal_orphaned_plans() - Orphan healing logic """ +# ruff: noqa: E402 from pathlib import Path _PKG_ROOT = Path(__file__).resolve().parents[4] # Standard imports import json +import os import shutil +import time from datetime import datetime, timezone from typing import Dict, List, Any @@ -42,6 +45,35 @@ from aipass.prax.apps.modules.logger import system_logger as logger FLOW_ROOT = _PKG_ROOT / "flow" FLOW_JSON_DIR = FLOW_ROOT / "flow_json" +MODULE_NAME = "mbank_process" +_LOCK_RETRIES = 10 +_LOCK_BACKOFF_BASE = 0.05 + + +def _acquire_lock(lock_path: Path) -> bool: + """Atomically acquire a lockfile via O_CREAT|O_EXCL with retry+backoff.""" + for attempt in range(_LOCK_RETRIES): + try: + fd = os.open(str(lock_path), os.O_CREAT | os.O_EXCL | os.O_WRONLY) + os.write(fd, str(os.getpid()).encode()) + os.close(fd) + return True + except FileExistsError: + logger.info("[%s] Lock contention on %s, retry %d", MODULE_NAME, lock_path, attempt + 1) + time.sleep(_LOCK_BACKOFF_BASE * (2**attempt)) + except OSError as exc: + logger.warning("[%s] Lock creation failed for %s: %s", MODULE_NAME, lock_path, exc) + return False + return False + + +def _release_lock(lock_path: Path) -> None: + """Remove lockfile, tolerating already-removed.""" + try: + lock_path.unlink(missing_ok=True) + except OSError as exc: + logger.warning("[%s] Could not release lock %s: %s", MODULE_NAME, lock_path, exc) + def _find_repo_root() -> Path: """Walk up from this file to find the repo root (contains AIPASS_REGISTRY.json).""" @@ -97,12 +129,22 @@ def load_flow_registry(registry_file: str | None = None) -> Dict[str, Any]: def save_flow_registry(registry: Dict[str, Any], registry_file: str | None = None) -> None: - """Save a plan registry.""" + """Save a plan registry with lockfile + atomic write.""" target = FLOW_JSON_DIR / registry_file if registry_file else REGISTRY_FILE + lock_path = target.with_suffix(".lock") + try: - registry["last_updated"] = datetime.now(timezone.utc).isoformat() - with open(target, "w", encoding="utf-8") as f: - json.dump(registry, f, indent=2, ensure_ascii=False) + if not _acquire_lock(lock_path): + raise OSError(f"Could not acquire lock for {target}") + + try: + registry["last_updated"] = datetime.now(timezone.utc).isoformat() + tmp_path = target.with_suffix(".tmp") + with open(tmp_path, "w", encoding="utf-8") as f: + json.dump(registry, f, indent=2, ensure_ascii=False) + os.replace(str(tmp_path), str(target)) + finally: + _release_lock(lock_path) except Exception as e: raise Exception(f"Failed to save flow registry: {e}") @@ -361,7 +403,10 @@ def is_template_content(content: str) -> bool: # today = datetime.now().strftime("%Y%m%d") # plan_num = plan_path.stem.replace("FPLAN-", "") # template_suffix = "-TEMP" if is_template else "" -# filename = f"{folder_context}-{analysis['type']}-{analysis['category']}-{analysis['action']}-FPLAN-{plan_num}{template_suffix}-{today}.md" +# filename = ( +# f"{folder_context}-{analysis['type']}-{analysis['category']}" +# f"-{analysis['action']}-FPLAN-{plan_num}{template_suffix}-{today}.md" +# ) # # filename = re.sub(r'[<>:"|?*]', '-', filename) # filename = re.sub(r'-+', '-', filename) @@ -627,7 +672,7 @@ def process_closed_plans() -> Dict[str, Any]: if archive_success: processed_count += 1 - # Vector intake handled by close_ops.py via drone @memory process-plans + # Vector intake triggered by post_close_runner via direct import results.append({"plan": plan_label, "status": "archived", "correlation_id": correlation_id}) else: error_count += 1 diff --git a/src/aipass/flow/apps/handlers/registry/save_registry.py b/src/aipass/flow/apps/handlers/registry/save_registry.py index bb40d9cd..7756811c 100644 --- a/src/aipass/flow/apps/handlers/registry/save_registry.py +++ b/src/aipass/flow/apps/handlers/registry/save_registry.py @@ -25,6 +25,8 @@ Usage: """ import json +import os +import time from pathlib import Path from datetime import datetime, timezone from typing import Dict, Any @@ -44,6 +46,35 @@ MODULE_NAME = "save_registry" FLOW_JSON_DIR = FLOW_ROOT / "flow_json" REGISTRY_FILE = FLOW_JSON_DIR / "fplan_registry.json" +_LOCK_RETRIES = 10 +_LOCK_BACKOFF_BASE = 0.05 + + +def _acquire_lock(lock_path: Path) -> bool: + """Atomically acquire a lockfile via O_CREAT|O_EXCL with retry+backoff.""" + for attempt in range(_LOCK_RETRIES): + try: + fd = os.open(str(lock_path), os.O_CREAT | os.O_EXCL | os.O_WRONLY) + os.write(fd, str(os.getpid()).encode()) + os.close(fd) + return True + except FileExistsError: + logger.info("[%s] Lock contention on %s, retry %d", MODULE_NAME, lock_path, attempt + 1) + time.sleep(_LOCK_BACKOFF_BASE * (2**attempt)) + except OSError as exc: + logger.warning("[%s] Lock creation failed for %s: %s", MODULE_NAME, lock_path, exc) + return False + return False + + +def _release_lock(lock_path: Path) -> None: + """Remove lockfile, tolerating already-removed.""" + try: + lock_path.unlink(missing_ok=True) + except OSError as exc: + logger.warning("[%s] Could not release lock %s: %s", MODULE_NAME, lock_path, exc) + + # ============================================= # HANDLER FUNCTION # ============================================= @@ -64,15 +95,30 @@ def save_registry(registry: Dict[str, Any], registry_file: str | None = None) -> Automatically updates the last_updated timestamp before saving. Creates the flow_json directory if it doesn't exist. + Uses a lockfile to serialize concurrent writes and atomic + tempfile+rename to prevent torn reads. """ target = FLOW_JSON_DIR / registry_file if registry_file else REGISTRY_FILE + lock_path = target.with_suffix(".lock") try: FLOW_JSON_DIR.mkdir(parents=True, exist_ok=True) - registry["_notice"] = "DO NOT MANUALLY EDIT — managed by flow close pipeline" - registry["last_updated"] = datetime.now(timezone.utc).isoformat() - with open(target, "w", encoding="utf-8") as f: - json.dump(registry, f, indent=2, ensure_ascii=False) + + if not _acquire_lock(lock_path): + logger.error("[%s] Could not acquire lock for %s after %d retries", MODULE_NAME, target, _LOCK_RETRIES) + return False + + try: + registry["_notice"] = "DO NOT MANUALLY EDIT — managed by flow close pipeline" + registry["last_updated"] = datetime.now(timezone.utc).isoformat() + + tmp_path = target.with_suffix(".tmp") + with open(tmp_path, "w", encoding="utf-8") as f: + json.dump(registry, f, indent=2, ensure_ascii=False) + os.replace(str(tmp_path), str(target)) + finally: + _release_lock(lock_path) + json_handler.log_operation( "registry_saved", { diff --git a/src/aipass/flow/apps/modules/post_close_runner.py b/src/aipass/flow/apps/modules/post_close_runner.py index 7d5c3403..942ca92a 100644 --- a/src/aipass/flow/apps/modules/post_close_runner.py +++ b/src/aipass/flow/apps/modules/post_close_runner.py @@ -32,7 +32,7 @@ if sys.platform == "win32": from pathlib import Path -from aipass.cli.apps.modules import console, error, warning +from aipass.cli.apps.modules import console, error, success, warning from aipass.flow.apps.handlers.json import json_handler from aipass.flow.apps.handlers.mbank.process import process_closed_plans from aipass.flow.apps.handlers.runner.lock_ops import acquire_lock, release_lock @@ -80,7 +80,25 @@ def handle_command(command: str, args: list) -> bool: try: process_closed_plans() - console.print("[green]Processing complete[/green]") + + try: + import importlib + + _plans_mod = importlib.import_module("aipass.memory.apps.handlers.intake.plans_processor") + result = _plans_mod.process_plans() + if result.get("success"): + count = result.get("files_processed", 0) + chunks = result.get("total_chunks", 0) + if count > 0: + success(f"Vectorized {count} plan(s) ({chunks} chunks)") + logger.info("[%s] Plan vectorization: %s", MODULE_NAME, result) + else: + logger.error("[%s] Plan vectorization failed: %s", MODULE_NAME, result.get("error", "unknown")) + error(f"Vectorization failed: {result.get('error', 'unknown')}") + except Exception as e: + logger.error("[%s] Plan vectorization error: %s", MODULE_NAME, e) + error(f"Vectorization error: {e}") + except Exception as e: logger.error(f"[{MODULE_NAME}] Background processing failed: {e}") error(f"Processing failed: {e}") @@ -130,6 +148,19 @@ if __name__ == "__main__": try: process_closed_plans() + + try: + import importlib + + _plans_mod = importlib.import_module("aipass.memory.apps.handlers.intake.plans_processor") + result = _plans_mod.process_plans() + if result.get("success"): + logger.info("[%s] Plan vectorization: %s", MODULE_NAME, result) + else: + logger.error("[%s] Plan vectorization failed: %s", MODULE_NAME, result.get("error", "unknown")) + except Exception as e: + logger.error("[%s] Plan vectorization error: %s", MODULE_NAME, e) + except Exception as e: logger.error(f"[{MODULE_NAME}] Background processing failed: {e}") finally: