Class: Libuv::Reactor

Inherits:
Object show all
Extended by:
Accessors, ClassMethods
Includes:
Assertions, Resource
Defined in:
lib/libuv/reactor.rb

Defined Under Namespace

Modules: ClassMethods

Constant Summary collapse

REACTORS =
::Concurrent::Map.new
CRITICAL =
::Mutex.new

Constants included from Accessors

Accessors::Functions

Constants included from Assertions

Assertions::MSG_NO_PROC

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from Accessors

reactor

Methods included from ClassMethods

create, current, default, new

Methods included from Resource

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

Methods included from Assertions

#assert_block, #assert_boolean, #assert_type

Constructor Details

#initialize(pointer) ⇒ Reactor

Initialize a reactor using an FFI::Pointer to a libuv reactor



57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/libuv/reactor.rb', line 57

def initialize(pointer) # :notnew:
    @pointer = pointer
    @reactor = self
    @run_count = 0
    @ref_count = 0

    # Create an async call for scheduling work from other threads
    @run_queue = Queue.new
    @process_queue = @reactor.async method(:process_queue_cb)
    @process_queue.unref

    # Create a next tick timer
    @next_tick = @reactor.timer method(:next_tick_cb)
    @next_tick.unref

    # Create an async call for ending the reactor
    @stop_reactor = @reactor.async method(:stop_cb)
    @stop_reactor.unref

    # Libuv can prevent the application shutting down once the main thread has ended
    # The addition of a prepare function prevents this from happening.
    @reactor_prep = Libuv::Prepare.new(@reactor, method(:noop))
    @reactor_prep.unref
    @reactor_prep.start

    # LibUV ingnores program interrupt by default.
    # We provide normal behaviour and allow this to be overriden
    @on_signal = proc { stop_cb }
    sig_callback = method(:signal_cb)
    self.signal(:INT, sig_callback).unref
    self.signal(:HUP, sig_callback).unref
    self.signal(:TERM, sig_callback).unref

    # Notify of errors
    @throw_on_exit = nil
    @reactor_notify = proc { |error|
        @throw_on_exit = error
    }
end

Instance Attribute Details

#reactor_thread ⇒ Thread (readonly)

Exposed to allow joining on the thread, when run in a multithreaded environment. Performing other actions on the thread has undefined semantics (read: a dangerous endevor).

Returns:

  • (Thread)


516
517
518
# File 'lib/libuv/reactor.rb', line 516

def reactor_thread
  @reactor_thread
end

#run_count ⇒ Object (readonly)

Returns the value of attribute run_count.



97
98
99
# File 'lib/libuv/reactor.rb', line 97

def run_count
  @run_count
end

Instance Method Details

#all(*promises) ⇒ ::Libuv::Q::Promise

Combines multiple promises into a single promise that is resolved when all of the input promises are resolved. (thread safe)

Parameters:

  • *promises (::Libuv::Q::Promise) —

    a number of promises that will be combined into a single promise

Returns:

  • (::Libuv::Q::Promise) —

    Returns a single promise that will be resolved with an array of values, each value corresponding to the promise at the same index in the promises array. If any of the promises is resolved with a rejection, this resulting promise will be resolved with the same rejection.



242
243
244
# File 'lib/libuv/reactor.rb', line 242

def all(*promises)
    Q.all(@reactor, *promises)
end

#any(*promises) ⇒ ::Libuv::Q::Promise

Combines multiple promises into a single promise that is resolved when any of the input promises are resolved.

Parameters:

  • *promises (::Libuv::Q::Promise) —

    a number of promises that will be combined into a single promise

Returns:



252
253
254
# File 'lib/libuv/reactor.rb', line 252

def any(*promises)
    Q.any(@reactor, *promises)
end

#async(callback = nil, &block) ⇒ ::Libuv::Async

Get a new Async handle

Returns:



370
371
372
373
374
375
# File 'lib/libuv/reactor.rb', line 370

def async(callback = nil, &block)
    callback ||= block
    handle = Async.new(@reactor)
    handle.progress callback if callback
    handle
end

#check(callback = nil, &blk) ⇒ ::Libuv::Check

Get a new Check handle

Returns:



355
356
357
# File 'lib/libuv/reactor.rb', line 355

def check(callback = nil, &blk)
    Check.new(@reactor, callback || blk)
end

#defer ⇒ ::Libuv::Q::Deferred

Creates a deferred result object for where the result of an operation may only be returned at some point in the future or is being processed on a different thread (thread safe)



230
231
232
# File 'lib/libuv/reactor.rb', line 230

def defer
    Q.defer(@reactor)
end

#file(path, flags = 0, mode: 0, **opts, &blk) ⇒ ::Libuv::File

Opens a file and returns an object that can be used to manipulate it

Parameters:

  • path (String) —

    the path to the file or folder for watching

  • flags (Integer) (defaults to: 0) —

    see ruby File::Constants

  • mode (Integer) (defaults to: 0)

Returns:



437
438
439
440
441
442
# File 'lib/libuv/reactor.rb', line 437

def file(path, flags = 0, mode: 0, **opts, &blk)
    assert_type(String, path, "path must be a String")
    assert_type(Integer, flags, "flags must be an Integer")
    assert_type(Integer, mode, "mode must be an Integer")
    File.new(@reactor, path, flags, mode: mode, **opts, &blk)
end

#filesystem ⇒ ::Libuv::Filesystem

Returns an object for manipulating the filesystem

Returns:



447
448
449
# File 'lib/libuv/reactor.rb', line 447

def filesystem
    Filesystem.new(@reactor)
end

#finally(*promises) ⇒ ::Libuv::Q::Promise

Combines multiple promises into a single promise that is resolved when all of the input promises are resolved or rejected.

Parameters:

  • *promises (::Libuv::Q::Promise) —

    a number of promises that will be combined into a single promise

Returns:

  • (::Libuv::Q::Promise) —

    Returns a single promise that will be resolved with an array of values, each [result, wasResolved] value pair corresponding to a at the same index in the promises array.



263
264
265
# File 'lib/libuv/reactor.rb', line 263

def finally(*promises)
    Q.finally(@reactor, *promises)
end

#fs_event(path) ⇒ ::Libuv::FSEvent

Get a new FSEvent instance

Parameters:

  • path (String) —

    the path to the file or folder for watching

Returns:

Raises:

  • (ArgumentError) —

    if path is not a string



426
427
428
429
# File 'lib/libuv/reactor.rb', line 426

def fs_event(path)
    assert_type(String, path)
    FSEvent.new(@reactor, path)
end

#handle ⇒ Object



152
# File 'lib/libuv/reactor.rb', line 152

def handle; @pointer; end

#idle(callback = nil, &block) ⇒ ::Libuv::Idle

Get a new Idle handle

Parameters:

  • callback (Proc) (defaults to: nil) —

    the callback to be called on idle trigger

Returns:



363
364
365
# File 'lib/libuv/reactor.rb', line 363

def idle(callback = nil, &block)
    Idle.new(@reactor, callback || block)
end

#inspect ⇒ Object

Overwrite as errors in jRuby can literally hang VM when inspecting as many many classes will reference this class



147
148
149
# File 'lib/libuv/reactor.rb', line 147

def inspect
    "#<#{self.class}:0x#{self.__id__.to_s(16)} NT=#{@run_queue.length}>"
end

#log(error, msg = nil, trace = nil) ⇒ Object

Notifies the reactor there was an event that should be logged

Parameters:

  • error (Exception) —

    the error

  • msg (String|nil) (defaults to: nil) —

    optional context on the error

  • trace (Array<String>) (defaults to: nil) —

    optional additional trace of caller if async



497
498
499
# File 'lib/libuv/reactor.rb', line 497

def log(error, msg = nil, trace = nil)
    @reactor_notify.call(error, msg, trace)
end

#lookup(hostname, hint = :IPv4, port = 9, wait: true, &block) ⇒ ::Libuv::Dns

Lookup a hostname

Parameters:

  • hostname (String) —

    the domain name to lookup

  • port (Integer, String) (defaults to: 9) —

    the service being connected too

  • callback (Proc) —

    the callback to be called on success

Returns:



411
412
413
414
415
416
417
418
419
# File 'lib/libuv/reactor.rb', line 411

def lookup(hostname, hint = :IPv4, port = 9, wait: true, &block)
    dns = Dns.new(@reactor, hostname, port, hint, wait: wait)    # Work is a promise object
    if block_given? || !wait
        dns.then block
        dns
    else
        dns.results
    end
end

#lookup_error(err) ⇒ ::Libuv::Error

Lookup an error code and return is as an error object

Parameters:

  • err (Integer) —

    The error code to look up.

Returns:



287
288
289
290
291
292
293
294
295
296
297
298
299
300
# File 'lib/libuv/reactor.rb', line 287

def lookup_error(err)
    name = ::Libuv::Ext.err_name(err)

    if name
        msg  = ::Libuv::Ext.strerror(err)
        ::Libuv::Error.const_get(name.to_sym).new("#{msg}, #{name}:#{err}")
    else
        # We want a back-trace in this case
        raise "error lookup failed for code #{err}"
    end
rescue Exception => e
    @reactor.log e, 'performing error lookup'
    e
end

#next_tick(callback = nil, &block) ⇒ Object

Queue some work to be processed in the next iteration of the event reactor (thread safe)

Parameters:

  • callback (Proc) (defaults to: nil) —

    the callback to be called on the reactor thread

Raises:

  • (ArgumentError) —

    if block is not given



473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
# File 'lib/libuv/reactor.rb', line 473

def next_tick(callback = nil, &block)
    callback ||= block
    assert_block(callback)

    @run_queue << callback
    if reactor_thread?
        # Create a next tick timer
        if not @next_tick_scheduled
            @next_tick.start(0)
            @next_tick_scheduled = true
            @next_tick.ref
        end
    else
        @process_queue.call
    end

    self
end

#notifier(callback = nil, &blk) ⇒ ::Libuv::Q::Promise

Provides a promise notifier for receiving un-handled exceptions

Returns:



221
222
223
224
# File 'lib/libuv/reactor.rb', line 221

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

#now ⇒ Fixnum

Get current time in milliseconds

Returns:

  • (Fixnum)


279
280
281
# File 'lib/libuv/reactor.rb', line 279

def now
    ::Libuv::Ext.now(@pointer)
end

#on_program_interrupt(callback = nil, &block) ⇒ Object

Allows user defined behaviour when sig int is received



389
390
391
392
# File 'lib/libuv/reactor.rb', line 389

def on_program_interrupt(callback = nil, &block)
    @on_signal = callback || block
    self
end

#pipe(ipc = false) ⇒ ::Libuv::Pipe

Get a new Pipe instance

Parameters:

  • ipc (true, false) (defaults to: false) —

    indicate if a handle will be used for ipc, useful for sharing tcp socket between processes

Returns:



333
334
335
# File 'lib/libuv/reactor.rb', line 333

def pipe(ipc = false)
    Pipe.new(@reactor, ipc)
end

#prepare(callback = nil, &blk) ⇒ ::Libuv::Prepare

Get a new Prepare handle

Returns:



348
349
350
# File 'lib/libuv/reactor.rb', line 348

def prepare(callback = nil, &blk)
    Prepare.new(@reactor, callback || blk)
end

#reactor_running? ⇒ Boolean Also known as: running?

Tells you whether the Libuv reactor reactor is currently running.

Returns:

  • (Boolean)


521
522
523
# File 'lib/libuv/reactor.rb', line 521

def reactor_running?
    !@reactor_thread.nil?
end

#reactor_thread? ⇒ Boolean

True if the calling thread is the same thread as the reactor.

Returns:

  • (Boolean)


509
510
511
# File 'lib/libuv/reactor.rb', line 509

def reactor_thread?
    @reactor_thread == ::Thread.current
end

#ref ⇒ Object

Prevents the reactor loop from stopping



202
203
204
205
206
207
# File 'lib/libuv/reactor.rb', line 202

def ref
    if reactor_thread? && reactor_running?
        @process_queue.ref if @ref_count == 0
        @ref_count += 1
    end
end

#run(run_type = :UV_RUN_DEFAULT) {|promise| ... } ⇒ Object

Run the actual event reactor. This method will block until the reactor is stopped.

Parameters:

  • run_type (:UV_RUN_DEFAULT, :UV_RUN_ONCE, :UV_RUN_NOWAIT) (defaults to: :UV_RUN_DEFAULT)

Yield Parameters:

  • promise (::Libuv::Q::Promise) —

    Yields a promise that can be used for logging unhandled exceptions on the reactor.



159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
# File 'lib/libuv/reactor.rb', line 159

def run(run_type = :UV_RUN_DEFAULT)
    if @reactor_thread.nil?
        begin
            @reactor_thread = ::Thread.current
            raise 'only one reactor allowed per-thread' if REACTORS[@reactor_thread]

            REACTORS[@reactor_thread] = @reactor
            @throw_on_exit = nil

            if block_given?
                update_time
                ::Fiber.new {
                    begin
                        yield @reactor
                    rescue Exception => e
                        log(e, 'in reactor run block')
                    end
                }.resume
            end
            @run_count += 1
            ::Libuv::Ext.run(@pointer, run_type)  # This is blocking
        ensure
            REACTORS.delete(@reactor_thread)
            @reactor_thread = nil
            @run_queue.clear
        end

        # Raise the last unhandled error to occur on the reactor thread
        raise @throw_on_exit if @throw_on_exit

    elsif block_given?
        if reactor_thread?
            update_time
            yield @reactor
        else
            raise 'reactor already running on another thread'
        end
    end

    @reactor
end

#schedule(callback = nil, &block) ⇒ Object

Schedule some work to be processed on the event reactor as soon as possible (thread safe)

Parameters:

  • callback (Proc) (defaults to: nil) —

    the callback to be called on the reactor thread

Raises:

  • (ArgumentError) —

    if block is not given



455
456
457
458
459
460
461
462
463
464
465
466
467
# File 'lib/libuv/reactor.rb', line 455

def schedule(callback = nil, &block)
    callback ||= block
    assert_block(callback)

    if reactor_thread?
        callback.call
    else
        @run_queue << callback
        @process_queue.call
    end

    self
end

#signal(signum = nil, callback = nil, &block) ⇒ ::Libuv::Signal

Get a new signal handler

Returns:



380
381
382
383
384
385
386
# File 'lib/libuv/reactor.rb', line 380

def signal(signum = nil, callback = nil, &block)
    callback ||= block
    handle = Signal.new(@reactor)
    handle.progress callback if callback
    handle.start(signum) if signum
    handle
end

#stop ⇒ Object

Closes handles opened by the reactor class and completes the current reactor iteration (thread safe)



502
503
504
# File 'lib/libuv/reactor.rb', line 502

def stop
    @stop_reactor.call
end

#tcp(callback = nil, &blk) ⇒ ::Libuv::TCP

Get a new TCP instance

Returns:



305
306
307
308
# File 'lib/libuv/reactor.rb', line 305

def tcp(callback = nil, &blk)
    callback ||= blk
    TCP.new(@reactor, progress: callback)
end

#timer(callback = nil, &blk) ⇒ ::Libuv::Timer

Get a new timer instance

Parameters:

  • callback (Proc) (defaults to: nil) —

    the callback to be called on timer trigger

Returns:



341
342
343
# File 'lib/libuv/reactor.rb', line 341

def timer(callback = nil, &blk)
    Timer.new(@reactor, callback || blk)
end

#tty(fileno, readable = false) ⇒ ::Libuv::TTY

Get a new TTY instance

Parameters:

  • fileno (Integer) —

    Integer file descriptor of a tty device

  • readable (true, false) (defaults to: false) —

    Boolean indicating if TTY is readable

Returns:



323
324
325
326
327
# File 'lib/libuv/reactor.rb', line 323

def tty(fileno, readable = false)
    assert_type(Integer, fileno, "io#fileno must return an integer file descriptor, #{fileno.inspect} given")

    TTY.new(@reactor, fileno, readable)
end

#udp(callback = nil, &blk) ⇒ ::Libuv::UDP

Get a new UDP instance

Returns:



313
314
315
316
# File 'lib/libuv/reactor.rb', line 313

def udp(callback = nil, &blk)
    callback ||= blk
    UDP.new(@reactor, progress: callback)
end

#unref ⇒ Object

Allows the reactor loop to stop



210
211
212
213
214
215
# File 'lib/libuv/reactor.rb', line 210

def unref
    if reactor_thread? && reactor_running? && @ref_count > 0
        @ref_count -= 1
        @process_queue.unref if @ref_count == 0
    end
end

#update_time ⇒ Object

forces reactor time update, useful for getting more granular times

Returns:

  • nil



271
272
273
274
# File 'lib/libuv/reactor.rb', line 271

def update_time
    ::Libuv::Ext.update_time(@pointer)
    self
end

#work(callback = nil, &block) ⇒ ::Libuv::Work

Queue some work for processing in the libuv thread pool

Parameters:

  • callback (Proc) (defaults to: nil) —

    the callback to be called in the thread pool

Returns:

Raises:

  • (ArgumentError) —

    if block is not given



399
400
401
402
403
# File 'lib/libuv/reactor.rb', line 399

def work(callback = nil, &block)
    callback ||= block
    assert_block(callback)
    Work.new(@reactor, callback)    # Work is a promise object
end