From 5827d8b137ebe333432c68369a4ec0a1403859b0 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Tue, 9 Oct 2018 13:32:41 -0400 Subject: [PATCH 1/9] Initial commit --- parsedmarc/cli.py | 42 +++++++++++++++++++++++++++++--- parsedmarc/kafkaclient.py | 51 +++++++++++++++++++++++++++++++++++++++ requirements.txt | 1 + 3 files changed, 91 insertions(+), 3 deletions(-) create mode 100644 parsedmarc/kafkaclient.py diff --git a/parsedmarc/cli.py b/parsedmarc/cli.py index fcf3d625..65a88e7a 100644 --- a/parsedmarc/cli.py +++ b/parsedmarc/cli.py @@ -13,10 +13,9 @@ import json from elasticsearch.exceptions import ElasticsearchException from parsedmarc import logger, IMAPError, get_dmarc_reports_from_inbox, \ - parse_report_file, elastic, splunk, save_output, watch_inbox, \ + parse_report_file, elastic, kafkaclient, splunk, save_output, watch_inbox, \ email_results, SMTPError, ParserError, __version__ - def _main(): """Called when the module is executed""" def process_reports(reports_): @@ -25,6 +24,12 @@ def _main(): indent=2)) if not args.silent: print(output_str) + if args.kafka_hosts: + try: + kafkaClient = kafkaclient.KafkaClient(args.kafka_hosts) + # dont do this + except Exception as error: + logger.error("Kafka Error: {0}".format(error.__str__())) if args.save_aggregate: for report in reports_["aggregate_reports"]: try: @@ -37,6 +42,14 @@ def _main(): logger.error("Elasticsearch Error: {0}".format( error_.__str__())) exit(1) + try: + if args.kafka_hosts: + kafkaClient.save_aggregate_reports_to_kafka( + report, kafka_aggregate_topic) + # dont do this catch specific exceptions + except Exception as error_: + logger.error("Kafka Error: {0}".format( + error_.__str__())) if args.hec: try: aggregate_reports_ = reports_["aggregate_reports"] @@ -55,6 +68,14 @@ def _main(): except ElasticsearchException as error_: logger.error("Elasticsearch Error: {0}".format( error_.__str__())) + try: + if args.kafka_hosts: + kafkaClient.save_forensic_reports_to_kafka( + report, kafka_forensic_topics) + # dont do this + except Exception as error_: + logger.error("Kafka Error: {0}".format( + error_.__str__())) if args.hec: try: forensic_reports_ = reports_["forensic_reports"] @@ -122,6 +143,12 @@ def _main(): default=False, help="Skip certificate verification for Splunk " "HEC") + arg_parser.add_argument("-K", "--kafka-hosts", nargs="*", + help="A list of one or more Kafka hostnames or URLs") + arg_parser.add_argument("--kafka-aggregate-topic", + help="The Kafka topic to publish aggregate reports to.") + arg_parser.add_argument("--kafka-forensic_topic", + help="The Kafka topic to publish forensic reports to.") arg_parser.add_argument("--save-aggregate", action="store_true", default=False, help="Save aggregate reports to search indexes") @@ -191,7 +218,7 @@ def _main(): es_forensic_index = "{0}_{1}".format(es_forensic_index, suffix) if args.save_aggregate or args.save_forensic: - if args.elasticsearch_host is None and args.hec is None: + if args.elasticsearch_host is None and args.hec and args.kafka_hosts is None: args.elasticsearch_host = ["localhost:9200"] try: if args.elasticsearch_host: @@ -214,6 +241,15 @@ def _main(): args.hec_index, verify=verify) + kafka_aggregate_topic = "dmarc_aggrregate" + kafka_forensic_topic = "dmarc_forensic" + + if args.kafka_aggregate_topic: + kafka_aggregate_topic = args.kafka_aggregate_topic + + if args.kafka_forensic_topic: + kafka_forensic_topic = args.kafka_forensic_topic + file_paths = [] for file_path in args.file_path: file_paths += glob(file_path) diff --git a/parsedmarc/kafkaclient.py b/parsedmarc/kafkaclient.py new file mode 100644 index 00000000..7eaeb473 --- /dev/null +++ b/parsedmarc/kafkaclient.py @@ -0,0 +1,51 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- + +from kafka import KafkaProducer +import json + +class KafkaError(RuntimeError): + """Raised when a Kafka error occurs""" + +class KafkaClient(object): + def __init__(self, kafka_hosts): + """ Right now lets just do one host""" + self.host = kafka_hosts + self.producer = KafkaProducer(bootstrap_servers=kafka_hosts, value_serializer=lambda v: json.dumps(v).encode('utf-8')) + + def save_aggregate_reports_to_kafka(self, aggregate_reports, aggregate_topic): + """ + Saves aggregate DMARC reports to Splunk + + Args: + aggregate_reports (list): A list of aggregate report dictionaries + to save to kafka + + """ + if type(aggregate_reports) == dict: + aggregate_reports = [aggregate_reports] + + if len(aggregate_reports) < 1: + return + + for report in aggregate_reports: + self.producer.send(aggregate_topic, report) + + + def save_forensic_reports_to_kafka(self, forensic_reports, forensic_topic): + """ + Saves forensic DMARC reports to Kafka + + Args: + forensic_reports (list): A list of forensic report dictionaries + to save to kafka + + """ + if type(forensic_reports) == dict: + forensic_reports = [forensic_reports] + + if len(forensic_reports) < 1: + return + + for report in forensic_reports: + self.producer.send(forensic_topic, report) diff --git a/requirements.txt b/requirements.txt index b10e4576..023ed17d 100644 --- a/requirements.txt +++ b/requirements.txt @@ -15,3 +15,4 @@ sphinx_rtd_theme collective.checkdocs wheel rstcheck +kafka-python From d4cf4a7e5f0cc0965aaf4747fd7bd6f6f244c529 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Tue, 9 Oct 2018 14:08:02 -0400 Subject: [PATCH 2/9] forgot to flush --- parsedmarc/kafkaclient.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/parsedmarc/kafkaclient.py b/parsedmarc/kafkaclient.py index 7eaeb473..f81aff94 100644 --- a/parsedmarc/kafkaclient.py +++ b/parsedmarc/kafkaclient.py @@ -30,6 +30,7 @@ class KafkaClient(object): for report in aggregate_reports: self.producer.send(aggregate_topic, report) + self.producer.flush() def save_forensic_reports_to_kafka(self, forensic_reports, forensic_topic): @@ -49,3 +50,4 @@ class KafkaClient(object): for report in forensic_reports: self.producer.send(forensic_topic, report) + self.producer.flush() From a3ba85803a82a197508e4717845aa693271b25c4 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 08:07:44 -0400 Subject: [PATCH 3/9] Modified to send entire ordered dict to Kafka. Bug: would barf on reports larger than 10 megs --- parsedmarc/kafkaclient.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/parsedmarc/kafkaclient.py b/parsedmarc/kafkaclient.py index f81aff94..55f6dfad 100644 --- a/parsedmarc/kafkaclient.py +++ b/parsedmarc/kafkaclient.py @@ -11,7 +11,8 @@ class KafkaClient(object): def __init__(self, kafka_hosts): """ Right now lets just do one host""" self.host = kafka_hosts - self.producer = KafkaProducer(bootstrap_servers=kafka_hosts, value_serializer=lambda v: json.dumps(v).encode('utf-8')) + self.producer = KafkaProducer(value_serializer=lambda v: json.dumps(v).encode('utf-8'), + bootstrap_servers=kafka_hosts) def save_aggregate_reports_to_kafka(self, aggregate_reports, aggregate_topic): """ @@ -28,9 +29,10 @@ class KafkaClient(object): if len(aggregate_reports) < 1: return - for report in aggregate_reports: - self.producer.send(aggregate_topic, report) - self.producer.flush() + print("Report is {}".format(aggregate_reports)) + print("Report type is {}".format(type(aggregate_reports))) + self.producer.send(aggregate_topic, aggregate_reports) + self.producer.flush() def save_forensic_reports_to_kafka(self, forensic_reports, forensic_topic): @@ -49,5 +51,5 @@ class KafkaClient(object): return for report in forensic_reports: - self.producer.send(forensic_topic, report) + self.producer.send(forensic_topic, json.dumps(report)) self.producer.flush() From 687a44ee5803fd0e9717e529bd7ae9975291cad1 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 09:11:24 -0400 Subject: [PATCH 4/9] split out individual records. --- parsedmarc/cli.py | 3 +++ parsedmarc/kafkaclient.py | 15 ++++++++------- 2 files changed, 11 insertions(+), 7 deletions(-) diff --git a/parsedmarc/cli.py b/parsedmarc/cli.py index 65a88e7a..c155e752 100644 --- a/parsedmarc/cli.py +++ b/parsedmarc/cli.py @@ -16,6 +16,9 @@ from parsedmarc import logger, IMAPError, get_dmarc_reports_from_inbox, \ parse_report_file, elastic, kafkaclient, splunk, save_output, watch_inbox, \ email_results, SMTPError, ParserError, __version__ +import sys +import os + def _main(): """Called when the module is executed""" def process_reports(reports_): diff --git a/parsedmarc/kafkaclient.py b/parsedmarc/kafkaclient.py index 55f6dfad..b1e27730 100644 --- a/parsedmarc/kafkaclient.py +++ b/parsedmarc/kafkaclient.py @@ -2,6 +2,7 @@ # -*- coding: utf-8 -*- from kafka import KafkaProducer +from collections import OrderedDict import json class KafkaError(RuntimeError): @@ -9,14 +10,12 @@ class KafkaError(RuntimeError): class KafkaClient(object): def __init__(self, kafka_hosts): - """ Right now lets just do one host""" - self.host = kafka_hosts - self.producer = KafkaProducer(value_serializer=lambda v: json.dumps(v).encode('utf-8'), + self.producer = KafkaProducer(value_serializer=lambda v: json.dumps(v).encode('utf-8'), bootstrap_servers=kafka_hosts) def save_aggregate_reports_to_kafka(self, aggregate_reports, aggregate_topic): """ - Saves aggregate DMARC reports to Splunk + Saves aggregate DMARC reports to Kafka Args: aggregate_reports (list): A list of aggregate report dictionaries @@ -29,9 +28,11 @@ class KafkaClient(object): if len(aggregate_reports) < 1: return - print("Report is {}".format(aggregate_reports)) - print("Report type is {}".format(type(aggregate_reports))) - self.producer.send(aggregate_topic, aggregate_reports) + for record in aggregate_reports['records']: + buffer = OrderedDict([('xml_schema', aggregate_reports['xml_schema']), + ('report_metadata',aggregate_reports['report_metadata']), + ('records', record)]) + self.producer.send(aggregate_topic, buffer) self.producer.flush() From 19df7f65c42426fa6b34537a2bd2459ceafedb26 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 09:54:03 -0400 Subject: [PATCH 5/9] PEP8 fixes --- parsedmarc/cli.py | 49 ++++++++++++++++++++++++----------------------- 1 file changed, 25 insertions(+), 24 deletions(-) diff --git a/parsedmarc/cli.py b/parsedmarc/cli.py index c155e752..038cf6bd 100644 --- a/parsedmarc/cli.py +++ b/parsedmarc/cli.py @@ -13,14 +13,13 @@ import json from elasticsearch.exceptions import ElasticsearchException from parsedmarc import logger, IMAPError, get_dmarc_reports_from_inbox, \ - parse_report_file, elastic, kafkaclient, splunk, save_output, watch_inbox, \ - email_results, SMTPError, ParserError, __version__ + parse_report_file, elastic, kafkaclient, splunk, save_output, \ + watch_inbox, email_results, SMTPError, ParserError, __version__ -import sys -import os def _main(): """Called when the module is executed""" + def process_reports(reports_): output_str = "{0}\n".format(json.dumps(reports_, ensure_ascii=False, @@ -29,10 +28,9 @@ def _main(): print(output_str) if args.kafka_hosts: try: - kafkaClient = kafkaclient.KafkaClient(args.kafka_hosts) - # dont do this + kafkaClient = kafkaclient.KafkaClient(args.kafka_hosts) except Exception as error: - logger.error("Kafka Error: {0}".format(error.__str__())) + logger.error("Kafka Error: {0}".format(error.__str__())) if args.save_aggregate: for report in reports_["aggregate_reports"]: try: @@ -46,12 +44,11 @@ def _main(): error_.__str__())) exit(1) try: - if args.kafka_hosts: - kafkaClient.save_aggregate_reports_to_kafka( - report, kafka_aggregate_topic) - # dont do this catch specific exceptions + if args.kafka_hosts: + kafkaClient.save_aggregate_reports_to_kafka( + report, kafka_aggregate_topic) except Exception as error_: - logger.error("Kafka Error: {0}".format( + logger.error("Kafka Error: {0}".format( error_.__str__())) if args.hec: try: @@ -72,13 +69,13 @@ def _main(): logger.error("Elasticsearch Error: {0}".format( error_.__str__())) try: - if args.kafka_hosts: - kafkaClient.save_forensic_reports_to_kafka( - report, kafka_forensic_topics) - # dont do this + if args.kafka_hosts: + kafkaClient.save_forensic_reports_to_kafka( + report, kafka_forensic_topic) + except Exception as error_: - logger.error("Kafka Error: {0}".format( - error_.__str__())) + logger.error("Kafka Error: {0}".format( + error_.__str__())) if args.hec: try: forensic_reports_ = reports_["forensic_reports"] @@ -147,11 +144,14 @@ def _main(): help="Skip certificate verification for Splunk " "HEC") arg_parser.add_argument("-K", "--kafka-hosts", nargs="*", - help="A list of one or more Kafka hostnames or URLs") + help="A list of one or more Kafka hostnames" + " or URLs") arg_parser.add_argument("--kafka-aggregate-topic", - help="The Kafka topic to publish aggregate reports to.") + help="The Kafka topic to publish aggregate " + "reports to.") arg_parser.add_argument("--kafka-forensic_topic", - help="The Kafka topic to publish forensic reports to.") + help="The Kafka topic to publish forensic reports" + " to.") arg_parser.add_argument("--save-aggregate", action="store_true", default=False, help="Save aggregate reports to search indexes") @@ -221,7 +221,8 @@ def _main(): es_forensic_index = "{0}_{1}".format(es_forensic_index, suffix) if args.save_aggregate or args.save_forensic: - if args.elasticsearch_host is None and args.hec and args.kafka_hosts is None: + if (args.elasticsearch_host is None and args.hec + and args.kafka_hosts is None): args.elasticsearch_host = ["localhost:9200"] try: if args.elasticsearch_host: @@ -248,10 +249,10 @@ def _main(): kafka_forensic_topic = "dmarc_forensic" if args.kafka_aggregate_topic: - kafka_aggregate_topic = args.kafka_aggregate_topic + kafka_aggregate_topic = args.kafka_aggregate_topic if args.kafka_forensic_topic: - kafka_forensic_topic = args.kafka_forensic_topic + kafka_forensic_topic = args.kafka_forensic_topic file_paths = [] for file_path in args.file_path: From 966495a2a97921f36ca845a108a1c185c2375730 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 10:04:30 -0400 Subject: [PATCH 6/9] PEP8 changes --- parsedmarc/kafkaclient.py | 42 +++++++++++++++++++++++---------------- 1 file changed, 25 insertions(+), 17 deletions(-) diff --git a/parsedmarc/kafkaclient.py b/parsedmarc/kafkaclient.py index b1e27730..7f48c649 100644 --- a/parsedmarc/kafkaclient.py +++ b/parsedmarc/kafkaclient.py @@ -2,18 +2,26 @@ # -*- coding: utf-8 -*- from kafka import KafkaProducer -from collections import OrderedDict +from kafka.errors import NoBrokersAvailable, UnknownTopicOrPartitionError import json + class KafkaError(RuntimeError): """Raised when a Kafka error occurs""" -class KafkaClient(object): - def __init__(self, kafka_hosts): - self.producer = KafkaProducer(value_serializer=lambda v: json.dumps(v).encode('utf-8'), - bootstrap_servers=kafka_hosts) - def save_aggregate_reports_to_kafka(self, aggregate_reports, aggregate_topic): +class KafkaClient(object): + def __init__(self, kafka_hosts): + try: + def serializer(v): lambda v: json.dumps(v).encode('utf-8') + self.producer = KafkaProducer( + value_serializer=serializer, + bootstrap_servers=kafka_hosts) + except NoBrokersAvailable: + raise KafkaError("No Kafka brokers availabe") + + def save_aggregate_reports_to_kafka(self, aggregate_reports, + aggregate_topic): """ Saves aggregate DMARC reports to Kafka @@ -28,20 +36,18 @@ class KafkaClient(object): if len(aggregate_reports) < 1: return - for record in aggregate_reports['records']: - buffer = OrderedDict([('xml_schema', aggregate_reports['xml_schema']), - ('report_metadata',aggregate_reports['report_metadata']), - ('records', record)]) - self.producer.send(aggregate_topic, buffer) + try: + self.producer.send(aggregate_topic, aggregate_reports) + except UnknownTopicOrPartitionError: + raise KafkaError("Unknown topic or partition on broker") self.producer.flush() - - def save_forensic_reports_to_kafka(self, forensic_reports, forensic_topic): + def save_forensic_reports_to_kafka(self, forensic_reports, forensic_topic): """ Saves forensic DMARC reports to Kafka Args: - forensic_reports (list): A list of forensic report dictionaries + forensic_reports (list): A list of forensic report dicts to save to kafka """ @@ -51,6 +57,8 @@ class KafkaClient(object): if len(forensic_reports) < 1: return - for report in forensic_reports: - self.producer.send(forensic_topic, json.dumps(report)) - self.producer.flush() + try: + self.producer.send(forensic_topic, forensic_reports) + except UnknownTopicOrPartitionError: + raise KafkaError("Unknown topic or partition on broker") + self.producer.flush() From 66e707bfdff443477f7a88980737028a34899b83 Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 10:12:34 -0400 Subject: [PATCH 7/9] bumping version --- parsedmarc/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parsedmarc/__init__.py b/parsedmarc/__init__.py index cc075e3d..6f4162df 100644 --- a/parsedmarc/__init__.py +++ b/parsedmarc/__init__.py @@ -44,7 +44,7 @@ import imapclient.exceptions import dateparser import mailparser -__version__ = "4.2.0" +__version__ = "4.2.0k" logger = logging.getLogger(__name__) logger.setLevel(logging.ERROR) From fe611ac9df22fe9ce3409729df40a717534a4eab Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 11:57:41 -0400 Subject: [PATCH 8/9] added k version to setup.py --- setup.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/setup.py b/setup.py index fe803509..a4a1e69f 100644 --- a/setup.py +++ b/setup.py @@ -14,7 +14,7 @@ from setuptools import setup from codecs import open from os import path -__version__ = "4.2.0" +__version__ = "4.2.0k" description = "A Python package and CLI for parsing aggregate and " \ "forensic DMARC reports" From 074ce9b8155a4aa6fd5c43fba2adfd53ee24b40f Mon Sep 17 00:00:00 2001 From: Mike Siegel Date: Wed, 10 Oct 2018 13:20:28 -0400 Subject: [PATCH 9/9] Removed logger from import --- parsedmarc/cli.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/parsedmarc/cli.py b/parsedmarc/cli.py index 852b6ebc..148e572f 100644 --- a/parsedmarc/cli.py +++ b/parsedmarc/cli.py @@ -12,7 +12,7 @@ import json from elasticsearch.exceptions import ElasticsearchException -from parsedmarc import logger, IMAPError, get_dmarc_reports_from_inbox, \ +from parsedmarc import IMAPError, get_dmarc_reports_from_inbox, \ parse_report_file, elastic, kafkaclient, splunk, save_output, \ watch_inbox, email_results, SMTPError, ParserError, __version__