librarian: stop the DOI search from hanging on a chunk-count mismatch
CI / compile (pull_request) Successful in 1m26s
CI / unit (pull_request) Successful in 1m8s
CI / integration (pull_request) Failing after 10h21m8s
CI / compile (push) Successful in 12m43s
build / build (push) Failing after 13m30s
CI / unit (push) Successful in 2m32s
CI / integration (push) Failing after 1h54m10s
CI / compile (pull_request) Successful in 1m26s
CI / unit (pull_request) Successful in 1m8s
CI / integration (pull_request) Failing after 10h21m8s
CI / compile (push) Successful in 12m43s
build / build (push) Failing after 13m30s
CI / unit (push) Successful in 2m32s
CI / integration (push) Failing after 1h54m10s
search_bot conflated MAXTHREADS into two jobs at once - how many chunk files to read (files 0..MAXTHREADS-1) AND how many producer sentinels to wait for - so the two had to match exactly. Set too low it silently skipped trailing chunks; set too high (or with any chunk missing/unreadable) a producer crashed before emitting its sentinel, the consumers' count never reached the threshold, and search_for_doi hung on join() forever. The idle-timeout failsafe that was meant to break a starved consumer was dead code: `if empty_counter > 5: ... elif empty_counter > 10: break` - >10 implies >5, so the elif never ran. Fix, three layers: * auto-discover the chunk files present (discover_chunk_files: <n>_chunk.txt in numeric order) instead of range(0, MAXTHREADS). All files are read regardless of count, and no producer is ever pointed at a missing file; * the sentinel threshold is now the number of producers actually started, so it can't drift from what's emitted; * producers emit their sentinel in a finally, so even a crash (missing/unreadable chunk) can't starve the count; and the idle backstop is reordered so it can actually fire (>EMPTY_LIMIT seconds) as a last resort. MAXTHREADS is deprecated and unused (kept only so old env files don't break); docs/env updated to say chunk files are auto-discovered. For the reported case (MAXTHREADS=40, files 0..43): before, files 40-43 were silently never searched, and any run that referenced a missing chunk hung forever. After, all 44 are searched and it always terminates. Verified in a pytest-only venv (tests/unit/test_search_bot.py): DOI in a trailing chunk is found; an unreadable chunk still terminates; empty dir returns at once; discovery is numeric-sorted. Full unit job 27 passed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit was merged in pull request #1.
This commit is contained in:
@@ -20,6 +20,7 @@ Global Variables:
|
||||
|
||||
# TODO: Wpiemdolić to wszystko w klasę z loggerem przysłanym z góry
|
||||
import os
|
||||
import re
|
||||
from queue import Empty, Queue
|
||||
from threading import Thread
|
||||
import time
|
||||
@@ -28,13 +29,54 @@ q = Queue()
|
||||
|
||||
# Deployment data is environment-overridable so the local DOI database can live
|
||||
# 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\\")
|
||||
|
||||
ENCODING = os.getenv("CONJURER_ENCODING", "utf-8")
|
||||
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()
|
||||
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):
|
||||
@@ -47,36 +89,45 @@ def producer(out_q, control_q, filename, _logger):
|
||||
filename (str): Name of the file.
|
||||
_logger: Logger object for logging.
|
||||
"""
|
||||
with open(DATABASE_PATH + filename, "r", encoding=ENCODING) as operated_file:
|
||||
print(f"Worker {filename} ")
|
||||
line_no = 0
|
||||
while True:
|
||||
line = operated_file.readline()
|
||||
line_no += 1
|
||||
print(f"\t \t \t \t \t \t W{filename}{line_no}\r", end="")
|
||||
try:
|
||||
with open(DATABASE_PATH + filename, "r", encoding=ENCODING) as operated_file:
|
||||
print(f"Worker {filename} ")
|
||||
line_no = 0
|
||||
while True:
|
||||
line = operated_file.readline()
|
||||
line_no += 1
|
||||
print(f"\t \t \t \t \t \t W{filename}{line_no}\r", end="")
|
||||
|
||||
if not line:
|
||||
print(f"EOF {filename}")
|
||||
break
|
||||
if not line:
|
||||
print(f"EOF {filename}")
|
||||
break
|
||||
|
||||
if not line:
|
||||
break
|
||||
out_q.put(line)
|
||||
try:
|
||||
check = control_q.get(block=False)
|
||||
except Empty:
|
||||
check = False
|
||||
out_q.put(line)
|
||||
try:
|
||||
check = control_q.get(block=False)
|
||||
except Empty:
|
||||
check = False
|
||||
|
||||
if check is _sentinel:
|
||||
print("TERM signal received")
|
||||
control_q.put(check)
|
||||
break
|
||||
print(f"Worker finished: {filename}")
|
||||
if check is _sentinel:
|
||||
print("TERM signal received")
|
||||
control_q.put(check)
|
||||
break
|
||||
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)
|
||||
|
||||
|
||||
# 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.
|
||||
|
||||
@@ -117,14 +168,18 @@ def consumer(in_q, control_q, doi, live_results, result_list, control_dict, no,
|
||||
empty_counter += 1
|
||||
time.sleep(1)
|
||||
print(f"Consumer {no} empty")
|
||||
if empty_counter > 5:
|
||||
print("Consumer %s empty lvl 2", no)
|
||||
time.sleep(2)
|
||||
elif empty_counter > 10:
|
||||
print(f"Consumer thread finished {no}")
|
||||
# Order matters: the >EMPTY_LIMIT break must be checked BEFORE the
|
||||
# lesser threshold, otherwise (as in the original) the first branch
|
||||
# always wins and the break is dead code, leaving the sentinel count
|
||||
# as the only exit - which is exactly what used to hang the search.
|
||||
if empty_counter > EMPTY_LIMIT:
|
||||
print(f"Consumer thread finished {no} (idle backstop)")
|
||||
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}")
|
||||
break
|
||||
|
||||
@@ -148,18 +203,25 @@ def search_for_doi(doi, live_results, _logger):
|
||||
for item in doi:
|
||||
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):
|
||||
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")
|
||||
threads.append(t_cons)
|
||||
for i in range(0, MAXTHREADS):
|
||||
# TEST DATA
|
||||
# filename = "test" + str(i) + CHUNk
|
||||
# DEPLOYMENT DATA
|
||||
filename = str(i) + CHUNK
|
||||
_logger.info(f"Creating worker thread no: {i}")
|
||||
for filename in chunk_files:
|
||||
_logger.info("Creating worker thread for %s", filename)
|
||||
threads.append(
|
||||
Thread(target=producer, args=(work_q, control_q, filename, _logger))
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user