from typing import Dict import psycopg2 from psycopg2.extras import RealDictCursor import boto3 import hashlib import logging import os import json log = logging.getLogger() REGION = os.environ["AWS_REGION"] if "AWS_REGION" in os.environ else "eu-west-1" S3_BUCKET_DEVEL = "frontend-api-devel-filestore" def get_secret(secret_name): # Create a Secrets Manager client session = boto3.session.Session() client = session.client( service_name='secretsmanager', region_name=REGION ) get_secret_value_response = json.loads(client.get_secret_value(SecretId=secret_name)['SecretString']) return get_secret_value_response DB_CONN = {'fansifter-rds': None, 'fansifter-rds-test': None, 'fansifter-rds-live': None, 'sandbox': None} DB_PORT = {'fansifter-rds': 5432, 'fansifter-rds-test': 5433, 'fansifter-rds-live': 5434, 'sandbox': None} def get_rds_connection(conn_name): global DB_CONN if DB_CONN[conn_name] is None: try: secret = get_secret(conn_name) DB_CONN[conn_name] = psycopg2.connect(host=secret['host'], # 'localhost', port=DB_PORT[conn_name], # secret['port'], #DB_PORT[conn_name], database=secret['dbname'], user=secret['username'], password=secret['password'], cursor_factory=RealDictCursor) except Exception as e: log.exception(e) return DB_CONN[conn_name] def get_rds_connection_sandbox(): global DB_CONN if DB_CONN['sandbox'] is None: try: secret = get_secret('sandbox_pass') DB_CONN['sandbox'] = psycopg2.connect(host='localhost', # secret['host'], port='5431', # secret['port'], database=secret['dbname'], user=secret['username'], password=secret['password'], cursor_factory=RealDictCursor) except Exception as e: log.exception(e) return DB_CONN['sandbox'] def rds_query(connection_name, query, vars=None, fetch=False): with get_rds_connection(connection_name) as conn: with conn.cursor() as cur: cur.execute(query, vars=vars) if fetch: return cur.fetchall() databases = [ 'fansifter-rds', 'fansifter-rds-test', 'fansifter-rds-live', ] def make_company_hash(event: Dict) -> str: try: if 'userCompany' in event: company = event['userCompany'] else: # for cognito_post_confirmation company = event['request']['userAttributes']['custom:company'] return hashlib.sha224(company.encode()).hexdigest() except KeyError as e: log.exception(e) raise RuntimeError("Please provide payload.userCompany field in appsync RequestMappingTemplate") def make_schema_id(event: Dict) -> str: return 'c' + make_company_hash(event) def schema_id(company): event = {'userCompany': company} return make_schema_id(event) ddb = boto3.resource('dynamodb') #table = ddb.Table('frontend-api-devel-userdata') table = ddb.Table('frontend-api-test-userdata') for user in rds_query('fansifter-rds-test', "SELECT * FROM commons.user_company", fetch=True): res = table.put_item(Item={'UID': user['user_id'], 'IID': user['default_schema'], 'institutionName': user['company_name'], 'email': user['email'], 'firstName': user['user_name'], 'lastName': '', 'createdDateTime': '2021-02-17T14:49:44.754408Z'}) print(res)