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.
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
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 |