Class: Takagi::Server::UdpWorker

Inherits:
Object
  • Object
show all
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

Constructor Details

#initialize(socket:, middleware_stack:, **options) ⇒ UdpWorker

Returns a new instance of UdpWorker.

Parameters:

  • socket: (Object)
  • middleware_stack: (Object)
  • options (Object)


10
11
12
13
14
15
16
17
18
19
# File 'lib/takagi/server/udp_worker.rb', line 10

def initialize(socket:, middleware_stack:, **options)
  @socket = socket
  @middleware_stack = middleware_stack
  @router = options.fetch(:router, nil)
  @sender = options.fetch(:sender)
  @logger = options.fetch(:logger)
  @port = options.fetch(:port)
  @threads = options.fetch(:threads)
  @dedup_cache = Takagi::Message::DeduplicationCache.new
end

Instance Method Details

#build_response(inbound_request, result) ⇒ Object

Parameters:

  • inbound_request (Object)
  • result (Object)

Returns:

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

Parameters:

  • request (Object)
  • addr (Object)

Returns:

  • (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.message_id, source_endpoint)
    if cached_response
      @logger.debug "Duplicate CON detected (MID: #{inbound_request.message_id}), 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.message}"
end

#immediate_response(inbound_request) ⇒ untyped?

Parameters:

  • inbound_request (Object)

Returns:

  • (untyped, nil)


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.message_id,
      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.message_id,
    type: 2
  )
end

#log_inbound_request(inbound_request) ⇒ Object

Parameters:

  • inbound_request (Object)

Returns:

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

Parameters:

  • result (Object)

Returns:

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

Parameters:

  • queue (Object)

Returns:

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

#runObject

Returns:

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

Parameters:

  • queue (Object)

Returns:

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

Parameters:

  • response (Object)
  • addr (Object)

Returns:

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