#680: cross-platform _is_branch_occupied + _read_session_type (wake.py + daemon.py) — extracted _get_pid_cwd (Linux /proc, macOS lsof -Fn) + _read_session_type_darwin (ps -wwE); macOS occupancy no longer always-False so wake-back won't double-session an interactive branch. Fail-safe on unreadable cwd/env. +11 tests, seedgo 31/31 both files.
This commit is contained in:
@@ -63,6 +63,20 @@ PyPI version — not the changelog header.
|
||||
|
||||
### Fixed
|
||||
|
||||
- **Interactive-occupancy detection is now cross-platform — the wake-back guard
|
||||
no longer goes blind on macOS (issue #680).** `_is_branch_occupied()` and
|
||||
`_read_session_type()` (duplicated in `dispatch/wake.py` and `dispatch/daemon.py`)
|
||||
read `/proc/{pid}/cwd` + `/proc/{pid}/environ`, which do not exist on macOS —
|
||||
so occupancy always resolved `False` there and an external wake-back could spawn
|
||||
a *second* Claude session on an already-interactive branch (double-session,
|
||||
weakening the TDPLAN-0012/#678 interactive-dispatcher guard). The per-PID cwd
|
||||
and session-type probes are now extracted into platform helpers: `_get_pid_cwd`
|
||||
(Linux `/proc` readlink, macOS `lsof -a -p PID -d cwd -Fn`) and
|
||||
`_read_session_type_darwin` (`ps -p PID -wwE`), applied identically in both
|
||||
files. Fail-safe: an unreadable cwd/env logs at info and continues — never
|
||||
crashes the wake path. +11 tests (macOS cwd, macOS session type, zombie,
|
||||
unsupported platform), seedgo 31/31 on both files.
|
||||
|
||||
- **SubagentStop gate no longer runs its ~600ms seedgo check on every internal
|
||||
turn (issue #606).** Claude Code creates an internal agent per response turn
|
||||
with an empty `agent_type`, so the `subagent_gate` handler was firing its full
|
||||
|
||||
@@ -23,7 +23,7 @@
|
||||
{
|
||||
"file": "apps/handlers/dispatch/daemon.py",
|
||||
"standard": "deep_nesting",
|
||||
"reason": "3 functions: check_inbox_for_dispatch() depth 4 (priority scanning with business logic), run_daemon() depth 4 (main daemon loop), _check_lock() depth 4 (lock validation + PID liveness + cleanup)"
|
||||
"reason": "run_daemon() depth 5 (main daemon loop with retry + signal handling)"
|
||||
},
|
||||
{
|
||||
"file": "apps/handlers/dispatch/wake.py",
|
||||
|
||||
@@ -514,18 +514,74 @@ def is_protected_branch(branch_email: str) -> bool:
|
||||
return branch_email == "@devpulse"
|
||||
|
||||
|
||||
def _read_session_type(pid_str: str) -> str:
|
||||
"""Read AIPASS_SESSION_TYPE from /proc/{pid}/environ. Returns 'interactive' if unset."""
|
||||
if sys.platform != "linux":
|
||||
return "interactive"
|
||||
def _get_pid_cwd(pid_str: str) -> Optional[str]:
|
||||
"""Get the cwd of a process. Cross-platform: Linux /proc, macOS lsof."""
|
||||
if sys.platform == "linux":
|
||||
try:
|
||||
return os.readlink(f"/proc/{pid_str}/cwd")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("[daemon] Cannot read cwd for PID %s", pid_str)
|
||||
return None
|
||||
if sys.platform == "darwin":
|
||||
return _get_pid_cwd_darwin(pid_str)
|
||||
logger.info("[daemon] Cannot determine cwd for PID %s on %s", pid_str, sys.platform)
|
||||
return None
|
||||
|
||||
|
||||
def _get_pid_cwd_darwin(pid_str: str) -> Optional[str]:
|
||||
"""macOS: get process cwd via lsof."""
|
||||
try:
|
||||
with open(f"/proc/{pid_str}/environ", "rb") as f:
|
||||
data = f.read()
|
||||
for entry in data.split(b"\0"):
|
||||
if entry.startswith(b"AIPASS_SESSION_TYPE="):
|
||||
return entry.split(b"=", 1)[1].decode("utf-8")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("Cannot read session type for PID %s", pid_str)
|
||||
result = subprocess.run(
|
||||
["lsof", "-a", "-p", pid_str, "-d", "cwd", "-Fn"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
)
|
||||
except (subprocess.SubprocessError, OSError):
|
||||
logger.info("[daemon] Cannot read cwd for PID %s on macOS", pid_str)
|
||||
return None
|
||||
if result.returncode != 0:
|
||||
return None
|
||||
for line in result.stdout.strip().split("\n"):
|
||||
if line.startswith("n/"):
|
||||
return line[1:]
|
||||
return None
|
||||
|
||||
|
||||
def _read_session_type(pid_str: str) -> str:
|
||||
"""Read AIPASS_SESSION_TYPE from process environment. Returns 'interactive' if unset."""
|
||||
if sys.platform == "linux":
|
||||
try:
|
||||
with open(f"/proc/{pid_str}/environ", "rb") as f:
|
||||
data = f.read()
|
||||
for entry in data.split(b"\0"):
|
||||
if entry.startswith(b"AIPASS_SESSION_TYPE="):
|
||||
return entry.split(b"=", 1)[1].decode("utf-8")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("[daemon] Cannot read session type for PID %s", pid_str)
|
||||
return "interactive"
|
||||
if sys.platform == "darwin":
|
||||
return _read_session_type_darwin(pid_str)
|
||||
return "interactive"
|
||||
|
||||
|
||||
def _read_session_type_darwin(pid_str: str) -> str:
|
||||
"""macOS: read AIPASS_SESSION_TYPE from ps environment output."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["ps", "-p", pid_str, "-wwE", "-o", "command="],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
)
|
||||
except (subprocess.SubprocessError, OSError):
|
||||
logger.info("[daemon] Cannot read session type for PID %s on macOS", pid_str)
|
||||
return "interactive"
|
||||
if result.returncode != 0:
|
||||
return "interactive"
|
||||
for token in result.stdout.split():
|
||||
if token.startswith("AIPASS_SESSION_TYPE="):
|
||||
return token.split("=", 1)[1]
|
||||
return "interactive"
|
||||
|
||||
|
||||
@@ -534,35 +590,25 @@ _NON_BLOCKING_SESSION_TYPES = {"dispatched", "daemon"}
|
||||
|
||||
|
||||
def _is_branch_occupied(branch_path: Path) -> bool:
|
||||
"""
|
||||
Check if an interactive Claude session is running in this branch.
|
||||
|
||||
Only interactive sessions block dispatch. Telegram, dispatched, and daemon
|
||||
sessions are idle/background and should not prevent new agent spawns.
|
||||
"""
|
||||
resolved = branch_path.resolve()
|
||||
"""Check if an interactive Claude session is running in this branch."""
|
||||
resolved = str(branch_path.resolve())
|
||||
try:
|
||||
result = subprocess.run(["pgrep", "-x", "claude"], capture_output=True, text=True, timeout=5)
|
||||
if result.returncode != 0:
|
||||
return False
|
||||
|
||||
for pid_str in result.stdout.strip().split("\n"):
|
||||
pid_str = pid_str.strip()
|
||||
if not pid_str:
|
||||
continue
|
||||
try:
|
||||
if sys.platform != "linux":
|
||||
continue
|
||||
cwd = os.readlink(f"/proc/{pid_str}/cwd")
|
||||
if Path(cwd).resolve() == resolved:
|
||||
session_type = _read_session_type(pid_str)
|
||||
if session_type not in _NON_BLOCKING_SESSION_TYPES:
|
||||
return True
|
||||
except (OSError, PermissionError, ValueError):
|
||||
logger.info("Cannot read cwd for PID %s", pid_str)
|
||||
cwd = _get_pid_cwd(pid_str)
|
||||
if cwd is None:
|
||||
continue
|
||||
if str(Path(cwd).resolve()) == resolved:
|
||||
session_type = _read_session_type(pid_str)
|
||||
if session_type not in _NON_BLOCKING_SESSION_TYPES:
|
||||
return True
|
||||
except Exception:
|
||||
logger.info("Failed to check branch occupancy for %s", branch_path)
|
||||
logger.info("[daemon] Failed to check branch occupancy for %s", branch_path)
|
||||
return False
|
||||
|
||||
|
||||
|
||||
@@ -233,18 +233,74 @@ def _load_config() -> dict:
|
||||
return config
|
||||
|
||||
|
||||
def _read_session_type(pid_str: str) -> str:
|
||||
"""Read AIPASS_SESSION_TYPE from /proc/{pid}/environ. Returns 'interactive' if unset."""
|
||||
if sys.platform != "linux":
|
||||
return "interactive"
|
||||
def _get_pid_cwd(pid_str: str) -> Optional[str]:
|
||||
"""Get the cwd of a process. Cross-platform: Linux /proc, macOS lsof."""
|
||||
if sys.platform == "linux":
|
||||
try:
|
||||
return os.readlink(f"/proc/{pid_str}/cwd")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("[wake] Cannot read cwd for PID %s", pid_str)
|
||||
return None
|
||||
if sys.platform == "darwin":
|
||||
return _get_pid_cwd_darwin(pid_str)
|
||||
logger.info("[wake] Cannot determine cwd for PID %s on %s", pid_str, sys.platform)
|
||||
return None
|
||||
|
||||
|
||||
def _get_pid_cwd_darwin(pid_str: str) -> Optional[str]:
|
||||
"""macOS: get process cwd via lsof."""
|
||||
try:
|
||||
with open(f"/proc/{pid_str}/environ", "rb") as f:
|
||||
data = f.read()
|
||||
for entry in data.split(b"\0"):
|
||||
if entry.startswith(b"AIPASS_SESSION_TYPE="):
|
||||
return entry.split(b"=", 1)[1].decode("utf-8")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("[wake] Cannot read session type for PID %s", pid_str)
|
||||
result = subprocess.run(
|
||||
["lsof", "-a", "-p", pid_str, "-d", "cwd", "-Fn"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
)
|
||||
except (subprocess.SubprocessError, OSError):
|
||||
logger.info("[wake] Cannot read cwd for PID %s on macOS", pid_str)
|
||||
return None
|
||||
if result.returncode != 0:
|
||||
return None
|
||||
for line in result.stdout.strip().split("\n"):
|
||||
if line.startswith("n/"):
|
||||
return line[1:]
|
||||
return None
|
||||
|
||||
|
||||
def _read_session_type(pid_str: str) -> str:
|
||||
"""Read AIPASS_SESSION_TYPE from process environment. Returns 'interactive' if unset."""
|
||||
if sys.platform == "linux":
|
||||
try:
|
||||
with open(f"/proc/{pid_str}/environ", "rb") as f:
|
||||
data = f.read()
|
||||
for entry in data.split(b"\0"):
|
||||
if entry.startswith(b"AIPASS_SESSION_TYPE="):
|
||||
return entry.split(b"=", 1)[1].decode("utf-8")
|
||||
except (OSError, PermissionError):
|
||||
logger.info("[wake] Cannot read session type for PID %s", pid_str)
|
||||
return "interactive"
|
||||
if sys.platform == "darwin":
|
||||
return _read_session_type_darwin(pid_str)
|
||||
return "interactive"
|
||||
|
||||
|
||||
def _read_session_type_darwin(pid_str: str) -> str:
|
||||
"""macOS: read AIPASS_SESSION_TYPE from ps environment output."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["ps", "-p", pid_str, "-wwE", "-o", "command="],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=5,
|
||||
)
|
||||
except (subprocess.SubprocessError, OSError):
|
||||
logger.info("[wake] Cannot read session type for PID %s on macOS", pid_str)
|
||||
return "interactive"
|
||||
if result.returncode != 0:
|
||||
return "interactive"
|
||||
for token in result.stdout.split():
|
||||
if token.startswith("AIPASS_SESSION_TYPE="):
|
||||
return token.split("=", 1)[1]
|
||||
return "interactive"
|
||||
|
||||
|
||||
@@ -263,17 +319,13 @@ def _is_branch_occupied(branch_path: Path) -> bool:
|
||||
pid_str = pid_str.strip()
|
||||
if not pid_str:
|
||||
continue
|
||||
try:
|
||||
if sys.platform != "linux":
|
||||
continue
|
||||
cwd = os.readlink(f"/proc/{pid_str}/cwd")
|
||||
if str(Path(cwd).resolve()) == resolved:
|
||||
session_type = _read_session_type(pid_str)
|
||||
if session_type not in _NON_BLOCKING_SESSION_TYPES:
|
||||
return True
|
||||
except (OSError, PermissionError, ValueError):
|
||||
logger.info("[wake] Cannot read cwd for PID %s", pid_str)
|
||||
cwd = _get_pid_cwd(pid_str)
|
||||
if cwd is None:
|
||||
continue
|
||||
if str(Path(cwd).resolve()) == resolved:
|
||||
session_type = _read_session_type(pid_str)
|
||||
if session_type not in _NON_BLOCKING_SESSION_TYPES:
|
||||
return True
|
||||
except (subprocess.SubprocessError, OSError):
|
||||
logger.info("[wake] Failed to check branch occupancy")
|
||||
return False
|
||||
@@ -313,18 +365,23 @@ def _check_pid_alive(pid: int) -> bool:
|
||||
except OSError as exc:
|
||||
logger.warning("[wake] PID %s os.kill error (assuming dead): %s", pid, exc)
|
||||
return False
|
||||
if sys.platform == "linux":
|
||||
try:
|
||||
with open(f"/proc/{pid}/status", "r") as f:
|
||||
for line in f:
|
||||
if line.startswith("State:"):
|
||||
return "Z" not in line
|
||||
except FileNotFoundError as exc:
|
||||
logger.warning("[wake] PID %s /proc not found: %s", pid, exc)
|
||||
return False
|
||||
if sys.platform == "linux" and _is_zombie_linux(pid):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _is_zombie_linux(pid: int) -> bool:
|
||||
"""Return True if PID is a zombie (Linux /proc/status check)."""
|
||||
try:
|
||||
with open(f"/proc/{pid}/status", "r") as f:
|
||||
for line in f:
|
||||
if line.startswith("State:"):
|
||||
return "Z" in line
|
||||
except FileNotFoundError as exc:
|
||||
logger.warning("[wake] PID %s /proc not found: %s", pid, exc)
|
||||
return False
|
||||
|
||||
|
||||
def _spawn_in_systemd_scope(monitor_cmd, branch_path, spawn_env, branch_email, lock_file_path, custom_message, status):
|
||||
"""Spawn monitor in its own systemd unit to survive cgroup cleanup (td-48).
|
||||
|
||||
|
||||
@@ -20,7 +20,11 @@ from aipass.ai_mail.apps.handlers.dispatch.wake import (
|
||||
_read_json,
|
||||
_check_lock,
|
||||
_check_pid_alive,
|
||||
_get_pid_cwd,
|
||||
_get_pid_cwd_darwin,
|
||||
_read_session_type,
|
||||
_read_session_type_darwin,
|
||||
_is_zombie_linux,
|
||||
_clean_zombies,
|
||||
_find_claude_bin,
|
||||
resolve_branch,
|
||||
@@ -201,11 +205,112 @@ def test_read_session_type_not_set(monkeypatch, tmp_path):
|
||||
|
||||
|
||||
def test_read_session_type_non_linux(monkeypatch):
|
||||
"""Non-linux platform returns 'interactive' immediately."""
|
||||
monkeypatch.setattr("sys.platform", "darwin")
|
||||
"""Non-linux, non-darwin platform returns 'interactive' immediately."""
|
||||
monkeypatch.setattr("sys.platform", "win32")
|
||||
assert _read_session_type("999") == "interactive"
|
||||
|
||||
|
||||
def test_read_session_type_darwin_found(monkeypatch):
|
||||
"""macOS: reads AIPASS_SESSION_TYPE from ps -wwE output."""
|
||||
monkeypatch.setattr("sys.platform", "darwin")
|
||||
|
||||
class FakeResult:
|
||||
returncode = 0
|
||||
stdout = "/usr/bin/claude AIPASS_SESSION_TYPE=dispatched HOME=/Users/u"
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", lambda *a, **kw: FakeResult())
|
||||
assert _read_session_type("123") == "dispatched"
|
||||
|
||||
|
||||
def test_read_session_type_darwin_not_set(monkeypatch):
|
||||
"""macOS: missing env var returns 'interactive'."""
|
||||
monkeypatch.setattr("sys.platform", "darwin")
|
||||
|
||||
class FakeResult:
|
||||
returncode = 0
|
||||
stdout = "/usr/bin/claude HOME=/Users/u"
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", lambda *a, **kw: FakeResult())
|
||||
assert _read_session_type("456") == "interactive"
|
||||
|
||||
|
||||
def test_read_session_type_darwin_ps_failure(monkeypatch):
|
||||
"""macOS: ps failure returns 'interactive'."""
|
||||
assert _read_session_type_darwin("999") == "interactive"
|
||||
|
||||
|
||||
# --- _get_pid_cwd tests ------------------------------------------------
|
||||
|
||||
|
||||
def test_get_pid_cwd_linux(monkeypatch, tmp_path):
|
||||
"""Linux: reads /proc/{pid}/cwd via readlink."""
|
||||
monkeypatch.setattr("sys.platform", "linux")
|
||||
target = str(tmp_path / "project")
|
||||
monkeypatch.setattr(os, "readlink", lambda p: target)
|
||||
assert _get_pid_cwd("100") == target
|
||||
|
||||
|
||||
def test_get_pid_cwd_linux_oserror(monkeypatch):
|
||||
"""Linux: OSError returns None."""
|
||||
monkeypatch.setattr("sys.platform", "linux")
|
||||
monkeypatch.setattr(os, "readlink", lambda p: (_ for _ in ()).throw(OSError("no proc")))
|
||||
assert _get_pid_cwd("100") is None
|
||||
|
||||
|
||||
def test_get_pid_cwd_darwin(monkeypatch, tmp_path):
|
||||
"""macOS: reads cwd via lsof."""
|
||||
monkeypatch.setattr("sys.platform", "darwin")
|
||||
target = str(tmp_path / "project")
|
||||
|
||||
class FakeResult:
|
||||
returncode = 0
|
||||
stdout = f"p100\nn{target}\n"
|
||||
|
||||
monkeypatch.setattr(subprocess, "run", lambda *a, **kw: FakeResult())
|
||||
assert _get_pid_cwd("100") == target
|
||||
|
||||
|
||||
def test_get_pid_cwd_darwin_failure(monkeypatch):
|
||||
"""macOS: lsof failure returns None."""
|
||||
assert _get_pid_cwd_darwin("999") is None
|
||||
|
||||
|
||||
def test_get_pid_cwd_unsupported_platform(monkeypatch):
|
||||
"""Unsupported platform returns None."""
|
||||
monkeypatch.setattr("sys.platform", "win32")
|
||||
assert _get_pid_cwd("100") is None
|
||||
|
||||
|
||||
# --- _is_zombie_linux tests --------------------------------------------
|
||||
|
||||
|
||||
def test_is_zombie_linux_not_zombie(monkeypatch, tmp_path):
|
||||
"""Non-zombie process returns False."""
|
||||
status_file = tmp_path / "status"
|
||||
status_file.write_text("Name:\tclaude\nState:\tS (sleeping)\nPid:\t42\n")
|
||||
monkeypatch.setattr(
|
||||
"builtins.open",
|
||||
_fake_open_factory(str(status_file), {"/proc/42/status": str(status_file)}),
|
||||
)
|
||||
assert _is_zombie_linux(42) is False
|
||||
|
||||
|
||||
def test_is_zombie_linux_zombie(monkeypatch, tmp_path):
|
||||
"""Zombie process returns True."""
|
||||
status_file = tmp_path / "status"
|
||||
status_file.write_text("Name:\tclaude\nState:\tZ (zombie)\nPid:\t42\n")
|
||||
monkeypatch.setattr(
|
||||
"builtins.open",
|
||||
_fake_open_factory(str(status_file), {"/proc/42/status": str(status_file)}),
|
||||
)
|
||||
assert _is_zombie_linux(42) is True
|
||||
|
||||
|
||||
def test_is_zombie_linux_no_proc(monkeypatch):
|
||||
"""Missing /proc entry returns False (not zombie, just gone)."""
|
||||
assert _is_zombie_linux(99999) is False
|
||||
|
||||
|
||||
# --- _check_lock tests ------------------------------------------------
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user