Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a514ff652c | |||
| 7628faf551 | |||
| 2caa18a17f | |||
| d4eec4088d | |||
| 86e631926f | |||
| 4166520923 | |||
| d0819c2683 |
+21
-8
@@ -20,13 +20,6 @@ The bot now supports real-time video game information through IGDB (Internet Gam
|
|||||||
- **Category**: Select appropriate category
|
- **Category**: Select appropriate category
|
||||||
3. Note down your **Client ID**
|
3. Note down your **Client ID**
|
||||||
4. Generate a **Client Secret**
|
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
|
### 2. Configure the Bot
|
||||||
|
|
||||||
@@ -35,6 +28,25 @@ Update your `config.toml` file:
|
|||||||
```toml
|
```toml
|
||||||
# IGDB Configuration for game information
|
# IGDB Configuration for game information
|
||||||
igdb-client-id = "your_actual_client_id_here"
|
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"
|
igdb-access-token = "your_actual_access_token_here"
|
||||||
enable-game-info = true
|
enable-game-info = true
|
||||||
```
|
```
|
||||||
@@ -99,7 +111,8 @@ The integration provides two OpenAI functions:
|
|||||||
- Verify client ID and access token are set
|
- Verify client ID and access token are set
|
||||||
|
|
||||||
2. **Authentication errors**
|
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
|
- Verify client ID matches your Twitch app
|
||||||
|
|
||||||
3. **No game results**
|
3. **No game results**
|
||||||
|
|||||||
@@ -83,3 +83,6 @@ ci: install-dev all-checks ## Full CI pipeline (install deps and run all checks)
|
|||||||
# Deploy targets (SPEC-007)
|
# Deploy targets (SPEC-007)
|
||||||
deploy: ## Deploy a tag to a host: make deploy HOST=ggg TAG=v3.0.0
|
deploy: ## Deploy a tag to a host: make deploy HOST=ggg TAG=v3.0.0
|
||||||
bash deploy/deploy.sh $(HOST) $(TAG)
|
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)
|
||||||
|
|||||||
+26
-1
@@ -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 Configuration for game information
|
||||||
igdb-client-id = "YOUR_IGDB_CLIENT_ID"
|
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
|
enable-game-info = true
|
||||||
|
|
||||||
# --- operator / safety (SPEC-003, SPEC-006) ---
|
# --- operator / safety (SPEC-003, SPEC-006) ---
|
||||||
@@ -67,3 +71,24 @@ enable-game-info = true
|
|||||||
# idle-impulse-hours = 12
|
# idle-impulse-hours = 12
|
||||||
# taskgen-interval-hours = 6
|
# taskgen-interval-hours = 6
|
||||||
# Staff: !bot tasks | task-approve <id> | task-cancel <id>
|
# 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"],
|
||||||
|
# ]
|
||||||
|
|||||||
@@ -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())
|
||||||
@@ -89,6 +89,7 @@ class FjerkroaBot(commands.Bot):
|
|||||||
self.tasks_enabled = True
|
self.tasks_enabled = True
|
||||||
self.quiet_until = 0.0
|
self.quiet_until = 0.0
|
||||||
self._staff_alert_times: deque = deque()
|
self._staff_alert_times: deque = deque()
|
||||||
|
self._consecutive_api_errors = 0 # OPS-16
|
||||||
|
|
||||||
self.init_observer()
|
self.init_observer()
|
||||||
self.init_aichannels()
|
self.init_aichannels()
|
||||||
@@ -498,6 +499,14 @@ class FjerkroaBot(commands.Bot):
|
|||||||
return True, False
|
return True, False
|
||||||
return False, bool(verdict.get("factual", 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):
|
async def send_message_with_typing(self, airesponder, channel, message):
|
||||||
"""Send the user message to the AI responder with typing animation in discord"""
|
"""Send the user message to the AI responder with typing animation in discord"""
|
||||||
async with channel.typing():
|
async with channel.typing():
|
||||||
@@ -603,8 +612,15 @@ class FjerkroaBot(commands.Bot):
|
|||||||
# Get the AI responder based on the channel name
|
# Get the AI responder based on the channel name
|
||||||
airesponder = self.get_ai_responder(channel_name)
|
airesponder = self.get_ai_responder(channel_name)
|
||||||
|
|
||||||
# Send the user message to the AI responder, with typing indicators
|
# Send the user message to the AI responder, with typing indicators.
|
||||||
response = await self.send_message_with_typing(airesponder, channel, message)
|
# 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
|
# SAF/OPS gates between model proposal and delivery
|
||||||
await self._apply_response_gates(message, response)
|
await self._apply_response_gates(message, response)
|
||||||
|
|||||||
@@ -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
@@ -1,30 +1,67 @@
|
|||||||
import logging
|
import logging
|
||||||
|
import time
|
||||||
from functools import cache
|
from functools import cache
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
import requests
|
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):
|
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.client_id = client_id
|
||||||
self.igdb_api_key = igdb_api_key
|
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):
|
def send_igdb_request(self, endpoint, query_body):
|
||||||
igdb_url = f"https://api.igdb.com/v4/{endpoint}"
|
igdb_url = f"https://api.igdb.com/v4/{endpoint}"
|
||||||
headers = {"Client-ID": self.client_id, "Authorization": f"Bearer {self.igdb_api_key}"}
|
|
||||||
|
|
||||||
try:
|
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()
|
response.raise_for_status()
|
||||||
return response.json()
|
return response.json()
|
||||||
except requests.RequestException as e:
|
except requests.RequestException as e:
|
||||||
print(f"Error during IGDB API request: {e}")
|
print(f"Error during IGDB API request: {e}")
|
||||||
return None
|
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
|
@staticmethod
|
||||||
def build_query(fields, filters=None, limit=10, offset=None):
|
def build_query(fields, filters=None, limit=10, offset=None, search_term=None):
|
||||||
query = f"fields {','.join(fields) if fields is not None and len(fields) > 0 else '*'}; limit {limit};"
|
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:
|
if offset is not None:
|
||||||
query += f" offset {offset};"
|
query += f" offset {offset};"
|
||||||
if filters:
|
if filters:
|
||||||
@@ -32,12 +69,12 @@ class IGDBQuery(object):
|
|||||||
query += " where " + " & ".join(filter_statements) + ";"
|
query += " where " + " & ".join(filter_statements) + ";"
|
||||||
return query
|
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}
|
all_filters = {key: f'~ "{value}"*' for key, value in params.items() if value}
|
||||||
if additional_filters:
|
if additional_filters:
|
||||||
all_filters.update(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)
|
data = self.send_igdb_request(endpoint, query)
|
||||||
print(f"{endpoint}: {query} -> {data}")
|
print(f"{endpoint}: {query} -> {data}")
|
||||||
return data
|
return data
|
||||||
@@ -79,7 +116,7 @@ class IGDBQuery(object):
|
|||||||
"id",
|
"id",
|
||||||
"name",
|
"name",
|
||||||
"alternative_names",
|
"alternative_names",
|
||||||
"category",
|
"game_type",
|
||||||
"release_dates",
|
"release_dates",
|
||||||
"franchise",
|
"franchise",
|
||||||
"language_supports",
|
"language_supports",
|
||||||
@@ -101,9 +138,10 @@ class IGDBQuery(object):
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
try:
|
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(
|
games = self.generalized_igdb_query(
|
||||||
{"name": query.strip()},
|
{},
|
||||||
"games",
|
"games",
|
||||||
[
|
[
|
||||||
"id",
|
"id",
|
||||||
@@ -120,8 +158,9 @@ class IGDBQuery(object):
|
|||||||
"themes.name",
|
"themes.name",
|
||||||
"cover.url",
|
"cover.url",
|
||||||
],
|
],
|
||||||
additional_filters={"category": "= 0"}, # Main games only
|
additional_filters={"game_type": "= 0"}, # Main games only (IGDB renamed category -> game_type)
|
||||||
limit=limit,
|
limit=limit,
|
||||||
|
search_term=query.strip(),
|
||||||
)
|
)
|
||||||
|
|
||||||
if not games:
|
if not games:
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ from typing import Any, Callable, Dict, List, Optional
|
|||||||
|
|
||||||
import aiohttp
|
import aiohttp
|
||||||
|
|
||||||
|
from .httpread import read_capped
|
||||||
from .persistence import PersistentStore
|
from .persistence import PersistentStore
|
||||||
|
|
||||||
DEFAULT_CACHE_MB = 500
|
DEFAULT_CACHE_MB = 500
|
||||||
@@ -79,7 +80,9 @@ class ImageCache:
|
|||||||
async with aiohttp.ClientSession(timeout=timeout) as session:
|
async with aiohttp.ClientSession(timeout=timeout) as session:
|
||||||
async with session.get(url) as response:
|
async with session.get(url) as response:
|
||||||
response.raise_for_status()
|
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]:
|
def data_url(self, sha256: str, ext: str) -> Optional[str]:
|
||||||
path = self._path(sha256, ext)
|
path = self._path(sha256, ext)
|
||||||
|
|||||||
@@ -0,0 +1,291 @@
|
|||||||
|
"""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
|
||||||
|
|
||||||
|
|
||||||
|
DEFAULT_SEEN_CAP = 5000
|
||||||
|
DEFAULT_POST_PER_FEED = 5
|
||||||
|
DEFAULT_POST_MAX_PER_RUN = 8
|
||||||
|
|
||||||
|
|
||||||
|
def item_key(item: Dict[str, str]) -> str:
|
||||||
|
return item.get("link") or item.get("title") or ""
|
||||||
|
|
||||||
|
|
||||||
|
class NewsPoster:
|
||||||
|
"""Post NEW feed items to Discord channel webhooks (ggg model, SPEC-013 NEWS-04..06)."""
|
||||||
|
|
||||||
|
def __init__(self, guard, fetch_bytes, post_webhook) -> None:
|
||||||
|
self._guard = guard
|
||||||
|
self._fetch_bytes = fetch_bytes
|
||||||
|
self._post_webhook = post_webhook
|
||||||
|
|
||||||
|
async def run_post(
|
||||||
|
self,
|
||||||
|
feeds: List[Tuple[str, str, str]],
|
||||||
|
webhooks: Dict[str, str],
|
||||||
|
seen: set,
|
||||||
|
per_feed: int,
|
||||||
|
max_per_run: int,
|
||||||
|
seed_only: bool,
|
||||||
|
) -> Tuple[int, set]:
|
||||||
|
"""Returns (posted_count, updated_seen). seed_only marks new items seen without posting."""
|
||||||
|
posted = 0
|
||||||
|
for url, label, channel in feeds:
|
||||||
|
reason = self._guard(url)
|
||||||
|
if reason:
|
||||||
|
logging.warning(f"news-post: skipping feed {label} — {reason}")
|
||||||
|
continue
|
||||||
|
try:
|
||||||
|
data = await self._fetch_bytes(url)
|
||||||
|
except Exception as err:
|
||||||
|
logging.warning(f"news-post: fetch failed for {label}: {repr(err)}")
|
||||||
|
continue
|
||||||
|
for item in parse_feed(data, label)[:per_feed]:
|
||||||
|
key = item_key(item)
|
||||||
|
if not key or key in seen:
|
||||||
|
continue
|
||||||
|
seen.add(key)
|
||||||
|
may_post = not seed_only and posted < max_per_run
|
||||||
|
if may_post and await self._deliver(item, label, channel, webhooks):
|
||||||
|
posted += 1
|
||||||
|
return posted, seen
|
||||||
|
|
||||||
|
async def _deliver(self, item: Dict[str, str], label: str, channel: str, webhooks: Dict[str, str]) -> bool:
|
||||||
|
hook = webhooks.get(channel)
|
||||||
|
if not hook:
|
||||||
|
logging.warning(f"news-post: no webhook for channel {channel!r} ({label})")
|
||||||
|
return False
|
||||||
|
title = sanitize_external_text(item["title"], 300)
|
||||||
|
link = item.get("link", "")
|
||||||
|
content = f"**[{label}]** {title}" + (f"\n{link}" if link else "")
|
||||||
|
try:
|
||||||
|
await self._post_webhook(hook, content)
|
||||||
|
return True
|
||||||
|
except Exception as err:
|
||||||
|
logging.warning(f"news-post: webhook post failed ({label}): {repr(err)}")
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def load_seen(path: str) -> Tuple[set, bool]:
|
||||||
|
"""(seen-set, existed). Missing/broken state -> empty set, existed=False (seed run)."""
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
|
||||||
|
if not os.path.exists(path):
|
||||||
|
return set(), False
|
||||||
|
try:
|
||||||
|
with open(path, encoding="utf-8") as fd:
|
||||||
|
return set(json.load(fd)), True
|
||||||
|
except Exception as err:
|
||||||
|
logging.warning(f"news-post: unreadable state {path}: {err!r} — reseeding")
|
||||||
|
return set(), False
|
||||||
|
|
||||||
|
|
||||||
|
def save_seen(path: str, seen: set, cap: int = DEFAULT_SEEN_CAP) -> None:
|
||||||
|
import json
|
||||||
|
|
||||||
|
# keep the newest `cap` keys (insertion order preserved by Python sets? no — use a bounded slice)
|
||||||
|
keys = list(seen)[-cap:]
|
||||||
|
with open(path, "w", encoding="utf-8") as fd:
|
||||||
|
json.dump(keys, fd)
|
||||||
|
|
||||||
|
|
||||||
|
def _post_feeds_from_config(config: Dict[str, Any]) -> List[Tuple[str, str, str]]:
|
||||||
|
feeds = []
|
||||||
|
for entry in config.get("news-post-feeds", []):
|
||||||
|
if isinstance(entry, (list, tuple)) and len(entry) >= 3:
|
||||||
|
feeds.append((str(entry[0]), str(entry[1]), str(entry[2])))
|
||||||
|
return feeds
|
||||||
|
|
||||||
|
|
||||||
|
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 _aiohttp_post(hook: str, content: str) -> None:
|
||||||
|
import aiohttp
|
||||||
|
|
||||||
|
timeout = aiohttp.ClientTimeout(total=FETCH_TIMEOUT_S)
|
||||||
|
async with aiohttp.ClientSession(timeout=timeout) as session:
|
||||||
|
# allowed_mentions none: a headline can never ping the channel (SAF-02 spirit)
|
||||||
|
payload = {"content": content[:2000], "allowed_mentions": {"parse": []}}
|
||||||
|
async with session.post(hook, json=payload) as response:
|
||||||
|
response.raise_for_status()
|
||||||
|
|
||||||
|
|
||||||
|
async def run_post(config: Dict[str, Any]) -> int:
|
||||||
|
"""Webhook-posting mode (ggg): post new items to channels. Returns posted count."""
|
||||||
|
from .url_reader import guard_url
|
||||||
|
|
||||||
|
webhooks = dict(config.get("news-post-webhooks", {}))
|
||||||
|
feeds = _post_feeds_from_config(config)
|
||||||
|
state_path = config.get("news-post-state", "news_state.json")
|
||||||
|
if not webhooks or not feeds:
|
||||||
|
logging.error("news-post: need news-post-webhooks and news-post-feeds")
|
||||||
|
return 0
|
||||||
|
seen, existed = load_seen(state_path)
|
||||||
|
poster = NewsPoster(guard_url, _aiohttp_fetch, _aiohttp_post)
|
||||||
|
posted, seen = await poster.run_post(
|
||||||
|
feeds,
|
||||||
|
webhooks,
|
||||||
|
seen,
|
||||||
|
int(config.get("news-post-per-feed", DEFAULT_POST_PER_FEED)),
|
||||||
|
int(config.get("news-post-max-per-run", DEFAULT_POST_MAX_PER_RUN)),
|
||||||
|
seed_only=not existed, # first run seeds without flooding the channels
|
||||||
|
)
|
||||||
|
save_seen(state_path, seen, int(config.get("news-post-seen-cap", DEFAULT_SEEN_CAP)))
|
||||||
|
logging.info(f"news-post: posted {posted} item(s)" + (" (seed run — nothing posted)" if not existed else ""))
|
||||||
|
return posted
|
||||||
|
|
||||||
|
|
||||||
|
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: --post to channel webhooks (ggg) or default {news} digest file (kroa)"
|
||||||
|
)
|
||||||
|
parser.add_argument("--config", required=True)
|
||||||
|
parser.add_argument("--post", action="store_true", help="webhook-posting mode (post new items to Discord channels)")
|
||||||
|
args = parser.parse_args()
|
||||||
|
with open(args.config, encoding="utf-8") as fd:
|
||||||
|
config = tomlkit.load(fd)
|
||||||
|
if args.post:
|
||||||
|
asyncio.run(run_post(config))
|
||||||
|
return 0
|
||||||
|
result = asyncio.run(run(config))
|
||||||
|
return 0 if result else 1
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
sys.exit(main())
|
||||||
@@ -12,6 +12,7 @@ from .ai_responder import AIResponder, exponential_backoff, sanitize_external_te
|
|||||||
from .igdblib import IGDBQuery
|
from .igdblib import IGDBQuery
|
||||||
from .leonardo_draw import LeonardoAIDrawMixIn
|
from .leonardo_draw import LeonardoAIDrawMixIn
|
||||||
from .quota import QuotaLedger
|
from .quota import QuotaLedger
|
||||||
|
from .url_reader import FETCH_URL_TOOL, URLReader
|
||||||
|
|
||||||
# The response envelope, enforced server-side via structured outputs
|
# The response envelope, enforced server-side via structured outputs
|
||||||
# (ENV-19). All fields required, closed object, nullable where the
|
# (ENV-19). All fields required, closed object, nullable where the
|
||||||
@@ -136,16 +137,20 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
|
|
||||||
# Initialize IGDB if enabled
|
# Initialize IGDB if enabled
|
||||||
self.igdb = None
|
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("IGDB Configuration Check:")
|
||||||
logging.info(f" enable-game-info: {self.config.get('enable-game-info', 'NOT SET')}")
|
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-client-id: {'SET' if 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-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:
|
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("✅ 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())}")
|
logging.info(f" Available functions: {len(self.igdb.get_openai_functions())}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error(f"❌ Failed to initialize IGDB: {e}")
|
logging.error(f"❌ Failed to initialize IGDB: {e}")
|
||||||
@@ -153,6 +158,33 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
else:
|
else:
|
||||||
logging.warning("❌ IGDB integration DISABLED - missing configuration or disabled in config")
|
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]:
|
async def draw_openai(self, description: str, count: int = 1) -> List[BytesIO]:
|
||||||
if not self.ledger.budget_ok():
|
if not self.ledger.budget_ok():
|
||||||
raise RuntimeError("daily budget exhausted - refusing image call")
|
raise RuntimeError("daily budget exhausted - refusing image call")
|
||||||
@@ -241,24 +273,13 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
# hashed, never the raw Discord name (SAF-10)
|
# hashed, never the raw Discord name (SAF-10)
|
||||||
chat_kwargs["safety_identifier"] = "discord-" + hashlib.sha256(author.encode()).hexdigest()[:16]
|
chat_kwargs["safety_identifier"] = "discord-" + hashlib.sha256(author.encode()).hexdigest()[:16]
|
||||||
|
|
||||||
if self.igdb and self.config.get("enable-game-info", False):
|
available_tools = self._available_tools()
|
||||||
try:
|
if available_tools:
|
||||||
igdb_functions = self.igdb.get_openai_functions()
|
chat_kwargs["tools"] = [{"type": "function", "function": func} for func in available_tools]
|
||||||
if igdb_functions and isinstance(igdb_functions, list):
|
chat_kwargs["tool_choice"] = "auto"
|
||||||
chat_kwargs["tools"] = [{"type": "function", "function": func} for func in igdb_functions]
|
# gpt-5.6 rejects tools + reasoning on chat/completions (ENV-21)
|
||||||
chat_kwargs["tool_choice"] = "auto"
|
chat_kwargs["reasoning_effort"] = self.config.get("reasoning-effort", "none")
|
||||||
# gpt-5.6 rejects tools + reasoning on chat/completions (ENV-21)
|
logging.info(f"🔧 Tools available to AI: {[func['name'] for func in available_tools]}")
|
||||||
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)
|
|
||||||
)
|
|
||||||
)
|
|
||||||
|
|
||||||
result = await openai_chat(self.client, **chat_kwargs)
|
result = await openai_chat(self.client, **chat_kwargs)
|
||||||
self._record_usage(result)
|
self._record_usage(result)
|
||||||
@@ -278,10 +299,8 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
tool_names = [tc.function.name for tc in message.tool_calls]
|
tool_names = [tc.function.name for tc in message.tool_calls]
|
||||||
logging.info(f"🔧 OpenAI requested function calls: {tool_names}")
|
logging.info(f"🔧 OpenAI requested function calls: {tool_names}")
|
||||||
|
|
||||||
# Check if we have function/tool calls and IGDB is enabled
|
# Any offered tool may have been called (IGDB or fetch_url)
|
||||||
has_tool_calls = (
|
has_tool_calls = bool(hasattr(message, "tool_calls") and message.tool_calls and available_tools)
|
||||||
hasattr(message, "tool_calls") and message.tool_calls and self.igdb and self.config.get("enable-game-info", False)
|
|
||||||
)
|
|
||||||
|
|
||||||
# Clean up any existing tool messages in the history to avoid conflicts
|
# Clean up any existing tool messages in the history to avoid conflicts
|
||||||
if has_tool_calls:
|
if has_tool_calls:
|
||||||
@@ -304,12 +323,12 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
function_name = tool_call.function.name
|
function_name = tool_call.function.name
|
||||||
function_args = json.loads(tool_call.function.arguments)
|
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
|
# Route to the right provider (IGDB or URL reader)
|
||||||
function_result = await self._execute_igdb_function(function_name, function_args)
|
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(
|
messages.append(
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -0,0 +1,183 @@
|
|||||||
|
"""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"],
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
_META_REFRESH_URL = re.compile(r"url\s*=\s*['\"]?([^'\";\s]+)", re.I)
|
||||||
|
|
||||||
|
|
||||||
|
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
|
||||||
|
self.refresh_url: 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"]
|
||||||
|
# meta-refresh redirect (link shorteners, getnews stubs) — URL-04
|
||||||
|
content = attr.get("content")
|
||||||
|
if tag == "meta" and (attr.get("http-equiv") or "").lower() == "refresh" and content:
|
||||||
|
match = _META_REFRESH_URL.search(content)
|
||||||
|
if match and self.refresh_url is None:
|
||||||
|
self.refresh_url = match.group(1)
|
||||||
|
|
||||||
|
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)
|
||||||
|
# follow a meta-refresh redirect (link shorteners / getnews stubs), re-guarded — URL-04
|
||||||
|
for _ in range(2):
|
||||||
|
extractor = self._extract(body.decode("utf-8", "ignore"))
|
||||||
|
if not extractor.refresh_url:
|
||||||
|
break
|
||||||
|
target = urljoin(final_url, extractor.refresh_url)
|
||||||
|
if guard_url(target) is not None or target == final_url:
|
||||||
|
break
|
||||||
|
logging.info(f"url reader: following meta-refresh -> {target}")
|
||||||
|
final_url, body = await self._get(session, target, max_bytes)
|
||||||
|
except Exception as err:
|
||||||
|
return {"error": str(err)}
|
||||||
|
html = body.decode("utf-8", "ignore")
|
||||||
|
clean = sanitize_external_text(self._to_text(html), int(config.get("url-max-chars", DEFAULT_MAX_CHARS)))
|
||||||
|
images = await self._ingest_images(html, final_url, channel, user)
|
||||||
|
return {"url": final_url, "text": clean, "images_cached": images}
|
||||||
|
|
||||||
|
def _extract(self, html: str) -> "_Extractor":
|
||||||
|
extractor = _Extractor()
|
||||||
|
try:
|
||||||
|
extractor.feed(html)
|
||||||
|
except Exception as err:
|
||||||
|
logging.debug(f"html parse failed: {err!r}")
|
||||||
|
return extractor
|
||||||
|
|
||||||
|
def _to_text(self, html: str) -> str:
|
||||||
|
return re.sub(r"\s+\n", "\n", " ".join(self._extract(html).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 = self._extract(html)
|
||||||
|
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
|
||||||
@@ -12,3 +12,4 @@ with date + result.
|
|||||||
| DEP-04 | 2026-07-13 | Smoke gate exercised on ggg: RUNNING + fresh login line. |
|
| 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-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. |
|
| 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). |
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ dependencies = [
|
|||||||
"tomlkit>=0.13",
|
"tomlkit>=0.13",
|
||||||
"watchdog>=6",
|
"watchdog>=6",
|
||||||
"requests>=2.32",
|
"requests>=2.32",
|
||||||
|
"defusedxml>=0.7",
|
||||||
]
|
]
|
||||||
|
|
||||||
[project.scripts]
|
[project.scripts]
|
||||||
|
|||||||
@@ -0,0 +1,60 @@
|
|||||||
|
# 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. **HTML
|
||||||
|
meta-refresh** redirects (link shorteners, the old getnews stubs) are
|
||||||
|
also followed — the target is SSRF-re-guarded and fetched, so the
|
||||||
|
reader returns the real article, not the "Redirecting…" stub.
|
||||||
|
|
||||||
|
### 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.
|
||||||
@@ -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.
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
# 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.
|
||||||
|
|
||||||
|
## Webhook posting (ggg model)
|
||||||
|
|
||||||
|
`--post` mode fetches feeds mapped to channels and posts NEW items to
|
||||||
|
the channel's Discord webhook — replacing the py3.8 `getnews.py`
|
||||||
|
(dead play3 feed, 35 MB substring-scan state file, HTML-redirect
|
||||||
|
cruft). Config: `news-post-feeds = [[url, label, channel], …]`,
|
||||||
|
`news-post-webhooks = {channel = url}`, `news-post-state`.
|
||||||
|
|
||||||
|
### NEWS-04 — Only unseen items post, then are marked seen (coverage: test)
|
||||||
|
|
||||||
|
`NewsPoster.run_post` posts each item whose key (link, else title) is
|
||||||
|
not in the seen-set, adds it to the set, and posts to the mapped
|
||||||
|
channel's webhook. Re-runs over the same feed post nothing new.
|
||||||
|
|
||||||
|
### NEWS-05 — First run seeds without flooding (coverage: test)
|
||||||
|
|
||||||
|
With no prior state file (`seed_only`), every current item is marked
|
||||||
|
seen but nothing is posted — migrating off getnews.py never dumps a
|
||||||
|
backlog into the channels. `news-post-max-per-run` caps steady-state
|
||||||
|
posts per run.
|
||||||
|
|
||||||
|
### NEWS-06 — Post failures and bad channels are survived (coverage: test)
|
||||||
|
|
||||||
|
A feed the SSRF guard rejects, a feed that fails to fetch, an item
|
||||||
|
whose channel has no configured webhook, and a webhook POST that
|
||||||
|
raises are each logged and skipped — one failure never sinks the
|
||||||
|
run, and the seen-set still advances for successfully-processed
|
||||||
|
items.
|
||||||
@@ -28,7 +28,7 @@ class TestIGDBIntegration(unittest.IsolatedAsyncioTestCase):
|
|||||||
|
|
||||||
responder = OpenAIResponder(self.config_with_igdb)
|
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)
|
self.assertEqual(responder.igdb, mock_igdb_instance)
|
||||||
|
|
||||||
def test_igdb_initialization_disabled(self):
|
def test_igdb_initialization_disabled(self):
|
||||||
|
|||||||
+100
-1
@@ -160,7 +160,7 @@ class TestIGDBQuery(unittest.TestCase):
|
|||||||
"id",
|
"id",
|
||||||
"name",
|
"name",
|
||||||
"alternative_names",
|
"alternative_names",
|
||||||
"category",
|
"game_type",
|
||||||
"release_dates",
|
"release_dates",
|
||||||
"franchise",
|
"franchise",
|
||||||
"language_supports",
|
"language_supports",
|
||||||
@@ -174,5 +174,104 @@ class TestIGDBQuery(unittest.TestCase):
|
|||||||
self.assertEqual(result, [{"id": 1, "name": "Super Mario Bros"}])
|
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__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
@@ -45,6 +45,45 @@ class TestIngest(unittest.TestCase):
|
|||||||
self.assertIsNone(cache.ingest_bytes(PNG, "chat", "alice", "1"))
|
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):
|
class TestVisionDataUrls(OpsBase):
|
||||||
async def test_attachment_becomes_data_url(self):
|
async def test_attachment_becomes_data_url(self):
|
||||||
"""IMG-11: the model sees a data: URL, never the CDN link."""
|
"""IMG-11: the model sees a data: URL, never the CDN link."""
|
||||||
|
|||||||
@@ -0,0 +1,175 @@
|
|||||||
|
"""Unit coverage for SPEC-013 news digest (NEWS-01..03)."""
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from unittest.mock import AsyncMock
|
||||||
|
|
||||||
|
from fjerkroa_bot.news import NewsFetcher, NewsPoster, load_seen, parse_feed, render_digest, save_seen
|
||||||
|
|
||||||
|
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)
|
||||||
|
|
||||||
|
|
||||||
|
class TestPoster(unittest.IsolatedAsyncioTestCase):
|
||||||
|
def poster(self, posts):
|
||||||
|
async def fetch(url):
|
||||||
|
return RSS
|
||||||
|
|
||||||
|
async def post(hook, content):
|
||||||
|
posts.append((hook, content))
|
||||||
|
|
||||||
|
return NewsPoster(lambda u: None, fetch, post)
|
||||||
|
|
||||||
|
async def test_posts_unseen_then_dedups(self):
|
||||||
|
"""NEWS-04: unseen items post to the mapped webhook; re-run posts nothing."""
|
||||||
|
posts = []
|
||||||
|
poster = self.poster(posts)
|
||||||
|
feeds = [("https://a.com/feed", "PS", "news")]
|
||||||
|
hooks = {"news": "https://discord.com/api/webhooks/x"}
|
||||||
|
posted, seen = await poster.run_post(feeds, hooks, set(), per_feed=5, max_per_run=8, seed_only=False)
|
||||||
|
self.assertEqual(posted, 2)
|
||||||
|
self.assertIn("PS", posts[0][1])
|
||||||
|
self.assertIn("https://discord.com/api/webhooks/x", posts[0][0])
|
||||||
|
# re-run with the accumulated seen -> nothing new
|
||||||
|
posts.clear()
|
||||||
|
posted2, _ = await poster.run_post(feeds, hooks, seen, per_feed=5, max_per_run=8, seed_only=False)
|
||||||
|
self.assertEqual(posted2, 0)
|
||||||
|
self.assertEqual(posts, [])
|
||||||
|
|
||||||
|
async def test_seed_run_posts_nothing(self):
|
||||||
|
"""NEWS-05: seed_only marks items seen without posting."""
|
||||||
|
posts = []
|
||||||
|
poster = self.poster(posts)
|
||||||
|
feeds = [("https://a.com/feed", "PS", "news")]
|
||||||
|
posted, seen = await poster.run_post(feeds, {"news": "h"}, set(), 5, 8, seed_only=True)
|
||||||
|
self.assertEqual(posted, 0)
|
||||||
|
self.assertEqual(posts, [])
|
||||||
|
self.assertEqual(len(seen), 2) # both marked seen
|
||||||
|
|
||||||
|
async def test_max_per_run_caps(self):
|
||||||
|
"""NEWS-05: max-per-run caps posts; extras stay seen (not re-posted next run)."""
|
||||||
|
posts = []
|
||||||
|
poster = self.poster(posts)
|
||||||
|
feeds = [("https://a.com/feed", "PS", "news")]
|
||||||
|
posted, seen = await poster.run_post(feeds, {"news": "h"}, set(), per_feed=5, max_per_run=1, seed_only=False)
|
||||||
|
self.assertEqual(posted, 1)
|
||||||
|
self.assertEqual(len(seen), 2) # both seen, only one posted
|
||||||
|
|
||||||
|
async def test_failures_survived(self):
|
||||||
|
"""NEWS-06: SSRF-skip, fetch fail, missing webhook, post error each survive."""
|
||||||
|
posts = []
|
||||||
|
|
||||||
|
async def fetch(url):
|
||||||
|
if "boom" in url:
|
||||||
|
raise ValueError("boom")
|
||||||
|
return RSS
|
||||||
|
|
||||||
|
async def post(hook, content):
|
||||||
|
if hook == "bad":
|
||||||
|
raise RuntimeError("post failed")
|
||||||
|
posts.append((hook, content))
|
||||||
|
|
||||||
|
def guard(url):
|
||||||
|
return "refused" if "internal" in url else None
|
||||||
|
|
||||||
|
poster = NewsPoster(guard, fetch, post)
|
||||||
|
feeds = [
|
||||||
|
("https://internal/feed", "I", "news"), # SSRF-skipped
|
||||||
|
("https://boom.com/feed", "B", "news"), # fetch fails
|
||||||
|
("https://ok.com/feed", "OK", "nowhere"), # no webhook for channel
|
||||||
|
("https://ok2.com/feed", "OK2", "news"), # webhook raises
|
||||||
|
]
|
||||||
|
posted, seen = await poster.run_post(feeds, {"news": "bad"}, set(), 5, 8, seed_only=False)
|
||||||
|
self.assertEqual(posted, 0) # everything failed/skipped, no crash
|
||||||
|
|
||||||
|
|
||||||
|
class TestSeenState(unittest.TestCase):
|
||||||
|
def test_roundtrip_and_seed_detection(self):
|
||||||
|
"""NEWS-05: missing state -> (empty, existed=False); saved state reloads."""
|
||||||
|
import tempfile
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
|
path = str(Path(tmp) / "state.json")
|
||||||
|
seen, existed = load_seen(path)
|
||||||
|
self.assertEqual((seen, existed), (set(), False))
|
||||||
|
save_seen(path, {"a", "b", "c"}, cap=5000)
|
||||||
|
reloaded, existed2 = load_seen(path)
|
||||||
|
self.assertEqual(reloaded, {"a", "b", "c"})
|
||||||
|
self.assertTrue(existed2)
|
||||||
|
|
||||||
|
def test_cap_bounds_state(self):
|
||||||
|
"""NEWS-05: save keeps at most `cap` keys."""
|
||||||
|
import json
|
||||||
|
import tempfile
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
|
path = str(Path(tmp) / "state.json")
|
||||||
|
save_seen(path, {f"k{i}" for i in range(100)}, cap=10)
|
||||||
|
self.assertEqual(len(json.load(open(path))), 10)
|
||||||
@@ -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)
|
||||||
@@ -0,0 +1,242 @@
|
|||||||
|
"""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 TestMetaRefresh(unittest.IsolatedAsyncioTestCase):
|
||||||
|
async def test_follows_meta_refresh_to_real_article(self):
|
||||||
|
"""URL-04: a getnews-style meta-refresh stub is followed to the real article."""
|
||||||
|
reader = URLReader(lambda: {}, None)
|
||||||
|
stub = (
|
||||||
|
b'<html><head><meta http-equiv="refresh" content="0;url=https://pushsquare.com/real"></head><body>Redirecting...</body></html>'
|
||||||
|
)
|
||||||
|
article = b"<html><body><h1>MARVEL Tokon</h1><p>Full article text here</p></body></html>"
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
async def fake_get(session, url, max_bytes):
|
||||||
|
calls.append(url)
|
||||||
|
return (url, stub if "stub" in url else article)
|
||||||
|
|
||||||
|
reader._get = fake_get # type: ignore
|
||||||
|
with patch("fjerkroa_bot.url_reader.guard_url", return_value=None):
|
||||||
|
import fjerkroa_bot.url_reader as ur
|
||||||
|
|
||||||
|
# patch the session context so fetch() runs against fake_get
|
||||||
|
class FakeCM:
|
||||||
|
async def __aenter__(self):
|
||||||
|
return object()
|
||||||
|
|
||||||
|
async def __aexit__(self, *a):
|
||||||
|
return False
|
||||||
|
|
||||||
|
with patch.object(ur.aiohttp, "ClientSession", return_value=FakeCM()):
|
||||||
|
result = await reader.fetch("https://gggemein.de/url/stub.html", "chat", "alice")
|
||||||
|
self.assertIn("Full article text", result["text"])
|
||||||
|
self.assertEqual(result["url"], "https://pushsquare.com/real")
|
||||||
|
self.assertIn("https://pushsquare.com/real", calls)
|
||||||
|
|
||||||
|
async def test_meta_refresh_to_internal_is_not_followed(self):
|
||||||
|
"""URL-04: a meta-refresh pointing at an internal IP is refused (SSRF)."""
|
||||||
|
reader = URLReader(lambda: {}, None)
|
||||||
|
stub = b'<meta http-equiv="refresh" content="0; url=http://127.0.0.1/secret">Redirecting'
|
||||||
|
|
||||||
|
async def fake_get(session, url, max_bytes):
|
||||||
|
return (url, stub)
|
||||||
|
|
||||||
|
reader._get = fake_get # type: ignore
|
||||||
|
import fjerkroa_bot.url_reader as ur
|
||||||
|
|
||||||
|
class FakeCM:
|
||||||
|
async def __aenter__(self):
|
||||||
|
return object()
|
||||||
|
|
||||||
|
async def __aexit__(self, *a):
|
||||||
|
return False
|
||||||
|
|
||||||
|
def guard(u):
|
||||||
|
return "refused" if "127.0.0.1" in u else None
|
||||||
|
|
||||||
|
with patch("fjerkroa_bot.url_reader.guard_url", side_effect=guard):
|
||||||
|
with patch.object(ur.aiohttp, "ClientSession", return_value=FakeCM()):
|
||||||
|
result = await reader.fetch("https://safe.com/x", "chat", "alice")
|
||||||
|
self.assertEqual(result["url"], "https://safe.com/x") # did not follow to 127.0.0.1
|
||||||
|
|
||||||
|
|
||||||
|
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)
|
||||||
@@ -627,6 +627,7 @@ version = "3.0.0"
|
|||||||
source = { editable = "." }
|
source = { editable = "." }
|
||||||
dependencies = [
|
dependencies = [
|
||||||
{ name = "aiohttp" },
|
{ name = "aiohttp" },
|
||||||
|
{ name = "defusedxml" },
|
||||||
{ name = "discord-py" },
|
{ name = "discord-py" },
|
||||||
{ name = "openai" },
|
{ name = "openai" },
|
||||||
{ name = "requests" },
|
{ name = "requests" },
|
||||||
@@ -656,6 +657,7 @@ dev = [
|
|||||||
[package.metadata]
|
[package.metadata]
|
||||||
requires-dist = [
|
requires-dist = [
|
||||||
{ name = "aiohttp", specifier = ">=3.12" },
|
{ name = "aiohttp", specifier = ">=3.12" },
|
||||||
|
{ name = "defusedxml", specifier = ">=0.7" },
|
||||||
{ name = "discord-py", specifier = ">=2.5,<3" },
|
{ name = "discord-py", specifier = ">=2.5,<3" },
|
||||||
{ name = "openai", specifier = ">=2.45" },
|
{ name = "openai", specifier = ">=2.45" },
|
||||||
{ name = "requests", specifier = ">=2.32" },
|
{ name = "requests", specifier = ">=2.32" },
|
||||||
|
|||||||
Reference in New Issue
Block a user