Merge branch 'refs/heads/swaglog-to-cloudwatch' into master-dev-c3

# Conflicts:
#	CHANGELOGS.md
This commit is contained in:
DevTekVE
2024-05-28 14:58:54 +02:00
3 changed files with 128 additions and 22 deletions
+69 -17
View File
@@ -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
@@ -550,7 +551,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 +559,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 +571,66 @@ 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.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())
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
@@ -580,7 +640,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,18 +651,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
@@ -619,7 +671,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:
+27 -2
View File
@@ -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,15 +20,19 @@ 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:
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'),
@@ -37,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,), 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}')
@@ -48,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:
@@ -113,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()]
+32 -3
View File
@@ -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,25 @@ 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)
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}")