java-topology/defects/rabbitmq/patch/rabbitmq-0002-queue-consumers-blocked-set.patch

24 lines
1.4 KiB
Diff

# UNDF: UNDF-2026-000000240
--- a/deps/rabbit/src/rabbit_queue_consumers.erl
+++ b/deps/rabbit/src/rabbit_queue_consumers.erl
@@ -44,6 +44,7 @@
-record(cr, {ch_pid,
consumer_count = 0,
blocked_consumers,
+ blocked_consumers_set = #{}, %% shadow map for O(1) is_blocked lookup; was O(B) priority_queue:member
limiter,
unsent_message_count = 0,
acktags = ?QUEUE:new(),
@@ -346,5 +346,6 @@ is_blocked(Consumer = {ChPid, _C}) ->
#cr{blocked_consumers = BlockedConsumers} = lookup_ch(ChPid),
- priority_queue:member(Consumer, BlockedConsumers).
+ #cr{blocked_consumers_set = BSet} = lookup_ch(ChPid),
+ maps:is_key(Consumer, BSet). %% O(1); was O(B) priority_queue:member scan
%% block_consumer: add to BlockedConsumers and set in shadow map
-block_consumer(C = #cr{blocked_consumers = BlockedQ}, QEntry = {_ChPid, #consumer{cfg = #consumer_cfg{priority = Priority}}}) ->
- update_ch_record(C#cr{blocked_consumers = priority_queue:in(QEntry, Priority, BlockedQ)});
+block_consumer(C = #cr{blocked_consumers = BlockedQ, blocked_consumers_set = BSet},
+ QEntry = {_ChPid, #consumer{cfg = #consumer_cfg{priority = Priority}}}) ->
+ update_ch_record(C#cr{blocked_consumers = priority_queue:in(QEntry, Priority, BlockedQ),
+ blocked_consumers_set = maps:put(QEntry, true, BSet)});