Module: RSMP::Proxy::Modules::Receive

Included in:
RSMP::Proxy
Defined in:
lib/rsmp/proxy/modules/receive.rb

Overview

Message processing functionality. Peer-controlled invalid input is represented by Result::Failure; unexpected implementation errors escape.

Instance Method Summary collapse

Instance Method Details

#expect_version_message(message) ⇒ Object

Raises:



159
160
161
162
163
# File 'lib/rsmp/proxy/modules/receive.rb', line 159

def expect_version_message(message)
  return if message.is_a?(Version) || message.is_a?(MessageAck) || message.is_a?(MessageNotAck)

  raise HandshakeError, 'Version must be received first'
end

#peer_failure(code, text, message: nil, **context) ⇒ Object



117
118
119
120
121
122
123
124
# File 'lib/rsmp/proxy/modules/receive.rb', line 117

def peer_failure(code, text, message: nil, **context)
  Failure.new(
    code: code,
    message: text,
    source: :peer,
    context: context.merge(message: message).compact
  )
end

#process_deferredObject



18
19
20
# File 'lib/rsmp/proxy/modules/receive.rb', line 18

def process_deferred
  @node.process_deferred
end

#process_incoming_message(message) ⇒ Object



42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
# File 'lib/rsmp/proxy/modules/receive.rb', line 42

def process_incoming_message(message)
  validate = should_validate_ingoing_message?(message)
  if validate
    validation = message.validate(schemas)
    return reject_invalid_message(message, validation) if validation.invalid?
  end

  verify_sequence(message)
  message.decode_for(schemas) if validate
  with_deferred_distribution do
    distribute(message)
    process_message(message)
  end
  process_deferred
  Result.success(message)
end

#process_message(message) ⇒ Object



138
139
140
141
142
143
144
145
146
147
148
149
150
151
# File 'lib/rsmp/proxy/modules/receive.rb', line 138

def process_message(message)
  case message
  when MessageAck
    process_ack(message)
  when MessageNotAck
    process_not_ack(message)
  when Version
    process_version(message)
  when RSMP::Watchdog
    process_watchdog(message)
  else
    dont_acknowledge(message, 'Received', "unknown message (#{message.type})")
  end
end

#process_packet(json) ⇒ Object



26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
# File 'lib/rsmp/proxy/modules/receive.rb', line 26

def process_packet(json)
  attributes = Message.parse_attributes(json)
  message = Message.build(attributes, json)
  process_incoming_message(message)
rescue InvalidPacket => e
  reject_invalid_packet(json, e)
rescue MalformedMessage => e
  reject_malformed_packet(json, attributes, e)
rescue PeerMessageError => e
  reject_processed_message(message, e)
rescue HandshakeError, FatalError => e
  reject_fatal_message(message, e)
ensure
  @node&.clear_deferred
end

#publish_peer_failure(failure, message:) ⇒ Object



126
127
128
129
130
131
132
133
134
135
136
# File 'lib/rsmp/proxy/modules/receive.rb', line 126

def publish_peer_failure(failure, message:)
  distribute_event(
    Event.new(
      type: :invalid_message,
      source: self,
      session_id: session_id,
      message: message,
      failure: failure
    )
  )
end

#reject_fatal_message(message, error) ⇒ Object



108
109
110
111
112
113
114
115
# File 'lib/rsmp/proxy/modules/receive.rb', line 108

def reject_fatal_message(message, error)
  text = "Rejected #{message&.type || 'message'}: #{error.message}"
  failure = peer_failure(:protocol_failure, text, message: message)
  publish_peer_failure(failure, message: message)
  dont_acknowledge(message, "Rejected #{message.type}", error.message.to_s) if message
  close(reason: :protocol_failure, failure: failure)
  Result.failure(failure: failure)
end

#reject_invalid_message(message, validation) ⇒ Object



76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
# File 'lib/rsmp/proxy/modules/receive.rb', line 76

def reject_invalid_message(message, validation)
  schemas_string = schemas.map { |schema| "#{schema.first}: #{schema.last}" }.join(', ')
  reason = "schema errors (#{schemas_string}): #{validation.message}"
  text = "Received invalid #{message.type}"
  failure = peer_failure(
    :invalid_peer_message,
    "#{text}: #{validation.message}",
    message: message,
    validation: validation
  )
  publish_peer_failure(failure, message: message)
  log(text, message: message, level: :warning)
  dont_acknowledge(message, text, reason)
  Result.failure(failure: failure)
end

#reject_invalid_packet(json, error) ⇒ Object



59
60
61
62
63
64
65
# File 'lib/rsmp/proxy/modules/receive.rb', line 59

def reject_invalid_packet(json, error)
  reject_packet(
    :invalid_packet,
    "Received invalid packet, expected JSON but got #{json.size} bytes: #{error.message}",
    raw: json
  )
end

#reject_malformed_packet(json, attributes, error) ⇒ Object



67
68
69
70
71
72
73
74
# File 'lib/rsmp/proxy/modules/receive.rb', line 67

def reject_malformed_packet(json, attributes, error)
  reject_packet(
    :malformed_message,
    "Received malformed message: #{error.message}",
    raw: json,
    message: Malformed.new(attributes || {})
  )
end

#reject_packet(code, text, raw:, message: nil) ⇒ Object



92
93
94
95
96
97
# File 'lib/rsmp/proxy/modules/receive.rb', line 92

def reject_packet(code, text, raw:, message: nil)
  failure = peer_failure(code, text, message: message, raw: raw)
  publish_peer_failure(failure, message: message)
  log(text, message: message, level: :warning)
  Result.failure(failure: failure)
end

#reject_processed_message(message, error) ⇒ Object



99
100
101
102
103
104
105
106
# File 'lib/rsmp/proxy/modules/receive.rb', line 99

def reject_processed_message(message, error)
  code = error.is_a?(MessageRejected) ? :message_rejected : :invalid_peer_message
  text = "Received invalid #{message.type}: #{error.message}"
  failure = peer_failure(code, text, message: message)
  publish_peer_failure(failure, message: message)
  dont_acknowledge(message, "Received invalid #{message.type}", error.message.to_s)
  Result.failure(failure: failure)
end

#should_validate_ingoing_message?(message) ⇒ Boolean

Returns:

  • (Boolean)


7
8
9
10
11
12
13
14
15
16
# File 'lib/rsmp/proxy/modules/receive.rb', line 7

def should_validate_ingoing_message?(message)
  return false if message.is_a?(Version) && !@version_determined
  return true unless @site_settings

  skip = @site_settings['skip_validation']
  return true unless skip

  klass = message.class.name.split('::').last
  !skip.include?(klass)
end

#verify_sequence(message) ⇒ Object



22
23
24
# File 'lib/rsmp/proxy/modules/receive.rb', line 22

def verify_sequence(message)
  expect_version_message(message) unless @version_determined
end

#will_not_handle(message) ⇒ Object



153
154
155
156
157
# File 'lib/rsmp/proxy/modules/receive.rb', line 153

def will_not_handle(message)
  reason = "since we're a #{self.class.name.downcase}"
  log "Ignoring #{message.type}, #{reason}", message: message, level: :warning
  dont_acknowledge(message, nil, reason)
end