Compare commits

..

7 Commits

Author SHA1 Message Date
Oleksandr Kozachuk 2caa18a17f igdb search + capped http reads: native fulltext, full-page fetch
search_games used `where name ~ "<q>"*`: prefix-only, diacritic- and
word-order-sensitive -- "MARVEL Tōkon" found nothing though IGDB has
it. Switch to IGDB's native `search` clause (relevance-ranked,
diacritic-insensitive).

fetch_url and image downloads read bodies with content.read(n), which
returns only the first buffered chunk (~7 KB): pages collapsed to
their <title>. httpread.read_capped collects chunks up to the byte
cap; image downloads read limit+1 so over-limit files are still
rejected instead of cached truncated (IMG-10).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-13 19:29:16 +02:00
Oleksandr Kozachuk d4eec4088d news digest (spec-013): rss/atom -> {news} file via cron
Replaces the broken pre-1.0-openai news_feed.py. Stdlib parsing with
defusedxml (feeds are untrusted XML), titles sanitized (SAF-03), feed
URLs SSRF-guarded. CLI: python -m fjerkroa_bot.news --config <cfg>.
2026-07-13 19:28:41 +02:00
Oleksandr Kozachuk 86e631926f igdb: auto-refresh twitch token; category -> game_type filter
Static app tokens expire after ~60 days -> every lookup failed with
401. With igdb-client-secret set, the bot fetches the token via
client-credentials OAuth itself, refreshes a day before expiry and
retries once on 401; a static igdb-access-token still works.

IGDB renamed games.category to game_type: the category = 0 filter in
search_games silently matched nothing even with a valid token.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-13 18:41:31 +02:00
Oleksandr Kozachuk 4166520923 ops: consistent rotated db backups + cron, consecutive-api-error staff alert 2026-07-13 18:25:01 +02:00
Oleksandr Kozachuk d0819c2683 url reading tool: fetch_url with ssrf guard, html->text, page images to vision cache 2026-07-13 18:10:54 +02:00
Oleksandr Kozachuk df1924bb80 self-tasking engine: persistent queue, idle-impulse + follow-up generators, approval mode 2026-07-13 17:49:39 +02:00
Oleksandr Kozachuk e7e51e4230 img-17: cache image-only posts; news pipeline notes 2026-07-13 16:50:58 +02:00
29 changed files with 1775 additions and 107 deletions
+21 -8
View File
@@ -20,13 +20,6 @@ The bot now supports real-time video game information through IGDB (Internet Gam
- **Category**: Select appropriate category
3. Note down your **Client ID**
4. Generate a **Client Secret**
5. Get an access token using this curl command:
```bash
curl -X POST 'https://id.twitch.tv/oauth2/token' \
-H 'Content-Type: application/x-www-form-urlencoded' \
-d 'client_id=YOUR_CLIENT_ID&client_secret=YOUR_CLIENT_SECRET&grant_type=client_credentials'
```
6. Save the `access_token` from the response
### 2. Configure the Bot
@@ -35,6 +28,25 @@ Update your `config.toml` file:
```toml
# IGDB Configuration for game information
igdb-client-id = "your_actual_client_id_here"
igdb-client-secret = "your_actual_client_secret_here"
enable-game-info = true
```
With the client secret configured, the bot fetches an app access token from
Twitch itself and refreshes it automatically before it expires (Twitch app
tokens live ~60 days) — no manual token handling needed.
Alternatively, a static token still works (legacy setup — it expires after
~60 days and then game lookups fail with 401 until you replace it):
```bash
curl -X POST 'https://id.twitch.tv/oauth2/token' \
-H 'Content-Type: application/x-www-form-urlencoded' \
-d 'client_id=YOUR_CLIENT_ID&client_secret=YOUR_CLIENT_SECRET&grant_type=client_credentials'
```
```toml
igdb-client-id = "your_actual_client_id_here"
igdb-access-token = "your_actual_access_token_here"
enable-game-info = true
```
@@ -99,7 +111,8 @@ The integration provides two OpenAI functions:
- Verify client ID and access token are set
2. **Authentication errors**
- Regenerate access token (they expire)
- Prefer `igdb-client-secret` — the bot then refreshes tokens itself
- With a static `igdb-access-token`: regenerate it (they expire)
- Verify client ID matches your Twitch app
3. **No game results**
+3
View File
@@ -83,3 +83,6 @@ ci: install-dev all-checks ## Full CI pipeline (install deps and run all checks)
# Deploy targets (SPEC-007)
deploy: ## Deploy a tag to a host: make deploy HOST=ggg TAG=v3.0.0
bash deploy/deploy.sh $(HOST) $(TAG)
backup: ## Back up a local bot.db: make backup DB=history/bot.db DIR=backups
uv run python deploy/backup_db.py $(DB) $(DIR) $(or $(KEEP),14)
+34 -1
View File
@@ -14,7 +14,11 @@ system = "You are a smart AI assistant with access to real-time video game infor
# IGDB Configuration for game information
igdb-client-id = "YOUR_IGDB_CLIENT_ID"
igdb-access-token = "YOUR_IGDB_ACCESS_TOKEN"
# With the Twitch app client secret set, the bot fetches and refreshes the
# access token itself (recommended). A static igdb-access-token still works
# but expires after ~60 days.
igdb-client-secret = "YOUR_IGDB_CLIENT_SECRET"
# igdb-access-token = "YOUR_IGDB_ACCESS_TOKEN"
enable-game-info = true
# --- operator / safety (SPEC-003, SPEC-006) ---
@@ -59,3 +63,32 @@ enable-game-info = true
# image-cache-mb = 500 # LRU cap (ggg: consider 2000 — screenshots)
# image-cache-ttl-days = 90
# image-max-bytes = 8388608 # 8 MB upload cap
# Self-tasking (SPEC-005) — experimental, DEFAULT OFF:
# tasks-enabled = true
# tasks-generators = ["idle-impulse", "follow-up"]
# tasks-max-per-channel-per-day = 2
# tasks-approval = false # true: neue Tasks brauchen !bot task-approve
# idle-impulse-hours = 12
# taskgen-interval-hours = 6
# Staff: !bot tasks | task-approve <id> | task-cancel <id>
# URL reading (SPEC-011, FDB-018) — DEFAULT OFF; web pages are hostile input:
# enable-url-reading = true
# url-max-bytes = 2097152 # 2 MB fetch cap
# url-max-chars = 6000 # text handed to the model
# url-max-images = 2 # page images into the vision cache
# url-daily-per-user = 20
# Ops (SPEC-012): consecutive OpenAI failures before a staff alert
# api-error-alert-threshold = 5
# Backups: cron runs deploy/backup_db.py daily -> ~/backups/<bot>/ (keep 14)
# News digest (SPEC-013) — `python -m fjerkroa_bot.news --config X.toml` via cron;
# writes the {news} file. Feeds are [url, label] pairs (RSS or Atom):
# news = "news_feed.txt"
# news-per-feed = 3
# news-max-items = 15
# news-feeds = [
# ["https://blog.playstation.com/feed/", "PS"],
# ["https://kotaku.com/rss", "Kotaku"],
# ["https://www.pushsquare.com/feeds/latest", "Push"],
# ["https://mein-mmo.de/feed/", "MeinMMO"],
# ]
+80
View File
@@ -0,0 +1,80 @@
#!/usr/bin/env python3
"""Consistent, rotated bot.db backups (SPEC-012 OPS-13).
Run from cron on each host. Uses the sqlite3 online-backup API so the
snapshot is consistent even while the bot writes (WAL-safe), gzips it,
and keeps the newest N. Stdlib only.
Usage: python3 backup_db.py <bot.db> <backup-dir> [keep]
"""
import gzip
import os
import shutil
import sqlite3
import sys
import tempfile
import time
from pathlib import Path
DEFAULT_KEEP = 14
BACKUP_GLOB = "bot-*.db.gz"
def snapshot(src: Path, dest_gz: Path) -> None:
"""Write a consistent gzipped snapshot of src to dest_gz (OPS-13)."""
fd, tmp_path = tempfile.mkstemp(suffix=".db", dir=str(dest_gz.parent))
os.close(fd)
tmp = Path(tmp_path)
try:
source = sqlite3.connect(str(src))
try:
target = sqlite3.connect(str(tmp))
try:
source.backup(target) # atomic, WAL-safe online backup
finally:
target.close()
finally:
source.close()
with open(tmp, "rb") as raw, gzip.open(str(dest_gz), "wb") as gz:
shutil.copyfileobj(raw, gz)
os.chmod(dest_gz, 0o600) # conversation data
finally:
tmp.unlink(missing_ok=True)
def victims(existing: list, keep: int) -> list:
"""Given backup paths (any order), return the ones to delete, oldest first (OPS-14)."""
ordered = sorted(existing) # timestamped names sort chronologically
return ordered[: max(0, len(ordered) - keep)]
def rotate(backup_dir: Path, keep: int) -> int:
removed = 0
for path in victims(list(backup_dir.glob(BACKUP_GLOB)), keep):
Path(path).unlink(missing_ok=True)
removed += 1
return removed
def main() -> int:
if len(sys.argv) < 3:
print("usage: backup_db.py <bot.db> <backup-dir> [keep]", file=sys.stderr)
return 2
src = Path(sys.argv[1]).expanduser()
backup_dir = Path(sys.argv[2]).expanduser()
keep = int(sys.argv[3]) if len(sys.argv) > 3 else DEFAULT_KEEP
if not src.exists():
print(f"backup: source {src} missing", file=sys.stderr)
return 1
backup_dir.mkdir(parents=True, exist_ok=True)
stamp = time.strftime("%Y%m%d-%H%M%S", time.gmtime())
dest = backup_dir / f"bot-{stamp}.db.gz"
snapshot(src, dest)
removed = rotate(backup_dir, keep)
print(f"backup: wrote {dest.name} ({dest.stat().st_size} bytes), rotated {removed} old")
return 0
if __name__ == "__main__":
sys.exit(main())
+97 -50
View File
@@ -1,7 +1,6 @@
import argparse
import asyncio
import logging
import math
import random
import re
import sys
@@ -18,6 +17,7 @@ from watchdog.observers import Observer
from .ai_responder import AIMessage
from .openai_responder import OpenAIResponder
from .tasks import TaskEngine
DEFAULT_PRIVACY_NOTICE = (
"I keep recent channel messages and a short conversation summary to answer better. "
@@ -89,6 +89,7 @@ class FjerkroaBot(commands.Bot):
self.tasks_enabled = True
self.quiet_until = 0.0
self._staff_alert_times: deque = deque()
self._consecutive_api_errors = 0 # OPS-16
self.init_observer()
self.init_aichannels()
@@ -114,38 +115,42 @@ class FjerkroaBot(commands.Bot):
self.staff_channel = self.channel_by_name(self.config["staff-channel"], no_ignore=True)
self.welcome_channel = self.channel_by_name(self.config["welcome-channel"], no_ignore=True)
def init_boreness(self):
if "chat-channel" not in self.config:
return
def init_tasks(self):
"""Task engine replaces the sigmoid boreness loop (TSK-07)."""
self.last_activity_time = time.monotonic()
self.loop.create_task(self.on_boreness())
logging.info("Boreness initialised.")
self.task_engine = TaskEngine(
store=self.airesponder.store,
ledger=self.airesponder.ledger,
config_getter=lambda: self.config,
execute=self._execute_task,
propose=self.airesponder.propose_task,
staff_alert=self.send_staff_alert,
allowed=self.bot_initiated_allowed,
idle_seconds=lambda: time.monotonic() - self.last_activity_time,
observe=self.airesponder.observe_event,
)
self.loop.create_task(self.task_loop())
logging.info("Task engine initialised.")
async def on_boreness(self):
logging.info(f"Boreness started on channel: {repr(self.chat_channel)}")
async def task_loop(self):
while True:
if self.chat_channel is None or not self.bot_initiated_allowed():
await asyncio.sleep(7)
continue
boreness_interval = float(self.config.get("boreness-interval", 12.0))
elapsed_time = (time.monotonic() - self.last_activity_time) / 3600.0
probability = 1 / (1 + math.exp(-1 * (elapsed_time - (boreness_interval / 2.0)) + math.log(1 / 0.2 - 1)))
if random.random() < probability:
prev_messages = [msg async for msg in self.chat_channel.history(limit=2)]
last_author = prev_messages[1].author.id if len(prev_messages) > 1 else None
if last_author and last_author != self.user.id:
logging.info(f"Borred with {probability} probability after {elapsed_time}")
boreness_prompt = self.config.get("boreness-prompt", "Pretend that you just now thought of something, be creative.")
message = AIMessage("system", boreness_prompt, self.config.get("chat-channel", "chat"), True, False)
try:
await self.respond(message, self.chat_channel)
except Exception as err:
logging.warning(f"Failed to activate borringness: {repr(err)}")
await asyncio.sleep(7)
await asyncio.sleep(60)
try:
await self.task_engine.tick()
except Exception as err:
logging.warning(f"task 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)
if channel is None:
raise RuntimeError(f"task channel {channel_name!r} not resolvable")
message = AIMessage("system", prompt, channel_name, True, False)
await self.respond(message, channel)
async def on_ready(self):
self.init_channels()
self.init_boreness()
self.init_tasks()
logging.info(
f"We have logged in as {self.user}" f" ({repr(self.staff_channel)}, {repr(self.welcome_channel)}, {repr(self.chat_channel)})"
)
@@ -241,13 +246,31 @@ class FjerkroaBot(commands.Bot):
return "\n".join(f"{pin['id']} [{pin['channel'] or 'global'}]: {pin['fact']}" for pin in pins) or "No pins."
return None
def _task_command(self, args) -> Optional[str]:
"""Task queue surface (OPS-12, TSK-05)."""
is_list = args[:1] == ["tasks"] and len(args) == 1
if args[:1] not in (["task-approve"], ["task-cancel"]) and not is_list:
return None
store = self.airesponder.store
if store is None:
return "No store configured - task commands unavailable."
if is_list:
tasks = store.tasks_open()
return "\n".join(f"{t['id']} [{t['state']}] {t['kind']} #{t['channel']} due {t['due_at']}" for t in tasks) or "No open tasks."
if args[1:2] and args[1].isdigit():
if args[0] == "task-approve":
return f"Approved {store.task_set_state(int(args[1]), 'queued')} task(s)."
return f"Cancelled {store.task_set_state(int(args[1]), 'cancelled')} task(s)."
return None
async def handle_staff_command(self, message: Message) -> None:
"""Operator kill-switches, staff channel only (OPS-01..05, OPS-09, MEM-07)."""
args = str(message.content).split()[1:]
memory_reply = self._memory_command(args)
if memory_reply is not None:
await message.channel.send(memory_reply, suppress_embeds=True)
return
for handler in (self._memory_command, self._task_command):
reply = handler(args)
if reply is not None:
await message.channel.send(reply, suppress_embeds=True)
return
reply = "Commands: pause, resume, images on|off, tasks on|off, quiet <minutes>, status, spend, memory <user>, forget-fact <id>, pin <channel|global> <fact>, unpin <id>"
if args[:1] == ["pause"]:
self.replies_enabled = False
@@ -401,35 +424,44 @@ class FjerkroaBot(commands.Bot):
def get_ai_responder(self, channel_name):
return self.aichannels[channel_name] if channel_name in self.aichannels else self.airesponder
async def _ingest_attachments(self, message, channel_name: str, airesponder) -> list:
"""Cache-first attachment handling; CDN URLs never travel further (IMG-10/11)."""
urls = []
for attachment in message.attachments:
if airesponder.image_cache is None:
urls.append(attachment.url)
continue
sha = await airesponder.image_cache.ingest_url(attachment.url, channel_name, message.author.name, str(message.id))
if sha is not None:
recent = airesponder.image_cache.recent(channel_name, 8)
ext = next((row["ext"] for row in recent if row["sha256"] == sha), "png")
data_url = airesponder.image_cache.data_url(sha, ext)
if data_url:
urls.append(data_url)
return urls
async def handle_message_through_responder(self, message):
"""Handle a message through the AI responder"""
message_content = str(message.content).strip()
if message.reference and message.reference.resolved and isinstance(message.reference.resolved.content, str):
reference_content = str(message.reference.resolved.content).replace("\n", "> \n")
message_content = f"> {reference_content}\n\n{message_content}"
channel_name = self.get_channel_name(message.channel)
airesponder = self.get_ai_responder(channel_name)
attachment_urls = []
if message.attachments:
attachment_urls = await self._ingest_attachments(message, channel_name, airesponder)
if len(message_content) < 1:
# image-only posts: cached + observed, no reply (IMG-17)
if attachment_urls:
await airesponder.observe_event(message.author.name, "image", f"posted {len(attachment_urls)} image(s)")
return
message_content = self._resolve_mentions(message_content)
channel_name = self.get_channel_name(message.channel)
msg = AIMessage(
message.author.name, message_content, channel_name, self.user in message.mentions or isinstance(message.channel, DMChannel)
)
airesponder = self.get_ai_responder(channel_name)
if message.attachments:
for attachment in message.attachments:
if not msg.urls:
msg.urls = []
if airesponder.image_cache is not None:
# cache-first: CDN URLs never travel further (IMG-10/11)
sha = await airesponder.image_cache.ingest_url(attachment.url, channel_name, message.author.name, str(message.id))
if sha is not None:
recent = airesponder.image_cache.recent(channel_name, 8)
ext = next((row["ext"] for row in recent if row["sha256"] == sha), "png")
data_url = airesponder.image_cache.data_url(sha, ext)
if data_url:
msg.urls.append(data_url)
else:
msg.urls.append(attachment.url)
if attachment_urls:
msg.urls = attachment_urls
# Reply/ignore classifier gate — direct messages bypass (BEH-01/02/03/07)
handled, factual = await self._classifier_gate(message, msg, airesponder, channel_name)
@@ -467,6 +499,14 @@ class FjerkroaBot(commands.Bot):
return True, False
return False, bool(verdict.get("factual", False))
async def _note_api_error(self, err: Exception) -> None:
"""Count consecutive failures; alert staff once at threshold (OPS-16)."""
self._consecutive_api_errors += 1
logging.warning(f"responder call failed ({self._consecutive_api_errors} in a row): {repr(err)}")
threshold = int(self.config.get("api-error-alert-threshold", 5))
if self._consecutive_api_errors == threshold:
await self.send_staff_alert(f"⚠️ {threshold} consecutive API errors — the bot may be down. Last: {str(err)[:200]}")
async def send_message_with_typing(self, airesponder, channel, message):
"""Send the user message to the AI responder with typing animation in discord"""
async with channel.typing():
@@ -572,8 +612,15 @@ class FjerkroaBot(commands.Bot):
# Get the AI responder based on the channel name
airesponder = self.get_ai_responder(channel_name)
# Send the user message to the AI responder, with typing indicators
response = await self.send_message_with_typing(airesponder, channel, message)
# Send the user message to the AI responder, with typing indicators.
# A raised call = a broken API path (cf. the gpt-5.6 tools incident):
# count it, alert staff at threshold, never crash the handler (OPS-16).
try:
response = await self.send_message_with_typing(airesponder, channel, message)
except Exception as err:
await self._note_api_error(err)
return
self._consecutive_api_errors = 0
# SAF/OPS gates between model proposal and delivery
await self._apply_response_gates(message, response)
+18
View File
@@ -0,0 +1,18 @@
"""Bounded HTTP body read (leaf module, no intra-package imports).
`response.content.read(n)` returns whatever is buffered, not n bytes,
so it silently truncates large or chunked bodies (and web feeds/pages
parse to garbage). This accumulates decompressed chunks up to a hard
cap instead.
"""
CHUNK = 65536
async def read_capped(response, max_bytes: int) -> bytes:
buf = bytearray()
async for chunk in response.content.iter_chunked(CHUNK):
buf.extend(chunk)
if len(buf) > max_bytes:
break
return bytes(buf[:max_bytes])
+50 -11
View File
@@ -1,30 +1,67 @@
import logging
import time
from functools import cache
from typing import Any, Dict, List, Optional
import requests
TWITCH_OAUTH_URL = "https://id.twitch.tv/oauth2/token"
# Refresh this long before Twitch expires the token (app tokens live ~60 days)
TOKEN_REFRESH_MARGIN = 86400
class IGDBQuery(object):
def __init__(self, client_id, igdb_api_key):
def __init__(self, client_id, igdb_api_key=None, client_secret=None):
self.client_id = client_id
self.igdb_api_key = igdb_api_key
self.client_secret = client_secret
# Unknown for statically configured tokens; set after each refresh
self._token_expires_at = None
def _refresh_token(self):
response = requests.post(
TWITCH_OAUTH_URL,
params={"client_id": self.client_id, "client_secret": self.client_secret, "grant_type": "client_credentials"},
)
response.raise_for_status()
data = response.json()
self.igdb_api_key = data["access_token"]
self._token_expires_at = time.time() + data.get("expires_in", 0) - TOKEN_REFRESH_MARGIN
logging.info("IGDB: refreshed Twitch app access token")
def _ensure_token(self):
if not self.client_secret:
return
if not self.igdb_api_key or (self._token_expires_at is not None and time.time() >= self._token_expires_at):
self._refresh_token()
def send_igdb_request(self, endpoint, query_body):
igdb_url = f"https://api.igdb.com/v4/{endpoint}"
headers = {"Client-ID": self.client_id, "Authorization": f"Bearer {self.igdb_api_key}"}
try:
response = requests.post(igdb_url, headers=headers, data=query_body)
self._ensure_token()
response = self._post_igdb(igdb_url, query_body)
if self.client_secret and response.status_code == 401:
# Token expired server-side (e.g. statically configured) — refresh and retry once
self._refresh_token()
response = self._post_igdb(igdb_url, query_body)
response.raise_for_status()
return response.json()
except requests.RequestException as e:
print(f"Error during IGDB API request: {e}")
return None
def _post_igdb(self, igdb_url, query_body):
headers = {"Client-ID": self.client_id, "Authorization": f"Bearer {self.igdb_api_key}"}
return requests.post(igdb_url, headers=headers, data=query_body)
@staticmethod
def build_query(fields, filters=None, limit=10, offset=None):
query = f"fields {','.join(fields) if fields is not None and len(fields) > 0 else '*'}; limit {limit};"
def build_query(fields, filters=None, limit=10, offset=None, search_term=None):
query = ""
if search_term:
escaped = search_term.replace("\\", "\\\\").replace('"', '\\"')
query += f'search "{escaped}"; '
query += f"fields {','.join(fields) if fields is not None and len(fields) > 0 else '*'}; limit {limit};"
if offset is not None:
query += f" offset {offset};"
if filters:
@@ -32,12 +69,12 @@ class IGDBQuery(object):
query += " where " + " & ".join(filter_statements) + ";"
return query
def generalized_igdb_query(self, params, endpoint, fields, additional_filters=None, limit=10, offset=None):
def generalized_igdb_query(self, params, endpoint, fields, additional_filters=None, limit=10, offset=None, search_term=None):
all_filters = {key: f'~ "{value}"*' for key, value in params.items() if value}
if additional_filters:
all_filters.update(additional_filters)
query = self.build_query(fields, all_filters, limit, offset)
query = self.build_query(fields, all_filters, limit, offset, search_term)
data = self.send_igdb_request(endpoint, query)
print(f"{endpoint}: {query} -> {data}")
return data
@@ -79,7 +116,7 @@ class IGDBQuery(object):
"id",
"name",
"alternative_names",
"category",
"game_type",
"release_dates",
"franchise",
"language_supports",
@@ -101,9 +138,10 @@ class IGDBQuery(object):
return None
try:
# Search for games with fuzzy matching
# IGDB native full-text search: diacritic- and word-order-insensitive,
# unlike a `name ~ "..."*` prefix filter
games = self.generalized_igdb_query(
{"name": query.strip()},
{},
"games",
[
"id",
@@ -120,8 +158,9 @@ class IGDBQuery(object):
"themes.name",
"cover.url",
],
additional_filters={"category": "= 0"}, # Main games only
additional_filters={"game_type": "= 0"}, # Main games only (IGDB renamed category -> game_type)
limit=limit,
search_term=query.strip(),
)
if not games:
+4 -1
View File
@@ -14,6 +14,7 @@ from typing import Any, Callable, Dict, List, Optional
import aiohttp
from .httpread import read_capped
from .persistence import PersistentStore
DEFAULT_CACHE_MB = 500
@@ -79,7 +80,9 @@ class ImageCache:
async with aiohttp.ClientSession(timeout=timeout) as session:
async with session.get(url) as response:
response.raise_for_status()
return await response.content.read(limit + 1)
# limit + 1: an over-limit body must stay over-limit so
# ingest_bytes rejects it instead of caching it truncated
return await read_capped(response, limit + 1)
def data_url(self, sha256: str, ext: str) -> Optional[str]:
path = self._path(sha256, ext)
+153
View File
@@ -0,0 +1,153 @@
"""News digest fetcher (SPEC-013, FDB-012 news rewrite).
Replaces the broken pre-1.0-openai `news_feed.py`. Fetches configured
RSS/Atom feeds (stdlib, no feedparser dep), builds a compact sanitized
headline digest, and writes it to the `{news}` file the responder
injects (AIResponder.message). Feeds are external input: titles are
sanitized (SAF-03) and each feed URL is SSRF-guarded before fetching.
CLI: python -m fjerkroa_bot.news --config kroa.toml
"""
import argparse
import logging
import sys
import time
from typing import Any, Dict, List, Optional, Tuple
import defusedxml.ElementTree as ElementTree # hardened XML: feeds are untrusted (XXE/billion-laughs)
from .ai_responder import sanitize_external_text
DEFAULT_PER_FEED = 3
DEFAULT_MAX_ITEMS = 15
FETCH_TIMEOUT_S = 15
_ATOM = "{http://www.w3.org/2005/Atom}"
def parse_feed(data: bytes, source: str = "") -> List[Dict[str, str]]:
"""Parse RSS or Atom bytes into [{title, link, source}] (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/>
for item in root.iter("item"):
title = (item.findtext("title") or "").strip()
link = (item.findtext("link") or "").strip()
if title:
items.append({"title": title, "link": link, "source": source})
# Atom: <feed><entry><title/><link href=/>
for entry in root.iter(f"{_ATOM}entry"):
title = (entry.findtext(f"{_ATOM}title") or "").strip()
link_el = entry.find(f"{_ATOM}link")
link = link_el.get("href", "") if link_el is not None else ""
if title:
items.append({"title": title, "link": link, "source": source})
return items
def render_digest(items: List[Dict[str, str]], max_items: int = DEFAULT_MAX_ITEMS) -> str:
"""Compact sanitized digest for the {news} prompt slot."""
lines = []
for item in items[:max_items]:
title = sanitize_external_text(item["title"], 200)
source = item.get("source", "")
link = item.get("link", "")
prefix = f"[{source}] " if source else ""
lines.append(f"- {prefix}{title}" + (f" ({link})" if link else ""))
return "\n".join(lines)
class NewsFetcher:
def __init__(self, guard, fetch_bytes) -> None:
# injected so tests need no network; production wires aiohttp + guard_url
self._guard = guard
self._fetch_bytes = fetch_bytes
async def collect(self, feeds: List[Tuple[str, str]], per_feed: int) -> List[Dict[str, str]]:
"""feeds = [(url, label)]; returns deduped items, order preserved."""
seen = set()
out: List[Dict[str, str]] = []
for url, label in feeds:
reason = self._guard(url)
if reason:
logging.warning(f"news: skipping feed {label}{reason}")
continue
try:
data = await self._fetch_bytes(url)
except Exception as err:
logging.warning(f"news: fetch failed for {label}: {repr(err)}")
continue
for item in parse_feed(data, label)[:per_feed]:
key = item["title"]
if key not in seen:
seen.add(key)
out.append(item)
return out
def _feeds_from_config(config: Dict[str, Any]) -> List[Tuple[str, str]]:
"""news-feeds = [["url", "label"], ...] or ["url", ...]."""
feeds = []
for entry in config.get("news-feeds", []):
if isinstance(entry, (list, tuple)):
feeds.append((str(entry[0]), str(entry[1]) if len(entry) > 1 else ""))
else:
feeds.append((str(entry), ""))
return feeds
async def _aiohttp_fetch(url: str) -> bytes:
import aiohttp
from .httpread import read_capped
timeout = aiohttp.ClientTimeout(total=FETCH_TIMEOUT_S)
async with aiohttp.ClientSession(timeout=timeout, headers={"User-Agent": "Mozilla/5.0 (compatible; FjerkroaBot-news/1.0)"}) as session:
async with session.get(url) as response:
response.raise_for_status()
return await read_capped(response, 4 * 1024 * 1024)
async def run(config: Dict[str, Any]) -> Optional[str]:
from .url_reader import guard_url
out_path = config.get("news")
if not out_path:
logging.error("news: no `news` output path in config")
return None
feeds = _feeds_from_config(config)
if not feeds:
logging.error("news: no `news-feeds` configured")
return None
fetcher = NewsFetcher(guard_url, _aiohttp_fetch)
items = await fetcher.collect(feeds, int(config.get("news-per-feed", DEFAULT_PER_FEED)))
digest = render_digest(items, int(config.get("news-max-items", DEFAULT_MAX_ITEMS)))
header = f"News as of {time.strftime('%Y-%m-%d %H:%M UTC', time.gmtime())}:\n"
with open(out_path, "w", encoding="utf-8") as fd:
fd.write(header + digest + "\n")
logging.info(f"news: wrote {len(items)} items to {out_path}")
return out_path
def main() -> int:
import asyncio
import tomlkit
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
parser = argparse.ArgumentParser(description="Fetch RSS/Atom feeds into the {news} digest file")
parser.add_argument("--config", required=True)
args = parser.parse_args()
with open(args.config, encoding="utf-8") as fd:
config = tomlkit.load(fd)
result = asyncio.run(run(config))
return 0 if result else 1
if __name__ == "__main__":
sys.exit(main())
+98 -31
View File
@@ -12,6 +12,7 @@ from .ai_responder import AIResponder, exponential_backoff, sanitize_external_te
from .igdblib import IGDBQuery
from .leonardo_draw import LeonardoAIDrawMixIn
from .quota import QuotaLedger
from .url_reader import FETCH_URL_TOOL, URLReader
# The response envelope, enforced server-side via structured outputs
# (ENV-19). All fields required, closed object, nullable where the
@@ -58,6 +59,31 @@ CONSOLIDATION_RESPONSE_FORMAT = {
"type": "json_schema",
"json_schema": {"name": "consolidation", "strict": True, "schema": CONSOLIDATION_SCHEMA},
}
# Follow-up task proposal (SPEC-005 TSK-08): one task or null
TASKGEN_SCHEMA = {
"type": "object",
"properties": {
"task": {
"type": ["object", "null"],
"properties": {
"channel": {"type": ["string", "null"], "description": "Target channel, or null for the main chat channel."},
"prompt": {"type": "string", "description": "Instruction the assistant will act on when the task runs."},
"due_hours": {"type": "number", "description": "Hours from now until the task should run (0 = now)."},
},
"required": ["channel", "prompt", "due_hours"],
"additionalProperties": False,
}
},
"required": ["task"],
"additionalProperties": False,
}
TASKGEN_RESPONSE_FORMAT = {"type": "json_schema", "json_schema": {"name": "task_proposal", "strict": True, "schema": TASKGEN_SCHEMA}}
TASKGEN_SYSTEM = (
"You plan the self-initiated actions of a Discord assistant. Given recent conversation summaries, propose AT MOST ONE follow-up"
" worth doing on the assistant's own initiative (ask how something announced went, revisit an open question, congratulate on an"
" event). Only propose something genuinely worthwhile — when in doubt, return a null task."
)
# Reply/ignore + factual pre-pass (SPEC-010 BEH-01): one cheap call
CLASSIFIER_SCHEMA = {
"type": "object",
@@ -111,16 +137,20 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
# Initialize IGDB if enabled
self.igdb = None
igdb_client_id = self.config.get("igdb-client-id")
igdb_client_secret = self.config.get("igdb-client-secret")
igdb_access_token = self.config.get("igdb-access-token")
logging.info("IGDB Configuration Check:")
logging.info(f" enable-game-info: {self.config.get('enable-game-info', 'NOT SET')}")
logging.info(f" igdb-client-id: {'SET' if self.config.get('igdb-client-id') else 'NOT SET'}")
logging.info(f" igdb-access-token: {'SET' if self.config.get('igdb-access-token') else 'NOT SET'}")
logging.info(f" igdb-client-id: {'SET' if igdb_client_id else 'NOT SET'}")
logging.info(f" igdb-client-secret: {'SET' if igdb_client_secret else 'NOT SET'}")
logging.info(f" igdb-access-token: {'SET' if igdb_access_token else 'NOT SET'}")
if self.config.get("enable-game-info", False) and self.config.get("igdb-client-id") and self.config.get("igdb-access-token"):
if self.config.get("enable-game-info", False) and igdb_client_id and (igdb_client_secret or igdb_access_token):
try:
self.igdb = IGDBQuery(self.config["igdb-client-id"], self.config["igdb-access-token"])
self.igdb = IGDBQuery(igdb_client_id, igdb_access_token, client_secret=igdb_client_secret)
logging.info("✅ IGDB integration SUCCESSFULLY enabled for game information")
logging.info(f" Client ID: {self.config['igdb-client-id'][:8]}...")
logging.info(f" Client ID: {igdb_client_id[:8]}...")
logging.info(f" Available functions: {len(self.igdb.get_openai_functions())}")
except Exception as e:
logging.error(f"❌ Failed to initialize IGDB: {e}")
@@ -128,6 +158,33 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
else:
logging.warning("❌ IGDB integration DISABLED - missing configuration or disabled in config")
# URL reading tool (SPEC-011); shares the image cache for page images
self.url_reader = URLReader(lambda: self.config, self.image_cache)
def _available_tools(self) -> List[Dict[str, Any]]:
"""Assemble the function-tool list from every enabled provider (URL-01)."""
functions: List[Dict[str, Any]] = []
if self.igdb and self.config.get("enable-game-info", False):
try:
igdb_functions = self.igdb.get_openai_functions()
if isinstance(igdb_functions, list):
functions.extend(igdb_functions)
except (TypeError, AttributeError) as err:
logging.warning(f"Error setting up IGDB functions: {err}")
if self.url_reader.enabled():
functions.append(FETCH_URL_TOOL)
return functions
async def _dispatch_tool(self, name: str, args: Dict[str, Any], author: str) -> Any:
"""Route a tool call to its provider (IGDB or URL reader)."""
if name == "fetch_url":
per_user_cap = int(self.config.get("url-daily-per-user", 20))
if self.ledger._get(f"url-fetch:{author}") >= per_user_cap: # URL-07
return {"error": "daily URL fetch limit reached"}
self.ledger._add(f"url-fetch:{author}", 1)
return await self.url_reader.fetch(str(args.get("url", "")), self.channel, author or "user")
return await self._execute_igdb_function(name, args)
async def draw_openai(self, description: str, count: int = 1) -> List[BytesIO]:
if not self.ledger.budget_ok():
raise RuntimeError("daily budget exhausted - refusing image call")
@@ -216,24 +273,13 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
# hashed, never the raw Discord name (SAF-10)
chat_kwargs["safety_identifier"] = "discord-" + hashlib.sha256(author.encode()).hexdigest()[:16]
if self.igdb and self.config.get("enable-game-info", False):
try:
igdb_functions = self.igdb.get_openai_functions()
if igdb_functions and isinstance(igdb_functions, list):
chat_kwargs["tools"] = [{"type": "function", "function": func} for func in igdb_functions]
chat_kwargs["tool_choice"] = "auto"
# gpt-5.6 rejects tools + reasoning on chat/completions (ENV-21)
chat_kwargs["reasoning_effort"] = self.config.get("reasoning-effort", "none")
logging.info(f"🎮 IGDB functions available to AI: {[f['name'] for f in igdb_functions]}")
logging.debug(f" Full chat_kwargs with tools: {list(chat_kwargs.keys())}")
except (TypeError, AttributeError) as e:
logging.warning(f"Error setting up IGDB functions: {e}")
else:
logging.debug(
"🎮 IGDB not available for this request (igdb={}, enabled={})".format(
self.igdb is not None, self.config.get("enable-game-info", False)
)
)
available_tools = self._available_tools()
if available_tools:
chat_kwargs["tools"] = [{"type": "function", "function": func} for func in available_tools]
chat_kwargs["tool_choice"] = "auto"
# gpt-5.6 rejects tools + reasoning on chat/completions (ENV-21)
chat_kwargs["reasoning_effort"] = self.config.get("reasoning-effort", "none")
logging.info(f"🔧 Tools available to AI: {[func['name'] for func in available_tools]}")
result = await openai_chat(self.client, **chat_kwargs)
self._record_usage(result)
@@ -253,10 +299,8 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
tool_names = [tc.function.name for tc in message.tool_calls]
logging.info(f"🔧 OpenAI requested function calls: {tool_names}")
# Check if we have function/tool calls and IGDB is enabled
has_tool_calls = (
hasattr(message, "tool_calls") and message.tool_calls and self.igdb and self.config.get("enable-game-info", False)
)
# Any offered tool may have been called (IGDB or fetch_url)
has_tool_calls = bool(hasattr(message, "tool_calls") and message.tool_calls and available_tools)
# Clean up any existing tool messages in the history to avoid conflicts
if has_tool_calls:
@@ -279,12 +323,12 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
function_name = tool_call.function.name
function_args = json.loads(tool_call.function.arguments)
logging.info(f"🎮 Executing IGDB function: {function_name} with args: {function_args}")
logging.info(f"🔧 Executing tool: {function_name} with args: {function_args}")
# Execute IGDB function
function_result = await self._execute_igdb_function(function_name, function_args)
# Route to the right provider (IGDB or URL reader)
function_result = await self._dispatch_tool(function_name, function_args, self._last_author(messages) or "")
logging.info(f"🎮 IGDB function result: {type(function_result)} - {str(function_result)[:200]}...")
logging.info(f"🔧 Tool result: {type(function_result)} - {str(function_result)[:200]}...")
messages.append(
{
@@ -381,6 +425,29 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
logging.info(f"edited {len(buffers)} image(s) on {model} from {len(handles)} input(s)")
return buffers
async def propose_task(self) -> Optional[Dict[str, Any]]:
"""One follow-up proposal from recent episodes on memory-model (TSK-08)."""
if "memory-model" not in self.config or self.store is None or not self.ledger.budget_ok():
return None
channel = self.config.get("chat-channel", "chat")
episodes = await asyncio.to_thread(self.store.recent_episodes, channel, 5)
if not episodes:
return None
episode_lines = "\n".join(f"- {episode}" for episode in episodes)
messages = [
{"role": "system", "content": TASKGEN_SYSTEM},
{"role": "user", "content": f"Recent conversation summaries in #{channel}:\n{episode_lines}"},
]
try:
result = await openai_chat(
self.client, model=self.config["memory-model"], messages=messages, response_format=TASKGEN_RESPONSE_FORMAT
)
self._record_usage(result)
return json.loads(result.choices[0].message.content)
except Exception as err:
logging.warning(f"task proposal failed: {repr(err)}")
return None
async def classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
"""~100-token reply/factual/emoji verdict on classifier-model (BEH-01/03)."""
if "classifier-model" not in self.config or not self.ledger.budget_ok():
+40 -1
View File
@@ -13,7 +13,7 @@ from contextlib import closing
from pathlib import Path
from typing import Any, Dict, List, Optional
SCHEMA_VERSION = 4
SCHEMA_VERSION = 5
class PersistentStore:
@@ -66,6 +66,12 @@ class PersistentStore:
" user TEXT NOT NULL, message_id TEXT, ext TEXT NOT NULL, bytes INTEGER NOT NULL,"
" created_at TEXT NOT NULL DEFAULT (datetime('now')))"
)
if version < 5:
conn.execute(
"CREATE TABLE IF NOT EXISTS tasks (id INTEGER PRIMARY KEY, kind TEXT NOT NULL, channel TEXT NOT NULL,"
" due_at TEXT NOT NULL, payload TEXT NOT NULL, state TEXT NOT NULL DEFAULT 'queued',"
" created_at TEXT NOT NULL DEFAULT (datetime('now')), executed_at TEXT)"
)
if version < SCHEMA_VERSION:
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
os.chmod(self.db_path, 0o600) # conversation data (PER-04)
@@ -195,6 +201,39 @@ class PersistentStore:
)
return cursor.rowcount
# --- task queue (SPEC-005, FDB-011) ---
def task_add(self, kind: str, channel: str, due_at: str, payload: str, state: str = "queued") -> int:
with closing(self._connect()) as conn, conn:
cursor = conn.execute(
"INSERT INTO tasks (kind, channel, due_at, payload, state) VALUES (?, ?, ?, ?, ?)",
(kind, channel, due_at, payload, state),
)
return int(cursor.lastrowid or 0)
def tasks_due(self) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute(
"SELECT id, kind, channel, payload FROM tasks WHERE state = 'queued' AND due_at <= datetime('now') ORDER BY due_at"
).fetchall()
return [{"id": row[0], "kind": row[1], "channel": row[2], "payload": row[3]} for row in rows]
def tasks_open(self) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute(
"SELECT id, kind, channel, due_at, state FROM tasks WHERE state IN ('queued', 'approval') ORDER BY id"
).fetchall()
return [{"id": row[0], "kind": row[1], "channel": row[2], "due_at": row[3], "state": row[4]} for row in rows]
def task_set_state(self, task_id: int, state: str) -> int:
with closing(self._connect()) as conn, conn:
return conn.execute("UPDATE tasks SET state = ?, executed_at = datetime('now') WHERE id = ?", (state, task_id)).rowcount
def tasks_pending_of_kind(self, kind: str) -> int:
with closing(self._connect()) as conn:
row = conn.execute("SELECT COUNT(*) FROM tasks WHERE kind = ? AND state IN ('queued', 'approval')", (kind,)).fetchone()
return int(row[0])
# --- image cache index (SPEC-004, FDB-010) ---
def image_add(self, sha256: str, channel: str, user: str, message_id: Optional[str], ext: str, nbytes: int) -> None:
+135
View File
@@ -0,0 +1,135 @@
"""Self-tasking engine (SPEC-005, FDB-011).
Generators propose, the scheduler executes through the injected
execute callback (= the normal responder path), so budget, gates and
kill-switches all apply. Everything is off unless `tasks-enabled` is
true (TSK-03).
"""
import asyncio
import logging
import time
from typing import Any, Awaitable, Callable, Dict, List, Optional
from .persistence import PersistentStore
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."
ExecuteCallback = Callable[[str, str], Awaitable[None]]
ProposeCallback = Callable[[], Awaitable[Optional[Dict[str, Any]]]]
AlertCallback = Callable[[str], Awaitable[None]]
class TaskEngine:
def __init__(
self,
store: Optional[PersistentStore],
ledger: QuotaLedger,
config_getter: Callable[[], Dict[str, Any]],
execute: ExecuteCallback,
propose: ProposeCallback,
staff_alert: AlertCallback,
allowed: Callable[[], bool],
idle_seconds: Callable[[], float],
observe: Callable[[str, str, str], Awaitable[None]],
) -> None:
self.store = store
self.ledger = ledger
self._config = config_getter
self._execute = execute
self._propose = propose
self._staff_alert = staff_alert
self._allowed = allowed
self._idle_seconds = idle_seconds
self._observe = observe
self._last_generation = 0.0
self._lock = asyncio.Lock()
def active(self) -> bool:
return self.store is not None and bool(self._config().get("tasks-enabled", False)) # TSK-03
def generators(self) -> List[str]:
return list(self._config().get("tasks-generators", ["idle-impulse", "follow-up"]))
async def enqueue(self, kind: str, channel: str, payload: str, due_hours: float = 0.0) -> int:
assert self.store is not None
state = "approval" if self._config().get("tasks-approval", False) else "queued" # TSK-05
due_at = time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime(time.time() + due_hours * 3600.0))
task_id = int(await asyncio.to_thread(self.store.task_add, kind, channel, due_at, payload, state) or 0)
if state == "approval":
await self._staff_alert(f"Task #{task_id} proposed ({kind}, #{channel}): {payload[:180]} — !bot task-approve {task_id}")
return task_id
def _channel_cap_ok(self, channel: str) -> bool:
cap = int(self._config().get("tasks-max-per-channel-per-day", DEFAULT_MAX_PER_CHANNEL_PER_DAY))
return self.ledger._get(f"task-runs:{channel}") < cap # TSK-04
async def tick(self) -> None:
"""One scheduler pass: execute due tasks, then maybe generate (TSK-02/06)."""
if not self.active() or not self._allowed() or self._lock.locked():
return
assert self.store is not None
async with self._lock:
await self._maybe_generate()
for task in await asyncio.to_thread(self.store.tasks_due):
if not self._channel_cap_ok(task["channel"]):
continue # stays queued for tomorrow (TSK-04)
try:
await self._execute(task["channel"], task["payload"])
await asyncio.to_thread(self.store.task_set_state, task["id"], "done")
self.ledger._add(f"task-runs:{task['channel']}", 1)
await self._observe("system", "task", f"executed {task['kind']}: {task['payload'][:200]}")
except Exception as err:
logging.warning(f"task {task['id']} failed: {repr(err)}")
await asyncio.to_thread(self.store.task_set_state, task["id"], "failed")
async def _maybe_generate(self) -> None:
config = self._config()
interval = float(config.get("taskgen-interval-hours", DEFAULT_TASKGEN_INTERVAL_HOURS)) * 3600.0
now = time.monotonic()
if self._last_generation and now - self._last_generation < interval:
return
self._last_generation = now
generators = self.generators()
if "idle-impulse" in generators:
await self._generate_idle_impulse()
if "follow-up" in generators:
await self._generate_follow_up()
async def _generate_idle_impulse(self) -> None:
"""Boreness, demoted to a deterministic generator (TSK-07)."""
assert self.store is not None
config = self._config()
channel = config.get("chat-channel")
if not channel:
return
idle_threshold = float(config.get("idle-impulse-hours", DEFAULT_IDLE_IMPULSE_HOURS)) * 3600.0
if self._idle_seconds() < idle_threshold:
return
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "idle-impulse"):
return # one pending impulse is enough
prompt = config.get("boreness-prompt", DEFAULT_BORENESS_PROMPT)
await self.enqueue("idle-impulse", channel, prompt)
async def _generate_follow_up(self) -> None:
"""Ask the model for one follow-up worth doing (TSK-08)."""
assert self.store is not None
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "follow-up"):
return
proposal = await self._propose()
if not proposal or not isinstance(proposal.get("task"), dict):
return
task = proposal["task"]
channel = str(task.get("channel") or self._config().get("chat-channel") or "")
prompt = str(task.get("prompt") or "").strip()
if not channel or not prompt:
return
try:
due_hours = max(0.0, float(task.get("due_hours") or 0.0))
except (TypeError, ValueError):
due_hours = 0.0
await self.enqueue("follow-up", channel, prompt, due_hours)
+164
View File
@@ -0,0 +1,164 @@
"""URL reading tool (SPEC-011, FDB-018).
The model calls `fetch_url`; this module fetches safely and returns
readable text plus prominent image URLs. Web pages are hostile input:
every fetch is SSRF-guarded (no private/loopback/link-local targets,
http/https only, redirects re-validated) and every byte of text is
sanitized before it can reach the prompt.
"""
import ipaddress
import logging
import re
import socket
from html.parser import HTMLParser
from typing import Any, Callable, Dict, List, Optional, Tuple
from urllib.parse import urljoin, urlparse
import aiohttp
from .ai_responder import sanitize_external_text
from .httpread import read_capped
DEFAULT_MAX_BYTES = 2 * 1024 * 1024
DEFAULT_MAX_CHARS = 6000
DEFAULT_MAX_IMAGES = 2
FETCH_TIMEOUT_S = 15
MAX_REDIRECTS = 5
FETCH_URL_TOOL = {
"name": "fetch_url",
"description": "Fetch a public web page and return its readable text plus prominent image links. "
"Use when the user shares a URL and asks about it, or to get details behind a news link.",
"parameters": {
"type": "object",
"properties": {"url": {"type": "string", "description": "The http/https URL to read."}},
"required": ["url"],
},
}
class _Extractor(HTMLParser):
def __init__(self) -> None:
super().__init__()
self._skip = 0
self.parts: List[str] = []
self.images: List[str] = []
self.og_image: Optional[str] = None
def handle_starttag(self, tag: str, attrs) -> None:
if tag in ("script", "style", "noscript", "svg"):
self._skip += 1
attr = dict(attrs)
src = attr.get("src")
if tag == "img" and src:
self.images.append(src)
if tag == "meta" and attr.get("property") == "og:image" and attr.get("content"):
self.og_image = attr["content"]
def handle_endtag(self, tag: str) -> None:
if tag in ("script", "style", "noscript", "svg") and self._skip > 0:
self._skip -= 1
def handle_data(self, data: str) -> None:
if self._skip == 0 and data.strip():
self.parts.append(data.strip())
def _ip_is_public(ip_str: str) -> bool:
try:
ip = ipaddress.ip_address(ip_str)
except ValueError:
return False
return not (ip.is_private or ip.is_loopback or ip.is_link_local or ip.is_multicast or ip.is_reserved or ip.is_unspecified)
def guard_url(url: str) -> Optional[str]:
"""Return None if safe to fetch, else a human-readable refusal reason (URL-02/03)."""
parsed = urlparse(url)
if parsed.scheme not in ("http", "https"):
return f"refused scheme {parsed.scheme!r} (only http/https)"
host = parsed.hostname
if not host:
return "refused: no host"
try:
literal = ipaddress.ip_address(host)
return None if _ip_is_public(str(literal)) else f"refused non-public address {host}"
except ValueError:
pass
try:
infos = socket.getaddrinfo(host, None)
except socket.gaierror:
return f"refused: cannot resolve {host}"
for info in infos:
if not _ip_is_public(str(info[4][0])):
return f"refused: {host} resolves to non-public address"
return None
class URLReader:
def __init__(self, config_getter: Callable[[], Dict[str, Any]], image_cache) -> None:
self._config = config_getter
self.image_cache = image_cache
def enabled(self) -> bool:
return bool(self._config().get("enable-url-reading", False))
async def _get(self, session, url: str, max_bytes: int) -> Tuple[str, bytes]:
"""Manual redirect handling so every hop is re-guarded (URL-04)."""
current = url
for _ in range(MAX_REDIRECTS):
reason = guard_url(current)
if reason:
raise ValueError(reason)
async with session.get(current, allow_redirects=False) as response:
if response.status in (301, 302, 303, 307, 308) and response.headers.get("Location"):
current = urljoin(current, response.headers["Location"])
continue
response.raise_for_status()
return str(response.url), await read_capped(response, max_bytes)
raise ValueError("too many redirects")
async def fetch(self, url: str, channel: str, user: str) -> Dict[str, Any]:
config = self._config()
max_bytes = int(config.get("url-max-bytes", DEFAULT_MAX_BYTES))
timeout = aiohttp.ClientTimeout(total=FETCH_TIMEOUT_S)
try:
async with aiohttp.ClientSession(timeout=timeout, headers={"User-Agent": "FjerkroaBot/1.0"}) as session:
final_url, body = await self._get(session, url, max_bytes)
except Exception as err:
return {"error": str(err)}
text = self._to_text(body.decode("utf-8", "ignore"))
clean = sanitize_external_text(text, int(config.get("url-max-chars", DEFAULT_MAX_CHARS)))
images = await self._ingest_images(body.decode("utf-8", "ignore"), final_url, channel, user)
return {"url": final_url, "text": clean, "images_cached": images}
def _to_text(self, html: str) -> str:
extractor = _Extractor()
try:
extractor.feed(html)
except Exception as err:
logging.debug(f"html parse (text) failed: {err!r}")
return re.sub(r"\s+\n", "\n", " ".join(extractor.parts))
async def _ingest_images(self, html: str, base_url: str, channel: str, user: str) -> int:
if self.image_cache is None:
return 0
extractor = _Extractor()
try:
extractor.feed(html)
except Exception as err:
logging.debug(f"html parse (images) failed: {err!r}")
candidates = ([extractor.og_image] if extractor.og_image else []) + extractor.images
limit = int(self._config().get("url-max-images", DEFAULT_MAX_IMAGES))
cached = 0
for src in candidates:
if cached >= limit:
break
absolute = urljoin(base_url, src)
if guard_url(absolute) is not None:
continue
sha = await self.image_cache.ingest_url(absolute, channel, user, None)
if sha is not None:
cached += 1
return cached
+3 -2
View File
@@ -6,9 +6,10 @@ with date + result.
| ID | Date | Result |
| --- | --- | --- |
| DEP-01 | 2026-07-13 | Verified with the v3.0.0 ggg deploy: tag-only refusal + untracked config/state survived. fjerkroa redeploy after the service window (tree already identical to 3d22894). |
| DEP-02 | 2026-07-13 | Service map exercised: luma restart via script (v3.0.0); kroa mapping code-reviewed, exercised on its next deploy. |
| DEP-01 | 2026-07-13 | Verified on both hosts (ggg v3.0.0..v3.3.2, fjerkroa v3.3.2): tag-only refusal + untracked config/state survived every deploy. |
| DEP-02 | 2026-07-13 | Service map exercised on both hosts: luma (v3.0.0..v3.3.2) and kroa (v3.3.2, DEPLOY_FORCE per operator order). |
| DEP-03 | 2026-07-13 | Exercised with the v3.1.0 ggg deploy: bot.db.pre-v3.1.0 confirmed on the host. (v3.0.0 note: no pre-existing db in the pickle era.) |
| DEP-04 | 2026-07-13 | Smoke gate exercised on ggg: RUNNING + fresh login line. |
| DEP-05 | 2026-07-13 | Live-verified: kroa deploy attempt ~15h Oslo refused without DEPLOY_FORCE=1. |
| DEP-06 | 2026-07-13 | Rollback documented (older tag + db backup restore); live drill pending — next release. |
| OPS-15 | 2026-07-13 | Backup cron installed on both hosts (daily 03:17 UTC → ~/backups/<bot>/, keep 14); first snapshots written + verified 0600 (kroa 10965 B, luma 25871 B). |
+1
View File
@@ -15,6 +15,7 @@ dependencies = [
"tomlkit>=0.13",
"watchdog>=6",
"requests>=2.32",
"defusedxml>=0.7",
]
[project.scripts]
+8
View File
@@ -82,6 +82,14 @@ SAF-08/MEM-09).
Bot-generated images are ingested like uploads (user `assistant`), so
"make a variant of that" remix chains work on the bot's own output.
### IMG-17 — Image-only messages are cached (coverage: test)
A message consisting only of attachments (no text) is ingested into
the cache and recorded as an observation, even though no reply is
produced — the image must be available for later `picture_edit` and
vision follow-ups. (Previously the empty-text early-return dropped
such posts entirely.)
### IMG-16 — The prompt announces editable images (coverage: test)
When the answer channel has cached images, the context suffix states
+59
View File
@@ -0,0 +1,59 @@
# SPEC-005 — Self-tasking
The sigmoid "boreness" loop becomes a persistent task queue (schema
v5): generators propose, a scheduler executes — through the normal
responder path, so every SAF gate, quota and kill-switch applies.
**Default off** (`tasks-enabled`, kitchen stays quiet unless opted
in); generators are individually selectable per persona
(`tasks-generators`). Marked experimental per plan v4.
### TSK-01 — Tasks are persistent queue rows (coverage: test)
A task is a store row: kind, channel, due-at, payload (the prompt the
responder will run), state (`queued`/`approval`/`done`/`cancelled`/
`failed`). Enqueued tasks survive restarts.
### TSK-02 — Due tasks run through the responder path (coverage: test)
The scheduler executes due queued tasks as system messages via the
normal respond flow (inheriting budget, gates, envelope), marks them
`done`/`failed`, and records the outcome as an observation so memory
learns what the bot did on its own.
### TSK-03 — Off by default (coverage: test)
Without `tasks-enabled = true` nothing is generated and nothing is
executed. A restaurant server does not improvise unless asked to.
### TSK-04 — Per-channel daily cap (coverage: test)
At most `tasks-max-per-channel-per-day` (default 2) task executions
per channel per day, counted in the ledger; further due tasks stay
queued for the next day.
### TSK-05 — Approval mode (coverage: test)
With `tasks-approval = true`, generated tasks enter state `approval`
and a staff alert announces them; `!bot task-approve <id>` moves them
to the queue, `!bot task-cancel <id>` kills them (works for queued
tasks too). Staged-rollout path from the review.
### TSK-06 — Kill-switches govern execution (coverage: test)
Task execution respects `bot_initiated_allowed()``!bot tasks off`,
pause, quiet mode and quiet hours all stop the scheduler tick.
### TSK-07 — Idle-impulse generator (boreness, demoted) (coverage: test)
When the chat channel has been idle longer than
`idle-impulse-hours` (default 12), the generator enqueues one
impulse task with the configured boreness prompt — deterministic
threshold instead of the old 7-second sigmoid dice loop, bounded by
TSK-04. The `on_boreness` loop is gone.
### TSK-08 — Follow-up generator proposes from memory (coverage: test)
Periodically (`taskgen-interval-hours`, default 6) the follow-up
generator asks `memory-model` (strict structured outputs) whether the
recent episodes/facts warrant one follow-up task (channel, prompt,
due-in-hours). A null/failed proposal enqueues nothing.
+6
View File
@@ -62,6 +62,12 @@ Bot-initiated posts also respect pause/quiet.
`!bot pins` answers with all pinned facts and their ids (global +
per-channel) — without it, `!bot unpin <id>` required guessing ids.
### OPS-12 — Task queue surface (coverage: test)
`!bot tasks` lists open tasks (queued + awaiting approval) with ids;
`!bot task-approve <id>` and `!bot task-cancel <id>` manage them
(TSK-05). `!bot tasks on|off` stays the kill-switch (OPS-09).
### OPS-10 — Spend report (coverage: test)
`!bot spend` answers in the staff channel with today's estimated
+57
View File
@@ -0,0 +1,57 @@
# SPEC-011 — URL reading
A `fetch_url` tool alongside IGDB: the model decides when to read a
link (user pastes a URL + question; a news item links an article
Luma wants details on). Web pages are the number-one injection
vector, so everything fetched is sanitized (SAF-03) and the fetch
itself is SSRF-guarded — the bot runs on shared hosting. Active only
when `enable-url-reading = true`.
### URL-01 — fetch_url is offered as a tool (coverage: test)
When `enable-url-reading` is true, the chat call's `tools` list
includes a `fetch_url` function (url string param) next to any IGDB
tools. When false, it is absent.
### URL-02 — Only http/https are fetched (coverage: test)
`file:`, `ftp:`, `data:`, `gopher:` and schemeless inputs are
refused before any network call, with an error result the model can
relay.
### URL-03 — SSRF guard blocks non-public addresses (coverage: test)
Before fetching, the host is resolved and every resulting IP is
checked; the fetch is refused when any is private, loopback,
link-local, or otherwise non-global (RFC1918, 127/8, 169.254/16,
::1, fc00::/7, etc.). A URL literal that is already such an IP is
refused without DNS.
### URL-04 — Redirects are re-validated (coverage: test)
Redirects are followed manually; each hop's target passes URL-02 and
URL-03 again. A public URL that 302-redirects to `localhost` or an
internal IP is refused at the redirect, not fetched.
### URL-05 — Fetched text is bounded and sanitized (coverage: test)
Responses are capped at `url-max-bytes` (default 2 MB) with a
download timeout; HTML is reduced to readable text (script/style
dropped, tags stripped, whitespace collapsed) and passed through
`sanitize_external_text` before it reaches the model, truncated to
`url-max-chars` (default 6000).
### URL-06 — Page images feed the cache (coverage: test)
Up to `url-max-images` (default 2) prominent images (og:image, then
large `<img>`) are ingested into the ImageCache for the requesting
channel (SSRF-guarded like the page), so the model can see them and
`picture_edit` can remix them. Ingestion failures are skipped, never
fatal to the text result.
### URL-07 — Fetches are metered and capped (coverage: test)
Each fetch increments a per-user daily counter; over
`url-daily-per-user` (default 20) `fetch_url` refuses with an error
result. The budget gate (SAF-04) still applies to the surrounding
model calls.
+32
View File
@@ -0,0 +1,32 @@
# SPEC-012 — Operations hardening
Runtime + host operability (FDB-012). Backups and host wiring are
`manual` coverage; the in-process alerting is `test`.
### OPS-13 — Consistent DB backups (coverage: test)
`deploy/backup_db.py` writes a gzipped snapshot of `bot.db` using the
sqlite3 online-backup API — consistent even while the bot writes
(WAL-safe) — with 0600 permissions. Restoring a snapshot yields a
readable database with the same rows.
### OPS-14 — Backups are rotated (coverage: test)
The newest `backup-keep` (default 14) snapshots are kept; older ones
are deleted. Timestamped names sort chronologically so rotation is a
pure list operation.
### OPS-15 — Backup cron on each host (coverage: manual)
Each host runs `backup_db.py` daily via cron, writing to
`~/backups/<bot>/` (outside `~/fjerkroa_bot`, so deploys and service
restarts never touch it). Verified by presence of the cron line and a
fresh snapshot.
### OPS-16 — Repeated API errors alert staff (coverage: test)
The responder counts consecutive OpenAI request failures; at
`api-error-alert-threshold` (default 5) in a row it fires one staff
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.
+25
View File
@@ -0,0 +1,25 @@
# SPEC-013 — News digest
Replaces the broken pre-1.0-openai `news_feed.py`. A CLI
(`python -m fjerkroa_bot.news --config <cfg>`) fetches the
`news-feeds` and writes a compact digest to the `news` file that
`AIResponder.message` injects into the `{news}` slot. Feeds are
external input and operator-configured.
### NEWS-01 — RSS 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
XML (returns an empty list, logs), and never raises.
### NEWS-02 — Digest is sanitized and bounded (coverage: test)
`render_digest` caps at `news-max-items`, and every headline passes
`sanitize_external_text` (SAF-03) — a feed cannot inject `@everyone`
or control characters into the prompt via a headline.
### NEWS-03 — Feeds are SSRF-guarded and deduped (coverage: test)
`NewsFetcher.collect` skips any feed URL the SSRF guard rejects,
skips feeds that fail to fetch (one bad feed never sinks the run),
and drops duplicate headlines across feeds.
+1 -1
View File
@@ -28,7 +28,7 @@ class TestIGDBIntegration(unittest.IsolatedAsyncioTestCase):
responder = OpenAIResponder(self.config_with_igdb)
mock_igdb.assert_called_once_with("test_client", "test_token")
mock_igdb.assert_called_once_with("test_client", "test_token", client_secret=None)
self.assertEqual(responder.igdb, mock_igdb_instance)
def test_igdb_initialization_disabled(self):
+100 -1
View File
@@ -160,7 +160,7 @@ class TestIGDBQuery(unittest.TestCase):
"id",
"name",
"alternative_names",
"category",
"game_type",
"release_dates",
"franchise",
"language_supports",
@@ -174,5 +174,104 @@ class TestIGDBQuery(unittest.TestCase):
self.assertEqual(result, [{"id": 1, "name": "Super Mario Bros"}])
class TestIGDBNativeSearch(unittest.TestCase):
def test_build_query_with_search_term(self):
"""search_games uses IGDB full-text search, not a name prefix filter."""
query = IGDBQuery.build_query(["name"], {"game_type": "= 0"}, limit=5, search_term="Marvel Tōkon")
self.assertEqual(query, 'search "Marvel Tōkon"; fields name; limit 5; where game_type = 0;')
def test_search_term_escapes_quotes_and_backslashes(self):
query = IGDBQuery.build_query(["name"], search_term='say "hi" \\ bye')
self.assertIn('search "say \\"hi\\" \\\\ bye";', query)
@patch.object(IGDBQuery, "generalized_igdb_query")
def test_search_games_passes_search_term(self, mock_query):
mock_query.return_value = []
IGDBQuery("cid", "token").search_games("Elden Ring", limit=3)
_, kwargs = mock_query.call_args
self.assertEqual(kwargs["search_term"], "Elden Ring")
self.assertEqual(mock_query.call_args.args[0], {})
class TestIGDBTokenRefresh(unittest.TestCase):
@staticmethod
def _oauth_response(token="fresh_token", expires_in=5_000_000):
response = Mock()
response.json.return_value = {"access_token": token, "expires_in": expires_in}
response.raise_for_status.return_value = None
return response
@staticmethod
def _api_response(payload, status_code=200):
response = Mock()
response.status_code = status_code
response.json.return_value = payload
response.raise_for_status.return_value = None
return response
@patch("fjerkroa_bot.igdblib.requests.post")
def test_fetches_token_when_only_secret_configured(self, mock_post):
"""Without a static token, the first request fetches one via Twitch OAuth."""
mock_post.side_effect = [self._oauth_response(), self._api_response([{"id": 1}])]
igdb = IGDBQuery("cid", client_secret="secret")
result = igdb.send_igdb_request("games", "fields name; limit 1;")
self.assertEqual(result, [{"id": 1}])
oauth_call, api_call = mock_post.call_args_list
self.assertEqual(oauth_call.args[0], "https://id.twitch.tv/oauth2/token")
self.assertEqual(
oauth_call.kwargs["params"],
{"client_id": "cid", "client_secret": "secret", "grant_type": "client_credentials"},
)
self.assertEqual(api_call.kwargs["headers"]["Authorization"], "Bearer fresh_token")
@patch("fjerkroa_bot.igdblib.requests.post")
def test_refreshes_and_retries_on_401(self, mock_post):
"""A 401 with a configured secret triggers one refresh and retry."""
mock_post.side_effect = [
self._api_response(None, status_code=401),
self._oauth_response(),
self._api_response([{"id": 2}]),
]
igdb = IGDBQuery("cid", "expired_token", client_secret="secret")
result = igdb.send_igdb_request("games", "fields name; limit 1;")
self.assertEqual(result, [{"id": 2}])
self.assertEqual(igdb.igdb_api_key, "fresh_token")
self.assertEqual(mock_post.call_args_list[2].kwargs["headers"]["Authorization"], "Bearer fresh_token")
@patch("fjerkroa_bot.igdblib.time.time")
@patch("fjerkroa_bot.igdblib.requests.post")
def test_proactive_refresh_before_expiry(self, mock_post, mock_time):
"""An expired self-fetched token is refreshed before the request."""
mock_time.return_value = 1_000_000.0
mock_post.side_effect = [self._oauth_response("token_a", expires_in=5_000_000), self._api_response([])]
igdb = IGDBQuery("cid", client_secret="secret")
igdb.send_igdb_request("games", "fields name;")
# jump past the token expiry -> next request refreshes first
mock_time.return_value = 1_000_000.0 + 5_000_000
mock_post.side_effect = [self._oauth_response("token_b"), self._api_response([])]
igdb.send_igdb_request("games", "fields name;")
self.assertEqual(igdb.igdb_api_key, "token_b")
@patch("fjerkroa_bot.igdblib.requests.post")
def test_no_refresh_without_secret(self, mock_post):
"""Static-token setups keep the old behavior: no OAuth calls, error -> None."""
response = Mock()
response.status_code = 401
response.raise_for_status.side_effect = requests.RequestException("401 Client Error")
mock_post.return_value = response
igdb = IGDBQuery("cid", "expired_token")
result = igdb.send_igdb_request("games", "fields name; limit 1;")
self.assertIsNone(result)
mock_post.assert_called_once()
if __name__ == "__main__":
unittest.main()
+61
View File
@@ -45,6 +45,45 @@ class TestIngest(unittest.TestCase):
self.assertIsNone(cache.ingest_bytes(PNG, "chat", "alice", "1"))
class TestOversizedDownloadRejected(unittest.IsolatedAsyncioTestCase):
async def test_download_stays_over_limit_and_is_rejected(self):
"""IMG-10: an over-limit download must be rejected, not cached truncated."""
with tempfile.TemporaryDirectory() as tmp:
_, cache = make_cache(tmp, {"image-max-bytes": 32})
class FakeContent:
@staticmethod
async def iter_chunked(size):
yield PNG # 72 bytes > 32
class FakeResp:
content = FakeContent()
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
def raise_for_status(self):
pass
class FakeSession:
def get(self, url):
return FakeResp()
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
with patch("fjerkroa_bot.images.aiohttp.ClientSession", return_value=FakeSession()):
data = await cache._download("http://x.com/big.png")
self.assertEqual(len(data), 33) # limit + 1, not silently capped to limit
self.assertIsNone(await cache.ingest_url("http://x.com/big.png", "chat", "alice", "1"))
class TestVisionDataUrls(OpsBase):
async def test_attachment_becomes_data_url(self):
"""IMG-11: the model sees a data: URL, never the CDN link."""
@@ -158,6 +197,28 @@ class TestGeneratedImagesCached(OpsBase):
self.assertEqual(rows[0]["user"], "assistant")
class TestImageOnlyMessages(OpsBase):
async def test_image_only_post_cached_no_reply(self):
"""IMG-17: attachment without text -> cached + observed, no reply."""
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp)
self.bot.airesponder.image_cache = cache
self.bot.airesponder.observe_event = AsyncMock()
self.bot.respond = AsyncMock()
message = self.public_msg("")
message.content = ""
message.channel.name = "chat"
attachment = Mock()
attachment.url = "https://cdn.discordapp.com/attachments/1/2/silent.png"
message.attachments = [attachment]
message.id = 77
with patch.object(ImageCache, "_download", new_callable=AsyncMock, return_value=PNG):
await self.bot.on_message(message)
self.assertEqual(len(store.images_recent("chat", 5)), 1)
self.bot.airesponder.observe_event.assert_awaited_once()
self.bot.respond.assert_not_awaited()
class TestContextAnnouncesImages(unittest.IsolatedAsyncioTestCase):
def test_suffix_mentions_picture_edit(self):
"""IMG-16: cached channel images are announced in the context suffix."""
+75
View File
@@ -0,0 +1,75 @@
"""Unit coverage for SPEC-013 news digest (NEWS-01..03)."""
import unittest
from unittest.mock import AsyncMock
from fjerkroa_bot.news import NewsFetcher, parse_feed, render_digest
RSS = b"""<?xml version="1.0"?><rss><channel>
<item><title>Game X released</title><link>https://ex.com/x</link></item>
<item><title>Patch Y notes</title><link>https://ex.com/y</link></item>
</channel></rss>"""
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>"""
class TestParse(unittest.TestCase):
def test_rss(self):
"""NEWS-01: RSS items parsed with title + link."""
items = parse_feed(RSS, "Src")
self.assertEqual([i["title"] for i in items], ["Game X released", "Patch Y notes"])
self.assertEqual(items[0]["link"], "https://ex.com/x")
self.assertEqual(items[0]["source"], "Src")
def test_atom(self):
"""NEWS-01: Atom entries parsed with href link."""
items = parse_feed(ATOM, "A")
self.assertEqual(items[0]["title"], "Atom headline")
self.assertEqual(items[0]["link"], "https://ex.com/a")
def test_malformed_never_raises(self):
"""NEWS-01: garbage XML returns [] without raising."""
self.assertEqual(parse_feed(b"<not xml", "bad"), [])
self.assertEqual(parse_feed(b"", "empty"), [])
class TestDigest(unittest.TestCase):
def test_sanitized_and_capped(self):
"""NEWS-02: headlines sanitized, item count capped."""
items = [{"title": "@everyone big news \x00", "link": "", "source": "S"} for _ in range(20)]
digest = render_digest(items, max_items=5)
self.assertEqual(digest.count("\n"), 4) # 5 lines
self.assertNotIn("@everyone", digest)
self.assertNotIn("\x00", digest)
class TestCollect(unittest.IsolatedAsyncioTestCase):
async def test_ssrf_skip_and_dedup(self):
"""NEWS-03: guarded feed skipped, dup titles dropped, bad fetch survived."""
def guard(url):
return "refused" if "internal" in url else None
async def fetch(url):
if "boom" in url:
raise ValueError("boom")
return RSS # same content from two feeds -> dedup
fetcher = NewsFetcher(guard, fetch)
feeds = [
("https://a.com/feed", "A"),
("https://internal/feed", "Internal"), # SSRF-skipped
("https://boom.com/feed", "Boom"), # fetch fails
("https://b.com/feed", "B"), # same RSS -> dup titles dropped
]
items = await fetcher.collect(feeds, per_feed=5)
titles = [i["title"] for i in items]
self.assertEqual(titles, ["Game X released", "Patch Y notes"]) # deduped, internal+boom skipped
async def test_per_feed_limit(self):
"""NEWS-03: per-feed cap honored."""
fetcher = NewsFetcher(lambda u: None, AsyncMock(return_value=RSS))
items = await fetcher.collect([("https://a.com", "A")], per_feed=1)
self.assertEqual(len(items), 1)
+100
View File
@@ -0,0 +1,100 @@
"""Unit coverage for SPEC-012 ops hardening (OPS-13/14/16)."""
import gzip
import sqlite3
import stat
import sys
import tempfile
import unittest
from pathlib import Path
from unittest.mock import AsyncMock, MagicMock
from fjerkroa_bot.persistence import PersistentStore
from .test_spec_ops import OpsBase
sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "deploy"))
import backup_db # noqa: E402
class TestSnapshotConsistency(unittest.TestCase):
def test_snapshot_roundtrips(self):
"""OPS-13: a gzipped snapshot restores to a readable DB with the same rows, 0600."""
with tempfile.TemporaryDirectory() as tmp:
db = Path(tmp) / "bot.db"
store = PersistentStore(db)
store.save_history("chat", [{"role": "user", "content": "hei"}])
store.add_user_fact("alice", "likes espresso", "self")
dest = Path(tmp) / "snap.db.gz"
backup_db.snapshot(db, dest)
self.assertEqual(stat.S_IMODE(dest.stat().st_mode), 0o600)
restored = Path(tmp) / "restored.db"
with gzip.open(dest, "rb") as gz, open(restored, "wb") as out:
out.write(gz.read())
conn = sqlite3.connect(restored)
try:
rows = conn.execute("SELECT content FROM history WHERE channel='chat'").fetchall()
facts = conn.execute("SELECT fact FROM user_facts").fetchall()
finally:
conn.close()
self.assertEqual(rows, [("hei",)])
self.assertEqual(facts, [("likes espresso",)])
def test_snapshot_during_writes(self):
"""OPS-13: snapshot succeeds while another connection holds the DB open (WAL)."""
with tempfile.TemporaryDirectory() as tmp:
db = Path(tmp) / "bot.db"
store = PersistentStore(db)
store.save_history("chat", [{"role": "user", "content": "x"}])
live = sqlite3.connect(db) # simulate the running bot's open handle
live.execute("PRAGMA journal_mode=WAL")
try:
dest = Path(tmp) / "snap.db.gz"
backup_db.snapshot(db, dest) # must not raise
self.assertTrue(dest.exists())
finally:
live.close()
class TestRotation(unittest.TestCase):
def test_victims_keeps_newest(self):
"""OPS-14: only the oldest beyond `keep` are selected for deletion."""
names = [f"bot-2026070{d}-000000.db.gz" for d in range(1, 8)] # 7 chronological
victims = backup_db.victims(list(reversed(names)), keep=3)
self.assertEqual(victims, names[:4]) # oldest 4 removed, newest 3 kept
def test_victims_under_keep_deletes_nothing(self):
"""OPS-14: fewer than `keep` backups -> nothing deleted."""
self.assertEqual(backup_db.victims(["bot-20260701-000000.db.gz"], keep=14), [])
def test_rotate_on_disk(self):
"""OPS-14: rotate removes the right files from a real dir."""
with tempfile.TemporaryDirectory() as tmp:
for d in range(1, 6):
(Path(tmp) / f"bot-2026070{d}-000000.db.gz").write_bytes(b"x")
removed = backup_db.rotate(Path(tmp), keep=2)
self.assertEqual(removed, 3)
self.assertEqual(len(list(Path(tmp).glob("bot-*.db.gz"))), 2)
class TestApiErrorAlert(OpsBase):
async def test_threshold_alert_and_reset(self):
"""OPS-16: N consecutive failures fire one staff alert; success resets."""
self.bot.config["api-error-alert-threshold"] = 3
self.bot.send_message_with_typing = AsyncMock(side_effect=RuntimeError("boom"))
origin = MagicMock()
from fjerkroa_bot.ai_responder import AIMessage
for _ in range(3):
await self.bot.respond(AIMessage("alice", "hei", "chat"), origin)
self.assertEqual(self.bot.staff_channel.send.await_count, 1) # exactly one alert at threshold
self.assertEqual(self.bot._consecutive_api_errors, 3)
# a success resets the counter
from fjerkroa_bot.ai_responder import AIResponse
self.bot.send_message_with_typing = AsyncMock(return_value=AIResponse(None, False, "chat", None, None, False, False))
await self.bot.respond(AIMessage("alice", "hei", "chat"), origin)
self.assertEqual(self.bot._consecutive_api_errors, 0)
+165
View File
@@ -0,0 +1,165 @@
"""Unit coverage for SPEC-005 self-tasking (TSK-01..08) + OPS-12."""
import tempfile
import unittest
from pathlib import Path
from unittest.mock import AsyncMock
from fjerkroa_bot.persistence import PersistentStore
from fjerkroa_bot.quota import QuotaLedger
from fjerkroa_bot.tasks import TaskEngine
from .test_spec_ops import OpsBase
class EngineBase(unittest.IsolatedAsyncioTestCase):
def setUp(self):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
self.store = PersistentStore(Path(self.tmp.name) / "bot.db")
self.config = {"tasks-enabled": True, "chat-channel": "chat", "taskgen-interval-hours": 0}
self.ledger = QuotaLedger(self.store, lambda: self.config)
self.execute = AsyncMock()
self.propose = AsyncMock(return_value={"task": None})
self.alert = AsyncMock()
self.observe = AsyncMock()
self.idle = lambda: 0.0
self.engine = TaskEngine(
self.store,
self.ledger,
lambda: self.config,
self.execute,
self.propose,
self.alert,
lambda: True,
lambda: self.idle(),
self.observe,
)
class TestQueuePersistence(EngineBase):
async def test_tasks_survive_restart(self):
"""TSK-01: enqueued tasks are store rows and reload in a fresh store."""
await self.engine.enqueue("follow-up", "chat", "frag bob nach der pruefung", due_hours=1)
reborn = PersistentStore(Path(self.tmp.name) / "bot.db")
tasks = reborn.tasks_open()
self.assertEqual(len(tasks), 1)
self.assertEqual((tasks[0]["kind"], tasks[0]["channel"], tasks[0]["state"]), ("follow-up", "chat", "queued"))
class TestExecution(EngineBase):
async def test_due_task_executes_and_completes(self):
"""TSK-02: due task runs via execute callback, marked done, observed."""
await self.engine.enqueue("idle-impulse", "chat", "sag was nettes")
await self.engine.tick()
self.execute.assert_awaited_once_with("chat", "sag was nettes")
self.assertEqual(self.store.tasks_open(), [])
self.observe.assert_awaited()
async def test_failed_execution_marked_failed(self):
"""TSK-02: execute raising marks the task failed, not done."""
self.execute.side_effect = RuntimeError("channel gone")
await self.engine.enqueue("idle-impulse", "chat", "x")
await self.engine.tick()
self.assertEqual(self.store.tasks_open(), []) # not queued anymore
rows = [t for t in self.store.tasks_due()]
self.assertEqual(rows, [])
class TestDefaultOff(EngineBase):
async def test_disabled_engine_is_dormant(self):
"""TSK-03: without tasks-enabled nothing executes or generates."""
self.store.task_add("idle-impulse", "chat", "2020-01-01 00:00:00", "x")
self.config.pop("tasks-enabled")
await self.engine.tick()
self.execute.assert_not_awaited()
self.propose.assert_not_awaited()
class TestDailyCap(EngineBase):
async def test_cap_defers_excess_tasks(self):
"""TSK-04: over the per-channel cap tasks stay queued."""
self.config["tasks-max-per-channel-per-day"] = 1
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "one")
self.store.task_add("b", "chat", "2020-01-01 00:00:00", "two")
await self.engine.tick()
self.assertEqual(self.execute.await_count, 1)
self.assertEqual(len(self.store.tasks_open()), 1)
class TestApproval(EngineBase):
async def test_approval_flow(self):
"""TSK-05: approval mode holds tasks until approved; cancel kills them."""
self.config["tasks-approval"] = True
task_id = await self.engine.enqueue("follow-up", "chat", "frag nach")
self.alert.assert_awaited_once()
await self.engine.tick()
self.execute.assert_not_awaited() # approval != due
self.store.task_set_state(task_id, "queued")
await self.engine.tick()
self.execute.assert_awaited_once()
class TestKillSwitch(EngineBase):
async def test_not_allowed_blocks_tick(self):
"""TSK-06: bot_initiated_allowed()=False stops execution."""
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "x")
engine = TaskEngine(
self.store, self.ledger, lambda: self.config, self.execute, self.propose, self.alert, lambda: False, lambda: 0.0, self.observe
)
await engine.tick()
self.execute.assert_not_awaited()
class TestIdleImpulse(EngineBase):
async def test_idle_enqueues_once(self):
"""TSK-07: long idle enqueues one impulse; pending impulse dedupes."""
self.config["idle-impulse-hours"] = 1
self.config["boreness-prompt"] = "denk dir was aus"
self.idle = lambda: 2 * 3600.0
await self.engine.tick()
open_tasks = self.store.tasks_open()
impulse = [t for t in open_tasks if t["kind"] == "idle-impulse"]
self.assertEqual(len(impulse), 0) # executed immediately (due now)
self.execute.assert_awaited_once_with("chat", "denk dir was aus")
async def test_short_idle_no_impulse(self):
"""TSK-07: below the idle threshold nothing is generated."""
self.config["idle-impulse-hours"] = 12
self.idle = lambda: 60.0
await self.engine.tick()
self.execute.assert_not_awaited()
class TestFollowUpGenerator(EngineBase):
async def test_proposal_becomes_task(self):
"""TSK-08: a proposed task is enqueued with its due offset."""
self.config.pop("chat-channel")
self.propose.return_value = {"task": {"channel": "chat", "prompt": "frag bob", "due_hours": 24}}
await self.engine.tick()
tasks = self.store.tasks_open()
self.assertEqual(len(tasks), 1)
self.assertEqual(tasks[0]["kind"], "follow-up")
async def test_null_proposal_no_task(self):
"""TSK-08: null proposal enqueues nothing."""
self.propose.return_value = {"task": None}
await self.engine.tick()
self.assertEqual(self.store.tasks_open(), [])
class TestTaskCommands(OpsBase):
async def test_list_approve_cancel(self):
"""OPS-12: !bot tasks lists; task-approve/task-cancel manage states."""
with tempfile.TemporaryDirectory() as tmp:
store = PersistentStore(Path(tmp) / "bot.db")
self.bot.airesponder.store = store
task_id = store.task_add("follow-up", "chat", "2099-01-01 00:00:00", "frag bob", "approval")
await self.bot.on_message(self.staff_msg("!bot tasks"))
listing = self.bot.staff_channel.send.await_args.args[0]
self.assertIn("follow-up", listing)
self.assertIn("approval", listing)
await self.bot.on_message(self.staff_msg(f"!bot task-approve {task_id}"))
self.assertEqual(store.tasks_open()[0]["state"], "queued")
await self.bot.on_message(self.staff_msg(f"!bot task-cancel {task_id}"))
self.assertEqual(store.tasks_open(), [])
+183
View File
@@ -0,0 +1,183 @@
"""Unit coverage for SPEC-011 URL reading (URL-01..07)."""
import unittest
from unittest.mock import AsyncMock, patch
from fjerkroa_bot.openai_responder import OpenAIResponder
from fjerkroa_bot.url_reader import FETCH_URL_TOOL, URLReader, guard_url
CONFIG = {"openai-token": "t", "model": "m", "system": "s", "history-limit": 5}
class TestToolOffered(unittest.TestCase):
def test_tool_present_only_when_enabled(self):
"""URL-01: fetch_url appears in the tool list only with enable-url-reading."""
off = OpenAIResponder(CONFIG, "chat")
self.assertNotIn("fetch_url", [f["name"] for f in off._available_tools()])
on = OpenAIResponder(dict(CONFIG, **{"enable-url-reading": True}), "chat")
self.assertIn("fetch_url", [f["name"] for f in on._available_tools()])
self.assertEqual(FETCH_URL_TOOL["name"], "fetch_url")
class TestSchemeGuard(unittest.TestCase):
def test_non_http_schemes_refused(self):
"""URL-02: only http/https pass the guard."""
self.assertIsNone(guard_url("https://example.com/article"))
for bad in ("file:///etc/passwd", "ftp://host/x", "data:text/html,x", "gopher://h", "no-scheme.com/x"):
self.assertIsNotNone(guard_url(bad))
class TestSSRFGuard(unittest.TestCase):
def test_private_and_loopback_refused(self):
"""URL-03: private/loopback/link-local literals are refused without DNS."""
for bad in (
"http://127.0.0.1/admin",
"http://localhost/x", # resolves to loopback
"http://10.0.0.5/x",
"http://192.168.1.1/x",
"http://169.254.169.254/latest/meta-data", # cloud metadata
"http://[::1]/x",
):
self.assertIsNotNone(guard_url(bad), f"{bad} should be refused")
def test_public_ip_allowed(self):
"""URL-03: a public IP literal passes."""
self.assertIsNone(guard_url("http://93.184.216.34/"))
@patch("fjerkroa_bot.url_reader.socket.getaddrinfo")
def test_dns_to_private_refused(self, getaddrinfo):
"""URL-03: a hostname resolving to a private IP is refused."""
getaddrinfo.return_value = [(2, 1, 6, "", ("10.1.2.3", 0))]
self.assertIsNotNone(guard_url("http://evil.example.com/x"))
class TestRedirectRevalidation(unittest.IsolatedAsyncioTestCase):
async def test_redirect_to_internal_refused(self):
"""URL-04: a public URL redirecting to localhost is refused at the hop."""
reader = URLReader(lambda: {}, None)
class FakeResp:
status = 302
headers = {"Location": "http://127.0.0.1/secret"}
url = "http://safe.example.com"
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
def raise_for_status(self):
pass
class FakeSession:
def get(self, url, allow_redirects=False):
return FakeResp()
with patch("fjerkroa_bot.url_reader.guard_url", side_effect=[None, "refused internal"]):
with self.assertRaises(ValueError):
await reader._get(FakeSession(), "http://safe.example.com", 1000)
class TestTextExtraction(unittest.TestCase):
def test_html_reduced_to_text(self):
"""URL-05: scripts/styles dropped, tags stripped."""
reader = URLReader(lambda: {}, None)
html = "<html><head><style>x{}</style></head><body><h1>Titel</h1><script>evil()</script><p>Inhalt hier</p></body></html>"
text = reader._to_text(html)
self.assertIn("Titel", text)
self.assertIn("Inhalt hier", text)
self.assertNotIn("evil", text)
self.assertNotIn("x{}", text)
class TestBodyReadCollectsAllChunks(unittest.IsolatedAsyncioTestCase):
async def test_get_reads_past_first_chunk(self):
"""URL-05 regression: body arrives in many chunks; all are collected up to the cap."""
reader = URLReader(lambda: {}, None)
chunks = [b"<title>t</title>", b"<p>middle</p>", b"<p>end</p>"]
class FakeContent:
@staticmethod
async def iter_chunked(size):
for chunk in chunks:
yield chunk
class FakeResp:
status = 200
headers = {}
url = "http://safe.example.com"
content = FakeContent()
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
def raise_for_status(self):
pass
class FakeSession:
def get(self, url, allow_redirects=False):
return FakeResp()
with patch("fjerkroa_bot.url_reader.guard_url", return_value=None):
_, body = await reader._get(FakeSession(), "http://safe.example.com", 1000)
self.assertEqual(body, b"".join(chunks))
with patch("fjerkroa_bot.url_reader.guard_url", return_value=None):
_, body = await reader._get(FakeSession(), "http://safe.example.com", 20)
self.assertEqual(body, b"".join(chunks)[:20])
class TestFetchSanitizes(unittest.IsolatedAsyncioTestCase):
async def test_fetch_result_is_sanitized_and_capped(self):
"""URL-05: fetch output is length-capped and @everyone-neutralized."""
reader = URLReader(lambda: {"url-max-chars": 50}, None)
payload = ("<p>@everyone " + "x" * 5000 + "</p>").encode()
with patch.object(reader, "_get", new=AsyncMock(return_value=("http://x.com", payload))):
result = await reader.fetch("http://x.com", "chat", "alice")
self.assertLessEqual(len(result["text"]), 50)
self.assertNotIn("@everyone", result["text"])
async def test_fetch_error_is_reported_not_raised(self):
"""URL-05: a fetch failure returns an error dict the model can relay."""
reader = URLReader(lambda: {}, None)
with patch.object(reader, "_get", new=AsyncMock(side_effect=ValueError("refused non-public address"))):
result = await reader.fetch("http://10.0.0.1", "chat", "alice")
self.assertIn("error", result)
class TestImageIngest(unittest.IsolatedAsyncioTestCase):
async def test_page_images_go_to_cache_ssrf_guarded(self):
"""URL-06: og:image + <img> ingested (cap honored), internal srcs skipped."""
cache = type("C", (), {})()
cache.ingest_url = AsyncMock(side_effect=["sha1", "sha2", "sha3"])
reader = URLReader(lambda: {"url-max-images": 2}, cache)
html = (
'<meta property="og:image" content="https://cdn.example.com/hero.jpg">'
'<img src="https://cdn.example.com/a.png"><img src="http://127.0.0.1/internal.png">'
)
# guard by scheme/loopback only, no real DNS in the test
def fake_guard(url):
return "refused" if "127.0.0.1" in url else None
with patch("fjerkroa_bot.url_reader.guard_url", side_effect=fake_guard):
count = await reader._ingest_images(html, "https://example.com", "chat", "alice")
self.assertEqual(count, 2) # og:image + first public img, cap 2
ingested = [call.args[0] for call in cache.ingest_url.await_args_list]
self.assertNotIn("http://127.0.0.1/internal.png", ingested)
class TestPerUserCap(unittest.IsolatedAsyncioTestCase):
async def test_dispatch_caps_fetches(self):
"""URL-07: over url-daily-per-user, fetch_url returns an error without fetching."""
responder = OpenAIResponder(dict(CONFIG, **{"enable-url-reading": True, "url-daily-per-user": 2}), "chat")
responder.url_reader.fetch = AsyncMock(return_value={"url": "x", "text": "ok"})
for _ in range(2):
await responder._dispatch_tool("fetch_url", {"url": "http://x.com"}, "alice")
blocked = await responder._dispatch_tool("fetch_url", {"url": "http://x.com"}, "alice")
self.assertIn("error", blocked)
self.assertEqual(responder.url_reader.fetch.await_count, 2)
Generated
+2
View File
@@ -627,6 +627,7 @@ version = "3.0.0"
source = { editable = "." }
dependencies = [
{ name = "aiohttp" },
{ name = "defusedxml" },
{ name = "discord-py" },
{ name = "openai" },
{ name = "requests" },
@@ -656,6 +657,7 @@ dev = [
[package.metadata]
requires-dist = [
{ name = "aiohttp", specifier = ">=3.12" },
{ name = "defusedxml", specifier = ">=0.7" },
{ name = "discord-py", specifier = ">=2.5,<3" },
{ name = "openai", specifier = ">=2.45" },
{ name = "requests", specifier = ">=2.32" },