Class: SimpleJob::Client

Inherits:
Object
  • Object
show all
Defined in:
lib/simplejob/client.rb

Direct Known Subclasses

Worker

Class Method Summary collapse

Instance Method Summary collapse

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, message)
  @exchange.publish(message, :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