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.



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

#id ⇒ Object (readonly)

Returns the value of attribute id.



298
299
300
# File 'lib/mnet.rb', line 298

def id
  @id
end

#peer_addr ⇒ Object (readonly)

Returns the value of attribute peer_addr.



298
299
300
# File 'lib/mnet.rb', line 298

def peer_addr
  @peer_addr
end

#state ⇒ Object (readonly)

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

Returns:

  • (Boolean)


354
355
356
# File 'lib/mnet.rb', line 354

def closed?
  @closed
end

#eof? ⇒ Boolean

Returns:

  • (Boolean)


358
359
360
# File 'lib/mnet.rb', line 358

def eof?
  @eof
end

#established? ⇒ Boolean

Returns:

  • (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
      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

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



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