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
-
#after_read ⇒ Object
Hook overridden by Session to advertise a reopened receive window.
-
#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)). -
#close_bridge ⇒ Object
Close the session-side end of the bridge so the @up pump thread gets EOF.
-
#pump_recv ⇒ Object
Hook overridden by KcpSession to pull decoded bytes out of the engine.
- #read(length = nil, outbuf = nil) ⇒ Object
- #read_nonblock(maxlen, outbuf = nil, exception: true) ⇒ Object
- #readpartial(maxlen, outbuf = nil) ⇒ Object
-
#remote_address ⇒ Object
返回对端真实地址(Addrinfo)。Session/KcpSession 没有内核 socket,远端地址 由 @peer_addr 跟踪(含切网迁移后的新地址),供 LSocket#io.remote_address 等使用。.
- #sync ⇒ Object
- #sync=(_value) ⇒ Object
- #sysread(maxlen, outbuf = nil) ⇒ Object
- #syswrite(data) ⇒ Object
- #to_io ⇒ Object
- #wait_readable(timeout = nil) ⇒ Object
- #wait_writable(_timeout = nil) ⇒ Object
- #write_nonblock(data, exception: true) ⇒ Object
Instance Method Details
#after_read ⇒ Object
Hook overridden by Session to advertise a reopened receive window.
171 172 |
# File 'lib/mnet.rb', line 171 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.
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_bridge ⇒ Object
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_recv ⇒ Object
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
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_address ⇒ Object
返回对端真实地址(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 |
#sync ⇒ Object
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_io ⇒ Object
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 |