self-tasking engine: persistent queue, idle-impulse + follow-up generators, approval mode
This commit is contained in:
+53
-31
@@ -1,7 +1,6 @@
|
||||
import argparse
|
||||
import asyncio
|
||||
import logging
|
||||
import math
|
||||
import random
|
||||
import re
|
||||
import sys
|
||||
@@ -18,6 +17,7 @@ from watchdog.observers import Observer
|
||||
|
||||
from .ai_responder import AIMessage
|
||||
from .openai_responder import OpenAIResponder
|
||||
from .tasks import TaskEngine
|
||||
|
||||
DEFAULT_PRIVACY_NOTICE = (
|
||||
"I keep recent channel messages and a short conversation summary to answer better. "
|
||||
@@ -114,38 +114,42 @@ class FjerkroaBot(commands.Bot):
|
||||
self.staff_channel = self.channel_by_name(self.config["staff-channel"], no_ignore=True)
|
||||
self.welcome_channel = self.channel_by_name(self.config["welcome-channel"], no_ignore=True)
|
||||
|
||||
def init_boreness(self):
|
||||
if "chat-channel" not in self.config:
|
||||
return
|
||||
def init_tasks(self):
|
||||
"""Task engine replaces the sigmoid boreness loop (TSK-07)."""
|
||||
self.last_activity_time = time.monotonic()
|
||||
self.loop.create_task(self.on_boreness())
|
||||
logging.info("Boreness initialised.")
|
||||
self.task_engine = TaskEngine(
|
||||
store=self.airesponder.store,
|
||||
ledger=self.airesponder.ledger,
|
||||
config_getter=lambda: self.config,
|
||||
execute=self._execute_task,
|
||||
propose=self.airesponder.propose_task,
|
||||
staff_alert=self.send_staff_alert,
|
||||
allowed=self.bot_initiated_allowed,
|
||||
idle_seconds=lambda: time.monotonic() - self.last_activity_time,
|
||||
observe=self.airesponder.observe_event,
|
||||
)
|
||||
self.loop.create_task(self.task_loop())
|
||||
logging.info("Task engine initialised.")
|
||||
|
||||
async def on_boreness(self):
|
||||
logging.info(f"Boreness started on channel: {repr(self.chat_channel)}")
|
||||
async def task_loop(self):
|
||||
while True:
|
||||
if self.chat_channel is None or not self.bot_initiated_allowed():
|
||||
await asyncio.sleep(7)
|
||||
continue
|
||||
boreness_interval = float(self.config.get("boreness-interval", 12.0))
|
||||
elapsed_time = (time.monotonic() - self.last_activity_time) / 3600.0
|
||||
probability = 1 / (1 + math.exp(-1 * (elapsed_time - (boreness_interval / 2.0)) + math.log(1 / 0.2 - 1)))
|
||||
if random.random() < probability:
|
||||
prev_messages = [msg async for msg in self.chat_channel.history(limit=2)]
|
||||
last_author = prev_messages[1].author.id if len(prev_messages) > 1 else None
|
||||
if last_author and last_author != self.user.id:
|
||||
logging.info(f"Borred with {probability} probability after {elapsed_time}")
|
||||
boreness_prompt = self.config.get("boreness-prompt", "Pretend that you just now thought of something, be creative.")
|
||||
message = AIMessage("system", boreness_prompt, self.config.get("chat-channel", "chat"), True, False)
|
||||
try:
|
||||
await self.respond(message, self.chat_channel)
|
||||
except Exception as err:
|
||||
logging.warning(f"Failed to activate borringness: {repr(err)}")
|
||||
await asyncio.sleep(7)
|
||||
await asyncio.sleep(60)
|
||||
try:
|
||||
await self.task_engine.tick()
|
||||
except Exception as err:
|
||||
logging.warning(f"task tick failed: {repr(err)}")
|
||||
|
||||
async def _execute_task(self, channel_name: str, prompt: str) -> None:
|
||||
"""Run a due task through the normal responder path (TSK-02)."""
|
||||
channel = self.channel_by_name(channel_name, getattr(self, "chat_channel", None), no_ignore=True)
|
||||
if channel is None:
|
||||
raise RuntimeError(f"task channel {channel_name!r} not resolvable")
|
||||
message = AIMessage("system", prompt, channel_name, True, False)
|
||||
await self.respond(message, channel)
|
||||
|
||||
async def on_ready(self):
|
||||
self.init_channels()
|
||||
self.init_boreness()
|
||||
self.init_tasks()
|
||||
logging.info(
|
||||
f"We have logged in as {self.user}" f" ({repr(self.staff_channel)}, {repr(self.welcome_channel)}, {repr(self.chat_channel)})"
|
||||
)
|
||||
@@ -241,13 +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 None
|
||||
|
||||
def _task_command(self, args) -> Optional[str]:
|
||||
"""Task queue surface (OPS-12, TSK-05)."""
|
||||
is_list = args[:1] == ["tasks"] and len(args) == 1
|
||||
if args[:1] not in (["task-approve"], ["task-cancel"]) and not is_list:
|
||||
return None
|
||||
store = self.airesponder.store
|
||||
if store is None:
|
||||
return "No store configured - task commands unavailable."
|
||||
if is_list:
|
||||
tasks = store.tasks_open()
|
||||
return "\n".join(f"{t['id']} [{t['state']}] {t['kind']} #{t['channel']} due {t['due_at']}" for t in tasks) or "No open tasks."
|
||||
if args[1:2] and args[1].isdigit():
|
||||
if args[0] == "task-approve":
|
||||
return f"Approved {store.task_set_state(int(args[1]), 'queued')} task(s)."
|
||||
return f"Cancelled {store.task_set_state(int(args[1]), 'cancelled')} task(s)."
|
||||
return None
|
||||
|
||||
async def handle_staff_command(self, message: Message) -> None:
|
||||
"""Operator kill-switches, staff channel only (OPS-01..05, OPS-09, MEM-07)."""
|
||||
args = str(message.content).split()[1:]
|
||||
memory_reply = self._memory_command(args)
|
||||
if memory_reply is not None:
|
||||
await message.channel.send(memory_reply, suppress_embeds=True)
|
||||
return
|
||||
for handler in (self._memory_command, self._task_command):
|
||||
reply = handler(args)
|
||||
if reply is not None:
|
||||
await message.channel.send(reply, suppress_embeds=True)
|
||||
return
|
||||
reply = "Commands: pause, resume, images on|off, tasks on|off, quiet <minutes>, status, spend, memory <user>, forget-fact <id>, pin <channel|global> <fact>, unpin <id>"
|
||||
if args[:1] == ["pause"]:
|
||||
self.replies_enabled = False
|
||||
|
||||
@@ -58,6 +58,31 @@ CONSOLIDATION_RESPONSE_FORMAT = {
|
||||
"type": "json_schema",
|
||||
"json_schema": {"name": "consolidation", "strict": True, "schema": CONSOLIDATION_SCHEMA},
|
||||
}
|
||||
# Follow-up task proposal (SPEC-005 TSK-08): one task or null
|
||||
TASKGEN_SCHEMA = {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"task": {
|
||||
"type": ["object", "null"],
|
||||
"properties": {
|
||||
"channel": {"type": ["string", "null"], "description": "Target channel, or null for the main chat channel."},
|
||||
"prompt": {"type": "string", "description": "Instruction the assistant will act on when the task runs."},
|
||||
"due_hours": {"type": "number", "description": "Hours from now until the task should run (0 = now)."},
|
||||
},
|
||||
"required": ["channel", "prompt", "due_hours"],
|
||||
"additionalProperties": False,
|
||||
}
|
||||
},
|
||||
"required": ["task"],
|
||||
"additionalProperties": False,
|
||||
}
|
||||
TASKGEN_RESPONSE_FORMAT = {"type": "json_schema", "json_schema": {"name": "task_proposal", "strict": True, "schema": TASKGEN_SCHEMA}}
|
||||
TASKGEN_SYSTEM = (
|
||||
"You plan the self-initiated actions of a Discord assistant. Given recent conversation summaries, propose AT MOST ONE follow-up"
|
||||
" worth doing on the assistant's own initiative (ask how something announced went, revisit an open question, congratulate on an"
|
||||
" event). Only propose something genuinely worthwhile — when in doubt, return a null task."
|
||||
)
|
||||
|
||||
# Reply/ignore + factual pre-pass (SPEC-010 BEH-01): one cheap call
|
||||
CLASSIFIER_SCHEMA = {
|
||||
"type": "object",
|
||||
@@ -381,6 +406,29 @@ class OpenAIResponder(AIResponder, LeonardoAIDrawMixIn):
|
||||
logging.info(f"edited {len(buffers)} image(s) on {model} from {len(handles)} input(s)")
|
||||
return buffers
|
||||
|
||||
async def propose_task(self) -> Optional[Dict[str, Any]]:
|
||||
"""One follow-up proposal from recent episodes on memory-model (TSK-08)."""
|
||||
if "memory-model" not in self.config or self.store is None or not self.ledger.budget_ok():
|
||||
return None
|
||||
channel = self.config.get("chat-channel", "chat")
|
||||
episodes = await asyncio.to_thread(self.store.recent_episodes, channel, 5)
|
||||
if not episodes:
|
||||
return None
|
||||
episode_lines = "\n".join(f"- {episode}" for episode in episodes)
|
||||
messages = [
|
||||
{"role": "system", "content": TASKGEN_SYSTEM},
|
||||
{"role": "user", "content": f"Recent conversation summaries in #{channel}:\n{episode_lines}"},
|
||||
]
|
||||
try:
|
||||
result = await openai_chat(
|
||||
self.client, model=self.config["memory-model"], messages=messages, response_format=TASKGEN_RESPONSE_FORMAT
|
||||
)
|
||||
self._record_usage(result)
|
||||
return json.loads(result.choices[0].message.content)
|
||||
except Exception as err:
|
||||
logging.warning(f"task proposal failed: {repr(err)}")
|
||||
return None
|
||||
|
||||
async def classify(self, message: Any, history_tail: List[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
|
||||
"""~100-token reply/factual/emoji verdict on classifier-model (BEH-01/03)."""
|
||||
if "classifier-model" not in self.config or not self.ledger.budget_ok():
|
||||
|
||||
@@ -13,7 +13,7 @@ from contextlib import closing
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
SCHEMA_VERSION = 4
|
||||
SCHEMA_VERSION = 5
|
||||
|
||||
|
||||
class PersistentStore:
|
||||
@@ -66,6 +66,12 @@ class PersistentStore:
|
||||
" user TEXT NOT NULL, message_id TEXT, ext TEXT NOT NULL, bytes INTEGER NOT NULL,"
|
||||
" created_at TEXT NOT NULL DEFAULT (datetime('now')))"
|
||||
)
|
||||
if version < 5:
|
||||
conn.execute(
|
||||
"CREATE TABLE IF NOT EXISTS tasks (id INTEGER PRIMARY KEY, kind TEXT NOT NULL, channel TEXT NOT NULL,"
|
||||
" due_at TEXT NOT NULL, payload TEXT NOT NULL, state TEXT NOT NULL DEFAULT 'queued',"
|
||||
" created_at TEXT NOT NULL DEFAULT (datetime('now')), executed_at TEXT)"
|
||||
)
|
||||
if version < SCHEMA_VERSION:
|
||||
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
|
||||
os.chmod(self.db_path, 0o600) # conversation data (PER-04)
|
||||
@@ -195,6 +201,39 @@ class PersistentStore:
|
||||
)
|
||||
return cursor.rowcount
|
||||
|
||||
# --- task queue (SPEC-005, FDB-011) ---
|
||||
|
||||
def task_add(self, kind: str, channel: str, due_at: str, payload: str, state: str = "queued") -> int:
|
||||
with closing(self._connect()) as conn, conn:
|
||||
cursor = conn.execute(
|
||||
"INSERT INTO tasks (kind, channel, due_at, payload, state) VALUES (?, ?, ?, ?, ?)",
|
||||
(kind, channel, due_at, payload, state),
|
||||
)
|
||||
return int(cursor.lastrowid or 0)
|
||||
|
||||
def tasks_due(self) -> List[Dict[str, Any]]:
|
||||
with closing(self._connect()) as conn:
|
||||
rows = conn.execute(
|
||||
"SELECT id, kind, channel, payload FROM tasks WHERE state = 'queued' AND due_at <= datetime('now') ORDER BY due_at"
|
||||
).fetchall()
|
||||
return [{"id": row[0], "kind": row[1], "channel": row[2], "payload": row[3]} for row in rows]
|
||||
|
||||
def tasks_open(self) -> List[Dict[str, Any]]:
|
||||
with closing(self._connect()) as conn:
|
||||
rows = conn.execute(
|
||||
"SELECT id, kind, channel, due_at, state FROM tasks WHERE state IN ('queued', 'approval') ORDER BY id"
|
||||
).fetchall()
|
||||
return [{"id": row[0], "kind": row[1], "channel": row[2], "due_at": row[3], "state": row[4]} for row in rows]
|
||||
|
||||
def task_set_state(self, task_id: int, state: str) -> int:
|
||||
with closing(self._connect()) as conn, conn:
|
||||
return conn.execute("UPDATE tasks SET state = ?, executed_at = datetime('now') WHERE id = ?", (state, task_id)).rowcount
|
||||
|
||||
def tasks_pending_of_kind(self, kind: str) -> int:
|
||||
with closing(self._connect()) as conn:
|
||||
row = conn.execute("SELECT COUNT(*) FROM tasks WHERE kind = ? AND state IN ('queued', 'approval')", (kind,)).fetchone()
|
||||
return int(row[0])
|
||||
|
||||
# --- image cache index (SPEC-004, FDB-010) ---
|
||||
|
||||
def image_add(self, sha256: str, channel: str, user: str, message_id: Optional[str], ext: str, nbytes: int) -> None:
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user