133 lines
4.1 KiB
Dart
133 lines
4.1 KiB
Dart
import 'dart:async';
|
||
import 'package:flutter_bloc/flutter_bloc.dart';
|
||
import 'package:fpdart/fpdart.dart';
|
||
import '../../../../core/di/injection.dart';
|
||
import '../../../../core/error/failure.dart';
|
||
import '../../../../core/logging/i_logger_service.dart';
|
||
import '../../../../core/network/base/i_connection.dart';
|
||
import '../../../../core/network/enums/control_mode.dart';
|
||
import '../../../../core/protocol/packet_entity.dart';
|
||
import '../../../../core/protocol/packet_parser_transformer.dart';
|
||
import '../../../../core/protocol/protocol_constants.dart';
|
||
import 'connectivity_state.dart';
|
||
|
||
class ConnectivityCubit extends Cubit<ConnectivityState> {
|
||
IConnection? _activeConnection;
|
||
StreamSubscription? _statusSub;
|
||
StreamSubscription<PacketEntity>? _packetSub;
|
||
Timer? _watchdogTimer;
|
||
|
||
// 核心:这是唯一的流入口
|
||
Stream<PacketEntity>? _broadcastPacketStream;
|
||
|
||
ConnectivityCubit() : super(ConnectivityState.initial());
|
||
|
||
void setConnection(IConnection? connection, ControlMode mode) {
|
||
// 1. 彻底清理(按顺序来)
|
||
_stopMonitoring();
|
||
_statusSub?.cancel();
|
||
_activeConnection?.disconnect();
|
||
|
||
// 2. 关键:切断旧流的所有引用
|
||
_broadcastPacketStream = null;
|
||
|
||
_activeConnection = connection;
|
||
|
||
if (connection != null) {
|
||
_statusSub = connection.statusStream.listen((status) {
|
||
emit(state.copyWith(activeMode: mode, status: status));
|
||
|
||
if (status == ConnectionStatus.connected) {
|
||
// 3. 只有连接成功时,才初始化唯一的广播流
|
||
_ensureStreamInitialized();
|
||
_startMonitoring();
|
||
} else {
|
||
_stopMonitoring();
|
||
}
|
||
});
|
||
connection.connect();
|
||
} else {
|
||
emit(state.copyWith(activeMode: ControlMode.NONE, status: ConnectionStatus.disconnected));
|
||
}
|
||
}
|
||
|
||
// 核心:强制初始化,确保 transform 只跑一次
|
||
void _ensureStreamInitialized() {
|
||
if (_activeConnection != null && _broadcastPacketStream == null) {
|
||
print("DEBUG: 正在创建唯一的广播流转换器");
|
||
|
||
// 确保 rawStream 是有效的
|
||
if (_activeConnection!.rawStream != null) {
|
||
_broadcastPacketStream = _activeConnection!.rawStream
|
||
.transform(PacketParserTransformer())
|
||
.asBroadcastStream();
|
||
} else {
|
||
print("DEBUG: rawStream is null,无法初始化广播流");
|
||
}
|
||
}
|
||
}
|
||
|
||
|
||
void _startMonitoring() {
|
||
_stopMonitoring();
|
||
_resetWatchdog();
|
||
|
||
// 确保流已存在
|
||
_ensureStreamInitialized();
|
||
|
||
// 此时 listen 的是已经转换好的广播流,绝对不会触发 rawStream.transform
|
||
_packetSub = _broadcastPacketStream?.listen((packet) {
|
||
_resetWatchdog();
|
||
if (packet.cmdType == 0xFF) {
|
||
sendCommand(0xFF, null);
|
||
sl<ILoggerService>().log("receive heartbeat,has response...");
|
||
}
|
||
}, onError: (e) => print("流监听错误: $e"));
|
||
}
|
||
|
||
// Getter 也要改,保证它只返回已经创建好的流
|
||
Stream<PacketEntity> get activePacketStream {
|
||
_ensureStreamInitialized();
|
||
return _broadcastPacketStream ?? const Stream.empty();
|
||
}
|
||
|
||
// --- 其他逻辑保持不变 ---
|
||
void _resetWatchdog() {
|
||
_watchdogTimer?.cancel();
|
||
_watchdogTimer = Timer(const Duration(seconds: 10), () {
|
||
print("警告:10秒内数据,触发断开逻辑");
|
||
emit(state.copyWith(status: ConnectionStatus.disconnected));
|
||
});
|
||
}
|
||
|
||
void _stopMonitoring() {
|
||
_packetSub?.cancel();
|
||
_packetSub = null;
|
||
_watchdogTimer?.cancel();
|
||
_watchdogTimer = null;
|
||
}
|
||
|
||
Future<Either<Failure, Unit>> sendCommand(int cmd, List<int>? payload) async {
|
||
if (_activeConnection == null) return left(NetworkFailure("无连接"));
|
||
|
||
final frame = <int>[
|
||
ProtocolConstants.header1,
|
||
ProtocolConstants.header2,
|
||
cmd,
|
||
...?payload,
|
||
if (payload != null && payload.isNotEmpty) ...[0x00, 0x00],
|
||
ProtocolConstants.endFlag1,
|
||
ProtocolConstants.endFlag2
|
||
];
|
||
|
||
return await _activeConnection!.send(frame);
|
||
}
|
||
|
||
@override
|
||
Future<void> close() {
|
||
_stopMonitoring();
|
||
_statusSub?.cancel();
|
||
_activeConnection?.disconnect();
|
||
return super.close();
|
||
}
|
||
} |