fix: CI green pass on PR#696 — lint (ruff format on trigger runaway tests), seedgo 100% both red branches (@hooks: new json_handler + persistent_alert/alert_dismiss wired through it, cc_sessions introspection gate, _menu_live nesting extraction, 1071 tests; @prax: rate_tracker DI refactor — module layer injects logs_dir + trigger.fire via configure(), 1028 tests), navmap trim 9.3k→7.9k under injection cap. All owner-fixed via dispatch, devpulse-verified: both audits 100% independently re-run, full runaway chain re-proven live post-refactor (332 lines/min storm → WARNING at 130s → trigger → @aipass triage → alert banner rendered in-session → dismiss clears). Burst-evasion design finding (pre-existing, not regression) filed to @prax by mail.
This commit is contained in:
@@ -0,0 +1,3 @@
|
||||
"""JSON Handler — Hooks Branch."""
|
||||
|
||||
__all__ = []
|
||||
@@ -0,0 +1,165 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: json_handler.py
|
||||
# Description: JSON auto-creating handler for hooks data files
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-07-15
|
||||
# Modified: 2026-07-15
|
||||
# =============================================
|
||||
|
||||
"""JSON auto-creating handler for hooks data files."""
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
import inspect
|
||||
|
||||
from aipass.prax.apps.modules.logger import system_logger as logger
|
||||
|
||||
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")
|
||||
|
||||
_BRANCH_ROOT = Path(__file__).resolve().parents[3]
|
||||
_BRANCH_NAME = _BRANCH_ROOT.name
|
||||
JSON_DIR = _BRANCH_ROOT / f"{_BRANCH_NAME}_json"
|
||||
|
||||
|
||||
def _get_caller_module_name() -> str:
|
||||
"""Auto-detect calling module name from call stack."""
|
||||
stack = inspect.stack()
|
||||
if len(stack) > 2:
|
||||
caller_frame = stack[2]
|
||||
caller_path = Path(caller_frame.filename)
|
||||
module_name = caller_path.stem
|
||||
if module_name and not module_name.startswith("_"):
|
||||
return module_name
|
||||
return "unknown"
|
||||
|
||||
|
||||
def _create_default(json_type: str, module_name: str) -> Any:
|
||||
"""Create default JSON structure from inline code defaults."""
|
||||
today = datetime.now().date().isoformat()
|
||||
if json_type == "config":
|
||||
return {
|
||||
"module_name": module_name,
|
||||
"version": "1.0.0",
|
||||
"config": {"max_log_entries": 100},
|
||||
"created": today,
|
||||
}
|
||||
elif json_type == "data":
|
||||
return {
|
||||
"module_name": module_name,
|
||||
"created": today,
|
||||
"last_updated": today,
|
||||
}
|
||||
elif json_type == "log":
|
||||
return []
|
||||
raise ValueError(f"Unknown json_type: {json_type}")
|
||||
|
||||
|
||||
def validate_json_structure(data: Any, json_type: str) -> bool:
|
||||
"""Validate JSON structure matches expected type."""
|
||||
if json_type == "config":
|
||||
return isinstance(data, dict) and all(k in data for k in ["module_name", "version", "config"])
|
||||
elif json_type == "data":
|
||||
return isinstance(data, dict) and all(k in data for k in ["created", "last_updated"])
|
||||
elif json_type == "log":
|
||||
return isinstance(data, list)
|
||||
return False
|
||||
|
||||
|
||||
def get_json_path(module_name: str, json_type: str) -> Path:
|
||||
"""Get path for module JSON file."""
|
||||
return JSON_DIR / f"{module_name}_{json_type}.json"
|
||||
|
||||
|
||||
def ensure_json_exists(module_name: str, json_type: str) -> bool:
|
||||
"""Ensure JSON file exists, create from template if missing."""
|
||||
JSON_DIR.mkdir(parents=True, exist_ok=True)
|
||||
json_path = get_json_path(module_name, json_type)
|
||||
if json_path.exists():
|
||||
try:
|
||||
data = json.loads(json_path.read_text(encoding="utf-8"))
|
||||
if validate_json_structure(data, json_type):
|
||||
return True
|
||||
except Exception as exc:
|
||||
logger.warning("[HOOKS] json_handler: ensure_json_exists failed for %s_%s: %s", module_name, json_type, exc)
|
||||
template = _create_default(json_type, module_name)
|
||||
json_path.write_text(json.dumps(template, indent=2, ensure_ascii=False) + "\n", encoding="utf-8")
|
||||
return True
|
||||
|
||||
|
||||
def load_json(module_name: str, json_type: str) -> Any | None:
|
||||
"""Load JSON file, auto-create if missing."""
|
||||
if not ensure_json_exists(module_name, json_type):
|
||||
return None
|
||||
json_path = get_json_path(module_name, json_type)
|
||||
return json.loads(json_path.read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def save_json(module_name: str, json_type: str, data: Any) -> bool:
|
||||
"""Save JSON file."""
|
||||
json_path = get_json_path(module_name, json_type)
|
||||
if not validate_json_structure(data, json_type):
|
||||
raise ValueError(f"Invalid structure for {json_type} JSON")
|
||||
if json_type == "data" and isinstance(data, dict):
|
||||
data["last_updated"] = datetime.now().date().isoformat()
|
||||
json_path.write_text(json.dumps(data, indent=2, ensure_ascii=False) + "\n", encoding="utf-8")
|
||||
return True
|
||||
|
||||
|
||||
def ensure_module_jsons(module_name: str) -> bool:
|
||||
"""Ensure all 3 JSON files exist for a module."""
|
||||
ensure_json_exists(module_name, "config")
|
||||
ensure_json_exists(module_name, "data")
|
||||
ensure_json_exists(module_name, "log")
|
||||
return True
|
||||
|
||||
|
||||
def log_operation(
|
||||
operation: str,
|
||||
data: dict[str, Any] | None = None,
|
||||
module_name: str | None = None,
|
||||
) -> bool:
|
||||
"""Add entry to module log with automatic rotation.
|
||||
|
||||
Auto-detects calling module if module_name not provided.
|
||||
"""
|
||||
if module_name is None:
|
||||
module_name = _get_caller_module_name()
|
||||
ensure_module_jsons(module_name)
|
||||
|
||||
config = load_json(module_name, "config")
|
||||
max_entries = 100
|
||||
if config and "config" in config:
|
||||
max_entries = config["config"].get("max_log_entries", 100)
|
||||
|
||||
log = load_json(module_name, "log")
|
||||
if log is None:
|
||||
log = []
|
||||
|
||||
entry: dict[str, Any] = {"timestamp": datetime.now().isoformat(), "operation": operation}
|
||||
if data:
|
||||
entry["data"] = data
|
||||
|
||||
log.append(entry)
|
||||
if len(log) > max_entries:
|
||||
log = log[-max_entries:]
|
||||
|
||||
return save_json(module_name, "log", log)
|
||||
|
||||
|
||||
def read_json_file(path: Path) -> Any:
|
||||
"""Read and parse a JSON file at an arbitrary path."""
|
||||
return json.loads(path.read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def write_json_file(path: Path, data: Any) -> None:
|
||||
"""Write data as JSON to an arbitrary path."""
|
||||
path.write_text(json.dumps(data, indent=2) + "\n", encoding="utf-8")
|
||||
@@ -423,6 +423,71 @@ def _is_session_file_present(pid: int | None) -> bool:
|
||||
return session_file.exists()
|
||||
|
||||
|
||||
def _menu_single_session(
|
||||
session: dict,
|
||||
branch: str,
|
||||
claude_bin: str,
|
||||
defaults: list[str],
|
||||
extra_args: list[str] | None,
|
||||
) -> dict:
|
||||
"""Handle menu for a single live session."""
|
||||
label = _session_label(session, branch)
|
||||
is_bg = session.get("kind") in ("bg", "background")
|
||||
sys.stderr.write(f"\n{branch} — live chat: {label}\n")
|
||||
if is_bg:
|
||||
sys.stderr.write(" [Enter] resume this chat (stops bg, reopens as normal chat)\n")
|
||||
sys.stderr.write(" [n] start new chat (stops bg first)\n")
|
||||
sys.stderr.write(" [c] close it and exit (stops bg)\n\n")
|
||||
else:
|
||||
sys.stderr.write(" [Enter] resume this chat\n")
|
||||
sys.stderr.write(" [n] start new chat (closes the one above first)\n")
|
||||
sys.stderr.write(" [c] close it and exit\n\n")
|
||||
|
||||
choice = _read_choice()
|
||||
|
||||
if choice in ("", "r"):
|
||||
return _resume_session(session, branch, claude_bin, defaults, extra_args)
|
||||
if choice == "n":
|
||||
return _menu_single_new(session, is_bg, branch, claude_bin, defaults, extra_args)
|
||||
if choice == "c":
|
||||
return _menu_single_close(session, is_bg, branch, claude_bin)
|
||||
if choice in ("exit", "q", "quit"):
|
||||
return {"exit_code": 0, "action": "quit"}
|
||||
sys.stderr.write(" Unknown choice. Exiting.\n")
|
||||
return {"exit_code": 1, "error": "unknown choice"}
|
||||
|
||||
|
||||
def _menu_single_new(
|
||||
session: dict,
|
||||
is_bg: bool,
|
||||
branch: str,
|
||||
claude_bin: str,
|
||||
defaults: list[str],
|
||||
extra_args: list[str] | None,
|
||||
) -> dict:
|
||||
"""Handle 'n' choice for single session — stop current, start fresh."""
|
||||
if is_bg:
|
||||
stop = _daemon_stop(claude_bin, branch, session.get("pid"))
|
||||
if not stop["ok"]:
|
||||
return {"exit_code": 1, "error": stop["error"]}
|
||||
else:
|
||||
_stop_session(session, claude_bin)
|
||||
return _start_fresh(branch, claude_bin, defaults, extra_args)
|
||||
|
||||
|
||||
def _menu_single_close(session: dict, is_bg: bool, branch: str, claude_bin: str) -> dict:
|
||||
"""Handle 'c' choice for single session — close and exit."""
|
||||
if is_bg:
|
||||
stop = _daemon_stop(claude_bin, branch, session.get("pid"))
|
||||
if not stop["ok"]:
|
||||
return {"exit_code": 1, "error": stop["error"]}
|
||||
sys.stderr.write(f" Stopped bg session PID {session.get('pid')}.\n")
|
||||
else:
|
||||
result = _stop_session(session, claude_bin)
|
||||
sys.stderr.write(f" {result}\n")
|
||||
return {"exit_code": 0, "action": "closed"}
|
||||
|
||||
|
||||
def _menu_live(
|
||||
live: list[dict],
|
||||
branch: str,
|
||||
@@ -432,46 +497,7 @@ def _menu_live(
|
||||
) -> dict:
|
||||
"""Display menu when live session(s) exist."""
|
||||
if len(live) == 1:
|
||||
session = live[0]
|
||||
label = _session_label(session, branch)
|
||||
is_bg = session.get("kind") in ("bg", "background")
|
||||
sys.stderr.write(f"\n{branch} — live chat: {label}\n")
|
||||
if is_bg:
|
||||
sys.stderr.write(" [Enter] resume this chat (stops bg, reopens as normal chat)\n")
|
||||
sys.stderr.write(" [n] start new chat (stops bg first)\n")
|
||||
sys.stderr.write(" [c] close it and exit (stops bg)\n\n")
|
||||
else:
|
||||
sys.stderr.write(" [Enter] resume this chat\n")
|
||||
sys.stderr.write(" [n] start new chat (closes the one above first)\n")
|
||||
sys.stderr.write(" [c] close it and exit\n\n")
|
||||
|
||||
choice = _read_choice()
|
||||
|
||||
if choice in ("", "r"):
|
||||
return _resume_session(session, branch, claude_bin, defaults, extra_args)
|
||||
elif choice == "n":
|
||||
if is_bg:
|
||||
stop = _daemon_stop(claude_bin, branch, session.get("pid"))
|
||||
if not stop["ok"]:
|
||||
return {"exit_code": 1, "error": stop["error"]}
|
||||
else:
|
||||
_stop_session(session, claude_bin)
|
||||
return _start_fresh(branch, claude_bin, defaults, extra_args)
|
||||
elif choice == "c":
|
||||
if is_bg:
|
||||
stop = _daemon_stop(claude_bin, branch, session.get("pid"))
|
||||
if not stop["ok"]:
|
||||
return {"exit_code": 1, "error": stop["error"]}
|
||||
sys.stderr.write(f" Stopped bg session PID {session.get('pid')}.\n")
|
||||
else:
|
||||
result = _stop_session(session, claude_bin)
|
||||
sys.stderr.write(f" {result}\n")
|
||||
return {"exit_code": 0, "action": "closed"}
|
||||
elif choice in ("exit", "q", "quit"):
|
||||
return {"exit_code": 0, "action": "quit"}
|
||||
else:
|
||||
sys.stderr.write(" Unknown choice. Exiting.\n")
|
||||
return {"exit_code": 1, "error": "unknown choice"}
|
||||
return _menu_single_session(live[0], branch, claude_bin, defaults, extra_args)
|
||||
|
||||
sys.stderr.write(f"\n{branch} — {len(live)} live sessions:\n")
|
||||
for i, session in enumerate(live, 1):
|
||||
|
||||
@@ -14,6 +14,7 @@ import json
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from aipass.hooks.apps.handlers.json import json_handler
|
||||
from aipass.prax.apps.modules.logger import system_logger as logger
|
||||
|
||||
_announced: set[str] = set()
|
||||
@@ -124,6 +125,7 @@ def handle(hook_data: dict) -> dict:
|
||||
plural = "s" if count != 1 else ""
|
||||
sound = f"alert: {count} active alert{plural}"
|
||||
|
||||
json_handler.log_operation("inject_alerts", {"count": len(alerts)})
|
||||
logger.info("[HOOKS] persistent_alert: %d active alerts injected", len(alerts))
|
||||
result = {"stdout": banner, "exit_code": 0}
|
||||
if sound:
|
||||
|
||||
@@ -14,6 +14,7 @@ import json
|
||||
from pathlib import Path
|
||||
|
||||
from aipass.cli.apps.modules import err_console
|
||||
from aipass.hooks.apps.handlers.json import json_handler
|
||||
from aipass.prax.apps.modules.logger import system_logger as logger
|
||||
|
||||
CONSOLE = err_console
|
||||
@@ -52,7 +53,7 @@ def _dismiss_alert(alert_id: str) -> bool:
|
||||
return False
|
||||
|
||||
try:
|
||||
data = json.loads(alerts_path.read_text(encoding="utf-8"))
|
||||
data = json_handler.read_json_file(alerts_path)
|
||||
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]")
|
||||
@@ -67,15 +68,13 @@ def _dismiss_alert(alert_id: str) -> bool:
|
||||
return False
|
||||
|
||||
try:
|
||||
alerts_path.write_text(
|
||||
json.dumps({"alerts": remaining}, indent=2) + "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
json_handler.write_json_file(alerts_path, {"alerts": remaining})
|
||||
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
|
||||
|
||||
json_handler.log_operation("dismiss_alert", {"alert_id": alert_id})
|
||||
logger.info("[HOOKS] dismiss: removed alert %s", alert_id)
|
||||
CONSOLE.print(f"[green]Dismissed alert {alert_id}[/green]")
|
||||
return True
|
||||
|
||||
@@ -264,9 +264,7 @@ def handle_command(command: str, args: list) -> bool:
|
||||
|
||||
if command in ("sessions", "cc_sessions"):
|
||||
if not args:
|
||||
CONSOLE.print("[bold cyan]sessions[/bold cyan]")
|
||||
CONSOLE.print(f" Sessions dir: {CC_SESSIONS_DIR}")
|
||||
_print_sessions_list()
|
||||
print_introspection()
|
||||
return True
|
||||
|
||||
if args[0] == "reclaim":
|
||||
|
||||
@@ -27,20 +27,10 @@ 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
|
||||
|
||||
@@ -103,13 +93,32 @@ _tracked: Dict[str, FileRateState] = {}
|
||||
|
||||
_suppressed_files: set = set()
|
||||
|
||||
_logs_dir: Optional[Path] = None
|
||||
|
||||
_EVENT_CALLBACK = None
|
||||
|
||||
_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 configure(
|
||||
logs_dir: Optional[Path] = None,
|
||||
event_callback=None,
|
||||
suppressed_files: Optional[set] = None,
|
||||
) -> None:
|
||||
"""Inject dependencies from the module layer.
|
||||
|
||||
Args:
|
||||
logs_dir: Path to system_logs/ directory to scan.
|
||||
event_callback: Callable(event_name, **kwargs) for firing events.
|
||||
suppressed_files: Set of log file names to skip during detection.
|
||||
"""
|
||||
global _logs_dir, _EVENT_CALLBACK, _suppressed_files
|
||||
if logs_dir is not None:
|
||||
_logs_dir = logs_dir
|
||||
if event_callback is not None:
|
||||
_EVENT_CALLBACK = event_callback
|
||||
if suppressed_files is not None:
|
||||
_suppressed_files = suppressed_files
|
||||
|
||||
|
||||
def _load_state() -> None:
|
||||
@@ -160,15 +169,14 @@ def scan_rates() -> list:
|
||||
"""
|
||||
_load_state()
|
||||
|
||||
logs_dir = get_system_logs_dir()
|
||||
if not logs_dir.exists():
|
||||
if _logs_dir is None or not _logs_dir.exists():
|
||||
return []
|
||||
|
||||
now = time.time()
|
||||
results = []
|
||||
|
||||
current_files = set()
|
||||
for log_file in logs_dir.glob("*.log"):
|
||||
for log_file in _logs_dir.glob("*.log"):
|
||||
file_key = str(log_file)
|
||||
current_files.add(file_key)
|
||||
|
||||
@@ -302,8 +310,8 @@ def _fire_event(
|
||||
},
|
||||
)
|
||||
|
||||
if _HAS_TRIGGER and trigger is not None:
|
||||
trigger.fire(
|
||||
if _EVENT_CALLBACK is not None:
|
||||
_EVENT_CALLBACK(
|
||||
"runaway_log_detected",
|
||||
file_path=file_path,
|
||||
rate_lines_per_min=rate_lines_per_min,
|
||||
@@ -358,11 +366,3 @@ def get_snapshot() -> list:
|
||||
}
|
||||
)
|
||||
return results
|
||||
|
||||
|
||||
def reset() -> None:
|
||||
"""Clear all tracking state. Used in tests."""
|
||||
global _state_loaded
|
||||
_tracked.clear()
|
||||
_suppressed_files.clear()
|
||||
_state_loaded = False
|
||||
|
||||
@@ -121,6 +121,17 @@ def _display_rates(results: list, is_scan: bool) -> None:
|
||||
console.print()
|
||||
|
||||
|
||||
def _get_event_callback():
|
||||
"""Return trigger.fire if available, else None."""
|
||||
try:
|
||||
from aipass.trigger.apps.modules.core import trigger
|
||||
|
||||
return trigger.fire
|
||||
except ImportError as exc:
|
||||
logger.info("[log-health] trigger not available: %s", exc)
|
||||
return None
|
||||
|
||||
|
||||
def handle_command(command: str, args: List[str]) -> bool:
|
||||
"""Handle log-health command.
|
||||
|
||||
@@ -142,7 +153,10 @@ def handle_command(command: str, args: List[str]) -> bool:
|
||||
print_help()
|
||||
return True
|
||||
|
||||
from aipass.prax.apps.handlers.monitoring.rate_tracker import scan_rates, get_snapshot
|
||||
from aipass.prax.apps.handlers.monitoring.rate_tracker import scan_rates, get_snapshot, configure
|
||||
from aipass.prax.apps.handlers.config.load import get_system_logs_dir
|
||||
|
||||
configure(logs_dir=get_system_logs_dir(), event_callback=_get_event_callback())
|
||||
|
||||
subcmd = args[0]
|
||||
logger.info("[log-health] %s", subcmd)
|
||||
|
||||
@@ -490,7 +490,18 @@ def _log_watcher_worker():
|
||||
|
||||
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
|
||||
from aipass.prax.apps.handlers.monitoring.rate_tracker import scan_rates, configure, SCAN_INTERVAL
|
||||
from aipass.prax.apps.handlers.config.load import get_system_logs_dir
|
||||
|
||||
try:
|
||||
from aipass.trigger.apps.modules.core import trigger
|
||||
|
||||
event_cb = trigger.fire
|
||||
except ImportError as exc:
|
||||
logger.info("[monitor] trigger not available for rate tracker: %s", exc)
|
||||
event_cb = None
|
||||
|
||||
configure(logs_dir=get_system_logs_dir(), event_callback=event_cb)
|
||||
|
||||
while not _stop_event.is_set():
|
||||
try:
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_rate_tracker.py
|
||||
# Description: Tests for the rate tracker runaway-log detector
|
||||
# Version: 1.0.0
|
||||
# Version: 1.1.0
|
||||
# Created: 2026-07-14
|
||||
# Modified: 2026-07-14
|
||||
# Modified: 2026-07-15
|
||||
# =============================================
|
||||
|
||||
"""Tests for apps/handlers/monitoring/rate_tracker.py
|
||||
@@ -13,9 +13,10 @@ Covers:
|
||||
- Sustained threshold detection (WARNING and CRITICAL)
|
||||
- Subsidence reset when rate drops
|
||||
- Per-file suppression
|
||||
- Event firing via trigger
|
||||
- Event firing via callback
|
||||
- File disappearance handling
|
||||
- get_snapshot() and reset()
|
||||
- get_snapshot() and configure()
|
||||
- Disk persistence
|
||||
"""
|
||||
|
||||
import sys
|
||||
@@ -29,7 +30,7 @@ _HANDLER_MOCKS = {
|
||||
}
|
||||
|
||||
|
||||
def _import_tracker(monkeypatch):
|
||||
def _import_tracker(monkeypatch, logs_dir=None):
|
||||
"""Import (or reload) rate_tracker with handler mocks."""
|
||||
monkeypatch.delenv("PYTEST_CURRENT_TEST", raising=False)
|
||||
fresh = {k: MagicMock() for k in _HANDLER_MOCKS}
|
||||
@@ -41,12 +42,14 @@ def _import_tracker(monkeypatch):
|
||||
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
|
||||
event_mock = MagicMock()
|
||||
mod._tracked.clear()
|
||||
mod._suppressed_files.clear()
|
||||
setattr(mod, "_state_loaded", False)
|
||||
setattr(mod, "_logs_dir", None)
|
||||
setattr(mod, "_EVENT_CALLBACK", None)
|
||||
mod.configure(logs_dir=logs_dir, event_callback=event_mock)
|
||||
return mod, event_mock
|
||||
|
||||
|
||||
class TestRateCalculation:
|
||||
@@ -54,32 +57,28 @@ class TestRateCalculation:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "test_module.log").write_text("line1\n" * 10)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
results = mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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),
|
||||
):
|
||||
with patch.object(mod.time, "time", return_value=time.time() + 10.0):
|
||||
results = mod.scan_rates()
|
||||
|
||||
assert len(results) == 1
|
||||
@@ -87,18 +86,14 @@ class TestRateCalculation:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "test_module.log").write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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),
|
||||
):
|
||||
with patch.object(mod.time, "time", return_value=time.time() + 10.0):
|
||||
results = mod.scan_rates()
|
||||
|
||||
assert len(results) == 1
|
||||
@@ -106,20 +101,17 @@ class TestRateCalculation:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 10000)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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),
|
||||
):
|
||||
with patch.object(mod.time, "time", return_value=time.time() + 10.0):
|
||||
results = mod.scan_rates()
|
||||
|
||||
assert results == []
|
||||
@@ -135,85 +127,82 @@ class TestSustainedThresholds:
|
||||
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,
|
||||
),
|
||||
with 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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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
|
||||
event_mock.assert_called_once()
|
||||
call_args = event_mock.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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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()
|
||||
event_mock.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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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 event_mock.call_count == 1
|
||||
call_args = event_mock.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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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
|
||||
assert event_mock.call_count == 1
|
||||
|
||||
|
||||
class TestSubsidence:
|
||||
@@ -221,56 +210,48 @@ class TestSubsidence:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
|
||||
base_time = time.time()
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
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,
|
||||
),
|
||||
with patch.object(
|
||||
mod.time,
|
||||
"time",
|
||||
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
|
||||
):
|
||||
mod.scan_rates()
|
||||
|
||||
assert trigger_mock.fire.call_count == 1
|
||||
assert event_mock.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,
|
||||
),
|
||||
with 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,
|
||||
),
|
||||
with patch.object(
|
||||
mod.time,
|
||||
"time",
|
||||
return_value=base_time + offset * mod.SCAN_INTERVAL,
|
||||
):
|
||||
mod.scan_rates()
|
||||
|
||||
assert trigger_mock.fire.call_count == 2
|
||||
assert event_mock.call_count == 2
|
||||
|
||||
|
||||
class TestSuppression:
|
||||
@@ -278,30 +259,28 @@ class TestSuppression:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "noisy_module.log").write_text("x" * 10000)
|
||||
|
||||
mod.configure_suppression({"noisy_module.log"})
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
mod.configure(suppressed_files={"noisy_module.log"})
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod.scan_rates()
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "normal_module.log").write_text("x" * 100)
|
||||
|
||||
mod.configure_suppression({"other_module.log"})
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
mod.configure(suppressed_files={"other_module.log"})
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod.scan_rates()
|
||||
|
||||
assert any("normal_module.log" in k for k in mod._tracked)
|
||||
|
||||
@@ -311,30 +290,28 @@ class TestEventPayload:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
mod, event_mock = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
|
||||
base_time = time.time()
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
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,
|
||||
),
|
||||
with 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]
|
||||
call_kwargs = event_mock.call_args[1]
|
||||
assert "file_path" in call_kwargs
|
||||
assert "rate_lines_per_min" in call_kwargs
|
||||
assert "sustained_duration_sec" in call_kwargs
|
||||
@@ -348,20 +325,19 @@ class TestFileDisappearance:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "ephemeral.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
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()
|
||||
mod.scan_rates()
|
||||
|
||||
assert not any("ephemeral.log" in k for k in mod._tracked)
|
||||
|
||||
@@ -371,13 +347,12 @@ class TestSnapshot:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "test_module.log").write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
mod.scan_rates()
|
||||
|
||||
snapshot = mod.get_snapshot()
|
||||
assert len(snapshot) == 1
|
||||
@@ -389,53 +364,49 @@ class TestSnapshot:
|
||||
assert mod.get_snapshot() == []
|
||||
|
||||
|
||||
class TestReset:
|
||||
"""reset() clears all state."""
|
||||
class TestConfigure:
|
||||
"""configure() sets module-level dependencies."""
|
||||
|
||||
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)
|
||||
def test_configure_clears_tracked_on_reimport(self, tmp_path, monkeypatch):
|
||||
"""Fresh import via _import_tracker starts with empty tracked dict."""
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "test_module.log").write_text("x" * 100)
|
||||
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
mod.scan_rates()
|
||||
|
||||
assert len(mod._tracked) > 0
|
||||
mod.reset()
|
||||
|
||||
mod._tracked.clear()
|
||||
assert len(mod._tracked) == 0
|
||||
|
||||
|
||||
class TestNoTrigger:
|
||||
"""When trigger is unavailable, detection still works — just no event fired."""
|
||||
"""When no event callback is set, 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)
|
||||
def test_detection_without_callback(self, tmp_path, monkeypatch):
|
||||
"""Rate tracking and threshold detection work without event callback."""
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 100)
|
||||
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
mod.configure(event_callback=None)
|
||||
|
||||
base_time = time.time()
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
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,
|
||||
),
|
||||
with patch.object(
|
||||
mod.time,
|
||||
"time",
|
||||
return_value=base_time + (i + 1) * mod.SCAN_INTERVAL,
|
||||
):
|
||||
results = mod.scan_rates()
|
||||
|
||||
@@ -447,14 +418,14 @@ class TestPersistence:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
(logs_dir / "test_module.log").write_text("x" * 100)
|
||||
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
|
||||
json_handler_mock = mod.json_handler
|
||||
with patch.object(mod, "get_system_logs_dir", return_value=log_file.parent):
|
||||
mod.scan_rates()
|
||||
mod.scan_rates()
|
||||
|
||||
json_handler_mock.save_json.assert_called()
|
||||
call_args = json_handler_mock.save_json.call_args
|
||||
@@ -466,12 +437,17 @@ class TestPersistence:
|
||||
|
||||
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)
|
||||
logs_dir = tmp_path / "system"
|
||||
logs_dir.mkdir(parents=True)
|
||||
log_file = logs_dir / "test_module.log"
|
||||
log_file.write_text("x" * 1300)
|
||||
|
||||
mod, _ = _import_tracker(monkeypatch, logs_dir=logs_dir)
|
||||
|
||||
persisted = {
|
||||
"module_name": "rate_tracker",
|
||||
"files": {
|
||||
str(tmp_path / "system" / "test_module.log"): {
|
||||
str(log_file): {
|
||||
"last_offset": 100,
|
||||
"last_check": time.time() - 15.0,
|
||||
"warning_sustained": 0,
|
||||
@@ -483,12 +459,7 @@ class TestPersistence:
|
||||
}
|
||||
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()
|
||||
results = mod.scan_rates()
|
||||
|
||||
assert len(results) == 1
|
||||
assert results[0]["rate_lines_per_min"] > 0
|
||||
@@ -501,9 +472,11 @@ class TestPersistence:
|
||||
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."""
|
||||
def test_reimport_clears_state_loaded_flag(self, monkeypatch):
|
||||
"""Fresh _import_tracker resets _state_loaded so next scan reloads from disk."""
|
||||
mod, _ = _import_tracker(monkeypatch)
|
||||
setattr(mod, "_state_loaded", True)
|
||||
mod.reset()
|
||||
assert mod._state_loaded is True
|
||||
|
||||
mod, _ = _import_tracker(monkeypatch)
|
||||
assert mod._state_loaded is False
|
||||
|
||||
@@ -221,11 +221,7 @@ class TestNoEmailCallback:
|
||||
)
|
||||
|
||||
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"
|
||||
]
|
||||
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"]
|
||||
|
||||
@@ -410,11 +406,7 @@ class TestSuppressionLog:
|
||||
)
|
||||
|
||||
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"
|
||||
]
|
||||
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
|
||||
|
||||
@@ -435,11 +427,7 @@ class TestSuppressionLog:
|
||||
)
|
||||
|
||||
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"
|
||||
]
|
||||
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"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user