Module: Mnet::SessionIO

Included in:
KcpSession, Session
Defined in:
lib/mnet.rb

Overview

IO-compatible interface so a Session can be passed directly to OpenSSL::SSL::SSLSocket -- this removes the socketpair bridge (and its two thread hops per message), which is the main Ruby-side throughput cost.

Instance Method Summary collapse

Instance Method Details

#after_readObject

Hook overridden by Session to advertise a reopened receive window.



171
172
# File 'lib/mnet.rb', line 171

def after_read
end

#bridgeObject

A real kernel IO (socketpair bridge) wrapping this session, for the few consumers that require a genuine File/IO -- notably OpenSSL::SSL::SSLSocket (its C init does Check_Type(io, T_FILE)). Only used when needed; the session itself is already IO-compatible via the methods above.



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
# File 'lib/mnet.rb', line 267

def bridge
  return @bridge_io if @bridge_io

  app, transport = Mnet.socket_pair
  @bridge_transport = transport
  [app, transport].each do |s|
    s.sync = true
    s.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_NODELAY, 1) rescue nil
  end

  @down = Thread.new do
    begin
      while (data = read)
        transport.write(data)
      end
    rescue IOError, EOFError, Errno::ECONNRESET, Errno::EPIPE
    ensure
      transport.close_write rescue nil
      close rescue nil
    end
  end

  @up = Thread.new do
    begin
      while (data = transport.readpartial(65_536))
        write(data)
      end
    rescue EOFError, IOError, Errno::ECONNRESET, Errno::EPIPE
    ensure
      close rescue nil
    end
  end

  # 桥接 socketpair 本身拿不到远端地址,把 remote_address 委托回 session,
  # 这样 SSLServer#io.remote_address 等仍能取到对端真实 ip:port。
  session = self
  app.define_singleton_method(:remote_address) { session.remote_address }

  @bridge_io = app
end

#close_bridgeObject

Close the session-side end of the bridge so the @up pump thread gets EOF.



257
258
259
260
261
# File 'lib/mnet.rb', line 257

def close_bridge
  return unless @bridge_transport
  @bridge_transport.shutdown rescue nil
  @bridge_transport.close rescue nil
end

#pump_recvObject

Hook overridden by KcpSession to pull decoded bytes out of the engine.



167
168
# File 'lib/mnet.rb', line 167

def pump_recv
end

#read(length = nil, outbuf = nil) ⇒ Object



174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
# File 'lib/mnet.rb', line 174

def read(length = nil, outbuf = nil)
  @m.synchronize do
    loop do
      pump_recv
      unless @recv_buf.empty?
        n = length.nil? ? @recv_buf.bytesize : [length, @recv_buf.bytesize].min
        out = @recv_buf.byteslice(0, n)
        @recv_buf = @recv_buf.byteslice(n, @recv_buf.bytesize - n) || "".b
        after_read
        return outbuf ? outbuf.replace(out) : out
      end
      return nil if @eof || @closed
      @cv.wait
    end
  end
end

#read_nonblock(maxlen, outbuf = nil, exception: true) ⇒ Object



209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
# File 'lib/mnet.rb', line 209

def read_nonblock(maxlen, outbuf = nil, exception: true)
  @m.synchronize do
    pump_recv
    unless @recv_buf.empty?
      n = [maxlen, @recv_buf.bytesize].min
      out = @recv_buf.byteslice(0, n)
      @recv_buf = @recv_buf.byteslice(n, @recv_buf.bytesize - n) || "".b
      after_read
      return outbuf ? outbuf.replace(out) : out
    end
    return nil if @eof || @closed
    raise IO::WaitReadable if exception
    :wait_readable
  end
end

#readpartial(maxlen, outbuf = nil) ⇒ Object

Raises:

  • (ArgumentError)


191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
# File 'lib/mnet.rb', line 191

def readpartial(maxlen, outbuf = nil)
  raise ArgumentError, "non-positive maxlen" if maxlen <= 0
  @m.synchronize do
    loop do
      pump_recv
      unless @recv_buf.empty?
        n = [maxlen, @recv_buf.bytesize].min
        out = @recv_buf.byteslice(0, n)
        @recv_buf = @recv_buf.byteslice(n, @recv_buf.bytesize - n) || "".b
        after_read
        return outbuf ? outbuf.replace(out) : out
      end
      raise EOFError, "end of file reached" if @eof || @closed
      @cv.wait
    end
  end
end

#remote_addressObject

返回对端真实地址(Addrinfo)。Session/KcpSession 没有内核 socket,远端地址 由 @peer_addr 跟踪(含切网迁移后的新地址),供 LSocket#io.remote_address 等使用。



159
160
161
162
163
164
# File 'lib/mnet.rb', line 159

def remote_address
  return nil unless @peer_addr
  Addrinfo.udp(@peer_addr[0], @peer_addr[1])
rescue SocketError
  nil
end

#syncObject



150
151
152
# File 'lib/mnet.rb', line 150

def sync
  true
end

#sync=(_value) ⇒ Object



154
155
# File 'lib/mnet.rb', line 154

def sync=(_value)
end

#sysread(maxlen, outbuf = nil) ⇒ Object



229
230
231
# File 'lib/mnet.rb', line 229

def sysread(maxlen, outbuf = nil)
  readpartial(maxlen, outbuf)
end

#syswrite(data) ⇒ Object



233
234
235
# File 'lib/mnet.rb', line 233

def syswrite(data)
  write(data)
end

#to_ioObject



146
147
148
# File 'lib/mnet.rb', line 146

def to_io
  self
end

#wait_readable(timeout = nil) ⇒ Object



237
238
239
240
241
242
243
244
245
246
247
248
249
250
# File 'lib/mnet.rb', line 237

def wait_readable(timeout = nil)
  @m.synchronize do
    pump_recv
    return true unless @recv_buf.empty?
    return nil if @eof || @closed
    if timeout
      @cv.wait(timeout)
    else
      @cv.wait
    end
    pump_recv # data may have arrived via the engine while we waited
    !@recv_buf.empty?
  end
end

#wait_writable(_timeout = nil) ⇒ Object



252
253
254
# File 'lib/mnet.rb', line 252

def wait_writable(_timeout = nil)
  true
end

#write_nonblock(data, exception: true) ⇒ Object



225
226
227
# File 'lib/mnet.rb', line 225

def write_nonblock(data, exception: true)
  write(data)
end