Class: Takagi::Server::UdpWorker
- Inherits:
-
Object
- Object
- Takagi::Server::UdpWorker
- Defined in:
- lib/takagi/server/udp_worker.rb,
sig/takagi/server/udp_worker.rbs
Overview
Handles incoming UDP messages on behalf of the master Udp server.
Instance Method Summary collapse
- #build_response(inbound_request, result) ⇒ Object
- #handle_request(request, addr) ⇒ Object
- #immediate_response(inbound_request) ⇒ untyped?
-
#initialize(socket:, middleware_stack:, **options) ⇒ UdpWorker
constructor
A new instance of UdpWorker.
- #log_inbound_request(inbound_request) ⇒ Object
- #log_middleware_result(result) ⇒ nil, untyped
- #process_loop(queue) ⇒ Object
- #run ⇒ Object
- #spawn_thread(queue) ⇒ Object
- #transmit(response, addr) ⇒ Object
Constructor Details
#initialize(socket:, middleware_stack:, **options) ⇒ UdpWorker
Returns a new instance of UdpWorker.
10 11 12 13 14 15 16 17 18 19 |
# File 'lib/takagi/server/udp_worker.rb', line 10 def initialize(socket:, middleware_stack:, **) @socket = socket @middleware_stack = middleware_stack @router = .fetch(:router, nil) @sender = .fetch(:sender) @logger = .fetch(:logger) @port = .fetch(:port) @threads = .fetch(:threads) @dedup_cache = Takagi::Message::DeduplicationCache.new end |
Instance Method Details
#build_response(inbound_request, result) ⇒ Object
165 166 167 |
# File 'lib/takagi/server/udp_worker.rb', line 165 def build_response(inbound_request, result) ResponseBuilder.build(inbound_request, result, logger: @logger) end |
#handle_request(request, addr) ⇒ Object
59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 |
# File 'lib/takagi/server/udp_worker.rb', line 59 def handle_request(request, addr) inbound_request = Takagi::Message::Inbound.new(request) log_inbound_request(inbound_request) # RFC 7252 §4.4: Check for duplicate messages source_endpoint = "#{addr[3]}:#{addr[1]}" if inbound_request.type.zero? # CON message cached_response = @dedup_cache.check_duplicate(inbound_request., source_endpoint) if cached_response @logger.debug "Duplicate CON detected (MID: #{inbound_request.}), resending cached response" return @sender.transmit(cached_response, addr[3], addr[1]) end end immediate = immediate_response(inbound_request) return transmit(immediate, addr) if immediate # Delegate to controller's thread pool if available delegate_to_controller_pool(inbound_request, addr) rescue StandardError => e @logger.error "Handle_request failed: #{e.}" end |
#immediate_response(inbound_request) ⇒ untyped?
132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 |
# File 'lib/takagi/server/udp_worker.rb', line 132 def immediate_response(inbound_request) if inbound_request.type == 1 && inbound_request.method == 'EMPTY' return Takagi::Message::Outbound.new( code: '0.00', payload: '', token: '', message_id: inbound_request., type: 3 ) end return unless inbound_request.code.zero? && inbound_request.method == 'EMPTY' Takagi::Message::Outbound.new( code: '0.00', payload: '', token: inbound_request.token, message_id: inbound_request., type: 2 ) end |
#log_inbound_request(inbound_request) ⇒ Object
127 128 129 130 |
# File 'lib/takagi/server/udp_worker.rb', line 127 def log_inbound_request(inbound_request) @logger.debug "Code: #{inbound_request.code}" @logger.debug "Method: #{inbound_request.method}" end |
#log_middleware_result(result) ⇒ nil, untyped
154 155 156 157 158 159 160 161 162 163 |
# File 'lib/takagi/server/udp_worker.rb', line 154 def log_middleware_result(result) @logger.debug "Middleware result class: #{result.class}" @logger.debug "Middleware result inspect: #{result.inspect}" return unless result.is_a?(Hash) @logger.debug "Hash response keys: #{result.keys}" result.each do |key, value| @logger.debug "Key: #{key.inspect} => #{value.inspect} (#{value.class})" end end |
#process_loop(queue) ⇒ Object
35 36 37 38 39 40 41 42 43 44 45 46 47 48 |
# File 'lib/takagi/server/udp_worker.rb', line 35 def process_loop(queue) loop do break if @shutdown next unless @socket.wait_readable(0.1) queue << @socket.recvfrom(1024) end @logger.debug "[Worker PID: #{Process.pid}] Shutting down..." exit(0) rescue Interrupt, SignalException @logger.debug "[Worker PID: #{Process.pid}] Interrupted, shutting down..." exit(0) end |
#run ⇒ Object
21 22 23 24 25 26 27 28 29 30 31 |
# File 'lib/takagi/server/udp_worker.rb', line 21 def run @shutdown = false trap('TERM') { @shutdown = true } trap('INT') { @shutdown = true } queue = Queue.new Array.new(@threads) { spawn_thread(queue) } @logger.debug "[Worker PID: #{Process.pid}] Listening on CoAP://0.0.0.0:#{@port} with #{@threads} threads" process_loop(queue) end |
#spawn_thread(queue) ⇒ Object
50 51 52 53 54 55 56 57 |
# File 'lib/takagi/server/udp_worker.rb', line 50 def spawn_thread(queue) Thread.new do loop do request, addr = queue.pop handle_request(request, addr) end end end |
#transmit(response, addr) ⇒ Object
169 170 171 |
# File 'lib/takagi/server/udp_worker.rb', line 169 def transmit(response, addr) @sender.transmit(response, addr[3], addr[1]) end |