dart-app-core/lib/socket/protobuf_client.dart
2024-01-11 11:53:53 +09:00

195 lines
No EOL
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, s) {
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, s) {}
}
}
}
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;
}
}