187 lines
4.8 KiB
Dart
187 lines
4.8 KiB
Dart
// ignore_for_file: avoid_init_to_null, prefer_final_fields, avoid_print
|
|
|
|
import 'dart:io';
|
|
import 'dart:async';
|
|
import 'dart:typed_data';
|
|
|
|
import 'package:dart_framework/utils/system_util.dart';
|
|
import 'package:protobuf/protobuf.dart';
|
|
import 'package:dart_framework/communicator.dart';
|
|
import 'package:dart_framework/socket/packets/message_common.pb.dart';
|
|
|
|
abstract class ProtobufClient extends Communicator {
|
|
final int _headerSize = 4;
|
|
late int _heartbeatIntervalTime;
|
|
late int _heartbeatWaitTime;
|
|
final Socket _socket;
|
|
late ResponseChecker? _heartbeatChecker = null;
|
|
int? _length = null;
|
|
bool _isAlive = false;
|
|
late List<int> _arrivedData = [];
|
|
List<void Function(ProtobufClient)> _onDisconnectListenerList = [];
|
|
|
|
ProtobufClient(
|
|
this._socket,
|
|
this._heartbeatIntervalTime,
|
|
this._heartbeatWaitTime,
|
|
Map<String, GeneratedMessage Function(List<int>)> parserMap) {
|
|
print('Connected New Client');
|
|
_isAlive = true;
|
|
parserMap.addAll({
|
|
(TestData).toString(): TestData.fromBuffer,
|
|
(HeartBeat).toString(): HeartBeat.fromBuffer
|
|
});
|
|
super.initialize(parserMap);
|
|
_socket.listen(onData, onError: onError).asFuture().then(onDisconnected);
|
|
addListener(onHeartBeat);
|
|
sendHeartBeat();
|
|
}
|
|
|
|
void sendHeartBeat() {
|
|
if (_isAlive) {
|
|
_heartbeatChecker?.responsed();
|
|
_heartbeatChecker =
|
|
ResponseChecker.second(this, _heartbeatIntervalTime, (client) {
|
|
if (_isAlive) {
|
|
send(HeartBeat());
|
|
_heartbeatChecker =
|
|
ResponseChecker.second(this, _heartbeatWaitTime, (client) {
|
|
if (_isAlive) {
|
|
dispose();
|
|
onDisconnected(null);
|
|
}
|
|
});
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
void onHeartBeat(HeartBeat data) {
|
|
print('=== onHeartBeat');
|
|
}
|
|
|
|
void addDisconnectListener(void Function(ProtobufClient) handler) {
|
|
if (!_onDisconnectListenerList.contains(handler)) {
|
|
_onDisconnectListenerList.add(handler);
|
|
}
|
|
}
|
|
|
|
void removeDisconnectListener(void Function(ProtobufClient) handler) {
|
|
if (_onDisconnectListenerList.contains(handler)) {
|
|
_onDisconnectListenerList.remove(handler);
|
|
}
|
|
}
|
|
|
|
void onDisconnected(dynamic data) {
|
|
//client disconnected
|
|
if (_isAlive) {
|
|
for (var item in _onDisconnectListenerList) {
|
|
item.call(this);
|
|
}
|
|
_onDisconnectListenerList.clear();
|
|
}
|
|
}
|
|
|
|
void onError(dynamic e) {
|
|
print('=========> onError: $e');
|
|
}
|
|
|
|
void onData(Uint8List data) async {
|
|
try {
|
|
// printPacket('## Received', data);
|
|
if (_length == null) {
|
|
var header = data.sublist(0, _headerSize);
|
|
_length = header.buffer.asByteData().getInt32(0);
|
|
_arrivedData.addAll(data.sublist(_headerSize));
|
|
} else {
|
|
_arrivedData.addAll(data);
|
|
}
|
|
|
|
if (_arrivedData.length == _length) {
|
|
var common = PacketBase.fromBuffer(_arrivedData);
|
|
var nonce = common.nonce;
|
|
onReceivedData(common.typeName, common.data);
|
|
_length = null;
|
|
_arrivedData.clear();
|
|
|
|
sendHeartBeat();
|
|
}
|
|
} on Exception catch (e) {
|
|
print(e);
|
|
}
|
|
}
|
|
|
|
@override
|
|
Future send<T extends GeneratedMessage>(T data) async {
|
|
if (_isAlive) {
|
|
try {
|
|
List<int> packet = [];
|
|
var base = PacketBase();
|
|
base.typeName = data.info_.qualifiedMessageName;
|
|
base.nonce = 0;
|
|
base.data = data.writeToBuffer();
|
|
var baseBytes = base.writeToBuffer();
|
|
int length = baseBytes.length + _headerSize;
|
|
var header = Uint8List(_headerSize)
|
|
..buffer.asByteData().setInt32(0, length);
|
|
packet.addAll(header);
|
|
packet.addAll(baseBytes);
|
|
|
|
// printPacket('Send', packet);
|
|
|
|
_socket.add(packet);
|
|
await _socket.flush();
|
|
} catch (e) {
|
|
onDisconnected(null);
|
|
}
|
|
}
|
|
return simpleFuture;
|
|
}
|
|
|
|
void printPacket(String prefix, List<int> packet) {
|
|
var s = StringBuffer();
|
|
for (var item in packet) {
|
|
s.write('$item, ');
|
|
}
|
|
print('$prefix: ${s.toString()}');
|
|
}
|
|
|
|
void dispose() async {
|
|
if (_isAlive) {
|
|
_isAlive = false;
|
|
_heartbeatChecker?.responsed();
|
|
try {
|
|
await _socket.close();
|
|
_socket.destroy();
|
|
} catch (e) {}
|
|
}
|
|
}
|
|
}
|
|
|
|
class ResponseChecker<T> {
|
|
final T _responser;
|
|
late int _time;
|
|
late Timer? _timer = null;
|
|
late void Function(T) _notResponseListener;
|
|
T get responser => _responser;
|
|
ResponseChecker(this._responser, this._time, this._notResponseListener) {
|
|
startTimer();
|
|
}
|
|
|
|
ResponseChecker.second(this._responser, int time, this._notResponseListener) {
|
|
this._time = time * 1000;
|
|
startTimer();
|
|
}
|
|
|
|
void startTimer() {
|
|
_timer = Timer(Duration(milliseconds: _time), () {
|
|
_notResponseListener(_responser);
|
|
});
|
|
}
|
|
|
|
T responsed() {
|
|
_timer?.cancel();
|
|
_timer = null;
|
|
return _responser;
|
|
}
|
|
}
|