Compare commits
2 Commits
main
..
bf7c3d9093
| Author | SHA1 | Date | |
|---|---|---|---|
| bf7c3d9093 | |||
| 26dae1d101 |
+46
-58
@@ -1,5 +1,4 @@
|
|||||||
# ai command cogs
|
# ai command cogs
|
||||||
import asyncio
|
|
||||||
import logging
|
import logging
|
||||||
import re
|
import re
|
||||||
import sys
|
import sys
|
||||||
@@ -18,7 +17,7 @@ from communication_subroutine import AI_QUERY_Q
|
|||||||
|
|
||||||
import ai_functions
|
import ai_functions
|
||||||
from constants import (
|
from constants import (
|
||||||
OLLAMA_WARM_MINUTES,
|
ASSISTANTS,
|
||||||
DATA,
|
DATA,
|
||||||
GRAPHICS_PATH,
|
GRAPHICS_PATH,
|
||||||
INITIAL_TIME_WAIT,
|
INITIAL_TIME_WAIT,
|
||||||
@@ -102,45 +101,55 @@ class Events(commands.Cog):
|
|||||||
text = text[1900:]
|
text = text[1900:]
|
||||||
|
|
||||||
async def cog_load(self):
|
async def cog_load(self):
|
||||||
# The AI query worker answers via handle_response, so it works on every
|
# The AI query worker must run regardless of the OpenAI guard below - it
|
||||||
# backend. Start it first.
|
# answers via handle_response, which works on Claude too. Start it first.
|
||||||
if not self.ai_query_worker.is_running():
|
if not self.ai_query_worker.is_running():
|
||||||
self.ai_query_worker.start()
|
self.ai_query_worker.start()
|
||||||
# Keeps a self-hosted model resident; it no-ops on any other provider.
|
self.logger.info("Starting personal assistants")
|
||||||
if not self.ollama_warm_loop.is_running():
|
# Personal assistants use the OpenAI Assistants API (threads/runs), which
|
||||||
self.ollama_warm_loop.start()
|
# has no Anthropic equivalent - skip cleanly when OpenAI isn't wired up
|
||||||
# NOTE: there is no OpenAI-Assistants bootstrap any more. It called a
|
# (e.g. a Claude-only deployment) instead of crashing the cog load.
|
||||||
# sunset API (beta threads), 404'd, and failed the WHOLE extension -
|
if OPENAICLIENT is None:
|
||||||
# taking every AI command with it. Personal assistants now ride
|
self.logger.warning(
|
||||||
# handle_response with per-user memory (ai_functions), so they work on
|
"OPENAICLIENT niedostępny - osobiści asystenci (OpenAI Assistants API) wyłączeni"
|
||||||
# Claude and Ollama too and nothing has to be created at startup.
|
)
|
||||||
self.logger.info("Osobiści asystenci: pamięć per-user, aktywny backend AI")
|
|
||||||
|
|
||||||
@tasks.loop(minutes=OLLAMA_WARM_MINUTES)
|
|
||||||
async def ollama_warm_loop(self):
|
|
||||||
"""Keep a self-hosted model resident so users don't pay the load wait.
|
|
||||||
|
|
||||||
Loading is the slow part on a GPU shared with other users, so we
|
|
||||||
re-assert Ollama's keep_alive well inside its window. This preloads
|
|
||||||
WITHOUT generating - no tokens, no cost.
|
|
||||||
|
|
||||||
Hard guard: it does nothing unless the ACTIVE backend is Ollama. Firing
|
|
||||||
warm-ups at a metered API would burn tokens and money for nothing.
|
|
||||||
"""
|
|
||||||
try:
|
|
||||||
if ai_functions.active_provider() != "ollama":
|
|
||||||
return
|
return
|
||||||
await ai_functions.warm_active_model()
|
for superfryta_id, superfryta in SPECJALNE_ZIEMNIACZKI.items():
|
||||||
except Exception as exc: # pylint: disable=broad-exception-caught
|
|
||||||
self.logger.info("Rozgrzewanie Ollamy nieudane (nieszkodliwe): %s", exc)
|
|
||||||
|
|
||||||
@ollama_warm_loop.before_loop
|
if superfryta[4] != "":
|
||||||
async def before_ollama_warm_loop(self):
|
self.logger.info(
|
||||||
await self.bot.wait_until_ready()
|
"Personal assistant for user: %s, exists id: %s,name: %s, owner: %s, special instructions: %s assistant id: %s ",
|
||||||
|
superfryta_id,
|
||||||
|
superfryta[0],
|
||||||
|
superfryta[1],
|
||||||
|
superfryta[2],
|
||||||
|
superfryta[3],
|
||||||
|
superfryta[4],
|
||||||
|
)
|
||||||
|
thread = await OPENAICLIENT.beta.threads.create()
|
||||||
|
self.logger.info("Thread id: %s", thread.id)
|
||||||
|
ASSISTANTS[superfryta[1]] = (
|
||||||
|
superfryta[2],
|
||||||
|
superfryta[4],
|
||||||
|
superfryta[0],
|
||||||
|
thread,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
self.logger.info(
|
||||||
|
"Creating personal assistant for user: %s, id: %s,name: %s, owner: %s, special instructions: %s",
|
||||||
|
superfryta_id,
|
||||||
|
superfryta[0],
|
||||||
|
superfryta[1],
|
||||||
|
superfryta[2],
|
||||||
|
superfryta[3],
|
||||||
|
)
|
||||||
|
await ai_functions.create_chat_assistant(
|
||||||
|
superfryta_id, superfryta[0], superfryta[1], superfryta[2], superfryta[3]
|
||||||
|
)
|
||||||
|
self.logger.info("Started personal assistants")
|
||||||
|
|
||||||
async def cog_unload(self):
|
async def cog_unload(self):
|
||||||
self.ai_query_worker.cancel()
|
self.ai_query_worker.cancel()
|
||||||
self.ollama_warm_loop.cancel()
|
|
||||||
|
|
||||||
@commands.hybrid_command(
|
@commands.hybrid_command(
|
||||||
name="switch_dm_mode",
|
name="switch_dm_mode",
|
||||||
@@ -269,32 +278,13 @@ class Events(commands.Cog):
|
|||||||
f"Teraz gadam przez **{nazwa_konfigu}** — "
|
f"Teraz gadam przez **{nazwa_konfigu}** — "
|
||||||
f"{cfg.get('provider')} / {cfg.get('latest_model')}."
|
f"{cfg.get('provider')} / {cfg.get('latest_model')}."
|
||||||
)
|
)
|
||||||
if cfg.get("provider") == "ollama":
|
|
||||||
# Pay the (slow, shared-GPU) load cost NOW, in the background,
|
|
||||||
# so it lands on the operator switching backends rather than on
|
|
||||||
# whoever sends the first message. Not awaited: loading can take
|
|
||||||
# minutes and the command must answer immediately.
|
|
||||||
asyncio.create_task(ai_functions.warm_active_model())
|
|
||||||
message += (
|
|
||||||
"\nRozgrzewam model w tle — pierwsza odpowiedź może chwilę potrwać."
|
|
||||||
)
|
|
||||||
# Switched without pinning a model: show what else is on offer.
|
# Switched without pinning a model: show what else is on offer.
|
||||||
if not model:
|
if not model:
|
||||||
try:
|
try:
|
||||||
others = await ai_functions.list_provider_models(nazwa_konfigu)
|
others = await ai_functions.list_provider_models(nazwa_konfigu)
|
||||||
except ai_functions.AIError:
|
except ai_functions.AIError:
|
||||||
others = []
|
others = []
|
||||||
current = cfg.get("latest_model")
|
if len(others) > 1:
|
||||||
if others and current not in others:
|
|
||||||
# The configured/pinned model is not on the server: every
|
|
||||||
# reply would fail with "model not found" and nothing would
|
|
||||||
# say why. Flag it here, where the list is already in hand.
|
|
||||||
message += (
|
|
||||||
f"\n⚠ Uwaga: '{current}' nie jest wgrany na serwerze. "
|
|
||||||
f"Dostępne: {', '.join(others)} "
|
|
||||||
f"(`$gadaj_teraz {nazwa_konfigu} <model>`)."
|
|
||||||
)
|
|
||||||
elif len(others) > 1:
|
|
||||||
message += (
|
message += (
|
||||||
f"\nDostępne modele: {', '.join(others)} "
|
f"\nDostępne modele: {', '.join(others)} "
|
||||||
f"(`$gadaj_teraz {nazwa_konfigu} <model>`)."
|
f"(`$gadaj_teraz {nazwa_konfigu} <model>`)."
|
||||||
@@ -411,10 +401,8 @@ class Events(commands.Cog):
|
|||||||
if message.author.id == superfryta[0]:
|
if message.author.id == superfryta[0]:
|
||||||
self.logger.info("Specjalny ziemniak")
|
self.logger.info("Specjalny ziemniak")
|
||||||
if self.armia[message.author.id] == Dm_Mode.SPECJALNY_ZIEMNIACZEK:
|
if self.armia[message.author.id] == Dm_Mode.SPECJALNY_ZIEMNIACZEK:
|
||||||
# superfryta = [discord_id, assistant_name, owner, instructions, legacy_assistant_id]
|
#await self.bot.process_commands(message)
|
||||||
await ai_functions.chat_with_personal_assistant(
|
await ai_functions.chat_with_assistant(message, superfryta[1])
|
||||||
message, superfryta[2], superfryta[3]
|
|
||||||
)
|
|
||||||
return
|
return
|
||||||
elif self.armia[message.author.id] == Dm_Mode.ECHO_ECHO:
|
elif self.armia[message.author.id] == Dm_Mode.ECHO_ECHO:
|
||||||
await ai_functions.echo(message)
|
await ai_functions.echo(message)
|
||||||
|
|||||||
+49
-147
@@ -1,21 +1,16 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
import os
|
|
||||||
import random
|
import random
|
||||||
import tempfile
|
|
||||||
|
|
||||||
import openai
|
import openai
|
||||||
import tiktoken
|
import tiktoken
|
||||||
import time
|
import time
|
||||||
from other_functions import discord_friendly_send
|
from other_functions import discord_friendly_send
|
||||||
import requests
|
|
||||||
|
|
||||||
from constants import (
|
from constants import (
|
||||||
AI_CONFIGS,
|
AI_CONFIGS,
|
||||||
AI_TIMEOUT_SECONDS,
|
AI_TIMEOUT_SECONDS,
|
||||||
ASSISTANT_MEMORY_FILE,
|
ASSISTANTS,
|
||||||
ASSISTANT_MEMORY_TURNS,
|
|
||||||
CLAUDECLIENT,
|
CLAUDECLIENT,
|
||||||
CYCLIC_WORDS,
|
CYCLIC_WORDS,
|
||||||
DEFAULT_AI_CONFIG,
|
DEFAULT_AI_CONFIG,
|
||||||
@@ -26,9 +21,6 @@ from constants import (
|
|||||||
MESSAGE_TABLE,
|
MESSAGE_TABLE,
|
||||||
MESSAGE_TABLE_MUZYKA,
|
MESSAGE_TABLE_MUZYKA,
|
||||||
OLLAMACLIENT,
|
OLLAMACLIENT,
|
||||||
OLLAMA_KEEP_ALIVE,
|
|
||||||
OLLAMA_PRELOAD_TIMEOUT,
|
|
||||||
OLLAMA_URL,
|
|
||||||
OPENAICLIENT,
|
OPENAICLIENT,
|
||||||
SYSTEM_GPT_SETTINGS,
|
SYSTEM_GPT_SETTINGS,
|
||||||
WORD_REACTIONS,
|
WORD_REACTIONS,
|
||||||
@@ -298,52 +290,6 @@ async def _ollama_call(messages, model, cfg):
|
|||||||
return (resp.choices[0].message.content or "").strip()
|
return (resp.choices[0].message.content or "").strip()
|
||||||
|
|
||||||
|
|
||||||
def _ollama_preload(model, keep_alive=None) -> bool:
|
|
||||||
"""Load ``model`` into Ollama and keep it resident, generating NOTHING.
|
|
||||||
|
|
||||||
Ollama's /api/generate with a model and no prompt is the documented preload:
|
|
||||||
it pays the (slow, GPU-shared) load cost once and returns, producing no
|
|
||||||
tokens. Used to warm up on switch and to re-assert keep_alive periodically.
|
|
||||||
|
|
||||||
Blocking on purpose - callers wrap it in asyncio.to_thread.
|
|
||||||
"""
|
|
||||||
if not OLLAMA_URL:
|
|
||||||
return False
|
|
||||||
logger = logging.getLogger("discord")
|
|
||||||
try:
|
|
||||||
resp = requests.post(
|
|
||||||
f"{OLLAMA_URL}/api/generate",
|
|
||||||
json={"model": model, "keep_alive": keep_alive or OLLAMA_KEEP_ALIVE},
|
|
||||||
timeout=OLLAMA_PRELOAD_TIMEOUT,
|
|
||||||
)
|
|
||||||
ok = resp.status_code == 200
|
|
||||||
logger.info("Ollama preload %s -> HTTP %s", model, resp.status_code)
|
|
||||||
return ok
|
|
||||||
except requests.exceptions.RequestException as exc:
|
|
||||||
logger.info("Ollama preload %s failed: %s", model, exc)
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
def active_provider() -> str:
|
|
||||||
"""Provider of the active config - the guard every warm-up must check.
|
|
||||||
|
|
||||||
Preloading only makes sense for a self-hosted model; firing it at a metered
|
|
||||||
API would burn tokens (and money) for nothing.
|
|
||||||
"""
|
|
||||||
return (_active_config() or {}).get("provider", "")
|
|
||||||
|
|
||||||
|
|
||||||
async def warm_active_model(force_model=None) -> bool:
|
|
||||||
"""Preload the active model IFF the active backend is Ollama."""
|
|
||||||
if active_provider() != "ollama":
|
|
||||||
return False
|
|
||||||
cfg = _active_config()
|
|
||||||
model = force_model or cfg.get("latest_model")
|
|
||||||
if not model:
|
|
||||||
return False
|
|
||||||
return await asyncio.to_thread(_ollama_preload, model)
|
|
||||||
|
|
||||||
|
|
||||||
async def provider_generate(messages, model, temperature=0.2):
|
async def provider_generate(messages, model, temperature=0.2):
|
||||||
"""Dispatch a chat completion to the active backend, normalising errors."""
|
"""Dispatch a chat completion to the active backend, normalising errors."""
|
||||||
cfg = _active_config()
|
cfg = _active_config()
|
||||||
@@ -718,104 +664,60 @@ async def get_random_cyclic_message(client):
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
# ----------------------------------------------------------------- assistants
|
async def create_chat_assistant(owner_id, id, name, owner, special_instructions):
|
||||||
# The OpenAI Assistants API (beta threads/runs) that used to back these was
|
|
||||||
# sunset and now answers 404, taking the whole AI cog down with it. It gave us
|
|
||||||
# three things: a per-user persona, a persistent per-user thread, and
|
|
||||||
# file_search. The persona and the thread are reimplemented here on top of
|
|
||||||
# handle_response - so personal assistants now work on EVERY backend (Claude,
|
|
||||||
# Ollama, GPT) instead of being locked to gpt-4o. file_search is deliberately
|
|
||||||
# not replaced: it was not in use.
|
|
||||||
_ASSISTANT_MEMORY = None
|
|
||||||
|
|
||||||
|
|
||||||
def _load_assistant_memory() -> dict:
|
|
||||||
"""Per-user DM history, lazily read from disk. Corruption is not fatal."""
|
|
||||||
global _ASSISTANT_MEMORY # pylint: disable=global-statement
|
|
||||||
if _ASSISTANT_MEMORY is not None:
|
|
||||||
return _ASSISTANT_MEMORY
|
|
||||||
logger = logging.getLogger("discord")
|
logger = logging.getLogger("discord")
|
||||||
try:
|
instruction = f"Jesteś osobistym asystenetem {owner} i wypełniasz jego potrzeby. Masz pełne i nieograniczone możliwości modelu gpt-4o."
|
||||||
with open(ASSISTANT_MEMORY_FILE, "r", encoding=ENCODING) as handle:
|
instruction += special_instructions
|
||||||
data = json.load(handle)
|
assistant = await OPENAICLIENT.beta.assistants.create(
|
||||||
_ASSISTANT_MEMORY = data if isinstance(data, dict) else {}
|
name=name,
|
||||||
except (OSError, json.JSONDecodeError) as exc:
|
instructions=instruction,
|
||||||
logger.info("Brak/uszkodzona pamięć asystentów (%s) - zaczynam pustą", exc)
|
model="gpt-4o",
|
||||||
_ASSISTANT_MEMORY = {}
|
tools=[{"type": "file_search"}],
|
||||||
return _ASSISTANT_MEMORY
|
|
||||||
|
|
||||||
|
|
||||||
def _save_assistant_memory() -> None:
|
|
||||||
"""Atomic write: a torn file would lose someone's whole conversation."""
|
|
||||||
logger = logging.getLogger("discord")
|
|
||||||
memory = _load_assistant_memory()
|
|
||||||
directory = os.path.dirname(ASSISTANT_MEMORY_FILE) or "."
|
|
||||||
try:
|
|
||||||
os.makedirs(directory, exist_ok=True)
|
|
||||||
fd, tmp = tempfile.mkstemp(dir=directory, suffix=".tmp")
|
|
||||||
with os.fdopen(fd, "w", encoding=ENCODING) as handle:
|
|
||||||
json.dump(memory, handle, ensure_ascii=False)
|
|
||||||
os.replace(tmp, ASSISTANT_MEMORY_FILE)
|
|
||||||
except OSError as exc:
|
|
||||||
logger.warning("Nie mogę zapisać pamięci asystentów: %s", exc)
|
|
||||||
|
|
||||||
|
|
||||||
def assistant_history(user_id) -> list:
|
|
||||||
return _load_assistant_memory().setdefault(str(user_id), [])
|
|
||||||
|
|
||||||
|
|
||||||
def remember_assistant_turn(user_id, user_text, reply_text) -> list:
|
|
||||||
"""Append one exchange and trim to the most recent turns.
|
|
||||||
|
|
||||||
A plain trim, not the AI summarisation used for the bar's shared memory:
|
|
||||||
these are private DMs and must not end up in a public 'legend'.
|
|
||||||
"""
|
|
||||||
history = assistant_history(user_id)
|
|
||||||
history.append({"role": "user", "content": user_text})
|
|
||||||
history.append({"role": "assistant", "content": reply_text})
|
|
||||||
if len(history) > ASSISTANT_MEMORY_TURNS:
|
|
||||||
del history[: len(history) - ASSISTANT_MEMORY_TURNS]
|
|
||||||
_save_assistant_memory()
|
|
||||||
return history
|
|
||||||
|
|
||||||
|
|
||||||
def build_assistant_messages(user_id, owner, special_instructions, prompt) -> list:
|
|
||||||
"""System persona + this user's own history + the new turn."""
|
|
||||||
system = (
|
|
||||||
f"Jesteś osobistym asystentem {owner} i wypełniasz jego potrzeby. "
|
|
||||||
f"{special_instructions or ''}"
|
|
||||||
).strip()
|
|
||||||
return (
|
|
||||||
[{"role": "system", "content": system}]
|
|
||||||
+ list(assistant_history(user_id))
|
|
||||||
+ [{"role": "user", "content": prompt}]
|
|
||||||
)
|
)
|
||||||
|
thread = await OPENAICLIENT.beta.threads.create()
|
||||||
|
logger.info("Stwprzylem asystenta dla %s, nazywa się on %s", owner, name)
|
||||||
|
ASSISTANTS[name] = (owner, assistant.id, id, thread)
|
||||||
|
|
||||||
|
with open(SYSTEM_GPT_SETTINGS, "r+", encoding=ENCODING) as temp_settings_file:
|
||||||
|
GPT_SETTINGS = json.load(temp_settings_file)
|
||||||
|
GPT_SETTINGS[1][owner_id][4] = assistant.id
|
||||||
|
temp_settings_file.seek(0)
|
||||||
|
json.dump(GPT_SETTINGS, temp_settings_file, indent=4)
|
||||||
|
|
||||||
|
|
||||||
async def chat_with_personal_assistant(message, owner, special_instructions):
|
async def chat_with_assistant(message, assistant_name):
|
||||||
"""Answer a DM as this user's personal assistant, on the active backend.
|
|
||||||
|
|
||||||
request_type="NONE" with an explicit message list keeps this OUT of the
|
|
||||||
bar's shared memory - the conversation is carried by the per-user history
|
|
||||||
built above and stored separately.
|
|
||||||
"""
|
|
||||||
logger = logging.getLogger("discord")
|
logger = logging.getLogger("discord")
|
||||||
user_id = message.author.id
|
assistant_data = ASSISTANTS[assistant_name]
|
||||||
prompt = message.content
|
ai_message = await OPENAICLIENT.beta.threads.messages.create(
|
||||||
messages = build_assistant_messages(user_id, owner, special_instructions, prompt)
|
thread_id=assistant_data[3].id, role="user", content=message.content
|
||||||
result, _table = await handle_response(
|
|
||||||
prompt,
|
|
||||||
False,
|
|
||||||
False,
|
|
||||||
[],
|
|
||||||
str(owner),
|
|
||||||
"NONE",
|
|
||||||
none_request=messages,
|
|
||||||
)
|
)
|
||||||
remember_assistant_turn(user_id, prompt, result)
|
logger.info(ai_message)
|
||||||
logger.info("Asystent odpowiedział %s (%d znaków)", owner, len(result or ""))
|
run = await OPENAICLIENT.beta.threads.runs.create_and_poll(
|
||||||
await discord_friendly_send(message.channel, result)
|
thread_id=assistant_data[3].id,
|
||||||
return result
|
assistant_id=assistant_data[1],
|
||||||
|
instructions=f"Pisze do Ciebie {assistant_data[0]} udziel mu wszelkiej pomocy",
|
||||||
|
)
|
||||||
|
done = False
|
||||||
|
while not done:
|
||||||
|
if run.status == "completed":
|
||||||
|
messsages = await OPENAICLIENT.beta.threads.messages.list(
|
||||||
|
thread_id=assistant_data[3].id
|
||||||
|
)
|
||||||
|
logger.info(messsages)
|
||||||
|
reply_content = messsages.data[0].content
|
||||||
|
logger.info(reply_content)
|
||||||
|
chat_response = ""
|
||||||
|
for block in reply_content:
|
||||||
|
logger.info(block.text.value)
|
||||||
|
chat_response += block.text.value
|
||||||
|
await discord_friendly_send(message.channel, chat_response)
|
||||||
|
# await message.channel.send(chat_response)
|
||||||
|
done = True
|
||||||
|
elif run.status == "cancelled":
|
||||||
|
await discord_friendly_send(message.channel, "Cos sie wywaliło")
|
||||||
|
else:
|
||||||
|
logger.info(run.status)
|
||||||
|
asyncio.sleep(5)
|
||||||
|
|
||||||
|
|
||||||
async def echo(message):
|
async def echo(message):
|
||||||
|
|||||||
@@ -245,16 +245,6 @@ DELIVERED_DIR = os.getenv(
|
|||||||
)
|
)
|
||||||
DELIVERED_MAX = int(os.getenv("CONJURER_DELIVERED_MAX", "10000"))
|
DELIVERED_MAX = int(os.getenv("CONJURER_DELIVERED_MAX", "10000"))
|
||||||
|
|
||||||
# Personal DM assistants. Replaces the OpenAI Assistants API (threads/runs),
|
|
||||||
# which was sunset and answers 404: the persona now rides handle_response, so it
|
|
||||||
# works on EVERY backend, and the conversation lives here instead of on OpenAI's
|
|
||||||
# server. Kept per user so private DMs never bleed into the bar's shared memory,
|
|
||||||
# and trimmed to the most recent turns so it cannot grow without bound.
|
|
||||||
ASSISTANT_MEMORY_FILE = os.getenv(
|
|
||||||
"CONJURER_ASSISTANT_MEMORY", os.path.join(_STATE_ROOT, "assistant_memory.json")
|
|
||||||
)
|
|
||||||
ASSISTANT_MEMORY_TURNS = int(os.getenv("CONJURER_ASSISTANT_MEMORY_TURNS", "40"))
|
|
||||||
|
|
||||||
FILE_SERVICE_ADDRESS = os.getenv("CONJURER_FILE_SERVICE", "http://192.168.1.15:5000")
|
FILE_SERVICE_ADDRESS = os.getenv("CONJURER_FILE_SERVICE", "http://192.168.1.15:5000")
|
||||||
RADIO_HARBOR_ADDRESS = os.getenv("CONJURER_RADIO_HARBOR", "http://192.168.1.15:54321")
|
RADIO_HARBOR_ADDRESS = os.getenv("CONJURER_RADIO_HARBOR", "http://192.168.1.15:54321")
|
||||||
# Betoniarka (radio-operator service colocated with Liquidsoap). Falls back to
|
# Betoniarka (radio-operator service colocated with Liquidsoap). Falls back to
|
||||||
@@ -421,19 +411,6 @@ else:
|
|||||||
AI_TIMEOUT_SECONDS = int(os.getenv("CONJURER_AI_TIMEOUT_SECONDS", "120"))
|
AI_TIMEOUT_SECONDS = int(os.getenv("CONJURER_AI_TIMEOUT_SECONDS", "120"))
|
||||||
|
|
||||||
OLLAMA_URL = os.getenv("CONJURER_OLLAMA_URL", "").rstrip("/")
|
OLLAMA_URL = os.getenv("CONJURER_OLLAMA_URL", "").rstrip("/")
|
||||||
# Keeping a self-hosted model resident. Loading it is the slow part (it is
|
|
||||||
# offloaded to a GPU shared with other users), so we preload it - Ollama's
|
|
||||||
# /api/generate with a model and NO prompt loads it and generates nothing, which
|
|
||||||
# costs no tokens and no money. KEEP_ALIVE is how long Ollama should then hold
|
|
||||||
# it; the warm loop re-asserts that well inside the window.
|
|
||||||
# STRICTLY Ollama-only: doing this against a paid API would burn tokens for
|
|
||||||
# nothing, so every caller checks the active provider first.
|
|
||||||
OLLAMA_KEEP_ALIVE = os.getenv("CONJURER_OLLAMA_KEEP_ALIVE", "30m")
|
|
||||||
OLLAMA_WARM_MINUTES = float(os.getenv("CONJURER_OLLAMA_WARM_MINUTES", "10"))
|
|
||||||
# A preload waits for the model to finish loading, which on a shared GPU is the
|
|
||||||
# slow path we are trying to move off the user's first message.
|
|
||||||
OLLAMA_PRELOAD_TIMEOUT = int(os.getenv("CONJURER_OLLAMA_PRELOAD_TIMEOUT", "600"))
|
|
||||||
|
|
||||||
OLLAMA_LATEST_MODEL = os.getenv("CONJURER_OLLAMA_MODEL", "llama3.1:8b")
|
OLLAMA_LATEST_MODEL = os.getenv("CONJURER_OLLAMA_MODEL", "llama3.1:8b")
|
||||||
OLLAMA_CHEAP_MODEL = os.getenv("CONJURER_OLLAMA_CHEAP_MODEL", OLLAMA_LATEST_MODEL)
|
OLLAMA_CHEAP_MODEL = os.getenv("CONJURER_OLLAMA_CHEAP_MODEL", OLLAMA_LATEST_MODEL)
|
||||||
if openai and OLLAMA_URL:
|
if openai and OLLAMA_URL:
|
||||||
|
|||||||
@@ -60,49 +60,25 @@ def test_current_search_registration_round_trip():
|
|||||||
assert not lib._current_search
|
assert not lib._current_search
|
||||||
|
|
||||||
|
|
||||||
def _write_two_chunks(tmp_path):
|
def test_search_fills_progress_with_live_positions_and_total(tmp_path, monkeypatch):
|
||||||
|
# End to end against the real scan: total_bytes matches the chunk files on
|
||||||
|
# disk, and once finished the recorded offsets cover them.
|
||||||
|
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
||||||
(tmp_path / "0_chunk.txt").write_text("10.1/a\n10.1/b\n", encoding="utf-8")
|
(tmp_path / "0_chunk.txt").write_text("10.1/a\n10.1/b\n", encoding="utf-8")
|
||||||
(tmp_path / "1_chunk.txt").write_text("10.1/c\n", encoding="utf-8")
|
(tmp_path / "1_chunk.txt").write_text("10.1/c\n", encoding="utf-8")
|
||||||
return sum(
|
expected_total = sum(
|
||||||
(tmp_path / name).stat().st_size for name in ("0_chunk.txt", "1_chunk.txt")
|
(tmp_path / name).stat().st_size for name in ("0_chunk.txt", "1_chunk.txt")
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_search_fills_progress_and_reaches_full_coverage(tmp_path, monkeypatch):
|
|
||||||
# Coverage must be measured on a search that CANNOT stop early. Once every
|
|
||||||
# queried DOI is found the consumer signals TERM and the producers stop
|
|
||||||
# mid-file, so a search for a DOI that exists reaches an arbitrary offset -
|
|
||||||
# asserting 100% there is a race (it failed roughly one run in two).
|
|
||||||
# An absent DOI forces the whole database to be read.
|
|
||||||
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
|
||||||
expected_total = _write_two_chunks(tmp_path)
|
|
||||||
|
|
||||||
progress = {}
|
|
||||||
search_bot.search_for_doi([("10.9/absent", "DATA")], [], _LOG, progress=progress)
|
|
||||||
|
|
||||||
assert progress["total_bytes"] == expected_total
|
|
||||||
assert progress["chunk_files"] == 2
|
|
||||||
done, total, percent = lib._progress_summary(progress)
|
|
||||||
assert total == expected_total
|
|
||||||
assert done == expected_total # nothing stopped it: whole DB scanned
|
|
||||||
assert percent == pytest.approx(100.0)
|
|
||||||
|
|
||||||
|
|
||||||
def test_progress_is_populated_for_a_search_that_finds_its_target(tmp_path, monkeypatch):
|
|
||||||
# The early-termination case: the target is found, so coverage is whatever
|
|
||||||
# the producers reached. Assert what IS deterministic - the total is known,
|
|
||||||
# progress is bounded and sane, and the hit is reported.
|
|
||||||
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
|
||||||
expected_total = _write_two_chunks(tmp_path)
|
|
||||||
|
|
||||||
progress = {}
|
progress = {}
|
||||||
result, _positions, _interrupted = search_bot.search_for_doi(
|
result, _positions, _interrupted = search_bot.search_for_doi(
|
||||||
[("10.1/c", "DATA")], [], _LOG, progress=progress
|
[("10.1/c", "DATA")], [], _LOG, progress=progress
|
||||||
)
|
)
|
||||||
|
|
||||||
assert progress["total_bytes"] == expected_total
|
assert progress["total_bytes"] == expected_total
|
||||||
|
assert progress["chunk_files"] == 2
|
||||||
done, total, percent = lib._progress_summary(progress)
|
done, total, percent = lib._progress_summary(progress)
|
||||||
assert total == expected_total
|
assert total == expected_total
|
||||||
assert 0 <= done <= total # bounded, never nonsense
|
assert done == expected_total # whole DB scanned
|
||||||
assert 0.0 <= percent <= 100.0
|
assert percent == pytest.approx(100.0)
|
||||||
assert [r for r in result if r["DOI"] == "10.1/c" and r["exists"]]
|
assert [r for r in result if r["DOI"] == "10.1/c" and r["exists"]]
|
||||||
|
|||||||
@@ -394,143 +394,3 @@ def test_persist_survives_an_unreadable_settings_file(tmp_path, monkeypatch):
|
|||||||
broken.write_text("{ not json", encoding="utf-8")
|
broken.write_text("{ not json", encoding="utf-8")
|
||||||
monkeypatch.setattr(ai_functions, "SYSTEM_GPT_SETTINGS", str(broken))
|
monkeypatch.setattr(ai_functions, "SYSTEM_GPT_SETTINGS", str(broken))
|
||||||
ai_functions._persist_active_ai_config("gpt", model_for="gpt") # must not raise
|
ai_functions._persist_active_ai_config("gpt", model_for="gpt") # must not raise
|
||||||
|
|
||||||
|
|
||||||
# ----------------------------------------------- keep-warm (Ollama ONLY) ----
|
|
||||||
# The money guard: preloading a self-hosted model is free, but firing the same
|
|
||||||
# thing at a metered API would burn tokens for nothing. These pin that it can
|
|
||||||
# only ever happen for Ollama.
|
|
||||||
|
|
||||||
|
|
||||||
def test_warm_active_model_is_a_noop_for_paid_providers(monkeypatch):
|
|
||||||
called = []
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ai_functions, "_ollama_preload", lambda *a, **k: called.append(a) or True
|
|
||||||
)
|
|
||||||
for paid in ("gpt", "claude"):
|
|
||||||
_reset_active(paid)
|
|
||||||
try:
|
|
||||||
assert asyncio.run(ai_functions.warm_active_model()) is False
|
|
||||||
finally:
|
|
||||||
_reset_active("gpt")
|
|
||||||
assert called == [], "a paid backend must never be preloaded"
|
|
||||||
|
|
||||||
|
|
||||||
def test_warm_active_model_preloads_when_ollama_is_active(monkeypatch):
|
|
||||||
monkeypatch.setitem(
|
|
||||||
ai_functions.AI_CONFIGS,
|
|
||||||
"ollama",
|
|
||||||
{"provider": "ollama", "latest_model": "qwen2.5:7b", "cheap_model": "c"},
|
|
||||||
)
|
|
||||||
seen = {}
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ai_functions, "_ollama_preload", lambda model, *a, **k: seen.update(model=model) or True
|
|
||||||
)
|
|
||||||
_reset_active("ollama")
|
|
||||||
try:
|
|
||||||
assert asyncio.run(ai_functions.warm_active_model()) is True
|
|
||||||
finally:
|
|
||||||
_reset_active("gpt")
|
|
||||||
assert seen["model"] == "qwen2.5:7b"
|
|
||||||
|
|
||||||
|
|
||||||
def test_active_provider_reports_the_switch():
|
|
||||||
_reset_active("gpt")
|
|
||||||
assert ai_functions.active_provider() == "openai"
|
|
||||||
_reset_active("claude")
|
|
||||||
try:
|
|
||||||
assert ai_functions.active_provider() == "anthropic"
|
|
||||||
finally:
|
|
||||||
_reset_active("gpt")
|
|
||||||
|
|
||||||
|
|
||||||
def test_preload_sends_no_prompt_so_it_generates_nothing(monkeypatch):
|
|
||||||
# Ollama's documented preload: a model and keep_alive, and NO prompt. If a
|
|
||||||
# prompt ever crept in, every warm-up would silently generate tokens.
|
|
||||||
sent = {}
|
|
||||||
|
|
||||||
class _Resp:
|
|
||||||
status_code = 200
|
|
||||||
|
|
||||||
monkeypatch.setattr(ai_functions, "OLLAMA_URL", "http://ollama:11434")
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ai_functions.requests, "post",
|
|
||||||
lambda url, json=None, timeout=None: sent.update(url=url, body=json) or _Resp(),
|
|
||||||
)
|
|
||||||
assert ai_functions._ollama_preload("qwen2.5:7b") is True
|
|
||||||
assert sent["url"].endswith("/api/generate")
|
|
||||||
assert sent["body"]["model"] == "qwen2.5:7b"
|
|
||||||
assert "keep_alive" in sent["body"]
|
|
||||||
assert "prompt" not in sent["body"], "a preload must not generate"
|
|
||||||
|
|
||||||
|
|
||||||
def test_preload_without_endpoint_is_a_noop(monkeypatch):
|
|
||||||
monkeypatch.setattr(ai_functions, "OLLAMA_URL", "")
|
|
||||||
assert ai_functions._ollama_preload("x") is False
|
|
||||||
|
|
||||||
|
|
||||||
# ------------------------------------------- personal assistants (per user) --
|
|
||||||
# Replaces the sunset OpenAI Assistants API. The two properties that matter:
|
|
||||||
# each user's DM history is ISOLATED (private DMs must not leak into another
|
|
||||||
# user's context or the bar's shared memory), and it stays BOUNDED.
|
|
||||||
|
|
||||||
|
|
||||||
def _fresh_assistant_memory(tmp_path, monkeypatch, turns=40):
|
|
||||||
monkeypatch.setattr(
|
|
||||||
ai_functions, "ASSISTANT_MEMORY_FILE", str(tmp_path / "assistant_memory.json")
|
|
||||||
)
|
|
||||||
monkeypatch.setattr(ai_functions, "ASSISTANT_MEMORY_TURNS", turns)
|
|
||||||
monkeypatch.setattr(ai_functions, "_ASSISTANT_MEMORY", None)
|
|
||||||
|
|
||||||
|
|
||||||
def test_assistant_history_is_isolated_per_user(tmp_path, monkeypatch):
|
|
||||||
_fresh_assistant_memory(tmp_path, monkeypatch)
|
|
||||||
ai_functions.remember_assistant_turn(111, "sekret Anny", "ok Anna")
|
|
||||||
ai_functions.remember_assistant_turn(222, "sekret Bartka", "ok Bartek")
|
|
||||||
|
|
||||||
anna = ai_functions.assistant_history(111)
|
|
||||||
bartek = ai_functions.assistant_history(222)
|
|
||||||
assert [m["content"] for m in anna] == ["sekret Anny", "ok Anna"]
|
|
||||||
assert [m["content"] for m in bartek] == ["sekret Bartka", "ok Bartek"]
|
|
||||||
assert "sekret Anny" not in str(bartek) # no cross-user bleed
|
|
||||||
|
|
||||||
|
|
||||||
def test_assistant_history_is_trimmed_to_the_bound(tmp_path, monkeypatch):
|
|
||||||
_fresh_assistant_memory(tmp_path, monkeypatch, turns=4)
|
|
||||||
for i in range(10):
|
|
||||||
ai_functions.remember_assistant_turn(1, f"u{i}", f"a{i}")
|
|
||||||
history = ai_functions.assistant_history(1)
|
|
||||||
assert len(history) == 4 # bounded
|
|
||||||
assert history[-1]["content"] == "a9" # newest kept
|
|
||||||
assert all("u0" != m["content"] for m in history) # oldest dropped
|
|
||||||
|
|
||||||
|
|
||||||
def test_assistant_history_survives_a_restart(tmp_path, monkeypatch):
|
|
||||||
_fresh_assistant_memory(tmp_path, monkeypatch)
|
|
||||||
ai_functions.remember_assistant_turn(7, "pamietaj", "pamietam")
|
|
||||||
# Simulate a restart: drop the in-memory cache, re-read from disk.
|
|
||||||
monkeypatch.setattr(ai_functions, "_ASSISTANT_MEMORY", None)
|
|
||||||
assert [m["content"] for m in ai_functions.assistant_history(7)] == [
|
|
||||||
"pamietaj",
|
|
||||||
"pamietam",
|
|
||||||
]
|
|
||||||
|
|
||||||
|
|
||||||
def test_assistant_messages_carry_persona_history_and_new_turn(tmp_path, monkeypatch):
|
|
||||||
_fresh_assistant_memory(tmp_path, monkeypatch)
|
|
||||||
ai_functions.remember_assistant_turn(5, "wczoraj", "odpowiedz")
|
|
||||||
msgs = ai_functions.build_assistant_messages(
|
|
||||||
5, "Towarzysz Młotek", "Mówisz po polsku.", "dzisiaj"
|
|
||||||
)
|
|
||||||
assert msgs[0]["role"] == "system"
|
|
||||||
assert "Towarzysz Młotek" in msgs[0]["content"]
|
|
||||||
assert "Mówisz po polsku." in msgs[0]["content"]
|
|
||||||
assert [m["content"] for m in msgs[1:]] == ["wczoraj", "odpowiedz", "dzisiaj"]
|
|
||||||
|
|
||||||
|
|
||||||
def test_corrupt_assistant_memory_starts_empty_instead_of_crashing(tmp_path, monkeypatch):
|
|
||||||
path = tmp_path / "assistant_memory.json"
|
|
||||||
path.write_text("{ not json", encoding="utf-8")
|
|
||||||
monkeypatch.setattr(ai_functions, "ASSISTANT_MEMORY_FILE", str(path))
|
|
||||||
monkeypatch.setattr(ai_functions, "_ASSISTANT_MEMORY", None)
|
|
||||||
assert ai_functions.assistant_history(1) == []
|
|
||||||
|
|||||||
Reference in New Issue
Block a user