discoverOnce method

  1. @override
Future<List<PeerHandle>> discoverOnce(
  1. DiscoveryQuery query, {
  2. Duration timeout = const Duration(seconds: 2),
  3. int minPeers = 0,
  4. int maxPeers = 10,
})
override

Runs a single bounded discovery.

Returns as soon as minPeers have been found, or when timeout expires — so a query expected to match nothing costs the full timeout. Election relies on exactly that: "no better candidate" can only be concluded by waiting the query out.

Handles returned here are not owned by any discovery cycle, so PeerHandle.take on them is a no-op and the caller is responsible for disposing any it does not pass to a stream.

Implementation

@override
Future<List<PeerHandle>> discoverOnce(
  DiscoveryQuery query, {
  Duration timeout = const Duration(seconds: 2),
  int minPeers = 0,
  int maxPeers = 10,
}) async {
  // Election depends on this waiting out the timeout when nothing matches:
  // "no better candidate exists" is only knowable by asking and hearing
  // nothing back for long enough. Returning early on an empty result would
  // make every node elect itself and split the session.
  final deadline = DateTime.now().add(timeout);
  final results = <WsPeerHandle>[];
  final completer = Completer<void>();

  final qid = connection.query(query, continuous: true);
  final subscription = connection.control.listen((frame) {
    if (frame.type != WsControl.queryResult) return;
    if (frame.payload['qid'] != qid) return;
    results
      ..clear()
      ..addAll(_peersFrom(frame));
    if (results.length >= minPeers &&
        results.isNotEmpty &&
        !completer.isCompleted) {
      completer.complete();
    }
  });

  try {
    if (minPeers > 0) {
      final remaining = deadline.difference(DateTime.now());
      if (remaining > Duration.zero) {
        await completer.future.timeout(remaining, onTimeout: () {});
      }
    } else {
      // Give the hub one round trip to answer.
      await Future<void>.delayed(const Duration(milliseconds: 20));
    }
  } finally {
    await subscription.cancel();
    connection.unquery(qid);
  }

  if (results.length < minPeers) return const [];
  return results.take(maxPeers).toList(growable: false);
}