feat(flow): feat(flow): foreground archival, --dry-run, template auto-heal, vector verify pipeline

Co-Authored-By: @flow <flow@aipass>
This commit is contained in:
AIOSAI
2026-03-18 23:31:35 -07:00
co-authored by @flow
parent d57a7a1322
commit e30f54b355
11 changed files with 175 additions and 50 deletions
+14 -4
View File
@@ -29,6 +29,8 @@ from datetime import datetime, timezone
from typing import Dict, List, Optional, Any
from aipass.flow.apps.handlers.json import json_handler
from aipass.prax.apps.modules.logger import system_logger as logger
from aipass.cli.apps.modules import error as cli_error, warning as cli_warning
# AI summarization removed — OpenRouter API no longer needed here
# from aipass.api.apps.modules.openrouter_client import get_response
@@ -689,7 +691,10 @@ def verify_and_heal_orphaned_plans() -> Dict[str, Any]:
except Exception:
continue
for plan_num, plan_info in registry.get("plans", {}).items():
if plan_info.get("processed") is True and plan_info.get("cleanup_completed") is True:
# Heal ANY closed plan whose file still sits at its original location.
# Covers both fully-processed plans and auto-closed orphans that
# never entered the post-close pipeline.
if plan_info.get("status") == "closed":
original_path = Path(plan_info.get("file_path", ""))
if original_path.exists():
orphans_found += 1
@@ -757,12 +762,17 @@ def process_closed_plans() -> Dict[str, Any]:
if archive_success:
processed_count += 1
# Best-effort vector processing
# Vector processing — errors go to prax log + console
try:
from aipass.memory.apps.handlers.intake.plans_processor import process_plans # type: ignore[import-not-found]
process_plans()
except Exception:
pass
logger.info("[mbank] Vector intake completed for %s", plan_label)
except ImportError:
logger.error("[mbank] Vector intake FAILED for %s — plans_processor not found", plan_label)
cli_error(f"Vector intake unavailable — memory plans_processor not found ({plan_label})")
except Exception as vec_err:
logger.error("[mbank] Vector intake FAILED for %s: %s", plan_label, vec_err)
cli_error(f"Vector intake failed for {plan_label}: {vec_err}")
results.append({"plan": plan_label, "status": "archived", "correlation_id": correlation_id})
else:
error_count += 1
+96 -17
View File
@@ -108,6 +108,7 @@ def _spawn_background_runner():
def close_plan_impl(plan_num: Any = None, confirm: bool = False,
all_plans: bool = False, spawn_background: bool = True,
dry_run: bool = False,
# Dependencies injected from module
normalize_plan_number: Any = None,
load_registry: Any = None,
@@ -131,6 +132,7 @@ def close_plan_impl(plan_num: Any = None, confirm: bool = False,
all_plans: If True, close all open plans (default False)
spawn_background: Whether to spawn background post-processing (default True).
Set False when called from close_all_plans() to avoid race condition.
dry_run: If True, preview what would be closed without taking action (default False)
(remaining args): Handler/service dependencies injected by module
Returns:
@@ -141,7 +143,7 @@ def close_plan_impl(plan_num: Any = None, confirm: bool = False,
# Handle --all flag
if all_plans:
return close_all_plans_fn(confirm)
return close_all_plans_fn(confirm, dry_run=dry_run)
# Single plan closure
if not plan_num:
@@ -194,6 +196,24 @@ def close_plan_impl(plan_num: Any = None, confirm: bool = False,
# Extract prefix for display functions (e.g. "FPLAN", "DPLAN")
plan_prefix = _extract_prefix(plan_label) or "FPLAN"
# DRY RUN: Preview what would be closed, then return early
if dry_run:
location = plan_info.get("location", "unknown")
subject = plan_info.get("subject", "No subject")
status = plan_info.get("status", "unknown")
messages.append({"type": "dim", "text": f"[DRY RUN] Would close {plan_label}"})
messages.append({"type": "dim", "text": f" Location: {location}"})
messages.append({"type": "dim", "text": f" Subject: {subject}"})
messages.append({"type": "dim", "text": f" Status: {status}"})
messages.append({"type": "dim", "text": "No action taken."})
logger.info(f"[{MODULE_NAME}] Dry run: would close {plan_label}")
return {
"success": True,
"messages": messages,
"plan_key": plan_key,
"cancelled": False,
}
# 4. IDEMPOTENCY CHECK: Prevent double-closing (with orphan cleanup)
if plan_info['status'] == 'closed':
closed_date = plan_info.get('closed', 'unknown')
@@ -300,21 +320,57 @@ def close_plan_impl(plan_num: Any = None, confirm: bool = False,
"cancelled": False,
}
# --- Step 3/5: Background processing ---
if spawn_background:
messages.append({"type": "step", "text": "[3/5] Starting background processing..."})
try:
_spawn_background_runner()
logger.info(f"[{MODULE_NAME}] Spawned background post-processing for {plan_label}")
messages.append({"type": "dim", "text": " Summary generation and archival running in background"})
except FileNotFoundError as e:
logger.warning(f"[{MODULE_NAME}] Background runner not found: {e}")
messages.append({"type": "warning", "text": " Background runner not found - will retry on next close"})
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Failed to spawn background post-processing: {e}")
messages.append({"type": "warning", "text": " Background archival failed to start - will retry on next close"})
else:
messages.append({"type": "step", "text": "[3/5] Background processing deferred (batch mode)"})
# --- Step 3/5: Archive plan to processed_plans ---
messages.append({"type": "step", "text": "[3/5] Archiving plan..."})
try:
from aipass.flow.apps.handlers.mbank.process import archive_plan
archive_success = archive_plan(plan_file)
if archive_success:
# Set flags on same registry object we already have in memory
plan_info["processed"] = True
plan_info["processed_date"] = datetime.now(timezone.utc).isoformat()
plan_info["cleanup_completed"] = True
plan_info["cleanup_date"] = datetime.now(timezone.utc).isoformat()
if reg_file:
save_registry(registry, registry_file=reg_file)
else:
save_registry(registry)
logger.info(f"[{MODULE_NAME}] Archived {plan_label} to processed_plans")
messages.append({"type": "dim", "text": " Plan archived to processed_plans/"})
else:
logger.error(f"[{MODULE_NAME}] Failed to archive {plan_label}")
messages.append({"type": "warning", "text": " Archive failed — plan file not moved"})
except Exception as e:
logger.error(f"[{MODULE_NAME}] Archive error for {plan_label}: {e}")
messages.append({"type": "warning", "text": f" Archive error: {e}"})
# --- Vector intake + verification ---
# Trigger memory's plan processor via drone (no cross-branch imports)
try:
subprocess.run(
["drone", "@memory", "process-plans"],
capture_output=True, timeout=30,
)
except Exception:
pass # Best effort — verification below reports actual status
# Verify vectorization via memory's verify module
try:
from aipass.memory.apps.modules.verify import is_plan_vectorized # type: ignore[import-not-found]
result = is_plan_vectorized(plan_label)
if result.get("found"):
chunk_count = result.get("count", 0)
logger.info(f"[{MODULE_NAME}] Vectorized: {plan_label} ({chunk_count} chunks)")
messages.append({"type": "dim", "text": f" Vectorized: {chunk_count} chunks in chroma"})
else:
logger.warning(f"[{MODULE_NAME}] NOT vectorized: {plan_label}")
messages.append({"type": "warning", "text": " NOT vectorized — check drone @memory process-plans"})
except ImportError:
logger.warning(f"[{MODULE_NAME}] Vector verify unavailable — memory verify module not found")
messages.append({"type": "warning", "text": " Vector status: unknown (memory verify not available)"})
except Exception as vec_err:
logger.warning(f"[{MODULE_NAME}] Vector verify failed: {vec_err}")
messages.append({"type": "warning", "text": f" Vector status: unknown ({vec_err})"})
# --- Step 4/5: Update dashboards ---
messages.append({"type": "step", "text": "[4/5] Updating dashboards..."})
@@ -385,7 +441,7 @@ def close_plan_impl(plan_num: Any = None, confirm: bool = False,
}
def close_all_plans_impl(confirm: bool = False,
def close_all_plans_impl(confirm: bool = False, dry_run: bool = False,
# Dependencies injected from module
get_open_plans: Any = None,
close_plan_fn: Any = None) -> Dict[str, Any]:
@@ -394,6 +450,7 @@ def close_all_plans_impl(confirm: bool = False,
Args:
confirm: Whether to ask for bulk confirmation (default False, auto-confirms)
dry_run: If True, preview what would be closed without taking action (default False)
get_open_plans: Handler function to get open plans
close_plan_fn: Function to close a single plan (the module's close_plan)
@@ -416,6 +473,28 @@ def close_all_plans_impl(confirm: bool = False,
"total": 0,
}
# DRY RUN: Preview all plans that would be closed, then return early
if dry_run:
messages.append({"type": "dim", "text": f"[DRY RUN] Would close {len(open_plans)} plan(s):"})
for plan_num, plan_info in open_plans:
subject = plan_info.get("subject", "No subject")
location = plan_info.get("location", "unknown")
# Derive prefix from file_path if available
plan_file = Path(plan_info.get("file_path", ""))
plan_label = plan_file.stem if plan_file.name else f"PLAN-{plan_num}"
prefix = _extract_prefix(plan_label) or "FPLAN"
display_id = f"{prefix}-{plan_num}"
messages.append({"type": "dim", "text": f" {display_id:<14}{location:<14}{subject}"})
messages.append({"type": "dim", "text": "No action taken."})
logger.info(f"[{MODULE_NAME}] Dry run: would close {len(open_plans)} plan(s)")
return {
"success": True,
"messages": messages,
"success_count": 0,
"failure_count": 0,
"total": len(open_plans),
}
# Build plan list for display
plan_list = []
for plan_num, plan_info in open_plans:
@@ -99,45 +99,53 @@ def parse_delete_command_args(args: List[str]) -> Tuple[str | None, bool, str |
return plan_num, confirm, None
def parse_close_command_args(args: List[str]) -> Tuple[str | None, bool, bool, str | None]:
def parse_close_command_args(args: List[str]) -> Tuple[str | None, bool, bool, bool, str | None]:
"""
Parse arguments for close command
Auto-confirms by default (running 'close' IS the intent).
Use --confirm or --interactive to explicitly request a confirmation prompt.
--yes/-y kept for backwards compatibility (now redundant, already auto-confirms).
--dry-run or --preview previews what would be closed without taking action.
Args:
args: Command arguments
Returns:
Tuple of (plan_num, confirm, all_plans, error_message)
Tuple of (plan_num, confirm, all_plans, dry_run, error_message)
- plan_num: Plan number from first arg, or None if --all or missing
- confirm: True only if --confirm or --interactive flag present, False otherwise
- all_plans: True if --all flag present, False otherwise
- dry_run: True if --dry-run or --preview flag present, False otherwise
- error_message: None if valid, error string if invalid args
Examples:
>>> parse_close_command_args(["42"])
("42", False, False, None)
("42", False, False, False, None)
>>> parse_close_command_args(["42", "--yes"])
("42", False, False, None)
("42", False, False, False, None)
>>> parse_close_command_args(["42", "--confirm"])
("42", True, False, None)
("42", True, False, False, None)
>>> parse_close_command_args(["42", "--interactive"])
("42", True, False, None)
("42", True, False, False, None)
>>> parse_close_command_args(["--all"])
(None, False, True, None)
(None, False, True, False, None)
>>> parse_close_command_args(["--all", "--confirm"])
(None, True, True, None)
(None, True, True, False, None)
>>> parse_close_command_args(["42", "--dry-run"])
("42", False, False, True, None)
>>> parse_close_command_args(["--all", "--preview"])
(None, False, True, True, None)
>>> parse_close_command_args([])
(None, False, False, "Plan number or --all required")
(None, False, False, False, "Plan number or --all required")
"""
# Check for --all flag
all_plans = '--all' in args
@@ -147,18 +155,21 @@ def parse_close_command_args(args: List[str]) -> Tuple[str | None, bool, bool, s
# --yes/-y kept for backwards compat (redundant, already auto-confirms)
confirm = '--confirm' in args or '--interactive' in args
# Check for --dry-run or --preview flag
dry_run = '--dry-run' in args or '--preview' in args
# If --all, plan_num is None
if all_plans:
return None, confirm, True, None
return None, confirm, True, dry_run, None
# Otherwise, need plan number
# Filter out flag args to find the plan number
non_flag_args = [a for a in args if not a.startswith('--') and a not in ('-y',)]
if not non_flag_args:
return None, False, False, "Plan number or --all required"
return None, False, False, dry_run, "Plan number or --all required"
plan_num = non_flag_args[0]
return plan_num, confirm, False, None
return plan_num, confirm, False, dry_run, None
def parse_restore_command_args(args: List[str]) -> Tuple[str | None, str | None]:
@@ -147,6 +147,38 @@ def load_registry() -> Dict[str, Any]:
"type_count": len(data["types"]),
}
# Auto-heal: prune orphaned types (directory deleted but registry entry remains)
templates_dir = FLOW_ROOT / "templates"
orphaned = [
dir_name
for dir_name in data["types"]
if dir_name not in _PROTECTED_TYPES
and not (templates_dir / dir_name).is_dir()
]
if orphaned:
plan_registry_dir = FLOW_ROOT / "flow_json"
for dir_name in orphaned:
entry = data["types"][dir_name]
shorthand = entry.get("shorthand", entry.get("prefix", "").lower())
logger.info(
"[%s] Auto-pruning orphaned type '%s' (directory missing)",
MODULE_NAME,
dir_name,
)
del data["types"][dir_name]
# Clean up the per-type plan registry JSON
if shorthand:
plan_reg = plan_registry_dir / f"{shorthand}_registry.json"
if plan_reg.exists():
plan_reg.unlink()
logger.info(
"[%s] Removed orphaned plan registry: %s",
MODULE_NAME,
plan_reg.name,
)
save_registry(data)
return data
+10 -4
View File
@@ -185,11 +185,13 @@ def print_help():
console.print("[yellow]OPTIONS:[/yellow]")
console.print(" --all Close all open plans")
console.print(" --confirm Interactive confirmation prompt")
console.print(" --dry-run Preview what would be closed (no action taken)")
console.print()
console.print("[yellow]EXAMPLES:[/yellow]")
console.print(" [dim]drone @flow close FPLAN-0042[/dim] # Close specific plan")
console.print(" [dim]drone @flow close DPLAN-0005[/dim] # Close a DPLAN")
console.print(" [dim]drone @flow close --all[/dim] # Close all open plans")
console.print(" [dim]drone @flow close --all --dry-run[/dim] # Preview close-all")
console.print()
@@ -197,7 +199,7 @@ def print_help():
# CLOSE PLAN WORKFLOW (thin orchestrator)
# =============================================
def close_plan(plan_num: str | None = None, confirm: bool = False, all_plans: bool = False, spawn_background: bool = True) -> bool:
def close_plan(plan_num: str | None = None, confirm: bool = False, all_plans: bool = False, spawn_background: bool = True, dry_run: bool = False) -> bool:
"""
Orchestrate plan closure workflow (thin orchestrator)
@@ -218,6 +220,7 @@ def close_plan(plan_num: str | None = None, confirm: bool = False, all_plans: bo
all_plans: If True, close all open plans (default False)
spawn_background: Whether to spawn background post-processing (default True).
Set False when called from close_all_plans() to avoid race condition.
dry_run: If True, preview what would be closed without taking action (default False)
Returns:
True if successful, False otherwise
@@ -227,6 +230,7 @@ def close_plan(plan_num: str | None = None, confirm: bool = False, all_plans: bo
confirm=confirm,
all_plans=all_plans,
spawn_background=spawn_background,
dry_run=dry_run,
# Inject dependencies
normalize_plan_number=normalize_plan_number,
load_registry=load_registry,
@@ -249,18 +253,20 @@ def close_plan(plan_num: str | None = None, confirm: bool = False, all_plans: bo
return bool(result)
def close_all_plans(confirm: bool = False) -> bool:
def close_all_plans(confirm: bool = False, dry_run: bool = False) -> bool:
"""
Close all open plans in one operation (thin orchestrator)
Args:
confirm: Whether to ask for bulk confirmation (default False, auto-confirms)
dry_run: If True, preview what would be closed without taking action (default False)
Returns:
True if at least one plan closed successfully, False otherwise
"""
result = close_all_plans_impl(
confirm=confirm,
dry_run=dry_run,
get_open_plans=get_open_plans,
close_plan_fn=close_plan,
)
@@ -313,7 +319,7 @@ def handle_command(command: str, args: List[str]) -> bool:
)
# 1. PARSE ARGS: Use command_parser handler
plan_num, confirm, all_plans, error = parse_close_command_args(args)
plan_num, confirm, all_plans, dry_run, error = parse_close_command_args(args)
# 2. VALIDATE: Check for parsing errors
if error:
@@ -321,7 +327,7 @@ def handle_command(command: str, args: List[str]) -> bool:
return True # Command was handled (error already displayed)
# 3. EXECUTE: Run workflow orchestrator
close_plan(plan_num=plan_num, confirm=confirm, all_plans=all_plans)
close_plan(plan_num=plan_num, confirm=confirm, all_plans=all_plans, dry_run=dry_run)
# 4. RETURN: True = command was handled (even if the operation failed,
# the error has already been displayed -- returning False would cause
@@ -1,13 +0,0 @@
# {plan_number}: {subject}
Tag: {tag}
> Testing template — auto-discovered from filesystem
---
## Notes
---
*Created: {today}*
*Updated: {today}*