Class: MessageBus::Backends::Postgres::Client

Inherits:
Object
  • Object
show all
Defined in:
lib/message_bus/backends/postgres.rb

Defined Under Namespace

Classes: Listener

Constant Summary collapse

INHERITED_CONNECTIONS =
[]

Instance Method Summary collapse

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

Yields:

  • (listener)


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.do_message.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