Class: Takagi::Server::Udp

Inherits:
Object
  • Object
show all
Defined in:
lib/takagi/server/udp.rb,
sig/takagi/server/udp.rbs

Overview

UDP server for handling CoAP messages

Instance Method Summary collapse

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.

Parameters:

  • port: (::Integer) (defaults to: 5683)
  • worker_processes: (::Integer) (defaults to: 2)
  • worker_threads: (::Integer) (defaults to: 2)
  • middleware_stack: (Object, nil) (defaults to: nil)
  • router: (Object, nil) (defaults to: nil)
  • logger: (Object, nil) (defaults to: nil)
  • watcher: (Object, nil) (defaults to: nil)


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_socketObject

Returns:

  • (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_workerObject

Returns:

  • (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_detailsObject

Returns:

  • (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

Returns:

  • (Object)


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

Returns:

  • (nil, untyped)


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_workersObject

Returns:

  • (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_workersnil, untyped

Returns:

  • (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

Returns:

  • (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 }

Returns:

  • ({ 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