711ce8c0c1
Librarian.__init__ reads the Crossref contact from CONJURER_CROSSREF_MAILTO, then tries to override it from a 'crossref' netrc entry. When no netrc is mounted (the normal container setup - default /root/.netrc) the read raises FileNotFoundError and it logged 'Crossref credentials missing in netrc ...' on EVERY search, even though the env var was set and used. Pure noise. Only warn when there is genuinely no contact from either source (env unset AND netrc unreadable) - which is also the case that then raises. When the env var is set, a missing netrc is expected and logged at debug. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
782 lines
29 KiB
Python
782 lines
29 KiB
Python
"""
|
|
This module contains the implementation of the Librarian class and
|
|
related functions for searching and refining queries.
|
|
|
|
Classes:
|
|
- Librarian: Represents a librarian object that performs search and
|
|
refinement operations on queries
|
|
|
|
Functions:
|
|
- flask_debug: Starts a Flask application in debug mode without using the reloader.
|
|
- waitress_run: Serves the Flask application using the Waitress WSGI server.
|
|
- BackgroundTaskSearch: Represents a background task for running the
|
|
Librarian object asynchronously
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import os
|
|
import threading
|
|
import time
|
|
from json.decoder import JSONDecodeError
|
|
from logging import handlers
|
|
from pathlib import Path
|
|
from queue import Queue
|
|
from typing import Dict, Optional
|
|
|
|
import requests
|
|
import lib_paths
|
|
import scrape_bot
|
|
import search_bot
|
|
# import search_bot2 as search_bot
|
|
from durable_queue import DiskQueue
|
|
from flask import Flask, jsonify, request, abort
|
|
from habanero import Crossref
|
|
from waitress import serve
|
|
|
|
try:
|
|
import netrc
|
|
except ImportError: # pragma: no cover
|
|
netrc = None
|
|
|
|
# Constants
|
|
|
|
|
|
def _env(name: str, default: str) -> str:
|
|
return os.getenv(name, default)
|
|
|
|
|
|
def _env_path(name: str, default: str) -> Path:
|
|
return Path(os.getenv(name, default)).expanduser().resolve()
|
|
|
|
|
|
BASE_DIR = Path(
|
|
os.getenv("CONJURER_LIBRARIAN_BASE", str(Path(__file__).resolve().parent))
|
|
)
|
|
NETRC_FILE = _env_path("CONJURER_NETRC_FILE", str(Path.home() / ".netrc"))
|
|
HOST_ADDRESS = _env("CONJURER_LIBRARIAN_HOST", "0.0.0.0")
|
|
PORT_ADDRESS = int(_env("CONJURER_LIBRARIAN_PORT", "5001"))
|
|
MAIN_BOT_ADDRESS = _env("CONJURER_MAIN_BOT", "http://127.0.0.1:5000")
|
|
SEND_RESULTS = _env("CONJURER_LIBRARIAN_RESULTS_ENDPOINT", "/conjurer")
|
|
MAX_CR_RESULTS = int(_env("CONJURER_LIBRARIAN_MAX_RESULTS", "500"))
|
|
ENCODING = _env("CONJURER_ENCODING", "utf-8")
|
|
API_KEY = os.getenv("CONJURER_API_KEY")
|
|
LOGFILE_PATH = _env_path(
|
|
"CONJURER_LIBRARIAN_LOG", str(BASE_DIR / "librarian.log")
|
|
)
|
|
|
|
# Durable OUTBOX for finished results. A result is expensive (hours of compute),
|
|
# so it is written here and only removed once the bot ACKs it (HTTP 200). Lives
|
|
# on the librarian's persistent state volume, so it survives a librarian restart
|
|
# and a transient bot outage; the resender thread keeps retrying until delivered.
|
|
OUTBOX_DIR = _env("CONJURER_LIBRARIAN_OUTBOX", os.path.join(lib_paths.STATE_DIR, "outbox"))
|
|
RESULT_SEND_ATTEMPTS = int(_env("CONJURER_RESULT_SEND_ATTEMPTS", "3"))
|
|
RESULT_SEND_BACKOFF = float(_env("CONJURER_RESULT_SEND_BACKOFF", "2"))
|
|
OUTBOX_RESEND_SECONDS = int(_env("CONJURER_OUTBOX_RESEND_SECONDS", "60"))
|
|
_outbox = DiskQueue(OUTBOX_DIR)
|
|
|
|
app = Flask(__name__)
|
|
|
|
librarian_queue = Queue()
|
|
librarian_list = []
|
|
|
|
# Lifecycle of every real search uuid: "queued" (accepted, sitting in
|
|
# librarian_queue) -> "processing" (worker pulled it) -> removed (worker
|
|
# finished AND attempted to send the result). The bot's per-query watchdog polls
|
|
# /query_status against this: a uuid that VANISHES from here without its result
|
|
# reaching the bot is a lost result (finished-but-never-delivered) and gets
|
|
# flagged in chat.
|
|
active_queries: Dict[str, str] = {}
|
|
_active_lock = threading.Lock()
|
|
# Set while the worker is grinding a real search. A ping arriving during this
|
|
# pongs back immediately WITHOUT queueing - being busy is healthy (you can keep
|
|
# piling searches on), so "busy" must never look like "dead" to the health check.
|
|
worker_busy = threading.Event()
|
|
|
|
|
|
def _service_headers() -> Dict[str, str]:
|
|
if API_KEY:
|
|
return {"X-Conjurer-Api-Key": API_KEY}
|
|
return {}
|
|
|
|
|
|
def _authorize_request() -> None:
|
|
if API_KEY and request.headers.get("X-Conjurer-Api-Key") != API_KEY:
|
|
abort(401)
|
|
|
|
|
|
def _post_pong(app_logger, ping_uuid) -> None:
|
|
"""POST a pong for ``ping_uuid`` back to the bot. Non-fatal on failure.
|
|
|
|
This is the SAME return path a real result takes (bot's /conjurer), so a
|
|
delivered pong proves the librarian->bot leg works - the one thing the ping
|
|
needs to establish. Used directly by the /ping route (busy ping, which skips
|
|
the queue) and, wrapped in a thread, by the worker (idle ping). Synchronous
|
|
so the /ping route can stay a plain (non-async) view."""
|
|
try:
|
|
requests.post(
|
|
f"{MAIN_BOT_ADDRESS}{SEND_RESULTS}",
|
|
json={"__pong__": ping_uuid},
|
|
headers=_service_headers(),
|
|
timeout=5,
|
|
)
|
|
except requests.exceptions.RequestException as exc:
|
|
app_logger.warning("PING pong send failed for %s: %s", ping_uuid, exc)
|
|
|
|
|
|
def _deliver_result(uuid, payload, app_logger, attempts=RESULT_SEND_ATTEMPTS) -> bool:
|
|
"""POST one result to the bot, retrying with backoff. True only on HTTP 200.
|
|
|
|
The bot's /conjurer is idempotent (dedups by uuid), so re-POSTing a result
|
|
it already has is safe - it just answers 200 again. That is what lets the
|
|
OUTBOX keep retrying until the result is truly acknowledged, without ever
|
|
double-delivering to the user.
|
|
"""
|
|
target = f"{MAIN_BOT_ADDRESS}{SEND_RESULTS}"
|
|
for attempt in range(1, max(1, attempts) + 1):
|
|
try:
|
|
response = requests.post(
|
|
target, json=payload, headers=_service_headers(), timeout=60
|
|
)
|
|
if response.status_code == 200:
|
|
app_logger.info("Result %s delivered (HTTP 200) on attempt %d", uuid, attempt)
|
|
return True
|
|
app_logger.warning(
|
|
"Result %s: bot returned HTTP %s (attempt %d/%d): %s",
|
|
uuid, response.status_code, attempt, attempts, response.text[:300],
|
|
)
|
|
except requests.exceptions.RequestException as exc:
|
|
app_logger.warning(
|
|
"Result %s delivery failed (attempt %d/%d): %s", uuid, attempt, attempts, exc
|
|
)
|
|
if attempt < attempts:
|
|
time.sleep(RESULT_SEND_BACKOFF * attempt)
|
|
return False
|
|
|
|
|
|
def _resend_once(app_logger) -> None:
|
|
"""One sweep of the OUTBOX: try to deliver every un-acked result, once each.
|
|
|
|
Removes each entry only after a positive ACK, so nothing is dropped until
|
|
the bot has it. Corrupt/unreadable entries are skipped by DiskQueue.items().
|
|
"""
|
|
for uuid, payload, _ts in _outbox.items():
|
|
if _deliver_result(uuid, payload, app_logger, attempts=1):
|
|
_outbox.remove(uuid)
|
|
|
|
|
|
def outbox_resender(app_logger) -> None:
|
|
"""Background loop: periodically flush the OUTBOX until the bot is reachable.
|
|
|
|
This is what makes an expensive result survive a transient bot outage or a
|
|
librarian restart - on restart the persisted OUTBOX is simply resent."""
|
|
pending = len(_outbox)
|
|
if pending:
|
|
app_logger.info("OUTBOX has %d un-acked result(s) on startup - will resend", pending)
|
|
while True:
|
|
try:
|
|
_resend_once(app_logger)
|
|
except Exception as exc: # pylint: disable=broad-exception-caught
|
|
app_logger.exception("OUTBOX resend sweep failed: %s", exc)
|
|
time.sleep(OUTBOX_RESEND_SECONDS)
|
|
|
|
|
|
# trunk-ignore(pylint/R0902)
|
|
class Librarian(object):
|
|
"""
|
|
Represents a librarian object that performs search and refinement operations on queries.
|
|
"""
|
|
|
|
def __init__(self, _app, query, uuid, _deep_search) -> None:
|
|
"""
|
|
Initializes a Librarian object.
|
|
|
|
Args:
|
|
- _app: The Flask application object.
|
|
- query: The query to be searched.
|
|
- uuid: The unique identifier for the search.
|
|
|
|
Attributes:
|
|
- cr: The Crossref object for performing the search.
|
|
- query: The query to be searched.
|
|
- uuid: The unique identifier for the search.
|
|
- limit: The maximum number of search results to fetch.
|
|
- fetched: The number of search results fetched so far.
|
|
- hit: The number of search results that match the refinement criteria.
|
|
- total: The total number of search results.
|
|
- app: The Flask application object.
|
|
- live_results: A list to store live search results.
|
|
- final_result: A list to store the final refined search results.
|
|
- not_in_db: A list to store search results that are not in the local database.
|
|
- search_result_from_cr: A dictionary to store the search results from Crossref.
|
|
- done: A flag indicating if the search is done.
|
|
"""
|
|
# Crossref only needs a contact mailto. It can come from
|
|
# CONJURER_CROSSREF_MAILTO (the usual container setup) OR from a
|
|
# "crossref" entry in the netrc; netrc takes precedence when present.
|
|
mailto_contact: Optional[str] = os.getenv("CONJURER_CROSSREF_MAILTO")
|
|
if netrc:
|
|
try:
|
|
netrc_mod = netrc.netrc(str(NETRC_FILE))
|
|
auth_tokens = netrc_mod.authenticators("crossref")
|
|
if auth_tokens:
|
|
mailto_contact = auth_tokens[0]
|
|
except (FileNotFoundError, netrc.NetrcParseError) as exc:
|
|
# A missing/unreadable netrc is NORMAL when the mailto is set via
|
|
# env - don't cry wolf on every single search. Only warn when we
|
|
# genuinely have no contact from either source.
|
|
_log = logging.getLogger("conjurer_librarian")
|
|
if mailto_contact:
|
|
_log.debug(
|
|
"netrc %s not used (%s) - using CONJURER_CROSSREF_MAILTO",
|
|
NETRC_FILE, exc,
|
|
)
|
|
else:
|
|
_log.warning(
|
|
"Crossref contact not configured: netrc %s unreadable (%s) "
|
|
"and CONJURER_CROSSREF_MAILTO unset",
|
|
NETRC_FILE, exc,
|
|
)
|
|
if not mailto_contact:
|
|
raise RuntimeError(
|
|
"Crossref credentials not configured. Set CONJURER_CROSSREF_MAILTO or add to netrc."
|
|
)
|
|
self.cr = Crossref(
|
|
mailto=mailto_contact,
|
|
ua_string=f"Conjurer project. mailto:{mailto_contact}"
|
|
)
|
|
self.query = query
|
|
self.uuid = str(uuid)
|
|
self.limit = MAX_CR_RESULTS
|
|
self.fetched = 0
|
|
self.hit = 0
|
|
self.total = 0
|
|
self.app = _app
|
|
self.live_results = []
|
|
self.final_result = {}
|
|
self.not_in_db = {}
|
|
self.search_result_from_cr = {}
|
|
self.done = False
|
|
self.deep_search = _deep_search
|
|
|
|
async def search_crossref(self, query, deep_search=False):
|
|
"""
|
|
Performs a search on Crossref for the given query.
|
|
|
|
Args:
|
|
- query: The query to be searched.
|
|
|
|
Returns:
|
|
- result: The search result from Crossref.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
|
|
self.app.logger.info("STARTED SEARCH")
|
|
|
|
if not deep_search:
|
|
query_limit = MAX_CR_RESULTS if MAX_CR_RESULTS < 1000 else 1000
|
|
cr_result = self.cr.works(query=query, limit=query_limit)
|
|
self.search_result_from_cr.update(cr_result)
|
|
self.total = cr_result["message"]["total-results"]
|
|
self.fetched += len(cr_result["message"]["items"])
|
|
self.app.logger.info(self.total)
|
|
self.app.logger.info(self.fetched)
|
|
while self.total > self.fetched and self.limit > self.fetched:
|
|
tmp_result = self.cr.works(query=query, limit=query_limit, offset=self.fetched)
|
|
cr_result["message"]["items"].extend(tmp_result["message"]["items"])
|
|
self.total = tmp_result["message"]["total-results"]
|
|
self.fetched = len(cr_result["message"]["items"])
|
|
self.app.logger.info(self.total)
|
|
self.app.logger.info(self.fetched)
|
|
await asyncio.sleep(0.1)
|
|
|
|
else:
|
|
cr_result = self.cr.works(query=query, cursor_max=15000, cursor='*', progress_bar = True)
|
|
result = cr_result[0]
|
|
for item in cr_result[1:]:
|
|
result["message"]["items"].extend(item["message"]["items"])
|
|
self.total = item["message"]["total-results"]
|
|
self.fetched = len(result["message"]["items"])
|
|
self.app.logger.info(self.total)
|
|
self.app.logger.info(self.fetched)
|
|
cr_result = result
|
|
self.app.logger.info("Total, fetched:")
|
|
self.app.logger.info(self.total)
|
|
self.app.logger.info(self.fetched)
|
|
self.search_result_from_cr.update(cr_result)
|
|
self.total = cr_result["message"]["total-results"]
|
|
self.app.logger.info("CROSSREF DONE")
|
|
|
|
self.app.logger.info("CROSSREF DONE")
|
|
with open(lib_paths.CR_RESULTS, "r+", encoding="utf-8") as data_file:
|
|
# First we load existing data into a dict.
|
|
try:
|
|
file_data = json.load(data_file)
|
|
except JSONDecodeError:
|
|
file_data = {}
|
|
data_file.truncate(0)
|
|
data_file.seek(0)
|
|
tmp = {self.uuid : self.search_result_from_cr}
|
|
if file_data:
|
|
file_data.update(tmp)
|
|
else:
|
|
file_data = tmp
|
|
json.dump(file_data, data_file, indent=4)
|
|
return cr_result
|
|
|
|
|
|
async def refine_search(self, unrefined_result):
|
|
"""
|
|
Refines the search query based on the unrefined search result.
|
|
|
|
Args:
|
|
- unrefined_result: The unrefined search result.
|
|
|
|
Returns:
|
|
- refined_result: The refined search result.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
summarized_results = []
|
|
self.app.logger.info("REFINE: Removing all derived works from the list")
|
|
for item in unrefined_result["message"]["items"]:
|
|
summarized_results.append(
|
|
{
|
|
"DOI": item["DOI"],
|
|
"title": item["title"] if "title" in item else None,
|
|
"type": item["type"] if "type" in item else None,
|
|
}
|
|
)
|
|
|
|
partial_result = {
|
|
self.uuid: {
|
|
"total_results": self.search_result_from_cr["message"][
|
|
"total-results"
|
|
],
|
|
"on_page": 1,
|
|
"summary": summarized_results,
|
|
"results": self.search_result_from_cr["message"]["items"],
|
|
}
|
|
}
|
|
self.app.logger.info("REFINE: Dumping to file")
|
|
temp = []
|
|
refined_result = {}
|
|
for key in partial_result:
|
|
self.app.logger.info("KEY:")
|
|
self.app.logger.info(key)
|
|
|
|
for item in partial_result[self.uuid]["summary"]:
|
|
if item["title"]:
|
|
temp.append(item)
|
|
|
|
for item in temp:
|
|
refined_result[item["DOI"]]= item
|
|
with open(lib_paths.RR_RESULTS, "r+", encoding="utf-8") as data_file:
|
|
# First we load existing data into a dict.
|
|
try:
|
|
file_data = json.load(data_file)
|
|
except JSONDecodeError:
|
|
file_data = {}
|
|
data_file.truncate(0)
|
|
data_file.seek(0)
|
|
tmp = {self.uuid: refined_result}
|
|
if file_data:
|
|
file_data.update(tmp)
|
|
else:
|
|
file_data = tmp
|
|
json.dump(file_data, data_file, indent=4)
|
|
return refined_result
|
|
|
|
async def check_if_exists(self, refined_result):
|
|
"""
|
|
Checks if the given DOI exists.
|
|
|
|
Args:
|
|
- doi: The DOI to be checked.
|
|
- brute_force: A flag indicating if brute force method should be used.
|
|
|
|
Returns:
|
|
- result: The search result.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
result = {}
|
|
self.app.logger.info("REFINE: Running search in the backend app")
|
|
dois = []
|
|
for item, value in refined_result.items():
|
|
dois.append([item, value])
|
|
coro = asyncio.to_thread(
|
|
search_bot.search_for_doi, dois, self.live_results, self.app.logger
|
|
)
|
|
result = await coro
|
|
result_list = []
|
|
result_no_db = []
|
|
for item in result:
|
|
if item["exists"]:
|
|
result_list.append(item)
|
|
else:
|
|
result_no_db.append(item)
|
|
self.hit = len(result)
|
|
return result_list, result_no_db
|
|
|
|
|
|
async def answer_query(self, deep_search=False):
|
|
"""
|
|
Answers the search query.
|
|
|
|
Args:
|
|
- None.
|
|
|
|
Returns:
|
|
- result: The search result.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
self.app.logger.info(f"Search started {self.uuid}")
|
|
cr_result = await self.search_crossref(query=self.query, deep_search=deep_search)
|
|
refined_result = await self.refine_search(cr_result)
|
|
answer, negative_answer = await self.check_if_exists(refined_result)
|
|
|
|
self.app.logger.info("Returning result")
|
|
self.app.logger.info(answer)
|
|
self.app.logger.info(negative_answer)
|
|
|
|
for item in answer:
|
|
self.final_result[item["DOI"]] = {"Title": item["data"]["title"], "type": item["data"]["type"]}
|
|
for item in negative_answer:
|
|
self.not_in_db[item["DOI"]] = {"Title": item["data"]["title"], "type": item["data"]["type"]}
|
|
self.app.logger.info("Returning result case2")
|
|
self.app.logger.info(self.final_result)
|
|
return self.final_result
|
|
|
|
# ============================= FLASK INTERNALS===============================
|
|
|
|
|
|
def flask_debug():
|
|
"""
|
|
Starts a Flask application in debug mode without using the reloader.
|
|
|
|
Args:
|
|
- None.
|
|
|
|
Returns:
|
|
- None.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
# trunk-ignore(bandit/B201)
|
|
app.run(debug=True, use_reloader=False, host=HOST_ADDRESS, port=PORT_ADDRESS)
|
|
|
|
|
|
def waitress_run():
|
|
"""
|
|
Serves the Flask application using the Waitress WSGI server.
|
|
|
|
Args:
|
|
- None.
|
|
|
|
Returns:
|
|
- None.
|
|
|
|
Raises:
|
|
- None.
|
|
"""
|
|
serve(app, host=HOST_ADDRESS, port=PORT_ADDRESS)
|
|
|
|
|
|
class BackgroundTaskSearch(threading.Thread):
|
|
"""
|
|
A background task for searching and saving results to files.
|
|
|
|
This class extends the `threading.Thread` class and is responsible for running
|
|
the search task in the background. It retrieves queries from a queue, performs
|
|
the search, and saves the results to files.
|
|
|
|
Attributes:
|
|
app (App): The application instance.
|
|
"""
|
|
|
|
def run(self):
|
|
"""
|
|
Run the background task.
|
|
|
|
This method is called when the thread is started. It creates a new event loop,
|
|
runs the `_run` method, and closes the event loop.
|
|
"""
|
|
loop = asyncio.new_event_loop()
|
|
loop.run_until_complete(self._run())
|
|
loop.close()
|
|
|
|
async def _run(self):
|
|
"""
|
|
Perform the search task.
|
|
|
|
This method is an asynchronous coroutine that runs in a loop. It retrieves a
|
|
librarian from the queue, answers the query, and saves the results to files.
|
|
It also sends the results to a remote server.
|
|
|
|
The search task continues running indefinitely until the thread is stopped.
|
|
"""
|
|
while True:
|
|
database = None
|
|
ndb_database = None
|
|
item = librarian_queue.get()
|
|
# Health-check ping: it has flowed through the internal queue and is
|
|
# now pulled off it - that is the whole point. Pong it straight back
|
|
# with the same uuid and DO NOT run a search.
|
|
if isinstance(item, dict) and "__ping__" in item:
|
|
ping_uuid = item["__ping__"]
|
|
self.app.logger.info(
|
|
"PING %s pulled off internal queue - ponging back (no search)",
|
|
ping_uuid,
|
|
)
|
|
await asyncio.to_thread(_post_pong, self.app.logger, ping_uuid)
|
|
continue
|
|
librarian = item
|
|
# Mark busy + processing for the whole search, and ALWAYS clear both
|
|
# (even on a crash) in finally: worker_busy so a ping doesn't wait
|
|
# behind us, and active_queries so the bot's watchdog can tell a
|
|
# finished-and-gone query from one still in flight.
|
|
worker_busy.set()
|
|
with _active_lock:
|
|
active_queries[str(librarian.uuid)] = "processing"
|
|
try:
|
|
self.app.logger.info("STARTED")
|
|
result = await librarian.answer_query(librarian.deep_search)
|
|
result = {librarian.uuid: result}
|
|
self.app.logger.info("Saving to file")
|
|
|
|
# Save results to "not_in_db.json" file
|
|
with open(lib_paths.NOT_IN_DB, "r+", encoding="utf-8") as ndb_file:
|
|
ndb_database = {}
|
|
try:
|
|
ndb_database = json.load(ndb_file)
|
|
except JSONDecodeError:
|
|
pass
|
|
if ndb_database:
|
|
ndb_database.update(librarian.not_in_db)
|
|
else:
|
|
ndb_database = librarian.not_in_db
|
|
ndb_file.truncate(0)
|
|
ndb_file.seek(0)
|
|
json.dump(ndb_database, ndb_file)
|
|
|
|
# Save results to "s_results.json" file
|
|
with open(lib_paths.S_RESULTS, "r+", encoding="utf-8") as s_file:
|
|
database = {}
|
|
try:
|
|
database = json.load(s_file)
|
|
except JSONDecodeError:
|
|
pass
|
|
if database:
|
|
self.app.logger.info(database)
|
|
self.app.logger.info(result)
|
|
database.update(result)
|
|
else:
|
|
database = result
|
|
self.app.logger.info("DUMPING DATA")
|
|
s_file.truncate(0)
|
|
s_file.seek(0)
|
|
json.dump(database, s_file)
|
|
self.app.logger.info("FINISHED")
|
|
|
|
# Persist the result to the durable OUTBOX FIRST, then try to
|
|
# deliver it. Writing to disk before sending is the whole point:
|
|
# an expensive (hours-long) result now survives a failed send, a
|
|
# bot outage, or a librarian restart - the resender keeps
|
|
# retrying until the bot ACKs, and only then is it removed.
|
|
payload = result # shape: {uuid: {DOI: {"Title": ..., "type": ...}}}
|
|
uuid = str(librarian.uuid)
|
|
hits = payload.get(librarian.uuid, {}) if isinstance(payload, dict) else {}
|
|
_outbox.put(uuid, payload)
|
|
self.app.logger.info(
|
|
"SENDING result for %s: %d DOI(s): %s (queued to OUTBOX)",
|
|
uuid, len(hits), list(hits.keys()),
|
|
)
|
|
if await asyncio.to_thread(_deliver_result, uuid, payload, self.app.logger):
|
|
_outbox.remove(uuid)
|
|
else:
|
|
self.app.logger.warning(
|
|
"Result %s not acked yet - left in OUTBOX for the resender", uuid
|
|
)
|
|
except Exception as exc: # pylint: disable=broad-exception-caught
|
|
# A crashing search must not kill the worker thread (which would
|
|
# freeze the whole queue). Log and move on; finally still clears
|
|
# busy/active so the query is correctly seen as "gone".
|
|
self.app.logger.exception("Search %s crashed: %s", librarian.uuid, exc)
|
|
finally:
|
|
worker_busy.clear()
|
|
with _active_lock:
|
|
active_queries.pop(str(librarian.uuid), None)
|
|
await asyncio.sleep(1)
|
|
|
|
|
|
# ==================================SERVER ROUTES==========================================
|
|
@app.route("/query", methods=["POST"])
|
|
async def query_database():
|
|
_authorize_request()
|
|
"""
|
|
Endpoint for querying the database.
|
|
|
|
This function receives a POST request containing a JSON payload with a query and a UUID.
|
|
It creates a Librarian object with the query and UUID,
|
|
and adds it to the librarian_queue and librarian_list.
|
|
Finally, it returns a JSON response indicating the success
|
|
of the operation, along with the query, UUID,
|
|
and the current size of the librarian_queue.
|
|
|
|
Returns:
|
|
tuple: A tuple containing a JSON response and a status code.
|
|
"""
|
|
record = json.loads(request.data)
|
|
app.logger.info(record)
|
|
app.logger.info(record["query"])
|
|
app.logger.info(record["UUID"])
|
|
uuid = record["UUID"]
|
|
deep_search = record["deep_search"]
|
|
cl = Librarian(app, record["query"], uuid, deep_search)
|
|
librarian_queue.put(cl)
|
|
librarian_list.append(cl)
|
|
# The bot's per-query watchdog polls /query_status for this uuid; mark it
|
|
# "queued" now so it counts as known the moment we accept it.
|
|
with _active_lock:
|
|
active_queries[str(uuid)] = "queued"
|
|
answer_data = (record["query"], record["UUID"], librarian_queue.qsize())
|
|
return_data = (
|
|
jsonify(isError=False, message="Success", statusCode=200, data=answer_data),
|
|
200,
|
|
)
|
|
return return_data
|
|
|
|
|
|
@app.route("/ping", methods=["POST"])
|
|
def ping_roundtrip():
|
|
_authorize_request()
|
|
"""
|
|
Health-check round-trip.
|
|
|
|
Two cases, one guarantee - the pong always comes back over the librarian->bot
|
|
return path (the only thing the ping must prove):
|
|
|
|
* IDLE: put a ping marker onto the SAME internal ``librarian_queue`` real
|
|
searches use and return 200. The worker pulls it off and pongs it back,
|
|
so a successful pong proves the whole pipeline flows (queue + worker + the
|
|
return leg), not just that Flask is up.
|
|
* BUSY (a search is grinding): DO NOT queue - the ping would just wait behind
|
|
a possibly hours-long search and time out, making a perfectly healthy busy
|
|
librarian look dead. Pong back immediately instead. Being busy is fine; you
|
|
can keep piling searches on. The ping only needs to catch a BROKEN return
|
|
path, and the direct pong exercises exactly that.
|
|
"""
|
|
record = json.loads(request.data)
|
|
ping_uuid = record["UUID"]
|
|
if worker_busy.is_set():
|
|
app.logger.info("PING %s while busy grinding - direct pong (skip queue)", ping_uuid)
|
|
_post_pong(app.logger, ping_uuid)
|
|
else:
|
|
app.logger.info("PING received %s - queued for round-trip", ping_uuid)
|
|
librarian_queue.put({"__ping__": ping_uuid})
|
|
return (
|
|
jsonify(isError=False, message="ping-queued", statusCode=200, data=ping_uuid),
|
|
200,
|
|
)
|
|
|
|
|
|
@app.route("/query_status", methods=["POST"])
|
|
def query_status():
|
|
_authorize_request()
|
|
"""
|
|
Per-query watchdog probe.
|
|
|
|
Returns whether ``UUID`` is still known to the librarian (queued or being
|
|
processed). The bot polls this after dispatching a search: while the uuid is
|
|
known the search is progressing; once it VANISHES here without the result
|
|
ever reaching the bot, the result was lost in transit and the bot tells the
|
|
user. A busy/queued search is never mistaken for a lost one.
|
|
"""
|
|
record = json.loads(request.data)
|
|
uuid = str(record["UUID"])
|
|
with _active_lock:
|
|
state = active_queries.get(uuid, "unknown")
|
|
return (
|
|
jsonify(
|
|
isError=False,
|
|
message="Success",
|
|
statusCode=200,
|
|
data={"uuid": uuid, "known": state != "unknown", "state": state},
|
|
),
|
|
200,
|
|
)
|
|
|
|
|
|
@app.route("/get_partial_result", methods=["POST"])
|
|
async def get_partial():
|
|
_authorize_request()
|
|
"""
|
|
Retrieves the partial result for a given UUID.
|
|
|
|
Returns:
|
|
A JSON response containing the partial result.
|
|
"""
|
|
record = json.loads(request.data)
|
|
app.logger.info(record)
|
|
app.logger.info(record["UUID"])
|
|
for lib in librarian_list:
|
|
if lib.uuid == record["UUID"]:
|
|
answer_data = lib.live_results
|
|
break
|
|
return_data = (
|
|
jsonify(isError=False, message="Success", statusCode=200, data=answer_data),
|
|
200,
|
|
)
|
|
return return_data
|
|
|
|
|
|
# =======================================MAIN===================================================
|
|
if __name__ == "__main__":
|
|
app.logger.setLevel(logging.DEBUG)
|
|
LOGFILE_PATH.parent.mkdir(parents=True, exist_ok=True)
|
|
h1 = handlers.RotatingFileHandler(
|
|
filename=str(LOGFILE_PATH),
|
|
encoding=ENCODING,
|
|
mode="a",
|
|
maxBytes=6 * 1024 * 1024,
|
|
backupCount=6,
|
|
)
|
|
|
|
app.logger.addHandler(h1)
|
|
threads = []
|
|
threads.append(threading.Thread(target=waitress_run, daemon=True))
|
|
# threads.append(threading.Thread(target=flask_debug))
|
|
bgtask = BackgroundTaskSearch()
|
|
bgtask.app = app
|
|
bgtask.daemon = True
|
|
threads.append(bgtask)
|
|
threads.append(
|
|
threading.Thread(
|
|
target=scrape_bot.scraper, args=(app.logger,), daemon=True
|
|
)
|
|
)
|
|
# Durable delivery: keep flushing the OUTBOX so any result not yet acked by
|
|
# the bot (transient outage, or left over from before a restart) is resent.
|
|
threads.append(
|
|
threading.Thread(target=outbox_resender, args=(app.logger,), daemon=True)
|
|
)
|
|
i = 0
|
|
try:
|
|
for worker in threads:
|
|
app.logger.info("App number: %s", i)
|
|
i += 1
|
|
worker.start()
|
|
for worker in threads:
|
|
worker.join()
|
|
except KeyboardInterrupt:
|
|
app.logger.info("Shutdown requested - exiting librarian service")
|