#include "ptp_clock.h" #include #include #include #include "aes67_cfg.h" #include "esp_log.h" #include "esp_random.h" #include "esp_timer.h" #include "freertos/FreeRTOS.h" #include "freertos/semphr.h" #include "freertos/task.h" #include "lwip/sockets.h" #include "ptp_hw.h" #define PTP_MCAST "224.0.1.129" #define PTP_EVENT_PORT 319 #define PTP_GENERAL_PORT 320 #define HDR_LEN 34 #define FLAG_TWO_STEP 0x0200 #define STEP_NS 1000000 // re-step instead of slewing above 1 ms #define LOCK_NS 1000 // |offset| below this counts as good #define LOCK_GOOD 8 // consecutive good Syncs to lock #define LOCK_BAD 3 // consecutive bad Syncs to unlock #define MAX_DRIFT_PPB 500000.0 #define WINDOW 64 // samples for interval/delay statistics #define SUMMARY_US (60 * 1000000LL) static const char *TAG = "ptp"; static const uint8_t PTP_MCAST_MAC[6] = { 0x01, 0x00, 0x5e, 0x00, 0x01, 0x81 }; // Rolling window of int64 samples. typedef struct { int64_t v[WINDOW]; int n, next; } window_t; typedef struct { bool valid; uint8_t port_id[10]; // sourcePortIdentity of the GM (clockId + port) uint32_t ip; // source address, for hybrid mode later uint8_t p1, cls, acc, p2; uint16_t var; uint8_t gm_id[8]; uint16_t steps; uint8_t flags; // flagField octet 1: timeTraceable 0x10, frequencyTraceable 0x20 int8_t log_announce; int64_t last_us; // last Announce (esp_timer) } master_t; static struct { esp_netif_t *netif; int ev, gen; uint8_t domain, dscp, timeout; uint8_t port_id[10]; // our clockId (EUI-64 from MAC) + port 1 master_t gm; // Sync / Follow_Up uint16_t sync_seq; bool sync_pending; int64_t t1, t2, sync_corr; // Delay_Req / Delay_Resp uint16_t dreq_seq; bool dreq_pending; int64_t t3; int8_t log_dreq; // from Delay_Resp logMessageInterval int64_t next_dreq_us; int64_t delay_ns; // mean path delay, 0 = not measured yet // Drift of (t2 - t1) between Syncs: corrects the delay for the time between t2 and t3 // while the local clock is not yet syntonised. int64_t raw, prev_raw, prev_t2; double rate; int8_t log_sync; // GM's Sync interval (from Sync logMessageInterval) // Servo (PI, linuxptp-style gains) bool stepped; // clock set to GM time since this GM was selected double drift_ppb; // integral term double freq_ppb; // correction currently applied (positive = faster) int64_t offset_ns; int good, bad; bool locked; // Statistics for status.ptp (guarded by lock) SemaphoreHandle_t lock; window_t sync_iv, announce_iv, delays; int64_t last_announce_us; uint32_t delay_req, delay_resp; uint8_t own_class; int64_t sum_max_ns, sum_start_us; } s; static void win_add(window_t *w, int64_t v) { w->v[w->next] = v; w->next = (w->next + 1) % WINDOW; if (w->n < WINDOW) { w->n++; } } static void win_clear(window_t *w) { w->n = w->next = 0; } // mean, min, max and standard deviation of a window static void win_stats(const window_t *w, double *mean, double *min, double *max, double *sd) { double sum = 0, sq = 0, lo = 0, hi = 0; for (int i = 0; i < w->n; i++) { double x = (double)w->v[i]; sum += x; sq += x * x; lo = i ? fmin(lo, x) : x; hi = i ? fmax(hi, x) : x; } double m = w->n ? sum / w->n : 0; *mean = m; if (min) *min = lo; if (max) *max = hi; if (sd) *sd = w->n > 1 ? sqrt(fmax(0, sq / w->n - m * m)) : 0; } #define LOCKED(stmt) do { xSemaphoreTake(s.lock, portMAX_DELAY); stmt; xSemaphoreGive(s.lock); } while (0) static void set_locked(bool locked); /* ----- helpers ----- */ static uint16_t rd16(const uint8_t *p) { return (p[0] << 8) | p[1]; } static int64_t rd_ts(const uint8_t *p) // 48-bit seconds + 32-bit ns { uint64_t sec = ((uint64_t)rd16(p) << 32) | ((uint32_t)p[2] << 24) | (p[3] << 16) | (p[4] << 8) | p[5]; uint32_t ns = ((uint32_t)p[6] << 24) | (p[7] << 16) | (p[8] << 8) | p[9]; return (int64_t)sec * 1000000000LL + ns; } static int64_t rd_corr_ns(const uint8_t *p) // correctionField: ns * 2^16 { int64_t v = 0; for (int i = 0; i < 8; i++) { v = (v << 8) | p[i]; } return v >> 16; } static int64_t mac_ns(const eth_mac_time_t *t) { return (int64_t)t->seconds * 1000000000LL + t->nanoseconds; } static void fmt_id(char *out, const uint8_t *id) { sprintf(out, "%02X-%02X-%02X-%02X-%02X-%02X-%02X-%02X", id[0], id[1], id[2], id[3], id[4], id[5], id[6], id[7]); } // IEEE 1588 dataset comparison (without the topology part): <0 if a is better. static int compare(const master_t *a, const master_t *b) { if (a->p1 != b->p1) return a->p1 - b->p1; if (a->cls != b->cls) return a->cls - b->cls; if (a->acc != b->acc) return a->acc - b->acc; if (a->var != b->var) return a->var - b->var; if (a->p2 != b->p2) return a->p2 - b->p2; int c = memcmp(a->gm_id, b->gm_id, 8); if (c) return c; return a->steps - b->steps; } static int open_socket(uint16_t port, struct in_addr ifaddr) { int fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); int one = 1; setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &one, sizeof(one)); struct sockaddr_in a = { .sin_family = AF_INET, .sin_port = htons(port), .sin_addr.s_addr = htonl(INADDR_ANY) }; struct ip_mreq m = { .imr_interface = ifaddr }; inet_aton(PTP_MCAST, &m.imr_multiaddr); if (bind(fd, (struct sockaddr *)&a, sizeof(a)) < 0 || setsockopt(fd, IPPROTO_IP, IP_ADD_MEMBERSHIP, &m, sizeof(m)) < 0) { ESP_LOGE(TAG, "socket %u: bind/join failed (errno %d)", port, errno); close(fd); return -1; } return fd; } static bool from_gm(const uint8_t *b) { return s.gm.valid && memcmp(b + 20, s.gm.port_id, 10) == 0; } /* ----- message handling ----- */ static void on_announce(const uint8_t *b, int len, uint32_t src_ip) { if (len < 64) { return; } master_t m = { .valid = true, .ip = src_ip, .p1 = b[47], .cls = b[48], .acc = b[49], .var = rd16(b + 50), .p2 = b[52], .steps = rd16(b + 61), .flags = b[7], .log_announce = (int8_t)b[33], .last_us = esp_timer_get_time(), }; memcpy(m.port_id, b + 20, 10); memcpy(m.gm_id, b + 53, 8); if (memcmp(m.port_id, s.port_id, 10) == 0) { return; // our own } if (from_gm(b)) { LOCKED({ if (s.last_announce_us) { win_add(&s.announce_iv, m.last_us - s.last_announce_us); } s.last_announce_us = m.last_us; s.gm = m; // refresh dataset and timeout }); } else if (!s.gm.valid || compare(&m, &s.gm) < 0) { char id[24]; fmt_id(id, m.gm_id); ESP_LOGI(TAG, "TimeTransmitter %s (p1 %u class %u p2 %u, %u hops) from " IPSTR, id, m.p1, m.cls, m.p2, m.steps, IP2STR((esp_ip4_addr_t *)&src_ip)); LOCKED({ s.gm = m; s.last_announce_us = m.last_us; win_clear(&s.announce_iv); win_clear(&s.sync_iv); win_clear(&s.delays); }); s.sync_pending = s.dreq_pending = false; LOCKED({ s.delay_ns = s.prev_t2 = 0; s.stepped = false; set_locked(false); }); s.log_dreq = 0; s.next_dreq_us = 0; } } static void set_locked(bool locked) { if (locked != s.locked) { s.locked = locked; if (locked) { win_clear(&s.delays); // delay statistics cover the locked period, not the pull-in ESP_LOGI(TAG, "locked: offset %+lld ns, frequency %+.3f ppm, path delay %lld ns", s.offset_ns, s.freq_ppb / 1000, s.delay_ns); } else { ESP_LOGW(TAG, "unlocked"); } } if (!locked) { s.good = 0; } } // offset = local - GM (ns) static void servo(int64_t offset) { s.offset_ns = offset; if (!s.stepped || llabs(offset) > STEP_NS) { if (!s.stepped) { // Seed the integral with the measured rate error (measured with the current correction applied). s.drift_ppb += s.rate * 1e9; } s.freq_ppb = -s.drift_ppb; ptp_hw_adj_freq(s.freq_ppb); esp_err_t err = ptp_hw_step(offset); ESP_LOGI(TAG, "clock stepped by %+lld ns (%s), frequency %+.3f ppm", -offset, esp_err_to_name(err), s.freq_ppb / 1000); s.stepped = true; s.offset_ns = 0; // the pre-step offset is history: not for status or the summary s.sum_max_ns = 0; s.sum_start_us = 0; s.prev_t2 = 0; // rate across the step is meaningless s.bad = 0; set_locked(false); return; } // linuxptp PI gains for hardware timestamps, scaled by the Sync interval double iv = ldexp(1.0, s.log_sync); double kp = fmin(0.7 * pow(iv, -0.3), 0.7 / iv); double ki = fmin(0.3 * pow(iv, 0.4), 0.3 / iv); double ppb = kp * offset + s.drift_ppb; s.drift_ppb = fmax(-MAX_DRIFT_PPB, fmin(MAX_DRIFT_PPB, s.drift_ppb + ki * offset)); s.freq_ppb = -ppb; ptp_hw_adj_freq(s.freq_ppb); if (llabs(offset) < LOCK_NS) { s.bad = 0; if (++s.good >= LOCK_GOOD) { set_locked(true); } } else { s.good = 0; if (s.locked && ++s.bad >= LOCK_BAD) { set_locked(false); } } } static void sync_complete(void) { s.sync_pending = false; s.raw = s.t2 - s.t1 - s.sync_corr; if (s.prev_t2 && s.t2 > s.prev_t2) { s.rate = (double)(s.raw - s.prev_raw) / (double)(s.t2 - s.prev_t2); LOCKED(win_add(&s.sync_iv, s.t2 - s.prev_t2)); } s.prev_raw = s.raw; s.prev_t2 = s.t2; if (!s.delay_ns) { return; } // offset = t2 - t1 - corrections - mean path delay LOCKED(servo(s.raw - s.delay_ns)); ESP_LOGD(TAG, "seq %u: offset %+lld ns, freq %+.3f ppm, path delay %lld ns%s", s.sync_seq, s.offset_ns, s.freq_ppb / 1000, s.delay_ns, s.locked ? ", locked" : ""); // Info-level summary once a minute: worst offset in the period. int64_t now = esp_timer_get_time(); s.sum_max_ns = llabs(s.offset_ns) > s.sum_max_ns ? llabs(s.offset_ns) : s.sum_max_ns; if (!s.sum_start_us) { s.sum_start_us = now; } else if (now - s.sum_start_us >= SUMMARY_US) { ESP_LOGI(TAG, "%s: max |offset| %lld ns, freq %+.3f ppm, path delay %lld ns (last 60 s)", s.locked ? "locked" : "unlocked", s.sum_max_ns, s.freq_ppb / 1000, s.delay_ns); s.sum_max_ns = 0; s.sum_start_us = now; } } static void on_sync(const uint8_t *b, int len) { if (len < 44 || !from_gm(b)) { return; } eth_mac_time_t t2; uint16_t seq = rd16(b + 30); if (!ptp_hw_rx_ts(PTP_MSG_SYNC, seq, b + 20, &t2)) { ESP_LOGW(TAG, "Sync %u: no HW RX timestamp", seq); return; } s.sync_seq = seq; s.log_sync = (int8_t)b[33]; s.t2 = mac_ns(&t2); s.sync_corr = rd_corr_ns(b + 8); if (rd16(b + 6) & FLAG_TWO_STEP) { s.sync_pending = true; // wait for Follow_Up } else { s.t1 = rd_ts(b + 34); sync_complete(); } } static void on_follow_up(const uint8_t *b, int len) { if (len < 44 || !from_gm(b) || !s.sync_pending || rd16(b + 30) != s.sync_seq) { return; } s.t1 = rd_ts(b + 34); s.sync_corr += rd_corr_ns(b + 8); sync_complete(); } static void send_delay_req(void) { uint8_t m[44] = { 0 }; m[0] = PTP_MSG_DELAY_REQ; m[1] = 2; m[3] = sizeof(m); m[4] = s.domain; memcpy(m + 20, s.port_id, 10); s.dreq_seq++; m[30] = s.dreq_seq >> 8; m[31] = s.dreq_seq & 0xff; m[32] = 1; // controlField: Delay_Req m[33] = 0x7f; uint32_t dst; inet_aton(PTP_MCAST, (struct in_addr *)&dst); eth_mac_time_t t3; esp_err_t err = ptp_hw_send_event(PTP_MCAST_MAC, dst, s.dscp, m, sizeof(m), &t3); if (err != ESP_OK) { ESP_LOGW(TAG, "Delay_Req %u: %s", s.dreq_seq, esp_err_to_name(err)); s.dreq_pending = false; return; } s.t3 = mac_ns(&t3); s.dreq_pending = true; LOCKED(s.delay_req++); } static void on_delay_resp(const uint8_t *b, int len) { if (len < 54 || !from_gm(b) || !s.dreq_pending || rd16(b + 30) != s.dreq_seq || memcmp(b + 44, s.port_id, 10) != 0) { return; } s.dreq_pending = false; LOCKED(s.delay_resp++); int64_t t4 = rd_ts(b + 34) - rd_corr_ns(b + 8); s.log_dreq = (int8_t)b[33]; if (!s.prev_t2) { return; } // mean path delay = ((t2 - t1 - corr) + (t4 - t3)) / 2, plus the offset drift between t2 and t3 int64_t d = (s.raw + (t4 - s.t3) + (int64_t)(s.rate * (double)(s.t3 - s.t2))) / 2; LOCKED({ s.delay_ns = s.delay_ns ? (s.delay_ns * 7 + d) / 8 : d; // light smoothing for the servo win_add(&s.delays, d); // raw samples for status }); } /* ----- task ----- */ static void load_config(void) { cJSON *c = cfg_get("ptp"); s.domain = cJSON_GetObjectItem(c, "domain")->valueint; s.dscp = cJSON_GetObjectItem(c, "dscp")->valueint; s.timeout = cJSON_GetObjectItem(c, "announce_timeout")->valueint; // clockClass per role: slave-only 255, auto/master 248 (TimeTransmitter itself: step 6) s.own_class = strcmp(cJSON_GetObjectItem(c, "role")->valuestring, "slave") == 0 ? 255 : 248; cJSON_Delete(c); } static void ptp_task(void *arg) { esp_netif_ip_info_t ip = { 0 }; while (esp_netif_get_ip_info(s.netif, &ip) != ESP_OK || !ip.ip.addr) { vTaskDelay(pdMS_TO_TICKS(500)); } struct in_addr ifaddr = { .s_addr = ip.ip.addr }; s.ev = open_socket(PTP_EVENT_PORT, ifaddr); s.gen = open_socket(PTP_GENERAL_PORT, ifaddr); if (s.ev < 0 || s.gen < 0) { vTaskDelete(NULL); } char id[24]; fmt_id(id, s.port_id); ESP_LOGI(TAG, "TimeReceiver on " IPSTR ", domain %u, clock %s, listening", IP2STR(&ip.ip), s.domain, id); uint8_t b[128]; while (1) { fd_set fds; FD_ZERO(&fds); FD_SET(s.ev, &fds); FD_SET(s.gen, &fds); struct timeval tv = { .tv_sec = 0, .tv_usec = 100000 }; if (select((s.ev > s.gen ? s.ev : s.gen) + 1, &fds, NULL, NULL, &tv) > 0) { for (int k = 0; k < 2; k++) { int fd = k ? s.gen : s.ev; if (!FD_ISSET(fd, &fds)) { continue; } struct sockaddr_in src; socklen_t sl = sizeof(src); int len = recvfrom(fd, b, sizeof(b), 0, (struct sockaddr *)&src, &sl); if (len < HDR_LEN || (b[1] & 0x0f) != 2 || b[4] != s.domain) { continue; } switch (b[0] & 0x0f) { case PTP_MSG_ANNOUNCE: on_announce(b, len, src.sin_addr.s_addr); break; case PTP_MSG_SYNC: on_sync(b, len); break; case PTP_MSG_FOLLOW_UP: on_follow_up(b, len); break; case PTP_MSG_DELAY_RESP: on_delay_resp(b, len); break; default: break; } } } int64_t now = esp_timer_get_time(); if (s.gm.valid) { // announceReceiptTimeout x the GM's announce interval int64_t window = (int64_t)s.timeout * (s.gm.log_announce >= 0 ? 1000000LL << s.gm.log_announce : 1000000LL >> -s.gm.log_announce); if (now - s.gm.last_us > window) { ESP_LOGW(TAG, "TimeTransmitter lost (no Announce for %lld ms), listening", window / 1000); LOCKED({ memset(&s.gm, 0, sizeof(s.gm)); s.delay_ns = s.prev_t2 = 0; s.stepped = false; set_locked(false); // frequency correction stays (holdover) }); } else if (now >= s.next_dreq_us && s.prev_t2) { send_delay_req(); // Delay_Req interval from the GM's Delay_Resp; randomised 0.5..1.5x int64_t iv = s.log_dreq >= 0 ? 1000000LL << s.log_dreq : 1000000LL >> -s.log_dreq; s.next_dreq_us = now + iv / 2 + (esp_random() % (uint32_t)iv); } } } } bool ptp_clock_gm_id(char out[24]) { xSemaphoreTake(s.lock, portMAX_DELAY); bool valid = s.gm.valid; if (valid) { fmt_id(out, s.gm.gm_id); } xSemaphoreGive(s.lock); return valid; } bool ptp_clock_locked(void) { return s.locked; } void ptp_clock_status(cJSON *st) { cJSON *p = cJSON_AddObjectToObject(st, "ptp"); char id[24]; double mean, min, max, sd; xSemaphoreTake(s.lock, portMAX_DELAY); const char *state = !s.gm.valid ? "LISTENING" : s.locked ? "SLAVE" : "UNCALIBRATED"; cJSON_AddStringToObject(p, "state", state); cJSON_AddBoolToObject(p, "locked", s.locked); cJSON_AddNumberToObject(p, "version", 2); cJSON_AddNumberToObject(p, "own_class", s.own_class); fmt_id(id, s.port_id); cJSON_AddStringToObject(p, "clock_id", id); cJSON_AddBoolToObject(p, "hw_ts", true); cJSON_AddNumberToObject(p, "window", WINDOW); cJSON_AddNumberToObject(p, "delay_req", s.delay_req); cJSON_AddNumberToObject(p, "delay_resp", s.delay_resp); if (s.gm.valid) { fmt_id(id, s.gm.gm_id); cJSON_AddStringToObject(p, "gm_id", id); cJSON_AddNumberToObject(p, "gm_class", s.gm.cls); cJSON_AddNumberToObject(p, "gm_accuracy", s.gm.acc); cJSON_AddNumberToObject(p, "gm_p1", s.gm.p1); cJSON_AddNumberToObject(p, "gm_p2", s.gm.p2); cJSON_AddNumberToObject(p, "steps_removed", s.gm.steps); cJSON_AddBoolToObject(p, "gm_time_traceable", s.gm.flags & 0x10); cJSON_AddBoolToObject(p, "gm_freq_traceable", s.gm.flags & 0x20); if (s.stepped) { cJSON_AddNumberToObject(p, "offset_ns", s.offset_ns); cJSON_AddNumberToObject(p, "freq_ppb", round(s.freq_ppb)); } if (s.delays.n) { win_stats(&s.delays, &mean, NULL, NULL, &sd); cJSON_AddNumberToObject(p, "path_delay_ns", round(mean)); cJSON_AddNumberToObject(p, "path_delay_sd_ns", round(sd)); } if (s.sync_iv.n) { win_stats(&s.sync_iv, &mean, &min, &max, &sd); cJSON_AddNumberToObject(p, "sync_avg_ms", mean / 1e6); cJSON_AddNumberToObject(p, "sync_min_ms", min / 1e6); cJSON_AddNumberToObject(p, "sync_max_ms", max / 1e6); cJSON_AddNumberToObject(p, "sync_jitter_us", sd / 1e3); } if (s.announce_iv.n) { win_stats(&s.announce_iv, &mean, NULL, NULL, NULL); cJSON_AddNumberToObject(p, "announce_avg_ms", mean / 1e3); } } xSemaphoreGive(s.lock); } esp_err_t ptp_clock_start(esp_netif_t *netif) { s.netif = netif; s.lock = xSemaphoreCreateMutex(); load_config(); uint8_t mac[6]; esp_netif_get_mac(netif, mac); const uint8_t pid[10] = { mac[0], mac[1], mac[2], 0xff, 0xfe, mac[3], mac[4], mac[5], 0, 1 }; memcpy(s.port_id, pid, sizeof(pid)); return xTaskCreate(ptp_task, "ptp", 4096, NULL, 10, NULL) == pdPASS ? ESP_OK : ESP_ERR_NO_MEM; }