Compare commits

...

3 Commits

8 changed files with 315 additions and 16 deletions
+10
View File
@@ -60,6 +60,16 @@ Decisions inside the set architecture. D-NNN, never renumbered.
broken classifier must never mute the bot; the budget gate already
bounds spend. Its verdict gates BEFORE the main call, the
envelope's answer_needed still gates after — two independent nets.
- **D-021** — Health monitoring (FDB-012, SPEC-012 OPS-18/19): a
separate `monitor_loop` (own cadence, default 300 s) rather than
folding checks into the 60 s task loop — monitoring is coarse and
should not run every minute. Checks are edge-triggered (alert on the
rising edge, re-arm on recovery) so a standing condition never spams;
they reuse the existing rate-limited staff-alert path. Metrics are
the cheap, high-signal ones (spend vs budget, free disk, task-queue
depth); each is independently skippable when it has no data, so a
deployment without a budget or store still runs the others. Opt-in
(`enable-monitoring`) like every other operational rollout.
- **D-020** — Web search via Exa (FDB-022, SPEC-015): a `web_search`
tool alongside fetch_url/IGDB/codex/get_news, filling the "look it up
on the open web" gap. Exa (not a raw search-engine scrape) because it
+53 -9
View File
@@ -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
@@ -65,11 +68,29 @@ def split_answer(text: str, threshold: int, max_parts: int) -> list:
class ConfigFileHandler(FileSystemEventHandler):
def __init__(self, on_modified):
self._on_modified = on_modified
"""Rename-safe config watch (CFG-05).
def on_modified(self, event):
self._on_modified(event)
Editors and tools save atomically — write a temp file, then rename it
over the target — which fires a *moved*/*created* event (not
*modified*) and swaps the inode, so watching the file directly goes
deaf after the first save. We watch the config's *directory* and react
to any event whose src or dest path is the config file.
"""
def __init__(self, config_path: str, on_change):
self._config_path = str(Path(config_path).resolve())
self._on_change = on_change
def _hits_config(self, event) -> bool:
for attr in ("src_path", "dest_path"):
path = getattr(event, attr, "")
if path and str(Path(path).resolve()) == self._config_path:
return True
return False
def on_any_event(self, event):
if not event.is_directory and self._hits_config(event):
self._on_change()
class FjerkroaBot(commands.Bot):
@@ -99,8 +120,10 @@ class FjerkroaBot(commands.Bot):
def init_observer(self):
self.observer = Observer()
self.file_handler = ConfigFileHandler(self.on_config_file_modified)
self.observer.schedule(self.file_handler, path=self.config_file, recursive=False)
config_path = Path(self.config_file).resolve()
self.file_handler = ConfigFileHandler(str(config_path), self.on_config_file_changed)
# Watch the directory, not the file — atomic saves replace the inode (CFG-05)
self.observer.schedule(self.file_handler, path=str(config_path.parent), recursive=False)
self.observer.start()
def init_aichannels(self):
@@ -130,6 +153,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 +172,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)
@@ -389,12 +435,10 @@ class FjerkroaBot(commands.Bot):
airesponder.image_cache.purge_message(str(message.id)) # IMG-14
await airesponder.observe_event(message.author.name, "delete", f"deleted: {message.content}")
def on_config_file_modified(self, event):
def on_config_file_changed(self):
# Runs on the watchdog observer thread — the swap itself is
# scheduled onto the event loop so no request reads a
# half-swapped config (CFG-04 / D9)
if event.src_path != self.config_file:
return
new_config = self.load_config(self.config_file)
if repr(new_config) == repr(self.config):
return
+82
View File
@@ -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
+5 -1
View File
@@ -17,7 +17,11 @@ from .quota import QuotaLedger
DEFAULT_MAX_PER_CHANNEL_PER_DAY = 2
DEFAULT_IDLE_IMPULSE_HOURS = 12.0
DEFAULT_TASKGEN_INTERVAL_HOURS = 6.0
DEFAULT_BORENESS_PROMPT = "Pretend that you just now thought of something, be creative."
DEFAULT_BORENESS_PROMPT = (
"A thought just occurred to you. Anchor it to something real you know — recent news (use get_news), a game "
"releasing soon, the weather, or a regular you remember — not a generic musing. Share it briefly, in your own "
"voice, as an observation, a gentle question, or a joke; never an advertisement. Read the room and stay in character."
)
ExecuteCallback = Callable[[str, str], Awaitable[None]]
ProposeCallback = Callable[[], Awaitable[Optional[Dict[str, Any]]]]
+10
View File
@@ -31,3 +31,13 @@ and all responder `.config` references is scheduled onto the event
loop (`call_soon_threadsafe`), so no request ever reads a
half-swapped config (D9). Before the loop runs (startup), the swap
applies directly — there are no concurrent readers yet.
### CFG-05 — Hot-reload is rename-safe (coverage: test)
The watcher observes the config file's **directory**, not the file, and
reacts to any event (modified, created, moved) whose source or
destination path is the config file. This catches atomic saves — write
a temp file, then rename it over the target — which replace the inode
and fire a move/create rather than a modify; watching the file directly
would go deaf after the first such save. Events for other files in the
directory, and directory events themselves, are ignored.
+20
View File
@@ -30,3 +30,23 @@ The responder counts consecutive OpenAI request failures; at
alert (rate-limited like all staff alerts) so a silently-broken bot
(cf. the gpt-5.6 tools/reasoning incident) surfaces within minutes
instead of hours. A success resets the counter.
### OPS-18 — Health monitor watches spend, disk, task-queue (coverage: test)
When `enable-monitoring` is true, a loop wakes every `monitor-interval`
(default 300 s) and checks three thresholds, alerting the staff channel
when one is crossed: daily spend at or above `monitor-spend-alert-frac`
(default 0.8) of `daily-budget-usd`; free disk below `monitor-disk-min-mb`
(default 500 MB); open task-queue depth at or above `monitor-taskqueue-max`
(default 20). A check with no data to evaluate (no budget set, no store,
a failed disk read) is skipped, never fatal. With the flag off the loop
does nothing.
### OPS-19 — Alerts fire once per crossing and re-arm on recovery (coverage: test)
Each metric alerts only on the rising edge — the first tick that finds
it over its threshold — and stays silent while it remains over, so a
persistent condition does not repeat every interval. When the metric
falls back below the threshold the alert re-arms silently, ready to fire
again on the next crossing. All alerts still pass through the
rate-limited staff-alert path (OPS-07).
+106
View File
@@ -0,0 +1,106 @@
"""Unit coverage for SPEC-012 health monitoring (OPS-18/19)."""
import unittest
from unittest.mock import AsyncMock
from fjerkroa_bot.monitor import HealthMonitor
class FakeLedger:
def __init__(self, spent=0.0):
self._spent = spent
def spent_usd(self):
return self._spent
class FakeStore:
def __init__(self, open_tasks=0):
self._n = open_tasks
def tasks_open(self):
return list(range(self._n))
def _monitor(cfg, ledger=None, store=None, disk=1000.0):
alert = AsyncMock()
monitor = HealthMonitor(lambda: cfg, ledger or FakeLedger(), store, lambda: disk, alert)
return monitor, alert
class TestEnabled(unittest.TestCase):
def test_opt_in(self):
"""OPS-18: monitoring is opt-in via enable-monitoring."""
self.assertFalse(_monitor({})[0].enabled())
self.assertTrue(_monitor({"enable-monitoring": True})[0].enabled())
class TestChecks(unittest.IsolatedAsyncioTestCase):
async def test_spend_over_threshold_alerts(self):
"""OPS-18: spend at/above frac*budget alerts."""
monitor, alert = _monitor({"daily-budget-usd": 2.0, "monitor-spend-alert-frac": 0.8}, ledger=FakeLedger(1.8))
await monitor.tick()
alert.assert_awaited_once()
self.assertIn("Spend", alert.await_args.args[0])
async def test_spend_under_threshold_silent(self):
"""OPS-18: spend below threshold stays silent."""
monitor, alert = _monitor({"daily-budget-usd": 2.0}, ledger=FakeLedger(0.5))
await monitor.tick()
alert.assert_not_awaited()
async def test_no_budget_skips_spend(self):
"""OPS-18: no budget configured -> spend check skipped, never fatal."""
monitor, alert = _monitor({"enable-monitoring": True}, ledger=FakeLedger(99))
await monitor.tick()
alert.assert_not_awaited()
async def test_low_disk_alerts(self):
"""OPS-18: free disk below the floor alerts."""
monitor, alert = _monitor({"monitor-disk-min-mb": 500}, disk=100.0)
await monitor.tick()
self.assertTrue(any("Low disk" in call.args[0] for call in alert.await_args_list))
async def test_disk_read_failure_skipped(self):
"""OPS-18: a failing disk read is skipped, not fatal."""
def boom():
raise OSError("nope")
alert = AsyncMock()
monitor = HealthMonitor(lambda: {}, FakeLedger(), None, boom, alert)
await monitor.tick()
alert.assert_not_awaited()
async def test_deep_queue_alerts(self):
"""OPS-18: task-queue depth at/above max alerts."""
monitor, alert = _monitor({"monitor-taskqueue-max": 3}, store=FakeStore(5), disk=9999.0)
await monitor.tick()
self.assertTrue(any("Task queue" in call.args[0] for call in alert.await_args_list))
async def test_no_store_skips_queue(self):
"""OPS-18: no store -> queue check skipped."""
monitor, alert = _monitor({"monitor-taskqueue-max": 1}, store=None, disk=9999.0)
await monitor.tick()
alert.assert_not_awaited()
class TestEdgeArming(unittest.IsolatedAsyncioTestCase):
async def test_fires_once_per_crossing(self):
"""OPS-19: a persistent over-threshold condition alerts once, not every tick."""
monitor, alert = _monitor({"daily-budget-usd": 2.0}, ledger=FakeLedger(1.9))
await monitor.tick()
await monitor.tick()
await monitor.tick()
self.assertEqual(alert.await_count, 1)
async def test_rearms_on_recovery(self):
"""OPS-19: recovery re-arms silently; the next crossing alerts again."""
ledger = FakeLedger(1.9)
monitor, alert = _monitor({"daily-budget-usd": 2.0}, ledger=ledger, disk=9999.0)
await monitor.tick() # over -> alert (1)
ledger._spent = 0.5
await monitor.tick() # recovered -> silent, re-arm
ledger._spent = 1.95
await monitor.tick() # over again -> alert (2)
self.assertEqual(alert.await_count, 2)
+29 -6
View File
@@ -120,9 +120,7 @@ class TestConfigReloadRace(TestBotBase):
new_config["history-limit"] = 99
self.bot.load_config = lambda path: new_config
self.bot.loop = MagicMock()
event = MagicMock()
event.src_path = self.bot.config_file
self.bot.on_config_file_modified(event)
self.bot.on_config_file_changed()
self.bot.loop.call_soon_threadsafe.assert_called_once()
apply_fn = self.bot.loop.call_soon_threadsafe.call_args.args[0]
apply_fn()
@@ -137,7 +135,32 @@ class TestConfigReloadRace(TestBotBase):
loop = MagicMock()
loop.call_soon_threadsafe.side_effect = RuntimeError("no running loop")
self.bot.loop = loop
event = MagicMock()
event.src_path = self.bot.config_file
self.bot.on_config_file_modified(event)
self.bot.on_config_file_changed()
self.assertEqual(self.bot.config["history-limit"], 42)
class TestConfigReloadRenameSafe(unittest.TestCase):
def test_atomic_rename_and_modify_trigger_reload(self):
"""CFG-05: a modified OR a renamed-into-place config fires the reload; unrelated files do not."""
from fjerkroa_bot.discord_bot import ConfigFileHandler
with tempfile.TemporaryDirectory() as tmp:
config = Path(tmp) / "kroa.toml"
config.write_text("x = 1\n")
hits = []
handler = ConfigFileHandler(str(config), lambda: hits.append(1))
def evt(is_dir=False, src=None, dest=None):
event = MagicMock()
event.is_directory = is_dir
event.src_path = src if src is not None else ""
event.dest_path = dest if dest is not None else ""
return event
handler.on_any_event(evt(src=str(config))) # in-place modify
handler.on_any_event(evt(src=str(Path(tmp) / "kroa.toml.tmp"), dest=str(config))) # atomic rename over
self.assertEqual(len(hits), 2)
handler.on_any_event(evt(src=str(Path(tmp) / "other.txt"))) # unrelated file
handler.on_any_event(evt(is_dir=True, src=str(config))) # directory event
self.assertEqual(len(hits), 2) # neither fired