"""Lambda fivetran-sync function module.""" from typing import Any, Dict import sentry_sdk from lambdacommon.common_config import logger from pyfivetran import FivetranClient from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration import config from src.connectors.boto_clients import secrets_manager_client as sm if config.SENTRY_DSN: sentry_sdk.init( dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)], ) def get_fivetran_credentials() -> Dict[str, str]: """Fetch Fivetran credentials from AWS Secrets Manager.""" secrets = { "secret_id": config.FIVETRAN_SM_KEY, "secret_key": config.FIVETRAN_SM_KEY_SECRET, } try: return { key: sm.get_secret_value(SecretId=value)["SecretString"] for key, value in secrets.items() } except Exception as e: logger.error(f"Failed to retrieve secrets: {str(e)}") raise def handler(event: Dict[str, Any], context: Any) -> Dict[str, Any]: """Lambda entry point.""" try: connector_id = event.get("connector_id") requires_historical_sync = event.get("historical_sync", "false") if not connector_id: raise ValueError('Missing required input: "connector_id"') creds = get_fivetran_credentials() client = FivetranClient(creds["secret_id"], creds["secret_key"]) try: # Get connector details connector = client.connector_endpoint.get_connector(connector_id) logger.info("Connector details:") logger.info(f"ID: {connector.fivetran_id}") logger.info(f"Service: {connector.service}") logger.info(f"Schema: {connector.schema}") logger.info(f"Paused: {connector.paused}") logger.info(f"Sync Frequency: {connector.sync_frequency}") # Trigger sync if requires_historical_sync == "true": sync_response = connector.resync() message = f"Historical re-sync triggered \ for connector '{connector_id}'" else: sync_response = connector.sync() message = f"Sync triggered for connector '{connector_id}'" logger.info("Sync triggered:") logger.info(sync_response) return { "message": message, "fivetran_response": sync_response, } except Exception as e: logger.error(f"Error fetching connector or triggering sync: {e}") raise except Exception as e: logger.error(f"Error triggering sync: {str(e)}", exc_info=True) raise