mirror of
https://github.com/espressif/openthread.git
synced 2026-09-01 06:49:54 +00:00
[posix-app] enhance virtual time for posix app (#3016)
This commit is contained in:
@@ -27,19 +27,25 @@
|
||||
# POSSIBILITY OF SUCH DAMAGE.
|
||||
#
|
||||
|
||||
import binascii
|
||||
import bisect
|
||||
import cmd
|
||||
import os
|
||||
import socket
|
||||
import struct
|
||||
import time
|
||||
import sys
|
||||
import traceback
|
||||
import time
|
||||
|
||||
import io
|
||||
import config
|
||||
import message
|
||||
import pcap
|
||||
|
||||
def dbg_print(*args):
|
||||
if False:
|
||||
print args
|
||||
|
||||
class RealTime:
|
||||
|
||||
def __init__(self):
|
||||
@@ -60,9 +66,29 @@ class RealTime:
|
||||
|
||||
class VirtualTime:
|
||||
|
||||
OT_SIM_EVENT_ALARM_FIRED = 0
|
||||
OT_SIM_EVENT_RADIO_RECEIVED = 1
|
||||
OT_SIM_EVENT_UART_WRITE = 2
|
||||
OT_SIM_EVENT_RADIO_SPINEL_WRITE = 3
|
||||
OT_SIM_EVENT_POSTCMD = 4
|
||||
|
||||
EVENT_TIME = 0
|
||||
EVENT_SEQUENCE = 1
|
||||
EVENT_ADDR = 2
|
||||
EVENT_TYPE = 3
|
||||
EVENT_DATA_LENGTH = 4
|
||||
EVENT_DATA = 5
|
||||
|
||||
BASE_PORT = 9000
|
||||
MAX_NODES = 34
|
||||
PORT_OFFSET = int(os.getenv('PORT_OFFSET', "0"))
|
||||
MAX_MESSAGE = 1024
|
||||
END_OF_TIME = 0x7fffffff
|
||||
PORT_OFFSET = int(os.getenv('PORT_OFFSET', '0'))
|
||||
|
||||
BLOCK_TIMEOUT = 4
|
||||
|
||||
RADIO_ONLY = os.getenv('RADIO_DEVICE') != None
|
||||
NCP_SIM = os.getenv('NODE_TYPE', 'sim') == 'ncp-sim'
|
||||
|
||||
def __init__(self):
|
||||
self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
@@ -73,11 +99,17 @@ class VirtualTime:
|
||||
|
||||
self.devices = {}
|
||||
self.event_queue = []
|
||||
self.event_count = 0
|
||||
# there could be events scheduled at exactly the same time
|
||||
self.event_sequence = 0
|
||||
self.current_time = 0
|
||||
self.current_event = None;
|
||||
self.current_event = None
|
||||
self.awake_devices = set()
|
||||
|
||||
self._pcap = pcap.PcapCodec(os.getenv('TEST_NAME', 'current'))
|
||||
# the addr for spinel-cli sending OT_SIM_EVENT_POSTCMD
|
||||
self._spinel_cli_addr = (ip, self.BASE_PORT + self.port)
|
||||
self.current_nodeid = None
|
||||
self._pause_time = 0
|
||||
|
||||
self._message_factory = config.create_default_thread_message_factory()
|
||||
|
||||
@@ -102,7 +134,6 @@ class VirtualTime:
|
||||
except Exception as e:
|
||||
# Just print the exception to the console
|
||||
print("EXCEPTION: %s" % e)
|
||||
pass
|
||||
|
||||
def set_lowpan_context(self, cid, prefix):
|
||||
self._message_factory.set_lowpan_context(cid, prefix)
|
||||
@@ -125,24 +156,63 @@ class VirtualTime:
|
||||
|
||||
return message.MessagesSet(messages)
|
||||
|
||||
def receive_events(self):
|
||||
def _is_radio(self, addr):
|
||||
return addr[1] < self.BASE_PORT * 2
|
||||
|
||||
if self.current_event is None:
|
||||
self.sock.setblocking(0)
|
||||
def _to_core_addr(self, addr):
|
||||
assert self._is_radio(addr)
|
||||
return (addr[0], addr[1] + self.BASE_PORT)
|
||||
|
||||
def _to_radio_addr(self, addr):
|
||||
assert not self._is_radio(addr)
|
||||
return (addr[0], addr[1] - self.BASE_PORT)
|
||||
|
||||
def _core_addr_from(self, nodeid):
|
||||
if self.RADIO_ONLY:
|
||||
return ('127.0.0.1', self.BASE_PORT + self.port + nodeid)
|
||||
else:
|
||||
self.sock.setblocking(1)
|
||||
return ('127.0.0.1', self.port + nodeid)
|
||||
|
||||
def _next_event_time(self):
|
||||
if len(self.event_queue) == 0:
|
||||
return self.END_OF_TIME
|
||||
else:
|
||||
return self.event_queue[0][0]
|
||||
|
||||
def receive_events(self):
|
||||
""" Receive events until all devices are asleep. """
|
||||
while True:
|
||||
try:
|
||||
msg, addr = self.sock.recvfrom(1024)
|
||||
except socket.error:
|
||||
break
|
||||
if self.current_event or len(self.awake_devices) or (self._next_event_time() > self._pause_time and self.current_nodeid):
|
||||
self.sock.settimeout(self.BLOCK_TIMEOUT)
|
||||
try:
|
||||
msg, addr = self.sock.recvfrom(self.MAX_MESSAGE)
|
||||
except socket.error:
|
||||
# print debug information on failure
|
||||
print('Current nodeid:')
|
||||
print(self.current_nodeid)
|
||||
print('Current awake:')
|
||||
print(self.awake_devices)
|
||||
print('Current time:')
|
||||
print(self.current_time)
|
||||
print('Current event:')
|
||||
print(self.current_event)
|
||||
print('Events:')
|
||||
for event in self.event_queue:
|
||||
print(event)
|
||||
raise
|
||||
else:
|
||||
self.sock.settimeout(0)
|
||||
try:
|
||||
msg, addr = self.sock.recvfrom(self.MAX_MESSAGE)
|
||||
except socket.error:
|
||||
break
|
||||
|
||||
if addr not in self.devices:
|
||||
if addr != self._spinel_cli_addr and addr not in self.devices:
|
||||
self.devices[addr] = {}
|
||||
self.devices[addr]['alarm'] = None
|
||||
self.devices[addr]['msgs'] = []
|
||||
self.devices[addr]['time'] = self.current_time
|
||||
self.awake_devices.discard(addr)
|
||||
#print "New device:", addr, self.devices
|
||||
|
||||
delay, type, datalen = struct.unpack('=QBH', msg[:11])
|
||||
@@ -150,31 +220,37 @@ class VirtualTime:
|
||||
|
||||
event_time = self.current_time + delay
|
||||
|
||||
if type == 0:
|
||||
if data:
|
||||
dbg_print("New event: ", event_time, addr, type, datalen, binascii.hexlify(data))
|
||||
else:
|
||||
dbg_print("New event: ", event_time, addr, type, datalen)
|
||||
|
||||
if type == self.OT_SIM_EVENT_ALARM_FIRED:
|
||||
# remove any existing alarm event for device
|
||||
try:
|
||||
if self.devices[addr]['alarm']:
|
||||
self.event_queue.remove(self.devices[addr]['alarm'])
|
||||
#print "-- Remove\t", self.devices[addr]['alarm']
|
||||
self.devices[addr]['alarm'] = None
|
||||
except ValueError:
|
||||
pass
|
||||
|
||||
# add alarm event to event queue
|
||||
event = (event_time, addr, type, datalen)
|
||||
event = (event_time, self.event_sequence, addr, type, datalen)
|
||||
self.event_sequence += 1
|
||||
#print "-- Enqueue\t", event, delay, self.current_time
|
||||
bisect.insort(self.event_queue, event)
|
||||
self.devices[addr]['alarm'] = event
|
||||
|
||||
if self.current_event is not None and self.current_event[1] == addr:
|
||||
self.awake_devices.discard(addr)
|
||||
|
||||
if self.current_event and self.current_event[self.EVENT_ADDR] == addr:
|
||||
#print "Done\t", self.current_event
|
||||
self.current_event = None
|
||||
return
|
||||
|
||||
elif type == 1:
|
||||
elif type == self.OT_SIM_EVENT_RADIO_RECEIVED:
|
||||
assert self._is_radio(addr)
|
||||
# add radio receive events event queue
|
||||
for device in self.devices:
|
||||
if device != addr:
|
||||
event = (event_time, device, type, datalen, data)
|
||||
if device != addr and self._is_radio(device):
|
||||
event = (event_time, self.event_sequence, device, type, datalen, data)
|
||||
self.event_sequence += 1
|
||||
#print "-- Enqueue\t", event
|
||||
bisect.insort(self.event_queue, event)
|
||||
|
||||
@@ -182,34 +258,67 @@ class VirtualTime:
|
||||
self._add_message(addr[1] - self.port, data)
|
||||
|
||||
# add radio transmit done events to event queue
|
||||
event = (event_time, addr, type, datalen, data)
|
||||
event = (event_time, self.event_sequence, addr, type, datalen, data)
|
||||
self.event_sequence += 1
|
||||
bisect.insort(self.event_queue, event)
|
||||
|
||||
self.awake_devices.add(addr)
|
||||
|
||||
elif type == self.OT_SIM_EVENT_RADIO_SPINEL_WRITE:
|
||||
assert not self._is_radio(addr)
|
||||
radio_addr = self._to_radio_addr(addr)
|
||||
if not self.devices.has_key(radio_addr):
|
||||
self.awake_devices.add(radio_addr)
|
||||
|
||||
event = (event_time, self.event_sequence, radio_addr, self.OT_SIM_EVENT_UART_WRITE, datalen, data)
|
||||
self.event_sequence += 1
|
||||
bisect.insort(self.event_queue, event)
|
||||
|
||||
self.awake_devices.add(addr)
|
||||
|
||||
elif type == self.OT_SIM_EVENT_UART_WRITE:
|
||||
assert self._is_radio(addr)
|
||||
core_addr = self._to_core_addr(addr)
|
||||
if not self.devices.has_key(core_addr):
|
||||
self.awake_devices.add(core_addr)
|
||||
|
||||
event = (event_time, self.event_sequence, core_addr, self.OT_SIM_EVENT_RADIO_SPINEL_WRITE, datalen, data)
|
||||
self.event_sequence += 1
|
||||
bisect.insort(self.event_queue, event)
|
||||
|
||||
self.awake_devices.add(addr)
|
||||
|
||||
elif type == self.OT_SIM_EVENT_POSTCMD:
|
||||
assert self.current_time == self._pause_time
|
||||
nodeid = struct.unpack('=B', data)[0]
|
||||
if self.current_nodeid == nodeid:
|
||||
self.current_nodeid = None
|
||||
|
||||
def _send_message(self, message, addr):
|
||||
while True:
|
||||
try:
|
||||
sent = self.sock.sendto(message, addr)
|
||||
except socket.error:
|
||||
traceback.print_exc()
|
||||
time.sleep(0)
|
||||
else:
|
||||
break
|
||||
assert sent == len(message)
|
||||
|
||||
def process_next_event(self):
|
||||
|
||||
if self.current_event != None:
|
||||
return
|
||||
|
||||
#print "Events", len(self.event_queue)
|
||||
count = 0
|
||||
for event in self.event_queue:
|
||||
#print count, event
|
||||
count += 1
|
||||
assert self.current_event is None
|
||||
assert self._next_event_time() < self.END_OF_TIME
|
||||
|
||||
# process next event
|
||||
try:
|
||||
event = self.event_queue.pop(0)
|
||||
except IndexError:
|
||||
return
|
||||
event = self.event_queue.pop(0)
|
||||
|
||||
#print "Pop\t", event
|
||||
|
||||
if len(event) == 4:
|
||||
event_time, addr, type, datalen = event
|
||||
if len(event) == 5:
|
||||
event_time, sequence, addr, type, datalen = event
|
||||
dbg_print("Pop event: ", event_time, addr, type, datalen)
|
||||
else:
|
||||
event_time, addr, type, datalen, data = event
|
||||
event_time, sequence, addr, type, datalen, data = event
|
||||
dbg_print("Pop event: ", event_time, addr, type, datalen, binascii.hexlify(data))
|
||||
|
||||
self.event_count += 1
|
||||
self.current_event = event
|
||||
|
||||
assert(event_time >= self.current_time)
|
||||
@@ -220,33 +329,52 @@ class VirtualTime:
|
||||
|
||||
message = struct.pack('=QBH', elapsed, type, datalen)
|
||||
|
||||
if type == 0:
|
||||
if type == self.OT_SIM_EVENT_ALARM_FIRED:
|
||||
self.devices[addr]['alarm'] = None
|
||||
self.sock.sendto(message, addr)
|
||||
elif type == 1:
|
||||
self._send_message(message, addr)
|
||||
elif type == self.OT_SIM_EVENT_RADIO_RECEIVED:
|
||||
message += data
|
||||
self.sock.sendto(message, addr)
|
||||
self._send_message(message, addr)
|
||||
elif type == self.OT_SIM_EVENT_RADIO_SPINEL_WRITE:
|
||||
message += data
|
||||
self._send_message(message, addr)
|
||||
elif type == self.OT_SIM_EVENT_UART_WRITE:
|
||||
message += data
|
||||
self._send_message(message, addr)
|
||||
|
||||
def sync_devices(self):
|
||||
self.current_time = self._pause_time
|
||||
for addr in self.devices:
|
||||
elapsed = self.current_time - self.devices[addr]['time']
|
||||
if elapsed == 0:
|
||||
continue
|
||||
dbg_print('syncing', addr, elapsed)
|
||||
self.devices[addr]['time'] = self.current_time
|
||||
message = struct.pack('=QBH', elapsed, 0, 0)
|
||||
self.sock.sendto(message, addr)
|
||||
|
||||
def go(self, duration):
|
||||
message = struct.pack('=QBH', elapsed, self.OT_SIM_EVENT_ALARM_FIRED, 0)
|
||||
self._send_message(message, addr)
|
||||
self.awake_devices.add(addr)
|
||||
self.receive_events()
|
||||
self.awake_devices.clear()
|
||||
|
||||
def go(self, duration, nodeid=None):
|
||||
assert self.current_time == self._pause_time
|
||||
duration = int(duration) * 1000000
|
||||
|
||||
start_time = self.current_time
|
||||
self.current_event = None
|
||||
|
||||
print "running for %d us" % duration
|
||||
|
||||
dbg_print('running for %d us' % duration)
|
||||
self._pause_time += duration
|
||||
if nodeid:
|
||||
if self.NCP_SIM:
|
||||
self.current_nodeid = nodeid
|
||||
self.awake_devices.add(self._core_addr_from(nodeid))
|
||||
self.receive_events()
|
||||
|
||||
while (self.current_time - start_time) < duration:
|
||||
while self._next_event_time() <= self._pause_time:
|
||||
self.process_next_event()
|
||||
self.receive_events()
|
||||
if duration > 0:
|
||||
self.sync_devices()
|
||||
dbg_print('current time %d us' % self.current_time)
|
||||
|
||||
self.sync_devices()
|
||||
|
||||
if __name__ == '__main__':
|
||||
simulator = VirtualTime()
|
||||
while True:
|
||||
simulator.go(0)
|
||||
|
||||
Reference in New Issue
Block a user