diff --git a/components/bt/common/ble_log/README.md b/components/bt/common/ble_log/README.md index 67c75721077..c8506fef349 100644 --- a/components/bt/common/ble_log/README.md +++ b/components/bt/common/ble_log/README.md @@ -19,9 +19,12 @@ flowchart TD The shared pool contains `CONFIG_BLE_LOG_POOL_TRANS_CNT` transports. The last `CONFIG_BLE_LOG_POOL_NON_YIELD_RESERVE_CNT` transports are reserved for ISR and -other contexts that cannot yield. Every yieldable-context writer, from the -public API and claims to the controller LL task, waits for a shared transport; -only ISR and critical-section writers fail fast. +other contexts that cannot yield. Ordinary yieldable task writers, from the +public API and claims to the controller LL task, wait for a shared transport +unless a claim explicitly opts out. Writers on the shared ESP Timer task never +wait for a transport: its callbacks must return so dispatch can progress. +While yieldable, they use shared transports only; ISR and critical-section +writers retain reserve access and fail-fast behavior. The Internal Snapshot and UART0 redirection transports are not members of the bitmap pool. @@ -202,10 +205,11 @@ if (payload) { ``` `ble_log_claim()` reserves the hidden frame header, ESP Timer timestamp, and -checksum. It waits for a shared transport in yieldable contexts when -`wait_for_transport` is true (writer backpressure); pass false for a lossy -fast path from contexts that must not block, such as the shared ESP Timer -task. Non-yieldable contexts (ISR, critical section) fail fast either way. +checksum. It waits for a shared transport in ordinary yieldable tasks when +`wait_for_transport` is true (writer backpressure); pass false to opt out. +The shared ESP Timer task and non-yieldable contexts fail fast either way. +A busy-pool rejection leaves a Global SN gap and increments the source's loss +counter, just like other pool acquisition failures. Every successful claim must be committed exactly once before `ble_log_deinit()`; `ble_log_commit(handle, 0)` cancels it. Handles include a transport generation so a stale handle cannot commit a later claim in the same diff --git a/components/bt/common/ble_log/include/ble_log.h b/components/bt/common/ble_log/include/ble_log.h index 98a7410e7c6..06c4003bd25 100644 --- a/components/bt/common/ble_log/include/ble_log.h +++ b/components/bt/common/ble_log/include/ble_log.h @@ -67,12 +67,12 @@ void ble_log_deinit(void); bool ble_log_enable(bool enable); /* Blocking; call only from a caller-owned task, not an ISR or system callback. */ void ble_log_flush(void); -/* Waits for a shared transport in yieldable contexts; ISR and critical-section callers fail fast. */ +/* Waits for a shared transport in ordinary yieldable tasks. The shared ESP + * Timer task, ISR and critical-section callers fail fast when none is available. */ bool ble_log_write_hex(ble_log_src_t src_code, const uint8_t *addr, size_t len); -/* Same backpressure as ble_log_write_hex(): yieldable claims wait for a - * shared transport when wait_for_transport is true; pass false for a lossy - * fast path from contexts that must not block (system periodic output). - * Non-yieldable contexts (ISR, critical section) fail fast either way. */ +/* Same backpressure as ble_log_write_hex(): ordinary yieldable tasks wait + * when wait_for_transport is true; false opts out. The shared ESP Timer + * task and non-yieldable contexts fail fast regardless of this flag. */ uint8_t *ble_log_claim(ble_log_src_t src_code, size_t max_len, uint32_t *handle, bool wait_for_transport); void ble_log_commit(uint32_t handle, size_t actual_len); diff --git a/components/bt/common/ble_log/src/ble_log_lbm_v2.c b/components/bt/common/ble_log/src/ble_log_lbm_v2.c index 9ac016d0bf3..f25223e036a 100644 --- a/components/bt/common/ble_log/src/ble_log_lbm_v2.c +++ b/components/bt/common/ble_log/src/ble_log_lbm_v2.c @@ -26,6 +26,9 @@ #endif // CONFIG_BT_DUAL_MODE_ARCH #endif /* CONFIG_BLE_LOG_LL_ENABLED && CONFIG_SOC_ESP_NIMBLE_CONTROLLER */ +/* hint: private API; keep this declaration in sync with esp_timer_impl.h. */ +extern TaskHandle_t esp_timer_impl_get_timer_task_handle(void); + /* MACRO */ #define BLE_LOG_POOL_MASK(count) (0xFFFFFFFFu >> (32 - (count))) #define BLE_LOG_POOL_ALL_MASK BLE_LOG_POOL_MASK(BLE_LOG_POOL_TRANS_CNT) @@ -130,6 +133,11 @@ extern void ble_log_test_enable_before_lifecycle_lock_hook(void) __attribute__(( extern void ble_log_test_disable_before_wake_hook(void) __attribute__((weak)); extern void ble_log_test_claim_locked_hook(void) __attribute__((weak)); extern void ble_log_test_init_snapshot_before_acquire_hook(void) __attribute__((weak)); +extern void ble_log_test_flush_hint_seen_hook(uint8_t id) __attribute__((weak)); +extern void ble_log_test_flush_stale_locked_hook(void) __attribute__((weak)); +extern void ble_log_test_recycle_pre_bitmap_hook(ble_log_prph_trans_t *trans) __attribute__((weak)); +extern void ble_log_test_flush_drain_between_loads_hook(void) __attribute__((weak)); +extern void ble_log_test_acquire_after_unregister_hook(void) __attribute__((weak)); #endif BLE_LOG_IRAM_ATTR BLE_LOG_STATIC @@ -221,6 +229,31 @@ BLE_LOG_IRAM_ATTR BLE_LOG_STATIC void ble_log_pool_notify_waiter(uint8_t id) } } +/* Releases a candidate lock that was taken on a stale bitmap hint without + * changing the transport state. The transport may have become claimable + * while the lock was held (recycle does not take this lock), and a waiter + * whose scan bounced off this lock depends on this release for its wake: + * availability is re-checked only AFTER the release, inside the same + * SC-fenced window as the waiter-count read. SENDING/CLAIMED must not + * mint: unconditional notification would let a registered scanning waiter + * wake itself on a full pool. A FREE transport notifies only when its + * free-bitmap hint is already published: state alone does not make it + * claimable (a preempted recycler sits between the FREE store and the + * bitmap set), and notifying an unadvertised transport lets a waiter + * mint and consume its own wake tokens in a self-sustaining spin. */ +BLE_LOG_IRAM_ATTR BLE_LOG_STATIC void +ble_log_pool_release_candidate(ble_log_prph_trans_t *trans) +{ + BLE_LOG_CAS_RELEASE(&trans->atomic_lock); + __atomic_thread_fence(__ATOMIC_SEQ_CST); + ble_log_trans_state_t st = BLE_LOG_ATOMIC_LOAD_RELAXED(trans->state); + if (st == BLE_LOG_TRANS_STATE_OPEN || + (st == BLE_LOG_TRANS_STATE_FREE && + (BLE_LOG_ATOMIC_LOAD_ACQUIRE(g_pool.free_bitmap) & BIT(trans->id)))) { + ble_log_pool_notify_waiter(trans->id); + } +} + /* Publish OPEN and release ownership as one operation so notification can * never be moved before the lock release. */ BLE_LOG_IRAM_ATTR BLE_LOG_STATIC void @@ -286,6 +319,11 @@ BLE_LOG_IRAM_ATTR void ble_log_lbm_recycle_trans(ble_log_prph_trans_t *trans) /* Publish FREE state before advertising the pool bitmap hint. */ BLE_LOG_ATOMIC_STORE_RELEASE(trans->state, BLE_LOG_TRANS_STATE_FREE); +#if CONFIG_BLE_LOG_PRPH_TEST + if (ble_log_test_recycle_pre_bitmap_hook) { + ble_log_test_recycle_pre_bitmap_hook(trans); + } +#endif ble_log_pool_bitmap_set(&g_pool.free_bitmap, trans->id); ble_log_pool_notify_waiter(trans->id); } @@ -358,9 +396,14 @@ ble_log_prph_trans_t *ble_log_pool_try_claim_from(volatile uint32_t *bitmap, if (!BLE_LOG_CAS_ACQUIRE(&trans->atomic_lock)) { continue; } - if (BLE_LOG_ATOMIC_LOAD_RELAXED(trans->state) != expected_state) { + /* Acceptance load: pairs with the recycler's STORE_RELEASE(FREE), + * which publishes pos/pending_seal without holding this lock (the + * last lock release predates the recycle). The OPEN domain is + * already covered by that lock pairing; the acquire is required + * for FREE acceptance. */ + if (BLE_LOG_ATOMIC_LOAD_ACQUIRE(trans->state) != expected_state) { /* The bitmap is only a hint; another owner may have changed state. */ - BLE_LOG_CAS_RELEASE(&trans->atomic_lock); + ble_log_pool_release_candidate(trans); continue; } @@ -421,7 +464,7 @@ ble_log_prph_trans_t *ble_log_pool_try_claim_available(uint32_t frame_len, bool } ble_log_pool_seal_and_send(open_trans); /* releases the lock */ } else { - BLE_LOG_CAS_RELEASE(&open_trans->atomic_lock); + ble_log_pool_release_candidate(open_trans); } } } @@ -462,7 +505,10 @@ ble_log_prph_trans_t *ble_log_pool_acquire(size_t log_len, for (;;) { ble_log_prph_trans_t *trans = ble_log_pool_try_claim_available(frame_len, use_reserve); - if (trans || !wait || !BLE_LOG_ATOMIC_LOAD_ACQUIRE(lbm_enabled)) { + /* The shared ESP Timer task cannot wait for its own dispatcher. + * Check its identity only on the yieldable, would-wait path. */ + if (trans || !wait || !BLE_LOG_ATOMIC_LOAD_ACQUIRE(lbm_enabled) || + xTaskGetCurrentTaskHandle() == esp_timer_impl_get_timer_task_handle()) { return trans; } @@ -485,6 +531,11 @@ ble_log_prph_trans_t *ble_log_pool_acquire(size_t log_len, xSemaphoreTake(g_pool.sem, portMAX_DELAY); BLE_LOG_REF_COUNT_ACQUIRE_SEQ_CST(&lbm_ref_count); ble_log_pool_waiter_adjust(-1); +#if CONFIG_BLE_LOG_PRPH_TEST + if (ble_log_test_acquire_after_unregister_hook) { + ble_log_test_acquire_after_unregister_hook(); + } +#endif if (!BLE_LOG_ATOMIC_LOAD_ACQUIRE(lbm_enabled)) { return NULL; @@ -563,12 +614,9 @@ void ble_log_pool_finish_frame(ble_log_prph_trans_t *trans, uint16_t payload_len /* ---------------------------------------------- */ /* Claim with a caller-chosen wait policy. wait_for_transport=true applies - * backpressure in a yieldable context (the writer waits for a shared - * transport instead of dropping); false is a lossy fast path that returns - * NULL on a busy pool, for callers that must never block (e.g. system - * periodic output on the shared ESP timer task). Non-yieldable contexts - * (ISR, scheduler suspended) never wait whatever the policy: they keep - * their dedicated reserve and fail fast on contention. */ + * backpressure in ordinary yieldable tasks; false returns NULL on a busy + * pool. The shared ESP Timer task never waits, regardless of that policy. + * Non-yieldable contexts retain their reserve and fail-fast behavior. */ uint8_t *ble_log_claim(ble_log_src_t src_code, size_t max_len, uint32_t *handle, bool wait_for_transport) { @@ -756,10 +804,24 @@ void ble_log_lbm_begin_deinit(void) /* Wake any blocked task writers and wait until BOTH the reference count * and the waiting-task count drain to zero. Blocked writers hold no * reference while parked, so waiting on ref_count alone could free the - * pool while a woken task still touches it (use-after-free). */ + * pool while a woken task still touches it (use-after-free). The two + * counters are read separately: a writer waking between the loads + * re-acquires its reference before unregistering, so a zero waiter + * count alone is not a drain. The SC fence pairs with that acquire: + * once the waiter load has observed the unregister, the re-read below + * must observe the re-acquired reference. */ TickType_t ticks_waited = 0; - while ((BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) > 0) || - (BLE_LOG_ATOMIC_LOAD_ACQUIRE(g_pool.waiting_task_count) > 0)) { + while (true) { + bool drained = false; + if (BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) == 0) { + if (BLE_LOG_ATOMIC_LOAD_ACQUIRE(g_pool.waiting_task_count) == 0) { + __atomic_thread_fence(__ATOMIC_SEQ_CST); + drained = BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) == 0; + } + } + if (drained) { + break; + } ble_log_pool_wake_all(); vTaskDelay(1); BLE_LOG_ASSERT(ticks_waited++ < BLE_LOG_WAIT_TIMEOUT_TICKS); @@ -872,7 +934,9 @@ bool ble_log_internal_snapshot(uint16_t reason_flags, TickType_t start_tick = xTaskGetTickCount(); for (;;) { if (BLE_LOG_CAS_ACQUIRE(&internal_trans->atomic_lock)) { - if (BLE_LOG_ATOMIC_LOAD_RELAXED(internal_trans->state) == + /* Acceptance load: pairs with the dedicated transport's + * lock-free recycle publication (pos=0, STORE_RELEASE(FREE)). */ + if (BLE_LOG_ATOMIC_LOAD_ACQUIRE(internal_trans->state) == BLE_LOG_TRANS_STATE_FREE) { break; } @@ -963,6 +1027,13 @@ void ble_log_lbm_flush_open_trans(void) continue; } ble_log_prph_trans_t *trans = g_pool.trans[id]; +#if CONFIG_BLE_LOG_PRPH_TEST + /* Test point: the OPEN hint was just observed; a test can play the + * other core before this scan takes the candidate lock. */ + if (ble_log_test_flush_hint_seen_hook) { + ble_log_test_flush_hint_seen_hook((uint8_t)id); + } +#endif if (!BLE_LOG_CAS_ACQUIRE(&trans->atomic_lock)) { /* A writer holds the buffer: leave a pending-seal marker for * the next claim instead of waiting (the flusher never @@ -976,7 +1047,14 @@ void ble_log_lbm_flush_open_trans(void) trans->pos > 0) { ble_log_pool_seal_and_send(trans); } else { - BLE_LOG_CAS_RELEASE(&trans->atomic_lock); +#if CONFIG_BLE_LOG_PRPH_TEST + /* Test point: this scan holds a candidate lock taken on a + * stale hint; the transport state is not OPEN. */ + if (ble_log_test_flush_stale_locked_hook) { + ble_log_test_flush_stale_locked_hook(); + } +#endif + ble_log_pool_release_candidate(trans); } } @@ -1053,8 +1131,29 @@ void ble_log_flush(void) bool lbm_enabled_copy = BLE_LOG_ATOMIC_LOAD_ACQUIRE(lbm_enabled); ble_log_lbm_disable(); TickType_t start_tick = xTaskGetTickCount(); - while ((BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) > 1) || - (BLE_LOG_ATOMIC_LOAD_ACQUIRE(g_pool.waiting_task_count) > 0)) { + bool writers_drained = false; + while (!writers_drained) { + if (BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) <= 1) { +#if CONFIG_BLE_LOG_PRPH_TEST + if (ble_log_test_flush_drain_between_loads_hook) { + ble_log_test_flush_drain_between_loads_hook(); + } +#endif + if (BLE_LOG_ATOMIC_LOAD_ACQUIRE(g_pool.waiting_task_count) == 0) { + /* The refcount and the waiter count are two separate + * loads: a writer waking between them re-acquires its + * reference BEFORE unregistering, so a zero waiter count + * alone is not a drain. The SC fence pairs with that + * re-acquire: once this load has observed the unregister, + * the re-read must observe the re-acquired reference. */ + __atomic_thread_fence(__ATOMIC_SEQ_CST); + writers_drained = + BLE_LOG_ATOMIC_LOAD_SEQ_CST(lbm_ref_count) <= 1; + } + } + if (writers_drained) { + break; + } ble_log_pool_wake_all(); if ((xTaskGetTickCount() - start_tick) >= BLE_LOG_WAIT_TIMEOUT_TICKS) { BLE_LOG_CONSOLE("@EW: Timed out waiting for BLE Log writers\n"); diff --git a/components/bt/common/ble_log/src/ble_log_task_registry.c b/components/bt/common/ble_log/src/ble_log_task_registry.c index f9e8c92a82d..537833a594b 100644 --- a/components/bt/common/ble_log/src/ble_log_task_registry.c +++ b/components/bt/common/ble_log/src/ble_log_task_registry.c @@ -164,7 +164,9 @@ BLE_LOG_STATIC bool ble_log_task_binding_send(uint8_t cnt, uint32_t timestamp) if (!BLE_LOG_CAS_ACQUIRE(&binding_trans->atomic_lock)) { return false; } - if (BLE_LOG_ATOMIC_LOAD_RELAXED(binding_trans->state) != + /* Acceptance load: pairs with the dedicated transport's lock-free + * recycle publication (pos=0, STORE_RELEASE(FREE)). */ + if (BLE_LOG_ATOMIC_LOAD_ACQUIRE(binding_trans->state) != BLE_LOG_TRANS_STATE_FREE) { BLE_LOG_CAS_RELEASE(&binding_trans->atomic_lock); return false; diff --git a/components/bt/common/ble_log/src/internal_include/ble_log_lbm_v2.h b/components/bt/common/ble_log/src/internal_include/ble_log_lbm_v2.h index 0a057a0b93e..da8dc7047f5 100644 --- a/components/bt/common/ble_log/src/internal_include/ble_log_lbm_v2.h +++ b/components/bt/common/ble_log/src/internal_include/ble_log_lbm_v2.h @@ -246,9 +246,9 @@ bool ble_log_internal_snapshot(uint16_t reason_flags, * only BLE_LOG_SRC_ENCODE is supported; its frames are stamped with the * ENCODE source ID on the wire. wait_for_transport=false turns a busy pool * into a lossy fast path (NULL, as in a non-yieldable context) instead of - * waiting, for output that must not block its caller (e.g. the shared ESP - * timer task); true waits in yieldable contexts, never in non-yieldable - * ones. */ + * waiting. true requests backpressure in ordinary yieldable tasks; the + * shared ESP Timer task and non-yieldable contexts never wait, regardless + * of the requested policy. */ uint8_t *ble_log_claim(ble_log_src_t src_code, size_t max_len, uint32_t *handle, bool wait_for_transport); void ble_log_commit(uint32_t handle, size_t actual_len); 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 3ca6ae4ecd5..c976ad3c219 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 @@ -11,6 +11,7 @@ #include #include "esp_chip_info.h" +#include "esp_timer.h" #include "freertos/FreeRTOS.h" #include "freertos/semphr.h" #include "freertos/task.h" @@ -67,11 +68,21 @@ static volatile bool s_disable_hook_armed; static SemaphoreHandle_t s_disable_hook_entered; static SemaphoreHandle_t s_disable_hook_continue; +/* Stale-hint candidate flush (pool lost-wakeup interleaving). */ +static volatile bool s_flush_hint_hook_armed; +static volatile bool s_flush_stale_hook_armed; +static volatile bool s_flush_hint_drained; +static volatile int s_flush_hint_id = -1; +static SemaphoreHandle_t s_stale_writer_start; +static SemaphoreHandle_t s_stale_writer_started; + void ble_log_test_claim_pre_publish_hook(void); void ble_log_test_claim_locked_hook(void); void ble_log_test_init_snapshot_before_acquire_hook(void); void ble_log_test_enable_before_lifecycle_lock_hook(void); void ble_log_test_disable_before_wake_hook(void); +void ble_log_test_flush_hint_seen_hook(uint8_t id); +void ble_log_test_flush_stale_locked_hook(void); #if CONFIG_BLE_HOST_COMPRESSED_LOG_ENABLE extern int ble_log_compressed_hex_print(uint8_t source, uint32_t log_index, size_t args_cnt, ...); @@ -121,6 +132,42 @@ void ble_log_test_disable_before_wake_hook(void) } } +void ble_log_test_flush_hint_seen_hook(uint8_t id) +{ + if (!s_flush_hint_hook_armed) { + return; + } + s_flush_hint_hook_armed = false; + s_flush_hint_id = id; + /* Play the other core between the flusher's hint read and its candidate + * lock: seal, send and recycle the one OPEN transport. The recycle + * notification runs while no writer is registered yet, so only the + * flusher's own candidate release can re-advertise it. drain dispatches + * to the peripheral queue; the peripheral releases on read, and that + * queue holds only the hinted transport (the other shared ones are + * pinned by open claims), so one read recycles exactly it. */ + ble_log_lbm_flush_open_trans(); + s_flush_hint_drained = ble_log_rt_drain(); + if (s_flush_hint_drained) { + s_flush_hint_drained = + ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), 0, 0, NULL) > 0; + } +} + +void ble_log_test_flush_stale_locked_hook(void) +{ + if (!s_flush_stale_hook_armed) { + return; + } + s_flush_stale_hook_armed = false; + /* The flusher holds the recycled transport's candidate lock. Start the + * writer: with every other shared transport pinned by a claim, its two + * scans bounce off this lock and it parks before this hook returns. */ + xSemaphoreGive(s_stale_writer_start); + (void)xSemaphoreTake(s_stale_writer_started, portMAX_DELAY); + vTaskDelay(pdMS_TO_TICKS(20)); +} + /* A commit field is hex characters, zero-padded after a shorter value; * anything else (garbage, non-hex, zeros after data) is invalid. */ static bool commit_is_valid(const uint8_t *commit, size_t len) @@ -1501,6 +1548,291 @@ TEST_CASE("BLE Log task writers wait for a shared transport", "[ble_log][lbm]") } } +typedef struct { + SemaphoreHandle_t done; + bool write_result; +} stale_writer_ctx_t; + +static void stale_candidate_writer_task(void *arg) +{ + stale_writer_ctx_t *ctx = arg; + (void)xSemaphoreTake(s_stale_writer_start, portMAX_DELAY); + xSemaphoreGive(s_stale_writer_started); + const uint8_t marker = 0x61; + ctx->write_result = ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1); + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +TEST_CASE("BLE Log flush re-advertises a stale candidate for parked writers", + "[ble_log][lbm]") +{ + /* Deterministic lost-wakeup interleaving: the periodic flusher reads an + * OPEN hint, the hinted transport is sealed, sent and recycled before + * the flusher takes its candidate lock, and a task writer registers + * and parks while that lock is held. Only the flusher's release can + * then re-advertise the FREE transport; a bare unlock loses the + * writer. */ + TEST_ASSERT_TRUE(ble_log_enable(true)); + 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) { + } + + /* Pin every shared transport but one with an open claim: claimed + * transports are in no bitmap and never enter the runtime queue. */ + uint32_t pinned[BLE_LOG_POOL_TRANS_CNT] = {0}; + int pinned_cnt = 0; + for (; pinned_cnt < BLE_LOG_POOL_SHARED_CNT - 1; pinned_cnt++) { + uint8_t *payload = ble_log_claim(BLE_LOG_SRC_ENCODE, 1, + &pinned[pinned_cnt], false); + TEST_ASSERT_NOT_NULL(payload); + } + + /* The last shared transport stays OPEN, so the armed flush targets it. */ + const uint8_t marker = 0x61; + TEST_ASSERT_TRUE(ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1)); + + stale_writer_ctx_t ctx = {.done = xSemaphoreCreateBinary()}; + s_stale_writer_start = xSemaphoreCreateBinary(); + s_stale_writer_started = xSemaphoreCreateBinary(); + TEST_ASSERT_NOT_NULL(ctx.done); + TEST_ASSERT_NOT_NULL(s_stale_writer_start); + TEST_ASSERT_NOT_NULL(s_stale_writer_started); + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreate(stale_candidate_writer_task, "ble_log_stale", + TEST_LIFECYCLE_STACK_SIZE, &ctx, + TEST_LIFECYCLE_PRIO, NULL)); + + s_flush_hint_id = -1; + s_flush_hint_drained = false; + s_flush_hint_hook_armed = true; + s_flush_stale_hook_armed = true; + ble_log_lbm_flush_open_trans(); + TEST_ASSERT_TRUE(s_flush_hint_drained); + TEST_ASSERT_TRUE(s_flush_hint_id >= 0); + + bool completed = xSemaphoreTake(ctx.done, pdMS_TO_TICKS(1000)) == pdTRUE; + if (!completed) { + /* Lost wake: un-park the writer, release the pinned claims and + * restore the gate before failing, so later cases stay clean. */ + TEST_ASSERT_TRUE(ble_log_enable(false)); + TEST_ASSERT_EQUAL(pdTRUE, xSemaphoreTake(ctx.done, pdMS_TO_TICKS(1000))); + TEST_ASSERT_FALSE(ctx.write_result); + for (int i = 0; i < pinned_cnt; i++) { + ble_log_commit(pinned[i], 0); + } + TEST_ASSERT_TRUE(ble_log_enable(true)); + TEST_FAIL_MESSAGE("stale-candidate release lost the parked writer's wake"); + } + TEST_ASSERT_TRUE(ctx.write_result); + + for (int i = 0; i < pinned_cnt; i++) { + ble_log_commit(pinned[i], 0); + } + + /* The recovered writer's frame reaches the peripheral. */ + ble_log_lbm_flush_open_trans(); + TEST_ASSERT_TRUE(ble_log_rt_drain()); + TEST_ASSERT_GREATER_THAN_size_t( + 0, ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(TEST_READ_TIMEOUT_MS), 0, NULL)); + + vSemaphoreDelete(ctx.done); + vSemaphoreDelete(s_stale_writer_start); + vSemaphoreDelete(s_stale_writer_started); +} + +typedef struct { + SemaphoreHandle_t done; + uint8_t marker; + bool result; + bool finished; +} conc_writer_ctx_t; + +#define CONC_WRITER_CNT (4) +#define CONC_FRAMES_EACH (96) + +static void conc_writer_task(void *arg) +{ + conc_writer_ctx_t *ctx = arg; + for (uint16_t i = 0; i < CONC_FRAMES_EACH; i++) { + uint8_t payload[3] = {ctx->marker, (uint8_t)(i >> 8), (uint8_t)i}; + if (!ble_log_write_hex(BLE_LOG_SRC_CUSTOM, payload, sizeof(payload))) { + ctx->result = false; + break; + } + } + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +typedef struct { + size_t frames; + size_t per_marker[CONC_WRITER_CNT]; + uint32_t sn[CONC_WRITER_CNT * CONC_FRAMES_EACH]; + size_t sn_count; +} conc_capture_t; + +static void capture_conc_frame(const test_ble_log_frame_t *frame, void *ctx_) +{ + conc_capture_t *cap = ctx_; + if (frame->src != BLE_LOG_SRC_CUSTOM || frame->payload_len != 4 + 3) { + return; + } + cap->frames++; + uint8_t marker = frame->payload[4]; + for (int w = 0; w < CONC_WRITER_CNT; w++) { + if (marker == 0xa0 + w) { + cap->per_marker[w]++; + } + } + if (cap->sn_count < CONC_WRITER_CNT * CONC_FRAMES_EACH) { + cap->sn[cap->sn_count++] = frame->sn; + } +} + +typedef struct { + SemaphoreHandle_t done; + uint32_t hammer_cnt; + bool taken; +} hammer_ctx_t; + +static void hammer_claim_task(void *arg) +{ + hammer_ctx_t *ctx = arg; + uint32_t committed = 0; + const uint8_t marker = 0xb0; + for (int i = 0; i < 4000; i++) { + uint32_t handle; + uint8_t *payload = ble_log_claim(BLE_LOG_SRC_ENCODE, 1, + &handle, false); + if (payload) { + payload[0] = marker; + ble_log_commit(handle, 1); + committed++; + } + } + ctx->hammer_cnt = committed; + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +typedef struct { + uint32_t sn[CONC_WRITER_CNT * CONC_FRAMES_EACH + 16 * 1024]; + size_t sn_count; +} hammer_capture_t; + +static void capture_hammer_frame(const test_ble_log_frame_t *frame, void *ctx_) +{ + hammer_capture_t *cap = ctx_; + if (frame->src == BLE_LOG_SRC_ENCODE && + cap->sn_count < CONC_WRITER_CNT * CONC_FRAMES_EACH + 16 * 1024) { + cap->sn[cap->sn_count++] = frame->sn; + } +} + +static size_t count_unique_sn(const uint32_t *sn, size_t n) +{ + size_t unique = 0; + for (size_t i = 0; i < n; i++) { + bool first = true; + for (size_t j = 0; j < i; j++) { + if (sn[j] == sn[i]) { + first = false; + break; + } + } + unique += first; + } + return unique; +} + +typedef struct { + SemaphoreHandle_t done; + bool write_result; + bool claim_result; +} timer_writer_ctx_t; + +static void timer_writer_cb(void *arg) +{ + timer_writer_ctx_t *ctx = arg; + const uint8_t marker = 0x59; + ctx->write_result = ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1); + uint32_t handle; + uint8_t *payload = ble_log_claim(BLE_LOG_SRC_ENCODE, 1, &handle, true); + ctx->claim_result = payload != NULL; + if (payload) { + payload[0] = marker; + ble_log_commit(handle, 1); + } +#if CONFIG_BLE_LOG_LL_ENABLED + ble_log_write_hex_ll(1, &marker, 0, NULL, BIT(BLE_LOG_LL_FLAG_TASK)); +#endif +#if CONFIG_BLE_HOST_COMPRESSED_LOG_ENABLE + ble_log_compressed_hex_print(BLE_COMPRESSED_LOG_OUT_SOURCE_HOST, 0x730, 0); +#endif + xSemaphoreGive(ctx->done); +} + +TEST_CASE("BLE Log ESP Timer writers never wait for shared transports", + "[ble_log][lbm]") +{ + static const uint8_t full_payload[ + BLE_LOG_MAX_PAYLOAD_LEN - sizeof(uint32_t)] = {0}; + TEST_ASSERT_TRUE(ble_log_enable(true)); + 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) { + } + /* The test peripheral holds these transports until explicitly read. + * Reserve capacity remains available but must not be used by this callback. */ + for (int i = 0; i < BLE_LOG_POOL_SHARED_CNT; i++) { + TEST_ASSERT_TRUE(ble_log_write_hex(BLE_LOG_SRC_CUSTOM, full_payload, + sizeof(full_payload))); + } + + timer_writer_ctx_t ctx = {.done = xSemaphoreCreateBinary()}; + TEST_ASSERT_NOT_NULL(ctx.done); + esp_timer_handle_t timer; + const esp_timer_create_args_t args = { + .callback = timer_writer_cb, + .arg = &ctx, + .dispatch_method = ESP_TIMER_TASK, + .name = "ble_log_writer", + }; + TEST_ESP_OK(esp_timer_create(&args, &timer)); + + bool returned = true; + bool results_match = true; + /* First fire with shared capacity exhausted, then again after recycle: + * both callbacks must return, but only the second may write/claim. */ + for (int attempt = 0; attempt < 2; attempt++) { + TEST_ESP_OK(esp_timer_start_once(timer, 1)); + bool completed = xSemaphoreTake(ctx.done, pdMS_TO_TICKS(1000)) == pdTRUE; + returned &= completed; + if (!completed) { + /* A regressed writer may be parked. Wake it before waiting for + * callback completion; never delete a task inside a pool API. */ + TEST_ASSERT_TRUE(ble_log_enable(false)); + } + TEST_ESP_OK(esp_timer_stop_blocking(timer, portMAX_DELAY)); + (void)xSemaphoreTake(ctx.done, 0); + results_match &= ctx.write_result == (attempt != 0) && + ctx.claim_result == (attempt != 0); + + 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) { + } + TEST_ASSERT_TRUE(ble_log_enable(true)); + } + TEST_ESP_OK(esp_timer_delete(timer)); + vSemaphoreDelete(ctx.done); + TEST_ASSERT_TRUE_MESSAGE(returned, "ESP Timer writer waited for a shared transport"); + TEST_ASSERT_TRUE(results_match); +} + #define SNAPSHOT_CAPTURE_MAX 8 typedef struct { @@ -1777,3 +2109,411 @@ TEST_CASE("BLE Log disable keeps waiter semaphore alive during deinit", TEST_ASSERT_TRUE(ble_log_init()); } #endif /* CONFIG_BLE_LOG_LL_ENABLED */ + +/* ------------------------------------------------------------------ */ +/* Review blocker repro: pool waiter self-wake spin. */ +/* A cancelled claim publishes FREE state before the free-bitmap hint; */ +/* a waiter scanning that transport from the stale open-cursor hint */ +/* must PARK. Unconditional notification on the unlocked FREE state */ +/* lets the waiter mint and consume its own wake tokens in a tight */ +/* loop: the writer never parks, and the publisher suspended inside */ +/* the recycle window can never run to publish the bitmap. */ +/* ------------------------------------------------------------------ */ +static volatile bool s_spin_hook_armed; +static SemaphoreHandle_t s_spin_writer_start; +static SemaphoreHandle_t s_spin_writer_entered; +static SemaphoreHandle_t s_spin_canceller_hold; + +typedef struct { + SemaphoreHandle_t done; + volatile bool result; +} spin_writer_ctx_t; + +static void spin_writer_task(void *arg) +{ + spin_writer_ctx_t *ctx = arg; + (void)xSemaphoreTake(s_spin_writer_start, portMAX_DELAY); + xSemaphoreGive(s_spin_writer_entered); + const uint8_t marker = 0xC1; + ctx->result = ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1); + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +typedef struct { + SemaphoreHandle_t done; + uint32_t handle; +} spin_cancel_ctx_t; + +static void spin_canceller_task(void *arg) +{ + spin_cancel_ctx_t *ctx = arg; + s_spin_hook_armed = true; + ble_log_commit(ctx->handle, 0); + xSemaphoreGive(ctx->done); + vTaskDelete(NULL); +} + +void ble_log_test_recycle_pre_bitmap_hook(ble_log_prph_trans_t *trans); + +void ble_log_test_recycle_pre_bitmap_hook(ble_log_prph_trans_t *trans) +{ + (void)trans; + if (!s_spin_hook_armed) { + return; + } + s_spin_hook_armed = false; + /* Deterministic preemption between the FREE publication and the + * free-bitmap hint: release the waiting writer, then stay suspended + * exactly like a publisher preempted in this window. */ + xSemaphoreGive(s_spin_writer_start); + (void)xSemaphoreTake(s_spin_canceller_hold, portMAX_DELAY); +} + +typedef struct { + volatile bool stop; + volatile uint32_t loops; +} canary_ctx_t; + +static void spin_canary_task(void *arg) +{ + canary_ctx_t *ctx = arg; + while (!ctx->stop) { + ctx->loops++; + vTaskDelay(1); + } + vTaskDelete(NULL); +} + +TEST_CASE("BLE Log pool waiter parks instead of self-waking on an unpublished FREE", + "[ble_log][lbm][repro]") +{ + canary_ctx_t canary = {0}; + spin_writer_ctx_t wctx = {.done = xSemaphoreCreateBinary()}; + spin_cancel_ctx_t cctx = {.done = xSemaphoreCreateBinary()}; + s_spin_writer_start = xSemaphoreCreateBinary(); + s_spin_writer_entered = xSemaphoreCreateBinary(); + s_spin_canceller_hold = xSemaphoreCreateBinary(); + TEST_ASSERT_NOT_NULL(wctx.done); + TEST_ASSERT_NOT_NULL(cctx.done); + TEST_ASSERT_NOT_NULL(s_spin_writer_start); + TEST_ASSERT_NOT_NULL(s_spin_writer_entered); + TEST_ASSERT_NOT_NULL(s_spin_canceller_hold); + + /* Fresh pool: the first shared transport must land on id 0 so the + * stale open-cursor hint (still 0 from init) points at it. */ + ble_log_deinit(); + TEST_ASSERT_TRUE(ble_log_init()); + 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) { + } + } + + TEST_ASSERT_NOT_NULL(ble_log_claim(BLE_LOG_SRC_ENCODE, 1, + &cctx.handle, false)); + TEST_ASSERT_EQUAL(0, cctx.handle & 0xff); + uint32_t pinned[BLE_LOG_POOL_TRANS_CNT] = {0}; + int pinned_cnt = 0; + for (; pinned_cnt < BLE_LOG_POOL_SHARED_CNT - 1; pinned_cnt++) { + TEST_ASSERT_NOT_NULL(ble_log_claim(BLE_LOG_SRC_ENCODE, 1, + &pinned[pinned_cnt], false)); + } + + /* The canary shares the writer's core: it only progresses while the + * writer is parked; a spinning writer monopolizes the core. */ + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreatePinnedToCore(spin_canary_task, "ble_log_canary", + TEST_LIFECYCLE_STACK_SIZE, &canary, + 1, NULL, 1)); + vTaskDelay(pdMS_TO_TICKS(20)); + uint32_t canary_base = canary.loops; + TEST_ASSERT_GREATER_THAN(0, canary_base); + + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreatePinnedToCore(spin_writer_task, "ble_log_spin", + TEST_LIFECYCLE_STACK_SIZE, &wctx, + TEST_LIFECYCLE_PRIO + 1, NULL, 1)); + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreatePinnedToCore(spin_canceller_task, "ble_log_cancel", + TEST_LIFECYCLE_STACK_SIZE, &cctx, + TEST_LIFECYCLE_PRIO - 1, NULL, 0)); + + TEST_ASSERT_TRUE(xSemaphoreTake(s_spin_writer_entered, pdMS_TO_TICKS(1000))); + /* Let the writer reach its steady state (parked or spinning), then + * sample the canary across a second window: a parked writer leaves + * the core to the canary; a spinning writer monopolizes it. */ + vTaskDelay(pdMS_TO_TICKS(30)); + uint32_t canary_settled = canary.loops; + vTaskDelay(pdMS_TO_TICKS(100)); + uint32_t canary_now = canary.loops; + bool completed_early = xSemaphoreTake(wctx.done, 0) == pdTRUE; + printf("B1 sample: writer completed_early=%d canary %u -> %u during hold (core1: %s)\n", + completed_early, (unsigned)canary_settled, (unsigned)canary_now, + canary_now == canary_settled ? "MONOPOLIZED by spinning writer" : "progressing, writer parked"); + + /* Let the suspended publisher finish the recycle and observe the + * recovery: the writer must complete once the bitmap is published. */ + xSemaphoreGive(s_spin_canceller_hold); + TEST_ASSERT_TRUE(xSemaphoreTake(cctx.done, pdMS_TO_TICKS(1000))); + TEST_ASSERT_TRUE(xSemaphoreTake(wctx.done, pdMS_TO_TICKS(1000))); + TEST_ASSERT_TRUE(wctx.result); + + for (int i = 0; i < pinned_cnt; i++) { + ble_log_commit(pinned[i], 0); + } + 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) { + } + + canary.stop = true; + vTaskDelay(pdMS_TO_TICKS(5)); + vSemaphoreDelete(wctx.done); + vSemaphoreDelete(cctx.done); + vSemaphoreDelete(s_spin_writer_start); + vSemaphoreDelete(s_spin_writer_entered); + vSemaphoreDelete(s_spin_canceller_hold); + s_spin_writer_start = NULL; + s_spin_writer_entered = NULL; + s_spin_canceller_hold = NULL; + + /* Contract: while its only exit is suspended, the waiter must be + * PARKED (canary progresses across the hold window), never spinning + * on self-minted tokens (canary frozen = core monopolized). */ + TEST_ASSERT_GREATER_THAN(canary_settled, canary_now); +} + +/* ------------------------------------------------------------------ */ +/* Review blocker repro: flush drain straddle. */ +/* The drain loop reads the writer refcount and the waiter count as */ +/* two separate loads. A writer waking between them is in-flight (ref */ +/* re-acquired, waiter unregistered, gate not yet checked) when the */ +/* drain concludes; flush then snapshots and re-enables producers */ +/* while the pre-existing writer still owns its reference. */ +/* ------------------------------------------------------------------ */ +static volatile bool s_straddle_acquire_armed; +static volatile bool s_straddle_flush_armed; +static volatile int s_straddle_flush_evals; +static volatile bool s_straddle_hog_stop; +static volatile bool s_straddle_writer_frozen_at_return; +static volatile bool s_straddle_writer_result; +static volatile int64_t s_straddle_flush_us; +static SemaphoreHandle_t s_straddle_freeze; +static SemaphoreHandle_t s_straddle_inflight; +static SemaphoreHandle_t s_straddle_flush_start; +static SemaphoreHandle_t s_straddle_flush_done; +static SemaphoreHandle_t s_straddle_writer_started; +static SemaphoreHandle_t s_straddle_writer_done; +static SemaphoreHandle_t s_straddle_reader_start; + +void ble_log_test_acquire_after_unregister_hook(void); +void ble_log_test_flush_drain_between_loads_hook(void); +static void straddle_hog_task(void *arg); + +void ble_log_test_acquire_after_unregister_hook(void) +{ + if (!s_straddle_acquire_armed) { + return; + } + s_straddle_acquire_armed = false; + /* The writer is in-flight: reference re-acquired, waiter + * unregistered, producer gate not yet observed. Any preemption of + * this task leaves exactly this state behind. */ + xSemaphoreGive(s_straddle_inflight); + (void)xSemaphoreTake(s_straddle_freeze, portMAX_DELAY); +} + +void ble_log_test_flush_drain_between_loads_hook(void) +{ + if (!s_straddle_flush_armed) { + return; + } + switch (++s_straddle_flush_evals) { + case 1: + /* Evaluation #1, before this iteration's wake_all: start a hog on + * the writer's core (idle: the writer is parked) so the wake token + * minted moments later stays unconsumed. */ + (void)xTaskCreatePinnedToCore(straddle_hog_task, "ble_log_strh", + TEST_LIFECYCLE_STACK_SIZE, NULL, + TEST_LIFECYCLE_PRIO + 2, NULL, 0); + break; + case 2: + /* Evaluation #2, between the two condition loads: stop the hog so + * the parked writer consumes the wake and runs its real post-take + * handoff steps (re-acquire + unregister) inside this window. + * Only now may the link recycle transports: the writer is + * in-flight and unregistered, so recycles mint no wake tokens. */ + s_straddle_flush_armed = false; + s_straddle_hog_stop = true; + if (xSemaphoreTake(s_straddle_inflight, pdMS_TO_TICKS(2000)) == pdTRUE) { + xSemaphoreGive(s_straddle_reader_start); + } + break; + default: + break; + } +} + +static void straddle_writer_task(void *arg) +{ + (void)arg; + xSemaphoreGive(s_straddle_writer_started); + const uint8_t marker = 0xD2; + s_straddle_writer_result = ble_log_write_hex(BLE_LOG_SRC_CUSTOM, &marker, 1); + xSemaphoreGive(s_straddle_writer_done); + vTaskDelete(NULL); +} + +static void straddle_flusher_task(void *arg) +{ + (void)arg; + (void)xSemaphoreTake(s_straddle_flush_start, portMAX_DELAY); + int64_t t0 = esp_timer_get_time(); + ble_log_flush(); + s_straddle_flush_us = esp_timer_get_time() - t0; + s_straddle_writer_frozen_at_return = !s_straddle_acquire_armed; + xSemaphoreGive(s_straddle_flush_done); + vTaskDelete(NULL); +} + +static void straddle_hog_task(void *arg) +{ + (void)arg; + while (!s_straddle_hog_stop) { + } + vTaskDelete(NULL); +} + +static volatile bool s_straddle_reader_stop; + +static void straddle_reader_task(void *arg) +{ + (void)arg; + /* Recycle transports the way a real link does, so flush_all_trans + * completes and the measured flush duration reflects the writer + * drain rather than the transport wait. Blocked until the drain has + * handed the writer off: recycling earlier would wake the parked + * writer before the straddle point. */ + (void)xSemaphoreTake(s_straddle_reader_start, portMAX_DELAY); + while (!s_straddle_reader_stop) { + (void)ble_log_prph_test_read(s_read_buf, sizeof(s_read_buf), + pdMS_TO_TICKS(10), 0, NULL); + } + vTaskDelete(NULL); +} + +TEST_CASE("BLE Log flush drain does not straddle an in-flight parked writer", + "[ble_log][lbm][repro]") +{ + static const uint8_t full_payload[ + BLE_LOG_MAX_PAYLOAD_LEN - sizeof(uint32_t)] = {0}; + + s_straddle_freeze = xSemaphoreCreateBinary(); + s_straddle_inflight = xSemaphoreCreateBinary(); + s_straddle_flush_start = xSemaphoreCreateBinary(); + s_straddle_flush_done = xSemaphoreCreateBinary(); + s_straddle_writer_started = xSemaphoreCreateBinary(); + s_straddle_writer_done = xSemaphoreCreateBinary(); + s_straddle_reader_start = xSemaphoreCreateBinary(); + TEST_ASSERT_NOT_NULL(s_straddle_freeze); + TEST_ASSERT_NOT_NULL(s_straddle_inflight); + TEST_ASSERT_NOT_NULL(s_straddle_flush_start); + TEST_ASSERT_NOT_NULL(s_straddle_flush_done); + TEST_ASSERT_NOT_NULL(s_straddle_writer_started); + TEST_ASSERT_NOT_NULL(s_straddle_writer_done); + TEST_ASSERT_NOT_NULL(s_straddle_reader_start); + + 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) { + } + } + + /* Keep every task-usable transport SENDING in the test peripheral. */ + for (int i = 0; i < BLE_LOG_POOL_SHARED_CNT; i++) { + TEST_ASSERT_TRUE(ble_log_write_hex(BLE_LOG_SRC_CUSTOM, full_payload, + sizeof(full_payload))); + } + TEST_ASSERT_TRUE(ble_log_rt_drain()); + + /* The writer parks on the full shared pool. */ + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreatePinnedToCore(straddle_writer_task, "ble_log_strw", + TEST_LIFECYCLE_STACK_SIZE, NULL, + TEST_LIFECYCLE_PRIO, NULL, 0)); + TEST_ASSERT_TRUE(xSemaphoreTake(s_straddle_writer_started, pdMS_TO_TICKS(1000))); + vTaskDelay(pdMS_TO_TICKS(10)); + + /* The eval#1 hook starts the hog on the writer's core so the first + * wake stays unconsumed; the eval#2 hook stops it exactly between + * the two condition loads. */ + s_straddle_flush_evals = 0; + s_straddle_flush_armed = true; + s_straddle_acquire_armed = true; + s_straddle_hog_stop = false; + s_straddle_reader_stop = false; + s_straddle_writer_frozen_at_return = false; + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreate(straddle_reader_task, "ble_log_strr", + TEST_READER_STACK_SIZE, NULL, + TEST_READER_PRIO, NULL)); + TEST_ASSERT_EQUAL(pdTRUE, + xTaskCreatePinnedToCore(straddle_flusher_task, "ble_log_strf", + TEST_READER_STACK_SIZE, NULL, + TEST_READER_PRIO + 2, NULL, 1)); + xSemaphoreGive(s_straddle_flush_start); + + TEST_ASSERT_TRUE(xSemaphoreTake(s_straddle_flush_done, pdMS_TO_TICKS(10000))); + int64_t flush_ms = s_straddle_flush_us / 1000; + bool writer_done_early = xSemaphoreTake(s_straddle_writer_done, 0) == pdTRUE; + printf("B2 sample: flush returned in %lld ms, writer in-flight at return=%d, writer already done=%d\n", + (long long)flush_ms, s_straddle_writer_frozen_at_return, + writer_done_early); + + /* Transports are already claimable (the reader kept the link busy); + * release the frozen writer: on the straddle it observes the + * re-enabled gate and completes its pre-flush write after flush + * already returned. */ + xSemaphoreGive(s_straddle_freeze); + TEST_ASSERT_TRUE(xSemaphoreTake(s_straddle_writer_done, pdMS_TO_TICKS(1000))); + printf("B2 outcome: frozen writer passed the re-enabled gate and wrote=%d\n", + s_straddle_writer_result); + + s_straddle_reader_stop = true; + vTaskDelay(pdMS_TO_TICKS(30)); + 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) { + } + + vSemaphoreDelete(s_straddle_freeze); + vSemaphoreDelete(s_straddle_inflight); + vSemaphoreDelete(s_straddle_flush_start); + vSemaphoreDelete(s_straddle_flush_done); + vSemaphoreDelete(s_straddle_writer_started); + vSemaphoreDelete(s_straddle_writer_done); + vSemaphoreDelete(s_straddle_reader_start); + s_straddle_freeze = NULL; + s_straddle_inflight = NULL; + s_straddle_flush_start = NULL; + s_straddle_flush_done = NULL; + s_straddle_writer_started = NULL; + s_straddle_writer_done = NULL; + s_straddle_reader_start = NULL; + + /* Contract: when a writer was in-flight at flush return, flush must + * have held the drain for the full documented timeout (1s) instead + * of concluding between the two counter loads. */ + if (s_straddle_writer_frozen_at_return) { + TEST_ASSERT_GREATER_OR_EQUAL(950, (int)flush_ms); + } +}