Merge pull request #3614 from cesanta/bsd

fixes and additions for BSD
This commit is contained in:
Sergio R. Caprile
2026-07-03 09:52:47 -03:00
committed by GitHub
5 changed files with 316 additions and 68 deletions
+25
View File
@@ -201,6 +201,31 @@ jobs:
with: { fetch-depth: 2 }
- run: sudo apt -y update ; sudo apt -y install binfmt-support qemu-user-static && docker run --rm --privileged multiarch/qemu-user-static --reset -p yes && make -C test ${{ matrix.target }}
bsd:
runs-on: ubuntu-latest
strategy:
fail-fast: false
matrix:
cc: [gcc, clang, g++, clang++]
name: bsd CC=${{ matrix.cc }}
steps:
- uses: actions/checkout@v5
with: { fetch-depth: 2 }
- uses: webfactory/ssh-agent@v0.10.0
with:
ssh-private-key: ${{ secrets.HEALTH_TESTS_SSH_KEY }}
- run: |
source ./test/setup_ga_network.sh
make -C test/bsd shm_tap_bridge
test/bsd/shm_tap_bridge &
sleep 2
make -C test/bsd
- if: success() || failure()
run: |
cat log
test/health.awk < log > json
scp -o "StrictHostKeyChecking=no" json "root@176.9.217.245:/data/downloads/health/bsd_$(date +"%Y%m%d").json"
s390:
runs-on: ubuntu-latest
strategy:
+1
View File
@@ -157,6 +157,7 @@ jobs:
- path: stm32/nucleo-f746zg/minimal
- path: stm32/nucleo-h563zi/minimal
- path: stm32/nucleo-h723zg/minimal
- path: stm32/nucleo-h723zg/minimal-bsd
- path: stm32/nucleo-h743zi/minimal
- path: stm32/stm32h573i-dk/minimal
- path: ti/ek-tm4c1294xl-make-baremetal-builtin
+144 -33
View File
@@ -123,6 +123,62 @@ fail:
#if MG_ENABLE_BSD_SOCKETS
#ifndef MG_ENABLE_BSD_LOG
#define MG_ENABLE_BSD_LOG 0
#define bsd_log(type, tag, a, b, c, n)
#else
#define MG_BSD_LOG_SOCK 1
#define MG_BSD_LOG_ACCEPT 2
#define MG_BSD_LOG_CLOSE 3
#define MG_BSD_LOG_TRANSPORT 4
#define MG_BSD_LOG_RESULT 5
#define MG_BSD_LOG_CONNECT 6
// Type is always a constant, so optimised builds fold the switch to one log.
#define bsd_log(type, tag, a, b, c, n) \
do { \
switch (type) { \
case MG_BSD_LOG_SOCK: { \
struct mg_bsd_sock *log_s = (struct mg_bsd_sock *) (a); \
MG_DEBUG(("SOCK %s %ld %p", tag, (long) (n), log_s->t)); \
break; \
} \
case MG_BSD_LOG_ACCEPT: { \
struct mg_connection *log_mc = (struct mg_connection *) (a); \
if ((n) >= 0) { \
MG_DEBUG(("ACCEPT %lu %p %u %ld", log_mc->id, c, \
(unsigned) mg_ntohs(log_mc->rem.port), (long) (n))); \
} else { \
MG_DEBUG(("ACCEPT %lu %p %u", log_mc->id, c, \
(unsigned) mg_ntohs(log_mc->rem.port))); \
} \
break; \
} \
case MG_BSD_LOG_CLOSE: { \
struct mg_connection *log_mc = (struct mg_connection *) (a); \
if ((n) >= 0) MG_INFO(("CLOSE %s %lu %p %ld", tag, log_mc->id, b, \
(long) (n))); \
else MG_DEBUG(("CLOSE %s %lu %p", tag, log_mc->id, b)); \
break; \
} \
case MG_BSD_LOG_TRANSPORT: \
if ((n) >= 0) MG_INFO(("TRANSP %s %p %ld", tag, a, (long) (n))); \
else MG_DEBUG(("TRANSP %s %p", tag, a)); \
break; \
case MG_BSD_LOG_RESULT: \
MG_DEBUG(("RESULT %s %p %ld", tag, a, (long) (n))); \
break; \
case MG_BSD_LOG_CONNECT: \
if ((c) != NULL) { \
MG_DEBUG(("CONNECT %s %p %p %s %ld", tag, a, b, \
(const char *) (c), (long) (n))); \
} else { \
MG_DEBUG(("CONNECT %s %p %p %ld", tag, a, b, (long) (n))); \
} \
break; \
} \
} while (0)
#endif
struct mg_bsd_sock {
void *t; // opaque transport handle
int fd;
@@ -160,7 +216,6 @@ static void release_sock(int fd) {
if (*p) *p = (*p)->next;
}
int socket(int domain, int type, int proto) {
struct mg_bsd_sock *s = (struct mg_bsd_sock *) calloc(1, sizeof(*s));
if (!s) { errno = ENOMEM; return -1; }
@@ -191,7 +246,7 @@ int accept(int fd, struct sockaddr *addr, socklen_t *addrlen) {
void *t = mg_bsd_transport_accept(ls->t, &peer, ls->nonblock);
if (!t) { if (ls->nonblock) errno = EAGAIN; return -1; }
struct mg_bsd_sock *ns = (struct mg_bsd_sock *) calloc(1, sizeof(*ns));
if (!ns || alloc_sock(ns) < 0) { mg_bsd_transport_free(t); free(ns); errno = ENOMEM; return -1; }
if (!ns || alloc_sock(ns) < 0) { mg_bsd_transport_close(t); free(ns); errno = ENOMEM; return -1; }
ns->t = t; ns->domain = ls->domain; ns->type = ls->type; ns->peer = peer;
if (addr && addrlen) {
size_t sz = sizeof(peer) < *addrlen ? sizeof(peer) : *addrlen;
@@ -246,7 +301,9 @@ ssize_t read(int fd, void *buf, size_t len) { return recv(fd, buf, len, 0); }
int close(int fd) {
struct mg_bsd_sock *s = get(fd);
if (!s) return -1;
bsd_log(MG_BSD_LOG_SOCK, "CLOSE", s, NULL, NULL, (long) fd);
mg_bsd_transport_close(s->t);
bsd_log(MG_BSD_LOG_SOCK, "FREE", s, NULL, NULL, (long) fd);
release_sock(fd);
free(s);
return 0;
@@ -385,17 +442,37 @@ struct mg_bsd_cmd {
static QueueHandle_t s_cmd_q;
// Single-slot DNS resolve state (not reentrant, sufficient for demos)
static struct { struct mg_addr addr; bool done, error; TaskHandle_t caller; } s_resolve;
static bool bsd_qsend(QueueHandle_t q, const void *item, TickType_t ticks) {
BaseType_t ok = xQueueSend(q, item, ticks);
if (ok != pdTRUE) MG_ERROR(("%p", q));
return ok == pdTRUE;
}
// static DNS resolve state; not reentrant (see below)
static struct { struct mg_addr addr; bool error; TaskHandle_t caller; } s_resolve;
// gethostbyname statics (official isn't reentrant anyway, and is obsolete)
static struct hostent s_hostent;
static char *s_h_aliases[1];
static char *s_h_addr_list[2];
static uint32_t s_h_addr;
static char s_h_name[64];
static void resolve_cb(struct mg_connection *c, int ev, void *ev_data) {
if (ev == MG_EV_RESOLVE) { s_resolve.addr = c->rem; s_resolve.done = true; c->is_closing = 1; }
else if ((ev == MG_EV_ERROR || ev == MG_EV_CLOSE) && !s_resolve.done) s_resolve.error = true;
if ((s_resolve.done || s_resolve.error) && s_resolve.caller) {
TaskHandle_t h = s_resolve.caller;
s_resolve.caller = NULL; // prevent double-notify on subsequent MG_EV_CLOSE
xTaskNotifyGive(h);
bool notify = false;
if (ev == MG_EV_RESOLVE) {
s_resolve.addr = c->rem;
c->is_closing = 1;
MG_DEBUG(("%lu resolved", c->id));
notify = true;
} else if (ev == MG_EV_ERROR) {
s_resolve.error = true;
MG_DEBUG(("%lu failed", c->id));
notify = true;
} else if (ev == MG_EV_CLOSE) {
s_resolve.caller = NULL; // The resolver connection has fully unwound.
MG_DEBUG(("%lu done", c->id));
}
if (notify) xTaskNotifyGive(s_resolve.caller);
(void) ev_data;
}
@@ -415,6 +492,7 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
if (ev == MG_EV_ACCEPT) {
// c is the new accepted connection; x is the listening transport
bool ok;
struct mg_xport *nx = xport_alloc();
if (!nx) { c->is_closing = 1; return; }
nx->c = c;
@@ -422,7 +500,8 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
nx->peer.sin_port = c->rem.port;
memcpy(&nx->peer.sin_addr, &c->rem.addr.ip4, 4);
c->fn_data = nx;
xQueueSend(x->accept_q, &nx, 0);
ok = bsd_qsend(x->accept_q, &nx, 0);
bsd_log(MG_BSD_LOG_ACCEPT, NULL, c, x, nx, ok ? 1 : 0);
} else if (ev == MG_EV_READ && x->recv_q) {
// Drain c->recv into recv_q in fixed-size chunks; task1 owns c->recv
size_t off = 0;
@@ -432,7 +511,7 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
if (n > MG_BSD_CHUNK_SIZE) n = MG_BSD_CHUNK_SIZE;
memcpy(chunk.data, c->recv.buf + off, n);
chunk.len = (uint16_t) n;
xQueueSend(x->recv_q, &chunk, portMAX_DELAY);
bsd_qsend(x->recv_q, &chunk, portMAX_DELAY); // TODO(): handle failure
off += n;
}
mg_iobuf_del(&c->recv, 0, c->recv.len);
@@ -444,22 +523,42 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
} else if (ev == MG_EV_CONNECT) {
// Outgoing connection established: wake the task blocked in connect()
if (x->connect_waiter) {
bsd_log(MG_BSD_LOG_CONNECT, "OK", x, c, NULL, 0);
if (x->connect_result) *x->connect_result = 0;
TaskHandle_t h = x->connect_waiter;
x->connect_waiter = NULL; x->connect_result = NULL;
xTaskNotifyGive(h);
}
} else if (ev == MG_EV_CLOSE) {
bsd_log(MG_BSD_LOG_CLOSE, "IN", c, x, NULL, -1);
x->c = NULL; x->closed = true; c->fn_data = NULL;
// If connect() is still waiting, signal failure
if (x->connect_waiter) {
bsd_log(MG_BSD_LOG_CONNECT, "FAIL", x, c, NULL, -1);
if (x->connect_result) *x->connect_result = -1;
TaskHandle_t h = x->connect_waiter;
x->connect_waiter = NULL; x->connect_result = NULL;
xTaskNotifyGive(h);
// The notified task owns the transport and can close/free it immediately.
return;
}
if (x->recv_q) { struct mg_bsd_chunk eof = {.len = 0}; xQueueSend(x->recv_q, &eof, 0); }
if (x->accept_q) { struct mg_xport *nil = NULL; xQueueSend(x->accept_q, &nil, 0); }
if (x->recv_q) {
struct mg_bsd_chunk eof;
memset(&eof, 0, sizeof(eof));
eof.len = 0;
bsd_log(MG_BSD_LOG_CLOSE, "EOF>", c, x, NULL, -1);
bsd_qsend(x->recv_q, &eof, 0);
// recv() wakeup transfers control to the transport owner.
return;
}
if (x->accept_q) {
struct mg_xport *nil = NULL;
bsd_log(MG_BSD_LOG_CLOSE, "ACCEPT>", c, x, NULL, -1);
bsd_qsend(x->accept_q, &nil, 0);
// accept() wakeup transfers control to the transport owner.
return;
}
bsd_log(MG_BSD_LOG_CLOSE, "OUT", c, x, NULL, -1);
}
(void) ev_data;
}
@@ -478,25 +577,36 @@ void mg_bsd_poll(struct mg_mgr *mgr) {
cmd.x->c = c;
*cmd.result = c ? 0 : -1;
} else if (cmd.type == BSD_CMD_CLOSE) {
bsd_log(MG_BSD_LOG_TRANSPORT, "CMD", cmd.x, NULL, NULL, -1);
if (cmd.x->c) {
cmd.x->c->fn_data = NULL;
cmd.x->c->is_draining = 1;
cmd.x->c = NULL;
}
cmd.x->closed = true;
*cmd.result = 0;
} else if (cmd.type == BSD_CMD_CONNECT) {
cmd.x->connect_waiter = cmd.caller;
cmd.x->connect_result = cmd.result;
struct mg_connection *c = mg_connect(mgr, cmd.url, xport_ev, cmd.x);
bsd_log(MG_BSD_LOG_CONNECT, "NEW", cmd.x, c, cmd.url, c ? 1 : 0);
cmd.x->c = c;
if (!c) { *cmd.result = -1; cmd.x->connect_waiter = NULL; cmd.x->connect_result = NULL; }
else notify = false; // xport_ev notifies when connected or on error
} else if (cmd.type == BSD_CMD_RESOLVE) {
s_resolve.done = s_resolve.error = false;
s_resolve.error = false;
s_resolve.caller = cmd.caller;
char url[80];
snprintf(url, sizeof(url), "tcp://%s:0", cmd.url);
if (!mg_connect(mgr, url, resolve_cb, NULL)) s_resolve.error = true;
else notify = false; // resolve_cb notifies when done
struct mg_connection *c;
c = mg_alloc_conn(mgr);
if (c == NULL) {
} else {
c->fn = resolve_cb;
LIST_ADD_HEAD(struct mg_connection, &mgr->conns, c);
mg_call(c, MG_EV_OPEN, NULL);
MG_DEBUG(("%lu resolve %s", c->id, cmd.url));
mg_resolve(c, cmd.url);
notify = false; // resolve_cb notifies when done
}
}
if (notify) xTaskNotifyGive(cmd.caller);
}
@@ -515,6 +625,7 @@ void *mg_bsd_transport_new(int domain, int type, int proto) {
void mg_bsd_transport_free(void *t) {
struct mg_xport *x = (struct mg_xport *) t;
if (!x) return;
bsd_log(MG_BSD_LOG_TRANSPORT, "FREE", x, NULL, NULL, -1);
if (x->recv_q) vQueueDelete(x->recv_q);
if (x->send_q) vQueueDelete(x->send_q);
if (x->accept_q) vQueueDelete(x->accept_q);
@@ -526,7 +637,7 @@ int mg_bsd_transport_listen(void *t, const struct sockaddr_in *addr) {
int result = -1;
struct mg_bsd_cmd cmd = {BSD_CMD_LISTEN, x, {0}, xTaskGetCurrentTaskHandle(), &result};
snprintf(cmd.url, sizeof(cmd.url), "tcp://0.0.0.0:%d", mg_ntohs(addr->sin_port));
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return -1;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
return result;
}
@@ -565,7 +676,7 @@ ssize_t mg_bsd_transport_send(void *t, const void *buf, size_t len, bool nonbloc
if (n > MG_BSD_CHUNK_SIZE) n = MG_BSD_CHUNK_SIZE;
memcpy(chunk.data, (const uint8_t *) buf + sent, n);
chunk.len = (uint16_t) n;
if (xQueueSend(x->send_q, &chunk, ticks) != pdTRUE) break;
bsd_qsend(x->send_q, &chunk, ticks); // TODO(): handle failure
sent += n;
}
return sent > 0 ? (ssize_t) sent : (errno = EAGAIN, -1);
@@ -582,22 +693,17 @@ int mg_bsd_transport_connect(void *t, const struct sockaddr_in *addr, bool nonbl
uint8_t *ip = (uint8_t *) &addr->sin_addr.s_addr;
snprintf(cmd.url, sizeof(cmd.url), "tcp://%d.%d.%d.%d:%d",
ip[0], ip[1], ip[2], ip[3], mg_ntohs(addr->sin_port));
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
bsd_log(MG_BSD_LOG_CONNECT, "REQ", x, NULL, cmd.url, -1);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return -1;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
return result;
}
// gethostbyname: resolve via Mongoose DNS (not reentrant)
static struct hostent s_hostent;
static char *s_h_aliases[1];
static char *s_h_addr_list[2];
static uint32_t s_h_addr;
static char s_h_name[64];
struct hostent *gethostbyname(const char *name) {
struct mg_bsd_cmd cmd = {BSD_CMD_RESOLVE, NULL, {0}, xTaskGetCurrentTaskHandle(), NULL};
snprintf(cmd.url, sizeof(cmd.url), "%s", name);
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return NULL;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
if (s_resolve.error) return NULL;
s_h_addr = s_resolve.addr.addr.ip4;
@@ -642,12 +748,18 @@ void freeaddrinfo(struct addrinfo *res) {
void mg_bsd_transport_close(void *t) {
struct mg_xport *x = (struct mg_xport *) t;
bsd_log(MG_BSD_LOG_TRANSPORT, "CLOSE", x, NULL, NULL, -1);
if (!x->closed && x->c) {
int result = 0;
struct mg_bsd_cmd cmd = {BSD_CMD_CLOSE, x, {0}, xTaskGetCurrentTaskHandle(), &result};
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
bool ok = bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY);
bsd_log(MG_BSD_LOG_TRANSPORT, "CMD>", x, NULL, NULL, ok ? 1 : 0);
if (ok) {
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
bsd_log(MG_BSD_LOG_RESULT, "CMD", x, NULL, NULL, (long) result);
}
}
bsd_log(MG_BSD_LOG_TRANSPORT, "FREE_REQ", x, NULL, NULL, -1);
mg_bsd_transport_free(x);
}
@@ -720,9 +832,8 @@ bool mg_wakeup(struct mg_mgr *mgr, unsigned long conn_id, const void *buf,
m->id = conn_id;
m->len = len;
memcpy(m->data, buf, len);
if (xQueueSend((QueueHandle_t) mgr->pipe.q, &m, 0) != pdTRUE) {
if (!bsd_qsend((QueueHandle_t) mgr->pipe.q, &m, 0)) {
free(m);
MG_ERROR(("xQueueSend"));
return false;
}
return true;
+144 -33
View File
@@ -2,6 +2,62 @@
#if MG_ENABLE_BSD_SOCKETS
#ifndef MG_ENABLE_BSD_LOG
#define MG_ENABLE_BSD_LOG 0
#define bsd_log(type, tag, a, b, c, n)
#else
#define MG_BSD_LOG_SOCK 1
#define MG_BSD_LOG_ACCEPT 2
#define MG_BSD_LOG_CLOSE 3
#define MG_BSD_LOG_TRANSPORT 4
#define MG_BSD_LOG_RESULT 5
#define MG_BSD_LOG_CONNECT 6
// Type is always a constant, so optimised builds fold the switch to one log.
#define bsd_log(type, tag, a, b, c, n) \
do { \
switch (type) { \
case MG_BSD_LOG_SOCK: { \
struct mg_bsd_sock *log_s = (struct mg_bsd_sock *) (a); \
MG_DEBUG(("SOCK %s %ld %p", tag, (long) (n), log_s->t)); \
break; \
} \
case MG_BSD_LOG_ACCEPT: { \
struct mg_connection *log_mc = (struct mg_connection *) (a); \
if ((n) >= 0) { \
MG_DEBUG(("ACCEPT %lu %p %u %ld", log_mc->id, c, \
(unsigned) mg_ntohs(log_mc->rem.port), (long) (n))); \
} else { \
MG_DEBUG(("ACCEPT %lu %p %u", log_mc->id, c, \
(unsigned) mg_ntohs(log_mc->rem.port))); \
} \
break; \
} \
case MG_BSD_LOG_CLOSE: { \
struct mg_connection *log_mc = (struct mg_connection *) (a); \
if ((n) >= 0) MG_INFO(("CLOSE %s %lu %p %ld", tag, log_mc->id, b, \
(long) (n))); \
else MG_DEBUG(("CLOSE %s %lu %p", tag, log_mc->id, b)); \
break; \
} \
case MG_BSD_LOG_TRANSPORT: \
if ((n) >= 0) MG_INFO(("TRANSP %s %p %ld", tag, a, (long) (n))); \
else MG_DEBUG(("TRANSP %s %p", tag, a)); \
break; \
case MG_BSD_LOG_RESULT: \
MG_DEBUG(("RESULT %s %p %ld", tag, a, (long) (n))); \
break; \
case MG_BSD_LOG_CONNECT: \
if ((c) != NULL) { \
MG_DEBUG(("CONNECT %s %p %p %s %ld", tag, a, b, \
(const char *) (c), (long) (n))); \
} else { \
MG_DEBUG(("CONNECT %s %p %p %ld", tag, a, b, (long) (n))); \
} \
break; \
} \
} while (0)
#endif
struct mg_bsd_sock {
void *t; // opaque transport handle
int fd;
@@ -39,7 +95,6 @@ static void release_sock(int fd) {
if (*p) *p = (*p)->next;
}
int socket(int domain, int type, int proto) {
struct mg_bsd_sock *s = (struct mg_bsd_sock *) calloc(1, sizeof(*s));
if (!s) { errno = ENOMEM; return -1; }
@@ -70,7 +125,7 @@ int accept(int fd, struct sockaddr *addr, socklen_t *addrlen) {
void *t = mg_bsd_transport_accept(ls->t, &peer, ls->nonblock);
if (!t) { if (ls->nonblock) errno = EAGAIN; return -1; }
struct mg_bsd_sock *ns = (struct mg_bsd_sock *) calloc(1, sizeof(*ns));
if (!ns || alloc_sock(ns) < 0) { mg_bsd_transport_free(t); free(ns); errno = ENOMEM; return -1; }
if (!ns || alloc_sock(ns) < 0) { mg_bsd_transport_close(t); free(ns); errno = ENOMEM; return -1; }
ns->t = t; ns->domain = ls->domain; ns->type = ls->type; ns->peer = peer;
if (addr && addrlen) {
size_t sz = sizeof(peer) < *addrlen ? sizeof(peer) : *addrlen;
@@ -125,7 +180,9 @@ ssize_t read(int fd, void *buf, size_t len) { return recv(fd, buf, len, 0); }
int close(int fd) {
struct mg_bsd_sock *s = get(fd);
if (!s) return -1;
bsd_log(MG_BSD_LOG_SOCK, "CLOSE", s, NULL, NULL, (long) fd);
mg_bsd_transport_close(s->t);
bsd_log(MG_BSD_LOG_SOCK, "FREE", s, NULL, NULL, (long) fd);
release_sock(fd);
free(s);
return 0;
@@ -264,17 +321,37 @@ struct mg_bsd_cmd {
static QueueHandle_t s_cmd_q;
// Single-slot DNS resolve state (not reentrant, sufficient for demos)
static struct { struct mg_addr addr; bool done, error; TaskHandle_t caller; } s_resolve;
static bool bsd_qsend(QueueHandle_t q, const void *item, TickType_t ticks) {
BaseType_t ok = xQueueSend(q, item, ticks);
if (ok != pdTRUE) MG_ERROR(("%p", q));
return ok == pdTRUE;
}
// static DNS resolve state; not reentrant (see below)
static struct { struct mg_addr addr; bool error; TaskHandle_t caller; } s_resolve;
// gethostbyname statics (official isn't reentrant anyway, and is obsolete)
static struct hostent s_hostent;
static char *s_h_aliases[1];
static char *s_h_addr_list[2];
static uint32_t s_h_addr;
static char s_h_name[64];
static void resolve_cb(struct mg_connection *c, int ev, void *ev_data) {
if (ev == MG_EV_RESOLVE) { s_resolve.addr = c->rem; s_resolve.done = true; c->is_closing = 1; }
else if ((ev == MG_EV_ERROR || ev == MG_EV_CLOSE) && !s_resolve.done) s_resolve.error = true;
if ((s_resolve.done || s_resolve.error) && s_resolve.caller) {
TaskHandle_t h = s_resolve.caller;
s_resolve.caller = NULL; // prevent double-notify on subsequent MG_EV_CLOSE
xTaskNotifyGive(h);
bool notify = false;
if (ev == MG_EV_RESOLVE) {
s_resolve.addr = c->rem;
c->is_closing = 1;
MG_DEBUG(("%lu resolved", c->id));
notify = true;
} else if (ev == MG_EV_ERROR) {
s_resolve.error = true;
MG_DEBUG(("%lu failed", c->id));
notify = true;
} else if (ev == MG_EV_CLOSE) {
s_resolve.caller = NULL; // The resolver connection has fully unwound.
MG_DEBUG(("%lu done", c->id));
}
if (notify) xTaskNotifyGive(s_resolve.caller);
(void) ev_data;
}
@@ -294,6 +371,7 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
if (ev == MG_EV_ACCEPT) {
// c is the new accepted connection; x is the listening transport
bool ok;
struct mg_xport *nx = xport_alloc();
if (!nx) { c->is_closing = 1; return; }
nx->c = c;
@@ -301,7 +379,8 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
nx->peer.sin_port = c->rem.port;
memcpy(&nx->peer.sin_addr, &c->rem.addr.ip4, 4);
c->fn_data = nx;
xQueueSend(x->accept_q, &nx, 0);
ok = bsd_qsend(x->accept_q, &nx, 0);
bsd_log(MG_BSD_LOG_ACCEPT, NULL, c, x, nx, ok ? 1 : 0);
} else if (ev == MG_EV_READ && x->recv_q) {
// Drain c->recv into recv_q in fixed-size chunks; task1 owns c->recv
size_t off = 0;
@@ -311,7 +390,7 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
if (n > MG_BSD_CHUNK_SIZE) n = MG_BSD_CHUNK_SIZE;
memcpy(chunk.data, c->recv.buf + off, n);
chunk.len = (uint16_t) n;
xQueueSend(x->recv_q, &chunk, portMAX_DELAY);
bsd_qsend(x->recv_q, &chunk, portMAX_DELAY); // TODO(): handle failure
off += n;
}
mg_iobuf_del(&c->recv, 0, c->recv.len);
@@ -323,22 +402,42 @@ static void xport_ev(struct mg_connection *c, int ev, void *ev_data) {
} else if (ev == MG_EV_CONNECT) {
// Outgoing connection established: wake the task blocked in connect()
if (x->connect_waiter) {
bsd_log(MG_BSD_LOG_CONNECT, "OK", x, c, NULL, 0);
if (x->connect_result) *x->connect_result = 0;
TaskHandle_t h = x->connect_waiter;
x->connect_waiter = NULL; x->connect_result = NULL;
xTaskNotifyGive(h);
}
} else if (ev == MG_EV_CLOSE) {
bsd_log(MG_BSD_LOG_CLOSE, "IN", c, x, NULL, -1);
x->c = NULL; x->closed = true; c->fn_data = NULL;
// If connect() is still waiting, signal failure
if (x->connect_waiter) {
bsd_log(MG_BSD_LOG_CONNECT, "FAIL", x, c, NULL, -1);
if (x->connect_result) *x->connect_result = -1;
TaskHandle_t h = x->connect_waiter;
x->connect_waiter = NULL; x->connect_result = NULL;
xTaskNotifyGive(h);
// The notified task owns the transport and can close/free it immediately.
return;
}
if (x->recv_q) { struct mg_bsd_chunk eof = {.len = 0}; xQueueSend(x->recv_q, &eof, 0); }
if (x->accept_q) { struct mg_xport *nil = NULL; xQueueSend(x->accept_q, &nil, 0); }
if (x->recv_q) {
struct mg_bsd_chunk eof;
memset(&eof, 0, sizeof(eof));
eof.len = 0;
bsd_log(MG_BSD_LOG_CLOSE, "EOF>", c, x, NULL, -1);
bsd_qsend(x->recv_q, &eof, 0);
// recv() wakeup transfers control to the transport owner.
return;
}
if (x->accept_q) {
struct mg_xport *nil = NULL;
bsd_log(MG_BSD_LOG_CLOSE, "ACCEPT>", c, x, NULL, -1);
bsd_qsend(x->accept_q, &nil, 0);
// accept() wakeup transfers control to the transport owner.
return;
}
bsd_log(MG_BSD_LOG_CLOSE, "OUT", c, x, NULL, -1);
}
(void) ev_data;
}
@@ -357,25 +456,36 @@ void mg_bsd_poll(struct mg_mgr *mgr) {
cmd.x->c = c;
*cmd.result = c ? 0 : -1;
} else if (cmd.type == BSD_CMD_CLOSE) {
bsd_log(MG_BSD_LOG_TRANSPORT, "CMD", cmd.x, NULL, NULL, -1);
if (cmd.x->c) {
cmd.x->c->fn_data = NULL;
cmd.x->c->is_draining = 1;
cmd.x->c = NULL;
}
cmd.x->closed = true;
*cmd.result = 0;
} else if (cmd.type == BSD_CMD_CONNECT) {
cmd.x->connect_waiter = cmd.caller;
cmd.x->connect_result = cmd.result;
struct mg_connection *c = mg_connect(mgr, cmd.url, xport_ev, cmd.x);
bsd_log(MG_BSD_LOG_CONNECT, "NEW", cmd.x, c, cmd.url, c ? 1 : 0);
cmd.x->c = c;
if (!c) { *cmd.result = -1; cmd.x->connect_waiter = NULL; cmd.x->connect_result = NULL; }
else notify = false; // xport_ev notifies when connected or on error
} else if (cmd.type == BSD_CMD_RESOLVE) {
s_resolve.done = s_resolve.error = false;
s_resolve.error = false;
s_resolve.caller = cmd.caller;
char url[80];
snprintf(url, sizeof(url), "tcp://%s:0", cmd.url);
if (!mg_connect(mgr, url, resolve_cb, NULL)) s_resolve.error = true;
else notify = false; // resolve_cb notifies when done
struct mg_connection *c;
c = mg_alloc_conn(mgr);
if (c == NULL) {
} else {
c->fn = resolve_cb;
LIST_ADD_HEAD(struct mg_connection, &mgr->conns, c);
mg_call(c, MG_EV_OPEN, NULL);
MG_DEBUG(("%lu resolve %s", c->id, cmd.url));
mg_resolve(c, cmd.url);
notify = false; // resolve_cb notifies when done
}
}
if (notify) xTaskNotifyGive(cmd.caller);
}
@@ -394,6 +504,7 @@ void *mg_bsd_transport_new(int domain, int type, int proto) {
void mg_bsd_transport_free(void *t) {
struct mg_xport *x = (struct mg_xport *) t;
if (!x) return;
bsd_log(MG_BSD_LOG_TRANSPORT, "FREE", x, NULL, NULL, -1);
if (x->recv_q) vQueueDelete(x->recv_q);
if (x->send_q) vQueueDelete(x->send_q);
if (x->accept_q) vQueueDelete(x->accept_q);
@@ -405,7 +516,7 @@ int mg_bsd_transport_listen(void *t, const struct sockaddr_in *addr) {
int result = -1;
struct mg_bsd_cmd cmd = {BSD_CMD_LISTEN, x, {0}, xTaskGetCurrentTaskHandle(), &result};
snprintf(cmd.url, sizeof(cmd.url), "tcp://0.0.0.0:%d", mg_ntohs(addr->sin_port));
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return -1;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
return result;
}
@@ -444,7 +555,7 @@ ssize_t mg_bsd_transport_send(void *t, const void *buf, size_t len, bool nonbloc
if (n > MG_BSD_CHUNK_SIZE) n = MG_BSD_CHUNK_SIZE;
memcpy(chunk.data, (const uint8_t *) buf + sent, n);
chunk.len = (uint16_t) n;
if (xQueueSend(x->send_q, &chunk, ticks) != pdTRUE) break;
bsd_qsend(x->send_q, &chunk, ticks); // TODO(): handle failure
sent += n;
}
return sent > 0 ? (ssize_t) sent : (errno = EAGAIN, -1);
@@ -461,22 +572,17 @@ int mg_bsd_transport_connect(void *t, const struct sockaddr_in *addr, bool nonbl
uint8_t *ip = (uint8_t *) &addr->sin_addr.s_addr;
snprintf(cmd.url, sizeof(cmd.url), "tcp://%d.%d.%d.%d:%d",
ip[0], ip[1], ip[2], ip[3], mg_ntohs(addr->sin_port));
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
bsd_log(MG_BSD_LOG_CONNECT, "REQ", x, NULL, cmd.url, -1);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return -1;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
return result;
}
// gethostbyname: resolve via Mongoose DNS (not reentrant)
static struct hostent s_hostent;
static char *s_h_aliases[1];
static char *s_h_addr_list[2];
static uint32_t s_h_addr;
static char s_h_name[64];
struct hostent *gethostbyname(const char *name) {
struct mg_bsd_cmd cmd = {BSD_CMD_RESOLVE, NULL, {0}, xTaskGetCurrentTaskHandle(), NULL};
snprintf(cmd.url, sizeof(cmd.url), "%s", name);
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
if (!bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY)) return NULL;
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
if (s_resolve.error) return NULL;
s_h_addr = s_resolve.addr.addr.ip4;
@@ -521,12 +627,18 @@ void freeaddrinfo(struct addrinfo *res) {
void mg_bsd_transport_close(void *t) {
struct mg_xport *x = (struct mg_xport *) t;
bsd_log(MG_BSD_LOG_TRANSPORT, "CLOSE", x, NULL, NULL, -1);
if (!x->closed && x->c) {
int result = 0;
struct mg_bsd_cmd cmd = {BSD_CMD_CLOSE, x, {0}, xTaskGetCurrentTaskHandle(), &result};
xQueueSend(s_cmd_q, &cmd, portMAX_DELAY);
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
bool ok = bsd_qsend(s_cmd_q, &cmd, portMAX_DELAY);
bsd_log(MG_BSD_LOG_TRANSPORT, "CMD>", x, NULL, NULL, ok ? 1 : 0);
if (ok) {
ulTaskNotifyTake(pdTRUE, portMAX_DELAY);
bsd_log(MG_BSD_LOG_RESULT, "CMD", x, NULL, NULL, (long) result);
}
}
bsd_log(MG_BSD_LOG_TRANSPORT, "FREE_REQ", x, NULL, NULL, -1);
mg_bsd_transport_free(x);
}
@@ -599,9 +711,8 @@ bool mg_wakeup(struct mg_mgr *mgr, unsigned long conn_id, const void *buf,
m->id = conn_id;
m->len = len;
memcpy(m->data, buf, len);
if (xQueueSend((QueueHandle_t) mgr->pipe.q, &m, 0) != pdTRUE) {
if (!bsd_qsend((QueueHandle_t) mgr->pipe.q, &m, 0)) {
free(m);
MG_ERROR(("xQueueSend"));
return false;
}
return true;
+2 -2
View File
@@ -88,8 +88,8 @@ wget https://cesanta.com/robots.txt
echo robots.txt:
cat robots.txt
rm robots.txt
wget https://ipv6.google.com/
rm index.html
wget https://ipv6.google.com/ || true
rm index.html || true
echo
# Confirm OK