"""Application configuration.""" from distutils.util import strtobool import logging import os from cryptography.hazmat.backends import default_backend from cryptography.hazmat.primitives import serialization from dotenv import load_dotenv from secrets_manager.lambda_ext import LambdaSecretsManager from snowflake_connector.etl_connector import SnowflakeSQLExecutor from snowflake_connector.private_key import get_private_key # Load env file if it exists load_dotenv(verbose=True) # Service information SERVICE_NAME = 'neo4j-cypher-scheduler' SERVICE_VERSION = '1.1.1' # Logger Name LOGGER_NAME = 'neo4j-cypher-scheduler' LOGGER_LEVEL = logging.INFO LOGGER_DSN = os.environ.get('LOGGER_DSN') or None ENVIRONMENT = os.environ.get('Environment', 'qa') print('The environment is {}'.format(ENVIRONMENT)) ENVIRONMENTS = ['local', 'dev', 'qa', 'prod'] if ENVIRONMENT not in ENVIRONMENTS: raise Exception('The environment used is not properly named.') # Neo4j connection settings NEO4J_URL = os.environ.get( 'NEO4J_URL', 'neo4j+s://{}-neo4j-cluster.theorchard.io:7687'.format(ENVIRONMENT)) NEO4J_DATABASE_NAME = os.environ.get('NEO4J_DATABASE_NAME', 'graph.db') # The name of the file containing the query to execute NEO4J_SOURCE_NAME = os.environ.get('NEO4J_SOURCE_NAME', 'test') # The name of the folder containing the queries to execute # If set, NEO4J_SOURCE_NAME will be ignored NEO4J_SOURCE_FOLDER_NAME = os.environ.get('NEO4J_SOURCE_FOLDER_NAME') CYPHER_SOURCE = ( NEO4J_SOURCE_FOLDER_NAME if NEO4J_SOURCE_FOLDER_NAME else NEO4J_SOURCE_NAME ) # When True, runs Snowflake sync assertions instead of executing .cypher files ASSERT_SNOWFLAKE_SYNC = bool(strtobool(os.environ.get('ASSERT_SNOWFLAKE_SYNC', 'False'))) if ENVIRONMENT == 'local': NEO4J_CONNECTION_USER = os.environ.get('NEO4J_CONNECTION_USER') NEO4J_CONNECTION_PASSWORD = os.environ.get('NEO4J_CONNECTION_PASSWORD') SENTRY_DSN = os.environ.get('SENTRY_DSN', '') DATADOG_API_KEY = os.environ.get('DATADOG_API_KEY', '') DATADOG_APP_KEY = os.environ.get('DATADOG_APP_KEY', '') secrets_manager_service = None else: secrets_manager_service = LambdaSecretsManager( environment=ENVIRONMENT, service_name=SERVICE_NAME) secrets_manager_datadog = LambdaSecretsManager(environment=ENVIRONMENT, service_name='datadog') NEO4J_CONNECTION_USER = secrets_manager_service.get_cred('NEO4J_CONNECTION_USER') NEO4J_CONNECTION_PASSWORD = secrets_manager_service.get_cred('NEO4J_CONNECTION_PASSWORD') SENTRY_DSN = secrets_manager_service.get_cred('SENTRY_DSN') DATADOG_API_KEY = secrets_manager_datadog.get_cred('DD_API_KEY') DATADOG_APP_KEY = secrets_manager_datadog.get_cred('DD_APP_KEY') DATADOG_CYPHER_INFO_METRIC = 'info' DATADOG_CYPHER_FAILURE_METRIC = 'failure' DATADOG_CYPHER_SUCCESS_METRIC = 'success' SEND_QUERY_RESULTS_TO_DATADOG = os.environ.get( 'SEND_QUERY_RESULTS_TO_DATADOG', '').split(',') # Flag indicating if the results of the queries should be asserted or not ASSERT_RESULTS = bool(strtobool(os.environ.get('ASSERT_RESULTS', 'False'))) # Flag indicating if the queries should continue # to be executed after one of them failed CONTINUE_AFTER_ERRORS = bool(strtobool( os.environ.get('CONTINUE_AFTER_ERRORS', 'False'))) # Flag indicating if the queries should be run in multiple batches RUN_IN_BATCHES = bool(strtobool( os.environ.get('RUN_IN_BATCHES', 'False'))) BATCH_SIZE = int(os.environ.get('BATCH_SIZE', 100000)) DISABLE_SENTRY = bool(strtobool(os.environ.get('DISABLE_SENTRY', 'False'))) SNOWFLAKE_DATABASE = os.environ.get('SNOWFLAKE_DATABASE') SNOWFLAKE_SCHEMA = os.environ.get('SNOWFLAKE_SCHEMA') WINDOW_START_HOURS = int(os.environ.get('WINDOW_START_HOURS', 25)) WINDOW_END_HOURS = int(os.environ.get('WINDOW_END_HOURS', 1)) DEFAULT_SAMPLE_NUMBER = int(os.environ.get('DEFAULT_SAMPLE_NUMBER', 100)) def build_snowflake_config(): """Return sf_config dict for SnowflakeSQLExecutor.""" if ENVIRONMENT == 'local': private_key_path = os.path.expanduser(os.environ.get( 'SNOWFLAKE_PRIVATE_KEY_PATH', '~/.ssh/snowflake/rsa_key.p8', )) private_key_passphrase = os.environ.get('SNOWFLAKE_PRIVATE_KEY_PASSPHRASE', '') return { 'account': os.environ.get('SNOWFLAKE_ACCOUNT'), 'user': os.environ.get('SNOWFLAKE_USER'), 'db': os.environ.get('SNOWFLAKE_DATABASE'), 'schema': os.environ.get('SNOWFLAKE_SCHEMA'), 'warehouse': os.environ.get('SNOWFLAKE_WAREHOUSE'), 'role': os.environ.get('SNOWFLAKE_ROLE'), 'private_key': get_private_key(private_key_path, private_key_passphrase), } pem_key = os.environ.get('SNOWFLAKE_PRIVATE_KEY') passphrase = os.environ.get('SNOWFLAKE_KEY_PASSPHRASE') sf_private_key = serialization.load_pem_private_key( bytes(pem_key, 'utf-8'), password=bytes(passphrase, 'utf-8'), backend=default_backend(), ) private_key_der = sf_private_key.private_bytes( encoding=serialization.Encoding.DER, format=serialization.PrivateFormat.PKCS8, encryption_algorithm=serialization.NoEncryption(), ) return { 'account': os.environ.get('SNOWFLAKE_ACCOUNT'), 'user': os.environ.get('SNOWFLAKE_USER'), 'db': os.environ.get('SNOWFLAKE_DATABASE'), 'schema': os.environ.get('SNOWFLAKE_SCHEMA'), 'warehouse': os.environ.get('SNOWFLAKE_WAREHOUSE'), 'role': os.environ.get('SNOWFLAKE_ROLE'), 'private_key': private_key_der, } def create_snowflake_executor(): """Return a SnowflakeSQLExecutor context manager.""" return SnowflakeSQLExecutor(build_snowflake_config())