Add native USB CDC broker transport
This commit is contained in:
@@ -0,0 +1,989 @@
|
||||
#include "usb_cdc_transport.h"
|
||||
|
||||
#include <stdatomic.h>
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
|
||||
#include "esp_mac.h"
|
||||
#include "freertos/FreeRTOS.h"
|
||||
#include "freertos/queue.h"
|
||||
#include "freertos/stream_buffer.h"
|
||||
#include "freertos/task.h"
|
||||
#include "serial_config.h"
|
||||
#include "serial_service.h"
|
||||
#include "tinyusb.h"
|
||||
#include "tinyusb_cdc_acm.h"
|
||||
#include "tinyusb_default_config.h"
|
||||
|
||||
#define USB_CDC_HOST_RX_STREAM_SIZE 4096U
|
||||
#define USB_CDC_IO_CHUNK_SIZE 256U
|
||||
#define USB_CDC_CONTROL_QUEUE_LENGTH 8U
|
||||
#define USB_CDC_TASK_STACK_SIZE 5632U
|
||||
#define USB_CDC_TASK_PRIORITY 8U
|
||||
#define USB_CDC_POLL_MS 2U
|
||||
#define USB_CDC_SERVICE_RETRY_MS 250U
|
||||
#define USB_CDC_CONNECT_RETRY_MS 100U
|
||||
#define USB_CDC_BROKER_RECONCILE_MS 100U
|
||||
|
||||
#define USB_STATE_ATTACHED (1U << 0)
|
||||
#define USB_STATE_DTR (1U << 1)
|
||||
#define USB_STATE_RTS (1U << 2)
|
||||
|
||||
typedef enum {
|
||||
USB_CDC_CONTROL_REQUEST_WRITER = 0,
|
||||
USB_CDC_CONTROL_RELEASE_WRITER,
|
||||
} usb_cdc_control_t;
|
||||
|
||||
typedef struct {
|
||||
uint8_t data[USB_CDC_IO_CHUNK_SIZE];
|
||||
size_t size;
|
||||
size_t offset;
|
||||
} usb_cdc_pending_buffer_t;
|
||||
|
||||
static StreamBufferHandle_t s_host_rx_stream;
|
||||
static QueueHandle_t s_control_queue;
|
||||
static atomic_uintptr_t s_transport_task;
|
||||
static portMUX_TYPE s_state_lock = portMUX_INITIALIZER_UNLOCKED;
|
||||
|
||||
static atomic_bool s_initialized;
|
||||
static atomic_bool s_initializing;
|
||||
static atomic_uint s_usb_state;
|
||||
/* Changes on every effective CDC open/close boundary, even during one task poll. */
|
||||
static atomic_uint s_connection_generation;
|
||||
|
||||
static session_broker_client_id_t s_broker_client_id;
|
||||
static bool s_writer;
|
||||
static usb_cdc_transport_line_coding_t s_line_coding;
|
||||
static bool s_line_coding_pending;
|
||||
static usb_cdc_transport_counters_t s_counters;
|
||||
|
||||
static const char s_language_descriptor[] = {0x09, 0x04};
|
||||
static const char s_manufacturer[] = "ESP32 Serial Tools";
|
||||
/* esp_tinyusb's default UTF-16 conversion accepts at most 31 characters. */
|
||||
static const char s_product[] = "ESP32 Serial Swiss Army Knife";
|
||||
static char s_serial_number[13];
|
||||
static const char s_cdc_interface[] = "USB CDC";
|
||||
static const char *s_string_descriptors[] = {
|
||||
s_language_descriptor,
|
||||
s_manufacturer,
|
||||
s_product,
|
||||
s_serial_number,
|
||||
s_cdc_interface,
|
||||
};
|
||||
|
||||
static TickType_t milliseconds_to_ticks(uint32_t milliseconds)
|
||||
{
|
||||
TickType_t ticks = pdMS_TO_TICKS(milliseconds);
|
||||
return (milliseconds > 0U && ticks == 0U) ? 1U : ticks;
|
||||
}
|
||||
|
||||
static void add_counter(uint64_t *counter, uint64_t amount)
|
||||
{
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
*counter += amount;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
static void notify_transport_task(void)
|
||||
{
|
||||
TaskHandle_t task = (TaskHandle_t)atomic_load(&s_transport_task);
|
||||
if (task != NULL) {
|
||||
xTaskNotifyGive(task);
|
||||
}
|
||||
}
|
||||
|
||||
static bool usb_state_is_open(unsigned int state)
|
||||
{
|
||||
return (state & (USB_STATE_ATTACHED | USB_STATE_DTR)) ==
|
||||
(USB_STATE_ATTACHED | USB_STATE_DTR);
|
||||
}
|
||||
|
||||
static bool usb_host_port_open(void)
|
||||
{
|
||||
return usb_state_is_open(atomic_load(&s_usb_state));
|
||||
}
|
||||
|
||||
static void set_broker_state(session_broker_client_id_t client_id, bool writer)
|
||||
{
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_broker_client_id = client_id;
|
||||
s_writer = writer;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
static void set_writer_state(bool writer)
|
||||
{
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_writer = writer;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
static void device_event_callback(tinyusb_event_t *event, void *arg)
|
||||
{
|
||||
(void)arg;
|
||||
|
||||
if (event == NULL) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
switch (event->id) {
|
||||
case TINYUSB_EVENT_ATTACHED:
|
||||
atomic_fetch_or(&s_usb_state, USB_STATE_ATTACHED);
|
||||
notify_transport_task();
|
||||
break;
|
||||
case TINYUSB_EVENT_DETACHED: {
|
||||
/* A new attachment must receive fresh control state and line coding. */
|
||||
unsigned int old_state = atomic_exchange(&s_usb_state, 0U);
|
||||
if (usb_state_is_open(old_state)) {
|
||||
atomic_fetch_add(&s_connection_generation, 1U);
|
||||
}
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_line_coding_pending = false;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
notify_transport_task();
|
||||
break;
|
||||
}
|
||||
default:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
static void cdc_rx_callback(int itf, cdcacm_event_t *event)
|
||||
{
|
||||
if (itf != TINYUSB_CDC_ACM_0 || event == NULL || event->type != CDC_EVENT_RX) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
uint8_t data[USB_CDC_IO_CHUNK_SIZE];
|
||||
for (;;) {
|
||||
size_t received = 0U;
|
||||
esp_err_t result = tinyusb_cdcacm_read(TINYUSB_CDC_ACM_0,
|
||||
data,
|
||||
sizeof(data),
|
||||
&received);
|
||||
if (result != ESP_OK) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
break;
|
||||
}
|
||||
if (received == 0U) {
|
||||
break;
|
||||
}
|
||||
|
||||
size_t queued = 0U;
|
||||
if (s_host_rx_stream != NULL) {
|
||||
/* This callback is the stream's only writer and never waits. */
|
||||
queued = xStreamBufferSend(s_host_rx_stream, data, received, 0U);
|
||||
}
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_counters.host_rx_bytes += received;
|
||||
s_counters.host_rx_stream_dropped_bytes += received - queued;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
notify_transport_task();
|
||||
}
|
||||
|
||||
static void cdc_wanted_char_callback(int itf, cdcacm_event_t *event)
|
||||
{
|
||||
if (itf != TINYUSB_CDC_ACM_0 || event == NULL ||
|
||||
event->type != CDC_EVENT_RX_WANTED_CHAR) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
/* Wanted-character notifications are only a scheduling hint here. */
|
||||
notify_transport_task();
|
||||
}
|
||||
|
||||
static void cdc_line_state_callback(int itf, cdcacm_event_t *event)
|
||||
{
|
||||
if (itf != TINYUSB_CDC_ACM_0 || event == NULL ||
|
||||
event->type != CDC_EVENT_LINE_STATE_CHANGED) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
unsigned int old_state = atomic_load(&s_usb_state);
|
||||
for (;;) {
|
||||
/* A valid CDC class callback is also proof that this device is attached. */
|
||||
unsigned int new_state = (old_state | USB_STATE_ATTACHED) &
|
||||
~(USB_STATE_DTR | USB_STATE_RTS);
|
||||
if (event->line_state_changed_data.dtr) {
|
||||
new_state |= USB_STATE_DTR;
|
||||
}
|
||||
if (event->line_state_changed_data.rts) {
|
||||
new_state |= USB_STATE_RTS;
|
||||
}
|
||||
if (atomic_compare_exchange_weak(&s_usb_state, &old_state, new_state)) {
|
||||
if (usb_state_is_open(old_state) != usb_state_is_open(new_state)) {
|
||||
atomic_fetch_add(&s_connection_generation, 1U);
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if (!event->line_state_changed_data.dtr) {
|
||||
/* Do not apply a closed host session's deferred line coding after reopen. */
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_line_coding_pending = false;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
notify_transport_task();
|
||||
}
|
||||
|
||||
static void cdc_line_coding_callback(int itf, cdcacm_event_t *event)
|
||||
{
|
||||
if (itf != TINYUSB_CDC_ACM_0 || event == NULL ||
|
||||
event->type != CDC_EVENT_LINE_CODING_CHANGED ||
|
||||
event->line_coding_changed_data.p_line_coding == NULL) {
|
||||
add_counter(&s_counters.callback_drops, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
/* The TinyUSB value is packed and callback-owned, so copy it immediately. */
|
||||
cdc_line_coding_t coding;
|
||||
memcpy(&coding,
|
||||
event->line_coding_changed_data.p_line_coding,
|
||||
sizeof(coding));
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
if (s_line_coding_pending) {
|
||||
/* Preserve the latest complete setting and account for the superseded one. */
|
||||
++s_counters.callback_drops;
|
||||
}
|
||||
s_line_coding = (usb_cdc_transport_line_coding_t) {
|
||||
.baud_rate = coding.bit_rate,
|
||||
.stop_bits = coding.stop_bits,
|
||||
.parity = coding.parity,
|
||||
.data_bits = coding.data_bits,
|
||||
};
|
||||
s_line_coding_pending = true;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
|
||||
notify_transport_task();
|
||||
}
|
||||
|
||||
static bool take_pending_line_coding(usb_cdc_transport_line_coding_t *coding)
|
||||
{
|
||||
bool pending;
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
pending = s_line_coding_pending;
|
||||
if (pending) {
|
||||
*coding = s_line_coding;
|
||||
s_line_coding_pending = false;
|
||||
}
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
return pending;
|
||||
}
|
||||
|
||||
static bool serial_configs_equal(const serial_config_t *left,
|
||||
const serial_config_t *right)
|
||||
{
|
||||
return left->version == right->version &&
|
||||
left->baud_rate == right->baud_rate &&
|
||||
left->data_bits == right->data_bits &&
|
||||
left->parity == right->parity &&
|
||||
left->stop_bits == right->stop_bits &&
|
||||
left->flow_control == right->flow_control &&
|
||||
left->dtr_behavior == right->dtr_behavior &&
|
||||
left->rts_threshold == right->rts_threshold;
|
||||
}
|
||||
|
||||
static bool map_line_coding(const usb_cdc_transport_line_coding_t *coding,
|
||||
serial_config_t *config)
|
||||
{
|
||||
if (coding->baud_rate < SERIAL_CONFIG_MIN_BAUD_RATE ||
|
||||
coding->baud_rate > SERIAL_CONFIG_MAX_BAUD_RATE) {
|
||||
return false;
|
||||
}
|
||||
config->baud_rate = coding->baud_rate;
|
||||
|
||||
switch (coding->data_bits) {
|
||||
case 7U:
|
||||
config->data_bits = SERIAL_CONFIG_DATA_BITS_7;
|
||||
break;
|
||||
case 8U:
|
||||
config->data_bits = SERIAL_CONFIG_DATA_BITS_8;
|
||||
break;
|
||||
default:
|
||||
return false;
|
||||
}
|
||||
|
||||
switch (coding->parity) {
|
||||
case CDC_LINE_CODING_PARITY_NONE:
|
||||
config->parity = SERIAL_CONFIG_PARITY_NONE;
|
||||
break;
|
||||
case CDC_LINE_CODING_PARITY_ODD:
|
||||
config->parity = SERIAL_CONFIG_PARITY_ODD;
|
||||
break;
|
||||
case CDC_LINE_CODING_PARITY_EVEN:
|
||||
config->parity = SERIAL_CONFIG_PARITY_EVEN;
|
||||
break;
|
||||
default:
|
||||
/* Mark and space parity are intentionally not representable by UART policy. */
|
||||
return false;
|
||||
}
|
||||
|
||||
switch (coding->stop_bits) {
|
||||
case CDC_LINE_CODING_STOP_BITS_1:
|
||||
config->stop_bits = SERIAL_CONFIG_STOP_BITS_1;
|
||||
break;
|
||||
case CDC_LINE_CODING_STOP_BITS_2:
|
||||
config->stop_bits = SERIAL_CONFIG_STOP_BITS_2;
|
||||
break;
|
||||
default:
|
||||
/* This also rejects USB's 1.5-stop-bit encoding. */
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
static void apply_pending_line_coding(bool writer)
|
||||
{
|
||||
if (!writer || !serial_service_is_running()) {
|
||||
return;
|
||||
}
|
||||
|
||||
/* Restarting UART1 discards queued TX, so defer framing changes until idle. */
|
||||
if (serial_service_tx_pending() > 0U) {
|
||||
return;
|
||||
}
|
||||
|
||||
usb_cdc_transport_line_coding_t coding;
|
||||
if (!take_pending_line_coding(&coding)) {
|
||||
return;
|
||||
}
|
||||
|
||||
serial_config_t current;
|
||||
if (serial_service_get_config(¤t) != ESP_OK) {
|
||||
add_counter(&s_counters.line_coding_failed, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
serial_config_t desired = current;
|
||||
if (!map_line_coding(&coding, &desired)) {
|
||||
add_counter(&s_counters.line_coding_rejected, 1U);
|
||||
return;
|
||||
}
|
||||
|
||||
/* Flow control, DTR policy, and RTS threshold remain from current RAM state. */
|
||||
if (serial_configs_equal(¤t, &desired)) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (serial_service_apply_config(&desired) == ESP_OK) {
|
||||
add_counter(&s_counters.line_coding_applied, 1U);
|
||||
} else {
|
||||
add_counter(&s_counters.line_coding_failed, 1U);
|
||||
}
|
||||
}
|
||||
|
||||
static bool writer_event_type(session_broker_event_type_t type)
|
||||
{
|
||||
return type == SESSION_BROKER_EVENT_WRITER_GRANTED ||
|
||||
type == SESSION_BROKER_EVENT_WRITER_RELEASED ||
|
||||
type == SESSION_BROKER_EVENT_WRITER_REVOKED ||
|
||||
type == SESSION_BROKER_EVENT_WRITER_DENIED;
|
||||
}
|
||||
|
||||
static esp_err_t reconcile_broker_state(session_broker_client_id_t client_id,
|
||||
bool *writer)
|
||||
{
|
||||
session_broker_client_snapshot_t snapshot;
|
||||
esp_err_t result = session_broker_get_client_snapshot(client_id, &snapshot);
|
||||
if (result == ESP_OK) {
|
||||
*writer = snapshot.is_writer;
|
||||
set_writer_state(*writer);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
static esp_err_t drain_broker_events(session_broker_client_id_t client_id,
|
||||
bool *writer)
|
||||
{
|
||||
session_broker_event_t event;
|
||||
for (;;) {
|
||||
esp_err_t result = session_broker_pop_event(client_id, &event);
|
||||
if (result == ESP_ERR_TIMEOUT) {
|
||||
return ESP_OK;
|
||||
}
|
||||
if (result != ESP_OK) {
|
||||
return result;
|
||||
}
|
||||
|
||||
/* writer_id is ownership after this event, including forced changes. */
|
||||
*writer = event.writer_id == client_id;
|
||||
set_writer_state(*writer);
|
||||
|
||||
if (!writer_event_type(event.type)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
++s_counters.writer_events;
|
||||
if (event.type == SESSION_BROKER_EVENT_WRITER_GRANTED &&
|
||||
event.client_id == client_id) {
|
||||
++s_counters.writer_grants;
|
||||
} else if (event.type == SESSION_BROKER_EVENT_WRITER_DENIED &&
|
||||
event.client_id == client_id) {
|
||||
++s_counters.writer_denials;
|
||||
} else if (event.type == SESSION_BROKER_EVENT_WRITER_REVOKED &&
|
||||
event.client_id == client_id) {
|
||||
++s_counters.writer_revocations;
|
||||
}
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
}
|
||||
|
||||
static void discard_host_input(usb_cdc_pending_buffer_t *pending)
|
||||
{
|
||||
uint64_t discarded = pending->size - pending->offset;
|
||||
pending->size = 0U;
|
||||
pending->offset = 0U;
|
||||
|
||||
uint8_t data[USB_CDC_IO_CHUNK_SIZE];
|
||||
size_t drain_budget = USB_CDC_HOST_RX_STREAM_SIZE;
|
||||
while (drain_budget > 0U) {
|
||||
size_t request = drain_budget < sizeof(data) ? drain_budget : sizeof(data);
|
||||
size_t received = xStreamBufferReceive(s_host_rx_stream, data, request, 0U);
|
||||
if (received == 0U) {
|
||||
break;
|
||||
}
|
||||
discarded += received;
|
||||
drain_budget -= received;
|
||||
}
|
||||
|
||||
if (discarded > 0U) {
|
||||
add_counter(&s_counters.broker_rejected_bytes, discarded);
|
||||
}
|
||||
}
|
||||
|
||||
static void discard_usb_pending(usb_cdc_pending_buffer_t *pending)
|
||||
{
|
||||
size_t discarded = pending->size - pending->offset;
|
||||
pending->size = 0U;
|
||||
pending->offset = 0U;
|
||||
if (discarded > 0U) {
|
||||
add_counter(&s_counters.usb_tx_dropped_bytes, discarded);
|
||||
}
|
||||
}
|
||||
|
||||
static esp_err_t move_host_data_to_broker(session_broker_client_id_t client_id,
|
||||
usb_cdc_pending_buffer_t *pending)
|
||||
{
|
||||
if (pending->offset == pending->size) {
|
||||
pending->size = xStreamBufferReceive(s_host_rx_stream,
|
||||
pending->data,
|
||||
sizeof(pending->data),
|
||||
0U);
|
||||
pending->offset = 0U;
|
||||
}
|
||||
if (pending->size == 0U) {
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
size_t accepted = 0U;
|
||||
esp_err_t result = session_broker_write(client_id,
|
||||
pending->data + pending->offset,
|
||||
pending->size - pending->offset,
|
||||
&accepted);
|
||||
size_t remaining = pending->size - pending->offset;
|
||||
if (accepted > remaining) {
|
||||
accepted = remaining;
|
||||
}
|
||||
pending->offset += accepted;
|
||||
if (accepted > 0U) {
|
||||
add_counter(&s_counters.broker_accepted_bytes, accepted);
|
||||
}
|
||||
if (pending->offset == pending->size) {
|
||||
pending->size = 0U;
|
||||
pending->offset = 0U;
|
||||
}
|
||||
|
||||
/* Timeout, zero acceptance, and partial acceptance retain exact byte order. */
|
||||
return result;
|
||||
}
|
||||
|
||||
static esp_err_t move_broker_data_to_usb(session_broker_client_id_t client_id,
|
||||
usb_cdc_pending_buffer_t *pending)
|
||||
{
|
||||
if (pending->offset == pending->size) {
|
||||
size_t received = 0U;
|
||||
esp_err_t result = session_broker_read(client_id,
|
||||
pending->data,
|
||||
sizeof(pending->data),
|
||||
&received);
|
||||
if (result != ESP_OK) {
|
||||
return result;
|
||||
}
|
||||
pending->size = received;
|
||||
pending->offset = 0U;
|
||||
if (received > 0U) {
|
||||
add_counter(&s_counters.broker_to_usb_bytes, received);
|
||||
}
|
||||
}
|
||||
|
||||
if (pending->offset < pending->size) {
|
||||
size_t remaining = pending->size - pending->offset;
|
||||
size_t queued = tinyusb_cdcacm_write_queue(TINYUSB_CDC_ACM_0,
|
||||
pending->data + pending->offset,
|
||||
remaining);
|
||||
if (queued > remaining) {
|
||||
queued = remaining;
|
||||
}
|
||||
pending->offset += queued;
|
||||
if (queued > 0U) {
|
||||
add_counter(&s_counters.usb_tx_queued_bytes, queued);
|
||||
}
|
||||
if (pending->offset == pending->size) {
|
||||
pending->size = 0U;
|
||||
pending->offset = 0U;
|
||||
}
|
||||
}
|
||||
|
||||
/* Keep endpoint servicing nonblocking even when the host stops reading. */
|
||||
(void)tinyusb_cdcacm_write_flush(TINYUSB_CDC_ACM_0, 0U);
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
static void discard_control_requests(void)
|
||||
{
|
||||
usb_cdc_control_t control;
|
||||
uint64_t discarded = 0U;
|
||||
while (xQueueReceive(s_control_queue, &control, 0U) == pdTRUE) {
|
||||
++discarded;
|
||||
}
|
||||
if (discarded > 0U) {
|
||||
add_counter(&s_counters.control_drops, discarded);
|
||||
}
|
||||
}
|
||||
|
||||
static esp_err_t process_control_requests(session_broker_client_id_t client_id,
|
||||
bool *writer)
|
||||
{
|
||||
usb_cdc_control_t control;
|
||||
while (xQueueReceive(s_control_queue, &control, 0U) == pdTRUE) {
|
||||
esp_err_t result;
|
||||
if (control == USB_CDC_CONTROL_REQUEST_WRITER) {
|
||||
result = session_broker_request_writer(client_id);
|
||||
/* A competing writer is an observed denial, not a lost control. */
|
||||
if (result != ESP_OK && result != ESP_ERR_INVALID_STATE) {
|
||||
add_counter(&s_counters.control_drops, 1U);
|
||||
}
|
||||
} else {
|
||||
result = session_broker_release_writer(client_id);
|
||||
/* Releasing while already an observer is an idempotent no-op. */
|
||||
if (result != ESP_OK && result != ESP_ERR_INVALID_STATE) {
|
||||
add_counter(&s_counters.control_drops, 1U);
|
||||
}
|
||||
}
|
||||
|
||||
if (result == ESP_ERR_NOT_FOUND) {
|
||||
return result;
|
||||
}
|
||||
esp_err_t event_result = drain_broker_events(client_id, writer);
|
||||
if (event_result != ESP_OK) {
|
||||
return event_result;
|
||||
}
|
||||
}
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
static bool retry_due(TickType_t now,
|
||||
TickType_t interval,
|
||||
TickType_t *last_attempt,
|
||||
bool *attempted)
|
||||
{
|
||||
if (!*attempted || (TickType_t)(now - *last_attempt) >= interval) {
|
||||
*last_attempt = now;
|
||||
*attempted = true;
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
static void mark_client_disconnected(session_broker_client_id_t *client_id,
|
||||
bool *writer)
|
||||
{
|
||||
*client_id = SESSION_BROKER_NO_CLIENT;
|
||||
*writer = false;
|
||||
set_broker_state(SESSION_BROKER_NO_CLIENT, false);
|
||||
add_counter(&s_counters.disconnections, 1U);
|
||||
}
|
||||
|
||||
static void transport_task(void *context)
|
||||
{
|
||||
(void)context;
|
||||
|
||||
session_broker_client_id_t client_id = SESSION_BROKER_NO_CLIENT;
|
||||
bool writer = false;
|
||||
unsigned int observed_generation = atomic_load(&s_connection_generation);
|
||||
TickType_t last_service_attempt = 0U;
|
||||
TickType_t last_connect_attempt = 0U;
|
||||
TickType_t last_reconcile = 0U;
|
||||
bool service_attempted = false;
|
||||
bool connect_attempted = false;
|
||||
bool reconciled = false;
|
||||
usb_cdc_pending_buffer_t host_pending = {0};
|
||||
usb_cdc_pending_buffer_t usb_pending = {0};
|
||||
const TickType_t poll_ticks = milliseconds_to_ticks(USB_CDC_POLL_MS);
|
||||
const TickType_t service_retry_ticks =
|
||||
milliseconds_to_ticks(USB_CDC_SERVICE_RETRY_MS);
|
||||
const TickType_t connect_retry_ticks =
|
||||
milliseconds_to_ticks(USB_CDC_CONNECT_RETRY_MS);
|
||||
const TickType_t reconcile_ticks =
|
||||
milliseconds_to_ticks(USB_CDC_BROKER_RECONCILE_MS);
|
||||
|
||||
for (;;) {
|
||||
unsigned int generation = atomic_load(&s_connection_generation);
|
||||
if (generation != observed_generation) {
|
||||
/* Never let pending data or a broker identity cross a CDC session. */
|
||||
service_attempted = false;
|
||||
connect_attempted = false;
|
||||
reconciled = false;
|
||||
discard_control_requests();
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
|
||||
if (client_id != SESSION_BROKER_NO_CLIENT) {
|
||||
esp_err_t result = session_broker_disconnect(client_id);
|
||||
if (result == ESP_OK || result == ESP_ERR_NOT_FOUND) {
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
}
|
||||
if (client_id != SESSION_BROKER_NO_CLIENT) {
|
||||
/* Leave the generation unmatched so cleanup is retried. */
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
observed_generation = generation;
|
||||
}
|
||||
|
||||
bool open = usb_host_port_open();
|
||||
if (!open) {
|
||||
discard_control_requests();
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
|
||||
if (client_id != SESSION_BROKER_NO_CLIENT) {
|
||||
esp_err_t event_result = drain_broker_events(client_id, &writer);
|
||||
if (event_result == ESP_ERR_NOT_FOUND) {
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
} else {
|
||||
esp_err_t result = session_broker_disconnect(client_id);
|
||||
if (result == ESP_OK || result == ESP_ERR_NOT_FOUND) {
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
}
|
||||
/* Other failures leave the broker transaction intact for retry. */
|
||||
}
|
||||
}
|
||||
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
continue;
|
||||
}
|
||||
|
||||
TickType_t now = xTaskGetTickCount();
|
||||
if (!serial_service_is_running() &&
|
||||
retry_due(now,
|
||||
service_retry_ticks,
|
||||
&last_service_attempt,
|
||||
&service_attempted)) {
|
||||
esp_err_t result = serial_service_start();
|
||||
if (result != ESP_OK && !serial_service_is_running()) {
|
||||
add_counter(&s_counters.service_start_failures, 1U);
|
||||
}
|
||||
}
|
||||
|
||||
if (client_id == SESSION_BROKER_NO_CLIENT &&
|
||||
serial_service_is_running() &&
|
||||
retry_due(now,
|
||||
connect_retry_ticks,
|
||||
&last_connect_attempt,
|
||||
&connect_attempted)) {
|
||||
session_broker_client_id_t new_client_id = SESSION_BROKER_NO_CLIENT;
|
||||
esp_err_t result = session_broker_connect(SESSION_BROKER_CLIENT_USB,
|
||||
"usb-cdc",
|
||||
&new_client_id);
|
||||
if (result == ESP_OK) {
|
||||
client_id = new_client_id;
|
||||
writer = false;
|
||||
set_broker_state(client_id, false);
|
||||
add_counter(&s_counters.connections, 1U);
|
||||
|
||||
/* Initial ownership is opportunistic; denial leaves an observer. */
|
||||
(void)session_broker_request_writer(client_id);
|
||||
esp_err_t state_result = drain_broker_events(client_id, &writer);
|
||||
if (state_result == ESP_OK) {
|
||||
state_result = reconcile_broker_state(client_id, &writer);
|
||||
reconciled = state_result == ESP_OK;
|
||||
last_reconcile = now;
|
||||
}
|
||||
if (state_result == ESP_ERR_NOT_FOUND) {
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (client_id == SESSION_BROKER_NO_CLIENT) {
|
||||
discard_control_requests();
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
continue;
|
||||
}
|
||||
|
||||
esp_err_t result = drain_broker_events(client_id, &writer);
|
||||
if (result == ESP_OK) {
|
||||
result = process_control_requests(client_id, &writer);
|
||||
}
|
||||
now = xTaskGetTickCount();
|
||||
if (result == ESP_OK &&
|
||||
retry_due(now, reconcile_ticks, &last_reconcile, &reconciled)) {
|
||||
/* Event queues are bounded; the snapshot is authoritative after gaps. */
|
||||
result = reconcile_broker_state(client_id, &writer);
|
||||
}
|
||||
if (result == ESP_ERR_NOT_FOUND) {
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
continue;
|
||||
}
|
||||
|
||||
/* Re-check after broker calls that may have waited for another task. */
|
||||
if (atomic_load(&s_connection_generation) != observed_generation) {
|
||||
continue;
|
||||
}
|
||||
|
||||
apply_pending_line_coding(writer);
|
||||
|
||||
if (atomic_load(&s_connection_generation) != observed_generation) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (writer) {
|
||||
result = move_host_data_to_broker(client_id, &host_pending);
|
||||
if (result == ESP_ERR_INVALID_STATE) {
|
||||
/* Ownership may have changed after an event-queue overflow. */
|
||||
result = reconcile_broker_state(client_id, &writer);
|
||||
}
|
||||
if (result == ESP_ERR_NOT_FOUND) {
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
continue;
|
||||
}
|
||||
} else {
|
||||
/* Observers receive UART data but may not retain host-originated input. */
|
||||
discard_host_input(&host_pending);
|
||||
}
|
||||
|
||||
result = move_broker_data_to_usb(client_id, &usb_pending);
|
||||
if (result == ESP_ERR_NOT_FOUND) {
|
||||
discard_host_input(&host_pending);
|
||||
discard_usb_pending(&usb_pending);
|
||||
mark_client_disconnected(&client_id, &writer);
|
||||
}
|
||||
|
||||
(void)ulTaskNotifyTake(pdTRUE, poll_ticks);
|
||||
}
|
||||
}
|
||||
|
||||
static void reset_uninitialized_state(void)
|
||||
{
|
||||
atomic_store(&s_usb_state, 0U);
|
||||
atomic_store(&s_connection_generation, 0U);
|
||||
atomic_store(&s_transport_task, (uintptr_t)NULL);
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
s_broker_client_id = SESSION_BROKER_NO_CLIENT;
|
||||
s_writer = false;
|
||||
s_line_coding = (usb_cdc_transport_line_coding_t) {
|
||||
.baud_rate = 115200U,
|
||||
.stop_bits = USB_CDC_TRANSPORT_STOP_BITS_1,
|
||||
.parity = USB_CDC_TRANSPORT_PARITY_NONE,
|
||||
.data_bits = 8U,
|
||||
};
|
||||
s_line_coding_pending = false;
|
||||
memset(&s_counters, 0, sizeof(s_counters));
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
}
|
||||
|
||||
static void cleanup_init_allocations(bool cdc_initialized, bool driver_installed)
|
||||
{
|
||||
if (cdc_initialized) {
|
||||
(void)tinyusb_cdcacm_deinit(TINYUSB_CDC_ACM_0);
|
||||
}
|
||||
if (driver_installed) {
|
||||
(void)tinyusb_driver_uninstall();
|
||||
}
|
||||
if (s_control_queue != NULL) {
|
||||
vQueueDelete(s_control_queue);
|
||||
s_control_queue = NULL;
|
||||
}
|
||||
if (s_host_rx_stream != NULL) {
|
||||
vStreamBufferDelete(s_host_rx_stream);
|
||||
s_host_rx_stream = NULL;
|
||||
}
|
||||
reset_uninitialized_state();
|
||||
}
|
||||
|
||||
esp_err_t usb_cdc_transport_init(void)
|
||||
{
|
||||
bool expected = false;
|
||||
if (atomic_load(&s_initialized) ||
|
||||
!atomic_compare_exchange_strong(&s_initializing, &expected, true)) {
|
||||
return ESP_ERR_INVALID_STATE;
|
||||
}
|
||||
|
||||
reset_uninitialized_state();
|
||||
|
||||
uint8_t mac[6];
|
||||
esp_err_t result = esp_read_mac(mac, ESP_MAC_WIFI_STA);
|
||||
if (result != ESP_OK) {
|
||||
atomic_store(&s_initializing, false);
|
||||
return result;
|
||||
}
|
||||
(void)snprintf(s_serial_number,
|
||||
sizeof(s_serial_number),
|
||||
"%02X%02X%02X%02X%02X%02X",
|
||||
mac[0],
|
||||
mac[1],
|
||||
mac[2],
|
||||
mac[3],
|
||||
mac[4],
|
||||
mac[5]);
|
||||
|
||||
s_host_rx_stream = xStreamBufferCreate(USB_CDC_HOST_RX_STREAM_SIZE, 1U);
|
||||
if (s_host_rx_stream == NULL) {
|
||||
atomic_store(&s_initializing, false);
|
||||
return ESP_ERR_NO_MEM;
|
||||
}
|
||||
|
||||
s_control_queue = xQueueCreate(USB_CDC_CONTROL_QUEUE_LENGTH,
|
||||
sizeof(usb_cdc_control_t));
|
||||
if (s_control_queue == NULL) {
|
||||
cleanup_init_allocations(false, false);
|
||||
atomic_store(&s_initializing, false);
|
||||
return ESP_ERR_NO_MEM;
|
||||
}
|
||||
|
||||
/* ESP32-S3's default full-speed internal PHY is fixed to GPIO19/20. */
|
||||
tinyusb_config_t usb_config = TINYUSB_DEFAULT_CONFIG(device_event_callback);
|
||||
usb_config.descriptor.string = s_string_descriptors;
|
||||
usb_config.descriptor.string_count =
|
||||
(int)(sizeof(s_string_descriptors) / sizeof(s_string_descriptors[0]));
|
||||
/* NULL device/config descriptor fields deliberately select class defaults. */
|
||||
|
||||
result = tinyusb_driver_install(&usb_config);
|
||||
if (result != ESP_OK) {
|
||||
cleanup_init_allocations(false, false);
|
||||
atomic_store(&s_initializing, false);
|
||||
return result;
|
||||
}
|
||||
|
||||
const tinyusb_config_cdcacm_t cdc_config = {
|
||||
.cdc_port = TINYUSB_CDC_ACM_0,
|
||||
.callback_rx = cdc_rx_callback,
|
||||
.callback_rx_wanted_char = cdc_wanted_char_callback,
|
||||
.callback_line_state_changed = cdc_line_state_callback,
|
||||
.callback_line_coding_changed = cdc_line_coding_callback,
|
||||
};
|
||||
result = tinyusb_cdcacm_init(&cdc_config);
|
||||
if (result != ESP_OK) {
|
||||
cleanup_init_allocations(false, true);
|
||||
atomic_store(&s_initializing, false);
|
||||
return result;
|
||||
}
|
||||
|
||||
TaskHandle_t task = NULL;
|
||||
if (xTaskCreate(transport_task,
|
||||
"usb_cdc_transport",
|
||||
USB_CDC_TASK_STACK_SIZE,
|
||||
NULL,
|
||||
USB_CDC_TASK_PRIORITY,
|
||||
&task) != pdPASS) {
|
||||
cleanup_init_allocations(true, true);
|
||||
atomic_store(&s_initializing, false);
|
||||
return ESP_ERR_NO_MEM;
|
||||
}
|
||||
|
||||
atomic_store(&s_transport_task, (uintptr_t)task);
|
||||
atomic_store(&s_initialized, true);
|
||||
atomic_store(&s_initializing, false);
|
||||
notify_transport_task();
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
esp_err_t usb_cdc_transport_get_snapshot(usb_cdc_transport_snapshot_t *snapshot)
|
||||
{
|
||||
if (snapshot == NULL) {
|
||||
return ESP_ERR_INVALID_ARG;
|
||||
}
|
||||
|
||||
unsigned int usb_state = atomic_load(&s_usb_state);
|
||||
snapshot->initialized = atomic_load(&s_initialized);
|
||||
snapshot->attached = (usb_state & USB_STATE_ATTACHED) != 0U;
|
||||
snapshot->dtr = (usb_state & USB_STATE_DTR) != 0U;
|
||||
snapshot->rts = (usb_state & USB_STATE_RTS) != 0U;
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
snapshot->broker_client_id = s_broker_client_id;
|
||||
snapshot->writer = s_writer;
|
||||
snapshot->line_coding = s_line_coding;
|
||||
snapshot->counters = s_counters;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
static esp_err_t enqueue_control_request(usb_cdc_control_t control)
|
||||
{
|
||||
if (!atomic_load(&s_initialized)) {
|
||||
return ESP_ERR_INVALID_STATE;
|
||||
}
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
bool connected = s_broker_client_id != SESSION_BROKER_NO_CLIENT;
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
if (!connected) {
|
||||
add_counter(&s_counters.control_drops, 1U);
|
||||
return ESP_ERR_INVALID_STATE;
|
||||
}
|
||||
|
||||
if (xQueueSend(s_control_queue, &control, 0U) != pdTRUE) {
|
||||
add_counter(&s_counters.control_drops, 1U);
|
||||
return ESP_ERR_TIMEOUT;
|
||||
}
|
||||
|
||||
notify_transport_task();
|
||||
return ESP_OK;
|
||||
}
|
||||
|
||||
esp_err_t usb_cdc_transport_request_writer(void)
|
||||
{
|
||||
return enqueue_control_request(USB_CDC_CONTROL_REQUEST_WRITER);
|
||||
}
|
||||
|
||||
esp_err_t usb_cdc_transport_release_writer(void)
|
||||
{
|
||||
return enqueue_control_request(USB_CDC_CONTROL_RELEASE_WRITER);
|
||||
}
|
||||
|
||||
esp_err_t usb_cdc_transport_clear_counters(void)
|
||||
{
|
||||
if (!atomic_load(&s_initialized)) {
|
||||
return ESP_ERR_INVALID_STATE;
|
||||
}
|
||||
|
||||
taskENTER_CRITICAL(&s_state_lock);
|
||||
memset(&s_counters, 0, sizeof(s_counters));
|
||||
taskEXIT_CRITICAL(&s_state_lock);
|
||||
return ESP_OK;
|
||||
}
|
||||
Reference in New Issue
Block a user