forked from sockjs/sockjs-ruby
-
Notifications
You must be signed in to change notification settings - Fork 5
Expand file tree
/
Copy pathtransport.rb
More file actions
359 lines (301 loc) · 9.5 KB
/
Copy pathtransport.rb
File metadata and controls
359 lines (301 loc) · 9.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
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
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
# encoding: utf-8
require "sockjs/session"
require "sockjs/servers/request"
require "sockjs/servers/response"
require 'rack/mount'
module SockJS
class SessionUnavailableError < StandardError
end
#If I had it to do over again, Endpoint wouldn't have subclasses - we'd
#subclass Response and instances of Endpoint would know what kind of Response
#to create for their mount point
class Endpoint
class MethodMap
def initialize(map)
@method_map = map
end
attr_reader :method_map
def call(env)
app = @method_map.fetch(env[Rack::REQUEST_METHOD])
app.call(env)
rescue KeyError
::SockJS.debug "Method not supported!"
[405, {"Allow" => methods_map.keys.join(", ") }, []]
end
end
class MethodNotSupportedApp
def initialize(methods)
@allowed_methods = methods
end
def response
::SockJS.debug "Method not supported! (app)"
@response ||=
[405, {"Allow" => @allowed_methods.join(",")}, []].freeze
end
def call(env)
return response
end
end
#XXX Remove
# @deprecated: See response.rb
CONTENT_TYPES ||= {
plain: "text/plain; charset=UTF-8",
html: "text/html; charset=UTF-8",
javascript: "application/javascript; charset=UTF-8",
event_stream: "text/event-stream; charset=UTF-8"
}
module ClassMethods
def add_routes(route_set, connection, options)
method_catching = Hash.new{|h,k| h[k] = []}
endpoints.each do |endpoint_class|
endpoint_class.add_route(route_set, connection, options)
method_catching[endpoint_class.routing_prefix] << endpoint_class.method
end
method_catching.each_pair do |prefix, methods|
route_set.add_route(MethodNotSupportedApp.new(methods), {:path_info => prefix}, {})
end
end
def routing_prefix
case prefix
when String
"/" + prefix
when Regexp
prefix
end
end
def route_conditions
{
:request_method => self.method,
:path_info => self.routing_prefix
}
end
def add_route(route_set, connection, options)
#SockJS.debug "Adding route for #{self} on #{route_conditions.inspect}"
route_set.add_route(self.new(connection, options), route_conditions, {})
end
def endpoints
@endpoints ||= []
end
def register(method, prefix)
@prefix = prefix
@method = method
Endpoint.endpoints << self
end
attr_reader :prefix, :method
end
extend ClassMethods
# Instance methods.
attr_reader :connection, :options
def initialize(connection, options)
@connection, @options = connection, options
options[:websocket] = true unless options.has_key?(:websocket)
options[:cookie_needed] = true unless options.has_key?(:cookie_needed)
end
def inspect
"<<#{self.class.name} #{options.inspect}>>"
end
def response_class
SockJS::Response
end
# Used for pings.
def empty_string
"\n"
end
def format_frame(session, payload)
raise TypeError.new("Payload must not be nil!") if payload.nil?
"#{payload}\n"
end
#How is this used?
#Thread safety?
attr_reader :remote_addr, :http_origin
def call(env)
@remote_addr = env["REMOTE_ADDR"]
@http_origin = env["HTTP_ORIGIN"]
SockJS.debug "Request for #{self.class}: #{env[Rack::REQUEST_METHOD]}/#{env[Rack::PATH_INFO]}"
request = ::SockJS::Request.new(env)
EM.next_tick do
handle(request)
end
return Thin::Connection::AsyncResponse
end
def handle(request)
handle_request(request)
rescue SockJS::HttpError => error
SockJS.debug "HttpError while handling request: #{([error.inspect] + error.backtrace).join("\n")}"
handle_http_error(request, error)
rescue Object => error
SockJS.debug "Error while handling request: #{([error.inspect] + error.backtrace).join("\n")}"
begin
response = response_class.new(request, 500)
response.write(error.message)
response.finish
return response
rescue Object => ex
SockJS.debug "Error while trying to send error HTTP response: #{ex.inspect}"
end
end
def handle_request(request)
response = build_response(request)
response.finish
return response
end
def error_content_type
:plain
end
def handle_http_error(request, error)
response = build_error_response(request)
response.status = error.status
response.set_content_type(error_content_type)
SockJS::debug "Built error response: #{response.inspect}"
response.write(error.message)
response
end
def build_response(request)
response = response_class.new(request)
setup_response(request, response)
return response
end
def build_error_response(request)
build_response(request)
end
def setup_response(request, response)
response.status = 200
end
end
class SessionEndpoint < Endpoint
def self.routing_prefix
legal_key_regexp = %r{[^./]+}
::Rack::Mount::Strexp.new("/:server_key/:session_key/#{self.prefix}", {:server_key => legal_key_regexp, :session_key => legal_key_regexp})
end
end
class Transport < SessionEndpoint
def handle_request(request)
SockJS.debug({:Request => request, :Transport => self}.inspect)
response = build_response(request)
session = get_session(response)
process_session(session, response)
return response
rescue SockJS::InvalidJSON => error
exception_response(request, error, 500)
rescue SockJS::SessionUnavailableError => error
handle_session_unavailable(request)
end
def response_beginning(response)
end
def exception_response(request, error, status)
SockJS::debug("Handling error #{error.inspect}")
response = build_response(request)
response.status = status
response.set_content_type(:plain)
response.set_session_id(request.session_id)
response.write(error.message)
SockJS::debug("Error response: #{response.inspect}")
return response
end
def handle_session_unavailable(request)
SockJS::debug("Handling missing session for #{request.inspect}")
response = build_response(request)
response.status = 404
response.set_content_type(:plain)
response.set_session_id(request.session_id)
response.write("Session is not open!")
return response
end
def server_key(response)
request = response.request
(request.env['rack.routing_args'] || {})[:server_key]
end
def session_key(response)
request = response.request
(request.env['rack.routing_args'] || {})[:session_key]
end
def request_data(request)
request.data.string
end
end
class ConsumingTransport < Transport
def process_session(session, response)
session.attach_consumer(response, self)
response.request.on_close do
begin
request_closed(session)
rescue Object => ex
SockJS::debug "Exception when closing request: #{ex.inspect}"
end
end
end
def request_closed(session)
session.detach_consumer
end
def finish_response(response)
response.finish
end
def opening_frame(response)
send_data(response, format_frame(response, Protocol::OpeningFrame.instance))
end
def heartbeat_frame(response)
send_data(response, format_frame(response, Protocol::HeartbeatFrame.instance))
end
def messages_frame(response, messages)
send_data(response, format_frame(response, Protocol::ArrayFrame.new(messages)))
end
def closing_frame(response, status, message)
send_data(response, format_frame(response, Protocol::ClosingFrame.new(status, message)))
finish_response(response)
end
#TODO: Consider absorbing format_frame into send_data
def send_data(response, data)
response.write(data)
return data.length
end
def format_frame(response, frame)
frame.to_s + "\n"
end
def get_session(response)
begin
session = connection.get_session(session_key(response))
response_beginning(response)
return session
rescue KeyError
SockJS::debug("Missing session for #{session_key(response)} - creating new")
session = connection.create_session(session_key(response))
response_beginning(response)
opening_frame(response)
return session
end
end
end
class PollingConsumingTransport < ConsumingTransport
def process_session(session, response)
super
#response.finish
end
end
class DeliveryTransport < Transport
def process_session(session, response)
session.receive_message(extract_message(response.request))
successful_response(response)
end
def extract_message(request)
body = request.data.read
raise "Payload expected." if body.empty?
return body
end
def setup_response(response)
response.status = 204
end
def successful_response(response)
response.finish
end
def get_session(response)
begin
session = connection.get_session(session_key(response))
response_beginning(response)
return session
rescue KeyError
SockJS::debug("Missing session for #{session_key(response)} - invalid request")
raise SessionUnavailableError
end
end
end
end