1066 lines
34 KiB
Dart
1066 lines
34 KiB
Dart
import 'dart:async';
|
|
import 'dart:io';
|
|
import 'dart:typed_data';
|
|
|
|
import 'package:test/test.dart';
|
|
import 'package:proto_socket/proto_socket.dart';
|
|
|
|
const _testPort = 19090;
|
|
const _testPortSsl = 19091;
|
|
const _testPortWs = 19092;
|
|
const _testPortWss = 19093;
|
|
const _host = '127.0.0.1';
|
|
const _certPath = 'test/certs/server.crt';
|
|
const _keyPath = 'test/certs/server.key';
|
|
|
|
bool _acceptTestCertificate(X509Certificate certificate) =>
|
|
certificate.pem == File(_certPath).readAsStringSync();
|
|
|
|
// ── 테스트 전용 픽스처 ──────────────────────────────────────────
|
|
|
|
class _TestClient extends ProtobufClient {
|
|
_TestClient(Socket socket)
|
|
: super(socket, 5, 3, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
}
|
|
|
|
class _TestServer extends ProtobufServer {
|
|
final receivedMessages = <TestData>[];
|
|
final connectedClients = <ProtobufClient>[];
|
|
|
|
_TestServer() : super(_host, _testPort, (socket) => _TestClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addListener<TestData>((data) => receivedMessages.add(data));
|
|
}
|
|
}
|
|
|
|
class _TestServerSsl extends ProtobufServer {
|
|
final receivedMessages = <TestData>[];
|
|
final connectedClients = <ProtobufClient>[];
|
|
|
|
_TestServerSsl(SecurityContext ctx)
|
|
: super.secure(_host, _testPortSsl, ctx, (socket) => _TestClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addListener<TestData>((data) => receivedMessages.add(data));
|
|
}
|
|
}
|
|
|
|
class _TestWsClient extends WsProtobufClient {
|
|
_TestWsClient(WebSocket ws)
|
|
: super(ws, 5, 3, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
}
|
|
|
|
class _TestWsServer extends WsProtobufServer {
|
|
final receivedMessages = <TestData>[];
|
|
final connectedClients = <WsProtobufClient>[];
|
|
|
|
_TestWsServer() : super(_host, _testPortWs, (ws) => _TestWsClient(ws));
|
|
|
|
@override
|
|
void onClientConnected(WsProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addListener<TestData>((data) => receivedMessages.add(data));
|
|
}
|
|
}
|
|
|
|
class _TestWsServerSsl extends WsProtobufServer {
|
|
final receivedMessages = <TestData>[];
|
|
final connectedClients = <WsProtobufClient>[];
|
|
|
|
_TestWsServerSsl(SecurityContext ctx)
|
|
: super.secure(_host, _testPortWss, ctx, (ws) => _TestWsClient(ws));
|
|
|
|
@override
|
|
void onClientConnected(WsProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addListener<TestData>((data) => receivedMessages.add(data));
|
|
}
|
|
}
|
|
|
|
class _TestReqServer extends ProtobufServer {
|
|
_TestReqServer() : super(_host, _testPort, (socket) => _TestClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {
|
|
client.addRequestListener<TestData, TestData>((req) async {
|
|
return TestData()
|
|
..index = req.index * 2
|
|
..message = 'echo: ${req.message}';
|
|
});
|
|
}
|
|
}
|
|
|
|
class _TestWsReqServer extends WsProtobufServer {
|
|
_TestWsReqServer([int port = _testPortWs])
|
|
: super(_host, port, (ws) => _TestWsClient(ws));
|
|
|
|
@override
|
|
void onClientConnected(WsProtobufClient client) {
|
|
client.addRequestListener<TestData, TestData>((req) async {
|
|
return TestData()
|
|
..index = req.index * 2
|
|
..message = 'echo: ${req.message}';
|
|
});
|
|
}
|
|
}
|
|
|
|
class _FastHeartbeatClient extends ProtobufClient {
|
|
_FastHeartbeatClient(Socket socket)
|
|
: super(socket, 1, 1, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
}
|
|
|
|
class _NoHeartbeatClient extends ProtobufClient {
|
|
_NoHeartbeatClient(Socket socket)
|
|
: super(socket, 5, 3, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
|
|
@override
|
|
void onHeartBeat(HeartBeat data) {}
|
|
}
|
|
|
|
class _NoHeartbeatServer extends ProtobufServer {
|
|
_NoHeartbeatServer()
|
|
: super(_host, _testPort, (socket) => _NoHeartbeatClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {}
|
|
}
|
|
|
|
class _DisabledHeartbeatClient extends ProtobufClient {
|
|
_DisabledHeartbeatClient(Socket socket)
|
|
: super(socket, 0, 0, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
}
|
|
|
|
class _DisabledHeartbeatServer extends ProtobufServer {
|
|
final receivedMessages = <TestData>[];
|
|
final connectedClients = <ProtobufClient>[];
|
|
|
|
_DisabledHeartbeatServer()
|
|
: super(_host, _testPort, (socket) => _DisabledHeartbeatClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addListener<TestData>((data) => receivedMessages.add(data));
|
|
}
|
|
}
|
|
|
|
class _FastHeartbeatWsClient extends WsProtobufClient {
|
|
_FastHeartbeatWsClient(WebSocket ws)
|
|
: super(ws, 1, 1, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
}
|
|
|
|
class _NoHeartbeatWsClient extends WsProtobufClient {
|
|
_NoHeartbeatWsClient(WebSocket ws)
|
|
: super(ws, 5, 3, {
|
|
TestData.getDefault().info_.qualifiedMessageName: TestData.fromBuffer,
|
|
});
|
|
|
|
@override
|
|
void onHeartBeat(HeartBeat data) {}
|
|
}
|
|
|
|
class _NoHeartbeatWsServer extends WsProtobufServer {
|
|
_NoHeartbeatWsServer()
|
|
: super(_host, _testPortWs, (ws) => _NoHeartbeatWsClient(ws));
|
|
|
|
@override
|
|
void onClientConnected(WsProtobufClient client) {}
|
|
}
|
|
|
|
// ── 테스트 ──────────────────────────────────────────────────────
|
|
|
|
void main() {
|
|
// ── Plain TCP ────────────────────────────────────────────────
|
|
|
|
group('ProtobufServer (plain)', () {
|
|
late _TestServer server;
|
|
|
|
setUp(() async {
|
|
server = _TestServer();
|
|
await server.start();
|
|
});
|
|
|
|
tearDown(() async {
|
|
await server.stop();
|
|
});
|
|
|
|
test('서버가 정상 시작된다', () {
|
|
expect(server.started, isTrue);
|
|
expect(server.isSecure, isFalse);
|
|
});
|
|
});
|
|
|
|
group('ProtobufClient (plain)', () {
|
|
late _TestServer server;
|
|
late _TestClient client;
|
|
|
|
setUp(() async {
|
|
server = _TestServer();
|
|
await server.start();
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
client = _TestClient(socket);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('클라이언트가 서버에 연결된다', () {
|
|
expect(server.connectedClients, isNotEmpty);
|
|
});
|
|
|
|
test('TestData 메시지를 서버가 수신한다', () async {
|
|
await client.send(TestData()
|
|
..index = 42
|
|
..message = 'hello proto-socket');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(42));
|
|
expect(
|
|
server.receivedMessages.first.message, equals('hello proto-socket'));
|
|
});
|
|
|
|
test('여러 메시지를 순서대로 수신한다', () async {
|
|
for (var i = 0; i < 5; i++) {
|
|
await client.send(TestData()
|
|
..index = i
|
|
..message = 'msg$i');
|
|
}
|
|
await Future.delayed(const Duration(milliseconds: 300));
|
|
|
|
expect(server.receivedMessages, hasLength(5));
|
|
for (var i = 0; i < 5; i++) {
|
|
expect(server.receivedMessages[i].index, equals(i));
|
|
}
|
|
});
|
|
|
|
test('서버에서 클라이언트로 메시지를 전송한다', () async {
|
|
final completer = Completer<TestData>();
|
|
client.addListener<TestData>((data) => completer.complete(data));
|
|
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.connectedClients.first.send(TestData()
|
|
..index = 99
|
|
..message = 'from server');
|
|
|
|
final received =
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(received.index, equals(99));
|
|
expect(received.message, equals('from server'));
|
|
});
|
|
|
|
test('nonce가 송신마다 증가한다', () async {
|
|
for (var i = 1; i <= 3; i++) {
|
|
await client.send(TestData()..index = i);
|
|
}
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
expect(server.receivedMessages, hasLength(3));
|
|
});
|
|
|
|
test('클라이언트 disconnect 시 서버 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
server.connectedClients.first.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await client.close();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(completer.isCompleted, isTrue);
|
|
});
|
|
|
|
test('서버 stop 시 클라이언트 disconnect 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
client.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await server.stop();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
test('HeartBeat interval 동안 연결이 유지된다', () async {
|
|
await Future.delayed(const Duration(seconds: 2));
|
|
await client.send(TestData()
|
|
..index = 1
|
|
..message = 'alive check');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
expect(server.receivedMessages, hasLength(1));
|
|
});
|
|
|
|
test('close 후 isAlive가 false다', () async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
test('서버가 모든 클라이언트에게 브로드캐스트한다', () async {
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
final secondClient = _TestClient(socket);
|
|
try {
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
final firstReceived = Completer<TestData>();
|
|
final secondReceived = Completer<TestData>();
|
|
client.addListener<TestData>((data) {
|
|
if (!firstReceived.isCompleted) firstReceived.complete(data);
|
|
});
|
|
secondClient.addListener<TestData>((data) {
|
|
if (!secondReceived.isCompleted) secondReceived.complete(data);
|
|
});
|
|
|
|
await server.broadcast(TestData()
|
|
..index = 77
|
|
..message = 'broadcast');
|
|
|
|
final first =
|
|
await firstReceived.future.timeout(const Duration(seconds: 2));
|
|
final second =
|
|
await secondReceived.future.timeout(const Duration(seconds: 2));
|
|
expect(first.message, equals('broadcast'));
|
|
expect(second.message, equals('broadcast'));
|
|
} finally {
|
|
await secondClient.close();
|
|
}
|
|
});
|
|
|
|
test('close 후 send는 무시된다', () async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
|
|
await expectLater(
|
|
client.send(TestData()
|
|
..index = 1
|
|
..message = 'ignored'),
|
|
completes,
|
|
);
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
test('TCP closes on oversized packet length', () async {
|
|
final rawSocket = await Socket.connect(_host, _testPort);
|
|
try {
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
final serverClient = server.connectedClients.last;
|
|
final disconnected = Completer<void>();
|
|
serverClient.addDisconnectListener((_) {
|
|
if (!disconnected.isCompleted) disconnected.complete();
|
|
});
|
|
|
|
final header = Uint8List(4)
|
|
..buffer.asByteData().setInt32(0, maxPacketSize + 1);
|
|
rawSocket.add(header);
|
|
await rawSocket.flush();
|
|
|
|
await disconnected.future.timeout(const Duration(seconds: 2));
|
|
expect(serverClient.isAlive, isFalse);
|
|
} finally {
|
|
rawSocket.destroy();
|
|
}
|
|
});
|
|
});
|
|
|
|
group('Heartbeat timeout', () {
|
|
test('TCP heartbeat 타임아웃 시 disconnect 콜백이 호출된다', () async {
|
|
final server = _NoHeartbeatServer();
|
|
await server.start();
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
final client = _FastHeartbeatClient(socket);
|
|
final completer = Completer<void>();
|
|
client.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
try {
|
|
await completer.future.timeout(const Duration(seconds: 4));
|
|
expect(client.isAlive, isFalse);
|
|
} finally {
|
|
await client.close();
|
|
await server.stop();
|
|
}
|
|
});
|
|
|
|
test('WS heartbeat 타임아웃 시 disconnect 콜백이 호출된다', () async {
|
|
final server = _NoHeartbeatWsServer();
|
|
await server.start();
|
|
final ws = await WsProtobufClient.connect(_host, _testPortWs);
|
|
final client = _FastHeartbeatWsClient(ws);
|
|
final completer = Completer<void>();
|
|
client.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
try {
|
|
await completer.future.timeout(const Duration(seconds: 4));
|
|
expect(client.isAlive, isFalse);
|
|
} finally {
|
|
await client.close();
|
|
await server.stop();
|
|
}
|
|
});
|
|
});
|
|
|
|
group('Heartbeat disabled (interval 0)', () {
|
|
test('interval 0이면 timer churn 없이 연결이 유지되고 송수신된다', () async {
|
|
final server = _DisabledHeartbeatServer();
|
|
await server.start();
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
final client = _DisabledHeartbeatClient(socket);
|
|
|
|
try {
|
|
// heartbeat가 비활성이면 wait timeout으로 self-close되지 않는다.
|
|
await Future.delayed(const Duration(milliseconds: 600));
|
|
expect(client.isAlive, isTrue);
|
|
expect(server.connectedClients.single.isAlive, isTrue);
|
|
|
|
await client.send(TestData()
|
|
..index = 7
|
|
..message = 'disabled hb');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(7));
|
|
expect(server.receivedMessages.first.message, equals('disabled hb'));
|
|
} finally {
|
|
await client.close();
|
|
await server.stop();
|
|
}
|
|
});
|
|
});
|
|
|
|
// ── SSL/TLS ──────────────────────────────────────────────────
|
|
|
|
group('ProtobufServer (SSL)', () {
|
|
late _TestServerSsl server;
|
|
|
|
setUp(() async {
|
|
final ctx = SecurityContext()
|
|
..useCertificateChain(_certPath)
|
|
..usePrivateKey(_keyPath);
|
|
server = _TestServerSsl(ctx);
|
|
await server.start();
|
|
});
|
|
|
|
tearDown(() async {
|
|
await server.stop();
|
|
});
|
|
|
|
test('SSL 서버가 정상 시작된다', () {
|
|
expect(server.started, isTrue);
|
|
expect(server.isSecure, isTrue);
|
|
});
|
|
});
|
|
|
|
group('ProtobufClient (SSL)', () {
|
|
late _TestServerSsl server;
|
|
late _TestClient client;
|
|
|
|
setUp(() async {
|
|
final serverCtx = SecurityContext()
|
|
..useCertificateChain(_certPath)
|
|
..usePrivateKey(_keyPath);
|
|
server = _TestServerSsl(serverCtx);
|
|
await server.start();
|
|
|
|
final clientCtx = SecurityContext()..setTrustedCertificates(_certPath);
|
|
final socket = await ProtobufClient.connectSecure(
|
|
_host,
|
|
_testPortSsl,
|
|
context: clientCtx,
|
|
onBadCertificate: _acceptTestCertificate,
|
|
);
|
|
client = _TestClient(socket);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('SSL 클라이언트가 서버에 연결된다', () {
|
|
expect(server.connectedClients, isNotEmpty);
|
|
});
|
|
|
|
test('SSL TestData 메시지를 서버가 수신한다', () async {
|
|
await client.send(TestData()
|
|
..index = 7
|
|
..message = 'hello over ssl');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(7));
|
|
expect(server.receivedMessages.first.message, equals('hello over ssl'));
|
|
});
|
|
|
|
test('SSL 서버에서 클라이언트로 메시지를 전송한다', () async {
|
|
final completer = Completer<TestData>();
|
|
client.addListener<TestData>((data) => completer.complete(data));
|
|
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.connectedClients.first.send(TestData()
|
|
..index = 55
|
|
..message = 'from ssl server');
|
|
|
|
final received =
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(received.index, equals(55));
|
|
expect(received.message, equals('from ssl server'));
|
|
});
|
|
|
|
test('SSL disconnect 시 서버 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
server.connectedClients.first.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await client.close();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(completer.isCompleted, isTrue);
|
|
});
|
|
});
|
|
|
|
// ── WebSocket ────────────────────────────────────────────────
|
|
|
|
group('WsProtobufServer (plain)', () {
|
|
late _TestWsServer server;
|
|
|
|
setUp(() async {
|
|
server = _TestWsServer();
|
|
await server.start();
|
|
});
|
|
|
|
tearDown(() async {
|
|
await server.stop();
|
|
});
|
|
|
|
test('WS 서버가 정상 시작된다', () {
|
|
expect(server.started, isTrue);
|
|
expect(server.isSecure, isFalse);
|
|
});
|
|
});
|
|
|
|
group('WsProtobufClient (plain)', () {
|
|
late _TestWsServer server;
|
|
late _TestWsClient client;
|
|
|
|
setUp(() async {
|
|
server = _TestWsServer();
|
|
await server.start();
|
|
final ws = await WsProtobufClient.connect(_host, _testPortWs);
|
|
client = _TestWsClient(ws);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('WS 클라이언트가 서버에 연결된다', () {
|
|
expect(server.connectedClients, isNotEmpty);
|
|
});
|
|
|
|
test('WS TestData 메시지를 서버가 수신한다', () async {
|
|
await client.send(TestData()
|
|
..index = 42
|
|
..message = 'hello ws');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(42));
|
|
expect(server.receivedMessages.first.message, equals('hello ws'));
|
|
});
|
|
|
|
test('WS 서버에서 클라이언트로 메시지를 전송한다', () async {
|
|
final completer = Completer<TestData>();
|
|
client.addListener<TestData>((data) => completer.complete(data));
|
|
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.connectedClients.first.send(TestData()
|
|
..index = 99
|
|
..message = 'from ws server');
|
|
|
|
final received =
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(received.index, equals(99));
|
|
expect(received.message, equals('from ws server'));
|
|
});
|
|
|
|
test('WS 클라이언트 disconnect 시 서버 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
server.connectedClients.first.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await client.close();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(completer.isCompleted, isTrue);
|
|
});
|
|
|
|
test('WS 서버 stop 시 클라이언트 disconnect 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
client.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await server.stop();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
test('WS close 후 isAlive가 false다', () async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
test('WS 서버가 모든 클라이언트에게 브로드캐스트한다', () async {
|
|
final ws = await WsProtobufClient.connect(_host, _testPortWs);
|
|
final secondClient = _TestWsClient(ws);
|
|
try {
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
final firstReceived = Completer<TestData>();
|
|
final secondReceived = Completer<TestData>();
|
|
client.addListener<TestData>((data) {
|
|
if (!firstReceived.isCompleted) firstReceived.complete(data);
|
|
});
|
|
secondClient.addListener<TestData>((data) {
|
|
if (!secondReceived.isCompleted) secondReceived.complete(data);
|
|
});
|
|
|
|
await server.broadcast(TestData()
|
|
..index = 88
|
|
..message = 'ws broadcast');
|
|
|
|
final first =
|
|
await firstReceived.future.timeout(const Duration(seconds: 2));
|
|
final second =
|
|
await secondReceived.future.timeout(const Duration(seconds: 2));
|
|
expect(first.message, equals('ws broadcast'));
|
|
expect(second.message, equals('ws broadcast'));
|
|
} finally {
|
|
await secondClient.close();
|
|
}
|
|
});
|
|
|
|
test('WS close 후 send는 무시된다', () async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
|
|
await expectLater(
|
|
client.send(TestData()
|
|
..index = 1
|
|
..message = 'ignored'),
|
|
completes,
|
|
);
|
|
expect(client.isAlive, isFalse);
|
|
});
|
|
|
|
// regression: 과거 _onMessage는 dispatch 중(_isDispatching) 도착한 WS frame을
|
|
// early return으로 버려, burst 수신 시 frame이 누락됐다(perf-regression의
|
|
// burst/Dart/ws hard-gate FAIL). drop-free drain queue로 모든 frame이
|
|
// 순서대로 수신되어야 한다.
|
|
test('WS burst 수신 시 frame을 누락하지 않고 순서대로 수신한다', () async {
|
|
const burst = 500;
|
|
final futures = <Future<void>>[];
|
|
for (var i = 0; i < burst; i++) {
|
|
futures.add(client.send(TestData()
|
|
..index = i
|
|
..message = 'burst$i'));
|
|
}
|
|
await Future.wait(futures);
|
|
|
|
final deadline = DateTime.now().add(const Duration(seconds: 5));
|
|
while (server.receivedMessages.length < burst &&
|
|
DateTime.now().isBefore(deadline)) {
|
|
await Future<void>.delayed(const Duration(milliseconds: 20));
|
|
}
|
|
|
|
expect(server.receivedMessages, hasLength(burst));
|
|
for (var i = 0; i < burst; i++) {
|
|
expect(server.receivedMessages[i].index, equals(i));
|
|
}
|
|
});
|
|
|
|
test('WS 1MB payload roundtrip', () async {
|
|
final port = await _freePort();
|
|
final reqServer = _TestWsReqServer(port);
|
|
await reqServer.start();
|
|
final ws = await WsProtobufClient.connect(_host, port);
|
|
final reqClient = _TestWsClient(ws);
|
|
try {
|
|
await Future<void>.delayed(const Duration(milliseconds: 100));
|
|
|
|
final testData = TestData()
|
|
..index = 789
|
|
..message = 'W' * (1024 * 1024); // 1MB
|
|
|
|
final response = await reqClient
|
|
.sendRequest<TestData, TestData>(testData)
|
|
.timeout(const Duration(seconds: 10));
|
|
|
|
expect(response.message.length, equals(1024 * 1024 + 6));
|
|
expect(response.index, equals(789 * 2));
|
|
expect(response.message.startsWith('echo: WW'), isTrue);
|
|
expect(response.message.endsWith('WW'), isTrue);
|
|
} finally {
|
|
await reqClient.close();
|
|
await Future<void>.delayed(const Duration(milliseconds: 100));
|
|
await reqServer.stop();
|
|
}
|
|
});
|
|
});
|
|
|
|
// ── WebSocket SSL ─────────────────────────────────────────────
|
|
|
|
group('WsProtobufServer (SSL)', () {
|
|
late _TestWsServerSsl server;
|
|
|
|
setUp(() async {
|
|
final ctx = SecurityContext()
|
|
..useCertificateChain(_certPath)
|
|
..usePrivateKey(_keyPath);
|
|
server = _TestWsServerSsl(ctx);
|
|
await server.start();
|
|
});
|
|
|
|
tearDown(() async {
|
|
await server.stop();
|
|
});
|
|
|
|
test('WSS 서버가 정상 시작된다', () {
|
|
expect(server.started, isTrue);
|
|
expect(server.isSecure, isTrue);
|
|
});
|
|
});
|
|
|
|
group('WsProtobufClient (SSL)', () {
|
|
late _TestWsServerSsl server;
|
|
late _TestWsClient client;
|
|
|
|
setUp(() async {
|
|
final serverCtx = SecurityContext()
|
|
..useCertificateChain(_certPath)
|
|
..usePrivateKey(_keyPath);
|
|
server = _TestWsServerSsl(serverCtx);
|
|
await server.start();
|
|
|
|
final clientCtx = SecurityContext()..setTrustedCertificates(_certPath);
|
|
final ws = await WsProtobufClient.connectSecure(
|
|
_host,
|
|
_testPortWss,
|
|
context: clientCtx,
|
|
onBadCertificate: _acceptTestCertificate,
|
|
);
|
|
client = _TestWsClient(ws);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('WSS 클라이언트가 서버에 연결된다', () {
|
|
expect(server.connectedClients, isNotEmpty);
|
|
});
|
|
|
|
test('WSS TestData 메시지를 서버가 수신한다', () async {
|
|
await client.send(TestData()
|
|
..index = 7
|
|
..message = 'hello over wss');
|
|
await Future.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(7));
|
|
expect(server.receivedMessages.first.message, equals('hello over wss'));
|
|
});
|
|
|
|
test('WSS disconnect 시 서버 콜백이 호출된다', () async {
|
|
final completer = Completer<void>();
|
|
server.connectedClients.first.addDisconnectListener((_) {
|
|
if (!completer.isCompleted) completer.complete();
|
|
});
|
|
|
|
await client.close();
|
|
await completer.future.timeout(const Duration(seconds: 2));
|
|
expect(completer.isCompleted, isTrue);
|
|
});
|
|
});
|
|
|
|
// ── Request-Response (TCP) ────────────────────────────────────
|
|
|
|
group('Request-Response (TCP plain)', () {
|
|
late _TestReqServer server;
|
|
late _TestClient client;
|
|
|
|
setUp(() async {
|
|
server = _TestReqServer();
|
|
await server.start();
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
client = _TestClient(socket);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('sendRequest로 서버 응답을 받는다', () async {
|
|
final response = await client
|
|
.sendRequest<TestData, TestData>(TestData()
|
|
..index = 21
|
|
..message = 'hello')
|
|
.timeout(const Duration(seconds: 2));
|
|
|
|
expect(response.index, equals(42));
|
|
expect(response.message, equals('echo: hello'));
|
|
});
|
|
|
|
test('여러 sendRequest가 각각 올바른 응답을 받는다', () async {
|
|
final futures = [
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 1
|
|
..message = 'a'),
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 2
|
|
..message = 'b'),
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 3
|
|
..message = 'c'),
|
|
];
|
|
final results =
|
|
await Future.wait(futures).timeout(const Duration(seconds: 3));
|
|
|
|
expect(results[0].index, equals(2));
|
|
expect(results[1].index, equals(4));
|
|
expect(results[2].index, equals(6));
|
|
expect(results[0].message, equals('echo: a'));
|
|
expect(results[1].message, equals('echo: b'));
|
|
expect(results[2].message, equals('echo: c'));
|
|
});
|
|
});
|
|
|
|
// ── Request-Response (WS) ─────────────────────────────────────
|
|
|
|
group('Request-Response (WS plain)', () {
|
|
late _TestWsReqServer server;
|
|
late _TestWsClient client;
|
|
|
|
setUp(() async {
|
|
server = _TestWsReqServer();
|
|
await server.start();
|
|
final ws = await WsProtobufClient.connect(_host, _testPortWs);
|
|
client = _TestWsClient(ws);
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
});
|
|
|
|
tearDown(() async {
|
|
await client.close();
|
|
await Future.delayed(const Duration(milliseconds: 100));
|
|
await server.stop();
|
|
});
|
|
|
|
test('WS sendRequest로 서버 응답을 받는다', () async {
|
|
final response = await client
|
|
.sendRequest<TestData, TestData>(TestData()
|
|
..index = 21
|
|
..message = 'hello ws')
|
|
.timeout(const Duration(seconds: 2));
|
|
|
|
expect(response.index, equals(42));
|
|
expect(response.message, equals('echo: hello ws'));
|
|
});
|
|
|
|
test('WS 여러 sendRequest가 각각 올바른 응답을 받는다', () async {
|
|
final futures = [
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 1
|
|
..message = 'x'),
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 2
|
|
..message = 'y'),
|
|
client.sendRequest<TestData, TestData>(TestData()
|
|
..index = 3
|
|
..message = 'z'),
|
|
];
|
|
final results =
|
|
await Future.wait(futures).timeout(const Duration(seconds: 3));
|
|
|
|
expect(results[0].index, equals(2));
|
|
expect(results[1].index, equals(4));
|
|
expect(results[2].index, equals(6));
|
|
});
|
|
|
|
test('TCP inbound backpressure pauses subscription on heavy load',
|
|
() async {
|
|
final server = _TestBackpressureServer();
|
|
final completer = Completer<TestData>();
|
|
server.handlerCompleter = completer;
|
|
await server.start();
|
|
|
|
final rawSocket = await ProtobufClient.connect(_host, _testPort);
|
|
final client = _TestClient(rawSocket);
|
|
|
|
final requestFuture = client.sendRequest<TestData, TestData>(
|
|
TestData()
|
|
..index = 1
|
|
..message = 'first',
|
|
);
|
|
await Future<void>.delayed(Duration.zero);
|
|
|
|
for (var i = 2; i <= 65; i++) {
|
|
await client.send(TestData()
|
|
..index = i
|
|
..message = 'heavy');
|
|
}
|
|
|
|
await client.send(TestData()
|
|
..index = 66
|
|
..message = 'overflow');
|
|
|
|
await Future<void>.delayed(const Duration(milliseconds: 100));
|
|
|
|
final srvClient = server.connectedClients.single;
|
|
expect(srvClient.isSourcePaused, isTrue);
|
|
|
|
completer.complete(TestData()
|
|
..index = 100
|
|
..message = 'done');
|
|
await requestFuture;
|
|
|
|
await Future<void>.delayed(const Duration(milliseconds: 50));
|
|
expect(srvClient.isSourcePaused, isFalse);
|
|
|
|
await client.close();
|
|
await server.stop();
|
|
});
|
|
});
|
|
|
|
group('TCP fragmentation and large payload', () {
|
|
late _TestServer server;
|
|
|
|
setUp(() async {
|
|
server = _TestServer();
|
|
await server.start();
|
|
});
|
|
|
|
tearDown(() async {
|
|
await server.stop();
|
|
});
|
|
|
|
test('TCP receives fragmented large frame without reordering or truncation',
|
|
() async {
|
|
final rawSocket = await Socket.connect(_host, _testPort);
|
|
try {
|
|
final testData = TestData()
|
|
..index = 123
|
|
..message = 'A' * 10000; // 10KB message
|
|
final packet = PacketBase()
|
|
..typeName = testData.info_.qualifiedMessageName
|
|
..data = testData.writeToBuffer();
|
|
|
|
final baseBytes = packet.writeToBuffer();
|
|
final header = Uint8List(4)
|
|
..buffer.asByteData().setInt32(0, baseBytes.length);
|
|
|
|
final fullFrame = Uint8List(header.length + baseBytes.length);
|
|
fullFrame.setRange(0, header.length, header);
|
|
fullFrame.setRange(header.length, fullFrame.length, baseBytes);
|
|
|
|
// Send the frame in small chunks of 100 bytes
|
|
const chunkSize = 100;
|
|
for (int i = 0; i < fullFrame.length; i += chunkSize) {
|
|
final end = (i + chunkSize < fullFrame.length)
|
|
? i + chunkSize
|
|
: fullFrame.length;
|
|
rawSocket.add(fullFrame.sublist(i, end));
|
|
await rawSocket.flush();
|
|
// Small delay to simulate packet fragmentation
|
|
await Future<void>.delayed(const Duration(milliseconds: 1));
|
|
}
|
|
|
|
// Wait for server to receive it
|
|
await Future<void>.delayed(const Duration(milliseconds: 200));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(123));
|
|
expect(server.receivedMessages.first.message, equals('A' * 10000));
|
|
} finally {
|
|
await rawSocket.close();
|
|
}
|
|
});
|
|
|
|
test('TCP 1MB payload roundtrip', () async {
|
|
final socket = await ProtobufClient.connect(_host, _testPort);
|
|
final client = _TestClient(socket);
|
|
try {
|
|
await Future<void>.delayed(const Duration(milliseconds: 100));
|
|
|
|
final testData = TestData()
|
|
..index = 456
|
|
..message = 'B' * (1024 * 1024); // 1MB message
|
|
|
|
await client.send(testData);
|
|
|
|
await Future<void>.delayed(const Duration(milliseconds: 300));
|
|
|
|
expect(server.receivedMessages, hasLength(1));
|
|
expect(server.receivedMessages.first.index, equals(456));
|
|
expect(
|
|
server.receivedMessages.first.message.length, equals(1024 * 1024));
|
|
} finally {
|
|
await client.close();
|
|
}
|
|
});
|
|
});
|
|
}
|
|
|
|
class _TestBackpressureServer extends ProtobufServer {
|
|
final connectedClients = <ProtobufClient>[];
|
|
Completer<TestData>? handlerCompleter;
|
|
|
|
_TestBackpressureServer()
|
|
: super(_host, _testPort, (socket) => _TestClient(socket));
|
|
|
|
@override
|
|
void onClientConnected(ProtobufClient client) {
|
|
connectedClients.add(client);
|
|
client.addRequestListener<TestData, TestData>((req) async {
|
|
return handlerCompleter!.future;
|
|
});
|
|
}
|
|
}
|
|
|
|
Future<int> _freePort() async {
|
|
final socket = await ServerSocket.bind(_host, 0);
|
|
final port = socket.port;
|
|
await socket.close();
|
|
return port;
|
|
}
|