Class: Mnet::Endpoint
- Inherits:
-
Object
- Object
- Mnet::Endpoint
- Defined in:
- lib/mnet.rb
Constant Summary collapse
- TICK_INTERVAL =
0.05
Instance Attribute Summary collapse
-
#socket ⇒ Object
readonly
Returns the value of attribute socket.
Instance Method Summary collapse
- #accept ⇒ Object
- #close ⇒ Object
- #dial(host, port, timeout: 5, key: nil, proto: :mnet) ⇒ Object
-
#hop(local_host = "0.0.0.0") ⇒ Object
(also: #rebind)
mosh-style port hop: open a fresh socket (new source address) and switch to it, keeping the old socket alive briefly to catch delayed packets.
-
#initialize(logger: nil, keys: nil, **opts) ⇒ Endpoint
constructor
A new instance of Endpoint.
- #listen(host = "0.0.0.0", port = 0) ⇒ Object
- #local_addr ⇒ Object
- #remove_session(id) ⇒ Object
- #send_raw(data, ip, port) ⇒ Object
Constructor Details
#initialize(logger: nil, keys: nil, **opts) ⇒ Endpoint
Returns a new instance of Endpoint.
909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 |
# File 'lib/mnet.rb', line 909 def initialize(logger: nil, keys: nil, **opts) @opts = opts @logger = logger @socks = [] # [socket, created_at] pairs; newest last @send_mutex = Mutex.new @sessions = {} @sessions_mutex = Mutex.new @accept = Queue.new @running = false @bound = false # Receive window (flow control): capped dynamically to the kernel's # actual socket buffer size (Linux clamps via net.core.rmem_max). @recv_cap = opts.fetch(:recv_capacity, 256 * 1024) # Pre-shared session keys (optional): id -> key, derived like mosh. @key_by_id = {} (keys || []).each { |k| @key_by_id[Mnet.session_id_from_key(k)] = k } end |
Instance Attribute Details
#socket ⇒ Object (readonly)
Returns the value of attribute socket.
907 908 909 |
# File 'lib/mnet.rb', line 907 def socket @socket end |
Instance Method Details
#accept ⇒ Object
955 956 957 |
# File 'lib/mnet.rb', line 955 def accept @accept.pop end |
#close ⇒ Object
963 964 965 966 |
# File 'lib/mnet.rb', line 963 def close @running = false @socks.each { |s, _| s.close rescue nil } end |
#dial(host, port, timeout: 5, key: nil, proto: :mnet) ⇒ Object
939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 |
# File 'lib/mnet.rb', line 939 def dial(host, port, timeout: 5, key: nil, proto: :mnet) ensure_bound id = key ? Mnet.session_id_from_key(key) : SecureRandom.random_bytes(Mnet::SESSION_ID_LEN) klass = proto == :kcp ? KcpSession : Session sess = klass.new(self, id, role: :client, peer_addr: [host, port], key: key, logger: @logger, **@opts.merge(recv_capacity: @recv_cap)) register(sess) begin sess.send_syn sess.wait_established(timeout) rescue remove_session(id) raise end end |
#hop(local_host = "0.0.0.0") ⇒ Object Also known as: rebind
mosh-style port hop: open a fresh socket (new source address) and switch to it, keeping the old socket alive briefly to catch delayed packets. In real life you do NOT need to call this -- bind the client socket to 0.0.0.0 and the OS re-picks the source address when the route changes; this is the explicit fallback for testing/deterministic roaming.
973 974 975 976 977 978 979 980 981 982 983 984 |
# File 'lib/mnet.rb', line 973 def hop(local_host = "0.0.0.0") new_sock = nil @send_mutex.synchronize do new_sock = UDPSocket.new configure_socket(new_sock) new_sock.bind(local_host, 0) @socks << [new_sock, Mnet.now] end prune_sockets @sessions_mutex.synchronize { @sessions.each_value(&:reanchor) } new_sock end |
#listen(host = "0.0.0.0", port = 0) ⇒ Object
929 930 931 932 933 934 935 936 937 |
# File 'lib/mnet.rb', line 929 def listen(host = "0.0.0.0", port = 0) sock = UDPSocket.new configure_socket(sock) sock.bind(host, port) @socks << [sock, Mnet.now] @bound = true start self end |
#local_addr ⇒ Object
959 960 961 |
# File 'lib/mnet.rb', line 959 def local_addr @socks.last[0].addr end |
#remove_session(id) ⇒ Object
998 999 1000 |
# File 'lib/mnet.rb', line 998 def remove_session(id) @sessions_mutex.synchronize { @sessions.delete(id) } end |
#send_raw(data, ip, port) ⇒ Object
988 989 990 991 992 993 994 995 996 |
# File 'lib/mnet.rb', line 988 def send_raw(data, ip, port) sock = @socks.last[0] @send_mutex.synchronize do sock.send(data, Mnet::MSG_DONTWAIT, ip, port) true rescue IO::WaitWritable, IO::WaitReadable, SystemCallError, IOError false end end |