Skip to content

Instantly share code, notes, and snippets.

@bgentry
Created November 29, 2016 18:53
Show Gist options
  • Save bgentry/5a4592dbbcc398c0ad651c53af7da51f to your computer and use it in GitHub Desktop.
Save bgentry/5a4592dbbcc398c0ad651c53af7da51f to your computer and use it in GitHub Desktop.
ActionCable Sequel Postgres adapter
# forked from:
# https://github.com/rails/rails/blob/a373be9da45d4bee684ea03420212780ec1ef4b1/actioncable/lib/action_cable/subscription_adapter/postgresql.rb
gem 'pg', '~> 0.18'
require 'pg'
require 'thread'
module ActionCable
module SubscriptionAdapter
class PostgresqlSequel < Base # :nodoc:
def initialize(*)
super
@listener = nil
end
def broadcast(channel, payload)
with_connection do |pg_conn|
pg_conn.exec("NOTIFY #{pg_conn.escape_identifier(channel)}, '#{pg_conn.escape_string(payload)}'")
end
end
def subscribe(channel, callback, success_callback = nil)
listener.add_subscriber(channel, callback, success_callback)
end
def unsubscribe(channel, callback)
listener.remove_subscriber(channel, callback)
end
def shutdown
listener.shutdown
end
def with_connection(&block) # :nodoc:
Sequel::Model.db.synchronize do |pg_conn|
unless pg_conn.is_a?(PG::Connection)
raise "Sequel database must be Postgres in order to use the Postgres Sequel ActionCable storage adapter"
end
yield pg_conn
end
end
private
def listener
@listener || @server.mutex.synchronize { @listener ||= Listener.new(self, @server.event_loop) }
end
class Listener < SubscriberMap
def initialize(adapter, event_loop)
super()
@adapter = adapter
@event_loop = event_loop
@queue = Queue.new
@thread = Thread.new do
Thread.current.abort_on_exception = true
listen
end
end
def listen
@adapter.with_connection do |pg_conn|
catch :shutdown do
loop do
until @queue.empty?
action, channel, callback = @queue.pop(true)
case action
when :listen
pg_conn.exec("LISTEN #{pg_conn.escape_identifier channel}")
@event_loop.post(&callback) if callback
when :unlisten
pg_conn.exec("UNLISTEN #{pg_conn.escape_identifier channel}")
when :shutdown
throw :shutdown
end
end
pg_conn.wait_for_notify(1) do |chan, pid, message|
broadcast(chan, message)
end
end
end
end
end
def shutdown
@queue.push([:shutdown])
Thread.pass while @thread.alive?
end
def add_channel(channel, on_success)
@queue.push([:listen, channel, on_success])
end
def remove_channel(channel)
@queue.push([:unlisten, channel])
end
def invoke_callback(*)
@event_loop.post { super }
end
end
end
end
end
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment