diff --git a/.travis.yml b/.travis.yml index 3175e24..2dec53c 100644 --- a/.travis.yml +++ b/.travis.yml @@ -4,3 +4,4 @@ rvm: - 1.9.2 - 1.9.3 - 2.0.0 + - 2.1.0 diff --git a/CHANGELOG.md b/CHANGELOG.md new file mode 100644 index 0000000..a8f4422 --- /dev/null +++ b/CHANGELOG.md @@ -0,0 +1,15 @@ +# Arbiter + +## 3.0.1 + +- Remove required dependencies + +## 3.0.0 + +- Changed away from ruby 'Marshal' in favor of sending parameters via 'MultiJSON' when communicating with Majordomo in the Async Arbiter. + +## 2.0.0 + +- Removed plain ZeroMQ Push/Pull Arbiter implementation. We don't use it locally and to be honest it's not worth the effort to attempt to update it since we use the Majordomo pattern internally. Let us know via an issue if you use it and would see value in us adding it back. +- Added new ZeroMQ::MajordomoAsynchronousArbiter. See README for usage details. +- Bump ffi-rzmq dependency to ~>2.0. We use this internally. diff --git a/Gemfile b/Gemfile index 13a9169..304102c 100644 --- a/Gemfile +++ b/Gemfile @@ -1,10 +1,10 @@ source "http://rubygems.org" -gem 'resque', '~>1.21' -gem "ffi-rzmq", "~>1.0" - group :development do - gem 'rspec', '~>2.11' - gem 'oj', '~>2.0' - gem 'rake', '~>0.9' + gem 'rspec', '~> 2.99' + gem 'oj', '~> 2.0' + gem 'rake', '~> 0.9' + + gem 'arbiter' + gem 'ffi-rzmq' end diff --git a/Gemfile.lock b/Gemfile.lock index 5147363..a1bfe0e 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -1,46 +1,33 @@ GEM remote: http://rubygems.org/ specs: - diff-lcs (1.1.3) - ffi (1.9.0) - ffi-rzmq (1.0.1) + arbiter (3.0.0) + diff-lcs (1.3) + ffi (1.9.23) + ffi-rzmq (2.0.6) + ffi-rzmq-core (>= 1.0.6) + ffi-rzmq-core (1.0.6) ffi - multi_json (1.3.6) - oj (2.0.10) - rack (1.4.1) - rack-protection (1.2.0) - rack - rake (0.9.2.2) - redis (3.0.1) - redis-namespace (1.2.1) - redis (~> 3.0.0) - resque (1.21.0) - multi_json (~> 1.0) - redis-namespace (~> 1.0) - sinatra (>= 0.9.2) - vegas (~> 0.1.2) - rspec (2.11.0) - rspec-core (~> 2.11.0) - rspec-expectations (~> 2.11.0) - rspec-mocks (~> 2.11.0) - rspec-core (2.11.1) - rspec-expectations (2.11.2) - diff-lcs (~> 1.1.3) - rspec-mocks (2.11.1) - sinatra (1.3.2) - rack (~> 1.3, >= 1.3.6) - rack-protection (~> 1.2) - tilt (~> 1.3, >= 1.3.3) - tilt (1.3.3) - vegas (0.1.11) - rack (>= 1.0.0) + oj (2.18.5) + rake (0.9.6) + rspec (2.99.0) + rspec-core (~> 2.99.0) + rspec-expectations (~> 2.99.0) + rspec-mocks (~> 2.99.0) + rspec-core (2.99.2) + rspec-expectations (2.99.2) + diff-lcs (>= 1.1.3, < 2.0) + rspec-mocks (2.99.4) PLATFORMS ruby DEPENDENCIES - ffi-rzmq (~> 1.0) + arbiter + ffi-rzmq oj (~> 2.0) rake (~> 0.9) - resque (~> 1.21) - rspec (~> 2.11) + rspec (~> 2.99) + +BUNDLED WITH + 1.16.1 diff --git a/README.md b/README.md index fcac53d..671acb7 100644 --- a/README.md +++ b/README.md @@ -104,25 +104,31 @@ The in-memory arbiter is the default, simplest driver. This is an in-memory arbi This is an Arbiter implementation that uses a Resque backend. You'll need a Resque worker running to process events. See the resque manual for details on this. -### ZeroMQ (ZeromqArbiter) +### ZeroMQ Majordomo (Zeromq::Majordomo::AsynchronousArbiter) -This is an Arbiter that uses ZeroMQ to send it's messages. This is by far the most advanced and powerful implementation, as you can use the power of zmq to setup any kind of messaging architecture you want. +This is an Arbiter that uses ZeroMQ and the [Majordomo Protocol](http://rfc.zeromq.org/spec:7) to send it's messages. It submits asynchronously to dispatch work via a broker (which is not provided by this gem). You must also have majordomo workers receiving and processing the work (workers are also not provided). -#### Running the backend +#### Configuring Zeromq::Majordomo::AsynchronousArbiter -To run the backend, you'll need to start two processes: +You'll need to setup some configuration in your app to use this arbiter. - - `rake "arbiter:proxy[frontend_uri,backend_uri]"` - - `rake "arbiter:worker[backend_uri]"` +* Address: The address of the Majordomo broker +* ZMQ Context: the ZMQ context implementation +* Timeout: timeout in seconds +* MD Service: The service name to use with Majordomo. This is how messages are routed. -The `frontend_uri` and `backend_uri` values above should conform to standard zmq addresses. You can use tcp, udp, ipc, or anything else that zmq supports. You must start the proxy first, and then connect the workers to the proxy. This allows you to scale the workers up and down as you see fit. For a light application, you'll probably only need one worker, but for very busy applications, you might need much more than that. +Example: -#### Configuring ZeromqArbiter - -You'll need to setup some configuration in your app to use zeromq. - -Set the backend worker: `ZeromqArbiter.frontend = 'frontend_uri'` - -The address should be the frontend location of your proxy. - -You can also set a logger for it if you wish: `ZeromqArbiter.logger = Logger` +```ruby +require 'arbiter' +require 'ffi-rzmq' + +class ZeromqConfig + arbiter = Zeromq::Majordomo::AsynchronousArbiter.new( + "tcp://192.168.1.1:9292", + ZMQ::Context.new, + 5, + "some-service-name-v1" + ) +end +``` diff --git a/arbiter.gemspec b/arbiter.gemspec index 1ba3093..0de1a6d 100644 --- a/arbiter.gemspec +++ b/arbiter.gemspec @@ -3,8 +3,8 @@ $:.unshift lib unless $:.include?(lib) Gem::Specification.new do |s| s.name = 'arbiter' - s.version = '1.0.1' - s.authors = ['Sitter City'] + s.version = '3.0.1' + s.authors = ['Sittercity'] s.email = ['dev@sittercity.com'] s.homepage = 'https://github.com/sittercity/arbiter' s.summary = 'A simple eventing framework' diff --git a/lib/zeromq/majordomo/asynchronous_arbiter.rb b/lib/zeromq/majordomo/asynchronous_arbiter.rb new file mode 100644 index 0000000..dc54051 --- /dev/null +++ b/lib/zeromq/majordomo/asynchronous_arbiter.rb @@ -0,0 +1,45 @@ +require 'ffi-rzmq' +require 'multi_json' + +module Zeromq + module Majordomo + class AsynchronousArbiter + INFINITE = -1 + + def initialize(address, context, timeout_in_sec, md_service) + @address = address + @context = context + @timeout = timeout_in_sec.to_i * 1000 #ms + @md_service = md_service + end + + def publish(method, params) + sock = connect + assert_zmq_ok(sock.send_strings( + ['', 'MDPC01', @md_service, method.to_s, MultiJson.dump(params)] + )) + ensure + disconnect(sock) if sock + end + + private + + def connect + sock = @context.socket(ZMQ::DEALER) + assert_zmq_ok(sock.setsockopt(ZMQ::LINGER, INFINITE)) + assert_zmq_ok(sock.setsockopt(ZMQ::SNDTIMEO, @timeout)) + assert_zmq_ok(sock.connect(@address)) + sock + end + + def disconnect(sock) + assert_zmq_ok(sock.disconnect(@address)) + assert_zmq_ok(sock.close) + end + + def assert_zmq_ok(rc) + raise ZMQ::Util.error_string unless ZMQ::Util.resultcode_ok?(rc) + end + end + end +end diff --git a/lib/zeromq_arbiter.rb b/lib/zeromq_arbiter.rb deleted file mode 100644 index cf8b06c..0000000 --- a/lib/zeromq_arbiter.rb +++ /dev/null @@ -1,99 +0,0 @@ -require 'arbiter' -require 'ffi-rzmq' -require 'multi_json' - -class ZeromqArbiter < Arbiter - - class << self - attr_accessor :frontend, :logger - end - - def self.publish(message, metadata) - context = ZMQ::Context.new - - outbound = context.socket(ZMQ::PUSH) - outbound.connect(frontend) - - outbound.send_string( - MultiJson.dump( - :message => message, - :metadata => metadata - ) - ) - - outbound.close - context.terminate - end - - def listen(proxy) - raise 'Must provide proxy location!' unless proxy - - ctx = ZMQ::Context.new - socket = ctx.socket(ZMQ::PULL) - rc = socket.connect(proxy) - - raise "Could not connect to #{proxy}!" unless rc == 0 - - log :info, "Connected to #{proxy}" - - while true - msg = '' - rc = socket.recv_string(msg) - - if error_check(rc) - break - else - process_message(msg) - end - end - - socket.close - ctx.terminate - end - - protected - - def log(type, msg) - if self.class.logger - self.class.logger.send(type, msg) - end - end - - def process_message(message) - message = MultiJson.decode(message) - log :info, "Processing: #{message}" - - begin - self.class.perform(message['message'].to_sym, symbolize_nested_keys(message['metadata'])) - rescue Exception => e - log :error, e - end - end - - def error_check(rc) - if ZMQ::Util.resultcode_ok?(rc) - false - else - log :error, "Operation failed, errno [#{ZMQ::Util.errno}] description [#{ZMQ::Util.error_string}]" - caller(1).each { |callstack| log :error, callstack } - true - end - end - - def symbolize_nested_keys(data) - case data - when Hash - data.map {|k, v| - {k.to_sym => symbolize_nested_keys(v)} - }.inject({}) { |coll, symbol_key_hash| - coll.merge(symbol_key_hash) - } - when Array - data.collect { |a| - symbolize_nested_keys(a) - } - else - data - end - end -end diff --git a/spec/zeromq/majordomo/asynchronous_arbiter_spec.rb b/spec/zeromq/majordomo/asynchronous_arbiter_spec.rb new file mode 100644 index 0000000..12c701c --- /dev/null +++ b/spec/zeromq/majordomo/asynchronous_arbiter_spec.rb @@ -0,0 +1,55 @@ +require 'zeromq/majordomo/asynchronous_arbiter' + +describe Zeromq::Majordomo::AsynchronousArbiter do + let(:response_object) { {:foo => :bar} } + + let(:socket_uri) { 'inproc://server' } + let(:zmq_context) { double(:context, socket: socket) } + let(:version) { 'some-version' } + + let(:socket) { double(:socket) } + + subject { described_class.new(socket_uri, zmq_context, 5, version) } + + context 'successful connect' do + before :each do + zmq_context.should_receive(:socket).with(ZMQ::DEALER).and_return(socket) + socket.should_receive(:setsockopt).at_least(:once).with(ZMQ::LINGER, -1).and_return(0) + socket.should_receive(:setsockopt).at_least(:once).with(ZMQ::SNDTIMEO, 5000).and_return(0) + socket.should_receive(:connect).at_least(:once).with(socket_uri).and_return(0) + + socket.should_receive(:disconnect).at_least(:once).with(socket_uri).and_return(0) + socket.should_receive(:close).at_least(:once).and_return(0) + end + + it 'sends an MDP client request on the reply socket' do + socket.should_receive(:send_strings).with([ + '', 'MDPC01', version, 'rpc-method', MultiJson.dump(:foo => :body) + ]).and_return(0) + + subject.publish('rpc-method', :foo => :body) + end + + it 'sends twice without blocking' do + socket.should_receive(:send_strings).with([ + '', 'MDPC01', version, 'some_method', MultiJson.dump(:some_arg) + ]).at_least(2).times.and_return(0) + + subject.publish(:some_method, :some_arg) + subject.publish(:some_method, :some_arg) + end + + it 'works with a Symbol method' do + socket.should_receive(:send_strings).with([ + '', 'MDPC01', version, 'some_method', MultiJson.dump(:some_arg) + ]).times.and_return(0) + + subject.publish(:some_method, :some_arg) + end + end + + it 'does not attempt disconnect if socket connection failed' do + zmq_context.should_receive(:socket).and_raise(StandardError) + lambda { subject.publish(:some_method, :some_arg) }.should raise_error(StandardError) + end +end diff --git a/spec/zeromq_arbiter_spec.rb b/spec/zeromq_arbiter_spec.rb deleted file mode 100644 index 555d44d..0000000 --- a/spec/zeromq_arbiter_spec.rb +++ /dev/null @@ -1,104 +0,0 @@ -require 'zeromq_arbiter' - -describe ZeromqArbiter do - let(:message) { :test } - let(:metadata) { { :the => :data } } - - let(:zmq_context) { double(:zmq_context) } - let(:socket) { double(:zmq_socket) } - - let(:frontend) { 'tcp://0.0.0.0:9000' } - - subject { - described_class.frontend = frontend - - described_class - } - - before :each do - ZMQ::Context.stub(:new => zmq_context) - end - - it 'sends a message to a zmq socket' do - zmq_context.should_receive(:socket).with(ZMQ::PUSH).and_return(socket) - zmq_context.should_receive(:terminate) - socket.should_receive(:connect).with(frontend) - socket.should_receive(:close) - - socket.should_receive(:send_string).with(MultiJson.dump(:message => message, :metadata => metadata)) - subject.publish(message, metadata) - end - - context :listen do - subject { described_class.new } - - it 'receives a message and pushes it to the listener classes' do - zmq_context.should_receive(:socket).with(ZMQ::PULL).and_return(socket) - zmq_context.should_receive(:terminate) - socket.should_receive(:connect).with(frontend).and_return(0) - socket.should_receive(:close) - - msg = '' - return_values = [0, -1].to_enum - socket.stub(:recv_string) do |msg| - msg.concat(MultiJson.dump(:message => message, :metadata => metadata)) - return_values.next - end - - described_class.should_receive(:perform).with(message, {:the => 'data'}) - - subject.listen(frontend) - end - - it 'raises an error if connect fails' do - zmq_context.should_receive(:socket).with(ZMQ::PULL).and_return(socket) - # FIXME zmq_context.should_receive(:terminate) - socket.should_receive(:connect).with(frontend).and_return(-1) - - lambda { subject.listen(frontend) }.should raise_error {|e| - e.message.should == "Could not connect to #{frontend}!" - } - end - - it 'logs an error if recieving fails' do - logger = double(:logger, :info => true) - described_class.logger = logger - - logger.should_receive(:error).any_number_of_times - - zmq_context.should_receive(:socket).with(ZMQ::PULL).and_return(socket) - zmq_context.should_receive(:terminate) - socket.should_receive(:connect).with(frontend).and_return(0) - socket.should_receive(:close) - socket.stub(:recv_string).and_return(-1) - - subject.listen(frontend) - end - - it 'logs an error if task perform raises an exception' do - class PerformError < Exception; end - - logger = double(:logger, :info => true, :error => true) - described_class.logger = logger - - logger.should_receive(:error) do |e| - e.should be_a PerformError - end - - zmq_context.should_receive(:socket).with(ZMQ::PULL).and_return(socket) - zmq_context.should_receive(:terminate) - socket.should_receive(:connect).with(frontend).and_return(0) - socket.should_receive(:close) - - return_values = [0, -1].to_enum - socket.stub(:recv_string) do |msg| - msg.concat(MultiJson.dump(:message => message, :metadata => metadata)) - return_values.next - end - - described_class.stub(:perform).and_raise(PerformError) - - subject.listen(frontend) - end - end -end