하트비트 수정
This commit is contained in:
parent
a0d8e93d6c
commit
777fac5d84
2 changed files with 68 additions and 42 deletions
|
|
@ -13,10 +13,10 @@ import 'package:dart_framework/socket/packets/message_common.pb.dart';
|
|||
abstract class ProtobufClient extends Communicator
|
||||
{
|
||||
final int _headerSize = 4;
|
||||
final int _heartbeatIntervalTime;
|
||||
final int _heartbeatWaitTime;
|
||||
late int _heartbeatIntervalTime;
|
||||
late int _heartbeatWaitTime;
|
||||
final Socket _socket;
|
||||
late ResponseChecker? _heartbeatChecker;
|
||||
late ResponseChecker? _heartbeatChecker = null;
|
||||
int? _length = null;
|
||||
bool _isAlive = false;
|
||||
late List<int> _arrivedData = [];
|
||||
|
|
@ -31,25 +31,32 @@ abstract class ProtobufClient extends Communicator
|
|||
(HeartBeat).toString() : HeartBeat.fromBuffer
|
||||
});
|
||||
super.initialize(parserMap);
|
||||
addListener(onHeartBeat);
|
||||
_socket.listen(onData, onError: onError)
|
||||
.asFuture().then(onDisconnected);
|
||||
addListener(onHeartBeat);
|
||||
sendHeartBeat();
|
||||
}
|
||||
|
||||
void sendHeartBeat() {
|
||||
_heartbeatChecker?.responsed();
|
||||
_heartbeatChecker = ResponseChecker.seecond(this, _heartbeatIntervalTime, (client) {
|
||||
send(HeartBeat());
|
||||
_heartbeatChecker = ResponseChecker.seecond(this, _heartbeatWaitTime, (client) {
|
||||
dispose();
|
||||
onDisconnected(null);
|
||||
if(_isAlive) {
|
||||
_heartbeatChecker?.responsed();
|
||||
_heartbeatChecker = ResponseChecker.second(this, _heartbeatIntervalTime, (client) {
|
||||
if(_isAlive) {
|
||||
print('=== sendHartBeat');
|
||||
send(HeartBeat());
|
||||
_heartbeatChecker = ResponseChecker.second(this, _heartbeatWaitTime, (client) {
|
||||
if(_isAlive) {
|
||||
dispose();
|
||||
onDisconnected(null);
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
void onHeartBeat(HeartBeat data) {
|
||||
sendHeartBeat();
|
||||
print('=== onHeartBeat');
|
||||
}
|
||||
|
||||
void addDisconnectListener(void Function(ProtobufClient) handler) {
|
||||
|
|
@ -68,10 +75,13 @@ abstract class ProtobufClient extends Communicator
|
|||
|
||||
void onDisconnected(dynamic data) {
|
||||
//client disconnected
|
||||
for (var item in _onDisconnectListenerList) {
|
||||
item.call(this);
|
||||
if(_isAlive)
|
||||
{
|
||||
for (var item in _onDisconnectListenerList) {
|
||||
item.call(this);
|
||||
}
|
||||
_onDisconnectListenerList.clear();
|
||||
}
|
||||
_onDisconnectListenerList.clear();
|
||||
}
|
||||
|
||||
void onError(dynamic e) {
|
||||
|
|
@ -81,7 +91,6 @@ abstract class ProtobufClient extends Communicator
|
|||
void onData(Uint8List data) async {
|
||||
try
|
||||
{
|
||||
sendHeartBeat();
|
||||
// printPacket('## Received', data);
|
||||
if(_length == null)
|
||||
{
|
||||
|
|
@ -101,6 +110,8 @@ abstract class ProtobufClient extends Communicator
|
|||
onReceivedData(common.typeName, common.data);
|
||||
_length = null;
|
||||
_arrivedData.clear();
|
||||
|
||||
sendHeartBeat();
|
||||
}
|
||||
} on Exception catch(e) {
|
||||
print(e);
|
||||
|
|
@ -109,22 +120,28 @@ abstract class ProtobufClient extends Communicator
|
|||
|
||||
@override
|
||||
Future send<T extends GeneratedMessage>(T data) async {
|
||||
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);
|
||||
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);
|
||||
// printPacket('Send', packet);
|
||||
|
||||
_socket.add(packet);
|
||||
await _socket.flush();
|
||||
_socket.add(packet);
|
||||
await _socket.flush();
|
||||
} catch(e, s) {
|
||||
onDisconnected(null);
|
||||
}
|
||||
}
|
||||
return simpleFuture;
|
||||
}
|
||||
|
||||
|
|
@ -142,27 +159,37 @@ abstract class ProtobufClient extends Communicator
|
|||
{
|
||||
_isAlive = false;
|
||||
_heartbeatChecker?.responsed();
|
||||
await _socket.close();
|
||||
_socket.destroy();
|
||||
try {
|
||||
await _socket.close();
|
||||
_socket.destroy();
|
||||
} catch (e, s) {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
class ResponseChecker<T> {
|
||||
final T _responser;
|
||||
late Timer _timer;
|
||||
ResponseChecker(this._responser, int time, void Function(T) notResponseListener) {
|
||||
_timer = Timer(Duration(milliseconds: time), () {
|
||||
notResponseListener(_responser);
|
||||
late int _time;
|
||||
late Timer? _timer = null;
|
||||
late void Function(T) _notResponseListener;
|
||||
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);
|
||||
});
|
||||
}
|
||||
|
||||
ResponseChecker.seecond(this._responser, int time, void Function(T) notResponseListener) {
|
||||
ResponseChecker(this._responser, time * 1000, notResponseListener);
|
||||
}
|
||||
|
||||
T responsed() {
|
||||
_timer.cancel();
|
||||
_timer?.cancel();
|
||||
_timer = null;
|
||||
return _responser;
|
||||
}
|
||||
}
|
||||
|
|
@ -34,7 +34,6 @@ abstract class ProtobufServer
|
|||
|
||||
void onDisconnectedClient(ProtobufClient client)
|
||||
{
|
||||
client.removeDisconnectListener(onDisconnectedClient);
|
||||
_clientList.remove(client);
|
||||
client.dispose();
|
||||
print('Client disconnected');
|
||||
|
|
|
|||
Loading…
Reference in a new issue