addInlet method

  1. @override
Future<void> addInlet(
  1. PeerHandle handle
)
override

Subscribes to a peer found by discovery.

Implementations must take ownership of handle (see PeerHandle.take) before using it.

Implementation

@override
Future<void> addInlet(PeerHandle handle) async {
  if (_disposed) return;
  if (!handle.taken) handle.take();

  final producer = handle.descriptor.endpointId;
  // Deliberately no self-check here. Whether a node consumes its own
  // output is decided by StreamParticipationMode via
  // PeerSession.getProducersForStream — allNodes, coordinatorOnly and
  // sendAllReceiveCoordinator all legitimately include the local node — and
  // `consumeCoordinationStreamAsCoordinator` has the coordinator subscribe
  // to its own coordination stream on purpose. A transport that silently
  // skipped self would override those decisions.
  if (!_subscribedProducers.add(producer)) return;
  if (handle is WsPeerHandle && handle.slot != null) {
    _slotOwners[handle.slot!] = handle.descriptor.nodeUId;
  }
  _producerByNodeUId[handle.descriptor.nodeUId] = producer;

  connection.subscribe(
    streamName: config.name,
    subscriberEndpointId: endpointId,
    producerEndpointIds: [producer],
  );
  await handle.dispose();
}