fix(ble_log): harden the pool waiter handshake and claim ordering

Consolidates the concurrency review fixes for the unified pool:

- the shared ESP Timer task never waits for a transport: its identity
  is checked on the would-wait path only, so its callbacks always
  return and dispatch can progress (this also bounds claim()'s
  backpressure);
- transports accepted from the FREE bitmap are ACQUIRE-loaded: the
  recycler publishes pos/pending_seal with a STORE_RELEASE(FREE)
  without holding the candidate lock, and the acceptance load pairs
  with that publication;
- a waiter whose scan bounced off a candidate lock is re-advertised:
  the releasing side re-checks availability after the lock release,
  inside the same seq_cst window as the waiter-count read, so a
  transport that became claimable while locked cannot strand its wake;
- waiters never self-wake on an unpublished FREE transport:
  notification for a FREE transport fires only when its free-bitmap
  hint is already published, so a registered waiter cannot mint and
  consume its own wake tokens in a self-sustaining spin;
- flush and deinit drains re-check writer references after observing a
  zero waiter count: a writer waking between the two loads re-acquires
  its reference before unregistering, and the seq_cst-fenced re-read
  must observe it before the drain concludes.
This commit is contained in:
Zhou Xiao
2026-09-10 15:39:56 +08:00
committed by guozifan
parent 2cf0638686
commit 22812d893a
6 changed files with 878 additions and 33 deletions
+11 -7
View File
@@ -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
@@ -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);
+116 -17
View File
@@ -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");
@@ -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;
@@ -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);
@@ -11,6 +11,7 @@
#include <string.h>
#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);
}
}