diff --git a/.github/workflows/broker-check.yml b/.github/workflows/broker-check.yml index a6c60e71b..848a462da 100644 --- a/.github/workflows/broker-check.yml +++ b/.github/workflows/broker-check.yml @@ -19,21 +19,35 @@ jobs: - name: "Broker default (dynamic alloc)" cflags: "" wolfmqtt_opts: "--enable-broker" + wolfssl_opts: "--enable-enckeys" - name: "Broker static memory" cflags: "-DWOLFMQTT_STATIC_MEMORY" wolfmqtt_opts: "--enable-broker" + wolfssl_opts: "--enable-enckeys" - name: "Broker with TLS" cflags: "" wolfmqtt_opts: "--enable-broker --enable-tls" + wolfssl_opts: "--enable-enckeys" - name: "Broker with TLS (static memory)" cflags: "-DWOLFMQTT_STATIC_MEMORY" wolfmqtt_opts: "--enable-broker --enable-tls" + wolfssl_opts: "--enable-enckeys" + - name: "Broker with DTLS" + cflags: "" + wolfmqtt_opts: "--enable-broker --enable-tls --enable-dtls" + wolfssl_opts: "--enable-enckeys --enable-dtls --enable-dtls13" + - name: "Broker with DTLS (static memory)" + cflags: "-DWOLFMQTT_STATIC_MEMORY" + wolfmqtt_opts: "--enable-broker --enable-tls --enable-dtls" + wolfssl_opts: "--enable-enckeys --enable-dtls --enable-dtls13" - name: "Broker no logging" cflags: "" wolfmqtt_opts: "--enable-broker --disable-broker-log" + wolfssl_opts: "--enable-enckeys" - name: "Broker minimal (no log, no retained, no will, no wildcards, no auth)" cflags: "" wolfmqtt_opts: "--enable-broker --disable-broker-log --disable-broker-retained --disable-broker-will --disable-broker-wildcards --disable-broker-auth" + wolfssl_opts: "--enable-enckeys" steps: - name: Install dependencies @@ -51,7 +65,7 @@ jobs: run: ./autogen.sh - name: wolfssl configure working-directory: ./wolfssl - run: ./configure --enable-enckeys + run: ./configure ${{ matrix.wolfssl_opts }} - name: wolfssl make working-directory: ./wolfssl run: make diff --git a/CMakeLists.txt b/CMakeLists.txt index 3a2e4d4fb..43b4adaaa 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -91,6 +91,18 @@ if (WOLFMQTT_TLS) list(APPEND WOLFMQTT_DEFINITIONS "-DENABLE_MQTT_TLS") endif() +add_option("WOLFMQTT_DTLS" + "Enable DTLS support (requires wolfSSL DTLS)" + "no" "yes;no") +if (WOLFMQTT_DTLS) + if (NOT WOLFMQTT_TLS) + message(WARNING "DTLS requires TLS support -- disabling DTLS") + set(WOLFMQTT_DTLS "no" CACHE STRING "" FORCE) + else() + list(APPEND WOLFMQTT_DEFINITIONS "-DENABLE_MQTT_DTLS") + endif() +endif() + add_option("WOLFMQTT_SN" "Enable MQTT-SN support" "no" "yes;no") diff --git a/configure.ac b/configure.ac index 9aed67332..cd76208eb 100644 --- a/configure.ac +++ b/configure.ac @@ -158,7 +158,6 @@ then test "$enable_v5" = "" && enable_v5=yes test "$enable_discb" = "" && enable_discb=yes test "$enable_mt" = "" && enable_mt=yes - test "$enable_" = "" && enable_=yes fi # TLS Support with wolfSSL @@ -176,6 +175,24 @@ AC_CHECK_LIB([wolfssl],[wolfCrypt_Init],,[AC_MSG_ERROR([libwolfssl is required a fi +# DTLS Support (requires wolfSSL with DTLS) +AC_ARG_ENABLE([dtls], + [AS_HELP_STRING([--enable-dtls],[Enable DTLS support (requires wolfSSL DTLS) (default: disabled)])], + [ ENABLED_DTLS=$enableval ], + [ ENABLED_DTLS=no ] + ) + +if test "x$ENABLED_DTLS" = "xyes" +then + if test "x$ENABLED_TLS" != "xyes" + then + AC_MSG_WARN([DTLS requires TLS support -- disabling DTLS]) + ENABLED_DTLS=no + else + AM_CFLAGS="$AM_CFLAGS -DENABLE_MQTT_DTLS" + fi +fi + # Non-Blocking support AC_ARG_ENABLE([nonblock], @@ -480,6 +497,7 @@ AM_CONDITIONAL([BUILD_NONBLOCK], [test "x$ENABLED_NONBLOCK" = "xyes"]) AM_CONDITIONAL([BUILD_MULTITHREAD], [test "x$ENABLED_MULTITHREAD" = "xyes"]) AM_CONDITIONAL([BUILD_WEBSOCKET], [test "x$ENABLED_WEBSOCKET" = "xyes"]) AM_CONDITIONAL([BUILD_BROKER], [test "x$ENABLED_BROKER" = "xyes"]) +AM_CONDITIONAL([BUILD_DTLS], [test "x$ENABLED_DTLS" = "xyes"]) diff --git a/examples/mqttclient/mqttclient.c b/examples/mqttclient/mqttclient.c index 1fea238b7..bff152e38 100644 --- a/examples/mqttclient/mqttclient.c +++ b/examples/mqttclient/mqttclient.c @@ -268,6 +268,12 @@ int mqttclient_test(MQTTCtx *mqttCtx) MqttSocket_Connect for use by mqtt_tls_verify_cb */ mqttCtx->client.ctx = mqttCtx; +#ifdef ENABLE_MQTT_DTLS + if (mqttCtx->use_dtls) { + MqttClient_Flags(&mqttCtx->client, 0, MQTT_CLIENT_FLAG_IS_DTLS); + } +#endif + #ifdef WOLFMQTT_DISCONNECT_CB /* setup disconnect callback */ rc = MqttClient_SetDisconnectCallback(&mqttCtx->client, diff --git a/examples/mqttexample.c b/examples/mqttexample.c index f23ad715a..ac9496868 100644 --- a/examples/mqttexample.c +++ b/examples/mqttexample.c @@ -236,6 +236,9 @@ void mqtt_show_usage(MQTTCtx* mqttCtx) #endif /* !ENABLE_MQTT_CURL */ PRINTF("-p Port to connect on, default: %d", MQTT_DEFAULT_PORT); +#endif +#ifdef ENABLE_MQTT_DTLS + PRINTF("-D Enable DTLS (UDP)"); #endif PRINTF("-q Qos Level 0-2, default: %d", mqttCtx->qos); @@ -312,9 +315,14 @@ int mqtt_parse_args(MQTTCtx* mqttCtx, int argc, char** argv) #else #define MQTT_V5_ARGS "" #endif + #ifdef ENABLE_MQTT_DTLS + #define MQTT_DTLS_ARGS "D" + #else + #define MQTT_DTLS_ARGS "" + #endif while ((rc = mygetopt(argc, argv, "?h:p:q:sk:i:lu:w:m:n:C:Tf:rtd" \ - MQTT_TLS_ARGS MQTT_V5_ARGS)) != -1) { + MQTT_TLS_ARGS MQTT_V5_ARGS MQTT_DTLS_ARGS)) != -1) { switch ((char)rc) { case '?' : mqtt_show_usage(mqttCtx); @@ -394,6 +402,13 @@ int mqtt_parse_args(MQTTCtx* mqttCtx, int argc, char** argv) mqttCtx->debug_on = 1; break; + #ifdef ENABLE_MQTT_DTLS + case 'D': + mqttCtx->use_tls = 1; + mqttCtx->use_dtls = 1; + break; + #endif + #ifdef ENABLE_MQTT_TLS case 'A': mqttCtx->ca_file = myoptarg; @@ -640,6 +655,16 @@ int mqtt_tls_cb(MqttClient* client) /* Use highest available and allow downgrade. If wolfSSL is built with * old TLS support, it is possible for a server to force a downgrade to * an insecure version. */ +#ifdef ENABLE_MQTT_DTLS + if (MqttClient_Flags(client, 0, 0) & MQTT_CLIENT_FLAG_IS_DTLS) { + #ifdef WOLFSSL_DTLS13 + client->tls.ctx = wolfSSL_CTX_new(wolfDTLSv1_3_client_method()); + #else + client->tls.ctx = wolfSSL_CTX_new(wolfDTLSv1_2_client_method()); + #endif + } + else +#endif client->tls.ctx = wolfSSL_CTX_new(wolfSSLv23_client_method()); if (client->tls.ctx) { wolfSSL_CTX_set_verify(client->tls.ctx, WOLFSSL_VERIFY_PEER, @@ -766,7 +791,11 @@ int mqtt_dtls_cb(MqttClient* client) { int rc = WOLFSSL_FAILURE; SocketContext * sock = (SocketContext *)client->net->context; +#ifdef WOLFSSL_DTLS13 + client->tls.ctx = wolfSSL_CTX_new(wolfDTLSv1_3_client_method()); +#else client->tls.ctx = wolfSSL_CTX_new(wolfDTLSv1_2_client_method()); +#endif if (client->tls.ctx) { wolfSSL_CTX_set_verify(client->tls.ctx, WOLFSSL_VERIFY_PEER, mqtt_tls_verify_cb); diff --git a/examples/mqttexample.h b/examples/mqttexample.h index 6b0289a3d..c76f7cb43 100644 --- a/examples/mqttexample.h +++ b/examples/mqttexample.h @@ -170,6 +170,9 @@ typedef struct _MQTTCtx { byte *tx_buf, *rx_buf; int return_code; int use_tls; +#ifdef ENABLE_MQTT_DTLS + int use_dtls; +#endif int retain; int enable_lwt; #ifdef WOLFMQTT_V5 diff --git a/examples/mqttnet.c b/examples/mqttnet.c index 4786931e0..305825c14 100644 --- a/examples/mqttnet.c +++ b/examples/mqttnet.c @@ -1217,12 +1217,20 @@ static int NetConnect(void *context, const char* host, word16 port, { SocketContext *sock = (SocketContext*)context; int type = SOCK_STREAM; + int proto = IPPROTO_TCP; int rc = -1; SOERROR_T so_error = 0; struct addrinfo *result = NULL; struct addrinfo hints; MQTTCtx* mqttCtx = sock->mqttCtx; +#ifdef ENABLE_MQTT_DTLS + if (mqttCtx->use_dtls) { + type = SOCK_DGRAM; + proto = IPPROTO_UDP; + } +#endif + /* Get address information for host and locate IPv4 */ switch(sock->stat) { case SOCK_BEGIN: @@ -1234,8 +1242,8 @@ static int NetConnect(void *context, const char* host, word16 port, XMEMSET(&hints, 0, sizeof(hints)); hints.ai_family = AF_INET; - hints.ai_socktype = SOCK_STREAM; - hints.ai_protocol = IPPROTO_TCP; + hints.ai_socktype = type; + hints.ai_protocol = proto; XMEMSET(&sock->addr, 0, sizeof(sock->addr)); sock->addr.sin_family = AF_INET; diff --git a/scripts/broker.test b/scripts/broker.test index 17cbf1583..e8dd24e3b 100755 --- a/scripts/broker.test +++ b/scripts/broker.test @@ -320,6 +320,55 @@ else echo "SKIP: TLS tests (broker not built with TLS support)" fi +# --- DTLS Tests --- +has_dtls=no +echo "$broker_features" | grep -q "dtls" && has_dtls=yes + +if [ "$has_dtls" = "yes" ]; then + # --- Test 17: DTLS pub/sub --- + echo "" + echo "--- Test 17: DTLS pub/sub ---" + # DTLS uses UDP; nc -z only checks TCP, so use sleep for startup + if [ $broker_pid != $no_pid ]; then + kill $broker_pid 2>/dev/null + wait $broker_pid 2>/dev/null || true + broker_pid=$no_pid + fi + generate_port + ./$broker_bin -D \ + -c scripts/broker_test/server-cert.pem \ + -K scripts/broker_test/server-key.pem \ + -p $port & + broker_pid=$! + sleep 2 + + ./$client_bin -T -h 127.0.0.1 -p $port -n "test/dtls" -D \ + -A scripts/broker_test/ca-cert.pem -C 5000 \ + >"${TMP_DIR}/t17.log" 2>&1 + if [ $? -eq 0 ]; then + echo "PASS: DTLS pub/sub" + else + echo "FAIL: DTLS pub/sub" + FAIL=1 + fi + + # --- Test 18: DTLS QoS 2 --- + echo "" + echo "--- Test 18: DTLS QoS 2 ---" + ./$client_bin -T -h 127.0.0.1 -p $port -n "test/dtls_qos2" -q 2 -D \ + -A scripts/broker_test/ca-cert.pem -C 5000 \ + >"${TMP_DIR}/t18.log" 2>&1 + if [ $? -eq 0 ]; then + echo "PASS: DTLS QoS 2" + else + echo "FAIL: DTLS QoS 2" + FAIL=1 + fi +else + echo "" + echo "SKIP: DTLS tests (broker not built with DTLS support)" +fi + # --- Test 11: Multi-client QoS 2 (separate pub/sub) --- echo "" echo "--- Test 11: Multi-client QoS 2 ---" diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index 87ed88e4f..e65fb3ea8 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -330,6 +330,106 @@ static int BrokerPosix_Close(void* ctx, BROKER_SOCKET_T sock) return MQTT_CODE_SUCCESS; } +#ifdef ENABLE_MQTT_DTLS +static int BrokerPosix_ListenUDP(void* ctx, BROKER_SOCKET_T* sock, + word16 port, int backlog) +{ + struct sockaddr_in addr; + int opt = 1; + BROKER_SOCKET_T fd; + (void)backlog; + + fd = socket(AF_INET, SOCK_DGRAM, 0); + if (fd < 0) { + WBLOG_ERR((MqttBroker*)ctx, "broker: UDP socket failed (%d)", errno); + return MQTT_CODE_ERROR_NETWORK; + } + + (void)setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)); +#ifdef SO_REUSEPORT + (void)setsockopt(fd, SOL_SOCKET, SO_REUSEPORT, &opt, sizeof(opt)); +#endif + + if (BrokerPosix_SetNonBlocking(fd) != MQTT_CODE_SUCCESS) { + WBLOG_ERR((MqttBroker*)ctx, "broker: UDP set nonblocking failed (%d)", + errno); + close(fd); + return MQTT_CODE_ERROR_SYSTEM; + } + + XMEMSET(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_addr.s_addr = htonl(INADDR_ANY); + addr.sin_port = htons(port); + + if (bind(fd, (struct sockaddr*)&addr, sizeof(addr)) < 0) { + WBLOG_ERR((MqttBroker*)ctx, "broker: UDP bind failed (%d)", errno); + close(fd); + return MQTT_CODE_ERROR_NETWORK; + } + /* No listen() for UDP */ + + *sock = fd; + return MQTT_CODE_SUCCESS; +} + +static int BrokerPosix_AcceptUDP(void* ctx, BROKER_SOCKET_T listen_sock, + BROKER_SOCKET_T* client_sock) +{ + MqttBroker* broker = (MqttBroker*)ctx; + struct sockaddr_in peer_addr; + socklen_t addr_len = sizeof(peer_addr); + char buf[1]; + int rc; + + /* Peek to see if data is available and get peer address */ + rc = (int)recvfrom(listen_sock, buf, sizeof(buf), + MSG_PEEK | MSG_DONTWAIT, + (struct sockaddr*)&peer_addr, &addr_len); + if (rc <= 0) { + return MQTT_CODE_CONTINUE; /* No data / no new peer */ + } + + /* Check if we already have a client for this peer */ + { + int found = 0; + #ifdef WOLFMQTT_STATIC_MEMORY + int i; + for (i = 0; i < BROKER_MAX_CLIENTS; i++) { + BrokerClient* bc = &broker->clients[i]; + if (bc->in_use && bc->peer_addr_len == (int)addr_len && + XMEMCMP(bc->peer_addr, &peer_addr, addr_len) == 0) { + found = 1; + break; + } + } + #else + BrokerClient* bc = broker->clients; + while (bc != NULL) { + if (bc->peer_addr_len == (int)addr_len && + XMEMCMP(bc->peer_addr, &peer_addr, addr_len) == 0) { + found = 1; + break; + } + bc = bc->next; + } + #endif + if (found) { + return MQTT_CODE_CONTINUE; /* Existing client has data */ + } + } + + /* New peer: save address for BrokerClient_Add to copy */ + XMEMCPY(broker->pending_peer_addr, &peer_addr, addr_len); + broker->pending_peer_addr_len = (int)addr_len; + + /* All DTLS clients share the listen socket. Per-client demuxing is + * handled by peek-based read callbacks (BrokerNetReadDTLS). */ + *client_sock = listen_sock; + return MQTT_CODE_SUCCESS; +} +#endif /* ENABLE_MQTT_DTLS */ + int MqttBrokerNet_Init(MqttBrokerNet* net) { if (net == NULL) { @@ -357,16 +457,38 @@ static int BrokerTls_Init(MqttBroker* broker) rc = MQTT_CODE_ERROR_BAD_ARG; } - /* Select TLS method based on version preference */ + /* Select TLS/DTLS method based on version preference */ if (rc == WOLFSSL_SUCCESS) { - if (broker->tls_version == 12) { - ctx = wolfSSL_CTX_new(wolfTLSv1_2_server_method()); - } - else if (broker->tls_version == 13) { - ctx = wolfSSL_CTX_new(wolfTLSv1_3_server_method()); + #ifdef ENABLE_MQTT_DTLS + if (broker->use_dtls) { + if (broker->tls_version == 12) { + ctx = wolfSSL_CTX_new(wolfDTLSv1_2_server_method()); + } + #ifdef WOLFSSL_DTLS13 + else if (broker->tls_version == 13) { + ctx = wolfSSL_CTX_new(wolfDTLSv1_3_server_method()); + } + else { + ctx = wolfSSL_CTX_new(wolfDTLSv1_3_server_method()); + } + #else + else { + ctx = wolfSSL_CTX_new(wolfDTLSv1_2_server_method()); + } + #endif } - else { - ctx = wolfSSL_CTX_new(wolfSSLv23_server_method()); + else + #endif /* ENABLE_MQTT_DTLS */ + { + if (broker->tls_version == 12) { + ctx = wolfSSL_CTX_new(wolfTLSv1_2_server_method()); + } + else if (broker->tls_version == 13) { + ctx = wolfSSL_CTX_new(wolfTLSv1_3_server_method()); + } + else { + ctx = wolfSSL_CTX_new(wolfSSLv23_server_method()); + } } if (ctx == NULL) { WBLOG_ERR(broker, "broker: wolfSSL_CTX_new failed"); @@ -500,6 +622,101 @@ static int BrokerNetDisconnect(void* context) return MQTT_CODE_SUCCESS; } +/* -------------------------------------------------------------------------- */ +/* DTLS per-client IO callbacks (shared UDP socket with peek-based demux) */ +/* -------------------------------------------------------------------------- */ +#if defined(ENABLE_MQTT_DTLS) && !defined(WOLFMQTT_BROKER_CUSTOM_NET) +static int BrokerNetReadDTLS(void* context, byte* buf, int buf_len, + int timeout_ms) +{ + BrokerClient* bc = (BrokerClient*)context; + struct sockaddr_in peer_addr; + socklen_t addr_len = sizeof(peer_addr); + fd_set rfds; + struct timeval tv; + int rc; + + if (bc == NULL || bc->broker == NULL || buf == NULL || buf_len <= 0) { + return MQTT_CODE_ERROR_BAD_ARG; + } + + /* Wait for data on the shared listen socket */ + FD_ZERO(&rfds); + FD_SET(bc->sock, &rfds); + tv.tv_sec = timeout_ms / 1000; + tv.tv_usec = (timeout_ms % 1000) * 1000; + rc = select(bc->sock + 1, &rfds, NULL, NULL, &tv); + if (rc == 0) { + return MQTT_CODE_ERROR_TIMEOUT; + } + if (rc < 0) { + return MQTT_CODE_ERROR_NETWORK; + } + + /* Peek to check source address before consuming */ + addr_len = sizeof(peer_addr); + rc = (int)recvfrom(bc->sock, buf, (size_t)buf_len, + MSG_PEEK | MSG_DONTWAIT, + (struct sockaddr*)&peer_addr, &addr_len); + if (rc <= 0) { + return MQTT_CODE_CONTINUE; + } + + /* Verify the datagram is from this client's peer */ + if ((int)addr_len != bc->peer_addr_len || + XMEMCMP(&peer_addr, bc->peer_addr, addr_len) != 0) { + /* Not for this client - leave it on the socket */ + return MQTT_CODE_ERROR_TIMEOUT; + } + + /* Source matches: consume the datagram */ + addr_len = sizeof(peer_addr); + rc = (int)recvfrom(bc->sock, buf, (size_t)buf_len, 0, + (struct sockaddr*)&peer_addr, &addr_len); + if (rc <= 0) { + return MQTT_CODE_ERROR_NETWORK; + } + return rc; +} + +static int BrokerNetWriteDTLS(void* context, const byte* buf, int buf_len, + int timeout_ms) +{ + BrokerClient* bc = (BrokerClient*)context; + int rc; + (void)timeout_ms; + + if (bc == NULL || bc->broker == NULL || buf == NULL || buf_len <= 0) { + return MQTT_CODE_ERROR_BAD_ARG; + } + + rc = (int)sendto(bc->sock, buf, (size_t)buf_len, 0, + (struct sockaddr*)bc->peer_addr, bc->peer_addr_len); + if (rc <= 0) { + if (rc < 0 && (errno == EWOULDBLOCK || errno == EAGAIN)) { + return MQTT_CODE_CONTINUE; + } + WBLOG_ERR(bc->broker, "broker: DTLS sendto error sock=%d errno=%d", + (int)bc->sock, errno); + return MQTT_CODE_ERROR_NETWORK; + } + return rc; +} + +static int BrokerNetDisconnectDTLS(void* context) +{ + BrokerClient* bc = (BrokerClient*)context; + if (bc != NULL && bc->broker != NULL && + bc->sock != BROKER_SOCKET_INVALID) { + WBLOG_INFO(bc->broker, "broker: DTLS disconnect sock=%d", + (int)bc->sock); + /* Don't close the shared listen socket */ + bc->sock = BROKER_SOCKET_INVALID; + } + return MQTT_CODE_SUCCESS; +} +#endif /* ENABLE_MQTT_DTLS && !WOLFMQTT_BROKER_CUSTOM_NET */ + /* -------------------------------------------------------------------------- */ /* Client management */ /* -------------------------------------------------------------------------- */ @@ -508,10 +725,9 @@ static void BrokerClient_Free(BrokerClient* bc) if (bc == NULL) { return; } - (void)BrokerNetDisconnect(bc); #ifdef ENABLE_MQTT_TLS if (bc->client.tls.ssl) { - /* Only send close_notify if handshake completed successfully */ + /* Send close_notify before disconnect so the socket is still valid */ if (bc->tls_handshake_done) { wolfSSL_shutdown(bc->client.tls.ssl); } @@ -519,6 +735,13 @@ static void BrokerClient_Free(BrokerClient* bc) bc->client.tls.ssl = NULL; } #endif + /* Use the per-client disconnect (DTLS override skips socket close) */ + if (bc->net.disconnect) { + (void)bc->net.disconnect(bc->net.context); + } + else { + (void)BrokerNetDisconnect(bc); + } MqttClient_DeInit(&bc->client); #ifdef WOLFMQTT_STATIC_MEMORY XMEMSET(bc, 0, sizeof(*bc)); @@ -626,6 +849,26 @@ static BrokerClient* BrokerClient_Add(MqttBroker* broker, wolfSSL_SetIOReadCtx(bc->client.tls.ssl, &bc->client); wolfSSL_SetIOWriteCtx(bc->client.tls.ssl, &bc->client); MqttClient_Flags(&bc->client, 0, MQTT_CLIENT_FLAG_IS_TLS); + #ifdef ENABLE_MQTT_DTLS + if (broker->use_dtls) { + MqttClient_Flags(&bc->client, 0, + MQTT_CLIENT_FLAG_IS_DTLS); + /* Copy peer address from pending accept */ + XMEMCPY(bc->peer_addr, broker->pending_peer_addr, + broker->pending_peer_addr_len); + bc->peer_addr_len = broker->pending_peer_addr_len; + /* Tell wolfSSL which peer this object serves */ + wolfSSL_dtls_set_peer(bc->client.tls.ssl, + bc->peer_addr, bc->peer_addr_len); + #ifndef WOLFMQTT_BROKER_CUSTOM_NET + /* Use DTLS-specific IO (peek-based demux on + * shared UDP socket) */ + bc->net.read = BrokerNetReadDTLS; + bc->net.write = BrokerNetWriteDTLS; + bc->net.disconnect = BrokerNetDisconnectDTLS; + #endif + } + #endif bc->tls_handshake_done = 0; } } @@ -2709,6 +2952,16 @@ int MqttBroker_Run(MqttBroker* broker) return MQTT_CODE_ERROR_BAD_ARG; } + /* Swap to UDP network callbacks for DTLS mode */ +#ifdef ENABLE_MQTT_DTLS + if (broker->use_dtls) { + #if !defined(WOLFMQTT_WOLFIP) && !defined(WOLFMQTT_BROKER_CUSTOM_NET) + broker->net.listen = BrokerPosix_ListenUDP; + broker->net.accept = BrokerPosix_AcceptUDP; + #endif + } +#endif + /* Start listening */ rc = broker->net.listen(broker->net.ctx, &broker->listen_sock, broker->port, BROKER_LISTEN_BACKLOG); @@ -2724,12 +2977,23 @@ int MqttBroker_Run(MqttBroker* broker) WBLOG_ERR(broker, "broker: TLS init failed rc=%d", rc); return rc; } - WBLOG_INFO(broker, "broker: listening on port %d (TLS)", broker->port); + #ifdef ENABLE_MQTT_DTLS + if (broker->use_dtls) { + WBLOG_INFO(broker, "broker: listening on port %d (DTLS)", + broker->port); + } + else + #endif + { + WBLOG_INFO(broker, "broker: listening on port %d (TLS)", + broker->port); + } } else #endif { - WBLOG_INFO(broker, "broker: listening on port %d (no TLS)", broker->port); + WBLOG_INFO(broker, "broker: listening on port %d (no TLS)", + broker->port); } #ifdef WOLFMQTT_BROKER_AUTH if (broker->auth_user || broker->auth_pass) { @@ -2814,6 +3078,9 @@ static void BrokerUsage(const char* prog) #endif #ifdef ENABLE_MQTT_TLS " [-t] [-V ver] [-c cert] [-K key] [-A ca]" +#endif +#ifdef ENABLE_MQTT_DTLS + " [-D]" #endif , prog); PRINTF(" -v Log level: 1=error, 2=info (default), 3=debug"); @@ -2823,6 +3090,9 @@ static void BrokerUsage(const char* prog) PRINTF(" -c Server certificate file (PEM)"); PRINTF(" -K Server private key file (PEM)"); PRINTF(" -A CA certificate for mutual TLS (PEM)"); +#endif +#ifdef ENABLE_MQTT_DTLS + PRINTF(" -D Enable DTLS (UDP) instead of TLS (TCP)"); #endif PRINTF("Features:" #ifdef WOLFMQTT_BROKER_RETAINED @@ -2839,6 +3109,9 @@ static void BrokerUsage(const char* prog) #endif #ifdef ENABLE_MQTT_TLS " tls" +#endif +#ifdef ENABLE_MQTT_DTLS + " dtls" #endif ); } @@ -2919,6 +3192,15 @@ int wolfmqtt_broker(int argc, char** argv) else if (XSTRCMP(argv[i], "-A") == 0 && i + 1 < argc) { broker.tls_ca = argv[++i]; } +#endif +#ifdef ENABLE_MQTT_DTLS + else if (XSTRCMP(argv[i], "-D") == 0) { + broker.use_dtls = 1; + broker.use_tls = 1; /* DTLS implies TLS */ + if (broker.port == MQTT_DEFAULT_PORT) { + broker.port = MQTT_SECURE_PORT; + } + } #endif else if (XSTRCMP(argv[i], "-h") == 0) { BrokerUsage(argv[0]); diff --git a/src/mqtt_socket.c b/src/mqtt_socket.c index a39401558..c54251b20 100644 --- a/src/mqtt_socket.c +++ b/src/mqtt_socket.c @@ -437,7 +437,13 @@ int MqttSocket_Connect(MqttClient *client, const char* host, word16 port, } #ifdef WOLFSSL_DTLS else { - client->tls.ctx = wolfSSL_CTX_new(wolfDTLSv1_2_client_method()); + #ifdef WOLFSSL_DTLS13 + client->tls.ctx = wolfSSL_CTX_new( + wolfDTLSv1_3_client_method()); + #else + client->tls.ctx = wolfSSL_CTX_new( + wolfDTLSv1_2_client_method()); + #endif } #endif if (client->tls.ctx == NULL) { diff --git a/wolfmqtt/mqtt_broker.h b/wolfmqtt/mqtt_broker.h index 35273d77a..b68b688bb 100644 --- a/wolfmqtt/mqtt_broker.h +++ b/wolfmqtt/mqtt_broker.h @@ -72,6 +72,11 @@ #ifndef BROKER_TX_BUF_SZ #define BROKER_TX_BUF_SZ 4096 #endif +#ifdef ENABLE_MQTT_DTLS +#ifndef BROKER_PEER_ADDR_SZ + #define BROKER_PEER_ADDR_SZ 16 /* sizeof(struct sockaddr_in) IPv4 */ +#endif +#endif #ifndef BROKER_TIMEOUT_MS #define BROKER_TIMEOUT_MS 1000 #endif @@ -208,6 +213,10 @@ typedef struct BrokerClient { #ifdef ENABLE_MQTT_TLS byte tls_handshake_done; #endif +#ifdef ENABLE_MQTT_DTLS + byte peer_addr[BROKER_PEER_ADDR_SZ]; + int peer_addr_len; +#endif } BrokerClient; /* -------------------------------------------------------------------------- */ @@ -289,6 +298,11 @@ typedef struct MqttBroker { const char* tls_ca; /* CA cert for mutual auth (optional) */ byte use_tls; byte tls_version; /* 0=auto (v23), 12=TLS 1.2, 13=TLS 1.3 */ +#ifdef ENABLE_MQTT_DTLS + byte use_dtls; /* Use DTLS (UDP) instead of TLS (TCP) */ + byte pending_peer_addr[BROKER_PEER_ADDR_SZ]; + int pending_peer_addr_len; +#endif #endif #ifdef WOLFMQTT_STATIC_MEMORY BrokerClient clients[BROKER_MAX_CLIENTS];