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
| | # -*- encoding: binary -*-
# :enddoc:
class Rainbows::Coolio::Client < Coolio::IO
include Rainbows::EvCore
APP = Rainbows.server.app
CONN = Rainbows::Coolio::CONN
KATO = Rainbows::Coolio::KATO
LOOP = Coolio::Loop.default
def initialize(io)
CONN[self] = false
super(io)
post_init
@deferred = nil
end
def want_more
enable unless enabled?
end
def quit
super
close if nil == @deferred && @_write_buffer.empty?
end
# override the Coolio::IO#write method try to write directly to the
# kernel socket buffers to avoid an extra userspace copy if
# possible.
def write(buf)
if @_write_buffer.empty?
begin
case rv = @_io.kgio_trywrite(buf)
when nil
return enable_write_watcher
when :wait_writable
break # fall through to super(buf)
when String
buf = rv # retry, skb could grow or been drained
end
rescue => e
return handle_error(e)
end while true
end
super(buf)
end
def on_readable
buf = @_io.kgio_tryread(CLIENT_HEADER_BUFFER_SIZE, RBUF)
case buf
when :wait_readable
when nil # eof
close
else
on_read buf
end
rescue Errno::ECONNRESET
close
end
# allows enabling of write watcher even when read watcher is disabled
def evloop
LOOP
end
def next!
attached? or return
@deferred = nil
enable_write_watcher # trigger on_write_complete
end
def timeout?
if nil == @deferred && @_write_buffer.empty?
@_io.shutdown
true
else
false
end
end
# used for streaming sockets and pipes
def stream_response_body(body, io, chunk)
# we only want to attach to the Coolio::Loop belonging to the
# main thread in Ruby 1.9
(chunk ? Rainbows::Coolio::ResponseChunkPipe :
Rainbows::Coolio::ResponsePipe).new(io, self, body).attach(LOOP)
@deferred = true
end
def hijacked
CONN.delete(self)
detach
nil
end
def write_response_path(status, headers, body, alive)
io = body_to_io(body)
st = io.stat
if st.file?
defer_file(status, headers, body, alive, io, st)
elsif st.socket? || st.pipe?
chunk = stream_response_headers(status, headers, alive, body)
return hijacked if nil == chunk
stream_response_body(body, io, chunk)
else
# char or block device... WTF?
write_response(status, headers, body, alive)
end
end
def ev_write_response(status, headers, body, alive)
if body.respond_to?(:to_path)
body = write_response_path(status, headers, body, alive)
else
body = write_response(status, headers, body, alive)
end
return hijacked unless body
return quit unless alive && :close != @state
@state = :headers
end
def app_call input
KATO.delete(self)
disable if enabled?
@env['rack.input'] = input
@env['REMOTE_ADDR'] = @_io.kgio_addr
@env['async.callback'] = method(:write_async_response)
@hp.hijack_setup(@_io)
status, headers, body = catch(:async) {
APP.call(@env.merge!(RACK_DEFAULTS))
}
return hijacked if @hp.hijacked?
(nil == status || -1 == status) ? @deferred = true :
ev_write_response(status, headers, body, @hp.next?)
end
def on_write_complete
case @deferred
when true then return # #next! will clear this bit
when nil # fall through
else
return if stream_file_chunk(@deferred)
close_deferred # EOF, fall through
end
case @state
when :close
close if @_write_buffer.empty?
when :headers
if @buf.empty?
buf = @_io.kgio_tryread(CLIENT_HEADER_BUFFER_SIZE, RBUF) or return close
String === buf and return on_read(buf)
# buf == :wait_readable
unless enabled?
enable
KATO[self] = Time.now
end
else
on_read(''.freeze)
end
end
rescue => e
handle_error(e)
end
def handle_error(e)
close_deferred
if msg = Rainbows::Error.response(e)
@_io.kgio_trywrite(msg) rescue nil
end
@_write_buffer.clear
ensure
quit
end
def close_deferred
if @deferred
begin
@deferred.close if @deferred.respond_to?(:close)
rescue => e
Unicorn.log_error(Rainbows.server.logger,
"closing deferred=#{@deferred.inspect}", e)
end
@deferred = nil
end
end
def on_close
close_deferred
CONN.delete(self)
KATO.delete(self)
end
if IO.method_defined?(:trysendfile)
def defer_file(status, headers, body, alive, io, st)
if r = sendfile_range(status, headers)
status, headers, range = r
body = write_headers(status, headers, alive, body) or return hijacked
range and defer_file_stream(range[0], range[1], io, body)
else
write_headers(status, headers, alive, body) or return hijacked
defer_file_stream(0, st.size, io, body)
end
body
end
def stream_file_chunk(sf) # +sf+ is a Rainbows::StreamFile object
case n = @_io.trysendfile(sf, sf.offset, sf.count)
when Integer
sf.offset += n
return if 0 == (sf.count -= n)
when :wait_writable
return enable_write_watcher
else
return
end while true
end
else
def defer_file(status, headers, body, alive, io, st)
write_headers(status, headers, alive, body) or return hijacked
defer_file_stream(0, st.size, io, body)
body
end
def stream_file_chunk(body)
buf = body.to_io.read(0x4000) and write(buf)
end
end
def defer_file_stream(offset, count, io, body)
@deferred = Rainbows::StreamFile.new(offset, count, io, body)
enable_write_watcher
end
end
|