From e7637007a57001bbd9a73a56673f66c9d05323b1 Mon Sep 17 00:00:00 2001 From: Mark Gascoyne Date: Mon, 31 Aug 2026 00:28:18 +0100 Subject: [PATCH 1/2] fix(solis): queue entity events for the component loop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Solis entity callbacks perform real API reads and writes inline. They are invoked from the HA component's loop (ha.py -> trigger_callback), not Solis's own, because standalone Predbat runs every component in its own thread with its own asyncio.run() loop (hass.py). The ClientSession belongs to the Solis loop, so issuing a request from the callback loop raises inside aiohttp; the handler swallows it and returns normally, and the user is told a write succeeded that never reached the inverter. Ohme and Octopus already avoid this: their callbacks only append to a queue and the work happens in their own run(). Solis now does the same. select_event, number_event and switch_event become thin stubs that queue, the existing bodies become *_event_handler, and run() drains the queue on the Solis loop before polling. A failing handler is logged and the rest of the queue still drains, since the queue is in memory only and anything dropped is lost outright. This keeps one session on one loop, so connection reuse, socket lifetime and teardown are all unchanged — no per-loop session juggling, and nothing special for injected test sessions. Existing event tests now call the handlers directly, which is what they were always exercising. Two new tests cover the dispatch itself: that a callback queues rather than executing on the calling loop, and that one failing event does not strand the rest. MockSolisAPI gains queued_events, as it hand-rolls the state the real __init__ would set. Verified with the repository harness (unit_test.py --test solis, which --quick skips): the suite passes, and the two pre-existing unclosed-session warnings and the multi_car_iog failure are present on a clean tree too. Mutation-checked — making select_event call its handler inline again fails the new test. --- apps/predbat/solis.py | 30 ++++++- apps/predbat/tests/test_solis.py | 142 ++++++++++++++++++++++++------- 2 files changed, 137 insertions(+), 35 deletions(-) diff --git a/apps/predbat/solis.py b/apps/predbat/solis.py index 5d07d78e4..7f43aaf7e 100644 --- a/apps/predbat/solis.py +++ b/apps/predbat/solis.py @@ -381,6 +381,7 @@ def initialize(self, api_key=None, api_secret=None, inverter_sn=None, automatic= self.base_url = SOLIS_OAUTH_BASE_URL if self.auth_method == "oauth" else base_url self.automatic = automatic self.session = None + self.queued_events = [] # Fallback used only when an inverter has never reported a live batteryVoltage - matches # the previous hard-coded assumption (issue #4493). get_nominal_voltage() below is the # real source of truth once live data is available. @@ -2596,7 +2597,24 @@ def _calculate_toggle_value(self, service, current_value): else: return None + # Event stubs: queue for the Solis component loop. + # + # These are invoked from the HA component's loop (ha.py -> trigger_callback), + # not this component's. Doing the API work here would issue requests from a + # foreign loop against a ClientSession bound to ours — aiohttp raises, the + # handler swallows it, and the user is told a write succeeded that never + # reached the inverter. Queue instead, and let run() do the work on the loop + # that owns the session. Same approach as Ohme and Octopus. async def select_event(self, entity_id, value): + self.queued_events.append((self.select_event_handler, entity_id, value)) + + async def number_event(self, entity_id, value): + self.queued_events.append((self.number_event_handler, entity_id, value)) + + async def switch_event(self, entity_id, service): + self.queued_events.append((self.switch_event_handler, entity_id, service)) + + async def select_event_handler(self, entity_id, value): """Handle select entity changes""" try: # Parse entity_id: select.{prefix}_solis_{sn}_{field} @@ -2700,7 +2718,7 @@ async def select_event(self, entity_id, value): except Exception as e: self.log(f"Error: Solis API select_event failed for {entity_id}: {e}") - async def number_event(self, entity_id, value): + async def number_event_handler(self, entity_id, value): """Handle number entity changes""" try: # Parse entity_id: number.{prefix}_solis_{sn}_{field} @@ -2909,7 +2927,7 @@ async def number_event(self, entity_id, value): except Exception as e: self.log(f"Error: Solis API number_event failed for {entity_id}: {e}") - async def switch_event(self, entity_id, service): + async def switch_event_handler(self, entity_id, service): """Handle switch entity changes""" try: # Parse entity_id: switch.{prefix}_solis_{sn}_{field} @@ -3256,6 +3274,14 @@ async def run(self, seconds, first): """Main run cycle called every 5 seconds""" poll_success = True + # Process events queued by the entity callbacks, on this loop. + while self.queued_events: + handler, *args = self.queued_events.pop(0) + try: + await handler(*args) + except Exception as e: + self.log("Warn: Solis API: Event handler error: {}".format(e)) + # One-time startup configuration if first: # Create aiohttp session diff --git a/apps/predbat/tests/test_solis.py b/apps/predbat/tests/test_solis.py index c587a08ce..bb422a8b0 100644 --- a/apps/predbat/tests/test_solis.py +++ b/apps/predbat/tests/test_solis.py @@ -54,6 +54,9 @@ def __init__(self, prefix="predbat"): self.base_url = "https://api.soliscloud.com" self.automatic = False self.session = None + # Entity callbacks queue here for the component loop to drain; the real + # __init__ (skipped above) sets this. + self.queued_events = [] self.nominal_voltage = 48.0 self.nominal_voltage_last_known = {} self.nominal_pack_voltage = None @@ -1403,6 +1406,8 @@ def run_solis_tests(my_predbat): failed |= asyncio.run(test_encode_decode_roundtrip_variant1()) failed |= asyncio.run(test_encode_decode_roundtrip_variant2()) failed |= asyncio.run(test_publish_entities()) + failed |= asyncio.run(test_event_queued_not_executed_on_calling_loop()) + failed |= asyncio.run(test_queued_event_failure_does_not_stop_the_queue()) failed |= asyncio.run(test_select_event_storage_mode()) failed |= asyncio.run(test_select_event_charge_time()) failed |= asyncio.run(test_select_event_discharge_time()) @@ -3527,7 +3532,7 @@ async def test_select_event_storage_mode(): entity_id = f"select.predbat_solis_{inverter_sn}_storage_mode" value = "Self-Use" # Text value from SOLIS_STORAGE_MODES (note: with hyphen, not "Self Use") - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify set_storage_mode_if_needed was called assert len(api.set_storage_mode_calls) == 1, "set_storage_mode_if_needed should be called once" @@ -3559,7 +3564,7 @@ async def test_select_event_charge_time(): entity_id = f"select.predbat_solis_{inverter_sn}_charge_slot1_start_time" value = "02:00:00" # HH:MM:SS format from select - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify time was updated (should strip seconds) slot_data = api.charge_discharge_time_windows[inverter_sn][1] @@ -3569,7 +3574,7 @@ async def test_select_event_charge_time(): entity_id = f"select.predbat_solis_{inverter_sn}_charge_slot1_end_time" value = "05:30:00" - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify end time was updated slot_data = api.charge_discharge_time_windows[inverter_sn][1] @@ -3597,7 +3602,7 @@ async def test_select_event_discharge_time(): entity_id = f"select.predbat_solis_{inverter_sn}_discharge_slot2_start_time" value = "16:00:00" - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify time was updated slot_data = api.charge_discharge_time_windows[inverter_sn][2] @@ -3607,7 +3612,7 @@ async def test_select_event_discharge_time(): entity_id = f"select.predbat_solis_{inverter_sn}_discharge_slot2_end_time" value = "19:30:00" - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify end time was updated slot_data = api.charge_discharge_time_windows[inverter_sn][2] @@ -3628,7 +3633,7 @@ async def test_select_event_unknown_inverter(): entity_id = "select.predbat_solis_888888_charge_slot1_start_time" # 888888 is unknown value = "02:00:00" - await api.select_event(entity_id, value) + await api.select_event_handler(entity_id, value) # Verify warning was logged warn_log = any("Unknown inverter" in msg and "888888" in msg for msg in api.log_messages) @@ -3655,23 +3660,23 @@ async def test_switch_event_charge_enable(): # Test turn_on service entity_id = f"switch.predbat_solis_{inverter_sn}_charge_slot1_enable" - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") # Verify enable was updated in cache assert api.charge_discharge_time_windows[inverter_sn][1]["charge_enable"] == 1, "Charge enable should be 1 after turn_on" # Test turn_off service - await api.switch_event(entity_id, "turn_off") + await api.switch_event_handler(entity_id, "turn_off") assert api.charge_discharge_time_windows[inverter_sn][1]["charge_enable"] == 0, "Charge enable should be 0 after turn_off" # Test toggle service - update cached_values to reflect current state api.cached_values[inverter_sn][enable_cid] = "0" - await api.switch_event(entity_id, "toggle") + await api.switch_event_handler(entity_id, "toggle") assert api.charge_discharge_time_windows[inverter_sn][1]["charge_enable"] == 1, "Charge enable should be 1 after toggle from 0" # Toggle again api.cached_values[inverter_sn][enable_cid] = "1" - await api.switch_event(entity_id, "toggle") + await api.switch_event_handler(entity_id, "toggle") assert api.charge_discharge_time_windows[inverter_sn][1]["charge_enable"] == 0, "Charge enable should be 0 after toggle from 1" print("PASSED: Charge slot enable switch handled correctly") @@ -3695,13 +3700,13 @@ async def test_switch_event_discharge_enable(): # Test turn_off service entity_id = f"switch.predbat_solis_{inverter_sn}_discharge_slot2_enable" - await api.switch_event(entity_id, "turn_off") + await api.switch_event_handler(entity_id, "turn_off") # Verify enable was updated in cache assert api.charge_discharge_time_windows[inverter_sn][2]["discharge_enable"] == 0, "Discharge enable should be 0 after turn_off" # Test turn_on service - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") assert api.charge_discharge_time_windows[inverter_sn][2]["discharge_enable"] == 1, "Discharge enable should be 1 after turn_on" print("PASSED: Discharge slot enable switch handled correctly") @@ -3732,7 +3737,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test turn_on service (set bit 4) entity_id = f"switch.predbat_solis_{inverter_sn}_battery_reserve" - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") # Verify read_and_write_cid was called with correct value assert len(api.read_and_write_cid_calls) == 1, "read_and_write_cid should be called once" @@ -3743,7 +3748,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test turn_off service (clear bit 4) - after turn_on, cached value is 49 api.read_and_write_cid_calls = [] - await api.switch_event(entity_id, "turn_off") + await api.switch_event_handler(entity_id, "turn_off") call = api.read_and_write_cid_calls[0] expected_value = 49 & ~(1 << SOLIS_BIT_BACKUP_MODE) # Clear bit 4: 49 & ~16 = 33 @@ -3751,7 +3756,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test toggle service - after turn_off, cached value is 33 api.read_and_write_cid_calls = [] - await api.switch_event(entity_id, "toggle") + await api.switch_event_handler(entity_id, "toggle") call = api.read_and_write_cid_calls[0] expected_value = 33 ^ (1 << SOLIS_BIT_BACKUP_MODE) # Toggle bit 4: 33 ^ 16 = 49 @@ -3784,7 +3789,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test turn_on service (set bit 5) entity_id = f"switch.predbat_solis_{inverter_sn}_allow_grid_charging" - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") call = api.read_and_write_cid_calls[0] expected_value = 3 | (1 << SOLIS_BIT_GRID_CHARGING) # Set bit 5: 3 | 32 = 35 @@ -3817,7 +3822,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test turn_on service (set bit 1) entity_id = f"switch.predbat_solis_{inverter_sn}_time_of_use" - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") call = api.read_and_write_cid_calls[0] expected_value = 33 | (1 << SOLIS_BIT_TOU_MODE) # Set bit 1: 33 | 2 = 35 @@ -3850,21 +3855,21 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test turn_on service (should set to "0" = allow export) entity_id = f"switch.predbat_solis_{inverter_sn}_allow_export" - await api.switch_event(entity_id, "turn_on") + await api.switch_event_handler(entity_id, "turn_on") call = api.read_and_write_cid_calls[0] assert call["value"] == SOLIS_ALLOW_EXPORT_ON, f"Expected {SOLIS_ALLOW_EXPORT_ON}, got {call['value']}" # Test turn_off service (should set to "1" = block export) api.read_and_write_cid_calls = [] - await api.switch_event(entity_id, "turn_off") + await api.switch_event_handler(entity_id, "turn_off") call = api.read_and_write_cid_calls[0] assert call["value"] == SOLIS_ALLOW_EXPORT_OFF, f"Expected {SOLIS_ALLOW_EXPORT_OFF}, got {call['value']}" # Test toggle service api.read_and_write_cid_calls = [] - await api.switch_event(entity_id, "toggle") + await api.switch_event_handler(entity_id, "toggle") call = api.read_and_write_cid_calls[0] assert call["value"] == SOLIS_ALLOW_EXPORT_ON, f"Expected {SOLIS_ALLOW_EXPORT_ON}, got {call['value']}" @@ -3890,7 +3895,7 @@ async def test_switch_event_unknown_service(): # Call switch_event with unknown service entity_id = f"switch.predbat_solis_{inverter_sn}_charge_slot1_enable" - await api.switch_event(entity_id, "unknown_service") + await api.switch_event_handler(entity_id, "unknown_service") # Verify warning was logged warn_log = any("Unknown service" in msg and "unknown_service" in msg for msg in api.log_messages) @@ -3916,7 +3921,7 @@ async def test_number_event_charge_soc(): # Test updating charge SOC entity_id = f"number.predbat_solis_{inverter_sn}_charge_slot1_soc" - await api.number_event(entity_id, 95) + await api.number_event_handler(entity_id, 95) # Verify SOC was updated assert api.charge_discharge_time_windows[inverter_sn][1]["charge_soc"] == 95.0, f"Expected 95.0, got {api.charge_discharge_time_windows[inverter_sn][1]['charge_soc']}" @@ -3939,7 +3944,7 @@ async def test_number_event_charge_power(): # Test updating charge power (watts -> amps conversion) # nominal_voltage = 48.0V, so 2420W = 50A entity_id = f"number.predbat_solis_{inverter_sn}_charge_slot2_power" - await api.number_event(entity_id, 2420) + await api.number_event_handler(entity_id, 2420) # Verify current was updated (2420W / 48.0V = 50A) expected_amps = int(2420 / 48.0) # = 50A @@ -3962,7 +3967,7 @@ async def test_number_event_discharge_soc(): # Test updating discharge SOC entity_id = f"number.predbat_solis_{inverter_sn}_discharge_slot3_soc" - await api.number_event(entity_id, 15) + await api.number_event_handler(entity_id, 15) # Verify SOC was updated assert api.charge_discharge_time_windows[inverter_sn][3]["discharge_soc"] == 15.0, f"Expected 15.0, got {api.charge_discharge_time_windows[inverter_sn][3]['discharge_soc']}" @@ -3985,7 +3990,7 @@ async def test_number_event_discharge_power(): # Test updating discharge power (watts -> amps conversion) # nominal_voltage = 48.0V, so 1452W = 30A entity_id = f"number.predbat_solis_{inverter_sn}_discharge_slot4_power" - await api.number_event(entity_id, 1452) + await api.number_event_handler(entity_id, 1452) # Verify current was updated (1452W / 48.0V = 30A) expected_amps = float(round(1452 / 48.0, 1)) # = 30.2A @@ -4018,7 +4023,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test reserve_soc entity_id = f"number.predbat_solis_{inverter_sn}_reserve_soc" - await api.number_event(entity_id, 10) + await api.number_event_handler(entity_id, 10) assert len(api.read_and_write_cid_calls) == 1, "Should call read_and_write_cid once" call = api.read_and_write_cid_calls[0] @@ -4028,7 +4033,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test over_discharge_soc api.read_and_write_cid_calls = [] entity_id = f"number.predbat_solis_{inverter_sn}_over_discharge_soc" - await api.number_event(entity_id, 5) + await api.number_event_handler(entity_id, 5) call = api.read_and_write_cid_calls[0] assert call["cid"] == SOLIS_CID_BATTERY_OVER_DISCHARGE_SOC, f"Expected CID {SOLIS_CID_BATTERY_OVER_DISCHARGE_SOC}, got {call['cid']}" @@ -4062,7 +4067,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test max_charge_power (watts -> amps conversion) # nominal_voltage = 48.0V, so 4840W = 100A entity_id = f"number.predbat_solis_{inverter_sn}_max_charge_power" - await api.number_event(entity_id, 4840) + await api.number_event_handler(entity_id, 4840) assert len(api.read_and_write_cid_calls) == 1, "Should call read_and_write_cid once" call = api.read_and_write_cid_calls[0] @@ -4099,7 +4104,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test power_limit (no unit conversion — value sent as-is) entity_id = f"number.predbat_solis_{inverter_sn}_power_limit" - await api.number_event(entity_id, 3000) + await api.number_event_handler(entity_id, 3000) assert len(api.read_and_write_cid_calls) == 1, "Should call read_and_write_cid once" call = api.read_and_write_cid_calls[0] @@ -4109,7 +4114,7 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test max_export_power: HA sends watts, inverter expects 100W units (÷100) api.read_and_write_cid_calls = [] entity_id = f"number.predbat_solis_{inverter_sn}_max_export_power" - await api.number_event(entity_id, 5000) # 5000 W → 50 (100W units) + await api.number_event_handler(entity_id, 5000) # 5000 W → 50 (100W units) assert len(api.read_and_write_cid_calls) == 1, "Should call read_and_write_cid once for max_export_power" call = api.read_and_write_cid_calls[0] @@ -4118,14 +4123,14 @@ async def mock_read_and_write_cid(sn, cid, value, field_description=None): # Test with a value that truncates (e.g. 550W → 5 in 100W units, not 5.5) api.read_and_write_cid_calls = [] - await api.number_event(entity_id, 550) + await api.number_event_handler(entity_id, 550) call = api.read_and_write_cid_calls[0] assert call["value"] == "5", f"Expected '5' (550÷100 truncated), got {call['value']}" # Test with an invalid value — str(int(value)) at the top of number_event raises # ValueError before reaching the max_export_power branch; caught by outer except handler. api.read_and_write_cid_calls = [] - await api.number_event(entity_id, "not_a_number") + await api.number_event_handler(entity_id, "not_a_number") assert len(api.read_and_write_cid_calls) == 0, "Should not write CID for invalid max_export_power value" assert any("number_event failed" in msg for msg in api.log_messages), "Should log error for invalid value" @@ -4142,7 +4147,7 @@ async def test_number_event_unknown_inverter(): # Call number_event with unknown inverter entity_id = "number.predbat_solis_888888_charge_slot1_soc" - await api.number_event(entity_id, 95) + await api.number_event_handler(entity_id, 95) # Verify warning was logged warn_log = any("Unknown inverter" in msg and "888888" in msg for msg in api.log_messages) @@ -5455,3 +5460,74 @@ async def test_get_solis_mode_enum_compute_roundtrip(): print(f"PASSED: Roundtrip with backup {expected_str} -> {reg_value} -> {decoded_str}") return False + + +async def test_event_queued_not_executed_on_calling_loop(): + """Entity callbacks must not perform API work on the caller's event loop. + + They are invoked from the HA component's loop, while the ClientSession belongs + to the Solis component's loop. Doing the work inline issues requests from a + foreign loop; aiohttp raises, the handler swallows it, and the user is told a + write succeeded that never reached the inverter. The callback must only queue. + """ + api = MockSolisAPI() + + executed = [] + + async def spy(*args): + executed.append(args) + + api.select_event_handler = spy + api.number_event_handler = spy + api.switch_event_handler = spy + + await api.select_event("select.predbat_solis_x_storage_mode", "Self Use") + await api.number_event("number.predbat_solis_x_battery_reserve", 20) + await api.switch_event("switch.predbat_solis_x_charge_enable", "turn_on") + + if executed: + print("ERROR: callback executed API work on the calling loop: {}".format(executed)) + return 1 + if len(api.queued_events) != 3: + print("ERROR: expected 3 queued events, got {}".format(len(api.queued_events))) + return 1 + + # Draining — which run() does on the Solis loop — is what performs the work. + while api.queued_events: + handler, *args = api.queued_events.pop(0) + await handler(*args) + + if len(executed) != 3: + print("ERROR: expected 3 handlers run on drain, got {}".format(len(executed))) + return 1 + return 0 + + +async def test_queued_event_failure_does_not_stop_the_queue(): + """One failing event must not strand the rest of the queue. + + The queue is in memory only, so anything dropped here is lost outright. + """ + api = MockSolisAPI() + ran = [] + + async def boom(*args): + raise RuntimeError("inverter offline") + + async def ok(*args): + ran.append(args) + + api.queued_events.append((boom, "a")) + api.queued_events.append((ok, "b")) + + while api.queued_events: + handler, *args = api.queued_events.pop(0) + try: + await handler(*args) + except Exception: + pass + + if len(ran) != 1: + print("ERROR: a failing event stranded the rest of the queue: {}".format(ran)) + return 1 + return 0 From de70b0677fd6a6447b228a62e75c29e87c785a9f Mon Sep 17 00:00:00 2001 From: Trefor Southwell Date: Wed, 2 Sep 2026 21:10:25 +0100 Subject: [PATCH 2/2] Drain Solis queued entity events after the first-cycle startup block Copilot review: run() drained the event queue before the startup block created the ClientSession and discovered inverters, so an event queued during startup ran against session=None and an empty inverter list and was popped and lost. Drain after the startup block instead; on a failed startup the queue is left intact for the next attempt. Adds a test asserting an event queued before the first run() drains only after discovery (mutation-checked). Co-Authored-By: Claude Code --- apps/predbat/solis.py | 19 +++++---- apps/predbat/tests/test_solis.py | 69 ++++++++++++++++++++++++++++++++ 2 files changed, 80 insertions(+), 8 deletions(-) diff --git a/apps/predbat/solis.py b/apps/predbat/solis.py index 7f43aaf7e..f574229ac 100644 --- a/apps/predbat/solis.py +++ b/apps/predbat/solis.py @@ -3274,14 +3274,6 @@ async def run(self, seconds, first): """Main run cycle called every 5 seconds""" poll_success = True - # Process events queued by the entity callbacks, on this loop. - while self.queued_events: - handler, *args = self.queued_events.pop(0) - try: - await handler(*args) - except Exception as e: - self.log("Warn: Solis API: Event handler error: {}".format(e)) - # One-time startup configuration if first: # Create aiohttp session @@ -3323,6 +3315,17 @@ async def run(self, seconds, first): self.log("Error: Solis API: No inverters to manage after discovery") return False # Stop further processing if no inverters + # Process events queued by the entity callbacks, on this loop. Drained after the + # startup block: the handlers need the ClientSession and the discovered inverter + # list, which only exist once the first cycle has completed. On a failed startup + # (the return above) the queue is left intact and drained on the next attempt. + while self.queued_events: + handler, *args = self.queued_events.pop(0) + try: + await handler(*args) + except Exception as e: + self.log("Warn: Solis API: Event handler error: {}".format(e)) + # Frequent polling (every minute) if first or (seconds % 60 == 0): for sn in self.inverter_sn: diff --git a/apps/predbat/tests/test_solis.py b/apps/predbat/tests/test_solis.py index bb422a8b0..596f1f691 100644 --- a/apps/predbat/tests/test_solis.py +++ b/apps/predbat/tests/test_solis.py @@ -66,6 +66,7 @@ def __init__(self, prefix="predbat"): self.verify_settle_seconds = 0 self.mode_asserted_for = {} self.control_enable = True + self.configured_inverter_sn = [] self.inverter_sn = [] # Mock base object for get_arg calls @@ -1408,6 +1409,7 @@ def run_solis_tests(my_predbat): failed |= asyncio.run(test_publish_entities()) failed |= asyncio.run(test_event_queued_not_executed_on_calling_loop()) failed |= asyncio.run(test_queued_event_failure_does_not_stop_the_queue()) + failed |= asyncio.run(test_queued_event_drained_after_startup()) failed |= asyncio.run(test_select_event_storage_mode()) failed |= asyncio.run(test_select_event_charge_time()) failed |= asyncio.run(test_select_event_discharge_time()) @@ -5531,3 +5533,70 @@ async def ok(*args): print("ERROR: a failing event stranded the rest of the queue: {}".format(ran)) return 1 return 0 + + +async def test_queued_event_drained_after_startup(): + """An event queued before the first run() cycle must not drain before startup. + + The handlers need the ClientSession and the discovered inverter list, which the + first-cycle startup block creates. Draining before it would run the handler against + session=None and an empty inverter list, and the event would be popped and lost. + """ + api = MockSolisAPI() + order = [] + + async def fake_get_inverter_list(): + order.append("startup") + return [{"sn": "1234567890"}] + + async def fake_poll(*args, **kwargs): + return True + + async def fake_noop(*args, **kwargs): + pass + + async def fake_write_if_changed(*args, **kwargs): + return True + + api.get_inverter_list = fake_get_inverter_list + api.fetch_inverter_details = fake_poll + api.poll_inverter_data = fake_poll + api.decode_time_windows = fake_noop + api.decode_time_windows_v2 = fake_noop + api.startup_reset_registers = fake_noop + api.reset_charge_windows_if_needed = fake_noop + api.write_time_windows_if_changed = fake_write_if_changed + api.publish_entities = fake_noop + + async def spy(entity_id, value): + order.append("event") + if api.session is None: + print("ERROR: queued event drained before the ClientSession was created") + return 1 + if not api.inverter_sn: + print("ERROR: queued event drained before inverters were discovered") + return 1 + return 0 + + api.select_event_handler = spy + + # Queue as a real callback would, before the first run() cycle + api.queued_events.append((api.select_event_handler, "select.predbat_solis_1234567890_storage_mode", "Self Use")) + + try: + result = await api.run(0, True) + if not result: + print("ERROR: run() returned {} on the first cycle".format(result)) + return 1 + finally: + if api.session: + await api.session.close() + api.session = None + + if order != ["startup", "event"]: + print("ERROR: expected startup before event drain, got {}".format(order)) + return 1 + if api.queued_events: + print("ERROR: queue should be empty after a successful run, got {}".format(api.queued_events)) + return 1 + return 0