Merge pull request #482 from AIOSAI/work/system
feat(system): fix: DPLAN-0157 pre-commit hook + DPLAN-0158 watchdog reply detection
This commit is contained in:
@@ -197,7 +197,7 @@ Always work on main. Edit files in your branch directory on the main branch. Whe
|
||||
|
||||
**Blocked system-wide via `.claude/settings.json` permission gate:** `git checkout*` (any form — switch, discard, new branch), `git add -f*`, `git add --force*`. These are denied for every agent including devpulse. Use `drone @git sync` to switch to main, `drone @git fix` to recover from broken states.
|
||||
|
||||
**Culturally blocked (no permission gate yet, still don't use):** `git commit`, `git push`, `gh pr create`. Go through drone.
|
||||
**Mechanically blocked via `.git/hooks/pre-commit`:** `git commit` is rejected on any branch except main (also catches detached HEAD). `git push`, `gh pr create` — go through drone.
|
||||
|
||||
**Allowed read-only:** `git status`, `git diff`, `git log`, `git stash` (safe transient save).
|
||||
|
||||
|
||||
@@ -119,6 +119,11 @@
|
||||
"file": "apps/modules/integrations_manager.py",
|
||||
"standard": "unused_function",
|
||||
"reason": "fetch_contracts() and call_contract() are test API hooks called from tests/test_integrations.py. The unused_function checker excludes test dirs from its search corpus, so these appear unused despite being actively called. Encapsulation standard requires tests to go through modules, not import handlers directly — these functions serve that purpose."
|
||||
},
|
||||
{
|
||||
"file": "apps/modules/api_key.py",
|
||||
"standard": "unused_function",
|
||||
"reason": "fetch_api_key() and fetch_validate_key() are module-level wrappers called from tests/test_critical_paths.py. The unused_function checker excludes test dirs from its search corpus. Encapsulation standard requires tests to go through modules — these functions serve that purpose (DPLAN-0155)."
|
||||
}
|
||||
],
|
||||
"notes": {
|
||||
|
||||
@@ -29,6 +29,7 @@ import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
from aipass.prax.apps.modules.logger import system_logger as logger
|
||||
@@ -184,6 +185,51 @@ def _pid_alive(pid: int) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _get_jsonl_projects_dir(branch_path: Path) -> Path:
|
||||
"""Get Claude's JSONL projects directory for a branch path."""
|
||||
cwd = str(branch_path)
|
||||
encoded = cwd.replace("\\", "-").replace("/", "-").replace(":", "-").replace("_", "-").replace(".", "-")
|
||||
return Path.home() / ".claude" / "projects" / encoded
|
||||
|
||||
|
||||
def _snapshot_jsonl_sizes(projects_dir: Path) -> dict:
|
||||
"""Snapshot current sizes of all JSONL files."""
|
||||
sizes = {}
|
||||
if not projects_dir.exists():
|
||||
return sizes
|
||||
try:
|
||||
for f in projects_dir.glob("*.jsonl"):
|
||||
try:
|
||||
sizes[f.name] = f.stat().st_size
|
||||
except OSError as exc:
|
||||
logger.info("[watchdog.agent] jsonl snapshot stat failed for %s: %s", f.name, exc)
|
||||
except OSError as exc:
|
||||
logger.info("[watchdog.agent] jsonl snapshot glob failed: %s", exc)
|
||||
return sizes
|
||||
|
||||
|
||||
def _has_jsonl_activity(projects_dir: Path, baseline: dict) -> bool:
|
||||
"""Return True if any JSONL file grew or a new file appeared since baseline."""
|
||||
if not projects_dir.exists():
|
||||
return False
|
||||
try:
|
||||
files = list(projects_dir.glob("*.jsonl"))
|
||||
except OSError as exc:
|
||||
logger.info("[watchdog.agent] jsonl activity glob failed: %s", exc)
|
||||
return False
|
||||
for f in files:
|
||||
try:
|
||||
current_size = f.stat().st_size
|
||||
except OSError as exc:
|
||||
logger.info("[watchdog.agent] jsonl activity stat failed for %s: %s", f.name, exc)
|
||||
continue
|
||||
if f.name not in baseline:
|
||||
return current_size > 0
|
||||
if current_size > baseline[f.name]:
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def _read_lock(lock_file: Path) -> dict | None:
|
||||
"""Read lock file, return dict or None on miss/error."""
|
||||
if not lock_file.exists():
|
||||
@@ -195,10 +241,54 @@ def _read_lock(lock_file: Path) -> dict | None:
|
||||
return None
|
||||
|
||||
|
||||
def _classify_exit(branch_path: Path, lock_existed: bool) -> tuple[str, str, int | None]:
|
||||
"""After lock disappears, decide success vs crash.
|
||||
def _check_sent_for_reply(branch_path: Path, dispatch_ts: str | None) -> bool:
|
||||
"""Check if the agent sent any messages after the dispatch timestamp."""
|
||||
sent_dir = branch_path / ".ai_mail.local" / "sent"
|
||||
if not sent_dir.is_dir():
|
||||
return False
|
||||
|
||||
baseline = None
|
||||
if dispatch_ts:
|
||||
for fmt in ("%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M"):
|
||||
try:
|
||||
baseline = datetime.strptime(dispatch_ts, fmt)
|
||||
break
|
||||
except ValueError:
|
||||
logger.info("[watchdog.agent] dispatch_ts %r didn't match fmt %s", dispatch_ts, fmt)
|
||||
continue
|
||||
|
||||
for msg_file in sent_dir.iterdir():
|
||||
if not msg_file.suffix == ".json":
|
||||
continue
|
||||
try:
|
||||
data = json.loads(msg_file.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError) as exc:
|
||||
logger.info("[watchdog.agent] skipping unreadable sent file %s: %s", msg_file.name, exc)
|
||||
continue
|
||||
if baseline is None:
|
||||
return True
|
||||
msg_ts_str = data.get("timestamp", "")
|
||||
for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S"):
|
||||
try:
|
||||
msg_ts = datetime.strptime(msg_ts_str, fmt)
|
||||
if msg_ts >= baseline:
|
||||
return True
|
||||
break
|
||||
except ValueError:
|
||||
logger.info("[watchdog.agent] msg timestamp %r didn't match fmt %s", msg_ts_str, fmt)
|
||||
continue
|
||||
return False
|
||||
|
||||
|
||||
def _classify_exit(
|
||||
branch_path: Path,
|
||||
lock_existed: bool,
|
||||
dispatch_ts: str | None = None,
|
||||
) -> tuple[str, str, int | None]:
|
||||
"""After lock disappears, decide success vs crash vs silent finish.
|
||||
|
||||
Returns (agent_state, reason, exit_code).
|
||||
States: completed_replied, completed_silent, crashed.
|
||||
"""
|
||||
bounce_file = branch_path / ".ai_mail.local" / "last_bounce.json"
|
||||
if bounce_file.exists():
|
||||
@@ -214,9 +304,15 @@ def _classify_exit(branch_path: Path, lock_existed: bool) -> tuple[str, str, int
|
||||
logger.warning("[watchdog.agent] bounce file unreadable: %s", exc)
|
||||
return ("crashed", "agent crashed (bounce file present, unreadable)", None)
|
||||
|
||||
if not lock_existed:
|
||||
return ("completed", "agent finished cleanly", 0)
|
||||
return ("completed", "agent finished cleanly (lock removed)", 0)
|
||||
replied = _check_sent_for_reply(branch_path, dispatch_ts)
|
||||
if replied:
|
||||
reason = "agent finished cleanly (lock removed)" if lock_existed else "agent finished cleanly"
|
||||
return ("completed_replied", reason, 0)
|
||||
|
||||
reason = "agent exited without sending a reply"
|
||||
if lock_existed:
|
||||
reason += " (lock removed)"
|
||||
return ("completed_silent", reason, 0)
|
||||
|
||||
|
||||
def watch_agent(
|
||||
@@ -234,7 +330,7 @@ def watch_agent(
|
||||
|
||||
Returns:
|
||||
dict with keys: woke, reason, elapsed, agent_state, exit_code, agent_id.
|
||||
agent_state is one of: "completed", "crashed", "timeout".
|
||||
agent_state is one of: "completed_replied", "completed_silent", "crashed", "timeout".
|
||||
"""
|
||||
started_at = time.monotonic()
|
||||
_stderr(f"[watchdog.agent] watching {agent_id} (timeout={timeout_seconds}s)")
|
||||
@@ -263,12 +359,13 @@ def watch_agent(
|
||||
lock_file = branch_path / ".ai_mail.local" / ".dispatch.lock"
|
||||
initial_lock = _read_lock(lock_file)
|
||||
initial_pid = initial_lock.get("pid") if initial_lock else None
|
||||
dispatch_ts = initial_lock.get("timestamp") if initial_lock else None
|
||||
lock_existed_initially = initial_lock is not None
|
||||
|
||||
if not lock_existed_initially:
|
||||
_stderr(f"[watchdog.agent] {agent_id}: no active lock — agent already idle")
|
||||
elapsed = int(time.monotonic() - started_at)
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=False)
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=False, dispatch_ts=dispatch_ts)
|
||||
return {
|
||||
"woke": True,
|
||||
"reason": f"no active dispatch ({reason})",
|
||||
@@ -281,6 +378,12 @@ def watch_agent(
|
||||
|
||||
_stderr(f"[watchdog.agent] {agent_id}: lock present, monitor PID={initial_pid}")
|
||||
|
||||
jsonl_dir = _get_jsonl_projects_dir(branch_path)
|
||||
jsonl_baseline = _snapshot_jsonl_sizes(jsonl_dir)
|
||||
last_activity_at = time.monotonic()
|
||||
stall_reported = False
|
||||
stall_threshold = 120.0
|
||||
|
||||
while True:
|
||||
elapsed = time.monotonic() - started_at
|
||||
if elapsed >= timeout_seconds:
|
||||
@@ -299,7 +402,7 @@ def watch_agent(
|
||||
if not lock_file.exists():
|
||||
_stderr(f"[watchdog.agent] {agent_id}: lock removed — agent done")
|
||||
elapsed_int = int(time.monotonic() - started_at)
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=True)
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=True, dispatch_ts=dispatch_ts)
|
||||
logger.info("[watchdog.agent] wake agent_id=%s state=%s elapsed=%s", agent_id, state, elapsed_int)
|
||||
return {
|
||||
"woke": True,
|
||||
@@ -317,8 +420,8 @@ def watch_agent(
|
||||
f"but lock still present — treating as crash"
|
||||
)
|
||||
elapsed_int = int(time.monotonic() - started_at)
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=True)
|
||||
if state == "completed":
|
||||
state, reason, exit_code = _classify_exit(branch_path, lock_existed=True, dispatch_ts=dispatch_ts)
|
||||
if state.startswith("completed"):
|
||||
state = "crashed"
|
||||
reason = f"monitor PID {initial_pid} dead, lock still present"
|
||||
logger.info("[watchdog.agent] wake agent_id=%s state=%s elapsed=%s", agent_id, state, elapsed_int)
|
||||
@@ -332,6 +435,27 @@ def watch_agent(
|
||||
"handle": handle,
|
||||
}
|
||||
|
||||
if _has_jsonl_activity(jsonl_dir, jsonl_baseline):
|
||||
jsonl_baseline = _snapshot_jsonl_sizes(jsonl_dir)
|
||||
last_activity_at = time.monotonic()
|
||||
if stall_reported:
|
||||
_stderr(f"[watchdog.agent] {agent_id}: activity resumed")
|
||||
logger.info("[watchdog.agent] activity resumed agent_id=%s", agent_id)
|
||||
stall_reported = False
|
||||
elif not stall_reported and (time.monotonic() - last_activity_at) >= stall_threshold:
|
||||
idle_secs = int(time.monotonic() - last_activity_at)
|
||||
_stderr(
|
||||
f"[watchdog.agent] {agent_id}: STALLED — no JSONL activity for {idle_secs}s "
|
||||
f"(PID {initial_pid} still alive)"
|
||||
)
|
||||
logger.info(
|
||||
"[watchdog.agent] stall detected agent_id=%s idle=%ss pid=%s",
|
||||
agent_id,
|
||||
idle_secs,
|
||||
initial_pid,
|
||||
)
|
||||
stall_reported = True
|
||||
|
||||
time.sleep(poll_interval)
|
||||
finally:
|
||||
_registry.deregister(handle)
|
||||
|
||||
@@ -353,9 +353,18 @@ def _handle_agent(sub_args: List[str]) -> bool:
|
||||
reason = result.get("reason", "")
|
||||
elapsed = result.get("elapsed", 0)
|
||||
console.print(f"[bold]watchdog agent[/bold] {agent_id} -> state={state} elapsed={elapsed}s reason={reason}")
|
||||
console.print(
|
||||
f'watchdog: {agent_id} stopped (state={state}). Next: drone @ai_mail dispatch {agent_id} "check in" "..."'
|
||||
)
|
||||
if state == "completed_silent":
|
||||
console.print(
|
||||
f"watchdog: {agent_id} stopped (state={state}) -- CHECK DELIVERABLES. "
|
||||
f'Next: drone @ai_mail dispatch {agent_id} "check in" "You finished your last task but did not send a reply. '
|
||||
f'Please reply with your results now via drone @ai_mail email @devpulse."'
|
||||
)
|
||||
elif state == "completed_replied":
|
||||
console.print(f"watchdog: {agent_id} stopped (state={state}). Reply detected — check inbox.")
|
||||
else:
|
||||
console.print(
|
||||
f'watchdog: {agent_id} stopped (state={state}). Next: drone @ai_mail dispatch {agent_id} "check in" "..."'
|
||||
)
|
||||
return True
|
||||
|
||||
|
||||
|
||||
@@ -80,12 +80,12 @@ def test_watch_agent_no_active_lock(monkeypatch, tmp_path):
|
||||
result = agent_handler.watch_agent("@fakebranch", timeout_seconds=5)
|
||||
|
||||
assert result["woke"] is True
|
||||
assert result["agent_state"] == "completed"
|
||||
assert result["agent_state"] == "completed_silent"
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
|
||||
def test_watch_agent_completed_via_lock_removal(monkeypatch, tmp_path):
|
||||
"""Lock present at start, then removed -> wake with state=completed."""
|
||||
"""Lock present at start, then removed -> wake with state=completed_silent (no sent messages)."""
|
||||
branch_path = _build_fake_branch(tmp_path)
|
||||
lock_file = _write_lock(branch_path, pid=os.getpid())
|
||||
monkeypatch.setattr(agent_handler, "_find_repo_root", lambda *a, **kw: tmp_path)
|
||||
@@ -105,9 +105,38 @@ def test_watch_agent_completed_via_lock_removal(monkeypatch, tmp_path):
|
||||
result = agent_handler.watch_agent("@fakebranch", timeout_seconds=5, poll_interval=0.01)
|
||||
|
||||
assert result["woke"] is True
|
||||
assert result["agent_state"] == "completed"
|
||||
assert result["agent_state"] == "completed_silent"
|
||||
assert result["exit_code"] == 0
|
||||
|
||||
|
||||
def test_watch_agent_completed_replied_via_sent_folder(monkeypatch, tmp_path):
|
||||
"""Lock removed + sent message after dispatch -> completed_replied."""
|
||||
branch_path = _build_fake_branch(tmp_path)
|
||||
lock_file = _write_lock(branch_path, pid=os.getpid())
|
||||
monkeypatch.setattr(agent_handler, "_find_repo_root", lambda *a, **kw: tmp_path)
|
||||
|
||||
sent_dir = branch_path / ".ai_mail.local" / "sent"
|
||||
sent_dir.mkdir(parents=True, exist_ok=True)
|
||||
sent_msg = {"to": "@devpulse", "from": "@fakebranch", "subject": "Done", "timestamp": "2026-04-14 00:01:00"}
|
||||
(sent_dir / "reply.json").write_text(json.dumps(sent_msg), encoding="utf-8")
|
||||
|
||||
call_count = {"n": 0}
|
||||
real_sleep = time.sleep
|
||||
|
||||
def fake_sleep(seconds):
|
||||
"""Remove lock on first poll to simulate clean exit."""
|
||||
call_count["n"] += 1
|
||||
if call_count["n"] >= 1:
|
||||
lock_file.unlink(missing_ok=True)
|
||||
real_sleep(0.01)
|
||||
|
||||
monkeypatch.setattr(agent_handler.time, "sleep", fake_sleep)
|
||||
|
||||
result = agent_handler.watch_agent("@fakebranch", timeout_seconds=5, poll_interval=0.01)
|
||||
|
||||
assert result["woke"] is True
|
||||
assert result["agent_state"] == "completed_replied"
|
||||
assert result["exit_code"] == 0
|
||||
assert "clean" in result["reason"].lower() or "finished" in result["reason"].lower()
|
||||
|
||||
|
||||
def test_watch_agent_crashed_via_bounce_file(monkeypatch, tmp_path):
|
||||
@@ -206,7 +235,7 @@ def test_watch_agent_live_dispatch_completes():
|
||||
|
||||
result = agent_handler.watch_agent("@drone", timeout_seconds=300, poll_interval=2.0)
|
||||
assert result["woke"] is True
|
||||
assert result["agent_state"] in ("completed", "crashed")
|
||||
assert result["agent_state"] in ("completed_replied", "completed_silent", "crashed")
|
||||
|
||||
|
||||
@pytest.mark.integration
|
||||
|
||||
@@ -98,7 +98,7 @@
|
||||
"f014": {
|
||||
"path": ".trinity/passport.json",
|
||||
"name": "passport.json",
|
||||
"content_hash": "9b1d3a691272",
|
||||
"content_hash": "3f5a0c2037d2",
|
||||
"has_branch_placeholder": false
|
||||
},
|
||||
"f016": {
|
||||
@@ -155,7 +155,7 @@
|
||||
"content_hash": "a4cf0a8e3b4f",
|
||||
"has_branch_placeholder": false
|
||||
},
|
||||
"f015": {
|
||||
"f026": {
|
||||
"path": "apps/modules/__init__.py",
|
||||
"name": "__init__.py",
|
||||
"content_hash": "e3b0c44298fc",
|
||||
@@ -263,7 +263,7 @@
|
||||
"content_hash": "28e9ae373563",
|
||||
"has_branch_placeholder": false
|
||||
},
|
||||
"f026": {
|
||||
"f015": {
|
||||
"path": "apps/plugins/__init__.py",
|
||||
"name": "__init__.py",
|
||||
"content_hash": "e3b0c44298fc",
|
||||
|
||||
Reference in New Issue
Block a user