Class: Takagi::Message::RetransmissionManager

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

Overview

Implements CoAP retransmission logic as per RFC 7252 Section 4.2

For CON (Confirmable) messages, the client MUST retransmit the message until it receives an ACK, RST, or the transmission times out.

RFC 7252 §4.2: Retransmission uses exponential back-off with random factor RFC 7252 §4.8: Default transmission parameters

Constant Summary collapse

ACK_TIMEOUT =

RFC 7252 §4.8: Transmission parameters

Returns:

  • (::Float)
2.0
ACK_RANDOM_FACTOR =

Initial timeout in seconds

Returns:

  • (::Float)
1.5
MAX_RETRANSMIT =

Random factor for timeout calculation

Returns:

  • (4)
4
PendingTransmission =

Pending transmission tracking

Returns:

  • (Object)

Instance Method Summary collapse

Constructor Details

#initialize(logger: nil) ⇒ RetransmissionManager

Returns a new instance of RetransmissionManager.

Parameters:

  • logger: (Object, nil) (defaults to: nil)


35
36
37
38
39
40
41
# File 'lib/takagi/message/retransmission_manager.rb', line 35

def initialize(logger: nil)
  @pending = {}
  @mutex = Mutex.new
  @logger = logger || Takagi.logger
  @running = false
  @thread = nil
end

Instance Method Details

#calculate_timeout(attempt) ⇒ Object

Calculate timeout with exponential backoff and random factor RFC 7252 §4.2: timeout = ACK_TIMEOUT * (2 ** attempt) * random_factor where random_factor is between 1.0 and ACK_RANDOM_FACTOR

Parameters:

  • attempt (Object)

Returns:

  • (Object)


186
187
188
189
190
# File 'lib/takagi/message/retransmission_manager.rb', line 186

def calculate_timeout(attempt)
  base_timeout = ACK_TIMEOUT * (2**attempt)
  random_factor = 1.0 + (rand * (ACK_RANDOM_FACTOR - 1.0))
  base_timeout * random_factor
end

#handle_response(message_id, response_data = nil) ⇒ nil, untyped

Handle incoming ACK or RST to cancel retransmission

Parameters:

  • message_id (Integer)

    CoAP Message ID

  • response_data (String) (defaults to: nil)

    Response message data

Returns:

  • (nil, untyped)


94
95
96
97
98
99
100
101
102
103
104
105
# File 'lib/takagi/message/retransmission_manager.rb', line 94

def handle_response(message_id, response_data = nil)
  transmission = nil

  @mutex.synchronize do
    transmission = @pending.delete(message_id)
  end

  return unless transmission

  @logger.debug "Received response for MID: #{message_id}, canceling retransmission"
  transmission.callback&.call(response_data, nil)
end

#handle_timeout(transmission) ⇒ Object

Handle a transmission timeout

Parameters:

  • transmission (Object)

Returns:

  • (Object)


148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
# File 'lib/takagi/message/retransmission_manager.rb', line 148

def handle_timeout(transmission)
  if transmission.attempt >= MAX_RETRANSMIT
    # Max retries exceeded
    @mutex.synchronize do
      @pending.delete(transmission.message_id)
    end

    @logger.warn "Message #{transmission.message_id} failed after #{MAX_RETRANSMIT} retransmissions"
    transmission.callback&.call(nil, 'Timeout: Max retransmissions exceeded')
  else
    # Retransmit with exponential backoff
    transmission.attempt += 1
    transmission.next_timeout = calculate_timeout(transmission.attempt)
    transmission.timeout_at = Time.now.to_f + transmission.next_timeout

    transmit(transmission)

    @logger.debug "Retransmitting MID: #{transmission.message_id} " \
                  "(attempt #{transmission.attempt}/#{MAX_RETRANSMIT}, " \
                  "next timeout: #{transmission.next_timeout}s)"
  end
end

#process_timeoutsObject

Process timed out transmissions

Returns:

  • (Object)


132
133
134
135
136
137
138
139
140
141
142
143
144
145
# File 'lib/takagi/message/retransmission_manager.rb', line 132

def process_timeouts
  current_time = Time.now.to_f
  timed_out = []

  @mutex.synchronize do
    @pending.each_value do |transmission|
      timed_out << transmission if transmission.timed_out?(current_time)
    end
  end

  timed_out.each do |transmission|
    handle_timeout(transmission)
  end
end

#run_retransmission_loopObject

Main retransmission loop

Returns:

  • (Object)


120
121
122
123
124
125
126
127
128
129
# File 'lib/takagi/message/retransmission_manager.rb', line 120

def run_retransmission_loop
  while @running
    sleep 0.1 # Check every 100ms

    process_timeouts
  end
rescue StandardError => e
  @logger.error "Retransmission loop error: #{e.message}"
  @logger.error e.backtrace.join("\n")
end

#send_confirmable(message_id, message_data, socket, host, port, &callback) {|arg0| ... } ⇒ Object

Send a CON message with automatic retransmission

Parameters:

  • message_id (Integer)

    CoAP Message ID

  • message_data (String)

    Serialized message bytes

  • socket (UDPSocket)

    Socket to send on

  • host (String)

    Destination host

  • port (Integer)

    Destination port

  • callback (Proc)

    Called with response or timeout error

Yields:

Yield Parameters:

  • arg0

Yield Returns:

  • (Object)

Returns:

  • (Object)


66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
# File 'lib/takagi/message/retransmission_manager.rb', line 66

def send_confirmable(message_id, message_data, socket, host, port, &callback)
  initial_timeout = calculate_timeout(0)

  transmission = PendingTransmission.new(
    message_id,
    message_data,
    socket,
    host,
    port,
    0, # attempt
    initial_timeout,
    Time.now.to_f + initial_timeout,
    callback
  )

  @mutex.synchronize do
    @pending[message_id] = transmission
  end

  # Send initial transmission
  transmit(transmission)

  @logger.debug "Scheduled CON message (MID: #{message_id}) with timeout #{initial_timeout}s"
end

#startnil, untyped

Start the retransmission manager background thread

Returns:

  • (nil, untyped)


44
45
46
47
48
49
50
# File 'lib/takagi/message/retransmission_manager.rb', line 44

def start
  return if @running

  @running = true
  @thread = Thread.new { run_retransmission_loop }
  @logger.debug 'Retransmission manager started'
end

#statsObject

Get statistics about pending transmissions

Returns:

  • (Object)


108
109
110
111
112
113
114
115
# File 'lib/takagi/message/retransmission_manager.rb', line 108

def stats
  @mutex.synchronize do
    {
      pending_count: @pending.size,
      message_ids: @pending.keys
    }
  end
end

#stopObject

Stop the retransmission manager

Returns:

  • (Object)


53
54
55
56
57
# File 'lib/takagi/message/retransmission_manager.rb', line 53

def stop
  @running = false
  @thread&.join
  @logger.debug 'Retransmission manager stopped'
end

#transmit(transmission) ⇒ Object

Transmit a message

Parameters:

  • transmission (Object)

Returns:

  • (Object)


172
173
174
175
176
177
178
179
180
181
# File 'lib/takagi/message/retransmission_manager.rb', line 172

def transmit(transmission)
  transmission.socket.send(
    transmission.message_data,
    0,
    transmission.host,
    transmission.port
  )
rescue StandardError => e
  @logger.error "Transmission failed for MID #{transmission.message_id}: #{e.message}"
end