feat: TCP协议支持protobuf序列化,保留JSON兼容
- 新增 tcp_messages_pb2.py protobuf Python绑定 - msg_handler.cpp: 新增 make_packet_pb/raw bytes打包, parse_packet_raw/raw解析 - archery_netcore.cpp: 注册 make_packet_pb/parse_packet_raw - network.py: 新增 _use_proto 标志,_make_send_packet/_parse_recv 自动选择proto/JSON - 登录version加+proto后缀标识proto模式 - protobuf序列化/反序列化失败时自动回退JSON
This commit is contained in:
+114
-7
@@ -23,6 +23,14 @@ from logger_manager import logger_manager
|
||||
from wifi import wifi_manager
|
||||
import subprocess
|
||||
|
||||
# protobuf 支持
|
||||
try:
|
||||
import tcp_messages_pb2 as pb
|
||||
_HAS_PROTO = True
|
||||
except ImportError:
|
||||
_HAS_PROTO = False
|
||||
print("[NET] tcp_messages_pb2 not found, protobuf disabled")
|
||||
|
||||
|
||||
def _wifi_tls_would_block(exc):
|
||||
"""
|
||||
@@ -72,6 +80,9 @@ class NetworkManager:
|
||||
self._raw_line_data = []
|
||||
self._manual_trigger_flag = False
|
||||
|
||||
# protobuf 协议支持
|
||||
self._use_proto = _HAS_PROTO # 默认启用 proto(如果可用)
|
||||
|
||||
# 限制并发命令线程数
|
||||
self._cmd_thread_lock = threading.Lock()
|
||||
self._cmd_thread_count = 0
|
||||
@@ -711,6 +722,103 @@ class NetworkManager:
|
||||
"""线程安全地将消息加入队列(公共方法)"""
|
||||
self._enqueue((msg_type, data_dict), high)
|
||||
|
||||
def _make_send_packet(self, msg_type, data_dict):
|
||||
"""根据协议模式构造发送数据包"""
|
||||
if self._use_proto and _HAS_PROTO:
|
||||
return self._make_proto_packet(msg_type, data_dict)
|
||||
return self._netcore.make_packet(msg_type, data_dict)
|
||||
|
||||
def _make_proto_packet(self, msg_type, data_dict):
|
||||
"""使用 protobuf 序列化构造数据包"""
|
||||
try:
|
||||
if msg_type == 1:
|
||||
# 登录消息
|
||||
msg = pb.LoginRequest(
|
||||
device_id=data_dict.get("deviceId", ""),
|
||||
password=data_dict.get("password", ""),
|
||||
if_admin=data_dict.get("ifAdmin", False),
|
||||
version=data_dict.get("version", ""),
|
||||
vol=data_dict.get("vol", 0),
|
||||
vol_per=data_dict.get("vol_per", 0),
|
||||
iccid=data_dict.get("iccid", ""),
|
||||
)
|
||||
elif msg_type == 4:
|
||||
# 心跳消息
|
||||
msg = pb.Heartbeat(
|
||||
t=data_dict.get("t", 0),
|
||||
vol=data_dict.get("vol", 0),
|
||||
vol_per=data_dict.get("vol_per", 0),
|
||||
)
|
||||
elif msg_type == 2:
|
||||
# 业务逻辑消息
|
||||
cmd = data_dict.get("cmd", 0)
|
||||
inner_data = {k: v for k, v in data_dict.items() if k != "cmd"}
|
||||
data_bytes = json.dumps(inner_data).encode("utf-8") if inner_data else b""
|
||||
msg = pb.LogicBody(cmd=cmd, data=data_bytes)
|
||||
else:
|
||||
# 其他消息类型,回退到 JSON
|
||||
return self._netcore.make_packet(msg_type, data_dict)
|
||||
|
||||
body_bytes = msg.SerializeToString()
|
||||
return self._netcore.make_packet_pb(msg_type, body_bytes)
|
||||
except Exception as e:
|
||||
self.logger.error(f"[NET] protobuf 序列化失败,回退到 JSON: {e}")
|
||||
return self._netcore.make_packet(msg_type, data_dict)
|
||||
|
||||
def _parse_recv(self, payload):
|
||||
"""解析接收的数据包,返回 (msg_type, body_dict)"""
|
||||
if self._use_proto and _HAS_PROTO:
|
||||
msg_type, body_bytes = self._netcore.parse_packet_raw(payload)
|
||||
if msg_type is None:
|
||||
return None, None
|
||||
try:
|
||||
body_dict = self._parse_proto_body(msg_type, body_bytes)
|
||||
return msg_type, body_dict
|
||||
except Exception as e:
|
||||
self.logger.error(f"[NET] protobuf 反序列化失败: {e}")
|
||||
# 回退到 JSON 解析
|
||||
return self._netcore.parse_packet(payload)
|
||||
else:
|
||||
return self._netcore.parse_packet(payload)
|
||||
|
||||
def _parse_proto_body(self, msg_type, body_bytes):
|
||||
"""将 protobuf body bytes 反序列化为 dict"""
|
||||
if msg_type == 1:
|
||||
msg = pb.LoginResponse()
|
||||
msg.ParseFromString(body_bytes)
|
||||
return {"cmd": msg.cmd, "data": msg.data}
|
||||
elif msg_type == 4:
|
||||
# 心跳 ACK 通常无 body
|
||||
return {}
|
||||
elif msg_type == 2:
|
||||
msg = pb.LogicBody()
|
||||
msg.ParseFromString(body_bytes)
|
||||
result = {"cmd": msg.cmd}
|
||||
if msg.data:
|
||||
try:
|
||||
result["data"] = json.loads(msg.data.decode("utf-8"))
|
||||
except:
|
||||
result["data"] = {"raw": msg.data.hex()}
|
||||
return result
|
||||
elif msg_type == 40:
|
||||
msg = pb.OtaFragment()
|
||||
msg.ParseFromString(body_bytes)
|
||||
return {"l": msg.l, "d": msg.d, "t": msg.t, "v": msg.v}
|
||||
elif msg_type == 100:
|
||||
msg = pb.ImageUploadCommand()
|
||||
msg.ParseFromString(body_bytes)
|
||||
return {"uploadUrl": msg.upload_url, "token": msg.token, "shootId": msg.shoot_id, "outlink": msg.outlink}
|
||||
elif msg_type == 101:
|
||||
msg = pb.LogUploadCommand()
|
||||
msg.ParseFromString(body_bytes)
|
||||
return {"uploadUrl": msg.upload_url, "token": msg.token, "key": msg.key, "outlink": msg.outlink, "archive": msg.archive}
|
||||
else:
|
||||
# 未知类型,尝试 JSON 解析
|
||||
try:
|
||||
return json.loads(body_bytes.decode("utf-8"))
|
||||
except:
|
||||
return {"raw": body_bytes.hex()}
|
||||
|
||||
def connect_server(self):
|
||||
"""
|
||||
连接到服务器(自动选择WiFi或4G)
|
||||
@@ -1842,14 +1950,13 @@ class NetworkManager:
|
||||
login_data = {
|
||||
"deviceId": self.device_id,
|
||||
"password": self.password,
|
||||
"version": config.APP_VERSION,
|
||||
"version": config.APP_VERSION + ("+proto" if self._use_proto else ""),
|
||||
"vol": vol_val,
|
||||
"vol_per": voltage_to_percent(vol_val)
|
||||
}
|
||||
iccid_pending_marker = self._maybe_add_iccid_to_login(login_data)
|
||||
print(f"login_data: {login_data}")
|
||||
# if not self.tcp_send_raw(self.make_packet(1, login_data)):
|
||||
if not self.tcp_send_raw(self._netcore.make_packet(1, login_data)):
|
||||
if not self.tcp_send_raw(self._make_send_packet(1, login_data)):
|
||||
self._tcp_connected = False
|
||||
try:
|
||||
self.disconnect_server()
|
||||
@@ -1926,7 +2033,7 @@ class NetworkManager:
|
||||
pass
|
||||
|
||||
# msg_type, body = self.parse_packet(payload)
|
||||
msg_type, body = self._netcore.parse_packet(payload)
|
||||
msg_type, body = self._parse_recv(payload)
|
||||
|
||||
# 处理登录响应
|
||||
if not logged_in and msg_type == 1:
|
||||
@@ -2310,7 +2417,7 @@ class NetworkManager:
|
||||
|
||||
if item:
|
||||
msg_type, data_dict = item
|
||||
pkt = self._netcore.make_packet(msg_type, data_dict)
|
||||
pkt = self._make_send_packet(msg_type, data_dict)
|
||||
if not self.tcp_send_raw(pkt):
|
||||
# 发送失败:将消息放回队首(队列满则丢弃)
|
||||
with self.get_queue_lock():
|
||||
@@ -2339,8 +2446,8 @@ class NetworkManager:
|
||||
current_time = time.ticks_ms()
|
||||
if logged_in and current_time - last_heartbeat_send_time > config.HEARTBEAT_INTERVAL * 1000:
|
||||
vol_val = get_bus_voltage()
|
||||
if not self.tcp_send_raw(
|
||||
self._netcore.make_packet(4, {"vol": vol_val, "vol_per": voltage_to_percent(vol_val)})):
|
||||
heartbeat_pkt = self._make_send_packet(4, {"vol": vol_val, "vol_per": voltage_to_percent(vol_val)})
|
||||
if not self.tcp_send_raw(heartbeat_pkt):
|
||||
# if not self.tcp_send_raw(self.make_packet(4, {"vol": vol_val, "vol_per": voltage_to_percent(vol_val)})):
|
||||
send_hartbeat_fail_count += 1
|
||||
# 短暂波动可能导致一次发送失败:连续失败达到阈值才重连,避免重连风暴
|
||||
|
||||
Reference in New Issue
Block a user