mirror of
https://github.com/esphome/esphome.git
synced 2026-09-07 21:46:09 +00:00
Compare commits
16
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c9fb996758 | ||
|
|
6aa6b682df | ||
|
|
5cc758dea3 | ||
|
|
bf718a28b9 | ||
|
|
a595590386 | ||
|
|
a0dc5a8d14 | ||
|
|
edec6aaf5e | ||
|
|
55804e6e16 | ||
|
|
228f8894a9 | ||
|
|
866574ca14 | ||
|
|
5747c736c2 | ||
|
|
0f6c266cd7 | ||
|
|
05dbc5ee59 | ||
|
|
29f7439154 | ||
|
|
89cd183a9f | ||
|
|
370cfb8898 |
@@ -13,6 +13,7 @@ from esphome.components.logger import request_log_listener
|
||||
from esphome.components.noise import ( # noqa: F401
|
||||
ENCRYPTION_SCHEMA,
|
||||
decode_encryption_key,
|
||||
enable_spare_ephemeral,
|
||||
encryption_schema,
|
||||
new_psk_progmem,
|
||||
validate_encryption_key,
|
||||
@@ -603,6 +604,7 @@ async def to_code(config: ConfigType) -> None:
|
||||
# and plaintext disabled. Only a factory reset can remove it.
|
||||
cg.add_define("USE_API_PLAINTEXT")
|
||||
cg.add_define("USE_API_NOISE")
|
||||
enable_spare_ephemeral()
|
||||
else:
|
||||
cg.add_define("USE_API_PLAINTEXT")
|
||||
|
||||
|
||||
@@ -316,9 +316,16 @@ class APIConnection final : public APIServerConnectionBase {
|
||||
void on_noise_encryption_set_key_request(const NoiseEncryptionSetKeyRequest &msg);
|
||||
#endif
|
||||
|
||||
// How long a new connection holds off the spare ephemeral refill
|
||||
static constexpr uint32_t CONNECT_GRACE_MS = 1000;
|
||||
bool is_authenticated() {
|
||||
return static_cast<ConnectionState>(this->flags_.connection_state) == ConnectionState::AUTHENTICATED;
|
||||
}
|
||||
// Unauthenticated and within its grace period; an older unauthenticated
|
||||
// connection is a stale half open client and no longer counts
|
||||
bool is_still_connecting(uint32_t now) {
|
||||
return !this->is_authenticated() && now - this->last_traffic_ < CONNECT_GRACE_MS;
|
||||
}
|
||||
bool is_connection_setup() {
|
||||
return static_cast<ConnectionState>(this->flags_.connection_state) == ConnectionState::CONNECTED ||
|
||||
this->is_authenticated();
|
||||
|
||||
@@ -143,6 +143,15 @@ void APIServer::loop() {
|
||||
this->accept_new_connections_();
|
||||
}
|
||||
|
||||
// Checked once per pass for the refill and for the clients below
|
||||
const bool connected = network::is_connected();
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
// Only the flag test is inline; refilling is the rare path
|
||||
if (connected && !noise::has_spare_ephemeral()) {
|
||||
this->refill_spare_ephemeral_();
|
||||
}
|
||||
#endif
|
||||
|
||||
if (this->api_connection_count_ == 0) {
|
||||
// Check reboot timeout - done in loop to avoid scheduler heap churn
|
||||
// (cancelled scheduler items sit in heap memory until their scheduled time).
|
||||
@@ -159,8 +168,7 @@ void APIServer::loop() {
|
||||
}
|
||||
|
||||
// Process clients and remove disconnected ones in a single pass
|
||||
// Check network connectivity once for all clients
|
||||
if (!network::is_connected()) {
|
||||
if (!connected) {
|
||||
// Network is down - disconnect all clients
|
||||
for (auto &client : this->active_clients()) {
|
||||
client->on_fatal_error();
|
||||
@@ -188,6 +196,21 @@ void APIServer::loop() {
|
||||
}
|
||||
}
|
||||
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
// Called with the network up; refill only while no api client is still
|
||||
// connecting (an OTA handshake is not visible here and just pays the refill
|
||||
// it triggered).
|
||||
void APIServer::refill_spare_ephemeral_() {
|
||||
const uint32_t now = App.get_loop_component_start_time();
|
||||
for (auto &client : this->active_clients()) {
|
||||
if (client->is_still_connecting(now)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
noise::prepare_spare_ephemeral();
|
||||
}
|
||||
#endif
|
||||
|
||||
void APIServer::remove_client_(uint8_t client_index) {
|
||||
auto &client = this->clients_[client_index];
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
#include "api_buffer.h"
|
||||
// Must precede clients_ so APIConnection is complete for default_delete (libc++).
|
||||
#include "api_connection.h"
|
||||
#ifdef USE_API_NOISE
|
||||
#if defined(USE_API_NOISE) || defined(USE_NOISE_SPARE_EPHEMERAL)
|
||||
// Only present in the build when the noise component is loaded
|
||||
#include "esphome/components/noise/noise.h"
|
||||
#endif
|
||||
@@ -363,6 +363,9 @@ class APIServer final : public Component,
|
||||
uint8_t provisioning_source_{0};
|
||||
#endif
|
||||
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
void refill_spare_ephemeral_();
|
||||
#endif
|
||||
#ifdef USE_API_NOISE
|
||||
noise::NoiseContext noise_ctx_;
|
||||
#ifndef USE_API_NOISE_PSK_FROM_YAML
|
||||
|
||||
@@ -37,25 +37,6 @@ CONF_HANDSHAKE_PIN = "handshake_pin"
|
||||
CONF_SDIO_FREQUENCY = "sdio_frequency"
|
||||
CONF_SPI_MODE = "spi_mode"
|
||||
|
||||
# ESP-NOW-over-hosted shim (esp_now_hosted.cpp). esp-hosted proxies esp_wifi.h
|
||||
# but not esp_now.h (espressif/esp-hosted-mcu#19), and esp_wifi_remote injects
|
||||
# the esp_now.h header on the ESP32-P4 host with no implementation, leaving the
|
||||
# esp_now_* symbols undefined at link. On a P4 host, esp_now_hosted.cpp DEFINES
|
||||
# those symbols and forwards each call to the co-processor over esp-hosted's
|
||||
# CustomRpc "peer data transfer" channel, so ESPHome's `espnow` component links
|
||||
# and runs unchanged (proven on a Tab5, 2026-07-20). The .cpp is guarded to
|
||||
# CONFIG_IDF_TARGET_ESP32P4 so it compiles to nothing on hosts with a native
|
||||
# ESP-NOW stack. CustomRpc needs these two host-side Kconfig options. Host
|
||||
# registers 3 handlers (RESP, RECV, SEND); the coprocessor registers 1 (REQ);
|
||||
# we ask for 8 to leave room for other CustomRpc extensions alongside.
|
||||
#
|
||||
# The coprocessor must run the matching custom firmware (a parallel effort in
|
||||
# esphome/esp-hosted-firmware). esp_now_hosted_rpc.h here is the canonical copy
|
||||
# of the wire contract and MUST stay byte-identical to the copy that coprocessor
|
||||
# firmware uses — the packed structs are the on-wire layout, so any divergence
|
||||
# silently corrupts every ESP-NOW frame.
|
||||
_MAX_CUSTOM_MSG_HANDLERS = 8
|
||||
|
||||
# Shared fields for both transport modes
|
||||
BASE_SCHEMA = cv.Schema(
|
||||
{
|
||||
@@ -281,23 +262,6 @@ async def to_code(config: ConfigType) -> None:
|
||||
else:
|
||||
_configure_spi(config)
|
||||
|
||||
# ESP-NOW-over-hosted shim: only the radio-less ESP32-P4 host needs it (see
|
||||
# the note by _MAX_CUSTOM_MSG_HANDLERS). Enabled for every P4 host, not
|
||||
# gated on the `espnow` component being present: the shim is tiny and the
|
||||
# esp_now_* symbols/CustomRpc calls it defines require these Kconfig options
|
||||
# to link whenever esp_now_hosted.cpp compiles (which is on any P4 host), so
|
||||
# coupling the two keeps the build consistent. When `espnow` is absent the
|
||||
# symbols are simply unused and never register a callback at runtime.
|
||||
if esp32.get_esp32_variant() == esp32.VARIANT_ESP32P4:
|
||||
add_define("USE_ESP_NOW_HOSTED")
|
||||
# esp-hosted's CustomRpc ("peer data transfer") path — off by default.
|
||||
esp32.add_idf_sdkconfig_option(
|
||||
"CONFIG_ESP_HOSTED_ENABLE_PEER_DATA_TRANSFER", True
|
||||
)
|
||||
esp32.add_idf_sdkconfig_option(
|
||||
"CONFIG_ESP_HOSTED_MAX_CUSTOM_MSG_HANDLERS", _MAX_CUSTOM_MSG_HANDLERS
|
||||
)
|
||||
|
||||
# Place the transport mempool in PSRAM. Required on memory-tight host
|
||||
# configurations (e.g. P4 with a large LVGL UI) where the internal-RAM
|
||||
# mempool allocation fails at boot with `sdio_mempool_create` assert.
|
||||
|
||||
@@ -1,467 +0,0 @@
|
||||
/*
|
||||
* esp_now_hosted — host-side shim implementing <esp_now.h> over esp-hosted
|
||||
* CustomRpc, so ESPHome's `espnow` component can run on a radio-less host
|
||||
* (e.g. the ESP32-P4) whose radio lives on an esp-hosted co-processor.
|
||||
*
|
||||
* A radio-less host has no native ESP-NOW. esp_wifi_remote INJECTS the full
|
||||
* esp_now.h header (types + declarations) but ships NO implementation, so every
|
||||
* esp_now_* symbol is an undefined reference at link time. This translation
|
||||
* unit provides those definitions; each forwards to the co-processor over
|
||||
* CustomRpc (see esphome/esp-hosted-firmware for the matching coprocessor
|
||||
* handlers). No esp-hosted or esp_wifi_remote source is patched, and there is no
|
||||
* duplicate-symbol clash because nothing else defines these symbols here.
|
||||
*
|
||||
* See esp_now_hosted_rpc.h for the wire protocol.
|
||||
*/
|
||||
|
||||
#include "sdkconfig.h"
|
||||
|
||||
// Only build the shim on the radio-less host. On chips with a native ESP-NOW
|
||||
// stack (S3, C6, …) the real symbols exist and this file must stay empty to
|
||||
// avoid duplicate definitions.
|
||||
#if defined(CONFIG_IDF_TARGET_ESP32P4)
|
||||
|
||||
#include <cstring>
|
||||
|
||||
#include "freertos/FreeRTOS.h"
|
||||
#include "freertos/semphr.h"
|
||||
|
||||
#include "esp_idf_version.h"
|
||||
#include "esp_log.h"
|
||||
#include "esp_timer.h"
|
||||
|
||||
#include <esp_now.h> // injected declarations we are now DEFINING
|
||||
#include <esp_wifi_types.h> // wifi_pkt_rx_ctrl_t, wifi_tx_info_t
|
||||
|
||||
// esp_hosted_misc.h (host) ships WITHOUT an extern "C" guard, so including it
|
||||
// from C++ would give its declarations C++ linkage and the real C symbols in
|
||||
// libesp_hosted would go unresolved at link. Wrap it. (Verified vs
|
||||
// esp_hosted 2.12.9.)
|
||||
extern "C" {
|
||||
#include "esp_hosted_misc.h" // esp_hosted_{send_custom_data,register_custom_callback}
|
||||
}
|
||||
|
||||
#include "esp_now_hosted_rpc.h"
|
||||
|
||||
namespace {
|
||||
|
||||
const char *const TAG = "esp_now_hosted";
|
||||
|
||||
// One outstanding request at a time. ESPHome drives esp_now_* from the main
|
||||
// loop; the matching response and the async RECV/SEND events all arrive on the
|
||||
// single esp-hosted RPC RX thread. Serializing requests keeps the shared
|
||||
// response slot race-free; a sequence number stops a late/stale response from
|
||||
// being mistaken for ours.
|
||||
SemaphoreHandle_t g_req_mutex = nullptr;
|
||||
SemaphoreHandle_t g_resp_sem = nullptr; // given when the matching RESP lands
|
||||
bool g_setup_done = false; // set only after setup fully succeeds
|
||||
uint8_t g_seq = 0;
|
||||
volatile uint8_t g_expect_seq = 0;
|
||||
volatile int32_t g_resp_status = 0;
|
||||
uint8_t g_resp_ret[16];
|
||||
volatile uint16_t g_resp_ret_len = 0;
|
||||
|
||||
// Written from the main loop (register/unregister/deinit), read from the
|
||||
// esp-hosted RX thread (on_recv/on_send). volatile for the same reason the
|
||||
// g_resp_* globals are: force the RX thread to observe an updated pointer
|
||||
// (e.g. a nulling by esp_now_deinit) rather than a cached one.
|
||||
volatile esp_now_recv_cb_t g_recv_cb = nullptr;
|
||||
volatile esp_now_send_cb_t g_send_cb = nullptr;
|
||||
|
||||
// Local mirror of the co-processor's peer table. ESPHome's espnow component
|
||||
// calls esp_now_is_peer_exist() on the main loop for every received frame
|
||||
// (twice) and every send; forwarding each as a blocking RPC round-trip stalls
|
||||
// the loop. The shim is the only path that mutates the co-processor peer table
|
||||
// (add/del/deinit all go through here), so this mirror is authoritative and
|
||||
// esp_now_is_peer_exist() can answer from it with no round-trip.
|
||||
//
|
||||
// esp_now_* are public C symbols: any component or user lambda may call them,
|
||||
// and although ESPHome's espnow touches peers only from the main loop today
|
||||
// (its RX/TX callbacks merely enqueue), the shim cannot rely on that. A short
|
||||
// spinlock keeps the mirror consistent from any task/core, matching native
|
||||
// esp_now_*'s own internal thread-safety. The critical sections are a bounded
|
||||
// (<=20-entry) scan, so they stay tiny. ESP_NOW_MAX_TOTAL_PEER_NUM is 20.
|
||||
constexpr size_t ESP_NOW_HOSTED_MAX_PEERS = 20;
|
||||
uint8_t g_peer_cache[ESP_NOW_HOSTED_MAX_PEERS][6];
|
||||
size_t g_peer_count = 0;
|
||||
portMUX_TYPE g_peer_lock = portMUX_INITIALIZER_UNLOCKED;
|
||||
|
||||
// Caller must hold g_peer_lock.
|
||||
int peer_cache_find_locked(const uint8_t *mac) {
|
||||
for (size_t i = 0; i < g_peer_count; i++) {
|
||||
if (memcmp(g_peer_cache[i], mac, 6) == 0)
|
||||
return static_cast<int>(i);
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
bool peer_cache_contains(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
const bool found = peer_cache_find_locked(mac) >= 0;
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
return found;
|
||||
}
|
||||
|
||||
void peer_cache_add(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
if (peer_cache_find_locked(mac) < 0 && g_peer_count < ESP_NOW_HOSTED_MAX_PEERS)
|
||||
memcpy(g_peer_cache[g_peer_count++], mac, 6);
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
void peer_cache_remove(const uint8_t *mac) {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
const int idx = peer_cache_find_locked(mac);
|
||||
if (idx >= 0) {
|
||||
g_peer_count--;
|
||||
if (static_cast<size_t>(idx) != g_peer_count) // move the last entry into the gap
|
||||
memcpy(g_peer_cache[idx], g_peer_cache[g_peer_count], 6);
|
||||
}
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
void peer_cache_clear() {
|
||||
portENTER_CRITICAL(&g_peer_lock);
|
||||
g_peer_count = 0;
|
||||
portEXIT_CRITICAL(&g_peer_lock);
|
||||
}
|
||||
|
||||
// ── CustomRpc event handlers (run on the esp-hosted RPC RX thread) ──────────
|
||||
// Keep them short and non-blocking. In particular they MUST NOT call back into
|
||||
// any esp_now_* shim function: that would try to take g_req_mutex / wait on the
|
||||
// RX thread that delivers the response, and deadlock.
|
||||
|
||||
void on_resp(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
if (len < sizeof(esp_now_hosted_resp_t)) {
|
||||
ESP_LOGW(TAG, "RESP too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *r = reinterpret_cast<const esp_now_hosted_resp_t *>(data);
|
||||
if (r->seq != g_expect_seq) { // late response from a timed-out request (expected)
|
||||
ESP_LOGV(TAG, "dropping stale RESP seq %u (want %u)", r->seq, g_expect_seq);
|
||||
return;
|
||||
}
|
||||
g_resp_status = r->status;
|
||||
uint16_t rl = r->ret_len;
|
||||
if (rl > sizeof(g_resp_ret)) {
|
||||
// Larger than any real opcode return — a likely wire-format drift signal.
|
||||
ESP_LOGW(TAG, "RESP ret_len %u exceeds buffer, clamping (wire drift?)", rl);
|
||||
rl = sizeof(g_resp_ret);
|
||||
}
|
||||
if (len >= sizeof(esp_now_hosted_resp_t) + rl) {
|
||||
memcpy(g_resp_ret, r->ret, rl);
|
||||
} else {
|
||||
// Truncated frame: fail closed. Never hand the caller stale bytes left in
|
||||
// g_resp_ret by a previous response, and don't let request() report a
|
||||
// zeroed payload as success — override the status to an error.
|
||||
ESP_LOGW(TAG, "RESP truncated: claims %u ret bytes, frame too short", rl);
|
||||
rl = 0;
|
||||
g_resp_status = ESP_ERR_INVALID_RESPONSE;
|
||||
}
|
||||
g_resp_ret_len = rl;
|
||||
xSemaphoreGive(g_resp_sem);
|
||||
}
|
||||
|
||||
void on_recv(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
// Read the volatile pointer once: esp_now_unregister_recv_cb()/deinit() (via
|
||||
// the espnow component's disable()) can null it on the main loop between the
|
||||
// guard and the call, which would otherwise turn the call into a null-deref.
|
||||
const esp_now_recv_cb_t cb = g_recv_cb;
|
||||
if (cb == nullptr)
|
||||
return;
|
||||
if (len < sizeof(esp_now_hosted_recv_evt_t)) {
|
||||
ESP_LOGW(TAG, "RECV too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *e = reinterpret_cast<const esp_now_hosted_recv_evt_t *>(data);
|
||||
if (len < sizeof(esp_now_hosted_recv_evt_t) + e->data_len) {
|
||||
ESP_LOGW(TAG, "RECV data_len %u exceeds frame", e->data_len);
|
||||
return;
|
||||
}
|
||||
|
||||
// ESPHome dereferences info->rx_ctrl->{rssi,timestamp}; give it a real one.
|
||||
wifi_pkt_rx_ctrl_t rx_ctrl;
|
||||
memset(&rx_ctrl, 0, sizeof(rx_ctrl));
|
||||
rx_ctrl.rssi = e->rssi;
|
||||
rx_ctrl.channel = e->channel;
|
||||
rx_ctrl.timestamp = static_cast<uint32_t>(esp_timer_get_time());
|
||||
|
||||
esp_now_recv_info_t info;
|
||||
info.src_addr = const_cast<uint8_t *>(e->src_addr);
|
||||
info.des_addr = const_cast<uint8_t *>(e->des_addr);
|
||||
info.rx_ctrl = &rx_ctrl;
|
||||
cb(&info, e->data, static_cast<int>(e->data_len));
|
||||
}
|
||||
|
||||
void on_send(uint32_t /*msg_id*/, const uint8_t *data, size_t len, void * /*ctx*/) {
|
||||
// Read the volatile pointer once (see on_recv): disable()/deinit() can null it
|
||||
// on the main loop concurrently with this RX-thread callback.
|
||||
const esp_now_send_cb_t cb = g_send_cb;
|
||||
if (cb == nullptr)
|
||||
return;
|
||||
if (len < sizeof(esp_now_hosted_send_evt_t)) {
|
||||
ESP_LOGW(TAG, "SEND evt too short: %u bytes", static_cast<unsigned>(len));
|
||||
return;
|
||||
}
|
||||
const auto *e = reinterpret_cast<const esp_now_hosted_send_evt_t *>(data);
|
||||
#if ESP_IDF_VERSION >= ESP_IDF_VERSION_VAL(5, 5, 0)
|
||||
// IDF >= 5.5: esp_now_send_cb_t takes esp_now_send_info_t (== wifi_tx_info_t),
|
||||
// whose des_addr is a POINTER (not an inline array). Point it at the event's
|
||||
// MAC (valid for this callback) — do NOT memcpy into it (that writes NULL and
|
||||
// faults). ESPHome reads only info->des_addr.
|
||||
esp_now_send_info_t si;
|
||||
memset(&si, 0, sizeof(si));
|
||||
si.des_addr = const_cast<uint8_t *>(e->des_addr);
|
||||
cb(&si, static_cast<esp_now_send_status_t>(e->status));
|
||||
#else
|
||||
cb(e->des_addr, static_cast<esp_now_send_status_t>(e->status));
|
||||
#endif
|
||||
}
|
||||
|
||||
esp_err_t ensure_setup() {
|
||||
// Gate on g_setup_done, not on g_req_mutex: a failure part-way through (a
|
||||
// semaphore that did not allocate, a callback that did not register) must not
|
||||
// leave a later call thinking setup completed. Semaphore creation is guarded
|
||||
// so a retry after a partial failure does not leak the earlier handles.
|
||||
if (g_setup_done)
|
||||
return ESP_OK;
|
||||
if (g_req_mutex == nullptr)
|
||||
g_req_mutex = xSemaphoreCreateMutex();
|
||||
if (g_resp_sem == nullptr)
|
||||
g_resp_sem = xSemaphoreCreateBinary();
|
||||
if (g_req_mutex == nullptr || g_resp_sem == nullptr)
|
||||
return ESP_ERR_NO_MEM;
|
||||
esp_err_t err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_RESP, on_resp, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_RECV, on_recv, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
if ((err = esp_hosted_register_custom_callback(ESP_NOW_HOSTED_MSG_SEND, on_send, nullptr)) != ESP_OK)
|
||||
return err;
|
||||
g_setup_done = true;
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
// Send one request envelope. With wait=true (default) block until the matching
|
||||
// response (or timeout); with wait=false return as soon as the frame is handed
|
||||
// to the transport (fire-and-forget, used by esp_now_send).
|
||||
//
|
||||
// `tail` is an optional second chunk written straight after `payload`. Callers
|
||||
// with a fixed header plus a bulk body (esp_now_send) pass the two separately
|
||||
// so they never need a build buffer of their own: both chunks are laid into the
|
||||
// request buffer here, under g_req_mutex, which keeps concurrent callers from
|
||||
// racing and saves a full copy of the body on every transmit.
|
||||
esp_err_t request(uint8_t opcode, const void *payload, uint16_t plen, void *ret, uint16_t ret_cap, uint16_t *ret_len,
|
||||
bool wait = true, const void *tail = nullptr, uint16_t tail_len = 0) {
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
if (plen > ESP_NOW_HOSTED_MAX_PAYLOAD || tail_len > ESP_NOW_HOSTED_MAX_PAYLOAD - plen)
|
||||
return ESP_ERR_INVALID_SIZE;
|
||||
const uint16_t total_len = static_cast<uint16_t>(plen + tail_len);
|
||||
|
||||
if (xSemaphoreTake(g_req_mutex, portMAX_DELAY) != pdTRUE)
|
||||
return ESP_FAIL;
|
||||
|
||||
static uint8_t buf[sizeof(esp_now_hosted_req_t) + ESP_NOW_HOSTED_MAX_PAYLOAD]; // guarded by g_req_mutex
|
||||
auto *req = reinterpret_cast<esp_now_hosted_req_t *>(buf);
|
||||
req->opcode = opcode;
|
||||
req->seq = ++g_seq;
|
||||
req->payload_len = total_len;
|
||||
if (plen != 0)
|
||||
memcpy(req->payload, payload, plen);
|
||||
if (tail_len != 0)
|
||||
memcpy(req->payload + plen, tail, tail_len);
|
||||
g_expect_seq = req->seq;
|
||||
|
||||
xSemaphoreTake(g_resp_sem, 0); // drain any stale signal before sending
|
||||
err = esp_hosted_send_custom_data(ESP_NOW_HOSTED_MSG_REQ, buf, sizeof(esp_now_hosted_req_t) + total_len);
|
||||
if (err != ESP_OK) {
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return err;
|
||||
}
|
||||
if (!wait) {
|
||||
// Fire-and-forget (esp_now_send): the co-processor enqueues the frame and
|
||||
// reports the real TX result later via the async SEND event, exactly like
|
||||
// native esp_now_send. Returning here keeps the main loop off the ~100 ms+
|
||||
// RPC round-trip. The matching RESP is ignored (seq won't match the next
|
||||
// waited request, so on_resp drops it).
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return ESP_OK;
|
||||
}
|
||||
if (xSemaphoreTake(g_resp_sem, pdMS_TO_TICKS(ESP_NOW_HOSTED_TIMEOUT_MS)) != pdTRUE) {
|
||||
ESP_LOGW(TAG, "opcode %u timed out", opcode);
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return ESP_ERR_TIMEOUT;
|
||||
}
|
||||
|
||||
const int32_t status = g_resp_status;
|
||||
if (ret != nullptr && ret_cap != 0) {
|
||||
uint16_t n = g_resp_ret_len < ret_cap ? g_resp_ret_len : ret_cap;
|
||||
memcpy(ret, const_cast<const uint8_t *>(g_resp_ret), n);
|
||||
if (ret_len != nullptr)
|
||||
*ret_len = n;
|
||||
}
|
||||
xSemaphoreGive(g_req_mutex);
|
||||
return static_cast<esp_err_t>(status);
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
// ── The <esp_now.h> surface, defined for the radio-less host ────────────────
|
||||
extern "C" {
|
||||
|
||||
esp_err_t esp_now_init(void) { return request(ESP_NOW_HOSTED_OP_INIT, nullptr, 0, nullptr, 0, nullptr); }
|
||||
|
||||
esp_err_t esp_now_deinit(void) {
|
||||
g_recv_cb = nullptr;
|
||||
g_send_cb = nullptr;
|
||||
peer_cache_clear(); // the co-processor drops all peers on deinit
|
||||
return request(ESP_NOW_HOSTED_OP_DEINIT, nullptr, 0, nullptr, 0, nullptr);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_get_version(uint32_t *version) {
|
||||
uint32_t v = 0;
|
||||
uint16_t rl = 0;
|
||||
esp_err_t err = request(ESP_NOW_HOSTED_OP_GET_VERSION, nullptr, 0, &v, sizeof(v), &rl);
|
||||
if (version != nullptr)
|
||||
*version = v;
|
||||
return err;
|
||||
}
|
||||
|
||||
esp_err_t esp_now_register_recv_cb(esp_now_recv_cb_t cb) {
|
||||
// Only arm the callback once the CustomRpc handlers are actually registered,
|
||||
// so a failed setup leaves g_recv_cb null rather than falsely "registered".
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
g_recv_cb = cb;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_unregister_recv_cb(void) {
|
||||
g_recv_cb = nullptr;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_register_send_cb(esp_now_send_cb_t cb) {
|
||||
esp_err_t err = ensure_setup();
|
||||
if (err != ESP_OK)
|
||||
return err;
|
||||
g_send_cb = cb;
|
||||
return ESP_OK;
|
||||
}
|
||||
esp_err_t esp_now_unregister_send_cb(void) {
|
||||
g_send_cb = nullptr;
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
static esp_err_t add_or_mod_peer(uint8_t opcode, const esp_now_peer_info_t *peer, bool wait) {
|
||||
if (peer == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
esp_now_hosted_peer_t p;
|
||||
memset(&p, 0, sizeof(p));
|
||||
memcpy(p.peer_addr, peer->peer_addr, 6);
|
||||
memcpy(p.lmk, peer->lmk, 16);
|
||||
p.channel = peer->channel;
|
||||
p.ifidx = static_cast<uint8_t>(peer->ifidx);
|
||||
p.encrypt = peer->encrypt ? 1 : 0;
|
||||
return request(opcode, &p, sizeof(p), nullptr, 0, nullptr, wait);
|
||||
}
|
||||
esp_err_t esp_now_add_peer(const esp_now_peer_info_t *peer) {
|
||||
// Fire-and-forget (wait=false): adding a peer is a blocking RPC round-trip,
|
||||
// and ESPHome's espnow calls it on the main loop when a device joins the mesh
|
||||
// — under co-processor load that stalls the UI (peer-churn stutter). Issue it
|
||||
// without waiting and mirror it locally. Safe against a following
|
||||
// esp_now_send to the same peer: both ride the same in-order CustomRpc
|
||||
// channel (mutex-serialized on the host) and the co-processor processes REQs
|
||||
// FIFO, so ADD_PEER is applied before the SEND. Trade-off: a co-processor-side
|
||||
// failure (e.g. peer table full) is no longer reported synchronously — the
|
||||
// same limitation as esp_now_send — but ESPHome only adds peers it validated.
|
||||
esp_err_t err = add_or_mod_peer(ESP_NOW_HOSTED_OP_ADD_PEER, peer, /*wait=*/false);
|
||||
if (err == ESP_OK)
|
||||
peer_cache_add(peer->peer_addr); // keep the local mirror in sync
|
||||
return err;
|
||||
}
|
||||
esp_err_t esp_now_mod_peer(const esp_now_peer_info_t *peer) {
|
||||
// mod_peer changes a peer's parameters, not its existence, so the cache is
|
||||
// unaffected. Kept synchronous — it is not on any hot path (espnow never
|
||||
// calls it), so the extra round-trip does not matter and the status is useful.
|
||||
return add_or_mod_peer(ESP_NOW_HOSTED_OP_MOD_PEER, peer, /*wait=*/true);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_del_peer(const uint8_t *peer_addr) {
|
||||
if (peer_addr == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
// Fire-and-forget for the same reason as add_peer (peer churn on the main
|
||||
// loop). Removal is order-independent, so this is strictly safe.
|
||||
esp_err_t err = request(ESP_NOW_HOSTED_OP_DEL_PEER, peer_addr, 6, nullptr, 0, nullptr, /*wait=*/false);
|
||||
if (err == ESP_OK)
|
||||
peer_cache_remove(peer_addr); // keep the local mirror in sync
|
||||
return err;
|
||||
}
|
||||
|
||||
bool esp_now_is_peer_exist(const uint8_t *peer_addr) {
|
||||
if (peer_addr == nullptr)
|
||||
return false;
|
||||
// Answered from the local mirror — no RPC round-trip. ESPHome's espnow calls
|
||||
// this on the main loop for every received frame and every send, so a
|
||||
// blocking round-trip here would stall rendering under mesh traffic.
|
||||
return peer_cache_contains(peer_addr);
|
||||
}
|
||||
|
||||
esp_err_t esp_now_send(const uint8_t *peer_addr, const uint8_t *data, size_t len) {
|
||||
if (len > ESP_NOW_HOSTED_MAX_FRAME)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
if (data == nullptr && len != 0) // native esp_now_send treats this as an arg error
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
// Only the small fixed header is built here; the caller's frame goes over as
|
||||
// the request tail, so request() lays both into its own buffer under
|
||||
// g_req_mutex. esp_now_send is a public C symbol and may be called from any
|
||||
// task, and a shared build buffer here would let two callers corrupt each
|
||||
// other's frame. Passing the body through also drops a full-frame copy per
|
||||
// transmit, on the path this shim exists to keep quick.
|
||||
uint8_t hdr[sizeof(esp_now_hosted_send_req_t)];
|
||||
auto *s = reinterpret_cast<esp_now_hosted_send_req_t *>(hdr);
|
||||
s->has_addr = peer_addr != nullptr ? 1 : 0;
|
||||
if (peer_addr != nullptr)
|
||||
memcpy(s->peer_addr, peer_addr, 6);
|
||||
else
|
||||
memset(s->peer_addr, 0, 6);
|
||||
s->data_len = static_cast<uint16_t>(len);
|
||||
// Fire-and-forget (wait=false): native esp_now_send returns once the frame is
|
||||
// queued, with the real TX result delivered later through the send callback.
|
||||
// The co-processor mirrors that — it acks enqueue immediately and reports the
|
||||
// outcome via the async SEND event (on_send -> on_send_report). Waiting for
|
||||
// the RPC RESP here would block the main loop for the full round-trip on
|
||||
// every transmit.
|
||||
return request(ESP_NOW_HOSTED_OP_SEND, hdr, sizeof(hdr), nullptr, 0, nullptr, /*wait=*/false, data,
|
||||
static_cast<uint16_t>(len));
|
||||
}
|
||||
|
||||
esp_err_t esp_now_set_pmk(const uint8_t *pmk) {
|
||||
if (pmk == nullptr)
|
||||
return ESP_ERR_ESPNOW_ARG;
|
||||
return request(ESP_NOW_HOSTED_OP_SET_PMK, pmk, 16, nullptr, 0, nullptr);
|
||||
}
|
||||
|
||||
// Remainder of the <esp_now.h> surface. Not used by ESPHome's espnow component
|
||||
// today; provided so the whole header links and future callers get a defined
|
||||
// (if unimplemented) symbol rather than a link error. Wire them through
|
||||
// CustomRpc if a use case appears.
|
||||
esp_err_t esp_now_get_peer(const uint8_t * /*peer_addr*/, esp_now_peer_info_t * /*peer*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
esp_err_t esp_now_fetch_peer(bool /*from_head*/, esp_now_peer_info_t * /*peer*/) { return ESP_ERR_NOT_SUPPORTED; }
|
||||
esp_err_t esp_now_get_peer_num(esp_now_peer_num_t * /*num*/) { return ESP_ERR_NOT_SUPPORTED; }
|
||||
esp_err_t esp_now_set_wake_window(uint16_t /*window*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED; // power-save wake window is not forwarded; don't claim success
|
||||
}
|
||||
esp_err_t esp_now_set_peer_rate_config(const uint8_t * /*peer_addr*/, esp_now_rate_config_t * /*cfg*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
esp_err_t esp_wifi_config_espnow_rate(wifi_interface_t /*ifx*/, wifi_phy_rate_t /*rate*/) {
|
||||
return ESP_ERR_NOT_SUPPORTED;
|
||||
}
|
||||
|
||||
} // extern "C"
|
||||
|
||||
#endif // CONFIG_IDF_TARGET_ESP32P4
|
||||
@@ -1,128 +0,0 @@
|
||||
/*
|
||||
* esp_now_hosted — ESP-NOW-over-CustomRpc wire protocol.
|
||||
*
|
||||
* Shared, byte-for-byte-identical contract between:
|
||||
* - the host shim (esphome/components/esp32_hosted/esp_now_hosted.cpp)
|
||||
* - the coprocessor firmware (esphome/esp-hosted-firmware)
|
||||
*
|
||||
* It rides esp-hosted's CustomRpc channel (RPC ID 388, "peer data transfer",
|
||||
* available since esp-hosted v2.8.1), teaching the radio-less host <-> radio
|
||||
* co-processor link to carry esp_now.h, which esp-hosted itself does not proxy
|
||||
* (Espressif issue espressif/esp-hosted-mcu#19).
|
||||
*
|
||||
* KEEP THE TWO COPIES IN SYNC. The canonical copy lives here; the coprocessor
|
||||
* firmware uses a verbatim copy. Both sides are little-endian, so these packed
|
||||
* structs are wire-compatible with no byte-swapping.
|
||||
*/
|
||||
|
||||
#ifndef ESP_NOW_HOSTED_RPC_H
|
||||
#define ESP_NOW_HOSTED_RPC_H
|
||||
|
||||
#ifdef __cplusplus
|
||||
#include <cstdint>
|
||||
#else
|
||||
#include <stdint.h>
|
||||
#endif
|
||||
|
||||
#ifdef __cplusplus
|
||||
extern "C" {
|
||||
#endif
|
||||
|
||||
/* ── CustomRpc message IDs (any uint32_t except 0xFFFFFFFF) ──────────────────
|
||||
* One REQ handler slot on the device; three event handler slots on the host.
|
||||
* The bytes spell "now" + index, a private range unlikely to clash with other
|
||||
* CustomRpc users (e.g. the stock peer_data_transfer example's 1..6). */
|
||||
#define ESP_NOW_HOSTED_MSG_REQ 0x6E6F7701u /* host -> device : request envelope */
|
||||
#define ESP_NOW_HOSTED_MSG_RESP 0x6E6F7702u /* device -> host : reply to a REQ */
|
||||
#define ESP_NOW_HOSTED_MSG_RECV 0x6E6F7703u /* device -> host : async RX frame */
|
||||
#define ESP_NOW_HOSTED_MSG_SEND 0x6E6F7704u /* device -> host : async TX status */
|
||||
|
||||
/* ── Request opcodes ────────────────────────────────────────────────────── */
|
||||
enum {
|
||||
ESP_NOW_HOSTED_OP_INIT = 1, /* esp_now_init + register device recv/send cbs */
|
||||
ESP_NOW_HOSTED_OP_DEINIT = 2, /* unregister cbs + esp_now_deinit */
|
||||
ESP_NOW_HOSTED_OP_ADD_PEER = 3, /* payload: esp_now_hosted_peer_t */
|
||||
ESP_NOW_HOSTED_OP_DEL_PEER = 4, /* payload: 6-byte peer MAC */
|
||||
ESP_NOW_HOSTED_OP_IS_PEER_EXIST = 5, /* payload: 6-byte MAC; ret: 1 byte bool */
|
||||
ESP_NOW_HOSTED_OP_SEND = 6, /* payload: esp_now_hosted_send_req_t */
|
||||
ESP_NOW_HOSTED_OP_GET_VERSION = 7, /* ret: uint32 version */
|
||||
ESP_NOW_HOSTED_OP_SET_PMK = 8, /* payload: 16-byte PMK */
|
||||
ESP_NOW_HOSTED_OP_MOD_PEER = 9, /* payload: esp_now_hosted_peer_t */
|
||||
};
|
||||
|
||||
/* Largest ESP-NOW payload we forward. ESP-NOW v2 (IDF >= 5.4) is 1470 B; well
|
||||
* under esp-hosted's 8166 B CustomRpc cap, so the shim never truncates. */
|
||||
#define ESP_NOW_HOSTED_MAX_FRAME 1470u
|
||||
/* Envelope slack for the largest opcode payload (a SEND req wrapping a frame). */
|
||||
#define ESP_NOW_HOSTED_MAX_PAYLOAD (ESP_NOW_HOSTED_MAX_FRAME + 16u)
|
||||
/* Host request/response round-trip timeout over the transport. Generous:
|
||||
* normal RTT is sub-millisecond, but Wi-Fi/BLE contention on the co-processor
|
||||
* can stall the RX thread. */
|
||||
#define ESP_NOW_HOSTED_TIMEOUT_MS 2000
|
||||
|
||||
/* ── Envelopes ──────────────────────────────────────────────────────────── */
|
||||
|
||||
/* These payloads are shared verbatim with the C co-processor firmware, so they
|
||||
* use C's `typedef struct {...} name;` idiom rather than C++ `using` aliases,
|
||||
* which would not compile there. Silence clang-tidy's modernize-use-using for
|
||||
* the shared struct block. */
|
||||
// NOLINTBEGIN(modernize-use-using)
|
||||
typedef struct {
|
||||
uint8_t opcode; /* one of ESP_NOW_HOSTED_OP_* */
|
||||
uint8_t seq; /* wraps 0..255; echoed in the response for matching */
|
||||
uint16_t payload_len; /* bytes of opcode-specific payload that follow */
|
||||
uint8_t payload[]; /* flexible */
|
||||
} __attribute__((packed)) esp_now_hosted_req_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t opcode; /* echoes the request opcode */
|
||||
uint8_t seq; /* echoes the request seq */
|
||||
int32_t status; /* esp_err_t from the native call on the co-processor */
|
||||
uint16_t ret_len; /* bytes of return payload that follow */
|
||||
uint8_t ret[]; /* flexible (e.g. version u32, is_peer_exist bool) */
|
||||
} __attribute__((packed)) esp_now_hosted_resp_t;
|
||||
|
||||
/* ── Opcode payloads ────────────────────────────────────────────────────── */
|
||||
|
||||
/* esp_now_peer_info_t minus the host-only `priv` pointer, which is meaningless
|
||||
* across the transport and never set by ESPHome's espnow component. */
|
||||
typedef struct {
|
||||
uint8_t peer_addr[6];
|
||||
uint8_t lmk[16];
|
||||
uint8_t channel; /* 0 = current channel */
|
||||
uint8_t ifidx; /* wifi_interface_t (0=STA, 1=AP) */
|
||||
uint8_t encrypt; /* bool */
|
||||
} __attribute__((packed)) esp_now_hosted_peer_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t has_addr; /* 0 => peer_addr is NULL (broadcast to all peers) */
|
||||
uint8_t peer_addr[6];
|
||||
uint16_t data_len;
|
||||
uint8_t data[]; /* flexible, up to ESP_NOW_HOSTED_MAX_FRAME */
|
||||
} __attribute__((packed)) esp_now_hosted_send_req_t;
|
||||
|
||||
/* ── Async events (device -> host) ──────────────────────────────────────── */
|
||||
|
||||
/* Reconstructed on the host into an esp_now_recv_info_t + a minimal
|
||||
* wifi_pkt_rx_ctrl_t. ESPHome's espnow reads info->src_addr, info->des_addr,
|
||||
* info->rx_ctrl->rssi and info->rx_ctrl->timestamp. */
|
||||
typedef struct {
|
||||
uint8_t src_addr[6];
|
||||
uint8_t des_addr[6];
|
||||
int8_t rssi;
|
||||
uint8_t channel;
|
||||
uint16_t data_len;
|
||||
uint8_t data[]; /* flexible */
|
||||
} __attribute__((packed)) esp_now_hosted_recv_evt_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t des_addr[6];
|
||||
uint8_t status; /* esp_now_send_status_t (0 = success) */
|
||||
} __attribute__((packed)) esp_now_hosted_send_evt_t;
|
||||
// NOLINTEND(modernize-use-using)
|
||||
|
||||
#ifdef __cplusplus
|
||||
}
|
||||
#endif
|
||||
|
||||
#endif /* ESP_NOW_HOSTED_RPC_H */
|
||||
@@ -3,7 +3,6 @@ from typing import Any
|
||||
from esphome import automation, core
|
||||
import esphome.codegen as cg
|
||||
from esphome.components import wifi
|
||||
from esphome.components.esp32 import VARIANT_ESP32P4, get_esp32_variant
|
||||
from esphome.components.udp import CONF_ON_RECEIVE
|
||||
import esphome.config_validation as cv
|
||||
from esphome.const import (
|
||||
@@ -18,7 +17,6 @@ from esphome.const import (
|
||||
)
|
||||
from esphome.core import CORE, HexInt
|
||||
from esphome.cpp_generator import MockObj, TemplateArgsType
|
||||
import esphome.final_validate as fv
|
||||
from esphome.types import ConfigType
|
||||
|
||||
CODEOWNERS = ["@jesserockz"]
|
||||
@@ -134,24 +132,6 @@ CONFIG_SCHEMA = cv.All(
|
||||
)
|
||||
|
||||
|
||||
def _validate_variant(config: ConfigType) -> ConfigType:
|
||||
# ESP-NOW rides the Wi-Fi PHY. Radio-less esp32 variants have no native
|
||||
# ESP-NOW; only the ESP32-P4 has a path, via the esp32_hosted shim that
|
||||
# supplies the esp_now_* symbols. Fail here with a clear message instead of
|
||||
# letting the build reach an "undefined reference to esp_now_*" link error.
|
||||
variant = get_esp32_variant()
|
||||
if wifi.variant_has_wifi(variant):
|
||||
return config
|
||||
if variant != VARIANT_ESP32P4:
|
||||
raise cv.Invalid(f"ESP-NOW is not supported on {variant} (no Wi-Fi radio)")
|
||||
if "esp32_hosted" not in fv.full_config.get():
|
||||
raise cv.Invalid(f"ESP-NOW on {variant} requires the esp32_hosted component")
|
||||
return config
|
||||
|
||||
|
||||
FINAL_VALIDATE_SCHEMA = _validate_variant
|
||||
|
||||
|
||||
async def _trigger_to_code(config: ConfigType) -> MockObj:
|
||||
if address := config.get(CONF_ADDRESS):
|
||||
address = address.parts
|
||||
|
||||
@@ -192,8 +192,6 @@ async def to_code(config: ConfigType) -> None:
|
||||
if CORE.using_arduino:
|
||||
if CORE.is_esp8266:
|
||||
cg.add_library("ESP8266mDNS", None)
|
||||
# No MDNS global in the build; mdns_esp8266.cpp owns a guarded MDNSResponder
|
||||
cg.add_build_flag("-DNO_GLOBAL_MDNS")
|
||||
elif CORE.is_rp2:
|
||||
cg.add_library("LEAmDNS", None)
|
||||
|
||||
|
||||
@@ -13,47 +13,8 @@
|
||||
|
||||
namespace esphome::mdns {
|
||||
|
||||
// Main-loop calls into LEAmDNS that send (update() and close(); begin(), addService() and
|
||||
// the scheduled restart never reach a send) can yield inside UdpContext::sendTimeout(); a
|
||||
// packet arriving then re-enters LEAmDNS from lwIP on the same UdpContext and both sides
|
||||
// free the same tx pbufs (#18760). Received packets stay queued during such a call and are
|
||||
// processed from the main loop afterwards.
|
||||
class GuardedMDNSResponder : public ::esp8266::MDNSImplementation::MDNSResponder {
|
||||
public:
|
||||
void update_guarded() { this->run_guarded_(&GuardedMDNSResponder::update); }
|
||||
void close_guarded() { this->run_guarded_(&GuardedMDNSResponder::close); }
|
||||
|
||||
private:
|
||||
void run_guarded_(bool (GuardedMDNSResponder::*fn)()) {
|
||||
UdpContext *ctx = this->m_pUDPContext;
|
||||
if (ctx == nullptr) {
|
||||
(this->*fn)();
|
||||
return;
|
||||
}
|
||||
// Set every time: a restart replaces the context together with its stock handler. Only
|
||||
// begin() and the scheduled netif callback restart, never update() or close(), so the
|
||||
// context cannot change underneath this call.
|
||||
ctx->onRx([this]() {
|
||||
if (!this->in_loop_call_) {
|
||||
this->_callProcess();
|
||||
}
|
||||
});
|
||||
this->in_loop_call_ = true;
|
||||
(this->*fn)();
|
||||
// close() releases the context; a yield in here queues further packets for this loop too
|
||||
while (this->m_pUDPContext != nullptr && this->m_pUDPContext->next()) {
|
||||
this->_parseMessage();
|
||||
}
|
||||
this->in_loop_call_ = false;
|
||||
}
|
||||
|
||||
volatile bool in_loop_call_{false};
|
||||
};
|
||||
|
||||
static GuardedMDNSResponder mdns_responder; // NOLINT(cppcoreguidelines-avoid-non-const-global-variables)
|
||||
|
||||
static void register_esp8266(MDNSComponent *, StaticVector<MDNSService, MDNS_SERVICE_COUNT> &services) {
|
||||
mdns_responder.begin(App.get_name().c_str());
|
||||
MDNS.begin(App.get_name().c_str());
|
||||
|
||||
for (const auto &service : services) {
|
||||
// Strip the leading underscore from the proto and service_type. While it is
|
||||
@@ -69,10 +30,10 @@ static void register_esp8266(MDNSComponent *, StaticVector<MDNSService, MDNS_SER
|
||||
service_type++;
|
||||
}
|
||||
uint16_t port = service.port.value();
|
||||
mdns_responder.addService(FPSTR(service_type), FPSTR(proto), port);
|
||||
MDNS.addService(FPSTR(service_type), FPSTR(proto), port);
|
||||
for (const auto &record : service.txt_records) {
|
||||
mdns_responder.addServiceTxt(FPSTR(service_type), FPSTR(proto), FPSTR(MDNS_STR_ARG(record.key)),
|
||||
FPSTR(MDNS_STR_ARG(record.value)));
|
||||
MDNS.addServiceTxt(FPSTR(service_type), FPSTR(proto), FPSTR(MDNS_STR_ARG(record.key)),
|
||||
FPSTR(MDNS_STR_ARG(record.value)));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -91,7 +52,7 @@ void MDNSComponent::start_polling_window_() {
|
||||
if (wifi->is_roaming() || (!wifi->is_connected() && !wifi->is_ap_active()))
|
||||
return;
|
||||
#endif
|
||||
mdns_responder.update_guarded();
|
||||
MDNS.update();
|
||||
});
|
||||
this->set_timeout(MDNS_POLL_STOP_ID, MDNS_POLL_WINDOW_MS, [this]() { this->cancel_interval(MDNS_POLL_ID); });
|
||||
}
|
||||
@@ -120,7 +81,7 @@ void MDNSComponent::on_ip_state(const network::IPAddresses &ips, const network::
|
||||
#endif
|
||||
|
||||
void MDNSComponent::on_shutdown() {
|
||||
mdns_responder.close_guarded();
|
||||
MDNS.close();
|
||||
delay(10);
|
||||
}
|
||||
|
||||
|
||||
@@ -86,6 +86,11 @@ def encryption_schema(config: ConfigType | None) -> ConfigType:
|
||||
return ENCRYPTION_SCHEMA(config)
|
||||
|
||||
|
||||
def enable_spare_ephemeral() -> None:
|
||||
"""Compile the spare ephemeral key slot; the component that refills it calls this."""
|
||||
cg.add_define("USE_NOISE_SPARE_EPHEMERAL")
|
||||
|
||||
|
||||
async def to_code(config: ConfigType) -> None:
|
||||
cg.add_define("USE_NOISE")
|
||||
cg.add_library("esphome/noise-c", "0.1.24")
|
||||
|
||||
@@ -1,12 +1,14 @@
|
||||
#include "noise.h"
|
||||
#ifdef USE_NOISE
|
||||
#include "esphome/core/hal.h"
|
||||
#include "esphome/core/helpers.h"
|
||||
#include "esphome/core/log.h"
|
||||
|
||||
#include <algorithm>
|
||||
#include <cstring>
|
||||
|
||||
#include <noise/protocol.h>
|
||||
#include <sodium.h>
|
||||
|
||||
#ifdef USE_ESP8266
|
||||
#include <pgmspace.h>
|
||||
@@ -24,6 +26,39 @@ void NoiseContext::load_psk(psk_t &out) const {
|
||||
progmem_memcpy(out.data(), this->psk_, out.size());
|
||||
}
|
||||
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
static constexpr size_t PRIVATE_KEY_SIZE = SPARE_EPHEMERAL_KEY_SIZE;
|
||||
static constexpr size_t PUBLIC_KEY_SIZE = SPARE_EPHEMERAL_KEY_SIZE;
|
||||
uint8_t spare_ephemeral[SPARE_EPHEMERAL_SIZE]; // NOLINT(cppcoreguidelines-avoid-non-const-global-variables)
|
||||
|
||||
void prepare_spare_ephemeral() {
|
||||
uint8_t *private_key = spare_ephemeral;
|
||||
uint8_t *public_key = spare_ephemeral + PRIVATE_KEY_SIZE;
|
||||
// Same steps as noise-c's curve25519 keygen; the clamp sets the ready bit,
|
||||
// a failure wipes the slot so the handshake generates its own key
|
||||
if (!random_bytes(private_key, PRIVATE_KEY_SIZE)) {
|
||||
sodium_memzero(spare_ephemeral, sizeof(spare_ephemeral));
|
||||
return;
|
||||
}
|
||||
private_key[0] &= 0xF8;
|
||||
private_key[PRIVATE_KEY_SIZE - 1] = (private_key[PRIVATE_KEY_SIZE - 1] & 0x7F) | 0x40;
|
||||
if (crypto_scalarmult_curve25519_base(public_key, private_key) != 0) {
|
||||
sodium_memzero(spare_ephemeral, sizeof(spare_ephemeral));
|
||||
}
|
||||
}
|
||||
|
||||
int consume_spare_ephemeral(NoiseHandshakeState *state) {
|
||||
if (!has_spare_ephemeral()) {
|
||||
return 0;
|
||||
}
|
||||
// noise-c keeps its own copy, so the slot is wiped either way
|
||||
int err = noise_handshakestate_set_local_ephemeral(state, spare_ephemeral, PRIVATE_KEY_SIZE,
|
||||
spare_ephemeral + PRIVATE_KEY_SIZE, PUBLIC_KEY_SIZE);
|
||||
sodium_memzero(spare_ephemeral, sizeof(spare_ephemeral));
|
||||
return err;
|
||||
}
|
||||
#endif // USE_NOISE_SPARE_EPHEMERAL
|
||||
|
||||
const LogString *noise_err_to_logstr(int err) {
|
||||
if (err == NOISE_ERROR_NO_MEMORY)
|
||||
return LOG_STR("NO_MEMORY");
|
||||
|
||||
@@ -6,6 +6,9 @@
|
||||
#include <cstdint>
|
||||
#include "esphome/core/log.h"
|
||||
|
||||
// noise-c handshake state; the full definition lives in <noise/protocol.h>
|
||||
using NoiseHandshakeState = struct NoiseHandshakeState_s;
|
||||
|
||||
namespace esphome::noise {
|
||||
|
||||
using psk_t = std::array<uint8_t, 32>;
|
||||
@@ -38,6 +41,26 @@ class NoiseContext {
|
||||
/// Convert a noise error code to a readable error
|
||||
const LogString *noise_err_to_logstr(int err);
|
||||
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
// One responder ephemeral key pair generated ahead of time (about 60 ms on
|
||||
// ESP8266), refilled by the api server while idle and consumed by the next
|
||||
// handshake of any noise transport; an empty slot means the handshake
|
||||
// generates its own key. The private key stays in RAM until consumed; it is
|
||||
// not wiped on shutdown.
|
||||
// Private key then public key; zero when empty
|
||||
static constexpr size_t SPARE_EPHEMERAL_KEY_SIZE = 32;
|
||||
static constexpr size_t SPARE_EPHEMERAL_SIZE = 2 * SPARE_EPHEMERAL_KEY_SIZE;
|
||||
extern uint8_t spare_ephemeral[SPARE_EPHEMERAL_SIZE]; // NOLINT(cppcoreguidelines-avoid-non-const-global-variables)
|
||||
// Polled every api loop tick, so it must inline. A clamped X25519 private key
|
||||
// always has bit 254 set, so that byte doubles as the ready flag.
|
||||
inline bool has_spare_ephemeral() { return (spare_ephemeral[SPARE_EPHEMERAL_KEY_SIZE - 1] & 0x40) != 0; }
|
||||
/// Fill the slot; blocks for the base point multiply
|
||||
void prepare_spare_ephemeral();
|
||||
/// Hand the slot's key pair to a handshake that has not started and wipe the
|
||||
/// slot; 0 when the slot was empty or the key was taken, else the noise-c error
|
||||
int consume_spare_ephemeral(NoiseHandshakeState *state);
|
||||
#endif
|
||||
|
||||
// Shared wire format for the noise transports (api and ota): every frame is
|
||||
// FRAME_INDICATOR, a 16-bit big-endian payload length, then the payload.
|
||||
// Handshake payloads start with a status byte; transport payloads end with
|
||||
|
||||
@@ -57,6 +57,13 @@ int NoiseResponderHandshake::init(const NoiseContext &ctx, const uint8_t *prolog
|
||||
HANDSHAKE_STEP_LOG("noise_handshakestate_set_prologue", err);
|
||||
return this->fail_init_(err);
|
||||
}
|
||||
#ifdef USE_NOISE_SPARE_EPHEMERAL
|
||||
err = consume_spare_ephemeral(this->handshake_);
|
||||
// Not fatal: the handshake generates its own key instead
|
||||
if (err != 0) {
|
||||
HANDSHAKE_STEP_LOG("noise_handshakestate_set_local_ephemeral", err);
|
||||
}
|
||||
#endif
|
||||
err = noise_handshakestate_start(this->handshake_);
|
||||
if (err != 0) {
|
||||
HANDSHAKE_STEP_LOG("noise_handshakestate_start", err);
|
||||
|
||||
@@ -37,7 +37,8 @@ class NoiseResponderHandshake {
|
||||
NoiseResponderHandshake &operator=(const NoiseResponderHandshake &) = delete;
|
||||
|
||||
/// Create and start the handshake with the context's PSK and the prologue.
|
||||
/// A repeated call frees the previous handshake state and starts over.
|
||||
/// A repeated call frees the previous handshake state and starts over. A
|
||||
/// spare ephemeral key, when one is ready, is used instead of generating.
|
||||
[[nodiscard]] int init(const NoiseContext &ctx, const uint8_t *prologue, size_t prologue_len);
|
||||
/// ACTION_FAILED is the catch-all: returned before init(), after split()
|
||||
/// has released the state, and when noise-c reports a failed handshake.
|
||||
|
||||
@@ -71,7 +71,6 @@
|
||||
#define USE_ESP32_HOSTED
|
||||
#define USE_ESP32_HOSTED_HTTP_UPDATE
|
||||
#define USE_ESP32_IMPROV_STATE_CALLBACK
|
||||
#define USE_ESP_NOW_HOSTED
|
||||
#define USE_EVENT
|
||||
#define USE_FAN
|
||||
#define USE_GPIO_BINARY_SENSOR_INTERRUPT
|
||||
@@ -231,6 +230,7 @@
|
||||
#define USE_IMPROV_SERIAL_NEXT_URL
|
||||
#define USE_MD5
|
||||
#define USE_NOISE
|
||||
#define USE_NOISE_SPARE_EPHEMERAL
|
||||
#define USE_SHA256
|
||||
#ifndef USE_RP2 // no MQTT backend or esp_wireguard library on RP2
|
||||
#define USE_MQTT
|
||||
|
||||
@@ -23,7 +23,6 @@ from esphome.net_retry import (
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from filelock import FileLock
|
||||
import requests
|
||||
|
||||
PathType = str | os.PathLike
|
||||
@@ -910,61 +909,6 @@ def _part_path(dest: Path) -> Path:
|
||||
return dest.with_name(dest.name + ".part")
|
||||
|
||||
|
||||
def downloaded_bytes(dest: Path, size: int | None = None) -> int:
|
||||
"""Bytes of ``dest`` on disk (its ``.part`` while streaming), capped at ``size``."""
|
||||
done = 0
|
||||
for candidate in (_part_path(dest), dest):
|
||||
try:
|
||||
done = candidate.stat().st_size
|
||||
break
|
||||
except FileNotFoundError:
|
||||
continue
|
||||
return done if size is None else min(done, size)
|
||||
|
||||
|
||||
# Short lock-acquire slices so a waiting worker still observes Ctrl-C
|
||||
_DOWNLOAD_LOCK_POLL = 1
|
||||
|
||||
# Waiting on another process's download; past this the caller leaves the
|
||||
# file to its holder (the later sequential install waits on the same lock)
|
||||
DOWNLOAD_LOCK_TIMEOUT = 60
|
||||
|
||||
|
||||
class DownloadLockUnavailable(OSError):
|
||||
"""The lock file cannot be used at all (a lock-less filesystem)."""
|
||||
|
||||
|
||||
def wait_for_download_lock(
|
||||
lock: "FileLock",
|
||||
tracker: Callable[[int], None],
|
||||
on_disk: Callable[[], int],
|
||||
name: str,
|
||||
) -> None:
|
||||
"""Acquire ``lock``, reporting ``on_disk()`` to ``tracker`` each poll so the
|
||||
bar follows the holder's download. Raises filelock's ``Timeout`` once
|
||||
``DOWNLOAD_LOCK_TIMEOUT`` seconds pass."""
|
||||
from filelock import Timeout
|
||||
|
||||
deadline = time.monotonic() + DOWNLOAD_LOCK_TIMEOUT
|
||||
waiting = False
|
||||
while True:
|
||||
try:
|
||||
lock.acquire(timeout=_DOWNLOAD_LOCK_POLL)
|
||||
return
|
||||
except Timeout:
|
||||
pass
|
||||
except OSError as err:
|
||||
# Distinct from an OSError out of on_disk(), which must not
|
||||
# read as "locks unsupported"
|
||||
raise DownloadLockUnavailable(*err.args) from err
|
||||
if not waiting:
|
||||
waiting = True
|
||||
_LOGGER.info("Waiting for another process downloading %s", name)
|
||||
tracker(on_disk()) # raises when the batch is cancelled
|
||||
if time.monotonic() >= deadline:
|
||||
raise Timeout(lock.lock_file)
|
||||
|
||||
|
||||
def discard_partial_download(dest: Path) -> None:
|
||||
"""Remove ``dest`` and the resume sidecars of an abandoned download."""
|
||||
part = _part_path(dest)
|
||||
@@ -1375,7 +1319,10 @@ def download_from_mirrors(
|
||||
)
|
||||
# Tick with the bytes already on disk so a combined bar holds
|
||||
# steady during the backoff instead of rewinding to zero
|
||||
done = downloaded_bytes(path_target) if progress is not None else 0
|
||||
done = 0
|
||||
if progress is not None:
|
||||
part = _part_path(path_target)
|
||||
done = part.stat().st_size if part.is_file() else 0
|
||||
_cancellable_sleep(delay, progress, done)
|
||||
|
||||
# 3. Report every attempted URL if all mirrors failed. failures spans
|
||||
|
||||
@@ -33,14 +33,11 @@ import time
|
||||
from typing import Any, NamedTuple
|
||||
|
||||
from esphome.framework_helpers import (
|
||||
DownloadLockUnavailable,
|
||||
content_length,
|
||||
discard_partial_download,
|
||||
downloaded_bytes,
|
||||
failure_reason,
|
||||
resume_fetch_job,
|
||||
run_batch_downloads,
|
||||
wait_for_download_lock,
|
||||
warn_prefetch_failures,
|
||||
)
|
||||
from esphome.helpers import get_bool_env, get_usable_cpu_count, rmtree
|
||||
@@ -64,10 +61,16 @@ _RESOLVE_WORKERS = 8
|
||||
# A hung child must not block the build; downloads resume on the next run
|
||||
_PREFETCH_TIMEOUT = 20 * 60
|
||||
|
||||
# Waiting on another process's URL download; past this, leave it to pio
|
||||
_DOWNLOAD_LOCK_TIMEOUT = 60
|
||||
|
||||
# Child exit for a handled, already-warned failure; 1 would collide with
|
||||
# the interpreter's own import-failure exit
|
||||
_EXIT_HANDLED = 3
|
||||
|
||||
# Short lock-acquire slices so a waiting worker still observes Ctrl-C
|
||||
_URI_LOCK_POLL = 1
|
||||
|
||||
# Resolution errored (vs a clean skip); suppresses the warm sentinel
|
||||
_RESOLVE_FAILED = object()
|
||||
|
||||
@@ -459,26 +462,17 @@ def _uri_jobs(
|
||||
|
||||
|
||||
def _serialized_fetch_job(
|
||||
dl_path: Path,
|
||||
lock_path: str,
|
||||
body: Any,
|
||||
size: int,
|
||||
stream_dest: Path | None = None,
|
||||
unlocked_ok: bool = True,
|
||||
dl_path: Path, lock_path: str, body: Any, unlocked_ok: bool = True
|
||||
) -> Any:
|
||||
"""Wrap ``body`` so the shared destination is single-writer (interleaved
|
||||
writers truncate each other's ``.part``, see registry.py). A blown deadline
|
||||
is a clean skip. On a lock-less filesystem a sha256-verified body runs
|
||||
unlocked with one warning; a checksum-less one (``unlocked_ok=False``) fails.
|
||||
"""
|
||||
"""Wrap ``body`` so the shared destination is single-writer.
|
||||
|
||||
def on_disk() -> int:
|
||||
# A URL job's holder streams beside the staging path until it
|
||||
# promotes; after that only dl_path is left
|
||||
done = downloaded_bytes(dl_path, size)
|
||||
if not done and stream_dest is not None:
|
||||
done = downloaded_bytes(stream_dest, size)
|
||||
return done
|
||||
Interleaved writers truncate each other's ``.part`` bytes (see
|
||||
registry.py). The bounded poll observes Ctrl-C via the tracker; a
|
||||
blown deadline is a clean skip (the holder's copy is what the build
|
||||
needs). On a lock-less filesystem a sha256-verified body runs
|
||||
unlocked with one warning; a checksum-less one
|
||||
(``unlocked_ok=False``) is a counted failure instead.
|
||||
"""
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
from filelock import FileLock, Timeout
|
||||
@@ -486,27 +480,33 @@ def _serialized_fetch_job(
|
||||
# fallback_to_soft would leave a stale marker on lock-less
|
||||
# filesystems that blocks every later build (see git.py)
|
||||
lock = FileLock(lock_path, fallback_to_soft=False)
|
||||
try:
|
||||
wait_for_download_lock(lock, tracker, on_disk, dl_path.name)
|
||||
except Timeout:
|
||||
# The holder's copy is what the build needs (a large
|
||||
# framework archive can outlast this deadline)
|
||||
_LOGGER.debug("Leaving %s to its current downloader", dl_path.name)
|
||||
return
|
||||
except DownloadLockUnavailable as err:
|
||||
if not unlocked_ok:
|
||||
# A body with no checksum to catch interleaved corruption
|
||||
raise
|
||||
lock = None
|
||||
_LOGGER.warning(
|
||||
"Could not lock %s (%s); downloading unlocked",
|
||||
dl_path.name,
|
||||
err,
|
||||
)
|
||||
deadline = time.monotonic() + _DOWNLOAD_LOCK_TIMEOUT
|
||||
while True:
|
||||
try:
|
||||
lock.acquire(timeout=_URI_LOCK_POLL)
|
||||
break
|
||||
except Timeout:
|
||||
tracker(0) # raises when the batch is cancelled
|
||||
if time.monotonic() >= deadline:
|
||||
# Another process is fetching this same file; its copy
|
||||
# is what the build needs (a large framework archive
|
||||
# can hold the lock far longer than this deadline)
|
||||
_LOGGER.debug("Leaving %s to its current downloader", dl_path.name)
|
||||
return
|
||||
except OSError as err:
|
||||
if not unlocked_ok:
|
||||
# A body with no checksum to catch interleaved corruption
|
||||
raise
|
||||
lock = None
|
||||
_LOGGER.warning(
|
||||
"Could not lock %s (%s); downloading unlocked",
|
||||
dl_path.name,
|
||||
err,
|
||||
)
|
||||
break
|
||||
try:
|
||||
if dl_path.is_file():
|
||||
tracker(size) # another process finished it while we waited
|
||||
return
|
||||
return # another process finished it while we waited
|
||||
body(tracker)
|
||||
finally:
|
||||
if lock is not None:
|
||||
@@ -540,7 +540,6 @@ def _registry_fetch_job(
|
||||
dl_path,
|
||||
f"{dl_path}.esphome.lock",
|
||||
resume_fetch_job(url, dl_path, sha256=checksum, size=size),
|
||||
size,
|
||||
)
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
@@ -572,9 +571,9 @@ def _uri_fetch_job(manager: Any, url: str, dl_path: Path, size: int) -> Any:
|
||||
tmp.replace(dl_path)
|
||||
|
||||
def run(tracker: Any) -> None:
|
||||
_serialized_fetch_job(
|
||||
dl_path, f"{tmp}.lock", promote, size, tmp, unlocked_ok=False
|
||||
)(tracker)
|
||||
_serialized_fetch_job(dl_path, f"{tmp}.lock", promote, unlocked_ok=False)(
|
||||
tracker
|
||||
)
|
||||
if dl_path.is_file():
|
||||
# Won or lost, the race is over; staging files left behind
|
||||
# are dead weight PlatformIO's cache never prunes
|
||||
|
||||
@@ -17,10 +17,8 @@ from esphome.framework_helpers import (
|
||||
archive_extract_all,
|
||||
download_from_mirrors,
|
||||
download_with_resume,
|
||||
downloaded_bytes,
|
||||
rmdir,
|
||||
run_batch_downloads,
|
||||
wait_for_download_lock,
|
||||
)
|
||||
from esphome.net_retry import fetch_with_retry, http_request
|
||||
|
||||
@@ -166,17 +164,11 @@ class _PendingArchive(NamedTuple):
|
||||
name: str
|
||||
version: str
|
||||
dest: Path
|
||||
archive: Path
|
||||
url: str
|
||||
sha256: str
|
||||
size: int
|
||||
|
||||
|
||||
def _archive_path(downloads_dir: Path, name: str, version: str) -> Path:
|
||||
"""The one archive path the prefetch and the sequential install share."""
|
||||
return downloads_dir / f"{name}-{version}"
|
||||
|
||||
|
||||
def _already_installed(dest: Path) -> bool:
|
||||
"""Whether ``dest`` holds a completed install (extraction marker)."""
|
||||
return (dest / ".esphome_extracted").is_file()
|
||||
@@ -195,18 +187,18 @@ def prefetch_packages(
|
||||
lock as ``install_package``: the archive's ``.part`` file is shared, and
|
||||
two concurrent writers would truncate each other's bytes.
|
||||
"""
|
||||
from filelock import FileLock, Timeout
|
||||
from filelock import FileLock
|
||||
|
||||
pending: list[_PendingArchive] = []
|
||||
seen: set[Path] = set()
|
||||
seen: set[str] = set()
|
||||
for name, version, dest, mirrors in packages:
|
||||
if mirrors or (dest / ".esphome_extracted").is_file():
|
||||
continue
|
||||
archive = _archive_path(downloads_dir, name, version)
|
||||
if archive in seen:
|
||||
archive_name = f"{name}-{version}"
|
||||
if archive_name in seen:
|
||||
# A duplicate entry would race itself between two workers
|
||||
continue
|
||||
seen.add(archive)
|
||||
seen.add(archive_name)
|
||||
try:
|
||||
url, sha256, size = registry_download(name, version)
|
||||
except EsphomeError as err:
|
||||
@@ -215,9 +207,10 @@ def prefetch_packages(
|
||||
continue
|
||||
if not size:
|
||||
continue
|
||||
archive = downloads_dir / archive_name
|
||||
if archive.is_file() and archive.stat().st_size == size:
|
||||
continue
|
||||
pending.append(_PendingArchive(name, version, dest, archive, url, sha256, size))
|
||||
pending.append(_PendingArchive(name, version, dest, url, sha256, size))
|
||||
if len(pending) < 2:
|
||||
return
|
||||
downloads_dir.mkdir(parents=True, exist_ok=True)
|
||||
@@ -229,36 +222,20 @@ def prefetch_packages(
|
||||
|
||||
def _fetch(entry: _PendingArchive, tracker: Callable[[int], None]) -> None:
|
||||
entry.dest.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
def on_disk() -> int:
|
||||
if done := downloaded_bytes(entry.archive, entry.size):
|
||||
return done
|
||||
# The holder deletes the archive once it has installed it
|
||||
return entry.size if _already_installed(entry.dest) else 0
|
||||
|
||||
lock = FileLock(f"{entry.dest}.lock", fallback_to_soft=False)
|
||||
try:
|
||||
wait_for_download_lock(lock, tracker, on_disk, entry.name)
|
||||
except Timeout:
|
||||
# install_package waits on this same lock and verifies the
|
||||
# holder's copy
|
||||
_LOGGER.debug("Leaving %s to its current downloader", entry.name)
|
||||
return
|
||||
try:
|
||||
if _already_installed(entry.dest):
|
||||
# A concurrent build installed it while we waited; a
|
||||
# re-download would orphan a fresh copy in downloads_dir
|
||||
tracker(entry.size)
|
||||
return
|
||||
download_with_resume(
|
||||
entry.url,
|
||||
entry.archive,
|
||||
sha256=entry.sha256,
|
||||
size=entry.size,
|
||||
progress=tracker,
|
||||
)
|
||||
finally:
|
||||
lock.release()
|
||||
with FileLock(f"{entry.dest}.lock", fallback_to_soft=False):
|
||||
# Marker re-check: a concurrent build may have installed (and
|
||||
# deleted the archive of) this package while we waited;
|
||||
# re-downloading would orphan a fresh copy in downloads_dir
|
||||
# no branch: the thread tracer misses the skip edge; both
|
||||
# arms of _already_installed are pinned directly
|
||||
if not _already_installed(entry.dest): # pragma: no branch
|
||||
download_with_resume(
|
||||
entry.url,
|
||||
downloads_dir / f"{entry.name}-{entry.version}",
|
||||
sha256=entry.sha256,
|
||||
size=entry.size,
|
||||
progress=tracker,
|
||||
)
|
||||
|
||||
failures = run_batch_downloads(
|
||||
"Downloading packages",
|
||||
@@ -311,7 +288,7 @@ def install_package(
|
||||
rmdir(dest, msg=f"Clean up incomplete {name} install")
|
||||
# Persistent location so an interrupted download resumes across runs.
|
||||
downloads_dir.mkdir(parents=True, exist_ok=True)
|
||||
archive = _archive_path(downloads_dir, name, version)
|
||||
archive = downloads_dir / f"{name}-{version}"
|
||||
_LOGGER.info("Downloading %s %s ...", name, version)
|
||||
if mirrors:
|
||||
_LOGGER.warning(
|
||||
|
||||
+1
-16
@@ -294,9 +294,6 @@ def highlight(s):
|
||||
"esphome/components/socket/headers.h",
|
||||
"esphome/core/defines.h",
|
||||
"esphome/components/http_request/httplib.h",
|
||||
# Shared C wire header (byte-identical with the co-processor firmware);
|
||||
# these are protocol constants and constexpr is C++-only.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
def lint_no_defines(fname, match):
|
||||
@@ -819,10 +816,6 @@ def lint_relative_py_import(fname: Path, line, col, content):
|
||||
"esphome/components/host/helpers.cpp",
|
||||
"esphome/components/zephyr/helpers.cpp",
|
||||
"esphome/components/http_request/httplib.h",
|
||||
# Global extern "C" esp_now_* linker symbols + shared C wire header;
|
||||
# neither can live in a C++ namespace.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted.cpp",
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
def lint_namespace(fname: Path, content: str) -> str | None:
|
||||
@@ -848,15 +841,7 @@ def lint_esphome_h(fname, line, col, content):
|
||||
)
|
||||
|
||||
|
||||
@lint_content_check(
|
||||
include=["*.h"],
|
||||
exclude=[
|
||||
"esphome/core/entity_types.h",
|
||||
# Shared C wire header; uses a classic #ifndef guard for portability
|
||||
# across the co-processor firmware repo it stays byte-identical with.
|
||||
"esphome/components/esp32_hosted/esp_now_hosted_rpc.h",
|
||||
],
|
||||
)
|
||||
@lint_content_check(include=["*.h"], exclude=["esphome/core/entity_types.h"])
|
||||
def lint_pragma_once(fname, content):
|
||||
if "#pragma once" not in content:
|
||||
return (
|
||||
|
||||
@@ -1,5 +0,0 @@
|
||||
# Exercises the ESP-NOW-over-hosted shim: on the ESP32-P4 host, esp32_hosted
|
||||
# supplies the esp_now_* symbols that the espnow component links against.
|
||||
packages:
|
||||
esp32_hosted: !include common.yaml
|
||||
espnow: !include ../espnow/common.yaml
|
||||
@@ -1,3 +1,4 @@
|
||||
import esphome.codegen as cg
|
||||
from tests.testing_helpers import ComponentManifestOverride
|
||||
|
||||
|
||||
@@ -5,3 +6,10 @@ def override_manifest(manifest: ComponentManifestOverride) -> None:
|
||||
# to_code must run: it defines USE_NOISE and adds the noise-c library
|
||||
# the component sources under test need.
|
||||
manifest.enable_codegen()
|
||||
real_to_code = manifest.to_code
|
||||
|
||||
async def to_code_testing(config):
|
||||
await real_to_code(config)
|
||||
cg.add_define("USE_NOISE_SPARE_EPHEMERAL")
|
||||
|
||||
manifest.to_code = to_code_testing
|
||||
|
||||
@@ -157,6 +157,74 @@ TEST(NoiseResponderHandshakeTest, FullHandshakeAndTransportRoundTrip) {
|
||||
noise_cipherstate_free(recv_cipher);
|
||||
}
|
||||
|
||||
// Drive one full NNpsk0 handshake between a fresh initiator and responder;
|
||||
// responder_e receives the ephemeral public key the responder put on the
|
||||
// wire (the clear text start of its message, taken before the initiator
|
||||
// consumes the buffer in place)
|
||||
static void run_handshake(NoiseResponderHandshake &responder, uint8_t responder_e[SPARE_EPHEMERAL_KEY_SIZE]) {
|
||||
const psk_t psk = make_psk(7);
|
||||
ASSERT_EQ(responder.init(ctx_for(psk), PROLOGUE, sizeof(PROLOGUE)), 0);
|
||||
Initiator initiator(psk, PROLOGUE, sizeof(PROLOGUE));
|
||||
uint8_t msg[MAX_HANDSHAKE_SIZE];
|
||||
size_t msg_len = initiator.write_message(msg, sizeof(msg));
|
||||
ASSERT_EQ(responder.read_message(msg, msg_len), 0);
|
||||
size_t reply_len = 0;
|
||||
ASSERT_EQ(responder.write_message(msg, sizeof(msg), reply_len), 0);
|
||||
ASSERT_GE(reply_len, SPARE_EPHEMERAL_KEY_SIZE);
|
||||
std::memcpy(responder_e, msg, SPARE_EPHEMERAL_KEY_SIZE);
|
||||
ASSERT_EQ(initiator.read_message(msg, reply_len), 0);
|
||||
ASSERT_EQ(responder.action(), Action::ACTION_SPLIT);
|
||||
}
|
||||
|
||||
TEST(SpareEphemeralTest, EmptySlotLeavesHandshakeToGenerate) {
|
||||
ASSERT_FALSE(has_spare_ephemeral());
|
||||
NoiseResponderHandshake responder;
|
||||
uint8_t responder_e[SPARE_EPHEMERAL_KEY_SIZE];
|
||||
run_handshake(responder, responder_e);
|
||||
EXPECT_FALSE(has_spare_ephemeral());
|
||||
}
|
||||
|
||||
TEST(SpareEphemeralTest, ConsumeHandsTheKeyToANewState) {
|
||||
prepare_spare_ephemeral();
|
||||
ASSERT_TRUE(has_spare_ephemeral());
|
||||
const NoiseProtocolId nid = {
|
||||
.prefix_id = NOISE_PREFIX_STANDARD,
|
||||
.pattern_id = NOISE_PATTERN_NN,
|
||||
.modifier_ids = {NOISE_MODIFIER_PSK0},
|
||||
.dh_id = NOISE_DH_CURVE25519,
|
||||
.cipher_id = NOISE_CIPHER_CHACHAPOLY,
|
||||
.hash_id = NOISE_HASH_SHA256,
|
||||
.hybrid_id = NOISE_DH_NONE,
|
||||
};
|
||||
NoiseHandshakeState *state = nullptr;
|
||||
ASSERT_EQ(noise_handshakestate_new_by_id(&state, &nid, NOISE_ROLE_RESPONDER), 0);
|
||||
const psk_t psk = make_psk(7);
|
||||
ASSERT_EQ(noise_handshakestate_set_pre_shared_key(state, psk.data(), psk.size()), 0);
|
||||
ASSERT_EQ(noise_handshakestate_set_prologue(state, PROLOGUE, sizeof(PROLOGUE)), 0);
|
||||
EXPECT_EQ(consume_spare_ephemeral(state), 0);
|
||||
EXPECT_FALSE(has_spare_ephemeral());
|
||||
noise_handshakestate_free(state);
|
||||
}
|
||||
|
||||
TEST(SpareEphemeralTest, SlotKeyIsOnTheWireAndConsumedOnce) {
|
||||
prepare_spare_ephemeral();
|
||||
ASSERT_TRUE(has_spare_ephemeral());
|
||||
uint8_t expected_pub[SPARE_EPHEMERAL_KEY_SIZE];
|
||||
std::memcpy(expected_pub, spare_ephemeral + SPARE_EPHEMERAL_KEY_SIZE, sizeof(expected_pub));
|
||||
|
||||
NoiseResponderHandshake first;
|
||||
uint8_t responder_e[SPARE_EPHEMERAL_KEY_SIZE];
|
||||
run_handshake(first, responder_e);
|
||||
// The spare, not a generated key, went out; and it went out once
|
||||
EXPECT_EQ(std::memcmp(responder_e, expected_pub, sizeof(expected_pub)), 0);
|
||||
EXPECT_FALSE(has_spare_ephemeral());
|
||||
|
||||
NoiseResponderHandshake second;
|
||||
run_handshake(second, responder_e);
|
||||
EXPECT_NE(std::memcmp(responder_e, expected_pub, sizeof(expected_pub)), 0);
|
||||
EXPECT_FALSE(has_spare_ephemeral());
|
||||
}
|
||||
|
||||
TEST(NoiseResponderHandshakeTest, ReInitRestartsHandshake) {
|
||||
// The documented retry shape: a repeated init() frees the previous state
|
||||
// and starts over. The first message under the new key authenticating
|
||||
|
||||
@@ -1,48 +0,0 @@
|
||||
"""Tests for the espnow component's final validation."""
|
||||
|
||||
import pytest
|
||||
|
||||
from esphome.components.esp32.const import (
|
||||
VARIANT_ESP32C3,
|
||||
VARIANT_ESP32H2,
|
||||
VARIANT_ESP32P4,
|
||||
)
|
||||
from esphome.components.espnow import _validate_variant
|
||||
import esphome.config_validation as cv
|
||||
import esphome.final_validate as fv
|
||||
from esphome.types import ConfigType
|
||||
|
||||
|
||||
def _run(
|
||||
monkeypatch, variant: str, full_config: dict, config: ConfigType
|
||||
) -> ConfigType:
|
||||
monkeypatch.setattr("esphome.components.espnow.get_esp32_variant", lambda: variant)
|
||||
token = fv.full_config.set(full_config)
|
||||
try:
|
||||
return _validate_variant(config)
|
||||
finally:
|
||||
fv.full_config.reset(token)
|
||||
|
||||
|
||||
def test_variant_with_native_wifi_passes(monkeypatch) -> None:
|
||||
"""A variant with a native Wi-Fi PHY needs no shim; config passes through."""
|
||||
config = {"id": "espnow"}
|
||||
assert _run(monkeypatch, VARIANT_ESP32C3, {}, config) is config
|
||||
|
||||
|
||||
def test_radioless_non_p4_variant_rejected(monkeypatch) -> None:
|
||||
"""Radio-less variants without any ESP-NOW path are rejected outright."""
|
||||
with pytest.raises(cv.Invalid, match="not supported"):
|
||||
_run(monkeypatch, VARIANT_ESP32H2, {}, {})
|
||||
|
||||
|
||||
def test_p4_without_esp32_hosted_rejected(monkeypatch) -> None:
|
||||
"""The P4 needs the esp32_hosted shim to supply the esp_now_* symbols."""
|
||||
with pytest.raises(cv.Invalid, match="esp32_hosted"):
|
||||
_run(monkeypatch, VARIANT_ESP32P4, {}, {})
|
||||
|
||||
|
||||
def test_p4_with_esp32_hosted_passes(monkeypatch) -> None:
|
||||
"""The P4 with esp32_hosted present validates; config passes through."""
|
||||
config = {"id": "espnow"}
|
||||
assert _run(monkeypatch, VARIANT_ESP32P4, {"esp32_hosted": {}}, config) is config
|
||||
@@ -9,7 +9,7 @@ not be part of a unit test suite.
|
||||
|
||||
"""
|
||||
|
||||
from collections.abc import Callable, Generator
|
||||
from collections.abc import Generator
|
||||
import os
|
||||
from pathlib import Path
|
||||
import sys
|
||||
@@ -137,40 +137,3 @@ def mock_get_component() -> Generator[Mock, None, None]:
|
||||
"""Mock get_component for config module."""
|
||||
with patch("esphome.config.get_component") as mock:
|
||||
yield mock
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def held_lock() -> Callable[..., Callable[..., None]]:
|
||||
"""Factory for a ``FileLock.acquire`` fake held by another downloader.
|
||||
|
||||
Each poll writes the next chunk to ``part`` (or runs it, for a callable)
|
||||
and raises ``Timeout``; when the chunks run out the part is removed,
|
||||
``land()`` runs, and the acquire succeeds (also for any later job, so
|
||||
``land`` must be idempotent).
|
||||
"""
|
||||
from filelock import Timeout
|
||||
|
||||
def make(
|
||||
part: Path,
|
||||
chunks: list[bytes | Callable[[], None]],
|
||||
land: Callable[[], None],
|
||||
) -> Callable[..., None]:
|
||||
polls = iter(chunks)
|
||||
|
||||
def acquire(*args, **kwargs) -> None:
|
||||
try:
|
||||
chunk = next(polls)
|
||||
except StopIteration:
|
||||
part.unlink(missing_ok=True)
|
||||
land()
|
||||
return
|
||||
if callable(chunk):
|
||||
chunk()
|
||||
else:
|
||||
part.parent.mkdir(parents=True, exist_ok=True)
|
||||
part.write_bytes(chunk)
|
||||
raise Timeout("held")
|
||||
|
||||
return acquire
|
||||
|
||||
return make
|
||||
|
||||
@@ -2353,20 +2353,3 @@ def test_discard_partial_download_logs_undeletable(
|
||||
):
|
||||
framework_helpers.discard_partial_download(dest)
|
||||
assert "Could not remove" in caplog.text
|
||||
|
||||
|
||||
def test_downloaded_bytes_reports_what_is_on_disk(tmp_path: Path) -> None:
|
||||
"""Part file first, then the landed file, both capped at size; else 0."""
|
||||
dest = tmp_path / "archive"
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 0
|
||||
part = tmp_path / "archive.part"
|
||||
part.write_bytes(b"ab")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 2
|
||||
part.write_bytes(b"abcdef")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 4
|
||||
part.unlink()
|
||||
dest.write_bytes(b"abc")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 3
|
||||
assert framework_helpers.downloaded_bytes(dest) == 3
|
||||
dest.write_bytes(b"abcdef")
|
||||
assert framework_helpers.downloaded_bytes(dest, 4) == 4
|
||||
|
||||
@@ -454,96 +454,23 @@ def test_uri_fetch_job_waits_out_a_briefly_held_lock(tmp_path: Path) -> None:
|
||||
assert dl_path.read_bytes() == b"data"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("staged", [b"", b"ab"])
|
||||
def test_lock_deadline_leaves_download_to_the_holder(
|
||||
tmp_path: Path, staged: bytes
|
||||
) -> None:
|
||||
"""A lock held past the deadline is another process's download; skip
|
||||
cleanly, polling the tracker with what the holder has staged so far."""
|
||||
def test_lock_deadline_leaves_download_to_the_holder(tmp_path: Path) -> None:
|
||||
"""A lock held past the deadline means another process is fetching the
|
||||
same file; skipping cleanly beats a misleading failure warning. The
|
||||
tracker is still polled so a parked worker observes cancellation."""
|
||||
dl_path = tmp_path / "archive"
|
||||
(tmp_path / "archive.prefetch.part").write_bytes(staged)
|
||||
ticks: list[int] = []
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch.object(pf, "_DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
):
|
||||
pf._uri_fetch_job(MagicMock(), "https://x/a.zip", dl_path, 4)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == [len(staged)]
|
||||
assert ticks == [0]
|
||||
assert not dl_path.exists()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("job", "part_name", "chunks", "expected"),
|
||||
[
|
||||
(
|
||||
lambda dl_path: pf._registry_fetch_job(
|
||||
MagicMock(), "https://x/a.tar.gz", dl_path, "ab" * 32, 4
|
||||
),
|
||||
"archive.part",
|
||||
[b"a", b"abc"],
|
||||
[1, 3, 4],
|
||||
),
|
||||
(
|
||||
lambda dl_path: pf._uri_fetch_job(
|
||||
MagicMock(), "https://x/a.zip", dl_path, 4
|
||||
),
|
||||
"archive.prefetch.part",
|
||||
[b"ab"],
|
||||
[2, 4],
|
||||
),
|
||||
],
|
||||
ids=["registry", "uri"],
|
||||
)
|
||||
def test_lock_wait_reports_the_holders_progress(
|
||||
tmp_path: Path,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
held_lock,
|
||||
job,
|
||||
part_name: str,
|
||||
chunks: list[bytes],
|
||||
expected: list[int],
|
||||
) -> None:
|
||||
"""A waiting job reports the holder's part file (the staging one for a
|
||||
URL job), then the full size once the holder lands the archive."""
|
||||
dl_path = tmp_path / "archive"
|
||||
ticks: list[int] = []
|
||||
acquire = held_lock(
|
||||
tmp_path / part_name, chunks, lambda: dl_path.write_bytes(b"abcd")
|
||||
)
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
caplog.at_level(logging.INFO),
|
||||
):
|
||||
job(dl_path)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == expected
|
||||
assert caplog.text.count("Waiting for another process downloading archive") == 1
|
||||
|
||||
|
||||
def test_uri_lock_wait_prefers_the_landed_archive(tmp_path: Path, held_lock) -> None:
|
||||
"""Between the holder's promotion rename and its release the staging
|
||||
part is gone; the landed cache file is credited instead of 0."""
|
||||
dl_path = tmp_path / "archive"
|
||||
ticks: list[int] = []
|
||||
acquire = held_lock(
|
||||
tmp_path / "archive.prefetch.part",
|
||||
[b"ab", lambda: dl_path.write_bytes(b"abcd")],
|
||||
lambda: None,
|
||||
)
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
):
|
||||
pf._uri_fetch_job(MagicMock(), "https://x/a.zip", dl_path, 4)(ticks.append)
|
||||
mock_download.assert_not_called()
|
||||
assert ticks == [2, 4, 4]
|
||||
|
||||
|
||||
def test_registry_lock_deadline_skips_registration(tmp_path: Path) -> None:
|
||||
"""A registry job that lost the download race to another process
|
||||
must not stamp a nonexistent archive into pio's usage.db."""
|
||||
@@ -552,7 +479,7 @@ def test_registry_lock_deadline_skips_registration(tmp_path: Path) -> None:
|
||||
with (
|
||||
patch("esphome.framework_helpers.download_with_resume") as mock_download,
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch.object(pf, "_DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
):
|
||||
pf._registry_fetch_job(manager, "https://x/a.tar.gz", dl_path, "ab" * 32, 4)(
|
||||
lambda done: None
|
||||
|
||||
@@ -8,7 +8,6 @@ import os
|
||||
from pathlib import Path
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from filelock import Timeout
|
||||
import pytest
|
||||
|
||||
from esphome.core import EsphomeError
|
||||
@@ -541,13 +540,16 @@ def test_prefetch_packages_skips_freshly_installed_dest(tmp_path: Path) -> None:
|
||||
dest = tmp_path / "a"
|
||||
dest.mkdir()
|
||||
|
||||
def marker_appears_under_lock(*args, **kwargs):
|
||||
from contextlib import contextmanager
|
||||
|
||||
@contextmanager
|
||||
def marker_appears_under_lock(path, **kwargs):
|
||||
# Simulates the concurrent build finishing while we waited
|
||||
(dest / ".esphome_extracted").touch()
|
||||
yield
|
||||
|
||||
with (
|
||||
patch("filelock.FileLock.acquire", side_effect=marker_appears_under_lock),
|
||||
patch("filelock.FileLock.release"),
|
||||
patch("filelock.FileLock", side_effect=marker_appears_under_lock),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10})
|
||||
@@ -557,69 +559,6 @@ def test_prefetch_packages_skips_freshly_installed_dest(tmp_path: Path) -> None:
|
||||
mock_download.assert_not_called()
|
||||
|
||||
|
||||
def test_prefetch_packages_waits_with_the_holders_progress(
|
||||
tmp_path: Path, held_lock
|
||||
) -> None:
|
||||
"""A worker parked on another build's lock reports that build's part
|
||||
file, then the full size once the marker appears."""
|
||||
dest = tmp_path / "a"
|
||||
dest.mkdir()
|
||||
ticks: list[int] = []
|
||||
part = tmp_path / "dl" / "a-1.0.part"
|
||||
|
||||
def installed_and_pruned() -> None:
|
||||
# install_package touches the marker, then unlinks the archive
|
||||
(dest / ".esphome_extracted").touch()
|
||||
part.unlink()
|
||||
|
||||
acquire = held_lock(
|
||||
part,
|
||||
[lambda: None, b"abc", installed_and_pruned],
|
||||
(dest / ".esphome_extracted").touch,
|
||||
)
|
||||
|
||||
def fake_batch(header, jobs):
|
||||
for _name, _size, fetch in jobs:
|
||||
fetch(ticks.append)
|
||||
return []
|
||||
|
||||
with (
|
||||
patch("filelock.FileLock.acquire", side_effect=acquire),
|
||||
patch("filelock.FileLock.release"),
|
||||
patch.object(registry, "run_batch_downloads", side_effect=fake_batch),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10, "b": 5})
|
||||
),
|
||||
):
|
||||
registry.prefetch_packages(
|
||||
[("a", "1.0", dest, []), ("b", "2.0", tmp_path / "b", [])],
|
||||
tmp_path / "dl",
|
||||
)
|
||||
assert ticks == [0, 3, 10, 10]
|
||||
mock_download.assert_called_once()
|
||||
|
||||
|
||||
def test_prefetch_packages_leaves_a_long_held_lock_to_its_holder(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
"""Past the deadline the worker skips; install_package waits on the same
|
||||
lock later and verifies whatever the holder produced."""
|
||||
with (
|
||||
patch("filelock.FileLock.acquire", side_effect=Timeout("held")),
|
||||
patch("esphome.framework_helpers.DOWNLOAD_LOCK_TIMEOUT", 0),
|
||||
patch.object(registry, "download_with_resume") as mock_download,
|
||||
patch.object(
|
||||
registry, "registry_download", side_effect=_resolve_for({"a": 10, "b": 5})
|
||||
),
|
||||
):
|
||||
registry.prefetch_packages(
|
||||
[("a", "1.0", tmp_path / "a", []), ("b", "2.0", tmp_path / "b", [])],
|
||||
tmp_path / "dl",
|
||||
)
|
||||
mock_download.assert_not_called()
|
||||
|
||||
|
||||
def test_already_installed_probe(tmp_path: Path) -> None:
|
||||
"""Both arms of the marker probe the prefetch worker keys on."""
|
||||
dest = tmp_path / "pkg"
|
||||
|
||||
Reference in New Issue
Block a user