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
701
702
703
704
# 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)
  # KCP 默认 MTU 1400,但外层还要加 mnet header(42) + GCM nonce/tag(28)。
  # 按 IP_MTU(1280) 反推,KCP 包最多 1280 - 42 - 28 = 1210,留点余量取 1200,
  # 保证「KCP 包 + 封装」不超路径 MTU,避免被 IP 分片(分片常被 NAT/防火墙丢)。
  @engine.mtu = 1200

  @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



739
740
741
742
743
744
745
746
747
748
749
750
# File 'lib/mnet.rb', line 739

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)


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

def closed?
  @closed
end

#eof? ⇒ Boolean

Returns:

  • (Boolean)


714
715
716
# File 'lib/mnet.rb', line 714

def eof?
  @eof
end

#established? ⇒ Boolean

Returns:

  • (Boolean)


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

def established?
  @state == :established
end

#handle_packet(pkt, addr) ⇒ Object



752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
# File 'lib/mnet.rb', line 752

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



732
733
734
735
736
737
# File 'lib/mnet.rb', line 732

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

#reanchor ⇒ Object



795
796
797
798
799
# File 'lib/mnet.rb', line 795

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

#send_syn ⇒ Object



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

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

#tick(now) ⇒ Object



769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
# File 'lib/mnet.rb', line 769

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



789
790
791
792
793
# File 'lib/mnet.rb', line 789

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

#write(data) ⇒ Object



718
719
720
721
722
723
724
725
726
727
728
729
730
# File 'lib/mnet.rb', line 718

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