proto-socket/dart/lib/src/communicator.dart
toki 9cc1f1d58f sync: update communicator implementation across all languages
- Align protocol documentation (PROTOCOL.md, README.md, VERSIONING.md)
- Go: add nonce test, update communicator
- Kotlin: update Communicator, TcpClient, TcpServer, add TLS test
- Python: update all modules, add certificate test resources
- TypeScript: update communicator, tcp/ws clients and servers, add tests
- Dart: update communicator, heartbeat mixin, and tests
2026-04-26 05:31:56 +09:00

239 lines
7.4 KiB
Dart

// ignore_for_file: prefer_final_fields
import 'dart:async';
import 'package:meta/meta.dart';
import 'package:protobuf/protobuf.dart';
import 'packets/message_common.pb.dart';
import 'transport.dart';
abstract class Communicator {
static const int maxNonce = 2147483647;
Map<String, IDataHandler> _handlerDic = {};
Map<String, Future<void> Function(List<int>, int)> _requestHandlerDic = {};
Map<int, _PendingRequest> _pendingRequests = {};
late Map<String, GeneratedMessage Function(List<int>)> _instanceGenerator;
Future<void> _outboundWrite = Future.value();
Transport? _transport;
/// Monotonically increasing nonce. Shared by send / sendRequest / response.
int _nonce = 0;
int get nonce => _nonce;
@protected
set nonce(int value) => _nonce = value;
@protected
int nextNonce() {
if (_nonce >= maxNonce) {
_nonce = 0;
}
_nonce += 1;
return _nonce;
}
/// Whether the connection is alive. Set by each transport implementation.
bool _isAlive = false;
bool get isAlive => _isAlive;
@protected
set isAlive(bool value) => _isAlive = value;
Communicator();
void initialize(
Map<String, GeneratedMessage Function(List<int>)> instanceGenerator,
{required Transport transport}) {
_instanceGenerator = instanceGenerator;
_transport = transport;
}
T Function(List<int>) getGenerator<T extends GeneratedMessage>(String type) {
if (!_instanceGenerator.containsKey(type)) {
throw Exception(
'Must set protobuf packet creator before use it. Type: ${(T).toString()}');
}
return _instanceGenerator[type] as T Function(List<int>);
}
/// Serializes writes so stream transports do not interleave packets.
Future<void> queuePacket(PacketBase base) {
final write = _outboundWrite.then((_) => _transport!.writePacket(base));
_outboundWrite = write.catchError((_) {});
return write;
}
Future<void> send<T extends GeneratedMessage>(T data) async {
if (isAlive) {
await queuePacket(PacketBase()
..typeName = data.info_.qualifiedMessageName
..nonce = nextNonce()
..data = data.writeToBuffer());
}
}
@protected
void cancelPendingRequests() {
final snapshot = Map<int, _PendingRequest>.from(_pendingRequests);
_pendingRequests.clear();
for (final pending in snapshot.values) {
pending.completeError(
StateError('connection closed'), StackTrace.current);
}
}
/// Sends [data] as a request and waits for a typed response.
///
/// The remote side must have registered an [addRequestListener] for [Req].
/// [Res] must be registered in the parser map.
///
/// ```dart
/// final res = await client.sendRequest<GetUser, UserData>(GetUser()..id = 1);
/// ```
Future<Res>
sendRequest<Req extends GeneratedMessage, Res extends GeneratedMessage>(
Req data,
{Duration timeout = const Duration(seconds: 30)}) async {
if (!isAlive) return Future.error(StateError('not connected'));
final requestNonce = nextNonce();
final completer = Completer<Res>();
final expectedResponseType = Res.toString();
_pendingRequests[requestNonce] = _PendingRequest(
expectedTypeName: expectedResponseType,
complete: (bytes) {
completer.complete(getGenerator<Res>(expectedResponseType)(bytes));
},
completeError: (error, stackTrace) {
completer.completeError(error, stackTrace);
},
);
await queuePacket(PacketBase()
..typeName = data.info_.qualifiedMessageName
..nonce = requestNonce
..data = data.writeToBuffer())
.catchError((error, stackTrace) {
_pendingRequests.remove(requestNonce);
completer.completeError(error, stackTrace);
});
return completer.future.timeout(timeout, onTimeout: () {
_pendingRequests.remove(requestNonce);
throw TimeoutException(
'sendRequest timeout for nonce $requestNonce', timeout);
});
}
/// Registers a request handler for [Req] that returns [Res].
///
/// When a [Req] packet arrives the handler is called and the returned [Res]
/// is automatically sent back to the caller.
///
/// ```dart
/// server.addRequestListener<GetUser, UserData>((req) async {
/// return UserData()..name = db.getUser(req.id).name;
/// });
/// ```
void addRequestListener<Req extends GeneratedMessage,
Res extends GeneratedMessage>(Future<Res> Function(Req) handler) {
final reqType = Req.toString();
if (_handlerDic.containsKey(reqType)) {
throw StateError(
'Type $reqType is already registered with addListener and cannot also use addRequestListener.');
}
_requestHandlerDic[reqType] = (List<int> bytes, int requestNonce) async {
final req = getGenerator<Req>(reqType)(bytes);
final res = await handler(req);
if (isAlive) {
await queuePacket(PacketBase()
..typeName = res.info_.qualifiedMessageName
..nonce = nextNonce()
..responseNonce = requestNonce
..data = res.writeToBuffer());
}
};
}
void onReceivedData(String typeName, List<int> data,
{int incomingNonce = 0, int responseNonce = 0}) {
if (responseNonce > 0) {
final pending = _pendingRequests.remove(responseNonce);
if (pending == null) return;
if (typeName != pending.expectedTypeName) {
pending.completeError(
StateError(
'Response type mismatch for nonce $responseNonce: expected ${pending.expectedTypeName}, got $typeName'),
StackTrace.current,
);
return;
}
pending.complete(data);
return;
}
if (_requestHandlerDic.containsKey(typeName)) {
_requestHandlerDic[typeName]!(data, incomingNonce);
return;
}
if (_handlerDic.containsKey(typeName)) {
_handlerDic[typeName]?.onMessage(data);
}
}
void addListener<T extends GeneratedMessage>(void Function(T) listener) {
var type = T.toString();
if (_requestHandlerDic.containsKey(type)) {
throw StateError(
'Type $type is already registered with addRequestListener and cannot also use addListener.');
}
if (!_handlerDic.containsKey(type)) {
_handlerDic[type] = DataHandler<T>(getGenerator(type));
}
var handler = _handlerDic[type] as DataHandler<T>;
handler.addListener(listener);
}
void removeListener<T extends GeneratedMessage>(void Function(T) listener) {
var type = T.toString();
if (_handlerDic.containsKey(type)) {
var handler = _handlerDic[type] as DataHandler<T>;
handler.removeListener(listener);
}
}
}
abstract class IDataHandler {
void onMessage(List<int> data);
}
class DataHandler<T extends GeneratedMessage> implements IDataHandler {
T Function(List<int>) _generator;
List<void Function(T)> _listeners = [];
DataHandler(this._generator);
@override
void onMessage(List<int> data) {
for (var listener in _listeners) {
listener.call(_generator(data));
}
}
void addListener(void Function(T) handler) {
removeListener(handler); // 중복 리스너 불허
_listeners.add(handler);
}
void removeListener(void Function(T) handler) {
_listeners.remove(handler);
}
}
class _PendingRequest {
final String expectedTypeName;
final void Function(List<int>) complete;
final void Function(Object, StackTrace) completeError;
_PendingRequest({
required this.expectedTypeName,
required this.complete,
required this.completeError,
});
}