Class: Mnet::Session
Instance Attribute Summary collapse
Instance Method Summary
collapse
Methods included from SessionIO
#bridge, #close_bridge, #pump_recv, #read, #read_nonblock, #readpartial, #remote_address, #sync, #sync=, #sysread, #syswrite, #to_io, #wait_readable, #wait_writable, #write_nonblock
Constructor Details
#initialize(endpoint, id, role:, peer_addr: nil, key: nil, logger: nil, **opts) ⇒ Session
Returns a new instance of Session.
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
|
# File 'lib/mnet.rb', line 300
def initialize(endpoint, id, role:, peer_addr: nil, key: nil, logger: nil, **opts)
@endpoint = endpoint
@id = id
@role = role
@peer_addr = peer_addr
@logger = logger
@key = key
@m = Monitor.new
@cv = @m.new_cond
@mss = opts.fetch(:mss, Mnet::DEFAULT_MSS)
@recv_cap = opts.fetch(:recv_capacity, 256 * 1024)
@max_out = opts.fetch(:max_outstanding, 4 * 1024 * 1024)
@rto = opts.fetch(:rto, 0.3)
@ping_after = opts.fetch(:ping_after, 15.0)
@idle_timeout = opts.fetch(:idle_timeout, 60.0)
@syn_retry = opts.fetch(:syn_retry, 0.25)
@state = role == :client ? :connecting : :established
@next_seq = 0
@send_buf = "".b @inflight = [] @inflight_bytes = 0
@outstanding = 0
@peer_window = 0
@srtt = nil
@rttvar = nil
@next_exp = 0 @recv_buf = "".b @reasm = {} @reasm_bytes = 0
@last_window = 0
@last_recv = Mnet.now
@last_send = Mnet.now
@eof = false
@closed = false
@bridge_io = nil
end
|
Instance Attribute Details
#id ⇒ Object
Returns the value of attribute id.
298
299
300
|
# File 'lib/mnet.rb', line 298
def id
@id
end
|
#peer_addr ⇒ Object
Returns the value of attribute peer_addr.
298
299
300
|
# File 'lib/mnet.rb', line 298
def peer_addr
@peer_addr
end
|
#state ⇒ Object
Returns the value of attribute state.
298
299
300
|
# File 'lib/mnet.rb', line 298
def state
@state
end
|
Instance Method Details
#after_read ⇒ Object
378
379
380
381
|
# File 'lib/mnet.rb', line 378
def after_read
@cv.broadcast
maybe_advertise_window
end
|
#close ⇒ Object
383
384
385
386
387
388
389
390
391
392
393
|
# File 'lib/mnet.rb', line 383
def close
@m.synchronize do
return if @closed
send_packet(Mnet::TYPE_FIN, 0, @next_exp) if @state != :connecting
@closed = true
@eof = true
@cv.broadcast
close_bridge
end
@endpoint.remove_session(@id)
end
|
#closed? ⇒ Boolean
354
355
356
|
# File 'lib/mnet.rb', line 354
def closed?
@closed
end
|
#eof? ⇒ Boolean
358
359
360
|
# File 'lib/mnet.rb', line 358
def eof?
@eof
end
|
#established? ⇒ Boolean
350
351
352
|
# File 'lib/mnet.rb', line 350
def established?
@state == :established
end
|
#handle_packet(pkt, addr) ⇒ Object
---- Internal: called by Endpoint threads ----------------------------
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
|
# File 'lib/mnet.rb', line 397
def handle_packet(pkt, addr)
@m.synchronize do
return if @closed
@last_recv = Mnet.now
update_peer(addr)
if @key
= Mnet.pack(pkt.session_id, pkt.seq, pkt.ack, pkt.type, pkt.flags, pkt.window, "")
pkt.payload = decrypt_payload(, pkt.payload)
return if pkt.payload.nil? end
case pkt.type
when Mnet::TYPE_SYN then on_syn(pkt)
when Mnet::TYPE_SYNACK then on_synack(pkt)
when Mnet::TYPE_DATA then on_data(pkt)
when Mnet::TYPE_ACK then apply_ack(pkt.ack, pkt.window)
when Mnet::TYPE_FIN then on_fin
when Mnet::TYPE_PING then send_packet(Mnet::TYPE_PONG, 0, @next_exp)
when Mnet::TYPE_PONG then nil
end
end
end
|
#reanchor ⇒ Object
464
465
466
467
468
|
# File 'lib/mnet.rb', line 464
def reanchor
@m.synchronize do
send_packet(Mnet::TYPE_PING, 0, @next_exp) if @state == :established
end
end
|
#send_syn ⇒ Object
454
455
456
|
# File 'lib/mnet.rb', line 454
def send_syn
@m.synchronize { send_packet(Mnet::TYPE_SYN, 0, 0) }
end
|
#tick(now) ⇒ Object
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
|
# File 'lib/mnet.rb', line 421
def tick(now)
@m.synchronize do
return if @closed
if @state == :connecting
send_packet(Mnet::TYPE_SYN, 0, 0) if now - @last_send >= @syn_retry
return
end
if !@inflight.empty?
f = @inflight.first
if now - f.sent_at >= @rto
if f.retries >= Mnet::MAX_RETRIES
log("giving up after #{f.retries} retries")
teardown
return
end
@rto = [@rto * 2, Mnet::MAX_RTO].min
f.sent_at = now
f.retries += 1
send_packet(Mnet::TYPE_DATA, f.seq, @next_exp, f.payload)
end
elsif now - @last_send >= @ping_after
send_packet(Mnet::TYPE_PING, 0, @next_exp)
end
teardown if now - @last_recv >= @idle_timeout
pump unless @send_buf.empty?
end
end
|
#wait_established(timeout) ⇒ Object
458
459
460
461
462
|
# File 'lib/mnet.rb', line 458
def wait_established(timeout)
@m.synchronize { @cv.wait(timeout) if @state == :connecting }
raise "connect timed out" unless @state == :established
self
end
|
#write(data) ⇒ Object
---- Application API -------------------------------------------------
364
365
366
367
368
369
370
371
372
373
374
375
376
|
# File 'lib/mnet.rb', line 364
def write(data)
data = data.to_s.b
return 0 if data.empty?
@m.synchronize do
@cv.wait_while { !@closed && !@eof && (@outstanding >= @max_out || unsendable?) }
return 0 if @closed || @eof
@send_buf << data
@outstanding += data.bytesize
end
pump
data.bytesize
end
|