"""Lambda copy-recovery-point function module.""" import os import datetime from collections import namedtuple from typing import Optional import boto3 import botocore from lambdacommon.common_config import logger import config backup = boto3.client('backup', region_name=config.AWS_REGION) dynamodb = boto3.client('dynamodb', region_name=config.AWS_REGION) def handler(event, context): """Lambda entry point.""" try: if config.SOURCE_NAME: info = find_resource(config.BACKUP_VAULT_NAME, config.SOURCE_NAME) else: info = get_info(config.BACKUP_VAULT_NAME) recovery_point = info.recovery_point_arn status = info.backup_vault_status table_arn = info.resource_arn source_vault = info.backup_vault_arn logger.info(f'Received Recovery Point State Change {status} event for source_vault: {source_vault}') if table_arn == config.PROD_TABLE: logger.info(f'Deleting table {config.REFRESH_TABLE_NAME}') delete_table(config.REFRESH_TABLE_NAME) logger.info(f'Attempting start_restore_job for table: {config.REFRESH_TABLE_NAME}') start_recovery_job(recovery_point, config.REFRESH_TABLE_NAME, config.REFRESH_TABLE_KMS_KEY_ARN, config.REFRESH_IAM_ROLE) logger.info('Restore job started successfully') return {'status': 'OK'} return {'status': 'OK'} except Exception as e: logger.exception(str(e)) raise e def delete_table(table_name): """Delete DynamoDB table. Args: table_name (string): DynamoDB table name Returns: boolean : Returns True if the table deletion is completed or the table not exist. """ try: logger.info(f'Issuing delete for {table_name}') dynamodb.delete_table(TableName=table_name) dynamodb.get_waiter('table_not_exists').wait(TableName=table_name) logger.info(f'Table {table_name} delete complete') return True except dynamodb.exceptions.ResourceNotFoundException: print(f'Table {table_name} did not exist') return True except Exception as e: logger.error(f'Failed to delete table {table_name}: {str(e)}', exc_info=True) return False def start_recovery_job(recovery_point, refresh_table_name, kms_key_arn, iam_role): """Start recovery job. Args: recovery_point (sting): recovery point arn refresh_table_name (sting): table name for restore kms_key_arn (sting): Backup encryption KMS key iam_role (sting): IAM role that Backup uses to create the target resource Returns: Respond from the start_restore_job """ res = backup.start_restore_job( RecoveryPointArn=recovery_point, Metadata={ 'targetTableName': refresh_table_name, 'encryptionType': 'KMS', 'keyArn': kms_key_arn, }, ResourceType='DynamoDB', IamRoleArn=iam_role, IdempotencyToken=datetime.datetime.now().strftime('%Y-%m-%d-%H-%M-%S') ) logger.info(f'Looking for table {refresh_table_name} to see if it is available.') dynamodb.get_waiter('table_exists').wait( TableName=refresh_table_name, WaiterConfig={ 'Delay': 60, 'MaxAttempts': 10}) logger.info(f'start_restore_job result: {res}') return res BackupRecoveryInfo = namedtuple( 'Info', ['recovery_point_arn', 'backup_vault_status', 'resource_arn', 'backup_vault_arn']) def find_resource( backup_vault_name: str, needle: str, ) -> Optional[BackupRecoveryInfo]: """ Search for a recovery point in a backup vault by resource name and status 'COMPLETED'. Args: backup_vault_name: Name of the backup vault. needle: Resource name to search for. Returns: BackupRecoveryInfo object if found, otherwise None. """ logger.info(f'Searching for recovery point in vault {backup_vault_name} with resource name {needle}') next_token = None while True: try: params = { 'BackupVaultName': backup_vault_name, 'ByResourceType': 'DynamoDB', 'MaxResults': 100, } if next_token: params['NextToken'] = next_token res = backup.list_recovery_points_by_backup_vault(**params) except Exception as e: logger.error('Error calling list_recovery_points_by_backup_vault: %s', e) return None for point in res.get('RecoveryPoints', []): if point.get('ResourceName') == needle and point.get('Status') == 'COMPLETED': recovery_point_arn = point['RecoveryPointArn'] resource_arn = point['ResourceArn'] backup_vault_status = point['Status'] backup_vault_arn = point['BackupVaultArn'] logger.info(f'Recovery point ARN: {recovery_point_arn}') logger.info(f'Resource ARN: {resource_arn}') logger.info(f'Backup Vault status: {backup_vault_status}') logger.info(f'Backup Vault ARN: {backup_vault_arn}') return BackupRecoveryInfo(recovery_point_arn, backup_vault_status, resource_arn, backup_vault_arn) next_token = res.get('NextToken') if not next_token: break return None def get_info(backup_vault_name): """Collect info. Args: backup_vault_name (string): Backup Vault Name Returns: Multiple value from respond to pass in handler. """ try: res = backup.list_recovery_points_by_backup_vault( BackupVaultName=backup_vault_name, MaxResults=1) if 'RecoveryPoints' in res and len(res['RecoveryPoints']) > 0: logger.info(res['RecoveryPoints']) recovery_point_arn = res['RecoveryPoints'][0]['RecoveryPointArn'] resource_arn = res['RecoveryPoints'][0]['ResourceArn'] backup_vault_status = res['RecoveryPoints'][0]['Status'] backup_vault_arn = res['RecoveryPoints'][0]['BackupVaultArn'] logger.info(f'Recovery point ARN: {recovery_point_arn}') logger.info(f'Resource ARN: {resource_arn}') logger.info(f'Backup Vault status: {backup_vault_status}') logger.info(f'Backup Vault ARN: {backup_vault_arn}') return BackupRecoveryInfo(recovery_point_arn, backup_vault_status, resource_arn, backup_vault_arn) else: logger.warning('No recovery points found for backup vault: %s', backup_vault_name) return None except Exception as e: logger.exception('An error occurred while retrieving recovery points.') return e