完成页面对接tcp推送的设备状态数据

完成页面测试解析tcp推送的设备状态数据的UI展示
This commit is contained in:
2026-03-03 19:34:17 +08:00
parent 0a3a062c04
commit 34b4c2bcbc
12 changed files with 676 additions and 230 deletions

View File

@@ -17,7 +17,9 @@ import 'package:maibu_satabot_v2/features/remote_control/domain/usecase/diff_ste
import 'package:shared_preferences/shared_preferences.dart';
import '../../features/auth/data/datasources/auth_http_datasource.dart';
import '../../features/auth/data/datasources/auth_tcp_datasource.dart';
import '../../features/auth/data/datasources/impl/auth_http_datasource_impl.dart';
import '../../features/auth/data/datasources/impl/auth_tcp_datasource_impl.dart';
import '../../features/auth/data/repositories/auth_repository_impl.dart';
import '../../features/auth/domain/repositories/auth_repository.dart';
import '../../features/auth/domain/usecases/login_usecase.dart';
@@ -41,6 +43,8 @@ import '../../features/devices/domain/usecases/get_work_record_usecase.dart';
import '../../features/devices/domain/usecases/route_planning_usecase.dart';
import '../../features/devices/domain/usecases/save_work_record_usecase.dart';
import '../../features/devices/domain/usecases/select_work_record_usecase.dart';
import '../../features/devices/domain/usecases/switch_device_usecase.dart';
import '../../features/devices/presentation/bloc/device_status_bloc.dart';
import '../../features/remote_control/presentation/bloc/remote_control_cubit.dart';
import '../app/app_user_cubit.dart';
import '../network/dio_client.dart';
@@ -59,7 +63,9 @@ Future<void> init() async {
sl.registerLazySingleton<Dio>(() => DioClient.create());
/// 1.1.2 TcpClient:TCP客户端
sl.registerLazySingleton(() => TcpClient());
sl.registerLazySingleton(() => TcpClient( sl<UserStorage>(),
getUserDeviceUseCase: sl<GetUserDeviceUseCase>(),
switchDeviceUseCase: sl<SwitchDeviceUseCase>(),));
/// 1.1.3 NetMessageDispatcher:消息调度器,并将 TcpClient 注入给它
sl.registerLazySingleton(() => NetMessageDispatcher(sl()));
@@ -83,7 +89,7 @@ Future<void> init() async {
() => AuthHttpDataSourceImpl(sl()),
);
sl.registerLazySingleton<DeviceHttpDatasource>(
() => DeviceHttpDatasourceImpl(sl()),
() => DeviceHttpDatasourceImpl(sl(), sl<UserStorage>()), //
);
/// 3. 仓库 (Repository)
@@ -126,7 +132,16 @@ Future<void> init() async {
sl.registerLazySingleton(() => UnbindDeviceUseCase(sl()));
sl.registerLazySingleton(() => UpdateDevicenameUsecase(sl()));
sl.registerLazySingleton<SaveWorkRecordUseCase>(() => SaveWorkRecordUseCase(sl<PathRepository>()));
sl.registerLazySingleton<SwitchDeviceUseCase>(() => SwitchDeviceUseCase(sl()));
sl.registerLazySingleton<AuthTcpDatasource>(
() => AuthTcpDatasourceImpl(
sl<TcpClient>(),
sl<UserStorage>(),
getUserDeviceUseCase: sl<GetUserDeviceUseCase>(),
switchDeviceUseCase: sl<SwitchDeviceUseCase>(),
),
);
/// 5. 状态管理 (Cubit/Bloc)
@@ -136,6 +151,11 @@ Future<void> init() async {
() => DevicesCubit(sl(), sl(), sl(), sl(), sl(), sl(), sl(), sl(), sl(),sl(),sl()),
);
sl.registerFactory(() => RemoteControlCubit(sl()));
/* // 工厂模式(留存,每次获取新实例)
sl.registerFactory(() => DeviceStatusBloc(sl<NetMessageDispatcher>()));
*/
//单例模式(全局共享一个实例)
sl.registerLazySingleton(() => DeviceStatusBloc(sl<NetMessageDispatcher>()));
/// . 路径生成
sl.registerLazySingleton<PathHttpDatasource>(
@@ -162,6 +182,7 @@ Future<void> init() async {
sl<TcpClient>(),
sl<AppUserCubit>(),
sl<NetMessageDispatcher>(),
sl<AuthTcpDatasource>(),
),
);

View File

@@ -16,7 +16,19 @@ class NetMessageDispatcher {
/// 示例解析方法:将 0x02 指令解析为 String
Stream<String> onStringMessage() {
return onCommand(0x02).map((p) => utf8.decode(p.payload));
print("0x02--TCP拦截推送解析开始");
return onCommand(0x02).map((p) {
try {
// 尝试解码
final result = utf8.decode(p.payload, allowMalformed: true); // 允许乱码,防止报错中断流
print('✅ 解码 0x02 成功:$result'); // 🔥 关键日志:看这里打印了吗?
return result;
} catch (e) {
print('❌ 解码 0x02 失败:$e, 原始字节:${p.payload}');
return ''; // 返回空字符串,避免流中断
}
});
/* ///return onCommand(0x02).map((p) => utf8.decode(p.payload));*/
}
/// 示例解析方法:将 0x12 指令解析为 JSON 并转为模型

View File

@@ -6,9 +6,14 @@ import 'dart:typed_data';
import 'package:flutter/cupertino.dart';
import 'package:maibu_satabot_v2/features/devices/data/repositories/route_planning_repository_impl.dart';
import '../../../features/auth/data/datasources/auth_tcp_datasource.dart';
import '../../../features/devices/data/models/route_plan_send_entity.dart';
import '../../../features/devices/domain/entities/device_entity.dart';
import '../../../features/devices/domain/entities/gps_entity.dart';
import '../../../features/devices/domain/entities/running_status_entity.dart';
import '../../../features/devices/domain/usecases/get_user_device_usecase.dart';
import '../../../features/devices/domain/usecases/switch_device_usecase.dart';
import '../../storage/user_storage.dart';
import '../protocol_decoder.dart';
@@ -18,39 +23,111 @@ class TcpClient {
// 使用 StreamController 将原始数据流暴露给外部,以便多个 DataSource 监听
final _controller = StreamController<RawPacket>.broadcast();
Stream<RawPacket> get packetStream => _controller.stream;
final GetUserDeviceUseCase getUserDeviceUseCase;
final SwitchDeviceUseCase switchDeviceUseCase;
TcpClient(this._userStorage, {required this.getUserDeviceUseCase, required this.switchDeviceUseCase});
// 新增:心跳定时器
Timer? _heartbeatTimer;
Timer? _reconnectTimer; // 重连定时器
String? _lastHost;
int? _lastPort;
final UserStorage _userStorage;
bool get isConnected => _socket != null;
// 通过参数配置,不硬编码
Future<void> connect({required String host, required int port}) async {
debugPrint('🔌 [TCP] 开始连接:$host:$port'); // ✅ 必须看到这条
_lastHost = host;
_lastPort = port;
try {
_socket = await Socket.connect(
host,
port,
timeout: const Duration(seconds: 5),
);
debugPrint('✅ [TCP] 连接成功!'); // ✅ 必须看到这条
// 开始认证tcp
await _sendAuthPacket();
_socket!.listen((data) {
debugPrint('📥 [TCP] 收到原始数据:${data.length} 字节, 内容:$data');
// var packets = _decoder.decode(data);
// debugPrint('📦 [TCP] 解码成功,包数量:${packets.length}');
//for (var packet in packets) {
// // 🔥若收到服务端心跳(cmd == 0xFF),立即回复一个心跳包
// if (packet.command == 0xFF) {
// debugPrint('收到服务端心跳,自动回复...');
// sendHeartbeat(); // 回复 AB AA FF AA AB
// }
// }
// _controller.add(packet);
// }
try {
var packets = _decoder.decode(data);
debugPrint('📦 [TCP] 解码成功,包数量:${packets.length}');
for (var packet in packets) {
// 🔥若收到服务端心跳(cmd == 0xFF),立即回复一个心跳包
if (packet.command == 0xFF) {
debugPrint('收到服务端心跳,自动回复...');
sendHeartbeat(); // 回复 AB AA FF AA AB
if (!_controller.isClosed) {
_controller.add(packet);
debugPrint('➡️ [TCP] 已分发 CMD: 0x${packet.command.toRadixString(16)}');
}
_controller.add(packet);
if (packet.command == 0xFF) {
debugPrint('收到服务端心跳,自动回复...');
sendHeartbeat(); // 回复 AB AA FF AA AB
}
if (packet.command == 0x03) {
debugPrint('⚠️ 收到认证响应:${packet.payload}');
// 解析 payload 看是否有错误信息
}
}
} catch (e, stackTrace) {
debugPrint('❌ [TCP] 解码数据时发生异常:$e\n$stackTrace'); // 🔥 捕获解码异常
}
},
onDone: (){
// TODO: 断线重连
debugPrint('onDone❌ [TCP] 连接已断开!');
_handleDisconnect();
},
onError: (e) {
debugPrint('❌ [TCP] 发生错误:$e');
_handleDisconnect(); // 统一走重连逻辑,保护 Controller 不被关闭
},
onDone: disconnect,
onError: (e) => disconnect(),
);
} catch (e) {
rethrow; // 向上抛出连接异常
}
}
void _handleDisconnect() {
stopHeartbeat();
_socket = null;
// 注意:这里不要关闭 _controller!否则监听者会丢失数据流
// _controller?.close();
_scheduleReconnect();
}
// ✅ 新增:调度重连
Future<void> _scheduleReconnect() async {
if (_reconnectTimer != null) return;
debugPrint('⏳ 调度重连:Host=${_lastHost}, Port=${_lastPort}'); // ✅ 检查 Host/Port 是否为空
if (_lastHost == null || _lastPort == null) {
debugPrint('❌ 无法重连:Host 或 Port 为空!');
return;
}
_reconnectTimer = Timer(const Duration(seconds: 5), () {
_reconnectTimer = null;
debugPrint('⏰ 定时器触发,开始执行重连...');
connect(host: _lastHost!, port: _lastPort!);
});
await _sendAuthPacket();
}
/// 发送数据
void send(Map<String, dynamic> data) {
if (_socket == null) throw Exception("Socket not connected");
@@ -89,7 +166,7 @@ class TcpClient {
// 新增:启动心跳(每 4 秒发送一次 0xFF 指令)
void startHeartbeat({Duration interval = const Duration(seconds: 4)}) {
if (_heartbeatTimer != null) return; // 防止重复启动
debugPrint('收到服务端心跳,自动回复...');
_heartbeatTimer = Timer.periodic(interval, (_) {
// 发送心跳帧:AB AA FF 00 00 AA AB(与 sendRaw 一致)
//sendRaw(0xFF, []);
@@ -130,6 +207,88 @@ class TcpClient {
}
//
Future<void> _sendAuthPacket() async {
if (_socket == null) return;
String? username= "";
String? token= "";
final user = await _userStorage.getUser();
debugPrint('tcp使用用户信息:$user');
if (user == null || user.token == null) {
debugPrint('❌ [TCP] 认证失败:用户未登录或 Token 为空,无法发送认证包');
return; // 直接返回,不要发送无效包
}
username = user.username;
token = user.token;
final authString = '$username:app:$token';
debugPrint('🔑 [TCP] 认证字符串:$authString');
final authBytes = utf8.encode(authString);
// 构造包结构:Head(2) + Cmd(1) + Payload(N) + CRC(2) + Foot(2)
// 总长度 = 3 + N + 4
final builder = BytesBuilder()
..addByte(0xAB)
..addByte(0xAA)
..addByte(0x03) // 认证指令
..add(authBytes);
// 添加 CRC (这里简单用 0x00 0x00 占位,如果协议严格校验需计算真实 CRC)
// Android 代码中是 packet[tailStart] = 0x00; packet[tailStart+1] = 0x00;
builder.addByte(0x00);
builder.addByte(0x00);
builder.addByte(0xAA);
builder.addByte(0xAB);
_socket!.add(builder.takeBytes());
debugPrint('🔑 [TCP] 已发送认证包 (0x03): $authString');
try {
// 1. 获取 Either 结果
final eitherResult = await getUserDeviceUseCase.repository.getUserDevice(username);
// 2. 使用 fold 解包 Either
// left: 处理错误情况 (DeviceFailure)
// right: 处理成功情况 (List<DeviceEntity>)
await eitherResult.fold(
(failure) {
// 处理失败:打印日志或抛出异常
debugPrint('❌ [AuthTcp] 获取设备列表失败:$failure');
throw Exception('获取设备列表失败:$failure');
},
(devices) async {
// 处理成功:devices 现在是真正的 List<DeviceEntity>
if (devices.isEmpty) {
debugPrint('⚠️ [AuthTcp] 当前用户无可用设备,跳过切换步骤');
return;
}
// 取第一个设备
final DeviceEntity targetDevice = devices.first;
debugPrint('📱 [AuthTcp] 准备切换至默认设备:${targetDevice.deviceName}');
// 切换设备 (同样,如果 switchDeviceUseCase 也返回 Either,也需要 fold 处理)
final switchResult = await switchDeviceUseCase.deviceRepository.switchDevice("app",targetDevice.deviceName);
await switchResult.fold(
(failure) {
debugPrint('❌ [AuthTcp] 切换设备失败:$failure');
throw Exception('切换设备失败:$failure');
},
(success) {
debugPrint('✅ [AuthTcp] 设备切换成功,服务端应开始推送数据');
},
);
},
);
} catch (e) {
debugPrint('❌ [AuthTcp] 设备订阅流程异常:$e');
rethrow;
}
}
void sendRawBytes(Uint8List bytes) {
if (_socket == null) return;
_socket!.add(bytes);
}
}