import os import json import boto3 import psycopg2 from psycopg2.extras import RealDictCursor import logging log = logging.getLogger() log.setLevel(logging.INFO) 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], 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 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', ] other_buckets = set() for connection in databases: all_chemas = rds_query(connection, "select table_schema, table_name from information_schema.tables where table_name IN ('company' )", fetch=True) print(f"Found {len(all_chemas)} schemas in database {connection}") for schema, table in [sch.values() for sch in all_chemas]: print(f'Working on: {schema}.{table}') statements = [] vars = {} try: for wrk in rds_query(connection, f"SELECT id FROM {schema}.workspace", fetch=True): wrk_schema = wrk['id'] if wrk_schema != schema: management_schema = schema print(f" Workspace: {wrk_schema}") # get profiles from workspace: wrk_profiles = [p['profile_id'] for p in rds_query(connection, f"SELECT profile_id FROM {wrk_schema}.fan", fetch=True)] # insert workspace profiles into management_schema batch = 0 batch_size = 100 while len(wrk_profiles) > batch * batch_size: batch_start = batch * batch_size wrk_values = ','.join([f"('{wp}')" for wp in wrk_profiles[batch*batch_size:batch_start+batch_size]]) sql = f"INSERT INTO {management_schema}.fan (profile_id) VALUES {wrk_values} " \ f"ON CONFLICT (profile_id) DO NOTHING" statements.append(sql) batch += 1 # remap id's to main schema id's statements.append( f"UPDATE {wrk_schema}.fan_attribute cf SET fan_id = (SELECT mf.id FROM {wrk_schema}.fan f " f"JOIN {management_schema}.fan mf ON f.profile_id = mf.profile_id WHERE f.id = cf.fan_id)") statements.append( f"DELETE FROM {wrk_schema}.collection_fan") statements.append( f"INSERT INTO {wrk_schema}.collection_fan (collection_id, fan_id) " f"SELECT ca.collection_id, ca.fan_id FROM {wrk_schema}.fan_attribute ca " f"ON CONFLICT DO NOTHING") except Exception as e: print(f' Oops, {e} // in: {schema}.{table}') for statement in statements: try: rds_query(connection, statement, vars) except Exception as e: log.exception("oops", exc_info=e)