mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-08-23 01:04:01 +08:00
@@ -0,0 +1,2 @@
|
||||
*.bz2
|
||||
diff.txt
|
||||
@@ -0,0 +1,15 @@
|
||||
# process replay
|
||||
|
||||
Process replay is a regression test designed to identify any changes in the output of a process. This test replays a segment through individual processes and compares the output to a known good replay. Each make is represented in the test with a segment.
|
||||
|
||||
If the test fails, make sure that you didn't unintentionally change anything. If there are intentional changes, the reference logs will be updated.
|
||||
|
||||
Use `test_processes.py` to run the test locally.
|
||||
|
||||
Currently the following processes are tested:
|
||||
|
||||
* controlsd
|
||||
* radard
|
||||
* plannerd
|
||||
* calibrationd
|
||||
|
||||
+45
@@ -0,0 +1,45 @@
|
||||
#!/usr/bin/env python2
|
||||
import bz2
|
||||
import os
|
||||
import sys
|
||||
|
||||
import dictdiffer
|
||||
if "CI" in os.environ:
|
||||
tqdm = lambda x: x
|
||||
else:
|
||||
from tqdm import tqdm
|
||||
|
||||
from tools.lib.logreader import LogReader
|
||||
|
||||
|
||||
def save_log(dest, log_msgs):
|
||||
dat = ""
|
||||
for msg in log_msgs:
|
||||
dat += msg.as_builder().to_bytes()
|
||||
dat = bz2.compress(dat)
|
||||
|
||||
with open(dest, "w") as f:
|
||||
f.write(dat)
|
||||
|
||||
def compare_logs(log1, log2, ignore=[]):
|
||||
assert len(log1) == len(log2), "logs are not same length"
|
||||
|
||||
diff = []
|
||||
for msg1, msg2 in tqdm(zip(log1, log2)):
|
||||
assert msg1.which() == msg2.which(), "msgs not aligned between logs"
|
||||
|
||||
msg1_bytes = msg1.as_builder().to_bytes()
|
||||
msg2_bytes = msg2.as_builder().to_bytes()
|
||||
|
||||
if msg1_bytes != msg2_bytes:
|
||||
msg1_dict = msg1.to_dict(verbose=True)
|
||||
msg2_dict = msg2.to_dict(verbose=True)
|
||||
dd = dictdiffer.diff(msg1_dict, msg2_dict, ignore=ignore, tolerance=0)
|
||||
diff.extend(dd)
|
||||
return diff
|
||||
|
||||
if __name__ == "__main__":
|
||||
log1 = list(LogReader(sys.argv[1]))
|
||||
log2 = list(LogReader(sys.argv[2]))
|
||||
|
||||
compare_logs(log1, log2, sys.argv[3:])
|
||||
+185
@@ -0,0 +1,185 @@
|
||||
#!/usr/bin/env python2
|
||||
import gc
|
||||
import os
|
||||
import time
|
||||
|
||||
if "CI" in os.environ:
|
||||
tqdm = lambda x: x
|
||||
else:
|
||||
from tqdm import tqdm
|
||||
|
||||
from cereal import car
|
||||
from selfdrive.car.car_helpers import get_car
|
||||
import selfdrive.manager as manager
|
||||
import selfdrive.messaging as messaging
|
||||
from common.params import Params
|
||||
from selfdrive.services import service_list
|
||||
from collections import namedtuple
|
||||
|
||||
ProcessConfig = namedtuple('ProcessConfig', ['proc_name', 'pub_sub', 'ignore', 'init_callback', 'should_recv_callback'])
|
||||
|
||||
def fingerprint(msgs, pub_socks, sub_socks):
|
||||
print "start fingerprinting"
|
||||
manager.prepare_managed_process("logmessaged")
|
||||
manager.start_managed_process("logmessaged")
|
||||
|
||||
can = pub_socks["can"]
|
||||
logMessage = messaging.sub_sock(service_list["logMessage"].port)
|
||||
|
||||
time.sleep(1)
|
||||
messaging.drain_sock(logMessage)
|
||||
|
||||
# controlsd waits for a health packet before fingerprinting
|
||||
msg = messaging.new_message()
|
||||
msg.init("health")
|
||||
pub_socks["health"].send(msg.to_bytes())
|
||||
|
||||
canmsgs = filter(lambda msg: msg.which() == "can", msgs)
|
||||
for msg in canmsgs[:200]:
|
||||
can.send(msg.as_builder().to_bytes())
|
||||
|
||||
time.sleep(0.005)
|
||||
log = messaging.recv_one_or_none(logMessage)
|
||||
if log is not None and "fingerprinted" in log.logMessage:
|
||||
break
|
||||
manager.kill_managed_process("logmessaged")
|
||||
print "finished fingerprinting"
|
||||
|
||||
def get_car_params(msgs, pub_socks, sub_socks):
|
||||
sendcan = pub_socks.get("sendcan", None)
|
||||
if sendcan is None:
|
||||
sendcan = messaging.pub_sock(service_list["sendcan"].port)
|
||||
logcan = sub_socks.get("can", None)
|
||||
if logcan is None:
|
||||
logcan = messaging.sub_sock(service_list["can"].port)
|
||||
can = pub_socks.get("can", None)
|
||||
if can is None:
|
||||
can = messaging.pub_sock(service_list["can"].port)
|
||||
|
||||
time.sleep(0.5)
|
||||
|
||||
canmsgs = filter(lambda msg: msg.which() == "can", msgs)
|
||||
for m in canmsgs[:200]:
|
||||
can.send(m.as_builder().to_bytes())
|
||||
_, CP = get_car(logcan, sendcan)
|
||||
Params().put("CarParams", CP.to_bytes())
|
||||
time.sleep(0.5)
|
||||
messaging.drain_sock(logcan)
|
||||
|
||||
def radar_rcv_callback(msg, CP):
|
||||
if msg.which() != "can":
|
||||
return []
|
||||
|
||||
# hyundai and subaru don't have radar
|
||||
radar_msgs = {"honda": [0x445], "toyota": [0x19f, 0x22f], "gm": [0x475],
|
||||
"hyundai": [], "chrysler": [0x2d4], "subaru": []}.get(CP.carName, None)
|
||||
|
||||
if radar_msgs is None:
|
||||
raise NotImplementedError
|
||||
|
||||
for m in msg.can:
|
||||
if m.src == 1 and m.address in radar_msgs:
|
||||
return ["radarState", "liveTracks"]
|
||||
|
||||
return []
|
||||
|
||||
def plannerd_rcv_callback(msg, CP):
|
||||
if msg.which() in ["model", "radarState"]:
|
||||
time.sleep(0.005)
|
||||
else:
|
||||
time.sleep(0.002)
|
||||
return {"model": ["pathPlan"], "radarState": ["plan"]}.get(msg.which(), [])
|
||||
|
||||
CONFIGS = [
|
||||
ProcessConfig(
|
||||
proc_name="controlsd",
|
||||
pub_sub={
|
||||
"can": ["controlsState", "carState", "carControl", "sendcan"],
|
||||
"thermal": [], "health": [], "liveCalibration": [], "driverMonitoring": [], "plan": [], "pathPlan": []
|
||||
},
|
||||
ignore=["logMonoTime", "controlsState.startMonoTime", "controlsState.cumLagMs"],
|
||||
init_callback=fingerprint,
|
||||
should_recv_callback=None,
|
||||
),
|
||||
ProcessConfig(
|
||||
proc_name="radard",
|
||||
pub_sub={
|
||||
"can": ["radarState", "liveTracks"],
|
||||
"liveParameters": [], "controlsState": [], "model": [],
|
||||
},
|
||||
ignore=["logMonoTime", "radarState.cumLagMs"],
|
||||
init_callback=get_car_params,
|
||||
should_recv_callback=radar_rcv_callback,
|
||||
),
|
||||
ProcessConfig(
|
||||
proc_name="plannerd",
|
||||
pub_sub={
|
||||
"model": ["pathPlan"], "radarState": ["plan"],
|
||||
"carState": [], "controlsState": [], "liveParameters": [],
|
||||
},
|
||||
ignore=["logMonoTime", "valid", "plan.processingDelay"],
|
||||
init_callback=get_car_params,
|
||||
should_recv_callback=plannerd_rcv_callback,
|
||||
),
|
||||
ProcessConfig(
|
||||
proc_name="calibrationd",
|
||||
pub_sub={
|
||||
"cameraOdometry": ["liveCalibration"]
|
||||
},
|
||||
ignore=["logMonoTime"],
|
||||
init_callback=get_car_params,
|
||||
should_recv_callback=None,
|
||||
),
|
||||
]
|
||||
|
||||
def replay_process(cfg, lr):
|
||||
gc.disable() # gc can occasionally cause canparser to timeout
|
||||
|
||||
pub_socks, sub_socks = {}, {}
|
||||
for pub, sub in cfg.pub_sub.iteritems():
|
||||
pub_socks[pub] = messaging.pub_sock(service_list[pub].port)
|
||||
|
||||
for s in sub:
|
||||
sub_socks[s] = messaging.sub_sock(service_list[s].port)
|
||||
|
||||
all_msgs = sorted(lr, key=lambda msg: msg.logMonoTime)
|
||||
pub_msgs = filter(lambda msg: msg.which() in pub_socks.keys(), all_msgs)
|
||||
|
||||
params = Params()
|
||||
params.manager_start()
|
||||
params.put("Passive", "0")
|
||||
|
||||
manager.gctx = {}
|
||||
manager.prepare_managed_process(cfg.proc_name)
|
||||
manager.start_managed_process(cfg.proc_name)
|
||||
time.sleep(3) # Wait for started process to be ready
|
||||
|
||||
if cfg.init_callback is not None:
|
||||
cfg.init_callback(all_msgs, pub_socks, sub_socks)
|
||||
|
||||
CP = car.CarParams.from_bytes(params.get("CarParams", block=True))
|
||||
|
||||
log_msgs = []
|
||||
for msg in tqdm(pub_msgs):
|
||||
if cfg.should_recv_callback is not None:
|
||||
recv_socks = cfg.should_recv_callback(msg, CP)
|
||||
else:
|
||||
recv_socks = cfg.pub_sub[msg.which()]
|
||||
|
||||
pub_socks[msg.which()].send(msg.as_builder().to_bytes())
|
||||
|
||||
if len(recv_socks):
|
||||
# TODO: add timeout
|
||||
for sock in recv_socks:
|
||||
m = messaging.recv_one(sub_socks[sock])
|
||||
|
||||
# make these values fixed for faster comparison
|
||||
m_builder = m.as_builder()
|
||||
m_builder.logMonoTime = 0
|
||||
m_builder.valid = True
|
||||
log_msgs.append(m_builder.as_reader())
|
||||
|
||||
gc.enable()
|
||||
manager.kill_managed_process(cfg.proc_name)
|
||||
return log_msgs
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
e3388c62ffb80f4b3ca8721da56a581a93c44e79
|
||||
+118
@@ -0,0 +1,118 @@
|
||||
#!/usr/bin/env python2
|
||||
import os
|
||||
import requests
|
||||
import sys
|
||||
import tempfile
|
||||
|
||||
from selfdrive.test.tests.process_replay.compare_logs import compare_logs
|
||||
from selfdrive.test.tests.process_replay.process_replay import replay_process, CONFIGS
|
||||
from tools.lib.logreader import LogReader
|
||||
|
||||
segments = [
|
||||
"0375fdf7b1ce594d|2019-06-13--08-32-25--3", # HONDA.ACCORD
|
||||
"99c94dc769b5d96e|2019-08-03--14-19-59--2", # HONDA.CIVIC
|
||||
"cce908f7eb8db67d|2019-08-02--15-09-51--3", # TOYOTA.COROLLA_TSS2
|
||||
"7ad88f53d406b787|2019-07-09--10-18-56--8", # GM.VOLT
|
||||
"704b2230eb5190d6|2019-07-06--19-29-10--0", # HYUNDAI.KIA_SORENTO
|
||||
"b6e1317e1bfbefa6|2019-07-06--04-05-26--5", # CHRYSLER.JEEP_CHEROKEE
|
||||
"7873afaf022d36e2|2019-07-03--18-46-44--0", # SUBARU.IMPREZA
|
||||
]
|
||||
|
||||
def get_segment(segment_name):
|
||||
route_name, segment_num = segment_name.rsplit("--", 1)
|
||||
rlog_url = "https://commadataci.blob.core.windows.net/openpilotci/%s/%s/rlog.bz2" \
|
||||
% (route_name.replace("|", "/"), segment_num)
|
||||
r = requests.get(rlog_url)
|
||||
if r.status_code != 200:
|
||||
return None
|
||||
|
||||
with tempfile.NamedTemporaryFile(delete=False, suffix=".bz2") as f:
|
||||
f.write(r.content)
|
||||
return f.name
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
process_replay_dir = os.path.dirname(os.path.abspath(__file__))
|
||||
ref_commit_fn = os.path.join(process_replay_dir, "ref_commit")
|
||||
|
||||
if not os.path.isfile(ref_commit_fn):
|
||||
print "couldn't find reference commit"
|
||||
sys.exit(1)
|
||||
|
||||
ref_commit = open(ref_commit_fn).read().strip()
|
||||
print "***** testing against commit %s *****" % ref_commit
|
||||
|
||||
results = {}
|
||||
for segment in segments:
|
||||
print "***** testing route segment %s *****\n" % segment
|
||||
|
||||
results[segment] = {}
|
||||
|
||||
rlog_fn = get_segment(segment)
|
||||
|
||||
if rlog_fn is None:
|
||||
print "failed to get segment %s" % segment
|
||||
sys.exit(1)
|
||||
|
||||
lr = LogReader(rlog_fn)
|
||||
|
||||
for cfg in CONFIGS:
|
||||
log_msgs = replay_process(cfg, lr)
|
||||
|
||||
log_fn = os.path.join(process_replay_dir, "%s_%s_%s.bz2" % (segment, cfg.proc_name, ref_commit))
|
||||
|
||||
if not os.path.isfile(log_fn):
|
||||
url = "https://commadataci.blob.core.windows.net/openpilotci/"
|
||||
req = requests.get(url + os.path.basename(log_fn))
|
||||
if req.status_code != 200:
|
||||
results[segment][cfg.proc_name] = "failed to download comparison log"
|
||||
continue
|
||||
|
||||
with tempfile.NamedTemporaryFile(suffix=".bz2") as f:
|
||||
f.write(req.content)
|
||||
f.flush()
|
||||
f.seek(0)
|
||||
cmp_log_msgs = list(LogReader(f.name))
|
||||
else:
|
||||
cmp_log_msgs = list(LogReader(log_fn))
|
||||
|
||||
diff = compare_logs(cmp_log_msgs, log_msgs, cfg.ignore)
|
||||
results[segment][cfg.proc_name] = diff
|
||||
os.remove(rlog_fn)
|
||||
|
||||
failed = False
|
||||
with open(os.path.join(process_replay_dir, "diff.txt"), "w") as f:
|
||||
f.write("***** tested against commit %s *****\n" % ref_commit)
|
||||
|
||||
for segment, result in results.items():
|
||||
f.write("***** differences for segment %s *****\n" % segment)
|
||||
print "***** results for segment %s *****" % segment
|
||||
|
||||
for proc, diff in result.items():
|
||||
f.write("*** process: %s ***\n" % proc)
|
||||
print "\t%s" % proc
|
||||
|
||||
if isinstance(diff, str):
|
||||
print "\t\t%s" % diff
|
||||
failed = True
|
||||
elif len(diff):
|
||||
cnt = {}
|
||||
for d in diff:
|
||||
f.write("\t%s\n" % str(d))
|
||||
|
||||
k = str(d[1])
|
||||
cnt[k] = 1 if k not in cnt else cnt[k] + 1
|
||||
|
||||
for k, v in sorted(cnt.items()):
|
||||
print "\t\t%s: %s" % (k, v)
|
||||
failed = True
|
||||
|
||||
if failed:
|
||||
print "TEST FAILED"
|
||||
else:
|
||||
print "TEST SUCCEEDED"
|
||||
|
||||
print "\n\nTo update the reference logs for this test run:"
|
||||
print "./update_refs.py"
|
||||
|
||||
sys.exit(int(failed))
|
||||
+42
@@ -0,0 +1,42 @@
|
||||
#!/usr/bin/env python2
|
||||
import os
|
||||
import sys
|
||||
|
||||
from selfdrive.test.openpilotci_upload import upload_file
|
||||
from selfdrive.test.tests.process_replay.compare_logs import save_log
|
||||
from selfdrive.test.tests.process_replay.process_replay import replay_process, CONFIGS
|
||||
from selfdrive.test.tests.process_replay.test_processes import segments, get_segment
|
||||
from selfdrive.version import get_git_commit
|
||||
from tools.lib.logreader import LogReader
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
no_upload = "--no-upload" in sys.argv
|
||||
|
||||
process_replay_dir = os.path.dirname(os.path.abspath(__file__))
|
||||
ref_commit_fn = os.path.join(process_replay_dir, "ref_commit")
|
||||
|
||||
ref_commit = get_git_commit()
|
||||
with open(ref_commit_fn, "w") as f:
|
||||
f.write(ref_commit)
|
||||
|
||||
for segment in segments:
|
||||
rlog_fn = get_segment(segment)
|
||||
|
||||
if rlog_fn is None:
|
||||
print "failed to get segment %s" % segment
|
||||
sys.exit(1)
|
||||
|
||||
lr = LogReader(rlog_fn)
|
||||
|
||||
for cfg in CONFIGS:
|
||||
log_msgs = replay_process(cfg, lr)
|
||||
log_fn = os.path.join(process_replay_dir, "%s_%s_%s.bz2" % (segment, cfg.proc_name, ref_commit))
|
||||
save_log(log_fn, log_msgs)
|
||||
|
||||
if not no_upload:
|
||||
upload_file(log_fn, os.path.basename(log_fn))
|
||||
os.remove(log_fn)
|
||||
os.remove(rlog_fn)
|
||||
|
||||
print "done"
|
||||
Reference in New Issue
Block a user