feat(flow): feat(flow): S84 multi-type support — fix FPLAN-only hardcoding in dashboards, aggregate, restore, append, get_closed, lock race

Co-Authored-By: @flow <flow@aipass>
This commit is contained in:
AIOSAI
2026-04-10 13:55:02 -07:00
co-authored by @flow
parent 0c9a57f912
commit 119aaea9de
10 changed files with 233 additions and 81 deletions
@@ -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(
@@ -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]]:
@@ -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]]]:
@@ -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:
@@ -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", "")
@@ -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
@@ -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()
@@ -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():
+3 -2
View File
@@ -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")
+14 -20
View File
@@ -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: