rainbows.git  about / heads / tags
Unicorn for sleepy apps and slow clients
blob 12c1434182e87d3a6e650ea6369f27bdb90cb271 5779 bytes (raw)
$ git show HEAD:lib/rainbows/coolio/client.rb	# shows this blob on the CLI

  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] = Rainbows.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

git clone https://yhbt.net/rainbows.git