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



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

def initialize(redis_config = {}, max_backlog_size = 1000)
  @redis_config = redis_config
  @max_backlog_size = max_backlog_size
  # we can store a lot of messages, since only one queue
  @max_global_backlog_size = 100000
  @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
end

Instance Attribute Details

#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



45
46
47
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 45

def after_fork
  pub_redis.client.reconnect
end

#backlog(channel, last_id = nil) ⇒ Object



181
182
183
184
185
186
187
188
189
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 181

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



63
64
65
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 63

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

#backlog_key(channel) ⇒ Object



59
60
61
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 59

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

#ensure_backlog_flushedObject



141
142
143
144
145
146
147
148
149
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
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 141

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



209
210
211
212
213
214
215
216
217
218
219
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 209

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



191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 191

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



71
72
73
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 71

def global_backlog_key
  "__mb_global_backlog_n"
end

#global_id_keyObject



67
68
69
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 67

def global_id_key
  "__mb_global_id_n"
end

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

Raises:

  • (ArgumentError)


269
270
271
272
273
274
275
276
277
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
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 269

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



261
262
263
264
265
266
267
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 261

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



176
177
178
179
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 176

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

#new_redis_connectionObject



41
42
43
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 41

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

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



240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 240

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



55
56
57
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 55

def pub_redis
  @pub_redis ||= new_redis_connection
end

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



82
83
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
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 82

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.zadd backlog_key, backlog_id, payload
  redis.zadd global_backlog_key, global_id, backlog_id.to_s << "|" << channel

  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

  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



49
50
51
52
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 49

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

#reset!Object

use with extreme care, will nuke all of the data



76
77
78
79
80
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 76

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

#subscribe(channel, last_id = nil) ⇒ Object

Raises:

  • (ArgumentError)


221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
# File 'lib/message_bus/redis/reliable_pub_sub.rb', line 221

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