Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4166520923 | |||
| d0819c2683 | |||
| df1924bb80 |
@@ -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)
|
||||||
|
|||||||
+17
@@ -59,3 +59,20 @@ enable-game-info = true
|
|||||||
# image-cache-mb = 500 # LRU cap (ggg: consider 2000 — screenshots)
|
# image-cache-mb = 500 # LRU cap (ggg: consider 2000 — screenshots)
|
||||||
# image-cache-ttl-days = 90
|
# image-cache-ttl-days = 90
|
||||||
# image-max-bytes = 8388608 # 8 MB upload cap
|
# image-max-bytes = 8388608 # 8 MB upload cap
|
||||||
|
# Self-tasking (SPEC-005) — experimental, DEFAULT OFF:
|
||||||
|
# tasks-enabled = true
|
||||||
|
# tasks-generators = ["idle-impulse", "follow-up"]
|
||||||
|
# tasks-max-per-channel-per-day = 2
|
||||||
|
# tasks-approval = false # true: neue Tasks brauchen !bot task-approve
|
||||||
|
# idle-impulse-hours = 12
|
||||||
|
# taskgen-interval-hours = 6
|
||||||
|
# Staff: !bot tasks | task-approve <id> | task-cancel <id>
|
||||||
|
# URL reading (SPEC-011, FDB-018) — DEFAULT OFF; web pages are hostile input:
|
||||||
|
# enable-url-reading = true
|
||||||
|
# url-max-bytes = 2097152 # 2 MB fetch cap
|
||||||
|
# url-max-chars = 6000 # text handed to the model
|
||||||
|
# url-max-images = 2 # page images into the vision cache
|
||||||
|
# url-daily-per-user = 20
|
||||||
|
# Ops (SPEC-012): consecutive OpenAI failures before a staff alert
|
||||||
|
# api-error-alert-threshold = 5
|
||||||
|
# Backups: cron runs deploy/backup_db.py daily -> ~/backups/<bot>/ (keep 14)
|
||||||
|
|||||||
@@ -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())
|
||||||
+71
-33
@@ -1,7 +1,6 @@
|
|||||||
import argparse
|
import argparse
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
import math
|
|
||||||
import random
|
import random
|
||||||
import re
|
import re
|
||||||
import sys
|
import sys
|
||||||
@@ -18,6 +17,7 @@ from watchdog.observers import Observer
|
|||||||
|
|
||||||
from .ai_responder import AIMessage
|
from .ai_responder import AIMessage
|
||||||
from .openai_responder import OpenAIResponder
|
from .openai_responder import OpenAIResponder
|
||||||
|
from .tasks import TaskEngine
|
||||||
|
|
||||||
DEFAULT_PRIVACY_NOTICE = (
|
DEFAULT_PRIVACY_NOTICE = (
|
||||||
"I keep recent channel messages and a short conversation summary to answer better. "
|
"I keep recent channel messages and a short conversation summary to answer better. "
|
||||||
@@ -89,6 +89,7 @@ class FjerkroaBot(commands.Bot):
|
|||||||
self.tasks_enabled = True
|
self.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()
|
||||||
@@ -114,38 +115,42 @@ class FjerkroaBot(commands.Bot):
|
|||||||
self.staff_channel = self.channel_by_name(self.config["staff-channel"], no_ignore=True)
|
self.staff_channel = self.channel_by_name(self.config["staff-channel"], no_ignore=True)
|
||||||
self.welcome_channel = self.channel_by_name(self.config["welcome-channel"], no_ignore=True)
|
self.welcome_channel = self.channel_by_name(self.config["welcome-channel"], no_ignore=True)
|
||||||
|
|
||||||
def init_boreness(self):
|
def init_tasks(self):
|
||||||
if "chat-channel" not in self.config:
|
"""Task engine replaces the sigmoid boreness loop (TSK-07)."""
|
||||||
return
|
|
||||||
self.last_activity_time = time.monotonic()
|
self.last_activity_time = time.monotonic()
|
||||||
self.loop.create_task(self.on_boreness())
|
self.task_engine = TaskEngine(
|
||||||
logging.info("Boreness initialised.")
|
store=self.airesponder.store,
|
||||||
|
ledger=self.airesponder.ledger,
|
||||||
|
config_getter=lambda: self.config,
|
||||||
|
execute=self._execute_task,
|
||||||
|
propose=self.airesponder.propose_task,
|
||||||
|
staff_alert=self.send_staff_alert,
|
||||||
|
allowed=self.bot_initiated_allowed,
|
||||||
|
idle_seconds=lambda: time.monotonic() - self.last_activity_time,
|
||||||
|
observe=self.airesponder.observe_event,
|
||||||
|
)
|
||||||
|
self.loop.create_task(self.task_loop())
|
||||||
|
logging.info("Task engine initialised.")
|
||||||
|
|
||||||
async def on_boreness(self):
|
async def task_loop(self):
|
||||||
logging.info(f"Boreness started on channel: {repr(self.chat_channel)}")
|
|
||||||
while True:
|
while True:
|
||||||
if self.chat_channel is None or not self.bot_initiated_allowed():
|
await asyncio.sleep(60)
|
||||||
await asyncio.sleep(7)
|
try:
|
||||||
continue
|
await self.task_engine.tick()
|
||||||
boreness_interval = float(self.config.get("boreness-interval", 12.0))
|
except Exception as err:
|
||||||
elapsed_time = (time.monotonic() - self.last_activity_time) / 3600.0
|
logging.warning(f"task tick failed: {repr(err)}")
|
||||||
probability = 1 / (1 + math.exp(-1 * (elapsed_time - (boreness_interval / 2.0)) + math.log(1 / 0.2 - 1)))
|
|
||||||
if random.random() < probability:
|
async def _execute_task(self, channel_name: str, prompt: str) -> None:
|
||||||
prev_messages = [msg async for msg in self.chat_channel.history(limit=2)]
|
"""Run a due task through the normal responder path (TSK-02)."""
|
||||||
last_author = prev_messages[1].author.id if len(prev_messages) > 1 else None
|
channel = self.channel_by_name(channel_name, getattr(self, "chat_channel", None), no_ignore=True)
|
||||||
if last_author and last_author != self.user.id:
|
if channel is None:
|
||||||
logging.info(f"Borred with {probability} probability after {elapsed_time}")
|
raise RuntimeError(f"task channel {channel_name!r} not resolvable")
|
||||||
boreness_prompt = self.config.get("boreness-prompt", "Pretend that you just now thought of something, be creative.")
|
message = AIMessage("system", prompt, channel_name, True, False)
|
||||||
message = AIMessage("system", boreness_prompt, self.config.get("chat-channel", "chat"), True, False)
|
await self.respond(message, channel)
|
||||||
try:
|
|
||||||
await self.respond(message, self.chat_channel)
|
|
||||||
except Exception as err:
|
|
||||||
logging.warning(f"Failed to activate borringness: {repr(err)}")
|
|
||||||
await asyncio.sleep(7)
|
|
||||||
|
|
||||||
async def on_ready(self):
|
async def on_ready(self):
|
||||||
self.init_channels()
|
self.init_channels()
|
||||||
self.init_boreness()
|
self.init_tasks()
|
||||||
logging.info(
|
logging.info(
|
||||||
f"We have logged in as {self.user}" f" ({repr(self.staff_channel)}, {repr(self.welcome_channel)}, {repr(self.chat_channel)})"
|
f"We have logged in as {self.user}" f" ({repr(self.staff_channel)}, {repr(self.welcome_channel)}, {repr(self.chat_channel)})"
|
||||||
)
|
)
|
||||||
@@ -241,13 +246,31 @@ class FjerkroaBot(commands.Bot):
|
|||||||
return "\n".join(f"{pin['id']} [{pin['channel'] or 'global'}]: {pin['fact']}" for pin in pins) or "No pins."
|
return "\n".join(f"{pin['id']} [{pin['channel'] or 'global'}]: {pin['fact']}" for pin in pins) or "No pins."
|
||||||
return None
|
return None
|
||||||
|
|
||||||
|
def _task_command(self, args) -> Optional[str]:
|
||||||
|
"""Task queue surface (OPS-12, TSK-05)."""
|
||||||
|
is_list = args[:1] == ["tasks"] and len(args) == 1
|
||||||
|
if args[:1] not in (["task-approve"], ["task-cancel"]) and not is_list:
|
||||||
|
return None
|
||||||
|
store = self.airesponder.store
|
||||||
|
if store is None:
|
||||||
|
return "No store configured - task commands unavailable."
|
||||||
|
if is_list:
|
||||||
|
tasks = store.tasks_open()
|
||||||
|
return "\n".join(f"{t['id']} [{t['state']}] {t['kind']} #{t['channel']} due {t['due_at']}" for t in tasks) or "No open tasks."
|
||||||
|
if args[1:2] and args[1].isdigit():
|
||||||
|
if args[0] == "task-approve":
|
||||||
|
return f"Approved {store.task_set_state(int(args[1]), 'queued')} task(s)."
|
||||||
|
return f"Cancelled {store.task_set_state(int(args[1]), 'cancelled')} task(s)."
|
||||||
|
return None
|
||||||
|
|
||||||
async def handle_staff_command(self, message: Message) -> None:
|
async def handle_staff_command(self, message: Message) -> None:
|
||||||
"""Operator kill-switches, staff channel only (OPS-01..05, OPS-09, MEM-07)."""
|
"""Operator kill-switches, staff channel only (OPS-01..05, OPS-09, MEM-07)."""
|
||||||
args = str(message.content).split()[1:]
|
args = str(message.content).split()[1:]
|
||||||
memory_reply = self._memory_command(args)
|
for handler in (self._memory_command, self._task_command):
|
||||||
if memory_reply is not None:
|
reply = handler(args)
|
||||||
await message.channel.send(memory_reply, suppress_embeds=True)
|
if reply is not None:
|
||||||
return
|
await message.channel.send(reply, suppress_embeds=True)
|
||||||
|
return
|
||||||
reply = "Commands: pause, resume, images on|off, tasks on|off, quiet <minutes>, status, spend, memory <user>, forget-fact <id>, pin <channel|global> <fact>, unpin <id>"
|
reply = "Commands: pause, resume, images on|off, tasks on|off, quiet <minutes>, status, spend, memory <user>, forget-fact <id>, pin <channel|global> <fact>, unpin <id>"
|
||||||
if args[:1] == ["pause"]:
|
if args[:1] == ["pause"]:
|
||||||
self.replies_enabled = False
|
self.replies_enabled = False
|
||||||
@@ -476,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():
|
||||||
@@ -581,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)
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -58,6 +59,31 @@ CONSOLIDATION_RESPONSE_FORMAT = {
|
|||||||
"type": "json_schema",
|
"type": "json_schema",
|
||||||
"json_schema": {"name": "consolidation", "strict": True, "schema": CONSOLIDATION_SCHEMA},
|
"json_schema": {"name": "consolidation", "strict": True, "schema": CONSOLIDATION_SCHEMA},
|
||||||
}
|
}
|
||||||
|
# Follow-up task proposal (SPEC-005 TSK-08): one task or null
|
||||||
|
TASKGEN_SCHEMA = {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"task": {
|
||||||
|
"type": ["object", "null"],
|
||||||
|
"properties": {
|
||||||
|
"channel": {"type": ["string", "null"], "description": "Target channel, or null for the main chat channel."},
|
||||||
|
"prompt": {"type": "string", "description": "Instruction the assistant will act on when the task runs."},
|
||||||
|
"due_hours": {"type": "number", "description": "Hours from now until the task should run (0 = now)."},
|
||||||
|
},
|
||||||
|
"required": ["channel", "prompt", "due_hours"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"required": ["task"],
|
||||||
|
"additionalProperties": False,
|
||||||
|
}
|
||||||
|
TASKGEN_RESPONSE_FORMAT = {"type": "json_schema", "json_schema": {"name": "task_proposal", "strict": True, "schema": TASKGEN_SCHEMA}}
|
||||||
|
TASKGEN_SYSTEM = (
|
||||||
|
"You plan the self-initiated actions of a Discord assistant. Given recent conversation summaries, propose AT MOST ONE follow-up"
|
||||||
|
" worth doing on the assistant's own initiative (ask how something announced went, revisit an open question, congratulate on an"
|
||||||
|
" event). Only propose something genuinely worthwhile — when in doubt, return a null task."
|
||||||
|
)
|
||||||
|
|
||||||
# Reply/ignore + factual pre-pass (SPEC-010 BEH-01): one cheap call
|
# Reply/ignore + factual pre-pass (SPEC-010 BEH-01): one cheap call
|
||||||
CLASSIFIER_SCHEMA = {
|
CLASSIFIER_SCHEMA = {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
@@ -128,6 +154,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")
|
||||||
@@ -216,24 +269,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)
|
||||||
@@ -253,10 +295,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:
|
||||||
@@ -279,12 +319,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(
|
||||||
{
|
{
|
||||||
@@ -381,6 +421,29 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
|||||||
logging.info(f"edited {len(buffers)} image(s) on {model} from {len(handles)} input(s)")
|
logging.info(f"edited {len(buffers)} image(s) on {model} from {len(handles)} input(s)")
|
||||||
return buffers
|
return buffers
|
||||||
|
|
||||||
|
async def propose_task(self) -> Optional[Dict[str, Any]]:
|
||||||
|
"""One follow-up proposal from recent episodes on memory-model (TSK-08)."""
|
||||||
|
if "memory-model" not in self.config or self.store is None or not self.ledger.budget_ok():
|
||||||
|
return None
|
||||||
|
channel = self.config.get("chat-channel", "chat")
|
||||||
|
episodes = await asyncio.to_thread(self.store.recent_episodes, channel, 5)
|
||||||
|
if not episodes:
|
||||||
|
return None
|
||||||
|
episode_lines = "\n".join(f"- {episode}" for episode in episodes)
|
||||||
|
messages = [
|
||||||
|
{"role": "system", "content": TASKGEN_SYSTEM},
|
||||||
|
{"role": "user", "content": f"Recent conversation summaries in #{channel}:\n{episode_lines}"},
|
||||||
|
]
|
||||||
|
try:
|
||||||
|
result = await openai_chat(
|
||||||
|
self.client, model=self.config["memory-model"], messages=messages, response_format=TASKGEN_RESPONSE_FORMAT
|
||||||
|
)
|
||||||
|
self._record_usage(result)
|
||||||
|
return json.loads(result.choices[0].message.content)
|
||||||
|
except Exception as err:
|
||||||
|
logging.warning(f"task proposal failed: {repr(err)}")
|
||||||
|
return None
|
||||||
|
|
||||||
async def classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
|
async def classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
|
||||||
"""~100-token reply/factual/emoji verdict on classifier-model (BEH-01/03)."""
|
"""~100-token reply/factual/emoji verdict on classifier-model (BEH-01/03)."""
|
||||||
if "classifier-model" not in self.config or not self.ledger.budget_ok():
|
if "classifier-model" not in self.config or not self.ledger.budget_ok():
|
||||||
|
|||||||
@@ -13,7 +13,7 @@ from contextlib import closing
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
SCHEMA_VERSION = 4
|
SCHEMA_VERSION = 5
|
||||||
|
|
||||||
|
|
||||||
class PersistentStore:
|
class PersistentStore:
|
||||||
@@ -66,6 +66,12 @@ class PersistentStore:
|
|||||||
" user TEXT NOT NULL, message_id TEXT, ext TEXT NOT NULL, bytes INTEGER NOT NULL,"
|
" user TEXT NOT NULL, message_id TEXT, ext TEXT NOT NULL, bytes INTEGER NOT NULL,"
|
||||||
" created_at TEXT NOT NULL DEFAULT (datetime('now')))"
|
" created_at TEXT NOT NULL DEFAULT (datetime('now')))"
|
||||||
)
|
)
|
||||||
|
if version < 5:
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS tasks (id INTEGER PRIMARY KEY, kind TEXT NOT NULL, channel TEXT NOT NULL,"
|
||||||
|
" due_at TEXT NOT NULL, payload TEXT NOT NULL, state TEXT NOT NULL DEFAULT 'queued',"
|
||||||
|
" created_at TEXT NOT NULL DEFAULT (datetime('now')), executed_at TEXT)"
|
||||||
|
)
|
||||||
if version < SCHEMA_VERSION:
|
if version < SCHEMA_VERSION:
|
||||||
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
|
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
|
||||||
os.chmod(self.db_path, 0o600) # conversation data (PER-04)
|
os.chmod(self.db_path, 0o600) # conversation data (PER-04)
|
||||||
@@ -195,6 +201,39 @@ class PersistentStore:
|
|||||||
)
|
)
|
||||||
return cursor.rowcount
|
return cursor.rowcount
|
||||||
|
|
||||||
|
# --- task queue (SPEC-005, FDB-011) ---
|
||||||
|
|
||||||
|
def task_add(self, kind: str, channel: str, due_at: str, payload: str, state: str = "queued") -> int:
|
||||||
|
with closing(self._connect()) as conn, conn:
|
||||||
|
cursor = conn.execute(
|
||||||
|
"INSERT INTO tasks (kind, channel, due_at, payload, state) VALUES (?, ?, ?, ?, ?)",
|
||||||
|
(kind, channel, due_at, payload, state),
|
||||||
|
)
|
||||||
|
return int(cursor.lastrowid or 0)
|
||||||
|
|
||||||
|
def tasks_due(self) -> List[Dict[str, Any]]:
|
||||||
|
with closing(self._connect()) as conn:
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT id, kind, channel, payload FROM tasks WHERE state = 'queued' AND due_at <= datetime('now') ORDER BY due_at"
|
||||||
|
).fetchall()
|
||||||
|
return [{"id": row[0], "kind": row[1], "channel": row[2], "payload": row[3]} for row in rows]
|
||||||
|
|
||||||
|
def tasks_open(self) -> List[Dict[str, Any]]:
|
||||||
|
with closing(self._connect()) as conn:
|
||||||
|
rows = conn.execute(
|
||||||
|
"SELECT id, kind, channel, due_at, state FROM tasks WHERE state IN ('queued', 'approval') ORDER BY id"
|
||||||
|
).fetchall()
|
||||||
|
return [{"id": row[0], "kind": row[1], "channel": row[2], "due_at": row[3], "state": row[4]} for row in rows]
|
||||||
|
|
||||||
|
def task_set_state(self, task_id: int, state: str) -> int:
|
||||||
|
with closing(self._connect()) as conn, conn:
|
||||||
|
return conn.execute("UPDATE tasks SET state = ?, executed_at = datetime('now') WHERE id = ?", (state, task_id)).rowcount
|
||||||
|
|
||||||
|
def tasks_pending_of_kind(self, kind: str) -> int:
|
||||||
|
with closing(self._connect()) as conn:
|
||||||
|
row = conn.execute("SELECT COUNT(*) FROM tasks WHERE kind = ? AND state IN ('queued', 'approval')", (kind,)).fetchone()
|
||||||
|
return int(row[0])
|
||||||
|
|
||||||
# --- image cache index (SPEC-004, FDB-010) ---
|
# --- image cache index (SPEC-004, FDB-010) ---
|
||||||
|
|
||||||
def image_add(self, sha256: str, channel: str, user: str, message_id: Optional[str], ext: str, nbytes: int) -> None:
|
def image_add(self, sha256: str, channel: str, user: str, message_id: Optional[str], ext: str, nbytes: int) -> None:
|
||||||
|
|||||||
@@ -0,0 +1,135 @@
|
|||||||
|
"""Self-tasking engine (SPEC-005, FDB-011).
|
||||||
|
|
||||||
|
Generators propose, the scheduler executes — through the injected
|
||||||
|
execute callback (= the normal responder path), so budget, gates and
|
||||||
|
kill-switches all apply. Everything is off unless `tasks-enabled` is
|
||||||
|
true (TSK-03).
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import time
|
||||||
|
from typing import Any, Awaitable, Callable, Dict, List, Optional
|
||||||
|
|
||||||
|
from .persistence import PersistentStore
|
||||||
|
from .quota import QuotaLedger
|
||||||
|
|
||||||
|
DEFAULT_MAX_PER_CHANNEL_PER_DAY = 2
|
||||||
|
DEFAULT_IDLE_IMPULSE_HOURS = 12.0
|
||||||
|
DEFAULT_TASKGEN_INTERVAL_HOURS = 6.0
|
||||||
|
DEFAULT_BORENESS_PROMPT = "Pretend that you just now thought of something, be creative."
|
||||||
|
|
||||||
|
ExecuteCallback = Callable[[str, str], Awaitable[None]]
|
||||||
|
ProposeCallback = Callable[[], Awaitable[Optional[Dict[str, Any]]]]
|
||||||
|
AlertCallback = Callable[[str], Awaitable[None]]
|
||||||
|
|
||||||
|
|
||||||
|
class TaskEngine:
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
store: Optional[PersistentStore],
|
||||||
|
ledger: QuotaLedger,
|
||||||
|
config_getter: Callable[[], Dict[str, Any]],
|
||||||
|
execute: ExecuteCallback,
|
||||||
|
propose: ProposeCallback,
|
||||||
|
staff_alert: AlertCallback,
|
||||||
|
allowed: Callable[[], bool],
|
||||||
|
idle_seconds: Callable[[], float],
|
||||||
|
observe: Callable[[str, str, str], Awaitable[None]],
|
||||||
|
) -> None:
|
||||||
|
self.store = store
|
||||||
|
self.ledger = ledger
|
||||||
|
self._config = config_getter
|
||||||
|
self._execute = execute
|
||||||
|
self._propose = propose
|
||||||
|
self._staff_alert = staff_alert
|
||||||
|
self._allowed = allowed
|
||||||
|
self._idle_seconds = idle_seconds
|
||||||
|
self._observe = observe
|
||||||
|
self._last_generation = 0.0
|
||||||
|
self._lock = asyncio.Lock()
|
||||||
|
|
||||||
|
def active(self) -> bool:
|
||||||
|
return self.store is not None and bool(self._config().get("tasks-enabled", False)) # TSK-03
|
||||||
|
|
||||||
|
def generators(self) -> List[str]:
|
||||||
|
return list(self._config().get("tasks-generators", ["idle-impulse", "follow-up"]))
|
||||||
|
|
||||||
|
async def enqueue(self, kind: str, channel: str, payload: str, due_hours: float = 0.0) -> int:
|
||||||
|
assert self.store is not None
|
||||||
|
state = "approval" if self._config().get("tasks-approval", False) else "queued" # TSK-05
|
||||||
|
due_at = time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime(time.time() + due_hours * 3600.0))
|
||||||
|
task_id = int(await asyncio.to_thread(self.store.task_add, kind, channel, due_at, payload, state) or 0)
|
||||||
|
if state == "approval":
|
||||||
|
await self._staff_alert(f"Task #{task_id} proposed ({kind}, #{channel}): {payload[:180]} — !bot task-approve {task_id}")
|
||||||
|
return task_id
|
||||||
|
|
||||||
|
def _channel_cap_ok(self, channel: str) -> bool:
|
||||||
|
cap = int(self._config().get("tasks-max-per-channel-per-day", DEFAULT_MAX_PER_CHANNEL_PER_DAY))
|
||||||
|
return self.ledger._get(f"task-runs:{channel}") < cap # TSK-04
|
||||||
|
|
||||||
|
async def tick(self) -> None:
|
||||||
|
"""One scheduler pass: execute due tasks, then maybe generate (TSK-02/06)."""
|
||||||
|
if not self.active() or not self._allowed() or self._lock.locked():
|
||||||
|
return
|
||||||
|
assert self.store is not None
|
||||||
|
async with self._lock:
|
||||||
|
await self._maybe_generate()
|
||||||
|
for task in await asyncio.to_thread(self.store.tasks_due):
|
||||||
|
if not self._channel_cap_ok(task["channel"]):
|
||||||
|
continue # stays queued for tomorrow (TSK-04)
|
||||||
|
try:
|
||||||
|
await self._execute(task["channel"], task["payload"])
|
||||||
|
await asyncio.to_thread(self.store.task_set_state, task["id"], "done")
|
||||||
|
self.ledger._add(f"task-runs:{task['channel']}", 1)
|
||||||
|
await self._observe("system", "task", f"executed {task['kind']}: {task['payload'][:200]}")
|
||||||
|
except Exception as err:
|
||||||
|
logging.warning(f"task {task['id']} failed: {repr(err)}")
|
||||||
|
await asyncio.to_thread(self.store.task_set_state, task["id"], "failed")
|
||||||
|
|
||||||
|
async def _maybe_generate(self) -> None:
|
||||||
|
config = self._config()
|
||||||
|
interval = float(config.get("taskgen-interval-hours", DEFAULT_TASKGEN_INTERVAL_HOURS)) * 3600.0
|
||||||
|
now = time.monotonic()
|
||||||
|
if self._last_generation and now - self._last_generation < interval:
|
||||||
|
return
|
||||||
|
self._last_generation = now
|
||||||
|
generators = self.generators()
|
||||||
|
if "idle-impulse" in generators:
|
||||||
|
await self._generate_idle_impulse()
|
||||||
|
if "follow-up" in generators:
|
||||||
|
await self._generate_follow_up()
|
||||||
|
|
||||||
|
async def _generate_idle_impulse(self) -> None:
|
||||||
|
"""Boreness, demoted to a deterministic generator (TSK-07)."""
|
||||||
|
assert self.store is not None
|
||||||
|
config = self._config()
|
||||||
|
channel = config.get("chat-channel")
|
||||||
|
if not channel:
|
||||||
|
return
|
||||||
|
idle_threshold = float(config.get("idle-impulse-hours", DEFAULT_IDLE_IMPULSE_HOURS)) * 3600.0
|
||||||
|
if self._idle_seconds() < idle_threshold:
|
||||||
|
return
|
||||||
|
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "idle-impulse"):
|
||||||
|
return # one pending impulse is enough
|
||||||
|
prompt = config.get("boreness-prompt", DEFAULT_BORENESS_PROMPT)
|
||||||
|
await self.enqueue("idle-impulse", channel, prompt)
|
||||||
|
|
||||||
|
async def _generate_follow_up(self) -> None:
|
||||||
|
"""Ask the model for one follow-up worth doing (TSK-08)."""
|
||||||
|
assert self.store is not None
|
||||||
|
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "follow-up"):
|
||||||
|
return
|
||||||
|
proposal = await self._propose()
|
||||||
|
if not proposal or not isinstance(proposal.get("task"), dict):
|
||||||
|
return
|
||||||
|
task = proposal["task"]
|
||||||
|
channel = str(task.get("channel") or self._config().get("chat-channel") or "")
|
||||||
|
prompt = str(task.get("prompt") or "").strip()
|
||||||
|
if not channel or not prompt:
|
||||||
|
return
|
||||||
|
try:
|
||||||
|
due_hours = max(0.0, float(task.get("due_hours") or 0.0))
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
due_hours = 0.0
|
||||||
|
await self.enqueue("follow-up", channel, prompt, due_hours)
|
||||||
@@ -0,0 +1,163 @@
|
|||||||
|
"""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
|
||||||
|
|
||||||
|
DEFAULT_MAX_BYTES = 2 * 1024 * 1024
|
||||||
|
DEFAULT_MAX_CHARS = 6000
|
||||||
|
DEFAULT_MAX_IMAGES = 2
|
||||||
|
FETCH_TIMEOUT_S = 15
|
||||||
|
MAX_REDIRECTS = 5
|
||||||
|
|
||||||
|
FETCH_URL_TOOL = {
|
||||||
|
"name": "fetch_url",
|
||||||
|
"description": "Fetch a public web page and return its readable text plus prominent image links. "
|
||||||
|
"Use when the user shares a URL and asks about it, or to get details behind a news link.",
|
||||||
|
"parameters": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {"url": {"type": "string", "description": "The http/https URL to read."}},
|
||||||
|
"required": ["url"],
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
class _Extractor(HTMLParser):
|
||||||
|
def __init__(self) -> None:
|
||||||
|
super().__init__()
|
||||||
|
self._skip = 0
|
||||||
|
self.parts: List[str] = []
|
||||||
|
self.images: List[str] = []
|
||||||
|
self.og_image: Optional[str] = None
|
||||||
|
|
||||||
|
def handle_starttag(self, tag: str, attrs) -> None:
|
||||||
|
if tag in ("script", "style", "noscript", "svg"):
|
||||||
|
self._skip += 1
|
||||||
|
attr = dict(attrs)
|
||||||
|
src = attr.get("src")
|
||||||
|
if tag == "img" and src:
|
||||||
|
self.images.append(src)
|
||||||
|
if tag == "meta" and attr.get("property") == "og:image" and attr.get("content"):
|
||||||
|
self.og_image = attr["content"]
|
||||||
|
|
||||||
|
def handle_endtag(self, tag: str) -> None:
|
||||||
|
if tag in ("script", "style", "noscript", "svg") and self._skip > 0:
|
||||||
|
self._skip -= 1
|
||||||
|
|
||||||
|
def handle_data(self, data: str) -> None:
|
||||||
|
if self._skip == 0 and data.strip():
|
||||||
|
self.parts.append(data.strip())
|
||||||
|
|
||||||
|
|
||||||
|
def _ip_is_public(ip_str: str) -> bool:
|
||||||
|
try:
|
||||||
|
ip = ipaddress.ip_address(ip_str)
|
||||||
|
except ValueError:
|
||||||
|
return False
|
||||||
|
return not (ip.is_private or ip.is_loopback or ip.is_link_local or ip.is_multicast or ip.is_reserved or ip.is_unspecified)
|
||||||
|
|
||||||
|
|
||||||
|
def guard_url(url: str) -> Optional[str]:
|
||||||
|
"""Return None if safe to fetch, else a human-readable refusal reason (URL-02/03)."""
|
||||||
|
parsed = urlparse(url)
|
||||||
|
if parsed.scheme not in ("http", "https"):
|
||||||
|
return f"refused scheme {parsed.scheme!r} (only http/https)"
|
||||||
|
host = parsed.hostname
|
||||||
|
if not host:
|
||||||
|
return "refused: no host"
|
||||||
|
try:
|
||||||
|
literal = ipaddress.ip_address(host)
|
||||||
|
return None if _ip_is_public(str(literal)) else f"refused non-public address {host}"
|
||||||
|
except ValueError:
|
||||||
|
pass
|
||||||
|
try:
|
||||||
|
infos = socket.getaddrinfo(host, None)
|
||||||
|
except socket.gaierror:
|
||||||
|
return f"refused: cannot resolve {host}"
|
||||||
|
for info in infos:
|
||||||
|
if not _ip_is_public(str(info[4][0])):
|
||||||
|
return f"refused: {host} resolves to non-public address"
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
class URLReader:
|
||||||
|
def __init__(self, config_getter: Callable[[], Dict[str, Any]], image_cache) -> None:
|
||||||
|
self._config = config_getter
|
||||||
|
self.image_cache = image_cache
|
||||||
|
|
||||||
|
def enabled(self) -> bool:
|
||||||
|
return bool(self._config().get("enable-url-reading", False))
|
||||||
|
|
||||||
|
async def _get(self, session, url: str, max_bytes: int) -> Tuple[str, bytes]:
|
||||||
|
"""Manual redirect handling so every hop is re-guarded (URL-04)."""
|
||||||
|
current = url
|
||||||
|
for _ in range(MAX_REDIRECTS):
|
||||||
|
reason = guard_url(current)
|
||||||
|
if reason:
|
||||||
|
raise ValueError(reason)
|
||||||
|
async with session.get(current, allow_redirects=False) as response:
|
||||||
|
if response.status in (301, 302, 303, 307, 308) and response.headers.get("Location"):
|
||||||
|
current = urljoin(current, response.headers["Location"])
|
||||||
|
continue
|
||||||
|
response.raise_for_status()
|
||||||
|
return str(response.url), await response.content.read(max_bytes + 1)
|
||||||
|
raise ValueError("too many redirects")
|
||||||
|
|
||||||
|
async def fetch(self, url: str, channel: str, user: str) -> Dict[str, Any]:
|
||||||
|
config = self._config()
|
||||||
|
max_bytes = int(config.get("url-max-bytes", DEFAULT_MAX_BYTES))
|
||||||
|
timeout = aiohttp.ClientTimeout(total=FETCH_TIMEOUT_S)
|
||||||
|
try:
|
||||||
|
async with aiohttp.ClientSession(timeout=timeout, headers={"User-Agent": "FjerkroaBot/1.0"}) as session:
|
||||||
|
final_url, body = await self._get(session, url, max_bytes)
|
||||||
|
except Exception as err:
|
||||||
|
return {"error": str(err)}
|
||||||
|
text = self._to_text(body.decode("utf-8", "ignore"))
|
||||||
|
clean = sanitize_external_text(text, int(config.get("url-max-chars", DEFAULT_MAX_CHARS)))
|
||||||
|
images = await self._ingest_images(body.decode("utf-8", "ignore"), final_url, channel, user)
|
||||||
|
return {"url": final_url, "text": clean, "images_cached": images}
|
||||||
|
|
||||||
|
def _to_text(self, html: str) -> str:
|
||||||
|
extractor = _Extractor()
|
||||||
|
try:
|
||||||
|
extractor.feed(html)
|
||||||
|
except Exception as err:
|
||||||
|
logging.debug(f"html parse (text) failed: {err!r}")
|
||||||
|
return re.sub(r"\s+\n", "\n", " ".join(extractor.parts))
|
||||||
|
|
||||||
|
async def _ingest_images(self, html: str, base_url: str, channel: str, user: str) -> int:
|
||||||
|
if self.image_cache is None:
|
||||||
|
return 0
|
||||||
|
extractor = _Extractor()
|
||||||
|
try:
|
||||||
|
extractor.feed(html)
|
||||||
|
except Exception as err:
|
||||||
|
logging.debug(f"html parse (images) failed: {err!r}")
|
||||||
|
candidates = ([extractor.og_image] if extractor.og_image else []) + extractor.images
|
||||||
|
limit = int(self._config().get("url-max-images", DEFAULT_MAX_IMAGES))
|
||||||
|
cached = 0
|
||||||
|
for src in candidates:
|
||||||
|
if cached >= limit:
|
||||||
|
break
|
||||||
|
absolute = urljoin(base_url, src)
|
||||||
|
if guard_url(absolute) is not None:
|
||||||
|
continue
|
||||||
|
sha = await self.image_cache.ingest_url(absolute, channel, user, None)
|
||||||
|
if sha is not None:
|
||||||
|
cached += 1
|
||||||
|
return cached
|
||||||
@@ -6,9 +6,10 @@ with date + result.
|
|||||||
|
|
||||||
| ID | Date | Result |
|
| ID | Date | Result |
|
||||||
| --- | --- | --- |
|
| --- | --- | --- |
|
||||||
| DEP-01 | 2026-07-13 | Verified with the v3.0.0 ggg deploy: tag-only refusal + untracked config/state survived. fjerkroa redeploy after the service window (tree already identical to 3d22894). |
|
| DEP-01 | 2026-07-13 | Verified on both hosts (ggg v3.0.0..v3.3.2, fjerkroa v3.3.2): tag-only refusal + untracked config/state survived every deploy. |
|
||||||
| DEP-02 | 2026-07-13 | Service map exercised: luma restart via script (v3.0.0); kroa mapping code-reviewed, exercised on its next deploy. |
|
| DEP-02 | 2026-07-13 | Service map exercised on both hosts: luma (v3.0.0..v3.3.2) and kroa (v3.3.2, DEPLOY_FORCE per operator order). |
|
||||||
| DEP-03 | 2026-07-13 | Exercised with the v3.1.0 ggg deploy: bot.db.pre-v3.1.0 confirmed on the host. (v3.0.0 note: no pre-existing db in the pickle era.) |
|
| DEP-03 | 2026-07-13 | Exercised with the v3.1.0 ggg deploy: bot.db.pre-v3.1.0 confirmed on the host. (v3.0.0 note: no pre-existing db in the pickle era.) |
|
||||||
| DEP-04 | 2026-07-13 | Smoke gate exercised on ggg: RUNNING + fresh login line. |
|
| DEP-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). |
|
||||||
|
|||||||
@@ -0,0 +1,59 @@
|
|||||||
|
# SPEC-005 — Self-tasking
|
||||||
|
|
||||||
|
The sigmoid "boreness" loop becomes a persistent task queue (schema
|
||||||
|
v5): generators propose, a scheduler executes — through the normal
|
||||||
|
responder path, so every SAF gate, quota and kill-switch applies.
|
||||||
|
**Default off** (`tasks-enabled`, kitchen stays quiet unless opted
|
||||||
|
in); generators are individually selectable per persona
|
||||||
|
(`tasks-generators`). Marked experimental per plan v4.
|
||||||
|
|
||||||
|
### TSK-01 — Tasks are persistent queue rows (coverage: test)
|
||||||
|
|
||||||
|
A task is a store row: kind, channel, due-at, payload (the prompt the
|
||||||
|
responder will run), state (`queued`/`approval`/`done`/`cancelled`/
|
||||||
|
`failed`). Enqueued tasks survive restarts.
|
||||||
|
|
||||||
|
### TSK-02 — Due tasks run through the responder path (coverage: test)
|
||||||
|
|
||||||
|
The scheduler executes due queued tasks as system messages via the
|
||||||
|
normal respond flow (inheriting budget, gates, envelope), marks them
|
||||||
|
`done`/`failed`, and records the outcome as an observation so memory
|
||||||
|
learns what the bot did on its own.
|
||||||
|
|
||||||
|
### TSK-03 — Off by default (coverage: test)
|
||||||
|
|
||||||
|
Without `tasks-enabled = true` nothing is generated and nothing is
|
||||||
|
executed. A restaurant server does not improvise unless asked to.
|
||||||
|
|
||||||
|
### TSK-04 — Per-channel daily cap (coverage: test)
|
||||||
|
|
||||||
|
At most `tasks-max-per-channel-per-day` (default 2) task executions
|
||||||
|
per channel per day, counted in the ledger; further due tasks stay
|
||||||
|
queued for the next day.
|
||||||
|
|
||||||
|
### TSK-05 — Approval mode (coverage: test)
|
||||||
|
|
||||||
|
With `tasks-approval = true`, generated tasks enter state `approval`
|
||||||
|
and a staff alert announces them; `!bot task-approve <id>` moves them
|
||||||
|
to the queue, `!bot task-cancel <id>` kills them (works for queued
|
||||||
|
tasks too). Staged-rollout path from the review.
|
||||||
|
|
||||||
|
### TSK-06 — Kill-switches govern execution (coverage: test)
|
||||||
|
|
||||||
|
Task execution respects `bot_initiated_allowed()` — `!bot tasks off`,
|
||||||
|
pause, quiet mode and quiet hours all stop the scheduler tick.
|
||||||
|
|
||||||
|
### TSK-07 — Idle-impulse generator (boreness, demoted) (coverage: test)
|
||||||
|
|
||||||
|
When the chat channel has been idle longer than
|
||||||
|
`idle-impulse-hours` (default 12), the generator enqueues one
|
||||||
|
impulse task with the configured boreness prompt — deterministic
|
||||||
|
threshold instead of the old 7-second sigmoid dice loop, bounded by
|
||||||
|
TSK-04. The `on_boreness` loop is gone.
|
||||||
|
|
||||||
|
### TSK-08 — Follow-up generator proposes from memory (coverage: test)
|
||||||
|
|
||||||
|
Periodically (`taskgen-interval-hours`, default 6) the follow-up
|
||||||
|
generator asks `memory-model` (strict structured outputs) whether the
|
||||||
|
recent episodes/facts warrant one follow-up task (channel, prompt,
|
||||||
|
due-in-hours). A null/failed proposal enqueues nothing.
|
||||||
@@ -62,6 +62,12 @@ Bot-initiated posts also respect pause/quiet.
|
|||||||
`!bot pins` answers with all pinned facts and their ids (global +
|
`!bot pins` answers with all pinned facts and their ids (global +
|
||||||
per-channel) — without it, `!bot unpin <id>` required guessing ids.
|
per-channel) — without it, `!bot unpin <id>` required guessing ids.
|
||||||
|
|
||||||
|
### OPS-12 — Task queue surface (coverage: test)
|
||||||
|
|
||||||
|
`!bot tasks` lists open tasks (queued + awaiting approval) with ids;
|
||||||
|
`!bot task-approve <id>` and `!bot task-cancel <id>` manage them
|
||||||
|
(TSK-05). `!bot tasks on|off` stays the kill-switch (OPS-09).
|
||||||
|
|
||||||
### OPS-10 — Spend report (coverage: test)
|
### OPS-10 — Spend report (coverage: test)
|
||||||
|
|
||||||
`!bot spend` answers in the staff channel with today's estimated
|
`!bot spend` answers in the staff channel with today's estimated
|
||||||
|
|||||||
@@ -0,0 +1,57 @@
|
|||||||
|
# SPEC-011 — URL reading
|
||||||
|
|
||||||
|
A `fetch_url` tool alongside IGDB: the model decides when to read a
|
||||||
|
link (user pastes a URL + question; a news item links an article
|
||||||
|
Luma wants details on). Web pages are the number-one injection
|
||||||
|
vector, so everything fetched is sanitized (SAF-03) and the fetch
|
||||||
|
itself is SSRF-guarded — the bot runs on shared hosting. Active only
|
||||||
|
when `enable-url-reading = true`.
|
||||||
|
|
||||||
|
### URL-01 — fetch_url is offered as a tool (coverage: test)
|
||||||
|
|
||||||
|
When `enable-url-reading` is true, the chat call's `tools` list
|
||||||
|
includes a `fetch_url` function (url string param) next to any IGDB
|
||||||
|
tools. When false, it is absent.
|
||||||
|
|
||||||
|
### URL-02 — Only http/https are fetched (coverage: test)
|
||||||
|
|
||||||
|
`file:`, `ftp:`, `data:`, `gopher:` and schemeless inputs are
|
||||||
|
refused before any network call, with an error result the model can
|
||||||
|
relay.
|
||||||
|
|
||||||
|
### URL-03 — SSRF guard blocks non-public addresses (coverage: test)
|
||||||
|
|
||||||
|
Before fetching, the host is resolved and every resulting IP is
|
||||||
|
checked; the fetch is refused when any is private, loopback,
|
||||||
|
link-local, or otherwise non-global (RFC1918, 127/8, 169.254/16,
|
||||||
|
::1, fc00::/7, etc.). A URL literal that is already such an IP is
|
||||||
|
refused without DNS.
|
||||||
|
|
||||||
|
### URL-04 — Redirects are re-validated (coverage: test)
|
||||||
|
|
||||||
|
Redirects are followed manually; each hop's target passes URL-02 and
|
||||||
|
URL-03 again. A public URL that 302-redirects to `localhost` or an
|
||||||
|
internal IP is refused at the redirect, not fetched.
|
||||||
|
|
||||||
|
### URL-05 — Fetched text is bounded and sanitized (coverage: test)
|
||||||
|
|
||||||
|
Responses are capped at `url-max-bytes` (default 2 MB) with a
|
||||||
|
download timeout; HTML is reduced to readable text (script/style
|
||||||
|
dropped, tags stripped, whitespace collapsed) and passed through
|
||||||
|
`sanitize_external_text` before it reaches the model, truncated to
|
||||||
|
`url-max-chars` (default 6000).
|
||||||
|
|
||||||
|
### URL-06 — Page images feed the cache (coverage: test)
|
||||||
|
|
||||||
|
Up to `url-max-images` (default 2) prominent images (og:image, then
|
||||||
|
large `<img>`) are ingested into the ImageCache for the requesting
|
||||||
|
channel (SSRF-guarded like the page), so the model can see them and
|
||||||
|
`picture_edit` can remix them. Ingestion failures are skipped, never
|
||||||
|
fatal to the text result.
|
||||||
|
|
||||||
|
### URL-07 — Fetches are metered and capped (coverage: test)
|
||||||
|
|
||||||
|
Each fetch increments a per-user daily counter; over
|
||||||
|
`url-daily-per-user` (default 20) `fetch_url` refuses with an error
|
||||||
|
result. The budget gate (SAF-04) still applies to the surrounding
|
||||||
|
model calls.
|
||||||
@@ -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,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,165 @@
|
|||||||
|
"""Unit coverage for SPEC-005 self-tasking (TSK-01..08) + OPS-12."""
|
||||||
|
|
||||||
|
import tempfile
|
||||||
|
import unittest
|
||||||
|
from pathlib import Path
|
||||||
|
from unittest.mock import AsyncMock
|
||||||
|
|
||||||
|
from fjerkroa_bot.persistence import PersistentStore
|
||||||
|
from fjerkroa_bot.quota import QuotaLedger
|
||||||
|
from fjerkroa_bot.tasks import TaskEngine
|
||||||
|
|
||||||
|
from .test_spec_ops import OpsBase
|
||||||
|
|
||||||
|
|
||||||
|
class EngineBase(unittest.IsolatedAsyncioTestCase):
|
||||||
|
def setUp(self):
|
||||||
|
self.tmp = tempfile.TemporaryDirectory()
|
||||||
|
self.addCleanup(self.tmp.cleanup)
|
||||||
|
self.store = PersistentStore(Path(self.tmp.name) / "bot.db")
|
||||||
|
self.config = {"tasks-enabled": True, "chat-channel": "chat", "taskgen-interval-hours": 0}
|
||||||
|
self.ledger = QuotaLedger(self.store, lambda: self.config)
|
||||||
|
self.execute = AsyncMock()
|
||||||
|
self.propose = AsyncMock(return_value={"task": None})
|
||||||
|
self.alert = AsyncMock()
|
||||||
|
self.observe = AsyncMock()
|
||||||
|
self.idle = lambda: 0.0
|
||||||
|
self.engine = TaskEngine(
|
||||||
|
self.store,
|
||||||
|
self.ledger,
|
||||||
|
lambda: self.config,
|
||||||
|
self.execute,
|
||||||
|
self.propose,
|
||||||
|
self.alert,
|
||||||
|
lambda: True,
|
||||||
|
lambda: self.idle(),
|
||||||
|
self.observe,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class TestQueuePersistence(EngineBase):
|
||||||
|
async def test_tasks_survive_restart(self):
|
||||||
|
"""TSK-01: enqueued tasks are store rows and reload in a fresh store."""
|
||||||
|
await self.engine.enqueue("follow-up", "chat", "frag bob nach der pruefung", due_hours=1)
|
||||||
|
reborn = PersistentStore(Path(self.tmp.name) / "bot.db")
|
||||||
|
tasks = reborn.tasks_open()
|
||||||
|
self.assertEqual(len(tasks), 1)
|
||||||
|
self.assertEqual((tasks[0]["kind"], tasks[0]["channel"], tasks[0]["state"]), ("follow-up", "chat", "queued"))
|
||||||
|
|
||||||
|
|
||||||
|
class TestExecution(EngineBase):
|
||||||
|
async def test_due_task_executes_and_completes(self):
|
||||||
|
"""TSK-02: due task runs via execute callback, marked done, observed."""
|
||||||
|
await self.engine.enqueue("idle-impulse", "chat", "sag was nettes")
|
||||||
|
await self.engine.tick()
|
||||||
|
self.execute.assert_awaited_once_with("chat", "sag was nettes")
|
||||||
|
self.assertEqual(self.store.tasks_open(), [])
|
||||||
|
self.observe.assert_awaited()
|
||||||
|
|
||||||
|
async def test_failed_execution_marked_failed(self):
|
||||||
|
"""TSK-02: execute raising marks the task failed, not done."""
|
||||||
|
self.execute.side_effect = RuntimeError("channel gone")
|
||||||
|
await self.engine.enqueue("idle-impulse", "chat", "x")
|
||||||
|
await self.engine.tick()
|
||||||
|
self.assertEqual(self.store.tasks_open(), []) # not queued anymore
|
||||||
|
rows = [t for t in self.store.tasks_due()]
|
||||||
|
self.assertEqual(rows, [])
|
||||||
|
|
||||||
|
|
||||||
|
class TestDefaultOff(EngineBase):
|
||||||
|
async def test_disabled_engine_is_dormant(self):
|
||||||
|
"""TSK-03: without tasks-enabled nothing executes or generates."""
|
||||||
|
self.store.task_add("idle-impulse", "chat", "2020-01-01 00:00:00", "x")
|
||||||
|
self.config.pop("tasks-enabled")
|
||||||
|
await self.engine.tick()
|
||||||
|
self.execute.assert_not_awaited()
|
||||||
|
self.propose.assert_not_awaited()
|
||||||
|
|
||||||
|
|
||||||
|
class TestDailyCap(EngineBase):
|
||||||
|
async def test_cap_defers_excess_tasks(self):
|
||||||
|
"""TSK-04: over the per-channel cap tasks stay queued."""
|
||||||
|
self.config["tasks-max-per-channel-per-day"] = 1
|
||||||
|
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "one")
|
||||||
|
self.store.task_add("b", "chat", "2020-01-01 00:00:00", "two")
|
||||||
|
await self.engine.tick()
|
||||||
|
self.assertEqual(self.execute.await_count, 1)
|
||||||
|
self.assertEqual(len(self.store.tasks_open()), 1)
|
||||||
|
|
||||||
|
|
||||||
|
class TestApproval(EngineBase):
|
||||||
|
async def test_approval_flow(self):
|
||||||
|
"""TSK-05: approval mode holds tasks until approved; cancel kills them."""
|
||||||
|
self.config["tasks-approval"] = True
|
||||||
|
task_id = await self.engine.enqueue("follow-up", "chat", "frag nach")
|
||||||
|
self.alert.assert_awaited_once()
|
||||||
|
await self.engine.tick()
|
||||||
|
self.execute.assert_not_awaited() # approval != due
|
||||||
|
self.store.task_set_state(task_id, "queued")
|
||||||
|
await self.engine.tick()
|
||||||
|
self.execute.assert_awaited_once()
|
||||||
|
|
||||||
|
|
||||||
|
class TestKillSwitch(EngineBase):
|
||||||
|
async def test_not_allowed_blocks_tick(self):
|
||||||
|
"""TSK-06: bot_initiated_allowed()=False stops execution."""
|
||||||
|
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "x")
|
||||||
|
engine = TaskEngine(
|
||||||
|
self.store, self.ledger, lambda: self.config, self.execute, self.propose, self.alert, lambda: False, lambda: 0.0, self.observe
|
||||||
|
)
|
||||||
|
await engine.tick()
|
||||||
|
self.execute.assert_not_awaited()
|
||||||
|
|
||||||
|
|
||||||
|
class TestIdleImpulse(EngineBase):
|
||||||
|
async def test_idle_enqueues_once(self):
|
||||||
|
"""TSK-07: long idle enqueues one impulse; pending impulse dedupes."""
|
||||||
|
self.config["idle-impulse-hours"] = 1
|
||||||
|
self.config["boreness-prompt"] = "denk dir was aus"
|
||||||
|
self.idle = lambda: 2 * 3600.0
|
||||||
|
await self.engine.tick()
|
||||||
|
open_tasks = self.store.tasks_open()
|
||||||
|
impulse = [t for t in open_tasks if t["kind"] == "idle-impulse"]
|
||||||
|
self.assertEqual(len(impulse), 0) # executed immediately (due now)
|
||||||
|
self.execute.assert_awaited_once_with("chat", "denk dir was aus")
|
||||||
|
|
||||||
|
async def test_short_idle_no_impulse(self):
|
||||||
|
"""TSK-07: below the idle threshold nothing is generated."""
|
||||||
|
self.config["idle-impulse-hours"] = 12
|
||||||
|
self.idle = lambda: 60.0
|
||||||
|
await self.engine.tick()
|
||||||
|
self.execute.assert_not_awaited()
|
||||||
|
|
||||||
|
|
||||||
|
class TestFollowUpGenerator(EngineBase):
|
||||||
|
async def test_proposal_becomes_task(self):
|
||||||
|
"""TSK-08: a proposed task is enqueued with its due offset."""
|
||||||
|
self.config.pop("chat-channel")
|
||||||
|
self.propose.return_value = {"task": {"channel": "chat", "prompt": "frag bob", "due_hours": 24}}
|
||||||
|
await self.engine.tick()
|
||||||
|
tasks = self.store.tasks_open()
|
||||||
|
self.assertEqual(len(tasks), 1)
|
||||||
|
self.assertEqual(tasks[0]["kind"], "follow-up")
|
||||||
|
|
||||||
|
async def test_null_proposal_no_task(self):
|
||||||
|
"""TSK-08: null proposal enqueues nothing."""
|
||||||
|
self.propose.return_value = {"task": None}
|
||||||
|
await self.engine.tick()
|
||||||
|
self.assertEqual(self.store.tasks_open(), [])
|
||||||
|
|
||||||
|
|
||||||
|
class TestTaskCommands(OpsBase):
|
||||||
|
async def test_list_approve_cancel(self):
|
||||||
|
"""OPS-12: !bot tasks lists; task-approve/task-cancel manage states."""
|
||||||
|
with tempfile.TemporaryDirectory() as tmp:
|
||||||
|
store = PersistentStore(Path(tmp) / "bot.db")
|
||||||
|
self.bot.airesponder.store = store
|
||||||
|
task_id = store.task_add("follow-up", "chat", "2099-01-01 00:00:00", "frag bob", "approval")
|
||||||
|
await self.bot.on_message(self.staff_msg("!bot tasks"))
|
||||||
|
listing = self.bot.staff_channel.send.await_args.args[0]
|
||||||
|
self.assertIn("follow-up", listing)
|
||||||
|
self.assertIn("approval", listing)
|
||||||
|
await self.bot.on_message(self.staff_msg(f"!bot task-approve {task_id}"))
|
||||||
|
self.assertEqual(store.tasks_open()[0]["state"], "queued")
|
||||||
|
await self.bot.on_message(self.staff_msg(f"!bot task-cancel {task_id}"))
|
||||||
|
self.assertEqual(store.tasks_open(), [])
|
||||||
@@ -0,0 +1,143 @@
|
|||||||
|
"""Unit coverage for SPEC-011 URL reading (URL-01..07)."""
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from unittest.mock import AsyncMock, patch
|
||||||
|
|
||||||
|
from fjerkroa_bot.openai_responder import OpenAIResponder
|
||||||
|
from fjerkroa_bot.url_reader import FETCH_URL_TOOL, URLReader, guard_url
|
||||||
|
|
||||||
|
CONFIG = {"openai-token": "t", "model": "m", "system": "s", "history-limit": 5}
|
||||||
|
|
||||||
|
|
||||||
|
class TestToolOffered(unittest.TestCase):
|
||||||
|
def test_tool_present_only_when_enabled(self):
|
||||||
|
"""URL-01: fetch_url appears in the tool list only with enable-url-reading."""
|
||||||
|
off = OpenAIResponder(CONFIG, "chat")
|
||||||
|
self.assertNotIn("fetch_url", [f["name"] for f in off._available_tools()])
|
||||||
|
on = OpenAIResponder(dict(CONFIG, **{"enable-url-reading": True}), "chat")
|
||||||
|
self.assertIn("fetch_url", [f["name"] for f in on._available_tools()])
|
||||||
|
self.assertEqual(FETCH_URL_TOOL["name"], "fetch_url")
|
||||||
|
|
||||||
|
|
||||||
|
class TestSchemeGuard(unittest.TestCase):
|
||||||
|
def test_non_http_schemes_refused(self):
|
||||||
|
"""URL-02: only http/https pass the guard."""
|
||||||
|
self.assertIsNone(guard_url("https://example.com/article"))
|
||||||
|
for bad in ("file:///etc/passwd", "ftp://host/x", "data:text/html,x", "gopher://h", "no-scheme.com/x"):
|
||||||
|
self.assertIsNotNone(guard_url(bad))
|
||||||
|
|
||||||
|
|
||||||
|
class TestSSRFGuard(unittest.TestCase):
|
||||||
|
def test_private_and_loopback_refused(self):
|
||||||
|
"""URL-03: private/loopback/link-local literals are refused without DNS."""
|
||||||
|
for bad in (
|
||||||
|
"http://127.0.0.1/admin",
|
||||||
|
"http://localhost/x", # resolves to loopback
|
||||||
|
"http://10.0.0.5/x",
|
||||||
|
"http://192.168.1.1/x",
|
||||||
|
"http://169.254.169.254/latest/meta-data", # cloud metadata
|
||||||
|
"http://[::1]/x",
|
||||||
|
):
|
||||||
|
self.assertIsNotNone(guard_url(bad), f"{bad} should be refused")
|
||||||
|
|
||||||
|
def test_public_ip_allowed(self):
|
||||||
|
"""URL-03: a public IP literal passes."""
|
||||||
|
self.assertIsNone(guard_url("http://93.184.216.34/"))
|
||||||
|
|
||||||
|
@patch("fjerkroa_bot.url_reader.socket.getaddrinfo")
|
||||||
|
def test_dns_to_private_refused(self, getaddrinfo):
|
||||||
|
"""URL-03: a hostname resolving to a private IP is refused."""
|
||||||
|
getaddrinfo.return_value = [(2, 1, 6, "", ("10.1.2.3", 0))]
|
||||||
|
self.assertIsNotNone(guard_url("http://evil.example.com/x"))
|
||||||
|
|
||||||
|
|
||||||
|
class TestRedirectRevalidation(unittest.IsolatedAsyncioTestCase):
|
||||||
|
async def test_redirect_to_internal_refused(self):
|
||||||
|
"""URL-04: a public URL redirecting to localhost is refused at the hop."""
|
||||||
|
reader = URLReader(lambda: {}, None)
|
||||||
|
|
||||||
|
class FakeResp:
|
||||||
|
status = 302
|
||||||
|
headers = {"Location": "http://127.0.0.1/secret"}
|
||||||
|
url = "http://safe.example.com"
|
||||||
|
|
||||||
|
async def __aenter__(self):
|
||||||
|
return self
|
||||||
|
|
||||||
|
async def __aexit__(self, *a):
|
||||||
|
return False
|
||||||
|
|
||||||
|
def raise_for_status(self):
|
||||||
|
pass
|
||||||
|
|
||||||
|
class FakeSession:
|
||||||
|
def get(self, url, allow_redirects=False):
|
||||||
|
return FakeResp()
|
||||||
|
|
||||||
|
with patch("fjerkroa_bot.url_reader.guard_url", side_effect=[None, "refused internal"]):
|
||||||
|
with self.assertRaises(ValueError):
|
||||||
|
await reader._get(FakeSession(), "http://safe.example.com", 1000)
|
||||||
|
|
||||||
|
|
||||||
|
class TestTextExtraction(unittest.TestCase):
|
||||||
|
def test_html_reduced_to_text(self):
|
||||||
|
"""URL-05: scripts/styles dropped, tags stripped."""
|
||||||
|
reader = URLReader(lambda: {}, None)
|
||||||
|
html = "<html><head><style>x{}</style></head><body><h1>Titel</h1><script>evil()</script><p>Inhalt hier</p></body></html>"
|
||||||
|
text = reader._to_text(html)
|
||||||
|
self.assertIn("Titel", text)
|
||||||
|
self.assertIn("Inhalt hier", text)
|
||||||
|
self.assertNotIn("evil", text)
|
||||||
|
self.assertNotIn("x{}", text)
|
||||||
|
|
||||||
|
|
||||||
|
class 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)
|
||||||
Reference in New Issue
Block a user