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.
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
349
350
351
352
353
354
355
356
357
358
359
360
361
362
|
# File 'lib/mnet.rb', line 314
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.
312
313
314
|
# File 'lib/mnet.rb', line 312
def id
@id
end
|
#peer_addr ⇒ Object
Returns the value of attribute peer_addr.
312
313
314
|
# File 'lib/mnet.rb', line 312
def peer_addr
@peer_addr
end
|
#state ⇒ Object
Returns the value of attribute state.
312
313
314
|
# File 'lib/mnet.rb', line 312
def state
@state
end
|
Instance Method Details
#after_read ⇒ Object
392
393
394
395
|
# File 'lib/mnet.rb', line 392
def after_read
@cv.broadcast
maybe_advertise_window
end
|
#close ⇒ Object
397
398
399
400
401
402
403
404
405
406
407
|
# File 'lib/mnet.rb', line 397
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
368
369
370
|
# File 'lib/mnet.rb', line 368
def closed?
@closed
end
|
#eof? ⇒ Boolean
372
373
374
|
# File 'lib/mnet.rb', line 372
def eof?
@eof
end
|
#established? ⇒ Boolean
364
365
366
|
# File 'lib/mnet.rb', line 364
def established?
@state == :established
end
|
#handle_packet(pkt, addr) ⇒ Object
---- Internal: called by Endpoint threads ----------------------------
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
|
# File 'lib/mnet.rb', line 411
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
478
479
480
481
482
|
# File 'lib/mnet.rb', line 478
def reanchor
@m.synchronize do
send_packet(Mnet::TYPE_PING, 0, @next_exp) if @state == :established
end
end
|
#send_syn ⇒ Object
468
469
470
|
# File 'lib/mnet.rb', line 468
def send_syn
@m.synchronize { send_packet(Mnet::TYPE_SYN, 0, 0) }
end
|
#tick(now) ⇒ Object
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
|
# File 'lib/mnet.rb', line 435
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
472
473
474
475
476
|
# File 'lib/mnet.rb', line 472
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 -------------------------------------------------
378
379
380
381
382
383
384
385
386
387
388
389
390
|
# File 'lib/mnet.rb', line 378
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
|