Class: Takagi::UdpClient

Inherits:
ClientBase show all
Defined in:
lib/takagi/client.rb,
sig/takagi/client.rbs

Overview

UDP-specific client implementation (internal) Users should use Takagi::Client with protocol: :udp instead

Instance Attribute Summary

Attributes inherited from ClientBase

#callbacks, #server_uri, #timeout

Instance Method Summary collapse

Methods inherited from ClientBase

#close, #closed?, #delete, #deliver_response, #get, #get_json, #observe, #on, open, #parse_json_response, #post, #post_json, #put, #put_json

Constructor Details

#initialize(server_uri, timeout: 5, use_retransmission: true) ⇒ UdpClient

Initializes the UDP client

Parameters:

  • URL of the Takagi server

  • (defaults to: 5)

    Maximum time to wait for a response

  • (defaults to: true)

    Enable RFC 7252 §4.2 compliant retransmission (default: true)

  • (defaults to: 5)
  • (defaults to: true)


141
142
143
144
145
146
147
148
149
# File 'lib/takagi/client.rb', line 141

def initialize(server_uri, timeout: 5, use_retransmission: true)
  super(server_uri, timeout: timeout)
  @use_retransmission = use_retransmission

  return unless @use_retransmission

  @retransmission_manager = Takagi::Message::RetransmissionManager.new
  @retransmission_manager.start
end

Instance Method Details

#check_socket_for_response(message, socket, state) ⇒ Object

Parameters:

Returns:



229
230
231
232
233
234
235
236
237
# File 'lib/takagi/client.rb', line 229

def check_socket_for_response(message, socket, state)
  return unless socket.wait_readable(0.1)

  state[:response_data], = socket.recvfrom(1024)
  @retransmission_manager.handle_response(message.message_id, state[:response_data])
  state[:response_received] = true
rescue StandardError => e
  update_state(state, nil, e.message)
end

#cleanup_resourcesObject

Stops the retransmission manager thread

Returns:



154
155
156
157
# File 'lib/takagi/client.rb', line 154

def cleanup_resources
  @retransmission_manager&.stop
  super
end

#deliver_raw_response(response, &callback) ⇒ void

This method returns an undefined value.

Parameters:



255
256
257
258
259
260
261
262
263
# File 'lib/takagi/client.rb', line 255

def deliver_raw_response(response, &callback)
  if callback
    callback.call(response)
  elsif @callbacks[:response]
    @callbacks[:response].call(response)
  else
    puts response
  end
end

#handle_response_state(state) {|arg0| ... } ⇒ Object

Parameters:

Yields:

Yield Parameters:

  • arg0

Yield Returns:

  • (Object)

Returns:



245
246
247
248
249
250
251
252
253
# File 'lib/takagi/client.rb', line 245

def handle_response_state(state, &callback)
  if state[:response_received] && !state[:error]
    deliver_response(state[:response_data], &callback)
  elsif state[:error]
    puts "TakagiClient Error: #{state[:error]}"
  else
    puts 'TakagiClient Error: Request timeout'
  end
end

#request(method, path, payload = nil, options: {}, type: nil, &callback) {|arg0| ... } ⇒ Object

Executes a request to the server using Takagi::Message::Request

Parameters:

  • HTTP method (:get, :post, :put, :delete)

  • Resource path

  • (defaults to: nil)

    (optional) Data for POST/PUT requests

  • (optional) Callback function for processing the response

Yields:

Yield Parameters:

  • arg0

Yield Returns:

  • (Object)

Returns:



168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
# File 'lib/takagi/client.rb', line 168

def request(method, path, payload = nil, options: {}, type: nil, &callback)
  uri = URI.join(server_uri.to_s, path)
  message = Takagi::Message::Request.new(
    method: method,
    uri: uri,
    payload: payload,
    type: type,
    options: options
  )

  if @use_retransmission
    request_with_retransmission(message, uri, &callback)
  else
    request_simple(message, uri, &callback)
  end
end

#request_simple(message, uri) {|arg0| ... } ⇒ Object

Simple request without retransmission (legacy mode)

Parameters:

Yields:

Yield Parameters:

  • arg0

Yield Returns:

  • (Object)

Returns:



186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
# File 'lib/takagi/client.rb', line 186

def request_simple(message, uri, &callback)
  socket = UDPSocket.new
  socket.send(message.to_bytes, 0, uri.host, uri.port || 5683)

  unless socket.wait_readable(@timeout)
    puts 'TakagiClient Error: Request timeout'
    return
  end

  response, = socket.recvfrom(1024)
  deliver_raw_response(response, &callback)
rescue StandardError => e
  puts "TakagiClient Error: #{e.message}"
ensure
  socket&.close unless socket&.closed?
end

#request_with_retransmission(message, uri) {|arg0| ... } ⇒ Object

RFC 7252 §4.2 compliant request with automatic retransmission

Parameters:

Yields:

Yield Parameters:

  • arg0

Yield Returns:

  • (Object)

Returns:



204
205
206
207
208
209
210
211
# File 'lib/takagi/client.rb', line 204

def request_with_retransmission(message, uri, &callback)
  socket = UDPSocket.new
  state = send_with_retransmission(message, socket, uri)
  socket.close
  handle_response_state(state, &callback)
rescue StandardError => e
  puts "TakagiClient Error: #{e.message}"
end

#send_with_retransmission(message, socket, uri) ⇒ Object

Parameters:

Returns:



213
214
215
216
217
218
219
220
221
222
# File 'lib/takagi/client.rb', line 213

def send_with_retransmission(message, socket, uri)
  state = { response_received: false, response_data: nil, error: nil }

  @retransmission_manager.send_confirmable(
    message.message_id, message.to_bytes, socket, uri.host, uri.port || 5683
  ) { |resp, err| update_state(state, resp, err) }

  wait_for_response(message, socket, state)
  state
end

#update_state(state, response_data, error) ⇒ Object

Parameters:

Returns:



239
240
241
242
243
# File 'lib/takagi/client.rb', line 239

def update_state(state, response_data, error)
  state[:response_data] = response_data
  state[:error] = error
  state[:response_received] = true
end

#wait_for_response(message, socket, state) ⇒ Object

Parameters:

Returns:



224
225
226
227
# File 'lib/takagi/client.rb', line 224

def wait_for_response(message, socket, state)
  start_time = Time.now
  check_socket_for_response(message, socket, state) until state[:response_received] || (Time.now - start_time) > @timeout
end