Compare commits
3 Commits
efb671befa
...
491f957315
| Author | SHA1 | Date | |
|---|---|---|---|
| 491f957315 | |||
| b13a8afa01 | |||
| 1a59c9f6c5 |
+65
-1
@@ -9,8 +9,10 @@ from pathlib import Path
|
|||||||
import discord
|
import discord
|
||||||
import openai
|
import openai
|
||||||
import requests
|
import requests
|
||||||
from discord.ext import commands
|
from queue import Empty
|
||||||
|
from discord.ext import commands, tasks
|
||||||
from other_functions import discord_friendly_send, discord_friendly_reply
|
from other_functions import discord_friendly_send, discord_friendly_reply
|
||||||
|
from communication_subroutine import AI_QUERY_Q
|
||||||
|
|
||||||
|
|
||||||
import ai_functions
|
import ai_functions
|
||||||
@@ -43,7 +45,66 @@ class Events(commands.Cog):
|
|||||||
self.armia[superfryta[0]] = Dm_Mode.SPECJALNY_ZIEMNIACZEK
|
self.armia[superfryta[0]] = Dm_Mode.SPECJALNY_ZIEMNIACZEK
|
||||||
self.logger.info(self.armia)
|
self.logger.info(self.armia)
|
||||||
|
|
||||||
|
@tasks.loop(seconds=2)
|
||||||
|
async def ai_query_worker(self):
|
||||||
|
"""Drain AI_QUERY_Q one prompt at a time and answer with the active backend.
|
||||||
|
|
||||||
|
This is the bot-side half of the AI query interface: prompts arrive over
|
||||||
|
HTTP (POST /ai_query) or in-process (submit_ai_query), get queued, and are
|
||||||
|
answered here with handle_response - so they automatically use whichever
|
||||||
|
provider $gadaj_teraz currently selects. The answer is posted to the
|
||||||
|
channel the request named.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
item = AI_QUERY_Q.get(block=False)
|
||||||
|
except Empty:
|
||||||
|
return
|
||||||
|
prompt = item.get("prompt", "")
|
||||||
|
request_type = item.get("request_type", "NONE")
|
||||||
|
channel_id = item.get("channel_id")
|
||||||
|
username = item.get("username", "conjurer")
|
||||||
|
self.logger.info(
|
||||||
|
"AI query from %s (%s) -> channel %s", username, request_type, channel_id
|
||||||
|
)
|
||||||
|
global MESSAGE_TABLE # pylint: disable=global-statement
|
||||||
|
try:
|
||||||
|
if request_type == "NONE":
|
||||||
|
# Clean one-shot: no persona system prompt, no memory write.
|
||||||
|
result, _ = await ai_functions.handle_response(
|
||||||
|
"", True, True, [], username, "NONE", none_request=prompt
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
result, MESSAGE_TABLE = await ai_functions.handle_response(
|
||||||
|
prompt, True, True, MESSAGE_TABLE, username, request_type
|
||||||
|
)
|
||||||
|
except Exception as exc: # pylint: disable=broad-except
|
||||||
|
self.logger.exception("AI query failed: %s", exc)
|
||||||
|
result = "*Kondziu drapie się po głowie* Coś się zjebało przy pytaniu do AI."
|
||||||
|
if channel_id is None:
|
||||||
|
self.logger.warning("AI query had no channel_id - answer dropped")
|
||||||
|
return
|
||||||
|
channel = self.bot.get_channel(channel_id)
|
||||||
|
if channel is None:
|
||||||
|
self.logger.warning("AI query channel %s not found - answer dropped", channel_id)
|
||||||
|
return
|
||||||
|
await self._send_chunked(channel, result)
|
||||||
|
|
||||||
|
@ai_query_worker.before_loop
|
||||||
|
async def _before_ai_query_worker(self):
|
||||||
|
await self.bot.wait_until_ready()
|
||||||
|
|
||||||
|
async def _send_chunked(self, channel, text):
|
||||||
|
"""Send text in <=1900-char pieces (Discord caps messages at 2000)."""
|
||||||
|
text = text or ""
|
||||||
|
while text:
|
||||||
|
await discord_friendly_send(channel, text[:1900])
|
||||||
|
text = text[1900:]
|
||||||
|
|
||||||
async def cog_load(self):
|
async def cog_load(self):
|
||||||
|
# The AI query worker must run regardless of the OpenAI guard below - it
|
||||||
|
# answers via handle_response, which works on Claude too. Start it first.
|
||||||
|
if not self.ai_query_worker.is_running():
|
||||||
|
self.ai_query_worker.start()
|
||||||
self.logger.info("Starting personal assistants")
|
self.logger.info("Starting personal assistants")
|
||||||
# Personal assistants use the OpenAI Assistants API (threads/runs), which
|
# Personal assistants use the OpenAI Assistants API (threads/runs), which
|
||||||
# has no Anthropic equivalent - skip cleanly when OpenAI isn't wired up
|
# has no Anthropic equivalent - skip cleanly when OpenAI isn't wired up
|
||||||
@@ -87,6 +148,9 @@ class Events(commands.Cog):
|
|||||||
)
|
)
|
||||||
self.logger.info("Started personal assistants")
|
self.logger.info("Started personal assistants")
|
||||||
|
|
||||||
|
async def cog_unload(self):
|
||||||
|
self.ai_query_worker.cancel()
|
||||||
|
|
||||||
@commands.hybrid_command(
|
@commands.hybrid_command(
|
||||||
name="switch_dm_mode",
|
name="switch_dm_mode",
|
||||||
description="Jeśli nie wiesz jak użyć tej komendy to nawet nie próbuj",
|
description="Jeśli nie wiesz jak użyć tej komendy to nawet nie próbuj",
|
||||||
|
|||||||
@@ -17,6 +17,11 @@ ICECAST_ADDRESS = os.getenv("CONJURER_ICECAST", "http://192.168.1.12:8000")
|
|||||||
API_KEY = os.getenv("CONJURER_API_KEY")
|
API_KEY = os.getenv("CONJURER_API_KEY")
|
||||||
OUT_COMM_Q = Queue()
|
OUT_COMM_Q = Queue()
|
||||||
IN_COMM_Q = Queue()
|
IN_COMM_Q = Queue()
|
||||||
|
# AI request queue: prompts to be answered by the bot's own AI backend (whatever
|
||||||
|
# $gadaj_teraz currently points at). Drained by the AI cog's worker loop, which
|
||||||
|
# calls handle_response and delivers the answer to the requested channel. Fed
|
||||||
|
# either over HTTP (POST /ai_query) or in-process via submit_ai_query().
|
||||||
|
AI_QUERY_Q = Queue()
|
||||||
SRCHTITLE = re.compile(rb"StreamTitle=\\*(?P<title>[^;]*);").search
|
SRCHTITLE = re.compile(rb"StreamTitle=\\*(?P<title>[^;]*);").search
|
||||||
|
|
||||||
awaiting_q = []
|
awaiting_q = []
|
||||||
@@ -47,13 +52,17 @@ class QueryControl:
|
|||||||
content, logger, context, and replies.
|
content, logger, context, and replies.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self, query_author, query_uuid, query_content, ctx) -> None:
|
def __init__(self, query_author, query_uuid, query_content, ctx, ai_review=False) -> None:
|
||||||
self.author = query_author
|
self.author = query_author
|
||||||
self.uuid = query_uuid
|
self.uuid = query_uuid
|
||||||
self.content = query_content
|
self.content = query_content
|
||||||
self.logger = logging.getLogger("discord")
|
self.logger = logging.getLogger("discord")
|
||||||
self.stop = False
|
self.stop = False
|
||||||
self.ctx = ctx
|
self.ctx = ctx
|
||||||
|
# When True, once the librarian returns hits, the DOI list + the search
|
||||||
|
# phrase are sent to the AI backend for a weighted-relevance re-rank and
|
||||||
|
# source review (see librarian_commands.check_data_q).
|
||||||
|
self.ai_review = ai_review
|
||||||
self.logger.info(
|
self.logger.info(
|
||||||
f"Created Query control for {self.author}, {self.uuid}: {self.content}"
|
f"Created Query control for {self.author}, {self.uuid}: {self.content}"
|
||||||
)
|
)
|
||||||
@@ -105,6 +114,53 @@ def answer_external_command():
|
|||||||
return jsonify("SUCCESS")
|
return jsonify("SUCCESS")
|
||||||
|
|
||||||
|
|
||||||
|
def submit_ai_query(prompt, channel_id=None, request_type="NONE", username="conjurer", query_uuid=None):
|
||||||
|
"""Queue an AI prompt for the bot to answer with its configured backend.
|
||||||
|
|
||||||
|
In-process entry point (used by the librarian result handler). ``channel_id``
|
||||||
|
is the Discord channel the answer should be posted to; ``request_type`` is
|
||||||
|
passed through to handle_response ("NONE" keeps it a clean one-shot that does
|
||||||
|
not touch conversation memory).
|
||||||
|
"""
|
||||||
|
AI_QUERY_Q.put(
|
||||||
|
{
|
||||||
|
"uuid": query_uuid,
|
||||||
|
"prompt": prompt,
|
||||||
|
"channel_id": channel_id,
|
||||||
|
"request_type": request_type,
|
||||||
|
"username": username,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@app.route("/ai_query", methods=["POST"])
|
||||||
|
def ai_query():
|
||||||
|
"""Inbound AI query: queue a prompt to be answered by the bot's AI backend.
|
||||||
|
|
||||||
|
Payload: {"prompt": str, "channel_id": int, "request_type"?: str,
|
||||||
|
"username"?: str, "uuid"?: str}. The prompt is queued and answered
|
||||||
|
asynchronously by the AI cog's worker; the answer is posted to channel_id.
|
||||||
|
"""
|
||||||
|
_authorize_request()
|
||||||
|
logger = logging.getLogger("discord")
|
||||||
|
record = json.loads(request.data)
|
||||||
|
prompt = record.get("prompt", "")
|
||||||
|
if not prompt:
|
||||||
|
return jsonify(isError=True, message="missing 'prompt'", statusCode=400, data=[]), 400
|
||||||
|
submit_ai_query(
|
||||||
|
prompt=prompt,
|
||||||
|
channel_id=record.get("channel_id"),
|
||||||
|
request_type=record.get("request_type", "NONE"),
|
||||||
|
username=record.get("username", "external"),
|
||||||
|
query_uuid=record.get("uuid"),
|
||||||
|
)
|
||||||
|
logger.info("Queued AI query (qsize=%s)", AI_QUERY_Q.qsize())
|
||||||
|
return (
|
||||||
|
jsonify(isError=False, message="Queued", statusCode=200, data={"qsize": AI_QUERY_Q.qsize()}),
|
||||||
|
200,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@app.route("/conjurer", methods=["GET"])
|
@app.route("/conjurer", methods=["GET"])
|
||||||
def check_alive():
|
def check_alive():
|
||||||
"""
|
"""
|
||||||
|
|||||||
+6
-1
@@ -9,6 +9,11 @@ import sys
|
|||||||
|
|
||||||
_ROOT = os.path.dirname(os.path.abspath(__file__))
|
_ROOT = os.path.dirname(os.path.abspath(__file__))
|
||||||
|
|
||||||
for _path in (_ROOT, os.path.join(_ROOT, "conjurer_musician")):
|
for _path in (
|
||||||
|
_ROOT,
|
||||||
|
os.path.join(_ROOT, "conjurer_musician"),
|
||||||
|
os.path.join(_ROOT, "conjurer_librarian"),
|
||||||
|
os.path.join(_ROOT, "conjurer_betoniarka"),
|
||||||
|
):
|
||||||
if _path not in sys.path:
|
if _path not in sys.path:
|
||||||
sys.path.insert(0, _path)
|
sys.path.insert(0, _path)
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ Global Variables:
|
|||||||
|
|
||||||
# TODO: Wpiemdolić to wszystko w klasę z loggerem przysłanym z góry
|
# TODO: Wpiemdolić to wszystko w klasę z loggerem przysłanym z góry
|
||||||
import os
|
import os
|
||||||
|
import re
|
||||||
from queue import Empty, Queue
|
from queue import Empty, Queue
|
||||||
from threading import Thread
|
from threading import Thread
|
||||||
import time
|
import time
|
||||||
@@ -28,13 +29,54 @@ q = Queue()
|
|||||||
|
|
||||||
# Deployment data is environment-overridable so the local DOI database can live
|
# Deployment data is environment-overridable so the local DOI database can live
|
||||||
# on a mounted volume (Docker/Linux) instead of the hardcoded Windows path.
|
# on a mounted volume (Docker/Linux) instead of the hardcoded Windows path.
|
||||||
MAXTHREADS = int(os.getenv("CONJURER_LIBRARIAN_MAXTHREADS", "41"))
|
|
||||||
DATABASE_PATH = os.getenv("CONJURER_LIBRARIAN_DB_PATH", r"C:\\Database\\chunks\\")
|
DATABASE_PATH = os.getenv("CONJURER_LIBRARIAN_DB_PATH", r"C:\\Database\\chunks\\")
|
||||||
|
|
||||||
ENCODING = os.getenv("CONJURER_ENCODING", "utf-8")
|
ENCODING = os.getenv("CONJURER_ENCODING", "utf-8")
|
||||||
CHUNK = os.getenv("CONJURER_LIBRARIAN_CHUNK", "_chunk.txt")
|
CHUNK = os.getenv("CONJURER_LIBRARIAN_CHUNK", "_chunk.txt")
|
||||||
|
|
||||||
|
# DEPRECATED. The chunk files are now auto-discovered from DATABASE_PATH, so the
|
||||||
|
# thread count and the termination threshold both derive from what is actually
|
||||||
|
# on disk. This used to be BOTH "how many chunk files to read" AND "how many
|
||||||
|
# sentinels to wait for", which had to match exactly: set too low it silently
|
||||||
|
# skipped trailing chunks, set too high it referenced a nonexistent file whose
|
||||||
|
# producer crashed, starving the sentinel count and hanging the search forever.
|
||||||
|
# Kept only so old env files / references don't break; it no longer gates logic.
|
||||||
|
MAXTHREADS = int(os.getenv("CONJURER_LIBRARIAN_MAXTHREADS", "0"))
|
||||||
|
|
||||||
_sentinel = object()
|
_sentinel = object()
|
||||||
WORK_Q_SIZE = 35500000
|
WORK_Q_SIZE = 35500000
|
||||||
|
# Idle backstop: after this many consecutive empty seconds a consumer assumes
|
||||||
|
# the producers are done (or dead) and exits, so the search can never hang even
|
||||||
|
# if a sentinel were somehow lost. The primary, correct termination is still the
|
||||||
|
# sentinel count reaching the number of producers actually started.
|
||||||
|
EMPTY_LIMIT = 30
|
||||||
|
|
||||||
|
# Chunk files are named "<n>_chunk.txt" (suffix from CHUNK).
|
||||||
|
_CHUNK_RE = re.compile(r"^(\d+)" + re.escape(CHUNK) + r"$")
|
||||||
|
|
||||||
|
|
||||||
|
def discover_chunk_files(_logger):
|
||||||
|
"""Return the ``<n>_chunk.txt`` files present in DATABASE_PATH, numeric order.
|
||||||
|
|
||||||
|
Reading what exists (rather than files 0..MAXTHREADS-1) removes both historic
|
||||||
|
failure modes at once: no trailing chunk is ever silently skipped, and no
|
||||||
|
producer is ever pointed at a missing file, so it cannot crash before
|
||||||
|
emitting its sentinel and deadlock the consumers.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
names = os.listdir(DATABASE_PATH)
|
||||||
|
except OSError as exc:
|
||||||
|
_logger.error("Cannot list DOI database dir %s: %s", DATABASE_PATH, exc)
|
||||||
|
return []
|
||||||
|
indexed = []
|
||||||
|
for name in names:
|
||||||
|
match = _CHUNK_RE.match(name)
|
||||||
|
if match:
|
||||||
|
indexed.append((int(match.group(1)), name))
|
||||||
|
indexed.sort()
|
||||||
|
ordered = [name for _, name in indexed]
|
||||||
|
_logger.info("Discovered %d chunk files in %s", len(ordered), DATABASE_PATH)
|
||||||
|
return ordered
|
||||||
|
|
||||||
|
|
||||||
def producer(out_q, control_q, filename, _logger):
|
def producer(out_q, control_q, filename, _logger):
|
||||||
@@ -47,36 +89,45 @@ def producer(out_q, control_q, filename, _logger):
|
|||||||
filename (str): Name of the file.
|
filename (str): Name of the file.
|
||||||
_logger: Logger object for logging.
|
_logger: Logger object for logging.
|
||||||
"""
|
"""
|
||||||
with open(DATABASE_PATH + filename, "r", encoding=ENCODING) as operated_file:
|
try:
|
||||||
print(f"Worker {filename} ")
|
with open(DATABASE_PATH + filename, "r", encoding=ENCODING) as operated_file:
|
||||||
line_no = 0
|
print(f"Worker {filename} ")
|
||||||
while True:
|
line_no = 0
|
||||||
line = operated_file.readline()
|
while True:
|
||||||
line_no += 1
|
line = operated_file.readline()
|
||||||
print(f"\t \t \t \t \t \t W{filename}{line_no}\r", end="")
|
line_no += 1
|
||||||
|
print(f"\t \t \t \t \t \t W{filename}{line_no}\r", end="")
|
||||||
|
|
||||||
if not line:
|
if not line:
|
||||||
print(f"EOF {filename}")
|
print(f"EOF {filename}")
|
||||||
break
|
break
|
||||||
|
|
||||||
if not line:
|
out_q.put(line)
|
||||||
break
|
try:
|
||||||
out_q.put(line)
|
check = control_q.get(block=False)
|
||||||
try:
|
except Empty:
|
||||||
check = control_q.get(block=False)
|
check = False
|
||||||
except Empty:
|
|
||||||
check = False
|
|
||||||
|
|
||||||
if check is _sentinel:
|
if check is _sentinel:
|
||||||
print("TERM signal received")
|
print("TERM signal received")
|
||||||
control_q.put(check)
|
control_q.put(check)
|
||||||
break
|
break
|
||||||
print(f"Worker finished: {filename}")
|
print(f"Worker finished: {filename}")
|
||||||
|
except OSError as exc:
|
||||||
|
# A missing or unreadable chunk must not take the whole search down with
|
||||||
|
# it - log and move on. The sentinel below still fires (finally), so the
|
||||||
|
# consumers' count stays correct and nothing deadlocks.
|
||||||
|
_logger.warning("Chunk %s unreadable, skipping: %s", filename, exc)
|
||||||
|
print(f"Worker {filename} failed: {exc}")
|
||||||
|
finally:
|
||||||
|
# ALWAYS emit exactly one sentinel per producer, on every exit path (EOF,
|
||||||
|
# early TERM, or crash). This is what lets the consumers count producers
|
||||||
|
# deterministically instead of hanging on a lost sentinel.
|
||||||
out_q.put(_sentinel)
|
out_q.put(_sentinel)
|
||||||
|
|
||||||
|
|
||||||
# IMPORTANT!!! ONLY ONE CONSUMER THREAD AS WE ARE NOT PUTTING SENTINELS BACK
|
# IMPORTANT!!! ONLY ONE CONSUMER THREAD AS WE ARE NOT PUTTING SENTINELS BACK
|
||||||
def consumer(in_q, control_q, doi, live_results, result_list, control_dict, no, _logger):
|
def consumer(in_q, control_q, doi, live_results, result_list, control_dict, expected_sentinels, no, _logger):
|
||||||
"""
|
"""
|
||||||
Consumes items from an input queue and checks if DOI exists in the live results list.
|
Consumes items from an input queue and checks if DOI exists in the live results list.
|
||||||
|
|
||||||
@@ -117,14 +168,18 @@ def consumer(in_q, control_q, doi, live_results, result_list, control_dict, no,
|
|||||||
empty_counter += 1
|
empty_counter += 1
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
print(f"Consumer {no} empty")
|
print(f"Consumer {no} empty")
|
||||||
if empty_counter > 5:
|
# Order matters: the >EMPTY_LIMIT break must be checked BEFORE the
|
||||||
print("Consumer %s empty lvl 2", no)
|
# lesser threshold, otherwise (as in the original) the first branch
|
||||||
time.sleep(2)
|
# always wins and the break is dead code, leaving the sentinel count
|
||||||
elif empty_counter > 10:
|
# as the only exit - which is exactly what used to hang the search.
|
||||||
print(f"Consumer thread finished {no}")
|
if empty_counter > EMPTY_LIMIT:
|
||||||
|
print(f"Consumer thread finished {no} (idle backstop)")
|
||||||
break
|
break
|
||||||
|
if empty_counter > 5:
|
||||||
|
print(f"Consumer {no} empty lvl 2")
|
||||||
|
time.sleep(2)
|
||||||
|
|
||||||
if control_dict["sentinels"] >= MAXTHREADS:
|
if control_dict["sentinels"] >= expected_sentinels:
|
||||||
_logger.info(f"All workers finished {no}")
|
_logger.info(f"All workers finished {no}")
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -148,18 +203,25 @@ def search_for_doi(doi, live_results, _logger):
|
|||||||
for item in doi:
|
for item in doi:
|
||||||
result_list.append({"DOI": item[0], "exists": False, "data": item[1]})
|
result_list.append({"DOI": item[0], "exists": False, "data": item[1]})
|
||||||
|
|
||||||
|
# One producer per chunk file that actually exists; the sentinel threshold is
|
||||||
|
# that same count, so the two can never drift apart the way MAXTHREADS did.
|
||||||
|
chunk_files = discover_chunk_files(_logger)
|
||||||
|
expected = len(chunk_files)
|
||||||
|
if expected == 0:
|
||||||
|
_logger.error(
|
||||||
|
"No '<n>%s' chunk files in %s - DOI search cannot run", CHUNK, DATABASE_PATH
|
||||||
|
)
|
||||||
|
return result_list
|
||||||
|
|
||||||
for i in range (0, (len(doi)//1000)+2):
|
for i in range (0, (len(doi)//1000)+2):
|
||||||
t_cons = Thread(
|
t_cons = Thread(
|
||||||
target=consumer, args=(work_q, control_q, doi, live_results, result_list, control_dict, i, _logger)
|
target=consumer,
|
||||||
|
args=(work_q, control_q, doi, live_results, result_list, control_dict, expected, i, _logger),
|
||||||
)
|
)
|
||||||
_logger.info("Consumer thread created")
|
_logger.info("Consumer thread created")
|
||||||
threads.append(t_cons)
|
threads.append(t_cons)
|
||||||
for i in range(0, MAXTHREADS):
|
for filename in chunk_files:
|
||||||
# TEST DATA
|
_logger.info("Creating worker thread for %s", filename)
|
||||||
# filename = "test" + str(i) + CHUNk
|
|
||||||
# DEPLOYMENT DATA
|
|
||||||
filename = str(i) + CHUNK
|
|
||||||
_logger.info(f"Creating worker thread no: {i}")
|
|
||||||
threads.append(
|
threads.append(
|
||||||
Thread(target=producer, args=(work_q, control_q, filename, _logger))
|
Thread(target=producer, args=(work_q, control_q, filename, _logger))
|
||||||
)
|
)
|
||||||
|
|||||||
Vendored
+6
-2
@@ -12,10 +12,14 @@ CONJURER_MAIN_BOT=http://BOT_VM_IP:5000
|
|||||||
# Crossref polite-pool contact (or put credentials in netrc under "crossref").
|
# Crossref polite-pool contact (or put credentials in netrc under "crossref").
|
||||||
CONJURER_CROSSREF_MAILTO=you@example.com
|
CONJURER_CROSSREF_MAILTO=you@example.com
|
||||||
|
|
||||||
# Local DOI chunk database (mounted volume): expects 0_chunk.txt .. N_chunk.txt
|
# Local DOI chunk database (mounted volume): expects 0_chunk.txt .. N_chunk.txt.
|
||||||
|
# The chunk files are auto-discovered, so ALL of them are searched no matter how
|
||||||
|
# many there are - just drop them in this directory.
|
||||||
CONJURER_LIBRARIAN_DB_PATH=/doi/
|
CONJURER_LIBRARIAN_DB_PATH=/doi/
|
||||||
CONJURER_LIBRARIAN_MAXTHREADS=41
|
|
||||||
CONJURER_LIBRARIAN_CHUNK=_chunk.txt
|
CONJURER_LIBRARIAN_CHUNK=_chunk.txt
|
||||||
|
# DEPRECATED and unused: chunk files are now auto-discovered. It used to have to
|
||||||
|
# equal the file count exactly or the search would skip files / hang forever.
|
||||||
|
# CONJURER_LIBRARIAN_MAXTHREADS=41
|
||||||
|
|
||||||
# Runtime JSON state dir (mounted, persistent): cr_results/rr_results/
|
# Runtime JSON state dir (mounted, persistent): cr_results/rr_results/
|
||||||
# not_in_db/s_results are seeded here on first run.
|
# not_in_db/s_results are seeded here on first run.
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ The three talk to each other over HTTP on the Proxmox LAN. Direction of calls:
|
|||||||
bot --(/query)--------------------------------> librarian
|
bot --(/query)--------------------------------> librarian
|
||||||
musician --(/prepped_tracks)--------------------> bot
|
musician --(/prepped_tracks)--------------------> bot
|
||||||
librarian --(/conjurer results)-----------------> bot
|
librarian --(/conjurer results)-----------------> bot
|
||||||
|
* --(/ai_query)-----------------------------> bot (see 1c-ter)
|
||||||
```
|
```
|
||||||
|
|
||||||
Everything is configured through `CONJURER_*` environment variables (see the
|
Everything is configured through `CONJURER_*` environment variables (see the
|
||||||
@@ -117,6 +118,23 @@ generation (`imaginuje sobie:`) and personal assistants stay on OpenAI whatever
|
|||||||
the switch says (Anthropic has no equivalent) and degrade quietly if OpenAI is
|
the switch says (Anthropic has no equivalent) and degrade quietly if OpenAI is
|
||||||
not configured, so a Claude-only box still boots.
|
not configured, so a Claude-only box still boots.
|
||||||
|
|
||||||
|
### 1c-ter. AI query interface + librarian AI review
|
||||||
|
|
||||||
|
The bot exposes a queued AI interface on its comm layer: `POST /ai_query` with
|
||||||
|
`{"prompt": ..., "channel_id": <discord channel id>, "request_type"?: "NONE"}`
|
||||||
|
(same `X-Conjurer-Api-Key` auth as the other endpoints). The prompt is queued and
|
||||||
|
answered asynchronously by whatever backend `$gadaj_teraz` currently selects
|
||||||
|
(GPT or Claude), and the answer is posted to `channel_id`. `request_type: "NONE"`
|
||||||
|
(the default) keeps it a clean one-shot that doesn't touch the bar's
|
||||||
|
conversation memory.
|
||||||
|
|
||||||
|
The first consumer of this is the librarian command **`$wyszukaj_z_recenzja`**:
|
||||||
|
it works like `$wyszukaj_linki_do_dokumentow`, but when the DOI hits come back
|
||||||
|
the list (already sorted by Crossref relevance) plus the search phrase are handed
|
||||||
|
to the AI for a weighted-relevance re-rank and a short source review, delivered
|
||||||
|
to the same channel right after the raw results. No extra config — it uses the
|
||||||
|
active AI backend.
|
||||||
|
|
||||||
### 1d. Configure and launch
|
### 1d. Configure and launch
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
@@ -186,8 +204,11 @@ sudo mkdir -p /srv/librarian/doi /srv/librarian/secrets
|
|||||||
```
|
```
|
||||||
|
|
||||||
> The old Windows path `C:\Database\chunks\` is now `CONJURER_LIBRARIAN_DB_PATH`
|
> The old Windows path `C:\Database\chunks\` is now `CONJURER_LIBRARIAN_DB_PATH`
|
||||||
> (defaults to `/doi/` in the container). `CONJURER_LIBRARIAN_MAXTHREADS` (41)
|
> (defaults to `/doi/` in the container); `CONJURER_LIBRARIAN_CHUNK`
|
||||||
> and `CONJURER_LIBRARIAN_CHUNK` (`_chunk.txt`) are configurable too.
|
> (`_chunk.txt`) sets the suffix. Chunk files (`0_chunk.txt … N_chunk.txt`) are
|
||||||
|
> auto-discovered and all searched, so just drop them in — put as many as you
|
||||||
|
> like. (`CONJURER_LIBRARIAN_MAXTHREADS` is deprecated and ignored: it used to
|
||||||
|
> have to equal the file count exactly or the search would skip files or hang.)
|
||||||
|
|
||||||
### 2b. Configure and launch
|
### 2b. Configure and launch
|
||||||
|
|
||||||
|
|||||||
+94
-2
@@ -14,7 +14,7 @@ import requests
|
|||||||
from discord.ext import commands, tasks
|
from discord.ext import commands, tasks
|
||||||
|
|
||||||
from ai_functions import handle_response
|
from ai_functions import handle_response
|
||||||
from communication_subroutine import IN_COMM_Q, OUT_COMM_Q, QueryControl
|
from communication_subroutine import IN_COMM_Q, OUT_COMM_Q, QueryControl, submit_ai_query
|
||||||
from constants import DIR_PATH_SADOX, LIBRARIAN_SERVICE_ADDRESS, SEND_QUERY, service_headers
|
from constants import DIR_PATH_SADOX, LIBRARIAN_SERVICE_ADDRESS, SEND_QUERY, service_headers
|
||||||
|
|
||||||
SERVICE_HEADERS = service_headers()
|
SERVICE_HEADERS = service_headers()
|
||||||
@@ -97,14 +97,18 @@ class DataModule(commands.Cog):
|
|||||||
if fresh_data.stop:
|
if fresh_data.stop:
|
||||||
searcher = fresh_data.author
|
searcher = fresh_data.author
|
||||||
query = fresh_data.content
|
query = fresh_data.content
|
||||||
|
# ai_lines is a clean, plain rendering of the SAME list in the
|
||||||
|
# SAME (Crossref-relevance) order, for the optional AI review.
|
||||||
|
ai_lines = []
|
||||||
l_p = 1
|
l_p = 1
|
||||||
for doi in fresh_data.entries:
|
for doi in fresh_data.entries:
|
||||||
self.logger.info(doi)
|
self.logger.info(doi)
|
||||||
desc = fresh_data.entries[doi]
|
desc = fresh_data.entries[doi]
|
||||||
title = desc["Title"][0]
|
title = desc["Title"][0] if desc.get("Title") else "(bez tytułu)"
|
||||||
entries.append(
|
entries.append(
|
||||||
f"{l_p}. {title} pod linkiem https://www.sci-hub.se/{doi} i jest to {desc['type']}\n"
|
f"{l_p}. {title} pod linkiem https://www.sci-hub.se/{doi} i jest to {desc['type']}\n"
|
||||||
)
|
)
|
||||||
|
ai_lines.append(f"{l_p}. {title} (DOI: {doi}, typ: {desc['type']})")
|
||||||
l_p += 1
|
l_p += 1
|
||||||
message = "*Z podłogi wysuwa się winda na książki*"
|
message = "*Z podłogi wysuwa się winda na książki*"
|
||||||
if fresh_data.ctx is not None:
|
if fresh_data.ctx is not None:
|
||||||
@@ -131,6 +135,36 @@ class DataModule(commands.Cog):
|
|||||||
await ctx.send(message)
|
await ctx.send(message)
|
||||||
message = ""
|
message = ""
|
||||||
|
|
||||||
|
# Optional AI pass: re-rank the (already Crossref-relevance-
|
||||||
|
# sorted) DOI list and review the sources. Enqueued to the AI
|
||||||
|
# worker so it runs on whatever backend $gadaj_teraz selected;
|
||||||
|
# the answer lands in this same channel.
|
||||||
|
if getattr(fresh_data, "ai_review", False) and ai_lines:
|
||||||
|
target = getattr(ctx, "channel", ctx)
|
||||||
|
review_prompt = (
|
||||||
|
f'Poniżej lista źródeł naukowych znalezionych dla zapytania: "{query}".\n'
|
||||||
|
"Lista jest już wstępnie posortowana według trafności wg Crossref "
|
||||||
|
"(od najtrafniejszej).\n\n"
|
||||||
|
"Twoje zadania:\n"
|
||||||
|
"1. Przeważ i uporządkuj listę według RZECZYWISTEJ trafności do zapytania "
|
||||||
|
"(najtrafniejsze u góry).\n"
|
||||||
|
"2. Do każdej pozycji dopisz jedno-, dwuzdaniową recenzję: typ i wiarygodność "
|
||||||
|
"źródła oraz dlaczego (nie) pasuje do zapytania.\n"
|
||||||
|
"Odpowiedz zwięźle, numerowaną listą, po polsku.\n\n"
|
||||||
|
"Źródła:\n" + "\n".join(ai_lines)
|
||||||
|
)
|
||||||
|
submit_ai_query(
|
||||||
|
prompt=review_prompt,
|
||||||
|
channel_id=target.id,
|
||||||
|
request_type="NONE",
|
||||||
|
username=searcher,
|
||||||
|
query_uuid=str(fresh_data.uuid),
|
||||||
|
)
|
||||||
|
await ctx.send(
|
||||||
|
"*Conjurer podaje listę naszemu rezydentowi-mądrali od AI* "
|
||||||
|
"Za chwilę dorzuci recenzję i swoje przesortowanie wg trafności."
|
||||||
|
)
|
||||||
|
|
||||||
# Kept for sentimental reasons
|
# Kept for sentimental reasons
|
||||||
# await ctx.send(f"O. A tak będzie wyglądało coś ciekawego w przyszłości: {data}")
|
# await ctx.send(f"O. A tak będzie wyglądało coś ciekawego w przyszłości: {data}")
|
||||||
except Empty:
|
except Empty:
|
||||||
@@ -199,6 +233,64 @@ class DataModule(commands.Cog):
|
|||||||
+ " Zapytania obsługuje algorytm zasilany czterema chomikami zapierdalającymi w kołowrotku - więc wyniki najwcześniej za kilka godzi - ale mogą być też dni."
|
+ " Zapytania obsługuje algorytm zasilany czterema chomikami zapierdalającymi w kołowrotku - więc wyniki najwcześniej za kilka godzi - ale mogą być też dni."
|
||||||
)
|
)
|
||||||
|
|
||||||
|
@commands.hybrid_command(
|
||||||
|
name="wyszukaj_z_recenzja",
|
||||||
|
description="Jak wyszukaj_linki_do_dokumentow, ale wyniki przesortuje trafnością i zrecenzuje AI",
|
||||||
|
guild=discord.Object(id=664789470779932693),
|
||||||
|
)
|
||||||
|
async def wyszukaj_z_recenzja(self, ctx):
|
||||||
|
"""Same as wyszukaj_linki_do_dokumentow, but flags the search for an AI
|
||||||
|
review: when the DOI hits come back, the list + the search phrase are sent
|
||||||
|
to the bot's AI backend for a weighted-relevance re-rank and a source
|
||||||
|
review, delivered to this channel. The flag rides on the QueryControl so
|
||||||
|
it survives the round-trip and is matched back to this search by UUID.
|
||||||
|
"""
|
||||||
|
query = ctx.message.content
|
||||||
|
query_uuid = uuid.uuid4()
|
||||||
|
ctx.message.content = ctx.message.content.replace("$wyszukaj_z_recenzja", "")
|
||||||
|
|
||||||
|
json_query = {
|
||||||
|
"UUID": str(query_uuid),
|
||||||
|
"query": str(query),
|
||||||
|
"page": 1,
|
||||||
|
"deep_search": False,
|
||||||
|
}
|
||||||
|
coroutine = asyncio.to_thread(
|
||||||
|
requests.post,
|
||||||
|
f"{LIBRARIAN_SERVICE_ADDRESS}{SEND_QUERY}",
|
||||||
|
json=json_query,
|
||||||
|
headers=SERVICE_HEADERS,
|
||||||
|
timeout=360,
|
||||||
|
)
|
||||||
|
await ctx.send(
|
||||||
|
"*Conjurer notuje, wrzuca liścik do rury pneumatycznej i mruży oko* Tym razem jak coś"
|
||||||
|
+ " znajdę, przepuszczę wyniki jeszcze przez naszego rezydenta-mądralę od AI - przeważy"
|
||||||
|
+ " trafność i zrecenzuje źródła. Poczekaj kilka godzin - biblioteka to 3/4 stacji."
|
||||||
|
)
|
||||||
|
query_response = await coroutine
|
||||||
|
if not query_response.status_code == 200:
|
||||||
|
await ctx.send(
|
||||||
|
"*Z rury wydobywa się dym. Conjurer pryska w nią pierwszą cieczą pod ręką i wybucha"
|
||||||
|
+ " drobny pożar.* Wołaj szefa - mam wrażenie że się coś wyjebało"
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
|
query, query_uuid, queue_size = (
|
||||||
|
query_response.json()["data"][0],
|
||||||
|
query_response.json()["data"][1],
|
||||||
|
query_response.json()["data"][2],
|
||||||
|
)
|
||||||
|
if ctx.message.author.nick:
|
||||||
|
username = ctx.message.author.nick
|
||||||
|
else:
|
||||||
|
username = ctx.message.author.name
|
||||||
|
query_object = QueryControl(username, query_uuid, query, ctx, ai_review=True)
|
||||||
|
OUT_COMM_Q.put(query_object)
|
||||||
|
await ctx.send(
|
||||||
|
f"Poszło z recenzją AI. Identyfikator: {query_uuid}. Jesteś {queue_size} w kolejce."
|
||||||
|
+ " Najpierw dojadą surowe wyniki, a zaraz po nich przesortowanie i recenzja od AI."
|
||||||
|
)
|
||||||
|
|
||||||
@commands.hybrid_command(
|
@commands.hybrid_command(
|
||||||
name="glebokie_gardlo",
|
name="glebokie_gardlo",
|
||||||
description="Przygotowuje drinka o nazwie głębokie gardło",
|
description="Przygotowuje drinka o nazwie głębokie gardło",
|
||||||
|
|||||||
@@ -0,0 +1,70 @@
|
|||||||
|
"""Integration: the bot's /ai_query endpoint enforces the shared key, validates
|
||||||
|
the payload, and queues accepted prompts onto AI_QUERY_Q for the AI worker.
|
||||||
|
"""
|
||||||
|
import communication_subroutine as cs
|
||||||
|
|
||||||
|
|
||||||
|
def _client(key="test-secret"):
|
||||||
|
cs.API_KEY = key
|
||||||
|
return cs.app.test_client()
|
||||||
|
|
||||||
|
|
||||||
|
def _drain():
|
||||||
|
while not cs.AI_QUERY_Q.empty():
|
||||||
|
cs.AI_QUERY_Q.get()
|
||||||
|
|
||||||
|
|
||||||
|
def test_ai_query_rejected_without_key():
|
||||||
|
_drain()
|
||||||
|
client = _client()
|
||||||
|
resp = client.post("/ai_query", json={"prompt": "x", "channel_id": 1})
|
||||||
|
assert resp.status_code == 401
|
||||||
|
assert cs.AI_QUERY_Q.empty() # nothing queued on a rejected call
|
||||||
|
|
||||||
|
|
||||||
|
def test_ai_query_accepted_with_key_and_queued():
|
||||||
|
_drain()
|
||||||
|
client = _client()
|
||||||
|
resp = client.post(
|
||||||
|
"/ai_query",
|
||||||
|
json={"prompt": "posortuj DOI", "channel_id": 42, "username": "siara"},
|
||||||
|
headers={"X-Conjurer-Api-Key": "test-secret"},
|
||||||
|
)
|
||||||
|
assert resp.status_code == 200
|
||||||
|
item = cs.AI_QUERY_Q.get()
|
||||||
|
assert item["prompt"] == "posortuj DOI"
|
||||||
|
assert item["channel_id"] == 42
|
||||||
|
assert item["username"] == "siara"
|
||||||
|
|
||||||
|
|
||||||
|
def test_ai_query_missing_prompt_is_rejected():
|
||||||
|
_drain()
|
||||||
|
client = _client()
|
||||||
|
resp = client.post(
|
||||||
|
"/ai_query",
|
||||||
|
json={"channel_id": 1},
|
||||||
|
headers={"X-Conjurer-Api-Key": "test-secret"},
|
||||||
|
)
|
||||||
|
assert resp.status_code == 400
|
||||||
|
assert cs.AI_QUERY_Q.empty()
|
||||||
|
|
||||||
|
|
||||||
|
def test_ai_query_open_when_key_unset():
|
||||||
|
_drain()
|
||||||
|
client = _client(key=None)
|
||||||
|
resp = client.post("/ai_query", json={"prompt": "y", "channel_id": 1})
|
||||||
|
assert resp.status_code == 200
|
||||||
|
|
||||||
|
|
||||||
|
def test_submit_ai_query_enqueues_expected_shape():
|
||||||
|
_drain()
|
||||||
|
cs.submit_ai_query(
|
||||||
|
prompt="P", channel_id=7, request_type="NONE", username="u", query_uuid="uid-1"
|
||||||
|
)
|
||||||
|
assert cs.AI_QUERY_Q.get() == {
|
||||||
|
"uuid": "uid-1",
|
||||||
|
"prompt": "P",
|
||||||
|
"channel_id": 7,
|
||||||
|
"request_type": "NONE",
|
||||||
|
"username": "u",
|
||||||
|
}
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
"""Integration: betoniarka enforces the shared key on /clear_pr_pls (which
|
||||||
|
moved here from the musician during the radio split) while /ping stays open.
|
||||||
|
|
||||||
|
This is the same auth contract the musician test used to cover for
|
||||||
|
/clear_pr_pls, now exercised against the service that actually owns it.
|
||||||
|
"""
|
||||||
|
import betoniarka as b
|
||||||
|
|
||||||
|
|
||||||
|
def _client(tmp_path, key="test-secret"):
|
||||||
|
b.API_KEY = key
|
||||||
|
# /clear_pr_pls truncates this file; point it at a writable temp path so the
|
||||||
|
# authorised case reaches 200 instead of failing on the default
|
||||||
|
# /srv/betoniarka/data path that does not exist in CI.
|
||||||
|
b.PRIORITY_PLAYLIST_PATH = tmp_path / "priority_queue.playlist"
|
||||||
|
return b.app.test_client()
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_pr_pls_rejected_without_key(tmp_path):
|
||||||
|
client = _client(tmp_path)
|
||||||
|
assert client.get("/clear_pr_pls").status_code == 401
|
||||||
|
|
||||||
|
|
||||||
|
def test_clear_pr_pls_accepted_with_key(tmp_path):
|
||||||
|
client = _client(tmp_path)
|
||||||
|
resp = client.get(
|
||||||
|
"/clear_pr_pls", headers={"X-Conjurer-Api-Key": "test-secret"}
|
||||||
|
)
|
||||||
|
assert resp.status_code == 200
|
||||||
|
# The endpoint's job is to empty the priority playlist.
|
||||||
|
assert b.PRIORITY_PLAYLIST_PATH.read_text() == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_ping_is_open(tmp_path):
|
||||||
|
client = _client(tmp_path)
|
||||||
|
assert client.get("/ping").status_code == 200
|
||||||
|
|
||||||
|
|
||||||
|
def test_open_when_key_unset(tmp_path):
|
||||||
|
client = _client(tmp_path, key=None)
|
||||||
|
assert client.get("/clear_pr_pls").status_code == 200
|
||||||
@@ -1,23 +1,36 @@
|
|||||||
"""Integration: the musician Flask service enforces the shared key on its
|
"""Integration: the musician Flask service enforces the shared key on its
|
||||||
authenticated endpoints while leaving the open ones reachable.
|
authenticated endpoints while leaving the open ones reachable.
|
||||||
|
|
||||||
|
/clear_pr_pls used to live here but moved to betoniarka during the
|
||||||
|
musician/radio split - its auth contract is now covered in
|
||||||
|
test_betoniarka_auth.py. These tests exercise the same contract against an
|
||||||
|
endpoint the musician still serves.
|
||||||
"""
|
"""
|
||||||
import conjurer_musician as m
|
import conjurer_musician as m
|
||||||
|
|
||||||
|
# An authenticated endpoint the musician still owns. The body passes the
|
||||||
|
# endpoint's own validation, so a permitted request reaches 200 rather than a
|
||||||
|
# 400 that would not distinguish auth from a bad payload.
|
||||||
|
AUTHED_ENDPOINT = "/get_share_list"
|
||||||
|
VALID_BODY = {"entries": 1, "keywords": ["conjurer"]}
|
||||||
|
|
||||||
|
|
||||||
def _client(key="test-secret"):
|
def _client(key="test-secret"):
|
||||||
m.API_KEY = key
|
m.API_KEY = key
|
||||||
return m.app.test_client()
|
return m.app.test_client()
|
||||||
|
|
||||||
|
|
||||||
def test_clear_pr_pls_rejected_without_key():
|
def test_authed_endpoint_rejected_without_key():
|
||||||
client = _client()
|
client = _client()
|
||||||
assert client.get("/clear_pr_pls").status_code == 401
|
assert client.post(AUTHED_ENDPOINT, json=VALID_BODY).status_code == 401
|
||||||
|
|
||||||
|
|
||||||
def test_clear_pr_pls_accepted_with_key():
|
def test_authed_endpoint_accepted_with_key():
|
||||||
client = _client()
|
client = _client()
|
||||||
resp = client.get(
|
resp = client.post(
|
||||||
"/clear_pr_pls", headers={"X-Conjurer-Api-Key": "test-secret"}
|
AUTHED_ENDPOINT,
|
||||||
|
json=VALID_BODY,
|
||||||
|
headers={"X-Conjurer-Api-Key": "test-secret"},
|
||||||
)
|
)
|
||||||
assert resp.status_code == 200
|
assert resp.status_code == 200
|
||||||
|
|
||||||
@@ -31,4 +44,4 @@ def test_mp3_list_is_open():
|
|||||||
|
|
||||||
def test_open_when_key_unset():
|
def test_open_when_key_unset():
|
||||||
client = _client(key=None)
|
client = _client(key=None)
|
||||||
assert client.get("/clear_pr_pls").status_code == 200
|
assert client.post(AUTHED_ENDPOINT, json=VALID_BODY).status_code == 200
|
||||||
|
|||||||
@@ -0,0 +1,89 @@
|
|||||||
|
"""Unit tests for the librarian DOI search - specifically that it always
|
||||||
|
terminates.
|
||||||
|
|
||||||
|
The producer/consumer search used to hang whenever the number of chunk files it
|
||||||
|
was told to read (MAXTHREADS) did not exactly match the files on disk: too few
|
||||||
|
and it silently skipped trailing chunks, too many and a producer pointed at a
|
||||||
|
missing file crashed before emitting its sentinel, starving the consumers'
|
||||||
|
termination count forever. These tests pin down the fix: chunk files are
|
||||||
|
auto-discovered, the sentinel threshold equals the number of producers actually
|
||||||
|
started, and every producer emits its sentinel even on error.
|
||||||
|
"""
|
||||||
|
import logging
|
||||||
|
import threading
|
||||||
|
|
||||||
|
import search_bot
|
||||||
|
|
||||||
|
_LOG = logging.getLogger("test-search-bot")
|
||||||
|
_LOG.addHandler(logging.NullHandler())
|
||||||
|
|
||||||
|
|
||||||
|
def _write_chunks(directory, count, target=None, target_index=None):
|
||||||
|
"""Create <n>_chunk.txt files; optionally drop `target` into one of them."""
|
||||||
|
for n in range(count):
|
||||||
|
lines = [f"10.0000/decoy-{n}-a\n", f"10.0000/decoy-{n}-b\n"]
|
||||||
|
if target is not None and n == target_index:
|
||||||
|
lines.append(target + "\n")
|
||||||
|
(directory / f"{n}_chunk.txt").write_text("".join(lines), encoding="utf-8")
|
||||||
|
|
||||||
|
|
||||||
|
def _run_bounded(dois, timeout=20):
|
||||||
|
"""Run search_for_doi in a thread; return (finished_in_time, result)."""
|
||||||
|
box = {}
|
||||||
|
live = []
|
||||||
|
worker = threading.Thread(
|
||||||
|
target=lambda: box.update(result=search_bot.search_for_doi(dois, live, _LOG)),
|
||||||
|
daemon=True,
|
||||||
|
)
|
||||||
|
worker.start()
|
||||||
|
worker.join(timeout)
|
||||||
|
return (not worker.is_alive()), box.get("result")
|
||||||
|
|
||||||
|
|
||||||
|
def test_finds_doi_in_trailing_chunk(tmp_path, monkeypatch):
|
||||||
|
# Target lives in the LAST chunk - the one the old MAXTHREADS=N-too-low would
|
||||||
|
# never have read. Auto-discovery must read every chunk present.
|
||||||
|
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
||||||
|
target = "10.1234/target.in.trailing.chunk"
|
||||||
|
_write_chunks(tmp_path, count=6, target=target, target_index=5)
|
||||||
|
|
||||||
|
finished, result = _run_bounded([(target, "DATA"), ("10.9999/absent", "DATA")])
|
||||||
|
|
||||||
|
assert finished, "search hung instead of terminating"
|
||||||
|
hit = [r for r in result if r["DOI"] == target and r["exists"]]
|
||||||
|
assert hit, "DOI in the trailing chunk was not found"
|
||||||
|
|
||||||
|
|
||||||
|
def test_terminates_when_a_chunk_is_unreadable(tmp_path, monkeypatch):
|
||||||
|
# A chunk that exists at discovery time but cannot be opened (here: it is a
|
||||||
|
# directory) makes its producer raise. The finally-sentinel must still fire
|
||||||
|
# so the consumers' count completes and the search does not deadlock.
|
||||||
|
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
||||||
|
_write_chunks(tmp_path, count=3)
|
||||||
|
(tmp_path / "9_chunk.txt").mkdir() # discovered as a chunk, un-openable
|
||||||
|
|
||||||
|
finished, _ = _run_bounded([("10.0000/decoy-0-a", "DATA")])
|
||||||
|
|
||||||
|
assert finished, "an unreadable chunk deadlocked the search"
|
||||||
|
|
||||||
|
|
||||||
|
def test_no_chunks_returns_immediately(tmp_path, monkeypatch):
|
||||||
|
# Empty database dir: return an (all-not-found) result at once, never hang.
|
||||||
|
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
||||||
|
|
||||||
|
finished, result = _run_bounded([("10.0/x", "DATA")], timeout=10)
|
||||||
|
|
||||||
|
assert finished
|
||||||
|
assert result == [{"DOI": "10.0/x", "exists": False, "data": "DATA"}]
|
||||||
|
|
||||||
|
|
||||||
|
def test_discover_chunk_files_sorted_numerically(tmp_path, monkeypatch):
|
||||||
|
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
|
||||||
|
for n in (0, 2, 10, 1):
|
||||||
|
(tmp_path / f"{n}_chunk.txt").write_text("x\n", encoding="utf-8")
|
||||||
|
(tmp_path / "notes.txt").write_text("ignore me\n", encoding="utf-8")
|
||||||
|
|
||||||
|
found = search_bot.discover_chunk_files(_LOG)
|
||||||
|
|
||||||
|
# Numeric order (10 after 2, not lexicographic), and non-chunk files ignored.
|
||||||
|
assert found == ["0_chunk.txt", "1_chunk.txt", "2_chunk.txt", "10_chunk.txt"]
|
||||||
Reference in New Issue
Block a user