11import os
2+ from calendar import EPOCH
23from datetime import UTC , datetime
34from io import BytesIO
45from pathlib import Path
6+ from typing import Any
57
68import boto3
79import snowflake .connector
10+ from aws_lambda_powertools import Logger
11+ from aws_lambda_powertools .utilities .parameters import SSMProvider
12+ from aws_lambda_powertools .utilities .typing import LambdaContext
813from botocore .config import Config
914from botocore .exceptions import ClientError
1015from cryptography .hazmat .backends import default_backend
1116from cryptography .hazmat .primitives import serialization
1217from jinja2 import Environment , PackageLoader , StrictUndefined
1318from snowflake .connector import DictCursor
1419
15- from aws_lambda_powertools import Logger
16- from aws_lambda_powertools .utilities .typing import LambdaContext
17- from aws_lambda_powertools .utilities .parameters import SSMProvider
18-
1920REGION = os .environ .get ("AWS_CURRENT_REGION" , "us-east-1" )
2021BFD_ENV = os .environ .get ("BFD_ENV" , "prod" )
2122PARTNERS = os .environ .get ("PARTNERS" , "" )
4344logger = Logger ()
4445
4546
46- def execute_query (query : str ) -> list :
47+ def execute_query (query : str ) -> list [ dict [ str , Any ]] :
4748 """Execute the given query and return the resultant rows from Snowflake.
4849
4950 Args:
@@ -53,54 +54,51 @@ def execute_query(query: str) -> list:
5354 List: The results of the query.
5455 """
5556 try :
56- account = SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_account" , decrypt = True )
57- database = SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_database" , decrypt = True )
58- private_key_raw = SSM .get (
59- f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_private_key" , decrypt = True
57+ account = str (SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_account" , decrypt = True ))
58+ database = str (SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_database" , decrypt = True ))
59+ private_key_raw = str (
60+ SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_private_key" , decrypt = True )
61+ )
62+ schema = str (SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_schema" , decrypt = True ))
63+ user = str (SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_username" , decrypt = True ))
64+ warehouse = str (
65+ SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_warehouse" , decrypt = True )
6066 )
61- schema = SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_schema" , decrypt = True )
62- user = SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_username" , decrypt = True )
63- warehouse = SSM .get (f"/bfd/{ BFD_ENV } /idr-pipeline/sensitive/idr_warehouse" , decrypt = True )
6467 except Exception as exc :
65- raise ValueError (f"Missing snowflake configuration: { exc } " )
68+ raise ValueError (f"Missing snowflake configuration: { exc } " ) from exc
69+
70+ # Load and prepare private key for authentication
71+ private_key = serialization .load_pem_private_key (
72+ private_key_raw .encode (),
73+ password = None ,
74+ backend = default_backend (),
75+ )
76+ private_key_bytes = private_key .private_bytes (
77+ encoding = serialization .Encoding .DER ,
78+ format = serialization .PrivateFormat .PKCS8 ,
79+ encryption_algorithm = serialization .NoEncryption (),
80+ )
6681
6782 try :
68- # Load and prepare private key for authentication
69- private_key = serialization .load_pem_private_key (
70- private_key_raw .encode (),
71- password = None ,
72- backend = default_backend (),
73- )
74- private_key_bytes = private_key .private_bytes (
75- encoding = serialization .Encoding .DER ,
76- format = serialization .PrivateFormat .PKCS8 ,
77- encryption_algorithm = serialization .NoEncryption (),
78- )
79-
8083 # Establish connection
81- conn = snowflake .connector .connect (
84+ with snowflake .connector .connect (
8285 account = account ,
8386 user = user ,
8487 private_key = private_key_bytes ,
8588 warehouse = warehouse ,
8689 database = database ,
8790 schema = schema ,
88- )
89-
90- cursor = conn .cursor (DictCursor )
91- cursor .execute (query )
92- results = cursor .fetchall ()
93- cursor .close ()
94- return results
95-
96- except snowflake .connector .errors .Error as e :
91+ ) as conn :
92+ cursor = conn .cursor (DictCursor )
93+ cursor .execute (query )
94+ results = cursor .fetchall ()
95+ cursor .close ()
96+
97+ return results
98+ except snowflake .connector .errors .Error :
9799 logger .exception ("Snowflake error: {e}" )
98100 raise
99101
100- finally :
101- if "conn" in locals ():
102- conn .close ()
103-
104102
105103class PartnerPreferences :
106104 def __init__ (self , partner : str ) -> None :
@@ -122,9 +120,10 @@ def last_execution(self) -> str | None:
122120 return self ._execution
123121
124122 except ClientError as e :
125- self ._logger .exception (f"""
126- Error retrieving last execution: { e .response ["Error" ]["Message" ]}
127- """ )
123+ self ._logger .exception (
124+ "Error retrieving last execution: %s" ,
125+ e .response ["Error" ]["Message" ], # pyright: ignore[reportAttributeAccessIssue]
126+ )
128127 raise
129128
130129 return self ._execution
@@ -140,8 +139,8 @@ def partner_code(self) -> str:
140139 codes = {"ab2d" : "MED_AB2D" , "bcda" : "MED" , "dpc" : "MED_DPC" }
141140 try :
142141 code = codes [self .partner ]
143- except KeyError as exc :
144- self ._logger .exception (f "Invalid partner: { exc } " )
142+ except KeyError :
143+ self ._logger .exception ("Invalid partner: %s" , self . partner )
145144 raise
146145
147146 return code
@@ -161,13 +160,14 @@ def _set_last_execution(self, timestamp: str | None = None) -> None:
161160 )
162161
163162 self ._logger .info (
164- f "Updating last_exeuction from { self .last_execution } to { latest_timestamp } "
163+ "Updating last_execution from %s to %s" , self .last_execution , latest_timestamp
165164 )
166165 self ._execution = latest_timestamp
167166
168167 except ClientError as e :
169168 self ._logger .exception (
170- f"Error updating last execution: { e .response ['Error' ]['Message' ]} "
169+ "Error updating last execution: %s" ,
170+ e .response ["Error" ]["Message" ], # pyright: ignore[reportAttributeAccessIssue]
171171 )
172172 raise
173173
@@ -181,7 +181,7 @@ def _store_preferences(
181181 if store_local :
182182 local_file = Path (file_name ).name
183183 with Path .open (Path (local_file ), "wb" ) as local :
184- self ._logger .info (f "Storing local file:{ local_file } " )
184+ self ._logger .info ("Storing local file: %s" , local_file )
185185 local .write (preferences_data .encode ("utf-8" ))
186186
187187 if store_remote :
@@ -191,7 +191,7 @@ def _store_preferences(
191191
192192 bucket = SSM .get (f"/bfd/{ BFD_ENV } /bene-prefs/{ self .partner } /nonsensitive/bucket" )
193193
194- self ._logger .info (f "Storing remote file: s3://{ bucket } / { file_name } " )
194+ self ._logger .info ("Storing remote file: s3://%s/%s" , bucket , file_name )
195195 S3 .upload_fileobj (buffer , bucket , file_name )
196196
197197 def generate_preferences (
@@ -204,16 +204,21 @@ def generate_preferences(
204204 """Generate and store the preferences report in AWS S3.
205205
206206 Args:
207- timestamp_range (tuple): Tuple bounding (lower, upper) ISO 8601 timestamp strings. Default range between last execution and utc-now.
208- set_last_execution (bool): When `True`, update last_execution in DynamoDB . Default is `True`.
207+ timestamp_range (tuple): Tuple bounding (lower, upper) ISO 8601 timestamp strings.
208+ Default range between last execution and utc-now.
209+ set_last_execution (bool): When `True`, update last_execution in DynamoDB . Default is
210+ `True`.
209211 store_local (bool): When `True`, store preferences locally. Default is `False`.
210212 store_remote (bool): When `True`, write the preferences to S3. Default is `True`.
211213 """
212214 query_since_timestamp = timestamp_range [0 ] or self .last_execution
213215 query_until_timestamp = timestamp_range [1 ] or datetime .now (UTC ).isoformat ()
214216
215217 self ._logger .info (
216- f"Rendering query template { self .query_template .filename } with query_since_timestamp={ query_since_timestamp } , query_until_timestamp={ query_until_timestamp } "
218+ "Rendering query template %s with query_since_timestamp=%s, query_until_timestamp=%s" ,
219+ self .query_template .filename ,
220+ query_since_timestamp ,
221+ query_until_timestamp ,
217222 )
218223
219224 query = self .query_template .render (
@@ -227,14 +232,25 @@ def generate_preferences(
227232 results = execute_query (query )
228233
229234 self ._logger .info (
230- f"Rendering prefs template { self .prefs_template .filename } of { len (results )} records with data=<redacted>, extract_date={ YYYYMMDD } "
235+ "Rendering prefs template %s of %d records with data=<redacted>, extract_date=%s" ,
236+ self .prefs_template .filename ,
237+ len (results ),
238+ YYYYMMDD ,
231239 )
232240 preferences_data = self .prefs_template .render (data = results , extract_date = YYYYMMDD )
233241
234242 report_time = datetime .now (UTC ).strftime ("%H%M%S" )
235243
236244 self ._logger .info (
237- f"Rendering file_name template { self .file_name_template .filename } with partner={ self .partner } , env_indicator={ self .environment_indicator } , report_date={ YYYYMMDD } , report_time={ report_time } "
245+ (
246+ "Rendering file_name template %s with partner=%s, env_indicator=%s, report_date=%s,"
247+ " report_time=%s"
248+ ),
249+ self .file_name_template .filename ,
250+ self .partner ,
251+ self .environment_indicator ,
252+ YYYYMMDD ,
253+ report_time ,
238254 )
239255
240256 file_name = self .file_name_template .render (
@@ -250,17 +266,18 @@ def generate_preferences(
250266
251267 if set_last_execution :
252268 if results :
269+ max_insert_datetime : datetime | None = max (
270+ results , key = lambda x : x .get ("IDR_INSRT_TS" , datetime (EPOCH , 1 , 1 ))
271+ ).get ("IDR_INSRT_TS" )
253272 max_insert_timestamp = (
254- max (results , key = lambda x : x .get ("IDR_INSRT_TS" ))
255- .get ("IDR_INSRT_TS" )
256- .isoformat ()
273+ max_insert_datetime .isoformat () if max_insert_datetime else None
257274 )
258275 self ._set_last_execution (max_insert_timestamp )
259276 else :
260277 self ._logger .info ("Empty results. Skipping set of last execution." )
261278
262279
263- def handler (event : dict , context : LambdaContext ) -> None :
280+ def handler (event : dict [ str , Any ], context : LambdaContext ) -> None : # noqa: ARG001
264281 """Lambda event handler function.
265282
266283 Args:
@@ -270,6 +287,8 @@ def handler(event: dict, context: LambdaContext) -> None:
270287 Raises:
271288 RuntimeError: If any required environment variables are undefined
272289 """
290+ # PARTNERS is a comma-separated list of participating partners. The value of PARTNERS is derived
291+ # directly from the bene-prefs Terraservice
273292 partners = [p .strip () for p in PARTNERS .split ("," )]
274293
275294 if "bcda" in partners :
0 commit comments