Compare commits

...

4 Commits

13 changed files with 1059 additions and 48 deletions
+12
View File
@@ -55,3 +55,15 @@ enable-game-info = true
# image-model = "gpt-image-2" # default; dall-e-3 gets clamped to n=1 # image-model = "gpt-image-2" # default; dall-e-3 gets clamped to n=1
# image-size = "1024x1024" # image-size = "1024x1024"
# image-quality = "medium" # passed through only when set # image-quality = "medium" # passed through only when set
# Image input pipeline (SPEC-004, FDB-010) — active with history-directory:
# image-cache-mb = 500 # LRU cap (ggg: consider 2000 — screenshots)
# image-cache-ttl-days = 90
# image-max-bytes = 8388608 # 8 MB upload cap
# Self-tasking (SPEC-005) — experimental, DEFAULT OFF:
# tasks-enabled = true
# tasks-generators = ["idle-impulse", "follow-up"]
# tasks-max-per-channel-per-day = 2
# tasks-approval = false # true: neue Tasks brauchen !bot task-approve
# idle-impulse-hours = 12
# taskgen-interval-hours = 6
# Staff: !bot tasks | task-approve <id> | task-cancel <id>
+22 -6
View File
@@ -12,6 +12,7 @@ from pathlib import Path
from pprint import pformat from pprint import pformat
from typing import Any, Dict, List, Optional, Tuple, Union from typing import Any, Dict, List, Optional, Tuple, Union
from .images import ImageCache
from .memory import MemoryManager from .memory import MemoryManager
from .persistence import PersistentStore from .persistence import PersistentStore
@@ -153,17 +154,16 @@ class AIResponder(AIResponderBase):
if stored_memory is not None: if stored_memory is not None:
self.memory = stored_memory self.memory = stored_memory
self.memory_manager = MemoryManager(self.store, lambda: self.config, self.consolidate, self.channel) self.memory_manager = MemoryManager(self.store, lambda: self.config, self.consolidate, self.channel)
self.image_cache: Optional[ImageCache] = None
if self.store is not None:
self.image_cache = ImageCache(self.store, Path(self.config["history-directory"]).expanduser() / "images", lambda: self.config)
logging.info(f"memmory:\n{self.memory}") logging.info(f"memmory:\n{self.memory}")
# Dynamic values move to a context suffix so the persona prefix # Dynamic values move to a context suffix so the persona prefix
# stays byte-stable for the prompt cache (ENV-20) # stays byte-stable for the prompt cache (ENV-20)
DYNAMIC_PLACEHOLDERS = ("{date}", "{time}", "{news}", "{memory}") DYNAMIC_PLACEHOLDERS = ("{date}", "{time}", "{news}", "{memory}")
def message(self, message: AIMessage, limit: Optional[int] = None) -> List[Dict[str, Any]]: def _context_lines(self, message: AIMessage) -> List[str]:
messages = []
persona = self.config.get(self.channel, self.config["system"])
for placeholder in self.DYNAMIC_PLACEHOLDERS:
persona = persona.replace(placeholder, "")
context = [f"date: {time.strftime('%Y-%m-%d')} ({time.strftime('%A')})", f"time: {time.strftime('%H:%M:%S')}"] context = [f"date: {time.strftime('%Y-%m-%d')} ({time.strftime('%A')})", f"time: {time.strftime('%H:%M:%S')}"]
news_feed = self.config.get("news") news_feed = self.config.get("news")
if news_feed and os.path.exists(news_feed): if news_feed and os.path.exists(news_feed):
@@ -173,7 +173,23 @@ class AIResponder(AIResponderBase):
memory_block = self.memory_manager.memory_block(participants, self.memory) memory_block = self.memory_manager.memory_block(participants, self.memory)
if memory_block: if memory_block:
context.append("memory:\n" + memory_block) context.append("memory:\n" + memory_block)
system = persona.rstrip() + "\n\n## Context\n" + "\n".join(context) if self.image_cache is not None:
recent_images = self.image_cache.recent(message.channel, 4)
if recent_images:
# the model cannot use picture_edit unless told images exist (IMG-16)
context.append(
f"recent images in this channel: {len(recent_images)}. When the user asks to modify, reuse, combine or"
" include a previously shared image, you MUST set picture_edit=true — text-to-image cannot see earlier"
" images; only picture_edit passes them to the image model."
)
return context
def message(self, message: AIMessage, limit: Optional[int] = None) -> List[Dict[str, Any]]:
messages = []
persona = self.config.get(self.channel, self.config["system"])
for placeholder in self.DYNAMIC_PLACEHOLDERS:
persona = persona.replace(placeholder, "")
system = persona.rstrip() + "\n\n## Context\n" + "\n".join(self._context_lines(message))
messages.append({"role": "system", "content": system}) messages.append({"role": "system", "content": system})
if limit is not None: if limit is not None:
while len(self.history) > limit: while len(self.history) > limit:
+96 -39
View File
@@ -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. "
@@ -114,38 +114,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)})"
) )
@@ -197,6 +201,8 @@ class FjerkroaBot(commands.Bot):
removed += self.airesponder.store.delete_history_of_user(user) removed += self.airesponder.store.delete_history_of_user(user)
# facts + observations + episode traces (MEM-09) # facts + observations + episode traces (MEM-09)
removed += self.airesponder.store.purge_user_memory(user) removed += self.airesponder.store.purge_user_memory(user)
if self.airesponder.image_cache is not None:
removed += self.airesponder.image_cache.purge_user(user) # IMG-14
logging.info(f"forgetme: removed {removed} entries for {user}") logging.info(f"forgetme: removed {removed} entries for {user}")
await message.channel.send( await message.channel.send(
f"Removed your messages, facts and memory traces ({removed} entries).", f"Removed your messages, facts and memory traces ({removed} entries).",
@@ -239,13 +245,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
@@ -337,6 +361,8 @@ class FjerkroaBot(commands.Bot):
async def on_message_delete(self, message): async def on_message_delete(self, message):
airesponder = self.get_ai_responder(self.get_channel_name(message.channel)) airesponder = self.get_ai_responder(self.get_channel_name(message.channel))
if airesponder.image_cache is not None:
airesponder.image_cache.purge_message(str(message.id)) # IMG-14
await airesponder.observe_event(message.author.name, "delete", f"deleted: {message.content}") await airesponder.observe_event(message.author.name, "delete", f"deleted: {message.content}")
def on_config_file_modified(self, event): def on_config_file_modified(self, event):
@@ -397,27 +423,46 @@ class FjerkroaBot(commands.Bot):
def get_ai_responder(self, channel_name): def get_ai_responder(self, channel_name):
return self.aichannels[channel_name] if channel_name in self.aichannels else self.airesponder return self.aichannels[channel_name] if channel_name in self.aichannels else self.airesponder
async def _ingest_attachments(self, message, channel_name: str, airesponder) -> list:
"""Cache-first attachment handling; CDN URLs never travel further (IMG-10/11)."""
urls = []
for attachment in message.attachments:
if airesponder.image_cache is None:
urls.append(attachment.url)
continue
sha = await airesponder.image_cache.ingest_url(attachment.url, channel_name, message.author.name, str(message.id))
if sha is not None:
recent = airesponder.image_cache.recent(channel_name, 8)
ext = next((row["ext"] for row in recent if row["sha256"] == sha), "png")
data_url = airesponder.image_cache.data_url(sha, ext)
if data_url:
urls.append(data_url)
return urls
async def handle_message_through_responder(self, message): async def handle_message_through_responder(self, message):
"""Handle a message through the AI responder""" """Handle a message through the AI responder"""
message_content = str(message.content).strip() message_content = str(message.content).strip()
if message.reference and message.reference.resolved and isinstance(message.reference.resolved.content, str): if message.reference and message.reference.resolved and isinstance(message.reference.resolved.content, str):
reference_content = str(message.reference.resolved.content).replace("\n", "> \n") reference_content = str(message.reference.resolved.content).replace("\n", "> \n")
message_content = f"> {reference_content}\n\n{message_content}" message_content = f"> {reference_content}\n\n{message_content}"
channel_name = self.get_channel_name(message.channel)
airesponder = self.get_ai_responder(channel_name)
attachment_urls = []
if message.attachments:
attachment_urls = await self._ingest_attachments(message, channel_name, airesponder)
if len(message_content) < 1: if len(message_content) < 1:
# image-only posts: cached + observed, no reply (IMG-17)
if attachment_urls:
await airesponder.observe_event(message.author.name, "image", f"posted {len(attachment_urls)} image(s)")
return return
message_content = self._resolve_mentions(message_content) message_content = self._resolve_mentions(message_content)
channel_name = self.get_channel_name(message.channel)
msg = AIMessage( msg = AIMessage(
message.author.name, message_content, channel_name, self.user in message.mentions or isinstance(message.channel, DMChannel) message.author.name, message_content, channel_name, self.user in message.mentions or isinstance(message.channel, DMChannel)
) )
if message.attachments: if attachment_urls:
for attachment in message.attachments: msg.urls = attachment_urls
if not msg.urls:
msg.urls = []
msg.urls.append(attachment.url)
# Reply/ignore classifier gate — direct messages bypass (BEH-01/02/03/07) # Reply/ignore classifier gate — direct messages bypass (BEH-01/02/03/07)
airesponder = self.get_ai_responder(channel_name)
handled, factual = await self._classifier_gate(message, msg, airesponder, channel_name) handled, factual = await self._classifier_gate(message, msg, airesponder, channel_name)
if handled: if handled:
return return
@@ -462,7 +507,19 @@ class FjerkroaBot(commands.Bot):
"""Send the answer paced, split and with images on the last part (BEH-04/05/06)""" """Send the answer paced, split and with images on the last part (BEH-04/05/06)"""
files = None files = None
if response.picture is not None: if response.picture is not None:
buffers = await airesponder.draw(response.picture, getattr(response, "picture_count", 1)) count = getattr(response, "picture_count", 1)
channel_name = self.get_channel_name(answer_channel)
buffers = None
if getattr(response, "picture_edit", False) and airesponder.image_cache is not None:
sources = airesponder.image_cache.recent_paths(channel_name, 4)
if sources:
buffers = await airesponder.edit_openai(response.picture, sources, count)
if buffers is None:
# empty cache or no edit request: plain generation (IMG-13 fallback)
buffers = await airesponder.draw(response.picture, count)
if airesponder.image_cache is not None:
for buffer in buffers:
airesponder.image_cache.ingest_bytes(buffer.getvalue(), channel_name, "assistant", None) # IMG-15
files = [discord.File(fp=buffer, filename=f"image-{index}.png") for index, buffer in enumerate(buffers)] files = [discord.File(fp=buffer, filename=f"image-{index}.png") for index, buffer in enumerate(buffers)]
parts = split_answer(response.answer, int(self.config.get("split-threshold", 1200)), int(self.config.get("split-max-parts", 3))) parts = split_answer(response.answer, int(self.config.get("split-threshold", 1200)), int(self.config.get("split-max-parts", 3)))
pace = float(self.config.get("typing-chars-per-second", 0) or 0) pace = float(self.config.get("typing-chars-per-second", 0) or 0)
+124
View File
@@ -0,0 +1,124 @@
"""Content-hash image cache (SPEC-004, FDB-010).
Attachments are downloaded once, sniffed, stored under their sha256
and served to vision as data: URLs — Discord's expiring CDN links
never travel further (IMG-10/11). LRU + TTL keep the cache bounded
(IMG-12); deletions and !forgetme propagate here (IMG-14).
"""
import base64
import hashlib
import logging
from pathlib import Path
from typing import Any, Callable, Dict, List, Optional
import aiohttp
from .persistence import PersistentStore
DEFAULT_CACHE_MB = 500
DEFAULT_TTL_DAYS = 90
DEFAULT_MAX_BYTES = 8 * 1024 * 1024
DOWNLOAD_TIMEOUT_S = 20
MAGIC = [
(b"\x89PNG", "png"),
(b"\xff\xd8\xff", "jpg"),
(b"GIF87a", "gif"),
(b"GIF89a", "gif"),
]
def sniff_ext(data: bytes) -> Optional[str]:
"""Extension from magic bytes only — names and headers lie (IMG-10)."""
for magic, ext in MAGIC:
if data.startswith(magic):
return ext
if data[:4] == b"RIFF" and data[8:12] == b"WEBP":
return "webp"
return None
class ImageCache:
def __init__(self, store: PersistentStore, root: Path, config_getter: Callable[[], Dict[str, Any]]) -> None:
self.store = store
self.root = Path(root)
self._config = config_getter
self.root.mkdir(parents=True, exist_ok=True)
def _path(self, sha256: str, ext: str) -> Path:
return self.root / f"{sha256}.{ext}"
def ingest_bytes(self, data: bytes, channel: str, user: str, message_id: Optional[str]) -> Optional[str]:
ext = sniff_ext(data)
if ext is None:
logging.warning(f"image cache: rejected non-image bytes from {user} (IMG-10)")
return None
if len(data) > int(self._config().get("image-max-bytes", DEFAULT_MAX_BYTES)):
logging.warning(f"image cache: rejected oversized upload from {user} ({len(data)} bytes)")
return None
sha256 = hashlib.sha256(data).hexdigest()
path = self._path(sha256, ext)
if not path.exists():
path.write_bytes(data)
self.store.image_add(sha256, channel, user, message_id, ext, len(data))
self.evict()
return sha256
async def ingest_url(self, url: str, channel: str, user: str, message_id: Optional[str]) -> Optional[str]:
try:
data = await self._download(url)
except Exception as err:
logging.warning(f"image cache: download failed for {user}: {repr(err)}")
return None
return self.ingest_bytes(data, channel, user, message_id)
async def _download(self, url: str) -> bytes:
limit = int(self._config().get("image-max-bytes", DEFAULT_MAX_BYTES))
timeout = aiohttp.ClientTimeout(total=DOWNLOAD_TIMEOUT_S)
async with aiohttp.ClientSession(timeout=timeout) as session:
async with session.get(url) as response:
response.raise_for_status()
return await response.content.read(limit + 1)
def data_url(self, sha256: str, ext: str) -> Optional[str]:
path = self._path(sha256, ext)
if not path.exists():
return None
mime = "jpeg" if ext == "jpg" else ext
return f"data:image/{mime};base64," + base64.b64encode(path.read_bytes()).decode()
def recent(self, channel: str, count: int) -> List[Dict[str, Any]]:
return self.store.images_recent(channel, count)
def recent_paths(self, channel: str, count: int) -> List[Path]:
paths = [self._path(row["sha256"], row["ext"]) for row in self.recent(channel, count)]
return [path for path in paths if path.exists()]
def _remove(self, sha256: str, ext: str) -> None:
self._path(sha256, ext).unlink(missing_ok=True)
self.store.images_delete(sha256)
def evict(self) -> None:
"""TTL first, then LRU down to the byte cap (IMG-12)."""
config = self._config()
for row in self.store.images_expired(int(config.get("image-cache-ttl-days", DEFAULT_TTL_DAYS))):
self._remove(row["sha256"], row["ext"])
cap = int(config.get("image-cache-mb", DEFAULT_CACHE_MB)) * 1024 * 1024
while self.store.images_total_bytes() > cap:
victims = self.store.images_oldest(1)
if not victims:
break
self._remove(victims[0]["sha256"], victims[0]["ext"])
def purge_user(self, user: str) -> int:
rows = self.store.images_for_user(user)
for row in rows:
self._remove(row["sha256"], row["ext"])
return len(rows)
def purge_message(self, message_id: str) -> int:
rows = self.store.images_for_message(message_id)
for row in rows:
self._remove(row["sha256"], row["ext"])
return len(rows)
+75
View File
@@ -58,6 +58,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",
@@ -96,6 +121,10 @@ async def openai_image(client, *args, **kwargs):
return await client.images.generate(*args, **kwargs) return await client.images.generate(*args, **kwargs)
async def openai_image_edit(client, *args, **kwargs):
return await client.images.edit(*args, **kwargs)
class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn): class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
def __init__(self, config: Dict[str, Any], channel: Optional[str] = None) -> None: def __init__(self, config: Dict[str, Any], channel: Optional[str] = None) -> None:
super().__init__(config, channel) super().__init__(config, channel)
@@ -354,6 +383,52 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
logging.debug(f"Full traceback: {traceback.format_exc()}") logging.debug(f"Full traceback: {traceback.format_exc()}")
return None, limit return None, limit
async def edit_openai(self, description: str, paths: List[Any], count: int = 1) -> List[BytesIO]:
"""Edit/remix from cached inputs, ≤4 files (IMG-13)."""
if not self.ledger.budget_ok():
raise RuntimeError("daily budget exhausted - refusing image edit")
model = self.config.get("image-model", "gpt-image-2")
handles = [open(path, "rb") for path in paths[:4]]
try:
response = await openai_image_edit(
self.client,
model=model,
image=handles if len(handles) > 1 else handles[0],
prompt=description,
n=max(1, min(int(count), 4)),
size=self.config.get("image-size", "1024x1024"),
)
finally:
for handle in handles:
handle.close()
buffers = [BytesIO(base64.b64decode(item.b64_json)) for item in response.data]
self.ledger.add_images(len(buffers))
logging.info(f"edited {len(buffers)} image(s) on {model} from {len(handles)} input(s)")
return buffers
async def propose_task(self) -> Optional[Dict[str, Any]]:
"""One follow-up proposal from recent episodes on memory-model (TSK-08)."""
if "memory-model" not in self.config or self.store is None or not self.ledger.budget_ok():
return None
channel = self.config.get("chat-channel", "chat")
episodes = await asyncio.to_thread(self.store.recent_episodes, channel, 5)
if not episodes:
return None
episode_lines = "\n".join(f"- {episode}" for episode in episodes)
messages = [
{"role": "system", "content": TASKGEN_SYSTEM},
{"role": "user", "content": f"Recent conversation summaries in #{channel}:\n{episode_lines}"},
]
try:
result = await openai_chat(
self.client, model=self.config["memory-model"], messages=messages, response_format=TASKGEN_RESPONSE_FORMAT
)
self._record_usage(result)
return json.loads(result.choices[0].message.content)
except Exception as err:
logging.warning(f"task proposal failed: {repr(err)}")
return None
async def classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]: 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():
+93 -1
View File
@@ -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 = 3 SCHEMA_VERSION = 5
class PersistentStore: class PersistentStore:
@@ -60,6 +60,18 @@ class PersistentStore:
) )
# Legacy single-string memories carry over as one episode each (MEM-08) # Legacy single-string memories carry over as one episode each (MEM-08)
conn.execute("INSERT INTO episodes (channel, summary) SELECT channel, content FROM memory") conn.execute("INSERT INTO episodes (channel, summary) SELECT channel, content FROM memory")
if version < 4:
conn.execute(
"CREATE TABLE IF NOT EXISTS images (id INTEGER PRIMARY KEY, sha256 TEXT UNIQUE NOT NULL, channel TEXT 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')))"
)
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)
@@ -189,6 +201,86 @@ 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) ---
def image_add(self, sha256: str, channel: str, user: str, message_id: Optional[str], ext: str, nbytes: int) -> None:
with closing(self._connect()) as conn, conn:
conn.execute(
"INSERT OR IGNORE INTO images (sha256, channel, user, message_id, ext, bytes) VALUES (?, ?, ?, ?, ?, ?)",
(sha256, channel, user, message_id, ext, nbytes),
)
def images_recent(self, channel: str, count: int) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute(
"SELECT sha256, user, ext FROM images WHERE channel = ? ORDER BY id DESC LIMIT ?", (channel, count)
).fetchall()
return [{"sha256": row[0], "user": row[1], "ext": row[2]} for row in rows]
def images_total_bytes(self) -> int:
with closing(self._connect()) as conn:
row = conn.execute("SELECT COALESCE(SUM(bytes), 0) FROM images").fetchone()
return int(row[0])
def images_oldest(self, count: int) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute("SELECT sha256, ext, bytes FROM images ORDER BY id LIMIT ?", (count,)).fetchall()
return [{"sha256": row[0], "ext": row[1], "bytes": row[2]} for row in rows]
def images_expired(self, ttl_days: int) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute(
"SELECT sha256, ext FROM images WHERE created_at < datetime('now', ?)", (f"-{int(ttl_days)} days",)
).fetchall()
return [{"sha256": row[0], "ext": row[1]} for row in rows]
def images_delete(self, sha256: str) -> None:
with closing(self._connect()) as conn, conn:
conn.execute("DELETE FROM images WHERE sha256 = ?", (sha256,))
def images_for_user(self, user: str) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute("SELECT sha256, ext FROM images WHERE user = ?", (user,)).fetchall()
return [{"sha256": row[0], "ext": row[1]} for row in rows]
def images_for_message(self, message_id: str) -> List[Dict[str, Any]]:
with closing(self._connect()) as conn:
rows = conn.execute("SELECT sha256, ext FROM images WHERE message_id = ?", (message_id,)).fetchall()
return [{"sha256": row[0], "ext": row[1]} for row in rows]
def purge_user_memory(self, user: str) -> int: def purge_user_memory(self, user: str) -> int:
"""Facts, observations and episode traces of one user (MEM-09).""" """Facts, observations and episode traces of one user (MEM-09)."""
removed = 0 removed = 0
+135
View File
@@ -0,0 +1,135 @@
"""Self-tasking engine (SPEC-005, FDB-011).
Generators propose, the scheduler executes — through the injected
execute callback (= the normal responder path), so budget, gates and
kill-switches all apply. Everything is off unless `tasks-enabled` is
true (TSK-03).
"""
import asyncio
import logging
import time
from typing import Any, Awaitable, Callable, Dict, List, Optional
from .persistence import PersistentStore
from .quota import QuotaLedger
DEFAULT_MAX_PER_CHANNEL_PER_DAY = 2
DEFAULT_IDLE_IMPULSE_HOURS = 12.0
DEFAULT_TASKGEN_INTERVAL_HOURS = 6.0
DEFAULT_BORENESS_PROMPT = "Pretend that you just now thought of something, be creative."
ExecuteCallback = Callable[[str, str], Awaitable[None]]
ProposeCallback = Callable[[], Awaitable[Optional[Dict[str, Any]]]]
AlertCallback = Callable[[str], Awaitable[None]]
class TaskEngine:
def __init__(
self,
store: Optional[PersistentStore],
ledger: QuotaLedger,
config_getter: Callable[[], Dict[str, Any]],
execute: ExecuteCallback,
propose: ProposeCallback,
staff_alert: AlertCallback,
allowed: Callable[[], bool],
idle_seconds: Callable[[], float],
observe: Callable[[str, str, str], Awaitable[None]],
) -> None:
self.store = store
self.ledger = ledger
self._config = config_getter
self._execute = execute
self._propose = propose
self._staff_alert = staff_alert
self._allowed = allowed
self._idle_seconds = idle_seconds
self._observe = observe
self._last_generation = 0.0
self._lock = asyncio.Lock()
def active(self) -> bool:
return self.store is not None and bool(self._config().get("tasks-enabled", False)) # TSK-03
def generators(self) -> List[str]:
return list(self._config().get("tasks-generators", ["idle-impulse", "follow-up"]))
async def enqueue(self, kind: str, channel: str, payload: str, due_hours: float = 0.0) -> int:
assert self.store is not None
state = "approval" if self._config().get("tasks-approval", False) else "queued" # TSK-05
due_at = time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime(time.time() + due_hours * 3600.0))
task_id = int(await asyncio.to_thread(self.store.task_add, kind, channel, due_at, payload, state) or 0)
if state == "approval":
await self._staff_alert(f"Task #{task_id} proposed ({kind}, #{channel}): {payload[:180]} — !bot task-approve {task_id}")
return task_id
def _channel_cap_ok(self, channel: str) -> bool:
cap = int(self._config().get("tasks-max-per-channel-per-day", DEFAULT_MAX_PER_CHANNEL_PER_DAY))
return self.ledger._get(f"task-runs:{channel}") < cap # TSK-04
async def tick(self) -> None:
"""One scheduler pass: execute due tasks, then maybe generate (TSK-02/06)."""
if not self.active() or not self._allowed() or self._lock.locked():
return
assert self.store is not None
async with self._lock:
await self._maybe_generate()
for task in await asyncio.to_thread(self.store.tasks_due):
if not self._channel_cap_ok(task["channel"]):
continue # stays queued for tomorrow (TSK-04)
try:
await self._execute(task["channel"], task["payload"])
await asyncio.to_thread(self.store.task_set_state, task["id"], "done")
self.ledger._add(f"task-runs:{task['channel']}", 1)
await self._observe("system", "task", f"executed {task['kind']}: {task['payload'][:200]}")
except Exception as err:
logging.warning(f"task {task['id']} failed: {repr(err)}")
await asyncio.to_thread(self.store.task_set_state, task["id"], "failed")
async def _maybe_generate(self) -> None:
config = self._config()
interval = float(config.get("taskgen-interval-hours", DEFAULT_TASKGEN_INTERVAL_HOURS)) * 3600.0
now = time.monotonic()
if self._last_generation and now - self._last_generation < interval:
return
self._last_generation = now
generators = self.generators()
if "idle-impulse" in generators:
await self._generate_idle_impulse()
if "follow-up" in generators:
await self._generate_follow_up()
async def _generate_idle_impulse(self) -> None:
"""Boreness, demoted to a deterministic generator (TSK-07)."""
assert self.store is not None
config = self._config()
channel = config.get("chat-channel")
if not channel:
return
idle_threshold = float(config.get("idle-impulse-hours", DEFAULT_IDLE_IMPULSE_HOURS)) * 3600.0
if self._idle_seconds() < idle_threshold:
return
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "idle-impulse"):
return # one pending impulse is enough
prompt = config.get("boreness-prompt", DEFAULT_BORENESS_PROMPT)
await self.enqueue("idle-impulse", channel, prompt)
async def _generate_follow_up(self) -> None:
"""Ask the model for one follow-up worth doing (TSK-08)."""
assert self.store is not None
if await asyncio.to_thread(self.store.tasks_pending_of_kind, "follow-up"):
return
proposal = await self._propose()
if not proposal or not isinstance(proposal.get("task"), dict):
return
task = proposal["task"]
channel = str(task.get("channel") or self._config().get("chat-channel") or "")
prompt = str(task.get("prompt") or "").strip()
if not channel or not prompt:
return
try:
due_hours = max(0.0, float(task.get("due_hours") or 0.0))
except (TypeError, ValueError):
due_hours = 0.0
await self.enqueue("follow-up", channel, prompt, due_hours)
+2 -2
View File
@@ -6,8 +6,8 @@ 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. |
+58
View File
@@ -37,3 +37,61 @@ The translate-before-draw step is deleted: the model's picture prompt
reaches the image API verbatim (current image models handle reaches the image API verbatim (current image models handle
Norwegian/German natively). The `translate()` method and its Norwegian/German natively). The `translate()` method and its
`fix-model` dependency are gone (closes D-009). `fix-model` dependency are gone (closes D-009).
## Input pipeline (FDB-010)
Attachments live in a content-hash cache
(`<history-directory>/images/<sha256>.<ext>`, index in the store,
schema v4). Active only with a store; without one the legacy CDN-URL
path remains.
### IMG-10 — Attachments are ingested at message time (coverage: test)
Every image attachment is downloaded immediately (timeout, size cap
`image-max-bytes` default 8 MB) and stored under its content hash.
Only sniffed png/jpeg/gif/webp bytes are accepted — extension and
declared MIME are ignored (attacker-controlled). Rejected content is
dropped and logged (D11 root fix + cache-abuse hardening).
### IMG-11 — Vision reads from the cache, never CDN URLs (coverage: test)
Vision parts are `data:` URLs built from cached bytes. Discord's
signed, expiring CDN URLs never reach the model or the history.
### IMG-12 — The cache is capped and aged (coverage: test)
`image-cache-mb` (default 500) LRU-evicts oldest-first;
`image-cache-ttl-days` (default 90) ages entries out. Eviction always
removes file and index row together.
### IMG-13 — picture_edit edits the newest channel images (coverage: test)
`picture_edit=true` calls `images.edit` with up to the 4 newest
cached images of the answer channel as inputs (API max is 16; 4 keeps
prompts sane). An empty cache falls back to plain generation — the
flag alone must never fail a reply.
### IMG-14 — Deletion propagates to the cache (coverage: test)
Deleting a Discord message purges its cached images; `!forgetme`
purges all of the user's images — files and rows (extends
SAF-08/MEM-09).
### IMG-15 — Generated images join the cache (coverage: test)
Bot-generated images are ingested like uploads (user `assistant`), so
"make a variant of that" remix chains work on the bot's own output.
### IMG-17 — Image-only messages are cached (coverage: test)
A message consisting only of attachments (no text) is ingested into
the cache and recorded as an observation, even though no reply is
produced — the image must be available for later `picture_edit` and
vision follow-ups. (Previously the empty-text early-return dropped
such posts entirely.)
### IMG-16 — The prompt announces editable images (coverage: test)
When the answer channel has cached images, the context suffix states
how many and that `picture_edit=true` edits the newest — the model
cannot use a capability it does not know about.
+59
View File
@@ -0,0 +1,59 @@
# SPEC-005 — Self-tasking
The sigmoid "boreness" loop becomes a persistent task queue (schema
v5): generators propose, a scheduler executes — through the normal
responder path, so every SAF gate, quota and kill-switch applies.
**Default off** (`tasks-enabled`, kitchen stays quiet unless opted
in); generators are individually selectable per persona
(`tasks-generators`). Marked experimental per plan v4.
### TSK-01 — Tasks are persistent queue rows (coverage: test)
A task is a store row: kind, channel, due-at, payload (the prompt the
responder will run), state (`queued`/`approval`/`done`/`cancelled`/
`failed`). Enqueued tasks survive restarts.
### TSK-02 — Due tasks run through the responder path (coverage: test)
The scheduler executes due queued tasks as system messages via the
normal respond flow (inheriting budget, gates, envelope), marks them
`done`/`failed`, and records the outcome as an observation so memory
learns what the bot did on its own.
### TSK-03 — Off by default (coverage: test)
Without `tasks-enabled = true` nothing is generated and nothing is
executed. A restaurant server does not improvise unless asked to.
### TSK-04 — Per-channel daily cap (coverage: test)
At most `tasks-max-per-channel-per-day` (default 2) task executions
per channel per day, counted in the ledger; further due tasks stay
queued for the next day.
### TSK-05 — Approval mode (coverage: test)
With `tasks-approval = true`, generated tasks enter state `approval`
and a staff alert announces them; `!bot task-approve <id>` moves them
to the queue, `!bot task-cancel <id>` kills them (works for queued
tasks too). Staged-rollout path from the review.
### TSK-06 — Kill-switches govern execution (coverage: test)
Task execution respects `bot_initiated_allowed()``!bot tasks off`,
pause, quiet mode and quiet hours all stop the scheduler tick.
### TSK-07 — Idle-impulse generator (boreness, demoted) (coverage: test)
When the chat channel has been idle longer than
`idle-impulse-hours` (default 12), the generator enqueues one
impulse task with the configured boreness prompt — deterministic
threshold instead of the old 7-second sigmoid dice loop, bounded by
TSK-04. The `on_boreness` loop is gone.
### TSK-08 — Follow-up generator proposes from memory (coverage: test)
Periodically (`taskgen-interval-hours`, default 6) the follow-up
generator asks `memory-model` (strict structured outputs) whether the
recent episodes/facts warrant one follow-up task (channel, prompt,
due-in-hours). A null/failed proposal enqueues nothing.
+6
View File
@@ -62,6 +62,12 @@ Bot-initiated posts also respect pause/quiet.
`!bot pins` answers with all pinned facts and their ids (global + `!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
+212
View File
@@ -0,0 +1,212 @@
"""Unit coverage for SPEC-004 input pipeline (IMG-10..16)."""
import base64
import sqlite3
import tempfile
import unittest
from pathlib import Path
from unittest.mock import AsyncMock, MagicMock, Mock, patch
from fjerkroa_bot.ai_responder import AIMessage, AIResponse
from fjerkroa_bot.images import ImageCache, sniff_ext
from fjerkroa_bot.openai_responder import OpenAIResponder
from fjerkroa_bot.persistence import PersistentStore
from .test_bdd_envelope import FakeModelResponder
from .test_spec_ops import OpsBase
PNG = b"\x89PNG\r\n\x1a\n" + b"x" * 64
def make_cache(tmp, config=None):
store = PersistentStore(Path(tmp) / "bot.db")
cache = ImageCache(store, Path(tmp) / "images", lambda: config or {})
return store, cache
class TestIngest(unittest.TestCase):
def test_sniffed_types_only(self):
"""IMG-10: magic bytes decide; garbage and foreign types are rejected."""
self.assertEqual(sniff_ext(PNG), "png")
self.assertEqual(sniff_ext(b"\xff\xd8\xff\xe0rest"), "jpg")
self.assertIsNone(sniff_ext(b"MZ\x90\x00 definitely-an-exe"))
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp)
self.assertIsNone(cache.ingest_bytes(b"not an image", "chat", "alice", "1"))
sha = cache.ingest_bytes(PNG, "chat", "alice", "1")
self.assertIsNotNone(sha)
self.assertTrue((Path(tmp) / "images" / f"{sha}.png").exists())
self.assertEqual(store.images_recent("chat", 5)[0]["sha256"], sha)
def test_size_cap(self):
"""IMG-10: oversized uploads are dropped."""
with tempfile.TemporaryDirectory() as tmp:
_, cache = make_cache(tmp, {"image-max-bytes": 32})
self.assertIsNone(cache.ingest_bytes(PNG, "chat", "alice", "1"))
class TestVisionDataUrls(OpsBase):
async def test_attachment_becomes_data_url(self):
"""IMG-11: the model sees a data: URL, never the CDN link."""
with tempfile.TemporaryDirectory() as tmp:
_, cache = make_cache(tmp)
self.bot.airesponder.image_cache = cache
self.bot.respond = AsyncMock()
message = self.public_msg("look at this")
attachment = Mock()
attachment.url = "https://cdn.discordapp.com/attachments/1/2/cat.png?ex=deadbeef"
message.attachments = [attachment]
message.id = 42
with patch.object(ImageCache, "_download", new_callable=AsyncMock, return_value=PNG):
await self.bot.on_message(message)
sent_msg = self.bot.respond.await_args.args[0]
self.assertTrue(sent_msg.urls[0].startswith("data:image/png;base64,"))
self.assertNotIn("cdn.discordapp.com", sent_msg.urls[0])
class TestEviction(unittest.TestCase):
def test_lru_cap(self):
"""IMG-12: byte cap evicts oldest first, file + row together."""
big = b"\x89PNG\r\n\x1a\n" + b"a" * (700 * 1024)
big2 = b"\x89PNG\r\n\x1a\n" + b"b" * (700 * 1024)
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp, {"image-cache-mb": 1})
first = cache.ingest_bytes(big, "chat", "alice", "1")
second = cache.ingest_bytes(big2, "chat", "alice", "2")
shas = [row["sha256"] for row in store.images_recent("chat", 5)]
self.assertNotIn(first, shas)
self.assertIn(second, shas)
self.assertFalse((Path(tmp) / "images" / f"{first}.png").exists())
def test_ttl(self):
"""IMG-12: entries past image-cache-ttl-days age out."""
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp, {"image-cache-ttl-days": 30})
sha = cache.ingest_bytes(PNG, "chat", "alice", "1")
with sqlite3.connect(store.db_path) as conn:
conn.execute("UPDATE images SET created_at = datetime('now', '-60 days') WHERE sha256 = ?", (sha,))
cache.evict()
self.assertEqual(store.images_recent("chat", 5), [])
self.assertFalse((Path(tmp) / "images" / f"{sha}.png").exists())
class TestEditPath(OpsBase):
async def prepare(self, with_images):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
_, cache = make_cache(self.tmp.name)
self.bot.airesponder.image_cache = cache
if with_images:
cache.ingest_bytes(PNG, "chat", "alice", "1")
self.bot.airesponder.edit_openai = AsyncMock(return_value=[__import__("io").BytesIO(PNG)])
self.bot.airesponder.draw = AsyncMock(return_value=[__import__("io").BytesIO(PNG)])
response = AIResponse("her", True, "chat", None, "als wikinger", True, False)
channel = MagicMock()
channel.name = "chat"
channel.send = AsyncMock()
channel.typing = MagicMock(return_value=AsyncMock(__aenter__=AsyncMock(), __aexit__=AsyncMock()))
await self.bot.send_answer_with_typing(response, channel, self.bot.airesponder, factual=True)
async def test_edit_uses_cached_sources(self):
"""IMG-13: picture_edit + cached images -> images.edit path."""
await self.prepare(with_images=True)
self.bot.airesponder.edit_openai.assert_awaited_once()
self.bot.airesponder.draw.assert_not_awaited()
async def test_empty_cache_falls_back_to_generate(self):
"""IMG-13: empty cache -> plain generation, the flag never fails a reply."""
await self.prepare(with_images=False)
self.bot.airesponder.edit_openai.assert_not_awaited()
self.bot.airesponder.draw.assert_awaited_once()
class TestPurges(OpsBase):
async def test_message_delete_and_forgetme_purge_images(self):
"""IMG-14: message deletion and !forgetme remove files + rows."""
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp)
self.bot.airesponder.image_cache = cache
cache.ingest_bytes(PNG, "chat", "alice", "99")
deleted = MagicMock()
deleted.id = 99
deleted.content = "pic"
deleted.author.name = "alice"
deleted.channel = MagicMock()
await self.bot.on_message_delete(deleted)
self.assertEqual(store.images_recent("chat", 5), [])
cache.ingest_bytes(b"\x89PNG\r\n\x1a\n" + b"z" * 32, "chat", "alice", "100")
message = self.public_msg("!forgetme")
message.author.name = "alice"
await self.bot.on_message(message)
self.assertEqual(store.images_recent("chat", 5), [])
class TestGeneratedImagesCached(OpsBase):
async def test_bot_output_joins_cache(self):
"""IMG-15: generated images are ingested as user 'assistant'."""
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp)
self.bot.airesponder.image_cache = cache
self.bot.airesponder.draw = AsyncMock(return_value=[__import__("io").BytesIO(PNG)])
response = AIResponse("her", True, "chat", None, "en katt", False, False)
channel = MagicMock()
channel.name = "chat"
channel.send = AsyncMock()
await self.bot.send_answer_with_typing(response, channel, self.bot.airesponder, factual=True)
rows = store.images_recent("chat", 5)
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["user"], "assistant")
class TestImageOnlyMessages(OpsBase):
async def test_image_only_post_cached_no_reply(self):
"""IMG-17: attachment without text -> cached + observed, no reply."""
with tempfile.TemporaryDirectory() as tmp:
store, cache = make_cache(tmp)
self.bot.airesponder.image_cache = cache
self.bot.airesponder.observe_event = AsyncMock()
self.bot.respond = AsyncMock()
message = self.public_msg("")
message.content = ""
message.channel.name = "chat"
attachment = Mock()
attachment.url = "https://cdn.discordapp.com/attachments/1/2/silent.png"
message.attachments = [attachment]
message.id = 77
with patch.object(ImageCache, "_download", new_callable=AsyncMock, return_value=PNG):
await self.bot.on_message(message)
self.assertEqual(len(store.images_recent("chat", 5)), 1)
self.bot.airesponder.observe_event.assert_awaited_once()
self.bot.respond.assert_not_awaited()
class TestContextAnnouncesImages(unittest.IsolatedAsyncioTestCase):
def test_suffix_mentions_picture_edit(self):
"""IMG-16: cached channel images are announced in the context suffix."""
with tempfile.TemporaryDirectory() as tmp:
config = {"system": "s", "history-limit": 5, "history-directory": tmp}
responder = FakeModelResponder(config, "chat")
responder.image_cache.ingest_bytes(PNG, "chat", "alice", "1")
system = responder.message(AIMessage("alice", "hei", "chat"))[0]["content"]
self.assertIn("picture_edit", system)
self.assertIn("recent images in this channel: 1", system)
class TestEditOpenai(unittest.IsolatedAsyncioTestCase):
async def test_edit_call_shape_and_metering(self):
"""IMG-13: images.edit gets the file handles, n clamped, ledger counts."""
responder = OpenAIResponder({"openai-token": "t", "model": "m", "system": "s", "history-limit": 5}, "chat")
with tempfile.TemporaryDirectory() as tmp:
paths = []
for index in range(2):
path = Path(tmp) / f"in{index}.png"
path.write_bytes(PNG)
paths.append(path)
api_result = Mock(data=[Mock(b64_json=base64.b64encode(b"out").decode())])
with patch("fjerkroa_bot.openai_responder.openai_image_edit", new_callable=AsyncMock) as edit_mock:
edit_mock.return_value = api_result
buffers = await responder.edit_openai("wikinger", paths, 9)
self.assertEqual(buffers[0].read(), b"out")
self.assertEqual(edit_mock.await_args.kwargs["n"], 4)
self.assertEqual(len(edit_mock.await_args.kwargs["image"]), 2)
self.assertEqual(responder.ledger.images_today(), 1)
+165
View File
@@ -0,0 +1,165 @@
"""Unit coverage for SPEC-005 self-tasking (TSK-01..08) + OPS-12."""
import tempfile
import unittest
from pathlib import Path
from unittest.mock import AsyncMock
from fjerkroa_bot.persistence import PersistentStore
from fjerkroa_bot.quota import QuotaLedger
from fjerkroa_bot.tasks import TaskEngine
from .test_spec_ops import OpsBase
class EngineBase(unittest.IsolatedAsyncioTestCase):
def setUp(self):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
self.store = PersistentStore(Path(self.tmp.name) / "bot.db")
self.config = {"tasks-enabled": True, "chat-channel": "chat", "taskgen-interval-hours": 0}
self.ledger = QuotaLedger(self.store, lambda: self.config)
self.execute = AsyncMock()
self.propose = AsyncMock(return_value={"task": None})
self.alert = AsyncMock()
self.observe = AsyncMock()
self.idle = lambda: 0.0
self.engine = TaskEngine(
self.store,
self.ledger,
lambda: self.config,
self.execute,
self.propose,
self.alert,
lambda: True,
lambda: self.idle(),
self.observe,
)
class TestQueuePersistence(EngineBase):
async def test_tasks_survive_restart(self):
"""TSK-01: enqueued tasks are store rows and reload in a fresh store."""
await self.engine.enqueue("follow-up", "chat", "frag bob nach der pruefung", due_hours=1)
reborn = PersistentStore(Path(self.tmp.name) / "bot.db")
tasks = reborn.tasks_open()
self.assertEqual(len(tasks), 1)
self.assertEqual((tasks[0]["kind"], tasks[0]["channel"], tasks[0]["state"]), ("follow-up", "chat", "queued"))
class TestExecution(EngineBase):
async def test_due_task_executes_and_completes(self):
"""TSK-02: due task runs via execute callback, marked done, observed."""
await self.engine.enqueue("idle-impulse", "chat", "sag was nettes")
await self.engine.tick()
self.execute.assert_awaited_once_with("chat", "sag was nettes")
self.assertEqual(self.store.tasks_open(), [])
self.observe.assert_awaited()
async def test_failed_execution_marked_failed(self):
"""TSK-02: execute raising marks the task failed, not done."""
self.execute.side_effect = RuntimeError("channel gone")
await self.engine.enqueue("idle-impulse", "chat", "x")
await self.engine.tick()
self.assertEqual(self.store.tasks_open(), []) # not queued anymore
rows = [t for t in self.store.tasks_due()]
self.assertEqual(rows, [])
class TestDefaultOff(EngineBase):
async def test_disabled_engine_is_dormant(self):
"""TSK-03: without tasks-enabled nothing executes or generates."""
self.store.task_add("idle-impulse", "chat", "2020-01-01 00:00:00", "x")
self.config.pop("tasks-enabled")
await self.engine.tick()
self.execute.assert_not_awaited()
self.propose.assert_not_awaited()
class TestDailyCap(EngineBase):
async def test_cap_defers_excess_tasks(self):
"""TSK-04: over the per-channel cap tasks stay queued."""
self.config["tasks-max-per-channel-per-day"] = 1
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "one")
self.store.task_add("b", "chat", "2020-01-01 00:00:00", "two")
await self.engine.tick()
self.assertEqual(self.execute.await_count, 1)
self.assertEqual(len(self.store.tasks_open()), 1)
class TestApproval(EngineBase):
async def test_approval_flow(self):
"""TSK-05: approval mode holds tasks until approved; cancel kills them."""
self.config["tasks-approval"] = True
task_id = await self.engine.enqueue("follow-up", "chat", "frag nach")
self.alert.assert_awaited_once()
await self.engine.tick()
self.execute.assert_not_awaited() # approval != due
self.store.task_set_state(task_id, "queued")
await self.engine.tick()
self.execute.assert_awaited_once()
class TestKillSwitch(EngineBase):
async def test_not_allowed_blocks_tick(self):
"""TSK-06: bot_initiated_allowed()=False stops execution."""
self.store.task_add("a", "chat", "2020-01-01 00:00:00", "x")
engine = TaskEngine(
self.store, self.ledger, lambda: self.config, self.execute, self.propose, self.alert, lambda: False, lambda: 0.0, self.observe
)
await engine.tick()
self.execute.assert_not_awaited()
class TestIdleImpulse(EngineBase):
async def test_idle_enqueues_once(self):
"""TSK-07: long idle enqueues one impulse; pending impulse dedupes."""
self.config["idle-impulse-hours"] = 1
self.config["boreness-prompt"] = "denk dir was aus"
self.idle = lambda: 2 * 3600.0
await self.engine.tick()
open_tasks = self.store.tasks_open()
impulse = [t for t in open_tasks if t["kind"] == "idle-impulse"]
self.assertEqual(len(impulse), 0) # executed immediately (due now)
self.execute.assert_awaited_once_with("chat", "denk dir was aus")
async def test_short_idle_no_impulse(self):
"""TSK-07: below the idle threshold nothing is generated."""
self.config["idle-impulse-hours"] = 12
self.idle = lambda: 60.0
await self.engine.tick()
self.execute.assert_not_awaited()
class TestFollowUpGenerator(EngineBase):
async def test_proposal_becomes_task(self):
"""TSK-08: a proposed task is enqueued with its due offset."""
self.config.pop("chat-channel")
self.propose.return_value = {"task": {"channel": "chat", "prompt": "frag bob", "due_hours": 24}}
await self.engine.tick()
tasks = self.store.tasks_open()
self.assertEqual(len(tasks), 1)
self.assertEqual(tasks[0]["kind"], "follow-up")
async def test_null_proposal_no_task(self):
"""TSK-08: null proposal enqueues nothing."""
self.propose.return_value = {"task": None}
await self.engine.tick()
self.assertEqual(self.store.tasks_open(), [])
class TestTaskCommands(OpsBase):
async def test_list_approve_cancel(self):
"""OPS-12: !bot tasks lists; task-approve/task-cancel manage states."""
with tempfile.TemporaryDirectory() as tmp:
store = PersistentStore(Path(tmp) / "bot.db")
self.bot.airesponder.store = store
task_id = store.task_add("follow-up", "chat", "2099-01-01 00:00:00", "frag bob", "approval")
await self.bot.on_message(self.staff_msg("!bot tasks"))
listing = self.bot.staff_channel.send.await_args.args[0]
self.assertIn("follow-up", listing)
self.assertIn("approval", listing)
await self.bot.on_message(self.staff_msg(f"!bot task-approve {task_id}"))
self.assertEqual(store.tasks_open()[0]["state"], "queued")
await self.bot.on_message(self.staff_msg(f"!bot task-cancel {task_id}"))
self.assertEqual(store.tasks_open(), [])