Merge remote-tracking branch 'upstream/dev' into neutral-ble-client

# Conflicts:
#	esphome/components/ble_device_base/ble_gatt_client.h
#	esphome/components/bluetooth_connection/__init__.py
#	esphome/components/bluetooth_connection/bluetooth_connection.cpp
#	esphome/components/bluetooth_connection/bluetooth_connection.h
#	esphome/components/bluetooth_connection/bluetooth_connection_bluedroid.cpp
#	esphome/components/bluetooth_connection/bluetooth_connection_bluedroid.h
#	esphome/components/bluetooth_connection/bluetooth_connection_hub.cpp
#	esphome/components/bluetooth_connection/bluetooth_connection_hub.h
#	esphome/components/bluetooth_proxy/__init__.py
#	esphome/components/bluetooth_proxy/bluetooth_proxy.cpp
#	esphome/core/defines.h
#	tests/component_tests/bluetooth_proxy/test_platform_gates.py
#	tests/components/ble_device_base/__init__.py
#	tests/components/bluetooth_proxy/test-passive.esp32-c6-idf.yaml
This commit is contained in:
J. Nick Koston
2026-08-10 10:50:00 -05:00
59 changed files with 3478 additions and 634 deletions
@@ -52,7 +52,7 @@ def _prime_cache(yaml_path: Path) -> None:
Mirrors ``esphome compile``: ``read_config`` populates ``CORE.config``,
then ``update_storage_json`` writes both the StorageJSON sidecar and
the ``.validated.yaml`` compiled-config cache.
the ``.validated.json`` compiled-config cache.
"""
CORE.config_path = yaml_path
config = read_config({}, skip_external_update=True)
@@ -1,30 +1,62 @@
"""Tests for the wake button's wakeup_pin requirement in ld6002b."""
"""Tests for the ld6002b validators that reach across platforms.
wake needs a pin on its own hub, apply_area needs a select on its own hub, and
area_config needs both a button and a select on its own hub. Every one of them
is a same-instance check, which is the half that breaks quietly.
"""
from __future__ import annotations
import pytest
from esphome.components.ld6002b.button import CONFIG_SCHEMA, FINAL_VALIDATE_SCHEMA
from esphome.components.ld6002b.button import (
CONFIG_SCHEMA as BUTTON_CONFIG_SCHEMA,
FINAL_VALIDATE_SCHEMA as BUTTON_FINAL_VALIDATE_SCHEMA,
)
from esphome.components.ld6002b.const import CONF_AREA_CONFIG, CONF_Z_MIN
from esphome.components.ld6002b.number import (
CONFIG_SCHEMA as NUMBER_CONFIG_SCHEMA,
FINAL_VALIDATE_SCHEMA as NUMBER_FINAL_VALIDATE_SCHEMA,
)
from esphome.config import Config
import esphome.config_validation as cv
from esphome.const import CONF_ID, CONF_WAKEUP_PIN, PlatformFramework
from esphome.const import (
CONF_AREA_ID,
CONF_BUTTON,
CONF_ID,
CONF_WAKEUP_PIN,
PlatformFramework,
)
from esphome.core import ID
from esphome.types import ConfigType
from tests.component_tests.types import SetCoreConfigCallable
HUB_ID = "ld6002b_hub"
OTHER_HUB_ID = "ld6002b_other"
def _full_config(hub: ConfigType) -> Config:
def _full_config(
hub: ConfigType,
*,
selects: list[ConfigType] | None = None,
buttons: list[ConfigType] | None = None,
) -> Config:
"""A full config carrying one ld6002b hub, as the ID pass leaves it.
final_validate resolves the hub through get_path_for_id, so the declaring
path has to be registered the way validate_config registers it: the path of
the id value itself, whose parent is the hub's own config.
The platform lists are what the cross-platform validators scan, so a test can
say which of them exist and which hub each one names.
"""
full = Config()
full["ld6002b"] = [hub]
full.declare_ids.append((hub[CONF_ID], ["ld6002b", 0, CONF_ID]))
if selects is not None:
full["select"] = selects
if buttons is not None:
full[CONF_BUTTON] = buttons
return full
@@ -44,10 +76,33 @@ def _buttons(**buttons: str) -> ConfigType:
return config
def _select(*, hub_id: str = HUB_ID) -> ConfigType:
"""A select platform config naming area_id on the given hub."""
return {
"ld6002b_id": ID(hub_id, is_declaration=False, type="ld6002b"),
CONF_AREA_ID: {"name": "Area ID"},
}
def _area_numbers(*, hub_id: str = HUB_ID) -> ConfigType:
"""A number platform config carrying one area_config bound."""
return {
"ld6002b_id": ID(hub_id, is_declaration=False, type="ld6002b"),
CONF_AREA_CONFIG: {CONF_Z_MIN: {"name": "Area Z Min"}},
}
def _validated(config: ConfigType) -> ConfigType:
"""Run the button schema, then the final validation the hub is checked in."""
config = CONFIG_SCHEMA(config)
FINAL_VALIDATE_SCHEMA(config)
config = BUTTON_CONFIG_SCHEMA(config)
BUTTON_FINAL_VALIDATE_SCHEMA(config)
return config
def _validated_numbers(config: ConfigType) -> ConfigType:
"""The same two passes for the number platform."""
config = NUMBER_CONFIG_SCHEMA(config)
NUMBER_FINAL_VALIDATE_SCHEMA(config)
return config
@@ -80,3 +135,82 @@ def test_other_buttons_do_not_need_the_pin(
)
_validated(_buttons(get_delay="Get Delay"))
def test_apply_area_without_select_is_rejected(
set_core_config: SetCoreConfigCallable,
) -> None:
"""apply_area sends the staged bounds to whichever area the select names."""
set_core_config(
PlatformFramework.ESP32_IDF, full_config=_full_config(_hub(wakeup_pin=False))
)
with pytest.raises(
cv.Invalid,
match=(
r"^apply_area requires select\.area_id for the same ld6002b instance"
r" @ data\['apply_area'\]$"
),
):
_validated(_buttons(apply_area="Apply Area"))
def test_apply_area_select_on_another_hub_is_rejected(
set_core_config: SetCoreConfigCallable,
) -> None:
"""A select exists, but on a second ld6002b -- which cannot serve this one."""
set_core_config(
PlatformFramework.ESP32_IDF,
full_config=_full_config(
_hub(wakeup_pin=False), selects=[_select(hub_id=OTHER_HUB_ID)]
),
)
with pytest.raises(
cv.Invalid,
match=(
r"^apply_area requires select\.area_id for the same ld6002b instance"
r" @ data\['apply_area'\]$"
),
):
_validated(_buttons(apply_area="Apply Area"))
def test_area_config_without_apply_area_is_rejected(
set_core_config: SetCoreConfigCallable,
) -> None:
"""The six numbers only stage a write; apply_area is what sends it."""
set_core_config(
PlatformFramework.ESP32_IDF,
full_config=_full_config(_hub(wakeup_pin=False), selects=[_select()]),
)
with pytest.raises(
cv.Invalid,
match=(
r"^area_config requires button\.apply_area for the same ld6002b instance"
r" @ data\['area_config'\]$"
),
):
_validated_numbers(_area_numbers())
def test_area_config_without_select_is_rejected(
set_core_config: SetCoreConfigCallable,
) -> None:
"""The validator's other half: the staged bounds also need an area to land in."""
set_core_config(
PlatformFramework.ESP32_IDF,
full_config=_full_config(
_hub(wakeup_pin=False), buttons=[_buttons(apply_area="Apply Area")]
),
)
with pytest.raises(
cv.Invalid,
match=(
r"^area_config requires select\.area_id for the same ld6002b instance"
r" @ data\['area_config'\]$"
),
):
_validated_numbers(_area_numbers())
@@ -6,10 +6,14 @@ def override_manifest(manifest: ComponentManifestOverride) -> None:
# resolve_irk() is compiled only when a sensor configures irk:
# (request_irk_support() emits USE_BLE_DEVICE_IRK). The unit-test build has
# no sensors, so emit the define here to put the real IRK path under test.
# Likewise the scan-response merger (emitted by the split-report trackers)
# and the listener vector it dispatches into (codegen-sized by consumers).
async def to_code_testing(config):
cg.add_define("USE_BLE_DEVICE_IRK")
# The gatt contract test exercises the gated lookup helpers; compile
# their definitions (ble_gatt_client.cpp) into the test build.
cg.add_define("USE_BLE_GATT_CLIENT")
cg.add_define("USE_BLE_SCAN_RESPONSE_MERGER")
cg.add_define("ESPHOME_BLE_DEVICE_BASE_LISTENER_COUNT", 4)
manifest.to_code = to_code_testing
@@ -0,0 +1,183 @@
// The host test build gets this from the manifest override; clang-tidy does not.
#ifndef USE_BLE_SCAN_RESPONSE_MERGER
#define USE_BLE_SCAN_RESPONSE_MERGER
#endif
#include <gtest/gtest.h>
#include <cstdint>
#include <cstring>
#include <vector>
#include "esphome/components/ble_device_base/scan_response_merger.h"
namespace esphome::ble_device_base::testing {
namespace {
// Pins the merge policy three trackers share (ln882h, rp2, bk72xx): slot
// bookkeeping, the same-device reuse path, the table-full fallback, the
// 62-byte truncation, the advertisement-RSSI choice and the raw_only gate.
// Delivery is observed through a real AdvDispatcher: the raw callback sees
// every frame (including raw_only), a listener only the parsed ones.
struct DeliveredFrame {
uint64_t address;
std::vector<uint8_t> data;
int8_t rssi;
};
struct RawCapture {
std::vector<DeliveredFrame> frames;
static void trampoline(void *self, const RawAdvertisement &adv) {
auto *capture = static_cast<RawCapture *>(self);
capture->frames.push_back({adv.address, std::vector<uint8_t>(adv.data, adv.data + adv.data_len), adv.rssi});
}
};
class CountingListener : public ESPBTDeviceListener {
public:
bool parse_device(const ESPBTDevice &device) override {
this->parsed++;
return true; // claimed: keeps the discovered log quiet
}
int parsed{0};
};
class ScanResponseMergerTest : public ::testing::Test {
protected:
void SetUp() override {
this->dispatcher_.set_raw_advertisement_callback({&this->raw_, &RawCapture::trampoline});
this->dispatcher_.register_listener(&this->listener_);
this->merger_.bind(&this->dispatcher_, &this->scan_continuous_, "test");
}
void stash_(const uint8_t (&mac)[6], int8_t rssi, uint8_t data_len, uint8_t fill, uint32_t now = 0) {
std::vector<uint8_t> data(data_len, fill);
this->merger_.stash_adv(mac, rssi, 0, data.data(), data_len, now);
}
void scan_rsp_(const uint8_t (&mac)[6], int8_t rssi, uint8_t data_len, uint8_t fill) {
std::vector<uint8_t> data(data_len, fill);
this->merger_.submit_scan_rsp(mac, rssi, 0, data.data(), data_len);
}
ScanResponseMerger merger_;
AdvDispatcher dispatcher_;
RawCapture raw_;
CountingListener listener_;
bool scan_continuous_{true};
};
constexpr uint8_t MAC_A[6] = {0x01, 0x02, 0x03, 0x04, 0x05, 0x06};
constexpr uint8_t MAC_B[6] = {0x11, 0x12, 0x13, 0x14, 0x15, 0x16};
TEST_F(ScanResponseMergerTest, MatchedPairDeliversOneMergedFrameWithAdvRssi) {
this->stash_(MAC_A, -40, 20, 0xAA);
EXPECT_TRUE(this->raw_.frames.empty()); // held, not delivered
this->scan_rsp_(MAC_A, -70, 10, 0xBB);
ASSERT_EQ(this->raw_.frames.size(), 1u);
const auto &frame = this->raw_.frames[0];
ASSERT_EQ(frame.data.size(), 30u); // adv + response as ONE frame
EXPECT_EQ(frame.data[0], 0xAA);
EXPECT_EQ(frame.data[19], 0xAA);
EXPECT_EQ(frame.data[20], 0xBB);
// The advertisement's RSSI, never the scan response's.
EXPECT_EQ(frame.rssi, -40);
EXPECT_EQ(this->listener_.parsed, 1);
EXPECT_TRUE(this->merger_.empty());
}
TEST_F(ScanResponseMergerTest, ReAdvertisementDeliversHeldFrameAndReusesSlot) {
this->stash_(MAC_A, -40, 20, 0xAA);
this->stash_(MAC_A, -45, 22, 0xCC); // same device again: first frame is delivered
ASSERT_EQ(this->raw_.frames.size(), 1u);
EXPECT_EQ(this->raw_.frames[0].data.size(), 20u);
EXPECT_EQ(this->raw_.frames[0].rssi, -40);
EXPECT_FALSE(this->merger_.empty()); // the second advertisement now holds the slot
this->scan_rsp_(MAC_A, -70, 5, 0xBB);
ASSERT_EQ(this->raw_.frames.size(), 2u);
EXPECT_EQ(this->raw_.frames[1].data.size(), 27u); // 22 + 5, merged from the reused slot
EXPECT_EQ(this->raw_.frames[1].rssi, -45);
}
TEST_F(ScanResponseMergerTest, FullTableDegradesToUnmergedDelivery) {
uint8_t mac[6] = {0x20, 0x00, 0x00, 0x00, 0x00, 0x00};
for (uint8_t i = 0; i < 8; i++) {
mac[5] = i;
this->stash_(mac, -50, 10, i);
}
EXPECT_TRUE(this->raw_.frames.empty()); // 8 slots, all held
mac[5] = 8;
this->stash_(mac, -50, 10, 8); // 9th device: no slot left
ASSERT_EQ(this->raw_.frames.size(), 1u); // delivered immediately, unmerged
EXPECT_EQ(this->raw_.frames[0].data.size(), 10u);
this->merger_.flush(); // the 8 held frames are all still intact
EXPECT_EQ(this->raw_.frames.size(), 9u);
EXPECT_TRUE(this->merger_.empty());
}
TEST_F(ScanResponseMergerTest, MergeTruncatesAtBufferCapacity) {
this->stash_(MAC_A, -40, 31, 0xAA);
this->scan_rsp_(MAC_A, -70, 40, 0xBB); // only 31 bytes of room remain
ASSERT_EQ(this->raw_.frames.size(), 1u);
EXPECT_EQ(this->raw_.frames[0].data.size(), 62u);
EXPECT_EQ(this->raw_.frames[0].data[31], 0xBB);
EXPECT_EQ(this->raw_.frames[0].data[61], 0xBB);
}
TEST_F(ScanResponseMergerTest, UnmatchedScanResponseIsRawOnly) {
this->scan_rsp_(MAC_B, -60, 12, 0xDD);
ASSERT_EQ(this->raw_.frames.size(), 1u); // still forwarded on the raw path
EXPECT_EQ(this->raw_.frames[0].rssi, -60);
EXPECT_EQ(this->listener_.parsed, 0); // but never parsed for listeners
}
TEST_F(ScanResponseMergerTest, AddrTypeIsPartOfTheMatchKey) {
std::vector<uint8_t> adv(20, 0xAA);
this->merger_.stash_adv(MAC_A, -40, /*addr_type=*/0, adv.data(), adv.size(), 0);
std::vector<uint8_t> rsp(10, 0xBB);
this->merger_.submit_scan_rsp(MAC_A, -70, /*addr_type=*/1, rsp.data(), rsp.size());
// Same MAC, different addr_type: no merge — the response goes out raw_only.
ASSERT_EQ(this->raw_.frames.size(), 1u);
EXPECT_EQ(this->raw_.frames[0].data.size(), 10u);
EXPECT_EQ(this->listener_.parsed, 0);
EXPECT_FALSE(this->merger_.empty()); // the advertisement is still held
}
TEST_F(ScanResponseMergerTest, SweepDeliversOnlyPastTheTimeout) {
this->stash_(MAC_A, -40, 20, 0xAA, /*now=*/1000);
this->merger_.sweep(1300); // exactly 300 ms: not yet past the timeout
EXPECT_TRUE(this->raw_.frames.empty());
this->merger_.sweep(1301);
ASSERT_EQ(this->raw_.frames.size(), 1u);
EXPECT_EQ(this->raw_.frames[0].rssi, -40);
EXPECT_EQ(this->listener_.parsed, 1); // timeout delivery is a full parse, not raw_only
EXPECT_TRUE(this->merger_.empty());
}
TEST_F(ScanResponseMergerTest, FlushDeliversEverythingImmediately) {
this->stash_(MAC_A, -40, 20, 0xAA, /*now=*/1000);
this->stash_(MAC_B, -50, 15, 0xBB, /*now=*/1000);
this->merger_.flush();
EXPECT_EQ(this->raw_.frames.size(), 2u);
EXPECT_EQ(this->listener_.parsed, 2);
EXPECT_TRUE(this->merger_.empty());
}
TEST_F(ScanResponseMergerTest, UnboundMergerDropsInsteadOfCrashing) {
ScanResponseMerger unbound;
std::vector<uint8_t> data(20, 0xAA);
unbound.stash_adv(MAC_A, -40, 0, data.data(), data.size(), 0);
unbound.submit_scan_rsp(MAC_A, -70, 0, data.data(), data.size());
unbound.sweep(1000);
unbound.flush(); // no null jump anywhere
EXPECT_TRUE(unbound.empty());
}
} // namespace
} // namespace esphome::ble_device_base::testing
+53
View File
@@ -42,6 +42,32 @@ sensor:
name: Target-3 Dop
cluster_id:
name: Target-3 Cluster
interference_area_0:
x_min:
name: Interference-0 X Min
x_max:
name: Interference-0 X Max
y_min:
name: Interference-0 Y Min
y_max:
name: Interference-0 Y Max
z_min:
name: Interference-0 Z Min
z_max:
name: Interference-0 Z Max
detection_area_0:
x_min:
name: Detection-0 X Min
x_max:
name: Detection-0 X Max
y_min:
name: Detection-0 Y Min
y_max:
name: Detection-0 Y Max
z_min:
name: Detection-0 Z Min
z_max:
name: Detection-0 Z Max
binary_sensor:
- platform: ld6002b
@@ -50,6 +76,8 @@ binary_sensor:
name: Presence
target_1:
name: Target-1 Presence
detection_area_0:
name: Detection Area-0 Presence
text_sensor:
- platform: ld6002b
@@ -70,6 +98,19 @@ number:
name: Z Max
low_power_sleep_time:
name: Low Power Sleep
area_config:
x_min:
name: Area X Min
x_max:
name: Area X Max
y_min:
name: Area Y Min
y_max:
name: Area Y Max
z_min:
name: Area Z Min
z_max:
name: Area Z Max
select:
- platform: ld6002b
@@ -80,6 +121,8 @@ select:
name: Trigger Speed
installation_mode:
name: Installation
area_id:
name: Area ID
switch:
- platform: ld6002b
@@ -94,6 +137,16 @@ switch:
button:
- platform: ld6002b
ld6002b_id: ld6002b_radar
apply_area:
name: Apply Area
auto_interference:
name: Auto Interference
get_areas:
name: Get Areas
clear_interference:
name: Clear Interference
reset_detection_area:
name: Reset Detection
get_delay:
name: Get Delay
get_sensitivity:
+4 -4
View File
@@ -134,7 +134,7 @@ TEST(HeapProbe, QueueingTypicalCommandsIsAllocationFree) {
size_t total = 0;
for (int i = 0; i != n; i++) {
req[2] = static_cast<uint8_t>(i); // distinct start addresses: identical frames would dedup, not enqueue
total += sample([&] { device.send_pdu(req); }).count;
total += sample([&] { device.queue_pdu(req); }).count;
}
printf("HEAPPROBE queue_%d_typical_commands total_allocs=%zu\n", n, total);
EXPECT_EQ(total, 0u);
@@ -151,11 +151,11 @@ TEST(HeapProbe, WriteBehindQueuedReadsAppendsAllocationFree) {
req.assign(read_pdu, read_pdu + sizeof(read_pdu));
for (int i = 0; i != 3; i++) {
req[2] = static_cast<uint8_t>(i); // distinct start addresses: identical frames would dedup, not enqueue
device.send_pdu(req);
device.queue_pdu(req);
}
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
Sample append = sample([&] { device.send_pdu(write_pdu); });
Sample append = sample([&] { device.queue_pdu(write_pdu); });
printf("HEAPPROBE write_append count=%zu bytes=%zu\n", append.count, append.bytes);
EXPECT_EQ(append.count, 0u);
}
@@ -180,7 +180,7 @@ TEST(HeapProbe, ResponseHandlingIsAllocationFreeAfterWarmup) {
const uint8_t small_resp[] = {0x03, 0x04, 0x00, 0x2A, 0x01, 0x00};
auto round_trip = [&](std::span<const uint8_t> response_pdu) {
device.send_pdu(req);
device.queue_pdu(req);
hub.loop(); // transmit; the tx queue is empty during the measured receive below
uart.inject_frame(0x02, response_pdu);
return sample([&] { hub.loop(); }); // receive + parse + match + dispatch
+161 -125
View File
@@ -116,7 +116,7 @@ TEST(ModbusClientHubNoResponse, RetryRequeuesWaitingFrame) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/true);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
ASSERT_EQ(hub.queued_frames(), 1u);
hub.force_send_next();
ASSERT_EQ(hub.queued_frames(), 0u);
@@ -141,7 +141,7 @@ TEST(ModbusClientHubNoResponse, NoRetryDropsWaitingFrame) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/false);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.timeout_waiting();
@@ -157,7 +157,7 @@ TEST(ModbusClientHubNoResponse, DetachedDeviceIsNotNotified) {
NoResponseProbeHub hub;
{
RetryingDevice device(&hub, 0x02, /*retry=*/true);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
// device destructor clears its queue entries, including the waiting frame's device pointer
}
@@ -177,7 +177,7 @@ TEST(ModbusClientHubNoResponse, RetryBehindInterruptedShell) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/true);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
// A frame from the wrong address (0x07, expected 0x02) hits the unexpected-frame branch.
@@ -204,7 +204,7 @@ TEST(ModbusClientHubNoResponse, InterruptedShellDeclinedRetryRetiresOnRelease) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/false);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
const uint8_t stray_pdu[] = {0x03, 0x04, 0x00, 0x2A, 0x01, 0x00};
@@ -229,7 +229,7 @@ TEST(ModbusClientHubNoResponse, MidCallbackClearCancelsRetry) {
NoResponseProbeHub hub;
ClearingRetryDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.timeout_waiting();
@@ -246,9 +246,9 @@ TEST(ModbusClientHubPriority, WritesSendBeforeQueuedReads) {
const uint8_t read_a[] = {0x03, 0x01, 0x00, 0x00, 0x02};
const uint8_t read_b[] = {0x03, 0x02, 0x00, 0x00, 0x02};
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
device.send_pdu(read_a);
device.send_pdu(read_b);
device.send_pdu(write_pdu);
device.queue_pdu(read_a);
device.queue_pdu(read_b);
device.queue_pdu(write_pdu);
ASSERT_EQ(hub.queued_frames(), 3u);
hub.force_send_next();
@@ -266,8 +266,8 @@ TEST(ModbusClientHubPriority, DuplicateQueuedFrameAbsorbedNotDuplicated) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/false);
device.send_pdu(read_pdu());
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
device.queue_pdu(read_pdu());
ASSERT_EQ(hub.queued_frames(), 1u);
EXPECT_EQ(hub.queued(0).pending, 2u); // one entry standing for two accepted requests
@@ -280,9 +280,9 @@ TEST(ModbusClientHubPriority, InFlightDuplicateRunsOnceMore) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/false);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
device.send_pdu(read_pdu()); // duplicate of the waiting frame
device.queue_pdu(read_pdu()); // duplicate of the waiting frame
EXPECT_EQ(hub.queued_frames(), 0u); // not queued twice
EXPECT_EQ(hub.waiting_command().pending, 2u);
@@ -303,9 +303,9 @@ TEST(ModbusClientHubPriority, AbsorbedRequestSurvivesDeviceRetry) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/true);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
device.send_pdu(read_pdu()); // duplicate of the waiting frame -> absorbed
device.queue_pdu(read_pdu()); // duplicate of the waiting frame -> absorbed
ASSERT_EQ(hub.waiting_command().pending, 2u);
hub.timeout_waiting(); // no response; the device requests a retry
@@ -459,8 +459,8 @@ TEST(ModbusClientHubPriority, WritesThenOneShotReadsThenContinuousPolls) {
device.read_holding_registers(0x100, 2, {.continuous = true});
const uint8_t one_shot[] = {0x03, 0x02, 0x00, 0x00, 0x01};
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
device.send_pdu(one_shot);
device.send_pdu(write_pdu);
device.queue_pdu(one_shot);
device.queue_pdu(write_pdu);
ASSERT_EQ(hub.queued_frames(), 3u);
EXPECT_EQ(hub.queued(0).priority(), CommandPriority::CONTINUOUS);
EXPECT_EQ(hub.queued(1).priority(), CommandPriority::READ);
@@ -482,7 +482,7 @@ TEST(ModbusClientHubPriority, ContinuousIgnoredForWrites) {
RetryingDevice device(&hub, 0x02, /*retry=*/false);
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
device.send_pdu(write_pdu, {.continuous = true});
device.queue_pdu(write_pdu, {.continuous = true});
ASSERT_EQ(hub.queued_frames(), 1u);
EXPECT_EQ(hub.queued(0).priority(), CommandPriority::WRITE);
EXPECT_FALSE(hub.queued(0).continuous);
@@ -530,8 +530,8 @@ TEST(ModbusClientHubPriority, DuplicateQueuedWriteRefused) {
SentCountingDevice device(&hub, 0x02);
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
EXPECT_TRUE(device.send_pdu(write_pdu));
EXPECT_FALSE(device.send_pdu(write_pdu)); // duplicate write: refused
EXPECT_TRUE(device.queue_pdu(write_pdu));
EXPECT_FALSE(device.queue_pdu(write_pdu)); // duplicate write: refused
hub.sweep_for_test();
ASSERT_EQ(hub.queued_frames(), 1u);
@@ -547,8 +547,8 @@ TEST(ModbusClientHubPriority, DuplicateCustomFunctionCodeRefused) {
SentCountingDevice device(&hub, 0x02);
const uint8_t custom_pdu[] = {0x41, 0x01, 0x02}; // user-defined function code
EXPECT_TRUE(device.send_pdu(custom_pdu));
EXPECT_FALSE(device.send_pdu(custom_pdu)); // duplicate custom command: refused
EXPECT_TRUE(device.queue_pdu(custom_pdu));
EXPECT_FALSE(device.queue_pdu(custom_pdu)); // duplicate custom command: refused
hub.sweep_for_test();
ASSERT_EQ(hub.queued_frames(), 1u);
@@ -562,8 +562,8 @@ TEST(ModbusClientHubPriority, AnonymousDuplicateDroppedNotPromoted) {
NoResponseProbeHub hub;
const uint8_t read[] = {0x03, 0x01, 0x00, 0x00, 0x02};
hub.send_pdu(0x02, read);
hub.send_pdu(0x02, read); // anonymous duplicate: dropped
hub.queue_pdu(0x02, read);
hub.queue_pdu(0x02, read); // anonymous duplicate: dropped
ASSERT_EQ(hub.queued_frames(), 1u);
EXPECT_EQ(hub.queued(0).pending, 1u); // never absorbed for a null owner
@@ -575,13 +575,13 @@ TEST(ModbusClientHubPriority, RetriedReadGoesBehindFreshReads) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/true);
device.send_pdu(read_pdu());
hub.force_send_next(); // the frame that will time out and retry
device.send_pdu(read_pdu()); // waiting duplicate: absorbed into the waiting entry
device.queue_pdu(read_pdu());
hub.force_send_next(); // the frame that will time out and retry
device.queue_pdu(read_pdu()); // waiting duplicate: absorbed into the waiting entry
const uint8_t fresh_a[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t fresh_b[] = {0x03, 0x00, 0x20, 0x00, 0x01};
device.send_pdu(fresh_a);
device.send_pdu(fresh_b);
device.queue_pdu(fresh_a);
device.queue_pdu(fresh_b);
ASSERT_EQ(hub.queued_frames(), 2u);
hub.timeout_waiting(); // device retries; the entry returns to READY behind the fresh reads
@@ -608,9 +608,9 @@ TEST(ModbusClientHubPriority, AbsorbedDuplicateKeepsPlaceInLine) {
const uint8_t read_a[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t read_b[] = {0x03, 0x00, 0x20, 0x00, 0x01};
device.send_pdu(read_a);
device.send_pdu(read_b);
device.send_pdu(read_a); // duplicate of the older entry: absorbed, place unchanged
device.queue_pdu(read_a);
device.queue_pdu(read_b);
device.queue_pdu(read_a); // duplicate of the older entry: absorbed, place unchanged
ASSERT_EQ(hub.queued_frames(), 2u);
const ModbusDeviceCommand *next = hub.next_ready();
@@ -626,14 +626,14 @@ TEST(ModbusClientHubPriority, RetriedWriteKeepsWritePriorityAndStaysNonRequeueab
RetryingDevice device(&hub, 0x02, /*retry=*/true);
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
device.send_pdu(write_pdu);
device.queue_pdu(write_pdu);
hub.force_send_next();
hub.timeout_waiting(); // no response -> device requests retry -> back to READY
ASSERT_EQ(hub.queued_frames(), 1u);
EXPECT_EQ(hub.queued(0).priority(), CommandPriority::WRITE); // retry preserves the WRITE class
device.send_pdu(write_pdu); // duplicate of the retried write
device.queue_pdu(write_pdu); // duplicate of the retried write
ASSERT_EQ(hub.queued_frames(), 1u); // still not queued twice...
EXPECT_EQ(hub.queued(0).priority(), CommandPriority::WRITE);
hub.sweep_for_test();
@@ -655,7 +655,7 @@ TEST(ModbusClientHubSent, BlockedHubDefersInsteadOfFailing) {
AlwaysBlockedHub hub;
SentCountingDevice device(&hub, 0x02);
EXPECT_TRUE(device.send_pdu(read_pdu()));
EXPECT_TRUE(device.queue_pdu(read_pdu()));
hub.send_next_for_test();
EXPECT_EQ(device.sent_count_, 0);
@@ -673,7 +673,7 @@ TEST(ModbusClientHubSent, FiresOnWireNotOnQueue) {
hub.setup(); // frame timing derives from the baud rate
SentCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
EXPECT_EQ(device.sent_count_, 0); // queued only - nothing sent yet
hub.send_next_for_test();
@@ -703,7 +703,7 @@ TEST(ModbusClientHubSent, SendRejectedAfterDelayLeavesFrameReady) {
hub.setup();
SentCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.send_next_for_test(); // gate passes, send_frame_ rejects on the post-delay re-check
EXPECT_EQ(device.sent_count_, 0); // nothing transmitted
@@ -766,7 +766,7 @@ TEST(ModbusClientHubCallbackCount, SingleReadSingleCallback) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
drain_with_responses(hub, OK_RESPONSE);
EXPECT_EQ(device.data_count_, 1);
@@ -781,8 +781,8 @@ TEST(ModbusClientHubCallbackCount, DuplicateReadExactlyTwoCallbacks) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
device.queue_pdu(read_pdu());
int cycles = drain_with_responses(hub, OK_RESPONSE);
EXPECT_EQ(cycles, 2);
@@ -812,8 +812,8 @@ TEST(ModbusClientHubCallbackCount, ClearFromResponseResolvesDuplicateWithNotSent
NoResponseProbeHub hub;
ClearOnFirstResponseDevice device(&hub, 0x02);
EXPECT_TRUE(device.send_pdu(read_pdu()));
EXPECT_TRUE(device.send_pdu(read_pdu())); // absorbed: one entry, pending 2
EXPECT_TRUE(device.queue_pdu(read_pdu()));
EXPECT_TRUE(device.queue_pdu(read_pdu())); // absorbed: one entry, pending 2
hub.force_send_next();
hub.receive_frame_for_test(0x02, OK_RESPONSE); // response -> on_response -> clear, then sweep
@@ -829,9 +829,9 @@ TEST(ModbusClientHubCallbackCount, TripleReadRefusesTheThird) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
EXPECT_TRUE(device.send_pdu(read_pdu()));
EXPECT_TRUE(device.send_pdu(read_pdu()));
EXPECT_FALSE(device.send_pdu(read_pdu())); // the entry is already at its cap
EXPECT_TRUE(device.queue_pdu(read_pdu()));
EXPECT_TRUE(device.queue_pdu(read_pdu()));
EXPECT_FALSE(device.queue_pdu(read_pdu())); // the entry is already at its cap
hub.sweep_for_test();
EXPECT_EQ(device.not_sent_count_, 0); // refused synchronously, nothing owed
int cycles = drain_with_responses(hub, OK_RESPONSE);
@@ -849,8 +849,8 @@ TEST(ModbusClientHubCallbackCount, DuplicateWriteRefusedWithoutLifecycle) {
DataCountingDevice device(&hub, 0x02);
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
EXPECT_TRUE(device.send_pdu(write_pdu));
EXPECT_FALSE(device.send_pdu(write_pdu));
EXPECT_TRUE(device.queue_pdu(write_pdu));
EXPECT_FALSE(device.queue_pdu(write_pdu));
hub.sweep_for_test();
EXPECT_EQ(device.terminals(), 0); // the accepted write has not resolved; the other never existed
@@ -867,7 +867,7 @@ TEST(ModbusClientHubCallbackCount, ErrorResponseIsSoleTerminal) {
hub.setup();
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.send_next_for_test();
const uint8_t exception_response[] = {0x83, 0x02};
hub.receive_frame_for_test(0x02, exception_response);
@@ -886,7 +886,7 @@ TEST(ModbusClientHubCallbackCount, NoResponseIsSoleTerminalAndNotSentHasNoSent)
hub.setup();
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.send_next_for_test();
hub.timeout_waiting();
EXPECT_EQ(device.no_response_count_, 1);
@@ -896,8 +896,8 @@ TEST(ModbusClientHubCallbackCount, NoResponseIsSoleTerminalAndNotSentHasNoSent)
// Unabsorbable duplicate: the second identical write is refused at the door - no lifecycle, no
// terminal, nothing sent.
const uint8_t write_pdu[] = {0x06, 0x00, 0x10, 0xBE, 0xEF};
EXPECT_TRUE(device.send_pdu(write_pdu));
EXPECT_FALSE(device.send_pdu(write_pdu));
EXPECT_TRUE(device.queue_pdu(write_pdu));
EXPECT_FALSE(device.queue_pdu(write_pdu));
hub.sweep_for_test();
EXPECT_EQ(device.not_sent_count_, 0);
EXPECT_EQ(device.terminals(), 1); // still just the read's timeout
@@ -920,7 +920,7 @@ TEST(ModbusClientHubCallbackCount, RetryLifecyclesEachGetSentAndTerminal) {
DataCountingDevice device(&hub, 0x02);
device.retries_ = 1; // ask for exactly one retry
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.send_next_for_test();
hub.timeout_waiting(); // lifecycle 1: sent + no_response (retry requested -> re-queued)
ASSERT_EQ(hub.queued_frames(), 1u);
@@ -946,13 +946,13 @@ TEST(ModbusClientHubCallbackCount, RetryIsNeverRefusedByFullQueue) {
device.retries_ = 1;
SentCountingDevice filler(&hub, 0x05);
device.send_pdu(read_pdu());
hub.force_send_next(); // waiting
device.send_pdu(read_pdu()); // absorbed: two requests pending
device.queue_pdu(read_pdu());
hub.force_send_next(); // waiting
device.queue_pdu(read_pdu()); // absorbed: two requests pending
// Fill the remaining live capacity with distinct frames.
for (uint16_t i = 0; hub.entries() < MODBUS_TX_BUFFER_SIZE; i++) {
const uint8_t fill[] = {0x03, static_cast<uint8_t>(i >> 8), static_cast<uint8_t>(i & 0xFF), 0x00, 0x01};
filler.send_pdu(fill);
filler.queue_pdu(fill);
}
hub.timeout_waiting(); // retry requested; the entry flips back to READY regardless of capacity
@@ -991,10 +991,10 @@ TEST(ModbusClientHubQueue, SendRawTooShortIsRefusedAtTheDoor) {
NotSentCountingRawDevice device(&hub, 0x02);
#pragma GCC diagnostic push
#pragma GCC diagnostic ignored "-Wdeprecated-declarations"
EXPECT_FALSE(device.send_raw({})); // too short to contain a PDU
device.send_raw({}); // too short to contain a PDU; the deprecated void spelling cannot report it
#pragma GCC diagnostic pop
EXPECT_EQ(device.not_sent_count_, 0); // refusals are returned, never delivered
EXPECT_TRUE(hub.tx_buffer_empty());
EXPECT_EQ(device.not_sent_count_, 0); // refused at the door: no callback delivered
EXPECT_TRUE(hub.tx_buffer_empty()); // the only evidence of the refusal is that nothing queued
}
// A continuous read: every wire transmission pairs one sent with one terminal, ending on the error.
@@ -1040,7 +1040,7 @@ class ChainOnSentDevice : public ModbusClientDevice {
if (!this->chained_) {
this->chained_ = true;
const uint8_t follow[] = {0x03, 0x00, 0x09, 0x00, 0x01}; // read holding 0x0009 x1
this->send_pdu(follow);
this->queue_pdu(follow);
}
}
bool chained_{false};
@@ -1059,9 +1059,9 @@ TEST(ModbusClientHubQueue, ClearAddressQueueNotifiesEveryOwner) {
const uint8_t read_a[] = {0x03, 0x01, 0x00, 0x00, 0x02};
const uint8_t read_b[] = {0x03, 0x02, 0x00, 0x00, 0x02};
const uint8_t read_c[] = {0x03, 0x03, 0x00, 0x00, 0x02};
controller_like.send_pdu(read_a);
bystander_same.send_pdu(read_b);
bystander_other.send_pdu(read_c);
controller_like.queue_pdu(read_a);
bystander_same.queue_pdu(read_b);
bystander_other.queue_pdu(read_c);
ASSERT_EQ(hub.queued_frames(), 3u);
controller_like.clear_tx_queue_for_address();
@@ -1083,8 +1083,8 @@ TEST(ModbusClientHubQueue, ClearAddressDeliversOneTerminalPerAcceptedRequest) {
SentCountingDevice device(&hub, 0x02);
const uint8_t read[] = {0x03, 0x01, 0x00, 0x00, 0x02};
device.send_pdu(read);
device.send_pdu(read); // duplicate: absorbed into the queued entry
device.queue_pdu(read);
device.queue_pdu(read); // duplicate: absorbed into the queued entry
ASSERT_EQ(hub.queued_frames(), 1u);
ASSERT_EQ(hub.queued(0).pending, 2u);
@@ -1103,8 +1103,8 @@ TEST(ModbusClientHubQueue, ClearSentOnInFlightDuplicateStillNotifiesTheDuplicate
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.send_pdu(read_pdu()); // absorbed: one entry, pending 2
device.queue_pdu(read_pdu());
device.queue_pdu(read_pdu()); // absorbed: one entry, pending 2
ASSERT_EQ(hub.queued(0).pending, 2u);
hub.force_send_next(); // the frame is sent (WAITING); pending still 2
ASSERT_TRUE(hub.waiting());
@@ -1122,7 +1122,7 @@ TEST(ModbusClientHubQueue, ClearWhileInFlightStillDeliversTheResponse) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next(); // sent, now WAITING
ASSERT_TRUE(hub.waiting());
@@ -1151,7 +1151,7 @@ class ResendOnNotSentDevice : public ModbusClientDevice {
this->not_sent_count_++;
if (this->not_sent_count_ == 1) {
const uint8_t again[] = {0x06, 0x00, 0x40, 0x00, 0x01}; // a write: ranked first at selection, not by position
this->send_pdu(again);
this->queue_pdu(again);
}
}
int not_sent_count_{0};
@@ -1165,7 +1165,7 @@ TEST(ModbusClientHubQueue, ClearAddressReentrantResendSurvives) {
ResendOnNotSentDevice device(&hub, 0x02);
const uint8_t read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
device.send_pdu(read);
device.queue_pdu(read);
ASSERT_EQ(hub.queued_frames(), 1u);
hub.clear_tx_queue_for_address(0x02);
@@ -1187,8 +1187,8 @@ TEST(ModbusClientHubQueue, ClearAddressReentrantResendNotSwept) {
const uint8_t read_victim[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t read_other[] = {0x03, 0x00, 0x20, 0x00, 0x01};
resender.send_pdu(read_victim);
bystander_other.send_pdu(read_other);
resender.queue_pdu(read_victim);
bystander_other.queue_pdu(read_other);
ASSERT_EQ(hub.queued_frames(), 2u);
hub.clear_tx_queue_for_address(0x02);
@@ -1212,13 +1212,14 @@ class AlwaysResendDevice : public ModbusClientDevice {
void on_not_sent(std::span<const uint8_t> request_pdu) override {
this->not_sent_count_++;
const uint8_t again[] = {0x03, 0x00, 0x50, 0x00, 0x01};
this->send_pdu(again);
this->queue_pdu(again);
}
int not_sent_count_{0};
};
// From inside on_not_sent, clears ANOTHER address - those victims must still be notified (the per-device
// guard suppresses deliveries only to a device already inside its own on_not_sent()).
// From inside on_not_sent, clears ANOTHER address - those victims must still be notified. Nothing
// suppresses that: a re-entrant clear only flips states, retire() is a no-op on an already-retired
// entry, and each entry still owes one notification per un-run request until pending reaches zero.
class ClearOtherOnNotSentDevice : public ModbusClientDevice {
public:
ClearOtherOnNotSentDevice(ModbusClientHub *hub, uint8_t address) : ModbusClientDevice(hub, address) {}
@@ -1238,9 +1239,9 @@ TEST(ModbusClientHubQueue, PendingNeverExceedsTheServableCap) {
AlwaysResendDevice device(&hub, 0x02);
const uint8_t read[] = {0x03, 0x00, 0x50, 0x00, 0x01};
EXPECT_TRUE(device.send_pdu(read));
EXPECT_TRUE(device.send_pdu(read));
EXPECT_FALSE(device.send_pdu(read)); // at the cap: refused
EXPECT_TRUE(device.queue_pdu(read));
EXPECT_TRUE(device.queue_pdu(read));
EXPECT_FALSE(device.queue_pdu(read)); // at the cap: refused
ASSERT_EQ(hub.queued_frames(), 1u);
EXPECT_EQ(hub.queued(0).pending, 2u);
@@ -1260,12 +1261,12 @@ TEST(ModbusClientHubQueue, FullQueueRefusesWithoutCallbacks) {
// Fill the queue with distinct frames (distinct start addresses keep the dedup from absorbing them).
for (uint16_t i = 0; i < MODBUS_TX_BUFFER_SIZE; i++) {
const uint8_t fill[] = {0x03, static_cast<uint8_t>(i >> 8), static_cast<uint8_t>(i & 0xFF), 0x00, 0x01};
filler.send_pdu(fill);
filler.queue_pdu(fill);
}
ASSERT_EQ(hub.queued_frames(), MODBUS_TX_BUFFER_SIZE);
const uint8_t read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
EXPECT_FALSE(device.send_pdu(read)); // refused synchronously
EXPECT_FALSE(device.queue_pdu(read)); // refused synchronously
hub.sweep_for_test();
EXPECT_EQ(device.not_sent_count_, 0); // nothing was accepted, so nothing is owed
@@ -1298,13 +1299,13 @@ TEST(ModbusClientHubQueue, SelfClearFromNotSentResolvesEveryRequest) {
const uint8_t read_a[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t read_b[] = {0x03, 0x00, 0x20, 0x00, 0x01};
const uint8_t read_c[] = {0x03, 0x00, 0x30, 0x00, 0x01};
clearer.send_pdu(read_a);
clearer.send_pdu(read_b);
bystander.send_pdu(read_c);
clearer.queue_pdu(read_a);
clearer.queue_pdu(read_b);
bystander.queue_pdu(read_c);
ASSERT_EQ(hub.queued_frames(), 3u);
EXPECT_FALSE(clearer.send_pdu(std::span<const uint8_t>{})); // empty: refused, no callback
clearer.clear_tx_queue_for_address(); // the clear the handler used to make
EXPECT_FALSE(clearer.queue_pdu(std::span<const uint8_t>{})); // empty: refused, no callback
clearer.clear_tx_queue_for_address(); // the clear the handler used to make
hub.sweep_for_test();
@@ -1323,8 +1324,8 @@ TEST(ModbusClientHubQueue, NestedClearFromNotSentStillNotifiesVictims) {
const uint8_t read_a[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t read_b[] = {0x03, 0x00, 0x20, 0x00, 0x01};
clearer.send_pdu(read_a);
victim.send_pdu(read_b);
clearer.queue_pdu(read_a);
victim.queue_pdu(read_b);
ASSERT_EQ(hub.queued_frames(), 2u);
hub.clear_tx_queue_for_address(0x02); // clearer's on_not_sent clears address 0x03 in turn
@@ -1345,7 +1346,7 @@ class ResendSecondFrameDevice : public ModbusClientDevice {
this->not_sent_count_++;
if (this->not_sent_count_ == 1) {
const uint8_t same_as_r2[] = {0x03, 0x00, 0x22, 0x00, 0x01};
this->send_pdu(same_as_r2);
this->queue_pdu(same_as_r2);
}
}
int not_sent_count_{0};
@@ -1360,8 +1361,8 @@ TEST(ModbusClientHubQueue, SweepDedupSkipsDeletedFrames) {
const uint8_t r1[] = {0x03, 0x00, 0x21, 0x00, 0x01};
const uint8_t r2[] = {0x03, 0x00, 0x22, 0x00, 0x01};
device.send_pdu(r1);
device.send_pdu(r2);
device.queue_pdu(r1);
device.queue_pdu(r2);
ASSERT_EQ(hub.queued_frames(), 2u);
hub.clear_tx_queue_for_address(0x02);
@@ -1383,7 +1384,7 @@ class ResendAndClearOnNotSentDevice : public ModbusClientDevice {
void on_not_sent(std::span<const uint8_t> request_pdu) override {
this->not_sent_count_++;
const uint8_t again[] = {0x03, 0x00, 0x70, 0x00, 0x01};
this->send_pdu(again);
this->queue_pdu(again);
this->clear_tx_queue_for_address();
}
int not_sent_count_{0};
@@ -1397,7 +1398,7 @@ TEST(ModbusClientHubQueue, ResendAndClearFromNotSentCannotExtendTheSweep) {
ResendAndClearOnNotSentDevice device(&hub, 0x02);
const uint8_t read[] = {0x03, 0x00, 0x70, 0x00, 0x01};
device.send_pdu(read);
device.queue_pdu(read);
hub.clear_tx_queue_for_address(0x02);
hub.sweep_for_test();
@@ -1421,8 +1422,8 @@ TEST(ModbusClientHubQueue, ClearDeviceQueueDropsSilently) {
const uint8_t read_a[] = {0x03, 0x01, 0x00, 0x00, 0x02};
const uint8_t read_b[] = {0x03, 0x02, 0x00, 0x00, 0x02};
device.send_pdu(read_a);
device.send_pdu(read_b);
device.queue_pdu(read_a);
device.queue_pdu(read_b);
ASSERT_EQ(hub.queued_frames(), 2u);
device.clear_tx_queue_for_device();
@@ -1431,7 +1432,7 @@ TEST(ModbusClientHubQueue, ClearDeviceQueueDropsSilently) {
EXPECT_EQ(device.not_sent_count_, 0); // silent drop: no terminal callback
}
// A send_pdu() from inside on_sent() enqueues behind the waiting frame rather than sending
// A queue_pdu() from inside on_sent() enqueues behind the waiting frame rather than sending
// immediately or corrupting the waiting transaction.
TEST(ModbusClientHubSent, ReentrantSendFromOnSentQueues) {
NullUART uart;
@@ -1440,7 +1441,7 @@ TEST(ModbusClientHubSent, ReentrantSendFromOnSentQueues) {
hub.setup();
ChainOnSentDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.send_next_for_test(); // first frame is sent -> on_sent chains a follow-up
EXPECT_TRUE(hub.waiting()); // first frame is waiting
@@ -1507,34 +1508,69 @@ class LegacyNameDevice : public ModbusClientDevice {
#pragma GCC diagnostic pop
} // namespace
// send_pdu() was renamed queue_pdu() because the call queues a request rather than transmitting one.
// The old spelling stays for the deprecation window with the signature 2026.7.4 shipped - void, no
// CommandOptions - so a component built against a real release still compiles and still queues.
#pragma GCC diagnostic push
#pragma GCC diagnostic ignored "-Wdeprecated-declarations"
TEST(ModbusClientHubCompat, DeprecatedSendPduStillQueues) {
NoResponseProbeHub hub;
RetryingDevice device(&hub, 0x02, /*retry=*/false);
const uint8_t read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
device.send_pdu(read); // deprecated device spelling: void, as 2026.7.4 shipped it
EXPECT_EQ(hub.queued_frames(), 1u);
// A refusal is invisible to this spelling - no return value and no callback - so the only evidence
// is that nothing was queued. Reporting the refusal is exactly what moving to queue_pdu() buys.
device.send_pdu(std::span<const uint8_t>());
EXPECT_EQ(hub.queued_frames(), 1u);
// The deprecated hub spelling queues the same way, addressed explicitly.
const uint8_t other[] = {0x03, 0x00, 0x20, 0x00, 0x01};
hub.send_pdu(0x03, other, &device);
EXPECT_EQ(hub.queued_frames(), 2u);
// Both frames resolve to the same owner. Drain them in turn: the device-spelling frame first (FIFO),
// then the hub-spelling frame - addressed to 0x03 yet owned by &device, so reaching device's
// on_no_response proves the request routes by owner pointer, not by address.
hub.force_send_next();
hub.timeout_waiting();
EXPECT_EQ(device.no_response_count_, 1); // device-spelling frame (address 0x02)
hub.force_send_next();
hub.timeout_waiting();
EXPECT_EQ(device.no_response_count_, 2); // hub-spelling frame (address 0x03, &device routing)
}
#pragma GCC diagnostic pop
TEST(ModbusClientHubCompat, LegacyCallbackNamesStillForward) {
NoResponseProbeHub hub;
LegacyNameDevice device(&hub, 0x02);
const uint8_t read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
device.send_pdu(read);
device.queue_pdu(read);
hub.force_send_next();
hub.timeout_waiting(); // no reply -> on_no_response -> forwards to on_modbus_no_response
EXPECT_EQ(device.legacy_no_response_, 1);
// A refused send returns false with no callback, so exercise the forward through an accepted
// request instead: a cleared queue entry delivers on_not_sent(), which forwards to the old name.
EXPECT_FALSE(device.send_pdu(std::span<const uint8_t>())); // empty PDU: refused at the door
EXPECT_FALSE(device.queue_pdu(std::span<const uint8_t>())); // empty PDU: refused at the door
EXPECT_EQ(device.legacy_not_sent_, 0);
const uint8_t queued[] = {0x03, 0x00, 0x11, 0x00, 0x01};
EXPECT_TRUE(device.send_pdu(queued));
EXPECT_TRUE(device.queue_pdu(queued));
hub.clear_tx_queue_for_address(0x02);
hub.sweep_for_test();
EXPECT_EQ(device.legacy_not_sent_, 1);
}
// The send_pdu() capacity bound: a PDU larger than MAX_PDU_SIZE would build a frame past the RTU
// The queue_pdu() capacity bound: a PDU larger than MAX_PDU_SIZE would build a frame past the RTU
// 256-byte limit, so it is refused up front - false at the call site, no entry, no callback.
TEST(ModbusClientHub, OversizedPduIsRefusedAtTheDoor) {
NoResponseProbeHub hub;
LegacyNameDevice device(&hub, 0x02);
std::vector<uint8_t> big(MAX_PDU_SIZE + 1, 0x41);
EXPECT_FALSE(device.send_pdu(big));
EXPECT_FALSE(device.queue_pdu(big));
EXPECT_EQ(device.legacy_not_sent_, 0); // refusals are returned, never delivered
EXPECT_TRUE(hub.tx_buffer_empty());
EXPECT_EQ(hub.entries(), 0u);
@@ -1568,7 +1604,7 @@ TEST(ModbusDeviceShim, LegacyCallbacksReceiveTheOldShapes) {
// Read response: on_modbus_data() historically received the payload after the function code and
// the byte-count byte, as an owning vector.
const uint8_t read_req[] = {0x03, 0x00, 0x10, 0x00, 0x02};
device.send_pdu(read_req);
device.queue_pdu(read_req);
hub.force_send_next();
const uint8_t response[] = {0x03, 0x04, 0x00, 0x2A, 0x01, 0x00};
hub.receive_frame_for_test(0x02, response);
@@ -1577,14 +1613,14 @@ TEST(ModbusDeviceShim, LegacyCallbacksReceiveTheOldShapes) {
// Write echo: no byte-count byte, so the payload is everything after the function code.
const uint8_t write_req[] = {0x06, 0x00, 0x10, 0x00, 0x2A};
device.send_pdu(write_req);
device.queue_pdu(write_req);
hub.force_send_next();
hub.receive_frame_for_test(0x02, write_req); // single-write responses echo the request
const std::vector<uint8_t> expected_echo{0x00, 0x10, 0x00, 0x2A};
EXPECT_EQ(device.last_data_, expected_echo);
// Exception response: on_modbus_error() received the masked function code and the exception code.
device.send_pdu(read_req);
device.queue_pdu(read_req);
hub.force_send_next();
const uint8_t error[] = {0x83, 0x02};
hub.receive_frame_for_test(0x02, error);
@@ -1675,9 +1711,9 @@ class ResendOnDataDevice : public ModbusClientDevice {
public:
ResendOnDataDevice(ModbusClientHub *hub, uint8_t address) : ModbusClientDevice(hub, address) {}
void on_response(std::span<const uint8_t> request_pdu, std::span<const uint8_t> response_pdu) override {
this->send_pdu(std::vector<uint8_t>(request_pdu.begin(), request_pdu.end()));
this->queue_pdu(std::vector<uint8_t>(request_pdu.begin(), request_pdu.end()));
}
void send_pdu(const std::vector<uint8_t> &pdu) { ModbusClientDevice::send_pdu(pdu); }
void queue_pdu(const std::vector<uint8_t> &pdu) { ModbusClientDevice::queue_pdu(pdu); }
};
} // namespace
@@ -1704,8 +1740,8 @@ TEST(ModbusClientHubPriority, ExceptionFlaggedDuplicateDroppedNotPromoted) {
SentCountingDevice device(&hub, 0x02);
const uint8_t weird[] = {0x83, 0x01, 0x00, 0x00, 0x02}; // read-shaped but exception-flagged
EXPECT_TRUE(device.send_pdu(weird));
EXPECT_FALSE(device.send_pdu(weird)); // non-requeueable: cap of one, so the duplicate is refused
EXPECT_TRUE(device.queue_pdu(weird));
EXPECT_FALSE(device.queue_pdu(weird)); // non-requeueable: cap of one, so the duplicate is refused
hub.sweep_for_test();
ASSERT_EQ(hub.queued_frames(), 1u);
@@ -1715,7 +1751,7 @@ TEST(ModbusClientHubPriority, ExceptionFlaggedDuplicateDroppedNotPromoted) {
// The write-shaped twin (0x86 masks to WRITE_SINGLE_REGISTER) must not take WRITE-class
// ordering either: exception-flagged codes are excluded from the mutates classification.
const uint8_t weird_write[] = {0x86, 0x00, 0x10, 0xBE, 0xEF};
device.send_pdu(weird_write);
device.queue_pdu(weird_write);
ASSERT_EQ(hub.queued_frames(), 2u);
EXPECT_EQ(hub.queued(1).priority(), CommandPriority::READ); // not WRITE
const ModbusDeviceCommand *next = hub.next_ready();
@@ -1732,7 +1768,7 @@ class ResendInFlightOnNotSentDevice : public ModbusClientDevice {
this->not_sent_count_++;
if (this->not_sent_count_ == 1) {
const uint8_t same_as_waiting[] = {0x03, 0x01, 0x00, 0x00, 0x02}; // == READ_PDU
this->send_pdu(same_as_waiting);
this->queue_pdu(same_as_waiting);
}
}
int not_sent_count_{0};
@@ -1746,10 +1782,10 @@ TEST(ModbusClientHubQueue, SweepResendAfterClearQueuesFreshNotAbsorbedIntoShell)
NoResponseProbeHub hub;
ResendInFlightOnNotSentDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next(); // READ_PDU now waiting
const uint8_t queued_read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
device.send_pdu(queued_read); // a queued frame for the sweep to notify
device.queue_pdu(queued_read); // a queued frame for the sweep to notify
ASSERT_EQ(hub.queued_frames(), 1u);
hub.clear_tx_queue_for_address(0x02);
@@ -1787,7 +1823,7 @@ TEST(ModbusClientHubNoResponse, SelfClearFromNoResponseDoesNotDoubleResolve) {
NoResponseProbeHub hub;
ClearAddressOnNoResponseDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.timeout_waiting();
@@ -1802,8 +1838,8 @@ TEST(ModbusClientHubNoResponse, SelfClearFromNoResponseResolvesTheAbsorbedReques
NoResponseProbeHub hub;
ClearAddressOnNoResponseDevice device(&hub, 0x02);
EXPECT_TRUE(device.send_pdu(read_pdu()));
EXPECT_TRUE(device.send_pdu(read_pdu())); // absorbed: one entry, two requests
EXPECT_TRUE(device.queue_pdu(read_pdu()));
EXPECT_TRUE(device.queue_pdu(read_pdu())); // absorbed: one entry, two requests
hub.force_send_next();
hub.timeout_waiting();
@@ -1818,7 +1854,7 @@ TEST(ModbusClientHubQueue, ClearedShellReleasesTheBusOnLateResponse) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.clear_tx_queue_for_address(0x02);
ASSERT_EQ(hub.waiting_command().state, FrameState::WAITING_RETIRED);
@@ -1836,7 +1872,7 @@ TEST(ModbusClientHubQueue, ClearedShellReleasesTheBusOnTimeout) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.clear_tx_queue_for_address(0x02);
ASSERT_EQ(hub.waiting_command().state, FrameState::WAITING_RETIRED);
@@ -1856,7 +1892,7 @@ TEST(ModbusClientHubQueue, ClearInterruptedFrameGetsNoResponseAtTimeout) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02); // declines the retry (retries_ == 0)
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
const uint8_t stray_pdu[] = {0x03, 0x04, 0x00, 0x2A, 0x01, 0x00};
hub.receive_frame_for_test(0x07, stray_pdu); // wrong address: interrupts the transaction
@@ -1887,7 +1923,7 @@ TEST(ModbusClientHubQueue, InterruptAfterClearStillDistrusts) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
hub.clear_tx_queue_for_address(0x02);
ASSERT_EQ(hub.waiting_command().state, FrameState::WAITING_RETIRED);
@@ -1912,8 +1948,8 @@ TEST(ModbusClientHubQueue, ClearedInFlightDuplicateTimesOutWithoutRerunning) {
NoResponseProbeHub hub;
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.send_pdu(read_pdu()); // absorbed: one entry, pending 2
device.queue_pdu(read_pdu());
device.queue_pdu(read_pdu()); // absorbed: one entry, pending 2
ASSERT_EQ(hub.queued(0).pending, 2u);
hub.force_send_next(); // sent, pending still 2
hub.clear_tx_queue_for_address(0x02);
@@ -1936,9 +1972,9 @@ TEST(ModbusClientHubCallbackCount, AbsorbedRequestRunsAfterErrorResponse) {
hub.setup();
DataCountingDevice device(&hub, 0x02);
device.send_pdu(read_pdu());
device.queue_pdu(read_pdu());
hub.force_send_next();
device.send_pdu(read_pdu()); // waiting duplicate: absorbed
device.queue_pdu(read_pdu()); // waiting duplicate: absorbed
const uint8_t exception_response[] = {0x83, 0x02};
hub.receive_frame_for_test(0x02, exception_response); // error terminal for request 1
@@ -1954,8 +1990,8 @@ TEST(ModbusClientHubPriority, ReadModifyWritesRankAsWrites) {
const uint8_t read[] = {0x03, 0x00, 0x10, 0x00, 0x01};
const uint8_t mask_write[] = {0x16, 0x00, 0x10, 0x00, 0xFF, 0x00, 0x01};
device.send_pdu(read);
device.send_pdu(mask_write);
device.queue_pdu(read);
device.queue_pdu(mask_write);
ASSERT_EQ(hub.queued_frames(), 2u);
const ModbusDeviceCommand *next = hub.next_ready();
@@ -3,11 +3,13 @@
from __future__ import annotations
import ast
from collections.abc import Callable
import importlib.util
import json
from pathlib import Path
import subprocess
import sys
from typing import Any
import pytest
@@ -205,6 +207,47 @@ def test_convert_keys_no_marker_for_non_sensitive_field() -> None:
assert "sensitive_source" not in entry
def _wildcard_validator(value: Any) -> Any:
return value
def test_convert_keys_marker_wrapped_callable_key_normalizes() -> None:
converted: dict = {}
_bls.convert_keys(converted, {cv.Optional(_wildcard_validator): cv.string}, "/root")
config_vars = converted["schema"]["config_vars"]
assert set(config_vars) == {"string"}
assert config_vars["string"]["key"] == "Optional"
assert config_vars["string"]["key_type"] == "_wildcard_validator"
def test_convert_keys_marker_wrapped_callable_beside_fixed_keys() -> None:
converted: dict = {}
_bls.convert_keys(
converted,
{cv.Required("id"): cv.string, cv.Optional(_wildcard_validator): cv.string},
"/root",
)
assert set(converted["schema"]["config_vars"]) == {"id", "string"}
def test_convert_keys_bare_callable_dotted_qualname() -> None:
def make_validator() -> Callable[[Any], Any]:
def validator(value: Any) -> Any:
return value
return validator
converted: dict = {}
_bls.convert_keys(converted, {make_validator(): cv.string}, "/root")
assert converted["key"] == "String"
assert converted["key_type"].endswith("make_validator.<locals>.validator")
assert "at 0x" not in converted["key_type"]
assert set(converted["schema"]["config_vars"]) == {"string"}
# ---------------------------------------------------------------------------
# Regression tests for the lvgl schema dump.
#
+14
View File
@@ -10,6 +10,7 @@ not be part of a unit test suite.
"""
from collections.abc import Generator
import os
from pathlib import Path
import sys
from unittest.mock import Mock, patch
@@ -40,6 +41,19 @@ def fixture_path() -> Path:
return here / "fixtures"
@pytest.fixture
def probe_env() -> dict[str, str]:
"""Environment for running fixture probe scripts as subprocesses.
Running a script file drops the cwd from sys.path, so prepend the
repo root for the child.
"""
python_path = str(package_root)
if ambient := os.environ.get("PYTHONPATH"):
python_path = os.pathsep.join((python_path, ambient))
return os.environ | {"PYTHONPATH": python_path}
@pytest.fixture
def setup_core(tmp_path: Path) -> Path:
"""Set up CORE with test paths."""
@@ -2,12 +2,13 @@
Executed as a subprocess by test_lazy_imports.py: heavy module names come
in on argv, the ones found in sys.modules afterwards go out on stdout.
Covers both fast-path claims: the bundle suffix check in run_esphome reads
BUNDLE_EXTENSION from esphome.const without importing esphome.bundle, and
the real validated-config cache parse, include resolution included, stays
voluptuous free.
Covers three fast-path claims: the bundle suffix check in run_esphome reads
BUNDLE_EXTENSION from esphome.const without importing esphome.bundle, the
validated-config cache parse stays voluptuous free, and the JSON cache
(lambda sentinel included) resolves without pyyaml or esphome.yaml_util.
"""
import json
import os
from pathlib import Path
import sys
@@ -16,7 +17,6 @@ from unittest.mock import patch
from _leak_report import print_leaked_modules
from _storage import make_storage
import yaml
# Everything imported past this point is the code under test; the pop
# below must only drop what the setup itself preloaded, or it would
@@ -24,8 +24,10 @@ import yaml
_FIXTURE_PRELOADED = frozenset(sys.modules)
from esphome import __main__ as main_mod # noqa: E402
from esphome.const import __version__ as ESPHOME_VERSION # noqa: E402
CONFIG_TEXT = "esphome:\n name: t\n"
LAMBDA_BODY = 'ESP_LOGD("t", "x");'
# An ambient data-dir override would relocate the storage tree away
# from the tmp config dir this fixture builds.
@@ -39,13 +41,23 @@ with tempfile.TemporaryDirectory() as _td:
storage_dir = tmp / ".esphome" / "storage"
storage_dir.mkdir(parents=True)
# The cache is a top-level !include so loading it resolves an
# IncludeFile for real on the fast path. The sidecar is written to the
# layout ext_storage_path resolves once run_esphome sets
# CORE.config_path; going through CORE here would be circular.
(storage_dir / "inc.yaml").write_text(CONFIG_TEXT)
cache_path = storage_dir / "test.yaml.validated.yaml"
cache_path.write_text("!include inc.yaml\n")
# The cache carries a lambda sentinel so loading revives a real Lambda
# on the fast path. The sidecar is written to the layout
# ext_storage_path resolves once run_esphome sets CORE.config_path;
# going through CORE here would be circular.
cache_path = storage_dir / "test.yaml.validated.json"
cache_path.write_text(
json.dumps(
{
"v": 1,
"esphome": ESPHOME_VERSION,
"config": {
"esphome": {"name": "t"},
"script": [{"lambda": {"__esphome_lambda__": LAMBDA_BODY}}],
},
}
)
)
os.utime(cache_path) # keep the cache at least as fresh as the source
make_storage().save(storage_dir / "test.yaml.json")
@@ -76,7 +88,13 @@ with tempfile.TemporaryDirectory() as _td:
# asserts so PYTHONOPTIMIZE in the ambient environment can't strip them.
if exit_code != 0:
sys.exit(f"run_esphome exited {exit_code} before dispatching upload")
if dispatched.get("config") != yaml.safe_load(CONFIG_TEXT):
sys.exit(f"cache include did not resolve through the fast path: {dispatched!r}")
config = dispatched.get("config")
if config is None or config.get("esphome") != {"name": "t"}:
sys.exit(f"cache did not resolve through the fast path: {dispatched!r}")
from esphome.core import Lambda
revived = config["script"][0]["lambda"]
if not isinstance(revived, Lambda) or revived.value != LAMBDA_BODY:
sys.exit(f"lambda sentinel did not revive: {revived!r}")
print_leaked_modules()
@@ -0,0 +1,21 @@
"""Report whether setup_log() pulled in colorama, then print a colored line.
Executed as a subprocess by test_log.py because module imports are
process-global: the parent prints ``colorama_loaded=True/False`` plus an
ANSI colored line so the caller can observe whether the codes survive to
the stream. Pass ``--dashboard`` to simulate a dashboard-spawned run.
"""
import sys
from esphome.core import CORE
from esphome.log import setup_log
if "--dashboard" in sys.argv:
CORE.dashboard = True
setup_log()
print(f"colorama_loaded={'colorama' in sys.modules}")
print("\033[31mred\033[0m end")
sys.stdout.flush()
+289 -50
View File
@@ -2,15 +2,20 @@
from __future__ import annotations
from ipaddress import IPv4Address, IPv4Network
import json
import os
from pathlib import Path
from typing import Any
from unittest.mock import patch
from uuid import UUID
import pytest
from esphome import const, yaml_util
from esphome.__main__ import run_esphome
from esphome.compiled_config import (
_LAMBDA_KEY,
compiled_config_path,
load_compiled_config,
save_compiled_config,
@@ -24,30 +29,26 @@ from esphome.const import (
KEY_TARGET_FRAMEWORK,
KEY_TARGET_PLATFORM,
KEY_VARIANT,
Toolchain,
)
from esphome.core import CORE
from esphome.yaml_util import ESPHomeDataBase
from esphome.core import CORE, ID, HexInt, Lambda, MACAddress, TimePeriodMilliseconds
from esphome.util import OrderedDict
_VALIDATED_CONFIG_YAML = """\
esphome:
name: lite_test
friendly_name: Lite Test Device
esp32:
board: nodemcu-32s
logger:
baud_rate: 115200
api:
port: 6053
encryption:
key: 6dGhpcyBpcyBhIHRlc3Q=
ota:
- platform: esphome
port: 3232
password: secret
wifi:
ssid: ssid
use_address: 192.168.1.42
"""
_VALIDATED_CONFIG = {
"esphome": {"name": "lite_test", "friendly_name": "Lite Test Device"},
"esp32": {"board": "nodemcu-32s"},
"logger": {"baud_rate": 115200},
"api": {"port": 6053, "encryption": {"key": "6dGhpcyBpcyBhIHRlc3Q="}},
"ota": [{"platform": "esphome", "port": 3232, "password": "secret"}],
"wifi": {"ssid": "ssid", "use_address": "192.168.1.42"},
}
def _cache_body(config: dict | None = None) -> str:
"""Render the JSON envelope the production save writes."""
return json.dumps(
{"v": 1, "esphome": const.__version__, "config": config or _VALIDATED_CONFIG}
)
def _write_storage(
@@ -79,10 +80,10 @@ def _write_storage(
storage_path.write_text(json.dumps(data), encoding="utf-8")
def _write_cache(cache_path: Path, body: str = _VALIDATED_CONFIG_YAML) -> Path:
def _write_cache(cache_path: Path, body: str | None = None) -> Path:
"""Write the cache file and return it."""
cache_path.parent.mkdir(parents=True, exist_ok=True)
cache_path.write_text(body, encoding="utf-8")
cache_path.write_text(body if body is not None else _cache_body(), encoding="utf-8")
return cache_path
@@ -96,24 +97,28 @@ def _set_cache_mtime(cache_path: Path, yaml_path: Path, *, offset: int) -> None:
@pytest.fixture
def fresh_cache_files(tmp_path: Path) -> Path:
"""YAML + StorageJSON + cache, all consistent and fresh."""
def primed_storage(tmp_path: Path) -> Path:
"""YAML + StorageJSON sidecar, no cache yet."""
yaml_path = tmp_path / "lite_test.yaml"
yaml_path.write_text("esphome:\n name: lite_test\n")
CORE.config_path = yaml_path
storage_dir = tmp_path / ".esphome" / "storage"
_write_storage(storage_dir / "lite_test.yaml.json")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.yaml")
_set_cache_mtime(cache, yaml_path, offset=5)
_write_storage(tmp_path / ".esphome" / "storage" / "lite_test.yaml.json")
return yaml_path
@pytest.fixture
def fresh_cache_files(primed_storage: Path) -> Path:
"""YAML + StorageJSON + cache, all consistent and fresh."""
storage_dir = primed_storage.parent / ".esphome" / "storage"
cache = _write_cache(storage_dir / "lite_test.yaml.validated.json")
_set_cache_mtime(cache, primed_storage, offset=5)
return primed_storage
def test_compiled_config_path_lives_alongside_sidecar(setup_core: Path) -> None:
"""The cache file shape is predictable from the YAML filename."""
path = compiled_config_path("device.yaml")
assert path.name == "device.yaml.validated.yaml"
assert path.name == "device.yaml.validated.json"
assert path.parent.name == "storage"
@@ -126,9 +131,8 @@ def test_load_compiled_config_happy_path(fresh_cache_files: Path) -> None:
assert config[CONF_API]["encryption"]["key"] == "6dGhpcyBpcyBhIHRlc3Q="
assert config["ota"][0]["password"] == "secret"
# The fast path loads without per-node source ranges (the full
# contract lives in test_yaml_util; this checks the flag is wired up).
assert not isinstance(config[CONF_ESPHOME][CONF_NAME], ESPHomeDataBase)
# The fast path loads plain scalars; no per-node source ranges exist.
assert type(config[CONF_ESPHOME][CONF_NAME]) is str
# apply_to_core populated exactly what upload/logs read off CORE.
assert CORE.name == "lite_test"
@@ -147,7 +151,7 @@ def test_load_compiled_config_populates_esp32_variant(tmp_path: Path) -> None:
storage_dir = tmp_path / ".esphome" / "storage"
_write_storage(storage_dir / "lite_test.yaml.json", esp_platform="ESP32S3")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.yaml")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.json")
_set_cache_mtime(cache, yaml_path, offset=5)
assert load_compiled_config(yaml_path) is not None
@@ -168,7 +172,7 @@ def test_load_compiled_config_skips_esp32_block_for_other_platforms(
esp_platform="ESP8266",
core_platform="esp8266",
)
cache = _write_cache(storage_dir / "lite_test.yaml.validated.yaml")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.json")
_set_cache_mtime(cache, yaml_path, offset=5)
assert load_compiled_config(yaml_path) is not None
@@ -185,7 +189,7 @@ def test_load_compiled_config_falls_back(tmp_path: Path, scenario: str) -> None:
yaml_path.write_text("esphome:\n name: lite_test\n")
CORE.config_path = yaml_path
storage_dir = tmp_path / ".esphome" / "storage"
cache_path = storage_dir / "lite_test.yaml.validated.yaml"
cache_path = storage_dir / "lite_test.yaml.validated.json"
sidecar_path = storage_dir / "lite_test.yaml.json"
if scenario == "missing_cache":
@@ -196,7 +200,7 @@ def test_load_compiled_config_falls_back(tmp_path: Path, scenario: str) -> None:
elif scenario == "corrupt_cache":
_write_storage(sidecar_path)
_set_cache_mtime(
_write_cache(cache_path, "not: valid: yaml: ["), yaml_path, offset=5
_write_cache(cache_path, '{"v": 1, "config": {'), yaml_path, offset=5
)
elif scenario == "missing_sidecar":
# Cache fresh + parseable, but no StorageJSON → can't populate CORE.
@@ -205,6 +209,108 @@ def test_load_compiled_config_falls_back(tmp_path: Path, scenario: str) -> None:
assert load_compiled_config(yaml_path) is None
@pytest.mark.parametrize(
"body",
[
pytest.param(
json.dumps(
{"v": 999, "esphome": const.__version__, "config": {"esphome": {}}}
),
id="wrong_version",
),
pytest.param(
json.dumps({"esphome": const.__version__, "config": {"esphome": {}}}),
id="missing_version",
),
pytest.param(
json.dumps({"v": 1, "esphome": "2020.1.0", "config": {"esphome": {}}}),
id="other_esphome_version",
),
pytest.param(
json.dumps({"v": 1, "config": {"esphome": {}}}),
id="missing_esphome_version",
),
pytest.param(
json.dumps(
{
"v": 1,
"esphome": const.__version__,
"config": ["not", "a", "dict"],
}
),
id="non_dict_config",
),
pytest.param(
json.dumps({"v": 1, "esphome": const.__version__}), id="missing_config"
),
pytest.param(json.dumps(["not", "an", "envelope"]), id="non_dict_envelope"),
],
)
def test_load_compiled_config_rejects_bad_envelope(
primed_storage: Path, body: str
) -> None:
"""A foreign or future cache shape falls back instead of half-loading."""
storage_dir = primed_storage.parent / ".esphome" / "storage"
cache = _write_cache(storage_dir / "lite_test.yaml.validated.json", body)
_set_cache_mtime(cache, primed_storage, offset=5)
assert load_compiled_config(primed_storage) is None
def test_load_ignores_legacy_yaml_cache(primed_storage: Path) -> None:
"""A fresh pre-JSON ``.validated.yaml`` alone can't drive the fast path."""
storage_dir = primed_storage.parent / ".esphome" / "storage"
legacy = _write_cache(
storage_dir / "lite_test.yaml.validated.yaml", "esphome:\n name: lite_test\n"
)
_set_cache_mtime(legacy, primed_storage, offset=5)
assert load_compiled_config(primed_storage) is None
def test_save_removes_stale_legacy_yaml_cache(tmp_path: Path) -> None:
"""A successful save leaves only the JSON cache behind."""
CORE.config_path = tmp_path / "lite_test.yaml"
legacy = tmp_path / ".esphome" / "storage" / "lite_test.yaml.validated.yaml"
legacy.parent.mkdir(parents=True, exist_ok=True)
legacy.write_text("esphome:\n name: lite_test\n")
save_compiled_config({"esphome": {"name": "lite_test"}})
assert compiled_config_path("lite_test.yaml").is_file()
assert not legacy.exists()
def test_save_removes_legacy_yaml_even_when_write_fails(tmp_path: Path) -> None:
"""The secret-bearing legacy cache goes away regardless of write outcome."""
CORE.config_path = tmp_path / "lite_test.yaml"
legacy = tmp_path / ".esphome" / "storage" / "lite_test.yaml.validated.yaml"
legacy.parent.mkdir(parents=True, exist_ok=True)
legacy.write_text("esphome:\n name: lite_test\n")
with patch("esphome.compiled_config.write_file", side_effect=RuntimeError("boom")):
save_compiled_config({"esphome": {"name": "lite_test"}})
assert not legacy.exists()
assert not compiled_config_path("lite_test.yaml").exists()
def test_save_warns_when_legacy_cache_unremovable(
tmp_path: Path, caplog: pytest.LogCaptureFixture
) -> None:
"""A secret-bearing legacy file that won't unlink warns; the write proceeds."""
CORE.config_path = tmp_path / "lite_test.yaml"
legacy = tmp_path / ".esphome" / "storage" / "lite_test.yaml.validated.yaml"
legacy.parent.mkdir(parents=True, exist_ok=True)
legacy.mkdir() # unlink() on a directory raises OSError
with caplog.at_level("WARNING", logger="esphome.compiled_config"):
save_compiled_config({"esphome": {"name": "lite_test"}})
assert "legacy validated-config cache" in caplog.text
assert compiled_config_path("lite_test.yaml").is_file()
@pytest.mark.parametrize("command", ["upload", "logs"])
def test_run_esphome_upload_and_logs_use_cache_when_fresh(
command: str,
@@ -258,7 +364,7 @@ def test_run_esphome_upload_does_not_refresh_cache_without_sidecar(
) -> None:
"""Without a StorageJSON sidecar (no compile has run), the fallback
skips the cache write -- load_compiled_config requires the sidecar,
so writing the rendered (secret-resolved) YAML would be inert and
so writing the rendered (secret-resolved) config would be inert and
leak secrets to disk for nothing."""
yaml_path = tmp_path / "lite_test.yaml"
yaml_path.write_text("esphome:\n name: lite_test\n")
@@ -293,7 +399,7 @@ def test_run_esphome_upload_and_logs_refresh_cache_on_fallback(
storage_dir = tmp_path / ".esphome" / "storage"
_write_storage(storage_dir / "lite_test.yaml.json")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.yaml")
cache = _write_cache(storage_dir / "lite_test.yaml.validated.json")
_set_cache_mtime(cache, yaml_path, offset=-60) # stale
fresh_config = {"esphome": {"name": "lite_test"}, "logger": {}}
@@ -386,28 +492,161 @@ def test_run_esphome_compile_does_not_use_cache(fresh_cache_files: Path) -> None
def test_save_compiled_config_writes_cache(tmp_path: Path) -> None:
"""`save_compiled_config` writes the dumped YAML next to the sidecar."""
"""`save_compiled_config` writes the JSON envelope next to the sidecar."""
CORE.config_path = tmp_path / "lite_test.yaml"
save_compiled_config({"esphome": {"name": "lite_test"}, "logger": {}})
cache_path = compiled_config_path("lite_test.yaml")
assert cache_path.is_file()
body = cache_path.read_text()
assert "name: lite_test" in body
assert "logger:" in body
envelope = json.loads(cache_path.read_text())
assert envelope["v"] == 1
assert envelope["esphome"] == const.__version__
assert envelope["config"] == {"esphome": {"name": "lite_test"}, "logger": {}}
def test_save_compiled_config_swallows_dump_errors(
def test_save_compiled_config_swallows_write_errors(
tmp_path: Path, caplog: pytest.LogCaptureFixture
) -> None:
"""Failures during the dump are non-fatal -- a bad cache just means
"""Failures during the write are non-fatal -- a bad cache just means
the next fast path falls back to read_config()."""
CORE.config_path = tmp_path / "lite_test.yaml"
with patch("esphome.yaml_util.dump", side_effect=RuntimeError("boom")):
with patch("esphome.compiled_config.write_file", side_effect=RuntimeError("boom")):
save_compiled_config({"esphome": {"name": "lite_test"}})
assert not compiled_config_path("lite_test.yaml").exists()
def test_save_stringifies_unknown_values(tmp_path: Path) -> None:
"""A type with no dedicated encoding stores its string form."""
class Weird:
def __str__(self) -> str:
return "weird-str"
CORE.config_path = tmp_path / "lite_test.yaml"
save_compiled_config({"esphome": {"name": "lite_test", "weird": Weird()}})
envelope = json.loads(compiled_config_path("lite_test.yaml").read_text())
assert envelope["config"]["esphome"]["weird"] == "weird-str"
def test_save_skips_cache_on_unserializable_key(tmp_path: Path) -> None:
"""A non-basic dict key aborts the write; the fast path falls back."""
CORE.config_path = tmp_path / "lite_test.yaml"
save_compiled_config({"esphome": {("a", "b"): "lite_test"}})
assert not compiled_config_path("lite_test.yaml").exists()
def _normalize(value: Any) -> Any:
"""Make Lambda comparable; everything else compares by value already."""
if isinstance(value, Lambda):
return ("__lambda__", value.value)
if isinstance(value, dict):
return {k: _normalize(v) for k, v in value.items()}
if isinstance(value, (list, tuple)):
return [_normalize(v) for v in value]
return value
def _round_trip_config() -> OrderedDict:
"""A post-validation shaped config exercising every representer type."""
return OrderedDict(
{
"esphome": OrderedDict(
{
"name": "lite_test",
"build_path": Path("/build/lite_test"),
"on_boot": [
OrderedDict(
{
"trigger_id": ID("trigger_1", type="Trigger"),
"then": [{"lambda": Lambda('ESP_LOGD("t", "x");')}],
}
)
],
}
),
"wifi": OrderedDict(
{
"id": ID("wifi_id", type="WiFiComponent"),
"reboot_timeout": TimePeriodMilliseconds(milliseconds=900000),
"use_address": IPv4Address("192.168.1.42"),
"subnet": IPv4Network("192.168.1.0/24"),
"mac": MACAddress(0xDE, 0xAD, 0xBE, 0xEF, 0x00, 0x01),
}
),
"misc": OrderedDict(
{
"uuid": UUID("12345678-1234-5678-1234-567812345678"),
"toolchain": Toolchain.PLATFORMIO,
"hex": HexInt(0x1234),
"levels": (1, 2.5, True, None),
"empty": {},
}
),
}
)
def test_cache_round_trip_matches_yaml_cache(primed_storage: Path) -> None:
"""The JSON cache loads the same tree the YAML cache used to."""
config = _round_trip_config()
save_compiled_config(config)
from_json = load_compiled_config(primed_storage)
assert from_json is not None
yaml_cache = primed_storage.parent / "dumped.yaml"
yaml_cache.write_text(yaml_util.dump(config, show_secrets=True))
from_yaml = yaml_util.load_yaml(
yaml_cache, clear_secrets=False, track_document_range=False
)
assert _normalize(from_json) == _normalize(from_yaml)
def test_lambda_sentinel_round_trips(primed_storage: Path) -> None:
"""A !lambda body comes back as a Lambda with the same source."""
body = 'id(sensor_1).publish_state(42);\nreturn "multi\\nline";'
save_compiled_config(
{
"esphome": {"name": "lite_test"},
"script": [{"then": [{"lambda": Lambda(body)}]}],
}
)
config = load_compiled_config(primed_storage)
assert config is not None
revived = config["script"][0]["then"][0]["lambda"]
assert isinstance(revived, Lambda)
assert revived.value == body
def test_object_hook_requires_exact_shape(primed_storage: Path) -> None:
"""Only the exact one-key string-valued sentinel revives a Lambda."""
storage_dir = primed_storage.parent / ".esphome" / "storage"
config = {
"esphome": {"name": "lite_test"},
"extra_key": {_LAMBDA_KEY: "x", "y": 1},
"non_str": {_LAMBDA_KEY: 5},
}
cache = _write_cache(
storage_dir / "lite_test.yaml.validated.json", _cache_body(config)
)
_set_cache_mtime(cache, primed_storage, offset=5)
loaded = load_compiled_config(primed_storage)
assert loaded is not None
assert loaded["extra_key"] == {_LAMBDA_KEY: "x", "y": 1}
assert loaded["non_str"] == {_LAMBDA_KEY: 5}
def test_int_keys_coerce_to_strings(primed_storage: Path) -> None:
"""Non-str basic keys stringify; validated configs only use string keys."""
save_compiled_config({"esphome": {"name": "lite_test"}, "table": {1: "a", 2: "b"}})
config = load_compiled_config(primed_storage)
assert config is not None
assert config["table"] == {"1": "a", "2": "b"}
def test_load_compiled_config_rejects_wizard_only_sidecar(tmp_path: Path) -> None:
"""A wizard-only sidecar (no compile -- no core_platform / target_platform)
can't drive upload/logs, so the fast path falls back."""
@@ -426,7 +665,7 @@ def test_load_compiled_config_rejects_wizard_only_sidecar(tmp_path: Path) -> Non
'"loaded_integrations": [], "loaded_platforms": [], "no_mdns": false, '
'"framework": null, "core_platform": null}'
)
cache_path = _write_cache(storage_dir / "lite_test.yaml.validated.yaml")
cache_path = _write_cache(storage_dir / "lite_test.yaml.validated.json")
_set_cache_mtime(cache_path, yaml_path, offset=5)
assert load_compiled_config(yaml_path) is None
+27
View File
@@ -1,5 +1,7 @@
import os
from pathlib import Path
import subprocess
import sys
from unittest.mock import patch
from hypothesis import given
@@ -213,6 +215,31 @@ class TestLambda:
assert str(target) is value.value
def test_init__expression_initializer(self):
from esphome.cpp_generator import RawExpression
target = core.Lambda(RawExpression("foo()"))
assert target.value == "foo();"
def test_init__other_initializer(self):
target = core.Lambda(123)
assert target.value == 123
def test_init_from_str_does_not_import_codegen(self):
"""The validated-config cache revives Lambdas on the upload fast path."""
# sys.exit rather than assert so ambient PYTHONOPTIMIZE can't strip it.
check = (
"import sys; from esphome.core import Lambda; "
"Lambda('return 1;'); "
"sys.exit('codegen leaked' if 'esphome.cpp_generator' in sys.modules else 0)"
)
result = subprocess.run(
[sys.executable, "-c", check], capture_output=True, text=True, check=False
)
assert result.returncode == 0, result.stderr
def test_parts(self):
target = core.Lambda(SAMPLE_LAMBDA.strip())
+325
View File
@@ -0,0 +1,325 @@
"""Tests for the Happy Eyeballs urllib3 shim."""
from __future__ import annotations
import asyncio
from collections.abc import Generator
import socket
from typing import Any
from unittest.mock import Mock, patch
import pytest
from esphome.happy_eyeballs import _make_create_connection, ensure_happy_eyeballs
def _addr_info(host: str, port: int) -> tuple[Any, ...]:
"""Build a getaddrinfo-style result tuple for an IPv4 address."""
return (socket.AF_INET, socket.SOCK_STREAM, socket.IPPROTO_TCP, "", (host, port))
@pytest.fixture
def create_connection() -> Any:
"""A freshly built Happy Eyeballs create_connection replacement."""
return _make_create_connection()
@pytest.fixture
def listener() -> Generator[tuple[str, int]]:
"""A listening TCP socket on localhost; yields its address."""
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server.bind(("127.0.0.1", 0))
server.listen(5)
yield server.getsockname()
server.close()
@pytest.fixture
def mock_gai(listener: tuple[str, int]) -> Generator[Any]:
"""Resolve every host to two copies of the listener's address."""
with patch("socket.getaddrinfo", return_value=[_addr_info(*listener)] * 2) as mock:
yield mock
def test_ensure_happy_eyeballs_patches_and_is_idempotent(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""The shim replaces urllib3's create_connection exactly once."""
import urllib3.util.connection
def stock(*args: Any, **kwargs: Any) -> None:
pass
monkeypatch.setattr(urllib3.util.connection, "create_connection", stock)
ensure_happy_eyeballs()
patched = urllib3.util.connection.create_connection
assert patched is not stock
assert patched._esphome_patched
ensure_happy_eyeballs()
assert urllib3.util.connection.create_connection is patched
def test_connects_and_restores_socket_state(
create_connection: Any, listener: tuple[str, int], mock_gai: Any
) -> None:
"""The winning socket comes back blocking, with timeout and options set."""
sock = create_connection(
("example.com", listener[1]),
timeout=5,
socket_options=[(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)],
)
try:
assert sock.getpeername() == listener
assert sock.gettimeout() == 5
assert sock.getsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY) != 0
finally:
sock.close()
def test_single_address_connects(
create_connection: Any, listener: tuple[str, int]
) -> None:
"""A host resolving to one address connects through the same path."""
with patch("socket.getaddrinfo", return_value=[_addr_info(*listener)]):
sock = create_connection(("example.com", listener[1]), timeout=5)
try:
assert sock.getpeername() == listener
finally:
sock.close()
def test_falls_back_to_working_address(
create_connection: Any, listener: tuple[str, int], monkeypatch: pytest.MonkeyPatch
) -> None:
"""An unreachable first address does not block the working one."""
from esphome import happy_eyeballs
# 192.0.2.1 (TEST-NET-1) blackholes or fails fast depending on the
# network; either way the second address must win well within the
# timeout instead of waiting out the first. A short stagger keeps the
# test's duration network independent.
monkeypatch.setattr(happy_eyeballs, "HAPPY_EYEBALLS_DELAY", 0.01)
addr_infos = [_addr_info("192.0.2.1", 9), _addr_info(*listener)]
with patch("socket.getaddrinfo", return_value=addr_infos):
sock = create_connection(("example.com", listener[1]), timeout=10)
try:
assert sock.getpeername() == listener
finally:
sock.close()
def test_bracketed_ipv6_host_is_stripped(
create_connection: Any, listener: tuple[str, int], mock_gai: Any
) -> None:
"""A bracketed IPv6 literal is unbracketed before resolution."""
sock = create_connection(("[::1]", listener[1]), timeout=5)
try:
assert mock_gai.call_args[0][0] == "::1"
assert sock.getpeername() == listener
finally:
sock.close()
def test_source_address_is_bound(
create_connection: Any, listener: tuple[str, int], mock_gai: Any
) -> None:
"""The socket binds to the requested source address before connecting."""
sock = create_connection(
("example.com", listener[1]),
timeout=5,
source_address=("127.0.0.1", 0),
)
try:
assert sock.getsockname()[0] == "127.0.0.1"
finally:
sock.close()
def test_socket_factory_failure_closes_socket(
listener: tuple[str, int], mock_gai: Any
) -> None:
"""A socket-option failure fails the connect instead of leaking sockets.
Instrumented at ``_set_socket_options`` (which the factory calls with
the just-created socket) rather than by patching ``socket.socket``,
which is platform dependent: the event loop's internal socketpair use
differs between platforms.
"""
created: list[socket.socket] = []
def failing_set_options(sock: socket.socket, options: Any) -> None:
created.append(sock)
raise OSError("bad socket option")
# Patch before building the closure; it binds _set_socket_options at
# creation time.
with patch("urllib3.util.connection._set_socket_options", new=failing_set_options):
create_connection = _make_create_connection()
with pytest.raises(OSError):
create_connection(
("example.com", listener[1]),
timeout=5,
socket_options=[(999999, 999999, 1)],
)
assert created, "socket factory never ran"
assert all(sock.fileno() == -1 for sock in created), "socket leaked open"
def test_default_timeout_yields_blocking_socket(
create_connection: Any, listener: tuple[str, int], mock_gai: Any
) -> None:
"""Without an explicit timeout the socket follows the global default."""
sock = create_connection(("example.com", listener[1]))
try:
assert sock.gettimeout() is socket.getdefaulttimeout()
finally:
sock.close()
def test_settimeout_failure_closes_socket(
create_connection: Any, mock_gai: Any
) -> None:
"""A failure restoring socket state closes the winner instead of leaking."""
bad_sock = Mock()
bad_sock.settimeout.side_effect = OSError("bad timeout")
with (
patch("esphome.async_thread.run_async", return_value=bad_sock),
pytest.raises(OSError, match="bad timeout"),
):
create_connection(("example.com", 80), timeout=5)
bad_sock.close.assert_called_once()
def test_connect_timeout_raises() -> None:
"""A connect that never completes raises within the timeout."""
async def never(*args: Any, **kwargs: Any) -> None:
await asyncio.sleep(60)
addr_infos = [_addr_info("192.0.2.1", 9), _addr_info("192.0.2.2", 9)]
# Patch before building the closure; it binds start_connection at
# creation time.
with patch("aiohappyeyeballs.start_connection", new=never):
create_connection = _make_create_connection()
with (
patch("socket.getaddrinfo", return_value=addr_infos),
pytest.raises(TimeoutError),
):
create_connection(("example.com", 80), timeout=0.1)
def test_invalid_host_raises_location_parse_error(create_connection: Any) -> None:
"""Hostnames urllib3 would reject are still rejected."""
from urllib3.exceptions import LocationParseError
with pytest.raises(LocationParseError):
create_connection(("a" * 300, 80))
def test_empty_getaddrinfo_raises_oserror(create_connection: Any) -> None:
"""An empty resolution matches stock urllib3's OSError, not ValueError."""
with (
patch("socket.getaddrinfo", return_value=[]),
pytest.raises(OSError, match="empty"),
):
create_connection(("example.com", 80), timeout=5)
def test_ensure_falls_back_to_stock_when_internals_move(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""If urllib3 private names disappear, downloads keep the stock connect
and the warning is latched to fire once, not per download."""
import urllib3.util.connection
from esphome import happy_eyeballs
def stock(*args: Any, **kwargs: Any) -> None:
pass
factory = Mock(side_effect=ImportError("gone"))
monkeypatch.setattr(urllib3.util.connection, "create_connection", stock)
monkeypatch.setattr(happy_eyeballs, "_make_create_connection", factory)
ensure_happy_eyeballs()
ensure_happy_eyeballs()
assert urllib3.util.connection.create_connection is stock
assert factory.call_count == 1
assert caplog.text.count("Happy Eyeballs unavailable") == 1
def test_ensure_survives_missing_urllib3(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""An unimportable urllib3 degrades with a warning instead of raising."""
import sys
with patch.dict(sys.modules, {"urllib3.util.connection": None}):
ensure_happy_eyeballs()
assert "Happy Eyeballs unavailable" in caplog.text
def test_requests_routes_through_shim(monkeypatch: pytest.MonkeyPatch) -> None:
"""Patching urllib3's create_connection actually reroutes requests."""
from http.server import BaseHTTPRequestHandler, HTTPServer
import threading
import requests
import urllib3.util.connection
class Handler(BaseHTTPRequestHandler):
def do_GET(self) -> None:
self.send_response(200)
self.send_header("Content-Length", "2")
self.end_headers()
self.wfile.write(b"ok")
def log_message(self, *args: Any) -> None:
pass
server = HTTPServer(("127.0.0.1", 0), Handler)
threading.Thread(target=server.serve_forever, daemon=True).start()
host, port = server.server_address
calls: list[Any] = []
shim = _make_create_connection()
def counting(*args: Any, **kwargs: Any) -> Any:
calls.append(args)
return shim(*args, **kwargs)
counting._esphome_patched = True
monkeypatch.setattr(urllib3.util.connection, "create_connection", counting)
real_getaddrinfo = socket.getaddrinfo
def fake_getaddrinfo(h: str, p: int, *args: Any, **kwargs: Any) -> Any:
if h == "shim-test.invalid":
return [_addr_info(host, port), _addr_info(host, port)]
return real_getaddrinfo(h, p, *args, **kwargs)
monkeypatch.setattr(socket, "getaddrinfo", fake_getaddrinfo)
try:
with requests.Session() as session:
session.trust_env = False
resp = session.get(f"http://shim-test.invalid:{port}/", timeout=5)
assert resp.status_code == 200
assert resp.content == b"ok"
assert calls, "requests did not go through the patched create_connection"
finally:
server.shutdown()
server.server_close()
+25 -17
View File
@@ -15,7 +15,6 @@ test pins down *which* heavy modules must stay out entirely.
from __future__ import annotations
import importlib.util
import os
from pathlib import Path
import subprocess
import sys
@@ -46,6 +45,11 @@ API_HEAVY_MODULES = ("aioesphomeapi",)
# never pays for the bundle machinery and its tarfile chain.
BUNDLE_HEAVY_MODULES = ("esphome.bundle", "tarfile")
# Heavy only for a cache-hit upload/logs run: the JSON cache parse must
# not resolve pyyaml or the yaml_util chain (the read_config fallback
# still uses both).
CACHE_HIT_HEAVY_MODULES = ("esphome.yaml_util", "yaml")
# Stdlib modules deferred out of the dispatch fast path: a cache-hit
# upload/logs run never writes a file (tempfile), spawns a process
# (subprocess), parses a URL (urllib.parse), or prints a serial
@@ -56,8 +60,6 @@ STDLIB_FAST_PATH_MODULES = (
"tempfile",
"subprocess",
"getpass",
# Pins the module-level contract only: PyYAML's constructor loads
# datetime during the cache parse until the JSON cache lands.
"datetime",
*(("urllib.parse",) if sys.version_info >= (3, 13) else ()),
)
@@ -108,6 +110,7 @@ def test_watched_heavy_modules_exist() -> None:
FAST_PATH_HEAVY_MODULES
+ API_HEAVY_MODULES
+ BUNDLE_HEAVY_MODULES
+ CACHE_HIT_HEAVY_MODULES
+ STDLIB_FAST_PATH_MODULES
):
assert importlib.util.find_spec(module) is not None, (
@@ -116,18 +119,17 @@ def test_watched_heavy_modules_exist() -> None:
def _leaked_from_fixture(
fixture_path: Path, script_name: str, extra: tuple[str, ...] = ()
fixture_path: Path,
env: dict[str, str],
script_name: str,
extra: tuple[str, ...] = (),
) -> str:
"""Run a fixture script with the watched modules on argv.
Running a script file drops the cwd from sys.path, so prepend the
repo root for the child; a non-zero exit surfaces the child's stderr.
``env`` comes from the ``probe_env`` fixture so the child can import
the repo checkout; a non-zero exit surfaces the child's stderr.
"""
script = fixture_path / "lazy_imports" / script_name
python_path = str(Path(__file__).parents[2])
if ambient := os.environ.get("PYTHONPATH"):
python_path = os.pathsep.join((python_path, ambient))
env = os.environ | {"PYTHONPATH": python_path}
result = subprocess.run(
[sys.executable, str(script), *FAST_PATH_HEAVY_MODULES, *extra],
capture_output=True,
@@ -141,12 +143,13 @@ def _leaked_from_fixture(
def test_storage_json_fast_path_does_not_import_heavy_modules(
fixture_path: Path,
probe_env: dict[str, str],
) -> None:
"""``apply_to_core`` runs on the upload/logs fast path for every
platform; parsing the stored framework version must not drag in the
validation stack or the esp32 component package.
"""
leaked = _leaked_from_fixture(fixture_path, "storage_json_fast_path.py")
leaked = _leaked_from_fixture(fixture_path, probe_env, "storage_json_fast_path.py")
assert not leaked, (
f"storage_json.apply_to_core pulls in heavy modules: {leaked}. "
"The upload/logs fast path skips validation; importing the "
@@ -156,12 +159,15 @@ def test_storage_json_fast_path_does_not_import_heavy_modules(
def test_esptool_upload_fast_path_does_not_import_heavy_modules(
fixture_path: Path,
probe_env: dict[str, str],
) -> None:
"""The esptool serial upload reads the esp32 variant from CORE.data;
resolving it must not drag in the esp32 component package or the
validation stack.
"""
leaked = _leaked_from_fixture(fixture_path, "esptool_upload_fast_path.py")
leaked = _leaked_from_fixture(
fixture_path, probe_env, "esptool_upload_fast_path.py"
)
assert not leaked, (
f"upload_using_esptool pulls in heavy modules: {leaked}. "
"The upload fast path skips validation; importing the validation "
@@ -262,6 +268,7 @@ def test_yaml_util_does_not_import_heavy_modules() -> None:
def test_upload_command_path_does_not_import_heavy_modules(
fixture_path: Path,
probe_env: dict[str, str],
) -> None:
"""The single-config dispatch path checks the bundle suffix on every
run; reading it from esphome.const must not drag in esphome.bundle
@@ -269,14 +276,15 @@ def test_upload_command_path_does_not_import_heavy_modules(
"""
leaked = _leaked_from_fixture(
fixture_path,
probe_env,
"upload_command_fast_path.py",
extra=BUNDLE_HEAVY_MODULES + STDLIB_FAST_PATH_MODULES,
extra=BUNDLE_HEAVY_MODULES + CACHE_HIT_HEAVY_MODULES + STDLIB_FAST_PATH_MODULES,
)
assert not leaked, (
f"the upload dispatch path pulls in heavy modules: {leaked}. "
"An ordinary run only needs the bundle suffix constant, and the "
"cache parse must not resolve voluptuous; keep the esphome.bundle "
"import inside the branch that extracts one, the Invalid import "
"inside the branch that raises it, and the deferred stdlib "
"imports inside the write/spawn/serial helpers that use them."
"JSON cache parse must not resolve voluptuous or pyyaml; keep the "
"esphome.bundle import inside the branch that extracts one, the "
"yaml_util imports inside the read_config fallback, and the "
"deferred stdlib imports inside the write/spawn/serial helpers."
)
+266 -1
View File
@@ -1,6 +1,44 @@
from collections.abc import Generator
import errno
import io
import logging
import os
from pathlib import Path
import select
import subprocess
import sys
import time
import pytest
from esphome.log import AnsiFore, AnsiStyle, color
from esphome.core import CORE
from esphome.log import AnsiFore, AnsiStyle, color, setup_log
class _FakeTty(io.StringIO):
def isatty(self) -> bool:
return True
@pytest.fixture
def restore_logging_state() -> Generator[None, None, None]:
"""Undo the global logging changes setup_log() makes."""
root = logging.getLogger()
handlers = root.handlers[:]
formatters = [handler.formatter for handler in handlers]
level = root.level
urllib3_level = logging.getLogger("urllib3").level
yield
root.handlers[:] = handlers
for handler, formatter in zip(handlers, formatters, strict=True):
handler.setFormatter(formatter)
root.setLevel(level)
logging.getLogger("urllib3").setLevel(urllib3_level)
def _probe_command(fixture_path: Path, *args: str) -> list[str]:
"""Build the command line for the setup_log probe fixture script."""
return [sys.executable, str(fixture_path / "log" / "setup_log_probe.py"), *args]
def test_color_keep_returns_unchanged_message() -> None:
@@ -78,3 +116,230 @@ def test_ansi_fore_keep_is_enum_member() -> None:
assert bool(AnsiFore.KEEP) is True
# But the value itself is still an empty string
assert AnsiFore.KEEP.value == ""
@pytest.mark.skipif(
sys.platform == "win32", reason="colorama always initializes on Windows"
)
def test_setup_log_redirected_output_strips_ansi(
fixture_path: Path, probe_env: dict[str, str]
) -> None:
"""A redirected run must keep colorama so ANSI codes are stripped."""
result = subprocess.run(
_probe_command(fixture_path),
capture_output=True,
text=True,
timeout=60,
check=False,
env=probe_env,
)
assert result.returncode == 0, result.stderr
assert "colorama_loaded=True" in result.stdout
assert "red end" in result.stdout
assert "\033" not in result.stdout
@pytest.mark.skipif(
sys.platform == "win32", reason="colorama always initializes on Windows"
)
def test_setup_log_dashboard_skips_colorama(
fixture_path: Path, probe_env: dict[str, str]
) -> None:
"""Dashboard runs escape their color codes, so colorama must not load."""
result = subprocess.run(
_probe_command(fixture_path, "--dashboard"),
capture_output=True,
text=True,
timeout=60,
check=False,
env=probe_env,
)
assert result.returncode == 0, result.stderr
assert "colorama_loaded=False" in result.stdout
# Codes pass through untouched for the dashboard to handle.
assert "\033[31mred\033[0m end" in result.stdout
def _run_probe_on_pty(
fixture_path: Path, probe_env: dict[str, str], *, stderr_to_pty: bool
) -> str:
"""Run the probe with stdout on a pty and return the decoded pty output.
With ``stderr_to_pty=False`` stderr goes to a pipe instead, giving the
mixed tty/redirect stream combination while keeping any traceback
available for the exit assertion.
"""
# Unix-only; a module-level import would break test collection on
# Windows, where all the callers are skipped anyway.
import pty
controller, follower = pty.openpty()
proc = None
output = b""
deadline = time.monotonic() + 60
try:
try:
proc = subprocess.Popen(
_probe_command(fixture_path),
stdout=follower,
stderr=follower if stderr_to_pty else subprocess.PIPE,
stdin=follower,
env=probe_env,
)
finally:
os.close(follower)
while True:
timeout = deadline - time.monotonic()
if timeout <= 0 or not select.select([controller], [], [], timeout)[0]:
pytest.fail(f"pty probe produced no EOF in time; got {output!r}")
try:
chunk = os.read(controller, 1024)
except OSError as err:
# macOS raises EIO once the child closes its end of the pty;
# anything else is a real failure, not end-of-stream.
if err.errno != errno.EIO:
raise
break
if not chunk:
break
output += chunk
stderr_text = ""
if proc.stderr is not None:
stderr_text = proc.stderr.read().decode(errors="replace")
proc.stderr.close()
assert proc.wait(60) == 0, stderr_text
finally:
os.close(controller)
if proc is not None and proc.poll() is None:
proc.kill()
proc.wait()
return output.decode()
@pytest.mark.skipif(
sys.platform == "win32", reason="pty is POSIX-only; colorama loads on Windows"
)
def test_setup_log_tty_skips_colorama(
fixture_path: Path, probe_env: dict[str, str]
) -> None:
"""A terminal run must skip colorama and keep ANSI codes intact."""
text = _run_probe_on_pty(fixture_path, probe_env, stderr_to_pty=True)
assert "colorama_loaded=False" in text
assert "\033[31mred\033[0m end" in text
@pytest.mark.skipif(
sys.platform == "win32", reason="pty is POSIX-only; colorama loads on Windows"
)
def test_setup_log_mixed_streams_init_colorama(
fixture_path: Path, probe_env: dict[str, str]
) -> None:
"""A tty stdout with a redirected stderr must still initialize colorama.
The guard requires both streams to be a tty; collapsing it to a
single-stream check would stop stripping ANSI from a redirected
stderr while stdout is a terminal.
"""
text = _run_probe_on_pty(fixture_path, probe_env, stderr_to_pty=False)
assert "colorama_loaded=True" in text
# stdout is a tty, so colorama leaves its codes alone.
assert "\033[31mred\033[0m end" in text
@pytest.fixture
def colorama_probe(
monkeypatch: pytest.MonkeyPatch, restore_logging_state: None
) -> Generator[None, None, None]:
"""Shared preamble for the in-process guard-branch tests.
Clears colorama from sys.modules so the assertions prove what
setup_log() itself did, and snapshots CORE.verbose/quiet, which is
not a no-op: CORE.reset() does not restore them, so without the
snapshot setup_log()'s log-level side effects would leak into later
tests.
"""
monkeypatch.delitem(sys.modules, "colorama", raising=False)
monkeypatch.setattr(CORE, "verbose", CORE.verbose)
monkeypatch.setattr(CORE, "quiet", CORE.quiet)
yield
# init() rebinds sys.stdout/stderr; restore them before monkeypatch
# puts the originals back.
if (colorama := sys.modules.get("colorama")) is not None:
colorama.deinit()
@pytest.mark.skipif(
sys.platform == "win32", reason="colorama always initializes on Windows"
)
def test_setup_log_dashboard_branch_skips_colorama_import(
monkeypatch: pytest.MonkeyPatch, colorama_probe: None
) -> None:
"""The dashboard side of the guard must not import colorama."""
monkeypatch.setattr(CORE, "dashboard", True)
setup_log()
assert "colorama" not in sys.modules
@pytest.mark.skipif(
sys.platform == "win32", reason="colorama always initializes on Windows"
)
def test_setup_log_tty_branch_skips_colorama_import(
monkeypatch: pytest.MonkeyPatch, colorama_probe: None
) -> None:
"""The tty side of the guard must not import colorama."""
monkeypatch.setattr(sys, "stdout", _FakeTty())
monkeypatch.setattr(sys, "stderr", _FakeTty())
setup_log()
assert "colorama" not in sys.modules
@pytest.mark.skipif(
sys.platform == "win32", reason="colorama always initializes on Windows"
)
def test_setup_log_redirected_branch_imports_colorama(
monkeypatch: pytest.MonkeyPatch, colorama_probe: None
) -> None:
"""Redirected streams must keep importing and initializing colorama."""
monkeypatch.setattr(sys, "stdout", io.StringIO())
monkeypatch.setattr(sys, "stderr", io.StringIO())
setup_log()
assert "colorama" in sys.modules
@pytest.mark.parametrize("broken", ["missing", "closed"])
def test_setup_log_broken_streams_import_colorama(
broken: str, monkeypatch: pytest.MonkeyPatch, colorama_probe: None
) -> None:
"""A missing or closed stream counts as a redirect and must not crash.
colorama tolerates both, so setup_log() has to reach its init rather
than raise inside the tty probe.
"""
if broken == "missing":
stream = None
else:
stream = io.StringIO()
stream.close()
monkeypatch.setattr(sys, "stdout", stream)
monkeypatch.setattr(sys, "stderr", stream)
setup_log()
assert "colorama" in sys.modules
def test_setup_log_win32_always_imports_colorama(
monkeypatch: pytest.MonkeyPatch, colorama_probe: None
) -> None:
"""The Windows clause must init colorama even when both streams are ttys.
Old Windows consoles need colorama to translate ANSI escapes, so the
platform check has to win over the tty check. colorama itself keys
off os.name, so on a POSIX host its init/deinit pair is a
passthrough.
"""
monkeypatch.setattr(sys, "platform", "win32")
# Both streams are ttys: without the platform clause this combination
# would skip colorama.
monkeypatch.setattr(sys, "stdout", _FakeTty())
monkeypatch.setattr(sys, "stderr", _FakeTty())
setup_log()
assert "colorama" in sys.modules