diff --git a/lib/socket/protobuf_client.dart b/lib/socket/protobuf_client.dart index 7630d81..b845243 100644 --- a/lib/socket/protobuf_client.dart +++ b/lib/socket/protobuf_client.dart @@ -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 _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 data) async { - List 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 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 { 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; } } \ No newline at end of file diff --git a/lib/socket/protobuf_server.dart b/lib/socket/protobuf_server.dart index 220ec5d..0add61b 100644 --- a/lib/socket/protobuf_server.dart +++ b/lib/socket/protobuf_server.dart @@ -34,7 +34,6 @@ abstract class ProtobufServer void onDisconnectedClient(ProtobufClient client) { - client.removeDisconnectListener(onDisconnectedClient); _clientList.remove(client); client.dispose(); print('Client disconnected');