Class: MessageBus::Redis::ReliablePubSub

Inherits:
Object
  • Object
show all
Defined in:
lib/message_bus/redis/reliable_pub_sub.rb

Defined Under Namespace

Classes: BackLogOutOfOrder, NoMoreRetries

Constant Summary collapse

UNSUB_MESSAGE =
"$$UNSUBSCRIBE"

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(redis_config = {}, max_backlog_size = 1000) ⇒ ReliablePubSub

max_backlog_size is per multiplexed channel



29
30
31
32
33
34
35
36
37
38
39
40
41
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 29

def initialize(redis_config = {}, max_backlog_size = 1000)
  @redis_config = redis_config
  @max_backlog_size = max_backlog_size
  @max_global_backlog_size = 2000
  @max_publish_retries = 10
  @max_publish_wait = 500 #ms
  @max_in_memory_publish_backlog = 1000
  @in_memory_backlog = []
  @lock = Mutex.new
  @flush_backlog_thread = nil
  # after 7 days inactive backlogs will be removed
  @max_backlog_age = 604800
end

Instance Attribute Details

#max_backlog_ageObject

Returns the value of attribute max_backlog_age.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_backlog_age
  @max_backlog_age
end

#max_backlog_sizeObject

Returns the value of attribute max_backlog_size.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_backlog_size
  @max_backlog_size
end

#max_global_backlog_sizeObject

Returns the value of attribute max_global_backlog_size.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_global_backlog_size
  @max_global_backlog_size
end

#max_in_memory_publish_backlogObject

Returns the value of attribute max_in_memory_publish_backlog.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_in_memory_publish_backlog
  @max_in_memory_publish_backlog
end

#max_publish_retriesObject

Returns the value of attribute max_publish_retries.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_publish_retries
  @max_publish_retries
end

#max_publish_waitObject

Returns the value of attribute max_publish_wait.



13
14
15
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 13

def max_publish_wait
  @max_publish_wait
end

#subscribedObject (readonly)

Returns the value of attribute subscribed.



12
13
14
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 12

def subscribed
  @subscribed
end

Instance Method Details

#after_forkObject



47
48
49
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 47

def after_fork
  pub_redis.client.reconnect
end

#backlog(channel, last_id = nil) ⇒ Object



190
191
192
193
194
195
196
197
198
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 190

def backlog(channel, last_id = nil)
  redis = pub_redis
  backlog_key = backlog_key(channel)
  items = redis.zrangebyscore backlog_key, last_id.to_i + 1, "+inf"

  items.map do |i|
    MessageBus::Message.decode(i)
  end
end

#backlog_id_key(channel) ⇒ Object



65
66
67
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 65

def backlog_id_key(channel)
  "__mb_backlog_id_n_#{channel}"
end

#backlog_key(channel) ⇒ Object



61
62
63
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 61

def backlog_key(channel)
  "__mb_backlog_n_#{channel}"
end

#ensure_backlog_flushedObject



150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 150

def ensure_backlog_flushed
  while true
    try_again = false

    @lock.synchronize do
      break if @in_memory_backlog.length == 0

      begin
        publish(*@in_memory_backlog[0],false)
      rescue Redis::CommandError => e
        if e.message =~ /^READONLY/
          try_again = true
        else
          MessageBus.logger.warn("Dropping undeliverable message #{e}")
        end
      rescue => e
        MessageBus.logger.warn("Dropping undeliverable message #{e}")
      end

      @in_memory_backlog.delete_at(0) unless try_again
    end

    if try_again
      sleep 0.005
      # in case we are not connected to the correct server
      # which can happen when sharing ips
      pub_redis.client.reconnect
    end
  end
ensure
  @lock.synchronize do
    @flush_backlog_thread = nil
  end
end

#get_message(channel, message_id) ⇒ Object



218
219
220
221
222
223
224
225
226
227
228
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 218

def get_message(channel, message_id)
  redis = pub_redis
  backlog_key = backlog_key(channel)

  items = redis.zrangebyscore backlog_key, message_id, message_id
  if items && items[0]
    MessageBus::Message.decode(items[0])
  else
    nil
  end
end

#global_backlog(last_id = nil) ⇒ Object



200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 200

def global_backlog(last_id = nil)
  last_id = last_id.to_i
  redis = pub_redis

  items = redis.zrangebyscore global_backlog_key, last_id.to_i + 1, "+inf"

  items.map! do |i|
    pipe = i.index "|"
    message_id = i[0..pipe].to_i
    channel = i[pipe+1..-1]
    m = get_message(channel, message_id)
    m
  end

  items.compact!
  items
end

#global_backlog_keyObject



73
74
75
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 73

def global_backlog_key
  "__mb_global_backlog_n"
end

#global_id_keyObject



69
70
71
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 69

def global_id_key
  "__mb_global_id_n"
end

#global_subscribe(last_id = nil, &blk) ⇒ Object

Raises:

  • (ArgumentError)


278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 278

def global_subscribe(last_id=nil, &blk)
  raise ArgumentError unless block_given?
  highest_id = last_id

  clear_backlog = lambda do
    retries = 4
    begin
      highest_id = process_global_backlog(highest_id, retries > 0, &blk)
    rescue BackLogOutOfOrder => e
      highest_id = e.highest_id
      retries -= 1
      sleep(rand(50) / 1000.0)
      retry
    end
  end


  begin
    @redis_global = new_redis_connection

    if highest_id
      clear_backlog.call(&blk)
    end

    @redis_global.subscribe(redis_channel_name) do |on|
      on.subscribe do
        if highest_id
          clear_backlog.call(&blk)
        end
        @subscribed = true
      end

      on.unsubscribe do
        @subscribed = false
      end

      on.message do |c,m|
        if m == UNSUB_MESSAGE
          @redis_global.unsubscribe
          return
        end
        m = MessageBus::Message.decode m

        # we have 3 options
        #
        # 1. message came in the correct order GREAT, just deal with it
        # 2. message came in the incorrect order COMPLICATED, wait a tiny bit and clear backlog
        # 3. message came in the incorrect order and is lowest than current highest id, reset

        if highest_id.nil? || m.global_id == highest_id + 1
          highest_id = m.global_id
          yield m
        else
          clear_backlog.call(&blk)
        end
      end
    end
  rescue => error
    MessageBus.logger.warn "#{error} subscribe failed, reconnecting in 1 second. Call stack #{error.backtrace}"
    sleep 1
    retry
  end
end

#global_unsubscribeObject



270
271
272
273
274
275
276
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 270

def global_unsubscribe
  if @redis_global
    pub_redis.publish(redis_channel_name, UNSUB_MESSAGE)
    @redis_global.disconnect
    @redis_global = nil
  end
end

#last_id(channel) ⇒ Object



185
186
187
188
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 185

def last_id(channel)
  backlog_id_key = backlog_id_key(channel)
  pub_redis.get(backlog_id_key).to_i
end

#new_redis_connectionObject



43
44
45
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 43

def new_redis_connection
  ::Redis.new(@redis_config)
end

#process_global_backlog(highest_id, raise_error, &blk) ⇒ Object



249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 249

def process_global_backlog(highest_id, raise_error, &blk)
  if highest_id > pub_redis.get(global_id_key).to_i
    highest_id = 0
  end

  global_backlog(highest_id).each do |old|
    if highest_id + 1 == old.global_id
      yield old
      highest_id = old.global_id
    else
      raise BackLogOutOfOrder.new(highest_id) if raise_error
      if old.global_id > highest_id
        yield old
        highest_id = old.global_id
      end
    end
  end

  highest_id
end

#pub_redisObject

redis connection used for publishing messages



57
58
59
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 57

def pub_redis
  @pub_redis ||= new_redis_connection
end

#publish(channel, data, queue_in_memory = true) ⇒ Object



84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 84

def publish(channel, data, queue_in_memory=true)
  redis = pub_redis
  backlog_id_key = backlog_id_key(channel)
  backlog_key = backlog_key(channel)

  global_id = nil
  backlog_id = nil

  redis.multi do |m|
    global_id = m.incr(global_id_key)
    backlog_id = m.incr(backlog_id_key)
  end

  global_id = global_id.value
  backlog_id = backlog_id.value

  msg = MessageBus::Message.new global_id, backlog_id, channel, data
  payload = msg.encode

  redis.multi do |m|

    redis.zadd backlog_key, backlog_id, payload
    redis.expire backlog_key, @max_backlog_age

    redis.zadd global_backlog_key, global_id, backlog_id.to_s << "|" << channel
    redis.expire global_backlog_key, @max_backlog_age

    redis.publish redis_channel_name, payload

    if backlog_id > @max_backlog_size
      redis.zremrangebyscore backlog_key, 1, backlog_id - @max_backlog_size
    end

    if global_id > @max_global_backlog_size
      redis.zremrangebyscore global_backlog_key, 1, global_id - @max_global_backlog_size
    end

  end

  backlog_id

rescue Redis::CommandError => e
  if queue_in_memory &&
        e.message =~ /^READONLY/

    @lock.synchronize do
      @in_memory_backlog << [channel,data]
      if @in_memory_backlog.length > @max_in_memory_publish_backlog
        @in_memory_backlog.delete_at(0)
        MessageBus.logger.warn("Dropping old message cause max_in_memory_publish_backlog is full")
      end
    end

    if @flush_backlog_thread == nil
      @lock.synchronize do
        if @flush_backlog_thread == nil
          @flush_backlog_thread = Thread.new{ensure_backlog_flushed}
        end
      end
    end
    nil
  else
    raise
  end
end

#redis_channel_nameObject



51
52
53
54
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 51

def redis_channel_name
  db = @redis_config[:db] || 0
  "_message_bus_#{db}"
end

#reset!Object

use with extreme care, will nuke all of the data



78
79
80
81
82
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 78

def reset!
  pub_redis.keys("__mb_*").each do |k|
    pub_redis.del k
  end
end

#subscribe(channel, last_id = nil) ⇒ Object

Raises:

  • (ArgumentError)


230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 230

def subscribe(channel, last_id = nil)
  # trivial implementation for now,
  #   can cut down on connections if we only have one global subscriber
  raise ArgumentError unless block_given?

  if last_id
    # we need to translate this to a global id, at least give it a shot
    #   we are subscribing on global and global is always going to be bigger than local
    #   so worst case is a replay of a few messages
    message = get_message(channel, last_id)
    if message
      last_id = message.global_id
    end
  end
  global_subscribe(last_id) do |m|
    yield m if m.channel == channel
  end
end