Module: RSMP::Supervisor::Modules::Connection

Included in:
RSMP::Supervisor
Defined in:
lib/rsmp/node/supervisor/modules/connection.rb

Overview

Handles incoming connections from sites

Instance Method Summary collapse

Instance Method Details

#accept?(_socket, _info) ⇒ Boolean

Returns:

  • (Boolean)


41
42
43
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 41

def accept?(_socket, _info)
  true
end

#accept_connection(socket, info) ⇒ Object



126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 126

def accept_connection(socket, info)
  log "Site connected from #{format_ip_and_port(info)}",
      ip: info[:ip],
      port: info[:port],
      level: :info,
      timestamp: Clock.now

  authorize_ip info[:ip]

  settings = build_proxy_settings(socket, info)
  id = settings[:site_id]
  proxy = setup_proxy(find_site(id), settings, id)

  validate_and_start_proxy(proxy, settings[:protocol])
ensure
  site_ids_changed
  stop if @supervisor_settings['one_shot']
end

#authorize_ip(ip) ⇒ Object

Raises:



53
54
55
56
57
58
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 53

def authorize_ip(ip)
  return if @supervisor_settings['ips'] == 'all'
  return if @supervisor_settings['ips'].include? ip

  raise ConnectionError, 'default ip not allowed'
end

#build_proxy_settings(socket, info) ⇒ Object



74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 74

def build_proxy_settings(socket, info)
  stream = IO::Stream::Buffered.new(socket)
  protocol = RSMP::Protocol.new stream
  site_id = retrieve_site_id(protocol)
  site_settings = site_id_to_site_setting site_id

  {
    supervisor: self,
    ip: info[:ip],
    port: info[:port],
    task: @task,
    collect: @collect,
    socket: socket,
    stream: stream,
    protocol: protocol,
    info: info,
    logger: @logger,
    archive: @archive,
    site_id: site_id,
    site_settings: site_settings
  }
end

#check_max_sitesObject

Raises:



60
61
62
63
64
65
66
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 60

def check_max_sites
  max = @supervisor_settings['max_sites']
  return unless max
  return unless @proxies.size >= max

  raise ConnectionError, "maximum of #{max} sites already connected"
end

#close(socket, info) ⇒ Object



149
150
151
152
153
154
155
156
157
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 149

def close(socket, info)
  if info
    log "Connection to #{format_ip_and_port(info)} closed", ip: info[:ip], level: :info, timestamp: Clock.now
  else
    log 'Connection closed', level: :info, timestamp: Clock.now
  end

  socket.close
end

#format_ip_and_port(info) ⇒ Object



45
46
47
48
49
50
51
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 45

def format_ip_and_port(info)
  if @logger.settings['hide_ip_and_port']
    '********'
  else
    "#{info[:ip]}:#{info[:port]}"
  end
end

#handle_connection(socket) ⇒ Object



6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 6

def handle_connection(socket)
  remote_port = socket.remote_address.ip_port
  remote_hostname = socket.remote_address.ip_address
  remote_ip = socket.remote_address.ip_address

  info = { ip: remote_ip, port: remote_port, hostname: remote_hostname, now: Clock.now }
  if accept? socket, info
    accept_connection socket, info
  else
    reject_connection socket, info
  end
rescue ConnectionError, HandshakeError => e
  report_connection_rejection(e, remote_ip, remote_port)
ensure
  close socket, info
end

#peek_version_message(protocol) ⇒ Object



68
69
70
71
72
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 68

def peek_version_message(protocol)
  json = protocol.peek_line
  attributes = Message.parse_attributes json
  Message.build attributes, json
end

#reject_connection(_socket, info) ⇒ Object



145
146
147
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 145

def reject_connection(_socket, info)
  log 'Site rejected', ip: info[:ip], level: :info
end

#report_connection_rejection(error, remote_ip, remote_port) ⇒ Object



23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 23

def report_connection_rejection(error, remote_ip, remote_port)
  log "Rejected connection from #{remote_ip}:#{remote_port}, #{error}", level: :warning
  failure = Failure.new(
    code: error.is_a?(HandshakeError) ? :handshake_failed : :connection_rejected,
    message: error.message,
    source: :peer,
    context: { ip: remote_ip, port: remote_port }
  )
  publish_event(
    Event.new(
      type: :connection_rejected,
      source: self,
      failure: failure,
      context: { ip: remote_ip, port: remote_port }
    )
  )
end

#retrieve_site_id(protocol) ⇒ Object



97
98
99
100
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 97

def retrieve_site_id(protocol)
  version_message = peek_version_message protocol
  version_message.attribute('siteId').first['sId']
end

#setup_proxy(proxy, settings, id) ⇒ Object



102
103
104
105
106
107
108
109
110
111
112
113
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 102

def setup_proxy(proxy, settings, id)
  if proxy
    raise ConnectionError, "Site #{id} already connected from port #{proxy.port}" if proxy.connected?

    proxy.revive settings
  else
    check_max_sites
    proxy = build_proxy settings
    @proxies.push proxy
  end
  proxy
end

#validate_and_start_proxy(proxy, protocol) ⇒ Object



115
116
117
118
119
120
121
122
123
124
# File 'lib/rsmp/node/supervisor/modules/connection.rb', line 115

def validate_and_start_proxy(proxy, protocol)
  proxy.setup_site_settings
  proxy.check_core_version peek_version_message(protocol)
  log "Validating using core version #{proxy.core_version}", level: :debug
  proxy.start
  proxy.wait
ensure
  proxy_type = proxy ? proxy.class.name.split('::').last : 'SiteProxy'
  log "Created #{proxy_type} for site #{proxy&.site_id}", level: :debug
end