Skip to content
Open
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion action_subscriber.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ Gem::Specification.new do |spec|
spec.add_development_dependency "bundler", ">= 1.6"
spec.add_development_dependency "pry-coolline"
spec.add_development_dependency "pry-nav"
spec.add_development_dependency "rabbitmq_http_api_client", "~> 1.2.0"
spec.add_development_dependency "rabbitmq_http_api_client", "~> 1.9.0"
spec.add_development_dependency "rspec", "~> 3.0"
spec.add_development_dependency "rake"
end
20 changes: 17 additions & 3 deletions lib/action_subscriber/bunny/subscriber.rb
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ def bunny_consumers
@bunny_consumers ||= []
end

def set_bunny_consumers(consumers)
@bunny_consumers = consumers
end

def cancel_consumers!
bunny_consumers.each(&:cancel)
::ActionSubscriber::ThreadPools.threadpools.each do |name, threadpool|
Expand All @@ -24,14 +28,19 @@ def setup_subscriptions!
end
end

def start_subscribers!
subscriptions.each do |subscription|
route = subscription[:route]
def start_subscription!(subscription)
route = subscription[:route]
queue = subscription[:queue]
channel = queue.channel
threadpool = ::ActionSubscriber::ThreadPools.threadpools.fetch(route.threadpool_name)
channel.prefetch(route.prefetch) if route.acknowledgements?
consumer = ::Bunny::Consumer.new(channel, queue, channel.generate_consumer_tag, !route.acknowledgements?)
consumer.on_cancellation do |_basic_cancel|
Comment thread
hugh-j marked this conversation as resolved.
Outdated
set_bunny_consumers(bunny_consumers.reject { |bunny_consumer| bunny_consumer == consumer })
subscription[:queue] = setup_queue(route)
start_subscription!(subscription)
end

consumer.on_delivery do |delivery_info, properties, encoded_payload|
::ActiveSupport::Notifications.instrument "received_event.action_subscriber", :payload_size => encoded_payload.bytesize, :queue => queue.name
properties = {
Expand All @@ -51,6 +60,11 @@ def start_subscribers!
end
bunny_consumers << consumer
queue.subscribe_with(consumer)
end

def start_subscribers!
Comment thread
liveh2o marked this conversation as resolved.
subscriptions.each do |subscription|
start_subscription!(subscription)
end
end

Expand Down
27 changes: 23 additions & 4 deletions lib/action_subscriber/march_hare/subscriber.rb
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,10 @@ def cancel_consumers!
def march_hare_consumers
@march_hare_consumers ||= []
end

def set_march_hare_consumers(consumers)
@march_hare_consumers = consumers
end

def setup_subscriptions!
fail ::RuntimeError, "you cannot setup queues multiple times, this should only happen once at startup" unless subscriptions.empty?
Expand All @@ -23,14 +27,23 @@ def setup_subscriptions!
}
end
end

def start_subscribers!
subscriptions.each do |subscription|

def start_subscription!(subscription)
route = subscription[:route]
queue = subscription[:queue]
queue.channel.prefetch = route.prefetch if route.acknowledgements?
threadpool = ::ActionSubscriber::ThreadPools.threadpools.fetch(route.threadpool_name)
consumer = queue.subscribe(route.queue_subscription_options) do |metadata, encoded_payload|

cancel = {
:on_cancellation => Proc.new { |channel, consumer|
channel.close
set_march_hare_consumers(march_hare_consumers.reject { |march_hare_consumer| march_hare_consumer == consumer })

subscription[:queue] = setup_queue(route)
start_subscription!(subscription)
}
}
consumer = queue.subscribe(cancel.merge(route.queue_subscription_options)) do |metadata, encoded_payload|
::ActiveSupport::Notifications.instrument "received_event.action_subscriber", :payload_size => encoded_payload.bytesize, :queue => queue.name
properties = {
:action => route.action,
Expand All @@ -49,6 +62,12 @@ def start_subscribers!
end

march_hare_consumers << consumer
end

def start_subscribers!
route_set = self
subscriptions.each do |subscription|
start_subscription! subscription
end
end

Expand Down
34 changes: 34 additions & 0 deletions spec/integration/consumer_cancel_notify_spec.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
require "rabbitmq/http/client"

class ZombieSubscriber < ActionSubscriber::Base
def groan
$messages << payload
end
end

describe "Rebuilds subscription after receiving consumer_cancel_notify", :integration => true, :slow => true do
let(:draw_routes) do
::ActionSubscriber.draw_routes do
default_routes_for ZombieSubscriber
end
end
let(:http_client) { RabbitMQ::HTTP::Client.new("http://127.0.0.1:15672") }
let(:subscriber) { ZombieSubscriber }

it "continues to receive messages following consumer_cancel_notify message" do
::ActionSubscriber::start_subscribers!
::ActivePublisher.publish("zombie.groan", "uuunngg", "events")

delete_queue!
sleep 5.0

::ActivePublisher.publish("zombie.groan", "meeuuhhh", "events")
verify_expectation_within(5.0) do
expect($messages).to eq(Set.new(["uuunngg", "meeuuhhh"]))
end
end

def delete_queue!
http_client.delete_queue("/","alice.zombie.groan")
end
end