Class: Karait::Queue

Inherits:
Object
  • Object
show all
Includes:
Karait
Defined in:
lib/queue.rb

Constant Summary collapse

MESSAGES_READ =
10
NO_OBJECT_FOUND_ERROR =
'No matching object found'

Instance Method Summary collapse

Constructor Details

#initialize(opts = {}) ⇒ Queue

Returns a new instance of Queue.



11
12
13
14
# File 'lib/queue.rb', line 11

def initialize(opts={})
  set_instance_variables opts
  create_mongo_connection
end

Instance Method Details

#delete_messages(messages) ⇒ Object



105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
# File 'lib/queue.rb', line 105

def delete_messages(messages)
  ids = []
  messages.each {|message| ids << message._get_id}
  @queue_collection.update(
      {
          '_id' => {
            '$in' => ids
          }
      },
      {
          '$set' => {
              '_meta.expired' => true
          }
      },
      :multi => true,
      :safe => true
  )
end

#read(opts = {}) ⇒ Object



35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
# File 'lib/queue.rb', line 35

def read(opts={})
  opts = {
    :messages_read => 10,
    :visibility_timeout => -1.0,
    :routing_key => nil
  }.update(opts)
  
  current_time = Time.new.to_f
  messages = []
  
  query = {
      '_meta.expired' => false,
      '_meta.visible_after' => {
        '$lt' => current_time
      }
  }
  if opts[:routing_key] != nil
    query['_meta.routing_key'] = opts[:routing_key]
  else
    query['_meta.routing_key'] = {
      '$exists' => false
    }
  end
  
  update = false
  if opts[:visibility_timeout] > -1.0
    update = {
        '$set' => {
          '_meta.visible_after' => current_time + opts[:visibility_timeout]
        }
    }
  end
  
  raw_messages = []
  
  if update
    (0..opts[:messages_read]).each do
      begin
        
        raw_message = @queue_collection.find_and_modify(:query => query, :update => update)
        
        if raw_message
          raw_messages << raw_message
        else
          break
        end
      
      rescue Mongo::OperationFailure => operation_failure
        if not  operation_failure.to_s.match(Queue::NO_OBJECT_FOUND_ERROR)
          raise operation_failure
        end
      end
      
    end
  else
    @queue_collection.find(query).limit(opts[:messages_read]).each do |raw_message|
      raw_messages << raw_message
    end
  end
  
  raw_messages.each do |raw_message|
    message = Karait::Message.new(raw_message=raw_message, queue_collection=@queue_collection)
    if not message.expired?
      messages << message
    end
  end
  
  return messages
end

#write(message, opts = {}) ⇒ Object



16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
# File 'lib/queue.rb', line 16

def write(message, opts={})
  if message.class == Hash
    message_dict = message
  else
    message_dict = message.to_hash
  end
  
  message_dict[:_meta] = {
    :expire => opts.fetch(:expire, -1.0),
    :timestamp => Time.now().to_f,
    :expired => false,
    :visible_after => -1.0
  }
  
  message_dict[:_meta][:routing_key] = opts.fetch(:routing_key) if opts[:routing_key]
  
  @queue_collection.insert(message_dict, :safe => true)
end