import json import logging import botocore.config from airflow.providers.amazon.aws.hooks.lambda_function import AwsLambdaHook from fansifter.constants import AWS_REGION class FanSifterLambdaHook: default_aws_config = dict(read_timeout=900, connect_timeout=900, retries={"max_attempts": 0}) def __init__(self, function_name, invocation_type="RequestResponse", aws_config=None, region_name=AWS_REGION): self.function_name = function_name self.invocation_type = invocation_type self.region_name = region_name self.aws_config = ( aws_config if aws_config is not None else botocore.config.Config(**dict(read_timeout=900, connect_timeout=900, retries={"max_attempts": 0})) ) def invoke_lambda(self, payload): lambda_hook = AwsLambdaHook( function_name=self.function_name, invocation_type=self.invocation_type, region_name=self.region_name, config=self.aws_config, ) payload = payload if isinstance(payload, str) else json.dumps(payload) response = lambda_hook.invoke_lambda(payload) logging.info(f"Response from {self.function_name}:") logging.info(response) lambda_result = json.loads(response["Payload"].read()) if lambda_result is None: return None else: logging.info("Lambda response payload:") logging.info(lambda_result) error = lambda_result.get("error") or lambda_result.get("errorMessage") if error is not None: raise RuntimeError( f"Lambda call to {self.function_name} with payload {json.dumps(payload)} has FAILED and returned an error: {error}" ) return lambda_result