This is a default queue implementation that ships with Notifications. It just pushes events to all registered log subscribers.
This class is thread safe. All methods are reentrant.
Namespace
- MODULE ActiveSupport::Notifications::Fanout::Subscribers
- CLASS ActiveSupport::Notifications::Fanout::Handle
Methods
- A
- B
- F
- L
- N
- P
- S
- U
- W
Class Public methods
new() Link
Instance Public methods
all_listeners_for(name) Link
# File activesupport/lib/active_support/notifications/fanout.rb, line 398 def all_listeners_for(name) # this is correctly done double-checked locking (Concurrent::Map's lookups have volatile semantics) @all_listeners_for[name] || @mutex.synchronize do # use synchronisation when accessing @registry @all_listeners_for[name] ||= @registry[name] + @registry.other_subscribers.select { |s| s.subscribed_to?(name) } end end
build_handle(name, id, payload) Link
# File activesupport/lib/active_support/notifications/fanout.rb, line 365 def build_handle(name, id, payload) groups = groups_for(name).map do |group_klass, grouped_listeners| group_klass.new(grouped_listeners, name, id, payload) end if groups.empty? NullHandle else Handle.new(name, id, groups, payload) end end
finish(name, id, payload, listeners = nil) Link
listeners_for(name) Link
listening?(name) Link
publish(name, ...) Link
publish_event(event) Link
start(name, id, payload) Link
subscribe(pattern = nil, callable = nil, monotonic: false, prepend: false, &block) Link
# File activesupport/lib/active_support/notifications/fanout.rb, line 144 def subscribe(pattern = nil, callable = nil, monotonic: false, prepend: false, &block) block = ActiveSupport::Ractors.try_shareable_proc(block) if block subscriber = Subscribers.new(pattern, callable || block, monotonic) @mutex.synchronize do case pattern when String @registry = @registry.add(subscriber, prepend: prepend) clear_cache(pattern) when NilClass, Regexp if prepend raise ArgumentError, "Cannot prepend Regex subscribers" end @registry = @registry.add(subscriber) clear_cache else raise ArgumentError, "pattern must be specified as a String, Regexp or empty" end end subscriber end
unsubscribe(subscriber_or_name) Link
# File activesupport/lib/active_support/notifications/fanout.rb, line 166 def unsubscribe(subscriber_or_name) @mutex.synchronize do case subscriber_or_name when String @registry = @registry.delete_pattern(subscriber_or_name) @registry.other_subscribers.each { |sub| sub.unsubscribe!(subscriber_or_name) } clear_cache(subscriber_or_name) else pattern = subscriber_or_name.try(:pattern) pattern = nil unless String === pattern @registry = @registry.delete(subscriber_or_name) clear_cache(pattern) end end end