9dfa06ffee
Docker image / Build (linux/amd64) (push) Has been cancelled
Docker image / Build (linux/arm64) (push) Has been cancelled
Docker image / Merge release multi-arch manifest (push) Has been cancelled
Docker image / Merge debug multi-arch manifest (push) Has been cancelled
Docker image / Build public push gateway (linux/amd64) (push) Has been cancelled
Docker image / Build public push gateway (linux/arm64) (push) Has been cancelled
Docker image / Publish public push gateway image (push) Has been cancelled
Sprig image / Build (linux/amd64) (push) Has been cancelled
Sprig image / Build (linux/arm64) (push) Has been cancelled
Sprig image / Merge multi-arch manifest (push) Has been cancelled
Harbor Buzz Orchestra / Python tests and lint (push) Has been cancelled
CI / Detect Changed Paths (push) Has been cancelled
CI / Rust Lint (push) Has been cancelled
CI / Unit Tests (push) Has been cancelled
CI / Desktop Core (push) Has been cancelled
CI / Desktop Smoke E2E (1) (push) Has been cancelled
CI / Desktop Smoke E2E (2) (push) Has been cancelled
CI / Desktop Smoke E2E (3) (push) Has been cancelled
CI / Desktop Smoke E2E (4) (push) Has been cancelled
CI / Desktop (push) Has been cancelled
CI / Desktop E2E Relay (push) Has been cancelled
CI / Desktop E2E Integration (1/2) (push) Has been cancelled
CI / Desktop E2E Integration (2/2) (push) Has been cancelled
CI / Desktop E2E Integration (push) Has been cancelled
CI / Backend Integration (relay e2e) (push) Has been cancelled
CI / Relay E2E (push) Has been cancelled
CI / Web (push) Has been cancelled
CI / Mobile (push) Has been cancelled
CI / Security (push) Has been cancelled
CI / Dead Token Reference Guard (push) Has been cancelled
CI / Server Cross-Compile (aarch64-unknown-linux-musl) (push) Has been cancelled
CI / Server Cross-Compile (x86_64-unknown-linux-musl) (push) Has been cancelled
CI / Windows Rust (x86_64-pc-windows-msvc) (push) Has been cancelled
CI / Desktop Build (macOS) (push) Has been cancelled
helm chart / lint + unittest + render matrix (push) Has been cancelled
helm chart / install on kind (gated) (push) Has been cancelled
helm chart / publish chart to GHCR (push) Has been cancelled
Mesh Lifecycle / Relay-Driven Mesh Lifecycle Smoke (push) Has been cancelled
Sprig / Build (aarch64-unknown-linux-musl) (push) Has been cancelled
Sprig / Build (x86_64-unknown-linux-musl) (push) Has been cancelled
Sprig / Publish rolling release (push) Has been cancelled
Sprig / Publish tagged release (push) Has been cancelled
Signed-off-by: cls_宁波本机 <908705107@qq.com>
274 lines
8.2 KiB
Dart
274 lines
8.2 KiB
Dart
import 'dart:async';
|
|
import 'dart:convert';
|
|
import 'dart:io';
|
|
|
|
import 'package:buzz/features/pairing/pairing_socket.dart';
|
|
import 'package:flutter_test/flutter_test.dart';
|
|
import 'package:web_socket_channel/web_socket_channel.dart';
|
|
|
|
const _privateKey =
|
|
'09b3065e3570a3a4054660dccd66e12774a99a904fdb0ca02dbc6c3136249506';
|
|
|
|
void main() {
|
|
group('PairingSocket', () {
|
|
test('connects when the pairing relay sends no AUTH challenge', () async {
|
|
final server = await _TestRelay.start((_) {});
|
|
addTearDown(server.close);
|
|
final socket = _socket(
|
|
server.url,
|
|
authChallengeTimeout: const Duration(milliseconds: 30),
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await socket.connect();
|
|
|
|
expect(socket.isConnected, isTrue);
|
|
});
|
|
|
|
test('answers an AUTH challenge and requires an accepted OK', () async {
|
|
final authReceived = Completer<List<dynamic>>();
|
|
final server = await _TestRelay.start((webSocket) async {
|
|
webSocket.add(jsonEncode(['AUTH', 'challenge']));
|
|
final auth =
|
|
jsonDecode(await webSocket.first as String) as List<dynamic>;
|
|
authReceived.complete(auth);
|
|
final event = auth[1] as Map<String, dynamic>;
|
|
webSocket.add(jsonEncode(['OK', event['id'], true, 'authenticated']));
|
|
});
|
|
addTearDown(server.close);
|
|
final socket = _socket(server.url);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await socket.connect();
|
|
|
|
expect(socket.isConnected, isTrue);
|
|
expect((await authReceived.future).first, 'AUTH');
|
|
});
|
|
|
|
test('fails when the pairing relay rejects AUTH', () async {
|
|
final server = await _TestRelay.start((webSocket) async {
|
|
webSocket.add(jsonEncode(['AUTH', 'challenge']));
|
|
final auth =
|
|
jsonDecode(await webSocket.first as String) as List<dynamic>;
|
|
final event = auth[1] as Map<String, dynamic>;
|
|
webSocket.add(jsonEncode(['OK', event['id'], false, 'bad auth']));
|
|
});
|
|
addTearDown(server.close);
|
|
var disconnectCount = 0;
|
|
final socket = _socket(
|
|
server.url,
|
|
onDisconnected: (_) => disconnectCount++,
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await expectLater(socket.connect(), throwsA(isA<PairingAuthException>()));
|
|
|
|
expect(socket.isConnected, isFalse);
|
|
expect(disconnectCount, 1);
|
|
});
|
|
|
|
test(
|
|
'answers a challenge after the optional AUTH wait completes',
|
|
() async {
|
|
final authReceived = Completer<void>();
|
|
final server = await _TestRelay.start((webSocket) async {
|
|
await Future<void>.delayed(const Duration(milliseconds: 80));
|
|
webSocket.add(jsonEncode(['AUTH', 'late-challenge']));
|
|
final auth =
|
|
jsonDecode(await webSocket.first as String) as List<dynamic>;
|
|
final event = auth[1] as Map<String, dynamic>;
|
|
webSocket.add(jsonEncode(['OK', event['id'], true, 'authenticated']));
|
|
authReceived.complete();
|
|
});
|
|
addTearDown(server.close);
|
|
final socket = _socket(
|
|
server.url,
|
|
authChallengeTimeout: const Duration(milliseconds: 30),
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await socket.connect();
|
|
await authReceived.future;
|
|
|
|
expect(socket.isConnected, isTrue);
|
|
},
|
|
);
|
|
|
|
test('fails when AUTH receives no OK response', () async {
|
|
final server = await _TestRelay.start((webSocket) {
|
|
webSocket.add(jsonEncode(['AUTH', 'challenge']));
|
|
});
|
|
addTearDown(server.close);
|
|
final socket = _socket(
|
|
server.url,
|
|
authResponseTimeout: const Duration(milliseconds: 100),
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await expectLater(socket.connect(), throwsA(isA<PairingAuthException>()));
|
|
|
|
expect(socket.isConnected, isFalse);
|
|
});
|
|
|
|
test('notifies once when a connected stream emits an error', () async {
|
|
var disconnectCount = 0;
|
|
final channel = _ControlledWebSocketChannel();
|
|
final socket = _socket(
|
|
'ws://unused',
|
|
onDisconnected: (_) => disconnectCount++,
|
|
channelFactory: (_) => channel,
|
|
authChallengeTimeout: Duration.zero,
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await socket.connect();
|
|
channel.emitError(Exception('stream failed'));
|
|
await Future<void>.delayed(Duration.zero);
|
|
|
|
expect(disconnectCount, 1);
|
|
});
|
|
|
|
test('notifies once when a connected stream closes', () async {
|
|
var disconnectCount = 0;
|
|
final channel = _ControlledWebSocketChannel();
|
|
final socket = _socket(
|
|
'ws://unused',
|
|
onDisconnected: (_) => disconnectCount++,
|
|
channelFactory: (_) => channel,
|
|
authChallengeTimeout: Duration.zero,
|
|
);
|
|
addTearDown(socket.disconnect);
|
|
|
|
await socket.connect();
|
|
await channel.closeStream();
|
|
|
|
expect(disconnectCount, 1);
|
|
});
|
|
|
|
test(
|
|
'does not notify when deliberately disconnected or disposed',
|
|
() async {
|
|
var disconnectCount = 0;
|
|
final disconnectChannel = _ControlledWebSocketChannel();
|
|
final disconnectingSocket = _socket(
|
|
'ws://unused',
|
|
onDisconnected: (_) => disconnectCount++,
|
|
channelFactory: (_) => disconnectChannel,
|
|
authChallengeTimeout: Duration.zero,
|
|
);
|
|
await disconnectingSocket.connect();
|
|
|
|
await disconnectingSocket.disconnect();
|
|
|
|
final disposeChannel = _ControlledWebSocketChannel();
|
|
final disposingSocket = _socket(
|
|
'ws://unused',
|
|
onDisconnected: (_) => disconnectCount++,
|
|
channelFactory: (_) => disposeChannel,
|
|
authChallengeTimeout: Duration.zero,
|
|
);
|
|
await disposingSocket.connect();
|
|
|
|
disposingSocket.dispose();
|
|
await Future<void>.delayed(Duration.zero);
|
|
|
|
expect(disconnectCount, 0);
|
|
},
|
|
);
|
|
});
|
|
}
|
|
|
|
PairingSocket _socket(
|
|
String url, {
|
|
Duration authChallengeTimeout = const Duration(milliseconds: 500),
|
|
Duration authResponseTimeout = const Duration(seconds: 10),
|
|
void Function(Object? error)? onDisconnected,
|
|
WebSocketChannel Function(Uri uri)? channelFactory,
|
|
}) => PairingSocket(
|
|
wsUrl: url,
|
|
ephemeralPrivkey: _privateKey,
|
|
onMessage: (_) {},
|
|
onDisconnected: onDisconnected ?? (_) {},
|
|
authChallengeTimeout: authChallengeTimeout,
|
|
authResponseTimeout: authResponseTimeout,
|
|
channelFactory: channelFactory ?? WebSocketChannel.connect,
|
|
);
|
|
|
|
class _ControlledWebSocketChannel implements WebSocketChannel {
|
|
final StreamController<dynamic> _streamController = StreamController();
|
|
final WebSocketSink _sink = _ControlledWebSocketSink();
|
|
|
|
void emitError(Object error) => _streamController.addError(error);
|
|
|
|
Future<void> closeStream() => _streamController.close();
|
|
|
|
@override
|
|
Future<void> get ready => Future.value();
|
|
|
|
@override
|
|
Stream<dynamic> get stream => _streamController.stream;
|
|
|
|
@override
|
|
WebSocketSink get sink => _sink;
|
|
|
|
@override
|
|
int? get closeCode => null;
|
|
|
|
@override
|
|
String? get closeReason => null;
|
|
|
|
@override
|
|
String? get protocol => null;
|
|
|
|
@override
|
|
dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation);
|
|
}
|
|
|
|
class _ControlledWebSocketSink implements WebSocketSink {
|
|
@override
|
|
void add(dynamic event) {}
|
|
|
|
@override
|
|
void addError(Object error, [StackTrace? stackTrace]) {}
|
|
|
|
@override
|
|
Future<void> addStream(Stream<dynamic> stream) async {
|
|
await stream.drain<void>();
|
|
}
|
|
|
|
@override
|
|
Future<void> close([int? closeCode, String? closeReason]) async {}
|
|
|
|
@override
|
|
Future<void> get done => Future.value();
|
|
}
|
|
|
|
class _TestRelay {
|
|
final HttpServer _server;
|
|
final List<WebSocket> _sockets = [];
|
|
|
|
_TestRelay._(this._server);
|
|
|
|
String get url => 'ws://${_server.address.host}:${_server.port}';
|
|
|
|
static Future<_TestRelay> start(
|
|
FutureOr<void> Function(WebSocket socket) onConnected,
|
|
) async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
final relay = _TestRelay._(server);
|
|
server.listen((request) async {
|
|
final socket = await WebSocketTransformer.upgrade(request);
|
|
relay._sockets.add(socket);
|
|
await onConnected(socket);
|
|
});
|
|
return relay;
|
|
}
|
|
|
|
Future<void> close() async {
|
|
for (final socket in _sockets) {
|
|
await socket.close();
|
|
}
|
|
await _server.close(force: true);
|
|
}
|
|
}
|