diff --git a/setup.py b/setup.py index e8a5be7..5c31b79 100755 --- a/setup.py +++ b/setup.py @@ -11,5 +11,6 @@ data_files=[('etc', ['src/etc/openbmp-forwarder.yml'])], package_dir={'': 'src/site-packages'}, packages=['openbmp', 'openbmp.parsed'], - scripts=['src/bin/openbmp-forwarder'] - ) + scripts=['src/bin/openbmp-forwarder'], + install_requires=['ipaddress~=1.0.16', 'PyYAML~=3.11', 'python-snappy~=0.5', 'kafka-python~=1.2.2'] + ) diff --git a/src/bin/openbmp-forwarder b/src/bin/openbmp-forwarder index 53c5917..a3a886a 100755 --- a/src/bin/openbmp-forwarder +++ b/src/bin/openbmp-forwarder @@ -15,6 +15,8 @@ import logging import yaml import time import signal +import re +import ipaddress from multiprocessing import Queue, Manager from openbmp.logger import LoggerThread @@ -46,6 +48,14 @@ def signal_handler(signum, frame): RUNNING = False +def log_print(content, isError=False): + if LOG: + if isError: + LOG.error(content) + else: + LOG.debug(content) + else: + print (content) def load_config(cfg_filename, LOG): """ Load and validate the configuration from YAML @@ -80,34 +90,12 @@ def load_config(cfg_filename, LOG): cfg['kafka']['group_id'] = APP_NAME cfg['kafka']['offset_reset_largest'] = False - if 'collector' in cfg: - if 'host' not in cfg['collector']: - if LOG: - LOG.error("Configuration is missing 'host' in collector section") - else: - print ("Configuration is missing 'host' in collector section") - sys.exit(2) - - if 'port' not in cfg['collector']: - if LOG: - LOG.error("Configuration is missing 'port' in collector section, using default of 5000") - else: - print ("Configuration is missing 'port' in collector section, using default of 5000") - - cfg['collector']['port'] = 5000 - - else: - if LOG: - LOG.error("Configuration is missing 'collector' section.") - else: - print ("Configuration is missing 'collector' section.") + if 'dest_peer_groups' not in cfg: + log_print("Configuration is missing 'dest_peer_groups' section.", True) sys.exit(2) if 'logging' not in cfg: - if LOG: - LOG.error("Configuration is missing 'logging' section.") - else: - print ("Configuration is missing 'logging' section.") + log_print("Configuration is missing 'logging' section.", True) sys.exit(2) except (IOError, yaml.YAMLError), e: @@ -174,6 +162,51 @@ def parse_cmd_args(argv): return cfg +def parse_peer_group(cfg): + """ + Parse dest_peer_group in config. + """ + peer_groups = cfg['dest_peer_groups'] + for i, peer_exp in enumerate(peer_groups): + if 'name' not in peer_exp: + log_print('dest_peer_group with index {} is missing \'name\' section.'.format(i), True) + sys.exit(2) + name = peer_exp['name'] + log_print('Config: dest_peer_group name = {}'.format(name)) + if 'collector' not in peer_exp: + log_print('dest_peer_group {} is missing \'collector\' section.'.format(name), True) + sys.exit(2) + if type(peer_exp['regexp_hostname'] is list): + try: + peer_exp['regexp_hostname'] = [re.compile(regex) for regex in peer_exp['regexp_hostname']] + log_print('Config: compiled regexp hostname: {}'.format(peer_exp['regexp_hostname'])) + except re.error as err: + log_print('Invalid regular expression pattern: {}'.format(err), True) + sys.exit(2) + else: + log_print('Invalid regexp_hostname config: should be list.', True) + if type(peer_exp['prefix_range'] is list): + try: + peer_exp['prefix_range'] = [ipaddress.ip_network(unicode(prefix_exp)) for prefix_exp in + peer_exp['prefix_range']] + log_print('Config: parse prefix_range successful: {}'.format(peer_exp['prefix_range'])) + except Exception as err: + log_print('Invalid prefix range given: {}'.format(err), True) + sys.exit(2) + else: + log_print('Invalid prefix_range config: should be list.', True) + sys.exit(2) + if type(peer_exp['asn'] is list): + for asn in peer_exp['asn']: + if type(asn) is not int: + log_print("Invalid asn, must be int", True) + sys.exit(2) + else: + log_print('Invalid asn config: should be list.', True) + sys.exit(2) + cfg['dest_peer_groups'] = peer_groups + + def main(): """ Main entry point """ global LOG, RUNNING @@ -187,7 +220,8 @@ def main(): cfg_dict['max_queue_size'] = cfg['max_queue_size'] cfg_dict['logging'] = cfg['logging'] cfg_dict['kafka'] = cfg['kafka'] - cfg_dict['collector'] = cfg['collector'] + cfg_dict['dest_peer_groups'] = cfg['dest_peer_groups'] + cfg_dict['collector_heartbeat_interval'] = cfg['collector_heartbeat_interval'] # Setup signal handers signal.signal(signal.SIGTERM, signal_handler) @@ -199,8 +233,11 @@ def main(): thread_logger = LoggerThread(log_queue, cfg_dict['logging']) thread_logger.start() + logging.basicConfig() LOG = logging.getLogger() + parse_peer_group(cfg_dict) + # Use manager queue to ensure no duplicates forward_queue = manager.Queue(cfg_dict['max_queue_size']) diff --git a/src/etc/openbmp-forwarder.yml b/src/etc/openbmp-forwarder.yml index a4d3d22..24a2f70 100644 --- a/src/etc/openbmp-forwarder.yml +++ b/src/etc/openbmp-forwarder.yml @@ -15,11 +15,59 @@ kafka: offset_reset_largest: False # -# Collector - Where to send the BMP forwarded messages +# The number of seconds after last collector heartbeat to determine if the collector is dead +# Collector is considered dead after 1.1*THIS VALUE(int) after last heartbeat +# If a collector is dead, a series of PEER_DOWN will be generated to be sent to all corresponding dest_peer_group # -collector: - host: 10.1.1.1 - port: 5000 +collector_heartbeat_interval: 5000 + +# Peer group settings - Allow sending to multiple BMP destinations based on selected peers +# Order of matching +# Matching order is performed in the following sequence. The first match found is used. +# +# regexp_hostname - Hostname/regular expression is used first +# prefix_range - Prefix range is used second +# asn - Peer asn list +dest_peer_groups: + # name defines the value that is substituted for the variable. This provides a consistent + # mapping for different IP's and hostnames + - name: "lab" + + # You can specify which collector receives message about matched peers + collector: + host: 10.1.1.1 + port: 5000 + + # You can define a list of regexp's that match for hostname to group mapping + regexp_hostname: + - .*\.lab\..* + + # You can also define a list of prefixes that match for ip to group mapping + prefix_range: + - 10.100.100.0/24 + - 10.100.104.0/24 + + # You can define the matching to look at the peer asn. + asn: + - 100 + - 65000 + - 65001 + + # Keep this, it's the default entry + - name: "default" + + collector: + host: 10.1.1.1 + port: 5000 + + regexp_hostname: + - .* + + prefix_range: + - 0.0.0.0/0 + + asn: + - 0 # # Log settings @@ -61,7 +109,6 @@ logging: handlers: [file] propagate: no - # General/main program messages root: level: INFO diff --git a/src/site-packages/openbmp/bmp.py b/src/site-packages/openbmp/bmp.py new file mode 100644 index 0000000..ed5ebde --- /dev/null +++ b/src/site-packages/openbmp/bmp.py @@ -0,0 +1,138 @@ +"""OpenBMP MRT + + Copyright (c) 2013-2016 Cisco Systems, Inc. and others. All rights reserved. + This program and the accompanying materials are made available under the + terms of the Eclipse Public License v1.0 which accompanies this distribution, + and is available at http://www.eclipse.org/legal/epl-v10.html + + .. moduleauthor:: Tim Evens +""" +import socket + +from struct import unpack + +def bmp_parse_peerhdr(data): + """ Parse BMP peer header + + :param data: BMP raw data - should start at peer header + + :return: dictionary defined as:: + { + type: , + dist_id: , + addr: , + asn: , + bgp_id: , + isIPv4: , + isPrePolicy + is2ByteASN + ts_secs: , + ts_usecs: + } + """ + hdr = { 'type': None, + 'raw_type': 0, + 'flags': None, + 'dist_id': 0, + 'addr': None, + 'asn': 0, + 'bgp_id': 0, + 'isIPv4': True, + 'isPrePolicy': True, + 'is2ByteASN': False, + 'ts_secs': 0, + 'ts_usecs': 0} + + (hdr['raw_type'], hdr['flags'], hdr['dist_id']) = unpack('>BBQ', data[:10]) + + if hdr['raw_type'] == 0: + hdr['type'] = 'GLOBAL' + else: + hdr['type'] = 'L3VPN' + + if hdr['flags'] & 0x80: # V flag + hdr['isIPv4'] = False + else: + hdr['isIPv4'] = True + + if hdr['flags'] & 0x40: # L flag + hdr['isPrePolicy'] = False + else: + hdr['isPrePolicy'] = True + + if hdr['flags'] & 0x20: # A flag + hdr['is2ByteASN'] = True + else: + hdr['is2ByteASN'] = False + + if hdr['isIPv4']: + hdr['addr'] = socket.inet_ntop(socket.AF_INET, data[22:26]) + else: + hdr['addr'] = socket.inet_ntop(socket.AF_INET6, data[10:26]) + + (hdr['asn'],) = unpack('>I', data[26:30]) + + hdr['bgp_id'] = socket.inet_ntop(socket.AF_INET, data[30:34]) + + (hdr['ts_secs'], hdr['ts_usecs']) = unpack('>II', data[34:42]) + + return hdr + + +def bmp_parse_bmphdr(data): + """ Parse BMP header from message string + + :param data: RAW BMP message - should start at bmp header + + :return: dictionary defined as:: + { + version: , + length: , + type: , + } + """ + hdr = { 'version': None, + 'length': 0, + 'type': None } + + if not data: + return None + + (hdr['version'],) = unpack('B', data[:1]) + + if hdr['version'] == 3: + (hdr['length'], type) = unpack('>IB', data[1:6]) + hdr['length'] -= 6 # remove the bytes of the common header + + if type == 0: + hdr['type'] = 'ROUTE_MON' + elif type == 1: + hdr['type'] = 'STATS_REPORT' + elif type == 2: + hdr['type'] = 'PEER_DOWN' + elif type == 3: + hdr['type'] = 'PEER_UP' + elif type == 4: + hdr['type'] = 'INIT' + elif type == 5: + hdr['type'] = 'TERM' + else: + hdr['type'] = "UNKNOWN=%d" % type + + else: + self.LOG.error("Unsupported BMP version type of %d, cannot proceed" % hdr['version']) + return None + + return hdr + + +def resolveIp(addr): + """ Resolves an IP address to FQDN. + + :param addr: IPv4/v6 address to resovle + :return: FQDN or IP address if FQDN unknown/not found + """ + try: + return socket.gethostbyaddr(addr)[0] + except: + return addr diff --git a/src/site-packages/openbmp/bmp_consumer.py b/src/site-packages/openbmp/bmp_consumer.py index 78c8453..14e7bfe 100644 --- a/src/site-packages/openbmp/bmp_consumer.py +++ b/src/site-packages/openbmp/bmp_consumer.py @@ -12,16 +12,18 @@ import socket import kafka import kafka.common -import copy import time import re - +from datetime import datetime +from struct import pack +from threading import Timer from openbmp.logger import init_mp_logger from openbmp.parsed.headers import headers as parsed_headers from openbmp.parsed.collector import collector from openbmp.parsed.router import router - +from openbmp.parsed.peer import peer +from openbmp.bmp import bmp_parse_bmphdr, bmp_parse_peerhdr class BMPMessageObject(object): """ OpenBMP consumer message object which is used between internal queues """ @@ -37,6 +39,12 @@ class BMPMessageObject(object): #: Router Name ROUTER_NAME = "" + DEST_CONFIG_NAME = "" + + DEST_COLLECTOR_HOST = "" + + DEST_COLLECTOR_PORT = 0 + class BMPConsumer(multiprocessing.Process): """ OpenBMP consumer for forwarding @@ -51,6 +59,9 @@ class BMPConsumer(multiprocessing.Process): # Memory cache of routers ROUTERS = {} + # Memory cache of peers + PEERS = {} + def __init__(self, cfg, forward_queue, log_queue): """ Constructor @@ -71,7 +82,7 @@ def run(self): self.LOG = init_mp_logger("bmp_consumer", self._log_queue) # Enable to topics/feeds - topics = ['openbmp.parsed.collector', 'openbmp.parsed.router', 'openbmp.bmp_raw'] + topics = ['openbmp.parsed.collector', 'openbmp.parsed.router', 'openbmp.parsed.peer', 'openbmp.bmp_raw'] self.LOG.info("Running bmp_consumer") @@ -108,6 +119,7 @@ def run(self): except kafka.common.KafkaUnavailableError as err: self.LOG.error("Kafka Error: %s" % str(err)) + self._fwd_queue.put('DISCONNECT ALL') except KeyboardInterrupt: pass @@ -136,6 +148,9 @@ def process_msg(self, msg): if msg.topic == 'openbmp.parsed.router': self.process_router_msg(hdr.getCollectorHashId(), data) + if msg.topic == 'openbmp.parsed.peer': + self.process_peer_msg(hdr.getCollectorHashId(), data) + elif msg.topic == 'openbmp.bmp_raw': self.process_bmp_raw_msg(hdr.getCollectorHashId(), hdr.getRouterHashId(), hdr.getRouterIp(), data) @@ -160,7 +175,13 @@ def process_collector_msg(self, c_hash, data): # Update collector hash/cache if obj.getHashId() not in self.COLLECTORS: self.COLLECTORS[obj.getHashId()] = {'admin_id': obj.getAdminId()} - + if obj.getAction() == 'heartbeat': + cur_collector = self.COLLECTORS[obj.getHashId()] + expire_time = self._cfg['collector_heartbeat_interval'] * 1.1 + if 'timer' in cur_collector: + cur_collector['timer'].cancel() + cur_collector['timer'] = Timer(expire_time, self.handle_collector_down, [obj.getHashId()]) + cur_collector['timer'].start() else: self.COLLECTORS.pop(obj.getHashId(), None) @@ -188,14 +209,59 @@ def process_router_msg(self, c_hash, data): # Update the router hash/cache if obj.getHashId() not in self.ROUTERS: self.ROUTERS[obj.getHashId()] = {'ip': obj.getIpAddress(), - 'name': obj.getName()} - # else: - # self.ROUTERS.pop(obj.getHashId(), None) + 'name': obj.getName(), + 'hash': obj.getHashId(), + 'collector_hash': c_hash} except: self.LOG.debug("router parse error"); pass + def process_peer_msg(self, c_hash, data): + """ Process Peer message + + :param c_hash: Collector Hash ID + :param data: Message data to be consumed (should not contain headers) + """ + obj = peer() + + # Log messages + for row in data.split('\n'): + if len(row): + try: + obj.parse(row) + + peer_key = obj.getRouterHashId() + '_' + obj.getRemoteBgpId() + '_' + str(obj.getPeerRd()) + + self.LOG.info("peer: [%s] %s %s", obj.getAction(), + obj.getName(), obj.getLocalIp()) + + if obj.getAction() in ('up'): + + asn_len = 2 + if ('4 Octet ASN' in obj.getAdvCapabilities() and + '4 Octet ASN' in obj.getRecvCapabilities()): + asn_len = 4 + + # Update the peer hash/cache + # replace raw data with this one contains peer_name information + if obj.getAction() != 'down': + if peer_key in self.PEERS: + self.PEERS[peer_key] = dict(({'peer_name': obj.getName(), + 'remote_ip': obj.getRemoteIp(), + 'remote_asn': obj.getRemoteAsn(), + 'local_ip': obj.getLocalIp(), + 'local_asn': obj.getLocalAsn(), + 'router_hash': obj.getRouterHashId(), + 'asn_len': asn_len, + 'rd': obj.getPeerRd() + }.get(k, k), v) for (k, v) in + self.PEERS[peer_key].items()) + + except: + self.LOG.debug("peer parse error") + pass + def process_bmp_raw_msg(self, c_hash, r_hash, r_ip, data): """ Process BMP RAW message @@ -206,13 +272,165 @@ def process_bmp_raw_msg(self, c_hash, r_hash, r_ip, data): """ msg = BMPMessageObject() + # Parse the BMP headers + bmp_hdrs = bmp_parse_bmphdr(data) + + if bmp_hdrs['type'] in ('INIT', 'TERM'): + # Send to ALL + for dest_peer_group in self._cfg['dest_peer_groups']: + if c_hash in self.COLLECTORS: + msg.COLLECTOR_ADMIN_ID = self.COLLECTORS[c_hash]['admin_id'] + + if r_hash in self.ROUTERS: + msg.ROUTER_IP = self.ROUTERS[r_hash]['ip'] + msg.ROUTER_NAME = self.ROUTERS[r_hash]['name'] + + msg.DEST_CONFIG_NAME = dest_peer_group['name'] + msg.DEST_COLLECTOR_HOST = dest_peer_group['collector']['host'] + msg.DEST_COLLECTOR_PORT = dest_peer_group['collector']['port'] + msg.BMP_MSG = data + self._fwd_queue.put(msg) + + if bmp_hdrs['type'] == 'TERM': + self.handle_router_down(c_hash, r_hash) + + else: + peer_hdr = bmp_parse_peerhdr(data[6:]) + + peer_key = r_hash + '_' + peer_hdr['bgp_id'] + '_' + str(peer_hdr['dist_id']) + + if c_hash in self.COLLECTORS: + msg.COLLECTOR_ADMIN_ID = self.COLLECTORS[c_hash]['admin_id'] + + if r_hash in self.ROUTERS: + msg.ROUTER_IP = self.ROUTERS[r_hash]['ip'] + msg.ROUTER_NAME = self.ROUTERS[r_hash]['name'] + + if peer_key not in self.PEERS: + self.PEERS[peer_key] = {'type': peer_hdr['raw_type'], + 'flags': peer_hdr['flags'], + 'bgp_id': peer_hdr['bgp_id'], + 'isIPv4': peer_hdr['isIPv4'], + 'peer_name': '', + 'remote_ip': '', + 'remote_asn': '', + 'local_ip': peer_hdr['addr'], + 'local_asn': peer_hdr['asn'], + 'router_hash': r_hash, + 'asn_len': 0, + 'rd': peer_hdr['dist_id']} + + current_peer = self.PEERS[peer_key] + + # Send to forward queue + self.LOG.debug('Doing match for PEER KEY: {}'.format(peer_key)) + + dest_peer_group = self.match_peer_group(current_peer) + if dest_peer_group is not None: + self.LOG.debug('PEER KEY: {} matched dest_peer_group {}'.format(peer_key, dest_peer_group['name'])) + msg.DEST_CONFIG_NAME = dest_peer_group['name'] + msg.DEST_COLLECTOR_HOST = dest_peer_group['collector']['host'] + msg.DEST_COLLECTOR_PORT = dest_peer_group['collector']['port'] + msg.BMP_MSG = data + self._fwd_queue.put(msg) + else: + self.LOG.error('PEER KEY: {} did not match any dest_peer_group!'.format(peer_key)) + + def match_peer_group(self, peer_dict): + """ + + :param peer_dict: the self.PEERS['key'] dictionary + :return: the dest_group_config matching this peer + """ + peer_groups = self._cfg['dest_peer_groups'] try: - msg.COLLECTOR_ADMIN_ID = self.COLLECTORS['admin_id'] - msg.ROUTER_IP = self.ROUTERS['ip'] - msg.ROUTER_NAME = self.ROUTERS['name'] + for peer_exp in peer_groups: + for regex in peer_exp['regexp_hostname']: + if re.match(regex, peer_dict['peer_name']) is not None: + return peer_exp + for prefix in peer_exp['prefix_range']: + if peer_dict['local_ip'] in prefix.hosts(): + return peer_exp + if peer_dict['local_asn'] in peer_exp['asn']: + return peer_exp + return None + except Exception() as err: + self.LOG.error('Something went wrong: {}'.format(err)) + + def handle_collector_down(self, c_hash): + """ Generate PEER DOWN messages for the peers under routers which ones are under this collector + + :param c_hash: The collector which is down + """ + for key in self.ROUTERS: + if self.ROUTERS[key]['collector_hash'] == c_hash: + self.handle_router_down(c_hash, key) + self.COLLECTORS.pop(c_hash, None) - except: - pass + def handle_router_down(self, c_hash, r_hash): + """ Generate PEER DOWN messages for the peers under a certain router + + :param c_hash: the collector hash which the router TERM comes from + :param r_hash: the TERMed router + """ + # Generate Peer Down notifications to send to peer groups + peer_list = [] + + for key, peer in self.PEERS.iteritems(): + if peer['router_hash'] == r_hash: + peer_list.append(key) + + for peer_key in peer_list: + if peer_key in self.PEERS: + current_peer = self.PEERS[peer_key] + msg = BMPMessageObject() + + peer_header = pack('>BBQ', current_peer['type'], current_peer['flags'], current_peer['rd']) + + if current_peer['isIPv4']: + peer_header += socket.inet_pton(socket.AF_INET, current_peer['local_ip']) + pack( + '>BBBBBBBBBBBB', 0, 0, + 0, 0, 0, 0, 0, 0, 0, 0, + 0, 0) + else: + peer_header += socket.inet_pton(socket.AF_INET6, current_peer['local_ip']) + + peer_header += pack('>I', current_peer['local_asn']) + + peer_header += socket.inet_pton(socket.AF_INET, current_peer['bgp_id']) + + peer_header += pack('>II', + int(time.time()), + datetime.now().microsecond) + + msg_body = pack('>B', 4) + + common_header = pack('>BIB', 3, len(peer_header) + len(msg_body) + 6, 2) + + if c_hash in self.COLLECTORS: + msg.COLLECTOR_ADMIN_ID = self.COLLECTORS[c_hash]['admin_id'] + + if r_hash in self.ROUTERS: + msg.ROUTER_IP = self.ROUTERS[r_hash]['ip'] + msg.ROUTER_NAME = self.ROUTERS[r_hash]['name'] + + # Send to forward queue + self.LOG.debug( + 'Doing match for generated peer down notification with PEER KEY: {}'.format(peer_key)) + + dest_peer_group = self.match_peer_group(current_peer) + if dest_peer_group is not None: + self.LOG.debug( + 'GENERATED PEER DOWN: {} matched dest_peer_group {}'.format(peer_key, + dest_peer_group['name'])) + msg.DEST_CONFIG_NAME = dest_peer_group['name'] + msg.DEST_COLLECTOR_HOST = dest_peer_group['collector']['host'] + msg.DEST_COLLECTOR_PORT = dest_peer_group['collector']['port'] + msg.BMP_MSG = common_header + peer_header + msg_body + self._fwd_queue.put(msg) + else: + self.LOG.error('GENERATED PEER DOWN: {} did not match any dest_peer_group!'.format(peer_key)) + + self.PEERS.pop(peer_key, None) - msg.BMP_MSG = data - self._fwd_queue.put(msg) + self.ROUTERS.pop(r_hash, None) \ No newline at end of file diff --git a/src/site-packages/openbmp/forwarder_bmp.py b/src/site-packages/openbmp/forwarder_bmp.py index 989bb6f..c5d3386 100644 --- a/src/site-packages/openbmp/forwarder_bmp.py +++ b/src/site-packages/openbmp/forwarder_bmp.py @@ -21,6 +21,8 @@ class BMPWriter(multiprocessing.Process): Pops messages from forwarder queue and transmits them to remote bmp collector. """ + SOCKET_TIMEOUT = 5 + def __init__(self, cfg, forward_queue, log_queue): """ Constructor @@ -35,9 +37,7 @@ def __init__(self, cfg, forward_queue, log_queue): self._fwd_queue = forward_queue self._log_queue = log_queue self.LOG = None - self._isConnected = False - - self._sock = None + self._dest_socks = {} def run(self): """ Override """ @@ -57,42 +57,58 @@ def run(self): while not self.stopped(): # Do not pop any message unless connected - if self._isConnected: - msg = self._fwd_queue.get() + msg = self._fwd_queue.get() - sent = False - while not sent: - sent = self.send(msg.BMP_MSG) + if msg == 'DISCONNECT ALL': + self.disconnect() - self.LOG.debug("Received bmp message: %s %s %s", msg.COLLECTOR_ADMIN_ID, - msg.ROUTER_IP, msg.ROUTER_NAME) else: - self.LOG.info("Not connected, attempting to reconnect") - sleep(1) - self.connect() + if msg.DEST_CONFIG_NAME not in self._dest_socks: + self.connect(msg.DEST_CONFIG_NAME) + + sock_arr = self._dest_socks[msg.DEST_CONFIG_NAME] + if sock_arr[1]: + sent = False + while not sent: + sent = self.send(msg.BMP_MSG, sock_arr[0]) + + self.LOG.debug("Received bmp message: %s %s %s", msg.COLLECTOR_ADMIN_ID, + msg.ROUTER_IP, msg.ROUTER_NAME) + else: + self.LOG.info("Peer group \'{}\' not connected, attempting to reconnect".format(msg.DEST_CONFIG_NAME)) + sock_arr[1] = self.connect(msg.DEST_CONFIG_NAME) except KeyboardInterrupt: pass self.LOG.info("rewrite stopped") - def connect(self): - """ Connect to remote collector + def connect(self, peer_group_name=None): + """ Connect to remote controller - :return: True if connected, False otherwise/error + :param peer_group_name: the name of peer_group config. If this argument is not passed, connect all socks found + in config. + :return: if a peer_group_name is given, return connection status on that one """ - try: - self._sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - self._sock.connect((self._cfg['collector']['host'], self._cfg['collector']['port'])) - self._isConnected = True - self.LOG.info("Connected to remote collector: %s:%d", self._cfg['collector']['host'], - self._cfg['collector']['port']) - - except socket.error as msg: - self.LOG.error("Failed to connect to remote collector: %r", msg) - self._isConnected = False - - def send(self, msg): + for peer_group in self._cfg['dest_peer_groups']: + if peer_group_name is None or peer_group_name == peer_group['name']: + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + is_connected = None + sock.settimeout(self.SOCKET_TIMEOUT) + try: + sock.connect((peer_group['collector']['host'], peer_group['collector']['port'])) + is_connected = True + self.LOG.info("Connected to remote collector: %s:%d", peer_group['collector']['host'], + peer_group['collector']['port']) + except socket.error as msg: + self.LOG.error("Failed to connect to remote collector: %r", msg) + is_connected = False + finally: + self._dest_socks[peer_group['name']] = [sock, is_connected] + if peer_group_name is not None: + return is_connected + + def send(self, msg, sock): """ Send BMP message to socket. :param msg: Message to send/write @@ -102,25 +118,20 @@ def send(self, msg): sent = False try: - self._sock.sendall(msg) + sock.sendall(msg) sent = True except socket.error as msg: self.LOG.error("Failed to send message to collector: %r", msg) - self.disconnect() - sleep(1) - self.connect() return sent def disconnect(self): """ Disconnect from remote collector """ - if self._sock: - self._sock.close() - self._sock = None - - self._isConnected = False + for key, sock_arr in self._dest_socks: + sock_arr[0].close() + del self._dest_socks[key] def stop(self): self._stop.set()