[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.
This commit is contained in:
Yakun Xu
2019-07-16 09:03:42 -07:00
committed by Jonathan Hui
parent b289c94122
commit e80f46a993
7 changed files with 143 additions and 84 deletions
+4 -2
View File
@@ -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.
+105 -57
View File
@@ -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;
+1 -1
View File
@@ -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
+18 -12
View File
@@ -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
+2 -2
View File
@@ -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():
+2 -4
View File
@@ -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()
+11 -6
View File
@@ -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()