import os import json import boto3 import psycopg2 from psycopg2._psycopg import ProgrammingError from psycopg2.errors import UndefinedTable 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], #secret['port'], , 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: if isinstance(query, list): res = [] for statement in query: cur.execute(statement, vars=vars) try: res.append(cur.fetchall()) except ProgrammingError as e: res.append([]) return res else: cur.execute(query, vars=vars) try: return cur.fetchall() except ProgrammingError as e: return [] MANAGEMENT_ROLES = ['owner', 'member', 'guest', 'invitation:member', 'invitation:guest'] ROLE_OWNER, ROLE_MEMBER, ROLE_GUEST, PENDING_MEMBER, PENDING_GUEST = MANAGEMENT_ROLES MANAGEMENT_ROLE_SET = set(MANAGEMENT_ROLES) def set_roles(user_id, caw_schema, management_schema, role_names, schema_type): """role_names - list of role names. Empty list removes previous roles record.""" if MANAGEMENT_ROLE_SET.issuperset(role_names): if role_names: rds_query(f"INSERT INTO {management_schema}.user_roles (id,company_alliance_workspace, user_id, roles) " f"VALUES (%(caw_schema)s, %(schema_type)s, %(user_id)s, %(roles)s) " f"ON CONFLICT (id, user_id) DO UPDATE SET roles = EXCLUDED.roles;", {'user_id': user_id, 'caw_schema': caw_schema, 'schema_type': schema_type, 'roles': json.dumps(role_names)}) else: rds_query(f"DELETE FROM {management_schema}.user_roles " f"WHERE id = %(caw_schema)s AND company_alliance_workspace = %(schema_type)s AND " f"user_id = %(user_id)s;", {'user_id': user_id, 'caw_schema': caw_schema, 'schema_type': schema_type}) else: raise RuntimeError('Invalid role') databases = [ # 'fansifter-rds', 'fansifter-rds-test', #'fansifter-rds-live', ] other_buckets = set() for connection in databases: all_chemas = rds_query(connection, "select DISTINCT default_schema from commons.user_company", fetch=True) print(f"Found {len(all_chemas)} schemas in database {connection}") for management_company_id in [sch['default_schema'] for sch in all_chemas]: print(f'Working on: {management_company_id}') try: for line in rds_query(connection, f"SELECT current_package_id from {management_company_id}.company"): break else: raise RuntimeError("Doesn't have company record") except UndefinedTable as e: pass except Exception as e: log.exception("wut?", exc_info=e) try: statements = [ # f"""alter table {management_company_id}.company add current_package_id int;""", # f"""alter table {management_company_id}.company add subscription_valid_until_date date;""", f"""INSERT INTO {management_company_id}.company (company_id, name, current_package_id) VALUES ('{management_company_id}', (SELECT company_name FROM commons.user_company WHERE default_schema = '{management_company_id}'), 1); """ # f"""CREATE TABLE IF NOT EXISTS {management_company_id}.meta_data ( # key varchar(128), # val text # );""", # f"""CREATE UNIQUE INDEX IF NOT EXISTS idx_meta_data_key ON {management_company_id}.meta_data(key);""", # f""" # INSERT INTO {management_company_id}.workspace (id) VALUES ('{management_company_id}') # ON CONFLICT (id) DO NOTHING; # """, # f"""DELETE FROM {management_company_id}.workspace WHERE id = '{{schema_name}}' """ # f"""INSERT INTO {management_company_id}.meta_data (key,val) # VALUES ('management_version', 'mvp 2.0') # ON CONFLICT (key) DO NOTHING;""" ] # statements = [ # f"ALTER TABLE {schema}.{table} SET (autovacuum_vacuum_scale_factor=0.3);", # f"ALTER TABLE {schema}.{table} SET (autovacuum_analyze_scale_factor=0.4);", # f"ALTER TABLE {schema}.{table} SET (autovacuum_vacuum_threshold=10000);", # f"ALTER TABLE {schema}.{table} SET (autovacuum_analyze_threshold=15000);" # ] try: rds_query(connection, statements, {}) except Exception as e: log.exception("oops", exc_info=e) except Exception as e: print(f"Oh, {management_company_id} doesn't even have meta_data")