dart-app-core/lib/socket/protobuf_client.dart
build@lguplus.co.kr a0d8e93d6c 하트비트 추가
2023-12-16 00:30:06 +09:00

168 lines
No EOL
4.4 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;
final int _heartbeatIntervalTime;
final int _heartbeatWaitTime;
final Socket _socket;
late ResponseChecker? _heartbeatChecker;
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);
addListener(onHeartBeat);
_socket.listen(onData, onError: onError)
.asFuture().then(onDisconnected);
sendHeartBeat();
}
void sendHeartBeat() {
_heartbeatChecker?.responsed();
_heartbeatChecker = ResponseChecker.seecond(this, _heartbeatIntervalTime, (client) {
send(HeartBeat());
_heartbeatChecker = ResponseChecker.seecond(this, _heartbeatWaitTime, (client) {
dispose();
onDisconnected(null);
});
});
}
void onHeartBeat(HeartBeat data) {
sendHeartBeat();
}
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
for (var item in _onDisconnectListenerList) {
item.call(this);
}
_onDisconnectListenerList.clear();
}
void onError(dynamic e) {
print('=========> onError: $e');
}
void onData(Uint8List data) async {
try
{
sendHeartBeat();
// 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();
}
} on Exception catch(e) {
print(e);
}
}
@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);
// printPacket('Send', packet);
_socket.add(packet);
await _socket.flush();
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();
await _socket.close();
_socket.destroy();
}
}
}
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);
});
}
ResponseChecker.seecond(this._responser, int time, void Function(T) notResponseListener) {
ResponseChecker(this._responser, time * 1000, notResponseListener);
}
T responsed() {
_timer.cancel();
return _responser;
}
}