Class: Mnet::Session

Inherits:
Object
  • Object
show all
Includes:
SessionIO
Defined in:
lib/mnet.rb

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  # optional 32-byte session key (AES-256-GCM); nil = plaintext

  @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

  # Send side (byte-offset sequence numbers, like TCP).
  @next_seq       = 0
  @send_buf       = "".b          # bytes queued but not yet transmitted
  @inflight       = []            # ordered list of Frame
  @inflight_bytes = 0
  @outstanding    = 0             # queued + inflight bytes (backpressure)

  @peer_window = 0                # last receive window advertised by peer

  # RTT estimation (RFC 6298 SRTT/RTTVAR), driving RTO.
  @srtt    = nil
  @rttvar  = nil

  # Receive side.
  @next_exp    = 0                # next expected byte offset (cumulative ack)
  @recv_buf    = "".b             # delivered, ordered, not yet read by app
  @reasm       = {}               # seq => payload (out of order)
  @reasm_bytes = 0
  @last_window = 0

  # Liveness.
  @last_recv = Mnet.now
  @last_send = Mnet.now

  @eof    = false
  @closed = false
  @bridge_io = nil
end

Instance Attribute Details

#idObject (readonly)

Returns the value of attribute id.



312
313
314
# File 'lib/mnet.rb', line 312

def id
  @id
end

#peer_addrObject (readonly)

Returns the value of attribute peer_addr.



312
313
314
# File 'lib/mnet.rb', line 312

def peer_addr
  @peer_addr
end

#stateObject (readonly)

Returns the value of attribute state.



312
313
314
# File 'lib/mnet.rb', line 312

def state
  @state
end

Instance Method Details

#after_readObject



392
393
394
395
# File 'lib/mnet.rb', line 392

def after_read
  @cv.broadcast
  maybe_advertise_window
end

#closeObject



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

Returns:

  • (Boolean)


368
369
370
# File 'lib/mnet.rb', line 368

def closed?
  @closed
end

#eof?Boolean

Returns:

  • (Boolean)


372
373
374
# File 'lib/mnet.rb', line 372

def eof?
  @eof
end

#established?Boolean

Returns:

  • (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
      header = Mnet.pack(pkt.session_id, pkt.seq, pkt.ack, pkt.type, pkt.flags, pkt.window, "")
      pkt.payload = decrypt_payload(header, pkt.payload)
      return if pkt.payload.nil? # failed authentication / wrong key
    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

#reanchorObject



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_synObject



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
      # Handshake is not yet reliable: keep retrying SYN until SYNACK.
      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