From 119aaea9decd42b13df13567e3be763edd4f8ccf Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Fri, 10 Apr 2026 13:55:02 -0700 Subject: [PATCH] =?UTF-8?q?feat(flow):=20feat(flow):=20S84=20multi-type=20?= =?UTF-8?q?support=20=E2=80=94=20fix=20FPLAN-only=20hardcoding=20in=20dash?= =?UTF-8?q?boards,=20aggregate,=20restore,=20append,=20get=5Fclosed,=20loc?= =?UTF-8?q?k=20race?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: @flow --- .../dashboard/push_branch_dashboard.py | 56 +++++++++++++++---- .../apps/handlers/dashboard/push_central.py | 53 ++++++++++++++---- .../apps/handlers/dashboard/update_local.py | 56 +++++++++++++++---- .../flow/apps/handlers/plan/aggregate_ops.py | 8 +-- .../apps/handlers/plan/append_closed_plan.py | 10 +++- .../apps/handlers/plan/get_closed_plans.py | 51 ++++++++++++++--- .../flow/apps/handlers/plan/restore_ops.py | 10 ++-- .../flow/apps/modules/post_close_runner.py | 31 ++++++++-- src/aipass/flow/tests/test_aggregate_ops.py | 5 +- src/aipass/flow/tests/test_plan_handlers.py | 34 +++++------ 10 files changed, 233 insertions(+), 81 deletions(-) diff --git a/src/aipass/flow/apps/handlers/dashboard/push_branch_dashboard.py b/src/aipass/flow/apps/handlers/dashboard/push_branch_dashboard.py index 33dbb689..d57e5049 100644 --- a/src/aipass/flow/apps/handlers/dashboard/push_branch_dashboard.py +++ b/src/aipass/flow/apps/handlers/dashboard/push_branch_dashboard.py @@ -54,8 +54,9 @@ from aipass.prax.apps.modules.logger import system_logger as logger _PKG_ROOT = Path(__file__).resolve().parents[4] FLOW_ROOT = _PKG_ROOT / "flow" -# Registry location -REGISTRY_FILE = FLOW_ROOT / "flow_json" / "fplan_registry.json" +# Registry location (fallback default) +FLOW_JSON_DIR = FLOW_ROOT / "flow_json" +REGISTRY_FILE = FLOW_JSON_DIR / "fplan_registry.json" # Dashboard template path (package-relative) DASHBOARD_TEMPLATE_FILE = _PKG_ROOT / "devpulse" / "templates" / "DASHBOARD.template.json" @@ -194,25 +195,56 @@ def _calculate_quick_status(sections: Dict[str, Any]) -> Dict[str, Any]: } +# ============================================= +# PLAN TYPE HELPERS +# ============================================= + + +def _get_all_registry_files() -> List[str]: + """Return per-type registry filenames via plan-type discovery.""" + try: + from aipass.flow.apps.handlers.template.plan_type_loader import discover_plan_types + + files: List[str] = [] + for _key, config in discover_plan_types().items(): + rf = config.get("registry_file") + if rf and rf not in files: + files.append(rf) + if files: + return files + except Exception as exc: + logger.warning("[push_branch_dashboard] Failed to discover plan types, falling back to default registry: %s", exc) + return [REGISTRY_FILE.name] + + # ============================================= # HELPER FUNCTIONS # ============================================= def _load_registry() -> Dict[str, Any]: """ - Load fplan_registry.json. + Load all per-type plan registries and merge into a single dict. Returns: - Registry dict or empty structure if unavailable + Merged registry dict or empty structure if unavailable """ - try: - if not REGISTRY_FILE.exists(): - return {"plans": {}, "next_number": 1} - with open(REGISTRY_FILE, 'r', encoding='utf-8') as f: - return json.load(f) - except Exception as exc: - logger.warning("Failed to load fplan registry '%s': %s", REGISTRY_FILE, exc) - return {"plans": {}, "next_number": 1} + merged: Dict[str, Any] = {"plans": {}, "next_number": 1} + for registry_file in _get_all_registry_files(): + target = FLOW_JSON_DIR / registry_file + try: + if not target.exists(): + continue + with open(target, 'r', encoding='utf-8') as f: + data = json.load(f) + for plan_num, plan_data in data.get("plans", {}).items(): + merged["plans"][plan_num] = plan_data + # Keep the highest next_number across registries + nn = data.get("next_number", 1) + if nn > merged["next_number"]: + merged["next_number"] = nn + except Exception as exc: + logger.warning("Failed to load registry '%s': %s", target, exc) + return merged def _filter_branch_plans( diff --git a/src/aipass/flow/apps/handlers/dashboard/push_central.py b/src/aipass/flow/apps/handlers/dashboard/push_central.py index b5b42851..6df9deed 100644 --- a/src/aipass/flow/apps/handlers/dashboard/push_central.py +++ b/src/aipass/flow/apps/handlers/dashboard/push_central.py @@ -49,6 +49,8 @@ from aipass.flow.apps.handlers.plan.aggregate_ops import aggregate_central_impl MODULE_NAME = "push_central" FLOW_JSON_DIR = FLOW_ROOT / "flow_json" REGISTRY_FILE = FLOW_JSON_DIR / "fplan_registry.json" + + def _find_repo_root() -> Path: """Walk up from this file to find the repo root (contains AIPASS_REGISTRY.json).""" current = Path(__file__).resolve().parent @@ -62,25 +64,54 @@ _REPO_ROOT = _find_repo_root() AI_CENTRAL_DIR = _REPO_ROOT / ".ai_central" CENTRAL_FILE = AI_CENTRAL_DIR / "PLANS.central.json" +# ============================================= +# PLAN TYPE HELPERS +# ============================================= + + +def _get_all_registry_files() -> List[str]: + """Return per-type registry filenames via plan-type discovery.""" + try: + from aipass.flow.apps.handlers.template.plan_type_loader import discover_plan_types + + files: List[str] = [] + for _key, config in discover_plan_types().items(): + rf = config.get("registry_file") + if rf and rf not in files: + files.append(rf) + if files: + return files + except Exception as exc: + logger.warning("[push_central] Failed to discover plan types, falling back to default registry: %s", exc) + return [REGISTRY_FILE.name] + + # ============================================= # HELPER FUNCTIONS # ============================================= def _load_registry() -> Dict[str, Any]: - """Load fplan_registry.json + """Load all per-type plan registries and merge into a single dict. Returns: - Registry dict or empty structure if file doesn't exist + Merged registry dict or empty structure if no files found """ - if not REGISTRY_FILE.exists(): - return {"plans": {}, "next_number": 1} - - try: - with open(REGISTRY_FILE, 'r', encoding='utf-8') as f: - return json.load(f) - except Exception as exc: - logger.warning("Failed to load fplan registry '%s': %s", REGISTRY_FILE, exc) - return {"plans": {}, "next_number": 1} + merged: Dict[str, Any] = {"plans": {}, "next_number": 1} + for registry_file in _get_all_registry_files(): + target = FLOW_JSON_DIR / registry_file + try: + if not target.exists(): + continue + with open(target, 'r', encoding='utf-8') as f: + data = json.load(f) + for plan_num, plan_data in data.get("plans", {}).items(): + merged["plans"][plan_num] = plan_data + nn = data.get("next_number", 1) + if nn > merged["next_number"]: + merged["next_number"] = nn + except Exception as exc: + logger.warning("Failed to load registry '%s': %s", target, exc) + return merged def _extract_flow_plans(registry: Dict[str, Any]) -> tuple[List[Dict], List[Dict]]: diff --git a/src/aipass/flow/apps/handlers/dashboard/update_local.py b/src/aipass/flow/apps/handlers/dashboard/update_local.py index 1923340f..619afd9c 100644 --- a/src/aipass/flow/apps/handlers/dashboard/update_local.py +++ b/src/aipass/flow/apps/handlers/dashboard/update_local.py @@ -82,28 +82,64 @@ FLOW_ROOT = _PKG_ROOT / "flow" # CONFIGURATION # ============================================= -REGISTRY_FILE = FLOW_ROOT / "flow_json" / "fplan_registry.json" +FLOW_JSON_DIR = FLOW_ROOT / "flow_json" +REGISTRY_FILE = FLOW_JSON_DIR / "fplan_registry.json" DASHBOARD_FILE = FLOW_ROOT / "DASHBOARD.local.json" + +# ============================================= +# PLAN TYPE HELPERS +# ============================================= + + +def _get_all_registry_files() -> List[str]: + """Return per-type registry filenames via plan-type discovery.""" + try: + from aipass.flow.apps.handlers.template.plan_type_loader import discover_plan_types + + files: List[str] = [] + for _key, config in discover_plan_types().items(): + rf = config.get("registry_file") + if rf and rf not in files: + files.append(rf) + if files: + return files + except Exception as exc: + logger.warning("[update_local] Failed to discover plan types, falling back to default registry: %s", exc) + return [REGISTRY_FILE.name] + + # ============================================= # HELPER FUNCTIONS # ============================================= def _read_registry() -> Optional[Dict[str, Any]]: """ - Read fplan_registry.json. + Read all per-type plan registries and merge into a single dict. Returns: - Registry dict or None if error + Merged registry dict or None if no registries found """ - try: - if not REGISTRY_FILE.exists(): - return None - with open(REGISTRY_FILE, 'r', encoding='utf-8') as f: - return json.load(f) - except Exception as exc: - logger.warning("Failed to read fplan registry '%s': %s", REGISTRY_FILE, exc) + merged: Dict[str, Any] = {"plans": {}, "next_number": 1} + found_any = False + for registry_file in _get_all_registry_files(): + target = FLOW_JSON_DIR / registry_file + try: + if not target.exists(): + continue + with open(target, 'r', encoding='utf-8') as f: + data = json.load(f) + found_any = True + for plan_num, plan_data in data.get("plans", {}).items(): + merged["plans"][plan_num] = plan_data + nn = data.get("next_number", 1) + if nn > merged["next_number"]: + merged["next_number"] = nn + except Exception as exc: + logger.warning("Failed to read registry '%s': %s", target, exc) + if not found_any: return None + return merged def _extract_flow_plans(registry: Dict[str, Any]) -> tuple[List[Dict[str, Any]], List[Dict[str, Any]]]: diff --git a/src/aipass/flow/apps/handlers/plan/aggregate_ops.py b/src/aipass/flow/apps/handlers/plan/aggregate_ops.py index c15d0c36..39628c8e 100644 --- a/src/aipass/flow/apps/handlers/plan/aggregate_ops.py +++ b/src/aipass/flow/apps/handlers/plan/aggregate_ops.py @@ -110,17 +110,17 @@ def save_branch_registry(registry_path: Path, registry: Dict[str, Any]) -> bool: def extract_plan_number(plan_id: str) -> Optional[str]: - """Extract plan number from plan_id (e.g., 'FPLAN-0148' -> '0148') + """Extract plan number from plan_id (e.g., 'FPLAN-0148' -> '0148', 'DPLAN-0004' -> '0004') Args: - plan_id: Plan ID string (e.g., 'FPLAN-0148') + plan_id: Plan ID string with any prefix (e.g., 'FPLAN-0148', 'DPLAN-0004') Returns: Plan number string or None if invalid format """ - if not plan_id or not plan_id.startswith("FPLAN-"): + if not plan_id or '-' not in plan_id: return None - return plan_id[6:] # Skip 'FPLAN-' prefix + return plan_id.split('-', 1)[1] def auto_close_plan(registry_path: Path, plan_id: str, branch_name: str) -> bool: diff --git a/src/aipass/flow/apps/handlers/plan/append_closed_plan.py b/src/aipass/flow/apps/handlers/plan/append_closed_plan.py index 57f17ccf..f8847795 100644 --- a/src/aipass/flow/apps/handlers/plan/append_closed_plan.py +++ b/src/aipass/flow/apps/handlers/plan/append_closed_plan.py @@ -14,6 +14,7 @@ Creates the file if it doesn't exist. """ import json +import re from pathlib import Path # INFRASTRUCTURE IMPORT PATTERN @@ -40,7 +41,12 @@ def append_to_closed_plans(plan_key: str, plan_info: dict, plan_location: Path) True on success, False on failure """ try: - plan_id = f"FPLAN-{plan_key}" + # Extract prefix from plan_info's file_path (e.g., FPLAN, DPLAN) + file_path = plan_info.get("file_path", "") + filename = Path(file_path).name if file_path else "" + prefix_match = re.match(r'^([A-Z]+PLAN)', filename) + prefix = prefix_match.group(1) if prefix_match else "FPLAN" + plan_id = f"{prefix}-{plan_key}" # Extract date (YYYY-MM-DD) from the closed ISO timestamp closed_raw = plan_info.get("closed", "") @@ -49,7 +55,7 @@ def append_to_closed_plans(plan_key: str, plan_info: dict, plan_location: Path) # Build the entry entry = { "plan_id": plan_id, - "type": "FPLAN", + "type": prefix, "subject": plan_info.get("subject", ""), "date_closed": date_closed, "location": plan_info.get("relative_path", "") diff --git a/src/aipass/flow/apps/handlers/plan/get_closed_plans.py b/src/aipass/flow/apps/handlers/plan/get_closed_plans.py index 17d93e39..2d6b3e6d 100644 --- a/src/aipass/flow/apps/handlers/plan/get_closed_plans.py +++ b/src/aipass/flow/apps/handlers/plan/get_closed_plans.py @@ -19,6 +19,8 @@ Usage: from pathlib import Path from typing import List, Tuple, Dict, Any +from aipass.prax.apps.modules.logger import system_logger as logger + # INFRASTRUCTURE IMPORT PATTERN _PKG_ROOT = Path(__file__).resolve().parents[4] @@ -26,13 +28,41 @@ _PKG_ROOT = Path(__file__).resolve().parents[4] from aipass.flow.apps.handlers.registry.load_registry import load_registry from aipass.flow.apps.handlers.json import json_handler +# ============================================= +# MULTI-REGISTRY DISCOVERY +# ============================================= + +MODULE_NAME = "get_closed_plans" + + +def _get_all_registry_files() -> List[str]: + """Return per-type registry filenames via plan-type discovery. + + Uses the same pattern as mbank/process.py to discover all plan-type + registries (e.g. fplan_registry.json, dplan_registry.json). + Falls back to the default fplan_registry.json if discovery fails. + """ + try: + from aipass.flow.apps.handlers.template.plan_type_loader import discover_plan_types + files: List[str] = [] + for _key, config in discover_plan_types().items(): + rf = config.get("registry_file") + if rf and rf not in files: + files.append(rf) + if files: + return files + except Exception as exc: + logger.warning("[%s] Failed to discover plan types, falling back to default registry: %s", MODULE_NAME, exc) + return ["fplan_registry.json"] + + # ============================================= # HANDLER FUNCTION # ============================================= def get_closed_plans() -> List[Tuple[str, Dict[str, Any]]]: """ - Get all closed plans from registry + Get all closed plans from ALL discovered registries Returns: List of tuples: [(plan_num, plan_info), ...] @@ -43,15 +73,18 @@ def get_closed_plans() -> List[Tuple[str, Dict[str, Any]]]: >>> for plan_num, plan_info in plans: ... print(f"PLAN{plan_num}: {plan_info['subject']}") """ - # Load registry - registry = load_registry() + closed_plans: List[Tuple[str, Dict[str, Any]]] = [] - # Filter for closed plans - closed_plans = [ - (plan_num, plan_info) - for plan_num, plan_info in registry.get("plans", {}).items() - if plan_info.get("status") == "closed" - ] + for reg_file in _get_all_registry_files(): + try: + registry = load_registry(registry_file=reg_file) + except Exception as exc: + logger.warning("[%s] Failed to load registry '%s': %s", MODULE_NAME, reg_file, exc) + continue + + for plan_num, plan_info in registry.get("plans", {}).items(): + if plan_info.get("status") == "closed": + closed_plans.append((plan_num, plan_info)) json_handler.log_operation("closed_plans_retrieved", {"count": len(closed_plans)}) return closed_plans diff --git a/src/aipass/flow/apps/handlers/plan/restore_ops.py b/src/aipass/flow/apps/handlers/plan/restore_ops.py index 14092e7c..ced17e7f 100644 --- a/src/aipass/flow/apps/handlers/plan/restore_ops.py +++ b/src/aipass/flow/apps/handlers/plan/restore_ops.py @@ -71,12 +71,12 @@ def recover_plan_from_backup(plan_key: str, load_registry: Any = None, save_regi # Search for any prefix matching the plan key (FPLAN-, DPLAN-, etc.) variants = list(processed_plans.glob(f"*-{plan_key}*.md")) if processed_plans.exists() else [] - plan_file = processed_plans / f"FPLAN-{plan_key}.md" # fallback default + plan_file = None # No default -- use variant search if variants: # Sort by modification time, newest first variants.sort(key=lambda p: p.stat().st_mtime, reverse=True) plan_file = variants[0] # Use most recent backup - elif not plan_file.exists(): + if plan_file is None or not plan_file.exists(): return False, f"Plan {plan_key} not found in backups" # Read plan file to extract original location from header @@ -209,7 +209,7 @@ def restore_plan_impl( exists, error_msg = validate_plan_exists(plan_key, registry) if not exists: # AUTO-RECOVERY: Try to recover from processed_plans - messages.append({"type": "warning", "text": f"FPLAN-{plan_key} not in registry - attempting recovery..."}) + messages.append({"type": "warning", "text": f"PLAN-{plan_key} not in registry - attempting recovery..."}) recovered, recovery_msg = recover_plan_from_backup_fn(plan_key) if recovered: @@ -234,7 +234,7 @@ def restore_plan_impl( # 4. VALIDATE: Check plan is closed if plan_info.get("status") != "closed": - logger.warning(f"[{MODULE_NAME}] FPLAN-{plan_key} is already open") + logger.warning(f"[{MODULE_NAME}] Plan {plan_key} is already open") messages.append({"type": "error", "error_type": "already_open", "plan_key": plan_key}) return { "success": False, @@ -268,7 +268,7 @@ def restore_plan_impl( plan_info.pop('memory_file', None) save_registry(registry) - logger.info(f"[{MODULE_NAME}] Restored FPLAN-{plan_key} to open status") + logger.info(f"[{MODULE_NAME}] Restored plan {plan_key} to open status") # 8. UPDATE DASHBOARDS: Sync dashboard files (handlers) dashboard_success = update_dashboard_local() diff --git a/src/aipass/flow/apps/modules/post_close_runner.py b/src/aipass/flow/apps/modules/post_close_runner.py index a6d96a3e..b8ed593e 100644 --- a/src/aipass/flow/apps/modules/post_close_runner.py +++ b/src/aipass/flow/apps/modules/post_close_runner.py @@ -93,18 +93,37 @@ def handle_command(command: str, args: list) -> bool: def _acquire_lock() -> bool: - """Try to acquire lock file. Returns True if acquired, False if another instance is running.""" - if LOCK_FILE.exists(): + """Try to acquire lock file. Returns True if acquired, False if another instance is running. + + Uses atomic O_CREAT | O_EXCL to avoid TOCTOU race between existence check and write. + """ + try: + fd = os.open(str(LOCK_FILE), os.O_CREAT | os.O_EXCL | os.O_WRONLY) + os.write(fd, str(os.getpid()).encode()) + os.close(fd) + return True + except FileExistsError: + # Lock file exists — check if owner process is alive try: - pid = int(LOCK_FILE.read_text().strip()) + pid = int(LOCK_FILE.read_text(encoding="utf-8").strip()) os.kill(pid, 0) # Signal 0 = check if process exists logger.info(f"[{MODULE_NAME}] Another instance running (PID {pid}), exiting") return False except (ValueError, ProcessLookupError, PermissionError): + # Stale lock — remove and retry once logger.info(f"[{MODULE_NAME}] Stale lock found, taking over") - - LOCK_FILE.write_text(str(os.getpid())) - return True + try: + LOCK_FILE.unlink() + except OSError: + return False + # Retry atomic creation + try: + fd = os.open(str(LOCK_FILE), os.O_CREAT | os.O_EXCL | os.O_WRONLY) + os.write(fd, str(os.getpid()).encode()) + os.close(fd) + return True + except FileExistsError: + return False # Another process grabbed it def _release_lock(): diff --git a/src/aipass/flow/tests/test_aggregate_ops.py b/src/aipass/flow/tests/test_aggregate_ops.py index 72e6c4a0..6ad1e676 100644 --- a/src/aipass/flow/tests/test_aggregate_ops.py +++ b/src/aipass/flow/tests/test_aggregate_ops.py @@ -152,9 +152,10 @@ class TestExtractPlanNumber: extract_plan_number = _import("extract_plan_number") assert extract_plan_number(None) is None - def test_wrong_prefix(self): + def test_other_prefix(self): + """Any PREFIX-number pattern is accepted (DPLAN, FPLAN, etc.)""" extract_plan_number = _import("extract_plan_number") - assert extract_plan_number("DPLAN-0001") is None + assert extract_plan_number("DPLAN-0001") == "0001" def test_no_dash(self): extract_plan_number = _import("extract_plan_number") diff --git a/src/aipass/flow/tests/test_plan_handlers.py b/src/aipass/flow/tests/test_plan_handlers.py index dba304d3..d3fcab38 100644 --- a/src/aipass/flow/tests/test_plan_handlers.py +++ b/src/aipass/flow/tests/test_plan_handlers.py @@ -393,12 +393,14 @@ class TestAutoCloseOrphanedPlans: class TestGetClosedPlans: """Tests for get_closed_plans().""" + _SINGLE_REG = ["fplan_registry.json"] + _DISCOVERY_PATH = "aipass.flow.apps.handlers.plan.get_closed_plans._get_all_registry_files" + _LOAD_PATH = "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry" + def test_returns_only_closed_plans(self, mock_registry): _, registry = mock_registry - with patch( - "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry", - return_value=registry, - ): + with patch(self._DISCOVERY_PATH, return_value=self._SINGLE_REG), \ + patch(self._LOAD_PATH, return_value=registry): result = get_closed_plans() assert len(result) == 1 @@ -412,19 +414,15 @@ class TestGetClosedPlans: "1": {"status": "open", "subject": "active"}, } } - with patch( - "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry", - return_value=registry, - ): + with patch(self._DISCOVERY_PATH, return_value=self._SINGLE_REG), \ + patch(self._LOAD_PATH, return_value=registry): result = get_closed_plans() assert result == [] def test_returns_empty_on_empty_registry(self): - with patch( - "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry", - return_value={"plans": {}}, - ): + with patch(self._DISCOVERY_PATH, return_value=self._SINGLE_REG), \ + patch(self._LOAD_PATH, return_value={"plans": {}}): result = get_closed_plans() assert result == [] @@ -437,10 +435,8 @@ class TestGetClosedPlans: "3": {"status": "open", "subject": "still going"}, } } - with patch( - "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry", - return_value=registry, - ): + with patch(self._DISCOVERY_PATH, return_value=self._SINGLE_REG), \ + patch(self._LOAD_PATH, return_value=registry): result = get_closed_plans() assert len(result) == 2 @@ -449,10 +445,8 @@ class TestGetClosedPlans: def test_result_tuples_contain_plan_num_and_info(self, mock_registry): _, registry = mock_registry - with patch( - "aipass.flow.apps.handlers.plan.get_closed_plans.load_registry", - return_value=registry, - ): + with patch(self._DISCOVERY_PATH, return_value=self._SINGLE_REG), \ + patch(self._LOAD_PATH, return_value=registry): result = get_closed_plans() for plan_num, plan_info in result: