#include "aes67_tx.h" #include #include #include #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 AES67_RATE 48000 // the only sample rate (all sources are 48 kHz) #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[] = { AES67_RATE }; // the sources (player, test signals) are 48 kHz 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, 1, 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) bool mono_sum; // 1-channel stream carries (L + R) / 2, else L 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->mono_sum = cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(a, "mono_sum")); 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; } // Pink noise: xorshift32 white noise through Paul Kellet's refined pink filter (within 0.05 dB // above 9 Hz). The scale gives -18 dBFS RMS (measured over 10 s; peaks about -6 dBFS). #define PINK_SCALE 0.0714f size_t aes67_tx_pink_noise(int32_t *buf, size_t frames) { static uint32_t rnd = 0x12345678; static float b0, b1, b2, b3, b4, b5, b6; for (size_t i = 0; i < frames; i++) { rnd ^= rnd << 13; rnd ^= rnd >> 17; rnd ^= rnd << 5; float w = (int32_t)rnd * (1.0f / 2147483648.0f); // -1..1 b0 = 0.99886f * b0 + w * 0.0555179f; b1 = 0.99332f * b1 + w * 0.0750759f; b2 = 0.96900f * b2 + w * 0.1538520f; b3 = 0.86650f * b3 + w * 0.3104856f; b4 = 0.55000f * b4 + w * 0.5329522f; b5 = -0.7616f * b5 - w * 0.0168980f; float pink = b0 + b1 + b2 + b3 + b4 + b5 + b6 + w * 0.5362f; b6 = w * 0.115926f; float x = pink * PINK_SCALE; x = x > 0.999f ? 0.999f : x < -0.999f ? -0.999f : x; // never reached in practice int32_t v = (int32_t)(x * 2147483647.0f); buf[i * AES67_TX_SRC_CHANNELS] = v; buf[i * AES67_TX_SRC_CHANNELS + 1] = v; } return frames; } // Stereo source -> 1-channel stream, in place. The sum is in 64 bit, so it can't clip; identical // L/R comes out at the same level. static void to_mono(int32_t *buf, int frames, bool sum) { for (int i = 0; i < frames; i++) { buf[i] = sum ? (int32_t)(((int64_t)buf[2 * i] + buf[2 * i + 1]) / 2) : buf[2 * i]; } } // 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): precomputed table. static void tone(int32_t *buf, int frames, int channels, int rate, int64_t s0) { static int32_t table[AES67_RATE / 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 > AES67_TX_SRC_CHANNELS ? MAX_CH : AES67_TX_SRC_CHANNELS)]; 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 * AES67_TX_SRC_CHANNELS, 0, (c.frames - got) * AES67_TX_SRC_CHANNELS * sizeof(int32_t)); s_underruns++; } if (c.channels == 1) { to_mono(pcm, c.frames, c.mono_sum); } } 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; } // Stored before 96 kHz was dropped: stored values aren't re-validated, so fix it here (SDP too). cJSON *a = cfg_get("aes67"); double rate = cJSON_GetObjectItemCaseSensitive(a, "rate")->valuedouble; cJSON_Delete(a); if (rate != AES67_RATE) { ESP_LOGW(TAG, "stored sample rate %.0f not supported: using %d", rate, AES67_RATE); cfg_set_number("aes67", "rate", AES67_RATE, false); } return status_register(tx_status); }