#include "aes67_syslog.h" #include #include #include #include "aes67_cfg.h" #include "aes67_net.h" #include "aes67_web.h" #include "esp_log.h" #include "freertos/FreeRTOS.h" #include "freertos/queue.h" #include "freertos/semphr.h" #include "freertos/task.h" #include "lwip/netdb.h" #include "lwip/sockets.h" #define QUEUE_LEN 48 #define MSG_MAX 240 // facility: 16 = local0 (PRI = facility * 8 + severity). static const char LOG_DEFAULTS[] = "{\"syslog\":false,\"host\":\"\",\"port\":514,\"level\":\"info\",\"facility\":16,\"format\":\"rfc5424\"}"; typedef struct { uint8_t severity; // 3 error, 4 warning, 6 info, 7 debug char text[MSG_MAX]; // "tag: message" without the level/timestamp prefix } msg_t; static QueueHandle_t s_queue; static SemaphoreHandle_t s_fmt_lock; static TaskHandle_t s_task; static vprintf_like_t s_uart; static volatile bool s_enabled; static volatile uint8_t s_min_severity = 6; static volatile bool s_reconfig = true; static volatile uint32_t s_dropped; static bool log_validate(const cJSON *g, char *err, size_t n) { static const char *const levels[] = { "error", "warn", "info", "debug", NULL }; static const char *const formats[] = { "rfc5424", "rfc3164", NULL }; static const double facilities[] = { 1, 3, 16, 17, 18, 19, 20, 21, 22, 23 }; bool on = cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(g, "syslog")); return cfg_check_str(g, "host", on ? 1 : 0, 63, err, n) && cfg_check_int(g, "port", 1, 65535, err, n) && cfg_check_enum(g, "level", levels, err, n) && cfg_check_num_in(g, "facility", facilities, 10, err, n) && cfg_check_enum(g, "format", formats, err, n); } static uint8_t severity_of(const char *level) { return level[0] == 'e' ? 3 : level[0] == 'w' ? 4 : level[0] == 'i' ? 6 : 7; } static void read_filter(void) { cJSON *l = cfg_get("log"); s_enabled = cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(l, "syslog")); s_min_severity = severity_of(cJSON_GetObjectItemCaseSensitive(l, "level")->valuestring); cJSON_Delete(l); } static void log_apply(const cJSON *g) { read_filter(); s_reconfig = true; if (s_task) { xTaskNotifyGive(s_task); } } /* ----- esp_log hook ----- */ // "\033[0;32mI (1234) tag: message\033[0m\n" -> severity + "tag: message" static bool parse_line(const char *in, msg_t *m) { char *o = m->text, *end = m->text + MSG_MAX - 1; const char *p = in; char level = 0; bool prefix_done = false; while (*p && o < end) { if (*p == '\033') { // skip ANSI colour codes while (*p && *p != 'm') { p++; } if (*p) { p++; } continue; } if (!level) { level = *p++; continue; } if (!prefix_done) { // skip " (1234) " if (*p == ')') { prefix_done = true; p++; if (*p == ' ') { p++; } } else { p++; } continue; } if (*p == '\n' || *p == '\r') { p++; continue; } *o++ = *p++; } *o = 0; m->severity = level == 'E' ? 3 : level == 'W' ? 4 : level == 'I' ? 6 : 7; return prefix_done && m->text[0]; } static int log_hook(const char *fmt, va_list ap) { va_list copy; va_copy(copy, ap); int ret = s_uart(fmt, ap); // UART output unchanged if (s_enabled && xTaskGetCurrentTaskHandle() != s_task && xSemaphoreTake(s_fmt_lock, 0) == pdTRUE) { static char line[MSG_MAX + 32]; static msg_t m; vsnprintf(line, sizeof(line), fmt, copy); if (parse_line(line, &m) && m.severity <= s_min_severity && xQueueSend(s_queue, &m, 0) != pdTRUE) { s_dropped++; } xSemaphoreGive(s_fmt_lock); } else if (s_enabled && xTaskGetCurrentTaskHandle() != s_task) { s_dropped++; // formatter busy: drop rather than block } va_end(copy); return ret; } /* ----- sender ----- */ typedef struct { struct sockaddr_in dst; bool have_dst; uint8_t facility; bool rfc5424; char host[64]; char hostname[64]; } sender_cfg_t; static void load_sender(sender_cfg_t *c) { cJSON *l = cfg_get("log"), *n = cfg_get("net"); strlcpy(c->host, cJSON_GetObjectItemCaseSensitive(l, "host")->valuestring, sizeof(c->host)); c->facility = (uint8_t)cJSON_GetObjectItemCaseSensitive(l, "facility")->valuedouble; c->rfc5424 = strcmp(cJSON_GetObjectItemCaseSensitive(l, "format")->valuestring, "rfc5424") == 0; uint16_t port = (uint16_t)cJSON_GetObjectItemCaseSensitive(l, "port")->valuedouble; strlcpy(c->hostname, cJSON_GetObjectItemCaseSensitive(n, "hostname")->valuestring, sizeof(c->hostname)); cJSON_Delete(l); cJSON_Delete(n); c->have_dst = false; memset(&c->dst, 0, sizeof(c->dst)); c->dst.sin_family = AF_INET; c->dst.sin_port = htons(port); if (inet_aton(c->host, &c->dst.sin_addr)) { c->have_dst = true; return; } struct addrinfo hints = { .ai_family = AF_INET, .ai_socktype = SOCK_DGRAM }, *res = NULL; if (c->host[0] && getaddrinfo(c->host, NULL, &hints, &res) == 0 && res) { c->dst.sin_addr = ((struct sockaddr_in *)res->ai_addr)->sin_addr; c->have_dst = true; } if (res) { freeaddrinfo(res); } } // RFC 5424: 1 - HOSTNAME APP-NAME - - - MSG (no wall-clock time yet: NILVALUE) // RFC 3164: HOSTNAME TAG: MSG (TIMESTAMP omitted; the server adds it) static int format_msg(char *out, size_t size, const sender_cfg_t *c, const msg_t *m) { const char *colon = strstr(m->text, ": "); char tag[33] = "-"; const char *msg = m->text; if (colon && colon - m->text < (int)sizeof(tag)) { memcpy(tag, m->text, colon - m->text); tag[colon - m->text] = 0; msg = colon + 2; } int pri = c->facility * 8 + m->severity; if (c->rfc5424) { return snprintf(out, size, "<%d>1 - %s %s - - - %s", pri, c->hostname, tag, msg); } return snprintf(out, size, "<%d>%s %s: %s", pri, c->hostname, tag, msg); } static void syslog_task(void *arg) { sender_cfg_t c = { 0 }; int fd = -1; // created once the interface is up: lwIP may not be running yet at init msg_t m; char out[MSG_MAX + 128]; while (1) { // Hold messages until the interface has an IP (early boot logs stay queued). esp_netif_ip_info_t ip; esp_netif_t *netif = aes67_net_netif(); if (!netif || esp_netif_get_ip_info(netif, &ip) != ESP_OK || !ip.ip.addr) { ulTaskNotifyTake(pdTRUE, pdMS_TO_TICKS(500)); continue; } if (fd < 0 && (fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP)) < 0) { vTaskDelay(pdMS_TO_TICKS(1000)); continue; } if (s_reconfig) { s_reconfig = false; load_sender(&c); if (s_enabled && !c.have_dst) { ESP_LOGW("syslog", "cannot resolve '%s'", c.host); // not sent: logged from this task } } if (xQueueReceive(s_queue, &m, pdMS_TO_TICKS(1000)) != pdTRUE) { continue; } if (s_enabled && c.have_dst) { int n = format_msg(out, sizeof(out), &c, &m); sendto(fd, out, n < (int)sizeof(out) ? n : (int)sizeof(out) - 1, 0, (struct sockaddr *)&c.dst, sizeof(c.dst)); } static uint32_t reported; if (s_dropped != reported) { reported = s_dropped; msg_t d = { .severity = 4 }; snprintf(d.text, sizeof(d.text), "syslog: %lu messages dropped so far (queue full)", (unsigned long)reported); xQueueSend(s_queue, &d, 0); } } } static esp_err_t test_post(httpd_req_t *req) { if (!s_enabled) { return httpd_resp_send_err(req, HTTPD_400_BAD_REQUEST, "syslog is disabled"); } msg_t m = { .severity = 6 }; // info, regardless of the level filter snprintf(m.text, sizeof(m.text), "syslog: test message from the web UI"); if (xQueueSend(s_queue, &m, 0) != pdTRUE) { return httpd_resp_send_err(req, HTTPD_500_INTERNAL_SERVER_ERROR, "syslog queue full"); } return httpd_resp_sendstr(req, ""); } esp_err_t aes67_syslog_init(void) { esp_err_t err = cfg_register("log", LOG_DEFAULTS, log_validate, log_apply); if (err != ESP_OK) { return err; } s_queue = xQueueCreate(QUEUE_LEN, sizeof(msg_t)); s_fmt_lock = xSemaphoreCreateMutex(); if (!s_queue || !s_fmt_lock) { return ESP_ERR_NO_MEM; } read_filter(); if (xTaskCreate(syslog_task, "syslog", 4096, NULL, 2, &s_task) != pdPASS) { return ESP_ERR_NO_MEM; } static const httpd_uri_t uri = { .uri = "/api/log/test", .method = HTTP_POST, .handler = test_post }; err = web_register_uri(&uri); s_uart = esp_log_set_vprintf(log_hook); return err; }