#!/usr/bin/env python3 import os import sys import csv import json import getopt from itertools import islice from datetime import datetime, timedelta import structlog from urllib.parse import urlencode from datadog_api_client import ApiClient, Configuration from datadog_api_client.v2.api.spans_api import SpansApi from datadog_api_client.v2.model.spans_aggregate_data import SpansAggregateData from datadog_api_client.v2.model.spans_aggregate_request import SpansAggregateRequest from datadog_api_client.v2.model.spans_aggregate_request_attributes import SpansAggregateRequestAttributes from datadog_api_client.v2.model.spans_aggregate_request_type import SpansAggregateRequestType from datadog_api_client.v2.model.spans_query_filter import SpansQueryFilter from datadog_api_client.v2.model.spans_group_by import SpansGroupBy from datadog_api_client.v2.model.spans_aggregate_sort import SpansAggregateSort from datadog_api_client.v2.model.spans_sort_order import SpansSortOrder from datadog_api_client.v2.model.spans_aggregate_sort_type import SpansAggregateSortType from datadog_api_client.v2.model.spans_aggregation_function import SpansAggregationFunction from datadog_api_client.v2.model.spans_list_request import SpansListRequest from datadog_api_client.v2.model.spans_list_request_attributes import SpansListRequestAttributes from datadog_api_client.v2.model.spans_list_request_data import SpansListRequestData from datadog_api_client.v2.model.spans_list_request_page import SpansListRequestPage from datadog_api_client.v2.model.spans_list_request_type import SpansListRequestType from datadog_api_client.v2.model.spans_sort import SpansSort logger = structlog.get_logger() DD_ENV_QA = "qa" DD_ENV_PROD = "prod" DD_SPAN_COUNT = 4000 USER_AGENT = "core-load-tester/1.0" ORCHARD_HEADERS_MAPPING = { "account_id": "Grass-Account-Id", "account_type": "Grass-Account-Type", "identity_id": "Orchard-Identity-Id", "identity_uuid": "Orchard-Identity-UUID", "profile_id": "Orchard-Profile-Id", "profile_type": "Orchard-Profile-Type", "user_id": "Orchard-User-Id", } class Configurator(Configuration): """ Responsible for collecting all the necessary configuration from the environment. """ # defaults no_header_row = False outfile_path = None result_limit = 1000 def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self._process_args(sys.argv) self._collect_env_vars(os.environ) self.api_key["apiKeyAuth"] = os.environ.get("DD_API_KEY") self.api_key["appKeyAuth"] = os.environ.get("DD_APP_KEY") self.service_name = os.environ.get("DD_SERVICE_NAME") # use for getting aggregate prod data for the formation of the weight urls self.start_period_prod = self.query_date # use for getting spans data from QA and getting first DD_SPAN_COUNT spans self.start_period_qa = self.query_date + timedelta(weeks=-1) self.end_period = self.query_date + timedelta(days=1, minutes=-1) logger.bind( service_name=self.service_name, query_date=self.query_date, start_period_prod=self.start_period_prod, start_period_qa=self.start_period_qa, end_period=self.end_period, outfile_path=self.outfile_path, span_count=DD_SPAN_COUNT, result_limit=self.result_limit, headers_mapping=ORCHARD_HEADERS_MAPPING, ).info( "Configuration..." ) def _process_args(self, argv): name = sys.argv[0] args = sys.argv[1:] usage_msg = " ".join([name, "-o -d "]) short_opts = "hHl:o:d:s:" long_opts = ["help", "no-header-row", "limit=", "outfile=", "date=", "shift="] try: opts, formal_args = getopt.gnu_getopt(args, short_opts, long_opts) except getopt.GetoptError: logger.error(usage_msg) sys.exit(2) for opt, arg in opts: if opt in ("-h", "--help"): sys.exit() elif opt in ("-H", "--no-header-row"): self.no_header_row = True elif opt in ("-l", "--limit"): self.result_limit = int(arg) elif opt in ("-o", "--outfile"): self.outfile_path = arg # query_date elif opt in ("-d", "--date"): date = arg self.query_date = datetime.strptime(date, '%Y-%m-%d').astimezone() if not self.outfile_path or not self.query_date: logger.error("Requires an argument of the date (-d) and outfile (-o).") sys.exit(2) def _collect_env_vars(self, environ): pass class OutputFormatter: """Output the Athena query results as a CSV.""" def __init__(self, config, results): self.config = config self.no_header_row = config.no_header_row self.headers = self._headers() self.rows = self._clean_rows(results) def output_csv(self, outfile_path): """ Write CSV to the given outfile path. If no path is given, defaults to STDOUT. """ logger.bind( service_name=self.config.service_name, query_date=self.config.query_date, headers=self.headers, rows_count=len(self.rows), ).info( "Output CSV" ) if outfile_path is None: file_handle = sys.stdout else: file_handle = open(outfile_path, "w") csv_writer = csv.writer(file_handle, quoting=csv.QUOTE_ALL) csv_writer.writerow(self.headers) for row in self.rows: csv_writer.writerow(row) @staticmethod def _headers(): return ["method", "pattern", "uri", "headers", "query_string", "count"] @staticmethod def _clean_rows(results): return [ ( item["method"], item["pattern"], item["url"], item["headers"], item["query_string"], item["count"], ) for item in results ] class DataDogSpanHandler: def __init__(self, config): self.config = config def get_resource_data(self): logger.bind( service_name=self.config.service_name, query_date=self.config.query_date, ).info( "Getting resource data" ) with ApiClient(self.config) as api_client: api_instance = SpansApi(api_client) resp_aggregate_spans_prod = api_instance.aggregate_spans( body=self._get_group_by_resource_name_query(env=DD_ENV_PROD), ) aggregate_spans_data_prod = self._get_aggregate_spans_data(resp_aggregate_spans_prod) resp_spans_qa = api_instance.list_spans_with_pagination( body=self._get_spans_with_pagination_query(env=DD_ENV_QA), ) spans_data_qa = self._get_spans_data(resp_spans_qa) group_spans_by_url_rule = self._group_spans_by_url_rule(spans_data_qa) resource_data = self._add_weight_urls(group_spans_by_url_rule, aggregate_spans_data_prod) return resource_data def _add_weight_urls(self, group_spans_by_url_rule, aggregate_spans_data_prod): data = [] if not (group_spans_by_url_rule or aggregate_spans_data_prod): return data frequently_url_rule = max(aggregate_spans_data_prod, key=aggregate_spans_data_prod.get) weight_coefficient_for_new_urls = ( aggregate_spans_data_prod[frequently_url_rule] // group_spans_by_url_rule[frequently_url_rule]["raw_count"] ) for url_rule, value in group_spans_by_url_rule.items(): prod_count = aggregate_spans_data_prod.get(url_rule) if prod_count: count = prod_count // len(value["urls"]) or 1 else: logger.bind( service_name=self.config.service_name, query_date=self.config.query_date, url_rule=url_rule, weight_coefficient_for_new_urls=weight_coefficient_for_new_urls, ).info( "URL rule doesn't exist in prod env" ) count = group_spans_by_url_rule.get(url_rule)["raw_count"] * weight_coefficient_for_new_urls for url, method, headers, query_string in value["urls"]: data.append( { "method": method, "pattern": url_rule, "url": url, "headers": headers, "query_string": query_string, "count": count, } ) return data def _get_spans_with_pagination_query(self, env=DD_ENV_QA): query = f'env:{env} service:{self.config.service_name} ' \ f'resource_name:(*GET*) -@http.status_code:(404 OR 403 OR 408) ' \ f'-@http.useragent:"{USER_AGENT}" type:web' body = SpansListRequest( data=SpansListRequestData( attributes=SpansListRequestAttributes( filter=SpansQueryFilter( _from=self.config.start_period_qa.isoformat(timespec="milliseconds"), query=query, to=self.config.end_period.isoformat(timespec="milliseconds"), ), page=SpansListRequestPage( limit=1000, ), sort=SpansSort.TIMESTAMP_DESCENDING, ), type=SpansListRequestType.SEARCH_REQUEST, ), ) return body def _get_group_by_resource_name_query(self, env=DD_ENV_QA): query_string = f"env:{env} service:{self.config.service_name} " \ "resource_name:(*GET*) -@http.status_code:(404 OR 403 OR 408) " \ f"-@http.useragent:'{USER_AGENT}' type:web" body = SpansAggregateRequest( data=SpansAggregateData( attributes=SpansAggregateRequestAttributes( filter=SpansQueryFilter( _from=self.config.start_period_prod.isoformat(timespec="milliseconds"), query=query_string, to=self.config.end_period.isoformat(timespec="milliseconds"), ), group_by=[ SpansGroupBy( facet="resource_name", limit=100, sort=SpansAggregateSort( type=SpansAggregateSortType.MEASURE, order=SpansSortOrder.DESCENDING, aggregation=SpansAggregationFunction.COUNT, metric="@duration", ) ), ], ), type=SpansAggregateRequestType.AGGREGATE_REQUEST, ), ) return body def _get_aggregate_spans_data(self, resp_data): data = {} for item in resp_data.data: resource_name = item.get("attributes").get("by").get("resource_name") url_rule = self._get_url_rule_by_resource_name(resource_name) count = int(item.get("attributes").get("compute").get("c0") or 0) data[url_rule] = count logger.bind( service_name=self.config.service_name, query_date=self.config.query_date, aggregated_spans_data_prod=data, ).info( "Aggregated spans data prod" ) return data @staticmethod def _get_url_rule_by_resource_name(resource_name): return resource_name.split(" ")[1] if resource_name else None def _get_spans_data(self, resp_data): data = [] for item in islice(resp_data, DD_SPAN_COUNT): http_data = item.get("attributes").get("custom").get("http") query_string = http_data.get("url_details").get("queryString") resource_name = item.get("attributes").get("resource_name") headers = {} request_context = item.get("attributes").get("custom").get("request_context", {}) for request_context_key, request_context_value in request_context.items(): header = ORCHARD_HEADERS_MAPPING.get(request_context_key) if header: headers[header] = ( str(request_context_value) if isinstance(request_context_value, int) else request_context_value ) request_headers = http_data.get("request", {}).get("headers", {}) for request_header, request_header_value in request_headers.items(): if request_header == "user-agent": continue headers[request_header] = ( str(request_header_value) if isinstance(request_header_value, int) else request_header_value ) data.append( { "method": http_data.get("method"), "url_rule": self._get_url_rule_by_resource_name(resource_name), "url": http_data.get("url_details").get("path"), "query_string": urlencode(query_string) if query_string else None, "headers": json.dumps(headers) if headers else None, } ) return data def _group_spans_by_url_rule(self, spans): url_rules = set([item["url_rule"] for item in spans]) data = {item: {"urls": []} for item in url_rules} for item in spans: data[item["url_rule"]]["urls"].append((item["url"], item["method"], item["headers"], item["query_string"])) for url_rule, value in data.items(): data[url_rule]["raw_count"] = len(data[url_rule]["urls"]) data[url_rule]["urls"] = set(data[url_rule]["urls"]) logger.bind( service_name=self.config.service_name, query_date=self.config.query_date, aggregated_spans_data_qa={k: v["raw_count"] for k, v in data.items()}, ).info( "Aggregated spans data qa" ) return data def main(): configuration = Configurator() span_handler = DataDogSpanHandler(configuration) resource_data = span_handler.get_resource_data() formatter = OutputFormatter(configuration, resource_data) formatter.output_csv(configuration.outfile_path) if __name__ == "__main__": main()