Compare commits

..

3 Commits

Author SHA1 Message Date
gitea ed8b271b4e Gate librarian cog on a full ping round-trip, not a bare GET
CI / compile (pull_request) Successful in 9s
CI / unit (pull_request) Successful in 18s
CI / integration (pull_request) Successful in 18s
The librarian health check was a plain GET to '/', which only proved
Flask was listening - not that the service could actually take a query,
run it through its internal queue+worker, and answer back. So the cog
could load against a librarian whose worker was wedged or that couldn't
reach the bot on the return leg.

Replace it with a ping that travels the SAME path a real search does, on
both sides:
  bot: QueryControl -> OUT_COMM_Q -> scan_queue -> awaiting_q
  librarian: POST /ping -> librarian_queue -> worker pulls it off
             (no Crossref/DOI search) -> pongs back with the same uuid
  bot: /conjurer -> incoming_q -> scan_incoming matches uuid, wakes waiter
The cog enables only when that whole loop closes within 3s. This also
proves the librarian->bot return path, which a GET never did.

Safety: uuid is random per ping; the wait and POST are both bounded so
startup can't stall; a pong that finds no waiter is dropped (never
orphaned into IN_COMM_Q, which would make the cog post a bogus 'no
results' message); and a ping whose pong never returns is swept out of
awaiting_q after PING_TTL_SECONDS so nothing leaks. All awaiting_q writes
stay within scan_queue (append) and scan_incoming (remove) - no locks,
no cross-thread mutation.

Integration tests cover: OK round-trip, timeout when accepted-but-no-pong,
unreachable, non-200, orphan-pong-dropped, and that real results still
reach IN_COMM_Q. Suite: 24 integration + 41 unit green.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-01 15:00:04 +02:00
gitea 442b8a2a60 Make service health-check self-diagnosing
build / build (push) Failing after 9s
CI / compile (push) Successful in 9s
CI / unit (push) Successful in 19s
CI / integration (push) Successful in 14s
When a service group's cog stays disabled the log now says WHY:
_service_health returns the actual connection error (ConnectionError
= refused/down, ConnectTimeout = firewall/slow, gaierror = DNS/wrong
host) instead of a bare 'unreachable', and on_ready logs the resolved
FILE/LIBRARIAN/RADIO addresses so an env-var that never reached the
process (address falls back to the 192.168.1.15:5000 default) is
obvious at a glance.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-01 11:52:04 +00:00
gitea defc482a22 fix: batch A - crash bugs, a leak, and startup/edge fragility
CI / compile (pull_request) Successful in 10s
CI / unit (pull_request) Successful in 19s
CI / integration (pull_request) Successful in 10s
build / build (push) Failing after 46m40s
CI / compile (push) Successful in 20s
CI / unit (push) Successful in 19s
CI / integration (push) Successful in 16s
Seven confirmed defects from the code audit, each small and low-risk.

* ai_functions.get_random_cyclic_message: random.randint(0, len(CYCLIC_WORDS))
  is inclusive -> could return len -> IndexError. Now randrange(len) + guard on
  an empty CYCLIC_WORDS.
* librarian_commands.get_image_sadox: random.randrange(0, len(res)-1) never
  picked the last comic and raised ValueError('empty range') on a single file.
  Now randrange(len) + an empty-dir guard.
* ai_commands image generation: every DALL-E error branch replied but did not
  return, so control fell through to `if response:` with response unbound ->
  UnboundLocalError right after the friendly message. Each branch now returns;
  response is pre-initialised; and PermissionDeniedError no longer passes a
  (message, text) tuple as a single arg.
* search_bot DOI match: `item["DOI"] in data` was a substring test, so a DOI
  that is a prefix of a longer one (10.1/1 vs 10.1/12) produced a false 'exists'
  hit. Now matches the line's first whitespace token exactly, via an O(1) dict
  index built once per consumer (also removes the O(queried-DOIs) per-line scan
  - a real win for large databases).
* communication_subroutine.scan_incoming: matched records were never removed
  from awaiting_q, so it grew unbounded over uptime and a reused UUID could
  re-match a stale record. Matched records are now dropped after dispatch.
* communication_subroutine.id3: (resp.headers.get("icy-name") or "").title()
  guards against a stream that omits headers (was AttributeError on None,
  500-ing the /prepped_tracks "next" handler).
* betoniarka.scan_tracks: waits for the radio logs to exist instead of dying
  with FileNotFoundError on a fresh deploy (which silently killed the
  now-playing forwarder until a restart).

Verified: tests/unit/test_search_bot.py gains exact-match and trailing-metadata
cases; full unit job 43 passed. Remaining observations (image-gen stale
/home/pi fallback paths + dead FileNotFoundError-after-OSError branch; tailer
still vulnerable to mid-run log rotation; DOI-first-token assumption) noted for
follow-up - none are crashes on the normal path.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-31 23:01:27 +02:00
8 changed files with 125 additions and 36 deletions
+14 -3
View File
@@ -427,6 +427,11 @@ class Events(commands.Cog):
return return
message_content_lower = message_content_lower.replace("imaginuje sobie: ", "") message_content_lower = message_content_lower.replace("imaginuje sobie: ", "")
self.logger.debug("Wywolanie obrazka: %s", message_content_lower) self.logger.debug("Wywolanie obrazka: %s", message_content_lower)
# Every error branch below must RETURN: otherwise control falls
# through to `if response:` with `response` unbound (the call
# raised) -> UnboundLocalError, crashing the handler right after
# the friendly message was already sent.
response = None
try: try:
response = await OPENAICLIENT.images.generate( response = await OPENAICLIENT.images.generate(
model="dall-e-3", model="dall-e-3",
@@ -440,10 +445,12 @@ class Events(commands.Cog):
await discord_friendly_reply( await discord_friendly_reply(
message, f"*Kondziu patrzy na terminal, czeka, czeka, czeka,.... Jeszcze chwile czeka Przypierdala w niego pięścią....* Nie mogę się połączyć z Openai spróbuj od nowa. *Na ekranie pojawia się*: {e}" message, f"*Kondziu patrzy na terminal, czeka, czeka, czeka,.... Jeszcze chwile czeka Przypierdala w niego pięścią....* Nie mogę się połączyć z Openai spróbuj od nowa. *Na ekranie pojawia się*: {e}"
) )
return
except openai.APIConnectionError as e: except openai.APIConnectionError as e:
await discord_friendly_reply( await discord_friendly_reply(
message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Nie mogę się połączyć z Openai. *Na ekranie pojawia się*: {e}" message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Nie mogę się połączyć z Openai. *Na ekranie pojawia się*: {e}"
) )
return
except openai.BadRequestError as e: except openai.BadRequestError as e:
# Handle invalid request error, e.g. validate parameters or log # Handle invalid request error, e.g. validate parameters or log
if message.author.nick: if message.author.nick:
@@ -461,27 +468,31 @@ class Events(commands.Cog):
await discord_friendly_reply( await discord_friendly_reply(
message, f"Sorki, cenzura: {resp}. Jak chcesz to są kanały na nudle #sexy-foteczky i #kanal-do-fapania *Na ekranie pojawia się: {e}" message, f"Sorki, cenzura: {resp}. Jak chcesz to są kanały na nudle #sexy-foteczky i #kanal-do-fapania *Na ekranie pojawia się: {e}"
) )
return
except openai.AuthenticationError as e: except openai.AuthenticationError as e:
# Handle authentication error, e.g. check credentials or log # Handle authentication error, e.g. check credentials or log
await discord_friendly_reply( await discord_friendly_reply(
message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Wołaj szefa - coś się z hasłem zjebało. *Na terminalu pojawia się:* {e}" message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Wołaj szefa - coś się z hasłem zjebało. *Na terminalu pojawia się:* {e}"
) )
return
except openai.PermissionDeniedError as e: except openai.PermissionDeniedError as e:
# Handle permission error, e.g. check scope or log # Handle permission error, e.g. check scope or log
# (was accidentally passing a (message, text) TUPLE as one arg)
await discord_friendly_reply( await discord_friendly_reply(
( message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Wołaj szefa - coś się z uprawnieniami zjebało. *Na terminalu pojawia się:* {e}"
message, f"*Kondziu patrzy na terminal, chwile się zastanawia. Przypierdala w niego pięścią....* Wołaj szefa - coś się z uprawnieniami zjebało. *Na terminalu pojawia się:* {e}"
)
) )
return
except openai.RateLimitError as e: except openai.RateLimitError as e:
await discord_friendly_reply( await discord_friendly_reply(
message, f"*Kondziu patrzy na terminal* Wołaj szefa. Zapłacić rachunki za AI trzeba. Jak chcesz to się na #zebranie dorzuć. {e}" message, f"*Kondziu patrzy na terminal* Wołaj szefa. Zapłacić rachunki za AI trzeba. Jak chcesz to się na #zebranie dorzuć. {e}"
) )
return
except openai.APIError as e: except openai.APIError as e:
# Handle API error, e.g. retry or log # Handle API error, e.g. retry or log
await discord_friendly_reply( await discord_friendly_reply(
message, f"*Kondziu nurkuje za bar, terminal wybucha. Przed tobą ląduje pergamin zapisany pięknym gotykiem a na nim*: {e}" message, f"*Kondziu nurkuje za bar, terminal wybucha. Przed tobą ląduje pergamin zapisany pięknym gotykiem a na nim*: {e}"
) )
return
if response: if response:
self.logger.info(response) self.logger.info(response)
image_url = response.data[0].url image_url = response.data[0].url
+4 -2
View File
@@ -541,10 +541,12 @@ async def get_random_cyclic_message(client):
# trunk-ignore(bandit/B311) # trunk-ignore(bandit/B311)
ai_check = random.randint(0, 10) ai_check = random.randint(0, 10)
logger.info("Losowa wypowiedź") logger.info("Losowa wypowiedź")
if ai_check < 2: if ai_check < 2 and CYCLIC_WORDS:
logger.info("Predefiniowana") logger.info("Predefiniowana")
# randrange(n) is 0..n-1; randint(0, n) was inclusive and could return n
# -> list(...)[n] IndexError. Guarded on empty CYCLIC_WORDS above.
# trunk-ignore(bandit/B311) # trunk-ignore(bandit/B311)
messnum = random.randint(0, len(CYCLIC_WORDS)) messnum = random.randrange(len(CYCLIC_WORDS))
logger.debug(messnum) logger.debug(messnum)
logger.debug(len(CYCLIC_WORDS)) logger.debug(len(CYCLIC_WORDS))
mess_key = list(CYCLIC_WORDS.keys())[messnum] mess_key = list(CYCLIC_WORDS.keys())[messnum]
+35 -14
View File
@@ -128,13 +128,15 @@ SERVICE_EXTENSION_GROUPS = {
SERVICE_RECHECK_SECONDS = 300 SERVICE_RECHECK_SECONDS = 300
def _service_alive(url: str) -> bool: def _service_health(url: str):
"""True when the service answers HTTP at all (any status code counts).""" """Return None when the service answers HTTP at all (any status counts),
otherwise the connection error explaining WHY it's unreachable (refused vs
timeout vs DNS - the difference points straight at the cause)."""
try: try:
requests.get(url, timeout=3) requests.get(url, timeout=3)
return True return None
except requests.exceptions.RequestException: except requests.exceptions.RequestException as exc:
return False return exc
async def _load_extension_safe(name: str) -> bool: async def _load_extension_safe(name: str) -> bool:
@@ -172,16 +174,25 @@ async def _load_service_groups() -> bool:
service_headers(), service_headers(),
LIBRARIAN_PING_TIMEOUT, LIBRARIAN_PING_TIMEOUT,
) )
if not alive:
logger.warning(
"Service 'librarian' ping round-trip failed (%s) - cogs stay disabled: %s",
group["health_url"],
", ".join(missing),
)
continue
else: else:
alive = await asyncio.to_thread(_service_alive, group["health_url"]) err = await asyncio.to_thread(_service_health, group["health_url"])
if not alive: if err is not None:
logger.warning( logger.warning(
"Service '%s' unreachable (%s) - cogs stay disabled: %s", "Service '%s' unreachable (%s) [%s: %s] - cogs stay disabled: %s",
service, service,
group["health_url"], group["health_url"],
", ".join(missing), type(err).__name__,
) err,
continue ", ".join(missing),
)
continue
logger.info("Service '%s' is alive - enabling: %s", service, ", ".join(missing)) logger.info("Service '%s' is alive - enabling: %s", service, ", ".join(missing))
for extension in missing: for extension in missing:
if await _load_extension_safe(extension): if await _load_extension_safe(extension):
@@ -225,6 +236,16 @@ async def on_ready():
for extension in CORE_EXTENSIONS: for extension in CORE_EXTENSIONS:
await _load_extension_safe(extension) await _load_extension_safe(extension)
# Log the ACTUALLY-resolved service addresses. When one shows the built-in
# default (192.168.1.15:5000) it means the matching CONJURER_* env var never
# reached the process - the single most common cause of "service unreachable"
# confusion. Printing them makes env-vs-default obvious at a glance.
logger.info(
"Resolved service addresses -> musician(file): %s | librarian: %s | radio: %s",
FILE_SERVICE_ADDRESS,
LIBRARIAN_SERVICE_ADDRESS,
RADIO_SERVICE_ADDRESS,
)
await _load_service_groups() await _load_service_groups()
logger.info("Sensors: online") logger.info("Sensors: online")
+14 -5
View File
@@ -269,14 +269,20 @@ def scan_incoming(stop_event: Optional[threading.Event] = None):
awaiting_q.remove(record) awaiting_q.remove(record)
logger.info("PONG matched for %s", pong_uuid) logger.info("PONG matched for %s", pong_uuid)
continue continue
record_stored = False # Collect matched records and drop them from awaiting_q afterwards -
# they used to stay forever (awaiting_q only ever grew), leaking
# memory over the bot's uptime and letting a reused UUID re-match a
# stale record.
matched = []
for record in awaiting_q: for record in awaiting_q:
if record.uuid in answer.keys(): if record.uuid in answer.keys():
record_stored = True
record.stop = True record.stop = True
record.entries = answer[record.uuid] record.entries = answer[record.uuid]
IN_COMM_Q.put(record) IN_COMM_Q.put(record)
if not record_stored: matched.append(record)
for record in matched:
awaiting_q.remove(record)
if not matched:
for key in answer.keys(): for key in answer.keys():
record = QueryControl("Orphaned", key, "Orphan", None) record = QueryControl("Orphaned", key, "Orphan", None)
record.stop = True record.stop = True
@@ -353,10 +359,13 @@ def id3(url: str) -> dict:
resp.read( resp.read(
metaint metaint
) # this isn't seekable so, arbitrarily read to the point we want ) # this isn't seekable so, arbitrarily read to the point we want
# Guard the headers: an Icecast stream that omits icy-name / icy-genre
# (e.g. while the radio is down) made `.title()` raise AttributeError on
# None, 500-ing the /prepped_tracks "next" handler that calls this.
tagdata = dict( tagdata = dict(
site_url=resp.headers.get("icy-url"), site_url=resp.headers.get("icy-url"),
name=resp.headers.get("icy-name").title(), name=(resp.headers.get("icy-name") or "").title(),
genre=resp.headers.get("icy-genre").title(), genre=(resp.headers.get("icy-genre") or "").title(),
title=get_stream_title(resp.read(255)), title=get_stream_title(resp.read(255)),
) )
return tagdata return tagdata
+6
View File
@@ -203,6 +203,12 @@ def wyszukaj(word_list, how_many, _logger=None, write_to=None):
# ---------------------------------------------------------------- tailer # ---------------------------------------------------------------- tailer
def scan_tracks(): def scan_tracks():
"""Tail the radio logs and forward play events to the bot.""" """Tail the radio logs and forward play events to the bot."""
# On a fresh deploy Liquidsoap may not have written its logs yet; wait for
# them instead of dying with FileNotFoundError, which used to silently kill
# the now-playing forwarder until the container was restarted.
while not (RADIOLOG_PATH.exists() and PERSISTENCE_PATH.exists()):
logger.info("Waiting for radio logs (%s, %s)...", RADIOLOG_PATH, PERSISTENCE_PATH)
time.sleep(5)
with open(RADIOLOG_PATH, "r", encoding=ENCODING) as log_file: with open(RADIOLOG_PATH, "r", encoding=ENCODING) as log_file:
log_file.seek(os.stat(RADIOLOG_PATH).st_size) log_file.seek(os.stat(RADIOLOG_PATH).st_size)
prev_size = os.stat(PERSISTENCE_PATH).st_size prev_size = os.stat(PERSISTENCE_PATH).st_size
+19 -11
View File
@@ -148,8 +148,11 @@ def consumer(in_q, control_q, doi, live_results, result_list, control_dict, expe
print(f"Consumer thread started: {no} no") print(f"Consumer thread started: {no} no")
empty_counter = 0 empty_counter = 0
alive_no = 0 alive_no = 0
# DOI -> result item, so a line is matched with one O(1) dict lookup instead
# of scanning every queried DOI. Items are shared with result_list, so
# setting exists here is seen by everyone.
doi_index = {item["DOI"]: item for item in result_list}
while True: while True:
done_check = True
try: try:
data = in_q.get(block=True, timeout = 1) data = in_q.get(block=True, timeout = 1)
if data is _sentinel: if data is _sentinel:
@@ -161,16 +164,21 @@ def consumer(in_q, control_q, doi, live_results, result_list, control_dict, expe
alive_no += 1 alive_no += 1
print(f"C{no}__{alive_no}\r", end="") print(f"C{no}__{alive_no}\r", end="")
for item in result_list: # Each DB line is a DOI (optionally followed by metadata). Match
if item["DOI"] in data and not item["exists"]: # the WHOLE first token exactly - the old `item["DOI"] in data`
print(f"HIT in {no} content {data[0]} line {data[1]} file {data[2]} {item['exists']}") # was a substring test, so a DOI that is a prefix of a longer one
_logger.info(data) # (10.1/1 vs 10.1/12) produced a false 'exists' hit.
_logger.info("HIT") parts = data.split()
item["exists"] = True line_doi = parts[0] if parts else ""
live_results.append(item) item = doi_index.get(line_doi)
done_check = done_check and item["exists"] if item is not None and not item["exists"]:
if done_check: print(f"HIT in {no}: {line_doi}")
control_q.put(_sentinel) _logger.info("HIT %s", line_doi)
item["exists"] = True
live_results.append(item)
# All found? Signal producers to stop early (rare -> cheap).
if all(it["exists"] for it in result_list):
control_q.put(_sentinel)
except Empty: except Empty:
empty_counter += 1 empty_counter += 1
time.sleep(1) time.sleep(1)
+6 -1
View File
@@ -52,8 +52,13 @@ class DataModule(commands.Cog):
# check if current path is a file # check if current path is a file
if os.path.isfile(os.path.join(DIR_PATH_SADOX, path)): if os.path.isfile(os.path.join(DIR_PATH_SADOX, path)):
res.append(path) res.append(path)
if not res:
await ctx.send("*Conjurer grzebie w pustej skrzyni* Nie ma dziś żadnych komiksów.")
return
# randrange(len) is 0..len-1; the old randrange(0, len-1) never picked
# the last file and raised ValueError('empty range') on a single file.
# trunk-ignore(bandit/B311) # trunk-ignore(bandit/B311)
filename = res[random.randrange(0, len(res) - 1)] filename = res[random.randrange(len(res))]
# select random page # select random page
file = open(DIR_PATH_SADOX + filename, "rb") file = open(DIR_PATH_SADOX + filename, "rb")
if True: if True:
+27
View File
@@ -94,6 +94,33 @@ def test_survives_invalid_utf8_byte_and_still_finds_later_doi(tmp_path, monkeypa
assert hit, "DOI after the bad byte was not found - the file was aborted mid-read" assert hit, "DOI after the bad byte was not found - the file was aborted mid-read"
def test_doi_match_is_exact_not_substring(tmp_path, monkeypatch):
# A DB line "10.1/12" must NOT satisfy a search for "10.1/1" (the old
# `doi in line` substring test did). The exact DOI must still be found.
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
(tmp_path / "0_chunk.txt").write_text(
"10.1/12\n10.1/1\n10.2/999\n", encoding="utf-8"
)
finished, result = _run_bounded([("10.1/1", "DATA"), ("10.9/absent", "DATA")])
assert finished
by_doi = {r["DOI"]: r["exists"] for r in result}
assert by_doi["10.1/1"] is True # exact line present -> found
assert by_doi["10.9/absent"] is False
def test_doi_match_handles_line_with_trailing_metadata(tmp_path, monkeypatch):
# Lines of the form "<DOI>\t<metadata>" still match on the first token.
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
(tmp_path / "0_chunk.txt").write_text("10.5/abc\tsome title here\n", encoding="utf-8")
finished, result = _run_bounded([("10.5/abc", "DATA")])
assert finished
assert result[0]["exists"] is True
def test_discover_chunk_files_sorted_numerically(tmp_path, monkeypatch): def test_discover_chunk_files_sorted_numerically(tmp_path, monkeypatch):
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/") monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
for n in (0, 2, 10, 1): for n in (0, 2, 10, 1):