consume method
Stream<Delivery>
consume(
- RoutingSubscription subscription, {
- int prefetch = 1,
- String? consumerGroup,
- 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;
}