From e80f46a9931daa26773307025dba380a7a74179c Mon Sep 17 00:00:00 2001 From: Yakun Xu Date: Wed, 17 Jul 2019 00:03:43 +0800 Subject: [PATCH] [posix-sim] use multicast to simulate radio (#3993) This commit uses multicast to simulate radio transmissions. This avoids sending a frame multiple times and avoids duplicated packets in captures. This implementation uses two separate sockets for tx and rx. The current implementation of thread-cert requires the sniffer to report the sender of each packet. This issue can be addressed by stateful tracking state of each node, but left for future work. --- examples/platforms/posix/platform-posix.h | 6 +- examples/platforms/posix/radio.c | 162 ++++++++++++------ examples/platforms/posix/sim/platform-sim.c | 2 +- examples/platforms/posix/system.c | 30 ++-- tests/scripts/thread-cert/config.py | 4 +- tests/scripts/thread-cert/sniffer.py | 6 +- .../scripts/thread-cert/sniffer_transport.py | 17 +- 7 files changed, 143 insertions(+), 84 deletions(-) diff --git a/examples/platforms/posix/platform-posix.h b/examples/platforms/posix/platform-posix.h index 244f8e490..4829bf0b7 100644 --- a/examples/platforms/posix/platform-posix.h +++ b/examples/platforms/posix/platform-posix.h @@ -170,10 +170,12 @@ void platformRadioUpdateFdSet(fd_set *aReadFdSet, fd_set *aWriteFdSet, int *aMax /** * This function performs radio driver processing. * - * @param[in] aInstance The OpenThread instance structure. + * @param[in] aInstance The OpenThread instance structure. + * @param[in] aReadFdSet A pointer to the read file descriptors. + * @param[in] aWriteFdSet A pointer to the write file descriptors. * */ -void platformRadioProcess(otInstance *aInstance); +void platformRadioProcess(otInstance *aInstance, const fd_set *aReadFdSet, const fd_set *aWriteFdSet); /** * This function initializes the random number service used by OpenThread. diff --git a/examples/platforms/posix/radio.c b/examples/platforms/posix/radio.c index 5c0c63c66..b42755826 100644 --- a/examples/platforms/posix/radio.c +++ b/examples/platforms/posix/radio.c @@ -38,6 +38,9 @@ #include "utils/code_utils.h" +// The IPv4 group for receiving packets of radio simulation +#define OT_RADIO_GROUP "224.0.0.116" + enum { IEEE802154_MIN_LENGTH = 5, @@ -94,7 +97,8 @@ enum extern int sSockFd; extern uint16_t sPortOffset; #else -static int sSockFd; +static int sTxFd = -1; +static int sRxFd = -1; static uint16_t sPortOffset = 0; #endif @@ -394,11 +398,77 @@ void otPlatRadioSetPromiscuous(otInstance *aInstance, bool aEnable) sPromiscuous = aEnable; } +#if OPENTHREAD_POSIX_VIRTUAL_TIME == 0 +static void initFds(void) +{ + int fd; + int one = 1; + struct sockaddr_in sockaddr; + + memset(&sockaddr, 0, sizeof(sockaddr)); + + otEXPECT_ACTION((fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP)) != -1, perror("socket(sTxFd)")); + + sockaddr.sin_family = AF_INET; + sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + gNodeId)); + sockaddr.sin_addr.s_addr = inet_addr("127.0.0.1"); + + otEXPECT_ACTION(setsockopt(fd, IPPROTO_IP, IP_MULTICAST_IF, &sockaddr.sin_addr, sizeof(sockaddr.sin_addr)) != -1, + perror("setsockopt(sTxFd, IP_MULTICAST_IF)")); + + otEXPECT_ACTION(setsockopt(fd, IPPROTO_IP, IP_MULTICAST_LOOP, &one, sizeof(one)) != -1, + perror("setsockopt(sRxFd, IP_MULTICAST_LOOP)")); + + otEXPECT_ACTION(bind(fd, (struct sockaddr *)&sockaddr, sizeof(sockaddr)) != -1, perror("bind(sTxFd)")); + + // Tx fd is successfully initialized. + sTxFd = fd; + + otEXPECT_ACTION((fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP)) != -1, perror("socket(sRxFd)")); + + otEXPECT_ACTION(setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &one, sizeof(one)) != -1, + perror("setsockopt(sRxFd, SO_REUSEADDR)")); + otEXPECT_ACTION(setsockopt(fd, SOL_SOCKET, SO_REUSEPORT, &one, sizeof(one)) != -1, + perror("setsockopt(sRxFd, SO_REUSEPORT)")); + + { + struct ip_mreqn mreq; + + memset(&mreq, 0, sizeof(mreq)); + inet_pton(AF_INET, OT_RADIO_GROUP, &mreq.imr_multiaddr); + + // Always use loopback device to send simulation packets. + mreq.imr_address.s_addr = inet_addr("127.0.0.1"); + + otEXPECT_ACTION(setsockopt(fd, IPPROTO_IP, IP_MULTICAST_IF, &mreq.imr_address, sizeof(mreq.imr_address)) != -1, + perror("setsockopt(sRxFd, IP_MULTICAST_IF)")); + otEXPECT_ACTION(setsockopt(fd, IPPROTO_IP, IP_ADD_MEMBERSHIP, &mreq, sizeof(mreq)) != -1, + perror("setsockopt(sRxFd, IP_ADD_MEMBERSHIP)")); + } + + sockaddr.sin_family = AF_INET; + sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + WELLKNOWN_NODE_ID)); + sockaddr.sin_addr.s_addr = inet_addr(OT_RADIO_GROUP); + + otEXPECT_ACTION(bind(fd, (struct sockaddr *)&sockaddr, sizeof(sockaddr)) != -1, perror("bind(sRxFd)")); + + // Rx fd is successfully initialized. + sRxFd = fd; + +exit: + if (sRxFd == -1 || sTxFd == -1) + { + exit(EXIT_FAILURE); + } +} +#endif // OPENTHREAD_POSIX_VIRTUAL_TIME == 0 + void platformRadioInit(void) { #if OPENTHREAD_POSIX_VIRTUAL_TIME == 0 struct sockaddr_in sockaddr; char * offset; + memset(&sockaddr, 0, sizeof(sockaddr)); sockaddr.sin_family = AF_INET; @@ -419,31 +489,9 @@ void platformRadioInit(void) sPortOffset *= WELLKNOWN_NODE_ID; } - if (sPromiscuous) - { - sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + WELLKNOWN_NODE_ID)); - } - else - { - sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + gNodeId)); - } - - sockaddr.sin_addr.s_addr = INADDR_ANY; - - sSockFd = (int)socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); - - if (sSockFd == -1) - { - perror("socket"); - exit(EXIT_FAILURE); - } - - if (bind(sSockFd, (struct sockaddr *)&sockaddr, sizeof(sockaddr)) == -1) - { - perror("bind"); - exit(EXIT_FAILURE); - } + initFds(); #endif // OPENTHREAD_POSIX_VIRTUAL_TIME == 0 + sReceiveFrame.mPsdu = sReceiveMessage.mPsdu; sTransmitFrame.mPsdu = sTransmitMessage.mPsdu; sAckFrame.mPsdu = sAckMessage.mPsdu; @@ -735,44 +783,53 @@ void platformRadioUpdateFdSet(fd_set *aReadFdSet, fd_set *aWriteFdSet, int *aMax { if (aReadFdSet != NULL && (sState != OT_RADIO_STATE_TRANSMIT || sTxWait)) { - FD_SET(sSockFd, aReadFdSet); + FD_SET(sRxFd, aReadFdSet); - if (aMaxFd != NULL && *aMaxFd < sSockFd) + if (aMaxFd != NULL && *aMaxFd < sRxFd) { - *aMaxFd = sSockFd; + *aMaxFd = sRxFd; } } if (aWriteFdSet != NULL && platformRadioIsTransmitPending()) { - FD_SET(sSockFd, aWriteFdSet); + FD_SET(sTxFd, aWriteFdSet); - if (aMaxFd != NULL && *aMaxFd < sSockFd) + if (aMaxFd != NULL && *aMaxFd < sTxFd) { - *aMaxFd = sSockFd; + *aMaxFd = sTxFd; } } } +// no need to close in virtual time mode. void platformRadioDeinit(void) { - close(sSockFd); + if (sRxFd != -1) + { + close(sRxFd); + } + + if (sTxFd != -1) + { + close(sTxFd); + } } #endif // OPENTHREAD_POSIX_VIRTUAL_TIME -void platformRadioProcess(otInstance *aInstance) +void platformRadioProcess(otInstance *aInstance, const fd_set *aReadFdSet, const fd_set *aWriteFdSet) { -#if OPENTHREAD_POSIX_VIRTUAL_TIME == 0 - const int flags = POLLIN | POLLRDNORM | POLLERR | POLLNVAL | POLLHUP; - struct pollfd pollfd = {sSockFd, flags, 0}; + OT_UNUSED_VARIABLE(aReadFdSet); + OT_UNUSED_VARIABLE(aWriteFdSet); - if (poll(&pollfd, 1, 0) > 0 && (pollfd.revents & flags) != 0) +#if OPENTHREAD_POSIX_VIRTUAL_TIME == 0 + if (FD_ISSET(sRxFd, aReadFdSet)) { - ssize_t rval = recvfrom(sSockFd, (char *)&sReceiveMessage, sizeof(sReceiveMessage), 0, NULL, NULL); + ssize_t rval = recvfrom(sRxFd, (char *)&sReceiveMessage, sizeof(sReceiveMessage), 0, NULL, NULL); if (rval < 0) { - perror("recvfrom"); + perror("recvfrom(sRxFd)"); exit(EXIT_FAILURE); } @@ -791,30 +848,21 @@ void platformRadioProcess(otInstance *aInstance) void radioTransmit(struct RadioMessage *aMessage, const struct otRadioFrame *aFrame) { #if OPENTHREAD_POSIX_VIRTUAL_TIME == 0 + ssize_t rval; struct sockaddr_in sockaddr; memset(&sockaddr, 0, sizeof(sockaddr)); sockaddr.sin_family = AF_INET; - inet_pton(AF_INET, "127.0.0.1", &sockaddr.sin_addr); + inet_pton(AF_INET, OT_RADIO_GROUP, &sockaddr.sin_addr); - for (uint32_t i = 1; i <= WELLKNOWN_NODE_ID; i++) + sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + WELLKNOWN_NODE_ID)); + rval = + sendto(sTxFd, (const char *)aMessage, 1 + aFrame->mLength, 0, (struct sockaddr *)&sockaddr, sizeof(sockaddr)); + + if (rval < 0) { - ssize_t rval; - - if (gNodeId == i) - { - continue; - } - - sockaddr.sin_port = htons((uint16_t)(9000 + sPortOffset + i)); - rval = sendto(sSockFd, (const char *)aMessage, 1 + aFrame->mLength, 0, (struct sockaddr *)&sockaddr, - sizeof(sockaddr)); - - if (rval < 0) - { - perror("sendto"); - exit(EXIT_FAILURE); - } + perror("sendto(sTxFd)"); + exit(EXIT_FAILURE); } #else // OPENTHREAD_POSIX_VIRTUAL_TIME == 0 struct Event event; diff --git a/examples/platforms/posix/sim/platform-sim.c b/examples/platforms/posix/sim/platform-sim.c index 8a2e16dd4..6d38e0824 100644 --- a/examples/platforms/posix/sim/platform-sim.c +++ b/examples/platforms/posix/sim/platform-sim.c @@ -301,7 +301,7 @@ void otSysProcessDrivers(otInstance *aInstance) } platformAlarmProcess(aInstance); - platformRadioProcess(aInstance); + platformRadioProcess(aInstance, &read_fds, &write_fds); #if OPENTHREAD_POSIX_VIRTUAL_TIME_UART == 0 platformUartProcess(); #endif diff --git a/examples/platforms/posix/system.c b/examples/platforms/posix/system.c index c41047f23..e509543de 100644 --- a/examples/platforms/posix/system.c +++ b/examples/platforms/posix/system.c @@ -135,25 +135,31 @@ void otSysProcessDrivers(otInstance *aInstance) platformRadioUpdateFdSet(&read_fds, &write_fds, &max_fd); platformAlarmUpdateTimeout(&timeout); - if (!otTaskletsArePending(aInstance)) + if (otTaskletsArePending(aInstance)) { - rval = select(max_fd + 1, &read_fds, &write_fds, &error_fds, &timeout); - - if ((rval < 0) && (errno != EINTR)) - { - perror("select"); - exit(EXIT_FAILURE); - } + timeout.tv_sec = 0; + timeout.tv_usec = 0; } + rval = select(max_fd + 1, &read_fds, &write_fds, &error_fds, &timeout); + + if (rval >= 0) + { + platformUartProcess(); + platformRadioProcess(aInstance, &read_fds, &write_fds); + } + else if (errno != EINTR) + { + perror("select"); + exit(EXIT_FAILURE); + } + + platformAlarmProcess(aInstance); + if (gTerminate) { exit(0); } - - platformUartProcess(); - platformRadioProcess(aInstance); - platformAlarmProcess(aInstance); } #endif // OPENTHREAD_POSIX_VIRTUAL_TIME == 0 diff --git a/tests/scripts/thread-cert/config.py b/tests/scripts/thread-cert/config.py index e6a49019a..382410b87 100644 --- a/tests/scripts/thread-cert/config.py +++ b/tests/scripts/thread-cert/config.py @@ -459,8 +459,8 @@ def create_default_thread_message_factory(master_key=DEFAULT_MASTER_KEY): return message.MessageFactory(lowpan_parser=lowpan_parser) -def create_default_thread_sniffer(nodeid=SNIFFER_ID): - return sniffer.Sniffer(nodeid, create_default_thread_message_factory()) +def create_default_thread_sniffer(): + return sniffer.Sniffer(create_default_thread_message_factory()) def create_default_simulator(): diff --git a/tests/scripts/thread-cert/sniffer.py b/tests/scripts/thread-cert/sniffer.py index bb546c20e..a181b5f34 100644 --- a/tests/scripts/thread-cert/sniffer.py +++ b/tests/scripts/thread-cert/sniffer.py @@ -54,21 +54,19 @@ class Sniffer: RECV_BUFFER_SIZE = 4096 - def __init__(self, nodeid, message_factory): + def __init__(self, message_factory): """ Args: - nodeid (int): Node identifier message_factory (MessageFactory): Class producing messages from data bytes. """ - self.nodeid = nodeid self._message_factory = message_factory self._pcap = pcap.PcapCodec(os.getenv('TEST_NAME', 'current')) # Create transport transport_factory = sniffer_transport.SnifferTransportFactory() - self._transport = transport_factory.create_transport(nodeid) + self._transport = transport_factory.create_transport() self._thread = None self._thread_alive = threading.Event() diff --git a/tests/scripts/thread-cert/sniffer_transport.py b/tests/scripts/thread-cert/sniffer_transport.py index d88820a28..3e041e66e 100644 --- a/tests/scripts/thread-cert/sniffer_transport.py +++ b/tests/scripts/thread-cert/sniffer_transport.py @@ -96,8 +96,9 @@ class SnifferSocketTransport(SnifferTransport): PORT_OFFSET = int(os.getenv('PORT_OFFSET', "0")) - def __init__(self, nodeid): - self._nodeid = nodeid + RADIO_GROUP = '224.0.0.116' + + def __init__(self): self._socket = None def __del__(self): @@ -106,7 +107,7 @@ class SnifferSocketTransport(SnifferTransport): self.close() - def _nodeid_to_address(self, nodeid, ip_address=""): + def _nodeid_to_address(self, nodeid, ip_address=''): return ( ip_address, self.BASE_PORT @@ -129,7 +130,11 @@ class SnifferSocketTransport(SnifferTransport): if not self.is_opened: raise RuntimeError("Transport opening failed.") - self._socket.bind(self._nodeid_to_address(self._nodeid)) + self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1) + self._socket.setsockopt(socket.IPPROTO_IP, socket.IP_ADD_MEMBERSHIP, + socket.inet_aton(self.RADIO_GROUP) + socket.inet_aton('127.0.0.1')) + self._socket.bind(self._nodeid_to_address(self.WELLKNOWN_NODE_ID)) def close(self): if not self.is_opened: @@ -164,5 +169,5 @@ class MacFrame(ctypes.Structure): class SnifferTransportFactory(object): - def create_transport(self, nodeid): - return SnifferSocketTransport(nodeid) + def create_transport(self): + return SnifferSocketTransport()