#662 flow: close no longer false-reports 'timed out after 30s' on a committed close. close_plan_impl now honors spawn_background — single close fires the detached _spawn_background_runner (same path as close_all) and returns right after archive; the synchronous 30s 'drone @memory process-plans' that drone was killing is gone. Removed cross-handler imports (archive/trigger injected). Fix built by @flow, verified by devpulse: 730 flow tests green (+2), seedgo 31/31 x3, live repro close=5.1s exit 0 (was 30s-timeout->false exit 1).

This commit is contained in:
AIOSAI
2026-07-09 21:15:18 -07:00
parent 26a5f3a2ee
commit bc0d403da9
5 changed files with 172 additions and 91 deletions
@@ -20,6 +20,7 @@ Usage:
_find_unregistered_plan_file,
_self_heal_unregistered_plan,
_spawn_background_runner,
_cleanup_orphaned_plan,
)
"""
@@ -253,6 +254,43 @@ def _self_heal_unregistered_plan(
return actual_key, registry
def _cleanup_orphaned_plan(
plan_file: Path,
plan_label: str,
plan_info: Dict[str, Any],
registry: Dict[str, Any],
save_registry: Any,
reg_file: Any,
messages: List[Dict[str, Any]],
archive_plan: Any = None,
) -> None:
"""Archive an orphaned .md file (registry-closed but file never moved)."""
if archive_plan is None:
logger.warning(f"[{MODULE_NAME}] archive_plan not injected, skipping orphan cleanup for {plan_label}")
messages.append({"type": "warning", "text": " Orphan cleanup skipped — archive_plan not available"})
return
messages.append({"type": "dim", "text": f" Cleaning up: moving {plan_file.name} to processed_plans/"})
try:
if archive_plan(plan_file):
logger.info(f"[{MODULE_NAME}] Cleaned up orphaned file for {plan_label}: {plan_file}")
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)
messages.append({"type": "success", "text": " Orphaned file archived successfully"})
else:
logger.warning(f"[{MODULE_NAME}] Failed to archive orphaned file for {plan_label}: {plan_file}")
messages.append({"type": "error_text", "text": " Failed to move orphaned file — manual cleanup required"})
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Error cleaning orphaned file for {plan_label}: {e}")
messages.append({"type": "error_text", "text": f" Error during cleanup: {e}"})
def _spawn_background_runner():
"""Spawn post_close_runner.py as a fully detached background process"""
bg_runner = FLOW_ROOT / "apps" / "modules" / "post_close_runner.py"
+32 -69
View File
@@ -19,7 +19,6 @@ Usage:
"""
import sys
import subprocess
from pathlib import Path
from datetime import datetime, timezone
from typing import Dict, Any, List
@@ -28,6 +27,7 @@ from aipass.prax import logger
from aipass.flow.apps.handlers.json import json_handler
from aipass.flow.apps.handlers.plan.close_helpers import (
PROCESSED_PLANS_DIR,
_extract_prefix,
_resolve_registry_file,
_find_plan_across_registries,
@@ -35,6 +35,7 @@ from aipass.flow.apps.handlers.plan.close_helpers import (
_find_unregistered_plan_file,
_self_heal_unregistered_plan,
_spawn_background_runner,
_cleanup_orphaned_plan,
)
MODULE_NAME = "close_plan"
@@ -62,6 +63,8 @@ def close_plan_impl(
push_to_plans_central: Any = None,
push_flow_to_branch_dashboard: Any = None,
close_all_plans_fn: Any = None,
archive_plan_fn: Any = None,
trigger_fire_fn: Any = None,
) -> Dict[str, Any]:
"""
Implement plan closure workflow
@@ -187,30 +190,16 @@ def close_plan_impl(
"text": f"{plan_label} already closed on {closed_date} — orphaned .md file detected",
}
)
messages.append({"type": "dim", "text": f" Cleaning up: moving {plan_file.name} to processed_plans/"})
try:
from aipass.flow.apps.handlers.mbank.process import archive_plan
if archive_plan(plan_file):
logger.info(f"[{MODULE_NAME}] Cleaned up orphaned file for {plan_label}: {plan_file}")
# Update registry flags that were missed on the failed first close
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)
messages.append({"type": "success", "text": " Orphaned file archived successfully"})
else:
logger.warning(f"[{MODULE_NAME}] Failed to archive orphaned file for {plan_label}: {plan_file}")
messages.append(
{"type": "error_text", "text": " Failed to move orphaned file — manual cleanup required"}
)
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Error cleaning orphaned file for {plan_label}: {e}")
messages.append({"type": "error_text", "text": f" Error during cleanup: {e}"})
_cleanup_orphaned_plan(
plan_file,
plan_label,
plan_info,
registry,
save_registry,
reg_file,
messages,
archive_plan=archive_plan_fn,
)
return {
"success": True,
"messages": messages,
@@ -320,15 +309,12 @@ def close_plan_impl(
# --- 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, PROCESSED_PLANS_DIR
# If file is already in processed_plans (found via relocation search), skip move
if plan_file.exists() and plan_file.parent == PROCESSED_PLANS_DIR:
archive_success = True
logger.info(f"[{MODULE_NAME}] {plan_label} already in processed_plans/, skipping move")
messages.append({"type": "dim", "text": " Already in processed_plans/ — skipping move"})
else:
archive_success = archive_plan(plan_file)
archive_success = archive_plan_fn(plan_file) if archive_plan_fn else False
if archive_success:
plan_info["processed"] = True
@@ -349,35 +335,17 @@ def close_plan_impl(
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 as e:
logger.warning(f"[{MODULE_NAME}] Best-effort drone @memory process-plans failed: {e}")
# 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})"})
# --- Vector intake (background) ---
if spawn_background:
try:
_spawn_background_runner()
logger.info(f"[{MODULE_NAME}] Spawned background vectorization for {plan_label}")
messages.append({"type": "dim", "text": " Vectorizing in background"})
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Background vectorization failed to start: {e}")
messages.append(
{"type": "warning", "text": " Background vectorization failed to start — will retry on next close"}
)
# --- Step 4/5: Update dashboards ---
messages.append({"type": "step", "text": "[4/5] Updating dashboards..."})
@@ -416,23 +384,18 @@ def close_plan_impl(
logger.warning(f"[{MODULE_NAME}] CLOSED_PLANS update failed (non-critical): {e}")
# Fire trigger event for plan closure
try:
from aipass.trigger.apps.modules.core import trigger
trigger.fire("plan_closed", plan_number=plan_key, location=str(plan_file.parent))
except ImportError:
logger.info(f"[{MODULE_NAME}] Trigger module not available, skipping event fire")
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Trigger fire failed (non-critical): {e}")
if trigger_fire_fn is not None:
try:
trigger_fire_fn("plan_closed", plan_number=plan_key, location=str(plan_file.parent))
except Exception as e:
logger.warning(f"[{MODULE_NAME}] Trigger fire failed (non-critical): {e}")
# --- VERIFY: Physical state check for self-healed plans ---
if plan_info.get("self_healed"):
messages.append({"type": "step", "text": "[VERIFY] Checking physical state..."})
try:
from aipass.flow.apps.handlers.mbank.process import PROCESSED_PLANS_DIR as _VERIFY_DIR
original_source = Path(plan_info.get("file_path", ""))
dest = _VERIFY_DIR / original_source.name
dest = PROCESSED_PLANS_DIR / original_source.name
if dest.exists():
messages.append({"type": "dim", "text": f" [OK] File in processed_plans/: {original_source.name}"})
else:
+13 -2
View File
@@ -67,12 +67,21 @@ from aipass.flow.apps.handlers.dashboard.update_local import update_dashboard_lo
from aipass.flow.apps.handlers.dashboard.push_central import push_to_plans_central
from aipass.flow.apps.handlers.dashboard.push_branch_dashboard import push_flow_to_branch_dashboard
# Internal: Memory template check (lightweight, no API calls)
from aipass.flow.apps.handlers.mbank.process import is_template_content
# Internal: Memory template check + archive (lightweight, no API calls)
from aipass.flow.apps.handlers.mbank.process import is_template_content, archive_plan
# Internal: Close operations handler (implementation)
from aipass.flow.apps.handlers.plan.close_ops import close_plan_impl, close_all_plans_impl
# Internal: Trigger (optional — may not be installed)
try:
from aipass.trigger.apps.modules.core import trigger as _trigger_module
_trigger_fire = _trigger_module.fire
except ImportError:
logger.info("[close_plan] Trigger module not available, plan events will be skipped")
_trigger_fire = None
# =============================================
# CONFIGURATION
# =============================================
@@ -264,6 +273,8 @@ def close_plan(
push_to_plans_central=push_to_plans_central,
push_flow_to_branch_dashboard=push_flow_to_branch_dashboard,
close_all_plans_fn=close_all_plans,
archive_plan_fn=archive_plan,
trigger_fire_fn=_trigger_fire,
)
# Handle dict result from handler
+74 -20
View File
@@ -49,6 +49,8 @@ def _make_deps(**overrides) -> dict:
"push_to_plans_central": MagicMock(return_value=True),
"push_flow_to_branch_dashboard": MagicMock(return_value=True),
"close_all_plans_fn": MagicMock(),
"archive_plan_fn": MagicMock(return_value=True),
"trigger_fire_fn": MagicMock(),
}
deps.update(overrides)
return deps
@@ -188,8 +190,7 @@ class TestClosePlanImplAlreadyClosedOrphan:
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops._find_plan_across_registries", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops.archive_plan", create=True)
def test_already_closed_orphan_cleanup(self, mock_archive, _mock_find, _mock_resolve, tmp_path):
def test_already_closed_orphan_cleanup(self, _mock_find, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
# Create orphan file on disk
@@ -209,9 +210,7 @@ class TestClosePlanImplAlreadyClosedOrphan:
deps["load_registry"].return_value = registry
deps["validate_plan_exists"].return_value = (True, None)
# Patch archive_plan inside the function (lazy import)
with patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True):
result = close_plan_impl(plan_num="2", **deps)
result = close_plan_impl(plan_num="2", **deps)
assert result["success"] is True
assert result["plan_key"] == "2"
@@ -257,7 +256,7 @@ class TestClosePlanImplSuccess:
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops._find_plan_across_registries", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops.subprocess")
@patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess")
def test_successful_close(self, mock_subprocess, _mock_find, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
@@ -279,7 +278,6 @@ class TestClosePlanImplSuccess:
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
patch("aipass.flow.apps.handlers.plan.append_closed_plan.append_to_closed_plans", create=True),
):
@@ -295,6 +293,67 @@ class TestClosePlanImplSuccess:
deps["push_to_plans_central"].assert_called_once()
class TestSpawnBackgroundBehavior:
"""#662: spawn_background controls whether vectorization runs inline or in background."""
@patch("aipass.flow.apps.handlers.plan.close_ops._spawn_background_runner")
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops._find_plan_across_registries", return_value=None)
def test_spawn_background_true_calls_background_runner(self, _mock_find, _mock_resolve, mock_runner, tmp_path):
close_plan_impl = _import_close_plan_impl()
plan_file = tmp_path / "FPLAN-0001_test_2026-03-20.md"
plan_file.write_text("# Real content\nNotes here.", encoding="utf-8")
registry = {
"plans": {
"1": {
"status": "open",
"subject": "Test plan",
"location": str(tmp_path),
"file_path": str(plan_file),
}
}
}
deps = _make_deps()
deps["load_registry"].return_value = registry
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
patch("aipass.flow.apps.handlers.plan.append_closed_plan.append_to_closed_plans", create=True),
):
result = close_plan_impl(plan_num="1", spawn_background=True, **deps)
assert result["success"] is True
mock_runner.assert_called_once()
assert any("background" in m.get("text", "").lower() for m in result["messages"])
@patch("aipass.flow.apps.handlers.plan.close_ops._spawn_background_runner")
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value=None)
@patch("aipass.flow.apps.handlers.plan.close_ops._find_plan_across_registries", return_value=None)
def test_spawn_background_false_skips_background_runner(self, _mock_find, _mock_resolve, mock_runner, tmp_path):
close_plan_impl = _import_close_plan_impl()
plan_file = tmp_path / "FPLAN-0001_test_2026-03-20.md"
plan_file.write_text("# Real content\nNotes here.", encoding="utf-8")
registry = {
"plans": {
"1": {
"status": "open",
"subject": "Test plan",
"location": str(tmp_path),
"file_path": str(plan_file),
}
}
}
deps = _make_deps()
deps["load_registry"].return_value = registry
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
patch("aipass.flow.apps.handlers.plan.append_closed_plan.append_to_closed_plans", create=True),
):
result = close_plan_impl(plan_num="1", spawn_background=False, **deps)
assert result["success"] is True
mock_runner.assert_not_called()
class TestClosePlanImplConfirmCancelled:
"""User cancels when confirm=True."""
@@ -688,7 +747,7 @@ class TestSelfHealCrossPrefixCollision:
class TestClosePlanImplSelfHeal:
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value="dplan_registry.json")
@patch("aipass.flow.apps.handlers.plan.close_ops.subprocess")
@patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess")
def test_triggers_self_heal_when_not_in_registry(self, mock_subprocess, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
@@ -723,7 +782,6 @@ class TestClosePlanImplSelfHeal:
},
),
) as mock_heal,
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
patch("aipass.flow.apps.handlers.plan.append_closed_plan.append_to_closed_plans", create=True),
):
@@ -752,7 +810,7 @@ class TestClosePlanImplSelfHeal:
class TestSelfHealVerifyBlock:
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value="fplan_registry.json")
@patch("aipass.flow.apps.handlers.plan.close_ops.subprocess")
@patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess")
def test_verify_all_pass(self, mock_subprocess, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
@@ -781,9 +839,8 @@ class TestSelfHealVerifyBlock:
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch(
"aipass.flow.apps.handlers.mbank.process.PROCESSED_PLANS_DIR",
"aipass.flow.apps.handlers.plan.close_ops.PROCESSED_PLANS_DIR",
processed_dir,
),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
@@ -803,7 +860,7 @@ class TestSelfHealVerifyBlock:
assert len(ok_msgs) >= 2
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value="fplan_registry.json")
@patch("aipass.flow.apps.handlers.plan.close_ops.subprocess")
@patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess")
def test_verify_fails_when_file_not_in_processed(self, mock_subprocess, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
@@ -830,9 +887,8 @@ class TestSelfHealVerifyBlock:
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch(
"aipass.flow.apps.handlers.mbank.process.PROCESSED_PLANS_DIR",
"aipass.flow.apps.handlers.plan.close_ops.PROCESSED_PLANS_DIR",
processed_dir,
),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
@@ -846,7 +902,7 @@ class TestSelfHealVerifyBlock:
assert any("NOT found in processed_plans" in m.get("text", "") for m in fail_msgs)
@patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value="fplan_registry.json")
@patch("aipass.flow.apps.handlers.plan.close_ops.subprocess")
@patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess")
def test_verify_fails_when_source_still_exists(self, mock_subprocess, _mock_resolve, tmp_path):
close_plan_impl = _import_close_plan_impl()
@@ -874,9 +930,8 @@ class TestSelfHealVerifyBlock:
deps["validate_plan_exists"].return_value = (True, None)
with (
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch(
"aipass.flow.apps.handlers.mbank.process.PROCESSED_PLANS_DIR",
"aipass.flow.apps.handlers.plan.close_ops.PROCESSED_PLANS_DIR",
processed_dir,
),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
@@ -913,8 +968,7 @@ class TestSelfHealVerifyBlock:
with (
patch("aipass.flow.apps.handlers.plan.close_ops._resolve_registry_file", return_value=None),
patch("aipass.flow.apps.handlers.plan.close_ops._find_plan_across_registries", return_value=None),
patch("aipass.flow.apps.handlers.plan.close_ops.subprocess"),
patch("aipass.flow.apps.handlers.mbank.process.archive_plan", return_value=True),
patch("aipass.flow.apps.handlers.plan.close_helpers.subprocess"),
patch("aipass.flow.apps.handlers.plan.close_ops.json_handler"),
patch("aipass.flow.apps.handlers.plan.append_closed_plan.append_to_closed_plans", create=True),
):