consume method

  1. @override
Stream<Delivery> consume(
  1. RoutingSubscription subscription, {
  2. int prefetch = 1,
  3. String? consumerGroup,
  4. String? consumerName,
})
override

Returns a stream of deliveries based on the supplied subscription.

The prefetch parameter specifies the number of messages to prefetch. consumerGroup and consumerName can be used for consumer identification.

Implementation

@override
Stream<Delivery> consume(
  RoutingSubscription subscription, {
  int prefetch = 1,
  String? consumerGroup,
  String? consumerName,
}) {
  if (subscription.queues.length > 1) {
    throw UnsupportedError(
      'InMemoryBroker currently supports consuming a single queue at a time.',
    );
  }
  final queue = subscription.queues.firstOrNull;
  final state = queue == null ? null : _state(queue);
  final consumer = consumerName ?? const Uuid().v7();
  final consumerKey = '${consumerGroup ?? 'default'}::$consumer';
  _BroadcastSubscription? broadcastSubscription;
  var active = true;

  late StreamController<Delivery> controller;
  controller = StreamController<Delivery>(
    onListen: () async {
      if (subscription.broadcastChannels.isNotEmpty) {
        broadcastSubscription = _broadcastHub.subscribe(
          consumer: consumerKey,
          channels: subscription.broadcastChannels,
          onDelivery: (delivery) {
            if (controller.isClosed) return;
            controller.add(delivery);
          },
        );
        _activeBroadcastSubscriptions.add(broadcastSubscription!);
      }

      if (state == null) {
        return;
      }

      state.resumeConsumer(consumer);
      try {
        while (active && !controller.isClosed) {
          final delivery = await state.nextDelivery(
            consumer: consumer,
            prefetch: prefetch,
            defaultVisibilityTimeout: defaultVisibilityTimeout,
          );
          if (!active || controller.isClosed || !controller.hasListener) {
            state.requeue(delivery.receipt);
            break;
          }
          controller.add(delivery);
        }
      } on _ConsumerCancelled {
        return;
      }
    },
    onCancel: () {
      active = false;
      state?.cancelWaiters(consumer);
      final broadcastSub = broadcastSubscription;
      if (broadcastSub != null) {
        broadcastSub.close();
        _activeBroadcastSubscriptions.remove(broadcastSub);
      }
    },
  );
  return controller.stream;
}