Merge pull request #489 from AIOSAI/work/trigger
feat(trigger): S117 stress test @trigger
This commit is contained in:
@@ -423,6 +423,57 @@
|
||||
"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_error_detected.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_error_detected.py",
|
||||
"standard": "encapsulation",
|
||||
"lines": [72],
|
||||
"pattern": "handler imported directly",
|
||||
"reason": "Test helper _import_module() must import the handler module directly to test it. All trigger test files follow this pattern."
|
||||
},
|
||||
{
|
||||
"file": "tests/test_event_handlers.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_event_handlers.py",
|
||||
"standard": "encapsulation",
|
||||
"lines": [51],
|
||||
"pattern": "handler imported directly",
|
||||
"reason": "Test helper _import_cli() must import the handler module directly to test it. All trigger test files follow this pattern."
|
||||
},
|
||||
{
|
||||
"file": "tests/test_core.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_core.py",
|
||||
"standard": "documentation",
|
||||
"pattern": "public functions missing docstrings",
|
||||
"reason": "Test helper functions (capture_handler, etc.) are internal test utilities, not public API. Docstrings on test fixtures add noise without value."
|
||||
},
|
||||
{
|
||||
"file": "tests/test_error_registry.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_error_registry.py",
|
||||
"standard": "encapsulation",
|
||||
"lines": [63],
|
||||
"pattern": "handler imported directly",
|
||||
"reason": "Test helper _import_registry() must import the handler module directly to test it. All trigger test files follow this pattern."
|
||||
}
|
||||
],
|
||||
"notes": {
|
||||
|
||||
@@ -776,6 +776,40 @@ def clear_resolved(days: int = 7) -> int:
|
||||
return 0
|
||||
|
||||
|
||||
def purge_stale(days: int = 30) -> int:
|
||||
"""Remove entries whose last_seen is older than N days regardless of status.
|
||||
|
||||
Args:
|
||||
days: Age threshold in days (default: 30)
|
||||
|
||||
Returns:
|
||||
Count of removed entries
|
||||
"""
|
||||
try:
|
||||
registry = _load_registry()
|
||||
cutoff = (datetime.now() - timedelta(days=days)).isoformat()
|
||||
removed = 0
|
||||
|
||||
fingerprints_to_remove = []
|
||||
for fp, entry in registry["errors"].items():
|
||||
last_seen = entry.get("last_seen", "")
|
||||
if last_seen and last_seen < cutoff:
|
||||
fingerprints_to_remove.append(fp)
|
||||
|
||||
for fp in fingerprints_to_remove:
|
||||
del registry["errors"][fp]
|
||||
removed += 1
|
||||
|
||||
if removed > 0:
|
||||
_save_registry(registry)
|
||||
|
||||
return removed
|
||||
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to purge stale entries: %s", exc)
|
||||
return 0
|
||||
|
||||
|
||||
def get_stats() -> dict:
|
||||
"""Get summary statistics from the error registry.
|
||||
|
||||
|
||||
@@ -512,7 +512,7 @@ def handle_error_detected(
|
||||
)
|
||||
|
||||
# Send via callback (set by module layer, trigger isn't a branch so PWD detection fails)
|
||||
_send_email(
|
||||
sent = _send_email(
|
||||
to_branch=recipient,
|
||||
subject=email_subject,
|
||||
message=notification_message,
|
||||
@@ -521,6 +521,10 @@ def handle_error_detected(
|
||||
from_branch="@trigger",
|
||||
)
|
||||
|
||||
if not sent:
|
||||
_log_warning(f"Email delivery failed for {recipient} (fingerprint={fingerprint})")
|
||||
return
|
||||
|
||||
# Wake the target branch so the email is processed immediately
|
||||
try:
|
||||
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
|
||||
|
||||
@@ -1,76 +0,0 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: error_logged.py
|
||||
# Description: Legacy error logged event handler — monitor-only (no dispatch)
|
||||
# Version: 3.0.0
|
||||
# Created: 2026-01-31
|
||||
# Modified: 2026-04-10
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Error Logged Event Handler (Monitor-Only)
|
||||
|
||||
Legacy handler for error_logged events. All dispatch now goes through
|
||||
error_detected.py (Medic v2) which provides circuit breaker, per-fingerprint
|
||||
backoff, and registry-based deduplication.
|
||||
|
||||
This handler logs event metadata for monitoring. No email, no wake_branch.
|
||||
"""
|
||||
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
_HANDLER_LOG = TRIGGER_ROOT / "logs" / "error_logged_handler.log"
|
||||
|
||||
|
||||
def _log_warning(message: str) -> None:
|
||||
"""Log warning to file (event handlers cannot import Prax logger - causes recursion)."""
|
||||
try:
|
||||
_HANDLER_LOG.parent.mkdir(parents=True, exist_ok=True)
|
||||
ts = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
|
||||
with open(_HANDLER_LOG, "a", encoding="utf-8") as f:
|
||||
f.write(f"{ts} | WARNING | {message}\n")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def handle_error_logged(
|
||||
branch: str | None = None,
|
||||
message: str | None = None,
|
||||
error_hash: str | None = None,
|
||||
source_module: str | None = None,
|
||||
module_name: str | None = None,
|
||||
**_kwargs: Any,
|
||||
) -> None:
|
||||
"""Handle error_logged event — monitor-only, no dispatch.
|
||||
|
||||
All dispatch now goes through error_detected.py (Medic v2).
|
||||
This handler logs event metadata for monitoring only.
|
||||
|
||||
Args:
|
||||
branch: Branch where error occurred
|
||||
message: Error message text
|
||||
error_hash: Unique error identifier
|
||||
source_module: Module that logged the error
|
||||
module_name: Deprecated alias for source_module
|
||||
**_kwargs: Additional event data (ignored)
|
||||
"""
|
||||
try:
|
||||
if not branch or not message or not error_hash:
|
||||
return
|
||||
|
||||
effective_module = source_module or module_name or "unknown"
|
||||
|
||||
json_handler.log_operation(
|
||||
"error_logged_event",
|
||||
{
|
||||
"branch": branch,
|
||||
"module": effective_module,
|
||||
"error_hash": error_hash,
|
||||
},
|
||||
)
|
||||
|
||||
except Exception as exc:
|
||||
_log_warning(f"handle_error_logged failed: {exc}")
|
||||
return
|
||||
@@ -1,29 +0,0 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: memory.py
|
||||
# Description: Memory event handler placeholder for future rollover triggers
|
||||
# Version: 0.1.0
|
||||
# Created: 2025-12-04
|
||||
# Modified: 2025-12-04
|
||||
# =============================================
|
||||
|
||||
"""Memory Event Handler - Handle memory-related events
|
||||
|
||||
Placeholder for future memory event handling.
|
||||
"""
|
||||
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
|
||||
def handle_memory_saved(**kwargs):
|
||||
"""Handle memory save events - placeholder for future
|
||||
|
||||
Will check line count and trigger rollover if needed.
|
||||
|
||||
Args:
|
||||
**kwargs: Event data (branch, lines, file_path, etc.)
|
||||
"""
|
||||
# Future: Check line count and trigger rollover
|
||||
# if lines > 600:
|
||||
# trigger_rollover(branch)
|
||||
json_handler.log_operation("memory_event", {"success": True})
|
||||
pass
|
||||
@@ -1,169 +0,0 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: memory_threshold_exceeded.py
|
||||
# Description: Memory threshold exceeded event handler for compression notifications
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-01-31
|
||||
# Modified: 2026-01-31
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Memory Threshold Exceeded Event Handler
|
||||
|
||||
Handles memory_threshold_exceeded events fired when a branch's memory file
|
||||
exceeds the configured threshold (600 lines by default).
|
||||
|
||||
Sends compression notification to the affected branch via AI_Mail.
|
||||
|
||||
Event data expected:
|
||||
- branch: Branch name where threshold exceeded
|
||||
- branch_path: Path to branch root
|
||||
- file_name: Memory file name (e.g., local.json, observations.json)
|
||||
- file_path: Full path to the memory file
|
||||
- line_count: Current line count
|
||||
- threshold: Threshold that was exceeded
|
||||
- timestamp: When detected
|
||||
"""
|
||||
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
from aipass.trigger.apps.config import TRIGGER_ROOT
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
|
||||
# Path resolution not needed - this handler uses only event data passed in kwargs
|
||||
|
||||
_HANDLER_LOG = TRIGGER_ROOT / "logs" / "memory_threshold_handler.log"
|
||||
|
||||
|
||||
def _log_warning(message: str) -> None:
|
||||
"""Log warning to file (event handlers cannot import Prax logger - causes recursion)."""
|
||||
try:
|
||||
_HANDLER_LOG.parent.mkdir(parents=True, exist_ok=True)
|
||||
ts = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
|
||||
with open(_HANDLER_LOG, "a", encoding="utf-8") as f:
|
||||
f.write(f"{ts} | WARNING | {message}\n")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _build_compression_message(branch: str, file_name: str, line_count: int, threshold: int) -> str:
|
||||
"""
|
||||
Build compression notification message.
|
||||
|
||||
Args:
|
||||
branch: Branch name
|
||||
file_name: Memory file that exceeded threshold
|
||||
line_count: Current line count
|
||||
threshold: Threshold that was exceeded
|
||||
|
||||
Returns:
|
||||
Formatted message string with compression instructions
|
||||
"""
|
||||
return f"""Memory file threshold exceeded - compression needed.
|
||||
|
||||
Branch: {branch}
|
||||
File: {file_name}
|
||||
Current lines: {line_count}
|
||||
Threshold: {threshold}
|
||||
|
||||
---
|
||||
COMPRESSION INSTRUCTIONS:
|
||||
|
||||
Target: Reduce to ~400 lines while preserving critical information.
|
||||
|
||||
Priority order (what to keep):
|
||||
1. Top 25% (most recent): Keep mostly intact
|
||||
2. Next 25%: Reduce slightly (combine related entries)
|
||||
3. Next 25%: Reduce more (summary format)
|
||||
4. Last 25% (oldest): Delete if needed for space
|
||||
|
||||
Always preserve:
|
||||
- Session headers and dates
|
||||
- Key achievements and milestones
|
||||
- Critical errors and resolutions
|
||||
- Important patterns and learnings
|
||||
|
||||
Safe to remove:
|
||||
- Routine status updates
|
||||
- Redundant information
|
||||
- Low-value details
|
||||
- Completed temporary tasks
|
||||
|
||||
Maintain chronological order (newest first).
|
||||
|
||||
---
|
||||
After compression, verify the file still loads correctly.
|
||||
"""
|
||||
|
||||
|
||||
def handle_memory_threshold_exceeded(
|
||||
branch: str | None = None,
|
||||
branch_path: str | None = None,
|
||||
file_name: str | None = None,
|
||||
file_path: str | None = None,
|
||||
line_count: int | None = None,
|
||||
threshold: int | None = None,
|
||||
timestamp: str | None = None,
|
||||
**_kwargs: Any,
|
||||
) -> None:
|
||||
"""
|
||||
Handle memory_threshold_exceeded event - send compression notification.
|
||||
|
||||
Sends AI_Mail to affected branch with compression instructions when
|
||||
their memory file exceeds the configured threshold.
|
||||
|
||||
Args:
|
||||
branch: Branch name where threshold exceeded - REQUIRED
|
||||
branch_path: Path to branch root
|
||||
file_name: Memory file name - REQUIRED
|
||||
file_path: Full path to the memory file
|
||||
line_count: Current line count - REQUIRED
|
||||
threshold: Threshold that was exceeded (defaults to 600)
|
||||
timestamp: When detected (defaults to now)
|
||||
**_kwargs: Additional event data (ignored)
|
||||
"""
|
||||
try:
|
||||
# Validate required fields
|
||||
if not branch or not file_name or line_count is None:
|
||||
return
|
||||
|
||||
# Import AI_Mail delivery (modules-level API, not handler-level)
|
||||
try:
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
except ImportError:
|
||||
return
|
||||
|
||||
# Set defaults
|
||||
if not timestamp:
|
||||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
if not threshold:
|
||||
threshold = 600
|
||||
|
||||
# Build target and message
|
||||
target_branch = f"@{branch.lower()}"
|
||||
subject = f"[MEMORY] {file_name} exceeded {threshold} lines - compress needed"
|
||||
|
||||
notification_message = _build_compression_message(
|
||||
branch=branch, file_name=file_name, line_count=line_count, threshold=threshold
|
||||
)
|
||||
|
||||
# Build and deliver email
|
||||
email_data = {
|
||||
"from": "@trigger",
|
||||
"from_name": "Trigger",
|
||||
"to": target_branch,
|
||||
"subject": subject,
|
||||
"message": notification_message,
|
||||
"timestamp": timestamp,
|
||||
"auto_execute": False,
|
||||
"priority": "normal",
|
||||
}
|
||||
|
||||
deliver_email_to_branch(target_branch, email_data)
|
||||
|
||||
json_handler.log_operation("memory_threshold_event", {"success": True})
|
||||
|
||||
except Exception as exc:
|
||||
_log_warning(f"handle memory threshold exceeded failed: {exc}")
|
||||
@@ -31,11 +31,9 @@ def setup_handlers():
|
||||
"""Register all event handlers on startup"""
|
||||
from aipass.trigger.apps.modules.core import trigger
|
||||
from .startup import handle_startup
|
||||
from .memory import handle_memory_saved
|
||||
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 .error_logged import handle_error_logged
|
||||
|
||||
# Wire up email send callback for error_detected handler (avoids handler importing from modules)
|
||||
try:
|
||||
@@ -64,21 +62,17 @@ def setup_handlers():
|
||||
_log_warning("ai_mail not available — error notifications won't send")
|
||||
from .warning_logged import handle_warning_logged
|
||||
from .bulletin_created import handle_bulletin_created
|
||||
from .memory_threshold_exceeded import handle_memory_threshold_exceeded
|
||||
from .memory_template_updated import handle_memory_template_updated
|
||||
from .pr_status_sync import handle_pr_created, handle_pr_merged
|
||||
|
||||
trigger.on("startup", handle_startup)
|
||||
trigger.on("memory_saved", handle_memory_saved)
|
||||
trigger.on("cli_header_displayed", handle_cli_header_displayed)
|
||||
trigger.on("plan_file_created", handle_plan_file_created)
|
||||
trigger.on("plan_file_deleted", handle_plan_file_deleted)
|
||||
trigger.on("plan_file_moved", handle_plan_file_moved)
|
||||
trigger.on("error_detected", handle_error_detected)
|
||||
trigger.on("error_logged", handle_error_logged)
|
||||
trigger.on("warning_logged", handle_warning_logged)
|
||||
trigger.on("bulletin_created", handle_bulletin_created)
|
||||
trigger.on("memory_threshold_exceeded", handle_memory_threshold_exceeded)
|
||||
trigger.on("memory_template_updated", handle_memory_template_updated)
|
||||
trigger.on("pr_created", handle_pr_created)
|
||||
trigger.on("pr_merged", handle_pr_merged)
|
||||
|
||||
@@ -281,8 +281,16 @@ class LogFileWatcher(WatchdogFileSystemEventHandler if WATCHDOG_AVAILABLE else o
|
||||
count=error_count,
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning("Registry unavailable, falling back to error_logged: %s", exc)
|
||||
trigger.fire("error_logged", **event_data)
|
||||
logger.warning("Registry unavailable, falling back to error_detected: %s", exc)
|
||||
trigger.fire(
|
||||
"error_detected",
|
||||
branch=branch,
|
||||
module=module_name,
|
||||
message=message,
|
||||
log_path=log_file,
|
||||
error_hash=error_hash,
|
||||
timestamp=timestamp,
|
||||
)
|
||||
json_handler.log_operation("system_log_event", {"level": level, "module": module_name})
|
||||
elif level == "warning":
|
||||
trigger.fire("warning_logged", **event_data)
|
||||
|
||||
@@ -48,8 +48,8 @@ class Trigger:
|
||||
_deferred_queue = [] # Queue for events fired during handling
|
||||
_draining_deferred = False # Prevents nested deferred processing
|
||||
_log_watcher_started = False # Lazy-start flag for log watcher
|
||||
_handler_failures = {} # handler -> consecutive failure count
|
||||
_disabled_handlers = set() # handlers auto-disabled after repeated failures
|
||||
_handler_failures = {} # (handler, branch) -> consecutive failure count
|
||||
_disabled_handlers = set() # (handler, branch) tuples auto-disabled
|
||||
_HANDLER_FAILURE_THRESHOLD = 5 # consecutive failures before auto-disable
|
||||
|
||||
@classmethod
|
||||
@@ -102,20 +102,22 @@ class Trigger:
|
||||
handlers = cls._handlers.get(event, [])
|
||||
data = dict(data) # Copy to avoid mutating caller's dict
|
||||
data["fire_event"] = cls.fire
|
||||
branch = data.get("branch", "__global__")
|
||||
for handler in handlers:
|
||||
if handler in cls._disabled_handlers:
|
||||
key = (handler, branch)
|
||||
if key in cls._disabled_handlers:
|
||||
continue
|
||||
try:
|
||||
handler(**data)
|
||||
cls._handler_failures.pop(handler, None) # Reset on success
|
||||
cls._handler_failures.pop(key, None) # Reset on success
|
||||
except Exception as e:
|
||||
count = cls._handler_failures.get(handler, 0) + 1
|
||||
cls._handler_failures[handler] = count
|
||||
count = cls._handler_failures.get(key, 0) + 1
|
||||
cls._handler_failures[key] = count
|
||||
if count >= cls._HANDLER_FAILURE_THRESHOLD:
|
||||
cls._disabled_handlers.add(handler)
|
||||
cls._disabled_handlers.add(key)
|
||||
logger.error(
|
||||
f"[TRIGGER] Handler {getattr(handler, '__name__', handler)} "
|
||||
f"disabled after {count} consecutive failures"
|
||||
f"disabled for branch '{branch}' after {count} consecutive failures"
|
||||
)
|
||||
else:
|
||||
logger.error(f"[TRIGGER] Handler error for {event}: {e}")
|
||||
|
||||
@@ -36,6 +36,7 @@ from aipass.trigger.apps.handlers.error_registry import (
|
||||
get_circuit_breaker_status,
|
||||
circuit_breaker_reset,
|
||||
update_source_fix_status,
|
||||
purge_stale,
|
||||
)
|
||||
from aipass.trigger.apps.handlers.error_reporter import ( # noqa: F401
|
||||
report_error,
|
||||
@@ -67,6 +68,7 @@ def print_introspection():
|
||||
console.print(" - error_registry.py (get_circuit_breaker_status — circuit breaker state)")
|
||||
console.print(" - error_registry.py (circuit_breaker_reset — reset circuit breaker)")
|
||||
console.print(" - error_registry.py (update_source_fix_status — update fix tracking)")
|
||||
console.print(" - error_registry.py (purge_stale — remove entries older than N days)")
|
||||
console.print(" - error_reporter.py (report_error — cross-branch push error reporting)")
|
||||
console.print(" - error_reporter.py (send_source_fix_email — notify branch to fix error)")
|
||||
console.print()
|
||||
@@ -132,6 +134,7 @@ def print_help() -> None:
|
||||
console.print(" [bold]suppress[/bold] <id> [reason] Mark error as suppressed")
|
||||
console.print(" [bold]resolve[/bold] <id> Mark error as resolved")
|
||||
console.print(" [bold]clear-resolved[/bold] Purge old resolved entries [dim](--days=7)[/dim]")
|
||||
console.print(" [bold]purge[/bold] Purge stale entries [dim](--days=30)[/dim]")
|
||||
console.print(" [bold]stats[/bold] Summary statistics + circuit breaker state")
|
||||
console.print(" [bold]circuit-breaker[/bold] Show or reset circuit breaker [dim](reset)[/dim]")
|
||||
console.print(" [bold]help[/bold] Show this help")
|
||||
@@ -178,6 +181,7 @@ def handle_command(command: str, args: list) -> bool:
|
||||
"suppress": _cmd_suppress,
|
||||
"resolve": _cmd_resolve,
|
||||
"clear-resolved": _cmd_clear_resolved,
|
||||
"purge": _cmd_purge,
|
||||
"stats": _cmd_stats,
|
||||
"circuit-breaker": _cmd_circuit_breaker,
|
||||
}
|
||||
@@ -285,7 +289,10 @@ def _cmd_detail(console, args: list) -> bool:
|
||||
f" [bold]First Seen:[/bold] {entry.get('first_seen', '?')}",
|
||||
f" [bold]Last Seen:[/bold] {entry.get('last_seen', '?')}",
|
||||
f" [bold]Log Path:[/bold] {entry.get('log_path', 'N/A') or 'N/A'}",
|
||||
f" [bold]Fix Status:[/bold] [{_FIX_STATUS_COLORS.get(entry.get('source_fix_status', 'none'), 'white')}]{entry.get('source_fix_status', 'none')}[/{_FIX_STATUS_COLORS.get(entry.get('source_fix_status', 'none'), 'white')}]",
|
||||
f" [bold]Fix Status:[/bold] "
|
||||
f"[{_FIX_STATUS_COLORS.get(entry.get('source_fix_status', 'none'), 'white')}]"
|
||||
f"{entry.get('source_fix_status', 'none')}"
|
||||
f"[/{_FIX_STATUS_COLORS.get(entry.get('source_fix_status', 'none'), 'white')}]",
|
||||
]
|
||||
if entry.get("suppress_reason"):
|
||||
lines.append(f" [bold]Suppress Reason:[/bold] {entry['suppress_reason']}")
|
||||
@@ -370,6 +377,16 @@ def _cmd_clear_resolved(console, args: list) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _cmd_purge(console, args: list) -> bool:
|
||||
"""Purge entries older than N days. Usage: purge [--days=30]"""
|
||||
parsed = _parse_args(args)
|
||||
days = int(parsed.get("days", "30"))
|
||||
removed = purge_stale(days=days)
|
||||
console.print(f"Purged {removed} entries older than {days} days")
|
||||
json_handler.log_operation("error_purge", {"days": days, "removed": removed})
|
||||
return True
|
||||
|
||||
|
||||
def _cmd_stats(console, args: list) -> bool:
|
||||
"""Show summary statistics and circuit breaker state."""
|
||||
stats = get_stats()
|
||||
|
||||
@@ -81,6 +81,8 @@ def trigger_cls():
|
||||
Trigger._deferred_queue = []
|
||||
Trigger._draining_deferred = False
|
||||
Trigger._log_watcher_started = False
|
||||
Trigger._handler_failures = {}
|
||||
Trigger._disabled_handlers = set()
|
||||
return Trigger
|
||||
|
||||
|
||||
@@ -458,3 +460,107 @@ def test_fire_event_kwarg_always_overwritten(trigger_cls):
|
||||
assert callback != "custom_value"
|
||||
assert callable(callback)
|
||||
assert getattr(callback, "__qualname__", "") == "Trigger.fire"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Tests -- per-(handler, branch) auto-disable
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_handler_disabled_after_threshold_failures(trigger_cls):
|
||||
"""A handler that fails 5 times for the same branch gets disabled for that branch."""
|
||||
fail_handler = MagicMock(side_effect=RuntimeError("boom"))
|
||||
trigger_cls.on("error_detected", fail_handler)
|
||||
|
||||
for _ in range(trigger_cls._HANDLER_FAILURE_THRESHOLD):
|
||||
trigger_cls.fire("error_detected", branch="DRONE")
|
||||
|
||||
# Handler should be disabled for DRONE now
|
||||
assert (fail_handler, "DRONE") in trigger_cls._disabled_handlers
|
||||
|
||||
# Fire again -- handler should NOT be called beyond the threshold invocations
|
||||
fail_handler.reset_mock()
|
||||
trigger_cls.fire("error_detected", branch="DRONE")
|
||||
fail_handler.assert_not_called()
|
||||
|
||||
|
||||
def test_handler_disabled_for_one_branch_still_fires_for_another(trigger_cls):
|
||||
"""Disabling a handler for branch A does not disable it for branch B."""
|
||||
call_count = 0
|
||||
|
||||
def flaky_handler(**kwargs):
|
||||
"""Raise only when branch is DRONE."""
|
||||
nonlocal call_count
|
||||
call_count += 1
|
||||
if kwargs.get("branch") == "DRONE":
|
||||
raise RuntimeError("noisy branch")
|
||||
|
||||
trigger_cls.on("error_detected", flaky_handler)
|
||||
|
||||
# Fail 5 times from DRONE to trigger auto-disable
|
||||
for _ in range(trigger_cls._HANDLER_FAILURE_THRESHOLD):
|
||||
trigger_cls.fire("error_detected", branch="DRONE")
|
||||
|
||||
assert (flaky_handler, "DRONE") in trigger_cls._disabled_handlers
|
||||
|
||||
# Fire from a different branch -- handler should still execute
|
||||
call_count = 0
|
||||
trigger_cls.fire("error_detected", branch="MEMORY")
|
||||
|
||||
assert call_count == 1
|
||||
assert (flaky_handler, "MEMORY") not in trigger_cls._disabled_handlers
|
||||
|
||||
|
||||
def test_handler_failure_count_resets_on_success(trigger_cls):
|
||||
"""A successful call resets the failure counter for that (handler, branch)."""
|
||||
attempt = 0
|
||||
|
||||
def sometimes_fails(**kwargs):
|
||||
"""Raise on the first 3 attempts, succeed on the 4th."""
|
||||
nonlocal attempt
|
||||
attempt += 1
|
||||
if attempt <= 3:
|
||||
raise RuntimeError("transient")
|
||||
|
||||
trigger_cls.on("check", sometimes_fails)
|
||||
|
||||
# Fail 3 times
|
||||
for _ in range(3):
|
||||
trigger_cls.fire("check", branch="API")
|
||||
|
||||
assert trigger_cls._handler_failures.get((sometimes_fails, "API")) == 3
|
||||
|
||||
# Succeed once -- counter should reset
|
||||
trigger_cls.fire("check", branch="API")
|
||||
assert (sometimes_fails, "API") not in trigger_cls._handler_failures
|
||||
|
||||
|
||||
def test_handler_failure_uses_global_when_no_branch(trigger_cls):
|
||||
"""When no branch is in data, the key uses '__global__' as the branch."""
|
||||
fail_handler = MagicMock(side_effect=RuntimeError("boom"))
|
||||
trigger_cls.on("some_event", fail_handler)
|
||||
|
||||
trigger_cls.fire("some_event")
|
||||
|
||||
assert (fail_handler, "__global__") in trigger_cls._handler_failures
|
||||
|
||||
|
||||
def test_different_branches_have_independent_failure_counts(trigger_cls):
|
||||
"""Failure counts are tracked independently per branch."""
|
||||
fail_handler = MagicMock(side_effect=RuntimeError("boom"))
|
||||
trigger_cls.on("error_detected", fail_handler)
|
||||
|
||||
# Fail 3 times from DRONE
|
||||
for _ in range(3):
|
||||
trigger_cls.fire("error_detected", branch="DRONE")
|
||||
|
||||
# Fail 2 times from MEMORY
|
||||
for _ in range(2):
|
||||
trigger_cls.fire("error_detected", branch="MEMORY")
|
||||
|
||||
assert trigger_cls._handler_failures[(fail_handler, "DRONE")] == 3
|
||||
assert trigger_cls._handler_failures[(fail_handler, "MEMORY")] == 2
|
||||
|
||||
# Neither should be disabled yet (threshold is 5)
|
||||
assert (fail_handler, "DRONE") not in trigger_cls._disabled_handlers
|
||||
assert (fail_handler, "MEMORY") not in trigger_cls._disabled_handlers
|
||||
|
||||
@@ -332,6 +332,30 @@ class TestHandleErrorDetectedHappyPath:
|
||||
fingerprint="fpX",
|
||||
)
|
||||
|
||||
def test_does_not_record_dispatch_when_send_fails(self) -> None:
|
||||
"""When _send_email returns False, dispatch is not recorded."""
|
||||
mod = _import_module()
|
||||
send = _setup_happy_path(mod)
|
||||
send.return_value = False
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_detected(
|
||||
branch="flow",
|
||||
module="cfg",
|
||||
message="err",
|
||||
error_hash="h1",
|
||||
count=2,
|
||||
fingerprint="fp_fail",
|
||||
)
|
||||
|
||||
# Email was attempted
|
||||
send.assert_called_once()
|
||||
# But nothing after it should have run
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
mod.registry_record_dispatch.assert_not_called() # type: ignore[attr-defined]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Fallback stubs (when error_registry import fails)
|
||||
|
||||
@@ -920,3 +920,95 @@ class TestEmptyJsonResilience:
|
||||
assert isinstance(result, dict)
|
||||
assert result["is_new"] is True
|
||||
assert result["count"] == 1
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# 16. purge_stale
|
||||
# ===========================================================================
|
||||
|
||||
|
||||
def test_purge_stale_removes_old_entries(tmp_path: Path) -> None:
|
||||
"""purge_stale removes entries older than N days regardless of status."""
|
||||
_seed_registry(tmp_path)
|
||||
er = _import_registry()
|
||||
|
||||
er.report("ImportError", "old error", "FLOW")
|
||||
er.report("IOError", "another old error", "BACKUP")
|
||||
|
||||
# Backdate last_seen to 60 days ago
|
||||
registry_file = tmp_path / "trigger_json" / "error_registry.json"
|
||||
data = json.loads(registry_file.read_text(encoding="utf-8"))
|
||||
for entry in data["errors"].values():
|
||||
entry["last_seen"] = "2025-01-01T00:00:00"
|
||||
registry_file.write_text(json.dumps(data, indent=2), encoding="utf-8")
|
||||
|
||||
removed = er.purge_stale(days=30)
|
||||
assert removed == 2
|
||||
|
||||
data_after = json.loads(registry_file.read_text(encoding="utf-8"))
|
||||
assert len(data_after["errors"]) == 0
|
||||
|
||||
|
||||
def test_purge_stale_keeps_recent_entries(tmp_path: Path) -> None:
|
||||
"""purge_stale keeps entries newer than the cutoff."""
|
||||
_seed_registry(tmp_path)
|
||||
er = _import_registry()
|
||||
|
||||
er.report("ImportError", "recent error", "FLOW")
|
||||
|
||||
# last_seen is set to now by report(), so it should survive a 30-day cutoff
|
||||
removed = er.purge_stale(days=30)
|
||||
assert removed == 0
|
||||
|
||||
|
||||
def test_purge_stale_ignores_status(tmp_path: Path) -> None:
|
||||
"""purge_stale removes old entries regardless of status (new, suppressed, etc.)."""
|
||||
_seed_registry(tmp_path)
|
||||
er = _import_registry()
|
||||
|
||||
er.report("ImportError", "old new error", "FLOW")
|
||||
r2 = er.report("IOError", "old suppressed error", "BACKUP")
|
||||
er.update_status(r2["fingerprint"], "suppressed", reason="known")
|
||||
|
||||
# Backdate both entries
|
||||
registry_file = tmp_path / "trigger_json" / "error_registry.json"
|
||||
data = json.loads(registry_file.read_text(encoding="utf-8"))
|
||||
for entry in data["errors"].values():
|
||||
entry["last_seen"] = "2025-01-01T00:00:00"
|
||||
registry_file.write_text(json.dumps(data, indent=2), encoding="utf-8")
|
||||
|
||||
removed = er.purge_stale(days=30)
|
||||
assert removed == 2
|
||||
|
||||
|
||||
def test_purge_stale_returns_zero_on_empty_registry(tmp_path: Path) -> None:
|
||||
"""purge_stale returns 0 when registry is empty."""
|
||||
_seed_registry(tmp_path)
|
||||
er = _import_registry()
|
||||
|
||||
removed = er.purge_stale(days=30)
|
||||
assert removed == 0
|
||||
|
||||
|
||||
def test_purge_stale_custom_days(tmp_path: Path) -> None:
|
||||
"""purge_stale respects custom days parameter."""
|
||||
_seed_registry(tmp_path)
|
||||
er = _import_registry()
|
||||
|
||||
er.report("ImportError", "semi-old error", "FLOW")
|
||||
|
||||
# Backdate to 10 days ago
|
||||
registry_file = tmp_path / "trigger_json" / "error_registry.json"
|
||||
data = json.loads(registry_file.read_text(encoding="utf-8"))
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
ten_days_ago = (datetime.now() - timedelta(days=10)).isoformat()
|
||||
for entry in data["errors"].values():
|
||||
entry["last_seen"] = ten_days_ago
|
||||
registry_file.write_text(json.dumps(data, indent=2), encoding="utf-8")
|
||||
|
||||
# 30-day cutoff should keep it
|
||||
assert er.purge_stale(days=30) == 0
|
||||
|
||||
# 7-day cutoff should remove it
|
||||
assert er.purge_stale(days=7) == 1
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
# Modified: 2026-04-25
|
||||
# =============================================
|
||||
|
||||
"""Tests for cli, error_logged, memory, memory_template_updated, warning_logged, and bulletin_created event handlers."""
|
||||
"""Tests for cli, memory_template_updated, warning_logged, and bulletin_created event handlers."""
|
||||
|
||||
import sys
|
||||
from pathlib import Path
|
||||
@@ -39,8 +39,6 @@ def _mock_infrastructure(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> Non
|
||||
|
||||
for mod_name in (
|
||||
"aipass.trigger.apps.handlers.events.cli",
|
||||
"aipass.trigger.apps.handlers.events.error_logged",
|
||||
"aipass.trigger.apps.handlers.events.memory",
|
||||
"aipass.trigger.apps.handlers.events.memory_template_updated",
|
||||
"aipass.trigger.apps.handlers.events.warning_logged",
|
||||
"aipass.trigger.apps.handlers.events.bulletin_created",
|
||||
@@ -55,20 +53,6 @@ def _import_cli():
|
||||
return m
|
||||
|
||||
|
||||
def _import_error_logged():
|
||||
"""Import error_logged handler module fresh after mocking."""
|
||||
import aipass.trigger.apps.handlers.events.error_logged as m
|
||||
|
||||
return m
|
||||
|
||||
|
||||
def _import_memory():
|
||||
"""Import memory handler module fresh after mocking."""
|
||||
import aipass.trigger.apps.handlers.events.memory as m
|
||||
|
||||
return m
|
||||
|
||||
|
||||
def _import_memory_template_updated():
|
||||
"""Import memory_template_updated handler module fresh after mocking."""
|
||||
import aipass.trigger.apps.handlers.events.memory_template_updated as m
|
||||
@@ -123,170 +107,6 @@ class TestHandleCliHeaderDisplayed:
|
||||
assert result is None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# error_logged.py -- handle_error_logged
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestHandleErrorLogged:
|
||||
"""Tests for handle_error_logged from error_logged.py."""
|
||||
|
||||
def test_happy_path_logs_operation(self) -> None:
|
||||
"""Logs error_logged_event with branch, module, and error_hash."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="flow", message="kaboom", error_hash="abc123")
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"error_logged_event",
|
||||
{"branch": "flow", "module": "unknown", "error_hash": "abc123"},
|
||||
)
|
||||
|
||||
def test_returns_early_missing_branch(self) -> None:
|
||||
"""Does not log when branch is None."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch=None, message="msg", error_hash="h")
|
||||
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_missing_message(self) -> None:
|
||||
"""Does not log when message is None."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="flow", message=None, error_hash="h")
|
||||
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_missing_error_hash(self) -> None:
|
||||
"""Does not log when error_hash is None."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="flow", message="msg", error_hash=None)
|
||||
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_empty_branch(self) -> None:
|
||||
"""Does not log when branch is empty string (falsy)."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="", message="msg", error_hash="h")
|
||||
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_prefers_source_module_over_module_name(self) -> None:
|
||||
"""Uses source_module when both source_module and module_name are given."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(
|
||||
branch="api",
|
||||
message="err",
|
||||
error_hash="xyz",
|
||||
source_module="config.py",
|
||||
module_name="old_name.py",
|
||||
)
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"error_logged_event",
|
||||
{"branch": "api", "module": "config.py", "error_hash": "xyz"},
|
||||
)
|
||||
|
||||
def test_falls_back_to_module_name(self) -> None:
|
||||
"""Uses module_name when source_module is not provided."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(
|
||||
branch="spawn",
|
||||
message="err",
|
||||
error_hash="h1",
|
||||
module_name="fallback.py",
|
||||
)
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"error_logged_event",
|
||||
{"branch": "spawn", "module": "fallback.py", "error_hash": "h1"},
|
||||
)
|
||||
|
||||
def test_falls_back_to_unknown(self) -> None:
|
||||
"""Uses 'unknown' when neither source_module nor module_name given."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="drone", message="oops", error_hash="h2")
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"error_logged_event",
|
||||
{"branch": "drone", "module": "unknown", "error_hash": "h2"},
|
||||
)
|
||||
|
||||
def test_exception_does_not_raise(self) -> None:
|
||||
"""Catches exception from log_operation without propagating."""
|
||||
mod = _import_error_logged()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.side_effect = RuntimeError("boom") # type: ignore[union-attr]
|
||||
|
||||
mod.handle_error_logged(branch="flow", message="msg", error_hash="h3")
|
||||
|
||||
json_handler.log_operation.side_effect = None # type: ignore[union-attr]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# memory.py -- handle_memory_saved
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestHandleMemorySaved:
|
||||
"""Tests for handle_memory_saved from memory.py."""
|
||||
|
||||
def test_calls_log_operation(self) -> None:
|
||||
"""Logs memory_event via json_handler."""
|
||||
mod = _import_memory()
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_saved()
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"memory_event", {"success": True}
|
||||
)
|
||||
|
||||
def test_accepts_kwargs(self) -> None:
|
||||
"""Does not crash when event data kwargs are passed."""
|
||||
mod = _import_memory()
|
||||
mod.handle_memory_saved(branch="flow", lines=150)
|
||||
|
||||
def test_returns_none(self) -> None:
|
||||
"""Handler returns None."""
|
||||
mod = _import_memory()
|
||||
result = mod.handle_memory_saved()
|
||||
assert result is None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# memory_template_updated.py -- handle_memory_template_updated
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -1,205 +0,0 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_memory_threshold_handler.py
|
||||
# Description: Tests for memory_threshold_exceeded event handler
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-04-25
|
||||
# Modified: 2026-04-25
|
||||
# =============================================
|
||||
|
||||
"""Tests for memory_threshold_exceeded event handler."""
|
||||
|
||||
import sys
|
||||
|
||||
import pytest
|
||||
from unittest.mock import MagicMock
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _mock_infrastructure(monkeypatch: pytest.MonkeyPatch, tmp_path: Path) -> None:
|
||||
"""Mock heavy infrastructure imports before importing the handler module."""
|
||||
from aipass.trigger.apps.config import atomic_write_json
|
||||
|
||||
mock_config = MagicMock()
|
||||
mock_config.TRIGGER_ROOT = tmp_path
|
||||
mock_config.atomic_write_json = atomic_write_json
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.config", mock_config)
|
||||
|
||||
mock_json_handler = MagicMock()
|
||||
mock_json_handler.log_operation = MagicMock(return_value=True)
|
||||
json_pkg = MagicMock()
|
||||
json_pkg.json_handler = mock_json_handler
|
||||
monkeypatch.setitem(sys.modules, "aipass.trigger.apps.handlers.json", json_pkg)
|
||||
monkeypatch.setitem(
|
||||
sys.modules,
|
||||
"aipass.trigger.apps.handlers.json.json_handler",
|
||||
mock_json_handler,
|
||||
)
|
||||
|
||||
# Mock ai_mail chain so the handler can import deliver_email_to_branch
|
||||
mock_email_send = MagicMock()
|
||||
mock_email_send.deliver_email_to_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.modules", MagicMock())
|
||||
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps.modules.email_send", mock_email_send)
|
||||
|
||||
monkeypatch.delitem(
|
||||
sys.modules,
|
||||
"aipass.trigger.apps.handlers.events.memory_threshold_exceeded",
|
||||
raising=False,
|
||||
)
|
||||
|
||||
|
||||
def _import_memory_threshold():
|
||||
"""Import fresh after mocking."""
|
||||
import aipass.trigger.apps.handlers.events.memory_threshold_exceeded as m
|
||||
|
||||
return m
|
||||
|
||||
|
||||
class TestHandleMemoryThresholdExceeded:
|
||||
"""Tests for handle_memory_threshold_exceeded."""
|
||||
|
||||
def test_returns_early_when_branch_missing(self) -> None:
|
||||
"""None branch skips email delivery."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(branch=None, file_name="local.json", line_count=700)
|
||||
|
||||
deliver_email_to_branch.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_when_file_name_missing(self) -> None:
|
||||
"""None file_name skips email delivery."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(branch="flow", file_name=None, line_count=700)
|
||||
|
||||
deliver_email_to_branch.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_when_line_count_is_none(self) -> None:
|
||||
"""None line_count skips email delivery."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(branch="flow", file_name="local.json", line_count=None)
|
||||
|
||||
deliver_email_to_branch.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_returns_early_when_ai_mail_import_fails(self, monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
"""No operation is logged when ai_mail import fails."""
|
||||
# Setting the module to None causes ImportError on `from ... import`
|
||||
monkeypatch.setitem(sys.modules, "aipass.ai_mail.apps.modules.email_send", None)
|
||||
|
||||
mod = _import_memory_threshold()
|
||||
mod.handle_memory_threshold_exceeded(branch="flow", file_name="local.json", line_count=700)
|
||||
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
|
||||
json_handler.log_operation.assert_not_called() # type: ignore[union-attr]
|
||||
|
||||
def test_happy_path_sends_email(self) -> None:
|
||||
"""Sends email with correct target and email data on valid input."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(
|
||||
branch="flow",
|
||||
file_name="local.json",
|
||||
line_count=700,
|
||||
threshold=600,
|
||||
timestamp="2026-04-25 12:00:00",
|
||||
)
|
||||
|
||||
deliver_email_to_branch.assert_called_once() # type: ignore[union-attr]
|
||||
target, email_data = deliver_email_to_branch.call_args[0] # type: ignore[union-attr]
|
||||
assert target == "@flow"
|
||||
assert email_data["to"] == "@flow"
|
||||
assert email_data["from"] == "@trigger"
|
||||
assert "local.json" in email_data["subject"]
|
||||
assert "600" in email_data["subject"]
|
||||
assert email_data["timestamp"] == "2026-04-25 12:00:00"
|
||||
|
||||
def test_uses_default_threshold_when_not_provided(self) -> None:
|
||||
"""Falls back to default threshold of 600 when not explicitly given."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(branch="drone", file_name="observations.json", line_count=800)
|
||||
|
||||
deliver_email_to_branch.assert_called_once() # type: ignore[union-attr]
|
||||
_, email_data = deliver_email_to_branch.call_args[0] # type: ignore[union-attr]
|
||||
assert "600" in email_data["subject"]
|
||||
|
||||
def test_uses_default_timestamp_when_not_provided(self) -> None:
|
||||
"""Generates a non-empty timestamp when none is supplied."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(branch="api", file_name="local.json", line_count=650)
|
||||
|
||||
deliver_email_to_branch.assert_called_once() # type: ignore[union-attr]
|
||||
_, email_data = deliver_email_to_branch.call_args[0] # type: ignore[union-attr]
|
||||
assert email_data["timestamp"] is not None
|
||||
assert len(email_data["timestamp"]) > 0
|
||||
|
||||
def test_does_not_raise_on_deliver_exception(self) -> None:
|
||||
"""Delivery failure is swallowed without propagating."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
deliver_email_to_branch.side_effect = RuntimeError( # type: ignore[union-attr]
|
||||
"delivery failed"
|
||||
)
|
||||
|
||||
mod.handle_memory_threshold_exceeded(
|
||||
branch="flow",
|
||||
file_name="local.json",
|
||||
line_count=700,
|
||||
threshold=600,
|
||||
)
|
||||
|
||||
deliver_email_to_branch.side_effect = None # type: ignore[union-attr]
|
||||
|
||||
def test_logs_operation_on_success(self) -> None:
|
||||
"""Logs memory_threshold_event via json_handler after successful send."""
|
||||
mod = _import_memory_threshold()
|
||||
|
||||
from aipass.trigger.apps.handlers.json import json_handler
|
||||
from aipass.ai_mail.apps.modules.email_send import deliver_email_to_branch
|
||||
|
||||
json_handler.log_operation.reset_mock() # type: ignore[union-attr]
|
||||
deliver_email_to_branch.reset_mock() # type: ignore[union-attr]
|
||||
|
||||
mod.handle_memory_threshold_exceeded(
|
||||
branch="system",
|
||||
file_name="local.json",
|
||||
line_count=900,
|
||||
threshold=600,
|
||||
)
|
||||
|
||||
json_handler.log_operation.assert_called_once_with( # type: ignore[union-attr]
|
||||
"memory_threshold_event", {"success": True}
|
||||
)
|
||||
Reference in New Issue
Block a user