Class: Firehose::Rack::Consumer::WebSocket::MultiplexingHandler
- Inherits:
-
Handler
- Object
- Handler
- Firehose::Rack::Consumer::WebSocket::MultiplexingHandler
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
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
|