Class: Firehose::Rack::Consumer::WebSocket::MultiplexingHandler

Inherits:
Handler
  • Object
show all
Defined in:
lib/firehose/rack/consumer/web_socket.rb

Defined Under Namespace

Classes: Subscription

Instance Method Summary collapse

Methods inherited from Handler

#error, #parse_message, #send_message

Constructor Details

#initialize(ws) ⇒ MultiplexingHandler

Returns a new instance of MultiplexingHandler.



117
118
119
120
121
# File 'lib/firehose/rack/consumer/web_socket.rb', line 117

def initialize(ws)
  super(ws)
  @subscriptions = {}
  subscribe_multiplexed Consumer.multiplex_subscriptions(@req)
end

Instance Method Details

#close(event) ⇒ Object



145
146
147
148
# File 'lib/firehose/rack/consumer/web_socket.rb', line 145

def close(event)
  @subscriptions.each_value(&:close)
  @subscriptions.clear
end

#message(event) ⇒ Object



123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
# File 'lib/firehose/rack/consumer/web_socket.rb', line 123

def message(event)
  msg = parse_message(event)

  if subscriptions = msg[:multiplex_subscribe]
    subscriptions = [subscriptions] unless subscriptions.is_a?(Array)
    return subscribe_multiplexed(subscriptions)
  end

  if channel_names = msg[:multiplex_unsubscribe]
    return unsubscribe(channel_names)
  end

  if msg[:ping] == 'PING'
    Firehose.logger.debug "WS ping received, sending pong"
    return send_message pong: "PONG"
  end
end

#open(event) ⇒ Object



141
142
143
# File 'lib/firehose/rack/consumer/web_socket.rb', line 141

def open(event)
  Firehose.logger.debug "Multiplexing Websocket connected: #{@req.path}"
end

#subscribe(channel_name, last_sequence) ⇒ Object

Subscribe the client to the channel on the server. Asks for the last sequence for clients that reconnect.



163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
# File 'lib/firehose/rack/consumer/web_socket.rb', line 163

def subscribe(channel_name, last_sequence)
  channel      = Server::Channel.new channel_name
  deferrable   = channel.next_messages last_sequence
  subscription = Subscription.new(channel, deferrable)

  @subscriptions[channel_name] = subscription

  deferrable.callback do |messages|
    messages.each do |message|
      send_message(
        channel: channel_name,
        message: message.payload,
        last_sequence: message.sequence
      )
      Firehose.logger.debug "WS sent `#{message.payload}` to `#{channel_name}` with sequence `#{message.sequence}`"
    end
    subscribe channel_name, messages.last.sequence
  end

  deferrable.errback do |e|
    EM.next_tick { raise e.inspect } unless e == :disconnect
  end
end

#subscribe_multiplexed(subscriptions) ⇒ Object



150
151
152
153
154
155
156
157
158
159
# File 'lib/firehose/rack/consumer/web_socket.rb', line 150

def subscribe_multiplexed(subscriptions)
  subscriptions.each do |sub|
    Firehose.logger.debug "Subscribing multiplexed to: #{sub}"

    channel, sequence = sub[:channel], sub[:message_sequence]
    next if channel.nil?

    subscribe(channel, sequence.to_i)
  end
end

#unsubscribe(channel_names) ⇒ Object



187
188
189
190
191
192
193
194
195
# File 'lib/firehose/rack/consumer/web_socket.rb', line 187

def unsubscribe(channel_names)
  Firehose.logger.debug "Unsubscribing from channels: #{channel_names}"
  Array(channel_names).each do |chan|
    if sub = @subscriptions[chan]
      sub.close
      @subscriptions.delete(chan)
    end
  end
end