import 'dart:async'; import 'dart:mirrors'; import 'package:test/test.dart'; import 'package:proto_socket/proto_socket.dart'; class _FakeCommunicator extends Communicator { final _FakeTransport transport = _FakeTransport(); _FakeCommunicator({bool packageQualifiedParser = false}) { isAlive = true; final testDataType = packageQualifiedParser ? 'example.TestData' : TestData.getDefault().info_.qualifiedMessageName; initialize({ testDataType: TestData.fromBuffer, HeartBeat.getDefault().info_.qualifiedMessageName: HeartBeat.fromBuffer, }, transport: transport); } void closeForTest() { isAlive = false; cancelPendingRequests(); } void setNonceForTest(int value) { nonce = value; } } class _FakeTransport implements Transport { final sentPackets = []; Object? error; @override Future writePacket(PacketBase base) async { final error = this.error; if (error != null) { throw error; } sentPackets.add(base); } @override Future close() async {} } Map _pendingRequestsOf(Communicator communicator) { final library = reflectClass(Communicator).owner as LibraryMirror; final symbol = MirrorSystem.getSymbol('_pendingRequests', library); return reflect(communicator).getField(symbol).reflectee as Map; } void main() { group('Communicator protocol guards', () { test('response typeName mismatch completes sendRequest with error', () async { final communicator = _FakeCommunicator(); final future = communicator.sendRequest( TestData() ..index = 1 ..message = 'hello', ); await Future.delayed(Duration.zero); final requestNonce = communicator.transport.sentPackets.single.nonce; communicator.onReceivedData( HeartBeat.getDefault().info_.qualifiedMessageName, HeartBeat().writeToBuffer(), responseNonce: requestNonce, ); await expectLater( future, throwsA( isA().having((error) => error.toString(), 'message', contains('Response type mismatch')), ), ); }); test('close 후 sendRequest는 StateError로 완료된다', () async { final communicator = _FakeCommunicator(); final future = communicator.sendRequest( TestData()..index = 1, ); await Future.delayed(Duration.zero); communicator.closeForTest(); await expectLater( future, throwsA( isA().having( (error) => error.message, 'message', contains('connection closed'), ), ), ); }); test('sendRequest가 timeout 내 응답 없으면 TimeoutException을 던진다', () async { final communicator = _FakeCommunicator(); final future = communicator.sendRequest( TestData() ..index = 1 ..message = 'wait', timeout: const Duration(milliseconds: 50), ); await Future.delayed(Duration.zero); expect(communicator.transport.sentPackets, hasLength(1)); await expectLater( future, throwsA(isA()), ); expect(_pendingRequestsOf(communicator), isEmpty); }); test('nonce wraps after int32 max without emitting zero', () async { final communicator = _FakeCommunicator(); communicator.setNonceForTest(Communicator.maxNonce - 1); await communicator.send(TestData()..index = 1); await communicator.send(TestData()..index = 2); expect( communicator.transport.sentPackets.map((packet) => packet.nonce), [Communicator.maxNonce, 1], ); expect( communicator.transport.sentPackets.map((packet) => packet.nonce), isNot(contains(0)), ); }); test('sendRequest matches response at nonce wrap boundary', () async { final communicator = _FakeCommunicator(); communicator.setNonceForTest(Communicator.maxNonce - 1); final maxNonceFuture = communicator.sendRequest( TestData() ..index = 1 ..message = 'max', ); await Future.delayed(Duration.zero); final maxNonceRequest = communicator.transport.sentPackets.single; expect(maxNonceRequest.nonce, Communicator.maxNonce); expect(maxNonceRequest.nonce, isNot(0)); communicator.onReceivedData( TestData.getDefault().info_.qualifiedMessageName, (TestData() ..index = 2 ..message = 'max response') .writeToBuffer(), responseNonce: maxNonceRequest.nonce, ); expect((await maxNonceFuture).message, 'max response'); final wrappedFuture = communicator.sendRequest( TestData() ..index = 3 ..message = 'wrapped', ); await Future.delayed(Duration.zero); final wrappedRequest = communicator.transport.sentPackets.last; expect(wrappedRequest.nonce, 1); expect(wrappedRequest.nonce, isNot(0)); communicator.onReceivedData( TestData.getDefault().info_.qualifiedMessageName, (TestData() ..index = 4 ..message = 'wrapped response') .writeToBuffer(), responseNonce: wrappedRequest.nonce, ); expect((await wrappedFuture).message, 'wrapped response'); }); test('sendRequest accepts package-qualified response typeName', () async { final communicator = _FakeCommunicator(packageQualifiedParser: true); final future = communicator.sendRequest( TestData() ..index = 1 ..message = 'request', ); await Future.delayed(Duration.zero); final requestNonce = communicator.transport.sentPackets.single.nonce; communicator.onReceivedData( 'example.TestData', (TestData() ..index = 2 ..message = 'qualified response') .writeToBuffer(), responseNonce: requestNonce, ); expect((await future).message, 'qualified response'); }); test('addListener accepts package-qualified message typeName', () { final communicator = _FakeCommunicator(packageQualifiedParser: true); final messages = []; communicator.addListener(messages.add); communicator.onReceivedData( 'example.TestData', (TestData() ..index = 3 ..message = 'qualified event') .writeToBuffer(), ); expect(messages, hasLength(1)); expect(messages.single.message, 'qualified event'); }); test('addRequestListener accepts package-qualified request typeName', () async { final communicator = _FakeCommunicator(packageQualifiedParser: true); communicator.addRequestListener((req) async { return TestData() ..index = req.index + 1 ..message = 'echo: ${req.message}'; }); communicator.onReceivedData( 'example.TestData', (TestData() ..index = 4 ..message = 'qualified request') .writeToBuffer(), incomingNonce: 12, ); await Future.delayed(Duration.zero); final response = communicator.transport.sentPackets.single; expect(response.responseNonce, 12); final data = TestData.fromBuffer(response.data); expect(data.index, 5); expect(data.message, 'echo: qualified request'); }); test( 'cannot register addRequestListener for a type already using addListener', () { final communicator = _FakeCommunicator(); communicator.addListener((_) {}); expect( () => communicator.addRequestListener( (req) async => TestData()..index = req.index, ), throwsA(isA()), ); }); test( 'cannot register addListener for a type already using addRequestListener', () { final communicator = _FakeCommunicator(); communicator.addRequestListener( (req) async => TestData()..index = req.index, ); expect( () => communicator.addListener((_) {}), throwsA(isA()), ); }); test('queuePacket uses injected transport', () async { final communicator = _FakeCommunicator(); final packet = PacketBase() ..typeName = TestData.getDefault().info_.qualifiedMessageName ..nonce = 1 ..data = (TestData()..index = 7).writeToBuffer(); await communicator.queuePacket(packet); expect(communicator.transport.sentPackets, hasLength(1)); expect(communicator.transport.sentPackets.single, same(packet)); final error = StateError('write failed'); communicator.transport.error = error; await expectLater( communicator.queuePacket(PacketBase()..typeName = 'TestData'), throwsA(same(error)), ); }); }); }