diff --git a/.github/workflows/nightly.yml b/.github/workflows/nightly.yml index 74487b7f..0e068487 100644 --- a/.github/workflows/nightly.yml +++ b/.github/workflows/nightly.yml @@ -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: diff --git a/.github/workflows/quicktest.yml b/.github/workflows/quicktest.yml index a94673f8..e01bfe97 100644 --- a/.github/workflows/quicktest.yml +++ b/.github/workflows/quicktest.yml @@ -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 diff --git a/mongoose.c b/mongoose.c index 8b7473fd..ed5b52ec 100644 --- a/mongoose.c +++ b/mongoose.c @@ -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; diff --git a/src/bsd.c b/src/bsd.c index 66a126ec..2fe25c0d 100644 --- a/src/bsd.c +++ b/src/bsd.c @@ -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; diff --git a/test/setup_ga_network.sh b/test/setup_ga_network.sh index b9e4a718..caac0363 100755 --- a/test/setup_ga_network.sh +++ b/test/setup_ga_network.sh @@ -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