From e3e3ac6541c411234d9704fab2beb2bb3e36efff Mon Sep 17 00:00:00 2001 From: Zhou Xiao Date: Sun, 6 Sep 2026 16:12:27 +0800 Subject: [PATCH] test(ble_log): soak the park/wake handshake and fix observer latching - concurrent-writer integrity and park/wake soak cases: many writers across shared and reserve transports keep frames intact and SNs unique through continuous park/wake turnover; - the deinit race case wakes parked writers before joining, so the join cannot deadlock on a task that must first observe the closed gate; - the perf observer latches its final counters only from FLUSH snapshots: periodic frames arriving after the FLUSH snapshot carry interval counters that the flush reset, and must not overwrite the final result; - README describes the timestamp prefix by writer path. --- components/bt/common/ble_log/README.md | 24 +- .../main/test_ble_log_perf.c | 17 +- .../main/test_ble_log_runtime.c | 11 +- .../ble_log_test/main/test_ble_log_rt.c | 211 ++++++++++++++++++ 4 files changed, 244 insertions(+), 19 deletions(-) diff --git a/components/bt/common/ble_log/README.md b/components/bt/common/ble_log/README.md index c8506fef349..fa86105ae0c 100644 --- a/components/bt/common/ble_log/README.md +++ b/components/bt/common/ble_log/README.md @@ -60,19 +60,21 @@ The checksum is `ble_log_fast_checksum()` over the six-byte header and exactly `payload length` bytes. It excludes the checksum field and peripheral-only DMA padding. -Core frames (all sources except REDIR) begin their payload with: +Core frame payloads follow the writer entry point, not a single rule: -```text -[4-byte low32 esp_timer_get_time() microseconds][source payload] -``` +- Public write / claim path (`write_hex`, compressed encoder): the API + prepends a `[4-byte low32 esp_timer_get_time() microseconds]` before + the source payload. The timestamp is captured at API entry before any + pool mutex wait. It wraps approximately every 71.6 minutes and must be + unwrapped modulo 2^32 by the receiver. +- LL callback path (`write_hex_ll`, covering the LL_TASK / LL_HCI / + LL_ISR sources): controller payloads are forwarded raw, with no ESP + timestamp prefix — the on-wire layout is exactly what the controller + handed over. -The timestamp is captured at API entry before any pool mutex wait. It wraps -approximately every 71.6 minutes and must be unwrapped modulo 2^32 by the -receiver. - -UART0 redirection payload is the raw console stream with no timestamp -prefix: the redirection stream keeps its own frame sequence, and its -receiver-side arrival time (aggregation delay bounded by the periodic +UART0 redirection payload is likewise the raw console stream with no +timestamp prefix: the redirection stream keeps its own frame sequence, and +its receiver-side arrival time (aggregation delay bounded by the periodic redirection flush) is the alignment reference against the core timeline. ### Sources diff --git a/components/bt/common/ble_log/test_apps/ble_log_perf_test/main/test_ble_log_perf.c b/components/bt/common/ble_log/test_apps/ble_log_perf_test/main/test_ble_log_perf.c index 95d02496988..33f7fb8af86 100644 --- a/components/bt/common/ble_log/test_apps/ble_log_perf_test/main/test_ble_log_perf.c +++ b/components/bt/common/ble_log/test_apps/ble_log_perf_test/main/test_ble_log_perf.c @@ -492,13 +492,18 @@ static void observe_perf_frame(const test_ble_log_frame_t *frame, void *ctx) if (snapshot.int_src_code != BLE_LOG_INT_SRC_SNAPSHOT) { return; } - for (int i = 0; i < BLE_LOG_SRC_CORE_COUNT; i++) { - int src = BLE_LOG_SRC_CORE_FIRST + i; - sink->stat_written[src] = snapshot.stats[i].written_frame_cnt; - sink->stat_lost[src] = snapshot.stats[i].lost_frame_cnt; + /* Latch the final counters only from the FLUSH snapshot: periodic + * snapshots keep arriving after it and carry interval counters + * that the flush reset; letting them through would overwrite the + * final result the run report and assertions rely on. */ + if (snapshot.reason_flags & BLE_LOG_SNAPSHOT_REASON_FLUSH) { + for (int i = 0; i < BLE_LOG_SRC_CORE_COUNT; i++) { + int src = BLE_LOG_SRC_CORE_FIRST + i; + sink->stat_written[src] = snapshot.stats[i].written_frame_cnt; + sink->stat_lost[src] = snapshot.stats[i].lost_frame_cnt; + } + sink->stat_final_seen = true; } - sink->stat_final_seen = - (snapshot.reason_flags & BLE_LOG_SNAPSHOT_REASON_FLUSH) != 0; } } diff --git a/components/bt/common/ble_log/test_apps/ble_log_rt_test/main/test_ble_log_runtime.c b/components/bt/common/ble_log/test_apps/ble_log_rt_test/main/test_ble_log_runtime.c index d1067cf1e9e..87217c87276 100644 --- a/components/bt/common/ble_log/test_apps/ble_log_rt_test/main/test_ble_log_runtime.c +++ b/components/bt/common/ble_log/test_apps/ble_log_rt_test/main/test_ble_log_runtime.c @@ -932,6 +932,13 @@ TEST_CASE("BLE Log runtime survives deinit racing submissions", /* Stop the writer and join with a bound before asserting so a Unity * longjmp cannot leak the writer into later tests. */ __atomic_store_n(&ctx->stop, true, __ATOMIC_RELEASE); + /* A writer parked inside ble_log_claim() cannot observe stop: only the + * deinit waiter wake returns it. Deinit here, before the join, so its + * drain closes the producer gate, wakes any parked writer and waits + * for it to leave the pool API. The vTaskDelete fallback below can + * then never cut a task down inside pool bookkeeping (which would + * leak waiting_task_count and hang the recovery deinit). */ + ble_log_deinit(); const int64_t join_deadline_us = esp_timer_get_time() + (int64_t)RT_JOIN_TIMEOUT_MS * 1000; while (!__atomic_load_n(&ctx->exited, __ATOMIC_ACQUIRE)) { @@ -949,8 +956,8 @@ TEST_CASE("BLE Log runtime survives deinit racing submissions", /* Recover module state before any assertion can abort the test: a * failed re-init leaves the module deinit-ed and tearDown does not - * restore it, which would cascade into every later test. */ - ble_log_deinit(); + * restore it, which would cascade into every later test. Deinit ran + * before the join (see above); only the init half remains. */ bool recovered = ble_log_init(); reinit_ok = reinit_ok && recovered; diff --git a/components/bt/common/ble_log/test_apps/ble_log_test/main/test_ble_log_rt.c b/components/bt/common/ble_log/test_apps/ble_log_test/main/test_ble_log_rt.c index c976ad3c219..02288e3a387 100644 --- a/components/bt/common/ble_log/test_apps/ble_log_test/main/test_ble_log_rt.c +++ b/components/bt/common/ble_log/test_apps/ble_log_test/main/test_ble_log_rt.c @@ -1748,6 +1748,217 @@ static size_t count_unique_sn(const uint32_t *sn, size_t n) return unique; } +TEST_CASE("BLE Log concurrent writers keep frames intact and SNs unique", + "[ble_log][lbm]") +{ + /* Scenario: 4 writers run in parallel across all cores, each emitting + * CONC_FRAMES_EACH marked frames while the pool recycles under them. + * Guards (mutation -> expected failure): + * - per-transport CAS lock weakened to plain test-and-set: two cores + * co-own one transport, frames interleave -> checksum walk fails. + * - Global SN fetch_add weakened to plain increment: concurrent RMWs + * lose updates -> duplicate SNs -> unique-count assertion fails. + * Every completed write must be present exactly once. */ + TEST_ASSERT_TRUE(ble_log_enable(true)); + ble_log_lbm_flush_open_trans(); + for (int round = 0; round < 2; round++) { + TEST_ASSERT_TRUE(ble_log_rt_drain()); + while (ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + 0, 0, NULL) > 0) { + } + } + + conc_writer_ctx_t ctx[CONC_WRITER_CNT]; + static conc_capture_t cap; + memset(&cap, 0, sizeof(cap)); + for (int w = 0; w < CONC_WRITER_CNT; w++) { + ctx[w] = (conc_writer_ctx_t){ + .done = xSemaphoreCreateBinary(), + .marker = (uint8_t)(0xa0 + w), + .result = true, + }; + TEST_ASSERT_NOT_NULL(ctx[w].done); + } + for (int w = 0; w < CONC_WRITER_CNT; w++) { + TEST_ASSERT_EQUAL(pdTRUE, xTaskCreatePinnedToCore( + conc_writer_task, "ble_log_conc", + TEST_LIFECYCLE_STACK_SIZE, &ctx[w], + TEST_LIFECYCLE_PRIO, NULL, + w % portNUM_PROCESSORS)); + } + /* Feed the pool while the writers run; the main task is the + * consumer. */ + for (;;) { + bool all_done = true; + for (int w = 0; w < CONC_WRITER_CNT; w++) { + if (!ctx[w].finished) { + if (xSemaphoreTake(ctx[w].done, 0) == pdTRUE) { + ctx[w].finished = true; + TEST_ASSERT_TRUE(ctx[w].result); + vSemaphoreDelete(ctx[w].done); + } else { + all_done = false; + } + } + } + if (all_done) { + break; + } + TEST_ASSERT_TRUE(ble_log_rt_drain()); + size_t len = ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(20), 0, NULL); + if (len) { + TEST_ASSERT_TRUE(test_ble_log_walk_frames( + s_read_buf, len, capture_conc_frame, &cap)); + } + } + + ble_log_lbm_flush_open_trans(); + TEST_ASSERT_TRUE(ble_log_rt_drain()); + for (;;) { + size_t len = ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(TEST_READ_TIMEOUT_MS), + 0, NULL); + if (!len) { + break; + } + TEST_ASSERT_TRUE(test_ble_log_walk_frames(s_read_buf, len, + capture_conc_frame, &cap)); + } + TEST_ASSERT_EQUAL_size_t(CONC_WRITER_CNT * CONC_FRAMES_EACH, cap.frames); + for (int w = 0; w < CONC_WRITER_CNT; w++) { + TEST_ASSERT_EQUAL_size_t(CONC_FRAMES_EACH, cap.per_marker[w]); + } + /* Unique Global SNs: a lost fetch_add update surfaces as two frames + * sharing one SN. */ + TEST_ASSERT_EQUAL_size_t(CONC_WRITER_CNT * CONC_FRAMES_EACH, + count_unique_sn(cap.sn, cap.sn_count)); + + /* Phase 2 - SN hammer: wait=false claims in tight loops across all + * cores. The SN fetch_add window dominates the failing-claim path, so + * a plain (non-atomic) increment loses updates at a high rate. Losses + * are expected and allowed; only SN uniqueness is asserted. */ + enum { HAMMER_TASKS = 4 }; + hammer_ctx_t hctx[HAMMER_TASKS]; + for (int t = 0; t < HAMMER_TASKS; t++) { + hctx[t] = (hammer_ctx_t){.done = xSemaphoreCreateBinary()}; + TEST_ASSERT_NOT_NULL(hctx[t].done); + TEST_ASSERT_EQUAL(pdTRUE, xTaskCreatePinnedToCore( + hammer_claim_task, "ble_log_hammer", + TEST_LIFECYCLE_STACK_SIZE, &hctx[t], + TEST_LIFECYCLE_PRIO, NULL, + t % portNUM_PROCESSORS)); + } + /* Feed the pool while the hammers run. */ + static hammer_capture_t hcap; + memset(&hcap, 0, sizeof(hcap)); + static bool taken[HAMMER_TASKS]; + memset(taken, 0, sizeof(taken)); + for (int finished = 0; finished < HAMMER_TASKS;) { + for (int t = 0; t < HAMMER_TASKS; t++) { + if (!taken[t] && + xSemaphoreTake(hctx[t].done, 0) == pdTRUE) { + taken[t] = true; + finished++; + vSemaphoreDelete(hctx[t].done); + } + } + TEST_ASSERT_TRUE(ble_log_rt_drain()); + size_t len = ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(2), 0, NULL); + if (len) { + TEST_ASSERT_TRUE(test_ble_log_walk_frames( + s_read_buf, len, capture_hammer_frame, &hcap)); + } + } + ble_log_lbm_flush_open_trans(); + TEST_ASSERT_TRUE(ble_log_rt_drain()); + for (;;) { + size_t len = ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(TEST_READ_TIMEOUT_MS), + 0, NULL); + if (!len) { + break; + } + TEST_ASSERT_TRUE(test_ble_log_walk_frames(s_read_buf, len, + capture_hammer_frame, &hcap)); + } + TEST_ASSERT_GREATER_THAN_size_t(0, hcap.sn_count); + TEST_ASSERT_EQUAL_size_t(hcap.sn_count, + count_unique_sn(hcap.sn, hcap.sn_count)); +} + +typedef struct { + SemaphoreHandle_t done; + bool result; +} soak_writer_ctx_t; + +static void soak_writer_task(void *arg) +{ + soak_writer_ctx_t *ctx = arg; + const uint8_t marker = 0x5c; + ctx->result = ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1); + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +TEST_CASE("BLE Log waiter handshake survives continuous park/wake turnover", + "[ble_log][lbm]") +{ + /* Scenario: each round fills every shared transport, parks two writers, + * then recycles transports one at a time. Each single recycle is the + * last-availability event for exactly one waiter: if the register/ + * publish SC-fence handshake is weakened, a waiter can miss the final + * notification and stay parked while capacity exists. + * Mutation (drop the SEQ_CST fences in the waiter/notify handshake) + * -> probabilistic: some round times out. This is a soak: it raises + * the hit probability, it cannot make the window deterministic. */ + TEST_ASSERT_TRUE(ble_log_enable(true)); + static const uint8_t full[ + BLE_LOG_MAX_PAYLOAD_LEN - sizeof(uint32_t)] = {0}; + + for (int round = 0; round < 64; round++) { + ble_log_lbm_flush_open_trans(); + TEST_ASSERT_TRUE(ble_log_rt_drain()); + while (ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + 0, 0, NULL) > 0) { + } + for (int i = 0; i < BLE_LOG_POOL_SHARED_CNT; i++) { + TEST_ASSERT_TRUE(ble_log_write_hex(BLE_LOG_SRC_CUSTOM, full, + sizeof(full))); + } + soak_writer_ctx_t w[2]; + for (int i = 0; i < 2; i++) { + w[i] = (soak_writer_ctx_t){ + .done = xSemaphoreCreateBinary(), + .result = false, + }; + TEST_ASSERT_NOT_NULL(w[i].done); + TEST_ASSERT_EQUAL(pdTRUE, xTaskCreatePinnedToCore( + soak_writer_task, "ble_log_soak", + TEST_LIFECYCLE_STACK_SIZE, &w[i], + TEST_LIFECYCLE_PRIO, NULL, + i % portNUM_PROCESSORS)); + } + /* Let both writers scan twice and park. */ + vTaskDelay(pdMS_TO_TICKS(20)); + /* Recycle exactly one transport at a time: each read is one + * waiter's last chance to be notified. */ + TEST_ASSERT_TRUE(ble_log_rt_drain()); + for (int i = 0; i < 2; i++) { + size_t len = ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(TEST_READ_TIMEOUT_MS), + 0, NULL); + TEST_ASSERT_GREATER_THAN_size_t(0, len); + } + for (int i = 0; i < 2; i++) { + TEST_ASSERT_TRUE(xSemaphoreTake(w[i].done, pdMS_TO_TICKS(2000))); + TEST_ASSERT_TRUE(w[i].result); + vSemaphoreDelete(w[i].done); + } + } +} + typedef struct { SemaphoreHandle_t done; bool write_result;