Class: Libuv::TCP

Inherits:
Handle show all
Includes:
Net, Stream
Defined in:
lib/libuv/tcp.rb

Defined Under Namespace

Classes: Socket4, Socket6, SocketBase

Constant Summary collapse

TLS_ERROR =
"TLS write failed".freeze

Constants included from Stream

Stream::BACKLOG_ERROR, Stream::CLOSED_HANDLE_ERROR, Stream::STREAM_CLOSED_ERROR, Stream::WRITE_ERROR

Constants included from Net

Net::INET6_ADDRSTRLEN, Net::INET_ADDRSTRLEN, Net::IP_ARGUMENT_ERROR, Net::PORT_ARGUMENT_ERROR

Constants included from Assertions

Assertions::MSG_NO_PROC

Constants inherited from Q::Promise

Q::Promise::MAKE_PROMISE

Instance Attribute Summary collapse

Attributes inherited from Handle

#closed, #reactor, #storage

Attributes inherited from Q::Promise

#trace

Instance Method Summary collapse

Methods included from Stream

#close_write, #flush, included, #listen, #progress, #read, #readable?, #start_read, #stop_read, #try_write, #writable?

Methods inherited from Handle

#active?, #closed?, #closing?, #ref, #unref

Methods included from Assertions

#assert_block, #assert_boolean, #assert_type

Methods included from Resource

#check_result, #check_result!, #resolve, #to_ptr

Methods included from Listener

included

Methods inherited from Q::DeferredPromise

#resolved?, #then

Methods inherited from Q::Promise

#catch, #finally, #progress, #ruby_catch, #value

Constructor Details

#initialize(reactor, acceptor = nil, progress: nil, flags: nil) ⇒ TCP

Returns a new instance of TCP.



25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
# File 'lib/libuv/tcp.rb', line 25

def initialize(reactor, acceptor = nil, progress: nil, flags: nil)
    @reactor = reactor
    @progress = progress

    tcp_ptr = ::Libuv::Ext.allocate_handle_tcp
    error = if flags
        check_result(::Libuv::Ext.tcp_init_ext(reactor.handle, tcp_ptr, flags))
    else
        check_result(::Libuv::Ext.tcp_init(reactor.handle, tcp_ptr))
    end

    if acceptor && error.nil?
        error = check_result(::Libuv::Ext.accept(acceptor, tcp_ptr))
        @connected = true
    else
        @connected = false
    end
    
    super(tcp_ptr, error)
end

Instance Attribute Details

#connected ⇒ Object (readonly)

Returns the value of attribute connected.



18
19
20
# File 'lib/libuv/tcp.rb', line 18

def connected
  @connected
end

#protocol ⇒ Object (readonly)

Returns the value of attribute protocol.



19
20
21
# File 'lib/libuv/tcp.rb', line 19

def protocol
  @protocol
end

Instance Method Details

#bind(ip, port, callback = nil, &blk) ⇒ Object

END TLS Abstraction ------------------



200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
# File 'lib/libuv/tcp.rb', line 200

def bind(ip, port, callback = nil, &blk)
    return self if @closed

    @on_accept = callback || blk
    @on_listen = method(:accept)

    assert_type(String, ip, IP_ARGUMENT_ERROR)
    assert_type(Integer, port, PORT_ARGUMENT_ERROR)

    begin
        @tcp_socket = create_socket(IPAddr.new(ip), port)
        @tcp_socket.bind
    rescue Exception => e
        reject(e)
    end

    self
end

#close ⇒ Object

overwrite the default close to ensure pending writes are rejected



135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
# File 'lib/libuv/tcp.rb', line 135

def close
    return self if @closed

    # Free tls memory
    # Next tick as may recieve data after closing
    if @tls
        @reactor.next_tick do
            @tls.cleanup
        end
    end
    @connected = false

    if @pending_writes
        @pending_writes.each do |deferred, data|
            deferred.reject(TLS_ERROR)
        end
        @pending_writes = nil
    end

    super
end

#close_cb ⇒ Object

Close can be called multiple times



110
111
112
113
114
115
116
117
118
# File 'lib/libuv/tcp.rb', line 110

def close_cb
    if @pending_write
        @pending_write.reject(TLS_ERROR)
        @pending_write = nil
    end

    # Shutdown the stream
    close
end

#connect(ip, port, callback = nil, &blk) ⇒ Object



236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
# File 'lib/libuv/tcp.rb', line 236

def connect(ip, port, callback = nil, &blk)
    return self if @closed

    @callback = callback || blk
    assert_type(String, ip, IP_ARGUMENT_ERROR)
    assert_type(Integer, port, PORT_ARGUMENT_ERROR)
    
    begin
        @tcp_socket = create_socket(IPAddr.new(ip), port)
        @tcp_socket.connect(callback(:on_connect, @tcp_socket.connect_req.address))
    rescue Exception => e
        reject(e)
    end

    if @callback.nil?
        @coroutine = @reactor.defer
        co @coroutine.promise
    end

    self
end

#direct_write ⇒ Object



163
# File 'lib/libuv/tcp.rb', line 163

alias_method :direct_write, :write

#disable_keepalive ⇒ Object



290
291
292
293
294
# File 'lib/libuv/tcp.rb', line 290

def disable_keepalive
    return self if @closed
    check_result ::Libuv::Ext.tcp_keepalive(handle, 0, 0)
    self
end

#disable_nodelay ⇒ Object



278
279
280
281
282
# File 'lib/libuv/tcp.rb', line 278

def disable_nodelay
    return self if @closed
    check_result ::Libuv::Ext.tcp_nodelay(handle, 0)
    self
end

#disable_simultaneous_accepts ⇒ Object



302
303
304
305
306
# File 'lib/libuv/tcp.rb', line 302

def disable_simultaneous_accepts
    return self if @closed
    check_result ::Libuv::Ext.tcp_simultaneous_accepts(handle, 0)
    self
end

#dispatch_cb(data) ⇒ Object

This is clear text data that has been decrypted Same as stream.rb on_read for clear text



90
91
92
93
94
95
96
# File 'lib/libuv/tcp.rb', line 90

def dispatch_cb(data)
    begin
        @progress.call data, self
    rescue Exception => e
        @reactor.log e, 'performing TLS read data callback'
    end
end

#do_shutdown ⇒ Object



186
# File 'lib/libuv/tcp.rb', line 186

alias_method :do_shutdown, :shutdown

#enable_keepalive(delay) ⇒ Object



284
285
286
287
288
# File 'lib/libuv/tcp.rb', line 284

def enable_keepalive(delay)
    return self if @closed                   # The to_i asserts integer
    check_result ::Libuv::Ext.tcp_keepalive(handle, 1, delay.to_i)
    self
end

#enable_nodelay ⇒ Object



272
273
274
275
276
# File 'lib/libuv/tcp.rb', line 272

def enable_nodelay
    return self if @closed
    check_result ::Libuv::Ext.tcp_nodelay(handle, 1)
    self
end

#enable_simultaneous_accepts ⇒ Object



296
297
298
299
300
# File 'lib/libuv/tcp.rb', line 296

def enable_simultaneous_accepts
    return self if @closed
    check_result ::Libuv::Ext.tcp_simultaneous_accepts(handle, 1)
    self
end

#handshake_cb(protocol = nil) ⇒ Object

Push through any pending writes when handshake has completed



64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
# File 'lib/libuv/tcp.rb', line 64

def handshake_cb(protocol = nil)
    @handshake = true
    @protocol = protocol

    writes = @pending_writes
    @pending_writes = nil
    writes.each do |deferred, data|
        @pending_write = deferred
        @tls.encrypt(data)
    end

    begin
        @on_handshake.call(self, protocol) if @on_handshake
    rescue => e
        @reactor.log e, 'performing TLS handshake callback'
    end
end

#on_handshake(callback = nil, &blk) ⇒ Object

Provide a callback once the TLS handshake has completed



83
84
85
86
# File 'lib/libuv/tcp.rb', line 83

def on_handshake(callback = nil, &blk)
    @on_handshake = callback || blk
    self
end

#open(fd, binding = true, callback = nil, &blk) ⇒ Object



219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
# File 'lib/libuv/tcp.rb', line 219

def open(fd, binding = true, callback = nil, &blk)
    return self if @closed

    if binding
        @on_listen = method(:accept)
        @on_accept = callback || blk
    else
        @callback = callback || blk
        @coroutine = @reactor.defer if @callback.nil?
    end
    error = check_result UV.tcp_open(handle, fd)
    reject(error) if error
    co @coroutine.promise if @coroutine

    self
end

#peername ⇒ Object



265
266
267
268
269
270
# File 'lib/libuv/tcp.rb', line 265

def peername
    return [] if @closed
    sockaddr, len = get_sockaddr_and_len
    check_result! ::Libuv::Ext.tcp_getpeername(handle, sockaddr, len)
    get_ip_and_port(::Libuv::Ext::Sockaddr.new(sockaddr), len.get_int(0))
end

#shutdown ⇒ Object



187
188
189
190
191
192
193
194
# File 'lib/libuv/tcp.rb', line 187

def shutdown
    if @pending_writes && @pending_writes.length > 0
        @pending_writes[-1][0].finally method(:do_shutdown)
    else
        do_shutdown
    end
    self
end

#sockname ⇒ Object



258
259
260
261
262
263
# File 'lib/libuv/tcp.rb', line 258

def sockname
    return [] if @closed
    sockaddr, len = get_sockaddr_and_len
    check_result! ::Libuv::Ext.tcp_getsockname(handle, sockaddr, len)
    get_ip_and_port(::Libuv::Ext::Sockaddr.new(sockaddr), len.get_int(0))
end

#start_tls(args = {}) ⇒ Object

TLS Abstraction ----------------------



51
52
53
54
55
56
57
58
59
60
61
# File 'lib/libuv/tcp.rb', line 51

def start_tls(args = {})
    return self unless @connected && @tls.nil?

    args[:verify_peer] = true if @on_verify

    @handshake = false
    @pending_writes = []
    @tls = ::RubyTls::SSL::Box.new(args[:server], self, args)
    @tls.start
    self
end

#tls? ⇒ Boolean

Check if tls active on the socket

Returns:

  • (Boolean)


22
# File 'lib/libuv/tcp.rb', line 22

def tls?; !@tls.nil?; end

#transmit_cb(data) ⇒ Object

We resolve the existing tls write promise with a the real writes promise (a close may have occurred)



100
101
102
103
104
105
106
107
# File 'lib/libuv/tcp.rb', line 100

def transmit_cb(data)
    if @pending_write
        @pending_write.resolve(direct_write(data))
        @pending_write = nil
    else
        direct_write(data)
    end
end

#verify_cb(cert) ⇒ Object



120
121
122
123
124
125
126
127
128
129
130
131
# File 'lib/libuv/tcp.rb', line 120

def verify_cb(cert)
    if @on_verify
        begin
            return @on_verify.call cert
        rescue => e
            @reactor.log e, 'performing TLS verify callback'
            return false
        end
    end

    true
end

#verify_peer(callback = nil, &blk) ⇒ Object

Verify peers will be called for each cert in the chain



158
159
160
161
# File 'lib/libuv/tcp.rb', line 158

def verify_peer(callback = nil, &blk)
    @on_verify = callback || blk
    self
end

#write(data, wait: false) ⇒ Object



164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
# File 'lib/libuv/tcp.rb', line 164

def write(data, wait: false)
    if @tls
        deferred = @reactor.defer
        
        if @handshake
            @pending_write = deferred
            @tls.encrypt(data)
        else
            @pending_writes << [deferred, data]
        end

        if wait
            return deferred.promise if wait == :promise
            co deferred.promise
        end

        self
    else
        direct_write(data, wait: wait)
    end
end