Module: RSMP::Proxy::Modules::Send

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

Overview

Message sending functionality. Expected operational failures are returned as Result::Failure; unexpected implementation errors are not rescued.

Instance Method Summary collapse

Instance Method Details

#apply_nts_message_attributes(message) ⇒ Object



137
138
139
140
141
142
# File 'lib/rsmp/proxy/modules/send.rb', line 137

def apply_nts_message_attributes(message)
  return if core_3_3?

  message.attributes['ntsOId'] = main && main.ntsoid ? main.ntsoid : ''
  message.attributes['xNId'] = main && main.xnid ? main.xnid : ''
end

#buffer_message(message) ⇒ Object

Base proxies do not buffer. SupervisorProxy's MessageBuffer override returns a cloned queued message when its explicit policy accepts it.



90
91
92
93
# File 'lib/rsmp/proxy/modules/send.rb', line 90

def buffer_message(message)
  log "Discarded #{message.type}; connection is #{@state}", message: message, level: :warning
  nil
end

#buffer_or_fail(message, buffer:) ⇒ Object



76
77
78
79
80
81
82
83
84
85
86
# File 'lib/rsmp/proxy/modules/send.rb', line 76

def buffer_or_fail(message, buffer:)
  buffered = buffer_message(message) if buffer
  return Result.success(Delivery.new(message: buffered, state: :buffered)) if buffered

  Result.failure(
    :not_ready,
    message: "Cannot send #{message.type}: connection is #{@state}",
    source: :connection,
    context: { message: message, state: @state, session_id: session_id }
  )
end

#disconnected_send_result(message) ⇒ Object



95
96
97
98
99
100
101
102
# File 'lib/rsmp/proxy/modules/send.rb', line 95

def disconnected_send_result(message)
  Result.failure(
    :disconnected,
    message: "Cannot send #{message.type}: connection transport is closed",
    source: :transport,
    context: { message: message, session_id: session_id }
  )
end

#invalid_outbound_result(message, detail, validation: nil, cause: nil) ⇒ Object



104
105
106
107
108
109
110
111
112
113
114
# File 'lib/rsmp/proxy/modules/send.rb', line 104

def invalid_outbound_result(message, detail, validation: nil, cause: nil)
  text = "Could not send #{message.type}: #{detail}"
  log(text, message: message, level: :error)
  Result.failure(
    :invalid_outbound_message,
    message: text,
    source: :local,
    context: { message: message, validation: validation }.compact,
    cause: cause
  )
end

#log_send(message, reason = nil) ⇒ Object



116
117
118
119
120
# File 'lib/rsmp/proxy/modules/send.rb', line 116

def log_send(message, reason = nil)
  text = reason ? "Sent #{message.type} #{reason}" : "Sent #{message.type}"
  level = message.type == 'MessageNotAck' ? :warning : :log
  log(text, message: message, level: level)
end

#prepare_message(message, validate:) ⇒ Object



41
42
43
44
45
46
47
48
49
50
51
52
# File 'lib/rsmp/proxy/modules/send.rb', line 41

def prepare_message(message, validate:)
  message.direction = :out
  message.encode_for(schemas) unless validate == false
  message.generate_json
  unless validate == false
    validation = message.validate(schemas)
    return invalid_outbound_result(message, validation.message, validation: validation) if validation.invalid?
  end
  Result.success(message)
rescue RSMP::Schema::UnknownMessageCodeError => e
  invalid_outbound_result(message, e.message, cause: e)
end

#recover_write_failure(message, result, buffer:) ⇒ Object



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

def recover_write_failure(message, result, buffer:)
  return result unless result.failure.code == :disconnected

  buffered = buffer_message(message) if buffer
  return Result.success(Delivery.new(message: buffered, state: :buffered)) if buffered

  result
end

#send_generated_messageObject

Messages constructed by RSMP itself should never fail local schema validation. Surface that as an implementation defect while preserving connection and transport failures as ordinary Result values.



26
27
28
29
30
# File 'lib/rsmp/proxy/modules/send.rb', line 26

def send_generated_message(...)
  result = send_message(...)
  result.value! if result.failure&.source == :local
  result
end

#send_message(message, reason = nil, validate: true, force: false, buffer: true) ⇒ Object



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

def send_message(message, reason = nil, validate: true, force: false, buffer: true)
  return buffer_or_fail(message, buffer: buffer) unless force || connected?

  written = write_message(message, validate: validate)
  return recover_write_failure(message, written, buffer: buffer) if written.failure?

  expect_acknowledgement(message)
  distribute(message)
  log_send(message, reason)
  Result.success(Delivery.new(message: message, state: :sent))
end

#send_message!Object



19
20
21
# File 'lib/rsmp/proxy/modules/send.rb', line 19

def send_message!(...)
  send_message(...).value!
end

#send_message_and_collect(message, collector, validate: true) ⇒ Object



122
123
124
125
126
127
128
129
130
131
# File 'lib/rsmp/proxy/modules/send.rb', line 122

def send_message_and_collect(message, collector, validate: true)
  collector.start
  delivery = send_message(message, validate: validate, buffer: false)
  collector.fail(delivery.failure) if delivery.failure?
  collector.wait.map do |collection|
    Exchange.new(request: message, collection: collection)
  end
ensure
  collector.stop
end

#send_message_and_collect!Object



133
134
135
# File 'lib/rsmp/proxy/modules/send.rb', line 133

def send_message_and_collect!(...)
  send_message_and_collect(...).value!
end

#write_message(message, validate:) ⇒ Object



32
33
34
35
36
37
38
39
# File 'lib/rsmp/proxy/modules/send.rb', line 32

def write_message(message, validate:)
  return disconnected_send_result(message) unless @protocol

  prepared = prepare_message(message, validate: validate)
  return prepared if prepared.failure?

  write_protocol(message)
end

#write_protocol(message) ⇒ Object



54
55
56
57
58
59
60
61
62
63
64
65
# File 'lib/rsmp/proxy/modules/send.rb', line 54

def write_protocol(message)
  @protocol.write_lines(message.json)
  Result.success(message)
rescue IOError, SystemCallError => e
  Result.failure(
    :disconnected,
    message: "Could not send #{message.type}: #{e.message}",
    source: :transport,
    context: { message: message, session_id: session_id },
    cause: e
  )
end