diff --git a/components/aes67_syslog/CMakeLists.txt b/components/aes67_syslog/CMakeLists.txt index a4b1dba..385541f 100644 --- a/components/aes67_syslog/CMakeLists.txt +++ b/components/aes67_syslog/CMakeLists.txt @@ -1,3 +1,3 @@ idf_component_register(SRCS "aes67_syslog.c" INCLUDE_DIRS "include" - PRIV_REQUIRES aes67_web) + PRIV_REQUIRES aes67_net aes67_web esp_netif log lwip) diff --git a/components/aes67_syslog/aes67_syslog.c b/components/aes67_syslog/aes67_syslog.c index 0d1d23a..41361ab 100644 --- a/components/aes67_syslog/aes67_syslog.c +++ b/components/aes67_syslog/aes67_syslog.c @@ -1,13 +1,41 @@ #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 }; @@ -22,7 +50,231 @@ static bool log_validate(const cJSON *g, char *err, size_t 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) { - return cfg_register("log", LOG_DEFAULTS, log_validate, NULL); + 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; } diff --git a/components/aes67_syslog/include/aes67_syslog.h b/components/aes67_syslog/include/aes67_syslog.h index 12d3731..fe7b77b 100644 --- a/components/aes67_syslog/include/aes67_syslog.h +++ b/components/aes67_syslog/include/aes67_syslog.h @@ -4,5 +4,6 @@ #include "esp_err.h" -// Registers the "log" config group. (Sender itself: step 5.) +// Registers the "log" config group, hooks esp_log (still printed on the UART) and starts the +// sender task. Messages logged before the interface has an IP are queued (up to the queue size). esp_err_t aes67_syslog_init(void); diff --git a/main/main.c b/main/main.c index af98ca4..d54712f 100644 --- a/main/main.c +++ b/main/main.c @@ -26,6 +26,8 @@ void app_main(void) chip.revision / 100, chip.revision % 100, chip.cores); project_cfg_defaults(); + // First, so boot messages are queued for syslog until the interface has an IP. + ESP_ERROR_CHECK(aes67_syslog_init()); esp_eth_handle_t eth; ESP_ERROR_CHECK(aes67_board_eth_init(ð)); @@ -34,7 +36,6 @@ void app_main(void) ESP_ERROR_CHECK(aes67_ptp_init()); ESP_ERROR_CHECK(aes67_ptp_start(eth)); ESP_ERROR_CHECK(aes67_tx_init()); - ESP_ERROR_CHECK(aes67_syslog_init()); project_cfg_register(); ESP_ERROR_CHECK(aes67_ota_init()); ESP_ERROR_CHECK(aes67_sdp_sap_init());