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
message_content_lower = message_content_lower.replace("imaginuje sobie: ", "")
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:
response = await OPENAICLIENT.images.generate(
model="dall-e-3",
@@ -440,10 +445,12 @@ class Events(commands.Cog):
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}"
)
return
except openai.APIConnectionError as e:
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}"
)
return
except openai.BadRequestError as e:
# Handle invalid request error, e.g. validate parameters or log
if message.author.nick:
@@ -461,27 +468,31 @@ class Events(commands.Cog):
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}"
)
return
except openai.AuthenticationError as e:
# Handle authentication error, e.g. check credentials or log
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}"
)
return
except openai.PermissionDeniedError as e:
# Handle permission error, e.g. check scope or log
# (was accidentally passing a (message, text) TUPLE as one arg)
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:
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}"
)
return
except openai.APIError as e:
# Handle API error, e.g. retry or log
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}"
)
return
if response:
self.logger.info(response)
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)
ai_check = random.randint(0, 10)
logger.info("Losowa wypowiedź")
if ai_check < 2:
if ai_check < 2 and CYCLIC_WORDS:
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)
messnum = random.randint(0, len(CYCLIC_WORDS))
messnum = random.randrange(len(CYCLIC_WORDS))
logger.debug(messnum)
logger.debug(len(CYCLIC_WORDS))
mess_key = list(CYCLIC_WORDS.keys())[messnum]
+35 -14
View File
@@ -128,13 +128,15 @@ SERVICE_EXTENSION_GROUPS = {
SERVICE_RECHECK_SECONDS = 300
def _service_alive(url: str) -> bool:
"""True when the service answers HTTP at all (any status code counts)."""
def _service_health(url: str):
"""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:
requests.get(url, timeout=3)
return True
except requests.exceptions.RequestException:
return False
return None
except requests.exceptions.RequestException as exc:
return exc
async def _load_extension_safe(name: str) -> bool:
@@ -172,16 +174,25 @@ async def _load_service_groups() -> bool:
service_headers(),
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:
alive = await asyncio.to_thread(_service_alive, group["health_url"])
if not alive:
logger.warning(
"Service '%s' unreachable (%s) - cogs stay disabled: %s",
service,
group["health_url"],
", ".join(missing),
)
continue
err = await asyncio.to_thread(_service_health, group["health_url"])
if err is not None:
logger.warning(
"Service '%s' unreachable (%s) [%s: %s] - cogs stay disabled: %s",
service,
group["health_url"],
type(err).__name__,
err,
", ".join(missing),
)
continue
logger.info("Service '%s' is alive - enabling: %s", service, ", ".join(missing))
for extension in missing:
if await _load_extension_safe(extension):
@@ -225,6 +236,16 @@ async def on_ready():
for extension in CORE_EXTENSIONS:
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()
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)
logger.info("PONG matched for %s", pong_uuid)
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:
if record.uuid in answer.keys():
record_stored = True
record.stop = True
record.entries = answer[record.uuid]
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():
record = QueryControl("Orphaned", key, "Orphan", None)
record.stop = True
@@ -353,10 +359,13 @@ def id3(url: str) -> dict:
resp.read(
metaint
) # 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(
site_url=resp.headers.get("icy-url"),
name=resp.headers.get("icy-name").title(),
genre=resp.headers.get("icy-genre").title(),
name=(resp.headers.get("icy-name") or "").title(),
genre=(resp.headers.get("icy-genre") or "").title(),
title=get_stream_title(resp.read(255)),
)
return tagdata
+6
View File
@@ -203,6 +203,12 @@ def wyszukaj(word_list, how_many, _logger=None, write_to=None):
# ---------------------------------------------------------------- tailer
def scan_tracks():
"""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:
log_file.seek(os.stat(RADIOLOG_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")
empty_counter = 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:
done_check = True
try:
data = in_q.get(block=True, timeout = 1)
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
print(f"C{no}__{alive_no}\r", end="")
for item in result_list:
if item["DOI"] in data and not item["exists"]:
print(f"HIT in {no} content {data[0]} line {data[1]} file {data[2]} {item['exists']}")
_logger.info(data)
_logger.info("HIT")
item["exists"] = True
live_results.append(item)
done_check = done_check and item["exists"]
if done_check:
control_q.put(_sentinel)
# Each DB line is a DOI (optionally followed by metadata). Match
# the WHOLE first token exactly - the old `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.
parts = data.split()
line_doi = parts[0] if parts else ""
item = doi_index.get(line_doi)
if item is not None and not item["exists"]:
print(f"HIT in {no}: {line_doi}")
_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:
empty_counter += 1
time.sleep(1)
+6 -1
View File
@@ -52,8 +52,13 @@ class DataModule(commands.Cog):
# check if current path is a file
if os.path.isfile(os.path.join(DIR_PATH_SADOX, 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)
filename = res[random.randrange(0, len(res) - 1)]
filename = res[random.randrange(len(res))]
# select random page
file = open(DIR_PATH_SADOX + filename, "rb")
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"
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):
monkeypatch.setattr(search_bot, "DATABASE_PATH", str(tmp_path) + "/")
for n in (0, 2, 10, 1):