Class: Mnet::KcpSession
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") @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
#id ⇒ Object
Returns the value of attribute id.
686
687
688
|
# File 'lib/mnet.rb', line 686
def id
@id
end
|
#peer_addr ⇒ Object
Returns the value of attribute peer_addr.
686
687
688
|
# File 'lib/mnet.rb', line 686
def peer_addr
@peer_addr
end
|
#state ⇒ Object
Returns the value of attribute state.
686
687
688
|
# File 'lib/mnet.rb', line 686
def state
@state
end
|
Instance Method Details
#close ⇒ Object
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 @cv.broadcast
close_bridge
end
@endpoint.remove_session(@id)
end
|
#closed? ⇒ Boolean
721
722
723
|
# File 'lib/mnet.rb', line 721
def closed?
@closed
end
|
#eof? ⇒ Boolean
725
726
727
|
# File 'lib/mnet.rb', line 725
def eof?
@eof
end
|
#established? ⇒ Boolean
717
718
719
|
# File 'lib/mnet.rb', line 717
def established?
@state == :established
end
|
#handle_packet(pkt, addr) ⇒ Object
#pump_recv ⇒ Object
754
755
756
757
758
759
|
# File 'lib/mnet.rb', line 754
def pump_recv
return if @closed while (chunk = @engine.recv)
@recv_buf << chunk
end
end
|
#reanchor ⇒ Object
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_syn ⇒ Object
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) 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
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
Thread.pass if @engine.waitsnd > 0
data.bytesize
end
|