mirror of
https://github.com/craigerl/aprsd.git
synced 2026-10-08 00:00:28 -04:00
fix: guard daemon threads against unhandled loop() exceptions (#294)
Stability fixes for the critical unhandled-exception findings: * APRSDThread.run() now catches exceptions raised by loop() and keeps the thread alive instead of silently killing it on the first transient error (with a 1s backoff to avoid a hot loop). * Collector.collect() and KeepAliveCollector.check()/log() no longer re-raise producer exceptions. A single failing stats or keepalive producer previously killed the KeepAliveThread, which is the very thread that reports dead threads. * APRSDProcessPacketThread.process_packet() no longer calls .lower() on a None addresse/to_call, which raised AttributeError and silently killed the ProcessPKT thread on malformed or third-party packets. The to_call None-guard now matches the existing guard in the MessagePacket branch.
This commit is contained in:
@@ -31,8 +31,10 @@ class Collector:
|
||||
serializable=serializable
|
||||
).copy()
|
||||
except Exception as e:
|
||||
# A single failing producer must not kill the caller (the
|
||||
# KeepAliveThread calls collect() in its loop). Log and
|
||||
# continue with the remaining producers.
|
||||
LOG.error(f'Error in producer {name} (stats): {e}')
|
||||
raise e
|
||||
return stats
|
||||
|
||||
def register_producer(self, producer_name: Callable):
|
||||
|
||||
+10
-1
@@ -83,7 +83,16 @@ class APRSDThread(threading.Thread, metaclass=abc.ABCMeta):
|
||||
self.wait(timeout=1)
|
||||
else:
|
||||
self.loop_count += 1
|
||||
can_loop = self.loop()
|
||||
try:
|
||||
can_loop = self.loop()
|
||||
except Exception as e:
|
||||
# A transient error in one loop() iteration must not
|
||||
# kill a long-running daemon thread. Log it, wait a
|
||||
# beat so we don't busy-loop on a persistent error,
|
||||
# and keep going.
|
||||
LOG.exception(f'Exception in thread {self.name} loop: {e}')
|
||||
self.wait(timeout=1)
|
||||
continue
|
||||
self._last_loop = datetime.datetime.now()
|
||||
if not can_loop:
|
||||
self.stop()
|
||||
|
||||
+3
-1
@@ -232,11 +232,13 @@ class APRSDProcessPacketThread(APRSDFilterThread):
|
||||
# plugins.
|
||||
if (
|
||||
isinstance(packet, packets.AckPacket)
|
||||
and packet.addresse
|
||||
and packet.addresse.lower() == our_call
|
||||
):
|
||||
self.process_ack_packet(packet)
|
||||
elif (
|
||||
isinstance(packet, packets.RejectPacket)
|
||||
and packet.addresse
|
||||
and packet.addresse.lower() == our_call
|
||||
):
|
||||
self.process_reject_packet(packet)
|
||||
@@ -277,7 +279,7 @@ class APRSDProcessPacketThread(APRSDFilterThread):
|
||||
else:
|
||||
self.process_other_packet(
|
||||
packet,
|
||||
for_us=(to_call.lower() == our_call),
|
||||
for_us=(bool(to_call) and to_call.lower() == our_call),
|
||||
)
|
||||
LOG.debug(f"Packet processing complete for pkt '{packet.key}'")
|
||||
return False
|
||||
|
||||
@@ -33,8 +33,10 @@ class KeepAliveCollector:
|
||||
try:
|
||||
cls.keepalive_check()
|
||||
except Exception as e:
|
||||
# A single failing producer must not kill the caller (the
|
||||
# KeepAliveThread calls check() in its loop). Log and
|
||||
# continue with the remaining producers.
|
||||
LOG.error(f'Error in producer {name} (check): {e}')
|
||||
raise e
|
||||
|
||||
def log(self) -> None:
|
||||
"""Log any relevant information during a KeepAlive check"""
|
||||
@@ -43,8 +45,9 @@ class KeepAliveCollector:
|
||||
try:
|
||||
cls.keepalive_log()
|
||||
except Exception as e:
|
||||
# Same as check(): swallow producer errors so the
|
||||
# KeepAliveThread loop cannot be killed by one bad producer.
|
||||
LOG.error(f'Error in producer {name} (check): {e}')
|
||||
raise e
|
||||
|
||||
def register(self, producer_name: Callable):
|
||||
if not isinstance(producer_name, KeepAliveProducer):
|
||||
|
||||
@@ -110,7 +110,12 @@ class TestStatsCollector(unittest.TestCase):
|
||||
self.assertIsInstance(stats, dict)
|
||||
|
||||
def test_collect_with_exception(self):
|
||||
"""Test collect() raises exception from producer."""
|
||||
"""Test collect() does not raise when a producer errors.
|
||||
|
||||
A failing producer must not propagate its exception out of
|
||||
collect(): the KeepAliveThread calls collect() and an unhandled
|
||||
exception there kills the monitoring thread itself.
|
||||
"""
|
||||
c = collector.Collector()
|
||||
|
||||
class FailingProducer:
|
||||
@@ -127,8 +132,50 @@ class TestStatsCollector(unittest.TestCase):
|
||||
producer = FailingProducer()
|
||||
c.register_producer(producer)
|
||||
|
||||
with self.assertRaises(RuntimeError):
|
||||
c.collect()
|
||||
# Should log the error but NOT raise
|
||||
with mock.patch('aprsd.stats.collector.LOG') as mock_log:
|
||||
stats = c.collect()
|
||||
self.assertEqual(stats, {})
|
||||
mock_log.error.assert_called_once()
|
||||
|
||||
def test_collect_continues_after_failing_producer(self):
|
||||
"""Test collect() still calls subsequent producers after an error."""
|
||||
c = collector.Collector()
|
||||
call_order = []
|
||||
|
||||
class FailingProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def stats(self, serializable=False):
|
||||
call_order.append('fail')
|
||||
raise RuntimeError('Stats error')
|
||||
|
||||
class GoodProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def stats(self, serializable=False):
|
||||
call_order.append('good')
|
||||
return {'ok': True}
|
||||
|
||||
c.register_producer(FailingProducer())
|
||||
c.register_producer(GoodProducer())
|
||||
|
||||
with mock.patch('aprsd.stats.collector.LOG'):
|
||||
stats = c.collect()
|
||||
|
||||
# The good producer after the failing one still ran
|
||||
self.assertEqual(call_order, ['fail', 'good'])
|
||||
self.assertIn('GoodProducer', stats)
|
||||
|
||||
def test_stop_all(self):
|
||||
"""Test stop_all() method."""
|
||||
|
||||
@@ -207,6 +207,34 @@ class TestAPRSDThread(unittest.TestCase):
|
||||
# Can't instantiate abstract class directly
|
||||
APRSDThread('AbstractThread')
|
||||
|
||||
def test_run_loop_exception_does_not_kill_thread(self):
|
||||
"""Test run() survives an exception raised in loop().
|
||||
|
||||
A transient error in one loop() iteration must not kill a
|
||||
long-running daemon thread: the exception is logged and the
|
||||
thread keeps looping.
|
||||
"""
|
||||
|
||||
class ExplodingThread(APRSDThread):
|
||||
def __init__(self, name):
|
||||
super().__init__(name)
|
||||
self.iterations = 0
|
||||
|
||||
def loop(self):
|
||||
self.iterations += 1
|
||||
if self.iterations == 1:
|
||||
raise RuntimeError('boom')
|
||||
return False
|
||||
|
||||
thread = ExplodingThread('ExplodeTest')
|
||||
thread.start()
|
||||
thread.join(timeout=2)
|
||||
|
||||
# loop() ran twice: first raised (caught), second returned False
|
||||
# which stops the thread normally.
|
||||
self.assertEqual(thread.iterations, 2)
|
||||
self.assertFalse(thread.is_alive())
|
||||
|
||||
|
||||
class TestAPRSDThreadList(unittest.TestCase):
|
||||
"""Unit tests for the APRSDThreadList class."""
|
||||
|
||||
@@ -432,6 +432,51 @@ class TestAPRSDProcessPacketThread(unittest.TestCase):
|
||||
self.process_thread.process_other_packet(packet, for_us=True)
|
||||
self.assertEqual(mock_log.info.call_count, 2)
|
||||
|
||||
def test_process_packet_no_destination_no_crash(self):
|
||||
"""A packet with no addresse/to_call must not raise AttributeError.
|
||||
|
||||
process_packet() dereferences addresse/to_call with .lower() to
|
||||
decide packet routing. A packet that carries neither field (e.g. a
|
||||
malformed or third-party frame) previously raised AttributeError,
|
||||
which -- since loop() only catches queue.Empty -- silently killed
|
||||
the ProcessPKT thread.
|
||||
"""
|
||||
from oslo_config import cfg
|
||||
|
||||
from aprsd.packets import core
|
||||
|
||||
CONF = cfg.CONF
|
||||
CONF.callsign = 'TEST'
|
||||
|
||||
packet = core.StatusPacket(
|
||||
from_call='KJ4ERJ',
|
||||
to_call=None,
|
||||
status='test status',
|
||||
)
|
||||
packet.addresse = None
|
||||
|
||||
# Must not raise
|
||||
self.process_thread.process_packet(packet)
|
||||
|
||||
def test_process_ack_packet_no_addresse_no_crash(self):
|
||||
"""An AckPacket with addresse=None must not raise AttributeError."""
|
||||
from oslo_config import cfg
|
||||
|
||||
from aprsd.packets import core
|
||||
|
||||
CONF = cfg.CONF
|
||||
CONF.callsign = 'TEST'
|
||||
|
||||
packet = core.AckPacket(
|
||||
from_call='KJ4ERJ',
|
||||
to_call=None,
|
||||
msgNo='5',
|
||||
)
|
||||
packet.addresse = None
|
||||
|
||||
# Must not raise
|
||||
self.process_thread.process_packet(packet)
|
||||
|
||||
|
||||
class TestPluginProcessPacketPiggybackAck(unittest.TestCase):
|
||||
"""Integration tests for Reply-Ack (piggyback ACK) in APRSDPluginProcessPacketThread."""
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import unittest
|
||||
from unittest import mock
|
||||
|
||||
from aprsd.utils.keepalive_collector import KeepAliveCollector
|
||||
|
||||
@@ -100,7 +101,12 @@ class TestKeepAliveCollector(unittest.TestCase):
|
||||
self.assertTrue(producer2().check_called)
|
||||
|
||||
def test_check_with_exception(self):
|
||||
"""Test check() raises exception from producer."""
|
||||
"""Test check() does not raise when a producer errors.
|
||||
|
||||
A failing producer must not propagate its exception out of
|
||||
check(): the KeepAliveThread calls it and an unhandled exception
|
||||
there kills the monitoring thread itself.
|
||||
"""
|
||||
collector = KeepAliveCollector()
|
||||
|
||||
class FailingProducer:
|
||||
@@ -120,8 +126,53 @@ class TestKeepAliveCollector(unittest.TestCase):
|
||||
producer = FailingProducer()
|
||||
collector.register(producer)
|
||||
|
||||
with self.assertRaises(RuntimeError):
|
||||
# Should log the error but NOT raise
|
||||
with mock.patch('aprsd.utils.keepalive_collector.LOG') as mock_log:
|
||||
collector.check()
|
||||
mock_log.error.assert_called_once()
|
||||
|
||||
def test_check_continues_after_failing_producer(self):
|
||||
"""Test check() still calls subsequent producers after an error."""
|
||||
collector = KeepAliveCollector()
|
||||
call_order = []
|
||||
|
||||
class FailingProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def keepalive_check(self):
|
||||
call_order.append('fail')
|
||||
raise RuntimeError('Check error')
|
||||
|
||||
def keepalive_log(self):
|
||||
pass
|
||||
|
||||
class GoodProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def keepalive_check(self):
|
||||
call_order.append('good')
|
||||
|
||||
def keepalive_log(self):
|
||||
pass
|
||||
|
||||
collector.register(FailingProducer())
|
||||
collector.register(GoodProducer())
|
||||
|
||||
with mock.patch('aprsd.utils.keepalive_collector.LOG'):
|
||||
collector.check()
|
||||
|
||||
# The good producer after the failing one still ran
|
||||
self.assertEqual(call_order, ['fail', 'good'])
|
||||
|
||||
def test_log(self):
|
||||
"""Test log() method."""
|
||||
@@ -137,7 +188,7 @@ class TestKeepAliveCollector(unittest.TestCase):
|
||||
self.assertTrue(producer2().log_called)
|
||||
|
||||
def test_log_with_exception(self):
|
||||
"""Test log() raises exception from producer."""
|
||||
"""Test log() does not raise when a producer errors."""
|
||||
collector = KeepAliveCollector()
|
||||
|
||||
class FailingProducer:
|
||||
@@ -157,8 +208,53 @@ class TestKeepAliveCollector(unittest.TestCase):
|
||||
producer = FailingProducer()
|
||||
collector.register(producer)
|
||||
|
||||
with self.assertRaises(RuntimeError):
|
||||
# Should log the error but NOT raise
|
||||
with mock.patch('aprsd.utils.keepalive_collector.LOG') as mock_log:
|
||||
collector.log()
|
||||
mock_log.error.assert_called_once()
|
||||
|
||||
def test_log_continues_after_failing_producer(self):
|
||||
"""Test log() still calls subsequent producers after an error."""
|
||||
collector = KeepAliveCollector()
|
||||
call_order = []
|
||||
|
||||
class FailingProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def keepalive_check(self):
|
||||
pass
|
||||
|
||||
def keepalive_log(self):
|
||||
call_order.append('fail')
|
||||
raise RuntimeError('Log error')
|
||||
|
||||
class GoodProducer:
|
||||
_instance = None
|
||||
|
||||
def __call__(self):
|
||||
if self._instance is None:
|
||||
self._instance = self
|
||||
return self._instance
|
||||
|
||||
def keepalive_check(self):
|
||||
pass
|
||||
|
||||
def keepalive_log(self):
|
||||
call_order.append('good')
|
||||
|
||||
collector.register(FailingProducer())
|
||||
collector.register(GoodProducer())
|
||||
|
||||
with mock.patch('aprsd.utils.keepalive_collector.LOG'):
|
||||
collector.log()
|
||||
|
||||
# The good producer after the failing one still ran
|
||||
self.assertEqual(call_order, ['fail', 'good'])
|
||||
|
||||
def test_multiple_producers(self):
|
||||
"""Test multiple producers are called."""
|
||||
|
||||
Reference in New Issue
Block a user