Class: Takagi::Server::Udp
- Inherits:
-
Object
- Object
- Takagi::Server::Udp
- Defined in:
- lib/takagi/server/udp.rb,
sig/takagi/server/udp.rbs
Overview
UDP server for handling CoAP messages
Instance Method Summary collapse
- #close_socket ⇒ Object
- #fork_worker ⇒ Object
-
#initialize(port: 5683, worker_processes: 2, worker_threads: 2, middleware_stack: nil, router: nil, logger: nil, watcher: nil) ⇒ Udp
constructor
A new instance of Udp.
- #log_boot_details ⇒ Object
-
#run! ⇒ Object
Starts the server with multiple worker processes.
-
#shutdown! ⇒ nil, untyped
Gracefully shuts down all workers.
- #spawn_workers ⇒ Object
- #terminate_workers ⇒ nil, untyped
- #test_environment? ⇒ Boolean
- #worker_configuration ⇒ { port: untyped, socket: untyped, middleware_stack: untyped, sender: untyped, logger: untyped, threads: untyped }
Constructor Details
#initialize(port: 5683, worker_processes: 2, worker_threads: 2, middleware_stack: nil, router: nil, logger: nil, watcher: nil) ⇒ Udp
Returns a new instance of Udp.
11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 |
# File 'lib/takagi/server/udp.rb', line 11 def initialize(port: 5683, worker_processes: 2, worker_threads: 2, middleware_stack: nil, router: nil, logger: nil, watcher: nil) @port = port @worker_processes = worker_processes @worker_threads = worker_threads @middleware_stack = middleware_stack || Takagi::MiddlewareStack.instance @router = router || Takagi::Router.instance @logger = logger || Takagi.logger @watcher = watcher || Takagi::Observer::Watcher.new(interval: 1) Initializer.run! @socket = UDPSocket.new @socket.bind('0.0.0.0', @port) Takagi::Network::UdpSender.instance.setup(socket: @socket) @sender = Takagi::Network::UdpSender.instance end |
Instance Method Details
#close_socket ⇒ Object
103 104 105 106 107 108 109 |
# File 'lib/takagi/server/udp.rb', line 103 def close_socket return unless @socket && !@socket.closed? @socket.close rescue StandardError nil end |
#fork_worker ⇒ Object
84 85 86 87 88 89 |
# File 'lib/takagi/server/udp.rb', line 84 def fork_worker fork do Process.setproctitle('takagi-worker') UdpWorker.new(**worker_configuration).run end end |
#log_boot_details ⇒ Object
71 72 73 74 75 |
# File 'lib/takagi/server/udp.rb', line 71 def log_boot_details @logger.info "Starting Takagi server with #{@worker_processes} processes and #{@worker_threads} threads per process..." @logger.info "Takagi server has version #{Takagi::VERSION} with name '#{Takagi::NAME}'" @logger.debug "run #{@router.all_routes}" end |
#run! ⇒ Object
Starts the server with multiple worker processes
30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 |
# File 'lib/takagi/server/udp.rb', line 30 def run! Takagi::Hooks.emit(:server_starting, protocol: :udp, port: @port) log_boot_details spawn_workers Takagi::Observable::Registry.start_all @watcher.start # Set flag instead of calling shutdown! directly from trap context # This avoids "can't be called from trap context" errors with logger trap('INT') { @shutdown_requested = true } # Wait for workers with periodic checks for shutdown until @shutdown_called || @shutdown_requested sleep 0.1 end # Call shutdown if it was requested by signal shutdown! if @shutdown_requested end |
#shutdown! ⇒ nil, untyped
Gracefully shuts down all workers
51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 |
# File 'lib/takagi/server/udp.rb', line 51 def shutdown! return if @shutdown_called @shutdown_called = true @watcher.stop close_socket terminate_workers Takagi::Observable::Registry.stop_all # Join the server thread if it was spawned if defined?(@server_thread) && @server_thread&.alive? @server_thread.join(5) # Wait up to 5 seconds end exit(0) unless test_environment? Takagi::Hooks.emit(:server_stopped, protocol: :udp, port: @port) end |
#spawn_workers ⇒ Object
77 78 79 80 81 82 |
# File 'lib/takagi/server/udp.rb', line 77 def spawn_workers @worker_pids = Array.new(@worker_processes) do @logger.debug "process with #{@router.all_routes}" fork_worker end end |
#terminate_workers ⇒ nil, untyped
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 |
# File 'lib/takagi/server/udp.rb', line 111 def terminate_workers return unless @worker_pids.is_a?(Array) @worker_pids.each do |pid| Process.kill('TERM', pid) rescue Errno::ESRCH # worker already exited end # Give workers a moment to shut down gracefully deadline = Time.now + 2 @worker_pids.each do |pid| begin timeout = [deadline - Time.now, 0].max Timeout.timeout(timeout) do Process.wait(pid) end rescue Timeout::Error, Errno::ECHILD, Errno::ESRCH # Worker didn't exit in time or already exited end end end |
#test_environment? ⇒ Boolean
134 135 136 |
# File 'lib/takagi/server/udp.rb', line 134 def test_environment? ENV['RACK_ENV'] == 'test' || defined?(RSpec) end |
#worker_configuration ⇒ { port: untyped, socket: untyped, middleware_stack: untyped, sender: untyped, logger: untyped, threads: untyped }
91 92 93 94 95 96 97 98 99 100 101 |
# File 'lib/takagi/server/udp.rb', line 91 def worker_configuration { port: @port, socket: @socket, middleware_stack: @middleware_stack, router: @router, sender: @sender, logger: @logger, threads: @worker_threads } end |