Add broker and web throughput diagnostics

This commit is contained in:
2026-09-08 22:40:58 +02:00
parent 36d41be422
commit 042499e4d6
20 changed files with 899 additions and 5 deletions
+122
View File
@@ -0,0 +1,122 @@
#pragma once
#include <assert.h>
#include <stdbool.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#include <setjmp.h>
typedef int esp_err_t;
#define ESP_OK 0
#define ESP_ERR_INVALID_STATE 1
#define ESP_ERR_INVALID_ARG 2
#define ESP_ERR_NO_MEM 3
#define ESP_ERR_NOT_FOUND 4
#define ESP_ERR_TIMEOUT 5
static const char *esp_err_to_name(esp_err_t e) { (void)e; return "fake-error"; }
#define MALLOC_CAP_SPIRAM 1
#define MALLOC_CAP_8BIT 2
#define MALLOC_CAP_INTERNAL 4
static size_t allocations;
static void *heap_caps_calloc_prefer(size_t n, size_t s, int choices, ...) {
(void)choices; ++allocations; return calloc(n, s);
}
#define heap_caps_free free
typedef unsigned TickType_t;
typedef unsigned UBaseType_t;
#define pdTRUE 1
#define pdPASS 1
#define portMAX_DELAY UINT32_MAX
#define pdMS_TO_TICKS(ms) (ms)
static int mutex;
typedef int *SemaphoreHandle_t;
static SemaphoreHandle_t xSemaphoreCreateMutex(void) { return &mutex; }
static int xSemaphoreTake(SemaphoreHandle_t m, TickType_t ticks) {
if (*m) { assert(ticks == 0); return 0; } *m = 1; return pdTRUE;
}
static void xSemaphoreGive(SemaphoreHandle_t m) { assert(*m); *m = 0; }
static void vSemaphoreDelete(SemaphoreHandle_t m) { assert(!*m); }
typedef struct { uint8_t *data; size_t capacity, used; } StaticStreamBuffer_t;
typedef StaticStreamBuffer_t *StreamBufferHandle_t;
static StreamBufferHandle_t xStreamBufferCreateStatic(size_t size, size_t trigger,
uint8_t *data, StaticStreamBuffer_t *s) {
assert(trigger == 1); *s = (StaticStreamBuffer_t){data, size - 1, 0}; return s;
}
static size_t xStreamBufferBytesAvailable(StreamBufferHandle_t s) {
assert(mutex); return s->used;
}
static size_t xStreamBufferSend(StreamBufferHandle_t s, const void *data, size_t size, TickType_t ticks) {
assert(mutex && ticks == 0);
if (size > s->capacity - s->used) size = s->capacity - s->used;
memcpy(s->data + s->used, data, size); s->used += size; return size;
}
static size_t xStreamBufferReceive(StreamBufferHandle_t s, void *data, size_t size, TickType_t ticks) {
assert(mutex && ticks == 0);
if (size > s->used) size = s->used;
memcpy(data, s->data, size); s->used -= size;
memmove(s->data, s->data + size, s->used); return size;
}
static void xStreamBufferReset(StreamBufferHandle_t s) { assert(mutex); s->used = 0; }
static void vStreamBufferDelete(StreamBufferHandle_t s) { (void)s; }
typedef struct { uint8_t *data; size_t capacity, item_size, used; } StaticQueue_t;
typedef StaticQueue_t *QueueHandle_t;
static QueueHandle_t xQueueCreateStatic(size_t capacity, size_t size, uint8_t *data, StaticQueue_t *q) {
*q = (StaticQueue_t){data, capacity, size, 0}; return q;
}
static int xQueueSend(QueueHandle_t q, const void *item, TickType_t ticks) {
assert(mutex && ticks == 0); if (q->used == q->capacity) return 0;
memcpy(q->data + q->used++ * q->item_size, item, q->item_size); return pdTRUE;
}
static int xQueueReceive(QueueHandle_t q, void *item, TickType_t ticks) {
assert(mutex && ticks == 0); if (!q->used) return 0;
memcpy(item, q->data, q->item_size); --q->used;
memmove(q->data, q->data + q->item_size, q->used * q->item_size); return pdTRUE;
}
static UBaseType_t uxQueueMessagesWaiting(QueueHandle_t q) { assert(mutex); return q->used; }
static void xQueueReset(QueueHandle_t q) { assert(mutex); q->used = 0; }
static void vQueueDelete(QueueHandle_t q) { (void)q; }
typedef void *TaskHandle_t;
static void (*task_entry)(void *);
static jmp_buf task_exit;
static const uint8_t *serial_input;
static size_t serial_remaining;
static int xTaskCreate(void (*entry)(void *), const char *name, unsigned stack,
void *context, unsigned priority, TaskHandle_t *handle) {
(void)name; (void)stack; (void)context; (void)priority;
task_entry = entry; *handle = &mutex; return pdPASS;
}
static void vTaskDelay(TickType_t ticks) {
assert(!mutex && ticks > 0); if (!serial_remaining) longjmp(task_exit, 1);
}
static size_t serial_service_read(uint8_t *data, size_t size) {
assert(mutex); if (size > serial_remaining) size = serial_remaining;
memcpy(data, serial_input, size); serial_input += size; serial_remaining -= size; return size;
}
static size_t serial_service_write(const uint8_t *data, size_t size) { (void)data; assert(mutex); return size; }
static esp_err_t serial_service_set_session_active(bool active) { (void)active; assert(mutex); return ESP_OK; }
typedef struct {
const char *command, *help, *hint;
int (*func)(int, char **);
void *argtable;
} esp_console_cmd_t;
static int (*registered_command)(int, char **);
static esp_err_t esp_console_cmd_register(const esp_console_cmd_t *cmd) {
registered_command = cmd->func; return ESP_OK;
}
static char console_output[16384];
static size_t console_used;
static int capture_printf(const char *format, ...) {
/* Console formatting must happen after releasing the broker mutex. */
assert(!mutex);
va_list args; va_start(args, format);
int n = vsnprintf(console_output + console_used, sizeof(console_output) - console_used, format, args);
va_end(args); assert(n >= 0 && (size_t)n < sizeof(console_output) - console_used);
console_used += (size_t)n; return n;
}
+26
View File
@@ -0,0 +1,26 @@
#!/usr/bin/env python3
"""Compile unchanged production broker/console with isolated deterministic doubles.
No SDK, device, network, or persistent generated files required.
"""
from pathlib import Path
import shutil
import subprocess
import tempfile
HERE = Path(__file__).resolve().parent
ROOT = HERE.parent.parent
with tempfile.TemporaryDirectory(prefix="broker-diagnostics-") as directory:
build = Path(directory)
(build / "freertos").mkdir()
for name in ("session_broker.c", "session_broker.h", "session_console.c", "session_console.h"):
shutil.copyfile(ROOT / "src" / name, build / name)
shutil.copyfile(HERE / "fake.h", build / "fake.h")
shutil.copyfile(HERE / "test.c", build / "test.c")
for name in ("esp_err.h", "esp_heap_caps.h", "esp_console.h", "serial_service.h",
"freertos/FreeRTOS.h", "freertos/queue.h", "freertos/semphr.h",
"freertos/stream_buffer.h", "freertos/task.h"):
(build / name).write_text('#include "fake.h"\n')
subprocess.run(["cc", "-std=c11", "-D_POSIX_C_SOURCE=200809L", "-Wall", "-Wextra",
"-Werror", "-I", str(build), str(build / "test.c"),
"-o", str(build / "test")], check=True)
subprocess.run([str(build / "test")], check=True)
+155
View File
@@ -0,0 +1,155 @@
#include "fake.h"
#include "session_broker.c"
#define printf capture_printf
#include "session_console.c"
#undef printf
static uint8_t payload[SESSION_BROKER_OUTPUT_SIZE + 256];
static void feed(size_t size) {
assert(size <= sizeof(payload));
serial_input = payload; serial_remaining = size;
if (setjmp(task_exit) == 0) task_entry(NULL);
assert(!mutex && serial_remaining == 0);
}
static session_broker_client_snapshot_t snapshot(session_broker_client_id_t id) {
session_broker_client_snapshot_t s;
assert(session_broker_get_client_snapshot(id, &s) == ESP_OK); return s;
}
static session_broker_global_snapshot_t global(void) {
session_broker_global_snapshot_t s;
assert(session_broker_get_global_snapshot(&s) == ESP_OK); return s;
}
static size_t drain(session_broker_client_id_t id, size_t size) {
uint8_t data[sizeof(payload)]; size_t received;
assert(size <= sizeof(data));
assert(session_broker_read(id, data, size, &received) == ESP_OK); return received;
}
static session_broker_client_id_t connect_type(session_broker_client_type_t type) {
session_broker_client_id_t id;
assert(session_broker_connect(type, "SECRET-NAME-\033[2J", &id) == ESP_OK); return id;
}
static int command(const char *operation, session_broker_client_id_t id, const char *size) {
char id_text[16]; snprintf(id_text, sizeof(id_text), "%u", id);
char *argv[] = {"broker", (char *)operation, id_text, (char *)size};
console_used = 0; console_output[0] = 0;
return registered_command(id ? (size ? 4 : 3) : 2, argv);
}
static void disconnect_all(void) {
session_broker_client_snapshot_t clients[SESSION_BROKER_MAX_CLIENTS];
size_t n = session_broker_list_clients(clients, SESSION_BROKER_MAX_CLIENTS);
for (size_t i = 0; i < n; ++i) assert(session_broker_disconnect(clients[i].id) == ESP_OK);
}
static void comparison(unsigned browsers) {
session_broker_client_id_t ids[4];
ids[0] = connect_type(SESSION_BROKER_CLIENT_USB);
ids[1] = connect_type(SESSION_BROKER_CLIENT_SSH);
for (unsigned i = 0; i < browsers; ++i) ids[2+i] = connect_type(SESSION_BROKER_CLIENT_WEB);
assert(session_broker_clear_counters() == ESP_OK);
size_t total = 0;
while (total < 71292) {
size_t n = 256;
if (n > 71292 - total) n = 71292 - total;
/* Deterministic slow-browser window, not a claim about real scheduling. */
if (browsers == 2 && total < 23412 && n > 23412 - total) n = 23412 - total;
feed(n); total += n;
for (unsigned i = 0; i < 2+browsers; ++i) {
if (browsers == 2 && i == 3 && total < 23412) continue;
drain(ids[i], sizeof(payload));
}
}
session_broker_global_snapshot_t g = global();
uint64_t expected = 71292U * (2+browsers);
uint64_t drops = browsers == 2 ? 19316 : 0;
assert(g.counters.uart_rx_bytes == 71292 && g.counters.disconnections == 0);
assert(g.counters.output_queued_bytes == expected-drops);
assert(g.counters.output_read_bytes == expected-drops && g.counters.output_dropped_bytes == drops);
for (unsigned i = 0; i < 2+browsers; ++i) {
session_broker_client_snapshot_t s = snapshot(ids[i]);
assert(s.counters.output_dropped_bytes == (i == 3 ? drops : 0));
assert(s.counters.output_high_water_bytes == (i == 3 ? 4096 : 256));
}
printf("PASS synthetic %u-browser comparison: UART=71292 copies=%" PRIu64 " queued=%" PRIu64 " read=%" PRIu64 " dropped=%" PRIu64 " disconnect=0\n",
browsers, expected, g.counters.output_queued_bytes, g.counters.output_read_bytes, drops);
disconnect_all();
}
int main(void) {
for (size_t i = 0; i < sizeof(payload); ++i) payload[i] = (uint8_t)i;
assert(session_broker_get_global_snapshot(&(session_broker_global_snapshot_t){0}) == ESP_ERR_INVALID_STATE);
assert(session_broker_init() == ESP_OK);
assert(session_console_register_commands() == ESP_OK);
size_t initial_allocations = allocations;
session_broker_client_id_t fast = connect_type(SESSION_BROKER_CLIENT_USB);
session_broker_client_id_t slow = connect_type(SESSION_BROKER_CLIENT_WEB);
assert(session_broker_request_writer(fast) == ESP_OK);
for (unsigned i = 0; i < 17; ++i) { feed(256); assert(drain(fast, 256) == 256); }
session_broker_client_snapshot_t s = snapshot(slow);
assert(s.output_bytes_pending == 4096 && s.counters.output_high_water_bytes == 4096);
assert(s.counters.output_queued_bytes == 4096 && s.counters.output_dropped_bytes == 256);
s = snapshot(fast);
assert(s.counters.output_read_bytes == 4352 && s.counters.output_high_water_bytes == 256);
assert(s.counters.output_dropped_bytes == 0 && s.is_writer);
puts("PASS isolated fanout overflow and exact 4096-byte high-water");
assert(drain(slow, 4000) == 4000);
assert(snapshot(slow).counters.output_high_water_bytes == 4096);
session_broker_global_snapshot_t before = global();
assert(session_broker_clear_client_counters(slow) == ESP_OK);
s = snapshot(slow);
assert(s.output_bytes_pending == 96 && s.counters.output_high_water_bytes == 96);
assert(s.counters.output_queued_bytes == 0 && s.counters.output_read_bytes == 0 && s.counters.output_dropped_bytes == 0);
assert(global().counters.output_queued_bytes == before.counters.output_queued_bytes);
feed(256); assert(snapshot(slow).counters.output_high_water_bytes == 352);
assert(session_broker_clear_counters() == ESP_OK);
s = snapshot(slow);
assert(s.output_bytes_pending == 352 && s.counters.output_high_water_bytes == 352);
assert(global().writer_id == fast && global().latest_event_sequence == before.latest_event_sequence);
assert(global().counters.uart_rx_bytes == 0 && s.counters.uart_rx_bytes == 0);
assert(drain(slow, 352) == 352);
assert(snapshot(slow).counters.output_read_bytes == 352);
feed(sizeof(payload));
assert(snapshot(slow).counters.output_dropped_bytes == 256);
puts("PASS per-client/global clears seed pending; subsequent read/drop/HWM accounting");
before = global();
assert(session_broker_disconnect(slow) == ESP_OK);
assert(global().counters.output_dropped_bytes == before.counters.output_dropped_bytes + 4096);
assert(session_broker_get_client_snapshot(slow, &s) == ESP_ERR_NOT_FOUND);
session_broker_client_id_t replacement = connect_type(SESSION_BROKER_CLIENT_WEB);
assert(replacement != slow && (replacement & 7) == (slow & 7));
s = snapshot(replacement);
assert(s.output_bytes_pending == 0 && s.counters.output_high_water_bytes == 0);
assert(s.counters.output_queued_bytes == 0 && s.counters.output_read_bytes == 0 && s.counters.output_dropped_bytes == 0);
assert(session_broker_clear_client_counters(slow) == ESP_ERR_NOT_FOUND);
puts("PASS disconnect discard retained globally and generation-safe reuse resets diagnostics");
while (session_broker_list_clients(NULL, 0) < SESSION_BROKER_MAX_CLIENTS)
connect_type(SESSION_BROKER_CLIENT_INTERNAL);
struct { session_broker_client_snapshot_t row; uint64_t guard; } bounded = {.guard = UINT64_MAX};
assert(session_broker_list_clients(&bounded.row, 1) == 1 && bounded.guard == UINT64_MAX);
before = global();
assert(command("counters", 0, NULL) == 0);
assert(strstr(console_output, "pending HWM") && !strstr(console_output, "SECRET-NAME") && !strchr(console_output, '\033'));
unsigned lines = 0;
for (const char *p = strstr(console_output, "ID type"); *p; ++p) lines += *p == '\n';
assert(lines == 1 + SESSION_BROKER_MAX_CLIENTS && console_used < 4096);
assert(global().counters.output_read_bytes == before.counters.output_read_bytes);
puts("PASS bounded eight-row metadata-only console counters; snapshots do not consume output");
feed(1024);
assert(command("read", replacement, "513") == 1);
assert(snapshot(replacement).output_bytes_pending == 1024);
assert(command("read", replacement, "0") == 1);
assert(command("read", replacement, NULL) == 0);
assert(strstr(console_output, "read 512 bytes: 00010203") && !strchr(console_output, '\033'));
assert(snapshot(replacement).output_bytes_pending == 512);
assert(command("read", replacement, "1") == 0);
assert(snapshot(replacement).output_bytes_pending == 511);
puts("PASS console read bounds and binary-safe hexadecimal rendering");
disconnect_all();
assert(command("counters", 0, NULL) == 0);
comparison(1); comparison(2);
assert(allocations == initial_allocations);
puts("PASS no post-init broker allocations; 7 diagnostic groups passed");
cleanup_allocations();
return 0;
}