c5643aa28f
CI / compile (pull_request) Successful in 9s
CI / unit (pull_request) Successful in 23s
CI / integration (pull_request) Successful in 27s
build / build (push) Successful in 33s
CI / compile (push) Successful in 13s
CI / unit (push) Successful in 33s
CI / integration (push) Successful in 27s
Two hygiene fixes on top of the work-queue OOM bound: Result dumps: cr_results / rr_results / s_results.json were write-only (nothing reads them) yet accumulated EVERY search forever and json.load'd the whole growing file on each write - unbounded RAM and PVC growth, and for a deep search the raw cr_results dump is hundreds of MB. They are now off by default (CONJURER_LIBRARIAN_DEBUG_DUMPS) and, when enabled, are overwritten with just the latest search - never loaded or accumulated. not_in_db.json is untouched: it's a real queue the scraper drains. Search logging: search_bot logged via print(), including a per-line carriage-return progress line that flooded stdout / the log file with millions of entries - fine for a desktop app, unreadable and bloating in a container. All of it is now proper logging at DEBUG (with coarse per-500k-line progress), so a normal run is quiet. The librarian log level is configurable (CONJURER_LIBRARIAN_LOG_LEVEL, default INFO) and a stdout handler is added so stays useful now that the search no longer prints straight to stdout. Set DEBUG for full verbosity. Also: make test_result_delivery_contract hermetic (point the durable spool at a temp dir so it can't pollute or be poisoned by the real result_inbox/ between runs) and gitignore the runtime spool dirs. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
773 lines
29 KiB
Python
773 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)
|
|
|
|
# cr_results/rr_results/s_results.json are write-only debug dumps (nothing reads
|
|
# them). They used to accumulate EVERY search forever AND json.load the whole
|
|
# growing file on each write - unbounded RAM + disk, and for a deep search the
|
|
# raw dump is hundreds of MB. Off by default now; when explicitly enabled they
|
|
# are overwritten with just the latest search (no load, no accumulation).
|
|
DEBUG_DUMPS = _env("CONJURER_LIBRARIAN_DEBUG_DUMPS", "0").lower() in ("1", "true", "yes")
|
|
# Log level: INFO keeps normal runs readable (the desktop-era per-line/per-file
|
|
# chatter is now DEBUG); set DEBUG to get the full verbosity back.
|
|
LOG_LEVEL = _env("CONJURER_LIBRARIAN_LOG_LEVEL", "INFO").upper()
|
|
|
|
|
|
def _dump_debug(path, uuid, data) -> None:
|
|
"""Optionally dump the latest search's data for debugging.
|
|
|
|
Overwrites (never accumulates) and does nothing unless DEBUG_DUMPS is on, so
|
|
it can't grow RAM or the state volume in normal operation."""
|
|
if not DEBUG_DUMPS:
|
|
return
|
|
try:
|
|
with open(path, "w", encoding="utf-8") as handle:
|
|
json.dump({uuid: data}, handle)
|
|
except OSError as exc:
|
|
logging.getLogger("conjurer_librarian").warning("Debug dump to %s failed: %s", path, exc)
|
|
|
|
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")
|
|
_dump_debug(lib_paths.CR_RESULTS, self.uuid, self.search_result_from_cr)
|
|
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
|
|
_dump_debug(lib_paths.RR_RESULTS, self.uuid, refined_result)
|
|
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)
|
|
|
|
# Optional debug dump of the final result (off by default).
|
|
_dump_debug(lib_paths.S_RESULTS, librarian.uuid, result[librarian.uuid])
|
|
self.app.logger.info("Search %s finished", librarian.uuid)
|
|
|
|
# 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__":
|
|
# Default INFO (readable). Set CONJURER_LIBRARIAN_LOG_LEVEL=DEBUG for the
|
|
# full per-file / per-line search chatter.
|
|
app.logger.setLevel(LOG_LEVEL)
|
|
_fmt = logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s")
|
|
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,
|
|
)
|
|
h1.setFormatter(_fmt)
|
|
app.logger.addHandler(h1)
|
|
# Console handler so `kubectl logs` shows what's happening (k8s reads stdout);
|
|
# the search internals no longer print() straight to stdout.
|
|
_console = logging.StreamHandler()
|
|
_console.setFormatter(_fmt)
|
|
app.logger.addHandler(_console)
|
|
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")
|