Class: RSMP::Site
Overview
RSMP site implementation that manages proxies and components.
Defined Under Namespace
Classes: Options
Instance Attribute Summary collapse
Attributes included from Components
#components, #main
Attributes inherited from Node
#archive, #clock, #collector, #deferred, #task
Attributes included from EventSource
#event_receivers
Attributes included from Task
#task
Attributes included from Logging
#archive
Class Method Summary
collapse
Instance Method Summary
collapse
#accept_supervisor_connection, #accept_supervisor_connections, #build_proxies, #connect_to_supervisor, #listen_for_supervisors
Methods included from Components
#add_component, #check_main_component, #clear_alarm_timestamps, #component_list, #find_component, #infer_component_type, #initialize_components, #natural_sort_key, #setup_components
Methods inherited from Node
#author, #check_required_settings, #clear_deferred, #defer, #do_deferred, #inspect, #now, #process_deferred
#add_event_receiver, #initialize_event_source, #publish_event, #remove_event_receiver
Methods included from Task
#initialize_task, #restart, #start, #stop_task, #task_status, #wait, #wait_for_condition, #wait_for_condition!, #wait_for_termination
Methods included from Logging
#author, #initialize_logging, #log
Constructor Details
#initialize(options = {}) ⇒ Site
Returns a new instance of Site.
13
14
15
16
17
18
19
20
21
22
|
# File 'lib/rsmp/node/site/site.rb', line 13
def initialize(options = {})
super
initialize_components
handle_site_settings options
@proxies = []
@sleep_condition = Async::Notification.new
@proxies_condition = Async::Notification.new
@ready_condition = Async::Notification.new
build_proxies
end
|
Instance Attribute Details
#core_version ⇒ Object
Returns the value of attribute core_version.
7
8
9
|
# File 'lib/rsmp/node/site/site.rb', line 7
def core_version
@core_version
end
|
#logger ⇒ Object
Returns the value of attribute logger.
7
8
9
|
# File 'lib/rsmp/node/site/site.rb', line 7
def logger
@logger
end
|
#proxies ⇒ Object
Returns the value of attribute proxies.
7
8
9
|
# File 'lib/rsmp/node/site/site.rb', line 7
def proxies
@proxies
end
|
#ready_condition ⇒ Object
Returns the value of attribute ready_condition.
7
8
9
|
# File 'lib/rsmp/node/site/site.rb', line 7
def ready_condition
@ready_condition
end
|
#site_settings ⇒ Object
Returns the value of attribute site_settings.
7
8
9
|
# File 'lib/rsmp/node/site/site.rb', line 7
def site_settings
@site_settings
end
|
Class Method Details
.options_class ⇒ Object
9
10
11
|
# File 'lib/rsmp/node/site/site.rb', line 9
def self.options_class
RSMP::Site::Options
end
|
Instance Method Details
#aggregated_status_changed(component, _options = {}) ⇒ Object
130
131
132
133
134
|
# File 'lib/rsmp/node/site/site.rb', line 130
def aggregated_status_changed(component, _options = {})
@proxies.each do |proxy|
proxy.send_aggregated_status component
end
end
|
#alarm_acknowledged(alarm_state) ⇒ Object
136
137
138
|
# File 'lib/rsmp/node/site/site.rb', line 136
def alarm_acknowledged(alarm_state)
send_alarm AlarmAcknowledged.new(alarm_state.to_hash)
end
|
#alarm_activated_or_deactivated(alarm_state) ⇒ Object
144
145
146
|
# File 'lib/rsmp/node/site/site.rb', line 144
def alarm_activated_or_deactivated(alarm_state)
send_alarm AlarmIssue.new(alarm_state.to_hash)
end
|
#alarm_suspended_or_resumed(alarm_state) ⇒ Object
140
141
142
|
# File 'lib/rsmp/node/site/site.rb', line 140
def alarm_suspended_or_resumed(alarm_state)
send_alarm AlarmSuspended.new(alarm_state.to_hash)
end
|
#build_component(id:, type:, settings:) ⇒ Object
223
224
225
226
227
228
229
230
231
|
# File 'lib/rsmp/node/site/site.rb', line 223
def build_component(id:, type:, settings:)
settings ||= {}
if type == 'main'
Component.new id: id, node: self, type: type, name: settings['name'], grouped: true,
ntsoid: settings['ntsOId'], xnid: settings['xNId']
else
Component.new id: id, node: self, type: type, name: settings['name'], grouped: false
end
end
|
#check_core_versions ⇒ Object
81
82
83
84
85
86
87
88
89
90
|
# File 'lib/rsmp/node/site/site.rb', line 81
def check_core_versions
version = @site_settings['core_version']
return unless version
return if %w[all latest].include? version
return if RSMP::Schema.normalize_core_version version
error_str = "Unknown core version: #{version}"
raise RSMP::ConfigurationError, error_str
end
|
#check_sxls ⇒ Object
69
70
71
72
73
74
75
76
77
78
79
|
# File 'lib/rsmp/node/site/site.rb', line 69
def check_sxls
raise RSMP::ConfigurationError, 'No SXLs specified' unless sxls
sxls.each do |sxl|
name = sxl['name']
version = sxl['version'].to_s
raise RSMP::ConfigurationError, 'SXL name cannot be core' if name.to_s == 'core'
RSMP::Schema.find_schema! name, version, lenient: true
end
end
|
#client_role? ⇒ Boolean
40
41
42
|
# File 'lib/rsmp/node/site/site.rb', line 40
def client_role?
@site_settings['connection_role'] != 'server'
end
|
#denormalize_sxls(settings) ⇒ Object
60
61
62
63
64
65
66
67
|
# File 'lib/rsmp/node/site/site.rb', line 60
def denormalize_sxls(settings)
sxls = settings['sxls']
return settings unless sxls.is_a?(Array)
settings.merge(
'sxls' => sxls.to_h { |sxl| [sxl['name'], sxl['version']] }
)
end
|
#find_supervisor(ip) ⇒ Object
216
217
218
219
220
221
|
# File 'lib/rsmp/node/site/site.rb', line 216
def find_supervisor(ip)
@proxies.each do |supervisor|
return supervisor if ip == :any || supervisor.ip == ip
end
nil
end
|
#handle_site_settings(options = {}) ⇒ Object
48
49
50
51
52
53
54
55
56
57
58
|
# File 'lib/rsmp/node/site/site.rb', line 48
def handle_site_settings(options = {})
options_class = self.class.options_class
settings = options[:site_settings] || {}
settings = denormalize_sxls(settings)
@site_options = options_class.new(settings)
@site_settings = @site_options.to_h
check_sxls
check_core_versions
setup_components @site_settings['components']
end
|
#log_site_starting ⇒ Object
96
97
98
99
100
101
102
103
104
105
106
|
# File 'lib/rsmp/node/site/site.rb', line 96
def log_site_starting
log "Starting #{site_type_name} #{@site_settings['site_id']}", level: :info, timestamp: @clock.now
sxl = "Using SXLs #{sxls.map { |item| "#{item['name']} #{item['version']}" }.join(', ')}"
version = @site_settings['core_version']
core = if version
"accepting only core version #{version}"
else
"accepting all core versions [#{RSMP::Schema.core_versions.join(', ')}]"
end
log "#{sxl}, #{core}", level: :info, timestamp: @clock.now
end
|
#primary_sxl ⇒ Object
28
29
30
|
# File 'lib/rsmp/node/site/site.rb', line 28
def primary_sxl
sxls.first
end
|
#run ⇒ Object
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
|
# File 'lib/rsmp/node/site/site.rb', line 108
def run
log_site_starting
barrier = Async::Barrier.new(parent: @task)
termination_task = barrier.async { wait_for_termination }
start_auxiliary_tasks(barrier)
if server_role?
listen_for_supervisors(parent: barrier)
else
@proxies.each { |proxy| proxy.start(parent: barrier) }
end
barrier.wait do |finished|
result = finished.wait
break result if finished.equal?(termination_task)
end
ensure
barrier.cancel if barrier && !barrier.empty?
end
|
#run_status_timer(task, interval) ⇒ Object
165
166
167
168
169
170
171
172
173
174
|
# File 'lib/rsmp/node/site/site.rb', line 165
def run_status_timer(task, interval)
next_time = Time.now.to_f
loop do
now = Clock.now
tick_status_subscriptions now
next_time += interval
duration = [next_time - Time.now.to_f, 0].max
task.sleep duration
end
end
|
#send_alarm(alarm) ⇒ Object
148
149
150
151
152
|
# File 'lib/rsmp/node/site/site.rb', line 148
def send_alarm(alarm)
@proxies.each do |proxy|
proxy.send_generated_message alarm if proxy.receive_alarms?
end
end
|
#server_role? ⇒ Boolean
44
45
46
|
# File 'lib/rsmp/node/site/site.rb', line 44
def server_role?
@site_settings['connection_role'] == 'server'
end
|
#site_id ⇒ Object
36
37
38
|
# File 'lib/rsmp/node/site/site.rb', line 36
def site_id
@site_settings['site_id']
end
|
#site_type_name ⇒ Object
92
93
94
|
# File 'lib/rsmp/node/site/site.rb', line 92
def site_type_name
'site'
end
|
#start_auxiliary_tasks(parent) ⇒ Object
126
127
128
|
# File 'lib/rsmp/node/site/site.rb', line 126
def start_auxiliary_tasks(parent)
start_status_timer(parent: parent)
end
|
#start_status_timer(parent:) ⇒ Object
154
155
156
157
158
159
160
161
162
163
|
# File 'lib/rsmp/node/site/site.rb', line 154
def start_status_timer(parent:)
return if @status_timer
interval = @site_settings['intervals']['timer'] || 1
log "Starting site status timer with interval #{interval} seconds", level: :debug
@status_timer = parent.async do |task|
task.annotate 'site status timer'
run_status_timer task, interval
end
end
|
#stop ⇒ Object
195
196
197
198
|
# File 'lib/rsmp/node/site/site.rb', line 195
def stop
log "Stopping site #{@site_settings['site_id']}", level: :info
super
end
|
#stop_status_timer ⇒ Object
180
181
182
183
184
|
# File 'lib/rsmp/node/site/site.rb', line 180
def stop_status_timer
@status_timer&.cancel if @status_timer&.running?
ensure
@status_timer = nil
end
|
#stop_subtasks ⇒ Object
186
187
188
189
190
191
192
|
# File 'lib/rsmp/node/site/site.rb', line 186
def stop_subtasks
stop_status_timer
@accept_task&.cancel if @accept_task&.running?
@accept_task = nil
@endpoint = nil
super
end
|
#sxl_version ⇒ Object
32
33
34
|
# File 'lib/rsmp/node/site/site.rb', line 32
def sxl_version
primary_sxl && primary_sxl['version']
end
|
#sxls ⇒ Object
24
25
26
|
# File 'lib/rsmp/node/site/site.rb', line 24
def sxls
@site_settings['sxls']
end
|
#tick_status_subscriptions(now) ⇒ Object
176
177
178
|
# File 'lib/rsmp/node/site/site.rb', line 176
def tick_status_subscriptions(now)
@proxies.each { |proxy| proxy.status_update_timer now }
end
|
#wait_for_supervisor(ip, timeout:) ⇒ Object
200
201
202
203
204
205
206
|
# File 'lib/rsmp/node/site/site.rb', line 200
def wait_for_supervisor(ip, timeout:)
supervisor = find_supervisor ip
return Result.success(supervisor) if supervisor
wait_for_condition(@proxies_condition, timeout: timeout) { find_supervisor ip }
.map { |proxy| proxy }
end
|
#wait_for_supervisor!(ip, timeout:) ⇒ Object
208
209
210
211
212
213
214
|
# File 'lib/rsmp/node/site/site.rb', line 208
def wait_for_supervisor!(ip, timeout:)
result = wait_for_supervisor(ip, timeout: timeout)
return result.value! if result.success?
failure = result.failure.with(message: "Supervisor '#{ip}' did not connect within #{timeout}s")
Result.failure(failure: failure).value!
end
|