From 81658ce0eab3dcdd7c1a1fb797afaef0b5cbd91d Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Wed, 15 Jul 2026 00:05:02 -0700 Subject: [PATCH] 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 --- .aipass/hooks.json | 5 + .aipass/tier1_navmap.md | 9 + CHANGELOG.md | 45 ++ README.md | 8 +- src/aipass/ai_mail/README.md | 7 +- .../handlers/dispatch/dispatch_monitor.py | 24 +- .../ai_mail/apps/handlers/dispatch/wake.py | 21 +- .../ai_mail/tests/test_dispatch_monitor.py | 159 ++---- src/aipass/ai_mail/tests/test_wake.py | 51 ++ src/aipass/hooks/README.md | 34 +- .../apps/handlers/prompt/persistent_alert.py | 131 +++++ .../hooks/apps/modules/alert_dismiss.py | 111 ++++ .../hooks/tests/test_persistent_alert.py | 425 +++++++++++++++ src/aipass/prax/CLOSED_PLANS.local.json | 7 + src/aipass/prax/README.md | 39 +- .../apps/handlers/monitoring/rate_tracker.py | 368 +++++++++++++ src/aipass/prax/apps/modules/log_health.py | 176 ++++++ src/aipass/prax/apps/modules/monitor.py | 25 +- src/aipass/prax/tests/test_monitor_module.py | 8 +- src/aipass/prax/tests/test_rate_tracker.py | 509 ++++++++++++++++++ src/aipass/trigger/.seedgo/bypass.json | 39 ++ src/aipass/trigger/README.md | 16 +- .../trigger/apps/handlers/events/registry.py | 3 + .../apps/handlers/events/runaway_handler.py | 284 ++++++++++ .../trigger/tests/test_runaway_handler.py | 467 ++++++++++++++++ 25 files changed, 2789 insertions(+), 182 deletions(-) create mode 100644 src/aipass/hooks/apps/handlers/prompt/persistent_alert.py create mode 100644 src/aipass/hooks/apps/modules/alert_dismiss.py create mode 100644 src/aipass/hooks/tests/test_persistent_alert.py create mode 100644 src/aipass/prax/apps/handlers/monitoring/rate_tracker.py create mode 100644 src/aipass/prax/apps/modules/log_health.py create mode 100644 src/aipass/prax/tests/test_rate_tracker.py create mode 100644 src/aipass/trigger/apps/handlers/events/runaway_handler.py create mode 100644 src/aipass/trigger/tests/test_runaway_handler.py diff --git a/.aipass/hooks.json b/.aipass/hooks.json index 3b32d2d4..96835ae7 100644 --- a/.aipass/hooks.json +++ b/.aipass/hooks.json @@ -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", diff --git a/.aipass/tier1_navmap.md b/.aipass/tier1_navmap.md index 8cfae959..106fefb1 100644 --- a/.aipass/tier1_navmap.md +++ b/.aipass/tier1_navmap.md @@ -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 diff --git a/CHANGELOG.md b/CHANGELOG.md index 7efd8a39..a12136ed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 `) — 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: diff --git a/README.md b/README.md index 597350a7..2879eab0 100644 --- a/README.md +++ b/README.md @@ -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 | diff --git a/src/aipass/ai_mail/README.md b/src/aipass/ai_mail/README.md index 446efb2d..97ed425e 100644 --- a/src/aipass/ai_mail/README.md +++ b/src/aipass/ai_mail/README.md @@ -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 diff --git a/src/aipass/ai_mail/apps/handlers/dispatch/dispatch_monitor.py b/src/aipass/ai_mail/apps/handlers/dispatch/dispatch_monitor.py index 406536b1..60d32d4e 100644 --- a/src/aipass/ai_mail/apps/handlers/dispatch/dispatch_monitor.py +++ b/src/aipass/ai_mail/apps/handlers/dispatch/dispatch_monitor.py @@ -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) diff --git a/src/aipass/ai_mail/apps/handlers/dispatch/wake.py b/src/aipass/ai_mail/apps/handlers/dispatch/wake.py index a76c477c..9501aac7 100644 --- a/src/aipass/ai_mail/apps/handlers/dispatch/wake.py +++ b/src/aipass/ai_mail/apps/handlers/dispatch/wake.py @@ -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) diff --git a/src/aipass/ai_mail/tests/test_dispatch_monitor.py b/src/aipass/ai_mail/tests/test_dispatch_monitor.py index 23aa3608..19bf9b92 100644 --- a/src/aipass/ai_mail/tests/test_dispatch_monitor.py +++ b/src/aipass/ai_mail/tests/test_dispatch_monitor.py @@ -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.""" diff --git a/src/aipass/ai_mail/tests/test_wake.py b/src/aipass/ai_mail/tests/test_wake.py index 00063866..e7cef1de 100644 --- a/src/aipass/ai_mail/tests/test_wake.py +++ b/src/aipass/ai_mail/tests/test_wake.py @@ -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) diff --git a/src/aipass/hooks/README.md b/src/aipass/hooks/README.md index 8f078815..52d4e2cd 100644 --- a/src/aipass/hooks/README.md +++ b/src/aipass/hooks/README.md @@ -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 ` | 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 ) │ │ ├── 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 ` 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* --- diff --git a/src/aipass/hooks/apps/handlers/prompt/persistent_alert.py b/src/aipass/hooks/apps/handlers/prompt/persistent_alert.py new file mode 100644 index 00000000..6a9848b6 --- /dev/null +++ b/src/aipass/hooks/apps/handlers/prompt/persistent_alert.py @@ -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 " + 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 diff --git a/src/aipass/hooks/apps/modules/alert_dismiss.py b/src/aipass/hooks/apps/modules/alert_dismiss.py new file mode 100644 index 00000000..4900ac92 --- /dev/null +++ b/src/aipass/hooks/apps/modules/alert_dismiss.py @@ -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 ", "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 ") + 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 ") + CONSOLE.print() + CONSOLE.print("Removes the alert with the given ID from .aipass/alerts.json.") + return True + + _dismiss_alert(args[0]) + return True diff --git a/src/aipass/hooks/tests/test_persistent_alert.py b/src/aipass/hooks/tests/test_persistent_alert.py new file mode 100644 index 00000000..782846f7 --- /dev/null +++ b/src/aipass/hooks/tests/test_persistent_alert.py @@ -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 diff --git a/src/aipass/prax/CLOSED_PLANS.local.json b/src/aipass/prax/CLOSED_PLANS.local.json index 1c4aefee..56074001 100644 --- a/src/aipass/prax/CLOSED_PLANS.local.json +++ b/src/aipass/prax/CLOSED_PLANS.local.json @@ -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": { diff --git a/src/aipass/prax/README.md b/src/aipass/prax/README.md index c0538621..ba1b8849 100644 --- a/src/aipass/prax/README.md +++ b/src/aipass/prax/README.md @@ -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 `/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) diff --git a/src/aipass/prax/apps/handlers/monitoring/rate_tracker.py b/src/aipass/prax/apps/handlers/monitoring/rate_tracker.py new file mode 100644 index 00000000..8c319121 --- /dev/null +++ b/src/aipass/prax/apps/handlers/monitoring/rate_tracker.py @@ -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 diff --git a/src/aipass/prax/apps/modules/log_health.py b/src/aipass/prax/apps/modules/log_health.py new file mode 100644 index 00000000..32cc0046 --- /dev/null +++ b/src/aipass/prax/apps/modules/log_health.py @@ -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) diff --git a/src/aipass/prax/apps/modules/monitor.py b/src/aipass/prax/apps/modules/monitor.py index bcdb6d13..3b4b0c71 100755 --- a/src/aipass/prax/apps/modules/monitor.py +++ b/src/aipass/prax/apps/modules/monitor.py @@ -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": diff --git a/src/aipass/prax/tests/test_monitor_module.py b/src/aipass/prax/tests/test_monitor_module.py index 7507759f..6d8200b3 100644 --- a/src/aipass/prax/tests/test_monitor_module.py +++ b/src/aipass/prax/tests/test_monitor_module.py @@ -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.""" diff --git a/src/aipass/prax/tests/test_rate_tracker.py b/src/aipass/prax/tests/test_rate_tracker.py new file mode 100644 index 00000000..0c786c5a --- /dev/null +++ b/src/aipass/prax/tests/test_rate_tracker.py @@ -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 diff --git a/src/aipass/trigger/.seedgo/bypass.json b/src/aipass/trigger/.seedgo/bypass.json index 2900bf23..23dbb43c 100644 --- a/src/aipass/trigger/.seedgo/bypass.json +++ b/src/aipass/trigger/.seedgo/bypass.json @@ -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": { diff --git a/src/aipass/trigger/README.md b/src/aipass/trigger/README.md index 41b5f112..30487dcc 100644 --- a/src/aipass/trigger/README.md +++ b/src/aipass/trigger/README.md @@ -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) diff --git a/src/aipass/trigger/apps/handlers/events/registry.py b/src/aipass/trigger/apps/handlers/events/registry.py index d39cb970..dbacd588 100644 --- a/src/aipass/trigger/apps/handlers/events/registry.py +++ b/src/aipass/trigger/apps/handlers/events/registry.py @@ -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}) diff --git a/src/aipass/trigger/apps/handlers/events/runaway_handler.py b/src/aipass/trigger/apps/handlers/events/runaway_handler.py new file mode 100644 index 00000000..5df7e371 --- /dev/null +++ b/src/aipass/trigger/apps/handlers/events/runaway_handler.py @@ -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}") diff --git a/src/aipass/trigger/tests/test_runaway_handler.py b/src/aipass/trigger/tests/test_runaway_handler.py new file mode 100644 index 00000000..30f23ef1 --- /dev/null +++ b/src/aipass/trigger/tests/test_runaway_handler.py @@ -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