From dc85aba74cd2c06838a945a2df2ca001f33a6015 Mon Sep 17 00:00:00 2001 From: Josh Hawkins <32435876+hawkeye217@users.noreply.github.com> Date: Sun, 16 Aug 2026 14:16:06 -0500 Subject: [PATCH] collapse in-flight retained values by topic and release the shutdown barrier from a finally --- frigate/comms/mqtt.py | 86 +++++++++++++++++------------ frigate/test/test_mqtt_lifecycle.py | 46 +++++++++++++++ 2 files changed, 96 insertions(+), 36 deletions(-) diff --git a/frigate/comms/mqtt.py b/frigate/comms/mqtt.py index c80f723c40..1bb4cea917 100644 --- a/frigate/comms/mqtt.py +++ b/frigate/comms/mqtt.py @@ -450,10 +450,14 @@ class MqttClient(Communicator): is clean, so the broker will not resume delivery on the new one. """ with self._retained_lock: - inflight = list(self._inflight_retained.values()) + # mids are insertion ordered, so collapsing by topic keeps the + # newest value when several updates to one topic were in flight + latest = { + topic: payload for topic, payload in self._inflight_retained.values() + } self._inflight_retained.clear() - for topic, payload in inflight: + for topic, payload in latest.items(): self._queue_retained(topic, payload, True, overwrite=False) def _buffer_undelivered( @@ -540,47 +544,57 @@ class MqttClient(Communicator): self._publish_direct(QueuedPublish(topic, payload, retain)) def _publish_direct(self, queued_publish: QueuedPublish) -> None: - """Publish a queued message from the worker thread's serialized context.""" - if self.client is None: - self._buffer_undelivered(queued_publish) - return + """Publish a queued message from the worker thread's serialized context. + The waiter is released however this exits. The message is already off + the queue by now, so nothing else can recover it for a stop() that is + blocked waiting on it. + """ try: - message_info = self.client.publish( - queued_publish.topic, - queued_publish.payload, - qos=self.config.mqtt.qos, - retain=queued_publish.retain, - ) - except (OSError, mqtt.WebsocketConnectionError) as err: - logger.warning("MQTT publish failed for %s: %s", queued_publish.topic, err) - # a newer buffered value for this topic wins over the failed one - self._buffer_undelivered(queued_publish, overwrite=False) - self._schedule_reconnect() - return + if self.client is None: + self._buffer_undelivered(queued_publish) + return - if message_info.rc != mqtt.MQTT_ERR_SUCCESS: - logger.error( - "Unable to publish to %s: mqtt error %s", - queued_publish.topic, - message_info.rc, - ) - self._buffer_undelivered(queued_publish, overwrite=False) - self._schedule_reconnect() - return - - # a successful rc only means paho accepted the message; above qos 0 it - # is not durable until the broker acks, so keep a copy for replay - if queued_publish.retain and not message_info.is_published(): - with self._retained_lock: - self._inflight_retained[message_info.mid] = ( + try: + message_info = self.client.publish( queued_publish.topic, queued_publish.payload, + qos=self.config.mqtt.qos, + retain=queued_publish.retain, ) + except (OSError, mqtt.WebsocketConnectionError) as err: + logger.warning( + "MQTT publish failed for %s: %s", queued_publish.topic, err + ) + # a newer buffered value for this topic wins over the failed one + self._buffer_undelivered(queued_publish, overwrite=False) + self._schedule_reconnect() + return - if queued_publish.done is not None: - self._wait_for_publish(message_info) - queued_publish.done.set() + if message_info.rc != mqtt.MQTT_ERR_SUCCESS: + logger.error( + "Unable to publish to %s: mqtt error %s", + queued_publish.topic, + message_info.rc, + ) + self._buffer_undelivered(queued_publish, overwrite=False) + self._schedule_reconnect() + return + + # a successful rc only means paho accepted the message; above qos 0 + # it is not durable until the broker acks, so keep a copy for replay + if queued_publish.retain and not message_info.is_published(): + with self._retained_lock: + self._inflight_retained[message_info.mid] = ( + queued_publish.topic, + queued_publish.payload, + ) + + if queued_publish.done is not None: + self._wait_for_publish(message_info) + finally: + if queued_publish.done is not None: + queued_publish.done.set() def _handle_publish_event(self, mid: int) -> None: """Drop the replay copy once the broker has acknowledged the message.""" diff --git a/frigate/test/test_mqtt_lifecycle.py b/frigate/test/test_mqtt_lifecycle.py index aa37870f47..878a78f61c 100644 --- a/frigate/test/test_mqtt_lifecycle.py +++ b/frigate/test/test_mqtt_lifecycle.py @@ -471,6 +471,52 @@ class TestMqttClientLifecycle(unittest.TestCase): self.assertTrue(barrier.is_set()) + def test_shutdown_barrier_releases_when_publish_raises_unexpectedly(self) -> None: + """The message is off the queue by the time this runs, so crash cleanup + cannot recover it and only _publish_direct can release the waiter.""" + self.client.client = MagicMock() + self.client.client.publish.side_effect = RuntimeError("unexpected bug") + barrier = threading.Event() + + with self.assertRaises(RuntimeError): + self.client._publish_direct( + QueuedPublish("frigate/available", "stopped", True, barrier) + ) + + self.assertTrue(barrier.is_set()) + + def test_newest_inflight_retained_value_wins(self) -> None: + """Several updates to one topic can be unacked at once above qos 0, and + the newest is the one subscribers should end up with.""" + mock_client = MagicMock() + self.client.client = mock_client + self.client.connected = True + + for mid, payload in ((1, "ON"), (2, "OFF")): + message_info = MagicMock(rc=mqtt.MQTT_ERR_SUCCESS, mid=mid) + message_info.is_published.return_value = False + mock_client.publish.return_value = message_info + self.client._publish_direct( + QueuedPublish("frigate/front/detect/state", payload, True) + ) + + self.client._requeue_inflight_retained() + + self.assertEqual( + self.client._pending_retained["frigate/front/detect/state"], ("OFF", True) + ) + + def test_inflight_retained_does_not_clobber_queued_value(self) -> None: + """Anything still queued was written later than anything in flight.""" + self.client._pending_retained = {"frigate/front/detect/state": ("OFF", True)} + self.client._inflight_retained = {1: ("frigate/front/detect/state", "ON")} + + self.client._requeue_inflight_retained() + + self.assertEqual( + self.client._pending_retained["frigate/front/detect/state"], ("OFF", True) + ) + def test_unacked_retained_publish_survives_reconnect(self) -> None: """Above qos 0 a successful rc only means paho queued the message, and dropping the client drops its outbound queue with it."""