From 8fada4e68a2bf8e9bbd3fd3840ca801351560853 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 11:41:42 +0200 Subject: [PATCH 01/15] Fix log configuration ID lifecycle (#577) --- cflib/crazyflie/log.py | 232 ++++++++++++++++++++++++++++--------- test/crazyflie/test_log.py | 215 ++++++++++++++++++++++++++++++++++ 2 files changed, 395 insertions(+), 52 deletions(-) create mode 100644 test/crazyflie/test_log.py diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 021e230e0..5e7659c9a 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -52,6 +52,8 @@ import errno import logging import struct +from collections import deque +from threading import Lock from .toc import Toc from .toc import TocFetcher @@ -60,7 +62,7 @@ from cflib.utils.callbacks import Caller __author__ = 'Bitcraze AB' -__all__ = ['Log', 'LogTocElement'] +__all__ = ['Log', 'LogConfigError', 'LogTocElement'] # Channels used for the logging port CHAN_TOC = 0 @@ -92,6 +94,10 @@ logger = logging.getLogger(__name__) +class LogConfigError(Exception): + """Raised when a log configuration cannot change lifecycle state.""" + + class LogVariable(): """A logging variable""" @@ -143,8 +149,8 @@ def __init__(self, name, period_in_ms): self.added_cb = Caller() self.err_no = 0 - # These 3 variables are set by the log subsystem when the bock is added - self.id = 0 + # These 3 variables are set by the log subsystem when the block is added + self.id = None self.cf = None self.useV2 = False @@ -152,12 +158,17 @@ def __init__(self, name, period_in_ms): self.period_in_ms = period_in_ms self._added = False self._started = False + self._delete_pending = False self.pending = False self.valid = False self.variables = [] self.default_fetch_as = [] + self._resolved_default_variables = [] self.name = name + def _get_effective_variables(self): + return self.variables + self._resolved_default_variables + def add_variable(self, name, fetch_as=None): """Add a new variable to the configuration. @@ -220,9 +231,10 @@ def _cmd_append_block(self): return CMD_APPEND_BLOCK def _setup_log_elements(self, pk, next_to_add): + variables = self._get_effective_variables() i = next_to_add - for i in range(next_to_add, len(self.variables)): - var = self.variables[i] + for i in range(next_to_add, len(variables)): + var = variables[i] if (var.is_toc_variable() is False): # Memory location logger.debug('Logging to raw memory %d, 0x%04X', var.get_storage_and_fetch_byte(), var.address) @@ -260,14 +272,15 @@ def create(self): for block in self.cf.log.log_blocks: if block.pending or block.added or block.started: pending += 1 - num_variables += len(block.variables) + num_variables += len(block._get_effective_variables()) if pending < Log.MAX_BLOCKS: # # The Crazyflie firmware can only handle 128 variables before # erroring out with ENOMEM. # - if num_variables + len(self.variables) > Log.MAX_VARIABLES: + if (num_variables + len(self._get_effective_variables()) > + Log.MAX_VARIABLES): raise AttributeError( ('Adding this configuration would exceed max number ' 'of variables (%d)' % Log.MAX_VARIABLES) @@ -291,7 +304,11 @@ def create(self): def start(self): """Start the logging for this entry""" - if (self.cf.link is not None): + cf = self.cf + if cf is None or self.id is None: + raise LogConfigError( + 'Log configuration must be added before it can be started') + if (cf.link is not None): if (self._added is False): self.create() logger.debug('First time block is started, add block') @@ -306,37 +323,30 @@ def start(self): def stop(self): """Stop the logging for this entry""" - if (self.cf.link is not None): - if (self.id is None): - logger.warning('Stopping block, but no block registered') - else: - logger.debug('Sending stop logging for block id=%d', self.id) - pk = CRTPPacket() - pk.set_header(5, CHAN_SETTINGS) - pk.data = (CMD_STOP_LOGGING, self.id) - self.cf.send_packet( - pk, expected_reply=(CMD_STOP_LOGGING, self.id)) + cf = self.cf + block_id = self.id + if cf is None or block_id is None: + return + if (cf.link is not None): + logger.debug('Sending stop logging for block id=%d', block_id) + pk = CRTPPacket() + pk.set_header(5, CHAN_SETTINGS) + pk.data = (CMD_STOP_LOGGING, block_id) + cf.send_packet( + pk, expected_reply=(CMD_STOP_LOGGING, block_id)) def delete(self): """Delete this entry in the Crazyflie""" - if (self.cf.link is not None): - if (self.id is None): - logger.warning('Delete block, but no block registered') - else: - logger.debug('LogEntry: Sending delete logging for block id=%d' - % self.id) - pk = CRTPPacket() - pk.set_header(5, CHAN_SETTINGS) - pk.data = (CMD_DELETE_BLOCK, self.id) - self.cf.send_packet( - pk, expected_reply=(CMD_DELETE_BLOCK, self.id)) + cf = self.cf + if cf is not None and self.id is not None: + cf.log._delete_config(self) def unpack_log_data(self, log_data, timestamp): """Unpack received logging data so it represent real values according to the configuration in the entry""" ret_data = {} data_index = 0 - for var in self.variables: + for var in self._get_effective_variables(): size = LogTocElement.get_size_from_id(var.fetch_as) name = var.name unpackstring = LogTocElement.get_unpack_string_from_id( @@ -415,6 +425,7 @@ class Log(): """Create log configuration""" MAX_BLOCKS = 16 + MAX_CONFIG_IDS = 256 MAX_VARIABLES = 128 # These codes can be decoded using os.stderror, but @@ -436,6 +447,7 @@ def __init__(self, crazyflie=None): self.cf = crazyflie self.toc = None self.cf.add_port_callback(CRTPPort.LOGGING, self._new_packet_cb) + self.cf.disconnected.add_callback(self._disconnected) self.toc_updated = Caller() self.state = IDLE @@ -444,7 +456,10 @@ def __init__(self, crazyflie=None): self._refresh_callback = None self._toc_cache = None - self._config_id_counter = 1 + self._registration_lock = Lock() + self._available_config_ids = deque() + self._ids_ready = False + self._reset_pending = False self._useV2 = False @@ -459,13 +474,21 @@ def add_config(self, logconf): connected when calling this method, otherwise it will fail.""" if not self.cf.link: - logger.error('Cannot add configs without being connected to a ' - 'Crazyflie!') - return + raise LogConfigError( + 'Cannot add log configurations without a connection') + + with self._registration_lock: + if logconf.id is not None or logconf.cf is not None: + raise LogConfigError( + 'Log configuration is already registered') + if not self._ids_ready: + raise LogConfigError( + 'Log configuration IDs are not ready') # If the log configuration contains variables that we added without # type (i.e we want the stored as type for fetching as well) then # resolve this now and add them to the block again. + resolved_default_variables = [] for name in logconf.default_fetch_as: var = self.toc.get_element_by_complete_name(name) if not var: @@ -473,15 +496,14 @@ def add_config(self, logconf): '%s not in TOC, this block cannot be used!', name) logconf.valid = False raise KeyError('Variable {} not in TOC'.format(name)) - # Now that we know what type this variable has, add it to the log - # config again with the correct type - logconf.add_variable(name, var.ctype) + resolved_default_variables.append(LogVariable(name, var.ctype)) # Now check that all the added variables are in the TOC and that # the total size constraint of a data packet with logging data is # not size = 0 - for var in logconf.variables: + effective_variables = logconf.variables + resolved_default_variables + for var in effective_variables: size += LogTocElement.get_size_from_id(var.fetch_as) # Check that we are able to find the variable in the TOC so # we can return error already now and not when the config is sent @@ -495,12 +517,23 @@ def add_config(self, logconf): if (size <= LogConfig.MAX_LEN and (logconf.period > 0 and logconf.period < 0xFF)): - logconf.valid = True - logconf.cf = self.cf - logconf.id = self._config_id_counter - logconf.useV2 = self._useV2 - self._config_id_counter = (self._config_id_counter + 1) % 255 - self.log_blocks.append(logconf) + with self._registration_lock: + if logconf.id is not None or logconf.cf is not None: + raise LogConfigError( + 'Log configuration is already registered') + if not self._ids_ready: + raise LogConfigError( + 'Log configuration IDs are not ready') + if not self._available_config_ids: + raise LogConfigError('No log configuration IDs available') + logconf.valid = True + logconf.cf = self.cf + logconf.id = self._available_config_ids.popleft() + logconf._delete_pending = False + logconf._resolved_default_variables = ( + resolved_default_variables) + logconf.useV2 = self._useV2 + self.log_blocks.append(logconf) self.block_added_cb.call(logconf) else: logconf.valid = False @@ -512,7 +545,6 @@ def reset(self): """ Reset the log system and remove all log blocks """ - self.log_blocks = [] self._send_reset_packet() def refresh_toc(self, refresh_done_callback, toc_cache): @@ -527,17 +559,101 @@ def refresh_toc(self, refresh_done_callback, toc_cache): self._send_reset_packet() def _send_reset_packet(self): + with self._registration_lock: + if self._reset_pending: + return + self._reset_pending = True + self._ids_ready = False + self._available_config_ids.clear() + pk = CRTPPacket() pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) pk.data = (CMD_RESET_LOGGING,) self.cf.send_packet(pk, expected_reply=(CMD_RESET_LOGGING,)) + def _detach_all_configs(self, restore_ids, require_reset_pending=False): + with self._registration_lock: + if require_reset_pending and not self._reset_pending: + return False + blocks = self.log_blocks + self.log_blocks = [] + if restore_ids: + self._available_config_ids = deque( + range(self.MAX_CONFIG_IDS)) + else: + self._available_config_ids.clear() + self._ids_ready = restore_ids + self._reset_pending = False + + callbacks = [] + for block in blocks: + callbacks.append((block, block.started, block.added)) + block._started = False + block._added = False + block._delete_pending = False + block.pending = False + block.id = None + block.cf = None + block._resolved_default_variables = [] + + for block, was_started, was_added in callbacks: + if was_started: + block.started_cb.call(block, False) + if was_added: + block.added_cb.call(block, False) + return True + + def _disconnected(self, uri): + self._detach_all_configs(restore_ids=False) + def _find_block(self, id): - for block in self.log_blocks: - if block.id == id: - return block + with self._registration_lock: + for block in self.log_blocks: + if block.id == id: + return block return None + def _retire_config(self, logconf, block_id): + with self._registration_lock: + if (logconf not in self.log_blocks or + logconf.id != block_id or + not logconf._delete_pending): + return False + + was_started = logconf.started + was_added = logconf.added + self.log_blocks.remove(logconf) + self._available_config_ids.append(block_id) + logconf._started = False + logconf._added = False + logconf._delete_pending = False + logconf.pending = False + logconf.id = None + logconf.cf = None + logconf._resolved_default_variables = [] + + if was_started: + logconf.started_cb.call(logconf, False) + if was_added: + logconf.added_cb.call(logconf, False) + return True + + def _delete_config(self, logconf): + with self._registration_lock: + if logconf not in self.log_blocks or logconf.id is None: + return + if logconf._delete_pending: + return + logconf._delete_pending = True + block_id = logconf.id + + logger.debug('LogEntry: Sending delete logging for block id=%d', + block_id) + pk = CRTPPacket() + pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) + pk.data = (CMD_DELETE_BLOCK, block_id) + self.cf.send_packet(pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) + def _new_packet_cb(self, packet): """Callback for newly arrived packets with TOC information""" chan = packet.channel @@ -604,15 +720,27 @@ def _new_packet_cb(self, packet): if error_status == 0x00 or error_status == errno.ENOENT: logger.info('Have successfully deleted id=%d', id) if block: - block.started = False - block.added = False + self._retire_config(block, id) + elif block: + with self._registration_lock: + if block.id != id or not block._delete_pending: + return + block._delete_pending = False + block.err_no = error_status + msg = self._err_codes[error_status] + block.error_cb.call(block, msg) if (cmd == CMD_RESET_LOGGING): + if error_status != 0x00: + return + + reset_completed = self._detach_all_configs( + restore_ids=True, require_reset_pending=True) + if not reset_completed: + return # Guard against multiple responses due to re-sending if not self.toc: logger.debug('Logging reset, continue with TOC download') - self.log_blocks = [] - self.toc = Toc() toc_fetcher = TocFetcher(self.cf, LogTocElement, CRTPPort.LOGGING, diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py new file mode 100644 index 000000000..bb9da780b --- /dev/null +++ b/test/crazyflie/test_log.py @@ -0,0 +1,215 @@ +# -*- coding: utf-8 -*- +# +# || ____ _ __ +# +------+ / __ )(_) /_______________ _____ ___ +# | 0xBC | / __ / / __/ ___/ ___/ __ `/_ / / _ \ +# +------+ / /_/ / / /_/ /__/ / / /_/ / / /_/ __/ +# || || /_____/\___/\___/_/ \__,_/ /___/\___/ +# +# Copyright (C) 2026 Bitcraze AB +# +# This program is free software; you can redistribute it and/or +# modify it under the terms of the GNU General Public License +# as published by the Free Software Foundation; either version 2 +# of the License, or (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# You should have received a copy of the GNU General Public License +# along with this program. If not, see . +import errno +import struct +import unittest +from unittest.mock import MagicMock + +from cflib.crazyflie import Crazyflie +from cflib.crazyflie.log import CHAN_SETTINGS +from cflib.crazyflie.log import CMD_DELETE_BLOCK +from cflib.crazyflie.log import CMD_RESET_LOGGING +from cflib.crazyflie.log import Log +from cflib.crazyflie.log import LogConfig +from cflib.crazyflie.log import LogConfigError +from cflib.crazyflie.toc import Toc +from cflib.crtp.crtpstack import CRTPPacket +from cflib.crtp.crtpstack import CRTPPort +from cflib.utils.callbacks import Caller + + +class LogTest(unittest.TestCase): + + def setUp(self): + self.cf = MagicMock(spec=Crazyflie) + self.cf.link = object() + self.cf.disconnected = Caller() + self.log = Log(self.cf) + self.cf.log = self.log + self.log.toc = Toc() + + def _acknowledge(self, command, block_id=0, error_status=0): + packet = CRTPPacket() + packet.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) + packet.data = (command, block_id, error_status) + self.log._new_packet_cb(packet) + + def _make_config(self, name): + config = LogConfig(name, 100) + config.add_memory('value', 'uint8_t', 'uint8_t', 0x1000) + return config + + def test_all_byte_values_are_available_as_log_config_ids(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + + configs = [self._make_config('config-{}'.format(i)) for i in range(256)] + for config in configs: + self.log.add_config(config) + + self.assertEqual(list(range(256)), [config.id for config in configs]) + with self.assertRaises(LogConfigError): + self.log.add_config(self._make_config('one-too-many')) + + def test_deleted_id_is_released_after_acknowledgement(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + deleted_config = self._make_config('deleted') + self.log.add_config(deleted_config) + deleted_config.delete() + for i in range(1, 256): + self.log.add_config(self._make_config('config-{}'.format(i))) + + with self.assertRaises(LogConfigError): + self.log.add_config(self._make_config('before-ack')) + + self._acknowledge(CMD_DELETE_BLOCK, deleted_config.id) + + self.assertIsNone(deleted_config.id) + self.assertIsNone(deleted_config.cf) + self.assertNotIn(deleted_config, self.log.log_blocks) + self.log.add_config(deleted_config) + self.assertEqual(0, deleted_config.id) + + def test_delete_is_idempotent_until_a_failed_acknowledgement(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + self.cf.send_packet.reset_mock() + + config.delete() + config.delete() + + self.assertEqual(1, self.cf.send_packet.call_count) + self._acknowledge(CMD_DELETE_BLOCK, config.id, errno.ENOMEM) + self.assertEqual(0, config.id) + self.assertIn(config, self.log.log_blocks) + + config.delete() + self.assertEqual(2, self.cf.send_packet.call_count) + + def test_reset_acknowledgement_detaches_configs_and_restores_ids(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('old-config') + self.log.add_config(config) + + self.log.reset() + + self.assertEqual(0, config.id) + with self.assertRaises(LogConfigError): + self.log.add_config(self._make_config('during-reset')) + + self._acknowledge(CMD_RESET_LOGGING) + + self.assertIsNone(config.id) + self.assertIsNone(config.cf) + self.assertEqual([], self.log.log_blocks) + new_config = self._make_config('new-config') + self.log.add_config(new_config) + self.assertEqual(0, new_config.id) + + def test_disconnect_detaches_configs_without_restoring_ids(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + + self.cf.disconnected.call('radio://test') + + self.assertIsNone(config.id) + self.assertIsNone(config.cf) + self.assertEqual([], self.log.log_blocks) + self.cf.link = None + with self.assertRaises(LogConfigError): + self.log.add_config(self._make_config('after-disconnect')) + + def test_detached_config_requires_registration_before_start(self): + config = self._make_config('config') + + config.stop() + config.delete() + with self.assertRaises(LogConfigError): + config.start() + + def test_config_cannot_be_registered_twice(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + + with self.assertRaises(LogConfigError): + self.log.add_config(config) + + other_config = self._make_config('other-config') + self.log.add_config(other_config) + self.assertEqual(1, other_config.id) + + def test_untyped_variables_are_resolved_fresh_when_reregistered(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + toc_element = MagicMock() + toc_element.ctype = 'uint8_t' + self.log.toc = MagicMock() + self.log.toc.get_element_by_complete_name.return_value = toc_element + config = LogConfig('config', 100) + config.add_variable('group.value') + self.log.add_config(config) + config.delete() + self._acknowledge(CMD_DELETE_BLOCK, config.id) + + toc_element.ctype = 'uint16_t' + self.log.add_config(config) + received = [] + config.data_received_cb.add_callback( + lambda timestamp, data, logconf: received.append(data)) + + config.unpack_log_data(struct.pack(' Date: Tue, 1 Sep 2026 11:45:53 +0200 Subject: [PATCH 02/15] Handle log lifecycle recovery edges --- cflib/crazyflie/log.py | 66 +++++++++++++++++++++----------------- test/crazyflie/test_log.py | 36 +++++++++++++++++++++ 2 files changed, 73 insertions(+), 29 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 5e7659c9a..aa10770f5 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -169,6 +169,23 @@ def __init__(self, name, period_in_ms): def _get_effective_variables(self): return self.variables + self._resolved_default_variables + def _detach(self): + previous_state = (self.started, self.added) + self._started = False + self._added = False + self._delete_pending = False + self.pending = False + self.id = None + self.cf = None + self._resolved_default_variables = [] + return previous_state + + def _call_detached_callbacks(self, was_started, was_added): + if was_started: + self.started_cb.call(self, False) + if was_added: + self.added_cb.call(self, False) + def add_variable(self, name, fetch_as=None): """Add a new variable to the configuration. @@ -587,20 +604,10 @@ def _detach_all_configs(self, restore_ids, require_reset_pending=False): callbacks = [] for block in blocks: - callbacks.append((block, block.started, block.added)) - block._started = False - block._added = False - block._delete_pending = False - block.pending = False - block.id = None - block.cf = None - block._resolved_default_variables = [] - - for block, was_started, was_added in callbacks: - if was_started: - block.started_cb.call(block, False) - if was_added: - block.added_cb.call(block, False) + callbacks.append((block, block._detach())) + + for block, previous_state in callbacks: + block._call_detached_callbacks(*previous_state) return True def _disconnected(self, uri): @@ -620,22 +627,12 @@ def _retire_config(self, logconf, block_id): not logconf._delete_pending): return False - was_started = logconf.started - was_added = logconf.added self.log_blocks.remove(logconf) - self._available_config_ids.append(block_id) - logconf._started = False - logconf._added = False - logconf._delete_pending = False - logconf.pending = False - logconf.id = None - logconf.cf = None - logconf._resolved_default_variables = [] + if self._ids_ready: + self._available_config_ids.append(block_id) + previous_state = logconf._detach() - if was_started: - logconf.started_cb.call(logconf, False) - if was_added: - logconf.added_cb.call(logconf, False) + logconf._call_detached_callbacks(*previous_state) return True def _delete_config(self, logconf): @@ -652,7 +649,15 @@ def _delete_config(self, logconf): pk = CRTPPacket() pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) pk.data = (CMD_DELETE_BLOCK, block_id) - self.cf.send_packet(pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) + try: + self.cf.send_packet( + pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) + except Exception: + with self._registration_lock: + if (logconf.id == block_id and + logconf._delete_pending): + logconf._delete_pending = False + raise def _new_packet_cb(self, packet): """Callback for newly arrived packets with TOC information""" @@ -732,6 +737,9 @@ def _new_packet_cb(self, packet): if (cmd == CMD_RESET_LOGGING): if error_status != 0x00: + with self._registration_lock: + if self._reset_pending: + self._reset_pending = False return reset_completed = self._detach_all_configs( diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index bb9da780b..5f73bdbf0 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -22,6 +22,7 @@ import errno import struct import unittest +from concurrent.futures import ThreadPoolExecutor from unittest.mock import MagicMock from cflib.crazyflie import Crazyflie @@ -210,6 +211,41 @@ def test_duplicate_reset_ack_does_not_detach_new_config(self): self.assertEqual(0, config.id) self.assertIn(config, self.log.log_blocks) + def test_delete_can_be_retried_when_sending_fails(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + self.cf.send_packet.reset_mock() + self.cf.send_packet.side_effect = [RuntimeError('send failed'), None] + + with self.assertRaises(RuntimeError): + config.delete() + config.delete() + + self.assertEqual(2, self.cf.send_packet.call_count) + + def test_reset_can_be_retried_after_failed_acknowledgement(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING, error_status=errno.ENOEXEC) + self.cf.send_packet.reset_mock() + + self.log.reset() + + self.assertEqual(1, self.cf.send_packet.call_count) + + def test_ids_are_unique_when_configs_are_registered_concurrently(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + configs = [self._make_config('config-{}'.format(i)) + for i in range(256)] + + with ThreadPoolExecutor(max_workers=16) as executor: + list(executor.map(self.log.add_config, configs)) + + self.assertEqual(list(range(256)), sorted( + config.id for config in configs)) + if __name__ == '__main__': unittest.main() From 17e4da787007f438e6dbea19acf15e3f6237e63e Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 11:49:55 +0200 Subject: [PATCH 03/15] Serialize log reset and delete commands --- cflib/crazyflie/log.py | 73 ++++++++++++++++++++++---------------- test/crazyflie/test_log.py | 45 +++++++++++++++++++++++ 2 files changed, 87 insertions(+), 31 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index aa10770f5..b5bf41a28 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -474,6 +474,7 @@ def __init__(self, crazyflie=None): self._toc_cache = None self._registration_lock = Lock() + self._command_lock = Lock() self._available_config_ids = deque() self._ids_ready = False self._reset_pending = False @@ -576,17 +577,24 @@ def refresh_toc(self, refresh_done_callback, toc_cache): self._send_reset_packet() def _send_reset_packet(self): - with self._registration_lock: - if self._reset_pending: - return - self._reset_pending = True - self._ids_ready = False - self._available_config_ids.clear() + with self._command_lock: + with self._registration_lock: + if self._reset_pending: + return + self._reset_pending = True + self._ids_ready = False + self._available_config_ids.clear() - pk = CRTPPacket() - pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) - pk.data = (CMD_RESET_LOGGING,) - self.cf.send_packet(pk, expected_reply=(CMD_RESET_LOGGING,)) + pk = CRTPPacket() + pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) + pk.data = (CMD_RESET_LOGGING,) + try: + self.cf.send_packet( + pk, expected_reply=(CMD_RESET_LOGGING,)) + except Exception: + with self._registration_lock: + self._reset_pending = False + raise def _detach_all_configs(self, restore_ids, require_reset_pending=False): with self._registration_lock: @@ -636,28 +644,31 @@ def _retire_config(self, logconf, block_id): return True def _delete_config(self, logconf): - with self._registration_lock: - if logconf not in self.log_blocks or logconf.id is None: - return - if logconf._delete_pending: - return - logconf._delete_pending = True - block_id = logconf.id - - logger.debug('LogEntry: Sending delete logging for block id=%d', - block_id) - pk = CRTPPacket() - pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) - pk.data = (CMD_DELETE_BLOCK, block_id) - try: - self.cf.send_packet( - pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) - except Exception: + with self._command_lock: with self._registration_lock: - if (logconf.id == block_id and - logconf._delete_pending): - logconf._delete_pending = False - raise + if (not self._ids_ready or + logconf not in self.log_blocks or + logconf.id is None): + return + if logconf._delete_pending: + return + logconf._delete_pending = True + block_id = logconf.id + + logger.debug('LogEntry: Sending delete logging for block id=%d', + block_id) + pk = CRTPPacket() + pk.set_header(CRTPPort.LOGGING, CHAN_SETTINGS) + pk.data = (CMD_DELETE_BLOCK, block_id) + try: + self.cf.send_packet( + pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) + except Exception: + with self._registration_lock: + if (logconf.id == block_id and + logconf._delete_pending): + logconf._delete_pending = False + raise def _new_packet_cb(self, packet): """Callback for newly arrived packets with TOC information""" diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 5f73bdbf0..a2a05a8b6 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -21,6 +21,7 @@ # along with this program. If not, see . import errno import struct +import threading import unittest from concurrent.futures import ThreadPoolExecutor from unittest.mock import MagicMock @@ -246,6 +247,50 @@ def test_ids_are_unique_when_configs_are_registered_concurrently(self): self.assertEqual(list(range(256)), sorted( config.id for config in configs)) + def test_reset_waits_for_in_flight_delete_command(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + delete_send_started = threading.Event() + allow_delete_send = threading.Event() + reset_finished = threading.Event() + + def send_packet(packet, expected_reply): + if packet.data[0] == CMD_DELETE_BLOCK: + delete_send_started.set() + allow_delete_send.wait() + + self.cf.send_packet.side_effect = send_packet + delete_thread = threading.Thread(target=config.delete) + delete_thread.start() + self.assertTrue(delete_send_started.wait(1.0)) + + def reset(): + self.log.reset() + reset_finished.set() + + reset_thread = threading.Thread(target=reset) + reset_thread.start() + + try: + self.assertFalse(reset_finished.wait(0.1)) + finally: + allow_delete_send.set() + delete_thread.join(1.0) + reset_thread.join(1.0) + self.assertFalse(delete_thread.is_alive()) + self.assertFalse(reset_thread.is_alive()) + + def test_reset_can_be_retried_when_sending_fails(self): + self.cf.send_packet.side_effect = [RuntimeError('send failed'), None] + + with self.assertRaises(RuntimeError): + self.log.reset() + self.log.reset() + + self.assertEqual(2, self.cf.send_packet.call_count) + if __name__ == '__main__': unittest.main() From 9cc01fee8fefac44e3feb2a44d6876d950f23a1e Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 11:53:41 +0200 Subject: [PATCH 04/15] Serialize log configuration commands --- cflib/crazyflie/log.py | 68 ++++++++++++++++++++++++++------------ test/crazyflie/test_log.py | 39 ++++++++++++++++++++++ 2 files changed, 86 insertions(+), 21 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index b5bf41a28..9769b4f0b 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -280,6 +280,15 @@ def _setup_log_elements(self, pk, next_to_add): def create(self): """Save the log configuration in the Crazyflie""" + cf = self.cf + if cf is None or self.id is None: + raise LogConfigError( + 'Log configuration must be added before it can be created') + with cf.log._command_lock: + cf.log._require_registered(self) + self._create() + + def _create(self): command = self._cmd_create_block() next_to_add = 0 is_done = False @@ -325,32 +334,38 @@ def start(self): if cf is None or self.id is None: raise LogConfigError( 'Log configuration must be added before it can be started') - if (cf.link is not None): - if (self._added is False): - self.create() - logger.debug('First time block is started, add block') - else: - logger.debug('Block already registered, starting logging' - ' for id=%d', self.id) - pk = CRTPPacket() - pk.set_header(5, CHAN_SETTINGS) - pk.data = (CMD_START_LOGGING, self.id, self.period) - self.cf.send_packet(pk, expected_reply=( - CMD_START_LOGGING, self.id)) + with cf.log._command_lock: + cf.log._require_registered(self) + if (cf.link is not None): + if (self._added is False): + self._create() + logger.debug('First time block is started, add block') + else: + logger.debug( + 'Block already registered, starting logging for id=%d', + self.id) + pk = CRTPPacket() + pk.set_header(5, CHAN_SETTINGS) + pk.data = (CMD_START_LOGGING, self.id, self.period) + cf.send_packet(pk, expected_reply=( + CMD_START_LOGGING, self.id)) def stop(self): """Stop the logging for this entry""" cf = self.cf - block_id = self.id - if cf is None or block_id is None: + if cf is None or self.id is None: return - if (cf.link is not None): - logger.debug('Sending stop logging for block id=%d', block_id) - pk = CRTPPacket() - pk.set_header(5, CHAN_SETTINGS) - pk.data = (CMD_STOP_LOGGING, block_id) - cf.send_packet( - pk, expected_reply=(CMD_STOP_LOGGING, block_id)) + with cf.log._command_lock: + if not cf.log._is_registered(self): + return + block_id = self.id + if (cf.link is not None): + logger.debug('Sending stop logging for block id=%d', block_id) + pk = CRTPPacket() + pk.set_header(5, CHAN_SETTINGS) + pk.data = (CMD_STOP_LOGGING, block_id) + cf.send_packet( + pk, expected_reply=(CMD_STOP_LOGGING, block_id)) def delete(self): """Delete this entry in the Crazyflie""" @@ -628,6 +643,17 @@ def _find_block(self, id): return block return None + def _is_registered(self, logconf): + with self._registration_lock: + return (self._ids_ready and + logconf in self.log_blocks and + logconf.cf is self.cf and + logconf.id is not None) + + def _require_registered(self, logconf): + if not self._is_registered(logconf): + raise LogConfigError('Log configuration is not registered') + def _retire_config(self, logconf, block_id): with self._registration_lock: if (logconf not in self.log_blocks or diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index a2a05a8b6..1d5dad180 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -28,6 +28,7 @@ from cflib.crazyflie import Crazyflie from cflib.crazyflie.log import CHAN_SETTINGS +from cflib.crazyflie.log import CMD_CREATE_BLOCK from cflib.crazyflie.log import CMD_DELETE_BLOCK from cflib.crazyflie.log import CMD_RESET_LOGGING from cflib.crazyflie.log import Log @@ -291,6 +292,44 @@ def test_reset_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_reset_waits_for_in_flight_start_command(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + self.log.toc = MagicMock() + self.log.toc.get_element_by_complete_name.return_value = MagicMock() + self.log.toc.get_element_id.return_value = 1 + config = LogConfig('config', 100) + config.add_variable('group.value', 'uint8_t') + self.log.add_config(config) + create_send_started = threading.Event() + allow_create_send = threading.Event() + reset_finished = threading.Event() + + def send_packet(packet, expected_reply): + if packet.data[0] == CMD_CREATE_BLOCK: + create_send_started.set() + allow_create_send.wait() + + def reset(): + self.log.reset() + reset_finished.set() + + self.cf.send_packet.side_effect = send_packet + start_thread = threading.Thread(target=config.start) + start_thread.start() + self.assertTrue(create_send_started.wait(1.0)) + reset_thread = threading.Thread(target=reset) + reset_thread.start() + + try: + self.assertFalse(reset_finished.wait(0.1)) + finally: + allow_create_send.set() + start_thread.join(1.0) + reset_thread.join(1.0) + self.assertFalse(start_thread.is_alive()) + self.assertFalse(reset_thread.is_alive()) + if __name__ == '__main__': unittest.main() From 62b635cfd72212aed7e893f185c79069f6a35b31 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 11:56:08 +0200 Subject: [PATCH 05/15] Guard log commands across lifecycle transitions --- cflib/crazyflie/log.py | 36 ++++++++++++++------------- test/crazyflie/test_log.py | 51 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 70 insertions(+), 17 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 9769b4f0b..b8e7b3cf0 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -53,6 +53,7 @@ import logging import struct from collections import deque +from contextlib import contextmanager from threading import Lock from .toc import Toc @@ -284,8 +285,7 @@ def create(self): if cf is None or self.id is None: raise LogConfigError( 'Log configuration must be added before it can be created') - with cf.log._command_lock: - cf.log._require_registered(self) + with cf.log._config_command(self): self._create() def _create(self): @@ -334,8 +334,7 @@ def start(self): if cf is None or self.id is None: raise LogConfigError( 'Log configuration must be added before it can be started') - with cf.log._command_lock: - cf.log._require_registered(self) + with cf.log._config_command(self): if (cf.link is not None): if (self._added is False): self._create() @@ -355,8 +354,8 @@ def stop(self): cf = self.cf if cf is None or self.id is None: return - with cf.log._command_lock: - if not cf.log._is_registered(self): + with cf.log._config_command(self, required=False) as registered: + if not registered: return block_id = self.id if (cf.link is not None): @@ -634,7 +633,8 @@ def _detach_all_configs(self, restore_ids, require_reset_pending=False): return True def _disconnected(self, uri): - self._detach_all_configs(restore_ids=False) + with self._command_lock: + self._detach_all_configs(restore_ids=False) def _find_block(self, id): with self._registration_lock: @@ -643,16 +643,18 @@ def _find_block(self, id): return block return None - def _is_registered(self, logconf): - with self._registration_lock: - return (self._ids_ready and - logconf in self.log_blocks and - logconf.cf is self.cf and - logconf.id is not None) - - def _require_registered(self, logconf): - if not self._is_registered(logconf): - raise LogConfigError('Log configuration is not registered') + @contextmanager + def _config_command(self, logconf, required=True): + with self._command_lock: + with self._registration_lock: + registered = (self._ids_ready and + logconf in self.log_blocks and + logconf.cf is self.cf and + logconf.id is not None and + not logconf._delete_pending) + if required and not registered: + raise LogConfigError('Log configuration is not registered') + yield registered def _retire_config(self, logconf, block_id): with self._registration_lock: diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 1d5dad180..8d66783db 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -330,6 +330,57 @@ def reset(): self.assertFalse(start_thread.is_alive()) self.assertFalse(reset_thread.is_alive()) + def test_start_is_rejected_while_delete_is_pending(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + config.delete() + + with self.assertRaises(LogConfigError): + config.start() + + self.assertEqual(2, self.cf.send_packet.call_count) + + def test_disconnect_waits_for_in_flight_start_command(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + self.log.toc = MagicMock() + self.log.toc.get_element_by_complete_name.return_value = MagicMock() + self.log.toc.get_element_id.return_value = 1 + config = LogConfig('config', 100) + config.add_variable('group.value', 'uint8_t') + self.log.add_config(config) + create_send_started = threading.Event() + allow_create_send = threading.Event() + disconnect_finished = threading.Event() + + def send_packet(packet, expected_reply): + if packet.data[0] == CMD_CREATE_BLOCK: + create_send_started.set() + allow_create_send.wait() + + def disconnect(): + self.cf.disconnected.call('radio://test') + disconnect_finished.set() + + self.cf.send_packet.side_effect = send_packet + start_thread = threading.Thread(target=config.start) + start_thread.start() + self.assertTrue(create_send_started.wait(1.0)) + disconnect_thread = threading.Thread(target=disconnect) + disconnect_thread.start() + + try: + self.assertFalse(disconnect_finished.wait(0.1)) + finally: + allow_create_send.set() + start_thread.join(1.0) + disconnect_thread.join(1.0) + self.assertFalse(start_thread.is_alive()) + self.assertFalse(disconnect_thread.is_alive()) + self.assertIsNone(config.id) + if __name__ == '__main__': unittest.main() From 7854889e8d72bd4fc017fba5b4a27852039eae85 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 11:59:06 +0200 Subject: [PATCH 06/15] Handle synchronous log disconnects --- cflib/crazyflie/log.py | 16 +++++++++++----- test/crazyflie/test_log.py | 29 +++++++++++++++++++++++++++++ 2 files changed, 40 insertions(+), 5 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index b8e7b3cf0..ae23a7845 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -55,6 +55,7 @@ from collections import deque from contextlib import contextmanager from threading import Lock +from threading import RLock from .toc import Toc from .toc import TocFetcher @@ -289,13 +290,15 @@ def create(self): self._create() def _create(self): + cf = self.cf + block_id = self.id command = self._cmd_create_block() next_to_add = 0 is_done = False num_variables = 0 pending = 0 - for block in self.cf.log.log_blocks: + for block in cf.log.log_blocks: if block.pending or block.added or block.started: pending += 1 num_variables += len(block._get_effective_variables()) @@ -319,11 +322,14 @@ def _create(self): while not is_done: pk = CRTPPacket() pk.set_header(5, CHAN_SETTINGS) - pk.data = (command, self.id) + pk.data = (command, block_id) is_done, next_to_add = self._setup_log_elements(pk, next_to_add) - logger.debug('Adding/appending log block id {}'.format(self.id)) - self.cf.send_packet(pk, expected_reply=(command, self.id)) + logger.debug('Adding/appending log block id {}'.format(block_id)) + cf.send_packet(pk, expected_reply=(command, block_id)) + if self.cf is not cf or self.id != block_id: + raise LogConfigError( + 'Log configuration was detached while being created') # Use append if we have to add more variables command = self._cmd_append_block() @@ -488,7 +494,7 @@ def __init__(self, crazyflie=None): self._toc_cache = None self._registration_lock = Lock() - self._command_lock = Lock() + self._command_lock = RLock() self._available_config_ids = deque() self._ids_ready = False self._reset_pending = False diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 8d66783db..4fd4be5c1 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -381,6 +381,35 @@ def disconnect(): self.assertFalse(disconnect_thread.is_alive()) self.assertIsNone(config.id) + def test_synchronous_disconnect_during_start_does_not_deadlock(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + self.log.toc = MagicMock() + self.log.toc.get_element_by_complete_name.return_value = MagicMock() + self.log.toc.get_element_id.return_value = 1 + config = LogConfig('config', 100) + config.add_variable('group.value', 'uint8_t') + self.log.add_config(config) + errors = [] + + def send_packet(packet, expected_reply): + self.cf.disconnected.call('radio://test') + + def start(): + try: + config.start() + except LogConfigError as error: + errors.append(error) + + self.cf.send_packet.side_effect = send_packet + start_thread = threading.Thread(target=start, daemon=True) + start_thread.start() + start_thread.join(1.0) + + self.assertFalse(start_thread.is_alive()) + self.assertEqual(1, len(errors)) + self.assertIsNone(config.id) + if __name__ == '__main__': unittest.main() From 510c4ab2c651b9561d79063f12c32771120f1cd9 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Tue, 1 Sep 2026 12:01:46 +0200 Subject: [PATCH 07/15] Validate log registration after sends --- cflib/crazyflie/log.py | 21 ++++++++++------ test/crazyflie/test_log.py | 50 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 63 insertions(+), 8 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index ae23a7845..1c6ee6c3d 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -327,9 +327,9 @@ def _create(self): logger.debug('Adding/appending log block id {}'.format(block_id)) cf.send_packet(pk, expected_reply=(command, block_id)) - if self.cf is not cf or self.id != block_id: + if not cf.log._is_current_registration(self, cf, block_id): raise LogConfigError( - 'Log configuration was detached while being created') + 'Log configuration changed while being created') # Use append if we have to add more variables command = self._cmd_append_block() @@ -652,16 +652,21 @@ def _find_block(self, id): @contextmanager def _config_command(self, logconf, required=True): with self._command_lock: - with self._registration_lock: - registered = (self._ids_ready and - logconf in self.log_blocks and - logconf.cf is self.cf and - logconf.id is not None and - not logconf._delete_pending) + registered = self._is_current_registration( + logconf, self.cf, logconf.id) if required and not registered: raise LogConfigError('Log configuration is not registered') yield registered + def _is_current_registration(self, logconf, cf, block_id): + with self._registration_lock: + return (self._ids_ready and + logconf in self.log_blocks and + logconf.cf is cf and + logconf.id == block_id and + block_id is not None and + not logconf._delete_pending) + def _retire_config(self, logconf, block_id): with self._registration_lock: if (logconf not in self.log_blocks or diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 4fd4be5c1..eae26789e 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -29,6 +29,7 @@ from cflib.crazyflie import Crazyflie from cflib.crazyflie.log import CHAN_SETTINGS from cflib.crazyflie.log import CMD_CREATE_BLOCK +from cflib.crazyflie.log import CMD_CREATE_BLOCK_V2 from cflib.crazyflie.log import CMD_DELETE_BLOCK from cflib.crazyflie.log import CMD_RESET_LOGGING from cflib.crazyflie.log import Log @@ -61,6 +62,17 @@ def _make_config(self, name): config.add_memory('value', 'uint8_t', 'uint8_t', 0x1000) return config + def _make_multi_packet_config(self): + self.log._useV2 = True + self.log.toc = MagicMock() + self.log.toc.get_element_by_complete_name.return_value = MagicMock() + self.log.toc.get_element_id.return_value = 1 + config = LogConfig('multi-packet', 100) + for i in range(20): + config.add_variable('group.value{}'.format(i), 'uint8_t') + self.log.add_config(config) + return config + def test_all_byte_values_are_available_as_log_config_ids(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) @@ -410,6 +422,44 @@ def start(): self.assertEqual(1, len(errors)) self.assertIsNone(config.id) + def test_synchronous_reset_during_create_prevents_append(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_multi_packet_config() + sent_commands = [] + + def send_packet(packet, expected_reply): + sent_commands.append(packet.data[0]) + if packet.data[0] == CMD_CREATE_BLOCK_V2: + self.log.reset() + + self.cf.send_packet.side_effect = send_packet + + with self.assertRaises(LogConfigError): + config.start() + + self.assertEqual( + [CMD_CREATE_BLOCK_V2, CMD_RESET_LOGGING], sent_commands) + + def test_synchronous_delete_during_create_prevents_append(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_multi_packet_config() + sent_commands = [] + + def send_packet(packet, expected_reply): + sent_commands.append(packet.data[0]) + if packet.data[0] == CMD_CREATE_BLOCK_V2: + config.delete() + + self.cf.send_packet.side_effect = send_packet + + with self.assertRaises(LogConfigError): + config.start() + + self.assertEqual( + [CMD_CREATE_BLOCK_V2, CMD_DELETE_BLOCK], sent_commands) + if __name__ == '__main__': unittest.main() From 2ae5f7dba61bb704aeaaa23c6cb5bb669126a957 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Thu, 3 Sep 2026 13:59:10 +0200 Subject: [PATCH 08/15] Restore log IDs when reset send fails --- cflib/crazyflie/log.py | 7 ++++++- test/crazyflie/test_log.py | 16 ++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 1c6ee6c3d..a1f660a72 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -601,6 +601,8 @@ def _send_reset_packet(self): with self._registration_lock: if self._reset_pending: return + ids_were_ready = self._ids_ready + available_config_ids = self._available_config_ids.copy() self._reset_pending = True self._ids_ready = False self._available_config_ids.clear() @@ -613,7 +615,10 @@ def _send_reset_packet(self): pk, expected_reply=(CMD_RESET_LOGGING,)) except Exception: with self._registration_lock: - self._reset_pending = False + if self._reset_pending: + self._reset_pending = False + self._ids_ready = ids_were_ready + self._available_config_ids = available_config_ids raise def _detach_all_configs(self, restore_ids, require_reset_pending=False): diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index eae26789e..9e3b0c413 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -304,6 +304,22 @@ def test_reset_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_failed_reset_send_restores_config_id_state(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + self.cf.send_packet.side_effect = RuntimeError('send failed') + + with self.assertRaises(RuntimeError): + self.log.reset() + + self.assertTrue(self.log._is_current_registration( + config, self.cf, config.id)) + other_config = self._make_config('other-config') + self.log.add_config(other_config) + self.assertEqual(1, other_config.id) + def test_reset_waits_for_in_flight_start_command(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) From 93e70a40076d35be3bac646f0413730fe98658cb Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Thu, 3 Sep 2026 14:10:20 +0200 Subject: [PATCH 09/15] Snapshot log registrations during create --- cflib/crazyflie/log.py | 4 +++- test/crazyflie/test_log.py | 29 +++++++++++++++++++++++++++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index a1f660a72..8c5ed5689 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -298,7 +298,9 @@ def _create(self): num_variables = 0 pending = 0 - for block in cf.log.log_blocks: + with cf.log._registration_lock: + log_blocks = list(cf.log.log_blocks) + for block in log_blocks: if block.pending or block.added or block.started: pending += 1 num_variables += len(block._get_effective_variables()) diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 9e3b0c413..abf5dbb8a 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -476,6 +476,35 @@ def send_packet(packet, expected_reply): self.assertEqual( [CMD_CREATE_BLOCK_V2, CMD_DELETE_BLOCK], sent_commands) + def test_create_counts_from_stable_registration_snapshot(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + existing_config = self._make_config('existing-config') + self.log.add_config(existing_config) + existing_config.pending = True + existing_config.variables *= Log.MAX_VARIABLES + config = self._make_config('config') + self.log.add_config(config) + + class RemovingBlock: + pending = False + added = False + started = False + + def __init__(self, log): + self.log = log + + def _get_effective_variables(self): + self.log.log_blocks.remove(self) + return [] + + removing_block = RemovingBlock(self.log) + removing_block.pending = True + self.log.log_blocks.insert(0, removing_block) + + with self.assertRaises(AttributeError): + config.create() + if __name__ == '__main__': unittest.main() From 68dab350d868b57c8437a254dd8dbb95f7e0aa54 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Thu, 3 Sep 2026 15:58:11 +0200 Subject: [PATCH 10/15] Harden log lifecycle concurrency --- cflib/crazyflie/log.py | 275 ++++++++++++++++++++++++----------- test/crazyflie/test_log.py | 289 +++++++++++++++++++++++++++++++++---- 2 files changed, 451 insertions(+), 113 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 8c5ed5689..88e5c9aa5 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -219,17 +219,25 @@ def add_memory(self, name, fetch_as, stored_as, address): stored_as, address)) def _set_added(self, added): - if added != self._added: + if self._set_added_state(added): self.added_cb.call(self, added) + + def _set_added_state(self, added): + changed = added != self._added self._added = added + return changed def _get_added(self): return self._added def _set_started(self, started): - if started != self._started: + if self._set_started_state(started): self.started_cb.call(self, started) + + def _set_started_state(self, started): + changed = started != self._started self._started = started + return changed def _get_started(self): return self._started @@ -296,14 +304,7 @@ def _create(self): next_to_add = 0 is_done = False - num_variables = 0 - pending = 0 - with cf.log._registration_lock: - log_blocks = list(cf.log.log_blocks) - for block in log_blocks: - if block.pending or block.added or block.started: - pending += 1 - num_variables += len(block._get_effective_variables()) + pending, num_variables = cf.log._get_active_config_usage() if pending < Log.MAX_BLOCKS: # @@ -497,9 +498,13 @@ def __init__(self, crazyflie=None): self._registration_lock = Lock() self._command_lock = RLock() + self._command_depth = 0 + self._deferred_calls = deque() + self._dispatching_deferred_calls = False self._available_config_ids = deque() self._ids_ready = False self._reset_pending = False + self._ids_before_reset = None self._useV2 = False @@ -513,6 +518,11 @@ def add_config(self, logconf): cannot be used. Since a valid TOC is required, a Crazyflie has to be connected when calling this method, otherwise it will fail.""" + with self._command_scope(): + self._add_config(logconf) + self._defer_call(self.block_added_cb.call, logconf) + + def _add_config(self, logconf): if not self.cf.link: raise LogConfigError( 'Cannot add log configurations without a connection') @@ -524,13 +534,17 @@ def add_config(self, logconf): if not self._ids_ready: raise LogConfigError( 'Log configuration IDs are not ready') + toc = self.toc + + if toc is None: + raise LogConfigError('Log TOC is not available') # If the log configuration contains variables that we added without # type (i.e we want the stored as type for fetching as well) then # resolve this now and add them to the block again. resolved_default_variables = [] for name in logconf.default_fetch_as: - var = self.toc.get_element_by_complete_name(name) + var = toc.get_element_by_complete_name(name) if not var: logger.warning( '%s not in TOC, this block cannot be used!', name) @@ -548,7 +562,7 @@ def add_config(self, logconf): # Check that we are able to find the variable in the TOC so # we can return error already now and not when the config is sent if var.is_toc_variable(): - if (self.toc.get_element_by_complete_name(var.name) is None): + if (toc.get_element_by_complete_name(var.name) is None): logger.warning( 'Log: %s not in TOC, this block cannot be used!', var.name) @@ -574,7 +588,6 @@ def add_config(self, logconf): resolved_default_variables) logconf.useV2 = self._useV2 self.log_blocks.append(logconf) - self.block_added_cb.call(logconf) else: logconf.valid = False raise AttributeError( @@ -590,21 +603,22 @@ def reset(self): def refresh_toc(self, refresh_done_callback, toc_cache): """Start refreshing the table of loggale variables""" - self._useV2 = self.cf.platform.get_protocol_version() >= 4 + with self._command_scope(): + self._useV2 = self.cf.platform.get_protocol_version() >= 4 - self._toc_cache = toc_cache - self._refresh_callback = refresh_done_callback - self.toc = None + self._toc_cache = toc_cache + self._refresh_callback = refresh_done_callback + self.toc = None - self._send_reset_packet() + self._send_reset_packet() def _send_reset_packet(self): - with self._command_lock: + with self._command_scope(): with self._registration_lock: if self._reset_pending: return - ids_were_ready = self._ids_ready - available_config_ids = self._available_config_ids.copy() + self._ids_before_reset = ( + self._ids_ready, self._available_config_ids.copy()) self._reset_pending = True self._ids_ready = False self._available_config_ids.clear() @@ -616,17 +630,24 @@ def _send_reset_packet(self): self.cf.send_packet( pk, expected_reply=(CMD_RESET_LOGGING,)) except Exception: - with self._registration_lock: - if self._reset_pending: - self._reset_pending = False - self._ids_ready = ids_were_ready - self._available_config_ids = available_config_ids + self._restore_ids_after_failed_reset() raise + def _restore_ids_after_failed_reset(self): + with self._registration_lock: + if not self._reset_pending: + return False + self._reset_pending = False + if self._ids_before_reset is not None: + self._ids_ready, self._available_config_ids = ( + self._ids_before_reset) + self._ids_before_reset = None + return True + def _detach_all_configs(self, restore_ids, require_reset_pending=False): with self._registration_lock: if require_reset_pending and not self._reset_pending: - return False + return None blocks = self.log_blocks self.log_blocks = [] if restore_ids: @@ -636,18 +657,72 @@ def _detach_all_configs(self, restore_ids, require_reset_pending=False): self._available_config_ids.clear() self._ids_ready = restore_ids self._reset_pending = False + self._ids_before_reset = None callbacks = [] for block in blocks: callbacks.append((block, block._detach())) - for block, previous_state in callbacks: - block._call_detached_callbacks(*previous_state) - return True + return callbacks def _disconnected(self, uri): + with self._command_scope(): + callbacks = self._detach_all_configs(restore_ids=False) + self._defer_detached_callbacks(callbacks) + + def _defer_detached_callbacks(self, callbacks): + for block, previous_state in callbacks: + self._defer_call( + block._call_detached_callbacks, *previous_state) + + @contextmanager + def _command_scope(self): + should_dispatch = False with self._command_lock: - self._detach_all_configs(restore_ids=False) + self._command_depth += 1 + try: + yield + finally: + self._command_depth -= 1 + if (self._command_depth == 0 and + self._deferred_calls and + not self._dispatching_deferred_calls): + self._dispatching_deferred_calls = True + should_dispatch = True + + if should_dispatch: + self._dispatch_deferred_calls() + + def _defer_call(self, callback, *args): + self._deferred_calls.append((callback, args)) + + def _dispatch_deferred_calls(self): + first_error = None + while True: + with self._command_lock: + if not self._deferred_calls: + self._dispatching_deferred_calls = False + break + callback, args = self._deferred_calls.popleft() + try: + callback(*args) + except Exception as error: + if first_error is None: + first_error = error + + if first_error is not None: + raise first_error + + def _get_active_config_usage(self): + with self._registration_lock: + log_blocks = list(self.log_blocks) + + active_blocks = [ + block for block in log_blocks + if block.pending or block.added or block.started] + variable_count = sum( + len(block._get_effective_variables()) for block in active_blocks) + return len(active_blocks), variable_count def _find_block(self, id): with self._registration_lock: @@ -658,7 +733,7 @@ def _find_block(self, id): @contextmanager def _config_command(self, logconf, required=True): - with self._command_lock: + with self._command_scope(): registered = self._is_current_registration( logconf, self.cf, logconf.id) if required and not registered: @@ -684,13 +759,16 @@ def _retire_config(self, logconf, block_id): self.log_blocks.remove(logconf) if self._ids_ready: self._available_config_ids.append(block_id) + elif self._reset_pending and self._ids_before_reset is not None: + ids_were_ready, available_config_ids = self._ids_before_reset + if ids_were_ready: + available_config_ids.append(block_id) previous_state = logconf._detach() - logconf._call_detached_callbacks(*previous_state) - return True + return logconf, previous_state def _delete_config(self, logconf): - with self._command_lock: + with self._command_scope(): with self._registration_lock: if (not self._ids_ready or logconf not in self.log_blocks or @@ -723,11 +801,34 @@ def _new_packet_cb(self, packet): payload = packet.data[1:] if (chan == CHAN_SETTINGS): + self._handle_settings_packet(cmd, payload) + + if (chan == CHAN_LOGDATA): + chan = packet.channel + id = packet.data[0] + block = self._find_block(id) + timestamps = struct.unpack(' Date: Thu, 3 Sep 2026 16:18:20 +0200 Subject: [PATCH 11/15] Fix log create failure state --- cflib/crazyflie/log.py | 42 +++++++++++++++++++------------ test/crazyflie/test_log.py | 51 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 16 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 88e5c9aa5..867ed2a97 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -321,21 +321,30 @@ def _create(self): raise AttributeError( 'Configuration has max number of blocks (%d)' % Log.MAX_BLOCKS ) - self.pending += 1 - while not is_done: - pk = CRTPPacket() - pk.set_header(5, CHAN_SETTINGS) - pk.data = (command, block_id) - is_done, next_to_add = self._setup_log_elements(pk, next_to_add) - - logger.debug('Adding/appending log block id {}'.format(block_id)) - cf.send_packet(pk, expected_reply=(command, block_id)) - if not cf.log._is_current_registration(self, cf, block_id): - raise LogConfigError( - 'Log configuration changed while being created') + self.pending = True + create_was_sent = False + try: + while not is_done: + pk = CRTPPacket() + pk.set_header(5, CHAN_SETTINGS) + pk.data = (command, block_id) + is_done, next_to_add = self._setup_log_elements( + pk, next_to_add) + + logger.debug( + 'Adding/appending log block id {}'.format(block_id)) + cf.send_packet(pk, expected_reply=(command, block_id)) + create_was_sent = True + if not cf.log._is_current_registration(self, cf, block_id): + raise LogConfigError( + 'Log configuration changed while being created') - # Use append if we have to add more variables - command = self._cmd_append_block() + # Use append if we have to add more variables + command = self._cmd_append_block() + except Exception: + if not create_was_sent: + self.pending = False + raise def start(self): """Start the logging for this entry""" @@ -850,7 +859,8 @@ def _handle_settings_packet(self, cmd, payload): logger.warning('Error %d when adding id=%d (%s)', error_status, id, msg) block.err_no = error_status - callbacks.append((block.added_cb, (False,))) + block.pending = False + callbacks.append((block.added_cb, (block, False))) callbacks.append((block.error_cb, (block, msg))) else: @@ -873,7 +883,7 @@ def _handle_settings_packet(self, cmd, payload): if (block is not None and self._is_current_registration( block, self.cf, id)): block.err_no = error_status - callbacks.append((block.started_cb, (self, False))) + callbacks.append((block.started_cb, (block, False))) # This is a temporary fix, we are adding a new issue # for this. For some reason we get an error back after # the block has been started and added. This will show diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index 400d9b4a4..e46fa9696 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -310,6 +310,57 @@ def test_delete_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_create_send_failure_clears_pending_state(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + self.cf.send_packet.side_effect = RuntimeError('send failed') + + with self.assertRaisesRegex(RuntimeError, 'send failed'): + config.create() + + self.assertFalse(config.pending) + self.assertEqual((0, 0), self.log._get_active_config_usage()) + + def test_create_error_clears_pending_and_reports_config(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + added_events = [] + error_events = [] + config.added_cb.add_callback( + lambda logconf, added: added_events.append((logconf, added))) + config.error_cb.add_callback( + lambda logconf, message: error_events.append((logconf, message))) + + config.create() + self.assertIs(config.pending, True) + self._acknowledge(CMD_CREATE_BLOCK, config.id, errno.ENOMEM) + + self.assertFalse(config.pending) + self.assertEqual([(config, False)], added_events) + self.assertEqual( + [(config, 'No more memory available')], error_events) + self.assertEqual((0, 0), self.log._get_active_config_usage()) + + def test_start_error_reports_config(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + started_events = [] + config.started_cb.add_callback( + lambda logconf, started: + started_events.append((logconf, started))) + + config.create() + self._acknowledge(CMD_CREATE_BLOCK, config.id) + self._acknowledge(CMD_START_LOGGING, config.id, errno.ENOEXEC) + + self.assertEqual([(config, False)], started_events) + def test_reset_can_be_retried_after_failed_acknowledgement(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) From 3204abf81d23e12a798a32927d32ea06f57738f2 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Fri, 4 Sep 2026 10:20:16 +0200 Subject: [PATCH 12/15] Handle log reset while disconnected --- cflib/crazyflie/log.py | 3 +++ test/crazyflie/test_log.py | 14 ++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 867ed2a97..ceeb4f4dd 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -623,6 +623,9 @@ def refresh_toc(self, refresh_done_callback, toc_cache): def _send_reset_packet(self): with self._command_scope(): + if self.cf.link is None: + return + with self._registration_lock: if self._reset_pending: return diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index e46fa9696..411762adf 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -459,6 +459,20 @@ def test_reset_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_reset_while_disconnected_does_not_block_later_reset(self): + self.cf.link = None + + self.log.reset() + + self.assertFalse(self.log._reset_pending) + self.cf.send_packet.assert_not_called() + + self.cf.link = object() + self.log.reset() + + self.assertTrue(self.log._reset_pending) + self.cf.send_packet.assert_called_once() + def test_failed_reset_send_restores_config_id_state(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) From 687e656d82cc20ddf7d957ffd32026172844ca25 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Fri, 4 Sep 2026 10:35:21 +0200 Subject: [PATCH 13/15] Serialize log data lifecycle handling --- cflib/crazyflie/log.py | 42 +++++++++++++++--------- test/crazyflie/test_log.py | 66 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 92 insertions(+), 16 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index ceeb4f4dd..444a0c44c 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -393,6 +393,10 @@ def delete(self): def unpack_log_data(self, log_data, timestamp): """Unpack received logging data so it represent real values according to the configuration in the entry""" + ret_data = self._unpack_log_data(log_data) + self.data_received_cb.call(timestamp, ret_data, self) + + def _unpack_log_data(self, log_data): ret_data = {} data_index = 0 for var in self._get_effective_variables(): @@ -404,7 +408,7 @@ def unpack_log_data(self, log_data, timestamp): unpackstring, log_data[data_index:data_index + size])[0] data_index += size ret_data[name] = value - self.data_received_cb.call(timestamp, ret_data, self) + return ret_data class LogTocElement: @@ -816,17 +820,22 @@ def _new_packet_cb(self, packet): self._handle_settings_packet(cmd, payload) if (chan == CHAN_LOGDATA): - chan = packet.channel - id = packet.data[0] - block = self._find_block(id) - timestamps = struct.unpack(' Date: Fri, 4 Sep 2026 10:51:56 +0200 Subject: [PATCH 14/15] Recover partial log configuration creates --- cflib/crazyflie/log.py | 53 ++++++++++------ test/crazyflie/test_log.py | 122 ++++++++++++++++++++++++++++++++++++- 2 files changed, 154 insertions(+), 21 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 444a0c44c..560f886cf 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -322,7 +322,7 @@ def _create(self): 'Configuration has max number of blocks (%d)' % Log.MAX_BLOCKS ) self.pending = True - create_was_sent = False + create_send_attempted = False try: while not is_done: pk = CRTPPacket() @@ -333,8 +333,8 @@ def _create(self): logger.debug( 'Adding/appending log block id {}'.format(block_id)) + create_send_attempted = True cf.send_packet(pk, expected_reply=(command, block_id)) - create_was_sent = True if not cf.log._is_current_registration(self, cf, block_id): raise LogConfigError( 'Log configuration changed while being created') @@ -342,8 +342,15 @@ def _create(self): # Use append if we have to add more variables command = self._cmd_append_block() except Exception: - if not create_was_sent: + if not create_send_attempted: self.pending = False + else: + try: + cf.log._delete_config(self) + except Exception: + logger.warning( + 'Failed to delete partial log block id=%d', + block_id, exc_info=True) raise def start(self): @@ -694,20 +701,30 @@ def _defer_detached_callbacks(self, callbacks): @contextmanager def _command_scope(self): should_dispatch = False - with self._command_lock: - self._command_depth += 1 - try: - yield - finally: - self._command_depth -= 1 - if (self._command_depth == 0 and - self._deferred_calls and - not self._dispatching_deferred_calls): - self._dispatching_deferred_calls = True - should_dispatch = True - - if should_dispatch: - self._dispatch_deferred_calls() + try: + with self._command_lock: + self._command_depth += 1 + try: + yield + finally: + self._command_depth -= 1 + if (self._command_depth == 0 and + self._deferred_calls and + not self._dispatching_deferred_calls): + self._dispatching_deferred_calls = True + should_dispatch = True + except BaseException: + if should_dispatch: + try: + self._dispatch_deferred_calls() + except BaseException: + logger.warning( + 'Deferred callback failed while handling command error', + exc_info=True) + raise + else: + if should_dispatch: + self._dispatch_deferred_calls() def _defer_call(self, callback, *args): self._deferred_calls.append((callback, args)) @@ -722,7 +739,7 @@ def _dispatch_deferred_calls(self): callback, args = self._deferred_calls.popleft() try: callback(*args) - except Exception as error: + except BaseException as error: if first_error is None: first_error = error diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index dce33085b..d0638ac7b 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -25,10 +25,12 @@ import unittest from concurrent.futures import ThreadPoolExecutor from unittest.mock import MagicMock +from unittest.mock import patch from cflib.crazyflie import Crazyflie from cflib.crazyflie.log import CHAN_LOGDATA from cflib.crazyflie.log import CHAN_SETTINGS +from cflib.crazyflie.log import CMD_APPEND_BLOCK_V2 from cflib.crazyflie.log import CMD_CREATE_BLOCK from cflib.crazyflie.log import CMD_CREATE_BLOCK_V2 from cflib.crazyflie.log import CMD_DELETE_BLOCK @@ -310,19 +312,133 @@ def test_delete_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) - def test_create_send_failure_clears_pending_state(self): + def test_create_send_failure_schedules_cleanup(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) config = self._make_toc_config('config') self.log.add_config(config) - self.cf.send_packet.side_effect = RuntimeError('send failed') + self.cf.send_packet.side_effect = [RuntimeError('send failed'), None] with self.assertRaisesRegex(RuntimeError, 'send failed'): config.create() - self.assertFalse(config.pending) + self.assertTrue(config.pending) + self.assertTrue(config._delete_pending) + self.assertEqual((1, 1), self.log._get_active_config_usage()) + + self._acknowledge(CMD_DELETE_BLOCK, config.id) + self.assertEqual((0, 0), self.log._get_active_config_usage()) + def test_create_send_exception_after_ack_deletes_firmware_block(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + commands = [] + + def send_packet(packet, expected_reply): + command = packet.data[0] + commands.append(command) + if command == CMD_CREATE_BLOCK: + self._acknowledge(CMD_CREATE_BLOCK, config.id) + raise RuntimeError('post-send callback failed') + + self.cf.send_packet.side_effect = send_packet + + with self.assertRaisesRegex(RuntimeError, 'post-send callback failed'): + config.create() + + self.assertEqual( + [CMD_CREATE_BLOCK, CMD_START_LOGGING, CMD_DELETE_BLOCK], commands) + self.assertTrue(config._delete_pending) + + def test_append_send_failure_deletes_partial_firmware_block(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_multi_packet_config() + commands = [] + + def send_packet(packet, expected_reply): + command = packet.data[0] + commands.append(command) + if command == CMD_APPEND_BLOCK_V2: + raise RuntimeError('append send failed') + + self.cf.send_packet.side_effect = send_packet + + with self.assertRaisesRegex(RuntimeError, 'append send failed'): + config.create() + + self.assertEqual( + [CMD_CREATE_BLOCK_V2, CMD_APPEND_BLOCK_V2, CMD_DELETE_BLOCK], + commands) + self.assertTrue(config._delete_pending) + block_id = config.id + + self._acknowledge(CMD_DELETE_BLOCK, block_id) + + self.assertIsNone(config.id) + self.assertNotIn(config, self.log.log_blocks) + + def test_append_failure_preserves_error_while_draining_callbacks(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_multi_packet_config() + added_states = [] + + class CallbackAbort(BaseException): + pass + + def added_callback(log_config, added): + added_states.append(added) + if added: + raise CallbackAbort('callback failed') + + config.added_cb.add_callback(added_callback) + + def send_packet(packet, expected_reply): + command = packet.data[0] + if command == CMD_CREATE_BLOCK_V2: + self._acknowledge(CMD_CREATE_BLOCK_V2, config.id) + elif command == CMD_APPEND_BLOCK_V2: + raise ValueError('append send failed') + elif command == CMD_DELETE_BLOCK: + self._acknowledge(CMD_DELETE_BLOCK, config.id) + + self.cf.send_packet.side_effect = send_packet + + with patch('cflib.crazyflie.log.logger') as log: + with self.assertRaisesRegex(ValueError, 'append send failed'): + config.create() + + log.warning.assert_called_once() + self.assertEqual([True, False], added_states) + self.assertEqual(0, len(self.log._deferred_calls)) + self.assertFalse(self.log._dispatching_deferred_calls) + + def test_deferred_base_exception_propagates_after_draining_callbacks(self): + calls = [] + + class CallbackAbort(BaseException): + pass + + def failing_callback(): + calls.append('failing') + raise CallbackAbort('callback failed') + + def later_callback(): + calls.append('later') + + with self.assertRaisesRegex(CallbackAbort, 'callback failed'): + with self.log._command_scope(): + self.log._defer_call(failing_callback) + self.log._defer_call(later_callback) + + self.assertEqual(['failing', 'later'], calls) + self.assertEqual(0, len(self.log._deferred_calls)) + self.assertFalse(self.log._dispatching_deferred_calls) + def test_create_error_clears_pending_and_reports_config(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) From 34552edfa531eba04f87742cb1068cb1a58e1ae6 Mon Sep 17 00:00:00 2001 From: Arnaud Taffanel Date: Fri, 4 Sep 2026 11:09:18 +0200 Subject: [PATCH 15/15] Harden log create failure recovery --- cflib/crazyflie/log.py | 30 +++++--- test/crazyflie/test_log.py | 145 ++++++++++++++++++++++++++++++++++++- 2 files changed, 163 insertions(+), 12 deletions(-) diff --git a/cflib/crazyflie/log.py b/cflib/crazyflie/log.py index 560f886cf..f39c23360 100644 --- a/cflib/crazyflie/log.py +++ b/cflib/crazyflie/log.py @@ -341,13 +341,13 @@ def _create(self): # Use append if we have to add more variables command = self._cmd_append_block() - except Exception: + except BaseException: if not create_send_attempted: self.pending = False else: try: cf.log._delete_config(self) - except Exception: + except BaseException: logger.warning( 'Failed to delete partial log block id=%d', block_id, exc_info=True) @@ -652,7 +652,7 @@ def _send_reset_packet(self): try: self.cf.send_packet( pk, expected_reply=(CMD_RESET_LOGGING,)) - except Exception: + except BaseException: self._restore_ids_after_failed_reset() raise @@ -820,7 +820,7 @@ def _delete_config(self, logconf): try: self.cf.send_packet( pk, expected_reply=(CMD_DELETE_BLOCK, block_id)) - except Exception: + except BaseException: with self._registration_lock: if (logconf.id == block_id and logconf._delete_pending): @@ -885,13 +885,21 @@ def _handle_settings_packet(self, cmd, payload): block, self.cf, id): return else: - msg = self._err_codes[error_status] + msg = self._err_codes.get( + error_status, 'Unknown error') logger.warning('Error %d when adding id=%d (%s)', error_status, id, msg) block.err_no = error_status - block.pending = False - callbacks.append((block.added_cb, (block, False))) - callbacks.append((block.error_cb, (block, msg))) + self._defer_call( + block.added_cb.call, block, False) + self._defer_call( + block.error_cb.call, block, msg) + try: + self._delete_config(block) + except Exception: + logger.warning( + 'Failed to delete rejected log block id=%d', + id, exc_info=True) else: logger.warning('No LogEntry to assign block to !!!') @@ -907,7 +915,8 @@ def _handle_settings_packet(self, cmd, payload): (block.started_cb, (block, True))) else: - msg = self._err_codes[error_status] + msg = self._err_codes.get( + error_status, 'Unknown error') logger.warning('Error %d when starting id=%d (%s)', error_status, id, msg) if (block is not None and self._is_current_registration( @@ -945,7 +954,8 @@ def _handle_settings_packet(self, cmd, payload): return block._delete_pending = False block.err_no = error_status - msg = self._err_codes[error_status] + msg = self._err_codes.get( + error_status, 'Unknown error') callbacks.append((block.error_cb, (block, msg))) if cmd == CMD_RESET_LOGGING: diff --git a/test/crazyflie/test_log.py b/test/crazyflie/test_log.py index d0638ac7b..bda1bdc4e 100644 --- a/test/crazyflie/test_log.py +++ b/test/crazyflie/test_log.py @@ -312,6 +312,45 @@ def test_delete_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_delete_unknown_error_is_reported_and_retryable(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + errors = [] + config.error_cb.add_callback( + lambda log_config, message: errors.append( + (log_config, message))) + + config.delete() + self._acknowledge(CMD_DELETE_BLOCK, config.id, 0xFF) + + self.assertEqual(0xFF, config.err_no) + self.assertEqual([(config, 'Unknown error')], errors) + self.assertFalse(config._delete_pending) + self.assertIn(config, self.log.log_blocks) + + config.delete() + + self.assertTrue(config._delete_pending) + + def test_delete_base_exception_restores_retryable_state(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_config('config') + self.log.add_config(config) + + class SendAbort(BaseException): + pass + + self.cf.send_packet.side_effect = SendAbort('send aborted') + + with self.assertRaisesRegex(SendAbort, 'send aborted'): + config.delete() + + self.assertFalse(config._delete_pending) + self.assertIn(config, self.log.log_blocks) + def test_create_send_failure_schedules_cleanup(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) @@ -353,6 +392,30 @@ def send_packet(packet, expected_reply): [CMD_CREATE_BLOCK, CMD_START_LOGGING, CMD_DELETE_BLOCK], commands) self.assertTrue(config._delete_pending) + def test_create_base_exception_still_schedules_cleanup(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + commands = [] + + class SendAbort(BaseException): + pass + + def send_packet(packet, expected_reply): + command = packet.data[0] + commands.append(command) + if command == CMD_CREATE_BLOCK: + raise SendAbort('send aborted') + + self.cf.send_packet.side_effect = send_packet + + with self.assertRaisesRegex(SendAbort, 'send aborted'): + config.create() + + self.assertEqual([CMD_CREATE_BLOCK, CMD_DELETE_BLOCK], commands) + self.assertTrue(config._delete_pending) + def test_append_send_failure_deletes_partial_firmware_block(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) @@ -439,7 +502,7 @@ def later_callback(): self.assertEqual(0, len(self.log._deferred_calls)) self.assertFalse(self.log._dispatching_deferred_calls) - def test_create_error_clears_pending_and_reports_config(self): + def test_create_error_schedules_cleanup_and_reports_config(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) config = self._make_toc_config('config') @@ -455,12 +518,57 @@ def test_create_error_clears_pending_and_reports_config(self): self.assertIs(config.pending, True) self._acknowledge(CMD_CREATE_BLOCK, config.id, errno.ENOMEM) - self.assertFalse(config.pending) + self.assertTrue(config.pending) self.assertEqual([(config, False)], added_events) self.assertEqual( [(config, 'No more memory available')], error_events) + self.assertEqual(0, config.id) + self.assertIs(self.cf, config.cf) + self.assertIn(config, self.log.log_blocks) + self.assertTrue(config._delete_pending) + self.assertEqual((1, 1), self.log._get_active_config_usage()) + self.assertEqual(CMD_DELETE_BLOCK, + self.cf.send_packet.call_args[0][0].data[0]) + with self.assertRaises(LogConfigError): + self.log.add_config(config) + + self._acknowledge(CMD_DELETE_BLOCK, config.id) + + self.assertIsNone(config.id) + self.assertIsNone(config.cf) + self.assertNotIn(config, self.log.log_blocks) self.assertEqual((0, 0), self.log._get_active_config_usage()) + self.log.add_config(config) + + self.assertEqual(1, config.id) + + def test_create_error_callbacks_survive_cleanup_base_exception(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + config.create() + added_events = [] + error_events = [] + config.added_cb.add_callback( + lambda log_config, added: added_events.append(added)) + config.error_cb.add_callback( + lambda log_config, message: error_events.append(message)) + + class SendAbort(BaseException): + pass + + self.cf.send_packet.side_effect = SendAbort('send aborted') + + with self.assertRaisesRegex(SendAbort, 'send aborted'): + self._acknowledge(CMD_CREATE_BLOCK, config.id, errno.ENOMEM) + + self.assertEqual([False], added_events) + self.assertEqual(['No more memory available'], error_events) + self.assertFalse(config._delete_pending) + self.assertIn(config, self.log.log_blocks) + def test_start_error_reports_config(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) @@ -477,6 +585,23 @@ def test_start_error_reports_config(self): self.assertEqual([(config, False)], started_events) + def test_unknown_start_error_reports_config(self): + self.log.reset() + self._acknowledge(CMD_RESET_LOGGING) + config = self._make_toc_config('config') + self.log.add_config(config) + started_events = [] + config.started_cb.add_callback( + lambda logconf, started: started_events.append( + (logconf, started))) + + config.start() + self._acknowledge(CMD_CREATE_BLOCK, config.id) + self._acknowledge(CMD_START_LOGGING, config.id, 0xFF) + + self.assertEqual(0xFF, config.err_no) + self.assertEqual([(config, False)], started_events) + def test_reset_can_be_retried_after_failed_acknowledgement(self): self.log.reset() self._acknowledge(CMD_RESET_LOGGING) @@ -575,6 +700,22 @@ def test_reset_can_be_retried_when_sending_fails(self): self.assertEqual(2, self.cf.send_packet.call_count) + def test_reset_base_exception_restores_retryable_state(self): + class SendAbort(BaseException): + pass + + self.cf.send_packet.side_effect = SendAbort('send aborted') + + with self.assertRaisesRegex(SendAbort, 'send aborted'): + self.log.reset() + + self.assertFalse(self.log._reset_pending) + + self.cf.send_packet.side_effect = None + self.log.reset() + + self.assertTrue(self.log._reset_pending) + def test_reset_while_disconnected_does_not_block_later_reset(self): self.cf.link = None