feat: auto-spawn watchdog after dispatch (FPLAN-0189)
This commit is contained in:
@@ -13,6 +13,8 @@ Orchestrates dispatch commands: status tracking and daemon management.
|
||||
Delegates all business logic to handlers.
|
||||
"""
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import List
|
||||
@@ -39,6 +41,7 @@ DISPATCH (send + wake):
|
||||
drone @ai_mail dispatch @branch "Subject" "Body" --fresh # Send + fresh wake
|
||||
drone @ai_mail dispatch @branch "Subject" "Body" --model opus # Send + wake with Opus
|
||||
drone @ai_mail dispatch @branch "Subject" "Body" --no-memory-save
|
||||
drone @ai_mail dispatch @branch "Subject" "Body" --no-watchdog # Skip auto-watchdog
|
||||
|
||||
WAKE ONLY:
|
||||
drone @ai_mail dispatch wake @branch # Wake with default inbox check
|
||||
@@ -216,6 +219,7 @@ def _orchestrate_dispatch_send(args: List[str]) -> bool:
|
||||
# Parse flags
|
||||
use_fresh = False
|
||||
no_memory_save = False
|
||||
no_watchdog = False
|
||||
from_branch = None
|
||||
use_model = None
|
||||
filtered = []
|
||||
@@ -229,6 +233,10 @@ def _orchestrate_dispatch_send(args: List[str]) -> bool:
|
||||
no_memory_save = True
|
||||
i += 1
|
||||
continue
|
||||
if args[i] == "--no-watchdog":
|
||||
no_watchdog = True
|
||||
i += 1
|
||||
continue
|
||||
if args[i] == "--from" and i + 1 < len(args):
|
||||
from_branch = args[i + 1]
|
||||
i += 2
|
||||
@@ -335,10 +343,57 @@ def _orchestrate_dispatch_send(args: List[str]) -> bool:
|
||||
if not wake_ok:
|
||||
logger.warning("[dispatch] Wake failed for %s — email was sent", target)
|
||||
error(f"Email sent but wake failed — retry: drone @ai_mail dispatch wake {target}")
|
||||
elif not no_watchdog:
|
||||
_spawn_watchdog(target)
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def _spawn_watchdog(target: str) -> None:
|
||||
"""Auto-spawn devpulse watchdog as a detached background process."""
|
||||
from aipass.ai_mail.apps.handlers.registry.read import get_branch_by_email
|
||||
|
||||
devpulse_info = get_branch_by_email("@devpulse")
|
||||
if not devpulse_info:
|
||||
logger.warning("[dispatch] Cannot spawn watchdog — @devpulse not in registry")
|
||||
return
|
||||
|
||||
_ai_mail_dir = Path(__file__).resolve().parents[2]
|
||||
_repo_root = _ai_mail_dir.parents[2]
|
||||
devpulse_path = devpulse_info.get("path", "")
|
||||
if not devpulse_path:
|
||||
logger.warning("[dispatch] Cannot spawn watchdog — @devpulse has no path")
|
||||
return
|
||||
|
||||
devpulse_dir = Path(devpulse_path)
|
||||
if not devpulse_dir.is_absolute():
|
||||
devpulse_dir = _repo_root / devpulse_dir
|
||||
|
||||
if not devpulse_dir.is_dir():
|
||||
logger.warning("[dispatch] Cannot spawn watchdog — devpulse dir not found: %s", devpulse_dir)
|
||||
return
|
||||
|
||||
cmd = ["drone", "@devpulse", "watchdog", "agent", target]
|
||||
|
||||
spawn_env = os.environ.copy()
|
||||
local_bin = str(Path.home() / ".local" / "bin")
|
||||
if local_bin not in spawn_env.get("PATH", ""):
|
||||
spawn_env["PATH"] = local_bin + ":" + spawn_env.get("PATH", "")
|
||||
|
||||
try:
|
||||
subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.DEVNULL,
|
||||
stderr=subprocess.DEVNULL,
|
||||
start_new_session=True,
|
||||
cwd=str(devpulse_dir),
|
||||
env=spawn_env,
|
||||
)
|
||||
console.print(f"[green]Watchdog armed for {target}[/green]")
|
||||
except Exception as e:
|
||||
logger.warning("[dispatch] Watchdog spawn failed for %s: %s", target, e)
|
||||
|
||||
|
||||
def _orchestrate_daemon() -> bool:
|
||||
"""Orchestrate daemon startup."""
|
||||
logger.info("[dispatch] Starting dispatch daemon")
|
||||
|
||||
@@ -968,3 +968,248 @@ class TestPrintIntrospection:
|
||||
assert "status.py" in combined
|
||||
assert "wake.py" in combined
|
||||
assert "daemon.py" in combined
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# _spawn_watchdog
|
||||
# ===========================================================================
|
||||
|
||||
|
||||
class TestSpawnWatchdog:
|
||||
"""Tests for _spawn_watchdog."""
|
||||
|
||||
def test_spawns_detached_subprocess(self, monkeypatch, tmp_path):
|
||||
"""Successful watchdog spawn calls Popen with correct args."""
|
||||
devpulse_dir = tmp_path / "src" / "aipass" / "devpulse"
|
||||
devpulse_dir.mkdir(parents=True)
|
||||
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
popen_calls: list[dict] = []
|
||||
mock_popen = MagicMock()
|
||||
|
||||
def tracking_popen(cmd, **kwargs):
|
||||
"""Capture Popen arguments."""
|
||||
popen_calls.append({"cmd": cmd, **kwargs})
|
||||
return mock_popen
|
||||
|
||||
with (
|
||||
patch(
|
||||
f"{_H_REG}.get_branch_by_email",
|
||||
return_value={"email": "@devpulse", "path": str(devpulse_dir)},
|
||||
),
|
||||
patch(f"{MOD}.subprocess.Popen", side_effect=tracking_popen),
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
assert len(popen_calls) == 1
|
||||
assert popen_calls[0]["cmd"] == ["drone", "@devpulse", "watchdog", "agent", "@flow"]
|
||||
assert popen_calls[0]["start_new_session"] is True
|
||||
assert popen_calls[0]["cwd"] == str(devpulse_dir)
|
||||
combined = " ".join(printed)
|
||||
assert "Watchdog armed for @flow" in combined
|
||||
|
||||
def test_devpulse_not_in_registry(self, monkeypatch):
|
||||
"""No spawn when @devpulse not found in registry."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
with (
|
||||
patch(f"{_H_REG}.get_branch_by_email", return_value=None),
|
||||
patch(f"{MOD}.subprocess.Popen") as mock_popen,
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
mock_popen.assert_not_called()
|
||||
|
||||
def test_devpulse_no_path(self, monkeypatch):
|
||||
"""No spawn when @devpulse has empty path."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
with (
|
||||
patch(
|
||||
f"{_H_REG}.get_branch_by_email",
|
||||
return_value={"email": "@devpulse", "path": ""},
|
||||
),
|
||||
patch(f"{MOD}.subprocess.Popen") as mock_popen,
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
mock_popen.assert_not_called()
|
||||
|
||||
def test_devpulse_dir_missing(self, monkeypatch, tmp_path):
|
||||
"""No spawn when devpulse directory doesn't exist."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
with (
|
||||
patch(
|
||||
f"{_H_REG}.get_branch_by_email",
|
||||
return_value={"email": "@devpulse", "path": str(tmp_path / "nonexistent")},
|
||||
),
|
||||
patch(f"{MOD}.subprocess.Popen") as mock_popen,
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
mock_popen.assert_not_called()
|
||||
|
||||
def test_popen_failure_warns_but_does_not_raise(self, monkeypatch, tmp_path):
|
||||
"""Popen failure logs warning but doesn't propagate."""
|
||||
devpulse_dir = tmp_path / "src" / "aipass" / "devpulse"
|
||||
devpulse_dir.mkdir(parents=True)
|
||||
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
with (
|
||||
patch(
|
||||
f"{_H_REG}.get_branch_by_email",
|
||||
return_value={"email": "@devpulse", "path": str(devpulse_dir)},
|
||||
),
|
||||
patch(f"{MOD}.subprocess.Popen", side_effect=FileNotFoundError("drone not found")),
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
# Should not raise — watchdog is optional
|
||||
|
||||
def test_relative_devpulse_path_resolved(self, monkeypatch, tmp_path):
|
||||
"""Relative path from registry is resolved against repo root."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
popen_calls: list[dict] = []
|
||||
|
||||
def tracking_popen(cmd, **kwargs):
|
||||
"""Capture Popen arguments."""
|
||||
popen_calls.append({"cmd": cmd, **kwargs})
|
||||
return MagicMock()
|
||||
|
||||
from aipass.ai_mail.apps.modules import dispatch as dispatch_mod
|
||||
|
||||
real_repo_root = dispatch_mod.Path(__file__).resolve().parents[4]
|
||||
devpulse_dir = real_repo_root / "src" / "aipass" / "devpulse"
|
||||
|
||||
with (
|
||||
patch(
|
||||
f"{_H_REG}.get_branch_by_email",
|
||||
return_value={"email": "@devpulse", "path": "src/aipass/devpulse"},
|
||||
),
|
||||
patch(f"{MOD}.subprocess.Popen", side_effect=tracking_popen),
|
||||
):
|
||||
from aipass.ai_mail.apps.modules.dispatch import _spawn_watchdog
|
||||
|
||||
_spawn_watchdog("@flow")
|
||||
|
||||
if devpulse_dir.is_dir():
|
||||
assert len(popen_calls) == 1
|
||||
assert "devpulse" in popen_calls[0]["cwd"]
|
||||
else:
|
||||
assert len(popen_calls) == 0
|
||||
|
||||
|
||||
class TestDispatchSendWatchdogIntegration:
|
||||
"""Tests for watchdog integration in _orchestrate_dispatch_send."""
|
||||
|
||||
def test_watchdog_spawned_after_successful_wake(self, monkeypatch):
|
||||
"""Watchdog is spawned after successful send + wake."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
watchdog_calls: list[str] = []
|
||||
monkeypatch.setattr(
|
||||
f"{MOD}._spawn_watchdog",
|
||||
lambda target: watchdog_calls.append(target),
|
||||
)
|
||||
|
||||
patches = _send_patches()
|
||||
with patches:
|
||||
from aipass.ai_mail.apps.modules.dispatch import _orchestrate_dispatch_send
|
||||
|
||||
result = _orchestrate_dispatch_send(["@target", "Subject", "Body"])
|
||||
|
||||
assert result is True
|
||||
assert watchdog_calls == ["@target"]
|
||||
|
||||
def test_watchdog_not_spawned_on_wake_failure(self, monkeypatch):
|
||||
"""Watchdog is NOT spawned when wake fails."""
|
||||
errors: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.error", lambda msg: errors.append(msg))
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
watchdog_calls: list[str] = []
|
||||
monkeypatch.setattr(
|
||||
f"{MOD}._spawn_watchdog",
|
||||
lambda target: watchdog_calls.append(target),
|
||||
)
|
||||
|
||||
mock_status = MagicMock()
|
||||
mock_status.format.return_value = "WAKE FAILED"
|
||||
patches = _send_patches(
|
||||
{
|
||||
f"{_H_WAKE}.wake_branch": MagicMock(return_value=(mock_status, False)),
|
||||
}
|
||||
)
|
||||
with patches:
|
||||
from aipass.ai_mail.apps.modules.dispatch import _orchestrate_dispatch_send
|
||||
|
||||
_orchestrate_dispatch_send(["@target", "Subject", "Body"])
|
||||
|
||||
assert watchdog_calls == []
|
||||
|
||||
def test_no_watchdog_flag_skips_spawn(self, monkeypatch):
|
||||
"""--no-watchdog flag prevents watchdog spawn."""
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
watchdog_calls: list[str] = []
|
||||
monkeypatch.setattr(
|
||||
f"{MOD}._spawn_watchdog",
|
||||
lambda target: watchdog_calls.append(target),
|
||||
)
|
||||
|
||||
patches = _send_patches()
|
||||
with patches:
|
||||
from aipass.ai_mail.apps.modules.dispatch import _orchestrate_dispatch_send
|
||||
|
||||
result = _orchestrate_dispatch_send(["@target", "Subject", "Body", "--no-watchdog"])
|
||||
|
||||
assert result is True
|
||||
assert watchdog_calls == []
|
||||
|
||||
def test_watchdog_not_spawned_on_send_failure(self, monkeypatch):
|
||||
"""Watchdog is NOT spawned when send fails."""
|
||||
errors: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.error", lambda msg: errors.append(msg))
|
||||
printed: list[str] = []
|
||||
monkeypatch.setattr(f"{MOD}.console", _mock_console(printed))
|
||||
|
||||
watchdog_calls: list[str] = []
|
||||
monkeypatch.setattr(
|
||||
f"{MOD}._spawn_watchdog",
|
||||
lambda target: watchdog_calls.append(target),
|
||||
)
|
||||
|
||||
patches = _send_patches(
|
||||
{
|
||||
f"{_H_SEND}.send_to_single": MagicMock(return_value=(False, "error")),
|
||||
}
|
||||
)
|
||||
with patches:
|
||||
from aipass.ai_mail.apps.modules.dispatch import _orchestrate_dispatch_send
|
||||
|
||||
_orchestrate_dispatch_send(["@target", "Subject", "Body"])
|
||||
|
||||
assert watchdog_calls == []
|
||||
|
||||
Reference in New Issue
Block a user