Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions recv_functions/sunstream_recv.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import click # Command line interface
from sunstream_util import print_in_hex_only
import time
from ..scripts.candata import non_can_data_to_mysql, can_data_to_mysql

class sunstream_receiver:
def __init__(self, mcast_grp: str, mcast_port: int, msg_callback, dbc_file: str, datagram_socket_timeout: int=2):
Expand Down Expand Up @@ -97,6 +98,9 @@ def compute_datagram_size(self) -> int:
def dump_test(msg: can.Message):
# Print out the can message
click.secho(msg, fg="blue")
# CAN msg to SQL db
non_can_data_to_mysql(msg)
can_data_to_mysql(msg)
pass

@click.command()
Expand Down
95 changes: 95 additions & 0 deletions scripts/candata.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import pandas as pd
from sqlalchemy import create_engine
import re

# local instance for now, will connect to official one later
cnx_string = 'mysql+mysqlconnector://root:@localhost:3306/midnight_sun'

engine = create_engine(cnx_string)

# sample can msgs from https://uwmidsun.atlassian.net/wiki/spaces/ELEC/pages/3157950471/Meeting+Notes+-+Strategy+x+FW
# msg1 = {
# "id": 0,
# "source": "BMS_CARRIER",
# "target": "CENTRE_CONSOLE",
# "msg_name": "bps heartbeat",
# "is_critical": True,
# "can_data": {
# "u8": {
# "field_name_1": "status"
# }
# }
# }

# msg2 = {
# "id": 1,
# "source": "CENTRE_CONSOLE",
# "target": "BMS_CARRIER, SOLAR_5_MPPTS, SOLAR_6_MPPTS, MOTOR_CONTROLLER",
# "msg_name": "set relay states",
# "is_critical": True,
# "can_data": {
# "u16": {
# "field_name_1": "relay_mask",
# "field_name_2": "relay_state"
# }
# }
# }

# msg3 = {
# "id": 33,
# "source": "BMS_CARRIER",
# "target": "TELEMETRY",
# "msg_name": "battery aggregate vc",
# "msg_readable_name": "battery aggregate voltage and current",
# "can_data": {
# "u32": {
# "field_name_1": "voltage",
# "field_name_2": "current"
# }
# }
# }

# msg4 = {
# "id": 55,
# "source": "POWER_DISTRIBUTION_REAR",
# "target": "CENTRE_CONSOLE",
# "msg_name": "rear current measurement",
# "can_data": {
# "u16": {
# "field_name_1": "current_id",
# "field_name_2": "current"
# }
# }
# }

# array of sample can msgs for testing
# msgs = [msg1, msg2, msg3, msg4]

# non can_data -------
def non_can_data_to_mysql(msgs):

df = pd.DataFrame(msgs, columns=["id", "source", "target", "msg_name", "is_critical", "msg_readable_name"])

# send to mysql db
df.to_sql('non_can_data', con=engine, if_exists='append', index=False)

def can_data_to_mysql(msgs):
# array to hold field values
can_data_values = []
# loop can msgs
for msg in msgs:
for key, value in msg["can_data"].items():
# Extract u{x} value from key using regular expression
match = re.search(r"u(\d+)", key)
u_value = int(match.group(1))

# loop through each nested field values and add it to the can_data array
can_data_values.extend([
{"id": msg["id"], "field_name": field_key, "field_value": field_value, "u_value": u_value}
for field_key, field_value in value.items()
])

df = pd.DataFrame(can_data_values)
df.to_sql('can_data', con=engine, if_exists='append', index=False)

engine.dispose()