/
in_http.rb
382 lines (328 loc) · 10.6 KB
/
in_http.rb
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
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
#
# Fluentd
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
module Fluent
class HttpInput < Input
Plugin.register_input('http', self)
include DetachMultiProcessMixin
require 'http/parser'
def initialize
require 'webrick/httputils'
require 'uri'
super
end
EMPTY_GIF_IMAGE = "GIF89a\u0001\u0000\u0001\u0000\x80\xFF\u0000\xFF\xFF\xFF\u0000\u0000\u0000,\u0000\u0000\u0000\u0000\u0001\u0000\u0001\u0000\u0000\u0002\u0002D\u0001\u0000;".force_encoding("UTF-8")
config_param :port, :integer, :default => 9880
config_param :bind, :string, :default => '0.0.0.0'
config_param :body_size_limit, :size, :default => 32*1024*1024 # TODO default
config_param :keepalive_timeout, :time, :default => 10 # TODO default
config_param :backlog, :integer, :default => nil
config_param :add_http_headers, :bool, :default => false
config_param :add_remote_addr, :bool, :default => false
config_param :format, :string, :default => 'default'
config_param :blocking_timeout, :time, :default => 0.5
config_param :cors_allow_origins, :array, :default => nil
config_param :respond_with_empty_img, :bool, :default => false
def configure(conf)
super
m = if @format == 'default'
method(:parse_params_default)
else
@parser = Plugin.new_parser(@format)
@parser.configure(conf)
method(:parse_params_with_parser)
end
(class << self; self; end).module_eval do
define_method(:parse_params, m)
end
end
class KeepaliveManager < Coolio::TimerWatcher
def initialize(timeout)
super(1, true)
@cons = {}
@timeout = timeout.to_i
end
def add(sock)
@cons[sock] = sock
end
def delete(sock)
@cons.delete(sock)
end
def on_timer
@cons.each_pair {|sock,val|
if sock.step_idle > @timeout
sock.close
end
}
end
end
def start
log.debug "listening http on #{@bind}:#{@port}"
lsock = TCPServer.new(@bind, @port)
detach_multi_process do
super
@km = KeepaliveManager.new(@keepalive_timeout)
#@lsock = Coolio::TCPServer.new(@bind, @port, Handler, @km, method(:on_request), @body_size_limit)
@lsock = Coolio::TCPServer.new(lsock, nil, Handler, @km, method(:on_request),
@body_size_limit, @format, log,
@cors_allow_origins)
@lsock.listen(@backlog) unless @backlog.nil?
@loop = Coolio::Loop.new
@loop.attach(@km)
@loop.attach(@lsock)
@thread = Thread.new(&method(:run))
end
end
def shutdown
@loop.watchers.each {|w| w.detach }
@loop.stop
@lsock.close
@thread.join
end
def run
@loop.run(@blocking_timeout)
rescue
log.error "unexpected error", :error=>$!.to_s
log.error_backtrace
end
def on_request(path_info, params)
begin
path = path_info[1..-1] # remove /
tag = path.split('/').join('.')
record_time, record = parse_params(params)
# Skip nil record
if record.nil?
if @respond_with_empty_img
return ["200 OK", {'Content-type'=>'image/gif; charset=utf-8'}, EMPTY_GIF_IMAGE]
else
return ["200 OK", {'Content-type'=>'text/plain'}, ""]
end
end
if @add_http_headers
params.each_pair { |k,v|
if k.start_with?("HTTP_")
record[k] = v
end
}
end
if @add_remote_addr
record['REMOTE_ADDR'] = params['REMOTE_ADDR']
end
time = if param_time = params['time']
param_time = param_time.to_i
param_time.zero? ? Engine.now : param_time
else
record_time.nil? ? Engine.now : record_time
end
rescue
return ["400 Bad Request", {'Content-type'=>'text/plain'}, "400 Bad Request\n#{$!}\n"]
end
# TODO server error
begin
# Support batched requests
if record.is_a?(Array)
mes = MultiEventStream.new
record.each do |single_record|
single_time = single_record.delete("time") || time
mes.add(single_time, single_record)
end
router.emit_stream(tag, mes)
else
router.emit(tag, time, record)
end
rescue
return ["500 Internal Server Error", {'Content-type'=>'text/plain'}, "500 Internal Server Error\n#{$!}\n"]
end
if @respond_with_empty_img
return ["200 OK", {'Content-type'=>'image/gif; charset=utf-8'}, EMPTY_GIF_IMAGE]
else
return ["200 OK", {'Content-type'=>'text/plain'}, ""]
end
end
private
def parse_params_default(params)
record = if msgpack = params['msgpack']
MessagePack.unpack(msgpack)
elsif js = params['json']
JSON.parse(js)
else
raise "'json' or 'msgpack' parameter is required"
end
return nil, record
end
EVENT_RECORD_PARAMETER = '_event_record'
def parse_params_with_parser(params)
if content = params[EVENT_RECORD_PARAMETER]
@parser.parse(content) { |time, record|
raise "Received event is not #{@format}: #{content}" if record.nil?
return time, record
}
else
raise "'#{EVENT_RECORD_PARAMETER}' parameter is required"
end
end
class Handler < Coolio::Socket
def initialize(io, km, callback, body_size_limit, format, log, cors_allow_origins)
super(io)
@km = km
@callback = callback
@body_size_limit = body_size_limit
@next_close = false
@format = format
@log = log
@cors_allow_origins = cors_allow_origins
@idle = 0
@km.add(self)
@remote_port, @remote_addr = *Socket.unpack_sockaddr_in(io.getpeername) rescue nil
end
def step_idle
@idle += 1
end
def on_close
@km.delete(self)
end
def on_connect
@parser = Http::Parser.new(self)
end
def on_read(data)
@idle = 0
@parser << data
rescue
@log.warn "unexpected error", :error=>$!.to_s
@log.warn_backtrace
close
end
def on_message_begin
@body = ''
end
def on_headers_complete(headers)
expect = nil
size = nil
if @parser.http_version == [1, 1]
@keep_alive = true
else
@keep_alive = false
end
@env = {}
@content_type = ""
headers.each_pair {|k,v|
@env["HTTP_#{k.gsub('-','_').upcase}"] = v
case k
when /Expect/i
expect = v
when /Content-Length/i
size = v.to_i
when /Content-Type/i
@content_type = v
when /Connection/i
if v =~ /close/i
@keep_alive = false
elsif v =~ /Keep-alive/i
@keep_alive = true
end
when /Origin/i
@origin = v
end
}
if expect
if expect == '100-continue'
if !size || size < @body_size_limit
send_response_nobody("100 Continue", {})
else
send_response_and_close("413 Request Entity Too Large", {}, "Too large")
end
else
send_response_and_close("417 Expectation Failed", {}, "")
end
end
end
def on_body(chunk)
if @body.bytesize + chunk.bytesize > @body_size_limit
unless closing?
send_response_and_close("413 Request Entity Too Large", {}, "Too large")
end
return
end
@body << chunk
end
def on_message_complete
return if closing?
# CORS check
# ==========
# For every incoming request, we check if we have some CORS
# restrictions and white listed origins through @cors_allow_origins.
unless @cors_allow_origins.nil?
unless @cors_allow_origins.include?(@origin)
send_response_and_close("403 Forbidden", {'Connection' => 'close'}, "")
return
end
end
@env['REMOTE_ADDR'] = @remote_addr if @remote_addr
uri = URI.parse(@parser.request_url)
params = WEBrick::HTTPUtils.parse_query(uri.query)
if @format != 'default'
params[EVENT_RECORD_PARAMETER] = @body
elsif @content_type =~ /^application\/x-www-form-urlencoded/
params.update WEBrick::HTTPUtils.parse_query(@body)
elsif @content_type =~ /^multipart\/form-data; boundary=(.+)/
boundary = WEBrick::HTTPUtils.dequote($1)
params.update WEBrick::HTTPUtils.parse_form_data(@body, boundary)
elsif @content_type =~ /^application\/json/
params['json'] = @body
end
path_info = uri.path
params.merge!(@env)
@env.clear
code, header, body = *@callback.call(path_info, params)
body = body.to_s
if @keep_alive
header['Connection'] = 'Keep-Alive'
send_response(code, header, body)
else
send_response_and_close(code, header, body)
end
end
def on_write_complete
close if @next_close
end
def send_response_and_close(code, header, body)
send_response(code, header, body)
@next_close = true
end
def closing?
@next_close
end
def send_response(code, header, body)
header['Content-length'] ||= body.bytesize
header['Content-type'] ||= 'text/plain'
data = %[HTTP/1.1 #{code}\r\n]
header.each_pair {|k,v|
data << "#{k}: #{v}\r\n"
}
data << "\r\n"
write data
write body
end
def send_response_nobody(code, header)
data = %[HTTP/1.1 #{code}\r\n]
header.each_pair {|k,v|
data << "#{k}: #{v}\r\n"
}
data << "\r\n"
write data
end
end
end
end