mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-10-03 22:07:04 +01:00
IRC health checks (#1151)
This commit is contained in:
@@ -25,6 +25,8 @@ logger = setup_logger(__name__)
|
||||
# Timing
|
||||
SOCKET_TIMEOUT = 300.0 # 5 minutes - long because we wait for DCC offers
|
||||
RECV_BUFFER = 4096
|
||||
# How often a deadline-bound read wakes up to re-check the clock
|
||||
POLL_INTERVAL = 2.0
|
||||
|
||||
# IRC channel user prefixes that indicate elevated status (ops, voice, etc.)
|
||||
# These are the download bots/servers
|
||||
@@ -248,11 +250,22 @@ class IRCClient:
|
||||
|
||||
# 366 = RPL_ENDOFNAMES - channel join is complete
|
||||
if msg.command == "366":
|
||||
logger.info(
|
||||
"Joined #%s - %s servers online",
|
||||
channel,
|
||||
len(self.online_servers),
|
||||
)
|
||||
if not self.online_servers:
|
||||
# Joining a channel that doesn't exist on this network
|
||||
# silently creates an empty one, so an empty name list is
|
||||
# the only hint that the channel name is wrong.
|
||||
logger.warning(
|
||||
"Joined #%s but no servers are online - the channel may "
|
||||
"be empty or not exist on %s",
|
||||
channel,
|
||||
self.server,
|
||||
)
|
||||
else:
|
||||
logger.info(
|
||||
"Joined #%s - %s servers online",
|
||||
channel,
|
||||
len(self.online_servers),
|
||||
)
|
||||
return
|
||||
|
||||
# Check for errors (e.g., banned, channel doesn't exist)
|
||||
@@ -296,27 +309,47 @@ class IRCClient:
|
||||
data = f"{message}\r\n".encode()
|
||||
self._socket.sendall(data)
|
||||
|
||||
def _recv_lines(self) -> Iterator[str]:
|
||||
"""Receive and yield complete CRLF-delimited IRC lines."""
|
||||
sock = self._require_socket()
|
||||
while True:
|
||||
# Check if we have a complete line in buffer
|
||||
while "\r\n" in self._buffer:
|
||||
line, self._buffer = self._buffer.split("\r\n", 1)
|
||||
if line:
|
||||
yield line
|
||||
def _recv_lines(self, deadline: float | None = None) -> Iterator[str]:
|
||||
"""Receive and yield complete CRLF-delimited IRC lines.
|
||||
|
||||
# Read more data
|
||||
try:
|
||||
data = sock.recv(RECV_BUFFER)
|
||||
if not data:
|
||||
return # Connection closed
|
||||
self._buffer += data.decode("utf-8", errors="replace")
|
||||
except TimeoutError:
|
||||
continue # Keep waiting
|
||||
except OSError as e:
|
||||
logger.warning("Socket error: %s", e)
|
||||
return # Connection error
|
||||
A deadline stops the read once it passes, even if nothing ever arrives.
|
||||
Callers time out by watching the messages they receive, so on a channel
|
||||
with no traffic at all there is nothing to watch: the recv would just
|
||||
keep blocking for SOCKET_TIMEOUT and retrying forever.
|
||||
"""
|
||||
sock = self._require_socket()
|
||||
original_timeout = sock.gettimeout()
|
||||
|
||||
try:
|
||||
while True:
|
||||
# Check if we have a complete line in buffer
|
||||
while "\r\n" in self._buffer:
|
||||
line, self._buffer = self._buffer.split("\r\n", 1)
|
||||
if line:
|
||||
yield line
|
||||
|
||||
if deadline is not None:
|
||||
remaining = deadline - time.time()
|
||||
if remaining <= 0:
|
||||
return
|
||||
# Wake up often enough to notice the deadline pass
|
||||
sock.settimeout(min(remaining, POLL_INTERVAL))
|
||||
|
||||
# Read more data
|
||||
try:
|
||||
data = sock.recv(RECV_BUFFER)
|
||||
if not data:
|
||||
return # Connection closed
|
||||
self._buffer += data.decode("utf-8", errors="replace")
|
||||
except TimeoutError:
|
||||
continue # Keep waiting (the deadline is re-checked above)
|
||||
except OSError as e:
|
||||
logger.warning("Socket error: %s", e)
|
||||
return # Connection error
|
||||
finally:
|
||||
if deadline is not None:
|
||||
with suppress(OSError):
|
||||
sock.settimeout(original_timeout)
|
||||
|
||||
def _parse_message(self, line: str) -> IRCMessage:
|
||||
"""Parse an IRC message line into components.
|
||||
@@ -427,9 +460,14 @@ class IRCClient:
|
||||
return False
|
||||
return True
|
||||
|
||||
def read_messages(self, *, auto_handle: bool = True) -> Iterator[IRCMessage]:
|
||||
def read_messages(
|
||||
self,
|
||||
*,
|
||||
auto_handle: bool = True,
|
||||
deadline: float | None = None,
|
||||
) -> Iterator[IRCMessage]:
|
||||
"""Read and yield IRC messages, optionally auto-handling PING/VERSION."""
|
||||
for line in self._recv_lines():
|
||||
for line in self._recv_lines(deadline):
|
||||
msg = self._parse_message(line)
|
||||
|
||||
# Auto-handle certain events
|
||||
@@ -453,13 +491,9 @@ class IRCClient:
|
||||
) -> DCCOffer | None:
|
||||
"""Wait for a DCC SEND offer. Returns None on timeout or no results."""
|
||||
target_event = IRCEvent.SEARCH_RESULT if result_type else IRCEvent.BOOK_RESULT
|
||||
start = time.time()
|
||||
|
||||
for msg in self.read_messages():
|
||||
if time.time() - start > timeout:
|
||||
logger.warning("Timeout waiting for DCC offer")
|
||||
return None
|
||||
deadline = time.time() + timeout
|
||||
|
||||
for msg in self.read_messages(deadline=deadline):
|
||||
if msg.event == target_event:
|
||||
if not self._is_allowed_dcc_sender(msg, expected_senders):
|
||||
continue
|
||||
@@ -491,6 +525,8 @@ class IRCClient:
|
||||
count = match.group(1)
|
||||
logger.info("Found %s matches", count)
|
||||
|
||||
if time.time() >= deadline:
|
||||
logger.warning("Timeout waiting for DCC offer")
|
||||
return None
|
||||
|
||||
@property
|
||||
|
||||
@@ -24,7 +24,7 @@ def test_wait_for_dcc_ignores_unexpected_sender(monkeypatch) -> None:
|
||||
event=IRCEvent.BOOK_RESULT,
|
||||
),
|
||||
]
|
||||
monkeypatch.setattr(client, "read_messages", lambda: iter(messages))
|
||||
monkeypatch.setattr(client, "read_messages", lambda **_: iter(messages))
|
||||
|
||||
offer = client.wait_for_dcc(timeout=1.0, result_type=False)
|
||||
|
||||
@@ -43,7 +43,7 @@ def test_wait_for_dcc_uses_expected_sender_over_online_server_list(monkeypatch)
|
||||
event=IRCEvent.BOOK_RESULT,
|
||||
)
|
||||
]
|
||||
monkeypatch.setattr(client, "read_messages", lambda: iter(messages))
|
||||
monkeypatch.setattr(client, "read_messages", lambda **_: iter(messages))
|
||||
|
||||
offer = client.wait_for_dcc(
|
||||
timeout=1.0,
|
||||
@@ -76,7 +76,7 @@ def test_wait_for_dcc_ignores_unsafe_offer_and_keeps_waiting(monkeypatch) -> Non
|
||||
event=IRCEvent.BOOK_RESULT,
|
||||
),
|
||||
]
|
||||
monkeypatch.setattr(client, "read_messages", lambda: iter(messages))
|
||||
monkeypatch.setattr(client, "read_messages", lambda **_: iter(messages))
|
||||
|
||||
offer = client.wait_for_dcc(timeout=1.0, result_type=False)
|
||||
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
"""A silent channel must still time out.
|
||||
|
||||
An empty or misnamed channel sends nothing at all, so a timeout that is only
|
||||
checked when a message arrives never fires and the search hangs forever.
|
||||
"""
|
||||
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from shelfmark.release_sources.irc.client import IRCClient
|
||||
|
||||
|
||||
class _SilentSocket:
|
||||
"""A socket that never delivers a line, only recv timeouts."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.timeout: float | None = 300.0
|
||||
self.recv_calls = 0
|
||||
|
||||
def gettimeout(self) -> float | None:
|
||||
return self.timeout
|
||||
|
||||
def settimeout(self, value: float | None) -> None:
|
||||
self.timeout = value
|
||||
|
||||
def recv(self, _bufsize: int) -> bytes:
|
||||
self.recv_calls += 1
|
||||
# A real socket blocks for the full timeout before raising; sleep a
|
||||
# sliver of it so the test stays fast but the deadline still advances.
|
||||
time.sleep(min(self.timeout or 0.0, 0.01))
|
||||
raise TimeoutError
|
||||
|
||||
|
||||
def test_wait_for_dcc_times_out_on_silent_channel() -> None:
|
||||
client = IRCClient(nick="reader", server="irc.example.test", port=6697)
|
||||
sock = _SilentSocket()
|
||||
client._socket = sock
|
||||
client._connected = True
|
||||
|
||||
start = time.time()
|
||||
offer = client.wait_for_dcc(timeout=0.5, result_type=True)
|
||||
elapsed = time.time() - start
|
||||
|
||||
assert offer is None
|
||||
assert elapsed == pytest.approx(0.5, abs=0.5)
|
||||
assert sock.recv_calls > 0
|
||||
# The long connection-wide timeout is put back for the next caller
|
||||
assert sock.timeout == 300.0
|
||||
|
||||
|
||||
def test_recv_lines_caps_socket_timeout_at_the_deadline() -> None:
|
||||
client = IRCClient(nick="reader", server="irc.example.test", port=6697)
|
||||
sock = _SilentSocket()
|
||||
client._socket = sock
|
||||
client._connected = True
|
||||
|
||||
lines = list(client._recv_lines(deadline=time.time() + 0.05))
|
||||
|
||||
assert lines == []
|
||||
# Never waits past the deadline on a single recv
|
||||
assert sock.timeout == 300.0
|
||||
Reference in New Issue
Block a user