Module: RSMP::SiteProxy::Modules::Status

Included in:
RSMP::SiteProxy
Defined in:
lib/rsmp/proxy/site/modules/status.rb

Overview

Handles status requests, responses, subscriptions and updates

Instance Method Summary collapse

Instance Method Details

#ensure_subscription_path(component_id, code, name) ⇒ Object



68
69
70
71
72
# File 'lib/rsmp/proxy/site/modules/status.rb', line 68

def ensure_subscription_path(component_id, code, name)
  @status_subscriptions[component_id] ||= {}
  @status_subscriptions[component_id][code] ||= {}
  @status_subscriptions[component_id][code][name] ||= {}
end

#process_status_response(message) ⇒ Object



61
62
63
64
65
66
# File 'lib/rsmp/proxy/site/modules/status.rb', line 61

def process_status_response(message)
  component = find_component message.attribute('cId')
  component.store_status message
  log "Received #{message.type}", message: message, level: :log
  acknowledge message
end

#process_status_update(message) ⇒ Object



183
184
185
186
187
188
189
# File 'lib/rsmp/proxy/site/modules/status.rb', line 183

def process_status_update(message)
  component = find_component message.attribute('cId')
  component.check_repeat_values message, @status_subscriptions
  component.store_status message
  log "Received #{message.type}", message: message, level: :log
  acknowledge message
end

#remove_subscription_item(component_id, code, name) ⇒ Object



139
140
141
142
143
144
145
# File 'lib/rsmp/proxy/site/modules/status.rb', line 139

def remove_subscription_item(component_id, code, name)
  return unless @status_subscriptions.dig(component_id, code, name)

  @status_subscriptions[component_id][code].delete name
  @status_subscriptions[component_id].delete(code) if @status_subscriptions[component_id][code].empty?
  @status_subscriptions.delete(component_id) if @status_subscriptions[component_id].empty?
end

#request_status(status_list, component: nil, m_id: nil, validate: true) ⇒ Object

Build and send a StatusRequest. Returns Result.



7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
# File 'lib/rsmp/proxy/site/modules/status.rb', line 7

def request_status(status_list, component: nil, m_id: nil, validate: true)
  readiness = validate_ready 'request status'
  return readiness if readiness.failure?

  component ||= main.c_id
  m_id ||= RSMP::Message.make_m_id

  list = RSMP::StatusList.new(status_list)

  # additional items can be used when verifying the response,
  # but must be removed from the request
  request_list = list.map { |item| item.slice('sCI', 'n') }

  message = RSMP::StatusRequest.new({
                                      'cId' => component,
                                      'sS' => request_list,
                                      'mId' => m_id
                                    })
  apply_nts_message_attributes message
  send_message(message, validate: validate).map(&:message)
end

#request_status!Object



29
30
31
# File 'lib/rsmp/proxy/site/modules/status.rb', line 29

def request_status!(...)
  request_status(...).value!
end

#request_status_and_collect(status_list, within:, component: nil, m_id: nil, validate: true) ⇒ Object

Build, send a StatusRequest and collect the StatusResponse. Returns Result.



34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
# File 'lib/rsmp/proxy/site/modules/status.rb', line 34

def request_status_and_collect(status_list, within:, component: nil, m_id: nil, validate: true)
  readiness = validate_ready 'request status'
  return readiness if readiness.failure?

  component ||= main.c_id
  m_id ||= RSMP::Message.make_m_id

  list = RSMP::StatusList.new(status_list)

  # additional items can be used when verifying the response,
  # but must be removed from the request
  request_list = list.map { |item| item.slice('sCI', 'n') }

  message = RSMP::StatusRequest.new({
                                      'cId' => component,
                                      'sS' => request_list,
                                      'mId' => m_id
                                    })
  apply_nts_message_attributes message
  collector = StatusCollector.new(self, list.to_a, timeout: within, m_id: m_id)
  send_message_and_collect(message, collector, validate: validate)
end

#request_status_and_collect!Object



57
58
59
# File 'lib/rsmp/proxy/site/modules/status.rb', line 57

def request_status_and_collect!(...)
  request_status_and_collect(...).value!
end

#subscribe_to_status(status_list, component: nil, m_id: nil, validate: true) ⇒ Object

Build and send a StatusSubscribe. Returns Result.



85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
# File 'lib/rsmp/proxy/site/modules/status.rb', line 85

def subscribe_to_status(status_list, component: nil, m_id: nil, validate: true)
  readiness = validate_ready 'subscribe to status'
  return readiness if readiness.failure?

  component ||= main.c_id
  m_id ||= RSMP::Message.make_m_id

  list = RSMP::StatusList.new(status_list)
  subscribe_list = list.map { |item| item.slice('sCI', 'n', 'uRt', 'sOc') }

  update_subscription(component, subscribe_list)
  find_component component

  message = RSMP::StatusSubscribe.new({
                                        'cId' => component,
                                        'sS' => subscribe_list,
                                        'mId' => m_id
                                      })
  apply_nts_message_attributes message
  send_message(message, validate: validate).map(&:message)
end

#subscribe_to_status!Object



107
108
109
# File 'lib/rsmp/proxy/site/modules/status.rb', line 107

def subscribe_to_status!(...)
  subscribe_to_status(...).value!
end

#subscribe_to_status_and_collect(status_list, within:, component: nil, m_id: nil, validate: true) ⇒ Object

Build, send a StatusSubscribe and collect the first matching status update. Returns Result.



112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
# File 'lib/rsmp/proxy/site/modules/status.rb', line 112

def subscribe_to_status_and_collect(status_list, within:, component: nil, m_id: nil, validate: true)
  readiness = validate_ready 'subscribe to status'
  return readiness if readiness.failure?

  component ||= main.c_id
  m_id ||= RSMP::Message.make_m_id

  list = RSMP::StatusList.new(status_list)
  subscribe_list = list.map { |item| item.slice('sCI', 'n', 'uRt', 'sOc') }

  update_subscription(component, subscribe_list)
  find_component component

  message = RSMP::StatusSubscribe.new({
                                        'cId' => component,
                                        'sS' => subscribe_list,
                                        'mId' => m_id
                                      })
  apply_nts_message_attributes message
  collector = StatusCollector.new(self, list.to_a, timeout: within, m_id: m_id)
  send_message_and_collect(message, collector, validate: validate)
end

#subscribe_to_status_and_collect!Object



135
136
137
# File 'lib/rsmp/proxy/site/modules/status.rb', line 135

def subscribe_to_status_and_collect!(...)
  subscribe_to_status_and_collect(...).value!
end

#unsubscribe_from_all(component: nil) ⇒ Object

unsubscribes to all statuses (with all attributes) defined in the used SXL



171
172
173
174
175
176
177
178
179
180
181
# File 'lib/rsmp/proxy/site/modules/status.rb', line 171

def unsubscribe_from_all(component: nil)
  component ||= main.c_id
  catalogue = accepted_sxls.each_with_object({}) do |sxl, memo|
    version = RSMP::Schema.sanitize_version(sxl['version'].to_s)
    memo.merge! RSMP::Schema.status_catalogue(sxl['name'], version)
  end
  status_list = catalogue.flat_map do |status_code_id, names|
    names.map { |name| { 'sCI' => status_code_id.to_s, 'n' => name.to_s } }
  end
  unsubscribe_to_status status_list, component: component
end

#unsubscribe_to_status(status_list, component: nil, validate: nil) ⇒ Object



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
# File 'lib/rsmp/proxy/site/modules/status.rb', line 147

def unsubscribe_to_status(status_list, component: nil, validate: nil)
  component ||= main.c_id

  status_list.each do |item|
    remove_subscription_item(component, item['sCI'], item['n'])
  end

  # Local subscription state must be cleaned up even when the peer has
  # already disconnected. In that case there is nothing left to send.
  return Result.success(nil) unless ready?

  message = RSMP::StatusUnsubscribe.new({
                                          'cId' => component,
                                          'sS' => status_list
                                        })
  apply_nts_message_attributes message
  send_message(message, validate: validate).map(&:message)
end

#unsubscribe_to_status!Object



166
167
168
# File 'lib/rsmp/proxy/site/modules/status.rb', line 166

def unsubscribe_to_status!(...)
  unsubscribe_to_status(...).value!
end

#update_subscription(component_id, subscribe_list) ⇒ Object



74
75
76
77
78
79
80
81
82
# File 'lib/rsmp/proxy/site/modules/status.rb', line 74

def update_subscription(component_id, subscribe_list)
  subscribe_list.each do |item|
    code = item['sCI']
    name = item['n']
    sub = ensure_subscription_path(component_id, code, name)
    sub['uRt'] = item['uRt']
    sub['sOc'] = item['sOc']
  end
end