Files
aes67-ESP32-P4/components/aes67_tx/aes67_tx.c
T
bsncubed 044e9913a7 AES67 TX: send each packet right after its last sample is due
The periodic TX timer had an arbitrary phase against packet boundaries,
so packets waited up to one packet time (different after every restart).
- One-shot esp_timer aimed at the next packet's due time from PTP, with
  the clock re-read after sending.
- 1 kHz tone from a per-rate lookup table (one period = rate/1000
  samples) instead of sinf() per sample.
Measured (RTP time - arrival, 4 ms vs 1 ms packets, expected -3.07 ms
incl. wire time): -3.3 ms extra before, now -3.14 ms. Tone within
1.0 LSB of the ideal PTP-phased sine, no gaps.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-25 06:48:47 +10:00

321 lines
12 KiB
C

#include "aes67_tx.h"
#include <math.h>
#include <stdio.h>
#include <string.h>
#include "aes67_cfg.h"
#include "aes67_net.h"
#include "aes67_ptp.h"
#include "aes67_web.h"
#include "esp_log.h"
#include "esp_random.h"
#include "esp_timer.h"
#include "freertos/FreeRTOS.h"
#include "freertos/task.h"
#include "lwip/inet.h"
#include "lwip/sockets.h"
#define MAX_FRAMES 192 // 4 ms at 48 kHz
#define MAX_CH 2
#define RTP_HDR 12
#define MAX_LAG_NS 20000000 // more than 20 ms behind: resync instead of bursting
#define TONE_HZ 1000
#define TONE_DBFS -18.0
#define WAKE_MARGIN_US 20 // wake just after a packet's last sample is due
static const char *TAG = "aes67_tx";
static TaskHandle_t s_task;
static volatile aes67_tx_read_cb_t s_read;
static volatile bool s_reconfig = true;
static volatile uint32_t s_packets, s_underruns;
// Defaults follow the Riedel Director 4-wire AES67 output; channels 2 (core default).
static const char AES67_DEFAULTS[] =
"{\"name\":\"AES67\",\"enabled\":true,\"discovery\":\"sap\",\"mcast\":\"239.69.1.10\","
"\"port\":5004,\"ttl\":32,\"dscp\":34,\"channels\":2,\"mono_sum\":false,\"encoding\":\"L24\","
"\"rate\":48000,\"ptime\":1,\"pt\":96,\"ssrc\":0,\"clk_offset\":0,"
"\"session_id\":1,\"session_ver\":1}";
// RTP multicast range 224.0.2.0 - 239.255.255.255 (keeps clear of 224.0.0.x/224.0.1.x, where PTP lives).
static bool check_mcast(const cJSON *g, char *err, size_t n)
{
const cJSON *v = cJSON_GetObjectItemCaseSensitive(g, "mcast");
struct in_addr a;
if (cJSON_IsString(v) && inet_aton(v->valuestring, &a)) {
uint32_t ip = ntohl(a.s_addr);
if (ip >= 0xE0000200 && ip <= 0xEFFFFFFF) {
return true;
}
}
snprintf(err, n, "mcast: must be 224.0.2.0 - 239.255.255.255");
return false;
}
static bool aes67_validate(const cJSON *g, char *err, size_t n)
{
static const char *const discovery[] = { "manual", "sap", NULL };
static const char *const encodings[] = { "L16", "L24", NULL };
static const double rates[] = { 48000, 96000 };
static const double ptimes[] = { 0.125, 0.25, 0.333, 1, 4 };
bool ok = cfg_check_str(g, "name", 1, 63, err, n) &&
cfg_check_enum(g, "discovery", discovery, err, n) &&
check_mcast(g, err, n) &&
cfg_check_int(g, "port", 1024, 65535, err, n) &&
cfg_check_int(g, "ttl", 1, 255, err, n) &&
cfg_check_int(g, "dscp", 0, 63, err, n) &&
cfg_check_int(g, "channels", 1, 2, err, n) &&
cfg_check_enum(g, "encoding", encodings, err, n) &&
cfg_check_num_in(g, "rate", rates, 2, err, n) &&
cfg_check_num_in(g, "ptime", ptimes, 5, err, n) &&
cfg_check_int(g, "pt", 96, 127, err, n) &&
cfg_check_int(g, "ssrc", 0, 4294967295.0, err, n) &&
cfg_check_int(g, "clk_offset", 0, 4294967295.0, err, n) &&
cfg_check_int(g, "session_id", 0, 4294967295.0, err, n) &&
cfg_check_int(g, "session_ver", 0, 4294967295.0, err, n);
if (ok && cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(g, "mono_sum")) &&
cJSON_GetObjectItemCaseSensitive(g, "channels")->valuedouble != 1) {
snprintf(err, n, "channels: must be 1 with mono_sum");
ok = false;
}
return ok;
}
static void tx_status(cJSON *st)
{
cJSON_AddNumberToObject(st, "tx_packets", s_packets);
cJSON_AddNumberToObject(st, "underruns", s_underruns);
}
static void aes67_apply(const cJSON *g)
{
s_reconfig = true; // the TX task picks up the new settings
}
void aes67_tx_set_source(aes67_tx_read_cb_t read)
{
s_read = read;
}
/* ----- sender ----- */
typedef struct {
bool enabled;
struct sockaddr_in dst;
uint8_t ttl, dscp, pt;
uint32_t ssrc, clk_offset;
int channels, rate, bytes; // bytes per sample: 3 (L24) or 2 (L16)
int frames; // samples per channel per packet
} tx_cfg_t;
static bool load(tx_cfg_t *c)
{
cJSON *a = cfg_get("aes67");
if (!a) {
return false;
}
#define NUM(k) cJSON_GetObjectItemCaseSensitive(a, k)->valuedouble
c->enabled = cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(a, "enabled"));
memset(&c->dst, 0, sizeof(c->dst));
c->dst.sin_family = AF_INET;
c->dst.sin_port = htons((uint16_t)NUM("port"));
inet_aton(cJSON_GetObjectItemCaseSensitive(a, "mcast")->valuestring, &c->dst.sin_addr);
c->ttl = (uint8_t)NUM("ttl");
c->dscp = (uint8_t)NUM("dscp");
c->pt = (uint8_t)NUM("pt");
c->ssrc = (uint32_t)NUM("ssrc");
c->clk_offset = (uint32_t)NUM("clk_offset");
c->channels = (int)NUM("channels");
c->rate = (int)NUM("rate");
c->bytes = strcmp(cJSON_GetObjectItemCaseSensitive(a, "encoding")->valuestring, "L16") == 0 ? 2 : 3;
c->frames = (int)lround(NUM("ptime") * c->rate / 1000.0); // 0.333 ms -> 16 at 48 kHz
#undef NUM
cJSON_Delete(a);
return c->frames > 0 && c->frames <= MAX_FRAMES && c->channels <= MAX_CH;
}
static int open_socket(const tx_cfg_t *c)
{
int fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP);
if (fd < 0) {
return -1;
}
int tos = c->dscp << 2;
uint8_t ttl = c->ttl, loop = 0;
setsockopt(fd, IPPROTO_IP, IP_TOS, &tos, sizeof(tos));
setsockopt(fd, IPPROTO_IP, IP_MULTICAST_TTL, &ttl, sizeof(ttl));
setsockopt(fd, IPPROTO_IP, IP_MULTICAST_LOOP, &loop, sizeof(loop));
esp_netif_ip_info_t ip;
if (esp_netif_get_ip_info(aes67_net_netif(), &ip) == ESP_OK) {
struct in_addr ifa = { .s_addr = ip.ip.addr };
setsockopt(fd, IPPROTO_IP, IP_MULTICAST_IF, &ifa, sizeof(ifa));
}
return fd;
}
// Sample index since the PTP epoch (no overflow: seconds * rate first).
static int64_t ns_to_samples(int64_t ns, int rate)
{
return (ns / 1000000000LL) * rate + (ns % 1000000000LL) * rate / 1000000000LL;
}
// PTP time at which sample index s starts (rounded up).
static int64_t samples_to_ns(int64_t s, int rate)
{
return (s / rate) * 1000000000LL + ((s % rate) * 1000000000LL + rate - 1) / rate;
}
// 1 kHz tone, phase from the PTP sample index, so every sender's tone lines up.
// One period is rate / 1000 samples (48 at 48 kHz, 96 at 96 kHz): precomputed table.
static void tone(int32_t *buf, int frames, int channels, int rate, int64_t s0)
{
static int32_t table[96000 / TONE_HZ];
static int table_rate;
int period = rate / TONE_HZ;
if (table_rate != rate) {
double amp = pow(10.0, TONE_DBFS / 20.0) * 2147483647.0;
for (int i = 0; i < period; i++) {
table[i] = (int32_t)(amp * sin(2.0 * M_PI * i / period));
}
table_rate = rate;
}
int ph = (int)(s0 % period);
for (int i = 0; i < frames; i++) {
int32_t v = table[ph];
ph = ph + 1 == period ? 0 : ph + 1;
for (int ch = 0; ch < channels; ch++) {
buf[i * channels + ch] = v;
}
}
}
static void timer_cb(void *arg)
{
xTaskNotifyGive(s_task);
}
static void tx_task(void *arg)
{
tx_cfg_t c = { 0 };
int fd = -1;
esp_timer_handle_t timer = NULL;
const esp_timer_create_args_t targs = { .callback = timer_cb, .name = "aes67_tx" };
esp_timer_create(&targs, &timer);
static int32_t pcm[MAX_FRAMES * MAX_CH];
static uint8_t pkt[RTP_HDR + MAX_FRAMES * MAX_CH * 3];
uint16_t seq = (uint16_t)esp_random();
int64_t next = 0; // sample index of the next packet; 0 = resync needed
const char *state = NULL, *new_state;
while (1) {
ulTaskNotifyTake(pdTRUE, pdMS_TO_TICKS(100));
if (s_reconfig) {
s_reconfig = false;
if (fd >= 0) {
close(fd);
fd = -1;
}
esp_timer_stop(timer);
if (!load(&c)) {
ESP_LOGE(TAG, "unsupported config (ptime/rate/channels)");
c.enabled = false;
}
next = 0;
if (c.enabled) {
fd = open_socket(&c);
}
}
int64_t now;
if (!c.enabled || fd < 0) {
new_state = "disabled";
} else if (!aes67_ptp_locked() || aes67_ptp_now_ns(&now) != ESP_OK) {
new_state = "waiting for PTP lock";
next = 0;
} else {
new_state = "sending";
int64_t due = ns_to_samples(now, c.rate);
// (Re)start aligned to a packet boundary, or after falling behind / a clock step back.
if (!next || due - next > (int64_t)c.rate * MAX_LAG_NS / 1000000000LL || next - due > 2 * c.frames) {
if (next) {
ESP_LOGW(TAG, "resync (%lld samples off)", due - next);
}
next = (due / c.frames + 1) * c.frames;
}
// Send every packet whose last sample is in the past.
while (samples_to_ns(next + c.frames, c.rate) <= now) {
aes67_tx_read_cb_t read = s_read;
if (read) {
size_t got = read(pcm, c.frames);
if (got < (size_t)c.frames) {
memset(pcm + got * c.channels, 0, (c.frames - got) * c.channels * sizeof(int32_t));
s_underruns++;
}
} else {
tone(pcm, c.frames, c.channels, c.rate, next);
}
uint32_t ts = (uint32_t)next + c.clk_offset;
pkt[0] = 0x80; // V=2
pkt[1] = c.pt & 0x7f;
pkt[2] = seq >> 8;
pkt[3] = seq & 0xff;
pkt[4] = ts >> 24;
pkt[5] = ts >> 16;
pkt[6] = ts >> 8;
pkt[7] = ts;
pkt[8] = c.ssrc >> 24;
pkt[9] = c.ssrc >> 16;
pkt[10] = c.ssrc >> 8;
pkt[11] = c.ssrc;
uint8_t *p = pkt + RTP_HDR;
for (int i = 0; i < c.frames * c.channels; i++) {
uint32_t v = (uint32_t)pcm[i]; // big-endian, top bytes of the 32-bit sample
*p++ = v >> 24;
*p++ = v >> 16;
if (c.bytes == 3) {
*p++ = v >> 8;
}
}
if (sendto(fd, pkt, p - pkt, 0, (struct sockaddr *)&c.dst, sizeof(c.dst)) > 0) {
s_packets++;
}
seq++;
next += c.frames;
}
// Wake when the next packet's last sample is due (PTP time), not on a free-running
// period, so each packet leaves right after it is complete. Re-read the clock:
// building and sending took time.
aes67_ptp_now_ns(&now);
int64_t wait_us = (samples_to_ns(next + c.frames, c.rate) - now) / 1000 + WAKE_MARGIN_US;
esp_timer_stop(timer);
esp_timer_start_once(timer, wait_us < 50 ? 50 : wait_us > 10000 ? 10000 : wait_us);
}
if (new_state != state) {
state = new_state;
if (c.enabled && state[0] == 's') {
ESP_LOGI(TAG, "%s: %s:%u, %s/%d/%d, %d samples per packet, %s", state,
inet_ntoa(c.dst.sin_addr), ntohs(c.dst.sin_port), c.bytes == 3 ? "L24" : "L16",
c.rate, c.channels, c.frames, s_read ? "source" : "1 kHz test tone");
} else {
ESP_LOGI(TAG, "%s", state);
}
}
}
}
esp_err_t aes67_tx_start(void)
{
return xTaskCreate(tx_task, "aes67_tx", 4096, NULL, 16, &s_task) == pdPASS ? ESP_OK : ESP_ERR_NO_MEM;
}
esp_err_t aes67_tx_init(void)
{
esp_err_t err = cfg_register("aes67", AES67_DEFAULTS, aes67_validate, aes67_apply);
if (err != ESP_OK) {
return err;
}
return status_register(tx_status);
}