Class: MessageBus::Backends::Postgres::Client
- Inherits:
-
Object
- Object
- MessageBus::Backends::Postgres::Client
- Defined in:
- lib/message_bus/backends/postgres.rb
Defined Under Namespace
Classes: Listener
Constant Summary collapse
- INHERITED_CONNECTIONS =
[]
Instance Method Summary collapse
- #add(channel, value) ⇒ Object
- #after_fork ⇒ Object
- #backlog(channel, backlog_id) ⇒ Object
- #clear_channel_backlog(channel, backlog_id, num_to_keep) ⇒ Object
- #clear_global_backlog(backlog_id, num_to_keep) ⇒ Object
- #destroy ⇒ Object
- #expire(max_backlog_age) ⇒ Object
-
#expire_all_backlogs! ⇒ Object
use with extreme care, will nuke all of the data.
- #get_value(channel, id) ⇒ Object
- #global_backlog(backlog_id) ⇒ Object
-
#initialize(config) ⇒ Client
constructor
A new instance of Client.
- #max_id(channel = nil) ⇒ Object
- #max_ids(*channels) ⇒ Object
- #publish(channel, data) ⇒ Object
-
#request_reconnect ⇒ Object
Sets a flag rather than closing the connection directly, since only the owning thread may touch a PGconn.
-
#reset! ⇒ Object
Dangerous, drops the message_bus table containing the backlog if it exists.
- #subscribe(channel) {|listener| ... } ⇒ Object
- #unsubscribe ⇒ Object
Constructor Details
#initialize(config) ⇒ Client
Returns a new instance of Client.
46 47 48 49 50 51 52 53 54 55 56 |
# File 'lib/message_bus/backends/postgres.rb', line 46 def initialize(config) @config = config @listening_on = {} @available = [] @allocated = {} @subscribe_connection = nil @reconnect_requested = false @subscribed = false @mutex = Mutex.new @pid = Process.pid end |
Instance Method Details
#add(channel, value) ⇒ Object
58 59 60 |
# File 'lib/message_bus/backends/postgres.rb', line 58 def add(channel, value) hold { |conn| exec_prepared(conn, 'insert_message', [channel, value]) { |r| r.getvalue(0, 0).to_i } } end |
#after_fork ⇒ Object
95 96 97 98 99 100 101 102 |
# File 'lib/message_bus/backends/postgres.rb', line 95 def after_fork sync do @pid = Process.pid INHERITED_CONNECTIONS.concat(@available) @available.clear @listening_on.clear end end |
#backlog(channel, backlog_id) ⇒ Object
79 80 81 82 83 |
# File 'lib/message_bus/backends/postgres.rb', line 79 def backlog(channel, backlog_id) hold do |conn| exec_prepared(conn, 'channel_backlog', [channel, backlog_id]) { |r| r.values.each { |a| a[0] = a[0].to_i } } end || [] end |
#clear_channel_backlog(channel, backlog_id, num_to_keep) ⇒ Object
69 70 71 72 |
# File 'lib/message_bus/backends/postgres.rb', line 69 def clear_channel_backlog(channel, backlog_id, num_to_keep) hold { |conn| exec_prepared(conn, 'clear_channel_backlog', [channel, backlog_id, num_to_keep]) } nil end |
#clear_global_backlog(backlog_id, num_to_keep) ⇒ Object
62 63 64 65 66 67 |
# File 'lib/message_bus/backends/postgres.rb', line 62 def clear_global_backlog(backlog_id, num_to_keep) if backlog_id > num_to_keep hold { |conn| exec_prepared(conn, 'clear_global_backlog', [backlog_id - num_to_keep]) } nil end end |
#destroy ⇒ Object
112 113 114 115 116 117 |
# File 'lib/message_bus/backends/postgres.rb', line 112 def destroy sync do @available.each(&:close) @available.clear end end |
#expire(max_backlog_age) ⇒ Object
74 75 76 77 |
# File 'lib/message_bus/backends/postgres.rb', line 74 def expire(max_backlog_age) hold { |conn| exec_prepared(conn, 'expire', [max_backlog_age]) } nil end |
#expire_all_backlogs! ⇒ Object
use with extreme care, will nuke all of the data
120 121 122 |
# File 'lib/message_bus/backends/postgres.rb', line 120 def expire_all_backlogs! reset! end |
#get_value(channel, id) ⇒ Object
91 92 93 |
# File 'lib/message_bus/backends/postgres.rb', line 91 def get_value(channel, id) hold { |conn| exec_prepared(conn, 'get_message', [channel, id]) { |r| r.getvalue(0, 0) } } end |
#global_backlog(backlog_id) ⇒ Object
85 86 87 88 89 |
# File 'lib/message_bus/backends/postgres.rb', line 85 def global_backlog(backlog_id) hold do |conn| exec_prepared(conn, 'global_backlog', [backlog_id]) { |r| r.values.each { |a| a[0] = a[0].to_i } } end || [] end |
#max_id(channel = nil) ⇒ Object
124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 |
# File 'lib/message_bus/backends/postgres.rb', line 124 def max_id(channel = nil) block = proc do |r| if r.ntuples > 0 r.getvalue(0, 0).to_i else 0 end end if channel hold { |conn| exec_prepared(conn, 'max_channel_id', [channel], &block) } else hold { |conn| exec_prepared(conn, 'max_id', &block) } end end |
#max_ids(*channels) ⇒ Object
140 141 142 143 144 145 146 147 148 149 150 151 152 153 |
# File 'lib/message_bus/backends/postgres.rb', line 140 def max_ids(*channels) block = proc do |pg_result| ids = Array.new(channels.size, 0) pg_result.ntuples.times do |i| channel = pg_result.getvalue(i, 0) max_id = pg_result.getvalue(i, 1) channel_index = channels.index(channel) ids[channel_index] = max_id.to_i end ids end hold { |conn| exec_prepared(conn, 'max_channel_ids', [PG::TextEncoder::Array.new.encode(channels)], &block) } end |
#publish(channel, data) ⇒ Object
155 156 157 |
# File 'lib/message_bus/backends/postgres.rb', line 155 def publish(channel, data) hold { |conn| exec_prepared(conn, 'publish', [channel, data]) } end |
#request_reconnect ⇒ Object
Sets a flag rather than closing the connection directly, since only the owning thread may touch a PGconn. Picked up by #subscribe's wait_for_notify loop, which raises to trigger the retry in MessageBus::Backends::Postgres#global_subscribe.
200 201 202 |
# File 'lib/message_bus/backends/postgres.rb', line 200 def request_reconnect sync { @reconnect_requested = true } end |
#reset! ⇒ Object
Dangerous, drops the message_bus table containing the backlog if it exists.
105 106 107 108 109 110 |
# File 'lib/message_bus/backends/postgres.rb', line 105 def reset! hold do |conn| conn.exec 'DROP TABLE IF EXISTS message_bus' create_table(conn) end end |
#subscribe(channel) {|listener| ... } ⇒ Object
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 184 185 186 187 188 189 190 |
# File 'lib/message_bus/backends/postgres.rb', line 159 def subscribe(channel) obj = Object.new sync do @listening_on[channel] = obj @reconnect_requested = false end listener = Listener.new yield listener conn = @subscribe_connection = raw_pg_connection begin conn.exec "LISTEN #{channel}" listener.do_sub.call while listening_on?(channel, obj) raise ReconnectRequested if reconnect_requested? conn.wait_for_notify(10) do |_, _, payload| break unless listening_on?(channel, obj) listener..call(nil, payload) end end listener.do_unsub.call conn.exec "UNLISTEN #{channel}" nil ensure @subscribe_connection&.close @subscribe_connection = nil end end |
#unsubscribe ⇒ Object
192 193 194 |
# File 'lib/message_bus/backends/postgres.rb', line 192 def unsubscribe sync { @listening_on.clear } end |