feat: runaway-log detection + escalation (DPLAN-0242, agent-designed) + citizen wake-back. Prax-led three-branch build via TDPLAN-0013: prax rate_tracker (disk-persisted volume detection, WARNING >100 l/min 2min / CRITICAL >10 l/s 1min, 4th monitor thread) + drone @prax log-health; trigger runaway_log_detected event + handler (per-file 30min cooldown independent of medic breaker, UNKNOWN->prax, writes .aipass/alerts.json); hooks persistent_alert banner + drone @hooks dismiss (devpulse fixed nearest-.aipass path bug in both + wired settings.json - registration is not deployment). Live-fire proven: planted 240 l/min storm -> detect 257 l/min -> event -> dispatch -> @aipass autonomous no-action triage -> banner -> dismiss. ai_mail wake ruling: owner-gate removed (citizen wake-back live-proven), devpulse structurally unwakeable (manager check all paths), self-wake loop found+fixed same night (guard + senderless wake-back sessions). Navmap comms section, 4 branch READMEs + root README + CHANGELOG. ~77 new tests, suites green: prax 1028, trigger 619, hooks 1071, ai_mail 765

This commit is contained in:
AIOSAI
2026-07-15 00:05:02 -07:00
parent de109846bf
commit 81658ce0ea
25 changed files with 2789 additions and 182 deletions
+5
View File
@@ -8,6 +8,11 @@
"handler": "aipass.hooks.apps.handlers.security.presence_gate.handle",
"matcher": ""
},
"persistent_alert": {
"enabled": true,
"handler": "aipass.hooks.apps.handlers.prompt.persistent_alert.handle",
"matcher": ""
},
"identity_injector": {
"enabled": true,
"handler": "aipass.hooks.apps.handlers.prompt.identity.handle",
+9
View File
@@ -70,6 +70,15 @@ drone @git status / diff / log # read-only git awareness
drone @memory search "query" # recall archived context
```
# Talking to other agents
Citizens dispatch each other directly — allowed and expected, no permission needed. Pick by one question: does the recipient need to ACT?
- Need an answer, input, or work from them → `dispatch` (send + wake). A sleeping agent never reads plain email — a question sent as `email` stalls unread.
- FYI only (status, steering an agent already awake) → `email` (no wake).
- Replies don't wake either — but wake-back does: when an agent you dispatched completes, YOU are woken automatically. Team mission: the lead dispatches each phase BEFORE sleeping, the worker replies normally, wake-back brings the lead back to verify and hand off the next phase.
- Exception — managers (`citizen_class: manager`, e.g. @devpulse) are never dispatched: they hold interactive sessions with the user, so the wake is blocked. `email` them — the mail lands and they see it live.
Always reply to dispatches — reply auto-closes. No silent completions.
# Plans — flow
+45
View File
@@ -40,6 +40,51 @@ PyPI version — not the changelog header.
### Added
- **Runaway-log detection + escalation — designed and built by the agents
themselves.** Patrick's mission brief went to @prax as lead ("I don't want it
to be you" — devpulse relayed requirements, not a design): prax researched,
collaborated with @hooks and @trigger by mail, wrote DPLAN-0242, and ran the
build as TDPLAN-0013 across three branches. The system: @prax `rate_tracker`
watches every log in `system_logs/` for volume (not content — orthogonal to
medic), disk-persisted state, WARNING at >100 lines/min sustained 2 min /
CRITICAL at >10 lines/sec 1 min, per-file suppression, runs as a 4th monitor
thread, plus `drone @prax log-health` for an at-a-glance rate overview.
@trigger registers the new `runaway_log_detected` event and dispatches to the
responsible branch with a per-file 30-min cooldown deliberately independent
of medic's circuit breaker (a storm can't silence both systems), UNKNOWN
attribution falls back to @prax, and every alert is written to
`.aipass/alerts.json`. @hooks `persistent_alert` injects an advisory banner
into every agent's prompt until the alert is fixed or dismissed
(`drone @hooks dismiss <id>`) — general-purpose, any agent can raise alerts.
Devpulse verification found and fixed the last-mile gaps: both hooks pieces
stopped at the first `.aipass/` dir walking up (every branch has one — the
banner could never render), and hooks.json registration alone isn't
deployment — the handler needed manual wiring into `~/.claude/settings.json`
(agents can't edit it; documented for future handlers). Live-fire acceptance:
a planted 240 lines/min storm was detected at 257 lines/min sustained 120s →
event → dispatch → **@aipass woke autonomously, root-caused the test writer
down to its PID and loop shape, triaged no-action** → alerts.json → banner
renders → dismiss clears. ~77 new tests across four branches, suites green
(prax 1028, trigger 619, hooks 1071), seedgo 98–100%.
- **Citizens wake each other freely — wake-back for everyone, devpulse
unwakeable by design.** Patrick's ruling after two team-mission stalls in one
evening (prax emailed sleeping collaborators; trigger replied instead of
dispatching back — replies never wake, and wake-back was owner-gated so
agent-to-agent dispatch never woke the sender). @ai_mail removed the
owner-gate: any citizen sender is woken when its dispatched agent completes
(proven live: "@trigger woken after @aipass completed" — first citizen
wake-back ever). @devpulse is now structurally unwakeable via a
`citizen_class: manager` check on every wake path — mail always lands, wake
always skips, no longer dependent on an interactive session happening to be
open. The gate removal exposed a self-wake loop within minutes (wake-back
sessions were attributed to @ai_mail as sender, so ai_mail kept waking
itself; the depth cap stopped it after one cycle) — fixed the same night:
self-wake guard + wake-back sessions carry no sender, so chains terminate at
the original dispatcher. 765 ai_mail tests green. The navmap gained a
"Talking to other agents" section: dispatch vs email semantics, team-relay
discipline, and the manager exception.
- **Medic is back on — and the loop is proven live.** Off since 2026-05-10 (a
pytest fixture storm flooded the error registry; the off switch was pulled to
stop the noise and forgotten for 65 days). Three fixes made re-enable safe:
+4 -4
View File
@@ -165,7 +165,7 @@ devpulse (orchestrator)
├── aipass — concierge + onboarding (aipass init, doctor, profile)
├── drone — command routing + @agent resolution
├── seedgo — automated quality standards
├── prax — real-time monitoring across all agents
├── prax — real-time monitoring + runaway-log detection across all agents
├── ai_mail — agent-to-agent communication + task dispatch
├── flow — plan lifecycle, templates, auto-archival
├── spawn — creates new agents anywhere on your filesystem
@@ -203,9 +203,9 @@ These agents work on the **same filesystem, same project, same time** — no san
| Agent | Role |
|-------|------|
| [**seedgo**](src/aipass/seedgo/README.md) | Automated quality standards, enforced across all agents |
| [**prax**](src/aipass/prax/README.md) | Real-time monitoring, logs, dashboards |
| [**prax**](src/aipass/prax/README.md) | Real-time monitoring, logs, dashboards, runaway-log detection |
| [**flow**](src/aipass/flow/README.md) | Plan lifecycle — multiple template types, auto-archival, vector verification |
| [**hooks**](src/aipass/hooks/README.md) | Hook engine — per-project config, sound control, event dispatch |
| [**hooks**](src/aipass/hooks/README.md) | Hook engine — per-project config, sound control, event dispatch, persistent alerts |
| [**trigger**](src/aipass/trigger/README.md) | Event-driven automation + self-healing |
| [**cli**](src/aipass/cli/README.md) | Terminal formatting and rich output |
| [**backup**](src/aipass/backup/README.md) | Local-first backups — snapshots, versioning, restore (optional Google Drive sync) |
@@ -242,7 +242,7 @@ The installer (`./aipass install`, powered by setup.sh) auto-detects which CLIs
| Metric | Value |
|--------|-------|
| Version | See [git tags](https://github.com/AIOSAI/AIPass/tags) |
| Agents | 13 core + user-created |
| Agents | 17 core + user-created |
| Quality | Automated standards enforced across every agent |
| Coverage | [![codecov](https://codecov.io/gh/AIOSAI/AIPass/graph/badge.svg)](https://codecov.io/gh/AIOSAI/AIPass) — 75% minimum, CI-gated |
| Tests | Extensive — every agent ships its own suite |
+5 -2
View File
@@ -62,19 +62,22 @@ The `dispatch` command sends an email and wakes the target branch in one step. D
### Wake Pipeline
1. `dispatch.py` orchestrates: send email via `send_to_single()`, then wake via `wake_branch()`
2. `wake.py` resolves the branch from the registry, finds the `claude` binary, spawns a subprocess
2. `wake.py` resolves the branch from the registry, checks `citizen_class` (managers are mail-only — wake skips), finds the `claude` binary, spawns a subprocess
3. `dispatch_monitor.py` wraps the claude process with safety features:
- **Startup health check** — monitors JSONL session files for 90s, kills if no activity
- **Auto-retry** — 3 strikes: attempt 1+2 resume, attempt 3 fresh (new session)
- **Bounce email** — on final failure, sends error report back to sender
- **Lock cleanup** — removes `.dispatch.lock` when agent exits
4. After wake, `_spawn_watchdog()` auto-launches `drone @devpulse watchdog agent @target` as a detached background process
- **Wake-back** — on agent exit, wakes the original sender so they can process the result. Wake-back sessions carry an empty sender, so chains terminate at the original dispatcher
### Safety Limits
- PID-based locking prevents concurrent agents per branch (`.dispatch.lock`)
- Max turns per wake, max dispatches per branch per day
- `WAKE_BLOCKLIST` protects `@devpulse` from cross-branch manual wakes
- **Manager structural block** — branches with `citizen_class: "manager"` in their passport (e.g. `@devpulse`) are unwakeable on all wake paths. Mail delivers, wake skips
- **Self-wake guard** — if sender equals target, wake-back is skipped (prevents self-loops)
- **Chain termination** — wake-back sessions carry an empty sender, so the chain always stops at the original dispatcher
- `dispatch_monitor.py` strips `AIPASS_CALLER_*` env vars to prevent parent context leaking into agent identity
- `AIPASS_BRANCH_NAME` env var set in spawn_env for CWD-independent identity
@@ -86,28 +86,22 @@ MAX_WAKE_DEPTH = 3
def _wake_sender(sender: str, branch_email: str, exit_code: int, lock_file: str) -> str:
"""Wake the dispatcher back after target completion.
Wake-back is owner-only: only the project owner (sealed registry)
gets woken. Non-owners silently skipped.
Any citizen sender gets woken back (same availability checks as
normal wake — interactive session, active lock, depth cap).
Returns a result tag for the dispatch_wake.log:
success, blocked_occupied, blocked_locked, blocked_depth,
skipped_sender, skipped_not_owner, failed
skipped_sender, skipped_self, failed
"""
if not sender or not sender.strip():
logger.info("[monitor] Wake-back skipped — no sender")
return "skipped_sender"
normalized = f"@{sender.lstrip('@').lower()}"
try:
from aipass.spawn.apps.handlers.registry import is_owner
except ImportError:
logger.warning("[monitor] Wake-back skipped — is_owner import failed")
return "failed"
if not is_owner(normalized):
logger.info("[monitor] Wake-back skipped — sender %s is not project owner", sender)
return "skipped_not_owner"
normalized_sender = f"@{sender.lstrip('@').lower()}"
normalized_target = f"@{branch_email.lstrip('@').lower()}"
if normalized_sender == normalized_target:
logger.info("[monitor] Wake-back skipped — sender %s is the completed agent (self-wake)", sender)
return "skipped_self"
depth = int(os.environ.get("AIPASS_WAKE_DEPTH", "0"))
if depth >= MAX_WAKE_DEPTH:
@@ -118,7 +112,7 @@ def _wake_sender(sender: str, branch_email: str, exit_code: int, lock_file: str)
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
os.environ["AIPASS_WAKE_DEPTH"] = str(depth + 1)
wake_status, success = wake_branch(sender, auto=True, sender="@ai_mail")
wake_status, success = wake_branch(sender, auto=True, sender="")
if success:
logger.info("[monitor] Wake-back: %s woken after %s completed (exit %d)", sender, branch_email, exit_code)
@@ -542,14 +542,27 @@ def wake_branch(
branch_path, email = result
status.ok("resolve", f"{email} → {branch_path}")
# Step 3: Zombie check (pre-flight)
# Step 3: Manager check — managers are never woken, mail only
passport_file = branch_path / ".trinity" / "passport.json"
try:
with open(passport_file, "r", encoding="utf-8") as f:
passport = json.load(f)
citizen_class = passport.get("identity", {}).get("citizen_class", "")
if citizen_class == "manager":
status.info("manager", f"{email} is a manager — mail only, wake skipped")
logger.info("[wake] %s is citizen_class=manager — wake skipped, mail delivered", email)
return status, True
except (FileNotFoundError, json.JSONDecodeError, OSError) as exc:
logger.info("[wake] Could not read passport for %s: %s", email, exc)
# Step 4: Zombie check (pre-flight)
zombie_count = _clean_zombies()
if zombie_count > 0:
status.warn("zombies", f"{zombie_count} zombie Claude process(es) detected")
else:
status.ok("pre-flight", "No zombie processes")
# Step 4: Lock check
# Step 5: Lock check
existing = _check_lock(branch_path)
if existing is not None:
pid = existing.get("pid", "?")
@@ -565,7 +578,7 @@ def wake_branch(
status.ok("lock", "No active lock — agent is sleeping")
# Step 5: Occupancy check
# Step 6: Occupancy check
if _is_branch_occupied(branch_path):
status.warn("occupancy", f"Interactive Claude session in {branch_path}")
status.fail("blocked", "Cannot spawn — interactive session running")
@@ -574,7 +587,7 @@ def wake_branch(
status.ok("occupancy", "No interactive session")
# Step 6: Build spawn command
# Step 7: Build spawn command
config = _load_config()
max_turns = config.get("max_turns_per_wake", 100)
+35 -124
View File
@@ -1840,15 +1840,7 @@ finally:
class TestWakeSender:
"""_wake_sender guards and owner-allowlist dispatch."""
@pytest.fixture(autouse=True)
def _mock_is_owner(self, monkeypatch):
"""Default: is_owner returns False (non-owner). Tests override as needed."""
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=False),
)
"""_wake_sender guards and wake-back dispatch."""
def test_skips_empty_sender(self, monkeypatch):
"""Empty sender returns skipped_sender."""
@@ -1862,44 +1854,22 @@ class TestWakeSender:
result = _wake_sender(" ", "@target", 0, "/fake/lock")
assert result == "skipped_sender"
def test_skips_non_owner_sender(self, monkeypatch):
"""Non-owner sender returns skipped_not_owner."""
def test_skips_self_wake(self, monkeypatch):
"""Sender equal to completed agent returns skipped_self."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=False),
)
result = _wake_sender("@someagent", "@target", 0, "/fake/lock")
assert result == "skipped_not_owner"
result = _wake_sender("@trigger", "@trigger", 0, "/fake/lock")
assert result == "skipped_self"
def test_skips_ai_mail_when_not_owner(self, monkeypatch):
"""@ai_mail is not owner — skipped."""
def test_skips_self_wake_case_insensitive(self, monkeypatch):
"""Self-wake guard is case-insensitive."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=False),
)
result = _wake_sender("@ai_mail", "@target", 0, "/fake/lock")
assert result == "skipped_not_owner"
result = _wake_sender("Trigger", "@TRIGGER", 0, "/fake/lock")
assert result == "skipped_self"
def test_skips_human_when_not_owner(self, monkeypatch):
"""@human is not owner — skipped."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=False),
)
result = _wake_sender("@human", "@target", 0, "/fake/lock")
assert result == "skipped_not_owner"
def test_owner_passes_guard(self, monkeypatch):
"""Owner sender passes the is_owner guard and reaches wake_branch."""
def test_wake_back_carries_empty_sender(self, monkeypatch):
"""Wake-back session carries empty sender to terminate the chain."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
mock_status = MagicMock()
mock_status.summary = "ok"
mock_wake = MagicMock(return_value=(mock_status, True))
@@ -1907,67 +1877,42 @@ class TestWakeSender:
"aipass.ai_mail.apps.handlers.dispatch.wake.wake_branch",
mock_wake,
)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
_wake_sender("@prax", "@trigger", 0, "/fake/lock")
mock_wake.assert_called_once_with("@prax", auto=True, sender="")
def test_any_citizen_reaches_wake_branch(self, monkeypatch):
"""Any citizen sender reaches wake_branch."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
mock_status = MagicMock()
mock_status.summary = "ok"
mock_wake = MagicMock(return_value=(mock_status, True))
monkeypatch.setattr(
"aipass.ai_mail.apps.handlers.dispatch.wake.wake_branch",
mock_wake,
)
result = _wake_sender("@prax", "@target", 0, "/fake/lock")
assert result == "success"
mock_wake.assert_called_once()
def test_is_owner_called_with_normalized_sender(self, monkeypatch):
"""is_owner receives normalized @-prefixed lowercase sender."""
monkeypatch.setattr(mod, "logger", MagicMock())
mock_is_owner = MagicMock(return_value=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
mock_is_owner,
)
_wake_sender("DevPulse", "@target", 0, "/fake/lock")
mock_is_owner.assert_called_once_with("@devpulse")
def test_is_owner_import_failure(self, monkeypatch):
"""ImportError from is_owner returns failed."""
monkeypatch.setattr(mod, "logger", MagicMock())
import builtins
real_import = builtins.__import__
def fail_import(name, *args, **kwargs):
if name == "aipass.spawn.apps.handlers.registry":
raise ImportError("no spawn")
return real_import(name, *args, **kwargs)
monkeypatch.setattr(builtins, "__import__", fail_import)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
assert result == "failed"
def test_depth_cap_blocks(self, monkeypatch):
"""AIPASS_WAKE_DEPTH >= MAX_WAKE_DEPTH returns blocked_depth."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
monkeypatch.setenv("AIPASS_WAKE_DEPTH", str(MAX_WAKE_DEPTH))
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@prax", "@target", 0, "/fake/lock")
assert result == "blocked_depth"
def test_depth_cap_over_max_blocks(self, monkeypatch):
"""Depth above max also blocks."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
monkeypatch.setenv("AIPASS_WAKE_DEPTH", str(MAX_WAKE_DEPTH + 5))
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@prax", "@target", 0, "/fake/lock")
assert result == "blocked_depth"
def test_success_on_wake(self, monkeypatch):
"""Successful wake_branch call returns success."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
mock_status = MagicMock()
mock_status.summary = "ok"
@@ -1977,18 +1922,14 @@ class TestWakeSender:
mock_wake,
)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@trigger", "@target", 0, "/fake/lock")
assert result == "success"
mock_wake.assert_called_once_with("@devpulse", auto=True, sender="@ai_mail")
mock_wake.assert_called_once_with("@trigger", auto=True, sender="")
def test_blocked_locked_on_lock_failure(self, monkeypatch):
"""wake_branch failing with lock-related message returns blocked_locked."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
mock_status = MagicMock()
mock_status.summary = "lock: Active agent (PID 1234)"
@@ -1998,17 +1939,13 @@ class TestWakeSender:
mock_wake,
)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@prax", "@target", 0, "/fake/lock")
assert result == "blocked_locked"
def test_blocked_occupied_on_interactive(self, monkeypatch):
"""wake_branch failing with occupancy message returns blocked_occupied."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
mock_status = MagicMock()
mock_status.summary = "blocked: Cannot spawn — interactive session running"
@@ -2018,34 +1955,26 @@ class TestWakeSender:
mock_wake,
)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@trigger", "@target", 0, "/fake/lock")
assert result == "blocked_occupied"
def test_failed_on_exception(self, monkeypatch):
"""Exception during wake returns failed."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
monkeypatch.setattr(
"aipass.ai_mail.apps.handlers.dispatch.wake.wake_branch",
MagicMock(side_effect=RuntimeError("broken")),
)
result = _wake_sender("@devpulse", "@target", 0, "/fake/lock")
result = _wake_sender("@prax", "@target", 0, "/fake/lock")
assert result == "failed"
def test_depth_incremented_before_wake(self, monkeypatch):
"""AIPASS_WAKE_DEPTH is incremented before calling wake_branch."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.setenv("AIPASS_WAKE_DEPTH", "1")
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
captured_depth = []
@@ -2060,17 +1989,13 @@ class TestWakeSender:
capture_wake,
)
_wake_sender("@devpulse", "@target", 0, "/fake/lock")
_wake_sender("@trigger", "@target", 0, "/fake/lock")
assert captured_depth == ["2"]
def test_wake_called_on_failure_exit(self, monkeypatch):
"""Wake fires on non-zero exit code too."""
monkeypatch.setattr(mod, "logger", MagicMock())
monkeypatch.delenv("AIPASS_WAKE_DEPTH", raising=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
MagicMock(return_value=True),
)
mock_status = MagicMock()
mock_status.summary = "ok"
@@ -2080,24 +2005,10 @@ class TestWakeSender:
mock_wake,
)
result = _wake_sender("@devpulse", "@target", 1, "/fake/lock")
result = _wake_sender("@prax", "@target", 1, "/fake/lock")
assert result == "success"
mock_wake.assert_called_once()
def test_sender_normalization_for_is_owner(self, monkeypatch):
"""Sender with or without @ prefix is normalized before is_owner call."""
monkeypatch.setattr(mod, "logger", MagicMock())
mock_is_owner = MagicMock(return_value=False)
monkeypatch.setattr(
"aipass.spawn.apps.handlers.registry.is_owner",
mock_is_owner,
)
_wake_sender("devpulse", "@target", 0, "/fake/lock")
_wake_sender("@devpulse", "@target", 0, "/fake/lock")
assert mock_is_owner.call_count == 2
for call in mock_is_owner.call_args_list:
assert call[0][0] == "@devpulse"
class TestLogWakeResult:
"""_log_wake_result writes to dispatch_wake.log."""
+51
View File
@@ -985,6 +985,57 @@ class TestWakeBranch:
assert ok is False
assert any(s[0] == "fail" and "resolve" in s[1] for s in status.steps)
# --- manager check ---
def test_manager_target_skips_wake(self, tmp_path, monkeypatch):
"""Target with citizen_class=manager returns True (mail only, no wake)."""
branch_path = _make_wake_fixtures(tmp_path, monkeypatch)
trinity = branch_path / ".trinity"
trinity.mkdir(parents=True, exist_ok=True)
(trinity / "passport.json").write_text(
json.dumps({"identity": {"citizen_class": "manager"}}),
encoding="utf-8",
)
status, ok = wake_branch("@testbranch")
assert ok is True
assert any(s[1] == "manager" for s in status.steps)
def test_non_manager_target_continues(self, tmp_path, monkeypatch):
"""Target with non-manager citizen_class proceeds to spawn."""
branch_path = _make_wake_fixtures(tmp_path, monkeypatch)
trinity = branch_path / ".trinity"
trinity.mkdir(parents=True, exist_ok=True)
(trinity / "passport.json").write_text(
json.dumps({"identity": {"citizen_class": "aipass_framework"}}),
encoding="utf-8",
)
_patch_wake_deps(monkeypatch, _clean_zombies=lambda: 0)
monkeypatch.setattr("subprocess.Popen", lambda *a, **kw: _FakeProc())
monkeypatch.setattr(
"aipass.ai_mail.apps.handlers.notify.send_notification",
lambda *a, **kw: None,
raising=False,
)
status, ok = wake_branch("@testbranch")
assert ok is True
assert not any(s[1] == "manager" for s in status.steps)
def test_missing_passport_continues(self, tmp_path, monkeypatch):
"""No passport.json — wake proceeds normally."""
_make_wake_fixtures(tmp_path, monkeypatch)
_patch_wake_deps(monkeypatch, _clean_zombies=lambda: 0)
monkeypatch.setattr("subprocess.Popen", lambda *a, **kw: _FakeProc())
monkeypatch.setattr(
"aipass.ai_mail.apps.handlers.notify.send_notification",
lambda *a, **kw: None,
raising=False,
)
status, ok = wake_branch("@testbranch")
assert ok is True
assert not any(s[1] == "manager" for s in status.steps)
# --- zombie check ---
def test_zombie_check_warns_but_continues(self, tmp_path, monkeypatch):
"""Zombie detected adds warning but dispatch continues."""
_make_wake_fixtures(tmp_path, monkeypatch)
+30 -4
View File
@@ -25,6 +25,7 @@ Every hook event flows through one engine. Platform bridges normalize the event
| `drone @hooks hooksound` | Show current sound mute status |
| `drone @hooks hooksound off` | Mute all hook sounds |
| `drone @hooks hooksound on` | Unmute all hook sounds |
| `drone @hooks dismiss <alert-id>` | Remove an alert from `.aipass/alerts.json` |
| `drone @hooks cadence` | Show prompt injection cadence config and state |
| `drone @hooks verify` | Cross-check provider settings vs project hook config |
| `drone @hooks --help` | Full help reference |
@@ -40,6 +41,8 @@ Hooks operate on two tiers:
**Why provider-only wiring?** Claude Code does not fire `PreToolUse`/`PostToolUse` hooks from project-level settings — only from user-level settings (DPLAN-0160 platform limitation). So all hook entries live in provider settings, and per-project control happens through `.aipass/hooks.json`.
**Deploying new handlers:** Registering a handler in `.aipass/hooks.json` is necessary but not sufficient. Each event type also needs a matching bridge command entry in `~/.claude/settings.json` — this is human-gated (agents cannot edit provider settings). After building a new handler, email @devpulse to wire the settings.json entry. Without it, the engine never receives the event and the handler never fires.
## Architecture
```
@@ -55,6 +58,7 @@ src/aipass/hooks/
│ │ ├── engine.py # Core dispatch — routes events to handlers
│ │ ├── hooksound.py # Sound control (drone @hooks hooksound on/off)
│ │ ├── hookstatus.py # Config viewer (drone @hooks status)
│ │ ├── alert_dismiss.py # Dismiss alerts (drone @hooks dismiss <id>)
│ │ ├── presence.py # Branch presence — claim/release/refresh for .ai_central/PRESENCE.central.json
│ │ ├── sandbox.py # Kernel sandbox — srt/bwrap wrapper + per-role policy generator
│ │ └── wire_verify.py # Wire verification — provider ↔ project hook wiring checker
@@ -66,7 +70,8 @@ src/aipass/hooks/
│ │ │ ├── branch_loader.py # Injects aipass_local_prompt.md
│ │ │ ├── tier0_kernel.py # Injects tier0 kernel prompt (every turn)
│ │ │ ├── navmap.py # Injects tier1 navmap prompt (periodic)
│ │ │ └── identity.py # Injects passport identity block
│ │ │ ├── identity.py # Injects passport identity block
│ │ │ └── persistent_alert.py # Injects advisory banners from .aipass/alerts.json
│ │ ├── security/ # Enforcement hooks
│ │ │ ├── edit_gate.py # Blocks unsafe edits (cross-branch, inbox, diagnostics)
│ │ │ ├── git_gate.py # Enforces git access tiers
@@ -91,7 +96,7 @@ src/aipass/hooks/
│ └── diagnostics.py # JSONL logging for hook execution
├── logs/
│ └── engine.jsonl # JSONL diagnostics (every hook execution)
└── tests/ # 913 tests across 28 test files
└── tests/ # 1071 tests across 29 test files
```
## How It Works
@@ -112,7 +117,7 @@ Handlers are called **dynamically at runtime** — the engine uses `importlib.im
| Event | Hooks | Description |
|---|---|---|
| UserPromptSubmit | presence_gate, identity, email, branch_loader, tier0_kernel, navmap | Presence gate + prompt injection + inbox check |
| UserPromptSubmit | presence_gate, persistent_alert, identity, email, branch_loader, tier0_kernel, navmap, auto_process, user_message_relay | Presence gate + alerts + prompt injection + inbox + auto-process + TG mirror |
| PreToolUse | tool_sound, edit_gate, git_gate, rm_gate, registry_gate | Security gates + guardrails + sound |
| PostToolUse | auto_fix, auto_watchdog | Diagnostics + watchdog |
| SubagentStop | subagent_gate | Seedgo validation |
@@ -140,6 +145,27 @@ The `git_gate` handler (`security/git_gate.py`) enforces git access via drone to
**Why it's on by default:** Agents reflexively reach for raw git, which causes state chaos in a multi-agent system. The gate redirects to `drone @git` which enforces access tiers (read-only for most branches, write-only for devpulse). External users who don't need multi-agent git orchestration can safely disable it.
## Persistent Alerts
The `persistent_alert` handler (`prompt/persistent_alert.py`) injects advisory banners into every prompt when active alerts exist. General-purpose — any agent can raise alerts (prax for runaway logs, trigger for medic, backup for sync failures).
**How it works:** Reads `.aipass/alerts.json` at the project root. Each alert has an ID, source, severity (`warning`/`critical`), title, body, and optional `expires_at`. Active alerts render as a banner every turn until dismissed or expired. Expired alerts are auto-cleaned on read.
**Sound:** Piper TTS fires on first injection per alert ID — subsequent turns are silent for known alerts. New alerts trigger a fresh announcement.
**Dismissing alerts:** `drone @hooks dismiss <alert-id>` removes an alert by ID from `alerts.json`.
**Schema:**
```json
{
"alerts": [{
"id": "uuid", "source": "prax", "severity": "warning",
"title": "High log rate", "body": "commons exceeds 50 lines/s",
"created_at": "iso", "expires_at": "iso or null"
}]
}
```
## Kernel Sandbox (srt/bwrap)
The sandbox module (`apps/modules/sandbox.py`) provides the kernel-level filesystem boundary for agent sessions. It wraps Anthropic's `@anthropic-ai/sandbox-runtime` (srt) library, which uses bubblewrap (bwrap) + Landlock + seccomp on Linux to enforce write/read restrictions at the OS level.
@@ -180,7 +206,7 @@ The @drone broker validates sandbox policy before agent launch. @ai_mail's dispa
- All branches via hook dispatch — every Claude Code session routes through the engine
- @ai_mail dispatch_monitor — sandbox_launch + build_policy for agent launch boundary
*Last Updated: 2026-06-29*
*Last Updated: 2026-07-15*
---
@@ -0,0 +1,131 @@
# =================== AIPass ====================
# Name: persistent_alert.py
# Version: 1.0.1
# Description: Injects advisory banners for active alerts on UserPromptSubmit
# Branch: hooks
# Layer: apps/handlers/prompt
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""Injects advisory banners for active alerts from .aipass/alerts.json."""
import json
from datetime import datetime, timezone
from pathlib import Path
from aipass.prax.apps.modules.logger import system_logger as logger
_announced: set[str] = set()
def _find_aipass_dir() -> Path | None:
"""Walk up from CWD; return the nearest .aipass/ that contains alerts.json.
Every branch has its own .aipass/ (branch prompt), so stopping at the first
.aipass directory would never reach the project root where alerts.json lives.
"""
search = Path.cwd()
home = Path.home()
while search != home and search.parent != search:
aipass_dir = search / ".aipass"
if (aipass_dir / "alerts.json").exists():
return aipass_dir
search = search.parent
return None
def _load_and_clean(alerts_path: Path) -> list[dict]:
"""Load alerts, remove expired, write back if cleaned. Returns active alerts."""
try:
data = json.loads(alerts_path.read_text(encoding="utf-8"))
except (json.JSONDecodeError, OSError) as exc:
logger.info("[HOOKS] persistent_alert: read error: %s", exc)
return []
alerts = data.get("alerts", []) if isinstance(data, dict) else []
if not alerts:
return []
now = datetime.now(timezone.utc)
active = []
cleaned = False
for alert in alerts:
expires = alert.get("expires_at")
if expires:
try:
exp_dt = datetime.fromisoformat(expires)
if exp_dt.tzinfo is None:
exp_dt = exp_dt.replace(tzinfo=timezone.utc)
if exp_dt < now:
cleaned = True
continue
except (ValueError, TypeError) as exc:
logger.info("[HOOKS] persistent_alert: bad expires_at: %s", exc)
active.append(alert)
if cleaned:
try:
alerts_path.write_text(
json.dumps({"alerts": active}, indent=2) + "\n",
encoding="utf-8",
)
except OSError as exc:
logger.info("[HOOKS] persistent_alert: cleanup write error: %s", exc)
return active
def _format_banner(alerts: list[dict]) -> str:
"""Format alert banners for prompt injection."""
lines = []
for alert in alerts:
severity = alert.get("severity", "warning").upper()
title = alert.get("title", "Untitled alert")
body = alert.get("body", "")
source = alert.get("source", "unknown")
alert_id = alert.get("id", "?")
lines.append(f"[{severity}] {title} (from @{source}, id: {alert_id})")
if body:
lines.append(f" {body}")
header = "# Active Alerts"
dismiss_hint = "Dismiss with: drone @hooks dismiss <alert-id>"
return "\n".join([header, ""] + lines + ["", dismiss_hint])
def handle(hook_data: dict) -> dict:
"""Inject advisory banners for active alerts.
Args:
hook_data: Parsed hook event dict from engine.
Returns:
Result dict with stdout (banner or empty) and exit_code.
"""
aipass_dir = _find_aipass_dir()
if not aipass_dir:
return {"stdout": "", "exit_code": 0}
alerts_path = aipass_dir / "alerts.json"
if not alerts_path.exists():
return {"stdout": "", "exit_code": 0}
alerts = _load_and_clean(alerts_path)
if not alerts:
return {"stdout": "", "exit_code": 0}
banner = _format_banner(alerts)
new_ids = [a["id"] for a in alerts if a.get("id") and a["id"] not in _announced]
sound = ""
if new_ids:
_announced.update(new_ids)
count = len(alerts)
plural = "s" if count != 1 else ""
sound = f"alert: {count} active alert{plural}"
logger.info("[HOOKS] persistent_alert: %d active alerts injected", len(alerts))
result = {"stdout": banner, "exit_code": 0}
if sound:
result["sound"] = sound
return result
@@ -0,0 +1,111 @@
# =================== AIPass ====================
# Name: alert_dismiss.py
# Version: 1.0.1
# Description: Dismiss alerts from .aipass/alerts.json via drone @hooks dismiss
# Branch: hooks
# Layer: apps/modules
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""Dismiss alerts from .aipass/alerts.json by ID."""
import json
from pathlib import Path
from aipass.cli.apps.modules import err_console
from aipass.prax.apps.modules.logger import system_logger as logger
CONSOLE = err_console
HELP_COMMANDS = [
("dismiss <alert-id>", "Remove an alert from .aipass/alerts.json"),
]
def _find_aipass_dir() -> Path | None:
"""Walk up from CWD; return the nearest .aipass/ that contains alerts.json.
Every branch has its own .aipass/ (branch prompt), so stopping at the first
.aipass directory would never reach the project root where alerts.json lives.
"""
search = Path.cwd()
home = Path.home()
while search != home and search.parent != search:
aipass_dir = search / ".aipass"
if (aipass_dir / "alerts.json").exists():
return aipass_dir
search = search.parent
return None
def _dismiss_alert(alert_id: str) -> bool:
"""Remove an alert by ID from alerts.json. Returns True if found and removed."""
aipass_dir = _find_aipass_dir()
if not aipass_dir:
CONSOLE.print("[red]No .aipass/ directory found in this directory tree.[/red]")
return False
alerts_path = aipass_dir / "alerts.json"
if not alerts_path.exists():
CONSOLE.print("[yellow]No alerts.json found — nothing to dismiss.[/yellow]")
return False
try:
data = json.loads(alerts_path.read_text(encoding="utf-8"))
except (json.JSONDecodeError, OSError) as exc:
logger.error("[HOOKS] dismiss: read error: %s", exc)
CONSOLE.print(f"[red]Failed to read alerts.json: {exc}[/red]")
return False
alerts = data.get("alerts", []) if isinstance(data, dict) else []
original_count = len(alerts)
remaining = [a for a in alerts if a.get("id") != alert_id]
if len(remaining) == original_count:
CONSOLE.print(f"[yellow]Alert {alert_id} not found.[/yellow]")
return False
try:
alerts_path.write_text(
json.dumps({"alerts": remaining}, indent=2) + "\n",
encoding="utf-8",
)
except OSError as exc:
logger.error("[HOOKS] dismiss: write error: %s", exc)
CONSOLE.print(f"[red]Failed to write alerts.json: {exc}[/red]")
return False
logger.info("[HOOKS] dismiss: removed alert %s", alert_id)
CONSOLE.print(f"[green]Dismissed alert {alert_id}[/green]")
return True
def print_introspection():
"""Print module structure for drone routing."""
CONSOLE.print("[bold cyan]alert_dismiss[/bold cyan] — Remove alerts from .aipass/alerts.json")
def handle_command(command: str, args: list) -> bool:
"""Route dismiss commands from drone @hooks."""
if command != "dismiss":
return False
if not args:
print_introspection()
CONSOLE.print()
CONSOLE.print(" drone @hooks dismiss <alert-id>")
CONSOLE.print()
CONSOLE.print("Removes the alert with the given ID from .aipass/alerts.json.")
return True
if args[0] in ("--help", "-h", "help"):
CONSOLE.print("[bold cyan]dismiss[/bold cyan] — Remove an alert")
CONSOLE.print()
CONSOLE.print(" drone @hooks dismiss <alert-id>")
CONSOLE.print()
CONSOLE.print("Removes the alert with the given ID from .aipass/alerts.json.")
return True
_dismiss_alert(args[0])
return True
@@ -0,0 +1,425 @@
# =================== AIPass ====================
# Name: test_persistent_alert.py
# Version: 1.0.0
# Description: Tests for persistent_alert handler and alert_dismiss module
# Branch: hooks
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""Tests for handlers/prompt/persistent_alert.py and modules/alert_dismiss.py."""
import json
from datetime import datetime, timedelta, timezone
from pathlib import Path
from unittest.mock import patch
def _make_alert(
alert_id="test-001",
source="prax",
severity="warning",
title="Test alert",
body="Something happened",
expires_at=None,
):
alert = {
"id": alert_id,
"source": source,
"severity": severity,
"title": title,
"body": body,
"created_at": datetime.now(timezone.utc).isoformat(),
"expires_at": expires_at,
}
return alert
def _write_alerts(aipass_dir: Path, alerts: list[dict]):
alerts_path = aipass_dir / "alerts.json"
alerts_path.write_text(
json.dumps({"alerts": alerts}, indent=2) + "\n",
encoding="utf-8",
)
class TestPersistentAlertHandler:
"""Banner injection behavior."""
def test_banner_injected_when_alerts_exist(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert()])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert result["exit_code"] == 0
assert "# Active Alerts" in result["stdout"]
assert "[WARNING] Test alert" in result["stdout"]
assert "drone @hooks dismiss" in result["stdout"]
def test_no_banner_when_no_alerts_file(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert result["stdout"] == ""
assert result["exit_code"] == 0
def test_no_banner_when_empty_alerts(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert result["stdout"] == ""
def test_no_banner_when_no_aipass_dir(self):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=None,
):
result = handle({})
assert result["stdout"] == ""
assert result["exit_code"] == 0
def test_multiple_alerts_all_shown(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(
aipass_dir,
[
_make_alert(alert_id="a1", title="First"),
_make_alert(alert_id="a2", title="Second", severity="critical"),
],
)
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert "[WARNING] First" in result["stdout"]
assert "[CRITICAL] Second" in result["stdout"]
def test_source_and_id_in_banner(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="abc-123", source="trigger")])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert "@trigger" in result["stdout"]
assert "abc-123" in result["stdout"]
def test_body_included_in_banner(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(body="Log rate exceeds 50/s")])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert "Log rate exceeds 50/s" in result["stdout"]
def test_no_body_line_when_body_empty(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(body="")])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
lines = result["stdout"].split("\n")
body_lines = [line for line in lines if line.startswith(" ")]
assert len(body_lines) == 0
class TestExpiredAlertCleanup:
"""Auto-cleaning of expired alerts."""
def test_expired_alerts_removed(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
past = (datetime.now(timezone.utc) - timedelta(hours=1)).isoformat()
_write_alerts(
aipass_dir,
[
_make_alert(alert_id="expired", expires_at=past),
_make_alert(alert_id="active", expires_at=None),
],
)
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert "active" in result["stdout"]
assert "expired" not in result["stdout"]
saved = json.loads((aipass_dir / "alerts.json").read_text())
assert len(saved["alerts"]) == 1
assert saved["alerts"][0]["id"] == "active"
def test_all_expired_returns_empty(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
past = (datetime.now(timezone.utc) - timedelta(hours=1)).isoformat()
_write_alerts(aipass_dir, [_make_alert(expires_at=past)])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert result["stdout"] == ""
def test_future_expiry_kept(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
future = (datetime.now(timezone.utc) + timedelta(hours=1)).isoformat()
_write_alerts(aipass_dir, [_make_alert(alert_id="still-valid", expires_at=future)])
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert "still-valid" in result["stdout"]
def test_corrupt_json_returns_empty(self, tmp_path):
from aipass.hooks.apps.handlers.prompt.persistent_alert import handle
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
(aipass_dir / "alerts.json").write_text("{bad json", encoding="utf-8")
with patch(
"aipass.hooks.apps.handlers.prompt.persistent_alert._find_aipass_dir",
return_value=aipass_dir,
):
result = handle({})
assert result["stdout"] == ""
assert result["exit_code"] == 0
class TestAlertSound:
"""Sound fires on first injection only."""
def test_sound_on_first_injection(self, tmp_path):
from aipass.hooks.apps.handlers.prompt import persistent_alert
persistent_alert._announced.clear()
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="snd-001")])
with patch.object(persistent_alert, "_find_aipass_dir", return_value=aipass_dir):
result = persistent_alert.handle({})
assert "sound" in result
assert "1 active alert" in result["sound"]
def test_no_sound_on_repeat_injection(self, tmp_path):
from aipass.hooks.apps.handlers.prompt import persistent_alert
persistent_alert._announced.clear()
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="snd-002")])
with patch.object(persistent_alert, "_find_aipass_dir", return_value=aipass_dir):
persistent_alert.handle({})
result = persistent_alert.handle({})
assert "sound" not in result
def test_sound_on_new_alert_added(self, tmp_path):
from aipass.hooks.apps.handlers.prompt import persistent_alert
persistent_alert._announced.clear()
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="snd-003")])
with patch.object(persistent_alert, "_find_aipass_dir", return_value=aipass_dir):
persistent_alert.handle({})
_write_alerts(
aipass_dir,
[
_make_alert(alert_id="snd-003"),
_make_alert(alert_id="snd-004"),
],
)
with patch.object(persistent_alert, "_find_aipass_dir", return_value=aipass_dir):
result = persistent_alert.handle({})
assert "sound" in result
assert "2 active alerts" in result["sound"]
def test_no_sound_when_no_alerts(self, tmp_path):
from aipass.hooks.apps.handlers.prompt import persistent_alert
persistent_alert._announced.clear()
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [])
with patch.object(persistent_alert, "_find_aipass_dir", return_value=aipass_dir):
result = persistent_alert.handle({})
assert "sound" not in result
class TestAlertDismiss:
"""drone @hooks dismiss behavior."""
def test_dismiss_removes_by_id(self, tmp_path):
from aipass.hooks.apps.modules.alert_dismiss import _dismiss_alert
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(
aipass_dir,
[
_make_alert(alert_id="keep"),
_make_alert(alert_id="remove"),
],
)
with patch(
"aipass.hooks.apps.modules.alert_dismiss._find_aipass_dir",
return_value=aipass_dir,
):
result = _dismiss_alert("remove")
assert result is True
saved = json.loads((aipass_dir / "alerts.json").read_text())
assert len(saved["alerts"]) == 1
assert saved["alerts"][0]["id"] == "keep"
def test_dismiss_nonexistent_returns_false(self, tmp_path):
from aipass.hooks.apps.modules.alert_dismiss import _dismiss_alert
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="exists")])
with patch(
"aipass.hooks.apps.modules.alert_dismiss._find_aipass_dir",
return_value=aipass_dir,
):
result = _dismiss_alert("nope")
assert result is False
def test_dismiss_no_alerts_file(self, tmp_path):
from aipass.hooks.apps.modules.alert_dismiss import _dismiss_alert
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
with patch(
"aipass.hooks.apps.modules.alert_dismiss._find_aipass_dir",
return_value=aipass_dir,
):
result = _dismiss_alert("any")
assert result is False
def test_dismiss_no_aipass_dir(self):
from aipass.hooks.apps.modules.alert_dismiss import _dismiss_alert
with patch(
"aipass.hooks.apps.modules.alert_dismiss._find_aipass_dir",
return_value=None,
):
result = _dismiss_alert("any")
assert result is False
def test_handle_command_routes_dismiss(self, tmp_path):
from aipass.hooks.apps.modules.alert_dismiss import handle_command
aipass_dir = tmp_path / ".aipass"
aipass_dir.mkdir()
_write_alerts(aipass_dir, [_make_alert(alert_id="cmd-test")])
with patch(
"aipass.hooks.apps.modules.alert_dismiss._find_aipass_dir",
return_value=aipass_dir,
):
result = handle_command("dismiss", ["cmd-test"])
assert result is True
saved = json.loads((aipass_dir / "alerts.json").read_text())
assert len(saved["alerts"]) == 0
def test_handle_command_ignores_other_commands(self):
from aipass.hooks.apps.modules.alert_dismiss import handle_command
assert handle_command("status", []) is False
def test_handle_command_help(self):
from aipass.hooks.apps.modules.alert_dismiss import handle_command
assert handle_command("dismiss", ["--help"]) is True
+7
View File
@@ -90,6 +90,13 @@
"subject": "TG relay send backoff: offline mode + log-once on network failures",
"date_closed": "2026-07-14",
"location": "prax"
},
{
"plan_id": "TDPLAN-0013",
"type": "TDPLAN",
"subject": "DPLAN-0242 build: runaway-log detection and escalation",
"date_closed": "2026-07-15",
"location": "prax"
}
],
"document_metadata": {
+27 -12
View File
@@ -5,7 +5,7 @@
**Purpose:** System-wide logging, real-time monitoring, and dashboard infrastructure for AIPass.
**Module:** `aipass.prax`
**Version:** 2.0.0
**Last Updated:** 2026-06-05
**Last Updated:** 2026-07-14
---
@@ -51,11 +51,23 @@ Real-time unified console showing:
- **Caller attribution** — `CALLER → TARGET` for drone commands
- **Model tags** — `[BRANCH/model]` (e.g., `[DEVPULSE/opus]`, `[DEVPULSE/gpt-5.4]`)
- **Multi-CLI** — Claude Code (JSONL), Codex (JSONL) session monitoring
- **Rate tracking** — 4th background thread scans `system_logs/` for runaway log growth every 10s
- **Polling fallback** — automatic fallback when inotify watches are exhausted
- **Soft start** — only shows new activity after launch (seeks to EOF on startup)
Interactive commands inside the monitor: `help`, `status`, `quit`/`exit`.
### Log Health
```bash
drone @prax log-health # Show module info
drone @prax log-health scan # Scan all log files, show current growth rates
drone @prax log-health snapshot # Show last known rates (no new scan)
drone @prax log-health --help # Log health usage
```
Quick overview of log file growth rates across `system_logs/`. Powered by the rate tracker handler — `scan` runs a fresh measurement, `snapshot` reads the last persisted state without scanning.
### Status
```bash
@@ -123,12 +135,13 @@ prax/
├── __init__.py # Public API: exports `logger` (NullLogger fallback)
├── apps/
│ ├── prax.py # Entry point — auto-discovers modules, routes commands
│ ├── modules/ # Business logic (5 command modules)
│ ├── modules/ # Business logic (6 command modules)
│ │ ├── logger.py # SystemLogger — auto-routing, two-tier logging
│ │ ├── monitor.py # Mission Control — 3-thread real-time monitoring
│ │ ├── monitor.py # Mission Control — 4-thread real-time monitoring
│ │ ├── dashboard.py # Dashboard — template management, refresh, write-through
│ │ ├── status.py # System status — health display (STATUS.md sync dormant)
│ │ └── log_audit.py # Log audit — scan, health summary, enforce limits
│ │ ├── log_audit.py # Log audit — scan, health summary, enforce limits
│ │ └── log_health.py # Log health — rate overview (scan/snapshot)
│ └── handlers/ # Implementation details (11 handler directories)
│ ├── central/ # Central file reader (.ai_central/*.central.json)
│ ├── config/ # Path resolution, log config, ignore patterns
@@ -137,13 +150,13 @@ prax/
│ ├── json/ # Auto-creating JSON handler (config/data/log per module)
│ ├── json_templates/ # Default JSON templates for auto-creation
│ ├── logging/ # Setup, rotation, introspection, override, direct logger
│ ├── monitoring/ # Event queue, branch detector, stream output, log watcher
│ ├── monitoring/ # Event queue, branch detector, stream output, log watcher, rate tracker
│ ├── registry/ # Module registry load/save
│ ├── status/ # STATUS.md sync handler (dormant — TDPLAN-0007)
│ └── watcher/ # Background system watchers
├── prax_json/ # Auto-created per-module config/data/log files
├── templates/ # Dashboard template schema (DASHBOARD.template.json)
└── tests/ # 1007 tests across 19 files
└── tests/ # 1028 tests across 20 files
```
### Design Pattern
@@ -164,14 +177,15 @@ drone @prax monitor run
1. **Auto-routing** — `logger.info()` inspects the call stack to identify the caller's module, branch, and file path, then routes the log entry to the correct per-module log file.
2. **Two-tier logging** — Each log entry goes to both `system_logs/` (central, all branches) and `<branch>/logs/` (branch-local), both with size-based rotation.
3. **Self-healing** — Auto-creates missing log directories, falls back to `system_logs/external/` for unknown modules, provides NullLogger if prax itself fails to import.
4. **Mission Control** — Three threads: display worker (pulls from event queue), file watcher (watchdog on branch `apps/` dirs), log watcher (tails `system_logs/*.log`). Falls back to polling when inotify is exhausted.
4. **Mission Control** — Four threads: display worker (pulls from event queue), file watcher (watchdog on branch `apps/` dirs), log watcher (tails `system_logs/*.log`), rate tracker (scans `system_logs/` for runaway growth every 10s). Falls back to polling when inotify is exhausted.
5. **Multi-CLI monitoring** — Watches Claude Code JSONL and Codex JSONL session files. Extracts agent activity (thinking, tool use, responses) with model detection and branch resolution.
6. **Dashboard** — Template-based per-branch dashboard files. Refreshes from central files (`*.central.json`). Write-through API for services to update sections directly.
7. **STATUS sync** — *(Dormant — TDPLAN-0007)* Previously scanned all branch `STATUS.local.md` files and built aggregated `STATUS.md`. Engine code intact but no longer triggered.
6. **Runaway-log detection** — Rate tracker measures byte growth per log file, estimates lines/min from byte deltas. Sustained thresholds: WARNING (>100 lines/min for 2 min), CRITICAL (>10 lines/sec for 1 min). Fires `runaway_log_detected` on the trigger event bus. State persists to disk across process restarts. Per-file suppression available.
7. **Dashboard** — Template-based per-branch dashboard files. Refreshes from central files (`*.central.json`). Write-through API for services to update sections directly.
8. **STATUS sync** — *(Dormant — TDPLAN-0007)* Previously scanned all branch `STATUS.local.md` files and built aggregated `STATUS.md`. Engine code intact but no longer triggered.
## Tests
1007 tests across 19 files, covering all major components:
1028 tests across 20 files, covering all major components:
| Test File | Tests | Coverage |
|-----------|-------|----------|
@@ -179,7 +193,7 @@ drone @prax monitor run
| test_monitoring_handlers.py | 139 | Branch detector, stream output, event handling |
| test_operations.py | 99 | Dashboard operations, write-through |
| test_log_watcher.py | 82 | Log file tailing, agent activity parsing |
| test_monitor_module.py | 73 | Monitor commands, thread lifecycle |
| test_monitor_module.py | 73 | Monitor commands, thread lifecycle (4-thread) |
| test_logging_handlers.py | 41 | Setup, rotation, introspection, direct logger |
| test_logging.py | 41 | Core logging system |
| test_logger_module.py | 40 | Logger init, routing, lifecycle |
@@ -193,6 +207,7 @@ drone @prax monitor run
| test_central.py | 14 | Central reader |
| test_devpulse_dashboard_plugin.py | 12 | Dashboard plugin (git, session, dispatch) |
| test_log_audit.py | 10 | Log audit |
| test_rate_tracker.py | 21 | Rate tracking, thresholds, persistence, suppression |
| test_status.py | 8 | Status commands |
## Integration Points
@@ -216,7 +231,7 @@ drone @prax monitor run
---
*Last Updated: 2026-06-05*
*Last Updated: 2026-07-14*
---
[← Back to AIPass](../../../README.md)
@@ -0,0 +1,368 @@
# =================== AIPass ====================
# Name: rate_tracker.py
# Description: Log file rate tracking for runaway detection
# Version: 1.1.0
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""
Rate Tracker — volume-based runaway-log detection.
Tracks byte growth rate per log file in system_logs/. When a file sustains
abnormal growth (lines/min above threshold for consecutive intervals), fires
a ``runaway_log_detected`` event on the trigger event bus.
Orthogonal to medic's content-based ERROR/CRITICAL detection — this catches
rate regardless of log level.
State persists to ``prax_json/rate_tracker_data.json`` so that rates survive
across process restarts and CLI invocations can display meaningful data.
"""
import time
from collections import deque
from pathlib import Path
from typing import Dict, Optional
from aipass.prax.apps.modules.logger import get_direct_logger
from aipass.prax.apps.handlers.json import json_handler
from aipass.prax.apps.handlers.config.load import get_system_logs_dir
from aipass.prax.apps.handlers.monitoring.branch_detector import detect_branch_from_log
logger = get_direct_logger()
try:
from aipass.trigger.apps.modules.core import trigger
_HAS_TRIGGER = True
except ImportError as exc:
logger.info("[rate_tracker] trigger module not available: %s", exc)
trigger = None # type: ignore[assignment]
_HAS_TRIGGER = False
SCAN_INTERVAL = 10.0
AVG_LINE_BYTES = 120
WARNING_LINES_PER_MIN = 100
WARNING_SUSTAINED_INTERVALS = 12 # 12 * 10s = 2 min
CRITICAL_LINES_PER_MIN = 600 # 10/sec * 60
CRITICAL_SUSTAINED_INTERVALS = 6 # 6 * 10s = 1 min
_RATE_HISTORY_SIZE = 30
_DATA_FILE = "rate_tracker"
class FileRateState:
"""Per-file tracking state."""
__slots__ = (
"last_offset",
"last_check",
"rates",
"warning_sustained",
"critical_sustained",
"fired_warning",
"fired_critical",
)
def __init__(self, offset: int, now: float) -> None:
self.last_offset: int = offset
self.last_check: float = now
self.rates: deque = deque(maxlen=_RATE_HISTORY_SIZE)
self.warning_sustained: int = 0
self.critical_sustained: int = 0
self.fired_warning: bool = False
self.fired_critical: bool = False
def to_dict(self) -> dict:
"""Serialize to a dict for disk persistence."""
return {
"last_offset": self.last_offset,
"last_check": self.last_check,
"warning_sustained": self.warning_sustained,
"critical_sustained": self.critical_sustained,
"fired_warning": self.fired_warning,
"fired_critical": self.fired_critical,
}
@classmethod
def from_dict(cls, d: dict) -> "FileRateState":
"""Restore from a persisted dict."""
state = cls(d.get("last_offset", 0), d.get("last_check", 0.0))
state.warning_sustained = d.get("warning_sustained", 0)
state.critical_sustained = d.get("critical_sustained", 0)
state.fired_warning = d.get("fired_warning", False)
state.fired_critical = d.get("fired_critical", False)
return state
_tracked: Dict[str, FileRateState] = {}
_suppressed_files: set = set()
_state_loaded: bool = False
def configure_suppression(file_names: Optional[set] = None) -> None:
"""Set the list of log file names to skip during detection."""
global _suppressed_files
_suppressed_files = file_names or set()
def _load_state() -> None:
"""Load persisted tracking state from disk on first scan."""
global _state_loaded
if _state_loaded:
return
_state_loaded = True
data = json_handler.load_json(_DATA_FILE, "data")
if data is None:
return
files = data.get("files", {})
for file_key, state_dict in files.items():
if not isinstance(state_dict, dict):
continue
_tracked[file_key] = FileRateState.from_dict(state_dict)
count = len(_tracked)
if count:
logger.info("[rate_tracker] Loaded %d file states from disk", count)
def _save_state() -> None:
"""Persist current tracking state to disk."""
files = {}
for file_key, state in _tracked.items():
files[file_key] = state.to_dict()
from datetime import date
today = date.today().isoformat()
data = {
"module_name": _DATA_FILE,
"created": today,
"last_updated": today,
"files": files,
}
json_handler.save_json(_DATA_FILE, "data", data)
def scan_rates() -> list:
"""Scan all .log files in system_logs/, update rates, fire events if thresholds met.
Returns a list of dicts describing each tracked file's current state,
suitable for display by the log-health command.
"""
_load_state()
logs_dir = get_system_logs_dir()
if not logs_dir.exists():
return []
now = time.time()
results = []
current_files = set()
for log_file in logs_dir.glob("*.log"):
file_key = str(log_file)
current_files.add(file_key)
if log_file.name in _suppressed_files:
continue
try:
size = log_file.stat().st_size
except OSError as exc:
logger.info("[rate_tracker] Cannot stat %s: %s", log_file.name, exc)
continue
state = _tracked.get(file_key)
if state is None:
_tracked[file_key] = FileRateState(size, now)
continue
elapsed = now - state.last_check
if elapsed < 1.0:
continue
if size < state.last_offset:
state.last_offset = size
state.last_check = now
state.warning_sustained = 0
state.critical_sustained = 0
continue
bytes_added = size - state.last_offset
lines_estimate = bytes_added / AVG_LINE_BYTES if bytes_added > 0 else 0.0
lines_per_min = (lines_estimate / elapsed) * 60.0
state.rates.append((now, lines_per_min))
state.last_offset = size
state.last_check = now
severity = _evaluate_thresholds(state, lines_per_min, file_key, log_file)
results.append(
{
"file": log_file.name,
"path": file_key,
"size_kb": round(size / 1024, 1),
"rate_lines_per_min": round(lines_per_min, 1),
"warning_sustained": state.warning_sustained,
"critical_sustained": state.critical_sustained,
"severity": severity,
"branch": detect_branch_from_log(file_key),
}
)
stale = [k for k in _tracked if k not in current_files]
for k in stale:
del _tracked[k]
_save_state()
return results
def _evaluate_thresholds(
state: FileRateState,
lines_per_min: float,
file_key: str,
log_file: Path,
) -> Optional[str]:
"""Update sustained counters and fire events when thresholds are crossed."""
severity = None
if lines_per_min >= CRITICAL_LINES_PER_MIN:
state.critical_sustained += 1
state.warning_sustained += 1
elif lines_per_min >= WARNING_LINES_PER_MIN:
state.critical_sustained = 0
state.warning_sustained += 1
else:
if state.fired_warning or state.fired_critical:
logger.info(
"[rate_tracker] %s rate subsided (%.0f lines/min)",
log_file.name,
lines_per_min,
)
state.warning_sustained = 0
state.critical_sustained = 0
state.fired_warning = False
state.fired_critical = False
return None
if state.critical_sustained >= CRITICAL_SUSTAINED_INTERVALS and not state.fired_critical:
severity = "critical"
state.fired_critical = True
duration = state.critical_sustained * SCAN_INTERVAL
_fire_event(file_key, lines_per_min, duration, "critical")
elif state.warning_sustained >= WARNING_SUSTAINED_INTERVALS and not state.fired_warning:
severity = "warning"
state.fired_warning = True
duration = state.warning_sustained * SCAN_INTERVAL
_fire_event(file_key, lines_per_min, duration, "warning")
else:
if state.critical_sustained > 0:
severity = "rising_critical"
elif state.warning_sustained > 0:
severity = "rising_warning"
return severity
def _fire_event(
file_path: str,
rate_lines_per_min: float,
sustained_duration_sec: float,
severity: str,
) -> None:
"""Fire runaway_log_detected on the trigger event bus."""
branch = detect_branch_from_log(file_path)
logger.warning(
"[rate_tracker] RUNAWAY %s: %s — %.0f lines/min sustained %.0fs (branch: %s)",
severity.upper(),
Path(file_path).name,
rate_lines_per_min,
sustained_duration_sec,
branch,
)
json_handler.log_operation(
"runaway_detected",
{
"file": file_path,
"rate": rate_lines_per_min,
"duration": sustained_duration_sec,
"severity": severity,
"branch": branch,
},
)
if _HAS_TRIGGER and trigger is not None:
trigger.fire(
"runaway_log_detected",
file_path=file_path,
rate_lines_per_min=rate_lines_per_min,
sustained_duration_sec=sustained_duration_sec,
severity=severity,
branch=branch,
)
def _resolve_severity(state: FileRateState) -> Optional[str]:
"""Derive display severity from a file's current state."""
if state.fired_critical:
return "critical"
if state.fired_warning:
return "warning"
if state.critical_sustained > 0:
return "rising_critical"
if state.warning_sustained > 0:
return "rising_warning"
return None
def _file_size(path: str) -> int:
"""Read file size, returning 0 on any OS error."""
try:
return Path(path).stat().st_size
except OSError:
logger.info("[rate_tracker] Cannot stat %s for snapshot", Path(path).name)
return 0
def get_snapshot() -> list:
"""Return the current tracking state without scanning (for display only).
Loads persisted state from disk if not already loaded, so CLI
invocations can display rates collected by the monitor process.
"""
_load_state()
results = []
for file_key, state in _tracked.items():
last_rate = state.rates[-1][1] if state.rates else 0.0
results.append(
{
"file": Path(file_key).name,
"path": file_key,
"size_kb": round(_file_size(file_key) / 1024, 1),
"rate_lines_per_min": round(last_rate, 1),
"warning_sustained": state.warning_sustained,
"critical_sustained": state.critical_sustained,
"severity": _resolve_severity(state),
"branch": detect_branch_from_log(file_key),
}
)
return results
def reset() -> None:
"""Clear all tracking state. Used in tests."""
global _state_loaded
_tracked.clear()
_suppressed_files.clear()
_state_loaded = False
+176
View File
@@ -0,0 +1,176 @@
# =================== AIPass ====================
# Name: log_health.py
# Description: PRAX Log Health Command — rate overview
# Version: 1.0.0
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""
PRAX Log Health Module
Implements the 'log-health' command showing current log file growth rates
across system_logs/. Powered by the rate_tracker handler.
"""
import os
import sys
from typing import List
if sys.platform == "win32":
os.environ.setdefault("PYTHONUTF8", "1")
for _stream in (sys.stdout, sys.stderr):
_reconfigure = getattr(_stream, "reconfigure", None)
if _reconfigure is not None:
_reconfigure(encoding="utf-8", errors="replace")
from aipass.prax.apps.modules.logger import system_logger as logger
from aipass.cli.apps.modules import console, error
from aipass.prax.apps.handlers.json import json_handler
def print_introspection():
"""Display module introspection."""
console.print()
console.print("[bold cyan]PRAX Log Health Module[/bold cyan]")
console.print()
console.print("[yellow]Purpose:[/yellow]")
console.print(" Show log file growth rates and detect runaway logs")
console.print()
console.print("[yellow]Connected Handlers:[/yellow]")
console.print()
console.print(" [cyan]prax/handlers/monitoring/[/cyan]")
console.print(" [dim]- rate_tracker.py (scan_rates, get_snapshot)[/dim]")
console.print()
console.print("[dim]Run 'drone @prax log-health --help' for usage[/dim]")
console.print()
def print_help():
"""Drone-compliant help output."""
console.print()
console.print("[bold cyan]PRAX Log Health[/bold cyan]")
console.print()
console.print("[yellow]Purpose:[/yellow]")
console.print(" Monitor log file growth rates across system_logs/")
console.print()
console.print("[yellow]Subcommands:[/yellow]")
console.print()
console.print(" [cyan]scan[/cyan] Scan all log files and show current rates")
console.print(" [cyan]snapshot[/cyan] Show last known rates (no new scan)")
console.print()
console.print("[yellow]Usage:[/yellow]")
console.print()
console.print(" [dim]# Scan and show current rates[/dim]")
console.print(" $ drone @prax log-health scan")
console.print()
console.print(" [dim]# Show last known rates without scanning[/dim]")
console.print(" $ drone @prax log-health snapshot")
console.print()
def _display_rates(results: list, is_scan: bool) -> None:
"""Display rate results in a formatted table."""
label = "Scan" if is_scan else "Snapshot"
console.print()
console.print(f"[bold cyan]Log Health {label}[/bold cyan] [dim](system_logs/)[/dim]")
if not results:
console.print(" [dim]No log files tracked yet[/dim]")
console.print()
return
active = [r for r in results if r["rate_lines_per_min"] > 0]
idle = [r for r in results if r["rate_lines_per_min"] == 0]
flagged = [r for r in results if r.get("severity")]
console.print(f" Files tracked: {len(results)}")
console.print(f" Active: {len(active)}, Idle: {len(idle)}")
if flagged:
console.print(f" [yellow]Flagged: {len(flagged)}[/yellow]")
if flagged:
console.print()
console.print("[yellow]Flagged files:[/yellow]")
for r in sorted(flagged, key=lambda x: x["rate_lines_per_min"], reverse=True):
sev = r["severity"] or ""
color = "red" if "critical" in sev else "yellow"
console.print(
f" [{color}]{sev.upper()}[/{color}] "
f"{r['file']}: {r['rate_lines_per_min']} lines/min "
f"({r['size_kb']} KB) [{r['branch']}]"
)
if active:
console.print()
console.print("[cyan]Active files:[/cyan]")
for r in sorted(active, key=lambda x: x["rate_lines_per_min"], reverse=True):
if r.get("severity"):
continue
console.print(f" {r['file']}: {r['rate_lines_per_min']} lines/min ({r['size_kb']} KB) [{r['branch']}]")
if idle and len(idle) <= 10:
console.print()
console.print("[dim]Idle files:[/dim]")
for r in sorted(idle, key=lambda x: x["file"]):
console.print(f" [dim]{r['file']}: {r['size_kb']} KB [{r['branch']}][/dim]")
elif idle:
console.print()
console.print(f" [dim]{len(idle)} idle files (0 lines/min)[/dim]")
console.print()
def handle_command(command: str, args: List[str]) -> bool:
"""Handle log-health command.
Args:
command: Command name
args: Command arguments
Returns:
True if command was handled
"""
if command != "log-health":
return False
if not args:
print_introspection()
return True
if args[0] in ("--help", "-h", "help"):
print_help()
return True
from aipass.prax.apps.handlers.monitoring.rate_tracker import scan_rates, get_snapshot
subcmd = args[0]
logger.info("[log-health] %s", subcmd)
json_handler.log_operation("log_health_executed", {"mode": subcmd})
if subcmd == "scan":
results = scan_rates()
_display_rates(results, is_scan=True)
return True
if subcmd == "snapshot":
results = get_snapshot()
_display_rates(results, is_scan=False)
return True
error(f"Unknown log-health subcommand: {subcmd}")
print_help()
return True
if __name__ == "__main__":
if len(sys.argv) == 1:
print_introspection()
sys.exit(0)
if "--help" in sys.argv:
print_help()
sys.exit(0)
args = [arg for arg in sys.argv[1:] if not arg.startswith("--")]
handle_command("log-health", args)
+19 -6
View File
@@ -63,6 +63,7 @@ _module_tracker: Optional[ModuleTracker] = None
_display_thread: Optional[threading.Thread] = None
_file_watcher_thread: Optional[threading.Thread] = None
_log_watcher_thread: Optional[threading.Thread] = None
_rate_tracker_thread: Optional[threading.Thread] = None
def print_introspection():
@@ -174,7 +175,7 @@ def _load_relay_config() -> Optional[dict]:
def _run_monitor(args: List[str]) -> bool:
"""Launch Mission Control live monitoring."""
global _event_queue, _module_tracker
global _display_thread, _file_watcher_thread, _log_watcher_thread
global _display_thread, _file_watcher_thread, _log_watcher_thread, _rate_tracker_thread
json_handler.log_operation("monitor_started", {"args": args})
logger.info(f"Starting unified monitoring (args: {args})")
@@ -223,20 +224,20 @@ def _run_monitor(args: List[str]) -> bool:
def _start_threads():
"""Start all monitoring threads"""
global _display_thread, _file_watcher_thread, _log_watcher_thread
global _display_thread, _file_watcher_thread, _log_watcher_thread, _rate_tracker_thread
# Display thread - pulls from event queue and displays
_display_thread = threading.Thread(target=_display_worker, daemon=True)
_display_thread.start()
# File watcher thread - watches filesystem changes
_file_watcher_thread = threading.Thread(target=_file_watcher_worker, daemon=True)
_file_watcher_thread.start()
# Log watcher thread - watches log files
_log_watcher_thread = threading.Thread(target=_log_watcher_worker, daemon=True)
_log_watcher_thread.start()
_rate_tracker_thread = threading.Thread(target=_rate_tracker_worker, daemon=True)
_rate_tracker_thread.start()
logger.info("All monitoring threads started")
@@ -251,7 +252,7 @@ def _stop_threads():
_event_queue.stop()
# Join all daemon threads with timeout
for t in (_display_thread, _file_watcher_thread, _log_watcher_thread):
for t in (_display_thread, _file_watcher_thread, _log_watcher_thread, _rate_tracker_thread):
if t is not None and t.is_alive():
t.join(timeout=2.0)
@@ -487,6 +488,18 @@ def _log_watcher_worker():
stop_log_watcher()
def _rate_tracker_worker():
"""Rate tracker thread — scans system_logs/ for runaway growth every SCAN_INTERVAL."""
from aipass.prax.apps.handlers.monitoring.rate_tracker import scan_rates, SCAN_INTERVAL
while not _stop_event.is_set():
try:
scan_rates()
except Exception as exc:
logger.info("[monitor] Rate tracker scan error: %s", exc)
_stop_event.wait(SCAN_INTERVAL)
def _handle_interactive_cmd(cmd: str, get_help_text) -> None:
"""Dispatch an interactive monitor command."""
if cmd == "help":
+4 -4
View File
@@ -317,14 +317,14 @@ class TestRenderEvent:
class TestThreadManagement:
"""Test thread start and stop functions."""
def test_start_threads_creates_three_threads(self):
"""_start_threads creates and starts display, file watcher, and log watcher threads."""
def test_start_threads_creates_four_threads(self):
"""_start_threads creates and starts display, file watcher, log watcher, and rate tracker threads."""
mod = _import_monitor()
mock_thread = MagicMock()
with patch("threading.Thread", return_value=mock_thread) as mock_cls:
mod._start_threads()
assert mock_cls.call_count == 3
assert mock_thread.start.call_count == 3
assert mock_cls.call_count == 4
assert mock_thread.start.call_count == 4
def test_stop_threads_sets_stop_event(self):
"""_stop_threads sets the stop event and stops the queue."""
+509
View File
@@ -0,0 +1,509 @@
# =================== AIPass ====================
# Name: test_rate_tracker.py
# Description: Tests for the rate tracker runaway-log detector
# Version: 1.0.0
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""Tests for apps/handlers/monitoring/rate_tracker.py
Covers:
- Rate calculation from byte offset changes
- Sustained threshold detection (WARNING and CRITICAL)
- Subsidence reset when rate drops
- Per-file suppression
- Event firing via trigger
- File disappearance handling
- get_snapshot() and reset()
"""
import sys
import time
from pathlib import Path
from unittest.mock import MagicMock, patch
_HANDLER_MOCKS = {
"aipass.prax.apps.handlers.json": MagicMock(),
"aipass.prax.apps.handlers.json.json_handler": MagicMock(),
}
def _import_tracker(monkeypatch):
"""Import (or reload) rate_tracker with handler mocks."""
monkeypatch.delenv("PYTEST_CURRENT_TEST", raising=False)
fresh = {k: MagicMock() for k in _HANDLER_MOCKS}
with patch.dict(sys.modules, fresh):
import importlib
if "aipass.prax.apps.handlers.monitoring.rate_tracker" in sys.modules:
mod = importlib.reload(sys.modules["aipass.prax.apps.handlers.monitoring.rate_tracker"])
else:
mod = importlib.import_module("aipass.prax.apps.handlers.monitoring.rate_tracker")
trigger_mock = MagicMock()
setattr(mod, "trigger", trigger_mock)
setattr(mod, "_HAS_TRIGGER", True)
mod.reset()
return mod, trigger_mock
class TestRateCalculation:
"""Rate is calculated from byte offset changes over elapsed time."""
def test_first_scan_initializes_no_rate(self, tmp_path, monkeypatch):
"""First scan seeds offsets — no rate calculated yet."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("line1\n" * 10)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
results = mod.scan_rates()
assert results == []
def test_second_scan_calculates_rate(self, tmp_path, monkeypatch):
"""Second scan with growth produces a rate."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
log_file.write_text("x" * 1300)
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(mod.time, "time", return_value=time.time() + 10.0),
):
results = mod.scan_rates()
assert len(results) == 1
assert results[0]["rate_lines_per_min"] > 0
def test_no_growth_produces_zero_rate(self, tmp_path, monkeypatch):
"""File that hasn't grown has rate 0."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(mod.time, "time", return_value=time.time() + 10.0),
):
results = mod.scan_rates()
assert len(results) == 1
assert results[0]["rate_lines_per_min"] == 0.0
def test_truncated_file_resets_offset(self, tmp_path, monkeypatch):
"""File that shrinks (rotation) resets offset without error."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 10000)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
log_file.write_text("x" * 100)
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(mod.time, "time", return_value=time.time() + 10.0),
):
results = mod.scan_rates()
assert results == []
class TestSustainedThresholds:
"""Events fire only after sustained intervals above threshold."""
def _grow_file(self, log_file, bytes_per_interval, mod, intervals):
"""Simulate N intervals of growth at a given rate."""
base_time = time.time()
results = []
for i in range(intervals):
log_file.write_bytes(b"x" * bytes_per_interval + log_file.read_bytes())
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
),
):
results = mod.scan_rates()
return results
def test_warning_fires_after_sustained_intervals(self, tmp_path, monkeypatch):
"""WARNING fires after WARNING_SUSTAINED_INTERVALS above WARNING_LINES_PER_MIN."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
self._grow_file(log_file, bytes_per_interval, mod, mod.WARNING_SUSTAINED_INTERVALS)
trigger_mock.fire.assert_called_once()
call_args = trigger_mock.fire.call_args
assert call_args[0][0] == "runaway_log_detected"
assert call_args[1]["severity"] == "warning"
def test_warning_does_not_fire_before_sustained(self, tmp_path, monkeypatch):
"""WARNING does not fire before reaching sustained count."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
self._grow_file(log_file, bytes_per_interval, mod, mod.WARNING_SUSTAINED_INTERVALS - 1)
trigger_mock.fire.assert_not_called()
def test_critical_fires_after_sustained_intervals(self, tmp_path, monkeypatch):
"""CRITICAL fires after CRITICAL_SUSTAINED_INTERVALS above CRITICAL_LINES_PER_MIN."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.CRITICAL_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
self._grow_file(log_file, bytes_per_interval, mod, mod.CRITICAL_SUSTAINED_INTERVALS)
assert trigger_mock.fire.call_count == 1
call_args = trigger_mock.fire.call_args
assert call_args[1]["severity"] == "critical"
def test_fires_only_once_until_subsides(self, tmp_path, monkeypatch):
"""Event fires once — not again on continued high rate."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
self._grow_file(log_file, bytes_per_interval, mod, mod.WARNING_SUSTAINED_INTERVALS + 5)
assert trigger_mock.fire.call_count == 1
class TestSubsidence:
"""Rate dropping below threshold resets sustained counters."""
def test_subsidence_resets_and_allows_refire(self, tmp_path, monkeypatch):
"""After rate drops and rises again, event can fire again."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
base_time = time.time()
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
for i in range(mod.WARNING_SUSTAINED_INTERVALS):
log_file.write_bytes(b"x" * bytes_per_interval + log_file.read_bytes())
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
),
):
mod.scan_rates()
assert trigger_mock.fire.call_count == 1
idle_offset = mod.WARNING_SUSTAINED_INTERVALS + 1
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + idle_offset * mod.SCAN_INTERVAL,
),
):
mod.scan_rates()
for i in range(mod.WARNING_SUSTAINED_INTERVALS):
log_file.write_bytes(b"x" * bytes_per_interval + log_file.read_bytes())
offset = idle_offset + i + 1
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + offset * mod.SCAN_INTERVAL,
),
):
mod.scan_rates()
assert trigger_mock.fire.call_count == 2
class TestSuppression:
"""Per-file suppression skips configured files."""
def test_suppressed_file_not_tracked(self, tmp_path, monkeypatch):
"""Suppressed files are skipped entirely."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "noisy_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 10000)
mod.configure_suppression({"noisy_module.log"})
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
mod.scan_rates()
assert "noisy_module.log" not in {Path(k).name for k in mod._tracked}
def test_non_suppressed_file_tracked(self, tmp_path, monkeypatch):
"""Non-suppressed files are tracked normally."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "normal_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
mod.configure_suppression({"other_module.log"})
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
assert any("normal_module.log" in k for k in mod._tracked)
class TestEventPayload:
"""Event payload carries correct fields."""
def test_event_payload_fields(self, tmp_path, monkeypatch):
"""Fired event includes all required fields."""
mod, trigger_mock = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
base_time = time.time()
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
for i in range(mod.WARNING_SUSTAINED_INTERVALS):
log_file.write_bytes(b"x" * bytes_per_interval + log_file.read_bytes())
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
),
):
mod.scan_rates()
call_kwargs = trigger_mock.fire.call_args[1]
assert "file_path" in call_kwargs
assert "rate_lines_per_min" in call_kwargs
assert "sustained_duration_sec" in call_kwargs
assert "severity" in call_kwargs
assert "branch" in call_kwargs
assert call_kwargs["sustained_duration_sec"] > 0
class TestFileDisappearance:
"""Deleted files are cleaned from tracking state."""
def test_deleted_file_removed_from_tracking(self, tmp_path, monkeypatch):
"""File removed between scans is cleaned from _tracked."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "ephemeral.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
assert any("ephemeral.log" in k for k in mod._tracked)
log_file.unlink()
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
assert not any("ephemeral.log" in k for k in mod._tracked)
class TestSnapshot:
"""get_snapshot() returns current state without scanning."""
def test_snapshot_returns_tracked_files(self, tmp_path, monkeypatch):
"""Snapshot includes files from a previous scan."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
snapshot = mod.get_snapshot()
assert len(snapshot) == 1
assert snapshot[0]["file"] == "test_module.log"
def test_snapshot_empty_before_any_scan(self, monkeypatch):
"""Snapshot is empty before any scan."""
mod, _ = _import_tracker(monkeypatch)
assert mod.get_snapshot() == []
class TestReset:
"""reset() clears all state."""
def test_reset_clears_tracked(self, tmp_path, monkeypatch):
"""After reset, tracked dict is empty."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
assert len(mod._tracked) > 0
mod.reset()
assert len(mod._tracked) == 0
class TestNoTrigger:
"""When trigger is unavailable, detection still works — just no event fired."""
def test_detection_without_trigger(self, tmp_path, monkeypatch):
"""Rate tracking and threshold detection work without trigger."""
mod, _ = _import_tracker(monkeypatch)
setattr(mod, "_HAS_TRIGGER", False)
setattr(mod, "trigger", None)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
base_time = time.time()
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
bytes_per_interval = int(mod.WARNING_LINES_PER_MIN * mod.AVG_LINE_BYTES * mod.SCAN_INTERVAL / 60 * 1.5)
results = []
for i in range(mod.WARNING_SUSTAINED_INTERVALS):
log_file.write_bytes(b"x" * bytes_per_interval + log_file.read_bytes())
with (
patch.object(mod, "get_system_logs_dir", return_value=log_file.parent),
patch.object(
mod.time,
"time",
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
),
):
results = mod.scan_rates()
assert any(r.get("severity") == "warning" for r in results)
class TestPersistence:
"""State persists to disk so CLI invocations and restarts work."""
def test_scan_saves_state_to_disk(self, tmp_path, monkeypatch):
"""scan_rates() calls save_json after scanning."""
mod, _ = _import_tracker(monkeypatch)
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 100)
json_handler_mock = mod.json_handler
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
mod.scan_rates()
json_handler_mock.save_json.assert_called()
call_args = json_handler_mock.save_json.call_args
assert call_args[0][0] == "rate_tracker"
assert call_args[0][1] == "data"
saved_data = call_args[0][2]
assert "files" in saved_data
assert any("test_module.log" in k for k in saved_data["files"])
def test_load_restores_offsets_from_disk(self, tmp_path, monkeypatch):
"""Loading persisted state restores file offsets so second scan can compute rates."""
mod, _ = _import_tracker(monkeypatch)
persisted = {
"module_name": "rate_tracker",
"files": {
str(tmp_path / "system" / "test_module.log"): {
"last_offset": 100,
"last_check": time.time() - 15.0,
"warning_sustained": 0,
"critical_sustained": 0,
"fired_warning": False,
"fired_critical": False,
}
},
}
mod.json_handler.load_json.return_value = persisted
log_file = tmp_path / "system" / "test_module.log"
log_file.parent.mkdir(parents=True)
log_file.write_text("x" * 1300)
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
results = mod.scan_rates()
assert len(results) == 1
assert results[0]["rate_lines_per_min"] > 0
def test_load_handles_missing_data(self, monkeypatch):
"""Loading when no data file exists is a no-op."""
mod, _ = _import_tracker(monkeypatch)
mod.json_handler.load_json.return_value = None
mod._load_state()
assert len(mod._tracked) == 0
def test_reset_clears_state_loaded_flag(self, monkeypatch):
"""reset() clears _state_loaded so next scan reloads from disk."""
mod, _ = _import_tracker(monkeypatch)
setattr(mod, "_state_loaded", True)
mod.reset()
assert mod._state_loaded is False
+39
View File
@@ -557,6 +557,45 @@
"standard": "unused_function",
"pattern": "get_muted_branches",
"reason": "Public API returning List[str] of active muted branches. Used by 15+ existing tests and part of the medic_state interface. get_muted_branches_detail() supplements it for status display — this is the simple accessor."
},
{
"file": "apps/handlers/events/runaway_handler.py",
"standard": "silent_catch",
"pattern": "_log_warning except",
"reason": "Meta-logging helper: _log_warning() writes directly to file. Its own except block cannot log — you cannot log a failure to log."
},
{
"file": "apps/handlers/events/runaway_handler.py",
"standard": "error_handling",
"lines": [52],
"pattern": "except Exception: pass",
"reason": "Meta-logging helper _log_warning() — cannot log a failure to log. Same pattern as silent_catch bypass."
},
{
"file": "apps/handlers/events/runaway_handler.py",
"standard": "handlers",
"lines": [272],
"pattern": "cross-handler import wake_branch",
"reason": "Inline import of wake_branch for dispatch — same pattern as error_detected.py line 572. Must wake target branch after email delivery."
},
{
"file": "apps/handlers/events/runaway_handler.py",
"standard": "encapsulation",
"lines": [272],
"pattern": "cross-handler import wake_branch",
"reason": "Inline import of wake_branch for dispatch — same pattern as error_detected.py. Handler must wake target branch after email delivery."
},
{
"file": "tests/test_runaway_handler.py",
"standard": "architecture",
"pattern": "3-layer structure",
"reason": "Test file — tests/ is the standard location for unit tests, not part of the apps/modules/handlers source tree."
},
{
"file": "tests/test_runaway_handler.py",
"standard": "encapsulation",
"pattern": "handler imported directly",
"reason": "Test must import the handler module directly to test it. All trigger test files follow this pattern."
}
],
"notes": {
+9 -7
View File
@@ -74,7 +74,7 @@ result = report_error(
## Events
15 events defined, 13 active (2 decommissioned by TDPLAN-0007). Registered via `handlers/events/registry.py` on first `Trigger.fire()`. All fire through the event bus.
16 events defined, 14 active (2 decommissioned by TDPLAN-0007). Registered via `handlers/events/registry.py` on first `Trigger.fire()`. All fire through the event bus.
| Event | Handler | Trigger | Action |
|-------|---------|---------|--------|
@@ -92,6 +92,7 @@ result = report_error(
| `cli_header_displayed` | `cli.py` | CLI displays headers | Registration hook |
| `pr_created` | `pr_status_sync.py` | PR opened on GitHub | ~~Runs `drone @prax status sync`~~ **Decommissioned** (TDPLAN-0007) |
| `pr_merged` | `pr_status_sync.py` | PR merged on GitHub | ~~Runs `drone @prax status sync`~~ **Decommissioned** (TDPLAN-0007) |
| `runaway_log_detected` | `runaway_handler.py` | Prax rate tracker detects sustained high log volume | Per-file cooldown dispatch to responsible branch; UNKNOWN attribution falls back to @prax; writes alert to `.aipass/alerts.json` |
| `memory_pool_auto_processed` | `memory_pool.py` | Hook engine runs `auto_process()` | Logs result; on failure fires `error_detected` for Medic dispatch |
## Medic
@@ -148,7 +149,7 @@ trigger/
│ ├── json/
│ │ └── json_handler.py # JSON structure logging
│ ├── events/
│ │ ├── registry.py # Auto-registers 13 active event handlers
│ │ ├── registry.py # Auto-registers 14 active event handlers
│ │ ├── startup.py # Startup catch-up scan
│ │ ├── error_detected.py # 8-gate Medic dispatch
│ │ ├── error_logged.py # Monitor-only (no dispatch)
@@ -159,11 +160,12 @@ trigger/
│ │ ├── memory_template_updated.py
│ │ ├── memory.py # memory_saved placeholder
│ │ ├── cli.py # cli_header_displayed hook
│ │ ├── runaway_handler.py # Runaway log dispatch (per-file cooldown, independent of Medic)
│ │ ├── pr_status_sync.py # PR → prax status sync (decommissioned TDPLAN-0007)
│ │ └── memory_pool.py # Pool auto-process observability
│ └── watchers/
│ └── log_watcher.py # System log watcher (system_logs/ dir)
├── tests/ # 563 tests across 19 modules
├── tests/ # 619 tests across 20 modules
├── trigger_json/ # Runtime state files
│ ├── trigger_config.json # Medic state, muted branches
│ ├── error_registry.json # All tracked errors
@@ -191,21 +193,21 @@ trigger/
## Testing
575 tests across 19 test modules, all passing. Coverage: 76/76 public functions (100%).
619 tests across 20 test modules, all passing. Coverage: 81/81 public functions (100%).
```bash
cd src/aipass/trigger && pytest # Run all tests
```
Test files: `test_core`, `test_errors`, `test_medic`, `test_error_registry`, `test_error_reporter`, `test_medic_state`, `test_log_watcher`, `test_watchers_log_watcher`, `test_branch_log_events`, `test_log_events`, `test_json_handler`, `test_pr_status_sync`, `test_error_detected`, `test_event_handlers`, `test_log_watcher_service`, `test_plan_file_handler`, `test_startup_handler`, `test_trigger_entry`, `test_memory_pool_handler`
Test files: `test_core`, `test_errors`, `test_medic`, `test_error_registry`, `test_error_reporter`, `test_medic_state`, `test_log_watcher`, `test_watchers_log_watcher`, `test_branch_log_events`, `test_log_events`, `test_json_handler`, `test_pr_status_sync`, `test_error_detected`, `test_event_handlers`, `test_log_watcher_service`, `test_plan_file_handler`, `test_startup_handler`, `test_trigger_entry`, `test_memory_pool_handler`, `test_runaway_handler`
## Compliance
Seedgo: 100% (34/34 standards). Zero type errors. All categories at 100%.
Seedgo: 100% (41/41 standards). Zero type errors. All categories at 100%.
---
*Last Updated: 2026-06-06*
*Last Updated: 2026-07-14*
---
[← Back to AIPass](../../../README.md)
@@ -36,6 +36,7 @@ def setup_handlers():
from .cli import handle_cli_header_displayed
from .plan_file import handle_plan_file_created, handle_plan_file_deleted, handle_plan_file_moved
from .error_detected import handle_error_detected, set_send_email_callback
from .runaway_handler import handle_runaway_log_detected, set_send_email_callback as set_runaway_email_callback
# Wire up email send callback for error_detected handler (avoids handler importing from modules)
try:
@@ -60,6 +61,7 @@ def setup_handlers():
return success
set_send_email_callback(_send_email_adapter)
set_runaway_email_callback(_send_email_adapter)
except ImportError:
_log_warning("ai_mail not available — error notifications won't send")
from .warning_logged import handle_warning_logged
@@ -79,5 +81,6 @@ def setup_handlers():
# trigger.on("pr_created", handle_pr_created) # TDPLAN-0007: status-sync decommissioned
# trigger.on("pr_merged", handle_pr_merged) # TDPLAN-0007: status-sync decommissioned
trigger.on("memory_pool_auto_processed", handle_memory_pool_auto_processed)
trigger.on("runaway_log_detected", handle_runaway_log_detected)
json_handler.log_operation("handlers_registered", {"success": True})
@@ -0,0 +1,284 @@
# =================== AIPass ====================
# Name: runaway_handler.py
# Description: Runaway log event handler with per-file cooldown gating
# Version: 1.0.0
# Created: 2026-07-14
# Modified: 2026-07-14
# =============================================
"""
Runaway Log Detected Event Handler
Handles runaway_log_detected events fired by prax's rate tracker.
Volume-based detection (rate of log output), orthogonal to error_detected
(content-based ERROR line matching).
Event payload from prax:
- file_path: Path to the runaway log file
- rate_lines_per_min: Current log rate
- sustained_duration_sec: How long the rate has been sustained
- severity: "warning" or "critical"
- branch: Responsible branch name
Gating:
- Per-file cooldown (30min default) — independent of medic circuit breaker
- Branch mute check (reuses TTL mute infrastructure from trigger_config.json)
- UNKNOWN/missing branch → dispatch to @prax as fallback
"""
import json
import time
import uuid
from datetime import datetime
from pathlib import Path
from typing import Any, Callable, Optional
from aipass.trigger.apps.config import TRIGGER_ROOT, atomic_write_json, json_file_lock
from aipass.trigger.apps.handlers.json import json_handler
try:
from aipass.prax import append_jsonl as _append_jsonl
except Exception:
_append_jsonl = None
_HANDLER_LOG = TRIGGER_ROOT / "logs" / "runaway_handler.jsonl"
def _log_warning(message: str) -> None:
"""Log warning to file (recursion-safe prax path)."""
if _append_jsonl is None:
return
try:
_append_jsonl(_HANDLER_LOG, {"level": "WARNING", "msg": message})
except Exception:
pass # seedgo:bypass meta-logging
def _find_repo_root() -> Path:
"""Walk up from this file to find the repo root (contains AIPASS_REGISTRY.json)."""
current = Path(__file__).resolve().parent
for parent in [current] + list(current.parents):
if (parent / "AIPASS_REGISTRY.json").exists():
return parent
return Path.cwd()
_REPO_ROOT = _find_repo_root()
ALERTS_FILE = _REPO_ROOT / ".aipass" / "alerts.json"
TRIGGER_CONFIG_FILE = TRIGGER_ROOT / "trigger_json" / "trigger_config.json"
_send_email: Optional[Callable[..., bool]] = None
_file_cooldowns: dict[str, float] = {}
COOLDOWN_SECONDS = 1800
def set_send_email_callback(callback: Callable[..., bool]) -> None:
"""Set the callback function for sending emails.
Must be called by the registry layer before events fire.
Args:
callback: Function matching deliver_email_to_branch adapter signature
"""
global _send_email
_send_email = callback
def _is_file_on_cooldown(file_path: str) -> bool:
"""Check if a file is still within its dispatch cooldown window.
Args:
file_path: Path to the log file
Returns:
True if cooldown has not expired
"""
last = _file_cooldowns.get(file_path, 0.0)
return (time.time() - last) < COOLDOWN_SECONDS
def _record_file_dispatch(file_path: str) -> None:
"""Record a dispatch timestamp for per-file cooldown.
Args:
file_path: Path to the log file that was dispatched
"""
_file_cooldowns[file_path] = time.time()
def _mute_entry_matches(entry, branch_lower: str, now: datetime) -> bool:
"""Check if a single mute entry matches the branch and is still active."""
if isinstance(entry, str):
return entry.lower() == branch_lower
if not isinstance(entry, dict):
return False
if entry.get("name", "").lower() != branch_lower:
return False
expires_at = entry.get("expires_at")
if expires_at is None:
return True
return datetime.fromisoformat(expires_at) > now
def _is_branch_muted(branch_name: str) -> bool:
"""Check if a branch is muted for dispatch.
Reads muted_branches from trigger_config.json. Supports both
plain-string entries (permanent) and dict entries with TTL.
Args:
branch_name: Branch name (case-insensitive)
Returns:
True if branch is actively muted
"""
try:
if not TRIGGER_CONFIG_FILE.exists():
return False
data = json.loads(TRIGGER_CONFIG_FILE.read_text(encoding="utf-8"))
muted = data.get("config", {}).get("muted_branches", [])
branch_lower = branch_name.lower()
now = datetime.now()
return any(_mute_entry_matches(e, branch_lower, now) for e in muted)
except Exception as exc:
_log_warning(f"_is_branch_muted config read failed: {exc}")
return False
def _write_suppression_log(reason: str, file_path: str, branch: str) -> None:
"""Write a line to the runaway suppression log."""
if _append_jsonl is None:
return
try:
suppressed_log = TRIGGER_ROOT / "logs" / "runaway_suppressed.jsonl"
entry = {
"ts": datetime.now().isoformat(),
"reason": reason,
"file": file_path,
"branch": branch,
}
_append_jsonl(suppressed_log, entry)
except Exception as exc:
_log_warning(f"suppression log write failed ({reason}): {exc}")
def _write_alert(file_path: str, severity: str, branch: str, rate: float, duration: float) -> None:
"""Write an alert entry to .aipass/alerts.json.
Args:
file_path: Path to the runaway log file
severity: "warning" or "critical"
branch: Responsible branch name
rate: Lines per minute
duration: Sustained duration in seconds
"""
try:
alert = {
"id": str(uuid.uuid4()),
"source": "prax",
"severity": severity,
"title": f"Runaway log: {Path(file_path).name}",
"body": (
f"Log file {file_path} producing {rate:.0f} lines/min sustained {duration:.0f}s. Branch: {branch}."
),
"created_at": datetime.now().isoformat(),
"expires_at": None,
}
ALERTS_FILE.parent.mkdir(parents=True, exist_ok=True)
with json_file_lock(ALERTS_FILE):
existing = {"alerts": []}
if ALERTS_FILE.exists():
raw = ALERTS_FILE.read_text(encoding="utf-8").strip()
if raw:
existing = json.loads(raw)
existing.setdefault("alerts", []).append(alert)
atomic_write_json(ALERTS_FILE, existing)
except Exception as exc:
_log_warning(f"_write_alert failed: {exc}")
def handle_runaway_log_detected(
file_path: str | None = None,
rate_lines_per_min: float = 0,
sustained_duration_sec: float = 0,
severity: str = "warning",
branch: str | None = None,
**kwargs: Any,
) -> None:
"""Handle runaway_log_detected event — dispatch to responsible branch.
Volume-based detection, independent of medic error_detected pipeline.
Uses per-file cooldown (30min) instead of the medic circuit breaker.
Args:
file_path: Path to the runaway log file — REQUIRED
rate_lines_per_min: Current log rate
sustained_duration_sec: How long the rate has been sustained
severity: "warning" or "critical"
branch: Responsible branch name (None/UNKNOWN → dispatch to @prax)
**kwargs: Additional event data (ignored)
"""
try:
if not file_path:
return
if _is_file_on_cooldown(file_path):
_write_suppression_log("cooldown", file_path, branch or "UNKNOWN")
return
is_unknown = not branch or branch.upper() == "UNKNOWN"
target_branch = branch or "UNKNOWN"
if not is_unknown and _is_branch_muted(target_branch):
_write_suppression_log("branch_muted", file_path, target_branch)
return
if _send_email is None:
_log_warning("No email callback — cannot dispatch runaway alert")
return
recipient = "@prax" if is_unknown else f"@{target_branch.lower()}"
subject = f"[RUNAWAY] {Path(file_path).name} — {severity.upper()}"
message = (
f"Runaway log detected.\n\n"
f"File: {file_path}\n"
f"Rate: {rate_lines_per_min:.0f} lines/min\n"
f"Sustained: {sustained_duration_sec:.0f}s\n"
f"Severity: {severity}\n"
f"Branch: {target_branch}\n\n"
f"---\n"
f"INVESTIGATION STEPS:\n"
f"1. Identify the process writing to this log\n"
f"2. Check for spin loops, retry storms, or misconfigured log levels\n"
f"3. Fix the root cause or kill the offending process\n"
f"4. Report to @devpulse\n"
)
sent = _send_email(
to_branch=recipient,
subject=subject,
message=message,
auto_execute=True,
reply_to="@devpulse",
from_branch="@trigger",
)
if not sent:
_log_warning(f"Email delivery failed for {recipient} ({file_path})")
return
try:
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
wake_branch(recipient, fresh=False, sender="@trigger")
except Exception:
pass # Email in inbox as fallback
_write_alert(file_path, severity, target_branch, rate_lines_per_min, sustained_duration_sec)
_record_file_dispatch(file_path)
json_handler.log_operation("runaway_dispatch_sent", {"recipient": recipient, "file": file_path})
except Exception as exc:
_log_warning(f"handle_runaway_log_detected failed: {exc}")
@@ -0,0 +1,467 @@
"""Tests for runaway_log_detected event handler."""
import json
import sys
from pathlib import Path
from unittest.mock import MagicMock, patch
import pytest
from aipass.trigger.apps.handlers.events import runaway_handler as mod
# ---------------------------------------------------------------------------
# Shared fixture: redirect file paths to tmp_path, mock _append_jsonl and
# wake_branch, clear cooldown state between tests.
# ---------------------------------------------------------------------------
@pytest.fixture(autouse=True)
def _reset_state(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): # type: ignore[misc]
"""Reset module state and redirect file paths to tmp_path."""
mod._file_cooldowns.clear()
mod._send_email = None
monkeypatch.setattr(mod, "TRIGGER_CONFIG_FILE", tmp_path / "trigger_config.json")
monkeypatch.setattr(mod, "ALERTS_FILE", tmp_path / "alerts.json")
monkeypatch.setattr(mod, "_append_jsonl", MagicMock())
# Mock wake_branch import chain so the in-function import succeeds
mock_wake_mod = MagicMock()
mock_wake_mod.wake_branch = MagicMock()
monkeypatch.setitem(sys.modules, "aipass.ai_mail", MagicMock())
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps", MagicMock())
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps.handlers", MagicMock())
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps.handlers.dispatch", MagicMock())
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps.handlers.dispatch.wake", mock_wake_mod)
yield
mod._file_cooldowns.clear()
def _setup_happy_path() -> MagicMock:
"""Set up a successful dispatch scenario and return the send_email mock."""
send_mock = MagicMock(return_value=True)
mod.set_send_email_callback(send_mock)
return send_mock
# ---------------------------------------------------------------------------
# 1. Missing file_path — returns without dispatch
# ---------------------------------------------------------------------------
class TestMissingFilePath:
"""Handler returns early when file_path is missing."""
def test_none_file_path_no_dispatch(self) -> None:
"""Returns without dispatch when file_path is None."""
send = _setup_happy_path()
mod.handle_runaway_log_detected(file_path=None, branch="flow")
send.assert_not_called()
def test_empty_file_path_no_dispatch(self) -> None:
"""Returns without dispatch when file_path is empty string."""
send = _setup_happy_path()
mod.handle_runaway_log_detected(file_path="", branch="flow")
send.assert_not_called()
# ---------------------------------------------------------------------------
# 2. Per-file cooldown — second call within 30min is suppressed
# ---------------------------------------------------------------------------
class TestPerFileCooldown:
"""Second call for the same file within 30min cooldown is suppressed."""
def test_second_call_within_cooldown_suppressed(self) -> None:
"""Second call for the same file is suppressed (no time mock needed)."""
send = _setup_happy_path()
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
assert send.call_count == 1
# Second call — same file, should be suppressed
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=120,
)
assert send.call_count == 1
# ---------------------------------------------------------------------------
# 3. Cooldown expired — call after cooldown passes dispatches again
# ---------------------------------------------------------------------------
class TestCooldownExpired:
"""Call after cooldown window expires dispatches again."""
@patch("aipass.trigger.apps.handlers.events.runaway_handler.time")
def test_dispatches_again_after_cooldown_expires(self, mock_time: MagicMock) -> None:
"""Dispatch succeeds again once the 1800s cooldown has elapsed."""
send = _setup_happy_path()
mock_time.time.return_value = 1_000_000.0
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
assert send.call_count == 1
# Advance past 1800s cooldown
mock_time.time.return_value = 1_000_000.0 + 1801
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=120,
)
assert send.call_count == 2
# ---------------------------------------------------------------------------
# 4. Branch muted — muted branch is suppressed
# ---------------------------------------------------------------------------
class TestBranchMuted:
"""Muted branch dispatch is suppressed."""
def test_muted_branch_suppressed(self, tmp_path: Path) -> None:
"""Branch listed in muted_branches is suppressed — no email sent."""
send = _setup_happy_path()
config_file = tmp_path / "trigger_config.json"
config_file.write_text(
json.dumps({"config": {"muted_branches": ["flow"]}}),
encoding="utf-8",
)
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
send.assert_not_called()
# ---------------------------------------------------------------------------
# 5. UNKNOWN branch — dispatches to @prax instead
# ---------------------------------------------------------------------------
class TestUnknownBranch:
"""UNKNOWN branch falls back to @prax."""
def test_unknown_branch_dispatches_to_prax(self) -> None:
"""Branch='UNKNOWN' dispatches email to @prax."""
send = _setup_happy_path()
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="UNKNOWN",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
send.assert_called_once()
assert send.call_args[1]["to_branch"] == "@prax"
# ---------------------------------------------------------------------------
# 6. None branch — dispatches to @prax instead
# ---------------------------------------------------------------------------
class TestNoneBranch:
"""None branch falls back to @prax."""
def test_none_branch_dispatches_to_prax(self) -> None:
"""Branch=None dispatches email to @prax."""
send = _setup_happy_path()
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch=None,
rate_lines_per_min=500,
sustained_duration_sec=60,
)
send.assert_called_once()
assert send.call_args[1]["to_branch"] == "@prax"
# ---------------------------------------------------------------------------
# 7. No email callback — logs warning, no dispatch
# ---------------------------------------------------------------------------
class TestNoEmailCallback:
"""Handler logs warning and returns when _send_email is None."""
def test_logs_warning_no_dispatch(self) -> None:
"""Logs warning via _append_jsonl when no callback set."""
# _send_email stays None (no set_send_email_callback call)
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
calls = mod._append_jsonl.call_args_list # type: ignore[union-attr]
warning_calls = [
c
for c in calls
if isinstance(c[0][1], dict) and c[0][1].get("level") == "WARNING"
]
assert len(warning_calls) >= 1
assert "No email callback" in warning_calls[0][0][1]["msg"]
# ---------------------------------------------------------------------------
# 8. Successful dispatch — email sent, wake called, alert written,
# cooldown recorded
# ---------------------------------------------------------------------------
class TestSuccessfulDispatch:
"""Full happy-path: email, wake, alert, cooldown."""
def test_full_dispatch(self, tmp_path: Path) -> None:
"""Email sent with correct kwargs, alert file exists, cooldown recorded."""
send = _setup_happy_path()
file_path = "/var/log/test.log"
mod.handle_runaway_log_detected(
file_path=file_path,
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
severity="critical",
)
# Email sent with expected kwargs
send.assert_called_once()
kwargs = send.call_args[1]
assert kwargs["to_branch"] == "@flow"
assert kwargs["auto_execute"] is True
assert kwargs["reply_to"] == "@devpulse"
assert kwargs["from_branch"] == "@trigger"
assert "[RUNAWAY]" in kwargs["subject"]
assert "CRITICAL" in kwargs["subject"]
# wake_branch called (via mocked import)
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
wake_branch.assert_called_once_with("@flow", fresh=False, sender="@trigger") # type: ignore[union-attr]
# Alert file written
alerts_file = tmp_path / "alerts.json"
assert alerts_file.exists()
# Cooldown recorded
assert file_path in mod._file_cooldowns
# ---------------------------------------------------------------------------
# 9. Alert file written — verify alerts.json schema
# ---------------------------------------------------------------------------
class TestAlertFileSchema:
"""Alert written to .aipass/alerts.json with correct schema."""
def test_alert_has_required_fields(self, tmp_path: Path) -> None:
"""Schema: {alerts: [{id, source, severity, title, body, created_at, expires_at}]}."""
_setup_happy_path()
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
severity="warning",
)
alerts_file = tmp_path / "alerts.json"
data = json.loads(alerts_file.read_text(encoding="utf-8"))
assert "alerts" in data
assert len(data["alerts"]) == 1
alert = data["alerts"][0]
required_keys = {"id", "source", "severity", "title", "body", "created_at", "expires_at"}
assert required_keys == set(alert.keys())
assert alert["source"] == "prax"
assert alert["severity"] == "warning"
assert "test.log" in alert["title"]
assert alert["body"]
assert alert["created_at"]
# ---------------------------------------------------------------------------
# 10. Alert appends — existing alerts preserved when new one appended
# ---------------------------------------------------------------------------
class TestAlertAppends:
"""Existing alerts are preserved when a new alert is appended."""
def test_existing_alerts_preserved(self, tmp_path: Path) -> None:
"""Pre-populated alerts.json keeps existing entries after append."""
alerts_file = tmp_path / "alerts.json"
existing_alert = {
"id": "existing-123",
"source": "medic",
"severity": "critical",
"title": "Existing alert",
"body": "Some body",
"created_at": "2026-01-01T00:00:00",
"expires_at": None,
}
alerts_file.write_text(
json.dumps({"alerts": [existing_alert]}),
encoding="utf-8",
)
_setup_happy_path()
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
data = json.loads(alerts_file.read_text(encoding="utf-8"))
assert len(data["alerts"]) == 2
assert data["alerts"][0]["id"] == "existing-123"
assert data["alerts"][1]["source"] == "prax"
# ---------------------------------------------------------------------------
# 11. Email send fails — returns early, no alert written, no cooldown
# ---------------------------------------------------------------------------
class TestEmailSendFails:
"""When _send_email returns False, no alert or cooldown is recorded."""
def test_returns_early_no_alert_no_cooldown(self, tmp_path: Path) -> None:
"""Failed email send means no alert file and no cooldown entry."""
send = MagicMock(return_value=False)
mod.set_send_email_callback(send)
file_path = "/var/log/test.log"
mod.handle_runaway_log_detected(
file_path=file_path,
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
send.assert_called_once()
alerts_file = tmp_path / "alerts.json"
assert not alerts_file.exists()
assert file_path not in mod._file_cooldowns
# ---------------------------------------------------------------------------
# 12. Suppression log written — cooldown and mute both write suppression log
# ---------------------------------------------------------------------------
class TestSuppressionLog:
"""Cooldown and mute suppressions write to the suppression log."""
def test_cooldown_writes_suppression_log(self) -> None:
"""Cooldown suppression writes reason='cooldown' via _append_jsonl."""
_setup_happy_path()
file_path = "/var/log/test.log"
# First call dispatches normally
mod.handle_runaway_log_detected(
file_path=file_path,
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
# Reset mock to isolate suppression log call
mod._append_jsonl.reset_mock() # type: ignore[union-attr]
# Second call is on cooldown — should write suppression log
mod.handle_runaway_log_detected(
file_path=file_path,
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=120,
)
calls = mod._append_jsonl.call_args_list # type: ignore[union-attr]
suppression_calls = [
c
for c in calls
if isinstance(c[0][1], dict) and c[0][1].get("reason") == "cooldown"
]
assert len(suppression_calls) == 1
assert suppression_calls[0][0][1]["file"] == file_path
def test_mute_writes_suppression_log(self, tmp_path: Path) -> None:
"""Branch mute suppression writes reason='branch_muted' via _append_jsonl."""
_setup_happy_path()
config_file = tmp_path / "trigger_config.json"
config_file.write_text(
json.dumps({"config": {"muted_branches": ["flow"]}}),
encoding="utf-8",
)
mod.handle_runaway_log_detected(
file_path="/var/log/test.log",
branch="flow",
rate_lines_per_min=500,
sustained_duration_sec=60,
)
calls = mod._append_jsonl.call_args_list # type: ignore[union-attr]
suppression_calls = [
c
for c in calls
if isinstance(c[0][1], dict) and c[0][1].get("reason") == "branch_muted"
]
assert len(suppression_calls) == 1
assert suppression_calls[0][0][1]["branch"] == "flow"
# ---------------------------------------------------------------------------
# 13. set_send_email_callback — sets the callback correctly
# ---------------------------------------------------------------------------
class TestSetSendEmailCallback:
"""Tests for set_send_email_callback."""
def test_sets_callback_correctly(self) -> None:
"""Stores the callback as module-level _send_email."""
callback = MagicMock()
mod.set_send_email_callback(callback)
assert mod._send_email is callback
def test_overwrites_previous_callback(self) -> None:
"""Second call replaces the first callback."""
first = MagicMock()
second = MagicMock()
mod.set_send_email_callback(first)
mod.set_send_email_callback(second)
assert mod._send_email is second