feat(driver_twai): add prioritized TX queue support

This commit is contained in:
wanckl
2026-06-23 20:18:04 +08:00
parent 671606b9b5
commit de63f3978b
13 changed files with 743 additions and 38 deletions
+1 -1
View File
@@ -9,7 +9,7 @@ set(public_include "include")
set(priv_req esp_driver_gpio esp_pm esp_timer)
if(CONFIG_SOC_TWAI_SUPPORTED)
list(APPEND srcs "esp_twai_onchip.c")
list(APPEND srcs "esp_twai_onchip.c" "twai_frame_queue.c")
endif()
idf_component_register(
+14 -34
View File
@@ -7,9 +7,10 @@
#include "esp_timer.h"
#include "esp_twai.h"
#include "esp_twai_onchip.h"
#include "twai_private.h"
#include "esp_private/twai_interface.h"
#include "esp_private/twai_utils.h"
#include "twai_private.h"
#include "esp_private/twai_frame_queue.h"
#include "hal/twai_periph.h"
#include "hal/twai_hal.h"
#if SOC_HAS(TWAI_FD)
@@ -46,7 +47,7 @@ typedef struct {
twai_hal_context_t *hal;
intr_handle_t intr_hdl;
intr_handle_t timer_intr_hdl;
QueueHandle_t tx_mount_queue;
twai_frame_queue_t tx_queue;
EventGroupHandle_t event_group;
twai_clock_source_t curr_clk_src;
uint32_t src_freq_hz;
@@ -183,7 +184,7 @@ static uint8_t _node_start_tx_batch_from_isr(twai_onchip_ctx_t *node, BaseType_t
// we should fill slot with same order to avoid data inverse
while (tx_idx < node->tx_slot_num) { // try best to fill all tx slots
const twai_frame_t *frame = NULL;
if (xQueueReceiveFromISR(node->tx_mount_queue, &frame, yield_required) != pdTRUE) {
if (twai_frame_queue_pop_from_isr(node->tx_queue, &frame, (bool *)yield_required) != ESP_OK) {
break;
}
node->p_curr_tx[tx_idx] = frame;
@@ -344,9 +345,7 @@ static void _node_destroy(twai_onchip_ctx_t *twai_ctx)
if (twai_ctx->timer_intr_hdl) {
esp_intr_free(twai_ctx->timer_intr_hdl);
}
if (twai_ctx->tx_mount_queue) {
vQueueDeleteWithCaps(twai_ctx->tx_mount_queue);
}
twai_frame_queue_del(twai_ctx->tx_queue);
if (twai_ctx->event_group) {
vEventGroupDeleteWithCaps(twai_ctx->event_group);
}
@@ -582,7 +581,7 @@ static esp_err_t _node_get_status(twai_node_handle_t node, twai_node_status_t *s
status_ret->state = atomic_load(&twai_ctx->state);
status_ret->tx_error_count = twai_hal_get_tec(twai_ctx->hal);
status_ret->rx_error_count = twai_hal_get_rec(twai_ctx->hal);
status_ret->tx_queue_remaining = uxQueueSpacesAvailable(twai_ctx->tx_mount_queue);
status_ret->tx_queue_remaining = twai_frame_queue_get_free_space(twai_ctx->tx_queue);
}
if (record_ret) {
*record_ret = twai_ctx->history;
@@ -606,7 +605,6 @@ static esp_err_t _node_queue_tx(twai_node_handle_t node, const twai_frame_t *fra
ESP_RETURN_ON_FALSE_ISR((!frame->header.brs) || (twai_ctx->valid_fd_timing), ESP_ERR_INVALID_ARG, TAG, "brs can't be used without config data_timing");
ESP_RETURN_ON_FALSE_ISR(!twai_ctx->hal->enable_listen_only, ESP_ERR_NOT_SUPPORTED, TAG, "node is config as listen only");
ESP_RETURN_ON_FALSE_ISR(atomic_load(&twai_ctx->state) != TWAI_ERROR_BUS_OFF, ESP_ERR_INVALID_STATE, TAG, "node is bus off");
TickType_t ticks_to_wait = (timeout == -1) ? portMAX_DELAY : pdMS_TO_TICKS(timeout);
xEventGroupClearBits(twai_ctx->event_group, TWAI_IDLE_EVENT_BIT); //going to send, clear the idle event
bool false_var = false;
@@ -615,38 +613,18 @@ static esp_err_t _node_queue_tx(twai_node_handle_t node, const twai_frame_t *fra
_node_start_trans(twai_ctx, 0);
} else {
// Hardware busy, need to queue the frame
BaseType_t is_isr_context = xPortInIsrContext();
BaseType_t yield_required = pdFALSE;
if (is_isr_context) {
// In ISR context - use ISR-safe queue operations
ESP_RETURN_ON_FALSE_ISR(xQueueSendFromISR(twai_ctx->tx_mount_queue, &frame, &yield_required), ESP_ERR_TIMEOUT, TAG, "tx queue full");
} else {
// In task context - use normal queue operations
ESP_RETURN_ON_FALSE(xQueueSend(twai_ctx->tx_mount_queue, &frame, ticks_to_wait), ESP_ERR_TIMEOUT, TAG, "tx queue full");
}
ESP_RETURN_ON_ERROR_ISR(twai_frame_queue_push_safe(twai_ctx->tx_queue, frame, frame->tx_queue_priority, timeout), TAG, "tx queue full");
// Second chance check for hardware availability
false_var = false;
if (atomic_compare_exchange_strong(&twai_ctx->hw_busy, &false_var, true)) {
BaseType_t dequeue_result;
if (is_isr_context) {
dequeue_result = xQueueReceiveFromISR(twai_ctx->tx_mount_queue, &twai_ctx->p_curr_tx[0], &yield_required);
} else {
dequeue_result = xQueueReceive(twai_ctx->tx_mount_queue, &twai_ctx->p_curr_tx[0], 0);
}
if (dequeue_result == pdTRUE) {
if (twai_frame_queue_pop_safe(twai_ctx->tx_queue, &twai_ctx->p_curr_tx[0]) == ESP_OK) {
_node_start_trans(twai_ctx, 0);
} else {
// any reason here means frame already taken and maybe finished by fast hardware, so back `hw_busy` to false
atomic_store(&twai_ctx->hw_busy, false);
}
}
// Handle ISR yield if required
if (is_isr_context && yield_required) {
portYIELD_FROM_ISR();
}
}
return ESP_OK;
}
@@ -657,9 +635,9 @@ static esp_err_t _node_wait_tx_all_done(twai_node_handle_t node, int timeout)
TickType_t ticks_to_wait = (timeout == -1) ? portMAX_DELAY : pdMS_TO_TICKS(timeout);
ESP_RETURN_ON_FALSE(atomic_load(&twai_ctx->state) != TWAI_ERROR_BUS_OFF, ESP_ERR_INVALID_STATE, TAG, "node is bus off");
// either hw_busy or tx_mount_queue is not empty, means tx is not finished
// either hw_busy or tx_queue is not empty, means tx is not finished
// otherwise, hardware is idle, return immediately
if (atomic_load(&twai_ctx->hw_busy) || uxQueueMessagesWaiting(twai_ctx->tx_mount_queue)) {
if (atomic_load(&twai_ctx->hw_busy) || twai_frame_queue_get_count(twai_ctx->tx_queue)) {
//wait for idle event bit but without clear it, every tasks block here can be waked up
if (TWAI_IDLE_EVENT_BIT != xEventGroupWaitBits(twai_ctx->event_group, TWAI_IDLE_EVENT_BIT, pdFALSE, pdFALSE, ticks_to_wait)) {
return ESP_ERR_TIMEOUT;
@@ -715,9 +693,11 @@ esp_err_t twai_new_node_onchip(const twai_onchip_node_config_t *node_config, twa
// state is in bus_off before enabled
atomic_store(&node->state, TWAI_ERROR_BUS_OFF);
node->tx_mount_queue = xQueueCreateWithCaps(node_config->tx_queue_depth, sizeof(twai_frame_t *), TWAI_MALLOC_CAPS);
if (!node_config->flags.enable_listen_only) {
ESP_GOTO_ON_ERROR(twai_frame_queue_new(&node->tx_queue, node_config->tx_queue_depth, TWAI_MALLOC_CAPS), err, TAG, "no_mem");
}
node->event_group = xEventGroupCreateWithCaps(TWAI_MALLOC_CAPS);
ESP_GOTO_ON_FALSE((node->tx_mount_queue && node->event_group) || node_config->flags.enable_listen_only, ESP_ERR_NO_MEM, err, TAG, "no_mem");
ESP_GOTO_ON_FALSE(node->event_group, ESP_ERR_NO_MEM, err, TAG, "no_mem");
uint32_t intr_flags = TWAI_INTR_ALLOC_FLAGS;
intr_flags |= (node_config->intr_priority > 0) ? BIT(node_config->intr_priority) : ESP_INTR_FLAG_LOWMED;
_lock_acquire(&s_platform.intr_mutex); // lock to prevent twai_intr and timer_intr registered to different cpu then triggered at the same time
@@ -0,0 +1,112 @@
/*
* SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
*
* SPDX-License-Identifier: Apache-2.0
*/
#pragma once
#include <stdbool.h>
#include <stddef.h>
#include <stdint.h>
#include "esp_err.h"
#include "esp_twai_types.h"
////////////////////////////////////////////////////////////////////
// !! This queue is ONLY for TWAI driver internal use //
////////////////////////////////////////////////////////////////////
#ifdef __cplusplus
extern "C" {
#endif
typedef struct twai_frame_queue_s *twai_frame_queue_t;
/**
* @brief Initialize a TWAI frame priority queue.
*
* Items are ordered by priority first. For equal priority values, items are popped in push order.
*/
esp_err_t twai_frame_queue_new(twai_frame_queue_t *queue, size_t capacity, uint32_t mem_caps);
/**
* @brief Release resources used by a TWAI frame priority queue.
*
* @return
* - ESP_OK: Queue was released successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
*/
esp_err_t twai_frame_queue_del(twai_frame_queue_t queue);
/**
* @brief Push an item with priority from task context.
*
* @return
* - ESP_OK: Item was queued successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_TIMEOUT: No free slot became available before timeout
*/
esp_err_t twai_frame_queue_push(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, int timeout_ms);
/**
* @brief Push an item with priority from ISR context.
*
* @return
* - ESP_OK: Item was queued successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_TIMEOUT: Queue is full
*/
esp_err_t twai_frame_queue_push_from_isr(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, bool *task_woken);
/**
* @brief Push an item with priority from task or ISR context.
*
* @return
* - ESP_OK: Item was queued successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_TIMEOUT: Queue is full or no free slot became available before timeout
*/
esp_err_t twai_frame_queue_push_safe(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, int timeout_ms);
/**
* @brief Pop the highest-priority item from task context.
*
* @return
* - ESP_OK: Item was popped successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_NOT_FOUND: Queue is empty
*/
esp_err_t twai_frame_queue_pop(twai_frame_queue_t queue, const twai_frame_t **data);
/**
* @brief Pop the highest-priority item from ISR context.
*
* @return
* - ESP_OK: Item was popped successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_NOT_FOUND: Queue is empty
*/
esp_err_t twai_frame_queue_pop_from_isr(twai_frame_queue_t queue, const twai_frame_t **data, bool *task_woken);
/**
* @brief Pop the highest-priority item from task or ISR context.
*
* @return
* - ESP_OK: Item was popped successfully
* - ESP_ERR_INVALID_ARG: Queue is invalid
* - ESP_ERR_NOT_FOUND: Queue is empty
*/
esp_err_t twai_frame_queue_pop_safe(twai_frame_queue_t queue, const twai_frame_t **data);
/**
* @brief Get the current number of queued items.
*/
size_t twai_frame_queue_get_count(twai_frame_queue_t queue);
/**
* @brief Get the number of available queue slots.
*/
size_t twai_frame_queue_get_free_space(twai_frame_queue_t queue);
#ifdef __cplusplus
}
#endif
@@ -33,6 +33,7 @@ typedef struct {
twai_frame_header_t header; /**< message attribute/metadata, exclude data buffer*/
uint8_t *buffer; /**< buffer address for tx and rx message data*/
size_t buffer_len; /**< buffer length of provided data buffer pointer, in bytes.*/
uint8_t tx_queue_priority; /**< Frame priority, range [0, 255], a frame with higher priority value will be picked first from the queue */
} twai_frame_t;
/**
+7
View File
@@ -9,10 +9,17 @@ entries:
esp_twai_onchip: _node_start_tx_batch_from_isr (noflash)
esp_twai_onchip: _node_mark_tx_idle (noflash)
esp_twai_onchip: _node_is_tx_all_done (noflash)
twai_frame_queue: twai_frame_queue_pop_from_isr (noflash)
if TWAI_IO_FUNC_IN_IRAM = y:
esp_twai_onchip: _node_queue_tx (noflash)
esp_twai: twai_node_transmit (noflash)
twai_frame_queue: twai_frame_queue_push (noflash)
twai_frame_queue: twai_frame_queue_push_from_isr (noflash)
twai_frame_queue: twai_frame_queue_push_safe (noflash)
twai_frame_queue: twai_frame_queue_pop (noflash)
twai_frame_queue: twai_frame_queue_pop_from_isr (noflash)
twai_frame_queue: twai_frame_queue_pop_safe (noflash)
[mapping:twai_hal]
archive: libesp_hal_twai.a
@@ -1,7 +1,7 @@
set(srcs "test_app_main.c")
if(CONFIG_SOC_TWAI_SUPPORTED)
list(APPEND srcs "test_twai_common.cpp" "test_twai_network.cpp")
list(APPEND srcs "test_twai_common.cpp" "test_twai_network.cpp" "test_twai_queue.cpp")
if(CONFIG_SOC_LIGHT_SLEEP_SUPPORTED)
list(APPEND srcs "test_twai_sleep.c")
endif()
@@ -0,0 +1,282 @@
/*
* SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
*
* SPDX-License-Identifier: Apache-2.0
*/
#include <stdbool.h>
#include <stddef.h>
#include <stdio.h>
#include <stdint.h>
#include <inttypes.h>
#include <stdatomic.h>
#include "esp_err.h"
#include "esp_heap_caps.h"
#include "esp_twai_types.h"
#include "freertos/FreeRTOS.h"
#include "freertos/semphr.h"
#include "freertos/task.h"
#include "unity.h"
#include "esp_private/twai_frame_queue.h"
typedef struct {
int value;
uint32_t priority;
} twai_frame_queue_test_item_t;
static void test_twai_frame_queue_print_header(const char *title)
{
printf("\n%s\n", title);
printf("+-------+----------+\n");
printf("| value | priority |\n");
printf("+-------+----------+\n");
}
static void test_twai_frame_queue_print_row(const twai_frame_queue_test_item_t *item)
{
printf("| %5d | %8" PRIu32 " |\n", item->value, item->priority);
}
static void test_twai_frame_queue_print_footer(void)
{
printf("+-------+----------+\n");
}
static void test_twai_frame_queue_pop_in_order(twai_frame_queue_t queue, size_t expected_count, const char *title)
{
const twai_frame_t *item = NULL;
const twai_frame_queue_test_item_t *last_item = NULL;
test_twai_frame_queue_print_header(title);
for (size_t i = 0; i < expected_count; i++) {
TEST_ESP_OK(twai_frame_queue_pop(queue, &item));
const twai_frame_queue_test_item_t *cur_item = (const twai_frame_queue_test_item_t *) item;
test_twai_frame_queue_print_row(cur_item);
if (last_item) {
TEST_ASSERT_LESS_OR_EQUAL(last_item->priority, cur_item->priority);
if (cur_item->priority == last_item->priority) {
TEST_ASSERT_GREATER_THAN(last_item->value, cur_item->value);
}
}
TEST_ASSERT_EQUAL(expected_count - i - 1, twai_frame_queue_get_count(queue));
last_item = cur_item;
}
test_twai_frame_queue_print_footer();
}
TEST_CASE("test twai frame priority queue", "[twai]")
{
twai_frame_queue_t test_q = NULL;
const twai_frame_t *item = NULL;
twai_frame_queue_test_item_t test_items[] = {
{ 0, 1 }, { 1, 4 }, { 2, 2 }, { 3, 5 },
{ 4, 3 }, { 5, 5 }, { 6, 1 }, { 7, 4 },
{ 8, 2 }, { 9, 3 }, { 10, 5 }, { 11, 0 },
{ 12, 4 }, { 13, 2 }, { 14, 3 }, { 15, 1 },
};
twai_frame_queue_test_item_t overflow = { 16, 6 };
const size_t capacity = sizeof(test_items) / sizeof(test_items[0]);
TEST_ESP_OK(twai_frame_queue_new(&test_q, capacity, MALLOC_CAP_DEFAULT));
TEST_ASSERT_EQUAL(0, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity, twai_frame_queue_get_free_space(test_q));
test_twai_frame_queue_print_header("Push all test data");
for (size_t i = 0; i < capacity; i++) {
TEST_ESP_OK(twai_frame_queue_push(test_q, (const twai_frame_t *) &test_items[i], test_items[i].priority, 0));
test_twai_frame_queue_print_row(&test_items[i]);
TEST_ASSERT_EQUAL(i + 1, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity - i - 1, twai_frame_queue_get_free_space(test_q));
}
test_twai_frame_queue_print_footer();
TEST_ASSERT_EQUAL(ESP_ERR_TIMEOUT, twai_frame_queue_push(test_q, (const twai_frame_t *) &overflow, 6, 0));
TEST_ASSERT_EQUAL(capacity, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(0, twai_frame_queue_get_free_space(test_q));
test_twai_frame_queue_pop_in_order(test_q, capacity, "Pop all test data");
TEST_ASSERT_EQUAL(ESP_ERR_NOT_FOUND, twai_frame_queue_pop(test_q, &item));
TEST_ASSERT_EQUAL(0, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity, twai_frame_queue_get_free_space(test_q));
TEST_ESP_OK(twai_frame_queue_del(test_q));
}
TEST_CASE("test twai queue mixed pop and push", "[twai]")
{
twai_frame_queue_t test_q = NULL;
const twai_frame_t *item = NULL;
twai_frame_queue_test_item_t test_items[] = {
{ 0, 1 },
{ 1, 5 },
{ 2, 3 },
{ 3, 5 },
{ 4, 2 },
{ 5, 4 },
{ 6, 6 },
{ 7, 3 },
{ 8, 4 },
};
const size_t first_push_count = 6;
const size_t capacity = 8;
const size_t remaining_count = sizeof(test_items) / sizeof(test_items[0]) - 2;
TEST_ESP_OK(twai_frame_queue_new(&test_q, capacity, MALLOC_CAP_DEFAULT));
test_twai_frame_queue_print_header("Initial push before refill");
for (size_t i = 0; i < first_push_count; i++) {
TEST_ESP_OK(twai_frame_queue_push(test_q, (const twai_frame_t *) &test_items[i], test_items[i].priority, 0));
test_twai_frame_queue_print_row(&test_items[i]);
}
test_twai_frame_queue_print_footer();
TEST_ASSERT_EQUAL(first_push_count, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity - first_push_count, twai_frame_queue_get_free_space(test_q));
test_twai_frame_queue_print_header("Pop before refill");
TEST_ESP_OK(twai_frame_queue_pop(test_q, &item));
test_twai_frame_queue_print_row((const twai_frame_queue_test_item_t *) item);
TEST_ASSERT_EQUAL_PTR(&test_items[1], item);
TEST_ESP_OK(twai_frame_queue_pop(test_q, &item));
test_twai_frame_queue_print_row((const twai_frame_queue_test_item_t *) item);
TEST_ASSERT_EQUAL_PTR(&test_items[3], item);
test_twai_frame_queue_print_footer();
TEST_ASSERT_EQUAL(4, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(4, twai_frame_queue_get_free_space(test_q));
test_twai_frame_queue_print_header("Push after partial pop");
for (size_t i = first_push_count; i < sizeof(test_items) / sizeof(test_items[0]); i++) {
TEST_ESP_OK(twai_frame_queue_push(test_q, (const twai_frame_t *) &test_items[i], test_items[i].priority, 0));
test_twai_frame_queue_print_row(&test_items[i]);
}
test_twai_frame_queue_print_footer();
TEST_ASSERT_EQUAL(remaining_count, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(1, twai_frame_queue_get_free_space(test_q));
test_twai_frame_queue_pop_in_order(test_q, remaining_count, "Pop after refill");
TEST_ASSERT_EQUAL(ESP_ERR_NOT_FOUND, twai_frame_queue_pop(test_q, &item));
TEST_ASSERT_EQUAL(0, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity, twai_frame_queue_get_free_space(test_q));
TEST_ESP_OK(twai_frame_queue_del(test_q));
}
typedef struct {
twai_frame_queue_t queue;
twai_frame_queue_test_item_t *items;
bool *seen;
size_t total;
atomic_int push_tasks_done;
atomic_int pop_count;
atomic_bool failed;
} twai_frame_queue_concurrent_ctx_t;
typedef struct {
twai_frame_queue_concurrent_ctx_t *ctx;
size_t begin;
size_t end;
} twai_frame_queue_push_args_t;
static void twai_frame_queue_push_task(void *arg)
{
twai_frame_queue_push_args_t *args = (twai_frame_queue_push_args_t *) arg;
twai_frame_queue_concurrent_ctx_t *ctx = args->ctx;
printf("%s started\n", pcTaskGetName(NULL));
for (size_t i = args->begin; i < args->end; i++) {
if (atomic_load(&ctx->failed)) {
vTaskDelete(NULL);
}
TEST_ESP_OK(twai_frame_queue_push(ctx->queue, (const twai_frame_t *) &ctx->items[i], ctx->items[i].priority, portMAX_DELAY));
vTaskDelay(1);
}
atomic_fetch_add(&ctx->push_tasks_done, 1);
vTaskDelete(NULL);
}
static void twai_frame_queue_pop_task(void *arg)
{
twai_frame_queue_concurrent_ctx_t *ctx = (twai_frame_queue_concurrent_ctx_t *) arg;
printf("%s started\n", pcTaskGetName(NULL));
const twai_frame_t *item = NULL;
while (atomic_load(&ctx->pop_count) < (int) ctx->total) {
if (twai_frame_queue_pop(ctx->queue, &item) == ESP_OK) {
const twai_frame_queue_test_item_t *cur = (const twai_frame_queue_test_item_t *) item;
int value = cur->value;
if (value < 0 || (size_t) value >= ctx->total || ctx->seen[value]) {
atomic_store(&ctx->failed, true);
vTaskDelete(NULL);
}
ctx->seen[value] = true;
atomic_fetch_add(&ctx->pop_count, 1);
} else if (atomic_load(&ctx->push_tasks_done) == 2 && twai_frame_queue_get_count(ctx->queue) == 0) {
break;
} else {
vTaskDelay(1);
}
}
vTaskDelete(NULL);
}
TEST_CASE("test twai queue concurrent push and pop", "[twai]")
{
twai_frame_queue_t test_q = NULL;
const size_t item_count = 500;
const size_t capacity = 4;
twai_frame_queue_test_item_t *test_items = (twai_frame_queue_test_item_t *) malloc(item_count * sizeof(twai_frame_queue_test_item_t));
bool seen[item_count] = {};
twai_frame_queue_concurrent_ctx_t ctx = {};
twai_frame_queue_push_args_t push_args0 = {};
twai_frame_queue_push_args_t push_args1 = {};
for (size_t i = 0; i < item_count; i++) {
test_items[i].value = (int) i;
test_items[i].priority = (uint32_t)((i * 7 + 3) % 6);
}
ctx.items = test_items;
ctx.seen = seen;
ctx.total = item_count;
push_args0.ctx = &ctx;
push_args0.begin = 0;
push_args0.end = item_count / 2;
push_args1.ctx = &ctx;
push_args1.begin = item_count / 2;
push_args1.end = item_count;
TEST_ESP_OK(twai_frame_queue_new(&test_q, capacity, MALLOC_CAP_DEFAULT));
ctx.queue = test_q;
TEST_ASSERT_EQUAL(pdPASS, xTaskCreate(twai_frame_queue_pop_task, "twai_q_pop", 4096, &ctx, 5, NULL));
TEST_ASSERT_EQUAL(pdPASS, xTaskCreate(twai_frame_queue_push_task, "twai_q_push0", 4096, &push_args0, 4, NULL));
TEST_ASSERT_EQUAL(pdPASS, xTaskCreate(twai_frame_queue_push_task, "twai_q_push1", 4096, &push_args1, 4, NULL));
TickType_t start = xTaskGetTickCount();
while (atomic_load(&ctx.push_tasks_done) < 2 || atomic_load(&ctx.pop_count) < (int) item_count) {
TEST_ASSERT_FALSE(atomic_load(&ctx.failed));
if ((xTaskGetTickCount() - start) > pdMS_TO_TICKS(5000)) {
TEST_FAIL_MESSAGE("concurrent push/pop timed out");
}
vTaskDelay(1);
}
TEST_ASSERT_FALSE(atomic_load(&ctx.failed));
TEST_ASSERT_EQUAL(2, atomic_load(&ctx.push_tasks_done));
printf("pop %d items\n", atomic_load(&ctx.pop_count));
TEST_ASSERT_EQUAL((int) item_count, atomic_load(&ctx.pop_count));
TEST_ASSERT_EQUAL(0, twai_frame_queue_get_count(test_q));
TEST_ASSERT_EQUAL(capacity, twai_frame_queue_get_free_space(test_q));
for (size_t i = 0; i < item_count; i++) {
TEST_ASSERT_TRUE(seen[i]);
}
printf("test finished\n");
free(test_items);
TEST_ESP_OK(twai_frame_queue_del(test_q));
vTaskDelay(10); // wait for tasks to be deleted
}
@@ -0,0 +1,277 @@
/*
* SPDX-FileCopyrightText: 2026 Espressif Systems (Shanghai) CO LTD
*
* SPDX-License-Identifier: Apache-2.0
*/
#include <string.h>
#include "esp_attr.h"
#include "esp_heap_caps.h"
#include "freertos/FreeRTOS.h"
#include "freertos/semphr.h"
#include "esp_private/twai_frame_queue.h"
/**
* @brief Priority queue item for TWAI frame pointers.
*
* The queue stores TWAI frame pointers and orders them by priority.
* Higher priority values are popped first; equal priorities keep FIFO order.
*/
typedef struct {
const twai_frame_t *data; /**< TWAI frame pointer stored in the queue */
uint32_t priority; /**< Higher value is dequeued first */
uint64_t seq; /**< Sequence number used to keep FIFO order for equal priority */
} twai_frame_queue_item_t;
/**
* @brief Priority queue used by the TWAI driver to schedule queued frames.
*/
struct twai_frame_queue_s {
portMUX_TYPE spinlock; /**< Protects heap storage in task and ISR context */
SemaphoreHandle_t spaces_sem; /**< Counts available queue slots */
twai_frame_queue_item_t *items; /**< Binary heap storage */
size_t capacity; /**< Maximum number of items */
size_t count; /**< Current number of queued items */
uint64_t next_seq; /**< Next sequence number assigned on push */
};
static inline IRAM_ATTR bool twai_frame_queue_item_greater(const twai_frame_queue_item_t *a, const twai_frame_queue_item_t *b)
{
if (a->priority != b->priority) {
return a->priority > b->priority;
}
return a->seq < b->seq;
}
static inline IRAM_ATTR void twai_frame_queue_swap_item(twai_frame_queue_item_t *a, twai_frame_queue_item_t *b)
{
twai_frame_queue_item_t tmp = *a;
*a = *b;
*b = tmp;
}
static IRAM_ATTR void twai_frame_queue_sift_up(twai_frame_queue_t queue, size_t index)
{
while (index > 0) {
size_t parent = (index - 1) / 2;
if (!twai_frame_queue_item_greater(&queue->items[index], &queue->items[parent])) {
break;
}
twai_frame_queue_swap_item(&queue->items[index], &queue->items[parent]);
index = parent;
}
}
static IRAM_ATTR void twai_frame_queue_sift_down(twai_frame_queue_t queue, size_t index)
{
while (true) {
size_t left = index * 2 + 1;
size_t right = left + 1;
size_t highest = index;
if ((left < queue->count) && twai_frame_queue_item_greater(&queue->items[left], &queue->items[highest])) {
highest = left;
}
if ((right < queue->count) && twai_frame_queue_item_greater(&queue->items[right], &queue->items[highest])) {
highest = right;
}
if (highest == index) {
break;
}
twai_frame_queue_swap_item(&queue->items[index], &queue->items[highest]);
index = highest;
}
}
static inline IRAM_ATTR void twai_frame_queue_push_locked(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority)
{
size_t index = queue->count++;
queue->items[index] = (twai_frame_queue_item_t) {
.data = data,
.priority = priority,
.seq = queue->next_seq++,
};
twai_frame_queue_sift_up(queue, index);
}
static inline IRAM_ATTR bool twai_frame_queue_pop_locked(twai_frame_queue_t queue, const twai_frame_t **data)
{
if (queue->count == 0) {
return false;
}
*data = queue->items[0].data;
queue->count--;
if (queue->count > 0) {
queue->items[0] = queue->items[queue->count];
twai_frame_queue_sift_down(queue, 0);
}
return true;
}
esp_err_t twai_frame_queue_new(twai_frame_queue_t *queue, size_t capacity, uint32_t mem_caps)
{
if (!queue || !capacity) {
return ESP_ERR_INVALID_ARG;
}
twai_frame_queue_t q_ctx = heap_caps_calloc(1, sizeof(struct twai_frame_queue_s), mem_caps);
if (!q_ctx) {
return ESP_ERR_NO_MEM;
}
q_ctx->items = heap_caps_calloc(capacity, sizeof(twai_frame_queue_item_t), mem_caps);
if (!q_ctx->items) {
heap_caps_free(q_ctx);
return ESP_ERR_NO_MEM;
}
q_ctx->spaces_sem = xSemaphoreCreateCountingWithCaps(capacity, capacity, mem_caps);
if (!q_ctx->spaces_sem) {
heap_caps_free(q_ctx->items);
heap_caps_free(q_ctx);
return ESP_ERR_NO_MEM;
}
q_ctx->spinlock = (portMUX_TYPE) portMUX_INITIALIZER_UNLOCKED;
q_ctx->capacity = capacity;
*queue = q_ctx;
return ESP_OK;
}
esp_err_t twai_frame_queue_del(twai_frame_queue_t queue)
{
if (!queue) {
return ESP_ERR_INVALID_ARG;
}
if (queue->spaces_sem) {
vSemaphoreDeleteWithCaps(queue->spaces_sem);
}
if (queue->items) {
heap_caps_free(queue->items);
}
heap_caps_free(queue);
return ESP_OK;
}
esp_err_t twai_frame_queue_push(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, int timeout_ms)
{
if (!queue || !queue->spaces_sem) {
return ESP_ERR_INVALID_ARG;
}
TickType_t ticks_to_wait = (timeout_ms == -1) ? portMAX_DELAY : pdMS_TO_TICKS(timeout_ms);
if (xSemaphoreTake(queue->spaces_sem, ticks_to_wait) != pdTRUE) {
return ESP_ERR_TIMEOUT;
}
portENTER_CRITICAL(&queue->spinlock);
twai_frame_queue_push_locked(queue, data, priority);
portEXIT_CRITICAL(&queue->spinlock);
return ESP_OK;
}
esp_err_t twai_frame_queue_push_from_isr(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, bool *task_woken)
{
if (!queue || !queue->spaces_sem) {
return ESP_ERR_INVALID_ARG;
}
if (xSemaphoreTakeFromISR(queue->spaces_sem, (BaseType_t *)task_woken) != pdTRUE) {
return ESP_ERR_TIMEOUT;
}
portENTER_CRITICAL_ISR(&queue->spinlock);
twai_frame_queue_push_locked(queue, data, priority);
portEXIT_CRITICAL_ISR(&queue->spinlock);
return ESP_OK;
}
esp_err_t twai_frame_queue_push_safe(twai_frame_queue_t queue, const twai_frame_t *data, uint32_t priority, int timeout_ms)
{
if (xPortInIsrContext()) {
bool task_woken = false;
esp_err_t ret = twai_frame_queue_push_from_isr(queue, data, priority, &task_woken);
if (task_woken) {
portYIELD_FROM_ISR();
}
return ret;
}
TickType_t ticks_to_wait = (timeout_ms == -1) ? portMAX_DELAY : pdMS_TO_TICKS(timeout_ms);
return twai_frame_queue_push(queue, data, priority, ticks_to_wait);
}
esp_err_t twai_frame_queue_pop(twai_frame_queue_t queue, const twai_frame_t **data)
{
bool ret;
if (!queue || !queue->spaces_sem) {
return ESP_ERR_INVALID_ARG;
}
portENTER_CRITICAL(&queue->spinlock);
ret = twai_frame_queue_pop_locked(queue, data);
portEXIT_CRITICAL(&queue->spinlock);
if (ret) {
xSemaphoreGive(queue->spaces_sem);
}
return ret ? ESP_OK : ESP_ERR_NOT_FOUND;
}
esp_err_t twai_frame_queue_pop_from_isr(twai_frame_queue_t queue, const twai_frame_t **data, bool *task_woken)
{
bool ret;
if (!queue || !queue->spaces_sem) {
return ESP_ERR_INVALID_ARG;
}
portENTER_CRITICAL_ISR(&queue->spinlock);
ret = twai_frame_queue_pop_locked(queue, data);
portEXIT_CRITICAL_ISR(&queue->spinlock);
if (ret) {
xSemaphoreGiveFromISR(queue->spaces_sem, (BaseType_t *)task_woken);
}
return ret ? ESP_OK : ESP_ERR_NOT_FOUND;
}
esp_err_t twai_frame_queue_pop_safe(twai_frame_queue_t queue, const twai_frame_t **data)
{
if (xPortInIsrContext()) {
bool task_woken = false;
esp_err_t ret = twai_frame_queue_pop_from_isr(queue, data, &task_woken);
if (task_woken) {
portYIELD_FROM_ISR();
}
return ret;
}
return twai_frame_queue_pop(queue, data);
}
size_t twai_frame_queue_get_count(twai_frame_queue_t queue)
{
size_t count;
if (!queue || !queue->spaces_sem) {
return 0;
}
portENTER_CRITICAL(&queue->spinlock);
count = queue->count;
portEXIT_CRITICAL(&queue->spinlock);
return count;
}
size_t twai_frame_queue_get_free_space(twai_frame_queue_t queue)
{
size_t count;
if (!queue || !queue->spaces_sem) {
return 0;
}
portENTER_CRITICAL(&queue->spinlock);
count = queue->count;
portEXIT_CRITICAL(&queue->spinlock);
return queue->capacity - count;
}
@@ -10,6 +10,7 @@
#include <string.h>
#include <sys/cdefs.h>
#include <sys/lock.h>
#include <sys/param.h>
#include <stdatomic.h>
#include "sdkconfig.h"
#if CONFIG_TWAI_ENABLE_DEBUG_LOG