fix(watch-together): authenticate retained room recovery

A transport reconnect could recover the wrong room or retain obsolete host authority.

Authenticate retained membership and recover the relay-authoritative role without falling back to create or join. Release retained membership through authenticated teardown and document the coordinated relay rollout.
This commit is contained in:
edde746
2026-09-06 20:21:11 +02:00
parent 99f4fda55a
commit 178634e943
8 changed files with 901 additions and 300 deletions
+4
View File
@@ -136,6 +136,10 @@ sudo moss it plezy
- Synchronized playback with friends
- Real-time play / pause / seek sync
- Host handoff when the relay and every connected peer support safe transfers
- Automatic reconnect authenticates retained room membership; it never silently joins a reused room code.
Self-hosted relays must support `authenticatedResume` for recovery. Deploy the updated relay before updating clients.
Older relays still accept explicit create/join, but failed recovery requires joining or creating a room again.
### <img src="assets/readme_icons/integrations.svg" height="20" alt="" align="center" /> Integrations
- Discord Rich Presence[^desktop]
@@ -6,6 +6,7 @@ abstract final class RelayProtocol {
static const String create = 'create';
static const String join = 'join';
static const String resume = 'resume';
static const String broadcast = 'broadcast';
static const String sendTo = 'sendTo';
static const String ping = 'ping';
@@ -14,6 +15,7 @@ abstract final class RelayProtocol {
static const String transferHost = 'transferHost';
static const String created = 'created';
static const String joined = 'joined';
static const String resumed = 'resumed';
static const String peerJoined = 'peerJoined';
static const String peerLeft = 'peerLeft';
static const String message = 'message';
@@ -36,6 +38,7 @@ abstract final class RelayProtocol {
static const String peerNotFoundCode = 'peer_not_found';
static const String hostTransferUnavailableCode = 'host_transfer_unavailable';
static const String atomicHostTransferFeature = 'atomicHostTransfer';
static const String authenticatedResumeFeature = 'authenticatedResume';
static const String hostTransferCapability = 'hostTransfer';
static const int maxRoomSize = 8;
@@ -69,6 +69,8 @@ class WatchTogetherPeerService with KeepaliveMixin {
StreamSubscription? _channelSubscription;
Completer<void>? _setupCompleter;
String? _setupRequestType;
bool _admitted = false;
bool _initialRequestMayHaveCommitted = false;
final Set<String> _connectedPeers = {};
String? _sessionId;
String? _myPeerId;
@@ -126,7 +128,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
/// Stream of connection state changes (true = connected, false = disconnected)
Stream<bool> get onConnectionStateChanged => _connectionStateController.stream;
/// Emitted when the host has durably ended the relay room.
/// Emitted when this session can no longer be continued safely, for either role.
Stream<void> get onSessionEnded => _sessionEndedController.stream;
/// Emitted with the new host's peer ID when the relay reassigns host
@@ -140,6 +142,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
bool canTransferHostTo(String peerId) =>
!_teardownInProgress &&
_admitted &&
_channel != null &&
_isHost &&
_relayEnforcesHostTransfer &&
@@ -163,16 +166,15 @@ class WatchTogetherPeerService with KeepaliveMixin {
/// Whether this peer is the host.
///
/// Derived, never stored: the relay is the authority on host identity and
/// names it in every admission and every `hostChanged`. Before it has
/// admitted us there is no authority yet, so the role is the one we
/// announced — which is what a release of a possibly-committed setup has to
/// go by.
/// names it in every admission and every `hostChanged`. Before admission,
/// the requested role is provisional; disconnected release always resumes
/// authenticated membership before choosing a terminal operation.
bool get _isHost => _hostPeerId == null ? _announcedAsHost : _hostPeerId == _myPeerId;
bool get isHost => _isHost;
/// Whether currently connected to a session
bool get isConnected => _channel != null && _connectedPeers.isNotEmpty;
bool get isConnected => _admitted && _channel != null && _connectedPeers.isNotEmpty;
/// List of connected peer IDs
List<String> get connectedPeers => _connectedPeers.toList();
@@ -222,6 +224,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
Duration? connectTimeout,
String connectOperation = 'WatchTogether connect',
}) async {
_requireCurrentConnection(epoch);
final channel = await _connectToRelay(timeout: connectTimeout, operation: connectOperation);
if (_disposed || epoch != _connectionEpoch || _sessionId == null) {
unawaited(channel.sink.close());
@@ -239,17 +242,23 @@ class WatchTogetherPeerService with KeepaliveMixin {
final completer = Completer<void>();
_setupCompleter = completer;
_setupRequestType = type;
if (type == RelayProtocol.create || type == RelayProtocol.join) {
final admission = type == RelayProtocol.create || type == RelayProtocol.join || type == RelayProtocol.resume;
if (admission) {
_admitted = false;
_clearHostTransferEligibility(resetFeature: true);
}
final reconnectToken = _reconnectToken;
if (type == RelayProtocol.create || type == RelayProtocol.join) {
// Enqueue failures cannot prove that the relay did not commit admission.
_initialRequestMayHaveCommitted = true;
}
_sendRaw({
'type': type,
'sessionId': _sessionId,
'peerId': _myPeerId,
'reconnectToken': ?reconnectToken,
'protocolVersion': _relayProtocolVersion,
if (type == RelayProtocol.create || type == RelayProtocol.join) ...{
if (admission) ...{
'syncProtocolVersion': SyncMessage.protocolVersion,
'capabilities': [RelayProtocol.hostTransferCapability],
},
@@ -301,6 +310,13 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
List<String> _acceptSetupResponse(Map<String, dynamic> msg, String type) {
final expectedType = switch (_setupRequestType) {
RelayProtocol.create => RelayProtocol.created,
RelayProtocol.join => RelayProtocol.joined,
RelayProtocol.resume => RelayProtocol.resumed,
_ => null,
};
if (_setupCompleter == null || type != expectedType) throw _invalidSetupResponse(type);
final responseSessionId = msg['sessionId'];
final hostPeerId = msg['hostPeerId'];
final reconnectToken = msg['reconnectToken'];
@@ -328,10 +344,15 @@ class WatchTogetherPeerService with KeepaliveMixin {
peers.add(peerId);
}
}
final features = msg['features'];
if (type == RelayProtocol.resumed &&
(features is! List || !features.contains(RelayProtocol.authenticatedResumeFeature))) {
throw _invalidSetupResponse(type);
}
_hostPeerId = hostPeerId;
_reconnectToken = reconnectToken;
final features = msg['features'];
_admitted = true;
_relayEnforcesHostTransfer = features is List && features.contains(RelayProtocol.atomicHostTransferFeature);
_clearHostTransferEligibility();
return peers;
@@ -347,6 +368,17 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
void _failSetup(PeerError error) {
if (_setupRequestType == RelayProtocol.resume && !_isRetryableInitialSetupError(error)) {
error = PeerError(
type: error.type,
message: t.watchTogether.errors.sessionUnavailable,
serverCode: error.serverCode,
);
if (!_initialSetupInProgress && !_teardownInProgress) {
_handleSessionEnded(error);
return;
}
}
_safeAdd(_errorController, error);
if (_setupCompleter case final completer? when !completer.isCompleted) {
_setupCompleter = null;
@@ -355,30 +387,20 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
}
bool _isExhaustedGuestReconnectRoomNotFound(String code) =>
code == RelayProtocol.roomNotFoundCode &&
!_isHost &&
!_initialSetupInProgress &&
!_teardownInProgress &&
_hostPeerId != null &&
_setupRequestType == RelayProtocol.join &&
_reconnectAttempts >= _maxReconnectAttempts;
void _handleGuestSessionEnded() {
_teardownInProgress = true;
++_connectionEpoch;
_reconnectTimer?.cancel();
_reconnectTimer = null;
stopKeepalive();
final error = PeerError(type: PeerErrorType.invalidSession, message: t.watchTogether.errors.sessionEnded);
if (_setupCompleter case final completer? when !completer.isCompleted) {
_setupCompleter = null;
_setupRequestType = null;
_safeAdd(_errorController, error);
completer.completeError(error);
void _handleSessionEnded([PeerError? error]) {
final endedError =
error ?? PeerError(type: PeerErrorType.invalidSession, message: t.watchTogether.errors.sessionEnded);
final completer = _setupCompleter;
_setupCompleter = null;
_setupRequestType = null;
if (completer != null && !completer.isCompleted) {
completer.completeError(endedError);
}
_safeAdd(_errorController, endedError);
// Clear credentials and fence the channel synchronously before notifying
// consumers, so terminal cleanup cannot recover into a replacement room.
unawaited(disconnect());
_safeAdd(_sessionEndedController, null);
_handleWebSocketClosed();
}
/// Handle an incoming server message (JSON string).
@@ -388,7 +410,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
final type = msg['type'] as String?;
switch (type) {
case RelayProtocol.created || RelayProtocol.joined:
case RelayProtocol.created || RelayProtocol.joined || RelayProtocol.resumed:
final previousHostPeerId = _hostPeerId;
late final List<String> peers;
try {
@@ -418,6 +440,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
case RelayProtocol.peerJoined:
if (!_admitted) break;
final peerId = msg['peerId'] as String;
appLogger.d('WatchTogether: Peer joined: $peerId');
_connectedPeers.add(peerId);
@@ -425,6 +448,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
_safeAdd(_connectionStateController, true);
case RelayProtocol.peerLeft:
if (!_admitted) break;
final peerId = msg['peerId'] as String;
appLogger.d('WatchTogether: Peer left: $peerId');
_connectedPeers.remove(peerId);
@@ -434,6 +458,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
case RelayProtocol.message:
if (!_admitted) break;
final payload = msg['payload'];
final serverFrom = msg['from'] as String?;
if (payload != null) {
@@ -452,12 +477,18 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
case RelayProtocol.left:
if (_setupRequestType != RelayProtocol.leave) {
_failSetup(_invalidSetupResponse(RelayProtocol.left));
break;
}
try {
_acceptTeardownResponse(msg, RelayProtocol.left);
} on PeerError catch (error) {
_failSetup(error);
break;
}
_admitted = false;
_reconnectToken = null;
if (_setupCompleter case final completer? when !completer.isCompleted) {
_setupCompleter = null;
_setupRequestType = null;
@@ -474,18 +505,23 @@ class WatchTogetherPeerService with KeepaliveMixin {
final expectedTeardown =
_setupRequestType == RelayProtocol.endSession || _setupRequestType == RelayProtocol.leave;
if (expectedTeardown) {
_admitted = false;
_reconnectToken = null;
if (_setupCompleter case final completer? when !completer.isCompleted) {
_setupCompleter = null;
_setupRequestType = null;
completer.complete();
}
} else if (!_isHost) {
_handleGuestSessionEnded();
} else if (_admitted ||
_setupRequestType == RelayProtocol.resume ||
_setupRequestType == RelayProtocol.join) {
_handleSessionEnded();
} else {
_failSetup(_invalidSetupResponse(RelayProtocol.ended));
}
case RelayProtocol.hostChanged:
if (!_admitted) break;
final newHostPeerId = msg['hostPeerId'];
if (msg['sessionId'] != _sessionId ||
newHostPeerId is! String ||
@@ -500,6 +536,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
_safeAdd(_hostChangedController, newHostPeerId);
case RelayProtocol.hostTransferEligibility:
if (!_admitted) break;
final targets = msg['hostTransferTargets'];
if (!_relayEnforcesHostTransfer ||
msg['sessionId'] != _sessionId ||
@@ -517,19 +554,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
case RelayProtocol.error:
final code = msg['code'] as String? ?? 'unknown';
final message = msg['message'] as String? ?? t.common.unknown;
if (_isExhaustedGuestReconnectRoomNotFound(code)) {
appLogger.d('WatchTogether: Room gone after reconnect retries; guest session ended');
_handleGuestSessionEnded();
break;
}
appLogger.e('WatchTogether: Server error: $code - $message');
final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message', serverCode: code);
_safeAdd(_errorController, error);
if (_setupCompleter case final completer? when !completer.isCompleted) {
_setupCompleter = null;
_setupRequestType = null;
completer.completeError(error);
}
_failSetup(error);
case RelayProtocol.pong:
// Handled by resetPongTimer() already
@@ -571,8 +598,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
/// Handle the WebSocket being closed unexpectedly — attempt reconnection.
void _handleWebSocketClosed() {
final channel = _channel;
final shouldReconnect = !_initialSetupInProgress && !_teardownInProgress;
final shouldReconnect = !_initialSetupInProgress && !_teardownInProgress && _reconnectToken != null;
if (shouldReconnect) ++_connectionEpoch;
_admitted = false;
stopKeepalive();
unawaited(_channelSubscription?.cancel());
_channelSubscription = null;
@@ -591,17 +619,13 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
}
/// Attempt to reconnect to the relay and re-join/re-create the room.
/// Recover only the authenticated membership retained by the relay.
void _attemptReconnect(int epoch) {
if (_disposed || epoch != _connectionEpoch || _sessionId == null) return;
if (_reconnectAttempts >= _maxReconnectAttempts) {
appLogger.e('WatchTogether: Max reconnect attempts reached');
_safeAdd(
_errorController,
const PeerError(
type: PeerErrorType.connectionFailed,
message: 'Lost connection to relay after multiple reconnect attempts',
),
_handleSessionEnded(
PeerError(type: PeerErrorType.connectionFailed, message: t.watchTogether.errors.sessionUnavailable),
);
return;
}
@@ -614,24 +638,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
_reconnectTimer = Timer(delay, () async {
if (_disposed || epoch != _connectionEpoch || _sessionId == null) return;
try {
final completer = await _connectAndAnnounce(RelayProtocol.join, epoch);
final completer = await _connectAndAnnounce(RelayProtocol.resume, epoch);
if (_disposed || epoch != _connectionEpoch) return;
try {
await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect');
} on PeerError catch (e) {
if (_disposed || epoch != _connectionEpoch) return;
if (_isHost && e.serverCode == RelayProtocol.roomNotFoundCode) {
appLogger.d('WatchTogether: Room gone, re-creating as host');
final createCompleter = _announce(RelayProtocol.create);
await createCompleter.future.namedTimeout(
const Duration(seconds: 10),
operation: 'WatchTogether reconnect create',
);
} else {
rethrow;
}
}
await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect');
await debugReconnectSetupSucceededBarrier?.call();
@@ -646,6 +655,10 @@ class WatchTogetherPeerService with KeepaliveMixin {
} catch (e) {
if (_disposed || epoch != _connectionEpoch) return;
appLogger.e('WatchTogether: Reconnect failed', error: e);
if (!_isRetryableInitialSetupError(e)) {
_handleSessionEnded(e is PeerError ? e : null);
return;
}
_handleWebSocketClosed();
}
});
@@ -659,12 +672,23 @@ class WatchTogetherPeerService with KeepaliveMixin {
error.type == PeerErrorType.timeout)) ||
error is! PeerError;
void _requireCurrentConnection(int epoch) {
if (_disposed || epoch != _connectionEpoch || _sessionId == null) {
throw StateError('Watch Together connection attempt became stale');
}
}
Future<void> _resetTransportForInitialRetry() async {
stopKeepalive();
final subscription = _channelSubscription;
final channel = _channel;
_channelSubscription = null;
_channel = null;
_admitted = false;
final completer = _setupCompleter;
if (completer != null && !completer.isCompleted) {
completer.completeError(StateError('Watch Together connection cancelled'));
}
_setupCompleter = null;
_setupRequestType = null;
_clearHostTransferEligibility(resetFeature: true);
@@ -681,12 +705,16 @@ class WatchTogetherPeerService with KeepaliveMixin {
try {
for (var attempt = 0; attempt < _maxReconnectAttempts; attempt++) {
try {
final completer = await _connectAndAnnounce(type, epoch);
_requireCurrentConnection(epoch);
final requestType = _initialRequestMayHaveCommitted ? RelayProtocol.resume : type;
final completer = await _connectAndAnnounce(requestType, epoch);
await completer.future.timeout(debugInitialSetupTimeout, onTimeout: () => throw timeoutError);
_requireCurrentConnection(epoch);
return;
} catch (error) {
if (_disposed || epoch != _connectionEpoch) rethrow;
await _resetTransportForInitialRetry();
_requireCurrentConnection(epoch);
if (!_isRetryableInitialSetupError(error) || attempt + 1 >= _maxReconnectAttempts) {
rethrow;
}
@@ -694,12 +722,14 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
}
} finally {
_initialSetupInProgress = false;
if (epoch == _connectionEpoch) _initialSetupInProgress = false;
}
}
Future<void> _bestEffortReleaseFailedSetup(Object setupError) async {
if (!_isRetryableInitialSetupError(setupError)) return;
Future<void> _bestEffortReleaseFailedSetup(Object setupError, int epoch) async {
if (epoch != _connectionEpoch || !_initialRequestMayHaveCommitted || !_isRetryableInitialSetupError(setupError)) {
return;
}
try {
await releaseSession();
} catch (releaseError) {
@@ -730,6 +760,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
_reconnectToken = _mintReconnectToken();
_reconnectAttempts = 0;
final epoch = ++_connectionEpoch;
final peerId = _myPeerId;
try {
await _performInitialSetup(
@@ -742,8 +773,13 @@ class WatchTogetherPeerService with KeepaliveMixin {
return _sessionId!;
} catch (e) {
appLogger.e('WatchTogether: Failed to create session', error: e);
await _bestEffortReleaseFailedSetup(e);
await disconnect();
if (epoch == _connectionEpoch) {
await _bestEffortReleaseFailedSetup(e, epoch);
// Release owns a new epoch; it must not clear a newer explicit entry.
if (_sessionId == resolvedSessionId && _myPeerId == peerId) {
await disconnect();
}
}
rethrow;
}
}
@@ -768,6 +804,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
_reconnectToken = _mintReconnectToken();
_reconnectAttempts = 0;
final epoch = ++_connectionEpoch;
final peerId = _myPeerId;
try {
await _performInitialSetup(
@@ -779,14 +816,19 @@ class WatchTogetherPeerService with KeepaliveMixin {
appLogger.d('WatchTogether: Joined session: $_sessionId');
} catch (e) {
appLogger.e('WatchTogether: Failed to join session', error: e);
await _bestEffortReleaseFailedSetup(e);
await disconnect();
if (epoch == _connectionEpoch) {
await _bestEffortReleaseFailedSetup(e, epoch);
if (_sessionId == resolvedSessionId && _myPeerId == peerId) {
await disconnect();
}
}
rethrow;
}
}
/// Broadcast a message to all connected peers
void broadcast(SyncMessage message) {
if (!_admitted) return;
final payload = message.toJson();
_sendRaw({'type': RelayProtocol.broadcast, 'payload': payload});
}
@@ -796,6 +838,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
if (!RelayProtocol.isValidPeerId(peerId)) {
throw ArgumentError.value(peerId, 'peerId', 'Must be 1${RelayProtocol.maxPeerIdLength} letters, digits, _ or -');
}
if (!_admitted) return;
final payload = message.toJson();
_sendRaw({'type': RelayProtocol.sendTo, 'to': peerId, 'payload': payload});
}
@@ -846,13 +889,16 @@ class WatchTogetherPeerService with KeepaliveMixin {
await _resetTransportForInitialRetry();
}
for (var attempt = 0; attempt < _maxReconnectAttempts; attempt++) {
_requireCurrentConnection(epoch);
// Which request a rejection below refers to. Re-admission failures
// are terminal answers about our identity; release failures are not.
var releaseRequested = false;
try {
if (_channel == null) {
if (_channel == null || !_admitted) {
if (_channel != null) await _resetTransportForInitialRetry();
_requireCurrentConnection(epoch);
final reconnectCompleter = await _connectAndAnnounce(
RelayProtocol.join,
RelayProtocol.resume,
epoch,
connectTimeout: debugReleaseConnectTimeout,
connectOperation: 'WatchTogether release reconnect',
@@ -863,6 +909,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
);
}
_requireCurrentConnection(epoch);
releaseRequested = true;
final releaseCompleter = _announce(_isHost ? RelayProtocol.endSession : RelayProtocol.leave);
await releaseCompleter.future.namedTimeout(
@@ -871,27 +918,38 @@ class WatchTogetherPeerService with KeepaliveMixin {
);
return;
} catch (error) {
_requireCurrentConnection(epoch);
if (error is PeerError) {
if (_isRoomAbsentCode(error.serverCode)) return;
if (_isRoomAbsentCode(error.serverCode)) {
_admitted = false;
_reconnectToken = null;
return;
}
if (error.serverCode == RelayProtocol.peerIdUnavailableCode) {
// From re-admission: the token no longer names an identity
// here (a processed leave whose ACK was lost, an expired
// reservation). Nothing is left to release.
if (!releaseRequested) return;
if (!releaseRequested) {
_admitted = false;
_reconnectToken = null;
return;
}
// From the release itself: the relay refused the operation we
// chose for the role we believed we had. A transfer that landed
// while we were tearing down (a guest promoted to host, a host
// demoted) makes exactly this rejection, and a role-mismatch
// rejection is not a release. Reauthenticate through the
// token: the `joined` admission names the current host, and
// token: the `resumed` admission names the current host, and
// the next pass sends the operation that role requires.
appLogger.d('WatchTogether: Release rejected for role host=$_isHost; re-authenticating ownership');
await _resetTransportForInitialRetry();
_requireCurrentConnection(epoch);
if (attempt + 1 >= _maxReconnectAttempts) rethrow;
continue;
}
}
await _resetTransportForInitialRetry();
_requireCurrentConnection(epoch);
if (!_isRetryableInitialSetupError(error) || attempt + 1 >= _maxReconnectAttempts) {
rethrow;
}
@@ -899,7 +957,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
}
}
} finally {
_teardownInProgress = false;
if (epoch == _connectionEpoch) _teardownInProgress = false;
}
}
@@ -907,7 +965,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
/// intentional exits call [releaseSession] before this cleanup step.
Future<void> disconnect() async {
appLogger.d('WatchTogether: Disconnecting...');
++_connectionEpoch;
final epoch = ++_connectionEpoch;
_reconnectTimer?.cancel();
_reconnectTimer = null;
stopKeepalive();
@@ -919,6 +977,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
final setupCompleter = _setupCompleter;
_setupCompleter = null;
_setupRequestType = null;
_admitted = false;
_initialSetupInProgress = false;
_initialRequestMayHaveCommitted = false;
if (setupCompleter != null && !setupCompleter.isCompleted) {
setupCompleter.completeError(StateError('Watch Together connection cancelled'));
}
@@ -938,7 +999,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
} catch (e) {
appLogger.d('WatchTogether: channel close ignored', error: e);
}
_safeAdd(_connectionStateController, false);
if (epoch == _connectionEpoch) _safeAdd(_connectionStateController, false);
}
/// Dispose all resources.
+4 -1
View File
@@ -4,6 +4,7 @@
"clientMessageTypes": {
"create": "create",
"join": "join",
"resume": "resume",
"broadcast": "broadcast",
"sendTo": "sendTo",
"ping": "ping",
@@ -14,6 +15,7 @@
"serverMessageTypes": {
"created": "created",
"joined": "joined",
"resumed": "resumed",
"peerJoined": "peerJoined",
"peerLeft": "peerLeft",
"message": "message",
@@ -39,7 +41,8 @@
"hostTransferUnavailable": "host_transfer_unavailable"
},
"features": {
"atomicHostTransfer": "atomicHostTransfer"
"atomicHostTransfer": "atomicHostTransfer",
"authenticatedResume": "authenticatedResume"
},
"capabilities": {
"hostTransfer": "hostTransfer"
+23 -11
View File
@@ -1936,7 +1936,7 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
ReconnectToken: reconnectToken,
ProtocolVersion: relayProtocolVersion,
Peers: existingPeers,
Features: []string{relayFeatureAtomicHostTransfer},
Features: []string{relayFeatureAtomicHostTransfer, relayFeatureAuthenticatedResume},
})
if hostWasAbsent {
existing.broadcastExceptLocked(msg.PeerID, serverMsg{
@@ -2005,17 +2005,19 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
HostPeerID: msg.PeerID,
ReconnectToken: reconnectToken,
ProtocolVersion: msg.ProtocolVersion,
Features: []string{relayFeatureAtomicHostTransfer},
Features: []string{relayFeatureAtomicHostTransfer, relayFeatureAuthenticatedResume},
})
room.publishHostTransferEligibilityLocked()
room.mu.Unlock()
case relayTypeJoin:
case relayTypeJoin, relayTypeResume:
resumeOnly := msg.Type == relayTypeResume
if !validRelayID(msg.SessionID, maxSessionIDLength) || !validRelayID(msg.PeerID, maxPeerIDLength) {
client.sendJSON(serverMsg{Type: relayTypeError, Code: relayErrorInvalidMessage, Message: "Invalid sessionId or peerId"})
continue
}
if !supportedRelayProtocolVersion(msg.ProtocolVersion) {
if !supportedRelayProtocolVersion(msg.ProtocolVersion) ||
(resumeOnly && msg.ProtocolVersion != relayProtocolVersion) {
client.sendJSON(serverMsg{
Type: relayTypeError,
Code: relayErrorProtocolMismatch,
@@ -2028,10 +2030,15 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
continue
}
newToken, newVerifier, err := mintReconnectToken()
if err != nil {
client.sendJSON(serverMsg{Type: relayTypeError, Code: relayErrorInvalidMessage, Message: "Unable to join room"})
continue
var newToken string
var newVerifier reconnectVerifier
if msg.ProtocolVersion == legacyRelayProtocolVersion {
var err error
newToken, newVerifier, err = mintReconnectToken()
if err != nil {
client.sendJSON(serverMsg{Type: relayTypeError, Code: relayErrorInvalidMessage, Message: "Unable to join room"})
continue
}
}
presentedVerifier, tokenValid := reconnectVerifierFromToken(msg.ReconnectToken)
@@ -2093,7 +2100,8 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
case occupied:
authorized = false
default:
authorized = tokenValid
// Only explicit initial admission may allocate an identity.
authorized = !resumeOnly && tokenValid
responseToken = msg.ReconnectToken
responseVerifier = presentedVerifier
}
@@ -2181,14 +2189,18 @@ func (s *Server) handleWS(w http.ResponseWriter, r *http.Request) {
if s.beforeJoinRoomAck != nil {
s.beforeJoinRoomAck()
}
responseType := relayTypeJoined
if resumeOnly {
responseType = relayTypeResumed
}
client.sendJSON(serverMsg{
Type: relayTypeJoined,
Type: responseType,
SessionID: msg.SessionID,
HostPeerID: room.HostPeerID,
ReconnectToken: responseToken,
ProtocolVersion: roomProtocolVersion,
Peers: existingPeers,
Features: []string{relayFeatureAtomicHostTransfer},
Features: []string{relayFeatureAtomicHostTransfer, relayFeatureAuthenticatedResume},
})
room.broadcastExceptLocked(msg.PeerID, serverMsg{Type: relayTypePeerJoined, PeerID: msg.PeerID})
room.publishHostTransferEligibilityLocked()
+271 -50
View File
@@ -2257,7 +2257,7 @@ func TestCreateNegotiatesModernProtocolWithClientKnownToken(t *testing.T) {
}
}
func TestModernCreateRetryAfterLostSetupResponseIsIdempotent(t *testing.T) {
func TestModernCreateLostSetupResponseResumesCommittedIdentity(t *testing.T) {
h := newRelayHarness(t)
hostToken, _ := mustReconnectToken(t)
create := clientMsg{
@@ -2282,8 +2282,10 @@ func TestModernCreateRetryAfterLostSetupResponseIsIdempotent(t *testing.T) {
}
retry := h.dial(t, "1.1.1.42")
retry.send(create)
created := retry.expectAuthority(relayTypeCreated, "H")
resume := create
resume.Type = relayTypeResume
retry.send(resume)
created := retry.expectAuthority(relayTypeResumed, "H")
if created.ReconnectToken != hostToken || created.ProtocolVersion != relayProtocolVersion {
t.Fatalf("retry authority changed: tokenMatch=%v protocol=%d", created.ReconnectToken == hostToken, created.ProtocolVersion)
}
@@ -3070,7 +3072,7 @@ func TestFullRoomAllowsOnlyAuthenticatedLiveReplacements(t *testing.T) {
unprovedHost := h.dial(t, "2.1.0.251")
unprovedHost.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "FULL",
PeerID: "H",
ProtocolVersion: relayProtocolVersion,
@@ -3078,7 +3080,7 @@ func TestFullRoomAllowsOnlyAuthenticatedLiveReplacements(t *testing.T) {
unprovedHost.expectError(relayErrorPeerIdUnavailable)
unprovedGuest := h.dial(t, "2.1.0.252")
unprovedGuest.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "FULL",
PeerID: "G1",
ProtocolVersion: relayProtocolVersion,
@@ -3086,7 +3088,7 @@ func TestFullRoomAllowsOnlyAuthenticatedLiveReplacements(t *testing.T) {
unprovedGuest.expectError(relayErrorPeerIdUnavailable)
wrongGuestToken, _ := mustReconnectToken(t)
unprovedGuest.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "FULL",
PeerID: "G1",
ReconnectToken: wrongGuestToken,
@@ -3096,13 +3098,13 @@ func TestFullRoomAllowsOnlyAuthenticatedLiveReplacements(t *testing.T) {
newHost := h.dial(t, "2.1.0.253")
newHost.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "FULL",
PeerID: "H",
ReconnectToken: created.ReconnectToken,
ProtocolVersion: relayProtocolVersion,
})
hostJoined := newHost.expectAuthority(relayTypeJoined, "H")
hostJoined := newHost.expectAuthority(relayTypeResumed, "H")
if len(hostJoined.Peers) != maxRoomSize-1 {
t.Fatalf("replacement host peers=%v, want %d peers", hostJoined.Peers, maxRoomSize-1)
}
@@ -3116,13 +3118,13 @@ func TestFullRoomAllowsOnlyAuthenticatedLiveReplacements(t *testing.T) {
newGuest := h.dial(t, "2.1.0.254")
newGuest.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "FULL",
PeerID: "G1",
ReconnectToken: guestTokens["G1"],
ProtocolVersion: relayProtocolVersion,
})
newGuest.expectAuthority(relayTypeJoined, "H")
newGuest.expectAuthority(relayTypeResumed, "H")
if messages, err := guests["G1"].recvUntilClosed(2 * time.Second); err != nil {
t.Fatalf("displaced guest did not close: %v (frames=%v)", err, messages)
}
@@ -3269,15 +3271,17 @@ func TestLegacySameSourceHostReconnectAndModernTokenEnforcement(t *testing.T) {
})
}
func TestJoinAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
func TestResumeAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
h := newRelayHarness(t)
_, hostVerifier := mustReconnectToken(t)
guestToken, guestVerifier := mustReconnectToken(t)
now := time.Now()
room := &Room{
SessionID: "ATOMIC_CLEANUP",
HostPeerID: "H",
hostVerifier: hostVerifier,
peerReservations: make(map[string]peerReservation),
ProtocolVersion: relayProtocolVersion,
peerReservations: map[string]peerReservation{"G1": {verifier: guestVerifier, absentSince: now}},
Peers: make(map[string]*Client),
CreatedAt: now.Add(-time.Hour),
LastActivityAt: now.Add(-emptyRoomMaxAge - time.Second),
@@ -3297,7 +3301,10 @@ func TestJoinAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
}
joiner := h.dial(t, "2.2.0.1")
joiner.send(clientMsg{Type: relayTypeJoin, SessionID: room.SessionID, PeerID: "G1"})
joiner.send(clientMsg{
Type: relayTypeResume, SessionID: room.SessionID, PeerID: "G1",
ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion,
})
<-reached
cleanupDone := make(chan struct{})
@@ -3312,7 +3319,7 @@ func TestJoinAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
}
close(release)
joiner.expectAuthority(relayTypeJoined, "H")
joiner.expectAuthority(relayTypeResumed, "H")
<-cleanupDone
h.srv.mu.RLock()
@@ -3329,7 +3336,11 @@ func TestJoinAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
}
second := h.dial(t, "2.2.0.2")
second.send(clientMsg{Type: relayTypeJoin, SessionID: room.SessionID, PeerID: "G2"})
secondToken, _ := mustReconnectToken(t)
second.send(clientMsg{
Type: relayTypeJoin, SessionID: room.SessionID, PeerID: "G2",
ReconnectToken: secondToken, ProtocolVersion: relayProtocolVersion,
})
second.expectAuthority(relayTypeJoined, "H")
joiner.expect(relayTypePeerJoined)
second.send(clientMsg{Type: relayTypeBroadcast, Payload: json.RawMessage(`{"atomic":true}`)})
@@ -3338,15 +3349,17 @@ func TestJoinAdmissionIsAtomicWithEmptyRoomCleanup(t *testing.T) {
}
}
func TestJoinAdmissionIsAtomicWithReservedRoomCreate(t *testing.T) {
func TestResumeAdmissionIsAtomicWithReservedRoomCreate(t *testing.T) {
h := newRelayHarness(t)
_, hostVerifier := mustReconnectToken(t)
guestToken, guestVerifier := mustReconnectToken(t)
now := time.Now()
room := &Room{
SessionID: "ATOMIC_CREATE",
HostPeerID: "H",
hostVerifier: hostVerifier,
peerReservations: make(map[string]peerReservation),
ProtocolVersion: relayProtocolVersion,
peerReservations: map[string]peerReservation{"G": {verifier: guestVerifier, absentSince: now}},
Peers: make(map[string]*Client),
CreatedAt: now,
LastActivityAt: now,
@@ -3366,11 +3379,18 @@ func TestJoinAdmissionIsAtomicWithReservedRoomCreate(t *testing.T) {
}
joiner := h.dial(t, "2.3.0.1")
joiner.send(clientMsg{Type: relayTypeJoin, SessionID: room.SessionID, PeerID: "G"})
joiner.send(clientMsg{
Type: relayTypeResume, SessionID: room.SessionID, PeerID: "G",
ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion,
})
<-reached
creator := h.dial(t, "2.3.0.2")
creator.send(clientMsg{Type: relayTypeCreate, SessionID: room.SessionID, PeerID: "OTHER"})
replacementToken, _ := mustReconnectToken(t)
creator.send(clientMsg{
Type: relayTypeCreate, SessionID: room.SessionID, PeerID: "OTHER",
ReconnectToken: replacementToken, ProtocolVersion: relayProtocolVersion,
})
type readResult struct {
message serverMsg
err error
@@ -3395,7 +3415,7 @@ func TestJoinAdmissionIsAtomicWithReservedRoomCreate(t *testing.T) {
}
close(release)
joiner.expectAuthority(relayTypeJoined, "H")
joiner.expectAuthority(relayTypeResumed, "H")
result := <-createResult
if result.err != nil {
t.Fatalf("read create result: %v", result.err)
@@ -3642,13 +3662,13 @@ func TestStalePeerSkipsCleanupBroadcast(t *testing.T) {
g2 := h.dial(t, "6.1.0.3")
g2.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "D2",
PeerID: "G",
ReconnectToken: guestToken,
ProtocolVersion: relayProtocolVersion,
})
g2.expectAuthority(relayTypeJoined, "H")
g2.expectAuthority(relayTypeResumed, "H")
host.expect(relayTypePeerJoined)
if messages, err := g1.recvUntilClosed(2 * time.Second); err != nil {
@@ -3670,6 +3690,195 @@ func TestStalePeerSkipsCleanupBroadcast(t *testing.T) {
}
}
func TestResumeCannotEnterReplacementRoom(t *testing.T) {
for _, replacementHost := range []string{"NEW_HOST", "H"} {
t.Run(replacementHost, func(t *testing.T) {
h := newRelayHarness(t)
host, guest, hostToken, guestToken := createModernRoomWithGuest(t, h, "REUSED", "6.5.0.1", "6.5.0.2")
if err := guest.conn.Close(); err != nil {
t.Fatal(err)
}
host.expect(relayTypePeerLeft)
host.send(clientMsg{Type: relayTypeEndSession, ReconnectToken: hostToken, ProtocolVersion: relayProtocolVersion})
// The old client may lose this terminal ACK; the committed room is gone.
host.expect(relayTypeEnded)
replacementToken, _ := mustReconnectToken(t)
replacement := h.dial(t, "6.5.0.3")
replacement.send(clientMsg{
Type: relayTypeCreate, SessionID: "REUSED", PeerID: replacementHost,
ReconnectToken: replacementToken, ProtocolVersion: relayProtocolVersion,
})
replacement.expectAuthority(relayTypeCreated, replacementHost)
for index, identity := range []struct{ peerID, token string }{{"H", hostToken}, {"G", guestToken}} {
stale := h.dial(t, fmt.Sprintf("6.5.1.%d", index+1))
stale.send(clientMsg{
Type: relayTypeResume, SessionID: "REUSED", PeerID: identity.peerID,
ReconnectToken: identity.token, ProtocolVersion: relayProtocolVersion,
})
stale.expectError(relayErrorPeerIdUnavailable)
stale.send(clientMsg{Type: relayTypeBroadcast, Payload: json.RawMessage(`{"stale":true}`)})
stale.expectError(relayErrorNotInRoom)
stale.send(clientMsg{Type: relayTypeEndSession, ReconnectToken: identity.token, ProtocolVersion: relayProtocolVersion})
stale.expectError(relayErrorNotInRoom)
}
// No peerJoined or payload may precede this barrier on the replacement.
replacement.send(clientMsg{Type: relayTypePing})
replacement.expect(relayTypePong)
h.srv.mu.RLock()
room := h.srv.rooms["REUSED"]
room.mu.RLock()
reserved, connected := len(room.peerReservations), len(room.Peers)
room.mu.RUnlock()
h.srv.mu.RUnlock()
if reserved != 0 || connected != 1 {
t.Fatalf("rejected resumes allocated membership: reserved=%d connected=%d", reserved, connected)
}
})
}
}
func TestResumeLostJoinAndLeaveAcknowledgements(t *testing.T) {
h := newRelayHarness(t)
hostToken, _ := mustReconnectToken(t)
host := h.dial(t, "6.5.2.1")
host.send(clientMsg{
Type: relayTypeCreate, SessionID: "LOST_GUEST_ACK", PeerID: "H",
ReconnectToken: hostToken, ProtocolVersion: relayProtocolVersion,
})
host.expectAuthority(relayTypeCreated, "H")
guestToken, _ := mustReconnectToken(t)
admission := clientMsg{
Type: relayTypeJoin, SessionID: "LOST_GUEST_ACK", PeerID: "G",
ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion,
}
guest := h.dial(t, "6.5.2.2")
guest.send(admission)
host.expect(relayTypePeerJoined)
// Drop the first transport without reading its committed join ACK.
if err := guest.conn.Close(); err != nil {
t.Fatal(err)
}
host.expect(relayTypePeerLeft)
resumed := h.dial(t, "6.5.2.3")
admission.Type = relayTypeResume
resumed.send(admission)
ack := resumed.expectAuthority(relayTypeResumed, "H")
if !slices.Contains(ack.Features, relayFeatureAuthenticatedResume) {
t.Fatal("resume did not advertise authenticated admission")
}
host.expect(relayTypePeerJoined)
resumed.send(clientMsg{Type: relayTypeLeave, ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion})
host.expect(relayTypePeerLeft)
// The leave committed, but the departing peer never reads its ACK.
if err := resumed.conn.Close(); err != nil {
t.Fatal(err)
}
cleanup := h.dial(t, "6.5.2.4")
cleanup.send(admission)
cleanup.expectError(relayErrorPeerIdUnavailable)
cleanup.send(clientMsg{Type: relayTypeLeave, ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion})
cleanup.expectError(relayErrorNotInRoom)
host.send(clientMsg{Type: relayTypePing})
host.expect(relayTypePong)
}
func TestResumeRequiresUnexpiredExistingMembership(t *testing.T) {
stateFile := filepath.Join(t.TempDir(), "rooms.json")
h := newRelayHarnessAt(t, t.TempDir(), stateFile)
host, guest, _, guestToken := createModernRoomWithGuest(t, h, "EXPIRING_RESUME", "6.5.3.1", "6.5.3.2")
if err := guest.conn.Close(); err != nil {
t.Fatal(err)
}
host.expect(relayTypePeerLeft)
h.srv.mu.RLock()
room := h.srv.rooms["EXPIRING_RESUME"]
h.srv.mu.RUnlock()
expireDisconnectedReservations(room, time.Now())
resume := clientMsg{
Type: relayTypeResume, SessionID: "EXPIRING_RESUME", PeerID: "G",
ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion,
}
probe := h.dial(t, "6.5.3.3")
probe.send(resume)
probe.expectError(relayErrorPeerIdUnavailable)
resume.PeerID = "NEVER_JOINED"
probe.send(resume)
probe.expectError(relayErrorPeerIdUnavailable)
resume.SessionID = "NEVER_CREATED"
probe.send(resume)
probe.expectError(relayErrorRoomNotFound)
resume.SessionID = "EXPIRING_RESUME"
resume.ProtocolVersion = legacyRelayProtocolVersion
probe.send(resume)
probe.expectError(relayErrorProtocolMismatch)
host.send(clientMsg{Type: relayTypePing})
host.expect(relayTypePong)
// Admission-time expiry is persisted; a restart cannot resurrect the guest.
if err := h.srv.snap.flushAndStop(2 * time.Second); err != nil {
t.Fatalf("flush reservation expiry: %v", err)
}
restarted := newRelayHarnessAt(t, t.TempDir(), copySnapshotForRestart(t, stateFile))
afterRestart := restarted.dial(t, "6.5.3.4")
resume.PeerID = "G"
resume.ProtocolVersion = relayProtocolVersion
afterRestart.send(resume)
afterRestart.expectError(relayErrorPeerIdUnavailable)
// An explicit new join is still allowed to claim the freed identity.
freshToken, _ := mustReconnectToken(t)
resume.Type = relayTypeJoin
resume.ReconnectToken = freshToken
probe.send(resume)
probe.expectAuthority(relayTypeJoined, "H")
host.expect(relayTypePeerJoined)
}
func TestOfflineGuestResumesAfterHostTransfer(t *testing.T) {
h := newRelayHarness(t)
host, promoted, _, _ := createTransferRoomWithGuest(t, h, "OFFLINE_TRANSFER", "6.5.4.1", "6.5.4.2")
token, _ := mustReconnectToken(t)
offline := h.dial(t, "6.5.4.3")
admission := currentAdmission(clientMsg{
Type: relayTypeJoin, SessionID: "OFFLINE_TRANSFER", PeerID: "OFFLINE",
ReconnectToken: token, ProtocolVersion: relayProtocolVersion,
})
offline.send(admission)
offline.expectAuthority(relayTypeJoined, "H")
offline.expectEligibility("OFFLINE_TRANSFER", "H", "G", "OFFLINE")
host.expect(relayTypePeerJoined)
host.expectEligibility("OFFLINE_TRANSFER", "H", "G", "OFFLINE")
promoted.expect(relayTypePeerJoined)
promoted.expectEligibility("OFFLINE_TRANSFER", "H", "G", "OFFLINE")
if err := offline.conn.Close(); err != nil {
t.Fatal(err)
}
host.expect(relayTypePeerLeft)
host.expectEligibility("OFFLINE_TRANSFER", "H", "G")
promoted.expect(relayTypePeerLeft)
promoted.expectEligibility("OFFLINE_TRANSFER", "H", "G")
host.send(clientMsg{Type: relayTypeTransferHost, To: "G", ProtocolVersion: relayProtocolVersion})
host.expect(relayTypeHostChanged)
host.expectEligibility("OFFLINE_TRANSFER", "G", "H")
promoted.expect(relayTypeHostChanged)
promoted.expectEligibility("OFFLINE_TRANSFER", "G", "H")
admission.Type = relayTypeResume
resumed := h.dial(t, "6.5.4.4")
resumed.send(admission)
resumed.expectAuthority(relayTypeResumed, "G")
resumed.expectEligibility("OFFLINE_TRANSFER", "G", "H", "OFFLINE")
host.expect(relayTypePeerJoined)
host.expectEligibility("OFFLINE_TRANSFER", "G", "H", "OFFLINE")
promoted.expect(relayTypePeerJoined)
promoted.expectEligibility("OFFLINE_TRANSFER", "G", "H", "OFFLINE")
resumed.send(clientMsg{Type: relayTypeSendTo, To: "G", Payload: json.RawMessage(`{"resumed":true}`)})
if received := promoted.expect(relayTypeMessage); received.From != "OFFLINE" {
t.Fatalf("resumed sender=%q, want OFFLINE", received.From)
}
}
func TestDisconnectedModernGuestIdentityRejectsTheftAndAcceptsRightfulReconnect(t *testing.T) {
h := newRelayHarness(t)
hostToken, _ := mustReconnectToken(t)
@@ -3718,7 +3927,7 @@ func TestDisconnectedModernGuestIdentityRejectsTheftAndAcceptsRightfulReconnect(
thiefToken, _ := mustReconnectToken(t)
thief := h.dial(t, "6.1.0.12")
thief.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "GUEST_RECONNECT",
PeerID: "G",
ReconnectToken: thiefToken,
@@ -3730,13 +3939,13 @@ func TestDisconnectedModernGuestIdentityRejectsTheftAndAcceptsRightfulReconnect(
rightful := h.dial(t, "6.1.0.13")
rightful.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "GUEST_RECONNECT",
PeerID: "G",
ReconnectToken: guestToken,
ProtocolVersion: relayProtocolVersion,
})
joined := rightful.expectAuthority(relayTypeJoined, "H")
joined := rightful.expectAuthority(relayTypeResumed, "H")
if joined.ReconnectToken != guestToken {
t.Fatal("rightful reconnect rotated the retained guest token")
}
@@ -3921,7 +4130,7 @@ func TestAdmissionPrunePersistsWhenJoinRejected(t *testing.T) {
wrongHostToken, _ := mustReconnectToken(t)
rejected := h.dial(t, "6.2.2.3")
rejected.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "REJECTED_AFTER_PRUNE",
PeerID: "H",
ReconnectToken: wrongHostToken,
@@ -3986,7 +4195,7 @@ func TestGuestReservationSnapshotV4Migration(t *testing.T) {
statePath := filepath.Join(t.TempDir(), "rooms.json")
now := time.Now().UTC()
_, hostVerifier := mustReconnectToken(t)
_, guestVerifier := mustReconnectToken(t)
guestToken, guestVerifier := mustReconnectToken(t)
legacy := stateSnapshot{
Version: 3,
SavedAt: now,
@@ -4039,6 +4248,12 @@ func TestGuestReservationSnapshotV4Migration(t *testing.T) {
if runtimeReservation.absentSince.IsZero() {
t.Fatal("legacy reservation was restored as connected")
}
guest := h.dial(t, "6.2.4.1")
guest.send(clientMsg{
Type: relayTypeResume, SessionID: "V3_MIGRATION", PeerID: "G",
ReconnectToken: guestToken, ProtocolVersion: relayProtocolVersion,
})
guest.expectAuthority(relayTypeResumed, "H")
}
func TestSnapshotV2LoadsAndRewritesV4(t *testing.T) {
@@ -4081,6 +4296,12 @@ func TestSnapshotV2LoadsAndRewritesV4(t *testing.T) {
if snapshot.Version != snapshotFormatVersion {
t.Fatalf("v2 rewrite version=%d, want %d", snapshot.Version, snapshotFormatVersion)
}
legacyHost := h.dial(t, "6.2.4.2")
legacyHost.send(clientMsg{
Type: relayTypeResume, SessionID: "V2_MIGRATION", PeerID: "H",
ProtocolVersion: legacyRelayProtocolVersion,
})
legacyHost.expectError(relayErrorProtocolMismatch)
}
func TestSnapshotV4RetainsGuestAbsenceAcrossRestart(t *testing.T) {
@@ -4722,7 +4943,7 @@ func TestLeavePendingReservationRejectsReplacement(t *testing.T) {
matching := h.dial(t, "6.3.1.3")
matching.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "PENDING_RELEASE",
PeerID: "G",
ReconnectToken: guestToken,
@@ -4732,7 +4953,7 @@ func TestLeavePendingReservationRejectsReplacement(t *testing.T) {
wrongToken, _ := mustReconnectToken(t)
wrong := h.dial(t, "6.3.1.4")
wrong.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "PENDING_RELEASE",
PeerID: "G",
ReconnectToken: wrongToken,
@@ -5587,9 +5808,9 @@ func TestTransferHostAdmissionDoesNotReuseConnectionCapability(t *testing.T) {
typ string
}{
{"host create retry", "H", relayTypeCreate},
{"host join", "H", relayTypeJoin},
{"target join", "G", relayTypeJoin},
{"bystander join", "B", relayTypeJoin},
{"host resume", "H", relayTypeResume},
{"target resume", "G", relayTypeResume},
{"bystander resume", "B", relayTypeResume},
} {
for _, disconnected := range []bool{false, true} {
t.Run(fmt.Sprintf("%s/disconnected=%v", admission.name, disconnected), func(t *testing.T) {
@@ -5637,7 +5858,7 @@ func TestTransferHostAdmissionDoesNotReuseConnectionCapability(t *testing.T) {
Type: admission.typ, SessionID: "XFER_REPLACE", PeerID: admission.peerID,
ReconnectToken: token, ProtocolVersion: relayProtocolVersion,
})
ackType := relayTypeJoined
ackType := relayTypeResumed
if admission.typ == relayTypeCreate {
ackType = relayTypeCreated
}
@@ -5645,7 +5866,7 @@ func TestTransferHostAdmissionDoesNotReuseConnectionCapability(t *testing.T) {
t.Fatalf("replacement changed reconnect identity: %+v", ack)
}
for _, observer := range observers {
if admission.typ == relayTypeJoin || disconnected {
if admission.typ == relayTypeResume || disconnected {
if joined := observer.expect(relayTypePeerJoined); joined.PeerID != admission.peerID {
t.Fatalf("wrong replacement identity: %+v", joined)
}
@@ -6030,7 +6251,7 @@ func TestTransferHostPreservesReconnectAuthority(t *testing.T) {
guest.expect(relayTypeHostChanged)
guest.expectEligibility("XFER_RECONNECT", "G", "H")
// The old host drops and rejoins with its original token as a guest.
// The old host resumes with its original token as a guest.
if err := host.conn.Close(); err != nil {
t.Fatalf("close old host: %v", err)
}
@@ -6041,13 +6262,13 @@ func TestTransferHostPreservesReconnectAuthority(t *testing.T) {
guest.expectEligibility("XFER_RECONNECT", "G")
oldHost := h.dial(t, "6.3.4.3")
oldHost.send(currentAdmission(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "XFER_RECONNECT",
PeerID: "H",
ReconnectToken: hostToken,
ProtocolVersion: relayProtocolVersion,
}))
rejoinedGuest := oldHost.expectAuthority(relayTypeJoined, "G")
rejoinedGuest := oldHost.expectAuthority(relayTypeResumed, "G")
if rejoinedGuest.ReconnectToken != hostToken {
t.Fatalf("old host reconnect token changed: %+v", rejoinedGuest)
}
@@ -6055,7 +6276,7 @@ func TestTransferHostPreservesReconnectAuthority(t *testing.T) {
guest.expect(relayTypePeerJoined)
guest.expectEligibility("XFER_RECONNECT", "G", "H")
// The new host drops and rejoins with its original token as the host.
// The new host resumes with its original token as the host.
if err := guest.conn.Close(); err != nil {
t.Fatalf("close new host: %v", err)
}
@@ -6066,13 +6287,13 @@ func TestTransferHostPreservesReconnectAuthority(t *testing.T) {
oldHost.expectEligibility("XFER_RECONNECT", "G")
newHost := h.dial(t, "6.3.4.4")
newHost.send(currentAdmission(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "XFER_RECONNECT",
PeerID: "G",
ReconnectToken: guestToken,
ProtocolVersion: relayProtocolVersion,
}))
rejoinedHost := newHost.expectAuthority(relayTypeJoined, "G")
rejoinedHost := newHost.expectAuthority(relayTypeResumed, "G")
if rejoinedHost.ReconnectToken != guestToken {
t.Fatalf("new host reconnect token changed: %+v", rejoinedHost)
}
@@ -6127,22 +6348,22 @@ func TestTransferHostSurvivesSnapshotRestart(t *testing.T) {
newHost := hB.dial(t, "6.3.5.3")
newHost.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "XFER_RESTART",
PeerID: "G",
ReconnectToken: guestToken,
ProtocolVersion: relayProtocolVersion,
})
newHost.expectAuthority(relayTypeJoined, "G")
newHost.expectAuthority(relayTypeResumed, "G")
oldHost := hB.dial(t, "6.3.5.4")
oldHost.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "XFER_RESTART",
PeerID: "H",
ReconnectToken: hostToken,
ProtocolVersion: relayProtocolVersion,
})
oldHost.expectAuthority(relayTypeJoined, "G")
oldHost.expectAuthority(relayTypeResumed, "G")
newHost.expect(relayTypePeerJoined)
// Persisted identity does not imply capability on these fresh connections.
newHost.send(clientMsg{Type: relayTypeTransferHost, To: "H", ProtocolVersion: relayProtocolVersion})
@@ -8745,7 +8966,7 @@ func TestSnapshotV4RetainsModernHostAndGuestReservationsAcrossRestart(t *testing
wrongHostToken, _ := mustReconnectToken(t)
hostThief := hB.dial(t, "8.0.1.3")
hostThief.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "V3_RESTART",
PeerID: "H",
ReconnectToken: wrongHostToken,
@@ -8755,13 +8976,13 @@ func TestSnapshotV4RetainsModernHostAndGuestReservationsAcrossRestart(t *testing
restartedHost := hB.dial(t, "8.0.1.4")
restartedHost.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "V3_RESTART",
PeerID: "H",
ReconnectToken: hostToken,
ProtocolVersion: relayProtocolVersion,
})
hostJoined := restartedHost.expectAuthority(relayTypeJoined, "H")
hostJoined := restartedHost.expectAuthority(relayTypeResumed, "H")
if hostJoined.ReconnectToken != hostToken || hostJoined.ProtocolVersion != relayProtocolVersion {
t.Fatalf("restored host authority changed: %+v", hostJoined)
}
@@ -8769,7 +8990,7 @@ func TestSnapshotV4RetainsModernHostAndGuestReservationsAcrossRestart(t *testing
wrongGuestToken, _ := mustReconnectToken(t)
guestThief := hB.dial(t, "8.0.1.5")
guestThief.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "V3_RESTART",
PeerID: "G",
ReconnectToken: wrongGuestToken,
@@ -8779,13 +9000,13 @@ func TestSnapshotV4RetainsModernHostAndGuestReservationsAcrossRestart(t *testing
restartedGuest := hB.dial(t, "8.0.1.6")
restartedGuest.send(clientMsg{
Type: relayTypeJoin,
Type: relayTypeResume,
SessionID: "V3_RESTART",
PeerID: "G",
ReconnectToken: guestToken,
ProtocolVersion: relayProtocolVersion,
})
guestJoined := restartedGuest.expectAuthority(relayTypeJoined, "H")
guestJoined := restartedGuest.expectAuthority(relayTypeResumed, "H")
if guestJoined.ReconnectToken != guestToken || guestJoined.ProtocolVersion != relayProtocolVersion {
t.Fatalf("restored guest authority changed: %+v", guestJoined)
}
+3
View File
@@ -8,6 +8,7 @@ const (
relayTypeCreate = "create"
relayTypeJoin = "join"
relayTypeResume = "resume"
relayTypeBroadcast = "broadcast"
relayTypeSendTo = "sendTo"
relayTypePing = "ping"
@@ -16,6 +17,7 @@ const (
relayTypeTransferHost = "transferHost"
relayTypeCreated = "created"
relayTypeJoined = "joined"
relayTypeResumed = "resumed"
relayTypePeerJoined = "peerJoined"
relayTypePeerLeft = "peerLeft"
relayTypeMessage = "message"
@@ -38,6 +40,7 @@ const (
relayErrorPeerNotFound = "peer_not_found"
relayErrorHostTransferUnavailable = "host_transfer_unavailable"
relayFeatureAtomicHostTransfer = "atomicHostTransfer"
relayFeatureAuthenticatedResume = "authenticatedResume"
relayCapabilityHostTransfer = "hostTransfer"
maxRoomSize = 8
@@ -5,7 +5,6 @@ import 'dart:io';
import 'package:flutter_test/flutter_test.dart';
import 'package:stream_channel/stream_channel.dart';
import 'package:web_socket_channel/web_socket_channel.dart';
import 'package:plezy/i18n/strings.g.dart';
import 'package:plezy/watch_together/services/watch_together_peer_service.dart';
import 'package:plezy/watch_together/services/watch_together_relay_endpoint.dart';
import 'package:plezy/watch_together/models/sync_message.dart';
@@ -311,9 +310,10 @@ void main() {
test('guest reconnect sends its retained capability', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'join') {
if (message['type'] == 'join' || message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
@@ -338,7 +338,7 @@ void main() {
final initialToken = relay.messages[0].single['reconnectToken'];
expect(relay.messages[1], [
{
'type': 'join',
'type': 'resume',
'sessionId': 'GUEST1',
'peerId': guestPeerId,
'reconnectToken': initialToken,
@@ -357,9 +357,10 @@ void main() {
// mirrors it and surfaces the change through the same onHostChanged path.
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (message['type'] == 'join') {
if (message['type'] == 'join' || message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': connection == 0 ? _relayHostId : 'new-host',
'reconnectToken': message['reconnectToken'],
@@ -393,10 +394,10 @@ void main() {
expect(service.hostPeerId, 'new-host');
expect(service.isHost, isFalse);
expect(errors, isEmpty);
expect(relay.messages[1].map((m) => m['type']), ['join'], reason: 'no leave: the admission stands');
expect(relay.messages[1].map((m) => m['type']), ['resume'], reason: 'authenticated continuity is required');
});
test('host reconnect proves ownership and re-creates with the retained authority', () async {
test('host membership loss is terminal and never recreates its room', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (connection == 0 && message['type'] == 'create') {
@@ -407,81 +408,44 @@ void main() {
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (connection == 1 && message['type'] == 'join') {
} else if (message['type'] == 'resume') {
relay.send(socket, {'type': 'error', 'code': 'room_not_found', 'message': 'Room not found'});
} else if (connection == 1 && message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
}
});
final service = serviceFor(relay);
final reconnected = Completer<void>();
var reconnectCallbacks = 0;
service.onReconnected = () {
reconnectCallbacks++;
reconnected.complete();
};
service.onReconnected = () => reconnectCallbacks++;
final ended = service.onSessionEnded.first;
await _withShortenedTimer(
original: const Duration(seconds: 2),
replacement: const Duration(milliseconds: 10),
body: () => service.createSession(sessionId: 'room2'),
body: () async {
await service.createSession(sessionId: 'room2');
await relay.sockets.single.close();
await ended.timeout(const Duration(seconds: 1));
await service.releaseSession();
},
);
final hostPeerId = service.myPeerId;
await relay.sockets.single.close();
await reconnected.future.timeout(const Duration(seconds: 6));
expect(reconnectCallbacks, 1);
expect(relay.sockets, hasLength(2));
final initialCreate = relay.messages[0].single;
final reconnectToken = initialCreate['reconnectToken'];
expect(initialCreate, {
'type': 'create',
'sessionId': 'ROOM2',
'peerId': hostPeerId,
'reconnectToken': matches(RegExp(r'^[A-Za-z0-9_-]{43}$')),
'protocolVersion': 2,
'syncProtocolVersion': SyncMessage.protocolVersion,
'capabilities': [RelayProtocol.hostTransferCapability],
});
expect(relay.messages[1], [
{
'type': 'join',
'sessionId': 'ROOM2',
'peerId': hostPeerId,
'reconnectToken': reconnectToken,
'protocolVersion': 2,
'syncProtocolVersion': SyncMessage.protocolVersion,
'capabilities': [RelayProtocol.hostTransferCapability],
},
{
'type': 'create',
'sessionId': 'ROOM2',
'peerId': hostPeerId,
'reconnectToken': reconnectToken,
'protocolVersion': 2,
'syncProtocolVersion': SyncMessage.protocolVersion,
'capabilities': [RelayProtocol.hostTransferCapability],
},
expect(reconnectCallbacks, 0);
expect(service.sessionId, isNull);
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()), [
['create'],
['resume'],
]);
expect(service.hostPeerId, hostPeerId);
});
test('initial create retry reuses its pre-minted identity and capability', () async {
test('lost create ACK resumes the pre-minted identity without another create', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) async {
if (message['type'] != 'create') return;
if (message['type'] != 'create' && message['type'] != 'resume') return;
if (connection == 0) {
await socket.close();
return;
}
relay.send(socket, {
'type': 'created',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
@@ -500,7 +464,7 @@ void main() {
final first = relay.messages[0].single;
final retry = relay.messages[1].single;
expect(first['reconnectToken'], matches(RegExp(r'^[A-Za-z0-9_-]{43}$')));
expect(retry, first);
expect(retry, {...first, 'type': 'resume'});
expect(service.myPeerId, first['peerId']);
expect(service.hostPeerId, first['peerId']);
});
@@ -517,11 +481,7 @@ void main() {
await expectLater(
_withRetryBackoffShortened(() => timeoutService.createSession(sessionId: 'slow1')),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.timeout)
.having((error) => error.message, 'message', t.watchTogether.errors.timedOut),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.timeout)),
);
late final _RelayServer errorRelay;
@@ -556,8 +516,7 @@ void main() {
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.serverError)
.having((error) => error.serverCode, 'serverCode', 'protocol_mismatch')
.having((error) => error.message, 'message', contains('Relay protocol version 2 is required')),
.having((error) => error.serverCode, 'serverCode', 'protocol_mismatch'),
),
);
});
@@ -579,11 +538,7 @@ void main() {
await expectLater(
service.createSession(sessionId: 'token1'),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.serverError)
.having((error) => error.message, 'message', t.watchTogether.errors.invalidRelayResponse),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.serverError)),
);
expect(service.hostPeerId, isNull);
});
@@ -597,11 +552,7 @@ void main() {
await expectLater(
oldRelayService.createSession(sessionId: 'old01'),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.serverError)
.having((error) => error.message, 'message', t.watchTogether.errors.invalidRelayResponse),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.serverError)),
);
expect(oldRelayService.hostPeerId, isNull);
@@ -619,11 +570,7 @@ void main() {
await expectLater(
malformedService.joinSession('bad01'),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.serverError)
.having((error) => error.message, 'message', t.watchTogether.errors.invalidRelayResponse),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.serverError)),
);
expect(malformedService.hostPeerId, isNull);
});
@@ -631,9 +578,10 @@ void main() {
test('exhausted create retries end a possibly committed room before clearing credentials', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (connection >= 3 && message['type'] == 'join') {
if (connection >= 3 && message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
@@ -651,18 +599,14 @@ void main() {
await expectLater(
_withRetryBackoffShortened(() => service.createSession(sessionId: 'lostc')),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.timeout)
.having((error) => error.message, 'message', t.watchTogether.errors.timedOut),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.timeout)),
);
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()).toList(), [
['create'],
['create'],
['create'],
['join', 'endSession'],
['resume'],
['resume'],
['resume', 'endSession'],
]);
final announcements = relay.messages.map((messages) => messages.first).toList();
expect(announcements.map((message) => message['peerId']).toSet(), hasLength(1));
@@ -674,9 +618,10 @@ void main() {
test('exhausted join retries leave a possibly committed guest reservation before clearing credentials', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (connection >= 3 && message['type'] == 'join') {
if (connection >= 3 && message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
@@ -696,18 +641,14 @@ void main() {
await expectLater(
_withRetryBackoffShortened(() => service.joinSession('lostj')),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.timeout)
.having((error) => error.message, 'message', isNotEmpty),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.timeout)),
);
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()).toList(), [
['join'],
['join'],
['join'],
['join', 'leave'],
['resume'],
['resume'],
['resume', 'leave'],
]);
final announcements = relay.messages.map((messages) => messages.first).toList();
expect(announcements.map((message) => message['peerId']).toSet(), hasLength(1));
@@ -728,11 +669,7 @@ void main() {
await expectLater(
service.joinSession('ended2'),
throwsA(
isA<PeerError>()
.having((error) => error.type, 'type', PeerErrorType.invalidSession)
.having((error) => error.message, 'message', t.watchTogether.errors.sessionEnded),
),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.invalidSession)),
);
await ended.timeout(const Duration(seconds: 1));
@@ -744,7 +681,7 @@ void main() {
test('guest reconnect treats ended as terminal and cancels further reconnects', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (message['type'] != 'join') return;
if (message['type'] != 'join' && message['type'] != 'resume') return;
if (connection == 0) {
relay.send(socket, {
'type': 'joined',
@@ -779,10 +716,10 @@ void main() {
expect(relay.sockets, hasLength(2));
});
test('guest reconnect converges to session ended after room-not-found retries are exhausted', () async {
test('guest membership loss is terminal without retrying initial admission', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (message['type'] != 'join') return;
if (message['type'] != 'join' && message['type'] != 'resume') return;
if (connection == 0) {
relay.send(socket, {
'type': 'joined',
@@ -828,10 +765,10 @@ void main() {
expect(reconnectCallbacks, 0);
expect(sessionEndedEvents, 1);
expect(service.isConnected, isFalse);
expect(relay.sockets, hasLength(4));
expect(relay.sockets, hasLength(2));
expect(
relay.messages.skip(1).map((messages) => messages.map((message) => message['type']).toList()),
everyElement(['join']),
everyElement(['resume']),
);
});
@@ -916,7 +853,7 @@ void main() {
'protocolVersion': 2,
'peers': [_relayHostId],
});
} else if (connection == 1 && message['type'] == 'join') {
} else if (connection == 1 && message['type'] == 'resume') {
relay.send(socket, {
'type': 'error',
'code': 'peer_id_unavailable',
@@ -931,7 +868,7 @@ void main() {
expect(relay.messages, hasLength(2));
expect(relay.messages[0].map((message) => message['type']), ['join', 'leave']);
expect(relay.messages[1].map((message) => message['type']), ['join']);
expect(relay.messages[1].map((message) => message['type']), ['resume']);
});
test('a guest promoted mid-teardown ends the room instead of reading the rejected leave as released', () async {
@@ -944,8 +881,10 @@ void main() {
relay = await relayWith((connection, socket, message) {
switch (message['type']) {
case 'join':
case 'resume':
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': connection == 0 ? _relayHostId : message['peerId'],
'reconnectToken': message['reconnectToken'],
@@ -970,7 +909,7 @@ void main() {
expect(relay.messages, hasLength(2));
expect(relay.messages[0].map((message) => message['type']), ['join', 'leave']);
expect(relay.messages[1].map((message) => message['type']), ['join', 'endSession']);
expect(relay.messages[1].map((message) => message['type']), ['resume', 'endSession']);
});
test('a host demoted mid-teardown leaves as the guest it became instead of failing the exit', () async {
@@ -987,9 +926,10 @@ void main() {
});
case 'endSession':
relay.send(socket, {'type': 'error', 'code': 'peer_id_unavailable', 'message': 'Unable to end room'});
case 'join':
case 'resume':
relay.send(socket, {
'type': 'joined',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': 'new-host',
'reconnectToken': message['reconnectToken'],
@@ -1016,7 +956,7 @@ void main() {
expect(relay.messages, hasLength(2));
expect(relay.messages[0].map((message) => message['type']), ['create', 'endSession']);
expect(relay.messages[1].map((message) => message['type']), ['join', 'leave']);
expect(relay.messages[1].map((message) => message['type']), ['resume', 'leave']);
});
test('a release the relay keeps refusing for an unchanged role stays bounded and surfaces', () async {
@@ -1024,8 +964,10 @@ void main() {
relay = await relayWith((connection, socket, message) {
switch (message['type']) {
case 'join':
case 'resume':
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
@@ -1045,8 +987,8 @@ void main() {
);
expect(relay.messages, hasLength(3));
for (final connection in relay.messages) {
expect(connection.map((message) => message['type']), ['join', 'leave']);
for (var i = 0; i < relay.messages.length; i++) {
expect(relay.messages[i].map((message) => message['type']), [i == 0 ? 'join' : 'resume', 'leave']);
}
});
@@ -1089,7 +1031,7 @@ void main() {
]);
});
test('sequential guest release accepts not-in-room as idempotent success', () async {
test('a concurrent host end completes guest release without requesting membership again', () async {
var leaveRequests = 0;
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
@@ -1104,16 +1046,12 @@ void main() {
});
} else if (message['type'] == 'leave') {
leaveRequests++;
if (leaveRequests == 1) {
relay.send(socket, {
'type': 'left',
'sessionId': message['sessionId'],
'peerId': message['peerId'],
'protocolVersion': 2,
});
} else {
relay.send(socket, {'type': 'error', 'code': 'not_in_room', 'message': 'Peer is not in the room'});
}
relay.send(socket, {
'type': 'ended',
'sessionId': message['sessionId'],
'peerId': message['peerId'],
'protocolVersion': 2,
});
}
});
final service = serviceFor(relay);
@@ -1122,8 +1060,8 @@ void main() {
await service.releaseSession();
await service.releaseSession();
expect(leaveRequests, 2);
expect(relay.messages.single.map((message) => message['type']), ['join', 'leave', 'leave']);
expect(leaveRequests, 1);
expect(relay.messages.single.map((message) => message['type']), ['join', 'leave']);
});
test('host end waits for a protocol-2 ended acknowledgement', () async {
@@ -1264,10 +1202,11 @@ void main() {
final service = serviceFor(relay);
final pending = service.createSession(sessionId: 'cancel1');
final cancelled = expectLater(pending, throwsStateError);
await announcementSeen.future.timeout(const Duration(seconds: 1));
await service.disconnect();
await expectLater(pending, throwsStateError);
await cancelled.timeout(const Duration(seconds: 1));
expect(service.sessionId, isNull);
expect(service.connectedPeers, isEmpty);
});
@@ -1285,9 +1224,10 @@ void main() {
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (connection == 1 && message['type'] == 'join') {
} else if (connection == 1 && message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
@@ -1460,14 +1400,17 @@ void main() {
test('reconnection must reestablish enforcing-feature and roster authorization', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (message['type'] == 'create' || message['type'] == 'join') {
if (message['type'] == 'create' || message['type'] == 'resume') {
relay.send(socket, {
'type': message['type'] == 'create' ? 'created' : 'joined',
'type': message['type'] == 'create' ? 'created' : 'resumed',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
if (connection != 1) 'features': [RelayProtocol.atomicHostTransferFeature],
'features': [
RelayProtocol.authenticatedResumeFeature,
if (connection != 1) RelayProtocol.atomicHostTransferFeature,
],
});
relay.send(socket, {
'type': 'hostTransferEligibility',
@@ -1489,7 +1432,7 @@ void main() {
service.onReconnected = () => reconnected.complete();
await relay.sockets.single.close();
await reconnected.future.timeout(const Duration(seconds: 5));
expect(service.canTransferHostTo('guest-1'), isFalse, reason: 'old relay cannot inherit prior ACK');
expect(service.canTransferHostTo('guest-1'), isFalse, reason: 'resume cannot inherit prior transfer ACK');
final restored = service.onHostTransferEligibilityChanged.firstWhere(
(_) => service.canTransferHostTo('guest-1'),
);
@@ -1576,9 +1519,10 @@ void main() {
test('guest reconnect accepts the host identity pinned by a transfer', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) {
if (message['type'] == 'join') {
if (message['type'] == 'join' || message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': connection == 0 ? _relayHostId : 'guest-2',
'reconnectToken': message['reconnectToken'],
@@ -1630,9 +1574,10 @@ void main() {
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (connection >= 1 && message['type'] == 'join') {
} else if (connection >= 1 && message['type'] == 'resume') {
relay.send(socket, {
'type': 'joined',
'type': 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
'hostPeerId': 'new-host',
'reconnectToken': message['reconnectToken'],
@@ -1672,8 +1617,10 @@ void main() {
relay = await relayWith((connection, socket, message) {
switch (message['type']) {
case 'join':
case 'resume':
relay.send(socket, {
'type': 'joined',
'type': message['type'] == 'join' ? 'joined' : 'resumed',
'features': [RelayProtocol.authenticatedResumeFeature],
'sessionId': message['sessionId'],
// The reconnect: the relay made this peer the host while it was
// offline, so the hostChanged broadcast never reached it.
@@ -1721,6 +1668,353 @@ void main() {
// The role the relay declared is the one the transport acts on: a host
// destroys the room instead of quietly leaving it behind.
await service.releaseSession();
expect(relay.messages[1].map((message) => message['type']), ['join', 'endSession']);
expect(relay.messages[1].map((message) => message['type']), ['resume', 'endSession']);
});
test('lost initial join ACK resumes a promotion before setup returns', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) async {
if (message['type'] == 'join') {
await socket.close();
} else if (message['type'] == 'resume') {
relay.send(socket, {
'type': 'resumed',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
'features': [RelayProtocol.authenticatedResumeFeature],
});
} else if (message['type'] == 'endSession') {
relay.send(socket, {'type': 'ended', 'sessionId': message['sessionId'], 'protocolVersion': 2});
}
});
final service = serviceFor(relay);
await _withRetryBackoffShortened(() => service.joinSession('promotedSetup'));
expect(service.isHost, isTrue);
expect(service.hostPeerId, service.myPeerId);
await service.releaseSession();
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()), [
['join'],
['resume', 'endSession'],
]);
});
test('uncommitted ambiguous setup fails closed and only explicit re-entry mints membership', () async {
late final _RelayServer relay;
relay = await relayWith((connection, socket, message) async {
if (connection == 0 && message['type'] == 'create') {
await socket.close();
} else if (message['type'] == 'resume') {
relay.send(socket, {'type': 'error', 'code': 'room_not_found', 'message': 'Room not found'});
} else if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
}
});
final service = serviceFor(relay);
await expectLater(
_withRetryBackoffShortened(() => service.createSession(sessionId: 'ambiguous')),
throwsA(isA<PeerError>().having((error) => error.serverCode, 'serverCode', 'room_not_found')),
);
expect(service.sessionId, isNull);
expect(relay.messages.map((messages) => messages.single['type']), ['create', 'resume']);
final original = relay.messages.first.single;
await service.createSession(sessionId: 'ambiguous');
final explicitEntry = relay.messages.last.single;
expect(explicitEntry['type'], 'create');
expect(explicitEntry['peerId'], isNot(original['peerId']));
expect(explicitEntry['reconnectToken'], isNot(original['reconnectToken']));
});
for (final invalid in ['joinedTokenEcho', 'missingFeature', 'malformedPeers', 'wrongSession']) {
test('resume rejects $invalid before publishing authority or room traffic', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'join') {
relay.send(socket, {
'type': 'joined',
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
'peers': [_relayHostId],
});
} else if (message['type'] == 'resume') {
relay.send(socket, {'type': 'hostChanged', 'sessionId': message['sessionId'], 'hostPeerId': 'replacement'});
relay.send(socket, {
'type': invalid == 'joinedTokenEcho' ? 'joined' : 'resumed',
'sessionId': invalid == 'wrongSession' ? 'OTHER_ROOM' : message['sessionId'],
'hostPeerId': 'replacement',
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
if (invalid != 'missingFeature') 'features': [RelayProtocol.authenticatedResumeFeature],
'peers': invalid == 'malformedPeers' ? ['replacement', 7] : ['replacement'],
});
relay.send(socket, {'type': 'peerJoined', 'peerId': 'replacement'});
relay.send(socket, {
'type': 'message',
'from': 'replacement',
'payload': SyncMessage.requestState().toJson(),
});
}
});
final service = serviceFor(relay);
final authorities = <String>[];
final connected = <bool>[];
final peers = <String>[];
final payloads = <SyncMessage>[];
var reconnectCallbacks = 0;
service.onReconnected = () => reconnectCallbacks++;
addTearDown(service.onHostChanged.listen(authorities.add).cancel);
addTearDown(service.onConnectionStateChanged.listen(connected.add).cancel);
addTearDown(service.onPeerConnected.listen(peers.add).cancel);
addTearDown(service.onMessageReceived.listen(payloads.add).cancel);
final ended = service.onSessionEnded.first;
await _withShortenedTimer(
original: const Duration(seconds: 2),
replacement: const Duration(milliseconds: 10),
body: () async {
await service.joinSession('validation');
await relay.sockets.single.close();
await ended.timeout(const Duration(seconds: 1));
},
);
expect(authorities, isEmpty);
expect(peers, [_relayHostId]);
expect(connected.where((value) => value), [true]);
expect(payloads, isEmpty);
expect(reconnectCallbacks, 0);
expect(service.hostPeerId, isNull);
expect(service.sessionId, isNull);
expect(relay.messages.map((messages) => messages.single['type']), ['join', 'resume']);
});
}
for (final releaseInstead in [false, true]) {
test(
'released relay compatibility fails closed on ${releaseInstead ? 'disconnected release' : 'resume'}',
() async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'join') {
// Released protocol-2 relays support initial admission without features.
relay.send(socket, {
'type': 'joined',
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
'peers': [_relayHostId],
});
} else {
relay.send(socket, {'type': 'error', 'code': 'invalid_message', 'message': 'Unknown message type'});
}
});
final service = serviceFor(relay);
var reconnectCallbacks = 0;
service.onReconnected = () => reconnectCallbacks++;
await _withShortenedTimer(
original: const Duration(seconds: 2),
replacement: releaseInstead ? const Duration(seconds: 2) : const Duration(milliseconds: 10),
body: () async {
await service.joinSession('oldRelay');
final disconnected = service.onConnectionStateChanged.firstWhere((connected) => !connected);
final ended = releaseInstead ? null : service.onSessionEnded.first;
await relay.sockets.single.close();
await disconnected.timeout(const Duration(seconds: 1));
if (releaseInstead) {
await expectLater(
service.releaseSession(),
throwsA(isA<PeerError>().having((error) => error.serverCode, 'serverCode', 'invalid_message')),
);
} else {
await ended!.timeout(const Duration(seconds: 1));
expect(service.sessionId, isNull);
}
},
);
expect(reconnectCallbacks, 0);
expect(relay.messages.map((messages) => messages.single['type']), ['join', 'resume']);
},
);
}
test('lost host end ACK cannot recreate or end a replacement room during cleanup', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (message['type'] == 'resume') {
// End committed and this reusable code now names unrelated membership.
relay.send(socket, {'type': 'error', 'code': 'peer_id_unavailable', 'message': 'Peer ID is unavailable'});
}
});
final service = serviceFor(relay, debugReleaseTimeout: const Duration(milliseconds: 10));
await service.createSession(sessionId: 'reusedCode');
await _withRetryBackoffShortened(service.releaseSession);
await service.releaseSession();
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()), [
['create', 'endSession'],
['resume'],
]);
});
test('transport failure before sending admission may retry the initial operation', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
}
});
var attempts = 0;
final service = serviceFor(
relay,
debugChannelFactory: (uri) {
if (attempts++ == 0) throw const SocketException('Connection refused before admission');
return WebSocketChannel.connect(uri);
},
);
await _withRetryBackoffShortened(() => service.createSession(sessionId: 'preSend'));
expect(service.isHost, isTrue);
expect(relay.messages.single.single['type'], 'create');
});
test('network failures exhaust bounded resume attempts without falling back to create', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) async {
if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (message['type'] == 'resume') {
await socket.close();
}
});
final service = serviceFor(relay);
var reconnectCallbacks = 0;
service.onReconnected = () => reconnectCallbacks++;
final ended = service.onSessionEnded.first;
await _withShortenedTimer(
original: const Duration(seconds: 2),
replacement: const Duration(milliseconds: 10),
body: () => _withShortenedTimer(
original: const Duration(seconds: 4),
replacement: const Duration(milliseconds: 10),
body: () => _withShortenedTimer(
original: const Duration(seconds: 6),
replacement: const Duration(milliseconds: 10),
body: () async {
await service.createSession(sessionId: 'networkLoss');
await relay.sockets.single.close();
},
),
),
);
// Only reconnect backoff belongs to the shortened timer zone.
await ended.timeout(const Duration(seconds: 2));
expect(reconnectCallbacks, 0);
expect(service.sessionId, isNull);
expect(relay.messages.map((messages) => messages.single['type']), ['create', 'resume', 'resume', 'resume']);
});
test('host release does not treat a guest left ACK as room destruction', () async {
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
} else if (message['type'] == 'endSession') {
relay.send(socket, {
'type': 'left',
'sessionId': message['sessionId'],
'peerId': message['peerId'],
'protocolVersion': 2,
});
}
});
final service = serviceFor(relay);
await service.createSession(sessionId: 'wrongReleaseAck');
await expectLater(
service.releaseSession(),
throwsA(isA<PeerError>().having((error) => error.type, 'type', PeerErrorType.serverError)),
);
expect(relay.messages.single.map((message) => message['type']), ['create', 'endSession']);
});
test('cancelling pending resume cannot publish or release a later explicit session', () async {
final resumeSeen = Completer<void>();
late final _RelayServer relay;
relay = await relayWith((_, socket, message) {
if (message['type'] == 'join') {
relay.send(socket, {
'type': 'joined',
'sessionId': message['sessionId'],
'hostPeerId': _relayHostId,
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
'peers': [_relayHostId],
});
} else if (message['type'] == 'resume') {
resumeSeen.complete();
} else if (message['type'] == 'create') {
relay.send(socket, {
'type': 'created',
'sessionId': message['sessionId'],
'hostPeerId': message['peerId'],
'reconnectToken': message['reconnectToken'],
'protocolVersion': 2,
});
}
});
final service = serviceFor(relay);
var reconnectCallbacks = 0;
service.onReconnected = () => reconnectCallbacks++;
await _withShortenedTimer(
original: const Duration(seconds: 2),
replacement: const Duration(milliseconds: 10),
body: () async {
await service.joinSession('cancelOld');
await relay.sockets.single.close();
await resumeSeen.future.timeout(const Duration(seconds: 1));
await service.disconnect();
await service.createSession(sessionId: 'explicitNew');
},
);
expect(service.sessionId, 'EXPLICITNEW');
expect(service.hostPeerId, service.myPeerId);
expect(service.isHost, isTrue);
expect(reconnectCallbacks, 0);
expect(relay.messages.map((messages) => messages.map((message) => message['type']).toList()), [
['join'],
['resume'],
['create'],
]);
});
}