From 83d1c5c68ae323ee0146f21a16336bc516508945 Mon Sep 17 00:00:00 2001 From: Oleg Kalachev Date: Tue, 30 Jun 2026 13:03:43 +0300 Subject: [PATCH] Implement entirely new logging subsystem Make topic-based logging mechanics. Print any log value to console. Possibility to expose any logging value to mavlink. --- flix/cli.ino | 18 ++- flix/flix.ino | 3 +- flix/log.ino | 297 +++++++++++++++++++++++++++++++++++-------- flix/mavlink.ino | 50 +++++++- flix/parameters.ino | 16 +++ flix/util.h | 43 +++++++ gazebo/Arduino.h | 1 + gazebo/flix.h | 11 +- gazebo/simulator.cpp | 1 + tools/csv_to_mcap.py | 11 +- tools/log_to_csv.py | 76 +++++++++++ 11 files changed, 464 insertions(+), 63 deletions(-) create mode 100755 tools/log_to_csv.py diff --git a/flix/cli.ino b/flix/cli.ino index 9eaa12c..9c82419 100644 --- a/flix/cli.ino +++ b/flix/cli.ino @@ -52,6 +52,11 @@ const char* motd = "mot - show motor output\n" "log [dump] - print log header [and data]\n" "mfr, mfl, mrr, mrl - test motor (remove props)\n" +"log - show log info\n" +"log header - show log header\n" +"log - setup log topic rate\n" +"l - show log values starting with str\n" +"l expose - expose log value to telemetry\n" "sys - show system info\n" "reset - reset drone's state\n" "reboot - reboot the drone\n"; @@ -153,9 +158,16 @@ void doCommand(String str, bool echo = false) { } else if (command == "mot") { print("front-right %g front-left %g rear-right %g rear-left %g\n", motors[MOTOR_FRONT_RIGHT], motors[MOTOR_FRONT_LEFT], motors[MOTOR_REAR_RIGHT], motors[MOTOR_REAR_LEFT]); - } else if (command == "log") { + } else if (command == "log" && arg0 == "") { + printLogInfo(); + } else if (command == "log" && arg1 != "") { + configLogThrottle(arg0.c_str(), arg1.toFloat()); + } else if (command == "log" && arg0 == "header") { printLogHeader(); - if (arg0 == "dump") printLogData(); + } else if (command == "l" && arg0 == "expose" && arg1 != "") { + exposeLogValue(arg1.c_str()); + } else if (command == "l") { + printLogValues(arg0.c_str()); } else if (command == "cr") { calibrateRC(); } else if (command == "ca") { @@ -174,6 +186,8 @@ void doCommand(String str, bool echo = false) { print("Temperature: %.1f °C\n", temperatureRead()); print("Total RAM: %d KB\n", ESP.getHeapSize() / 1024); print("Free heap: %d KB\n", ESP.getFreeHeap() / 1024); + print("PSRAM: %d KB\n", ESP.getPsramSize() / 1024); + print("Free PSRAM: %d KB\n", ESP.getFreePsram() / 1024); print("Firmware: " __DATE__ " " __TIME__ "\n"); // Print tasks table print("Num Task MinSt Prio Core CPU%%\n"); diff --git a/flix/flix.ino b/flix/flix.ino index 615664b..385babc 100644 --- a/flix/flix.ino +++ b/flix/flix.ino @@ -26,6 +26,7 @@ void setup() { setupWiFi(); setupIMU(); setupRC(); + setupLog(); setLED(false); print("Initializing complete\n"); } @@ -40,6 +41,6 @@ void loop() { handleInput(); processMavlink(); readVoltage(); - logData(); + loopLog(); syncParameters(); } diff --git a/flix/log.ino b/flix/log.ino index 6ccd6c6..7d4e758 100644 --- a/flix/log.ino +++ b/flix/log.ino @@ -1,77 +1,270 @@ // Copyright (c) 2023 Oleg Kalachev // Repository: https://github.com/okalachev/flix -// In-RAM logging +// Logging subsystem #include "vector.h" #include "util.h" -#define LOG_RATE 100 -#define LOG_DURATION 10 -#define LOG_SIZE LOG_DURATION * LOG_RATE +int logMemory = 0; // 0 - RAM, 1 - PSRAM, -1 - disabled +float logUsage = 0.5; // fraction of free memory to use for log -Vector attitudeEuler; -Vector attitudeTargetEuler; - -struct LogEntry { +struct LogValue { const char *name; - float *value; + Value value; + float lastValue = NAN; + bool logged = true; // if false, use only for triggering log update + LogValue() : name(nullptr), value() {}; // empty value constructor + template + LogValue(const char *name, T value, bool logged = true) : name(name), value(value), logged(logged) {}; }; -LogEntry logEntries[] = { - {"t", &t}, - {"rates.x", &rates.x}, - {"rates.y", &rates.y}, - {"rates.z", &rates.z}, - {"ratesTarget.x", &ratesTarget.x}, - {"ratesTarget.y", &ratesTarget.y}, - {"ratesTarget.z", &ratesTarget.z}, - {"attitude.x", &attitudeEuler.x}, - {"attitude.y", &attitudeEuler.y}, - {"attitude.z", &attitudeEuler.z}, - {"attitudeTarget.x", &attitudeTargetEuler.x}, - {"attitudeTarget.y", &attitudeTargetEuler.y}, - {"attitudeTarget.z", &attitudeTargetEuler.z}, - {"thrustTarget", &thrustTarget} +struct LogTopic { + LogValue values[10]; + int length = 0; // number of logged values + float throttle; // max update rate, Hz + float lastUpdate = -INFINITY; + + LogTopic(float throttle, LogValue v0, LogValue v1 = {}, LogValue v2 = {}, LogValue v3 = {}, LogValue v4 = {}, LogValue v5 = {}, LogValue v6 = {}, LogValue v7 = {}, LogValue v8 = {}, LogValue v9 = {}) : + throttle(throttle), values{v0, v1, v2, v3, v4, v5, v6, v7, v8, v9} { + // Count logged values + for (auto& v : values) { + if (v.name == nullptr) break; + if (v.logged) length++; + } + }; + + LogTopic(LogValue v0, LogValue v1 = {}, LogValue v2 = {}, LogValue v3 = {}, LogValue v4 = {}, LogValue v5 = {}, LogValue v6 = {}, LogValue v7 = {}, LogValue v8 = {}, LogValue v9 = {}) : + LogTopic(INFINITY, v0, v1, v2, v3, v4, v5, v6, v7, v8, v9) {}; }; -const int logColumns = sizeof(logEntries) / sizeof(logEntries[0]); -float logBuffer[LOG_SIZE][logColumns]; +LogTopic logTopics[] = { + // time + LogTopic({"t", &t}), // must be the first topic + LogTopic(1, {"loopRate", &loopRate}), -void prepareLogData() { - attitudeEuler = attitude.toEuler(); - attitudeTargetEuler = attitudeTarget.toEuler(); -} + // imu + LogTopic( + {"gyro.x", &gyro.x}, + {"gyro.y", &gyro.y}, + {"gyro.z", &gyro.z}), -void logData() { - if (!armed) return; - static int logPointer = 0; - static Rate period(LOG_RATE); - if (!period) return; + LogTopic(50, + {"acc.x", &acc.x}, + {"acc.y", &acc.y}, + {"acc.z", &acc.z}), - prepareLogData(); + LogTopic(10, + {"gyroBias.x", &gyroBias.x}, + {"gyroBias.y", &gyroBias.y}, + {"gyroBias.z", &gyroBias.z}), - for (int i = 0; i < logColumns; i++) { - logBuffer[logPointer][i] = *logEntries[i].value; - } + // estimation + LogTopic(50, + {"rates.x", &rates.x}, + {"rates.y", &rates.y}, + {"rates.z", &rates.z}, + {"attitude.roll", []() { return attitude.getRoll(); }}, + {"attitude.pitch", []() { return attitude.getPitch(); }}, + {"attitude.yaw", []() { return attitude.getYaw(); }}), - logPointer++; - if (logPointer >= LOG_SIZE) { - logPointer = 0; + // rc + LogTopic(10, + {"controlTime", &controlTime, false}, // trigger value + {"controlRoll", &controlRoll}, + {"controlPitch", &controlPitch}, + {"controlYaw", &controlYaw}, + {"controlThrottle", &controlThrottle}), + + // control + LogTopic({"armed", &armed}), + LogTopic({"mode", &mode}), + + LogTopic(10, + {"ratesTarget.x", &ratesTarget.x}, + {"ratesTarget.y", &ratesTarget.y}, + {"ratesTarget.z", &ratesTarget.z}, + {"attitudeTarget.roll", []() { return attitudeTarget.getRoll(); }}, + {"attitudeTarget.pitch", []() { return attitudeTarget.getPitch(); }}, + {"attitudeTarget.yaw", []() { return attitudeTarget.getYaw(); }}, + {"thrustTarget", &thrustTarget}), + + // motors + LogTopic( + {"motors[0]", &motors[0]}, + {"motors[1]", &motors[1]}, + {"motors[2]", &motors[2]}, + {"motors[3]", &motors[3]}), + + // misc + LogTopic(5, + {"voltage", &voltage}, + {"temp", &temperatureRead}, + {"imuTemp", []() { return imu.getTemp(); }}), +}; + +void *logBuffer; // buffer for log data +size_t logCapacity; +size_t logCursor = 0; +size_t logLength = 0; +LogValue *logExposed = nullptr; // log values exposed to telemetry + +void setupLog() { + print("Setup log\n"); + + free(logBuffer); // when reconfiguring + logBuffer = nullptr; + logCursor = 0; + logLength = 0; + + if (logMemory == 0) { + logCapacity = ESP.getFreeHeap() * logUsage; + logBuffer = (uint8_t *)calloc(logCapacity, 1); + } else if (logMemory == 1) { + logCapacity = ESP.getFreePsram() * logUsage; + logBuffer = (uint8_t *)heap_caps_calloc(logCapacity, 1, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); } } -void printLogHeader() { - for (int i = 0; i < logColumns; i++) { - print("%s%s", logEntries[i].name, i < logColumns - 1 ? "," : "\n"); - } -} +void loopLog() { + if (logBuffer == nullptr || !armed) return; -void printLogData() { - for (int i = 0; i < LOG_SIZE; i++) { - if (logBuffer[i][0] == 0) continue; // skip empty records - for (int j = 0; j < logColumns; j++) { - print("%g%s", logBuffer[i][j], j < logColumns - 1 ? "," : "\n"); + if (!logLength) resetLog(); // reset state on first log write + + static Rate sync(2); + if (sync) { + const uint8_t marker[] = {0x1A, 0x91, 0x4F, 0xF6, 0x7F}; + writeLog(&marker, sizeof(marker)); // write sync marker + } + + for (uint8_t i = 0; i < sizeof(logTopics) / sizeof(logTopics[0]); i++) { + LogTopic& topic = logTopics[i]; + if (t - topic.lastUpdate < 1 / topic.throttle) continue; // throttle topic + if (!isTopicUpdated(i)) continue; // skip if topic was't updated + + topic.lastUpdate = t; + writeLog(&i, sizeof(i)); // write topic index + + for (auto& value : topic.values) { + if (value.name == nullptr) break; + if (!value.logged) continue; + value.lastValue = value.value.get(); + writeLog(&value.lastValue, sizeof(float)); // write value } } } + +void resetLog() { + for (auto& topic : logTopics) { + topic.lastUpdate = -INFINITY; + for (auto& value : topic.values) { + value.lastValue = NAN; + } + } +} + +void writeLog(const void *data, size_t size) { + size_t first = min(size, logCapacity - logCursor); + size_t second = size - first; + memcpy(logBuffer + logCursor, data, first); + logCursor = (logCursor + first) % logCapacity; + if (second > 0) { + memcpy(logBuffer + logCursor, data + first, second); + logCursor = (logCursor + second) % logCapacity; + } + logLength = min(logLength + size, logCapacity); +} + +void readLog(void *data, size_t position, size_t size) { + if (logLength == logCapacity) { + position = (logCursor + position) % logCapacity; + } + size_t first = min(size, logCapacity - position); + size_t second = size - first; + memcpy(data, logBuffer + position, first); + if (second > 0) { + memcpy(data + first, logBuffer, second); + } +} + +bool isTopicUpdated(const uint8_t topic) { + LogTopic& logTopic = logTopics[topic]; + bool updated = false; + + for (auto& value : logTopic.values) { + if (value.name == nullptr) break; + float v = value.value.get(); + if (!floatEquals(value.lastValue, v)) { + value.lastValue = v; + updated = true; + } + } + return updated; +} + +void printLogInfo() { + if (logMemory == -1) return print("Log: disabled\n"); + + print("Memory: %s\n", logMemory == 0 ? "RAM" : "PSRAM"); + print("Usage: %.f%%\n", logUsage * 100); + print("Capacity: %u bytes\n", (unsigned)logCapacity); + print("Used: %u bytes\n", (unsigned)logLength); + print("Estimated duration: %d seconds\n", estimateLogDuration()); +} + +int estimateLogDuration() { + float bandwidth = 0; + for (LogTopic& topic : logTopics) { + float rate = isinf(topic.throttle) ? loopRate : topic.throttle; + bandwidth += rate * topic.length * sizeof(float); + } + return logCapacity / bandwidth; +} + +void printLogHeader() { + int i = 0; + for (auto& topic : logTopics) { + print("Topic #%d (%g Hz):\n", i++, topic.throttle); + for (auto& value : topic.values) { + if (value.name == nullptr) break; + print(" %s%s\n", value.name, value.logged ?"" : " (not logged)"); + } + } +} + +void printLogValues(const char *filter) { + for (LogTopic& topic : logTopics) { + for (LogValue& value : topic.values) { + if (value.name == nullptr) break; + if (strncasecmp(value.name, filter, strlen(filter))) continue; + print("%s = %g\n", value.name, value.value.get()); + } + } +} + +void configLogThrottle(const char *name, float throttle) { + for (LogTopic& topic : logTopics) { + for (LogValue& value : topic.values) { + if (value.name == nullptr) break; + if (strcasecmp(value.name, name) != 0) continue; + topic.throttle = throttle; + print("Log throttle for %s set to %.1f Hz\n", name, throttle); + return; + } + } + print("Log value not found: %s\n", name); +} + +void exposeLogValue(const char *name) { + for (int i = 0; i < sizeof(logTopics) / sizeof(logTopics[0]); i++) { + LogTopic& topic = logTopics[i]; + for (LogValue& value : topic.values) { + if (value.name == nullptr) break; + if (strcasecmp(value.name, name) != 0) continue; + logExposed = &value; + print("Log value %s exposed\n", name); + return; + } + } + print("Log value not found: %s\n", name); +} diff --git a/flix/mavlink.ino b/flix/mavlink.ino index d73203b..a8c343b 100644 --- a/flix/mavlink.ino +++ b/flix/mavlink.ino @@ -16,6 +16,7 @@ Rate telemetryAttitude(20); Rate telemetryRC(10); Rate telemetryMotors(10); Rate telemetryIMU(15); +Rate telemetryTopic(10); bool mavlinkConnected = false; String mavlinkPrintBuffer; @@ -84,6 +85,12 @@ void sendMavlink() { 0, 0, 0, 0); sendMessage(&msg); } + + if (telemetryTopic && logExposed != nullptr) { + mavlink_msg_named_value_float_pack(mavlinkSysId, MAV_COMP_ID_AUTOPILOT1, &msg, + time, logExposed->name, logExposed->value.get()); + sendMessage(&msg); + } } void sendMessage(const void *msg) { @@ -92,6 +99,26 @@ void sendMessage(const void *msg) { sendWiFi(buf, len); } +static uint8_t mavlinkBatch[ESP_NOW_MAX_DATA_LEN_V2]; +static int mavlinkBatchSize = 0; + +void batchMessage(const void *msg) { + uint8_t buf[MAVLINK_MAX_PACKET_LEN]; + int len = mavlink_msg_to_send_buffer(buf, (mavlink_message_t *)msg); + + if (mavlinkBatchSize + len > sizeof(mavlinkBatch)) { + sendWiFi(mavlinkBatch, mavlinkBatchSize); + mavlinkBatchSize = 0; + } + memcpy(mavlinkBatch + mavlinkBatchSize, buf, len); + mavlinkBatchSize += len; +} + +void flushBatchMessages() { + sendWiFi(mavlinkBatch, mavlinkBatchSize); + mavlinkBatchSize = 0; +} + void receiveMavlink() { uint8_t buf[MAVLINK_MAX_PACKET_LEN]; int len = receiveWiFi(buf, MAVLINK_MAX_PACKET_LEN); @@ -225,18 +252,29 @@ void handleMavlink(const void *_msg) { armed = motors[0] > 0 || motors[1] > 0 || motors[2] > 0 || motors[3] > 0; } + if (msg.msgid == MAVLINK_MSG_ID_LOG_REQUEST_LIST) { + const uint32_t qgcEpoch = 1262304000; // qgc accepts only timestamps after 2010-01-01 + mavlink_message_t response; + mavlink_msg_log_entry_pack(mavlinkSysId, MAV_COMP_ID_AUTOPILOT1, &response, + 0, 1, 0, qgcEpoch + t * 60, logLength); // put fake unique date to make qgc happy with saving logs + sendMessage(&response); + } + if (msg.msgid == MAVLINK_MSG_ID_LOG_REQUEST_DATA) { mavlink_log_request_data_t m; mavlink_msg_log_request_data_decode(&msg, &m); if (m.target_system && m.target_system != mavlinkSysId) return; - // Send all log records - for (int i = 0; i < sizeof(logBuffer) / sizeof(logBuffer[0]); i++) { - mavlink_message_t msg; - mavlink_msg_log_data_pack(mavlinkSysId, MAV_COMP_ID_AUTOPILOT1, &msg, 0, i, - sizeof(logBuffer[0]), (uint8_t *)logBuffer[i]); - sendMessage(&msg); + for (int i = 0; i < m.count; i += MAVLINK_MSG_LOG_DATA_FIELD_DATA_LEN) { + int chunkSize = min(MAVLINK_MSG_LOG_DATA_FIELD_DATA_LEN, (int)(m.count - i)); + mavlink_message_t response; + uint8_t data[MAVLINK_MSG_LOG_DATA_FIELD_DATA_LEN]; + readLog(data, m.ofs + i, chunkSize); + mavlink_msg_log_data_pack(mavlinkSysId, MAV_COMP_ID_AUTOPILOT1, &response, + m.id, m.ofs + i, chunkSize, data); + batchMessage(&response); } + flushBatchMessages(); } // Handle commands diff --git a/flix/parameters.ino b/flix/parameters.ino index 0d3e19a..bd0b9ea 100644 --- a/flix/parameters.ino +++ b/flix/parameters.ino @@ -121,6 +121,22 @@ Parameter parameters[] = { {"MAV_RATE_RC", &telemetryRC.rate}, {"MAV_RATE_MOT", &telemetryMotors.rate}, {"MAV_RATE_IMU", &telemetryIMU.rate}, + {"MAV_RATE_TOPIC", &telemetryTopic.rate}, + // log + {"LOG_MEMORY", &logMemory, setupLog}, + {"LOG_USAGE", &logUsage, setupLog}, + {"LOG_RATE_000", &logTopics[0].throttle}, + {"LOG_RATE_001", &logTopics[1].throttle}, + {"LOG_RATE_002", &logTopics[2].throttle}, + {"LOG_RATE_003", &logTopics[3].throttle}, + {"LOG_RATE_004", &logTopics[4].throttle}, + {"LOG_RATE_005", &logTopics[5].throttle}, + {"LOG_RATE_006", &logTopics[6].throttle}, + {"LOG_RATE_007", &logTopics[7].throttle}, + {"LOG_RATE_008", &logTopics[8].throttle}, + {"LOG_RATE_009", &logTopics[9].throttle}, + {"LOG_RATE_010", &logTopics[10].throttle}, + {"LOG_RATE_011", &logTopics[11].throttle}, // power {"PWR_VOLT_PIN", &voltagePin, setupPower}, {"PWR_VOLT_SCALE", &voltageScale}, diff --git a/flix/util.h b/flix/util.h index 4db65a4..f3971aa 100644 --- a/flix/util.h +++ b/flix/util.h @@ -64,6 +64,49 @@ public: } }; +// Simple variant type for logging and parameters +struct Value { + enum { EMPTY, FLOAT, INT, BOOL, FLOAT_FN, INT_FN, BOOL_FN } type; + union { + void *pointer; + float *_float; + int *_int; + bool *_bool; + float (*floatFn)(); + int (*intFn)(); + bool (*boolFn)(); + }; + + Value() : type(EMPTY), pointer(nullptr) {}; + Value(float *pt) : type(FLOAT), _float(pt) {}; + Value(int *pt) : type(INT), _int(pt) {}; + Value(bool *pt) : type(BOOL), _bool(pt) {}; + Value(float (*fn)()) : type(FLOAT_FN), floatFn(fn) {}; + Value(int (*fn)()) : type(INT_FN), intFn(fn) {}; + Value(bool (*fn)()) : type(BOOL_FN), boolFn(fn) {}; + + float get() const { + switch (type) { + case FLOAT: return *_float; + case INT: return *_int; + case BOOL: return *_bool ? 1 : 0; + case FLOAT_FN: return floatFn(); + case INT_FN: return intFn(); + case BOOL_FN: return boolFn() ? 1 : 0; + default: return NAN; + } + }; + + void set(float value) const { + switch (type) { + case FLOAT: *_float = value; break; + case INT: *_int = value; break; + case BOOL: *_bool = (value != 0); break; + default: break; + } + }; +}; + // Rate limiter class Rate { public: diff --git a/gazebo/Arduino.h b/gazebo/Arduino.h index 6122225..b9b03dd 100644 --- a/gazebo/Arduino.h +++ b/gazebo/Arduino.h @@ -156,6 +156,7 @@ HardwareSerial Serial, Serial1, Serial2; class EspClass { public: void restart() { Serial.println("Ignore reboot in simulation"); } + uint32_t getFreeHeap() { return 300 * 1024; } // assume 300 KB free heap } ESP; unsigned long __delayTime = 0; diff --git a/gazebo/flix.h b/gazebo/flix.h index e0ba9ab..f8ab659 100644 --- a/gazebo/flix.h +++ b/gazebo/flix.h @@ -48,8 +48,17 @@ void normalizeRC(); void calibrateRC(); void calibrateRCChannel(int*, uint16_t[16], uint16_t[16], const char*); void printRCCalibration(); +void loopLog(); +void resetLog(); +void writeLog(const void *data, size_t size); +void readLog(void *data, size_t position, size_t size); +bool isTopicUpdated(const uint8_t topic); +void printLogInfo(); +int estimateLogDuration(); void printLogHeader(); -void printLogData(); +void printLogValues(const char *filter); +void configLogThrottle(const char *name, float throttle); +void exposeLogValue(const char *name); void processMavlink(); void sendMavlink(); void sendMessage(const void *msg); diff --git a/gazebo/simulator.cpp b/gazebo/simulator.cpp index ca2776c..44b1826 100644 --- a/gazebo/simulator.cpp +++ b/gazebo/simulator.cpp @@ -55,6 +55,7 @@ public: initNode(); Serial.begin(0); setupParameters(); + setupLog(); gzmsg << "Flix plugin loaded" << endl; } diff --git a/tools/csv_to_mcap.py b/tools/csv_to_mcap.py index f881565..06898f1 100755 --- a/tools/csv_to_mcap.py +++ b/tools/csv_to_mcap.py @@ -9,6 +9,7 @@ Usage: import csv import json import docopt +import math from mcap.writer import Writer args = docopt.docopt(__doc__) @@ -39,7 +40,15 @@ channel_id = writer.register_channel( ) for row in csv_reader: - data = {key: float(value) for key, value in zip(header, row)} + if row[0] == '': continue + data = {} + for key, value in zip(header, row): + if value == '' or math.isnan(float(value)): + data[key] = None + else: + data[key] = float(value) + + data = {key: float(value) if value != '' else None for key, value in zip(header, row)} timestamp = round(float(row[0]) * 1e9) writer.add_message(channel_id=channel_id, log_time=timestamp, data=json.dumps(data).encode(), publish_time=timestamp,) diff --git a/tools/log_to_csv.py b/tools/log_to_csv.py new file mode 100755 index 0000000..b1b0361 --- /dev/null +++ b/tools/log_to_csv.py @@ -0,0 +1,76 @@ +#!/usr/bin/env python3 + +"""Convert flight from Flix format to CSV + +Usage: + log_to_csv.py +""" + +import os +from pyflix import Flix +import docopt +import struct +import csv + +DIR = os.path.dirname(os.path.realpath(__file__)) +HEADER_FILE = os.path.join(DIR, 'log/log_header.txt') + +# Read log header +try: + # read from file + header = open(HEADER_FILE, 'r').read() +except FileNotFoundError: + flix = Flix() + header = flix.cli('log header') # receive the log schema + open(HEADER_FILE, 'w').write(header) # save to file + +# Parse log header +topics = [] +for line in header.splitlines(): + if not line.startswith(' '): + topics.append([]) + elif 'not logged' not in line: + topics[-1].append(line.strip()) + +# Read log file +args = docopt.docopt(__doc__) +input_file = args[''] +outfile_file = input_file + '.csv' +data = open(input_file, 'rb').read() + +# Search for sync marker +SYNC_MARKER = bytes.fromhex('1A 91 4F F6 7F') +sync_offset = data.find(SYNC_MARKER) +if sync_offset == -1: + raise ValueError('Sync marker not found in log file') +data = data[sync_offset + len(SYNC_MARKER):] +data = data.replace(SYNC_MARKER, b'') # remove all other sync markers + +header_row = [f'{value}' for topic in topics for value in topic] +rows = [] + +offset = 0 +while offset < len(data): + try: + topic = struct.unpack_from('B', data, offset)[0] + offset += 1 + + if topic >= len(topics): + raise ValueError(f'Invalid topic {topic} at offset {offset}') + + if topic == 0 or not rows: + rows.append({}) + + for name in topics[topic]: + value = struct.unpack_from('