Compare commits

...

3 Commits

11 changed files with 289 additions and 19 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
+1 -1
View File
@@ -1,7 +1,7 @@
"""Codex Mechanicus search tool (SPEC-014, FDB-019).
Luma's own sacred archive — the Codex Mechanicus at binaric.tech — as a
function tool. She searches the codex index and answers Cult Mechanicus
function tool. He searches the codex index and answers Cult Mechanicus
lore from real, sourced inscriptions instead of inventing it. The index
is fetched over HTTPS (SSRF-guarded, size-bounded, cached in memory) and
every field returned to the model is sanitized (SAF-03), because even
+26
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
@@ -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)
+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
+20 -12
View File
@@ -27,6 +27,7 @@ DEFAULT_SUMMARY_CHARS = 200
DEFAULT_NEWS_KEEP = 400
FETCH_TIMEOUT_S = 15
_ATOM = "{http://www.w3.org/2005/Atom}"
_RSS1 = "{http://purl.org/rss/1.0/}" # RSS 1.0 / RDF (e.g. 4gamer.net) namespaces <item>/<title>/<link>
_TAG_RE = re.compile(r"<[^>]+>")
@@ -36,22 +37,28 @@ def _clean_summary(raw: str, max_len: int = 300) -> str:
return re.sub(r"\s+", " ", text).strip()[:max_len]
def _rss_items(root: Any, ns: str, source: str) -> List[Dict[str, str]]:
"""RSS 2.0 (ns='') and RSS 1.0/RDF (ns=_RSS1) both use <item><title><link><description>."""
out: List[Dict[str, str]] = []
for item in root.iter(f"{ns}item"):
title = (item.findtext(f"{ns}title") or "").strip()
link = (item.findtext(f"{ns}link") or "").strip()
summary = _clean_summary(item.findtext(f"{ns}description") or "")
if title:
out.append({"title": title, "link": link, "source": source, "summary": summary})
return out
def parse_feed(data: bytes, source: str = "") -> List[Dict[str, str]]:
"""Parse RSS or Atom bytes into [{title, link, source}] (tolerant)."""
"""Parse RSS 2.0, RSS 1.0/RDF, or Atom bytes into [{title, link, source, summary}] (tolerant)."""
try:
root = ElementTree.fromstring(data)
except Exception as err:
# malformed XML or a blocked entity/DTD attack — tolerate, never raise (NEWS-01)
logging.warning(f"news: unparseable/unsafe feed {source!r}: {err!r}")
return []
items: List[Dict[str, str]] = []
# RSS: <rss><channel><item><title/><link/><description/>
for item in root.iter("item"):
title = (item.findtext("title") or "").strip()
link = (item.findtext("link") or "").strip()
summary = _clean_summary(item.findtext("description") or "")
if title:
items.append({"title": title, "link": link, "source": source, "summary": summary})
# RSS 2.0 (unqualified) + RSS 1.0/RDF (namespaced, e.g. 4gamer) share <item><title><link><description>
items: List[Dict[str, str]] = _rss_items(root, "", source) + _rss_items(root, _RSS1, source)
# Atom: <feed><entry><title/><link href=/><summary|content/>
for entry in root.iter(f"{_ATOM}entry"):
title = (entry.findtext(f"{_ATOM}title") or "").strip()
@@ -184,9 +191,10 @@ class NewsPoster:
GET_NEWS_TOOL = {
"name": "get_news",
"description": "Fetch recent real-world news the bot has collected from its RSS feeds (local, national, world, sport, "
"culture). Use when someone asks what is new or what is happening, optionally about a topic or from a particular "
"source. Returns headlines with a short summary and a link to read more.",
"description": "Fetch news the bot has collected from its RSS feeds — this is the SAME news that gets posted in the "
"server's news channels (e.g. #news, #newsjp / ニュース). Use this FIRST, before web_search, for anything about "
"current news or about something someone saw in a news channel; filter by topic (a keyword, also matches the source "
"label) or by source. Returns headlines with a short summary and a link; follow up with fetch_url for the full text.",
"parameters": {
"type": "object",
"properties": {
+4 -3
View File
@@ -24,9 +24,10 @@ FETCH_TIMEOUT_S = 15
WEB_SEARCH_TOOL = {
"name": "web_search",
"description": "Search the open web for current information when the user asks you to look something up and it is not "
"covered by game info (IGDB), the Codex, the news store, or a URL they pasted. Returns result titles, URLs, and a short "
"snippet; follow up with fetch_url on a result link for the full article.",
"description": "Search the open web for general information. Use ONLY when the answer is not in your own sources: for "
"the server's news use get_news, for Adeptus Mechanicus / Warhammer 40k lore use codex_search, for video-game facts use "
"the game tools, for a specific URL someone pasted use fetch_url. Returns result titles, URLs, and a short snippet; "
"follow up with fetch_url on a result link for the full article.",
"parameters": {
"type": "object",
"properties": {
+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).
+1 -1
View File
@@ -6,7 +6,7 @@ Replaces the broken pre-1.0-openai `news_feed.py`. A CLI
`AIResponder.message` injects into the `{news}` slot. Feeds are
external input and operator-configured.
### NEWS-01 — RSS and Atom parse to items (coverage: test)
### NEWS-01 — RSS (2.0 and 1.0/RDF) and Atom parse to items (coverage: test)
`parse_feed(bytes, label)` extracts `{title, link, source}` from both
RSS (`<item>`) and Atom (`<entry>`) documents, tolerates malformed
+2 -2
View File
@@ -1,8 +1,8 @@
# SPEC-014 — Codex Mechanicus search
Luma is an Adeptus Mechanicus tech-priest; her lore has a real home —
Luma is an Adeptus Mechanicus tech-priest; his lore has a real home —
the priest's own Codex Mechanicus at `binaric.tech` (an Astro/MDX
archive, five tongues). A `codex_search` function tool lets her consult
archive, five tongues). A `codex_search` function tool lets him consult
that archive and answer from sourced inscriptions instead of inventing
lore. The index is public but still untrusted by the time it reaches a
prompt: the fetch is SSRF-guarded (SPEC-011 shares `guard_url`),
+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)
+17
View File
@@ -29,6 +29,16 @@ ATOM = b"""<?xml version="1.0"?><feed xmlns="http://www.w3.org/2005/Atom">
<entry><title>Atom headline</title><link href="https://ex.com/a"/></entry>
</feed>"""
RSS1 = (
'<?xml version="1.0" encoding="UTF-8"?>'
'<rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#" xmlns="http://purl.org/rss/1.0/">'
'<channel rdf:about="https://ex.jp"><title>Feed</title></channel>'
'<item rdf:about="https://ex.jp/1"><title>ゲームニュース</title><link>https://ex.jp/1</link>'
"<description>本文ここ</description></item>"
'<item rdf:about="https://ex.jp/2"><title>Second</title><link>https://ex.jp/2</link></item>'
"</rdf:RDF>"
).encode("utf-8")
class TestParse(unittest.TestCase):
def test_rss(self):
@@ -44,6 +54,13 @@ class TestParse(unittest.TestCase):
self.assertEqual(items[0]["title"], "Atom headline")
self.assertEqual(items[0]["link"], "https://ex.com/a")
def test_rss1_rdf(self):
"""NEWS-01: RSS 1.0/RDF (namespaced <item>, e.g. 4gamer.net) parses like RSS 2.0."""
items = parse_feed(RSS1, "JP")
self.assertEqual([i["title"] for i in items], ["ゲームニュース", "Second"])
self.assertEqual(items[0]["link"], "https://ex.jp/1")
self.assertEqual(items[0]["summary"], "本文ここ")
def test_malformed_never_raises(self):
"""NEWS-01: garbage XML returns [] without raising."""
self.assertEqual(parse_feed(b"<not xml", "bad"), [])