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.



688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
# File 'lib/mnet.rb', line 688

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)
  @engine.setmtu(opts.fetch(:kcp_mtu, Mnet.kcp_mtu(@key)))

  @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

#idObject (readonly)

Returns the value of attribute id.



686
687
688
# File 'lib/mnet.rb', line 686

def id
  @id
end

#peer_addrObject (readonly)

Returns the value of attribute peer_addr.



686
687
688
# File 'lib/mnet.rb', line 686

def peer_addr
  @peer_addr
end

#stateObject (readonly)

Returns the value of attribute state.



686
687
688
# File 'lib/mnet.rb', line 686

def state
  @state
end

Instance Method Details

#closeObject



761
762
763
764
765
766
767
768
769
770
771
772
# File 'lib/mnet.rb', line 761

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)


721
722
723
# File 'lib/mnet.rb', line 721

def closed?
  @closed
end

#eof?Boolean

Returns:

  • (Boolean)


725
726
727
# File 'lib/mnet.rb', line 725

def eof?
  @eof
end

#established?Boolean

Returns:

  • (Boolean)


717
718
719
# File 'lib/mnet.rb', line 717

def established?
  @state == :established
end

#handle_packet(pkt, addr) ⇒ Object



774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
# File 'lib/mnet.rb', line 774

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_recvObject



754
755
756
757
758
759
# File 'lib/mnet.rb', line 754

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

#reanchorObject



817
818
819
820
821
# File 'lib/mnet.rb', line 817

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

#send_synObject



807
808
809
# File 'lib/mnet.rb', line 807

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

#tick(now) ⇒ Object



791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
# File 'lib/mnet.rb', line 791

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



811
812
813
814
815
# File 'lib/mnet.rb', line 811

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

#write(data) ⇒ Object



729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
# File 'lib/mnet.rb', line 729

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

  total = data.bytesize
  @m.synchronize do
    return 0 if @closed || @eof

    # ikcp_send 单次最多接受 IKCP_WND_RCV(128) 个分段,超出会返回 -2 把整块
    # 数据静默丢弃。这里按 (mss * (窗口-1)) 切块喂入,避免大块写被丢。
    chunk = @engine.mss * (KCP_WND_RCV - 1)
    offset = 0
    while offset < total
      n = [chunk, total - offset].min
      queued = @engine.send(data.byteslice(offset, n))
      break if queued <= 0
      offset += queued
      pump_output
    end
  end
  # 窗口满(有积压)时让出 GVL,避免写线程紧循环饿死处理 ACK 的传输线程。
  Thread.pass if @engine.waitsnd > 0
  data.bytesize
end