feat(flow): close pipeline self-completing — auto-vectorization from post_close_runner + locked atomic registry writes (DPLAN-0245). E2E proven, 730 green
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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",
|
||||
{
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user