librarian: napraw nieskończoną pętlę w wyszukiwaniu DOI #1

Merged
gitea merged 1 commits from fix-librarian-search-hang into main 2026-07-30 09:26:45 +00:00
5 changed files with 204 additions and 42 deletions
+5 -1
View File
@@ -9,6 +9,10 @@ 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"),
):
if _path not in sys.path: if _path not in sys.path:
sys.path.insert(0, _path) sys.path.insert(0, _path)
+99 -37
View File
@@ -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))
) )
+6 -2
View File
@@ -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.
+5 -2
View File
@@ -186,8 +186,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
+89
View File
@@ -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"]