Files
flutterApp/lib/features/connectivity/presentation/bloc/connectivity_cubit.dart
2026-01-29 19:16:41 +08:00

127 lines
3.9 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;
static const int _watchdogTimerDuration = 10;
// 核心:这是唯一的流入口
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) {
_broadcastPacketStream = _activeConnection!.rawStream
.transform(PacketParserTransformer())
.asBroadcastStream();
}
}
void _startMonitoring() {
_stopMonitoring();
_resetWatchdog();
// 确保流已存在
_ensureStreamInitialized();
// 此时 listen 的是已经转换好的广播流,绝对不会触发 rawStream.transform
_packetSub = _broadcastPacketStream?.listen((packet) {
_resetWatchdog();
if (packet.cmdType == 0xFF) {
sendCommand(0xFF, null);
sl<ILoggerService>().log("收到心跳包,已回复...");
}
}, onError: (e) => sl<ILoggerService>().log("流监听错误: $e"));
}
// Getter 也要改,保证它只返回已经创建好的流
Stream<PacketEntity> get activePacketStream {
_ensureStreamInitialized();
return _broadcastPacketStream ?? const Stream.empty();
}
// --- 其他逻辑保持不变 ---
void _resetWatchdog() {
_watchdogTimer?.cancel();
_watchdogTimer = Timer(const Duration(seconds: _watchdogTimerDuration), () {
sl<ILoggerService>().log("$_watchdogTimerDuration秒内未收到心跳包...");
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();
}
}