From 5030860ab8c7af741de087a9265d988740446e6a Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 11 Jun 2024 06:52:11 +0000 Subject: [PATCH 1/3] [sunnylink] Bugfixes like sidebar & improved logging for error troubleshoot --- selfdrive/athena/athenad.py | 1 + selfdrive/athena/sunnylinkd.py | 8 ++++---- selfdrive/manager/manager.py | 1 + selfdrive/ui/qt/sidebar.cc | 13 +++++++++---- 4 files changed, 15 insertions(+), 8 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index d5b1e37ce1..443173e36d 100755 --- a/selfdrive/athena/athenad.py +++ b/selfdrive/athena/athenad.py @@ -178,6 +178,7 @@ def jsonrpc_handler(end_event: threading.Event) -> None: elif "id" in data and ("result" in data or "error" in data): log_recv_queue.put_nowait(data) else: + cloudlog.event("athena.jsonrpc_handler.invalid_request", error=True, data=data) raise Exception("not a valid request or response") except queue.Empty: pass diff --git a/selfdrive/athena/sunnylinkd.py b/selfdrive/athena/sunnylinkd.py index 9c26dd42a1..b0c6e8d4dc 100755 --- a/selfdrive/athena/sunnylinkd.py +++ b/selfdrive/athena/sunnylinkd.py @@ -28,6 +28,7 @@ LOCAL_PORT_WHITELIST = {8022} params = Params() sunnylink_api = SunnylinkApi(params.get("SunnylinkDongleId", encoding='utf-8')) def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: + cloudlog.info("sunnylinkd.handle_long_poll started") end_event = threading.Event() threads = [ @@ -55,7 +56,7 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: raise finally: for thread in threads: - cloudlog.debug(f"athena.joining {thread.name}") + cloudlog.info(f"sunnylinkd athena.joining {thread.name}") thread.join() @@ -92,7 +93,7 @@ def ws_ping(ws: WebSocket, end_event: threading.Event) -> None: except Exception: cloudlog.exception("sunnylinkd.ws_ping.exception") end_event.set() - time.sleep(RECONNECT_TIMEOUT_S * 0.8) # Sleep about 80% before a timeout + time.sleep(RECONNECT_TIMEOUT_S * 0.7) # Sleep about 70% before a timeout def ws_queue(end_event: threading.Event) -> None: resume_requested = False @@ -169,8 +170,7 @@ def main(exit_event: threading.Event = None): cloudlog.event("sunnylinkd.main.connecting_ws", ws_uri=ws_uri, retries=conn_retries) ws = create_connection(ws_uri, cookie="jwt=" + sunnylink_api.get_token(), - enable_multithread=True, - timeout=30.0) + enable_multithread=True) cloudlog.event("sunnylinkd.main.connected_ws", ws_uri=ws_uri, retries=conn_retries, duration=time.monotonic() - conn_start) conn_start = None diff --git a/selfdrive/manager/manager.py b/selfdrive/manager/manager.py index ef8cd286a9..07bd5f5848 100755 --- a/selfdrive/manager/manager.py +++ b/selfdrive/manager/manager.py @@ -107,6 +107,7 @@ def manager_init() -> None: ("SunnylinkEnabled", "1"), ("CustomDrivingModel", "0"), ("DrivingModelGeneration", "4"), + ("LastSunnylinkPingTime", "0"), ] if not PC: default_params.append(("LastUpdateTime", datetime.datetime.utcnow().isoformat().encode('utf8'))) diff --git a/selfdrive/ui/qt/sidebar.cc b/selfdrive/ui/qt/sidebar.cc index df9ebc45c0..5d8245acc5 100644 --- a/selfdrive/ui/qt/sidebar.cc +++ b/selfdrive/ui/qt/sidebar.cc @@ -155,17 +155,22 @@ void Sidebar::updateState(const UIState &s) { setProperty("pandaStatus", QVariant::fromValue(pandaStatus)); ItemStatus sunnylinkStatus; - auto last_sunnylink_ping = std::strtol(params.get("LastSunnylinkPingTime").c_str(), nullptr, 10); + auto last_sunnylink_ping_str = params.get("LastSunnylinkPingTime"); + auto last_sunnylink_ping = std::stoull(last_sunnylink_ping_str.empty() ? "0" : last_sunnylink_ping_str); + auto current_nanos = nanos_since_boot(); + auto elapsed_sunnylink_ping = current_nanos - last_sunnylink_ping; auto sunnylink_enabled = params.getBool("SunnylinkEnabled"); if (!sunnylink_enabled) { sunnylinkStatus = ItemStatus{{tr("SUNNYLINK"), tr("DISABLED")}, disabled_color}; - } else if (last_ping == 0) { + } else if (last_sunnylink_ping == 0) { sunnylinkStatus = ItemStatus{{tr("SUNNYLINK"), tr("OFFLINE")}, warning_color}; } else { - if (nanos_since_boot() - last_sunnylink_ping < 80e9) + if (elapsed_sunnylink_ping < 80000000000ULL) { sunnylinkStatus = ItemStatus{{tr("SUNNYLINK"), tr("ONLINE")}, good_color}; - else + } + else { sunnylinkStatus = ItemStatus{{tr("SUNNYLINK"), tr("ERROR")}, danger_color}; + } } setProperty("sunnylinkStatus", QVariant::fromValue(sunnylinkStatus)); } From 2b0dbb4e8e9a8a260aebb9347ca32ea5000051c9 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 11 Jun 2024 07:19:54 +0000 Subject: [PATCH 2/3] [sunnylink] Swaglogs to cloudwatch Enable log_handler in sunnylinkd.py --- selfdrive/athena/athenad.py | 87 +++++++++++++++++++++++++++------- selfdrive/athena/sunnylinkd.py | 31 ++++++++++-- system/loggerd/xattr_cache.py | 15 ++++-- 3 files changed, 109 insertions(+), 24 deletions(-) diff --git a/selfdrive/athena/athenad.py b/selfdrive/athena/athenad.py index 443173e36d..e8ef682026 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 @@ -551,7 +552,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 = [] @@ -559,7 +560,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): @@ -571,8 +572,68 @@ def get_logs_to_send_sorted() -> list[str]: return sorted(logs)[:-1] -def log_handler(end_event: threading.Event) -> None: +def add_log_to_queue(log_path, log_id, is_sunnylink = False): + MAX_SIZE_KB = 32 + MAX_SIZE_BYTES = MAX_SIZE_KB * 1024 + + with open(log_path, 'r') as f: + data = f.read() + + # Check if the file is empty + if not data: + cloudlog.warning(f"Log file {log_path} is empty.") + return + + # 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.debug(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()) + 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.debug(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.debug(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.debug(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: + cloudlog.debug("athena.log_handler: Not supported on PC") return log_files = [] @@ -581,7 +642,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 @@ -592,18 +653,10 @@ 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)) - 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 + setxattr(log_path, log_attr_name, int.to_bytes(curr_time, 4, sys.byteorder)) + + add_log_to_queue(log_path, log_entry, is_sunnylink) + curr_log = log_entry except OSError: pass # file could be deleted by log rotation @@ -620,7 +673,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 b0c6e8d4dc..0b784d0371 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) @@ -20,16 +20,20 @@ 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")) LOCAL_PORT_WHITELIST = {8022} +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: cloudlog.info("sunnylinkd.handle_long_poll started") + 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 +42,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=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,14 +53,28 @@ 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.debug(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.debug(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: - cloudlog.info(f"sunnylinkd athena.joining {thread.name}") + cloudlog.debug(f"sunnylinkd athena.joining {thread.name}") thread.join() @@ -114,6 +132,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()] diff --git a/system/loggerd/xattr_cache.py b/system/loggerd/xattr_cache.py index d3220118ac..88aa4f7916 100644 --- a/system/loggerd/xattr_cache.py +++ b/system/loggerd/xattr_cache.py @@ -1,5 +1,12 @@ -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 +14,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 +26,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) \ No newline at end of file From 10047307509582c3b9bfde008b6a49a3d6f95712 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Tue, 11 Jun 2024 07:21:45 +0000 Subject: [PATCH 3/3] [sunnylink] Add sunnylink_uploader to process_config and implement uploader --- selfdrive/manager/process_config.py | 13 +- system/loggerd/sunnylink_uploader.py | 302 +++++++++++++++++++++++++++ 2 files changed, 311 insertions(+), 4 deletions(-) create mode 100755 system/loggerd/sunnylink_uploader.py diff --git a/selfdrive/manager/process_config.py b/selfdrive/manager/process_config.py index 60c5b8ebc6..de84a43f16 100644 --- a/selfdrive/manager/process_config.py +++ b/selfdrive/manager/process_config.py @@ -104,10 +104,15 @@ procs = [ PythonProcess("webjoystick", "tools.bodyteleop.web", notcar), ] -if os.path.exists("../athena/manage_sunnylinkd.py") and Params().get_bool("SunnylinkEnabled"): - procs += [ - DaemonProcess("manage_sunnylinkd", "selfdrive.athena.manage_sunnylinkd", "SunnylinkdPid"), - ] +if Params().get_bool("SunnylinkEnabled"): + if os.path.exists("../athena/manage_sunnylinkd.py"): + procs += [ + DaemonProcess("manage_sunnylinkd", "selfdrive.athena.manage_sunnylinkd", "SunnylinkdPid"), + ] + if os.path.exists("../../system/loggerd/sunnylink_uploader.py"): + procs += [ + PythonProcess("sunnylink_uploader", "system.loggerd.sunnylink_uploader", always_run), + ] if os.path.exists("./gitlab_runner.sh") and True: # Of course and True is always true. Placeholder for a param :D # Only devs! diff --git a/system/loggerd/sunnylink_uploader.py b/system/loggerd/sunnylink_uploader.py new file mode 100755 index 0000000000..cb0e837901 --- /dev/null +++ b/system/loggerd/sunnylink_uploader.py @@ -0,0 +1,302 @@ +#!/usr/bin/env python3 +import bz2 +import io +import json +import os +import random +import requests +import threading +import time +import traceback +import datetime +from typing import BinaryIO +from collections.abc import Iterator + +from cereal import log +import cereal.messaging as messaging +from openpilot.common.api import SunnylinkApi +from openpilot.common.params import Params +from openpilot.common.realtime import set_core_affinity +from openpilot.system.hardware.hw import Paths +from openpilot.system.loggerd.xattr_cache import getxattr, setxattr +from openpilot.common.swaglog import cloudlog + +NetworkType = log.DeviceState.NetworkType +UPLOAD_ATTR_NAME = 'user.sunny.upload' + +UPLOAD_ATTR_VALUE = b'1' + +UPLOAD_QLOG_QCAM_MAX_SIZE = 5 * 1e6 # MB + +allow_sleep = bool(os.getenv("UPLOADER_SLEEP", "1")) +force_wifi = os.getenv("FORCEWIFI") is not None +fake_upload = os.getenv("FAKEUPLOAD") is not None + +OFFROAD_TRANSITION_TIMEOUT = 900. # wait until offroad for 15 minutes before allowing uploads + + +class FakeRequest: + def __init__(self): + self.headers = {"Content-Length": "0"} + + +class FakeResponse: + def __init__(self): + self.status_code = 200 + self.request = FakeRequest() + + +def get_directory_sort(d: str) -> list[str]: + # ensure old format is sorted sooner + o = ["0", ] if d.startswith("2024-") else ["1", ] + return o + [s.rjust(10, '0') for s in d.rsplit('--', 1)] + +def listdir_by_creation(d: str) -> list[str]: + if not os.path.isdir(d): + return [] + + try: + paths = [f for f in os.listdir(d) if os.path.isdir(os.path.join(d, f))] + paths = sorted(paths, key=get_directory_sort) + return paths + except OSError: + cloudlog.exception("listdir_by_creation failed") + return [] + +def clear_locks(root: str) -> None: + for logdir in os.listdir(root): + path = os.path.join(root, logdir) + try: + for fname in os.listdir(path): + if fname.endswith(".lock"): + os.unlink(os.path.join(path, fname)) + except OSError: + cloudlog.exception("clear_locks failed") + + +class Uploader: + def __init__(self, dongle_id: str, root: str): + self.dongle_id = dongle_id + self.api = SunnylinkApi(dongle_id) + self.root = root + + self.params = Params() + + # stats for last successfully uploaded file + self.last_filename = "" + + self.immediate_folders = ["crash/", "boot/"] + self.immediate_priority = {"qlog": 0, "qlog.bz2": 0, "qcamera.ts": 1} + + def list_upload_files(self, metered: bool) -> Iterator[tuple[str, str, str]]: + r = self.params.get("AthenadRecentlyViewedRoutes", encoding="utf8") + requested_routes = [] if r is None else r.split(",") + + for logdir in listdir_by_creation(self.root): + path = os.path.join(self.root, logdir) + try: + names = os.listdir(path) + except OSError: + continue + + if any(name.endswith(".lock") for name in names): + continue + + for name in sorted(names, key=lambda n: self.immediate_priority.get(n, 1000)): + key = os.path.join(logdir, name) + fn = os.path.join(path, name) + # skip files already uploaded + try: + ctime = os.path.getctime(fn) + is_uploaded = getxattr(fn, UPLOAD_ATTR_NAME) == UPLOAD_ATTR_VALUE + except OSError: + cloudlog.event("uploader_getxattr_failed", key=key, fn=fn) + # deleter could have deleted, so skip + continue + if is_uploaded: + continue + + # limit uploading on metered connections + if metered: + dt = datetime.timedelta(hours=12) + if logdir in self.immediate_folders and (datetime.datetime.now() - datetime.datetime.fromtimestamp(ctime)) < dt: + continue + + if name == "qcamera.ts" and not any(logdir.startswith(r.split('|')[-1]) for r in requested_routes): + continue + + yield name, key, fn + + def next_file_to_upload(self, metered: bool) -> tuple[str, str, str] | None: + upload_files = list(self.list_upload_files(metered)) + + for name, key, fn in upload_files: + if any(f in fn for f in self.immediate_folders): + return name, key, fn + + for name, key, fn in upload_files: + if name in self.immediate_priority: + return name, key, fn + + return None + + def do_upload(self, key: str, fn: str): + url_resp = self.api.get("device/" + self.dongle_id + "/upload_url/", timeout=10, path=key, access_token=self.api.get_token()) + if url_resp.status_code == 412: + return url_resp + + url_resp_json = json.loads(url_resp.text) + url = url_resp_json['url'] + headers = url_resp_json['headers'] + cloudlog.debug("sunnylink upload_url %s | Headers: %s", url, headers) + + if fake_upload: + return FakeResponse() + + with open(fn, "rb") as f: + data: BinaryIO + if key.endswith('.bz2') and not fn.endswith('.bz2'): + compressed = bz2.compress(f.read()) + data = io.BytesIO(compressed) + else: + data = f + + return requests.put(url, data=data, headers=headers, timeout=10) + + def upload(self, name: str, key: str, fn: str, network_type: int, metered: bool) -> bool: + try: + sz = os.path.getsize(fn) + except OSError: + cloudlog.exception("upload: getsize failed") + return False + + cloudlog.event("upload_start", key=key, fn=fn, sz=sz, network_type=network_type, metered=metered) + + if sz == 0: + # tag files of 0 size as uploaded + success = True + elif name in self.immediate_priority and sz > UPLOAD_QLOG_QCAM_MAX_SIZE: + cloudlog.event("uploader_too_large", key=key, fn=fn, sz=sz) + success = True + else: + start_time = time.monotonic() + + stat = None + last_exc = None + try: + stat = self.do_upload(key, fn) + except Exception as e: + last_exc = (e, traceback.format_exc()) + + if stat is not None and stat.status_code in (200, 201, 412): + self.last_filename = fn + dt = time.monotonic() - start_time + if stat.status_code == 412: + cloudlog.event("upload_ignored", key=key, fn=fn, sz=sz, network_type=network_type, metered=metered) + else: + content_length = int(stat.request.headers.get("Content-Length", 0)) + speed = (content_length / 1e6) / dt + cloudlog.event("upload_success", key=key, fn=fn, sz=sz, content_length=content_length, + network_type=network_type, metered=metered, speed=speed) + success = True + elif stat is not None: # 401, 403... Not sure why they were up to begin with + success = False + cloudlog.event("upload_failed with content", stat=stat, exc=last_exc, key=key, fn=fn, sz=sz, network_type=network_type, metered=metered, error=stat.content.decode("utf-8")) + else: + success = False + cloudlog.event("upload_failed", stat=stat, exc=last_exc, key=key, fn=fn, sz=sz, network_type=network_type, metered=metered) + + if success: + # tag file as uploaded + try: + setxattr(fn, UPLOAD_ATTR_NAME, UPLOAD_ATTR_VALUE) + except OSError: + cloudlog.event("uploader_setxattr_failed", exc=last_exc, key=key, fn=fn, sz=sz) + + return success + + + def step(self, network_type: int, metered: bool) -> bool | None: + d = self.next_file_to_upload(metered) + if d is None: + return None + + name, key, fn = d + + # qlogs and bootlogs need to be compressed before uploading + if key.endswith(('qlog', 'rlog')) or (key.startswith('boot/') and not key.endswith('.bz2')): + key += ".bz2" + + return self.upload(name, key, fn, network_type, metered) + + +def main(exit_event: threading.Event = None) -> None: + if exit_event is None: + exit_event = threading.Event() + + try: + set_core_affinity([0, 1, 2, 3]) + except Exception: + cloudlog.exception("failed to set core affinity") + + clear_locks(Paths.log_root()) + + params = Params() + dongle_id = params.get("SunnylinkDongleId", encoding='utf8') + + offroad_transition_prev = 0. + offroad_last = False + + if dongle_id is None: + cloudlog.info("uploader missing dongle_id") + raise Exception("uploader can't start without dongle id") + + sm = messaging.SubMaster(['deviceState']) + uploader = Uploader(dongle_id, Paths.log_root()) + + backoff = 0.1 + while not exit_event.is_set(): + sm.update(0) + + offroad = params.get_bool("IsOffroad") + t = time.monotonic() + if offroad and not offroad_last and t > 300.: + offroad_transition_prev = time.monotonic() + offroad_last = offroad + + network_type = sm['deviceState'].networkType if not force_wifi else NetworkType.wifi + if network_type == NetworkType.none: + if allow_sleep: + time.sleep(60 if offroad else 5) + continue + + if params.get_bool("DisableOnroadUploads"): + if not offroad or (offroad_transition_prev > 0. and t - offroad_transition_prev < OFFROAD_TRANSITION_TIMEOUT): + if not offroad: + cloudlog.info("not uploading: onroad uploads disabled") + else: + wait_minutes = int(OFFROAD_TRANSITION_TIMEOUT / 60) + time_left = OFFROAD_TRANSITION_TIMEOUT - (t - offroad_transition_prev) + if time_left / 60. > 2.: + time_left_str = f"{int(time_left / 60)} minute(s)" + else: + time_left_str = f"{int(time_left)} seconds(s)" + cloudlog.info(f"not uploading: waiting until offroad for {wait_minutes} minutes; {time_left_str} left") + if allow_sleep: + time.sleep(60) + continue + + success = uploader.step(sm['deviceState'].networkType.raw, sm['deviceState'].networkMetered) + if success is None: + backoff = 60 if offroad else 5 + elif success: + backoff = 0.1 + else: + cloudlog.info("upload backoff %r", backoff) + backoff = min(backoff*2, 120) + if allow_sleep: + time.sleep(backoff + random.uniform(0, backoff)) + + +if __name__ == "__main__": + main()