Class: Takagi::UdpClient
- Inherits:
-
ClientBase
- Object
- ClientBase
- Takagi::UdpClient
- 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
- #check_socket_for_response(message, socket, state) ⇒ Object
-
#cleanup_resources ⇒ Object
Stops the retransmission manager thread.
- #deliver_raw_response(response, &callback) ⇒ void
- #handle_response_state(state) {|arg0| ... } ⇒ Object
-
#initialize(server_uri, timeout: 5, use_retransmission: true) ⇒ UdpClient
constructor
Initializes the UDP client.
-
#request(method, path, payload = nil, options: {}, type: nil, &callback) {|arg0| ... } ⇒ Object
Executes a request to the server using Takagi::Message::Request.
-
#request_simple(message, uri) {|arg0| ... } ⇒ Object
Simple request without retransmission (legacy mode).
-
#request_with_retransmission(message, uri) {|arg0| ... } ⇒ Object
RFC 7252 §4.2 compliant request with automatic retransmission.
- #send_with_retransmission(message, socket, uri) ⇒ Object
- #update_state(state, response_data, error) ⇒ Object
- #wait_for_response(message, socket, state) ⇒ Object
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
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
229 230 231 232 233 234 235 236 237 |
# File 'lib/takagi/client.rb', line 229 def check_socket_for_response(, socket, state) return unless socket.wait_readable(0.1) state[:response_data], = socket.recvfrom(1024) @retransmission_manager.handle_response(., state[:response_data]) state[:response_received] = true rescue StandardError => e update_state(state, nil, e.) end |
#cleanup_resources ⇒ Object
Stops the retransmission manager thread
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.
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
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
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) = Takagi::Message::Request.new( method: method, uri: uri, payload: payload, type: type, options: ) if @use_retransmission request_with_retransmission(, uri, &callback) else request_simple(, uri, &callback) end end |
#request_simple(message, uri) {|arg0| ... } ⇒ Object
Simple request without retransmission (legacy mode)
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(, uri, &callback) socket = UDPSocket.new socket.send(.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
204 205 206 207 208 209 210 211 |
# File 'lib/takagi/client.rb', line 204 def request_with_retransmission(, uri, &callback) socket = UDPSocket.new state = send_with_retransmission(, 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
213 214 215 216 217 218 219 220 221 222 |
# File 'lib/takagi/client.rb', line 213 def send_with_retransmission(, socket, uri) state = { response_received: false, response_data: nil, error: nil } @retransmission_manager.send_confirmable( ., .to_bytes, socket, uri.host, uri.port || 5683 ) { |resp, err| update_state(state, resp, err) } wait_for_response(, socket, state) state end |
#update_state(state, response_data, error) ⇒ Object
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
224 225 226 227 |
# File 'lib/takagi/client.rb', line 224 def wait_for_response(, socket, state) start_time = Time.now check_socket_for_response(, socket, state) until state[:response_received] || (Time.now - start_time) > @timeout end |