From 060b8f084d47e9900eaa4f262dcf4f25f073fed1 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sun, 26 May 2024 21:26:39 +0200 Subject: [PATCH 01/14] Enable log_handler in sunnylinkd.py The 'log_handler' thread in sunnylinkd.py was previously commented out and has now been enabled for execution. This update will allow the 'log_handler' to perform its task in the thread execution sequence. --- selfdrive/athena/sunnylinkd.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index 9c26dd42a1..b69bc7a762 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -11,7 +11,7 @@ import threading import time from openpilot.selfdrive.athena.athenad import ws_send, jsonrpc_handler, \ - recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage + recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler from jsonrpc import dispatcher from websocket import (ABNF, WebSocket, WebSocketException, WebSocketTimeoutException, create_connection) @@ -37,7 +37,7 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: threading.Thread(target=ws_ping, args=(ws, end_event), name='ws_ping'), threading.Thread(target=ws_queue, args=(end_event,), name='ws_queue'), # threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler'), - # threading.Thread(target=log_handler, args=(end_event,), name='log_handler'), + threading.Thread(target=log_handler, args=(end_event,), name='log_handler'), # threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), ] + [ threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') From a366d5b873c709d244db20c374d853bc32173ad9 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sun, 26 May 2024 22:04:06 +0200 Subject: [PATCH 02/14] Add custom attribute name for Sunnylink logs The commit introduces a custom attribute name for Sunnylink log entries. The custom attribute name 'sunnylink.user.upload' is used for all operations in the athenad.py and sunnylinkd.py files that previously referred to the default LOG_ATTR_NAME. This will allow more flexibility in handling logs specific to Sunnylink. --- selfdrive/athena/athenad.py | 12 ++++++------ selfdrive/athena/sunnylinkd.py | 5 +++-- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index d5b1e37ce1..8e95e3a4aa 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -550,7 +550,7 @@ def takeSnapshot() -> str | dict[str, str] | None: raise Exception("not available while camerad is started") -def get_logs_to_send_sorted() -> list[str]: +def get_logs_to_send_sorted(log_attr_name=LOG_ATTR_NAME) -> list[str]: # TODO: scan once then use inotify to detect file creation/deletion curr_time = int(time.time()) logs = [] @@ -558,7 +558,7 @@ def get_logs_to_send_sorted() -> list[str]: log_path = os.path.join(Paths.swaglog_root(), log_entry) time_sent = 0 try: - value = getxattr(log_path, LOG_ATTR_NAME) + value = getxattr(log_path, log_attr_name) if value is not None: time_sent = int.from_bytes(value, sys.byteorder) except (ValueError, TypeError): @@ -570,7 +570,7 @@ def get_logs_to_send_sorted() -> list[str]: return sorted(logs)[:-1] -def log_handler(end_event: threading.Event) -> None: +def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None: if PC: return @@ -580,7 +580,7 @@ def log_handler(end_event: threading.Event) -> None: try: curr_scan = time.monotonic() if curr_scan - last_scan > 10: - log_files = get_logs_to_send_sorted() + log_files = get_logs_to_send_sorted(log_attr_name) last_scan = curr_scan # send one log @@ -591,7 +591,7 @@ def log_handler(end_event: threading.Event) -> None: try: curr_time = int(time.time()) log_path = os.path.join(Paths.swaglog_root(), log_entry) - setxattr(log_path, LOG_ATTR_NAME, int.to_bytes(curr_time, 4, sys.byteorder)) + setxattr(log_path, log_attr_name, int.to_bytes(curr_time, 4, sys.byteorder)) with open(log_path) as f: jsonrpc = { "method": "forwardLogs", @@ -619,7 +619,7 @@ def log_handler(end_event: threading.Event) -> None: if log_entry and log_success: log_path = os.path.join(Paths.swaglog_root(), log_entry) try: - setxattr(log_path, LOG_ATTR_NAME, LOG_ATTR_VALUE_MAX_UNIX_TIME) + setxattr(log_path, log_attr_name, LOG_ATTR_VALUE_MAX_UNIX_TIME) except OSError: pass # file could be deleted by log rotation if curr_log == log_entry: diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index b69bc7a762..6b9028bf20 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -11,7 +11,7 @@ import threading import time from openpilot.selfdrive.athena.athenad import ws_send, jsonrpc_handler, \ - recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler + recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler, LOG_ATTR_NAME from jsonrpc import dispatcher from websocket import (ABNF, WebSocket, WebSocketException, WebSocketTimeoutException, create_connection) @@ -24,6 +24,7 @@ from openpilot.common.swaglog import cloudlog SUNNYLINK_ATHENA_HOST = os.getenv('SUNNYLINK_ATHENA_HOST', 'wss://ws.stg.api.sunnypilot.ai') HANDLER_THREADS = int(os.getenv('HANDLER_THREADS', "4")) LOCAL_PORT_WHITELIST = {8022} +SUNNYLINK_LOG_ATTR_NAME = "sunnylink.user.upload" params = Params() sunnylink_api = SunnylinkApi(params.get("SunnylinkDongleId", encoding='utf-8')) @@ -37,7 +38,7 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: threading.Thread(target=ws_ping, args=(ws, end_event), name='ws_ping'), threading.Thread(target=ws_queue, args=(end_event,), name='ws_queue'), # threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler'), - threading.Thread(target=log_handler, args=(end_event,), name='log_handler'), + threading.Thread(target=log_handler, args=(end_event,SUNNYLINK_LOG_ATTR_NAME,), name='log_handler'), # threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), ] + [ threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') From f1a31495516a65c65cc8c1068df378192a1c93d3 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sun, 26 May 2024 22:18:24 +0200 Subject: [PATCH 03/14] Update attribute name for Sunnylink log upload The attribute name used for Sunnylink's log upload function was modified. The new attribute name "user.sunny.upload" replaces the previous "sunnylink.user.upload" to reflect recent changes in naming convention. --- selfdrive/athena/sunnylinkd.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index 6b9028bf20..1b461dfbbd 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -24,7 +24,7 @@ from openpilot.common.swaglog import cloudlog SUNNYLINK_ATHENA_HOST = os.getenv('SUNNYLINK_ATHENA_HOST', 'wss://ws.stg.api.sunnypilot.ai') HANDLER_THREADS = int(os.getenv('HANDLER_THREADS', "4")) LOCAL_PORT_WHITELIST = {8022} -SUNNYLINK_LOG_ATTR_NAME = "sunnylink.user.upload" +SUNNYLINK_LOG_ATTR_NAME = "user.sunny.upload" params = Params() sunnylink_api = SunnylinkApi(params.get("SunnylinkDongleId", encoding='utf-8')) From b56692bd2e8e4ed575904c177ecd32a621e3b8d2 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sun, 26 May 2024 22:33:34 +0200 Subject: [PATCH 04/14] Add platform-specific xattr handling in loggerd This commit adapts the xattr handling in the loggerd module to individual platforms, specifically macOS and others. It imports specific xattr functions based on the running system, enhancing compatibility and reducing the risk of errors. In specific, 'ENOATTR' error is now also taken into account which is relevant in some non-Linux platforms. --- system/loggerd/xattr_cache.py | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/system/loggerd/xattr_cache.py b/system/loggerd/xattr_cache.py index d3220118ac..0f7e8a8a52 100644 --- a/system/loggerd/xattr_cache.py +++ b/system/loggerd/xattr_cache.py @@ -1,5 +1,13 @@ import os import errno +import platform + +if platform.system() == 'Darwin': # macOS + from xattr import getxattr as _getxattr + from xattr import setxattr as _setxattr +else: + from os import getxattr as _getxattr + from os import setxattr as _setxattr _cached_attributes: dict[tuple, bytes | None] = {} @@ -7,10 +15,10 @@ def getxattr(path: str, attr_name: str) -> bytes | None: key = (path, attr_name) if key not in _cached_attributes: try: - response = os.getxattr(path, attr_name) + response = _getxattr(path, attr_name) except OSError as e: # ENODATA means attribute hasn't been set - if e.errno == errno.ENODATA: + if e.errno == errno.ENODATA or e.errno == errno.ENOATTR: response = None else: raise @@ -19,4 +27,4 @@ def getxattr(path: str, attr_name: str) -> bytes | None: def setxattr(path: str, attr_name: str, attr_value: bytes) -> None: _cached_attributes.pop((path, attr_name), None) - return os.setxattr(path, attr_name, attr_value) + return _setxattr(path, attr_name, attr_value) From 86ce2bc79135cd54ea94263eddbf5e6e138ca949 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sun, 26 May 2024 23:55:41 +0200 Subject: [PATCH 05/14] Remove unused import in sunnylinkd.py The unused import, LOG_ATTR_NAME, in the file selfdrive/athena/sunnylinkd.py has been removed to provide cleaner, simpler, and more readable code. This also helps follow good coding practices such as removing unnecessary or unused imports or resources. --- selfdrive/athena/sunnylinkd.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index 1b461dfbbd..6eaec38181 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -11,7 +11,7 @@ import threading import time from openpilot.selfdrive.athena.athenad import ws_send, jsonrpc_handler, \ - recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler, LOG_ATTR_NAME + recv_queue, RECONNECT_TIMEOUT_S, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler from jsonrpc import dispatcher from websocket import (ABNF, WebSocket, WebSocketException, WebSocketTimeoutException, create_connection) From e7b40f61a8022cdacccb2d276be7bc923fe79178 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Mon, 27 May 2024 19:44:41 +0200 Subject: [PATCH 06/14] Add support for network metering and PrimeType checking in sunnylinkd This commit introduces network metering and PrimeType checking in sunnylinkd. The implementation includes an additional threading event and a new log handler function, 'sunny_log_handler'. It defines PrimeType and metering conditions to set and clear the new threading event as part of the main processing loop in handle_long_poll function. --- selfdrive/athena/sunnylinkd.py | 26 +++++++++++++++++++++++++- 1 file changed, 25 insertions(+), 1 deletion(-) diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index 6eaec38181..6ebf6bf653 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -20,6 +20,7 @@ from openpilot.common.api import SunnylinkApi from openpilot.common.params import Params from openpilot.common.realtime import set_core_affinity from openpilot.common.swaglog import cloudlog +import cereal.messaging as messaging SUNNYLINK_ATHENA_HOST = os.getenv('SUNNYLINK_ATHENA_HOST', 'wss://ws.stg.api.sunnypilot.ai') HANDLER_THREADS = int(os.getenv('HANDLER_THREADS', "4")) @@ -29,7 +30,9 @@ SUNNYLINK_LOG_ATTR_NAME = "user.sunny.upload" params = Params() sunnylink_api = SunnylinkApi(params.get("SunnylinkDongleId", encoding='utf-8')) def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: + sm = messaging.SubMaster(['deviceState']) end_event = threading.Event() + comma_prime_cellular_end_event = threading.Event() threads = [ threading.Thread(target=ws_manage, args=(ws, end_event), name='ws_manage'), @@ -38,7 +41,7 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: threading.Thread(target=ws_ping, args=(ws, end_event), name='ws_ping'), threading.Thread(target=ws_queue, args=(end_event,), name='ws_queue'), # threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler'), - threading.Thread(target=log_handler, args=(end_event,SUNNYLINK_LOG_ATTR_NAME,), name='log_handler'), + threading.Thread(target=sunny_log_handler, args=(end_event, comma_prime_cellular_end_event), name='log_handler'), # threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), ] + [ threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') @@ -49,10 +52,24 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: thread.start() try: while not end_event.wait(0.1): + sm.update(0) if exit_event is not None and exit_event.is_set(): end_event.set() + comma_prime_cellular_end_event.set() + + prime_type = params.get("PrimeType", encoding='utf-8') + metered = sm['deviceState'].networkMetered + + if int(prime_type) > 2 and metered: + cloudlog.info(f"sunnylinkd.handle_long_poll: PrimeType({prime_type}) > 2 and networkMetered({metered})") + comma_prime_cellular_end_event.set() + elif comma_prime_cellular_end_event.is_set(): + cloudlog.info(f"sunnylinkd.handle_long_poll: comma_prime_cellular_end_event is set and not PrimeType({prime_type}) > 2 or not networkMetered({metered})") + comma_prime_cellular_end_event.clear() + except (KeyboardInterrupt, SystemExit): end_event.set() + comma_prime_cellular_end_event.set() raise finally: for thread in threads: @@ -114,6 +131,13 @@ def ws_queue(end_event: threading.Event) -> None: cloudlog.debug("Resume requested or end_event is set, exiting ws_queue thread") +def sunny_log_handler(end_event: threading.Event, comma_prime_cellular_end_event: threading.Event) -> None: + while not end_event.wait(0.1): + if not comma_prime_cellular_end_event.is_set(): + log_handler(comma_prime_cellular_end_event, SUNNYLINK_LOG_ATTR_NAME) + comma_prime_cellular_end_event.set() + + @dispatcher.add_method def getParamsAllKeys() -> list[str]: keys: list[str] = [k.decode('utf-8') for k in Params().all_keys()] From 9d79f243ff1068c46215e7b3149e0416ead56d02 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Mon, 27 May 2024 21:51:00 +0200 Subject: [PATCH 07/14] Refactor log handling and add exception handling The changes refactor the log handling in athenad.py for more efficient processing. The log is now split into maximum chunk sizes of 128KB to prevent oversized requests. Additionally, an invalid request or response exception has been introduced for better error handling. --- selfdrive/athena/athenad.py | 39 +++++++++++++++++++++++++------------ 1 file changed, 27 insertions(+), 12 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index 8e95e3a4aa..aaee9fca3c 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -177,7 +177,7 @@ def jsonrpc_handler(end_event: threading.Event) -> None: send_queue.put_nowait(response.json) elif "id" in data and ("result" in data or "error" in data): log_recv_queue.put_nowait(data) - else: + elif data: raise Exception("not a valid request or response") except queue.Empty: pass @@ -570,6 +570,28 @@ def get_logs_to_send_sorted(log_attr_name=LOG_ATTR_NAME) -> list[str]: return sorted(logs)[:-1] +def add_log_to_queue(log_path, id): + # Define the maximum size of a chunk in bytes + MAX_CHUNK_SIZE = 128 * 1024 # 128KB + + # Open the log file + with open(log_path, 'r') as f: + while True: + chunk = f.read(MAX_CHUNK_SIZE) + if not chunk: + break + + jsonrpc = { + "method": "forwardLogs", + "params": { + "logs": chunk + }, + "jsonrpc": "2.0", + "id": id + } + low_priority_send_queue.put_nowait(json.dumps(jsonrpc)) + + def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None: if PC: return @@ -592,17 +614,10 @@ def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None curr_time = int(time.time()) log_path = os.path.join(Paths.swaglog_root(), log_entry) setxattr(log_path, log_attr_name, int.to_bytes(curr_time, 4, sys.byteorder)) - with open(log_path) as f: - jsonrpc = { - "method": "forwardLogs", - "params": { - "logs": f.read() - }, - "jsonrpc": "2.0", - "id": log_entry - } - low_priority_send_queue.put_nowait(json.dumps(jsonrpc)) - curr_log = log_entry + # now we need to check the log size because we cant go larger than 128kb per request so we need to split by regex for example regex.compile(r'\{(?:[^{}]|(?R))*\}') + + add_log_to_queue(log_path, log_entry) + curr_log = log_entry except OSError: pass # file could be deleted by log rotation From 797e273ba7e16f42ff09225ddd17f776e53b050d Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Mon, 27 May 2024 22:19:39 +0200 Subject: [PATCH 08/14] Reduce maximum chunk size in athenad.py Updating the maximum chunk size constant in athenad.py to decrease its value from 128KB to 32KB. This change has been made to optimize the handling of log files and data transmission. --- selfdrive/athena/athenad.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index aaee9fca3c..5fecb253a8 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -572,7 +572,7 @@ def get_logs_to_send_sorted(log_attr_name=LOG_ATTR_NAME) -> list[str]: def add_log_to_queue(log_path, id): # Define the maximum size of a chunk in bytes - MAX_CHUNK_SIZE = 128 * 1024 # 128KB + MAX_CHUNK_SIZE = 32 * 1024 # 32KB # Open the log file with open(log_path, 'r') as f: From 091e5d1343f393cf9215a63e6c3838ee97063407 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Mon, 27 May 2024 22:51:27 +0200 Subject: [PATCH 09/14] Adjust maximum chunk size in athenad.py This commit decreases the maximum chunk size from 32KB to 28KB in the add_log_to_queue function in athenad.py. This size reduction will affect how the log files are loaded and processed. --- selfdrive/athena/athenad.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index 5fecb253a8..0935c8d218 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -572,7 +572,7 @@ def get_logs_to_send_sorted(log_attr_name=LOG_ATTR_NAME) -> list[str]: def add_log_to_queue(log_path, id): # Define the maximum size of a chunk in bytes - MAX_CHUNK_SIZE = 32 * 1024 # 32KB + MAX_CHUNK_SIZE = 28 * 1024 # 32KB # Open the log file with open(log_path, 'r') as f: From 23742d046f952c4b75edbcf80dd6b83cd4bbdc72 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 28 May 2024 02:26:26 +0000 Subject: [PATCH 10/14] Disable Screen Recorder --- CHANGELOGS.md | 2 ++ selfdrive/ui/SConscript | 2 +- .../ui/qt/offroad/sunnypilot/display_settings.cc | 2 ++ selfdrive/ui/qt/screenrecorder/omx_encoder.cc | 16 +++++++++++++++- 4 files changed, 20 insertions(+), 2 deletions(-) diff --git a/CHANGELOGS.md b/CHANGELOGS.md index 850013c029..2fc0388fa0 100644 --- a/CHANGELOGS.md +++ b/CHANGELOGS.md @@ -57,6 +57,8 @@ sunnypilot - 0.9.7.0 (2024-05-xx) * NEW❗: Metrics is now being displayed below the chevron instead of above * NEW❗: Display both Distance and Speed simultaneously * NEW❗: View sunnylink connectivity status on the left sidebar! +* REMOVED: Screen Recorder + * A better/equivalent version of Screen Recorder will be supported in the near future. Stay tuned! sunnypilot - 0.9.6.1 (2024-02-27) ======================== diff --git a/selfdrive/ui/SConscript b/selfdrive/ui/SConscript index 155e790e7c..3c7662e6ce 100644 --- a/selfdrive/ui/SConscript +++ b/selfdrive/ui/SConscript @@ -10,7 +10,7 @@ if arch == 'larch64': base_libs.append('EGL') maps = arch in ['larch64', 'aarch64', 'x86_64'] -screenrecorder = arch in ['larch64'] +screenrecorder = False # arch in ['larch64'] if arch == "Darwin": del base_libs[base_libs.index('OpenCL')] diff --git a/selfdrive/ui/qt/offroad/sunnypilot/display_settings.cc b/selfdrive/ui/qt/offroad/sunnypilot/display_settings.cc index 3d83fce252..a09bc22565 100644 --- a/selfdrive/ui/qt/offroad/sunnypilot/display_settings.cc +++ b/selfdrive/ui/qt/offroad/sunnypilot/display_settings.cc @@ -13,12 +13,14 @@ DisplayPanel::DisplayPanel(QWidget *parent) : ListWidget(parent, false) { .arg(tr("Enabled: Wake the brightness of the screen to display all events.")) .arg(tr("Disabled: Wake the brightness of the screen to display critical events.")), "../assets/offroad/icon_blank.png", +#ifdef ENABLE_DASHCAM }, { "ScreenRecorder", tr("Enable Screen Recorder"), tr("Enable this will display a button on the onroad screen to toggle on or off real-time screen recording with UI elements."), "../assets/offroad/icon_blank.png" +#endif } }; diff --git a/selfdrive/ui/qt/screenrecorder/omx_encoder.cc b/selfdrive/ui/qt/screenrecorder/omx_encoder.cc index d051b600f8..29607ed921 100644 --- a/selfdrive/ui/qt/screenrecorder/omx_encoder.cc +++ b/selfdrive/ui/qt/screenrecorder/omx_encoder.cc @@ -10,6 +10,8 @@ #include #include +#include + #include #include #include @@ -18,6 +20,7 @@ #include "msm_media_info.h" #include "common/swaglog.h" #include "common/util.h" +#include "common/watchdog.h" using namespace libyuv; @@ -327,8 +330,19 @@ OmxEncoder::OmxEncoder(const char* path, int width, int height, int fps, int bit int err = OMX_GetHandle(&this->handle, component, this, &omx_callbacks); if (err != OMX_ErrorNone) { LOGE("error getting codec: %x", err); + // TODO: We force quit Qt UI, and trigger fast restart for now + qApp->exit(18); + watchdog_kick(0); } - assert(err == OMX_ErrorNone); + // TODO: Investigate the insufficient memory issue in prebuilts, we force quit Qt UI, and trigger fast restart for now + /* + if (err == OMX_ErrorInsufficientResources) { + qApp->exit(18); + watchdog_kick(0); + } else { + assert(err == OMX_ErrorNone); + } + */ // printf("handle: %p\n", this->handle); // setup input port From 753bcce9a73919b4da4326433d3914edfbfd0583 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 28 May 2024 12:20:19 +0200 Subject: [PATCH 11/14] Implement compression and encoding for large log files The `add_log_to_queue` function has been updated to compress and base64 encode log files that exceed a certain size, specifically for the "sunnylink" scenario. This function also provides logging for various stages of the process, including initial file size, size post-compression/encoding, and whether the final payload is small enough to be sent in one request. The change will improve handle of larger log files and prevent payload size related issues. --- selfdrive/athena/athenad.py | 74 ++++++++++++++++++++++++++++--------- 1 file changed, 56 insertions(+), 18 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index 0935c8d218..ee47e53bea 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -16,6 +16,7 @@ import sys import tempfile import threading import time +import gzip from dataclasses import asdict, dataclass, replace from datetime import datetime from functools import partial @@ -570,29 +571,66 @@ def get_logs_to_send_sorted(log_attr_name=LOG_ATTR_NAME) -> list[str]: return sorted(logs)[:-1] -def add_log_to_queue(log_path, id): - # Define the maximum size of a chunk in bytes - MAX_CHUNK_SIZE = 28 * 1024 # 32KB +def add_log_to_queue(log_path, log_id, is_sunnylink = False): + MAX_SIZE_KB = 32 + MAX_SIZE_BYTES = MAX_SIZE_KB * 1024 - # Open the log file with open(log_path, 'r') as f: - while True: - chunk = f.read(MAX_CHUNK_SIZE) - if not chunk: - break + data = f.read() - jsonrpc = { - "method": "forwardLogs", - "params": { - "logs": chunk - }, - "jsonrpc": "2.0", - "id": id - } - low_priority_send_queue.put_nowait(json.dumps(jsonrpc)) + # Check if the file is empty + if not data: + cloudlog.warning(f"Log file {log_path} is empty.") + return + + # Log the current size of the file + current_size = os.path.getsize(log_path) + cloudlog.info(f"Current size of log file {log_path}: {current_size} bytes") + + # Initialize variables for encoding + payload = data + is_compressed = False + + if is_sunnylink and current_size > MAX_SIZE_BYTES: + # Compress and encode the data if it exceeds the maximum size + compressed_data = gzip.compress(data.encode()) + payload = base64.b64encode(compressed_data).decode() + is_compressed = True + + # Log the size after compression and encoding + compressed_size = len(compressed_data) + encoded_size = len(payload) + cloudlog.info(f"Size of log file {log_path} " + f"after compression: {compressed_size} bytes, " + f"after encoding: {encoded_size} bytes") + + jsonrpc = { + "method": "forwardLogs", + "params": { + "logs": payload + }, + "jsonrpc": "2.0", + "id": log_id + } + + if is_sunnylink and is_compressed: + jsonrpc["params"]["compressed"] = is_compressed + + jsonrpc_str = json.dumps(jsonrpc) + size_in_bytes = len(jsonrpc_str.encode('utf-8')) + + if is_sunnylink and size_in_bytes <= MAX_SIZE_BYTES: + cloudlog.info(f"Target is sunnylink and log file {log_path} is small enough to send in one request ({size_in_bytes} bytes).") + low_priority_send_queue.put_nowait(jsonrpc_str) + elif is_sunnylink: + cloudlog.warning(f"Target is sunnylink and log file {log_path} is too large to send in one request.") + else: + cloudlog.info(f"Target is not sunnylink, proceeding to send log file {log_path} in one request ({size_in_bytes} bytes).") + low_priority_send_queue.put_nowait(jsonrpc_str) def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None: + is_sunnylink = log_attr_name != LOG_ATTR_NAME if PC: return @@ -616,7 +654,7 @@ def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None setxattr(log_path, log_attr_name, int.to_bytes(curr_time, 4, sys.byteorder)) # now we need to check the log size because we cant go larger than 128kb per request so we need to split by regex for example regex.compile(r'\{(?:[^{}]|(?R))*\}') - add_log_to_queue(log_path, log_entry) + add_log_to_queue(log_path, log_entry, is_sunnylink) curr_log = log_entry except OSError: pass # file could be deleted by log rotation From 52ea16a93a4ef75698565a06fe28b191250dc8eb Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 28 May 2024 14:14:22 +0200 Subject: [PATCH 12/14] Update log file size calculation method The code for calculating the size of the log file has been revised. Instead of using os.path.getsize, which simply captures the size of the file on disk, we're now calculating the size of the payload after it has been serialized and encoded, plus an overhead of 100 bytes. This change provides a more accurate measure of the data that will be sent. --- selfdrive/athena/athenad.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index ee47e53bea..8ec1094124 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -583,14 +583,14 @@ def add_log_to_queue(log_path, log_id, is_sunnylink = False): cloudlog.warning(f"Log file {log_path} is empty.") return - # Log the current size of the file - current_size = os.path.getsize(log_path) - cloudlog.info(f"Current size of log file {log_path}: {current_size} bytes") - # Initialize variables for encoding payload = data is_compressed = False + # Log the current size of the file + current_size = len(json.dumps(payload).encode("utf-8")) + len(log_id.encode("utf-8")) + 100 # Add 100 bytes to account for encoding overhead + cloudlog.info(f"Current size of log file {log_path}: {current_size} bytes") + if is_sunnylink and current_size > MAX_SIZE_BYTES: # Compress and encode the data if it exceeds the maximum size compressed_data = gzip.compress(data.encode()) From 2d2a35c7e4b1111f012199a799a4e5b6fcee1601 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 28 May 2024 12:54:11 +0000 Subject: [PATCH 13/14] Apply 2 suggestion(s) to 1 file(s) --- selfdrive/athena/athenad.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index 8ec1094124..984b724aa6 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -178,7 +178,7 @@ def jsonrpc_handler(end_event: threading.Event) -> None: send_queue.put_nowait(response.json) elif "id" in data and ("result" in data or "error" in data): log_recv_queue.put_nowait(data) - elif data: + else: raise Exception("not a valid request or response") except queue.Empty: pass @@ -652,7 +652,6 @@ def log_handler(end_event: threading.Event, log_attr_name=LOG_ATTR_NAME) -> None curr_time = int(time.time()) log_path = os.path.join(Paths.swaglog_root(), log_entry) setxattr(log_path, log_attr_name, int.to_bytes(curr_time, 4, sys.byteorder)) - # now we need to check the log size because we cant go larger than 128kb per request so we need to split by regex for example regex.compile(r'\{(?:[^{}]|(?R))*\}') add_log_to_queue(log_path, log_entry, is_sunnylink) curr_log = log_entry From 0c909e98bbddd798b965bbb0c8aa366edce89bf3 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 28 May 2024 14:54:35 +0200 Subject: [PATCH 14/14] Spacing --- system/loggerd/xattr_cache.py | 25 +++++++++++++++++++++++-- 1 file changed, 23 insertions(+), 2 deletions(-) diff --git a/system/loggerd/xattr_cache.py b/system/loggerd/xattr_cache.py index 0f7e8a8a52..8d950f71a2 100644 --- a/system/loggerd/xattr_cache.py +++ b/system/loggerd/xattr_cache.py @@ -3,8 +3,8 @@ import errno import platform if platform.system() == 'Darwin': # macOS - from xattr import getxattr as _getxattr - from xattr import setxattr as _setxattr + from xattr import getxattr as _getxattr + from xattr import setxattr as _setxattr else: from os import getxattr as _getxattr from os import setxattr as _setxattr @@ -28,3 +28,24 @@ def getxattr(path: str, attr_name: str) -> bytes | None: def setxattr(path: str, attr_name: str, attr_value: bytes) -> None: _cached_attributes.pop((path, attr_name), None) return _setxattr(path, attr_name, attr_value) + + + + +import os + +# Get the list of all files in current directory +files = os.listdir('.') + +# Iterate over the list of files +for file_name in files: + # Construct full file path + file_path = os.path.join('.', file_name) + # Check if it's a regular file + if os.path.isfile(file_path): + # Try to remove the xattr, if it exists + try: + os.removexattr(file_path, "user.sunny.upload") + print(f"Removed xattr from {file_name}") + except OSError: + print(f"No such xattr on {file_name}") \ No newline at end of file