image input pipeline: content-hash cache, data-url vision, picture_edit real
This commit is contained in:
@@ -12,6 +12,7 @@ from pathlib import Path
|
||||
from pprint import pformat
|
||||
from typing import Any, Dict, List, Optional, Tuple, Union
|
||||
|
||||
from .images import ImageCache
|
||||
from .memory import MemoryManager
|
||||
from .persistence import PersistentStore
|
||||
|
||||
@@ -153,17 +154,16 @@ class AIResponder(AIResponderBase):
|
||||
if stored_memory is not None:
|
||||
self.memory = stored_memory
|
||||
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}")
|
||||
|
||||
# Dynamic values move to a context suffix so the persona prefix
|
||||
# stays byte-stable for the prompt cache (ENV-20)
|
||||
DYNAMIC_PLACEHOLDERS = ("{date}", "{time}", "{news}", "{memory}")
|
||||
|
||||
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, "")
|
||||
def _context_lines(self, message: AIMessage) -> List[str]:
|
||||
context = [f"date: {time.strftime('%Y-%m-%d')} ({time.strftime('%A')})", f"time: {time.strftime('%H:%M:%S')}"]
|
||||
news_feed = self.config.get("news")
|
||||
if news_feed and os.path.exists(news_feed):
|
||||
@@ -173,7 +173,19 @@ class AIResponder(AIResponderBase):
|
||||
memory_block = self.memory_manager.memory_block(participants, self.memory)
|
||||
if 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)} (set picture_edit=true to edit/remix the newest)")
|
||||
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})
|
||||
if limit is not None:
|
||||
while len(self.history) > limit:
|
||||
|
||||
@@ -197,6 +197,8 @@ class FjerkroaBot(commands.Bot):
|
||||
removed += self.airesponder.store.delete_history_of_user(user)
|
||||
# facts + observations + episode traces (MEM-09)
|
||||
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}")
|
||||
await message.channel.send(
|
||||
f"Removed your messages, facts and memory traces ({removed} entries).",
|
||||
@@ -337,6 +339,8 @@ class FjerkroaBot(commands.Bot):
|
||||
|
||||
async def on_message_delete(self, message):
|
||||
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}")
|
||||
|
||||
def on_config_file_modified(self, event):
|
||||
@@ -410,14 +414,24 @@ class FjerkroaBot(commands.Bot):
|
||||
msg = AIMessage(
|
||||
message.author.name, message_content, channel_name, self.user in message.mentions or isinstance(message.channel, DMChannel)
|
||||
)
|
||||
airesponder = self.get_ai_responder(channel_name)
|
||||
if message.attachments:
|
||||
for attachment in message.attachments:
|
||||
if not msg.urls:
|
||||
msg.urls = []
|
||||
msg.urls.append(attachment.url)
|
||||
if airesponder.image_cache is not None:
|
||||
# cache-first: CDN URLs never travel further (IMG-10/11)
|
||||
sha = await airesponder.image_cache.ingest_url(attachment.url, channel_name, message.author.name, str(message.id))
|
||||
if sha is not None:
|
||||
recent = airesponder.image_cache.recent(channel_name, 8)
|
||||
ext = next((row["ext"] for row in recent if row["sha256"] == sha), "png")
|
||||
data_url = airesponder.image_cache.data_url(sha, ext)
|
||||
if data_url:
|
||||
msg.urls.append(data_url)
|
||||
else:
|
||||
msg.urls.append(attachment.url)
|
||||
|
||||
# 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)
|
||||
if handled:
|
||||
return
|
||||
@@ -462,7 +476,19 @@ class FjerkroaBot(commands.Bot):
|
||||
"""Send the answer paced, split and with images on the last part (BEH-04/05/06)"""
|
||||
files = 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)]
|
||||
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)
|
||||
|
||||
@@ -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)
|
||||
@@ -96,6 +96,10 @@ async def openai_image(client, *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):
|
||||
def __init__(self, config: Dict[str, Any], channel: Optional[str] = None) -> None:
|
||||
super().__init__(config, channel)
|
||||
@@ -354,6 +358,29 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
||||
logging.debug(f"Full traceback: {traceback.format_exc()}")
|
||||
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 classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
|
||||
"""~100-token reply/factual/emoji verdict on classifier-model (BEH-01/03)."""
|
||||
if "classifier-model" not in self.config or not self.ledger.budget_ok():
|
||||
|
||||
@@ -13,7 +13,7 @@ from contextlib import closing
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
SCHEMA_VERSION = 3
|
||||
SCHEMA_VERSION = 4
|
||||
|
||||
|
||||
class PersistentStore:
|
||||
@@ -60,6 +60,12 @@ class PersistentStore:
|
||||
)
|
||||
# Legacy single-string memories carry over as one episode each (MEM-08)
|
||||
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 < SCHEMA_VERSION:
|
||||
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
|
||||
os.chmod(self.db_path, 0o600) # conversation data (PER-04)
|
||||
@@ -189,6 +195,53 @@ class PersistentStore:
|
||||
)
|
||||
return cursor.rowcount
|
||||
|
||||
# --- 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:
|
||||
"""Facts, observations and episode traces of one user (MEM-09)."""
|
||||
removed = 0
|
||||
|
||||
Reference in New Issue
Block a user