from typing import Dict from internal.helpers import get_schema_id, rds_query, get_email, get_id_and_management_schema_for_workspace, \ generate_schema_id from internal.commons import BUCKET import datetime import logging log = logging.getLogger() ## # Pre request filters # handler(event,request_params) -> request_params ## def default_arguments(event: Dict, req: Dict) -> Dict: if 'arguments' in event: req.update(event['arguments']) return req def async_source(event: Dict, req: Dict) -> Dict: req['async_source'] = f"appsync.{event['field']}" # TODO: req['async_response_target'] = api_name return req def appsync_field_param(event: Dict, req: Dict) -> Dict: if 'field' in event: req['appsync_field'] = event['field'] return req def make_clientid_param(event: Dict, req: Dict) -> Dict: try: default_schema = get_schema_id(event) if 'workspaceId' in req and req['workspaceId'] is not None and req['workspaceId'] != default_schema: workspace_schema, _ = get_id_and_management_schema_for_workspace(default_schema, req['workspaceId']) req['workspace_schema'] = workspace_schema req['management_schema'] = default_schema else: req['workspace_schema'] = default_schema req['management_schema'] = default_schema except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def generate_clientid_param(event: Dict, req: Dict) -> Dict: try: default_schema = generate_schema_id(event) req['new_schema'] = default_schema req['management_schema'] = default_schema except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def make_email_param(event: Dict, req: Dict) -> Dict: try: req['email'] = get_email(event) except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def make_user_id_param(event: Dict, req: Dict) -> Dict: try: req['user_id'] = event['userSub'] if 'userSub' in event else \ event['request']['userAttributes']['sub'] except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def make_company_schema_param(event: Dict, req: Dict) -> Dict: req['caw'] = 'company' return req def make_initiate_file_upload_params(event: Dict, req: Dict) -> Dict: try: req['bucket'] = BUCKET['user'] except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def make_s3csv_query_params(event: Dict, req: Dict) -> Dict: bucket_query = f"SELECT source FROM {req['workspace_schema']}.collection WHERE id = %(collectionId)s" for bk in rds_query(bucket_query, {'collectionId': req['collectionId']}): source_val = bk["source"] if source_val.startswith('klaviyo'): log.info('-------------- kalviyo make_s3csv_query_params -------------------------') req['klaviyo'] = {'collectionId': req['collectionId']} else: log.info(f"bucket_key: {source_val}") sl = source_val.find('/') if sl > -1: req['object_params'] = {'Bucket': source_val[:sl], 'Key': source_val[sl + 1:]} else: raise RuntimeError("Source for collection {req['collectionId']} was not in 'bucket/key' format, weird") return req else: raise RuntimeError('Invalid collectionId') def make_s3_query_params(event: Dict, req: Dict) -> Dict: try: key = req.get('key', req.get('fileName', req.get('fileKey', None))) if key is None: raise KeyError('Need one argument named key, fileName or fileKey') req['object_params'] = {'Bucket': BUCKET[req['bucket']], 'Key': key} except KeyError as e: log.exception(e) raise RuntimeError("Please provide correct parameters in appsync RequestMappingTemplate") return req def make_rds_params(event: Dict, req: Dict) -> Dict: if 'queries' in event: q = event['queries'][0] if 'query' in q: rq = q['query'] if '' in rq: default_schema = get_schema_id(event) if 'workspaceId' in req and req['workspaceId'] != default_schema: workspace_id, _ = get_id_and_management_schema_for_workspace(default_schema, req['workspaceId']) schema = workspace_id else: schema = default_schema rq = rq.replace('', schema) # sad hacky day req['query'] = rq if 'params' in q: req['params'] = q['params'] return req raise RuntimeError("Did not find query definition") def decode_alliance_workspace_id_params(event: Dict, req: Dict) -> Dict: for inparam in ['collectionId', 'audienceId', 'sourceAudienceId', 'parentId', 'setId', 'memberId', 'sourceId']: if inparam in req and req[inparam] is not None: try: parts = str(req[inparam]).split('-') req[inparam] = '-'.join(parts[1:]) if len(parts) >= 2: schema_marker = parts[0][0] if schema_marker == 'a': req['allianceId'] = parts[0] elif schema_marker in 'cw': req['workspaceId'] = parts[0] else: raise RuntimeError(f'Invalid {inparam} schema_marker selection') # pre-workspaces legacy helper TODO: this should be error after proper implementation. elif len(parts) == 1: req['workspaceId'] = get_schema_id(event) else: raise RuntimeError(f'Invalid {inparam}') except Exception as e: log.exception("decode_collection_id_param FAILED", exc_info=e) raise RuntimeError(f'Invalid {inparam}') # same for lists of id's for inparam in ['collectionIds']: if inparam in req and req[inparam] is not None and isinstance(req[inparam], list): try: decoded_list = [] for elem in req[inparam]: parts = str(elem).split('-') decoded_list.append(str(int(parts[-1]))) if len(parts) == 2: schema_marker = parts[0][0] if schema_marker == 'a': req['allianceId'] = parts[0] elif schema_marker in 'cw': req['workspaceId'] = parts[0] else: raise RuntimeError(f'Invalid {inparam} schema_marker') else: raise RuntimeError(f'Invalid {inparam}') req[inparam] = decoded_list except Exception as e: log.exception("decode_collection_id_param FAILED", exc_info=e) raise RuntimeError(f'Invalid {inparam}') return req def make_current_package(event: Dict, req: Dict) -> Dict: """ Check current package and make sure current subscription has not expired. This needs to happen after management_schema is already in request parameter set. Also make sure it is the current user management_schema, not alliance's management_schema. :param event: :param req: :return: """ try: management_schema = req['management_schema'] if 'management_schema' in req else get_schema_id(event) for company in rds_query(f"SELECT current_package_id, subscription_valid_until_date FROM {management_schema}.company"): if company['subscription_valid_until_date'] is not None and \ company['subscription_valid_until_date'] >= datetime.datetime.now().date(): # go on. if company['current_package_id'] is not None: req['current_package_id'] = company['current_package_id'] else: for default in rds_query(f"SELECT id FROM commons.packages WHERE package_name ='default'"): req['current_package_id'] = default['id'] else: # default rights for default in rds_query(f"SELECT id FROM commons.packages WHERE package_name ='default'"): req['current_package_id'] = default['id'] break else: raise RuntimeError('Missing company info') except Exception as e: log.exception("PACKAGE CHECK FAILED", exc_info=e) return req