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_read ⇒ Object

Hook overridden by Session to advertise a reopened receive window.



157
158
# File 'lib/mnet.rb', line 157

def after_read
end

#bridge ⇒ Object

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.



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

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_bridge ⇒ Object

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



243
244
245
246
247
# File 'lib/mnet.rb', line 243

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

#pump_recv ⇒ Object

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



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

def pump_recv
end

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



160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
# File 'lib/mnet.rb', line 160

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



195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
# File 'lib/mnet.rb', line 195

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)


177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
# File 'lib/mnet.rb', line 177

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_address ⇒ Object

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



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

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

#sync ⇒ Object



136
137
138
# File 'lib/mnet.rb', line 136

def sync
  true
end

#sync=(_value) ⇒ Object



140
141
# File 'lib/mnet.rb', line 140

def sync=(_value)
end

#sysread(maxlen, outbuf = nil) ⇒ Object



215
216
217
# File 'lib/mnet.rb', line 215

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

#syswrite(data) ⇒ Object



219
220
221
# File 'lib/mnet.rb', line 219

def syswrite(data)
  write(data)
end

#to_io ⇒ Object



132
133
134
# File 'lib/mnet.rb', line 132

def to_io
  self
end

#wait_readable(timeout = nil) ⇒ Object



223
224
225
226
227
228
229
230
231
232
233
234
235
236
# File 'lib/mnet.rb', line 223

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



238
239
240
# File 'lib/mnet.rb', line 238

def wait_writable(_timeout = nil)
  true
end

#write_nonblock(data, exception: true) ⇒ Object



211
212
213
# File 'lib/mnet.rb', line 211

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