From 713b3b2bc94b47b046e4ac00b3fcb07ea234d385 Mon Sep 17 00:00:00 2001 From: "J. Nick Koston" Date: Sun, 16 Aug 2026 20:38:23 -0500 Subject: [PATCH] [ld2420] Simplify after review rounds Model the listen settle window as its own state instead of two flags, pass reception into the state machine as a parameter instead of the rx_seen_ member, consolidate the duplicated button action guards into a shared precondition helper, extract the give-up path from the retry ladder, batch the receive drain reads, rename the command attempt counter, and deduplicate the newest integration tests with a shared log watcher and subscription helpers. --- esphome/components/ld2420/ld2420.cpp | 127 +++---- esphome/components/ld2420/ld2420.h | 13 +- tests/integration/test_uart_mock_ld2420.py | 393 ++++++++++----------- 3 files changed, 257 insertions(+), 276 deletions(-) diff --git a/esphome/components/ld2420/ld2420.cpp b/esphome/components/ld2420/ld2420.cpp index bb767d119b..b471008e41 100644 --- a/esphome/components/ld2420/ld2420.cpp +++ b/esphome/components/ld2420/ld2420.cpp @@ -73,8 +73,6 @@ static constexpr uint8_t CMD_MAX_RETRIES = 3; static constexpr uint32_t STARTUP_LISTEN_TIMEOUT_MS = 10000; static constexpr uint32_t STARTUP_RETRY_LISTEN_MS = 3000; static constexpr uint32_t STARTUP_LISTEN_SETTLE_MS = 500; -static constexpr uint32_t CMD_ACK_TIMEOUT_MS = 1000; -static constexpr uint8_t CMD_MAX_RETRIES = 3; static constexpr uint8_t STARTUP_SEQUENCE_MAX_RETRIES = 3; // Command sets @@ -241,16 +239,17 @@ void LD2420Component::begin_startup_() { } void LD2420Component::begin_listen_() { - this->rx_seen_ = false; - this->listen_drained_ = false; - this->buffer_pos_ = 0; this->phase_start_ms_ = millis(); - this->startup_state_ = StartupState::STARTUP_STATE_LISTEN; + this->startup_state_ = StartupState::STARTUP_STATE_LISTEN_SETTLE; } void LD2420Component::drain_rx_() { - while (this->available()) { - this->read(); + uint8_t buf[MAX_LINE_LENGTH]; + size_t avail; + while ((avail = this->available()) > 0) { + if (!this->read_array(buf, std::min(avail, sizeof(buf)))) { + break; + } } this->buffer_pos_ = 0; } @@ -305,7 +304,7 @@ void LD2420Component::send_startup_cmd_() { void LD2420Component::start_startup_cmd_(StartupState state) { this->startup_state_ = state; - this->startup_cmd_retries_ = 1; + this->startup_cmd_attempts_ = 1; this->send_startup_cmd_(); } @@ -325,8 +324,8 @@ bool LD2420Component::startup_ack_check_() { if (millis() - this->phase_start_ms_ <= CMD_ACK_TIMEOUT_MS) { return false; } - if (this->startup_cmd_retries_ < CMD_MAX_RETRIES) { - this->startup_cmd_retries_++; + if (this->startup_cmd_attempts_ < CMD_MAX_RETRIES) { + this->startup_cmd_attempts_++; ESP_LOGV(TAG, "No reply to startup command %2X; resending", this->startup_cmd_); this->send_startup_cmd_(); return false; @@ -339,8 +338,13 @@ bool LD2420Component::startup_ack_check_() { this->begin_listen_(); return false; } - // Give up on configuration but keep parsing the stream; a module that is - // still streaming keeps publishing sensor data even without a config read. + this->abandon_startup_(); + return false; +} + +// Gives up on configuration but keeps parsing the stream; a module that is +// still streaming keeps publishing sensor data even without a config read. +void LD2420Component::abandon_startup_() { ESP_LOGE(TAG, ESP_LOG_MSG_COMM_FAIL); if (ld2420::get_firmware_int(this->firmware_ver_) == 0) { // Old firmware streams text frames that are only parsed in simple mode; @@ -355,7 +359,6 @@ bool LD2420Component::startup_ack_check_() { #endif this->status_set_warning(ESP_LOG_MSG_COMM_FAIL); this->startup_state_ = StartupState::STARTUP_STATE_RUNNING; - return false; } void LD2420Component::abort_startup_cmd_() { @@ -370,36 +373,36 @@ void LD2420Component::abort_startup_cmd_() { this->write_cmd_frame_(frame); } -void LD2420Component::loop_startup_() { +void LD2420Component::loop_startup_(bool got_data) { switch (this->startup_state_) { - case StartupState::STARTUP_STATE_LISTEN: { - const uint32_t elapsed = millis() - this->phase_start_ms_; + case StartupState::STARTUP_STATE_LISTEN_SETTLE: // Bytes can already be in flight when the listen phase starts: the tail // of a frame the module was transmitting when it was told to restart, // stale data buffered before setup, or the ack to the blind config mode - // exit. Ignore reception during a short settle window (clearing the rx - // flag and the frame parser) so only data the module sends afterwards - // counts as proof that it is up and streaming. The listen_drained_ flag - // guarantees at least one such drain pass even when the main loop - // stalls past the whole window, so bytes that arrived before the listen - // phase can never be mistaken for fresh data. (A full-frame check - // cannot serve as that proof here: old-firmware text frames are only - // recognized once the operating mode is known, which requires the very - // handshake this phase gates.) - if (elapsed < STARTUP_LISTEN_SETTLE_MS || !this->listen_drained_) { - this->drain_rx_(); - this->rx_seen_ = false; - this->listen_drained_ = true; - return; + // exit. Discard everything received during this settle window so only + // data the module sends afterwards counts as proof that it is up and + // streaming. The state runs at least one drain pass even when the main + // loop stalls past the whole window, so bytes that arrived before the + // listen phase can never be mistaken for fresh data. + this->drain_rx_(); + if (millis() - this->phase_start_ms_ >= STARTUP_LISTEN_SETTLE_MS) { + this->phase_start_ms_ = millis(); + this->startup_state_ = StartupState::STARTUP_STATE_LISTEN; } + return; + + case StartupState::STARTUP_STATE_LISTEN: // The module locks up until power cycled if it receives data before it // has sent its first frame after powering on, so wait until it has - // provably transmitted before sending anything. A module stuck in some - // other state stays quiet, so fall through after the listen window. - if (!this->rx_seen_) { + // provably transmitted before sending anything. (A full-frame check + // cannot serve as that proof: old-firmware text frames are only + // recognized once the operating mode is known, which requires the very + // handshake this phase gates.) A module stuck in some other state + // stays quiet, so fall through after the listen window. + if (!got_data) { const uint32_t listen_timeout_ms = this->startup_sequence_retries_ == 0 ? STARTUP_LISTEN_TIMEOUT_MS : STARTUP_RETRY_LISTEN_MS; - if (elapsed < listen_timeout_ms) { + if (millis() - this->phase_start_ms_ < listen_timeout_ms) { return; } ESP_LOGW(TAG, "No data received from the module; attempting configuration anyway"); @@ -408,7 +411,6 @@ void LD2420Component::loop_startup_() { this->drain_rx_(); this->start_startup_cmd_(StartupState::STARTUP_STATE_ENTER_CONFIG); return; - } case StartupState::STARTUP_STATE_ENTER_CONFIG: if (!this->startup_ack_check_()) { @@ -495,15 +497,25 @@ void LD2420Component::loop_startup_() { } } -void LD2420Component::apply_config_action() { +// Common precondition for the button actions: the startup handshake must have +// finished, and actions that write configuration additionally require that the +// configuration was actually read (setup may have given up before the version +// read; writing the unread config to the module's NVM would wipe its stored +// thresholds). +bool LD2420Component::action_allowed_(bool needs_config) { if (this->startup_state_ != StartupState::STARTUP_STATE_RUNNING) { ESP_LOGW(TAG, "Module is still starting up; ignoring"); - return; + return false; } - if (ld2420::get_firmware_int(this->firmware_ver_) == 0) { - // Setup gave up before the configuration was ever read; writing the - // unread config to the module's NVM would wipe its stored thresholds + if (needs_config && ld2420::get_firmware_int(this->firmware_ver_) == 0) { ESP_LOGW(TAG, "Module configuration was never read; ignoring"); + return false; + } + return true; +} + +void LD2420Component::apply_config_action() { + if (!this->action_allowed_(true)) { return; } const uint8_t checksum = calc_checksum(&this->new_config, sizeof(this->new_config)); @@ -539,12 +551,7 @@ void LD2420Component::apply_config_action() { } void LD2420Component::factory_reset_action() { - if (this->startup_state_ != StartupState::STARTUP_STATE_RUNNING) { - ESP_LOGW(TAG, "Module is still starting up; ignoring"); - return; - } - if (ld2420::get_firmware_int(this->firmware_ver_) == 0) { - ESP_LOGW(TAG, "Module configuration was never read; ignoring"); + if (!this->action_allowed_(true)) { return; } ESP_LOGD(TAG, "Setting factory defaults"); @@ -579,8 +586,7 @@ void LD2420Component::factory_reset_action() { } void LD2420Component::restart_module_action() { - if (this->startup_state_ != StartupState::STARTUP_STATE_RUNNING) { - ESP_LOGW(TAG, "Module is still starting up; ignoring"); + if (!this->action_allowed_(false)) { return; } ESP_LOGD(TAG, "Restarting"); @@ -592,8 +598,7 @@ void LD2420Component::restart_module_action() { } void LD2420Component::revert_config_action() { - if (this->startup_state_ != StartupState::STARTUP_STATE_RUNNING) { - ESP_LOGW(TAG, "Module is still starting up; ignoring"); + if (!this->action_allowed_(false)) { return; } memcpy(&this->new_config, &this->current_config, sizeof(this->current_config)); @@ -610,10 +615,7 @@ void LD2420Component::loop() { } const bool got_data = this->read_batch_(this->buffer_data_); if (this->startup_state_ != StartupState::STARTUP_STATE_RUNNING) { - if (got_data) { - this->rx_seen_ = true; - } - this->loop_startup_(); + this->loop_startup_(got_data); } } @@ -1070,21 +1072,6 @@ void LD2420Component::build_system_mode_frame_(CmdFrameT &frame, uint16_t mode) ESP_LOGV(TAG, "Sending write system mode command: %2X", frame.command); } -void LD2420Component::build_system_mode_frame_(CmdFrameT &frame, uint16_t mode) { - uint16_t unknown_parm = 0x0000; - frame.data_length = 0; - frame.header = CMD_FRAME_HEADER; - frame.command = CMD_WRITE_SYS_PARAM; - memcpy(&frame.data[frame.data_length], &CMD_SYSTEM_MODE, sizeof(CMD_SYSTEM_MODE)); - frame.data_length += sizeof(CMD_SYSTEM_MODE); - memcpy(&frame.data[frame.data_length], &mode, sizeof(mode)); - frame.data_length += sizeof(mode); - memcpy(&frame.data[frame.data_length], &unknown_parm, sizeof(unknown_parm)); - frame.data_length += sizeof(unknown_parm); - frame.footer = CMD_FRAME_FOOTER; - ESP_LOGV(TAG, "Sending write system mode command: %2X", frame.command); -} - void LD2420Component::set_system_mode(uint16_t mode) { CmdFrameT cmd_frame; this->build_system_mode_frame_(cmd_frame, mode); diff --git a/esphome/components/ld2420/ld2420.h b/esphome/components/ld2420/ld2420.h index 24c2f0b62e..e94013412d 100644 --- a/esphome/components/ld2420/ld2420.h +++ b/esphome/components/ld2420/ld2420.h @@ -157,7 +157,8 @@ class LD2420Component final : public Component, public uart::UARTDevice { // first frame after powering on, so the state machine listens for data from // the module before transmitting anything. enum class StartupState : uint8_t { - STARTUP_STATE_LISTEN = 0, + STARTUP_STATE_LISTEN_SETTLE = 0, + STARTUP_STATE_LISTEN, STARTUP_STATE_ENTER_CONFIG, STARTUP_STATE_READ_LIMITS, STARTUP_STATE_READ_VERSION, @@ -169,11 +170,13 @@ class LD2420Component final : public Component, public uart::UARTDevice { void begin_startup_(); void begin_listen_(); - void loop_startup_(); + void loop_startup_(bool got_data); void start_startup_cmd_(StartupState state); void send_startup_cmd_(); void abort_startup_cmd_(); + void abandon_startup_(); bool startup_ack_check_(); + bool action_allowed_(bool needs_config); void drain_rx_(); void write_cmd_frame_(const CmdFrameT &frame); bool build_startup_frame_(CmdFrameT &frame); @@ -212,13 +215,11 @@ class LD2420Component final : public Component, public uart::UARTDevice { uint16_t system_mode_{0}; // Set to the energy mode default in begin_startup_() uint16_t gate_energy_[TOTAL_GATES]; uint32_t phase_start_ms_{0}; - StartupState startup_state_{StartupState::STARTUP_STATE_LISTEN}; + StartupState startup_state_{StartupState::STARTUP_STATE_LISTEN_SETTLE}; uint8_t startup_cmd_{0}; // Command byte of the in-flight startup command, for ack matching - uint8_t startup_cmd_retries_{0}; + uint8_t startup_cmd_attempts_{0}; uint8_t startup_sequence_retries_{0}; uint8_t startup_gate_{0}; - bool rx_seen_{false}; - bool listen_drained_{false}; uint8_t buffer_pos_{0}; // where to resume processing/populating buffer uint8_t buffer_data_[MAX_LINE_LENGTH]; char firmware_ver_[8]{"v0.0.0"}; diff --git a/tests/integration/test_uart_mock_ld2420.py b/tests/integration/test_uart_mock_ld2420.py index 753032b67b..fecec2b23b 100644 --- a/tests/integration/test_uart_mock_ld2420.py +++ b/tests/integration/test_uart_mock_ld2420.py @@ -55,7 +55,12 @@ from pathlib import Path from aioesphomeapi import ButtonInfo import pytest -from .state_utils import InitialStateHelper, SensorStateCollector, find_entity +from .state_utils import ( + InitialStateHelper, + SensorStateCollector, + find_entity, + require_entity, +) from .types import APIClientConnectedFactory, RunCompiledFunction @@ -190,6 +195,92 @@ async def test_uart_mock_ld2420( ) +SETUP_COMPLETE_LOG = "Module setup complete; firmware v2.0.0" + + +class _LogWatcher: + """Resolves futures when watched substrings appear in device log lines. + + Use as the run_compiled line_callback. watch() returns a future that + resolves once a line containing all given substrings has been seen `count` + times; `after` gates matching on another future being done, and `until` + stops matching once another future is done. collect() gathers every line + containing any of the given substrings into `self.collected`. + """ + + def __init__(self) -> None: + self._loop = asyncio.get_running_loop() + self._watches: list[dict] = [] + self._collect_substrings: tuple[str, ...] = () + self.collected: list[str] = [] + + def watch( + self, + substrings: str | list[str], + *, + count: int = 1, + after: asyncio.Future | None = None, + until: asyncio.Future | None = None, + ) -> asyncio.Future: + subs = [substrings] if isinstance(substrings, str) else substrings + watch = { + "subs": subs, + "count": count, + "after": after, + "until": until, + "future": self._loop.create_future(), + "seen": 0, + } + self._watches.append(watch) + return watch["future"] + + def collect(self, *substrings: str) -> None: + self._collect_substrings = substrings + + def __call__(self, line: str) -> None: + for watch in self._watches: + if watch["future"].done(): + continue + if watch["after"] is not None and not watch["after"].done(): + continue + if watch["until"] is not None and watch["until"].done(): + continue + if all(s in line for s in watch["subs"]): + watch["seen"] += 1 + if watch["seen"] >= watch["count"]: + watch["future"].set_result(True) + if any(s in line for s in self._collect_substrings): + self.collected.append(line) + + +async def _wait_or_fail(awaitable, timeout: float, message) -> None: + """Await with a timeout, translating TimeoutError into pytest.fail. + + `message` may be a string or a zero-argument callable evaluated at + failure time (for messages that embed the current collector state). + """ + try: + await asyncio.wait_for(awaitable, timeout=timeout) + except TimeoutError: + pytest.fail(message() if callable(message) else message) + + +async def _subscribe_and_wait(client, collector: SensorStateCollector | None = None): + """List entities, subscribe states, and wait for the initial state flood.""" + entities, _ = await client.list_entities_services() + if collector is not None: + collector.build_key_mapping(entities) + initial_state_helper = InitialStateHelper(entities) + on_state = collector.on_state if collector is not None else (lambda s: None) + client.subscribe_states(initial_state_helper.on_state_wrapper(on_state)) + await _wait_or_fail( + initial_state_helper.wait_for_initial_states(), + 11.0, + "Timeout waiting for initial states", + ) + return entities + + async def _run_listen_first_test( yaml_config: str, run_compiled: RunCompiledFunction, @@ -219,10 +310,7 @@ async def _run_listen_first_test( rx_seen = True elif "TX " in line and not rx_seen: tx_before_rx = True - if ( - "Module setup complete; firmware v2.0.0" in line - and not setup_complete.done() - ): + if SETUP_COMPLETE_LOG in line and not setup_complete.done(): setup_complete.set_result(True) if ( "marked FAILED" in line @@ -249,50 +337,40 @@ async def _run_listen_first_test( run_compiled(yaml_config, line_callback=line_callback), api_client_connected() as client, ): - entities, _ = await client.list_entities_services() - collector.build_key_mapping(entities) - - initial_state_helper = InitialStateHelper(entities) - client.subscribe_states( - initial_state_helper.on_state_wrapper(collector.on_state) - ) - - try: - await initial_state_helper.wait_for_initial_states() - except TimeoutError: - pytest.fail("Timeout waiting for initial states") + await _subscribe_and_wait(client, collector) # Setup handshake must complete once the module has talked - try: - await asyncio.wait_for(setup_complete, timeout=10.0) - except TimeoutError: - pytest.fail( - "Timeout waiting for 'Module setup complete' log line. " - "The startup state machine did not finish its handshake." - ) + await _wait_or_fail( + setup_complete, + 10.0, + "Timeout waiting for 'Module setup complete' log line. " + "The startup state machine did not finish its handshake.", + ) # Sensor data must flow from the stream - try: - await collector.wait_for_all(timeout=5.0) - except TimeoutError: - pytest.fail( + await _wait_or_fail( + collector.wait_for_all(timeout=5.0), + 6.0, + lambda: ( f"Timeout waiting for sensor data. Received:\n" f" sensor_states: {collector.sensor_states}\n" f" binary_states: {collector.binary_states}" - ) + ), + ) assert collector.sensor_states["moving_distance"][0] == pytest.approx(100.0) assert collector.binary_states["has_target"][0] is True if post_setup_received is not None: - try: - await asyncio.wait_for(post_setup_received, timeout=5.0) - except TimeoutError: - pytest.fail( + await _wait_or_fail( + post_setup_received, + 5.0, + lambda: ( f"Timeout waiting for post-setup frame " f"(distance={post_setup_distance}). Received:\n" f" moving_distance: {collector.sensor_states['moving_distance']}" - ) + ), + ) # The component must never transmit before the module has talked; # real hardware locks up until power cycled if it does. @@ -333,26 +411,10 @@ async def test_uart_mock_ld2420_cmd_retry( api_client_connected: APIClientConnectedFactory, ) -> None: """First config command gets no reply; the resend must recover.""" - loop = asyncio.get_running_loop() - - setup_complete = loop.create_future() - resend_seen = loop.create_future() - failure_lines: list[str] = [] - - def line_callback(line: str) -> None: - if "No reply to startup command" in line and not resend_seen.done(): - resend_seen.set_result(True) - if ( - "Module setup complete; firmware v2.0.0" in line - and not setup_complete.done() - ): - setup_complete.set_result(True) - if ( - "marked FAILED" in line - or "Communication failed" in line - or "Module setup attempt" in line - ): - failure_lines.append(line) + watcher = _LogWatcher() + resend_seen = watcher.watch("No reply to startup command") + setup_complete = watcher.watch(SETUP_COMPLETE_LOG) + watcher.collect("marked FAILED", "Communication failed", "Module setup attempt") collector = SensorStateCollector( sensor_names=["moving_distance"], @@ -360,47 +422,39 @@ async def test_uart_mock_ld2420_cmd_retry( ) async with ( - run_compiled(yaml_config, line_callback=line_callback), + run_compiled(yaml_config, line_callback=watcher), api_client_connected() as client, ): - entities, _ = await client.list_entities_services() - collector.build_key_mapping(entities) - - initial_state_helper = InitialStateHelper(entities) - client.subscribe_states( - initial_state_helper.on_state_wrapper(collector.on_state) - ) - - try: - await initial_state_helper.wait_for_initial_states() - except TimeoutError: - pytest.fail("Timeout waiting for initial states") + await _subscribe_and_wait(client, collector) # The first enable command is ignored, so a resend must happen - try: - await asyncio.wait_for(resend_seen, timeout=10.0) - except TimeoutError: - pytest.fail("Timeout waiting for the startup command resend log line") + await _wait_or_fail( + resend_seen, 10.0, "Timeout waiting for the startup command resend log line" + ) # The resend gets an ack and the handshake completes normally - try: - await asyncio.wait_for(setup_complete, timeout=10.0) - except TimeoutError: - pytest.fail("Timeout waiting for 'Module setup complete' after the resend") + await _wait_or_fail( + setup_complete, + 10.0, + "Timeout waiting for 'Module setup complete' after the resend", + ) - try: - await collector.wait_for_all(timeout=5.0) - except TimeoutError: - pytest.fail( + await _wait_or_fail( + collector.wait_for_all(timeout=5.0), + 6.0, + lambda: ( f"Timeout waiting for sensor data. Received:\n" f" sensor_states: {collector.sensor_states}" - ) + ), + ) assert collector.sensor_states["moving_distance"][0] == pytest.approx(100.0) # A single command resend must not burn a whole sequence retry or # produce any failure log line - assert not failure_lines, f"Unexpected failure log lines: {failure_lines}" + assert not watcher.collected, ( + f"Unexpected failure log lines: {watcher.collected}" + ) @pytest.mark.asyncio @@ -410,32 +464,15 @@ async def test_uart_mock_ld2420_give_up( api_client_connected: APIClientConnectedFactory, ) -> None: """Version read never answers; retries then give-up, stream keeps working.""" - loop = asyncio.get_running_loop() - - sequence_retry_seen = loop.create_future() - give_up_seen = loop.create_future() - parser_alive_after_give_up = loop.create_future() - marked_failed_lines: list[str] = [] - - def line_callback(line: str) -> None: - if "Module setup attempt 1 failed; retrying" in line and ( - not sequence_retry_seen.done() - ): - sequence_retry_seen.set_result(True) - if "Firmware version and operating mode were never read" in line and ( - not give_up_seen.done() - ): - give_up_seen.set_result(True) - # The overflow probe injected at t=22s (after the give-up) makes the - # parser log this warning only if it is still running - if ( - "Max command length exceeded" in line - and give_up_seen.done() - and not parser_alive_after_give_up.done() - ): - parser_alive_after_give_up.set_result(True) - if "marked FAILED" in line: - marked_failed_lines.append(line) + watcher = _LogWatcher() + sequence_retry_seen = watcher.watch("Module setup attempt 1 failed; retrying") + give_up_seen = watcher.watch("Firmware version and operating mode were never read") + # The overflow probe injected at t=22s (after the give-up) makes the + # parser log this warning only if it is still running + parser_alive_after_give_up = watcher.watch( + "Max command length exceeded", after=give_up_seen + ) + watcher.collect("marked FAILED") collector = SensorStateCollector( sensor_names=["moving_distance"], @@ -443,43 +480,28 @@ async def test_uart_mock_ld2420_give_up( ) async with ( - run_compiled(yaml_config, line_callback=line_callback), + run_compiled(yaml_config, line_callback=watcher), api_client_connected() as client, ): - entities, _ = await client.list_entities_services() - collector.build_key_mapping(entities) - - initial_state_helper = InitialStateHelper(entities) - client.subscribe_states( - initial_state_helper.on_state_wrapper(collector.on_state) - ) - - try: - await initial_state_helper.wait_for_initial_states() - except TimeoutError: - pytest.fail("Timeout waiting for initial states") + await _subscribe_and_wait(client, collector) # The version read times out three times, then the sequence retries - try: - await asyncio.wait_for(sequence_retry_seen, timeout=15.0) - except TimeoutError: - pytest.fail("Timeout waiting for the sequence retry log line") + await _wait_or_fail( + sequence_retry_seen, 15.0, "Timeout waiting for the sequence retry log line" + ) # After all sequence retries the component gives up with a warning - try: - await asyncio.wait_for(give_up_seen, timeout=30.0) - except TimeoutError: - pytest.fail("Timeout waiting for the give-up log line") + await _wait_or_fail( + give_up_seen, 30.0, "Timeout waiting for the give-up log line" + ) - # The stream must still be parsed after giving up; the overflow probe - # injected at t=22s only produces its warning if the parser runs - try: - await asyncio.wait_for(parser_alive_after_give_up, timeout=20.0) - except TimeoutError: - pytest.fail( - "No parser activity after the give-up; the stream parser " - "must keep running in the degraded state" - ) + # The stream must still be parsed after giving up + await _wait_or_fail( + parser_alive_after_give_up, + 20.0, + "No parser activity after the give-up; the stream parser " + "must keep running in the degraded state", + ) # The stream published sensor data while the handshake was failing assert pytest.approx(100.0) in collector.sensor_states["moving_distance"], ( @@ -488,8 +510,8 @@ async def test_uart_mock_ld2420_give_up( ) # The whole point of the degraded state: the component keeps running - assert not marked_failed_lines, ( - f"Component was marked failed: {marked_failed_lines}" + assert not watcher.collected, ( + f"Component was marked failed: {watcher.collected}" ) @@ -500,85 +522,56 @@ async def test_uart_mock_ld2420_restart_button( api_client_connected: APIClientConnectedFactory, ) -> None: """Restart action must not transmit into the module's boot window.""" - loop = asyncio.get_running_loop() - - setup_complete_count = 0 - first_setup_complete = loop.create_future() - second_setup_complete = loop.create_future() - restart_seen = False - module_frame_after_restart_seen = False - tx_into_boot_window = False - failure_lines: list[str] = [] - - def line_callback(line: str) -> None: - nonlocal setup_complete_count, restart_seen, module_frame_after_restart_seen - nonlocal tx_into_boot_window - if "Module setup complete; firmware v2.0.0" in line: - setup_complete_count += 1 - if setup_complete_count == 1 and not first_setup_complete.done(): - first_setup_complete.set_result(True) - elif setup_complete_count == 2 and not second_setup_complete.done(): - second_setup_complete.set_result(True) - if "[ld2420" in line and "Restarting" in line: - restart_seen = True - if restart_seen and "RX inject 45 bytes" in line: - # The module's first frame after its simulated 2 s boot - module_frame_after_restart_seen = True - if ( - restart_seen - and not module_frame_after_restart_seen - and "uart_mock" in line - and "TX " in line - and "FF:00:02:00" in line - ): - # Config mode enable transmitted before the module's first - # post-boot frame; on real hardware this locks the module up - tx_into_boot_window = True - if "marked FAILED" in line or "Communication failed" in line: - failure_lines.append(line) + watcher = _LogWatcher() + first_setup_complete = watcher.watch(SETUP_COMPLETE_LOG) + second_setup_complete = watcher.watch(SETUP_COMPLETE_LOG, count=2) + restart_seen = watcher.watch(["[ld2420", "Restarting"]) + # The module's first frame after its simulated 2 s boot + module_frame_after_restart = watcher.watch("RX inject 45 bytes", after=restart_seen) + # Config mode enable transmitted before the module's first post-boot + # frame; on real hardware this locks the module up + tx_into_boot_window = watcher.watch( + ["uart_mock", "TX ", "FF:00:02:00"], + after=restart_seen, + until=module_frame_after_restart, + ) + watcher.collect("marked FAILED", "Communication failed") async with ( - run_compiled(yaml_config, line_callback=line_callback), + run_compiled(yaml_config, line_callback=watcher), api_client_connected() as client, ): - entities, _ = await client.list_entities_services() - - initial_state_helper = InitialStateHelper(entities) - client.subscribe_states(initial_state_helper.on_state_wrapper(lambda s: None)) - - try: - await initial_state_helper.wait_for_initial_states() - except TimeoutError: - pytest.fail("Timeout waiting for initial states") + entities = await _subscribe_and_wait(client) # Wait for the initial startup handshake to finish - try: - await asyncio.wait_for(first_setup_complete, timeout=10.0) - except TimeoutError: - pytest.fail("Timeout waiting for the initial 'Module setup complete'") + await _wait_or_fail( + first_setup_complete, + 10.0, + "Timeout waiting for the initial 'Module setup complete'", + ) # Restart the module; the button automation also injects the in-flight # frame tail immediately and the module's first frame 2 s later - restart_btn = find_entity(entities, "restart_module", ButtonInfo) - assert restart_btn is not None, "Restart Module button not found" + restart_btn = require_entity(entities, "restart_module", ButtonInfo) client.button_command(restart_btn.key) # The handshake must complete again after the module comes back - try: - await asyncio.wait_for(second_setup_complete, timeout=15.0) - except TimeoutError: - pytest.fail( - "Timeout waiting for 'Module setup complete' after the restart. " - "The component did not recover from the module restart." - ) + await _wait_or_fail( + second_setup_complete, + 15.0, + "Timeout waiting for 'Module setup complete' after the restart. " + "The component did not recover from the module restart.", + ) - assert not tx_into_boot_window, ( + assert not tx_into_boot_window.done(), ( "Component transmitted the config handshake into the module's " "boot window after a restart; the in-flight frame tail bytes must " "not count as proof the module is up" ) - assert not failure_lines, f"Unexpected failure log lines: {failure_lines}" + assert not watcher.collected, ( + f"Unexpected failure log lines: {watcher.collected}" + ) @pytest.mark.asyncio