diff --git a/.github/workflows/linux.yml b/.github/workflows/linux.yml index bf34d250..5714211d 100644 --- a/.github/workflows/linux.yml +++ b/.github/workflows/linux.yml @@ -92,6 +92,12 @@ jobs: set -euo pipefail timeout --preserve-status 5m ./build/test-ttl-expired + - name: Run FreeRTOS BSD wrapper tests + timeout-minutes: 2 + run: | + set -euo pipefail + for t in build/test-freertos-*; do timeout --preserve-status 1m "$t"; done + - name: Testing ICMP socket by stealing system calls in ping timeout-minutes: 2 run: | diff --git a/docs/API.md b/docs/API.md index cc2e493b..ef7f2512 100644 --- a/docs/API.md +++ b/docs/API.md @@ -229,7 +229,13 @@ Processes pending network events. - Parameters: - s: wolfIP instance - now: Current timestamp -- Returns: Number of events processed +- Returns: Milliseconds until the stack next needs `wolfIP_poll()` for its own deadlines, at most `WOLFIP_POLL_MAX_WAIT_MS` (default 1000); 0 when work is still pending; negative on error. Received frames are not included: a caller that sleeps for the returned time must also wake when its link driver has a frame. + +```c +typedef void (*wolfIP_wake_cb)(void *arg); +void wolfIP_set_wake_cb(struct wolfIP *s, wolfIP_wake_cb cb, void *arg); +``` +Registers a callback the stack calls when a socket call (or `wolfIP_recv()` outside `wolfIP_poll()`) leaves work for the next `wolfIP_poll()`, such as a queued frame, a newly armed timer or a socket event raised again after a partial read. While a callback is set, a timer armed between polls starts at the next `wolfIP_poll()`, when its frame goes out, instead of at the time of the previous one. It is never called from inside `wolfIP_poll()`, and runs in the caller's context with whatever lock the caller holds, so it should only signal the thread that runs `wolfIP_poll()`. Pass `NULL` to unregister. Call it before other threads use the stack, or under the lock that serializes `wolfIP_poll()` and the socket calls. ```c void wolfIP_recv(struct wolfIP *s, void *buf, uint32_t len); diff --git a/docs/http_server_howto.md b/docs/http_server_howto.md index b15ae501..090f3007 100644 --- a/docs/http_server_howto.md +++ b/docs/http_server_howto.md @@ -316,8 +316,8 @@ httpd_register_handler(&httpd, "/status", status_handler); /* Main loop: the server runs entirely inside wolfIP_poll(). */ for (;;) { - uint32_t ms_next = wolfIP_poll(s, now_ms()); - /* sleep up to ms_next, service other work, then loop */ + int ms_next = wolfIP_poll(s, now_ms()); + /* sleep up to ms_next, or until the link driver has a frame */ } ``` diff --git a/docs/migrating_from_lwIP.md b/docs/migrating_from_lwIP.md index f1402c39..9ef9a529 100644 --- a/docs/migrating_from_lwIP.md +++ b/docs/migrating_from_lwIP.md @@ -1111,7 +1111,7 @@ Internally, the FreeRTOS port uses: * a poll task that calls `wolfIP_poll()`; * callbacks from wolfIP that wake blocked tasks by giving the socket’s semaphore. -The FreeRTOS poll task locks the stack, calls `wolfIP_poll(ipstack, now_ms)`, unlocks the stack, bounds the next sleep between a minimum and maximum, converts milliseconds to ticks, and then calls `vTaskDelay()`. The default wrapper constants include `WOLFIP_FREERTOS_BSD_MAX_FDS 16`, `WOLFIP_FREERTOS_POLL_MAX_MS 20`, and `WOLFIP_FREERTOS_POLL_MIN_MS 5`. +The FreeRTOS poll task locks the stack, calls `wolfIP_poll(ipstack, now_ms)`, unlocks the stack, bounds the returned time to the next deadline between a minimum and maximum, converts milliseconds to ticks, and then waits on a semaphore for that long. The stack gives the semaphore through `wolfIP_set_wake_cb()` whenever a socket call queues work, and a link driver can give it from its RX interrupt. The default wrapper constants include `WOLFIP_FREERTOS_BSD_MAX_FDS 16`, `WOLFIP_FREERTOS_POLL_MAX_MS 5`, and `WOLFIP_FREERTOS_POLL_MIN_MS 1`. ### 10.2 Blocking semantics in the wrapper @@ -1216,10 +1216,10 @@ static void wolfip_os_poll_task(void *arg) for (;;) { uint64_t now_ms = os_time_millis(); - uint32_t next_ms; + int next_ms; os_mutex_lock(&g_core_lock); - next_ms = (uint32_t)wolfIP_poll(ipstack, now_ms); + next_ms = wolfIP_poll(ipstack, now_ms); os_mutex_unlock(&g_core_lock); if (next_ms < OS_WOLFIP_POLL_MIN_MS) { @@ -1230,7 +1230,8 @@ static void wolfip_os_poll_task(void *arg) next_ms = OS_WOLFIP_POLL_MAX_MS; } - os_sleep_ms(next_ms); + /* Given by the wolfIP_set_wake_cb() callback and the RX interrupt. */ + os_sem_wait_timeout(&g_wake, next_ms); } } ``` @@ -1386,9 +1387,9 @@ Follow these rules in the new OS port: ### 11.5 Poll-task timing -The FreeRTOS wrapper bounds the poll delay between 5 ms and 20 ms by default. That is a reasonable starting point for an RTOS port because it prevents the poll task from spinning while still giving TCP timers, ACKs, retransmissions, and queued TX work regular progress. +The FreeRTOS wrapper sleeps until the deadline `wolfIP_poll()` returns, bounded between 1 ms and 5 ms by default, and socket calls wake it early through `wolfIP_set_wake_cb()`. The minimum prevents the poll task from spinning; the maximum bounds how long a received frame waits when the link driver has no RX interrupt. -For latency-sensitive products, reduce the maximum delay. For power-sensitive products, allow a larger maximum delay only after confirming that retransmission behavior, DNS, DHCP, and application latency still meet product requirements. +For latency-sensitive products without an RX interrupt, reduce the maximum delay, or wake the poll task from the RX interrupt. For power-sensitive products, allow a larger maximum delay only after confirming that receive latency and application behavior still meet product requirements. --- diff --git a/docs/migrating_from_lwIP_JP.md b/docs/migrating_from_lwIP_JP.md index e6a67828..f7f2991f 100644 --- a/docs/migrating_from_lwIP_JP.md +++ b/docs/migrating_from_lwIP_JP.md @@ -1110,7 +1110,7 @@ int close(int sockfd); * `wolfIP_poll()` を呼び出すポールタスク; * ソケットのセマフォを与えることでブロックされたタスクを起動する wolfIP からのコールバック。 -FreeRTOS ポールタスクはスタックをロックし、`wolfIP_poll(ipstack, now_ms)` を呼び出し、スタックをアンロックし、次のスリープを最小と最大の間に制限し、ミリ秒をティックに変換し、`vTaskDelay()` を呼び出します。デフォルトのラッパー定数には `WOLFIP_FREERTOS_BSD_MAX_FDS 16`、`WOLFIP_FREERTOS_POLL_MAX_MS 20`、`WOLFIP_FREERTOS_POLL_MIN_MS 5` が含まれます。 +FreeRTOS ポールタスクはスタックをロックし、`wolfIP_poll(ipstack, now_ms)` を呼び出し、スタックをアンロックし、戻り値の次の期限までの時間を最小と最大の間に制限し、ミリ秒をティックに変換し、その時間だけセマフォを待ちます。ソケット呼び出しが作業をキューに入れると、スタックは `wolfIP_set_wake_cb()` を通じてセマフォを与えます。リンクドライバは RX 割り込みからも与えることができます。デフォルトのラッパー定数には `WOLFIP_FREERTOS_BSD_MAX_FDS 16`、`WOLFIP_FREERTOS_POLL_MAX_MS 5`、`WOLFIP_FREERTOS_POLL_MIN_MS 1` が含まれます。 ### 10.2 ラッパーのブロッキング動作 @@ -1215,10 +1215,10 @@ static void wolfip_os_poll_task(void *arg) for (;;) { uint64_t now_ms = os_time_millis(); - uint32_t next_ms; + int next_ms; os_mutex_lock(&g_core_lock); - next_ms = (uint32_t)wolfIP_poll(ipstack, now_ms); + next_ms = wolfIP_poll(ipstack, now_ms); os_mutex_unlock(&g_core_lock); if (next_ms < OS_WOLFIP_POLL_MIN_MS) { @@ -1229,7 +1229,8 @@ static void wolfip_os_poll_task(void *arg) next_ms = OS_WOLFIP_POLL_MAX_MS; } - os_sleep_ms(next_ms); + /* wolfIP_set_wake_cb() のコールバックと RX 割り込みが与える */ + os_sem_wait_timeout(&g_wake, next_ms); } } ``` @@ -1385,9 +1386,9 @@ int recv(int public_fd, void *buf, size_t len, int flags) ### 11.5 ポールタスクのタイミング -FreeRTOS ラッパーはデフォルトでポール遅延を 5 ms から 20 ms の間に制限します。これは RTOS ポートの合理的な出発点です。ポールタスクがスピンするのを防ぎながら、TCP タイマー、ACK、再送信、キュー済み TX 作業に定期的なプログレスを与えるためです。 +FreeRTOS ラッパーは `wolfIP_poll()` が返す期限までスリープし、デフォルトではその時間を 1 ms から 5 ms の間に制限します。ソケット呼び出しは `wolfIP_set_wake_cb()` を通じてポールタスクを早期に起こします。最小値はポールタスクのスピンを防ぎ、最大値はリンクドライバに RX 割り込みがない場合に受信フレームが待つ時間を制限します。 -レイテンシが重要な製品では、最大遅延を小さくしてください。省電力が重要な製品では、再送信動作、DNS、DHCP、アプリケーションレイテンシが製品要件を満たすことを確認した後にのみ、より大きな最大遅延を許可してください。 +RX 割り込みのないレイテンシ重視の製品では、最大遅延を小さくするか、RX 割り込みからポールタスクを起こしてください。省電力が重要な製品では、受信レイテンシとアプリケーション動作が製品要件を満たすことを確認した後にのみ、より大きな最大遅延を許可してください。 --- diff --git a/docs/porting_guide.md b/docs/porting_guide.md index 970160aa..9836e259 100644 --- a/docs/porting_guide.md +++ b/docs/porting_guide.md @@ -260,7 +260,7 @@ static int tap_poll(struct wolfIP_ll_dev *ll, void *buf, uint32_t len) (void)ll; pfd.fd = tap_fd; pfd.events = POLLIN; - ret = poll(&pfd, 1, 2); + ret = poll(&pfd, 1, 0); if (ret < 0) { perror("poll"); return -1; /* driver error */ @@ -730,7 +730,7 @@ socket-offload model) for contrast. | Binary semaphore / event per socket | Sleep a task until its socket is ready. | | Task creation | Run the poll task. | | Millisecond clock | Provide `now_ms` to `wolfIP_poll()`. | -| Sleep/delay | Idle the poll task between cycles. | +| Timed wait on a wakeable object | Idle the poll task until the next deadline; socket calls and the RX interrupt end it early. | That is the whole dependency list. wolfIP needs no dynamic memory, no per-socket threads, and no timer callbacks from the OS. @@ -738,9 +738,9 @@ threads, and no timer callbacks from the OS. ### 6.2 The poll task: the heartbeat of the stack The poll task is a forever-loop that takes the core mutex, runs one poll cycle, -releases the mutex, and sleeps for a bounded interval. `wolfIP_poll()` returns -`>= 0` on success and a negative value on error; the FreeRTOS version clamps the -sleep to a `[MIN, MAX]` window so the task neither spins nor oversleeps: +releases the mutex, and sleeps until the deadline `wolfIP_poll()` returns. The +FreeRTOS version clamps the sleep to a `[MIN, MAX]` window so the task neither +spins nor oversleeps: ```c static void wolfip_bsd_poll_task(void *arg) @@ -748,14 +748,14 @@ static void wolfip_bsd_poll_task(void *arg) struct wolfIP *ipstack = (struct wolfIP *)arg; for (;;) { - uint32_t next_ms; + int next_ms; TickType_t delay_ticks; - uint64_t now_ms = (uint64_t)xTaskGetTickCount() * (uint64_t)portTICK_PERIOD_MS; + uint64_t now_ms = (uint64_t)xTaskGetTickCount() * 1000u / configTICK_RATE_HZ; /* One poll cycle under the global lock so socket ops and timer * processing see a consistent core state. */ xSemaphoreTake(g_lock, portMAX_DELAY); - next_ms = (uint32_t)wolfIP_poll(ipstack, now_ms); + next_ms = wolfIP_poll(ipstack, now_ms); xSemaphoreGive(g_lock); if (next_ms < WOLFIP_FREERTOS_POLL_MIN_MS) next_ms = WOLFIP_FREERTOS_POLL_MIN_MS; @@ -763,15 +763,28 @@ static void wolfip_bsd_poll_task(void *arg) delay_ticks = pdMS_TO_TICKS(next_ms); if (delay_ticks == 0) delay_ticks = 1; /* always yield at least 1 tick */ - vTaskDelay(delay_ticks); + (void)xSemaphoreTake(g_wake, delay_ticks); } } ``` -The default bounds are 5 ms minimum and 20 ms maximum. That floor stops the -task from busy-spinning; the ceiling guarantees TCP retransmit timers, delayed -ACKs, DHCP, and DNS still fire promptly. Lower the ceiling for latency, raise it -for power — but verify TCP behaviour after raising it. +`wolfIP_poll()` returns the milliseconds until its next timer, ARP retry or +other deadline. It cannot see frames the driver has not handed over yet, so the +ceiling (`WOLFIP_FREERTOS_POLL_MAX_MS`, default 5 ms) bounds how long a received +frame waits on a driver without an RX interrupt. The floor +(`WOLFIP_FREERTOS_POLL_MIN_MS`, default 1 ms, never less than one tick) stops +the task from spinning when `wolfIP_poll()` reports work still pending. Raise +the ceiling for power, but verify receive latency after raising it. + +The sleep ends early when `g_wake` is given. The port registers a callback with +`wolfIP_set_wake_cb()` that gives it; the core calls it whenever a socket call +queues a frame or arms a timer, so transmits leave at once. A link driver's RX +interrupt can do the same through `wolfip_freertos_notify_from_isr()`. Such an +interrupt-driven wake has no minimum sleep: if frames arrive faster than one +poll cycle runs, the poll task never blocks and starves lower-priority tasks. +Mask the RX interrupt in the ISR and re-enable it from `ll->poll` once the RX +ring is drained. Like any FreeRTOS `...FromISR` call, the ISR must run at or +below `configMAX_SYSCALL_INTERRUPT_PRIORITY`. ### 6.3 The core mutex diff --git a/src/port/freeRTOS/README.md b/src/port/freeRTOS/README.md index 2af56c96..e751d1c3 100644 --- a/src/port/freeRTOS/README.md +++ b/src/port/freeRTOS/README.md @@ -16,7 +16,7 @@ This directory provides a FreeRTOS integration layer for wolfIP with: ## Design 1. A global lock protects wolfIP core/socket operations. -2. A poll thread/task calls `wolfIP_poll()` periodically. +2. A poll thread/task calls `wolfIP_poll()` and sleeps until the deadline it returns, bounded by `WOLFIP_FREERTOS_POLL_MIN_MS`/`WOLFIP_FREERTOS_POLL_MAX_MS`. The core wakes it early when a socket call queues work, and a link driver can wake it from its RX interrupt. 3. Blocking socket operations: - Try the underlying non-blocking wolfIP socket call. - If `-WOLFIP_EAGAIN`, register a callback and block on a FreeRTOS semaphore. @@ -26,7 +26,7 @@ This gives application code standard blocking socket behavior while wolfIP remai ## Poll Task Example -The integration creates a dedicated task similar to: +The integration registers a wake callback with `wolfIP_set_wake_cb()` that gives the binary semaphore `g_wake`, and creates a dedicated task similar to: ```c static void wolfip_poll_task(void *arg) @@ -34,12 +34,12 @@ static void wolfip_poll_task(void *arg) struct wolfIP *ipstack = (struct wolfIP *)arg; for (;;) { - uint32_t next_ms; + int next_ms; TickType_t delay_ticks; - uint64_t now_ms = (uint64_t)xTaskGetTickCount() * (uint64_t)portTICK_PERIOD_MS; + uint64_t now_ms = (uint64_t)xTaskGetTickCount() * 1000u / configTICK_RATE_HZ; xSemaphoreTake(g_lock, portMAX_DELAY); - next_ms = (uint32_t)wolfIP_poll(ipstack, now_ms); + next_ms = wolfIP_poll(ipstack, now_ms); xSemaphoreGive(g_lock); if (next_ms < WOLFIP_FREERTOS_POLL_MIN_MS) { @@ -53,11 +53,13 @@ static void wolfip_poll_task(void *arg) if (delay_ticks == 0) { delay_ticks = 1; } - vTaskDelay(delay_ticks); + (void)xSemaphoreTake(g_wake, delay_ticks); } } ``` +`WOLFIP_FREERTOS_POLL_MAX_MS` bounds how long a received frame waits when the link driver has no RX interrupt. `WOLFIP_FREERTOS_POLL_MIN_MS` keeps the task from spinning when `wolfIP_poll()` reports work still pending. + ## Integration Steps 1. Include headers: @@ -91,6 +93,7 @@ close(fd); - `int wolfip_freertos_socket_init(struct wolfIP *ipstack, UBaseType_t poll_task_priority, uint16_t poll_task_stack_words);` - `int socket_last_error(void);` +- `void wolfip_freertos_notify_from_isr(void);` - call from a link driver's receive interrupt so an arriving frame is serviced at once instead of after up to `WOLFIP_FREERTOS_POLL_MAX_MS`. Does nothing before `wolfip_freertos_socket_init()`. Only call it from interrupts at or below `configMAX_SYSCALL_INTERRUPT_PRIORITY`. Each call ends the poll task's sleep, so under sustained receive load it can starve lower-priority tasks; mask the RX interrupt in the ISR and re-enable it from `ll->poll` once the RX ring is drained. - Socket calls: - `socket`, `bind`, `listen`, `accept`, `connect`, `close` - `send`, `sendto`, `recv`, `recvfrom` @@ -101,8 +104,9 @@ close(fd); Defined in `bsd_socket.c`: - `WOLFIP_FREERTOS_BSD_MAX_FDS` (default: `16`) -- `WOLFIP_FREERTOS_POLL_MIN_MS` (default: `5`) -- `WOLFIP_FREERTOS_POLL_MAX_MS` (default: `20`) +- `WOLFIP_FREERTOS_POLL_MIN_MS` (default: `1`) +- `WOLFIP_FREERTOS_POLL_MAX_MS` (default: `5`) +- `WOLFIP_BSD_DEBUG_CALLBACK` (default: `0`) - set to `1` to log socket callbacks from the poll task Override via compiler flags, for example: diff --git a/src/port/freeRTOS/bsd_socket.c b/src/port/freeRTOS/bsd_socket.c index 1b455b22..f14b76ca 100644 --- a/src/port/freeRTOS/bsd_socket.c +++ b/src/port/freeRTOS/bsd_socket.c @@ -38,11 +38,11 @@ #endif #ifndef WOLFIP_FREERTOS_POLL_MAX_MS -#define WOLFIP_FREERTOS_POLL_MAX_MS 20u +#define WOLFIP_FREERTOS_POLL_MAX_MS 5u #endif #ifndef WOLFIP_FREERTOS_POLL_MIN_MS -#define WOLFIP_FREERTOS_POLL_MIN_MS 5u +#define WOLFIP_FREERTOS_POLL_MIN_MS 1u #endif typedef struct { @@ -55,41 +55,68 @@ typedef struct { static struct wolfIP *g_ipstack; static SemaphoreHandle_t g_lock; +static SemaphoreHandle_t volatile g_wake; static wolfip_bsd_fd_entry g_fds[WOLFIP_FREERTOS_BSD_MAX_FDS]; static int g_last_error; + +#ifndef WOLFIP_BSD_DEBUG_CALLBACK +#define WOLFIP_BSD_DEBUG_CALLBACK 0 +#endif +#if WOLFIP_BSD_DEBUG_CALLBACK static volatile uint32_t g_cb_log_count; +#endif + +static void wolfip_bsd_wake(void *arg) +{ + (void)arg; + (void)xSemaphoreGive(g_wake); +} + +void wolfip_freertos_notify_from_isr(void) +{ + BaseType_t woken = pdFALSE; + SemaphoreHandle_t wake = g_wake; + + if (wake == NULL) { + return; + } + (void)xSemaphoreGiveFromISR(wake, &woken); + portYIELD_FROM_ISR(woken); +} static void wolfip_bsd_poll_task(void *arg) { struct wolfIP *ipstack = (struct wolfIP *)arg; for (;;) { - uint32_t next_ms; + int next_ms; TickType_t delay_ticks; - uint64_t now_ms = (uint64_t)xTaskGetTickCount() * (uint64_t)portTICK_PERIOD_MS; + uint64_t now_ms = (uint64_t)xTaskGetTickCount() * 1000u / + configTICK_RATE_HZ; /* Run one wolfIP poll cycle under the global lock so socket operations * and timer processing see a consistent core state. */ xSemaphoreTake(g_lock, portMAX_DELAY); - next_ms = (uint32_t)wolfIP_poll(ipstack, now_ms); + next_ms = wolfIP_poll(ipstack, now_ms); xSemaphoreGive(g_lock); - /* Bound sleep time to keep progress predictable and to avoid either - * spinning too fast or sleeping too long when no timers are pending. */ - if (next_ms < WOLFIP_FREERTOS_POLL_MIN_MS) { + /* Sleep until the next wolfIP deadline, bounded so received frames + * are still polled and the task never spins. */ + if (next_ms < (int)WOLFIP_FREERTOS_POLL_MIN_MS) { next_ms = WOLFIP_FREERTOS_POLL_MIN_MS; } - if (next_ms > WOLFIP_FREERTOS_POLL_MAX_MS) { + if (next_ms > (int)WOLFIP_FREERTOS_POLL_MAX_MS) { next_ms = WOLFIP_FREERTOS_POLL_MAX_MS; } - /* Convert milliseconds to RTOS ticks and always sleep at least one tick - * so the poll task yields CPU time to application tasks. */ + /* Convert milliseconds to RTOS ticks and always block at least one tick + * so the poll task yields CPU time to application tasks. A socket call + * or the RX interrupt giving g_wake ends the wait early. */ delay_ticks = pdMS_TO_TICKS(next_ms); if (delay_ticks == 0) { delay_ticks = 1; } - vTaskDelay(delay_ticks); + (void)xSemaphoreTake(g_wake, delay_ticks); } } @@ -162,6 +189,7 @@ static void wolfip_bsd_socket_cb(int internal_fd, uint16_t events, void *arg) return; } entry->seen_events |= events; +#if WOLFIP_BSD_DEBUG_CALLBACK g_cb_log_count++; if ((events & CB_EVENT_CLOSED) != 0u || (g_cb_log_count & 0x1Fu) == 0u) { printf("[sock_cb] ifd=%d events=0x%04x wait=0x%04x cb_count=%lu\n", @@ -170,6 +198,7 @@ static void wolfip_bsd_socket_cb(int internal_fd, uint16_t events, void *arg) (unsigned)entry->wait_events, (unsigned long)g_cb_log_count); } +#endif if ((events & entry->wait_events) != 0) { (void)xSemaphoreGive(entry->ready_sem); } @@ -230,6 +259,15 @@ int wolfip_freertos_socket_init(struct wolfIP *ipstack, if (g_lock == NULL) { return -WOLFIP_ENOMEM; } + /* Never deleted: an RX interrupt may use it as soon as it exists. */ + if (g_wake == NULL) { + g_wake = xSemaphoreCreateBinary(); + } + if (g_wake == NULL) { + vSemaphoreDelete(g_lock); + g_lock = NULL; + return -WOLFIP_ENOMEM; + } for (i = 0; i < WOLFIP_FREERTOS_BSD_MAX_FDS; i++) { g_fds[i].in_use = 0; @@ -240,10 +278,14 @@ int wolfip_freertos_socket_init(struct wolfIP *ipstack, g_ipstack = ipstack; g_last_error = 0; +#if WOLFIP_BSD_DEBUG_CALLBACK g_cb_log_count = 0; +#endif + wolfIP_set_wake_cb(g_ipstack, wolfip_bsd_wake, NULL); if (xTaskCreate(wolfip_bsd_poll_task, "wolfip_poll", poll_task_stack_words, g_ipstack, poll_task_priority, NULL) != pdPASS) { + wolfIP_set_wake_cb(g_ipstack, NULL, NULL); g_ipstack = NULL; vSemaphoreDelete(g_lock); g_lock = NULL; diff --git a/src/port/freeRTOS/bsd_socket.h b/src/port/freeRTOS/bsd_socket.h index 906e7a82..05be1286 100644 --- a/src/port/freeRTOS/bsd_socket.h +++ b/src/port/freeRTOS/bsd_socket.h @@ -61,6 +61,10 @@ int getsockopt(int sockfd, int level, int optname, int getsockname(int sockfd, struct wolfIP_sockaddr *addr, socklen_t *addrlen); int getpeername(int sockfd, struct wolfIP_sockaddr *addr, socklen_t *addrlen); +/* ISR-safe at or below configMAX_SYSCALL_INTERRUPT_PRIORITY: lets a link + * driver's receive interrupt wake the poll task. */ +void wolfip_freertos_notify_from_isr(void); + #ifdef __cplusplus } #endif diff --git a/src/port/posix/bsd_socket.c b/src/port/posix/bsd_socket.c index f57c0a36..2fc169a1 100644 --- a/src/port/posix/bsd_socket.c +++ b/src/port/posix/bsd_socket.c @@ -508,7 +508,7 @@ static int wolfip_wait_for_event_locked(struct wolfip_fd_entry *entry, short wai return -errno; } if (poll_ret == 0) { - return -ETIMEDOUT; + return -EAGAIN; } while (WOLFIP_HOST_CALL(read)(entry->public_fd, &c, 1) > 0) { wolfip_consume_token_locked(entry, c); @@ -1473,6 +1473,22 @@ int bind(int sockfd, const struct sockaddr *addr, socklen_t addrlen) { } int getsockopt(int sockfd, int level, int optname, void *optval, socklen_t *optlen) { + if (!in_the_stack && level == SOL_SOCKET && optname == SO_RCVTIMEO && + optval && optlen && *optlen >= (socklen_t)sizeof(struct timeval)) { + struct wolfip_fd_entry *entry; + pthread_mutex_lock(&wolfIP_mutex); + entry = wolfip_entry_from_public(sockfd); + if (entry) { + struct timeval *tv = (struct timeval *)optval; + int ms = (entry->rcv_timeout_ms > 0) ? entry->rcv_timeout_ms : 0; + tv->tv_sec = ms / 1000; + tv->tv_usec = (ms % 1000) * 1000; + *optlen = sizeof(struct timeval); + pthread_mutex_unlock(&wolfIP_mutex); + return 0; + } + pthread_mutex_unlock(&wolfIP_mutex); + } conditional_steal_call(getsockopt, sockfd, level, optname, optval, optlen); } @@ -1485,6 +1501,33 @@ int getsockname(int sockfd, struct sockaddr *addr, socklen_t *addrlen) { } int setsockopt(int sockfd, int level, int optname, const void *optval, socklen_t optlen) { + if (!in_the_stack && level == SOL_SOCKET && optname == SO_RCVTIMEO && + optval && optlen >= (socklen_t)sizeof(struct timeval)) { + struct wolfip_fd_entry *entry; + pthread_mutex_lock(&wolfIP_mutex); + entry = wolfip_entry_from_public(sockfd); + if (entry) { + const struct timeval *tv = (const struct timeval *)optval; + int ms; + if (tv->tv_usec < 0 || tv->tv_usec >= 1000000) { + pthread_mutex_unlock(&wolfIP_mutex); + errno = EDOM; + return -1; + } + /* As on Linux: {0, 0} or a huge value blocks forever, a negative one does not wait. */ + if (tv->tv_sec < 0) + ms = 0; + else if (tv->tv_sec >= INT_MAX / 1000 || + (tv->tv_sec == 0 && tv->tv_usec == 0)) + ms = -1; + else + ms = (int)(tv->tv_sec * 1000 + (tv->tv_usec + 999) / 1000); + entry->rcv_timeout_ms = ms; + pthread_mutex_unlock(&wolfIP_mutex); + return 0; + } + pthread_mutex_unlock(&wolfIP_mutex); + } conditional_steal_call(setsockopt, sockfd, level, optname, optval, optlen); } @@ -1978,18 +2021,55 @@ int poll(struct pollfd *fds, nfds_t nfds, int timeout) { * Implemented in port/posix/tap_*.c */ extern int tap_init(struct wolfIP_ll_dev *dev, const char *name, uint32_t host_ip); +extern int tap_get_fd(void); #endif +static int wolfip_wake_pipe[2] = { -1, -1 }; +static int wolfip_rx_fd = -1; + +/* Runs under wolfIP_mutex, so it must not go through the intercepted write(). */ +static void wolfip_posix_wake(void *arg) +{ + char c = 'w'; + + (void)arg; + (void)WOLFIP_HOST_CALL(write)(wolfip_wake_pipe[1], &c, 1); +} + void *wolfIP_sock_posix_ip_loop(void *arg) { struct wolfIP *ipstack = (struct wolfIP *)arg; - uint32_t ms_next; + int ms_next; struct timeval tv; + struct pollfd pfd[2]; + char drain[16]; while (1) { pthread_mutex_lock(&wolfIP_mutex); gettimeofday(&tv, NULL); ms_next = wolfIP_poll(ipstack, tv.tv_sec * 1000 + tv.tv_usec / 1000); pthread_mutex_unlock(&wolfIP_mutex); - usleep(ms_next * 1000); + if (ms_next < 0) + ms_next = 0; + /* Without a wake pipe or RX descriptor nothing can end the sleep early. */ + if ((wolfip_wake_pipe[0] < 0 || wolfip_rx_fd < 0) && ms_next > 1) + ms_next = 1; + /* Sleep until the next deadline, a received frame, or a socket call + * that queued data. poll() ignores a negative fd. */ + pfd[0].fd = wolfip_wake_pipe[0]; + pfd[0].events = POLLIN; + pfd[0].revents = 0; + pfd[1].fd = wolfip_rx_fd; + pfd[1].events = POLLIN; + pfd[1].revents = 0; + if (WOLFIP_HOST_CALL(poll)(pfd, 2, ms_next) > 0) { + if (pfd[0].revents & POLLIN) { + while (WOLFIP_HOST_CALL(read)(wolfip_wake_pipe[0], drain, + sizeof(drain)) > 0) + ; + } + /* A device that went away would end every poll() at once. */ + if (pfd[1].revents & (POLLERR | POLLHUP | POLLNVAL)) + wolfip_rx_fd = -1; + } in_the_stack = 1; } return NULL; @@ -2104,13 +2184,22 @@ void __attribute__((constructor)) init_wolfip_posix() { return; } } + wolfip_rx_fd = vde_get_fd(); #else if (tap_init(tapdev, "wtcp0", host_stack_ip.s_addr) < 0) { perror("tap init"); pthread_mutex_unlock(&wolfIP_mutex); return; } + wolfip_rx_fd = tap_get_fd(); #endif + if (pipe(wolfip_wake_pipe) == 0) { + WOLFIP_HOST_CALL(fcntl)(wolfip_wake_pipe[0], F_SETFD, FD_CLOEXEC); + WOLFIP_HOST_CALL(fcntl)(wolfip_wake_pipe[1], F_SETFD, FD_CLOEXEC); + WOLFIP_HOST_CALL(fcntl)(wolfip_wake_pipe[0], F_SETFL, O_NONBLOCK); + WOLFIP_HOST_CALL(fcntl)(wolfip_wake_pipe[1], F_SETFL, O_NONBLOCK); + wolfIP_set_wake_cb(IPSTACK, wolfip_posix_wake, NULL); + } #if WOLFIP_POSIX_TCPDUMP if (!tcpdump_atexit_registered) { atexit(wolfIP_stop_tcpdump_atexit); diff --git a/src/port/posix/tap_freebsd.c b/src/port/posix/tap_freebsd.c index 7ba9bfcb..8e6b4996 100644 --- a/src/port/posix/tap_freebsd.c +++ b/src/port/posix/tap_freebsd.c @@ -54,7 +54,7 @@ static int tap_poll(struct wolfIP_ll_dev *ll, void *buf, uint32_t len) pfd.fd = tap_fd; pfd.events = POLLIN; - ret = poll(&pfd, 1, 2); + ret = poll(&pfd, 1, 0); if (ret < 0) { perror("poll"); return -1; @@ -167,6 +167,11 @@ static void tap_fetch_mac(struct wolfIP_ll_dev *ll) freeifaddrs(ifas); } +int tap_get_fd(void) +{ + return tap_fd; +} + int tap_init(struct wolfIP_ll_dev *ll, const char *ifname, uint32_t host_ip) { char devpath[PATH_MAX]; diff --git a/src/port/posix/tap_linux.c b/src/port/posix/tap_linux.c index ea7440ad..901bdea3 100644 --- a/src/port/posix/tap_linux.c +++ b/src/port/posix/tap_linux.c @@ -57,7 +57,7 @@ static int tap_poll(struct wolfIP_ll_dev *ll, void *buf, uint32_t len) (void)ll; pfd.fd = tap_fd; pfd.events = POLLIN; - ret = poll(&pfd, 1, 2); + ret = poll(&pfd, 1, 0); if (ret < 0) { perror("poll"); return -1; @@ -77,6 +77,11 @@ static int tap_send(struct wolfIP_ll_dev *ll, void *buf, uint32_t len) return write(tap_fd, buf, len); } +int tap_get_fd(void) +{ + return tap_fd; +} + int tap_init(struct wolfIP_ll_dev *ll, const char *ifname, uint32_t host_ip) { struct ifreq ifr; @@ -120,6 +125,9 @@ int tap_init(struct wolfIP_ll_dev *ll, const char *ifname, uint32_t host_ip) close(tap_fd); return -1; } + /* Setting it marks it user-assigned, so udev's MACAddressPolicy leaves it alone. */ + if (ioctl(tap_fd, SIOCSIFHWADDR, &ifr) < 0) + perror("ioctl SIOCSIFHWADDR"); strncpy(ll->ifname, ifname, sizeof(ll->ifname) - 1); memcpy(ll->mac, ifr.ifr_hwaddr.sa_data, 6); ll->mac[5] ^= 1; diff --git a/src/port/posix/utun_darwin.c b/src/port/posix/utun_darwin.c index 3813556b..b7ef92c8 100644 --- a/src/port/posix/utun_darwin.c +++ b/src/port/posix/utun_darwin.c @@ -75,7 +75,7 @@ static int utun_poll(struct wolfIP_ll_dev *ll, void *buf, uint32_t len) pfd.fd = utun_fd; pfd.events = POLLIN; - if (poll(&pfd, 1, 2) <= 0) + if (poll(&pfd, 1, 0) <= 0) return 0; n = read(utun_fd, tmp, sizeof(tmp)); @@ -149,6 +149,11 @@ static int utun_setup_ipv4(const char *ifname, uint32_t host_ip, uint32_t peer_i return 0; } +int tap_get_fd(void) +{ + return utun_fd; +} + int tap_init(struct wolfIP_ll_dev *ll, const char *requested_ifname, uint32_t host_ip) { struct ctl_info info; diff --git a/src/port/vde2/vde_device.c b/src/port/vde2/vde_device.c index 893cf582..16710cef 100644 --- a/src/port/vde2/vde_device.c +++ b/src/port/vde2/vde_device.c @@ -235,6 +235,11 @@ void vde_cleanup(void) } } +int vde_get_fd(void) +{ + return vde_conn ? vde_datafd(vde_conn) : -1; +} + /** * Get random number for wolfIP (used for MAC generation) */ diff --git a/src/port/vde2/vde_device.h b/src/port/vde2/vde_device.h index 83a81add..00cdcdad 100644 --- a/src/port/vde2/vde_device.h +++ b/src/port/vde2/vde_device.h @@ -41,4 +41,6 @@ int vde_init(struct wolfIP_ll_dev *ll, const char *socket_path, */ void vde_cleanup(void); +int vde_get_fd(void); + #endif /* VDE_DEVICE_H */ diff --git a/src/test/esp/test_esp.c b/src/test/esp/test_esp.c index 4e2f0b0e..1e23f7ae 100644 --- a/src/test/esp/test_esp.c +++ b/src/test/esp/test_esp.c @@ -209,6 +209,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 1) @@ -359,6 +361,8 @@ static void *pt_echoclient(void *arg) /* Test code (host side). * Thread with echo server to test the client. */ +static volatile int host_listening; + static void *pt_echoserver(void *arg) { int fd, listen_fd, ret; @@ -373,6 +377,7 @@ static void *pt_echoserver(void *arg) fd = socket(AF_INET, IPSTACK_SOCK_STREAM, 0); if (fd < 0) { printf("test server socket: %d\n", fd); + host_listening = -1; return (void *)-1; } local_sock.sin_addr.s_addr = inet_addr(HOST_STACK_IP); @@ -381,15 +386,18 @@ static void *pt_echoserver(void *arg) if (ret < 0) { printf("test server bind: %d (%s)\n", ret, strerror(errno)); close(fd); + host_listening = -1; return (void *)-1; } ret = listen(fd, 1); if (ret < 0) { printf("test server listen: %d\n", ret); close(fd); + host_listening = -1; return (void *)-1; } listen_fd = fd; + host_listening = 1; printf("Waiting for client\n"); ret = accept(fd, NULL, NULL); if (ret < 0) { @@ -481,10 +489,19 @@ static void test_wolfip_echoclient(struct wolfIP *s) conn_fd = wolfIP_sock_socket(s, AF_INET, IPSTACK_SOCK_STREAM, 0); printf("client socket: %04x\n", conn_fd); wolfIP_register_callback(s, conn_fd, client_cb, s); + host_listening = 0; + if (pthread_create(&pt, NULL, pt_echoserver, (void*)1) != 0) + host_listening = -1; + /* A SYN that beats listen() is reset and never retried. */ + while (host_listening == 0) + usleep(1000); + if (host_listening < 0) { + printf("host server did not start\n"); + exit(1); + } printf("Connecting to %s:8\n", HOST_STACK_IP); wolfIP_sock_connect(s, conn_fd, (struct wolfIP_sockaddr *)&remote_sock, sizeof(remote_sock)); - pthread_create(&pt, NULL, pt_echoserver, (void*)1); printf("Starting test: echo client active close\n"); ret = test_loop(s, 1); printf("Test echo client active close: %d\n", ret); diff --git a/src/test/freertos_mocks/FreeRTOS.h b/src/test/freertos_mocks/FreeRTOS.h index 5c5bf8de..48b20d86 100644 --- a/src/test/freertos_mocks/FreeRTOS.h +++ b/src/test/freertos_mocks/FreeRTOS.h @@ -11,7 +11,11 @@ typedef uint32_t TickType_t; #define pdFALSE 0 #define pdPASS 1 #define portMAX_DELAY ((TickType_t)0xffffffffu) +#ifndef configTICK_RATE_HZ +#define configTICK_RATE_HZ 1000u +#endif #define portTICK_PERIOD_MS 1u #define pdMS_TO_TICKS(ms) ((TickType_t)(ms)) +#define portYIELD_FROM_ISR(woken) ((void)(woken)) #endif diff --git a/src/test/freertos_mocks/semphr.h b/src/test/freertos_mocks/semphr.h index a9d17b85..2913b2c2 100644 --- a/src/test/freertos_mocks/semphr.h +++ b/src/test/freertos_mocks/semphr.h @@ -10,5 +10,6 @@ SemaphoreHandle_t xSemaphoreCreateMutex(void); BaseType_t xSemaphoreTake(SemaphoreHandle_t sem, TickType_t ticks); BaseType_t xSemaphoreGive(SemaphoreHandle_t sem); void vSemaphoreDelete(SemaphoreHandle_t sem); +#define xSemaphoreGiveFromISR(sem, woken) ((void)(woken), xSemaphoreGive(sem)) #endif diff --git a/src/test/freertos_mocks/wolfip.h b/src/test/freertos_mocks/wolfip.h index 0c71b985..a80a417a 100644 --- a/src/test/freertos_mocks/wolfip.h +++ b/src/test/freertos_mocks/wolfip.h @@ -16,6 +16,7 @@ struct wolfIP_sockaddr { }; typedef void (*tsocket_cb)(int fd, uint16_t event, void *arg); +typedef void (*wolfIP_wake_cb)(void *arg); #define AF_INET 2 #define IPSTACK_SOCK_STREAM 1 @@ -54,5 +55,6 @@ int wolfIP_sock_can_write(struct wolfIP *s, int fd); int wolfIP_sock_can_read(struct wolfIP *s, int fd); int wolfIP_sock_close(struct wolfIP *s, int fd); void wolfIP_register_callback(struct wolfIP *s, int fd, tsocket_cb cb, void *arg); +void wolfIP_set_wake_cb(struct wolfIP *s, wolfIP_wake_cb cb, void *arg); #endif diff --git a/src/test/ipfilter_logger.c b/src/test/ipfilter_logger.c index 5cd3f4ce..b1fe4b78 100644 --- a/src/test/ipfilter_logger.c +++ b/src/test/ipfilter_logger.c @@ -243,6 +243,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) diff --git a/src/test/test_dhcp_dns.c b/src/test/test_dhcp_dns.c index 9e368b9d..c4638d73 100644 --- a/src/test/test_dhcp_dns.c +++ b/src/test/test_dhcp_dns.c @@ -125,6 +125,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) @@ -138,6 +140,8 @@ static int test_loop(struct wolfIP *s, int active_close) /* Test code (host side). * Provide an echo server to exercise the client path. */ +static volatile int host_listening; + static void *pt_echoserver(void *arg) { int fd, listen_fd, ret; @@ -152,6 +156,7 @@ static void *pt_echoserver(void *arg) fd = socket(AF_INET, IPSTACK_SOCK_STREAM, 0); if (fd < 0) { printf("test server socket: %d\n", fd); + host_listening = -1; return (void *)-1; } local_sock.sin_addr.s_addr = inet_addr(HOST_STACK_IP); @@ -160,15 +165,18 @@ static void *pt_echoserver(void *arg) if (ret < 0) { printf("test server bind: %d (%s)\n", ret, strerror(errno)); close(fd); + host_listening = -1; return (void *)-1; } ret = listen(fd, 1); if (ret < 0) { printf("test server listen: %d\n", ret); close(fd); + host_listening = -1; return (void *)-1; } listen_fd = fd; + host_listening = 1; printf("Waiting for client\n"); ret = accept(fd, NULL, NULL); if (ret < 0) { @@ -221,9 +229,18 @@ void test_wolfip_echoclient(struct wolfIP *s) conn_fd = wolfIP_sock_socket(s, AF_INET, IPSTACK_SOCK_STREAM, 0); printf("client socket: %04x\n", conn_fd); wolfIP_register_callback(s, conn_fd, client_cb, s); + host_listening = 0; + if (pthread_create(&pt, NULL, pt_echoserver, (void*)1) != 0) + host_listening = -1; + /* A SYN that beats listen() is reset and never retried. */ + while (host_listening == 0) + usleep(1000); + if (host_listening < 0) { + printf("host server did not start\n"); + exit(1); + } printf("Connecting to %s:8\n", HOST_STACK_IP); wolfIP_sock_connect(s, conn_fd, (struct wolfIP_sockaddr *)&remote_sock, sizeof(remote_sock)); - pthread_create(&pt, NULL, pt_echoserver, (void*)1); printf("Starting test: echo client active close\n"); ret = test_loop(s, 1); printf("Test echo client active close: %d\n", ret); diff --git a/src/test/test_eventloop.c b/src/test/test_eventloop.c index 4c4fcc18..b584e6df 100644 --- a/src/test/test_eventloop.c +++ b/src/test/test_eventloop.c @@ -192,6 +192,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) @@ -337,6 +339,8 @@ void *pt_echoclient(void *arg) /* Test code (host side). * Thread with echo server to test the client. */ +static volatile int host_listening; + static void *pt_echoserver(void *arg) { int fd, listen_fd, ret; @@ -351,6 +355,7 @@ static void *pt_echoserver(void *arg) fd = socket(AF_INET, IPSTACK_SOCK_STREAM, 0); if (fd < 0) { printf("test server socket: %d\n", fd); + host_listening = -1; return (void *)-1; } local_sock.sin_addr.s_addr = inet_addr(HOST_STACK_IP); @@ -359,15 +364,18 @@ static void *pt_echoserver(void *arg) if (ret < 0) { printf("test server bind: %d (%s)\n", ret, strerror(errno)); close(fd); + host_listening = -1; return (void *)-1; } ret = listen(fd, 1); if (ret < 0) { printf("test server listen: %d\n", ret); close(fd); + host_listening = -1; return (void *)-1; } listen_fd = fd; + host_listening = 1; printf("Waiting for client\n"); ret = accept(fd, NULL, NULL); if (ret < 0) { @@ -460,9 +468,18 @@ void test_wolfip_echoclient(struct wolfIP *s) conn_fd = wolfIP_sock_socket(s, AF_INET, IPSTACK_SOCK_STREAM, 0); printf("client socket: %04x\n", conn_fd); wolfIP_register_callback(s, conn_fd, client_cb, s); + host_listening = 0; + if (pthread_create(&pt, NULL, pt_echoserver, (void*)1) != 0) + host_listening = -1; + /* A SYN that beats listen() is reset and never retried. */ + while (host_listening == 0) + usleep(1000); + if (host_listening < 0) { + printf("host server did not start\n"); + exit(1); + } printf("Connecting to %s:8\n", HOST_STACK_IP); wolfIP_sock_connect(s, conn_fd, (struct wolfIP_sockaddr *)&remote_sock, sizeof(remote_sock)); - pthread_create(&pt, NULL, pt_echoserver, (void*)1); printf("Starting test: echo client active close\n"); ret = test_loop(s, 1); printf("Test echo client active close: %d\n", ret); diff --git a/src/test/test_eventloop_tun.c b/src/test/test_eventloop_tun.c index 01b3cbb3..34fb90c6 100644 --- a/src/test/test_eventloop_tun.c +++ b/src/test/test_eventloop_tun.c @@ -194,6 +194,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) @@ -341,6 +343,8 @@ void *pt_echoclient(void *arg) /* Test code (host side). * Thread with echo server to test the client. */ +static volatile int host_listening; + static void *pt_echoserver(void *arg) { int fd, listen_fd, ret; @@ -355,6 +359,7 @@ static void *pt_echoserver(void *arg) fd = socket(AF_INET, IPSTACK_SOCK_STREAM, 0); if (fd < 0) { printf("test server socket: %d\n", fd); + host_listening = -1; return (void *)-1; } local_sock.sin_addr.s_addr = inet_addr(HOST_STACK_IP); @@ -363,15 +368,18 @@ static void *pt_echoserver(void *arg) if (ret < 0) { printf("test server bind: %d (%s)\n", ret, strerror(errno)); close(fd); + host_listening = -1; return (void *)-1; } ret = listen(fd, 1); if (ret < 0) { printf("test server listen: %d\n", ret); close(fd); + host_listening = -1; return (void *)-1; } listen_fd = fd; + host_listening = 1; printf("Waiting for client\n"); ret = accept(fd, NULL, NULL); if (ret < 0) { @@ -463,9 +471,18 @@ void test_wolfip_echoclient(struct wolfIP *s) conn_fd = wolfIP_sock_socket(s, AF_INET, IPSTACK_SOCK_STREAM, 0); printf("client socket: %04x\n", conn_fd); wolfIP_register_callback(s, conn_fd, client_cb, s); + host_listening = 0; + if (pthread_create(&pt, NULL, pt_echoserver, (void*)1) != 0) + host_listening = -1; + /* A SYN that beats listen() is reset and never retried. */ + while (host_listening == 0) + usleep(1000); + if (host_listening < 0) { + printf("host server did not start\n"); + exit(1); + } printf("Connecting to %s:8\n", HOST_STACK_IP); wolfIP_sock_connect(s, conn_fd, (struct wolfIP_sockaddr *)&remote_sock, sizeof(remote_sock)); - pthread_create(&pt, NULL, pt_echoserver, (void*)1); printf("Starting test: echo client active close\n"); ret = test_loop(s, 1); printf("Test echo client active close: %d\n", ret); diff --git a/src/test/test_freertos_close_last_ack.c b/src/test/test_freertos_close_last_ack.c index da17c804..536f932d 100644 --- a/src/test/test_freertos_close_last_ack.c +++ b/src/test/test_freertos_close_last_ack.c @@ -20,6 +20,9 @@ static tsocket_cb registered_cb; static void *registered_arg; static int fake_internal_fd = MARK_TCP_SOCKET; static int close_calls; +static int sem_gives; +static int task_create_fail; +static wolfIP_wake_cb registered_wake_cb; SemaphoreHandle_t xSemaphoreCreateBinary(void) { @@ -55,9 +58,10 @@ BaseType_t xSemaphoreTake(SemaphoreHandle_t sem, TickType_t ticks) BaseType_t xSemaphoreGive(SemaphoreHandle_t sem) { - if (sem == NULL) + sem_gives++; + if (sem == NULL || sem->count > 0) return pdFALSE; - sem->count++; + sem->count = 1; return pdTRUE; } @@ -78,7 +82,7 @@ BaseType_t xTaskCreate(TaskFunction_t task, const char *name, (void)arg; (void)priority; (void)handle; - return pdPASS; + return task_create_fail ? pdFALSE : pdPASS; } void vTaskDelay(TickType_t ticks) @@ -219,19 +223,66 @@ void wolfIP_register_callback(struct wolfIP *s, int fd, tsocket_cb cb, void *arg registered_arg = arg; } +void wolfIP_set_wake_cb(struct wolfIP *s, wolfIP_wake_cb cb, void *arg) +{ + (void)s; + (void)arg; + registered_wake_cb = cb; +} + #include "../port/freeRTOS/bsd_socket.c" +/* Brings the wrapper up through a failed init first, leaving it initialized. */ +static int check_wake(struct wolfIP *stack) +{ + int fails = 0; + int gives; + + gives = sem_gives; + wolfip_freertos_notify_from_isr(); + if (sem_gives != gives) { + printf("ISR wake before init gave a semaphore\n"); + fails++; + } + + task_create_fail = 1; + if (wolfip_freertos_socket_init(stack, 1, 128) == 0 || + registered_wake_cb != NULL || g_lock != NULL || g_ipstack != NULL) { + printf("failed init left the wake callback or the lock behind\n"); + fails++; + } + task_create_fail = 0; + + if (wolfip_freertos_socket_init(stack, 1, 128) != 0 || + registered_wake_cb == NULL || g_wake == NULL) { + printf("init did not register a wake callback\n"); + return fails + 1; + } + registered_wake_cb(NULL); + if (g_wake->count != 1) { + printf("wake callback did not give the poll task's semaphore\n"); + fails++; + } + + (void)xSemaphoreTake(g_wake, 0); + wolfip_freertos_notify_from_isr(); + if (g_wake->count != 1) { + printf("ISR wake did not give the poll task's semaphore\n"); + fails++; + } + return fails; +} + int main(void) { struct wolfIP stack; int fd; int rc; + int frees; memset(&stack, 0, sizeof(stack)); - if (wolfip_freertos_socket_init(&stack, 1, 128) != 0) { - printf("init failed\n"); + if (check_wake(&stack) != 0) return 1; - } fd = socket(AF_INET, SOCK_STREAM, 0); if (fd < 0) { @@ -239,13 +290,14 @@ int main(void) return 1; } + frees = sem_frees; rc = close(fd); if (rc != 0) { printf("close failed rc=%d close_calls=%d sem_allocs=%d sem_frees=%d\n", rc, close_calls, sem_allocs, sem_frees); return 1; } - if (close_calls != 2 || sem_frees != 1) { + if (close_calls != 2 || sem_frees - frees != 1) { printf("unexpected state close_calls=%d sem_allocs=%d sem_frees=%d\n", close_calls, sem_allocs, sem_frees); return 1; diff --git a/src/test/test_httpd.c b/src/test/test_httpd.c index 8cabdd38..ed962643 100644 --- a/src/test/test_httpd.c +++ b/src/test/test_httpd.c @@ -55,6 +55,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) diff --git a/src/test/test_native_wolfssl.c b/src/test/test_native_wolfssl.c index 406a0c3e..97a97ed9 100644 --- a/src/test/test_native_wolfssl.c +++ b/src/test/test_native_wolfssl.c @@ -155,6 +155,8 @@ static int test_loop(struct wolfIP *s, int active_close) struct timeval tv; gettimeofday(&tv, NULL); ms_next = wolfIP_poll(s, tv.tv_sec * 1000 + tv.tv_usec / 1000); + if (ms_next > 1) + ms_next = 1; usleep(ms_next * 1000); if (exit_ok > 0) { if (exit_count++ < 10) diff --git a/src/test/unit/unit.c b/src/test/unit/unit.c index 25a6c90d..20e3cd7c 100644 --- a/src/test/unit/unit.c +++ b/src/test/unit/unit.c @@ -297,6 +297,7 @@ Suite *wolf_suite(void) tcase_add_test(tc_utils, test_udp_no_icmp_unreachable_for_multicast_dst); #ifdef IP_MULTICAST tcase_add_test(tc_utils, test_multicast_join_and_drop_reports); + tcase_add_test(tc_utils, test_multicast_join_wakes_poller); tcase_add_test(tc_utils, test_multicast_join_report_repeated); tcase_add_test(tc_utils, test_multicast_join_report_repeat_heap_full_rearmed_on_poll); tcase_add_test(tc_utils, test_multicast_join_validation_and_shared_refs); @@ -1513,6 +1514,26 @@ Suite *wolf_suite(void) #endif /* WOLFIP_PACKET_SOCKETS */ tcase_add_test(tc_core, test_poll_combined_timer_and_socket_cb_in_same_tick); tcase_add_test(tc_core, test_poll_no_timers_and_no_events_is_noop); + tcase_add_test(tc_core, test_poll_returns_ms_to_next_timer); + tcase_add_test(tc_core, test_poll_returns_zero_when_driver_defers_tx); + tcase_add_test(tc_core, test_poll_returns_arp_retry_deadline); + tcase_add_test(tc_core, test_poll_returns_zero_when_rx_budget_exhausted); + tcase_add_test(tc_core, test_poll_socket_events_pending_tcp); +#if WOLFIP_ENABLE_LOOPBACK + tcase_add_test(tc_core, test_poll_returns_zero_with_loopback_frame_queued); +#endif + tcase_add_test(tc_core, test_poll_returns_zero_while_flush_events_undelivered); + tcase_add_test(tc_core, test_wake_cb_on_socket_tx); + tcase_add_test(tc_core, test_wake_cb_on_register_with_pending_events); + tcase_add_test(tc_core, test_wake_cb_on_recv_and_loopback_outside_poll); + tcase_add_test(tc_core, test_wake_cb_not_fired_inside_poll); +#if WOLFIP_ENABLE_LOOPBACK + tcase_add_test(tc_core, test_wake_cb_on_ack_retry); +#endif + tcase_add_test(tc_core, test_wake_cb_on_accept_send_and_partial_read); + tcase_add_test(tc_core, test_wake_cb_on_icmp_raw_packet_sendto); + tcase_add_test(tc_core, test_wake_cb_starts_timers_armed_between_polls); + tcase_add_test(tc_core, test_wake_cb_deferred_timers_keep_heap_order); tcase_add_test(tc_core, test_poll_last_tick_updated); tcase_add_test(tc_core, test_poll_loopback_interface_iterated); tcase_add_test(tc_core, test_poll_multiple_udp_sockets_both_cbs_dispatched); diff --git a/src/test/unit/unit_tests_branches.c b/src/test/unit/unit_tests_branches.c index b3f2b987..802309f6 100644 --- a/src/test/unit/unit_tests_branches.c +++ b/src/test/unit/unit_tests_branches.c @@ -1141,7 +1141,7 @@ START_TEST(test_poll_dispatches_socket_callback) ts->events = CB_EVENT_READABLE; socket_cb_calls = 0; socket_cb_last_fd = -1; - ck_assert_int_eq(wolfIP_poll(&s, 1), 0); + ck_assert_int_eq(wolfIP_poll(&s, 1), WOLFIP_POLL_MAX_WAIT_MS); ck_assert_int_eq(socket_cb_calls, 1); ck_assert_int_eq(socket_cb_last_fd, SOCKET_UNMARK(udp_sd) | MARK_UDP_SOCKET); } @@ -1158,7 +1158,7 @@ START_TEST(test_poll_fires_expired_timer) tmr.cb = test_timer_cb; timers_binheap_insert(&s.timers, tmr); timer_cb_calls = 0; - ck_assert_int_eq(wolfIP_poll(&s, 200), 0); + ck_assert_int_eq(wolfIP_poll(&s, 200), WOLFIP_POLL_MAX_WAIT_MS); ck_assert_int_eq(timer_cb_calls, 1); } END_TEST @@ -1214,7 +1214,7 @@ START_TEST(test_poll_keeps_timer_armed_after_earlier_cancel) ck_assert_int_gt(id_second, 0); timer_binheap_cancel(&s.timers, (uint32_t)id_first); timer_cb_calls = 0; - ck_assert_int_eq(wolfIP_poll(&s, 250), 0); + ck_assert_int_eq(wolfIP_poll(&s, 250), WOLFIP_POLL_MAX_WAIT_MS); ck_assert_int_eq(timer_cb_calls, 1); } END_TEST @@ -1242,7 +1242,7 @@ START_TEST(test_poll_arp_pending_when_nexthop_unresolved) ck_assert_uint_gt(fifo_len(&ts->sock.udp.txbuf), 0U); last_frame_sent_size = 0; /* Use now > 1000 so the ARP rate-limit window has elapsed. */ - ck_assert_int_eq(wolfIP_poll(&s, 2000), 0); + ck_assert_int_eq(wolfIP_poll(&s, 2000), WOLFIP_POLL_MAX_WAIT_MS); /* Poll should have emitted an ARP request and left the datagram queued. */ ck_assert_uint_eq(last_frame_sent_size, sizeof(struct arp_packet)); ck_assert_uint_gt(fifo_len(&ts->sock.udp.txbuf), 0U); @@ -1284,7 +1284,7 @@ START_TEST(test_poll_filter_block_holds_tx) wolfIP_filter_set_callback(test_filter_cb_block, NULL); wolfIP_filter_set_udp_mask(WOLFIP_FILT_MASK(WOLFIP_FILT_SENDING)); last_frame_sent_size = 0; - ck_assert_int_eq(wolfIP_poll(&s, 2), 0); + ck_assert_int_eq(wolfIP_poll(&s, 2), WOLFIP_POLL_MAX_WAIT_MS); /* Filter blocked send: nothing transmitted, packet still in txbuf. */ ck_assert_uint_eq(last_frame_sent_size, 0U); ck_assert_uint_gt(fifo_len(&ts->sock.udp.txbuf), 0U); @@ -1322,7 +1322,7 @@ START_TEST(test_poll_drains_icmp_tx) ck_assert_int_eq(wolfIP_sock_sendto(&s, icmp_sd, payload, sizeof(payload), 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), (int)sizeof(payload)); last_frame_sent_size = 0; - ck_assert_int_eq(wolfIP_poll(&s, 2), 0); + ck_assert_int_eq(wolfIP_poll(&s, 2), WOLFIP_POLL_MAX_WAIT_MS); ck_assert_uint_gt(last_frame_sent_size, 0U); ck_assert_uint_eq(fifo_len(&ts->sock.udp.txbuf), 0U); } @@ -1617,7 +1617,7 @@ START_TEST(test_udp_send_and_receive_through_poll) ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, buf, sizeof(buf), 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), (int)sizeof(buf)); last_frame_sent_size = 0; - ck_assert_int_eq(wolfIP_poll(&s, 2), 0); + ck_assert_int_eq(wolfIP_poll(&s, 2), WOLFIP_POLL_MAX_WAIT_MS); ck_assert_uint_gt(last_frame_sent_size, 0U); ck_assert_uint_eq(fifo_len(&ts->sock.udp.txbuf), 0U); diff --git a/src/test/unit/unit_tests_multicast.c b/src/test/unit/unit_tests_multicast.c index a8cfbad4..1afb7f64 100644 --- a/src/test/unit/unit_tests_multicast.c +++ b/src/test/unit/unit_tests_multicast.c @@ -98,6 +98,35 @@ START_TEST(test_multicast_join_and_drop_reports) } END_TEST +static int mcast_wake_calls; + +static void mcast_wake_cb(void *arg) +{ + (void)arg; + mcast_wake_calls++; +} + +START_TEST(test_multicast_join_wakes_poller) +{ + struct wolfIP s; + int sd; + struct wolfIP_ip_mreq mreq; + + wolfIP_init(&s); + mock_link_init(&s); + wolfIP_ipconfig_set(&s, 0x0A000002U, 0xFFFFFF00U, 0); + sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(sd, 0); + mcast_wake_calls = 0; + wolfIP_set_wake_cb(&s, mcast_wake_cb, NULL); + + multicast_mreq(&mreq, 0xE9010203U, IPADDR_ANY); + ck_assert_int_eq(wolfIP_sock_setsockopt(&s, sd, WOLFIP_SOL_IP, + WOLFIP_IP_ADD_MEMBERSHIP, &mreq, sizeof(mreq)), 0); + ck_assert_int_eq(mcast_wake_calls, 1); +} +END_TEST + /* RFC 3376 §5.1: the unsolicited join report is repeated once after a short delay */ START_TEST(test_multicast_join_report_repeated) { @@ -290,7 +319,7 @@ START_TEST(test_multicast_udp_send_mac_ttl_loop_and_options) last_frame_sent_size = 0; ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, payload, sizeof(payload), 0, (struct wolfIP_sockaddr *)&dst, sizeof(dst)), (int)sizeof(payload)); - ck_assert_int_eq(wolfIP_poll(&s, 1), 0); + ck_assert_int_ge(wolfIP_poll(&s, 1), 0); ck_assert_uint_gt(last_frame_sent_size, 0); ck_assert_mem_eq(last_frame_sent, "\x01\x00\x5e\x01\x02\x06", 6); udp = (struct wolfIP_udp_datagram *)last_frame_sent; @@ -791,7 +820,7 @@ START_TEST(test_multicast_if_pins_egress_interface) ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, payload, sizeof(payload), 0, (struct wolfIP_sockaddr *)&dst, sizeof(dst)), (int)sizeof(payload)); - ck_assert_int_eq(wolfIP_poll(&s, 1), 0); + ck_assert_int_ge(wolfIP_poll(&s, 1), 0); ck_assert_uint_gt(last_frame_sent_size, 0); ck_assert_mem_eq(last_frame_sent + 6, secondary_mac, 6); @@ -809,7 +838,7 @@ START_TEST(test_multicast_if_pins_egress_interface) ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, payload, sizeof(payload), 0, (struct wolfIP_sockaddr *)&dst, sizeof(dst)), (int)sizeof(payload)); - ck_assert_int_eq(wolfIP_poll(&s, 1), 0); + ck_assert_int_ge(wolfIP_poll(&s, 1), 0); ck_assert_uint_gt(last_frame_sent_size, 0); ck_assert_mem_eq(last_frame_sent + 6, primary_mac, 6); diff --git a/src/test/unit/unit_tests_poll_dispatcher.c b/src/test/unit/unit_tests_poll_dispatcher.c index e088cc6a..e1dce78b 100644 --- a/src/test/unit/unit_tests_poll_dispatcher.c +++ b/src/test/unit/unit_tests_poll_dispatcher.c @@ -1724,9 +1724,582 @@ START_TEST(test_poll_no_timers_and_no_events_is_noop) wolfIP_init(&s); mock_link_init(&s); - /* Nothing registered; should return 0 silently */ - ck_assert_int_eq(wolfIP_poll(&s, 100), 0); - ck_assert_int_eq(wolfIP_poll(&s, 101), 0); + /* Nothing armed: the idle maximum. */ + ck_assert_int_eq(wolfIP_poll(&s, 100), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_int_eq(wolfIP_poll(&s, 101), WOLFIP_POLL_MAX_WAIT_MS); +} +END_TEST + +static void poll_neighbor_setup(struct wolfIP *s) +{ + uint8_t neighbor_mac[6] = {0x02, 0xAA, 0xBB, 0xCC, 0xDD, 0xEE}; + + wolfIP_init(s); + mock_link_init(s); + mock_link_capture_reset(); + wolfIP_ipconfig_set(s, 0x0A000001U, 0xFFFFFF00U, 0); + s->arp.neighbors[0].ip = 0x0A000002U; + s->arp.neighbors[0].if_idx = TEST_PRIMARY_IF; + s->arp.neighbors[0].ts = 1; + memcpy(s->arp.neighbors[0].mac, neighbor_mac, 6); + s->last_tick = 1; +} + +START_TEST(test_poll_returns_ms_to_next_timer) +{ + struct wolfIP s; + struct wolfIP_timer tmr = {0}; + + wolfIP_init(&s); + mock_link_init(&s); + tmr.expires = 250; + tmr.cb = test_timer_cb; + timers_binheap_insert(&s.timers, tmr); + timer_cb_calls = 0; + ck_assert_int_eq(wolfIP_poll(&s, 200), 50); + ck_assert_int_eq(timer_cb_calls, 0); + ck_assert_int_eq(wolfIP_poll(&s, 250), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_int_eq(timer_cb_calls, 1); +} +END_TEST + +START_TEST(test_poll_returns_zero_when_driver_defers_tx) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + uint8_t payload[4] = {0}; + struct tsocket *ts; + int udp_sd; + + poll_neighbor_setup(&s); + udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(udp_sd, 0); + ts = &s.udpsockets[SOCKET_UNMARK(udp_sd)]; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000002U); + ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, payload, sizeof(payload), + 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); + mock_send_eagain_armed = 1; + ck_assert_int_eq(wolfIP_poll(&s, 2), 0); + ck_assert_uint_gt(fifo_len(&ts->sock.udp.txbuf), 0U); + ck_assert_int_eq(wolfIP_poll(&s, 3), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_uint_eq(fifo_len(&ts->sock.udp.txbuf), 0U); +} +END_TEST + +START_TEST(test_poll_returns_arp_retry_deadline) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + uint8_t payload[4] = {0}; + int udp_sd; + + wolfIP_init(&s); + mock_link_init(&s); + wolfIP_ipconfig_set(&s, 0x0A000001U, 0xFFFFFF00U, 0); + udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(udp_sd, 0); + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000003U); + ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, payload, sizeof(payload), + 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); + /* The request goes out at 2000; the next one is allowed at 3001. */ + ck_assert_int_eq(wolfIP_poll(&s, 2000), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_int_eq(wolfIP_poll(&s, 2500), 501); +} +END_TEST + +static int endless_rx_poll(struct wolfIP_ll_dev *dev, void *frame, uint32_t len) +{ + (void)dev; + memset(frame, 0, len < 64U ? len : 64U); + return 64; +} + +START_TEST(test_poll_returns_zero_when_rx_budget_exhausted) +{ + struct wolfIP s; + + wolfIP_init(&s); + mock_link_init(&s); + wolfIP_getdev(&s)->poll = endless_rx_poll; + ck_assert_int_eq(wolfIP_poll(&s, 1), 0); +} +END_TEST + +START_TEST(test_poll_socket_events_pending_tcp) +{ + struct wolfIP s; + struct tsocket *ts; + int tcp_sd; + + wolfIP_init(&s); + mock_link_init(&s); + tcp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_STREAM, WI_IPPROTO_TCP); + ck_assert_int_gt(tcp_sd, 0); + ts = &s.tcpsockets[SOCKET_UNMARK(tcp_sd)]; + wolfIP_register_callback(&s, tcp_sd, test_socket_cb, NULL); + + /* handle_socket_callbacks() skips a closed socket without CLOSED. */ + ts->sock.tcp.state = TCP_CLOSED; + ts->events = CB_EVENT_READABLE; + ck_assert_int_eq(socket_events_pending(&s), 0); + ts->events |= CB_EVENT_CLOSED; + ck_assert_int_eq(socket_events_pending(&s), 1); + ts->events = 0; + ck_assert_int_eq(socket_events_pending(&s), 0); + ts->close_notify_pending = 1; + ck_assert_int_eq(socket_events_pending(&s), 1); +} +END_TEST + +#if WOLFIP_ENABLE_LOOPBACK +static struct wolfIP *loopback_timer_stack; + +static void loopback_send_from_timer(void *arg) +{ + struct wolfIP_ll_dev *loop; + uint8_t frame[16] = {0}; + + (void)arg; + loop = wolfIP_getdev_ex(loopback_timer_stack, TEST_LOOPBACK_IF); + ck_assert_int_eq(wolfIP_loopback_send(loop, frame, sizeof(frame)), + (int)sizeof(frame)); +} + +START_TEST(test_poll_returns_zero_with_loopback_frame_queued) +{ + struct wolfIP s; + struct wolfIP_timer tmr = {0}; + + wolfIP_init(&s); + mock_link_init(&s); + loopback_timer_stack = &s; + tmr.expires = 5; + tmr.cb = loopback_send_from_timer; + ck_assert_uint_ne(timers_binheap_insert(&s.timers, tmr), 0U); + + /* The timer runs after the loopback device was polled. */ + ck_assert_int_eq(wolfIP_poll(&s, 10), 0); + ck_assert_uint_eq(s.loopback_count, 1U); + ck_assert_int_eq(wolfIP_poll(&s, 10), WOLFIP_POLL_MAX_WAIT_MS); +} +END_TEST +#endif + +START_TEST(test_poll_returns_zero_while_flush_events_undelivered) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + uint8_t payload[4] = {0}; + struct tsocket *ts; + int udp_sd; + + poll_neighbor_setup(&s); + udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(udp_sd, 0); + ts = &s.udpsockets[SOCKET_UNMARK(udp_sd)]; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000002U); + ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, payload, sizeof(payload), + 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); + ts->events = 0; + wolfIP_register_callback(&s, udp_sd, test_socket_cb, NULL); + socket_cb_calls = 0; + + /* The flush drains the txbuf after callbacks ran: WRITABLE is pending. */ + ck_assert_int_eq(wolfIP_poll(&s, 2), 0); + ck_assert_int_eq(socket_cb_calls, 0); + ck_assert_uint_ne(ts->events & CB_EVENT_WRITABLE, 0U); + ck_assert_int_eq(wolfIP_poll(&s, 2), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_int_eq(socket_cb_calls, 1); +} +END_TEST + +static int wake_calls; +static struct wolfIP *wake_stack; +static int wake_udp_sd; + +static void test_wake_cb(void *arg) +{ + ck_assert_ptr_eq(arg, wake_stack); + wake_calls++; +} + +static void test_wake_send_from_timer(void *arg) +{ + struct wolfIP_sockaddr_in sin; + uint8_t payload[4] = {0}; + + (void)arg; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000002U); + ck_assert_int_eq(wolfIP_sock_sendto(wake_stack, wake_udp_sd, payload, + sizeof(payload), 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); +} + +static void wake_setup(struct wolfIP *s) +{ + poll_neighbor_setup(s); + wake_stack = s; + wake_calls = 0; + wolfIP_set_wake_cb(s, test_wake_cb, s); +} + +START_TEST(test_wake_cb_on_socket_tx) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + uint8_t payload[4] = {0}; + int udp_sd; + int tcp_sd; + + wake_setup(&s); + udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(udp_sd, 0); + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000002U); + ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, payload, sizeof(payload), + 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); + ck_assert_int_eq(wake_calls, 1); + + tcp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_STREAM, WI_IPPROTO_TCP); + ck_assert_int_gt(tcp_sd, 0); + (void)wolfIP_sock_connect(&s, tcp_sd, (struct wolfIP_sockaddr *)&sin, + sizeof(sin)); + ck_assert_int_eq(wake_calls, 2); + + /* Nothing queued: no wake. */ + ck_assert_int_eq(wolfIP_sock_recvfrom(&s, udp_sd, payload, + sizeof(payload), 0, NULL, NULL), -WOLFIP_EAGAIN); + ck_assert_int_eq(wake_calls, 2); + + wolfIP_set_wake_cb(&s, NULL, NULL); + ck_assert_int_eq(wolfIP_sock_sendto(&s, udp_sd, payload, sizeof(payload), + 0, (struct wolfIP_sockaddr *)&sin, sizeof(sin)), + (int)sizeof(payload)); + ck_assert_int_eq(wake_calls, 2); +} +END_TEST + +START_TEST(test_wake_cb_on_register_with_pending_events) +{ + struct wolfIP s; + struct tsocket *ts; + int udp_sd; + + wake_setup(&s); + udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_UDP); + ck_assert_int_gt(udp_sd, 0); + ts = &s.udpsockets[SOCKET_UNMARK(udp_sd)]; + + ts->events = 0; + wolfIP_register_callback(&s, udp_sd, test_socket_cb, NULL); + ck_assert_int_eq(wake_calls, 0); + ts->events = CB_EVENT_READABLE; + wolfIP_register_callback(&s, udp_sd, NULL, NULL); + ck_assert_int_eq(wake_calls, 0); + wolfIP_register_callback(&s, udp_sd, test_socket_cb, NULL); + ck_assert_int_eq(wake_calls, 1); +} +END_TEST + +START_TEST(test_wake_cb_on_recv_and_loopback_outside_poll) +{ + struct wolfIP s; + struct wolfIP_ll_dev *loop; + uint8_t frame[64] = {0}; + + wake_setup(&s); + wolfIP_recv(&s, frame, sizeof(frame)); + ck_assert_int_eq(wake_calls, 1); + wolfIP_recv_ex(&s, TEST_PRIMARY_IF, frame, sizeof(frame)); + ck_assert_int_eq(wake_calls, 2); + + loop = wolfIP_getdev_ex(&s, TEST_LOOPBACK_IF); + ck_assert_ptr_nonnull(loop); + ck_assert_int_eq(wolfIP_loopback_send(loop, frame, 16), 16); + ck_assert_int_eq(wake_calls, 3); +} +END_TEST + +START_TEST(test_wake_cb_not_fired_inside_poll) +{ + struct wolfIP s; + struct wolfIP_timer tmr = {0}; + struct tsocket *ts; + + wake_setup(&s); + wake_udp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, + WI_IPPROTO_UDP); + ck_assert_int_gt(wake_udp_sd, 0); + ts = &s.udpsockets[SOCKET_UNMARK(wake_udp_sd)]; + tmr.expires = 5; + tmr.cb = test_wake_send_from_timer; + timers_binheap_insert(&s.timers, tmr); + + ck_assert_int_eq(wolfIP_poll(&s, 10), WOLFIP_POLL_MAX_WAIT_MS); + ck_assert_int_eq(wake_calls, 0); + /* Queued during the poll, sent by the same poll. */ + ck_assert_uint_gt(last_frame_sent_size, 0U); + ck_assert_uint_eq(fifo_len(&ts->sock.udp.txbuf), 0U); +} +END_TEST + +#if WOLFIP_ENABLE_LOOPBACK +START_TEST(test_wake_cb_on_ack_retry) +{ + struct wolfIP s; + struct tsocket *ts; + struct wolfIP_ll_dev *loop; + uint8_t frame[16] = {0}; + unsigned int i; + + wake_setup(&s); + loop = wolfIP_getdev_ex(&s, TEST_LOOPBACK_IF); + ck_assert_ptr_nonnull(loop); + for (i = 0; i < WOLFIP_LOOPBACK_QUEUE_DEPTH; i++) { + ck_assert_int_eq(wolfIP_loopback_send(loop, frame, sizeof(frame)), + (int)sizeof(frame)); + } + + /* No TX FIFO and a full loopback queue: the ACK is deferred. */ + ts = &s.tcpsockets[0]; + memset(ts, 0, sizeof(*ts)); + ts->proto = WI_IPPROTO_TCP; + ts->S = &s; + ts->if_idx = TEST_LOOPBACK_IF; + ts->sock.tcp.state = TCP_ESTABLISHED; + ts->src_port = 1234; + ts->dst_port = 4321; + ts->local_ip = 0x7F000001U; + ts->remote_ip = 0x7F000001U; + wake_calls = 0; + tcp_send_ack(ts); + ck_assert_uint_eq(ts->sock.tcp.ack_retry_pending, 1U); + ck_assert_int_eq(wake_calls, 1); +} +END_TEST +#endif + +START_TEST(test_wake_cb_on_accept_send_and_partial_read) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in peer; + socklen_t peer_len = sizeof(peer); + struct tsocket *lsn; + struct tsocket *ts; + uint8_t buf[8] = {0}; + int fd; + int acc; + + wolfIP_init(&s); + mock_link_init(&s); + wolfIP_ipconfig_set(&s, LLK_LOCAL_IP, LLK_NET_MASK, 0); + wake_stack = &s; + wolfIP_set_wake_cb(&s, test_wake_cb, &s); + fd = llk_open_listener(&s); + lsn = &s.tcpsockets[SOCKET_UNMARK(fd)]; + llk_keep_arp_fresh(&s, LLK_ATT_IP); + llk_attacker_syn(&s, LLK_ATT_IP, 41000, 1, 0); + llk_complete_handshake(&s, lsn, LLK_ATT_IP, 41000, 1); + ck_assert_int_eq(lsn->sock.tcp.state, TCP_ESTABLISHED); + + wake_calls = 0; + memset(&peer, 0, sizeof(peer)); + acc = wolfIP_sock_accept(&s, fd, (struct wolfIP_sockaddr *)&peer, + &peer_len); + ck_assert_int_ge(acc, 0); + ck_assert_int_eq(wake_calls, 1); + + ck_assert_int_eq(wolfIP_sock_send(&s, acc, buf, sizeof(buf), 0), + (int)sizeof(buf)); + ck_assert_int_eq(wake_calls, 2); + + /* A partial read re-raises READABLE for the callback. */ + ts = &s.tcpsockets[SOCKET_UNMARK(acc)]; + ck_assert_int_eq(queue_insert(&ts->sock.tcp.rxbuf, buf, ts->sock.tcp.ack, + sizeof(buf)), 0); + wolfIP_register_callback(&s, acc, test_socket_cb, NULL); + ts->events = 0; + /* A scaled window hides the 4-byte update, so no ACK wakes the poller. */ + ts->sock.tcp.ws_enabled = 1; + ts->sock.tcp.rcv_wscale = 14; + wake_calls = 0; + ck_assert_int_eq(wolfIP_sock_recv(&s, acc, buf, 4, 0), 4); + ck_assert_uint_ne(ts->events & CB_EVENT_READABLE, 0U); + ck_assert_int_eq(wake_calls, 1); +} +END_TEST + +START_TEST(test_wake_cb_on_icmp_raw_packet_sendto) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + uint8_t payload[8] = {0}; + int sd; + + wake_setup(&s); + sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_DGRAM, WI_IPPROTO_ICMP); + ck_assert_int_gt(sd, 0); + payload[0] = ICMP_ECHO_REQUEST; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_addr.s_addr = ee32(0x0A000002U); + ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, payload, sizeof(payload), 0, + (struct wolfIP_sockaddr *)&sin, sizeof(sin)), (int)sizeof(payload)); + ck_assert_int_eq(wake_calls, 1); + +#if WOLFIP_RAWSOCKETS + sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_RAW, WI_IPPROTO_UDP); + ck_assert_int_ge(sd, 0); + ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, payload, sizeof(payload), 0, + (struct wolfIP_sockaddr *)&sin, sizeof(sin)), (int)sizeof(payload)); + ck_assert_int_eq(wake_calls, 2); +#endif + +#if WOLFIP_PACKET_SOCKETS + { + struct wolfIP_sockaddr_ll sll; + uint8_t frame[ETH_HEADER_LEN + 8] = {0}; + + sd = wolfIP_sock_socket(&s, AF_PACKET, IPSTACK_SOCK_RAW, + ee16(ETH_TYPE_IP)); + ck_assert_int_ge(sd, 0); + memset(&sll, 0, sizeof(sll)); + sll.sll_family = AF_PACKET; + sll.sll_protocol = ee16(ETH_TYPE_IP); + sll.sll_ifindex = TEST_PRIMARY_IF; + sll.sll_halen = 6; + memset(sll.sll_addr, 0xFF, 6); + ck_assert_int_eq(wolfIP_sock_bind(&s, sd, + (struct wolfIP_sockaddr *)&sll, sizeof(sll)), 0); + ck_assert_int_eq(wolfIP_sock_sendto(&s, sd, frame, sizeof(frame), 0, + (struct wolfIP_sockaddr *)&sll, sizeof(sll)), + (int)sizeof(frame)); + ck_assert_int_eq(wake_calls, 3); + } +#endif +} +END_TEST + +START_TEST(test_wake_cb_starts_timers_armed_between_polls) +{ + struct wolfIP s; + struct wolfIP_sockaddr_in sin; + struct tsocket *ts; + int tcp_sd; + + wake_setup(&s); + ck_assert_int_eq(wolfIP_poll(&s, 1), WOLFIP_POLL_MAX_WAIT_MS); + tcp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_STREAM, WI_IPPROTO_TCP); + ck_assert_int_gt(tcp_sd, 0); + ts = &s.tcpsockets[SOCKET_UNMARK(tcp_sd)]; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + sin.sin_port = ee16(1234); + sin.sin_addr.s_addr = ee32(0x0A000002U); + mock_link_capture_reset(); + (void)wolfIP_sock_connect(&s, tcp_sd, (struct wolfIP_sockaddr *)&sin, + sizeof(sin)); + + /* Connect 900 ms into an idle sleep: the SYN leaves at 901. */ + (void)wolfIP_poll(&s, 901); + ck_assert_uint_eq(last_frame_sent_count, 1U); + (void)wolfIP_poll(&s, 900 + TCP_RTO_MIN_MS); + ck_assert_uint_eq(last_frame_sent_count, 1U); + ck_assert_uint_eq(ts->sock.tcp.ctrl_rto_retries, 0U); + (void)wolfIP_poll(&s, 901 + TCP_RTO_MIN_MS); + ck_assert_uint_eq(last_frame_sent_count, 2U); + + /* Without a wake callback the caller polls on its own schedule. */ + wolfIP_set_wake_cb(&s, NULL, NULL); + ck_assert_int_eq(wolfIP_sock_close(&s, tcp_sd), 0); + (void)wolfIP_poll(&s, 2000); + tcp_sd = wolfIP_sock_socket(&s, AF_INET, IPSTACK_SOCK_STREAM, WI_IPPROTO_TCP); + ck_assert_int_gt(tcp_sd, 0); + ts = &s.tcpsockets[SOCKET_UNMARK(tcp_sd)]; + (void)wolfIP_sock_connect(&s, tcp_sd, (struct wolfIP_sockaddr *)&sin, + sizeof(sin)); + mock_link_capture_reset(); + (void)wolfIP_poll(&s, 2900); + ck_assert_uint_eq(last_frame_sent_count, 1U); + (void)wolfIP_poll(&s, 2000 + TCP_RTO_MIN_MS); + ck_assert_uint_eq(last_frame_sent_count, 2U); + ck_assert_uint_eq(ts->sock.tcp.ctrl_rto_retries, 1U); +} +END_TEST + +static struct wolfIP *deferred_stack; +static int deferred_fired[3]; +static int deferred_fired_n; + +static void deferred_record_cb(void *arg) +{ + deferred_fired[deferred_fired_n++] = (int)(uintptr_t)arg; +} + +static void deferred_arm_in_poll_cb(void *arg) +{ + struct wolfIP_timer tmr = {0}; + + (void)arg; + tmr.expires = deferred_stack->last_tick + 200; + tmr.cb = deferred_record_cb; + tmr.arg = (void *)1; + ck_assert_uint_ne(timers_binheap_insert(&deferred_stack->timers, tmr), 0U); +} + +START_TEST(test_wake_cb_deferred_timers_keep_heap_order) +{ + struct wolfIP s; + struct wolfIP_timer tmr = {0}; + + wake_setup(&s); + deferred_stack = &s; + deferred_fired_n = 0; + tmr.expires = 1000; + tmr.cb = deferred_arm_in_poll_cb; + ck_assert_uint_ne(timers_binheap_insert(&s.timers, tmr), 0U); + (void)wolfIP_poll(&s, 1000); + + /* Armed between polls: they start at 1150, after the one armed in poll. */ + tmr.cb = deferred_record_cb; + tmr.expires = 1100; + tmr.arg = (void *)2; + ck_assert_uint_ne(timers_binheap_insert(&s.timers, tmr), 0U); + tmr.expires = 1300; + tmr.arg = (void *)3; + ck_assert_uint_ne(timers_binheap_insert(&s.timers, tmr), 0U); + + ck_assert_int_eq(wolfIP_poll(&s, 1150), 50); + ck_assert_int_eq(deferred_fired_n, 0); + (void)wolfIP_poll(&s, 1200); + ck_assert_int_eq(deferred_fired_n, 1); + ck_assert_int_eq(deferred_fired[0], 1); + (void)wolfIP_poll(&s, 1249); + ck_assert_int_eq(deferred_fired_n, 1); + (void)wolfIP_poll(&s, 1250); + ck_assert_int_eq(deferred_fired_n, 2); + ck_assert_int_eq(deferred_fired[1], 2); + (void)wolfIP_poll(&s, 1450); + ck_assert_int_eq(deferred_fired_n, 3); + ck_assert_int_eq(deferred_fired[2], 3); } END_TEST diff --git a/src/wolfip.c b/src/wolfip.c index d0ebf728..fba45329 100644 --- a/src/wolfip.c +++ b/src/wolfip.c @@ -72,6 +72,7 @@ static inline int wolfIP_is_loopback_if(unsigned int if_idx) #if WOLFIP_ENABLE_LOOPBACK static int wolfIP_loopback_send(struct wolfIP_ll_dev *ll, void *buf, uint32_t len); #endif +static void wolfIP_wake(struct wolfIP *s); static void wolfIP_recv_on(struct wolfIP *s, unsigned int if_idx, void *buf, uint32_t len); struct wolfIP_eth_frame; @@ -1475,6 +1476,7 @@ struct wolfIP; struct wolfIP_timer { uint32_t id; + uint8_t deferred; uint64_t expires; void *arg; void (*cb)(void *arg); @@ -1493,6 +1495,7 @@ struct wolfIP_route_entry { struct timers_binheap { struct wolfIP_timer timers[MAX_TIMERS]; uint32_t size; + uint8_t defer; }; /* The main wolfip stack context structure. */ @@ -1542,6 +1545,10 @@ struct wolfIP { #endif uint16_t ipcounter; uint64_t last_tick; + uint64_t poll_next_at; + wolfIP_wake_cb wake_cb; + void *wake_arg; + uint8_t in_poll; #if WOLFIP_ENABLE_FORWARDING uint32_t route_generation; struct wolfIP_route_entry routes[WOLFIP_MAX_ROUTES]; @@ -1742,6 +1749,7 @@ static int wolfIP_loopback_send(struct wolfIP_ll_dev *ll, void *buf, uint32_t le s->loopback_pending_len[slot] = len; s->loopback_tail = (slot + 1U) % WOLFIP_LOOPBACK_QUEUE_DEPTH; s->loopback_count++; + wolfIP_wake(s); return (int)len; } @@ -1789,8 +1797,21 @@ static inline int wolfIP_ll_is_non_ethernet(struct wolfIP *s, unsigned int if_id return (ll && ll->non_ethernet) ? 1 : 0; } -static inline int wolfIP_ll_send_frame(struct wolfIP *s, unsigned int if_idx, - void *buf, uint32_t len) +static void wolfIP_wake(struct wolfIP *s) +{ + if (s && s->wake_cb && !s->in_poll) + s->wake_cb(s->wake_arg); +} + +/* Lowers the deadline the current wolfIP_poll() returns. */ +static void wolfIP_poll_by(struct wolfIP *s, uint64_t when) +{ + if (when < s->poll_next_at) + s->poll_next_at = when; +} + +static inline int wolfIP_ll_xmit(struct wolfIP *s, unsigned int if_idx, + void *buf, uint32_t len) { struct wolfIP_ll_dev *ll; uint32_t frame_mtu; @@ -1852,6 +1873,17 @@ static inline int wolfIP_ll_send_frame(struct wolfIP *s, unsigned int if_idx, return ll->send(ll, buf, len); } +static inline int wolfIP_ll_send_frame(struct wolfIP *s, unsigned int if_idx, + void *buf, uint32_t len) +{ + int ret = wolfIP_ll_xmit(s, if_idx, buf, len); + + /* A driver that is out of TX buffers wants the frame retried soon. */ + if (ret == -WOLFIP_EAGAIN) + wolfIP_poll_by(s, s->last_tick); + return ret; +} + static inline struct ipconf *wolfIP_ipconf_at(struct wolfIP *s, unsigned int if_idx) { if (!s || if_idx >= s->if_count) @@ -2789,6 +2821,7 @@ void wolfIP_register_callback(struct wolfIP *s, int sock_fd, tsocket_cb cb, void *arg) { struct tsocket *t; + uint16_t pending = 0; if (!s) return; if (sock_fd < 0) @@ -2799,18 +2832,21 @@ void wolfIP_register_callback(struct wolfIP *s, int sock_fd, tsocket_cb cb, t = &s->tcpsockets[SOCKET_UNMARK(sock_fd)]; t->callback = cb; t->callback_arg = arg; + pending = t->events; } else if (IS_SOCKET_UDP(sock_fd)) { if (SOCKET_UNMARK(sock_fd) >= MAX_UDPSOCKETS) return; t = &s->udpsockets[SOCKET_UNMARK(sock_fd)]; t->callback = cb; t->callback_arg = arg; + pending = t->events; } else if (IS_SOCKET_ICMP(sock_fd)) { if (SOCKET_UNMARK(sock_fd) >= MAX_ICMPSOCKETS) return; t = &s->icmpsockets[SOCKET_UNMARK(sock_fd)]; t->callback = cb; t->callback_arg = arg; + pending = t->events; } #if WOLFIP_RAWSOCKETS else if (IS_SOCKET_RAW(sock_fd)) { @@ -2818,6 +2854,7 @@ void wolfIP_register_callback(struct wolfIP *s, int sock_fd, tsocket_cb cb, return; s->rawsockets[SOCKET_UNMARK(sock_fd)].callback = cb; s->rawsockets[SOCKET_UNMARK(sock_fd)].callback_arg = arg; + pending = s->rawsockets[SOCKET_UNMARK(sock_fd)].events; } #endif #if WOLFIP_PACKET_SOCKETS @@ -2826,8 +2863,12 @@ void wolfIP_register_callback(struct wolfIP *s, int sock_fd, tsocket_cb cb, return; s->packetsockets[SOCKET_UNMARK(sock_fd)].callback = cb; s->packetsockets[SOCKET_UNMARK(sock_fd)].callback_arg = arg; + pending = s->packetsockets[SOCKET_UNMARK(sock_fd)].events; } #endif + /* Events raised while no callback was set are delivered by the next poll. */ + if (cb && pending) + wolfIP_wake(s); } /* Timers */ @@ -2869,6 +2910,7 @@ static int timers_binheap_insert(struct timers_binheap *heap, struct wolfIP_time if (heap->size >= MAX_TIMERS) return 0; /* heap full */ tmr.id = timer_id++; + tmr.deferred = heap->defer; /* Insert at the end */ heap->timers[heap->size] = tmr; heap->size++; @@ -2925,6 +2967,54 @@ static void timers_heap_rebase(struct timers_binheap *heap, uint64_t now) } } +static void timers_sift_down(struct timers_binheap *heap, uint32_t i) +{ + uint32_t n = heap->size; + + while (2 * i + 1 < n) { + uint32_t j = 2 * i + 1; + struct wolfIP_timer tmp; + if (j + 1 < n && heap->timers[j + 1].expires < heap->timers[j].expires) + j++; + if (heap->timers[i].expires <= heap->timers[j].expires) + break; + tmp = heap->timers[i]; + heap->timers[i] = heap->timers[j]; + heap->timers[j] = tmp; + i = j; + } +} + +/* With a wake callback the poll follows a call that armed a timer at once, + * so start the timer at that poll rather than at the stale last_tick. */ +static void timers_start_deferred(struct timers_binheap *heap, uint64_t delta) +{ + uint32_t i; + int moved = 0; + + for (i = 0; i < heap->size; i++) { + if (!heap->timers[i].deferred) + continue; + heap->timers[i].deferred = 0; + if (heap->timers[i].expires != 0 && delta != 0) { + heap->timers[i].expires += delta; + moved = 1; + } + } + if (!moved) + return; + for (i = heap->size / 2; i-- > 0; ) + timers_sift_down(heap, i); +} + +/* Earliest pending expiry, or 0 when no timer is armed. */ +static uint64_t timers_next_expiry(struct timers_binheap *heap) +{ + while (heap->size > 0 && heap->timers[0].expires == 0) + timers_binheap_pop(heap); + return (heap->size > 0) ? heap->timers[0].expires : 0; +} + static int is_timer_expired(struct timers_binheap *heap, uint64_t now) { while (heap->size > 0 && heap->timers[0].expires == 0) { @@ -4063,8 +4153,10 @@ static int tcp_send_empty(struct tsocket *t, uint8_t flags) return -WOLFIP_EINVAL; tcp = (struct wolfIP_tcp_seg *)buffer; frame_len = tcp_build_empty(t, tcp, flags); - if (fifo_push(&t->sock.tcp.txbuf, tcp, frame_len) == 0) + if (fifo_push(&t->sock.tcp.txbuf, tcp, frame_len) == 0) { + wolfIP_wake(t->S); return 0; + } /* Pure ACKs have no retransmission path, so do not drop them when the * shared data/control TX FIFO is saturated by already queued payload. */ @@ -4095,10 +4187,12 @@ static void tcp_send_ack(struct tsocket *t) if (!t) return; ret = tcp_send_empty(t, TCP_FLAG_ACK); - if (ret == -WOLFIP_EAGAIN) + if (ret == -WOLFIP_EAGAIN) { t->sock.tcp.ack_retry_pending = 1; - else if (ret >= 0) + wolfIP_wake(t->S); + } else if (ret >= 0) { t->sock.tcp.ack_retry_pending = 0; + } } static void tcp_send_reset_reply(struct wolfIP *s, unsigned int if_idx, @@ -4332,7 +4426,11 @@ static int tcp_send_syn(struct tsocket *t, uint8_t flags) opt_len++; } tcp->hlen = ((20 + opt_len) << 2) & 0xF0; - return fifo_push(&t->sock.tcp.txbuf, tcp, sizeof(struct wolfIP_tcp_seg) + opt_len); + if (fifo_push(&t->sock.tcp.txbuf, tcp, + sizeof(struct wolfIP_tcp_seg) + opt_len) < 0) + return -1; + wolfIP_wake(t->S); + return 0; } /* Returns true when handshake/teardown control traffic is outstanding and @@ -7436,6 +7534,7 @@ int wolfIP_sock_accept(struct wolfIP *s, int sockfd, struct wolfIP_sockaddr *add * peer waiting for its own retransmission timer. */ fifo_init(&newts->sock.tcp.txbuf, newts->txmem, TXBUF_SIZE); newts->sock.tcp.ack_retry_pending = 1; + wolfIP_wake(s); /* Readiness follows the child's own buffers. The listener's * CB_EVENT_READABLE means "a connection is pending accept" and * would otherwise dispatch a read callback on an accepted socket @@ -7652,6 +7751,7 @@ int wolfIP_sock_sendto(struct wolfIP *s, int sockfd, const void *buf, size_t len (struct wolfIP_tcp_seg *)((uint8_t *)last_desc + sizeof(*last_desc)); last_tcp->flags |= TCP_FLAG_PSH; } + wolfIP_wake(s); return sent; } } else if (IS_SOCKET_UDP(sockfd)) { @@ -7733,6 +7833,7 @@ int wolfIP_sock_sendto(struct wolfIP *s, int sockfd, const void *buf, size_t len (uint16_t)(frame_len - ETH_HEADER_LEN)); if (fifo_push(&ts->sock.udp.txbuf, udp, frame_len) < 0) return -WOLFIP_EAGAIN; + wolfIP_wake(s); if (sin && !ts->sock.udp.connected) { /* An unconnected socket adopts the explicit destination as * its last destination for subsequent plain sends (DHCP/DNS @@ -7814,6 +7915,7 @@ int wolfIP_sock_sendto(struct wolfIP *s, int sockfd, const void *buf, size_t len (uint16_t)(frame_len - ETH_HEADER_LEN)); if (fifo_push(&ts->sock.udp.txbuf, icmp, frame_len) < 0) return -WOLFIP_EAGAIN; + wolfIP_wake(s); return (int)payload_len; } #if WOLFIP_RAWSOCKETS @@ -7934,6 +8036,7 @@ int wolfIP_sock_sendto(struct wolfIP *s, int sockfd, const void *buf, size_t len return -WOLFIP_EAGAIN; if (fifo_push(&rs->txbuf, rip, total_len) < 0) return -WOLFIP_EAGAIN; + wolfIP_wake(s); return (int)len; } #endif @@ -7981,6 +8084,7 @@ int wolfIP_sock_sendto(struct wolfIP *s, int sockfd, const void *buf, size_t len return -WOLFIP_EAGAIN; if (fifo_push(&ps->txbuf, pkt_frame, (uint32_t)len) < 0) return -WOLFIP_EAGAIN; + wolfIP_wake(s); return (int)len; } #endif @@ -8038,8 +8142,11 @@ int wolfIP_sock_recvfrom(struct wolfIP *s, int sockfd, void *buf, size_t len, in int ret = queue_pop(&ts->sock.tcp.rxbuf, buf, len); if (ret > 0) { uint16_t win_after = tcp_adv_win(ts, 1); - if (queue_len(&ts->sock.tcp.rxbuf) > 0) + if (queue_len(&ts->sock.tcp.rxbuf) > 0) { ts->events |= CB_EVENT_READABLE; + if (ts->callback) + wolfIP_wake(s); + } if (win_after > win_before) tcp_send_ack(ts); } @@ -8052,8 +8159,11 @@ int wolfIP_sock_recvfrom(struct wolfIP *s, int sockfd, void *buf, size_t len, in int ret = queue_pop(&ts->sock.tcp.rxbuf, buf, len); if (ret > 0) { uint16_t win_after = tcp_adv_win(ts, 1); - if (queue_len(&ts->sock.tcp.rxbuf) > 0) + if (queue_len(&ts->sock.tcp.rxbuf) > 0) { ts->events |= CB_EVENT_READABLE; + if (ts->callback) + wolfIP_wake(s); + } if (win_after > win_before) tcp_send_ack(ts); } @@ -8318,6 +8428,7 @@ static int udp_mcast_join(struct wolfIP *s, struct tsocket *ts, ip4 group, (wolfIP_getrandom() % IGMP_UNSOLICITED_REPORT_MS) + 1U; m->tmr_unsol = igmp_arm_report(s, m, m->unsol_at, igmp_unsolicited_timer_cb); + wolfIP_wake(s); } return 0; } @@ -10814,11 +10925,13 @@ static void arp_request(struct wolfIP *s, unsigned int if_idx, ip4 tip) * so the first request is never held back by the window. */ if (s->arp.last_arp[if_idx] != 0 && s->arp.last_arp[if_idx] + 1000 > s->last_tick) { + wolfIP_poll_by(s, s->arp.last_arp[if_idx] + 1000); return; } /* Store tick+1 so a request sent at tick 0 is distinguishable from * "never sent" (last_arp == 0). */ s->arp.last_arp[if_idx] = s->last_tick + 1; + wolfIP_poll_by(s, s->arp.last_arp[if_idx] + 1000); memset(&arp, 0, sizeof(struct arp_packet)); eth_output_add_header(s, if_idx, NULL, &arp.eth, ETH_TYPE_ARP); arp.htype = ee16(1); /* Ethernet */ @@ -11898,9 +12011,12 @@ void wolfIP_recv(struct wolfIP *s, void *buf, uint32_t len) #endif memcpy(frame + ETH_HEADER_LEN, buf, len); wolfIP_recv_on(s, WOLFIP_PRIMARY_IF_IDX, frame, len + ETH_HEADER_LEN); + wolfIP_wake(s); return; } wolfIP_recv_on(s, WOLFIP_PRIMARY_IF_IDX, buf, len); + /* A frame handed in outside wolfIP_poll() leaves its events for the next. */ + wolfIP_wake(s); } void wolfIP_recv_ex(struct wolfIP *s, unsigned int if_idx, void *buf, uint32_t len) @@ -11915,9 +12031,11 @@ void wolfIP_recv_ex(struct wolfIP *s, unsigned int if_idx, void *buf, uint32_t l #endif memcpy(frame + ETH_HEADER_LEN, buf, len); wolfIP_recv_on(s, if_idx, frame, len + ETH_HEADER_LEN); + wolfIP_wake(s); return; } wolfIP_recv_on(s, if_idx, buf, len); + wolfIP_wake(s); } /* DNS Client */ @@ -12544,6 +12662,8 @@ static void poll_devices(struct wolfIP *s) budget--; } } while (len > 0 && budget > 0); + if (budget == 0) + wolfIP_poll_by(s, s->last_tick); } } @@ -12662,6 +12782,45 @@ static void handle_socket_callbacks(struct wolfIP *s) #endif } +/* True when handle_socket_callbacks() would dispatch something. */ +static int socket_events_pending(struct wolfIP *s) +{ + int i; + + for (i = 0; i < MAX_TCPSOCKETS; i++) { + struct tsocket *ts = &s->tcpsockets[i]; + if (ts->close_notify_pending) + return 1; + if (ts->callback && ts->events && + !((ts->sock.tcp.state == TCP_CLOSED) && + !(ts->events & CB_EVENT_CLOSED))) + return 1; + } + for (i = 0; i < MAX_UDPSOCKETS; i++) { + if (s->udpsockets[i].callback && s->udpsockets[i].events) + return 1; + } + for (i = 0; i < MAX_ICMPSOCKETS; i++) { + if (s->icmpsockets[i].callback && s->icmpsockets[i].events) + return 1; + } +#if WOLFIP_RAWSOCKETS + for (i = 0; i < WOLFIP_MAX_RAWSOCKETS; i++) { + struct rawsocket *r = &s->rawsockets[i]; + if (r->used && r->callback && r->events) + return 1; + } +#if WOLFIP_PACKET_SOCKETS + for (i = 0; i < WOLFIP_MAX_PACKETSOCKETS; i++) { + struct packetsocket *p = &s->packetsockets[i]; + if (p->used && p->callback && p->events) + return 1; + } +#endif +#endif + return 0; +} + static void flush_tcp_tx(struct wolfIP *s, uint64_t now) { int i; @@ -13120,11 +13279,14 @@ static void flush_packet_tx(struct wolfIP *s) * * This function also handles timers for all supported protocols. * - * TODO: Return the number of milliseconds to wait before - * calling it again. + * Returns the number of milliseconds until the next deadline, at most + * WOLFIP_POLL_MAX_WAIT_MS; received frames can need service sooner. */ int wolfIP_poll(struct wolfIP *s, uint64_t now) { + uint64_t next_tmr; + int32_t wait; + if (!s) return -WOLFIP_EINVAL; @@ -13173,7 +13335,12 @@ int wolfIP_poll(struct wolfIP *s, uint64_t now) #endif } + timers_start_deferred(&s->timers, + (now > s->last_tick) ? now - s->last_tick : 0); s->last_tick = now; + s->in_poll = 1; + s->timers.defer = 0; + s->poll_next_at = now + WOLFIP_POLL_MAX_WAIT_MS; /* Poll the device */ poll_devices(s); @@ -13197,7 +13364,29 @@ int wolfIP_poll(struct wolfIP *s, uint64_t now) flush_raw_tx(s); flush_packet_tx(s); - return 0; + /* The flushes can raise events after this poll dispatched callbacks. */ + if (socket_events_pending(s)) + wolfIP_poll_by(s, now); +#if WOLFIP_ENABLE_LOOPBACK + if (s->loopback_count > 0) + wolfIP_poll_by(s, now); +#endif + next_tmr = timers_next_expiry(&s->timers); + if (next_tmr != 0) + wolfIP_poll_by(s, next_tmr); + s->in_poll = 0; + s->timers.defer = (s->wake_cb != NULL); + + wait = (int32_t)((uint32_t)s->poll_next_at - (uint32_t)now); + return (wait > 0) ? (int)wait : 0; +} + +void wolfIP_set_wake_cb(struct wolfIP *s, wolfIP_wake_cb cb, void *arg) +{ + if (!s) + return; + s->wake_cb = cb; + s->wake_arg = arg; } void wolfIP_ipconfig_set(struct wolfIP *s, ip4 ip, ip4 mask, ip4 gw) diff --git a/wolfip.h b/wolfip.h index 41c9f778..a1078cb2 100644 --- a/wolfip.h +++ b/wolfip.h @@ -505,7 +505,15 @@ int wolfIP_register_eapol_handler(struct wolfIP *s, uint32_t len), void *ctx); size_t wolfIP_instance_size(void); +/* What wolfIP_poll() returns when no deadline is pending. */ +#ifndef WOLFIP_POLL_MAX_WAIT_MS +#define WOLFIP_POLL_MAX_WAIT_MS 1000U +#endif int wolfIP_poll(struct wolfIP *s, uint64_t now); +/* Called when a socket call queues work for the next wolfIP_poll(); never + * from inside it. Set it under the lock that serializes wolfIP calls. */ +typedef void (*wolfIP_wake_cb)(void *arg); +void wolfIP_set_wake_cb(struct wolfIP *s, wolfIP_wake_cb cb, void *arg); void wolfIP_recv(struct wolfIP *s, void *buf, uint32_t len); void wolfIP_recv_ex(struct wolfIP *s, unsigned int if_idx, void *buf, uint32_t len); void wolfIP_ipconfig_set(struct wolfIP *s, ip4 ip, ip4 mask, ip4 gw);