publish method

  1. @override
Future<void> publish(
  1. Envelope envelope, {
  2. RoutingInfo? routing,
})
override

Enqueues a message into the in-memory queue or delay set.

Implementation

@override
/// Enqueues a message into the in-memory queue or delay set.
Future<void> publish(Envelope envelope, {RoutingInfo? routing}) async {
  final resolvedRoute =
      routing ??
      RoutingInfo.queue(queue: envelope.queue, priority: envelope.priority);
  if (resolvedRoute.isBroadcast) {
    final channel = resolvedRoute.broadcastChannel ?? envelope.queue;
    final message = envelope.copyWith(queue: channel);
    _broadcastHub.publish(
      channel: channel,
      envelope: message,
      delivery: resolvedRoute.delivery ?? 'at-least-once',
    );
    return;
  }
  final targetQueue = resolvedRoute.queue ?? envelope.queue;
  final state = _state(targetQueue);
  final msg = envelope.copyWith(
    queue: targetQueue,
    priority: resolvedRoute.priority ?? envelope.priority,
  );

  if (msg.notBefore != null && msg.notBefore!.isAfter(stemNow())) {
    state.addDelayed(msg);
  } else {
    state.enqueue(msg);
  }
}