Class: SimpleJob::Client
- Inherits:
-
Object
- Object
- SimpleJob::Client
- Defined in:
- lib/simplejob/client.rb
Direct Known Subclasses
Class Method Summary collapse
Instance Method Summary collapse
- #bind(queue, key) ⇒ Object
- #exchange(name) ⇒ Object
- #orig_send ⇒ Object
- #publish(topic, message) ⇒ Object
- #queue(name) ⇒ Object
- #send(topic, props = {}) ⇒ Object
- #start(opts, &proc) ⇒ Object
- #stop ⇒ Object
- #subscribe(queue, &proc) ⇒ Object
Class Method Details
.start(opts = {}, &proc) ⇒ Object
5 6 7 8 9 |
# File 'lib/simplejob/client.rb', line 5 def self.start(opts = {}, &proc) instance = new instance.start(opts, &proc) instance end |
Instance Method Details
#bind(queue, key) ⇒ Object
36 37 38 |
# File 'lib/simplejob/client.rb', line 36 def bind(queue, key) queue.bind(@exchange, :key => key) end |
#exchange(name) ⇒ Object
28 29 30 |
# File 'lib/simplejob/client.rb', line 28 def exchange(name) @channel.topic(name, :durable => true, :auto_delete => false) end |
#orig_send ⇒ Object
48 |
# File 'lib/simplejob/client.rb', line 48 alias_method :orig_send, :send |
#publish(topic, message) ⇒ Object
40 41 42 |
# File 'lib/simplejob/client.rb', line 40 def publish(topic, ) @exchange.publish(, :routing_key => topic, :persistent => true) end |
#queue(name) ⇒ Object
32 33 34 |
# File 'lib/simplejob/client.rb', line 32 def queue(name) @channel.queue(name, :durable => true, :auto_delete => false) end |
#send(topic, props = {}) ⇒ Object
49 50 51 52 53 |
# File 'lib/simplejob/client.rb', line 49 def send(topic, props = {}) raise "Message properties should be a Hash" unless props.kind_of?(Hash) log.info "[simplejob] New job: #{topic}, #{props.truncated_inspect}" publish(topic, props.to_json) end |
#start(opts, &proc) ⇒ Object
11 12 13 14 15 16 17 18 19 20 21 22 |
# File 'lib/simplejob/client.rb', line 11 def start(opts, &proc) client = self AMQP.start(opts) do Signal.trap("INT") { puts; stop } Signal.trap("TERM") { puts; stop } @channel = AMQP::Channel.new @exchange = exchange(opts[:exchange_name] || DEFAULT_EXCHANGE_NAME) instance_exec(opts, &proc) end end |
#stop ⇒ Object
24 25 26 |
# File 'lib/simplejob/client.rb', line 24 def stop AMQP.stop { EM.stop } end |
#subscribe(queue, &proc) ⇒ Object
44 45 46 |
# File 'lib/simplejob/client.rb', line 44 def subscribe(queue, &proc) queue.subscribe({ :ack => true }, &proc) end |