Class: Mnet::KcpSession

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

Overview

Session with the reliable byte stream provided by the KCP (C) engine instead of the pure-Ruby ARQ. The migration layer (session token, peer address update, control handshake, socketpair bridge, optional AES-GCM) is identical to Session.

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from SessionIO

#after_read, #bridge, #close_bridge, #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) ⇒ KcpSession

Returns a new instance of KcpSession.



674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
# File 'lib/mnet.rb', line 674

def initialize(endpoint, id, role:, peer_addr: nil, key: nil, logger: nil, **opts)
  @endpoint = endpoint
  @id       = id
  @role     = role
  @peer_addr = peer_addr
  @key      = key
  @logger   = logger

  @conv   = id.unpack1("N") # KCP conv derived from the session id
  @engine = Kcp::Engine.new(@conv)

  @m  = Monitor.new
  @cv = @m.new_cond

  @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

  @recv_buf = "".b
  @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.



672
673
674
# File 'lib/mnet.rb', line 672

def id
  @id
end

#peer_addr ⇒ Object (readonly)

Returns the value of attribute peer_addr.



672
673
674
# File 'lib/mnet.rb', line 672

def peer_addr
  @peer_addr
end

#state ⇒ Object (readonly)

Returns the value of attribute state.



672
673
674
# File 'lib/mnet.rb', line 672

def state
  @state
end

Instance Method Details

#close ⇒ Object



735
736
737
738
739
740
741
742
743
744
745
746
# File 'lib/mnet.rb', line 735

def close
  @m.synchronize do
    return if @closed
    send_control(Mnet::TYPE_FIN) if @state != :connecting
    @closed = true
    @eof    = true
    @engine.close # free the C engine while holding @m (no race with read)
    @cv.broadcast
    close_bridge
  end
  @endpoint.remove_session(@id)
end

#closed? ⇒ Boolean

Returns:

  • (Boolean)


706
707
708
# File 'lib/mnet.rb', line 706

def closed?
  @closed
end

#eof? ⇒ Boolean

Returns:

  • (Boolean)


710
711
712
# File 'lib/mnet.rb', line 710

def eof?
  @eof
end

#established? ⇒ Boolean

Returns:

  • (Boolean)


702
703
704
# File 'lib/mnet.rb', line 702

def established?
  @state == :established
end

#handle_packet(pkt, addr) ⇒ Object



748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
# File 'lib/mnet.rb', line 748

def handle_packet(pkt, addr)
  @m.synchronize do
    return if @closed
    @last_recv = Mnet.now
    update_peer(addr)

    case pkt.type
    when Mnet::TYPE_SYN    then on_syn(pkt)
    when Mnet::TYPE_SYNACK then on_synack
    when Mnet::TYPE_DATA   then on_data(pkt)
    when Mnet::TYPE_PING   then send_control(Mnet::TYPE_PONG)
    when Mnet::TYPE_PONG   then nil
    when Mnet::TYPE_FIN    then on_fin
    end
  end
end

#pump_recv ⇒ Object



728
729
730
731
732
733
# File 'lib/mnet.rb', line 728

def pump_recv
  return if @closed # engine is freed in close
  while (chunk = @engine.recv)
    @recv_buf << chunk
  end
end

#reanchor ⇒ Object



791
792
793
794
795
# File 'lib/mnet.rb', line 791

def reanchor
  @m.synchronize do
    send_control(Mnet::TYPE_PING) if @state == :established
  end
end

#send_syn ⇒ Object



781
782
783
# File 'lib/mnet.rb', line 781

def send_syn
  @m.synchronize { send_control(Mnet::TYPE_SYN) }
end

#tick(now) ⇒ Object



765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
# File 'lib/mnet.rb', line 765

def tick(now)
  @m.synchronize do
    return if @closed
    @engine.update(now_ms) # timer-based retransmit
    pump_output

    if @state == :connecting
      send_syn if now - @last_send >= @syn_retry
    elsif now - @last_send >= @ping_after
      send_control(Mnet::TYPE_PING)
    end

    teardown if now - @last_recv >= @idle_timeout
  end
end

#wait_established(timeout) ⇒ Object



785
786
787
788
789
# File 'lib/mnet.rb', line 785

def wait_established(timeout)
  @m.synchronize { @cv.wait(timeout) if @state == :connecting }
  raise "connect timed out" unless @state == :established
  self
end

#write(data) ⇒ Object



714
715
716
717
718
719
720
721
722
723
724
725
726
# File 'lib/mnet.rb', line 714

def write(data)
  data = data.to_s.b
  return 0 if data.empty?

  @m.synchronize do
    return 0 if @closed || @eof
    @engine.send(data)
    pump_output
  end
  # 窗口满(有积压)时让出 GVL,避免写线程紧循环饿死处理 ACK 的传输线程。
  Thread.pass if @engine.waitsnd > 0
  data.bytesize
end