health monitor (ops-18/19): spend/disk/task-queue thresholds -> staff alerts, edge-triggered, opt-in
This commit is contained in:
@@ -3,9 +3,11 @@ import asyncio
|
||||
import logging
|
||||
import random
|
||||
import re
|
||||
import shutil
|
||||
import sys
|
||||
import time
|
||||
from collections import deque
|
||||
from pathlib import Path
|
||||
from typing import Optional, Union
|
||||
|
||||
import discord
|
||||
@@ -16,6 +18,7 @@ from watchdog.events import FileSystemEventHandler
|
||||
from watchdog.observers import Observer
|
||||
|
||||
from .ai_responder import AIMessage
|
||||
from .monitor import HealthMonitor
|
||||
from .openai_responder import OpenAIResponder
|
||||
from .tasks import TaskEngine
|
||||
|
||||
@@ -130,6 +133,15 @@ class FjerkroaBot(commands.Bot):
|
||||
observe=self.airesponder.observe_event,
|
||||
)
|
||||
self.loop.create_task(self.task_loop())
|
||||
# Proactive health monitoring -> staff alerts (OPS-18/19)
|
||||
self.health_monitor = HealthMonitor(
|
||||
config_getter=lambda: self.config,
|
||||
ledger=self.airesponder.ledger,
|
||||
store=self.airesponder.store,
|
||||
disk_free_mb=self._disk_free_mb,
|
||||
alert=self.send_staff_alert,
|
||||
)
|
||||
self.loop.create_task(self.monitor_loop())
|
||||
logging.info("Task engine initialised.")
|
||||
|
||||
async def task_loop(self):
|
||||
@@ -140,6 +152,20 @@ class FjerkroaBot(commands.Bot):
|
||||
except Exception as err:
|
||||
logging.warning(f"task tick failed: {repr(err)}")
|
||||
|
||||
def _disk_free_mb(self) -> float:
|
||||
directory = Path(self.config.get("history-directory", ".")).expanduser()
|
||||
target = directory if directory.exists() else Path.home()
|
||||
return shutil.disk_usage(target).free / (1024 * 1024)
|
||||
|
||||
async def monitor_loop(self):
|
||||
while True:
|
||||
await asyncio.sleep(int(self.config.get("monitor-interval", 300)))
|
||||
if self.health_monitor.enabled():
|
||||
try:
|
||||
await self.health_monitor.tick()
|
||||
except Exception as err:
|
||||
logging.warning(f"monitor tick failed: {repr(err)}")
|
||||
|
||||
async def _execute_task(self, channel_name: str, prompt: str) -> None:
|
||||
"""Run a due task through the normal responder path (TSK-02)."""
|
||||
channel = self.channel_by_name(channel_name, getattr(self, "chat_channel", None), no_ignore=True)
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
"""Proactive health monitoring -> staff alerts (SPEC-012, FDB-012).
|
||||
|
||||
A periodic check that watches daily spend against the budget, free disk,
|
||||
and task-queue depth, and posts a staff alert when a threshold is crossed
|
||||
— once per crossing, re-arming when the metric recovers, so a persistent
|
||||
condition never spams. Opt-in per deployment (`enable-monitoring`); it
|
||||
reuses the rate-limited staff-alert channel (OPS-07).
|
||||
"""
|
||||
|
||||
import logging
|
||||
from typing import Any, Callable, Dict, Optional, Tuple
|
||||
|
||||
# A check returns (metric-name, is-over-threshold, alert-message) or None when not applicable.
|
||||
Check = Optional[Tuple[str, bool, str]]
|
||||
|
||||
|
||||
class HealthMonitor:
|
||||
def __init__(
|
||||
self,
|
||||
config_getter: Callable[[], Dict[str, Any]],
|
||||
ledger: Any,
|
||||
store: Any,
|
||||
disk_free_mb: Callable[[], float],
|
||||
alert: Callable[[str], Any],
|
||||
) -> None:
|
||||
self._config = config_getter
|
||||
self._ledger = ledger
|
||||
self._store = store
|
||||
self._disk_free_mb = disk_free_mb
|
||||
self._alert = alert
|
||||
self._armed: Dict[str, bool] = {}
|
||||
|
||||
def enabled(self) -> bool:
|
||||
return bool(self._config().get("enable-monitoring", False))
|
||||
|
||||
def _check_spend(self) -> Check:
|
||||
config = self._config()
|
||||
if "daily-budget-usd" not in config:
|
||||
return None
|
||||
budget = float(config["daily-budget-usd"])
|
||||
if budget <= 0:
|
||||
return None
|
||||
spent = float(self._ledger.spent_usd())
|
||||
frac = spent / budget
|
||||
threshold = float(config.get("monitor-spend-alert-frac", 0.8))
|
||||
return ("spend", frac >= threshold, f"💸 Spend at ${spent:.2f} of ${budget:.2f} today ({frac:.0%}, alert ≥ {threshold:.0%}).")
|
||||
|
||||
def _check_disk(self) -> Check:
|
||||
try:
|
||||
free = float(self._disk_free_mb())
|
||||
except Exception as err:
|
||||
logging.debug(f"monitor: disk check failed: {err!r}")
|
||||
return None
|
||||
min_mb = float(self._config().get("monitor-disk-min-mb", 500))
|
||||
return ("disk", free < min_mb, f"💾 Low disk: {free:.0f} MB free (alert < {min_mb:.0f} MB).")
|
||||
|
||||
def _check_queue(self) -> Check:
|
||||
if self._store is None:
|
||||
return None
|
||||
try:
|
||||
depth = len(self._store.tasks_open())
|
||||
except Exception as err:
|
||||
logging.debug(f"monitor: queue check failed: {err!r}")
|
||||
return None
|
||||
limit = int(self._config().get("monitor-taskqueue-max", 20))
|
||||
return ("task-queue", depth >= limit, f"🗒️ Task queue deep: {depth} open (alert ≥ {limit}).")
|
||||
|
||||
async def tick(self) -> None:
|
||||
"""Evaluate every check; alert on a rising edge only (OPS-18/19)."""
|
||||
for check in (self._check_spend(), self._check_disk(), self._check_queue()):
|
||||
if check is None:
|
||||
continue
|
||||
metric, over, message = check
|
||||
await self._fire(metric, over, message)
|
||||
|
||||
async def _fire(self, metric: str, over: bool, message: str) -> None:
|
||||
was_over = self._armed.get(metric, False)
|
||||
if over and not was_over:
|
||||
self._armed[metric] = True
|
||||
await self._alert(message)
|
||||
elif not over and was_over:
|
||||
self._armed[metric] = False # recovered — re-arm silently for the next crossing
|
||||
Reference in New Issue
Block a user