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
- #expect_version_message(message) ⇒ Object
- #peer_failure(code, text, message: nil, **context) ⇒ Object
- #process_deferred ⇒ Object
- #process_incoming_message(message) ⇒ Object
- #process_message(message) ⇒ Object
- #process_packet(json) ⇒ Object
- #publish_peer_failure(failure, message:) ⇒ Object
- #reject_fatal_message(message, error) ⇒ Object
- #reject_invalid_message(message, validation) ⇒ Object
- #reject_invalid_packet(json, error) ⇒ Object
- #reject_malformed_packet(json, attributes, error) ⇒ Object
- #reject_packet(code, text, raw:, message: nil) ⇒ Object
- #reject_processed_message(message, error) ⇒ Object
- #should_validate_ingoing_message?(message) ⇒ Boolean
- #verify_sequence(message) ⇒ Object
- #will_not_handle(message) ⇒ Object
Instance Method Details
#expect_version_message(message) ⇒ Object
159 160 161 162 163 |
# File 'lib/rsmp/proxy/modules/receive.rb', line 159 def () return if .is_a?(Version) || .is_a?(MessageAck) || .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: ).compact ) end |
#process_deferred ⇒ Object
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 () validate = () if validate validation = .validate(schemas) return (, validation) if validation.invalid? end verify_sequence() .decode_for(schemas) if validate with_deferred_distribution do distribute() () end process_deferred Result.success() 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 () case when MessageAck process_ack() when MessageNotAck process_not_ack() when Version process_version() when RSMP::Watchdog process_watchdog() else dont_acknowledge(, 'Received', "unknown 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.build(attributes, json) () rescue InvalidPacket => e reject_invalid_packet(json, e) rescue MalformedMessage => e reject_malformed_packet(json, attributes, e) rescue PeerMessageError => e (, e) rescue HandshakeError, FatalError => e (, 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: , 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 (, error) text = "Rejected #{&.type || 'message'}: #{error.}" failure = peer_failure(:protocol_failure, text, message: ) publish_peer_failure(failure, message: ) dont_acknowledge(, "Rejected #{.type}", error..to_s) if 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 (, validation) schemas_string = schemas.map { |schema| "#{schema.first}: #{schema.last}" }.join(', ') reason = "schema errors (#{schemas_string}): #{validation.}" text = "Received invalid #{.type}" failure = peer_failure( :invalid_peer_message, "#{text}: #{validation.}", message: , validation: validation ) publish_peer_failure(failure, message: ) log(text, message: , level: :warning) dont_acknowledge(, 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.}", 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.}", 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: , raw: raw) publish_peer_failure(failure, message: ) log(text, 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 (, error) code = error.is_a?(MessageRejected) ? :message_rejected : :invalid_peer_message text = "Received invalid #{.type}: #{error.}" failure = peer_failure(code, text, message: ) publish_peer_failure(failure, message: ) dont_acknowledge(, "Received invalid #{.type}", error..to_s) Result.failure(failure: failure) end |
#should_validate_ingoing_message?(message) ⇒ Boolean
7 8 9 10 11 12 13 14 15 16 |
# File 'lib/rsmp/proxy/modules/receive.rb', line 7 def () return false if .is_a?(Version) && !@version_determined return true unless @site_settings skip = @site_settings['skip_validation'] return true unless skip klass = .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() () 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() reason = "since we're a #{self.class.name.downcase}" log "Ignoring #{.type}, #{reason}", message: , level: :warning dont_acknowledge(, nil, reason) end |