connectRelay method
Future<Tuple<bool, String> >
connectRelay({
- required String dirtyUrl,
- required ConnectionSource connectionSource,
- String? authPubkey,
- int connectTimeout = DEFAULT_WEB_SOCKET_CONNECT_TIMEOUT,
Connects to a relay to the relay pool. Returns a tuple with the first element being a boolean indicating success \ and the second element being a string with the error message if any.
Implementation
Future<Tuple<bool, String>> connectRelay({
required String dirtyUrl,
required ConnectionSource connectionSource,
String? authPubkey,
int connectTimeout = DEFAULT_WEB_SOCKET_CONNECT_TIMEOUT,
}) async {
String? url = cleanRelayUrl(dirtyUrl);
if (url == null) {
updateRelayConnectivity();
return Tuple(false, "unclean url");
}
if (globalState.blockedRelays.contains(url)) {
updateRelayConnectivity();
return Tuple(false, "relay is blocked");
}
final connectionKey = authPubkey == null
? RelayConnectionKey.anonymous(url)
: RelayConnectionKey.authenticated(url, authPubkey);
if (isConnectionOpen(connectionKey)) {
Logger.log.t(() => "relay already connected: $connectionKey");
updateRelayConnectivity();
return Tuple(true, "");
}
if (isConnectionConnecting(connectionKey) ||
_connectReadyCompleters.containsKey(connectionKey)) {
Logger.log.t(() => "relay is already connecting: $connectionKey");
final inFlightConnect = _connectReadyCompleters[connectionKey];
if (inFlightConnect != null) {
final connected = await inFlightConnect.future;
updateRelayConnectivity();
return Tuple(
connected,
connected
? "relay finished connecting"
: "relay failed while connecting",
);
}
updateRelayConnectivity();
return Tuple(false, "relay is still connecting");
}
RelayConnectivity? relayConnectivity = globalState.relays[connectionKey];
final connectCompleter = Completer<bool>();
_connectReadyCompleters[connectionKey] = connectCompleter;
NostrTransport? transport;
bool ownsTransport() =>
transport != null &&
identical(globalState.relays[connectionKey], relayConnectivity) &&
identical(relayConnectivity?.relayTransport, transport);
try {
if (relayConnectivity == null) {
relayConnectivity = RelayConnectivity<T>(
key: connectionKey,
relay: Relay(url: url, connectionSource: connectionSource),
specificEngineData: engineAdditionalDataFactory?.call(),
);
globalState.relays[connectionKey] = relayConnectivity;
}
relayConnectivity.relay.tryingToConnect();
Logger.log.i(() => "connecting to relay $dirtyUrl");
// a fresh socket for a key we may already know: nothing the previous one
// authenticated carries over
_forgetAuthState(connectionKey);
// A disconnected transport may still own reconnect timers and listeners.
// Retire it before replacement or it can later open an untracked socket.
await relayConnectivity.close();
if (!identical(globalState.relays[connectionKey], relayConnectivity)) {
throw StateError(
'Connection was removed while replacing its transport',
);
}
transport = nostrTransportFactory(
url,
onReconnect: () {
if (!ownsTransport()) return;
// the relay accepted our AUTH on the socket that just died, not on
// this one; the binding survives, the authentication does not
_forgetAuthState(connectionKey);
reSubscribeInFlightSubscriptions(relayConnectivity!);
updateRelayConnectivity();
},
onDisconnect: (code, error, reason) {
if (!ownsTransport()) return;
relayConnectivity!.stats.connectionErrors++;
// the transport reconnects under us and keeps its message stream
// open, so this is the only notice we get that the socket the relay
// authenticated died
_forgetAuthState(connectionKey);
// the requests died with that socket too; onReconnect replays them
relayConnectivity.stats.openRequestIds.clear();
updateRelayConnectivity();
},
);
relayConnectivity.relayTransport = transport;
// Start listening immediately so we don't miss early frames such as
// relay AUTH challenges that may arrive before the transport reports
// itself fully open.
_startListeningToSocket(relayConnectivity);
final opened = await _waitForTransportOpen(
transport,
timeoutSeconds: connectTimeout,
stillOwned: ownsTransport,
);
if (!opened || !ownsTransport()) {
throw TimeoutException(
"Future not completed",
Duration(seconds: connectTimeout),
);
}
Logger.log.i(() => "connected to relay: $url");
relayConnectivity.relay.succeededToConnect();
relayConnectivity.stats.connections++;
getRelayInfo(url).then((info) {
relayConnectivity!.relayInfo = info;
});
if (!connectCompleter.isCompleted) {
connectCompleter.complete(true);
}
if (identical(_connectReadyCompleters[connectionKey], connectCompleter)) {
_connectReadyCompleters.remove(connectionKey);
}
updateRelayConnectivity();
return Tuple(true, "");
} catch (e) {
Logger.log.e(() => "!! could not connect to $url -> $e");
try {
if (transport == null ||
identical(relayConnectivity!.relayTransport, transport)) {
await relayConnectivity!.close();
} else {
await transport.close();
}
} catch (closeError) {
Logger.log.w(() => "Error retiring transport for $url: $closeError");
}
}
final failedConnectivity = relayConnectivity;
if (failedConnectivity != null) {
failedConnectivity.relay.failedToConnect();
failedConnectivity.stats.connectionErrors++;
}
if (!connectCompleter.isCompleted) {
connectCompleter.complete(false);
}
if (identical(_connectReadyCompleters[connectionKey], connectCompleter)) {
_connectReadyCompleters.remove(connectionKey);
}
updateRelayConnectivity();
return Tuple(false, "could not connect to $url");
}