Class: MessageBus::Backends::Postgres
- Defined in:
- lib/message_bus/backends/postgres.rb
Overview
This backend diverges from the standard in Base in the following ways:
- Does not support in-memory buffering of messages on publication
- Does not expire backlogs until they are published to
The Postgres backend stores published messages in a single Postgres table
with only global IDs, and an index on channel name and ID for fast
per-channel lookup. All queries are implemented as prepared statements
to reduce the wire-chatter during use. In addition to storage in the
table, messages are published using pg_notify; this is used for
actively subscribed message_bus servers to consume published messages in
real-time while connected and forward them to subscribers, while catch-up
is performed from the backlog table.
Defined Under Namespace
Classes: Client
Constant Summary collapse
- ReconnectRequested =
Raised inside the subscriber thread to make it drop its LISTEN connection and reconnect. See #request_reconnect.
Class.new(StandardError)
Constants inherited from Base
Base::ConcreteClassMustImplementError, Base::UNSUB_MESSAGE
Instance Attribute Summary
Attributes inherited from Base
#clear_every, #max_backlog_age, #max_backlog_size, #max_global_backlog_size, #max_in_memory_publish_backlog, #subscribed
Class Method Summary collapse
Instance Method Summary collapse
-
#after_fork ⇒ Object
Reconnects to Postgres; used after a process fork, typically triggered by a forking webserver.
-
#backlog(channel, last_id = 0) ⇒ Array<MessageBus::Message>
Get messages from a channel backlog.
-
#destroy ⇒ Object
Closes all open connections to the storage.
-
#expire_all_backlogs! ⇒ Object
abstract
Deletes all backlogs and their data.
-
#get_message(channel, message_id) ⇒ MessageBus::Message?
Get a specific message from a channel.
-
#global_backlog(last_id = 0) ⇒ Array<MessageBus::Message>
Get messages from the global backlog.
-
#global_subscribe(last_id = nil) {|message| ... } ⇒ nil
Subscribe to messages on all channels.
-
#global_unsubscribe ⇒ Object
Causes all subscribers to the bus to unsubscribe, and terminates the local connection.
-
#initialize(config = {}, max_backlog_size = 1000) ⇒ Postgres
constructor
A new instance of Postgres.
-
#last_id(channel) ⇒ Integer
Get the ID of the last message published on a channel.
-
#last_ids(*channels) ⇒ Array<Integer>
Get the ID of the last message published on multiple channels.
-
#publish(channel, data, opts = nil) ⇒ Integer
Publishes a message to a channel.
-
#request_reconnect ⇒ Object
Asks the subscriber to drop its connection and re-establish it, so a wedged connection (a half-open socket, for example) recovers without a process restart.
-
#reset! ⇒ Object
Deletes all message_bus data from the backend.
-
#subscribe(channel, last_id = nil) {|message| ... } ⇒ nil
Subscribe to messages on a particular channel.
Constructor Details
#initialize(config = {}, max_backlog_size = 1000) ⇒ Postgres
Returns a new instance of Postgres.
299 300 301 302 303 304 305 306 307 308 |
# File 'lib/message_bus/backends/postgres.rb', line 299 def initialize(config = {}, max_backlog_size = 1000) @config = config @max_backlog_size = max_backlog_size @max_global_backlog_size = 2000 # after 7 days inactive backlogs will be removed @max_backlog_age = 604800 @clear_every = config[:clear_every] || 1 @mutex = Mutex.new @client = nil end |
Class Method Details
.reset!(config) ⇒ Object
290 291 292 |
# File 'lib/message_bus/backends/postgres.rb', line 290 def self.reset!(config) MessageBus::Postgres::Client.new(config).reset! end |
Instance Method Details
#after_fork ⇒ Object
Reconnects to Postgres; used after a process fork, typically triggered by a forking webserver
312 313 314 |
# File 'lib/message_bus/backends/postgres.rb', line 312 def after_fork client.after_fork end |
#backlog(channel, last_id = 0) ⇒ Array<MessageBus::Message>
Get messages from a channel backlog
363 364 365 366 367 368 369 |
# File 'lib/message_bus/backends/postgres.rb', line 363 def backlog(channel, last_id = 0) items = client.backlog channel, last_id.to_i items.map! do |id, data| MessageBus::Message.new id, id, channel, data end end |
#destroy ⇒ Object
Closes all open connections to the storage.
322 323 324 |
# File 'lib/message_bus/backends/postgres.rb', line 322 def destroy client.destroy end |
#expire_all_backlogs! ⇒ Object
Deletes all backlogs and their data. Does not delete non-backlog data that message_bus may persist, depending on the concrete backend implementation. Use with extreme caution.
327 328 329 |
# File 'lib/message_bus/backends/postgres.rb', line 327 def expire_all_backlogs! client.expire_all_backlogs! end |
#get_message(channel, message_id) ⇒ MessageBus::Message?
Get a specific message from a channel
381 382 383 384 385 386 387 |
# File 'lib/message_bus/backends/postgres.rb', line 381 def (channel, ) if data = client.get_value(channel, ) MessageBus::Message.new , , channel, data else nil end end |
#global_backlog(last_id = 0) ⇒ Array<MessageBus::Message>
Get messages from the global backlog
372 373 374 375 376 377 378 |
# File 'lib/message_bus/backends/postgres.rb', line 372 def global_backlog(last_id = 0) items = client.global_backlog last_id.to_i items.map! do |id, channel, data| MessageBus::Message.new id, id, channel, data end end |
#global_subscribe(last_id = nil) {|message| ... } ⇒ nil
Subscribe to messages on all channels. Each message since the last ID specified will be delivered by yielding to the passed block as soon as it is available. This will block until subscription is terminated.
417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 |
# File 'lib/message_bus/backends/postgres.rb', line 417 def global_subscribe(last_id = nil) raise ArgumentError unless block_given? highest_id = last_id begin # Seed a cursor if we don't already have one, so that a reconnect # (triggered by ReconnectRequested below) always has a replay # point and never treats "no explicit last_id" as "never replay". highest_id ||= client.max_id client.subscribe(postgresql_channel_name) do |on| h = {} on.subscribe do highest_id = process_global_backlog(highest_id) do |m| h[m.global_id] = true yield m end h = nil if h.empty? @subscribed = true end on.unsubscribe do @subscribed = false end on. do |_c, m| if m == UNSUB_MESSAGE @subscribed = false return # rubocop:disable Lint/NonLocalExitFromIterator end m = MessageBus::Message.decode m # If already yielded during the clear backlog when subscribing, # don't yield a duplicate copy. duplicate = h && h.delete(m.global_id) h = nil if h&.empty? unless duplicate highest_id = m.global_id if m.global_id > highest_id yield m end end end rescue => error @subscribed = false @config[:logger].warn "#{error} subscribe failed, reconnecting in 1 second. Call stack\n#{error.backtrace.join("\n")}" sleep 1 retry end end |
#global_unsubscribe ⇒ Object
Causes all subscribers to the bus to unsubscribe, and terminates the local connection. Typically used to reset tests.
401 402 403 404 |
# File 'lib/message_bus/backends/postgres.rb', line 401 def global_unsubscribe client.publish(postgresql_channel_name, UNSUB_MESSAGE) @subscribed = false end |
#last_id(channel) ⇒ Integer
Get the ID of the last message published on a channel
353 354 355 |
# File 'lib/message_bus/backends/postgres.rb', line 353 def last_id(channel) client.max_id(channel) end |
#last_ids(*channels) ⇒ Array<Integer>
Get the ID of the last message published on multiple channels
358 359 360 |
# File 'lib/message_bus/backends/postgres.rb', line 358 def last_ids(*channels) client.max_ids(*channels) end |
#publish(channel, data, opts = nil) ⇒ Integer
:queue_in_memory NOT SUPPORTED
Publishes a message to a channel
333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 |
# File 'lib/message_bus/backends/postgres.rb', line 333 def publish(channel, data, opts = nil) # TODO in memory queue? c = client backlog_id = c.add(channel, data) msg = MessageBus::Message.new backlog_id, backlog_id, channel, data payload = msg.encode c.publish postgresql_channel_name, payload if backlog_id && backlog_id % clear_every == 0 max_backlog_size = (opts && opts[:max_backlog_size]) || self.max_backlog_size max_backlog_age = (opts && opts[:max_backlog_age]) || self.max_backlog_age c.clear_global_backlog(backlog_id, @max_global_backlog_size) c.expire(max_backlog_age) c.clear_channel_backlog(channel, backlog_id, max_backlog_size) end backlog_id end |
#request_reconnect ⇒ Object
Asks the subscriber to drop its connection and re-establish it, so a wedged connection (a half-open socket, for example) recovers without a process restart. The backend's rescue/retry around #global_subscribe performs the actual reconnection.
Called from the keepalive watchdog thread, never from the subscriber thread; implementations must not touch anything requiring single-thread ownership (see the Postgres backend). Backends with no way to detach a stuck connection may leave this as the inherited no-op.
Flags the subscriber thread instead of closing its connection: a
PG::Connection may only be used by its owning thread. The reconnect
therefore takes effect when the subscriber's wait_for_notify call
next returns, so within 10 seconds.
412 413 414 |
# File 'lib/message_bus/backends/postgres.rb', line 412 def request_reconnect client.request_reconnect end |
#reset! ⇒ Object
Deletes all message_bus data from the backend. Use with extreme caution.
317 318 319 |
# File 'lib/message_bus/backends/postgres.rb', line 317 def reset! client.reset! end |
#subscribe(channel, last_id = nil) {|message| ... } ⇒ nil
Subscribe to messages on a particular channel. Each message since the last ID specified will be delivered by yielding to the passed block as soon as it is available. This will block until subscription is terminated.
390 391 392 393 394 395 396 397 398 |
# File 'lib/message_bus/backends/postgres.rb', line 390 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? global_subscribe(last_id) do |m| yield m if m.channel == channel end end |