diff --git a/apps/bfd-db-migrator-ng/migrations/V1_40_0__add_utn_pattern_idx.sql b/apps/bfd-db-migrator-ng/migrations/V1_40_0__add_utn_pattern_idx.sql new file mode 100644 index 0000000000..1f643f32f0 --- /dev/null +++ b/apps/bfd-db-migrator-ng/migrations/V1_40_0__add_utn_pattern_idx.sql @@ -0,0 +1,2 @@ +CREATE INDEX ON idr.prior_auth (utn) WHERE utn NOT LIKE '-%'; +CREATE INDEX ON idr.prior_auth_item (utn) WHERE utn NOT LIKE '-%'; diff --git a/apps/bfd-pipeline-idr/load-credentials.sh b/apps/bfd-pipeline-idr/load-credentials.sh index f52943fc62..32237a20a0 100755 --- a/apps/bfd-pipeline-idr/load-credentials.sh +++ b/apps/bfd-pipeline-idr/load-credentials.sh @@ -12,6 +12,8 @@ IDR_WAREHOUSE="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensi export IDR_WAREHOUSE IDR_DATABASE="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensitive/idr_database --with-decryption --query "Parameter.Value" --output text)" export IDR_DATABASE +IDR_EDP_DATABASE="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensitive/idr_edp_database --with-decryption --query "Parameter.Value" --output text)" +export IDR_EDP_DATABASE IDR_SCHEMA="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensitive/idr_schema --with-decryption --query "Parameter.Value" --output text)" export IDR_SCHEMA diff --git a/apps/bfd-pipeline-idr/load-synthetic-env.sh b/apps/bfd-pipeline-idr/load-synthetic-env.sh index 2da968b2c3..80f9509a3c 100755 --- a/apps/bfd-pipeline-idr/load-synthetic-env.sh +++ b/apps/bfd-pipeline-idr/load-synthetic-env.sh @@ -35,13 +35,10 @@ export IDR_WAREHOUSE IDR_DATABASE="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensitive/synthetic_env_database --with-decryption --query "Parameter.Value" --output text)" readonly IDR_DATABASE export IDR_DATABASE -IDR_SCHEMA="$(aws ssm get-parameter --name /bfd/${BFD_ENV}/idr-pipeline/sensitive/synthetic_env_schema --with-decryption --query "Parameter.Value" --output text)" -readonly IDR_SCHEMA -export IDR_SCHEMA args=('--load-type' 'initial' '--source' 'snowflake' '--load-mode' 'synthetic') if [[ -n "$1" ]]; then args+=('--seed-from' "$1") fi -IDR_ENABLE_DATE_PARTITIONS=0 IDR_ENABLE_PRIOR_AUTH=1 uv run idr-pipeline "${args[@]}" +IDR_ENABLE_DATE_PARTITIONS=0 uv run idr-pipeline "${args[@]}" diff --git a/apps/bfd-pipeline-idr/mock-idr.sql b/apps/bfd-pipeline-idr/mock-idr.sql index 54fc9aa3f3..499839fc6a 100644 --- a/apps/bfd-pipeline-idr/mock-idr.sql +++ b/apps/bfd-pipeline-idr/mock-idr.sql @@ -871,7 +871,5 @@ CREATE TABLE cms_edp_view_cvm_prau_prd.prauc ( mr_count_end_dt DATE, att_phy_npi VARCHAR(10) NOT NULL, rrb_excl_ind VARCHAR(1), - idr_insrt_ts TIMESTAMPTZ, - idr_updt_ts TIMESTAMPTZ, PRIMARY KEY(mbi_num, utn, current_segment) ); diff --git a/apps/bfd-pipeline-idr/run-db.sh b/apps/bfd-pipeline-idr/run-db.sh index d51a46bbe0..a9e59d751a 100755 --- a/apps/bfd-pipeline-idr/run-db.sh +++ b/apps/bfd-pipeline-idr/run-db.sh @@ -18,15 +18,14 @@ function do_load() { PGPASSWORD="$DB_PASSWORD" psql "host=$DB_ENDPOINT port=5432 dbname=fhirdb user=$DB_USERNAME" -f "$SCRIPT_DIR/mock-idr.sql" docker exec -u postgres bfd-idr-db psql fhirdb bfd -c "VACUUM FULL ANALYZE" BFD_DB_USERNAME="$DB_USERNAME" \ - BFD_DB_PASSWORD="$DB_PASSWORD" \ - BFD_DB_ENDPOINT="$DB_ENDPOINT" \ - IDR_ENABLE_DATE_PARTITIONS=0 \ - IDR_ENABLE_PRIOR_AUTH=1 \ - uv run idr-pipeline \ - --source postgres \ - --load-mode synthetic \ - --load-type initial \ - --seed-from "${1:-"${SCRIPT_DIR}/../bfd-model-idr/out"}" + BFD_DB_PASSWORD="$DB_PASSWORD" \ + BFD_DB_ENDPOINT="$DB_ENDPOINT" \ + IDR_ENABLE_DATE_PARTITIONS=0 \ + uv run idr-pipeline \ + --source postgres \ + --load-mode synthetic \ + --load-type initial \ + --seed-from "${1:-"${SCRIPT_DIR}/../bfd-model-idr/out"}" } image=postgres:16.6 diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/constants.py b/apps/bfd-pipeline-idr/src/idr_pipeline/constants.py index a2c061a46f..b04f04cbf3 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/constants.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/constants.py @@ -1,7 +1,12 @@ from dateutil.relativedelta import relativedelta from .load_partition import LoadPartition, LoadPartitionGroup, PartitionType -from .settings import PARTITION_TYPE +from .settings import IDR_DATABASE, IDR_EDP_DATABASE, PARTITION_TYPE + + +def _period_delimited(*segments: str) -> str: + return ".".join(segment for segment in segments if segment) + DEFAULT_MAX_DATE = "9999-12-31" DEFAULT_MIN_DATE = "0001-01-01" @@ -25,7 +30,7 @@ MCS_CLM_SOURCE = "22000" VMS_CLM_SOURCE = "23000" -IDR_PREFIX = "cms_vdm_view_mdcr_prd" +IDR_PREFIX = _period_delimited(IDR_DATABASE, "cms_vdm_view_mdcr_prd") IDR_BENE_HISTORY_TABLE = f"{IDR_PREFIX}.v2_mdcr_bene_hstry" IDR_BENE_MBI_TABLE = f"{IDR_PREFIX}.v2_mdcr_bene_mbi_id" IDR_BENE_XREF_TABLE = f"{IDR_PREFIX}.v2_mdcr_bene_xref" @@ -65,7 +70,7 @@ IDR_CONTRACT_PBP_CONTACT_TABLE = f"{IDR_PREFIX}.v2_mdcr_cntrct_pbp_cntct" IDR_CONTRACT_PBP_SEGMENT_TABLE = f"{IDR_PREFIX}.v2_mdcr_cntrct_pbp_sgmt" -IDR_PRIOR_AUTH_PREFIX = "cms_edp_view_cvm_prau_prd" +IDR_PRIOR_AUTH_PREFIX = _period_delimited(IDR_EDP_DATABASE, "cms_edp_view_cvm_prau_prd") IDR_PRIOR_AUTH_TABLE = f"{IDR_PRIOR_AUTH_PREFIX}.prauc" DEATH_DATE_CUTOFF_YEARS = 4 diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/extractor.py b/apps/bfd-pipeline-idr/src/idr_pipeline/extractor.py index 0f0dc8f6c9..ea41880a15 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/extractor.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/extractor.py @@ -33,9 +33,7 @@ BATCH_MULTIPLIER, ENABLE_DATE_PARTITIONS, IDR_ACCOUNT, - IDR_DATABASE, IDR_PRIVATE_KEY, - IDR_SCHEMA, IDR_USERNAME, IDR_WAREHOUSE, MIN_BATCH_COMPLETION_DATE, @@ -192,6 +190,14 @@ def extract_idr_data( {"timestamp": compare_timestamp}, ) + def extract_full_idr_data(self, source: Source) -> Iterator[list[T]]: + start_time = self.cls.model_type().min_transaction_date + fetch_query = self.get_query(start_time, source) + logger.info("extracting full {}", self.cls.table()) + return self.extract_many( + fetch_query.replace("{MIN_TS}", "%(timestamp)s"), {"timestamp": start_time} + ) + def _transform(self, batch: list[dict[str, DbType]]) -> list[T]: self.transform_timer.start() res = self.type_adapter.validate_python( @@ -318,8 +324,6 @@ def connect() -> SnowflakeConnection: private_key=private_key_bytes, account=IDR_ACCOUNT, warehouse=IDR_WAREHOUSE, - database=IDR_DATABASE, - schema=IDR_SCHEMA, ) @override @@ -379,8 +383,6 @@ def __init__(self) -> None: "user": IDR_USERNAME, "private_key": private_key_bytes, # type: ignore "warehouse": IDR_WAREHOUSE, - "database": IDR_DATABASE, - "schema": IDR_SCHEMA, } ).create() self.conn = SnowflakeExtractor.connect() diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/loader.py b/apps/bfd-pipeline-idr/src/idr_pipeline/loader.py index 97aebf7158..d0f477380f 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/loader.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/loader.py @@ -1,8 +1,9 @@ +import functools import itertools import operator from collections.abc import Awaitable, Callable, Iterator, Sequence from datetime import UTC, datetime -from typing import Any, cast +from typing import Any, Generic, cast, override import anyio import psycopg @@ -78,7 +79,8 @@ async def _async_load( timeout=600, ) as pool: await pool.wait() - return await BatchLoader( + loader_cls = FullSyncBatchLoader if model.should_delete_missing() else BatchLoader + return await loader_cls( fetch_results, model, pool, @@ -91,7 +93,7 @@ async def _async_load( ).load() -class BatchLoader: +class BatchLoader(Generic[T]): # noqa: UP046 def __init__( self, fetch_results: Iterator[list[T]], @@ -116,19 +118,18 @@ def __init__( self.batch_start = datetime.now(UTC) self.insert_cols = list(model.insert_keys()) self.insert_cols.sort() - self.immutable = not model.update_timestamp_col() self.meta_keys = ( - ["bfd_created_ts"] if self.immutable else ["bfd_created_ts", "bfd_updated_ts"] + ["bfd_created_ts"] if model.is_immutable() else ["bfd_created_ts", "bfd_updated_ts"] ) self.cols_str = ", ".join(self.insert_cols) self.meta_keys_str = ", ".join(self.meta_keys) self.ordered_pkeys = model.ordered_pkeys() self.primary_keys_str = ", ".join(self.ordered_pkeys) - self.update_set = [v for v in self.insert_cols if v not in model.ordered_pkeys()] - self.update_set_str = ", ".join([f"{v}=EXCLUDED.{v}" for v in self.update_set]) - self.on_conflict_where_clause = ( - f"WHERE ({', '.join(f't.{v}' for v in self.update_set)}) IS " - f"DISTINCT FROM ({', '.join(f'EXCLUDED.{v}' for v in self.update_set)})" + update_set = [v for v in self.insert_cols if v not in self.ordered_pkeys] + update_set_str = ", ".join([f"{v}=EXCLUDED.{v}" for v in update_set]) + on_conflict_where_clause = ( + f"WHERE ({', '.join(f't.{v}' for v in update_set)}) IS " + f"DISTINCT FROM ({', '.join(f'EXCLUDED.{v}' for v in update_set)})" ) # For immutable tables, we may still be attempting to re-load some data # due to a batch cancellation. @@ -137,24 +138,24 @@ def __init__( # Additionally, if there are no extra columns to update, we can skip it. self.on_conflict_clause = ( "DO NOTHING" - if self.immutable or not self.update_set + if model.is_immutable() or not update_set else ( - f"DO UPDATE SET {self.update_set_str}, bfd_updated_ts=%(timestamp)s " - f"{self.on_conflict_where_clause}" + f"DO UPDATE SET {update_set_str}, bfd_updated_ts=%(timestamp)s " + f"{on_conflict_where_clause}" ) ) # Used in _upsert so that relevant primary/last updated timestamp columns are returned for # rows that are actually updated during the upsert so that last updated can be ran for # just rows with changes during the load self.updated_keys_returning_str = ", ".join( - set( + { col for col in [ - *self.model.ordered_pkeys(), + *self.ordered_pkeys, self.model.last_updated_timestamp_col(), ] if col - ) + } ) self.timestamp_placeholders = ", ".join("%(timestamp)s" for _ in self.meta_keys) @@ -167,32 +168,20 @@ def __init__( self.full_batch_timer = Timer("full_batch", model, partition) self.full_load_timer = Timer("full_load", model, partition) self.load_type = load_type + self.load_mode = load_mode self.enable_load_progress = should_track_load_progress(load_mode) async def load(self) -> bool: timestamp = datetime.now(UTC) - self.full_load_timer.start() async with self.pool.connection() as conn, conn.cursor(binary=True) as cur: - self.progress_start_timer.start() - await self._insert_batch_start(cur) - await conn.commit() - self.progress_start_timer.stop() + await self._record_batch_start(conn, cur, commit=True) - data_loaded = False - num_rows = 0 batch_num = 1 - while True: - self.idr_query_timer.start() - # We unfortunately need to use a while true loop here since we need to wrap the - # iterator with the timer calls. - results = next(self.fetch_results, None) - self.idr_query_timer.stop() - if not results: - break + async def _process_batch(results: list[T]) -> None: + nonlocal batch_num self.full_batch_timer.start() - data_loaded = True logger.info( "{}-{}-{}: loading next {} results concurrently {} row(s) at a time", self.table, @@ -201,8 +190,6 @@ async def load(self) -> bool: len(results), PER_BATCH_CONCURRENT_ROWS, ) - num_rows += len(results) - self.sort_batch_timer.start() results.sort(key=operator.attrgetter(*self.ordered_pkeys)) self.sort_batch_timer.stop() @@ -258,6 +245,9 @@ async def _wrap_batch_chunk( batch_num += 1 self.full_batch_timer.stop() + num_rows = await self._stage_all_batches(_process_batch) + data_loaded = num_rows > 0 + # Wait until the background worker signals that all pending loading tasks are completed # for the current partition before marking it totally complete self.worker_client.wait_until_done(self.model, self.partition) @@ -342,7 +332,10 @@ async def _mark_batch_complete(self, cur: psycopg.AsyncCursor) -> None: ) async def _setup_temp_table( - self, cur: psycopg.AsyncCursor[Any], suffix: str | None = None + self, + cur: psycopg.AsyncCursor[Any], + suffix: str | None = None, + copy_primary_key: bool = False, ) -> str: # Load each batch into a temp table # This is necessary because we want to use COPY to quickly @@ -354,9 +347,12 @@ async def _setup_temp_table( # For simplicity's sake, we'll create our temp tables using the existing schema and # just drop the columns we need to ignore. full_tablename = f"{self.temp_table}_{suffix or ''}" + copy_primary_key_option = ( + f", PRIMARY KEY ({self.primary_keys_str})" if copy_primary_key else "" + ) await cur.execute( - f'CREATE TEMPORARY TABLE "{full_tablename}" (LIKE {self.table}) ' # type: ignore - "ON COMMIT DROP" + f'CREATE TEMPORARY TABLE "{full_tablename}" ' # type: ignore + f"(LIKE {self.table} {copy_primary_key_option}) ON COMMIT DROP" ) # Created/updated columns don't need to be loaded from the source. for col in self.meta_keys: @@ -417,6 +413,105 @@ async def _copy_data( [_remove_null_bytes(getattr(row, k)) for k in self.insert_cols] ) + async def _record_batch_start( + self, conn: psycopg.AsyncConnection, cur: psycopg.AsyncCursor[Any], commit: bool + ) -> None: + self.progress_start_timer.start() + await self._insert_batch_start(cur) + if commit: + await conn.commit() + self.progress_start_timer.stop() + + def _next_batch(self) -> list[T] | None: + self.idr_query_timer.start() + results = next(self.fetch_results, None) + self.idr_query_timer.stop() + return results + + async def _stage_all_batches(self, process_batch: Callable[[list[T]], Awaitable[None]]) -> int: + num_rows = 0 + + while True: + # We unfortunately need to use a while true loop here since we need to wrap the + # iterator with the timer calls. + self.idr_query_timer.start() + results = next(self.fetch_results, None) + self.idr_query_timer.stop() + if not results: + break + + num_rows += len(results) + await process_batch(results) + + return num_rows + + +class FullSyncBatchLoader(BatchLoader[T]): + @override + async def load(self) -> bool: + timestamp = datetime.now(UTC) + self.full_load_timer.start() + data_loaded = False + + async with self.pool.connection() as conn, conn.cursor(binary=True) as cur: + await self._record_batch_start(conn, cur, commit=False) + full_temp_table = await self._setup_temp_table(cur, "full_temp", copy_primary_key=True) + + num_rows = await self._stage_all_batches( + functools.partial(self._copy_data, cur, full_temp_table) + ) + data_loaded = num_rows > 0 + logger.info( + "{}-{}: staged {} row(s) for full sync", + self.table, + self.partition.name, + num_rows, + ) + + self.insert_batch_timer.start() + updated_keys = await self._upsert(cur, full_temp_table, timestamp) + deleted_count = await self._delete_missing(cur, full_temp_table) + self.insert_batch_timer.stop() + + logger.info( + "{}-{}: upserted {} new/changed row(s), deleted {} row(s) no longer present " + "upstream", + self.table, + self.partition.name, + len(updated_keys), + deleted_count, + ) + await self._mark_batch_complete(cur) + + self.full_load_timer.stop() + logger.info( + "{}-{}: finished full sync", + self.table, + self.partition.name, + ) + return data_loaded + + async def _delete_missing(self, cur: psycopg.AsyncCursor[Any], temp_tablename: str) -> int: + # We have to exclude our synthetic data that also exists in prod from deletion + synthetic_data_filter = self.model.synthetic_data_filter() + synthetic_where_clause = ( + f"WHERE {synthetic_data_filter}" + if synthetic_data_filter and self.load_mode != LoadMode.SYNTHETIC + else "" + ) + result = await cur.execute( # type: ignore + f''' + DELETE FROM {self.table} + WHERE ({self.primary_keys_str}) IN ( + SELECT {self.primary_keys_str} FROM {self.table} + {synthetic_where_clause} + EXCEPT + SELECT {self.primary_keys_str} FROM "{temp_tablename}" + ) + ''' # type: ignore + ) + return result.rowcount # type: ignore + def _remove_null_bytes(val: DbType) -> DbType: # Some IDR strings have null bytes. @@ -431,5 +526,5 @@ def _remove_null_bytes(val: DbType) -> DbType: def should_track_load_progress(load_mode: LoadMode) -> bool: - # Whether to read/write load progress, which is diabled for synthetic and testing loads. + # Whether to read/write load progress, which is disabled for synthetic and testing loads. return load_mode == LoadMode.PROD or force_load_progress() diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/base_model.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/base_model.py index f28ddf3c1b..80e1381766 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/base_model.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/base_model.py @@ -490,6 +490,20 @@ def should_replace() -> bool: """Whether to merge or replace data when loading this table.""" return False + @staticmethod + def should_delete_missing() -> bool: + """Whether upstream data deletion requires manual cleanup on our end. + + Upstream data can be deleted with no indicator like an obsolete timestamp, requiring + us to delete it on our end. + """ + return False + + @staticmethod + def synthetic_data_filter() -> str: + """Expression used to exclude synthetic data from being deleted in FullSyncBatchLoader.""" + return "" + @classmethod @abstractmethod def fetch_query( @@ -528,6 +542,10 @@ def batch_timestamp_col(cls, is_historical: bool) -> list[str]: def update_timestamp_col(cls) -> list[str]: return cls._extract_meta_keys(UPDATE_TIMESTAMP) + @classmethod + def is_immutable(cls) -> bool: + return not cls.update_timestamp_col() + @classmethod def batch_id_col_alias(cls) -> str | None: col = cls._single_or_default(BATCH_ID) diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_institutional_ss.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_institutional_ss.py index a7a0b10d55..5a86982638 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_institutional_ss.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_institutional_ss.py @@ -295,12 +295,12 @@ class IdrClaimInstitutionalSs(IdrBaseModel): # Columns from v2_mdcr_clm_fiss clm_crnt_stus_cd: Annotated[str, {ALIAS: ALIAS_FISS}, BeforeValidator(transform_default_string)] clm_pps_ind: Annotated[str, {ALIAS: ALIAS_FISS}, BeforeValidator(transform_default_string)] - idr_insrt_ts: Annotated[ + idr_insrt_ts_fiss: Annotated[ datetime, {ALIAS: ALIAS_FISS, **INSERT_FIELD}, BeforeValidator(transform_null_date_to_min), ] - idr_updt_ts: Annotated[ + idr_updt_ts_fiss: Annotated[ datetime, {ALIAS: ALIAS_FISS, **UPDATE_FIELD}, BeforeValidator(transform_null_date_to_min), diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_nch.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_nch.py index f6899b37d4..7423d8fadb 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_nch.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_nch.py @@ -97,12 +97,12 @@ class IdrClaimProfessionalNch(IdrBaseModel): clm_blg_prvdr_tax_num: Annotated[ str, {ALIAS: ALIAS_CLM}, BeforeValidator(transform_default_string) ] - idr_insrt_ts: Annotated[ + idr_insrt_ts_clm: Annotated[ datetime, {ALIAS: ALIAS_CLM, **INSERT_FIELD}, BeforeValidator(transform_null_date_to_min), ] - idr_updt_ts: Annotated[ + idr_updt_ts_clm: Annotated[ datetime, {ALIAS: ALIAS_CLM, **UPDATE_FIELD}, BeforeValidator(transform_null_date_to_min), diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_ss.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_ss.py index 7531b97568..cfc341177a 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_ss.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_claim_professional_ss.py @@ -113,12 +113,12 @@ class IdrClaimProfessionalSs(IdrBaseModel): clm_prvdr_rmng_due_amt: float | None clm_blood_ncvrd_chrg_amt: float | None clm_prvdr_intrst_pd_amt: float | None - idr_insrt_ts: Annotated[ + idr_insrt_ts_clm: Annotated[ datetime, {BATCH_TIMESTAMP: True, INSERT_EXCLUDE: True, ALIAS: ALIAS_CLM, COLUMN_MAP: "idr_insrt_ts"}, BeforeValidator(transform_null_date_to_min), ] - idr_updt_ts: Annotated[ + idr_updt_ts_clm: Annotated[ datetime, {UPDATE_TIMESTAMP: True, INSERT_EXCLUDE: True, ALIAS: ALIAS_CLM, COLUMN_MAP: "idr_updt_ts"}, BeforeValidator(transform_null_date_to_min), diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth.py index 7980bc066d..495eea5df0 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth.py @@ -10,11 +10,8 @@ ALIAS_PRVDR_ATT_PHY, ALIAS_PRVDR_ORDER_REFER, ALIAS_PRVDR_RENDER, - BATCH_TIMESTAMP, EXPR, - INSERT_EXCLUDE, PRIMARY_KEY_ORDER, - UPDATE_TIMESTAMP, IdrBaseModel, ModelType, Source, @@ -24,7 +21,6 @@ transform_null_date_to_max, transform_null_date_to_min, ) -from ..settings import MIN_PRIOR_AUTH_LOAD_DATE class IdrPriorAuth(IdrBaseModel): @@ -62,13 +58,6 @@ class IdrPriorAuth(IdrBaseModel): BeforeValidator(transform_default_string), ] bfd_att_phy_npi_type: Annotated[int | None, {EXPR: provider_npi_type_expr(ALIAS_PRVDR_ATT_PHY)}] - # TBD: might have to change the insert & update ts once IDR adds those - idr_insrt_ts: Annotated[datetime, {BATCH_TIMESTAMP: True, INSERT_EXCLUDE: True}] - idr_updt_ts: Annotated[ - datetime, - {UPDATE_TIMESTAMP: True, INSERT_EXCLUDE: True}, - BeforeValidator(transform_null_date_to_min), - ] @override @staticmethod @@ -85,6 +74,21 @@ def last_updated_date_column() -> list[str]: def model_type() -> ModelType: return ModelType.PRIOR_AUTH + @override + @staticmethod + def should_delete_missing() -> bool: + return True + + @override + @classmethod + def is_immutable(cls) -> bool: + return False + + @override + @staticmethod + def synthetic_data_filter() -> str: + return "utn NOT LIKE '-%'" + @override @classmethod def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Source) -> str: @@ -98,8 +102,8 @@ def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Sou SELECT *, ROW_NUMBER() OVER (PARTITION BY mbi_num, utn ORDER BY current_segment) as row_order FROM {IDR_PRIOR_AUTH_TABLE} - WHERE pa_req_rec_dt > '{MIN_PRIOR_AUTH_LOAD_DATE}' - ) + WHERE pa_req_rec_dt > {{MIN_TS}} + ) SELECT {{COLUMNS}} FROM distinct_prior_auths {prior_auth} LEFT JOIN {IDR_PROVIDER_HISTORY_TABLE} {prvdr_att_phy} ON {prvdr_att_phy}.prvdr_npi_num = {prior_auth}.att_phy_npi @@ -110,5 +114,5 @@ def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Sou LEFT JOIN {IDR_PROVIDER_HISTORY_TABLE} {prvdr_render} ON {prvdr_render}.prvdr_npi_num = {prior_auth}.render_npi AND {prvdr_render}.prvdr_hstry_obslt_dt >= '{DEFAULT_MAX_DATE}' - {{WHERE_CLAUSE}} AND row_order = 1; + WHERE row_order = 1; """ diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth_item.py b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth_item.py index c609e9412a..72e11cc5bc 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth_item.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/model/idr_prior_auth_item.py @@ -7,10 +7,7 @@ from ..load_partition import LoadPartition from ..model.base_model import ( ALIAS_PRIOR_AUTH, - BATCH_TIMESTAMP, - INSERT_EXCLUDE, PRIMARY_KEY_ORDER, - UPDATE_TIMESTAMP, IdrBaseModel, ModelType, Source, @@ -18,7 +15,6 @@ transform_null_date_to_max, transform_null_date_to_min, ) -from ..settings import MIN_PRIOR_AUTH_LOAD_DATE class IdrPriorAuthItem(IdrBaseModel): @@ -43,12 +39,6 @@ class IdrPriorAuthItem(IdrBaseModel): mr_count_st_dt: Annotated[date, BeforeValidator(transform_null_date_to_min)] mr_count_end_dt: Annotated[date, BeforeValidator(transform_null_date_to_max)] rrb_excl_ind: Annotated[str, BeforeValidator(transform_default_string)] - idr_insrt_ts: Annotated[datetime, {BATCH_TIMESTAMP: True, INSERT_EXCLUDE: True}] - idr_updt_ts: Annotated[ - datetime, - {UPDATE_TIMESTAMP: True, INSERT_EXCLUDE: True}, - BeforeValidator(transform_null_date_to_min), - ] @override @staticmethod @@ -65,6 +55,21 @@ def last_updated_date_column() -> list[str]: def model_type() -> ModelType: return ModelType.PRIOR_AUTH + @override + @staticmethod + def should_delete_missing() -> bool: + return True + + @override + @classmethod + def is_immutable(cls) -> bool: + return False + + @override + @staticmethod + def synthetic_data_filter() -> str: + return "utn NOT LIKE '-%'" + @override @classmethod def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Source) -> str: @@ -75,8 +80,8 @@ def fetch_query(cls, partition: LoadPartition, start_time: datetime, source: Sou SELECT *, ROW_NUMBER() OVER (PARTITION BY mbi_num, utn, current_segment ORDER BY mbi_num) as row_order FROM {IDR_PRIOR_AUTH_TABLE} - WHERE pa_req_rec_dt > '{MIN_PRIOR_AUTH_LOAD_DATE}' + WHERE pa_req_rec_dt > {{MIN_TS}} ) SELECT {{COLUMNS}} FROM distinct_prior_auths {prior_auth} - {{WHERE_CLAUSE}} AND row_order = 1; + WHERE row_order = 1; """ diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/pipeline_utils.py b/apps/bfd-pipeline-idr/src/idr_pipeline/pipeline_utils.py index 04125200c5..4dc07a0171 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/pipeline_utils.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/pipeline_utils.py @@ -110,7 +110,11 @@ def extract_and_load( else: logger.info("no previous progress for {} - {}", cls.table(), partition.name) - data_iter = data_extractor.extract_idr_data(progress, job_start, source) + data_iter = ( + data_extractor.extract_full_idr_data(source) + if cls.should_delete_missing() + else data_extractor.extract_idr_data(progress, job_start, source) + ) res = loader.load( data_iter, cls, diff --git a/apps/bfd-pipeline-idr/src/idr_pipeline/settings.py b/apps/bfd-pipeline-idr/src/idr_pipeline/settings.py index 8d68a49884..1d652712c0 100644 --- a/apps/bfd-pipeline-idr/src/idr_pipeline/settings.py +++ b/apps/bfd-pipeline-idr/src/idr_pipeline/settings.py @@ -27,7 +27,7 @@ def bfd_test_date() -> datetime | None: def enable_prior_auth_ingestion() -> bool: - return _parse_bool_default_false("IDR_ENABLE_PRIOR_AUTH") + return _parse_bool_default_true("IDR_ENABLE_PRIOR_AUTH") ENABLE_DATE_PARTITIONS = _parse_bool_default_true("IDR_ENABLE_DATE_PARTITIONS") @@ -113,7 +113,7 @@ def enable_prior_auth_ingestion() -> bool: IDR_ACCOUNT = getenv("IDR_ACCOUNT", "") IDR_WAREHOUSE = getenv("IDR_WAREHOUSE", "") IDR_DATABASE = getenv("IDR_DATABASE", "") -IDR_SCHEMA = getenv("IDR_SCHEMA", "") +IDR_EDP_DATABASE = getenv("IDR_EDP_DATABASE", "") or IDR_DATABASE # These need to be lazy-loaded since we override them in the tests diff --git a/apps/bfd-pipeline-idr/test/test_pipeline.py b/apps/bfd-pipeline-idr/test/test_pipeline.py index 9cf47a09a7..47a571b98d 100644 --- a/apps/bfd-pipeline-idr/test/test_pipeline.py +++ b/apps/bfd-pipeline-idr/test/test_pipeline.py @@ -26,7 +26,6 @@ from idr_pipeline.logger_config import configure_logger from idr_pipeline.model.base_model import LoadMode, Source from idr_pipeline.pydantic_utils import fields -from idr_pipeline.settings import enable_prior_auth_ingestion # ryuk throws a 500 or 404 error for some reason # seems to have issues with podman https://github.com/testcontainers/testcontainers-python/issues/753 @@ -82,16 +81,15 @@ def _do_test_pipeline(conn: Connection[DictRow], load_type: LoadType) -> None: rows = cur.fetchmany(1) assert rows[0]["bene_xref_efctv_sk"] == 353816021 - if enable_prior_auth_ingestion(): - cur = conn.execute("select * from idr.prior_auth order by mbi_num") - assert cur.rowcount == 21 - rows = cur.fetchmany(1) - assert rows[0]["mbi_num"] == "1OX4Y88RV68" + cur = conn.execute("select * from idr.prior_auth order by mbi_num") + assert cur.rowcount == 21 + rows = cur.fetchmany(1) + assert rows[0]["mbi_num"] == "1OX4Y88RV68" - cur = conn.execute("select * from idr.prior_auth_item order by mbi_num") - assert cur.rowcount == 64 - rows = cur.fetchmany(1) - assert rows[0]["mbi_num"] == "1OX4Y88RV68" + cur = conn.execute("select * from idr.prior_auth_item order by mbi_num") + assert cur.rowcount == 64 + rows = cur.fetchmany(1) + assert rows[0]["mbi_num"] == "1OX4Y88RV68" # Seed stale non-Part-D parent claims so the prune job has rows to delete. # CSVs cover item pruning because stale non-Part-D parents do not load. @@ -226,7 +224,7 @@ def _do_test_pipeline(conn: Connection[DictRow], load_type: LoadType) -> None: UPDATE {IDR_BENE_HISTORY_TABLE} SET bene_mbi_id = '1S000000000', idr_insrt_ts=%(timestamp)s, idr_updt_ts=%(timestamp)s WHERE bene_sk = 10464258 - """, + """, # type: ignore {"timestamp": datetime_now}, ) conn.commit() @@ -724,6 +722,74 @@ def _do_test_pipeline(conn: Connection[DictRow], load_type: LoadType) -> None: assert updated_ss_job.completion_time >= ss_clm_ts +def _do_test_prior_auth_update_and_delete(conn: Connection[DictRow], load_type: LoadType) -> None: + cur = conn.execute( + "select * from idr.prior_auth where mbi_num = '7ZM6HW2AT68' and utn = '-OTENCJLOQRAKA'" + ) + assert cur.rowcount == 1 + rows = cur.fetchmany(6) + assert rows[0]["mbi_num"] == "7ZM6HW2AT68" + original_updated_ts = rows[0]["bfd_updated_ts"] + original_name = rows[0]["name"] + + cur = conn.execute( + "select * from idr.prior_auth where mbi_num = '5OH0K85GU23' and utn = '-SC21YQR4UY4LI'" + ) + assert cur.rowcount == 1 + row = cur.fetchone() + assert row is not None + + prauc_table = sql.Identifier("cms_edp_view_cvm_prau_prd", "prauc") + conn.execute( + t""" + UPDATE {prauc_table:i} + SET name = 'BITE AID PHARMACY' + WHERE mbi_num = '7ZM6HW2AT68' + AND utn = '-OTENCJLOQRAKA' + """ + ) + + conn.execute( + t""" + DELETE FROM {prauc_table:i} + WHERE mbi_num = '5OH0K85GU23' + AND utn = '-SC21YQR4UY4LI' + """ + ) + conn.commit() + + _advance_time(datetime.now() + timedelta(days=1)) + run(Source.POSTGRES, LoadMode.SYNTHETIC, load_type) + + # verify that updated rows by upstream were updated + cur = conn.execute( + "select * from idr.prior_auth where mbi_num = '7ZM6HW2AT68' and utn = '-OTENCJLOQRAKA'" + ) + assert cur.rowcount == 1 + updated_row = cur.fetchone() + assert updated_row is not None + assert updated_row["name"] != original_name + assert updated_row["bfd_updated_ts"] > original_updated_ts + + # verify that deleted rows by upstream were deleted in header and item level for prior auth + cur = conn.execute( + "select * from idr.prior_auth where mbi_num = '5OH0K85GU23' and utn = '-SC21YQR4UY4LI'" + ) + assert cur.rowcount == 0 + + cur = conn.execute( + "select * from idr.prior_auth_item where mbi_num = '5OH0K85GU23' and utn = '-SC21YQR4UY4LI'" + ) + assert cur.rowcount == 0 + + # verify that untouched rows by upstream were not updated + cur = conn.execute( + "select * from idr.prior_auth where mbi_num = '7ZM6HW2AT68' and utn = '-RVUOWAUT5V5QZ'" + ) + rows = cur.fetchmany(2) + assert rows[0]["bfd_updated_ts"] < updated_row["bfd_updated_ts"] + + def _advance_time(timestamp: datetime) -> None: new_time = timestamp + timedelta(minutes=1) os.environ["BFD_TEST_DATE"] = new_time.isoformat() @@ -784,7 +850,6 @@ def _setup_pipeline_environment(info: psycopg.ConnectionInfo) -> None: os.environ["BFD_TEST_DATE"] = "2023-04-02" os.environ["IDR_PER_BATCH_MIN_CONNECTIONS"] = "1" os.environ["IDR_PER_BATCH_MAX_CONNECTIONS"] = "1" - os.environ["IDR_ENABLE_PRIOR_AUTH"] = "1" @pytest.fixture(scope="module") @@ -802,6 +867,7 @@ def _test_pipeline_load(postgres_db: tuple[PostgresContainer, str], load_type: L _reset_db(conn, sample_dir, postgres) _setup_pipeline_environment(conn.info) _do_test_pipeline(cast(Connection[DictRow], conn), load_type) + _do_test_prior_auth_update_and_delete(cast(Connection[DictRow], conn), load_type) logger.remove() diff --git a/apps/bfd-server-ng/src/test/java/gov/cms/bfd/server/ng/IntegrationTestConfiguration.java b/apps/bfd-server-ng/src/test/java/gov/cms/bfd/server/ng/IntegrationTestConfiguration.java index c488638353..11e2901942 100644 --- a/apps/bfd-server-ng/src/test/java/gov/cms/bfd/server/ng/IntegrationTestConfiguration.java +++ b/apps/bfd-server-ng/src/test/java/gov/cms/bfd/server/ng/IntegrationTestConfiguration.java @@ -132,7 +132,6 @@ private void runPython(PostgreSQLContainer container, Instant date, String... env.put("BFD_TEST_DATE", date.toString()); env.put("IDR_PER_BATCH_MIN_CONNECTIONS", "1"); env.put("IDR_PER_BATCH_MAX_CONNECTIONS", "1"); - env.put("IDR_ENABLE_PRIOR_AUTH", "1"); // Suppress noisy output unless the log level is explicitly set env.putIfAbsent("IDR_LOG_LEVEL", "WARNING"); diff --git a/ops/services/01-config/values/ephemeral.yaml b/ops/services/01-config/values/ephemeral.yaml index abdf6318ca..50d8143918 100644 --- a/ops/services/01-config/values/ephemeral.yaml +++ b/ops/services/01-config/values/ephemeral.yaml @@ -20,7 +20,7 @@ copy: - /bfd/${env}/idr-pipeline/sensitive/idr_account - /bfd/${env}/idr-pipeline/sensitive/idr_warehouse - /bfd/${env}/idr-pipeline/sensitive/idr_database - - /bfd/${env}/idr-pipeline/sensitive/idr_schema + - /bfd/${env}/idr-pipeline/sensitive/idr_edp_database - /bfd/${env}/idr-pipeline/sensitive/idr_private_key - /bfd/${env}/idr-pipeline/sensitive/synthetic_env_username - /bfd/${env}/idr-pipeline/sensitive/synthetic_env_account diff --git a/ops/services/01-config/values/prod.sopsw.yaml b/ops/services/01-config/values/prod.sopsw.yaml index 762a958ff4..1f61d5d589 100644 --- a/ops/services/01-config/values/prod.sopsw.yaml +++ b/ops/services/01-config/values/prod.sopsw.yaml @@ -36,8 +36,8 @@ /bfd/${env}/idr-pipeline/sensitive/idr_username: ENC[AES256_GCM,data:3+gL6zDI0hU+2XeaPhXe2Q==,iv:GgqF3nzKu8+500JqBhZsyCHDrrULlbpcZwDNFl6+A9g=,tag:p+kP4MvyjEbVIISSfTwG6A==,type:str] /bfd/${env}/idr-pipeline/sensitive/idr_account: ENC[AES256_GCM,data:4H6iwPfyiLJf7W4COWrJfVA9og==,iv:J7NvBn2Wf3YUDmGLVJhSNObyRQVfxdgcGC26XhpBnMM=,tag:1wf3r9XhfxAWkPDdIpFEwQ==,type:str] /bfd/${env}/idr-pipeline/sensitive/idr_database: ENC[AES256_GCM,data:gIhp6mUS2I4=,iv:pv079MXQlvrV1Mqgi7/IMAVP+grP/2JrsZlH4/EAM6g=,tag:AQKFbGjcubAHtGY+rnTfsA==,type:str] +/bfd/${env}/idr-pipeline/sensitive/idr_edp_database: ENC[AES256_GCM,data:xkf0u8/iWNV4Nfh3,iv:hwEP1kjB9CiabIZte3ZUSA9rPb38GUFm0wC6jjmT0Rw=,tag:VIBGYKkO1ROEgUiIJ6sqpg==,type:str] /bfd/${env}/idr-pipeline/sensitive/idr_warehouse: ENC[AES256_GCM,data:fsLtGL6QxPg4gxWZ,iv:gQsj0XqsFJmMDFEqLDoXe2fel5q/A6TMIgrlluwrV4c=,tag:+XxJMhD783hjWW7DKXDk+Q==,type:str] -/bfd/${env}/idr-pipeline/sensitive/idr_schema: ENC[AES256_GCM,data:f3/iwiAAhKV8FJ6xrx+CnAlDIrT5,iv:6GpQvT8Ws05Rrvq6ON6Qzcr+dTIbZNLxRTFFE/EHIps=,tag:rrhG+2C/qryS5q4Fjo+tAg==,type:str] /bfd/${env}/idr-pipeline/sensitive/idr_private_key: ENC[AES256_GCM,data:EHE2/UTtMoxrIYZHSaQO3+QRpkxc0crMRtxeOjjk4BD6m/ddWVmlOih1stnSmoHaOoSuZWVa196Y8sKyStqmsx/UjLM/JpdIQM68yxnp/iXiWF4+Q8V3FSajqov7yq/2uC5i+mD1UH6BZ9ywBM5hAko9L7lbQcIvSkvzG0S/ZWDOna0Y0BFb3d+il/3l45KxtWq6V8XbcGBFIHNSd/wbz5vbPxGZZG/8L84MwdQG8vzE7VZHyE4sxHCmYCVvdGg2T/K4OlU94u7cDoLuHueavVjocingNIeyrmUnyRLyBpeDo2r2qFgQ3ZaKBzV7T53Fsfh4UTP7btNS+EtczHQRiyjaTa2+G8qythmplpb8eKX4z5cfGJOl+KzK1eFcXJMQoYJpxUJKq5Oiy6kBYbH9bdr6pQtICvQbtJzFUq9wi3BGPYyOkk+k8UIyo12cZDyP73e1LNgDMsxhBE99i06CjJldgAx6FYcgu8fHZR8PsqqWwj+BLwhwKCxa1r6vsPXWO3mwsWt5dvUoxhTQbKCLLb9vMMjeiHKUvSiVVo5Ki4Lc0h3yuVIINTcdi5MfLOlj9W+mwI85Le+eEhdGjJqXeMlG1jzFyJAhHObnXUC4doiwfu5BUztVqGgv5Mh8kGSOlLXVevttgdP/Mjsvq9eqNwXi5OE7D8BX//JE51vHVpd+7J9FdGABw/BdMbgbpIv72UXLOAeyc6DHvhEpkgObrTHmGfNS8SBpZE/pkKDM2FuDXQXaOZpWxcBFC0Db4gzbKNrU1okFjrXuXJGYbu6X73CieFnZRYD9/Ilgq+sFy6NXu2m2lsrRfI3mlT7LltK1f2qzPvSWXlAmR22OQYzaSwYXxInHXh/q2dpx8FZHK/9ngThDnqlrYX5YW2zMsBsoFbwBM4E4ELPvCprnaVPvoyfp2HDts95lS8OuX4z+MG+eWFrvlqeXMRekl+522O0imt+OUeiRKWSxcgV0jKPDK5UoBGjlh6FxWhSlzt1+Z6c1wUoewEvN8ydHhBgpbccK+slyNp3dSgjDUqV/6Gw0LxHCVvaLknNyhNKVXYc6b766F0jtFxP3G2zrnq91o3obAujikpAawU5qCHYCXq5bwtoscoVMGqHcrg1mhe0MUpLX2hwqgWu0sV7gF7mCwBGv111R/MrMraLDtyOsrIa2+SeCCfOB/XeMKsE5270NQoriC2D1kOQTD+evRYBZCBTeFIKz/VemSYVbJvu5QS+m50LCIklpzl/K9LDvkqG5oANztG6gJ9lN/aEhLaWOVLi2aRNcwUatTeKkzhm4iQmVtTRcmjKFj6tkxfbep9fe01/6P3GIQtVdTcSWlsAaYdHbYFv95QkCQqK4dM3wVv1L2VAgwKz8yYcm8JCRP8TGe9VJNLUh0/yDOfzkN7kQO0D+EfnwKkA7ctRU+erGlaBq9wA2xxLAhvkdbI4ARgZqYbFmG7e/JxZvLq1J+iSDPF/TeaqSx/LFy1IU12wFo8zdnzQ9ePwaWCbOoRklLfFnjSLbI597GyOmPckY4tj23CeDD3YeOHoBh5dbMElDC1VxPR33pGbIFVwSYqzLedyJOjstWkzKSJECdIRa7phfk1o04wfFtj/UyqGzYvfNGwdD5PltKpeKzuFHLp7yf8oVwFIof+1K7AysNK9mqN+v/n13y5VZWDmhd9KZ6R4z8UXrcVmhVVZmuVjPqH0fsegS6MX7qfn025DChLFXP9YbyqmtSC4kuXq2YPu6Xy0z5LXcunoLLJ0Z6A9kYAgk/VFym036MfFU03ErNW2UsLiYXOHwcQwKebst67uGr7RF4azmDoFRqDBGtzmG2bglTx1FFS1B/gXXY1R+hbFb5Z2Ew9RH7mBzNj/P+qVDfrNS2nOihj7PT5/7WWMc/dYskdb+1deD+eD0sXbjXEWVrKp14bBzccTnvgLXd8fuPaGxGnNVsl2Kg49YPziGt/xBGEhQhgj9NS9ljzhzOmi9vqKmLS3LmbdFPw7+NXOAAKuH6ugyw77/XKlwtPDHYZ8JW+GjMjht7Ea1D0iRuRC8kJbywKcmlU7E32pRewtmnUV86sWz8c7KlrbE+AII1llAdPGVvZoS7MYj5Hfrv2cvWryJYESyPt0vBEOW4GpQKJWsY2wgSjmGbKfJ2WMsN2+I2RXa3HK+0sQmi4YUg35LMmCDBW1IyqGDFZGe3OXF9L7oO/CO2xN1Ol3BzAmU7Diib40rUuIecRHGudDrRveAn/D/C+BFHb0XkvXJ64nQWHy23ZXGMMn5HRPz6aM=,iv:bzDTkn40uAVkDjmo7ZpTYWBI6oKnOiiFH/fz30IGJcQ=,tag:nWtcMVlmvN8eNoEFDXebhw==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_username: ENC[AES256_GCM,data:Q9ZAGflOrN/LIrTc68Fx4SM=,iv:lhhjzq08tsGnGUm2+ffxToB/2j+U9gKoRmXqepY/Itc=,tag:iyWslH1zh1zf9t7xBaz5aQ==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_account: ENC[AES256_GCM,data:OLc5uFN4cA==,iv:N8FYKhkL/1lqK7XsTQwPpQ+om7gRiFshRqF9ixF9cYU=,tag:KTJJGT+Bdx1vKSUWCP9v6g==,type:str] diff --git a/ops/services/01-config/values/sandbox.sopsw.yaml b/ops/services/01-config/values/sandbox.sopsw.yaml index b1746425b0..a145cb1dc7 100644 --- a/ops/services/01-config/values/sandbox.sopsw.yaml +++ b/ops/services/01-config/values/sandbox.sopsw.yaml @@ -38,7 +38,6 @@ /bfd/${env}/idr-pipeline/sensitive/synthetic_env_account: ENC[AES256_GCM,data:pzUAvD89DQ==,iv:TfaEPTKk8leprcEXIbyVse/6j/JOi/MAtS4MI3bG9Ag=,tag:WAOeHgImuLoaoIufeePteA==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_database: ENC[AES256_GCM,data:O1XCauuSqUOdy2M=,iv:UKkTybopum/fRDmwNMrqL7Y09tM6O4GclpfHFGUzJFU=,tag:qG2FHUMlf7SM8h5UnF+EvA==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_warehouse: ENC[AES256_GCM,data:Q3NZ,iv:EhFDOBM1d5xTPsrcHvwPul4I5fc/09RC6vo9iqP+nw8=,tag:KXGYd0GUTtUDOVYpqNQmVA==,type:str] -/bfd/${env}/idr-pipeline/sensitive/synthetic_env_schema: ENC[AES256_GCM,data:ZGI3f4Z0EuH+nsCfPmf+sBST+mE8,iv:djGgV5qQZLQmslgakkLdxn/l1m1Dz1KYewmyokOito0=,tag:vhUfro1yNDHNKHcMRseb5w==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_private_key: ENC[AES256_GCM,data:Ew3Mp7+DPNyG8ngR1u1w2QzTywrFfpqqkr1HRMAfRoaGtsyE2YdCNADNcpzoAoW20Sq/4bBN224YVY4FEyia7yvmVfTJN2ICjKWBH6irQ5XeoddfjSxXWxcOfCDsAFE7C4hc8THQJ1lI8pK3QKcCIsNXd9K82FRWxHzLATbSLciQhMhqFJ0AK9NLK1LVMxctvhd6q9hvxlZC3hwp9OrfeLeS78MGgvGg31E0RzC5SCib0s9s5tZPg/hAawoZxuwMguXAMrLyu0EfizRQa6znycQTjlsM76OzqC4QrPCuM+NzquF12qAvmdYXByISejrvA1FnIgdg8AIPndoUj0gLjrzT+74nU7x18Rjer70IIaxM+39eXe7Mv3J2kHlxoJM6Ov240Hpj237KTWksBWHvIBh3Aqej2ugHnBZ+nhktO9FwPJfGZzmcT3w8ZpGy1d6WAB73YBP23jkjQE6EVRykvs9ri9CswImWUxGTg0DqzX65LkweUsUioAjaFyaLVrL85Z4/RifZq3jhsfdgOJy6QdL34ONWaoHf7pXqyfQt+41mIhB+oKdgrmAQbZkz31cHk6UzXnh3oO1ZXsHqdvxxZV8q0PtJjZ5ji7G/tZYp2iUIBYTJ3GkpDMiEq+vQ0+tM/HH9hWK99lFKHpB3cOKcnSZ76tgzr0nPszgTES4k0KcLRxrtpBrEnZqwNlUcZtzLmo5rQP6CnmKwPg7MIQkMZOJtbQcf4uzYt+5wD+y/fYx40vvrje2z9QfcpItWYHZHyPfIZtWMclR2UtkKGqhANn3kHmRWk2qMMI+mmpmXFP6Vt/DaFn3MVX2HgRhS7rpS8v+gm+qLSz2EfNRB52vgztzbWKbQVDOs5NhS2LpSX1EaHyeXWQj8rDT0lED8YxDaklmQ8deEwmcBNQUZ7v/eyw+CtmhahwkiVJhsTn5LOVVVbwZyEtTYaZZJb50gOS8d9q12akSei8bcgzh8wN22ySH9KUgUBUEqGXYGvR4D8keP09cfLC2JpIfpWCBQ/51Aoh4lwujqDQC9pvbB0QE59FWufI25WzmXbXZGF/4nyeWi6IgoS7DFyxQR1lWjDF6tkMAVDYS83KsX5IDEhOfUj6WQTl1cu4pgj7kkDocWgNmBFplQXf77oNMoIe3kBT+AdZ7znquKa8VeoT7kBfVJq6mWMxAk2B9sAQgwi7+xaGzUS/iw/fUVnoUJBomyEz1ePEpUtsYkYXgBt9yhRlLjDMbASyPr4kgSKOyPuKWuNQBqKE9EfoYKwENBDxn5GyRu7+uRm7dLw8NpY0JMeKc1ofcmxPnni9NIo53pM7YgTAaFmxW2yVy9+94LYHFd4nkdvQYDSTAMtarGakuwcNOWzzkMjDCThYl/q1HTR3xoh3suXr24tgGkQIT52AvDsIRwEykJGFK8jliM2HSapJsX1gsW4IrAAWJTl+tZ45lxFchp8N21fxTA+RnQT/YyTo+ghpv3CJ+tJl4RetobMhAxbSw9JYludPUC8UQvCxSxcV0OdR91ESXs2go29xHlcMCA7fb2IhUMnFcfxA2ntFMm2EuJzXLUabFoBql0ufbLmcok0JHAYsw0Us6A0RO9M/7MWCkO8FBY9ycK48fwyKBar0YSwdFEwha5m7U4vqRjfgVSWXuW6dWp1JhBG2+/u9DwYtCN6HCMbSlSNynlhCLmR/Dhs8TFFLlMJ9jEBVn1PyojXs9+deGZOXGqqb44ELIOqlwRg2ur5cDmFxJvlBLpbLptsenH1K1a/QXtESBmJiMtJRKzxH2gnPuXH2+arAu6Mccq+/zO+hgz/JUQQauQSDmWXJ4nzwTXk9sDkv9PdbNFfiav08q4CVwdH2FERenKy/4J6qd9BNpPDZSsah3dmZsRxLXTVVLFweimWlm++O+4pBwlvjGYxZ5v620oc94/MCZfcMqm6/uwlG69Hja1ws4NKsphfneRRn0WWWoyiKuicwchVUH4KjNqw6XSm8Fs3KG2Dj5YAnvzPy/t8muBrAt1Y4xLbkWeCm4pYBtaZWw1D2phVYObOww9Y+DijApeRAvwuAkuCrzWZWX/pPsFEO+ky3sjySlWspKMyp9QoPK6Ra7KmkGiPml+D0wS0pMYR+oT7TPtkF+FLZ1WOZ/XWEjRnGL80Lx1MWwNd4XAWGs1hAk1FoaQ0+okOTRp1nvy8cGbvQ8IMSwJQKsRFH/jNS5Uk7wbD+djY2CYlu5Qt4/UWYtB/jCjZOuUVx85YVw+ezDf3Yv0HOpfztdns31793T44vzm9r75,iv:mto+UK1H6tKWKajPJRSaOFZ477WlNXTHJRzjI8xh2tI=,tag:aRJnXhOFDX4/uUSP2bm8xg==,type:str] /bfd/${env}/idr-pipeline/nonsensitive/synthetic_env_public_key: | -----BEGIN PUBLIC KEY----- @@ -1575,20 +1574,20 @@ sops: kms: - arn: arn:aws:kms:us-east-1:${ACCOUNT_ID}:alias/bfd-prod-sbx-config-cmk + aws_profile: "" created_at: "2025-05-29T19:05:53Z" enc: AQICAHjP/Aj2Cj3mc/evGqWprczmF5nsbcARIepeXMoZGoJjRgFLoS1HtIVJis6KCi6IGgYtAAAAfjB8BgkqhkiG9w0BBwagbzBtAgEAMGgGCSqGSIb3DQEHATAeBglghkgBZQMEAS4wEQQMYXpFYeIHnLME7ja5AgEQgDsaaGnkO4KCbl3PxuQ/6xlsl42bBPosBcbDSZCjuEw5FiozdEJPIQKi2GnHBR3a4PJAgTMxZddz57uOWw== - aws_profile: "" - arn: arn:aws:kms:us-east-1:${ACCOUNT_ID}:alias/bfd-prod-sbx-cmk + aws_profile: "" created_at: "2025-05-29T19:05:53Z" enc: AQICAHhHJZYi+MFRKFTHfYVWYS6PhbzwLNq4isexRIBWUW6UVAHYOoM3Wmc9vHbogkQNmI5rAAAAfjB8BgkqhkiG9w0BBwagbzBtAgEAMGgGCSqGSIb3DQEHATAeBglghkgBZQMEAS4wEQQMvR9IxBZxsccEiPA2AgEQgDtS2TovK1kZfuk9NlVbJrTZqESM7m1Is19bKNTK5gfUMCyor8I3+MlnS5wsNnLBVvyA0j1n7nzQU6N3KA== - aws_profile: "" - arn: arn:aws:kms:us-east-1:${ACCOUNT_ID}:alias/bfd-sandbox-cmk + aws_profile: "" created_at: "2025-06-16T19:15:24Z" enc: AQICAHgeb6WVlAiQlv36YrB0vjY473KhX7b/JfaXaUcCtUQFkQEoS26161rCnQXywCnVnsASAAAAfjB8BgkqhkiG9w0BBwagbzBtAgEAMGgGCSqGSIb3DQEHATAeBglghkgBZQMEAS4wEQQMJg1DN2R+a6HG5UeZAgEQgDtHqLiAgHHpn07reZ3rEcUBwW6lMquhLGljV7LqOTlSWNakoMsmeFPOpxBtwTal8ceyfl7ufkWp8IUjng== - aws_profile: "" - arn: arn:aws:kms:us-west-2:${ACCOUNT_ID}:alias/bfd-sandbox-cmk + aws_profile: "" created_at: "2025-06-16T19:15:34Z" enc: AQICAHhJyyTYz7F9ZtNNOqGQey2/S5rpRtPB96w5RrdTRf7ycQFKTOb7uLgtYqFd+lCmkLF3AAAAfjB8BgkqhkiG9w0BBwagbzBtAgEAMGgGCSqGSIb3DQEHATAeBglghkgBZQMEAS4wEQQM2hGiKtGbX/DxiwP0AgEQgDvBGTkDcOudxsok1MRLWKiTtoPF+gvgaeTcTHaVK4hX2w0zxRbM3/llsarbb7f1u+7fAEPFZftQd+EUhA== - aws_profile: "" unencrypted_regex: /nonsensitive/ - version: 3.12.2 + version: 3.13.3 diff --git a/ops/services/01-config/values/test.sopsw.yaml b/ops/services/01-config/values/test.sopsw.yaml index b19de14c80..baf0efcf1e 100644 --- a/ops/services/01-config/values/test.sopsw.yaml +++ b/ops/services/01-config/values/test.sopsw.yaml @@ -38,7 +38,6 @@ /bfd/${env}/idr-pipeline/sensitive/synthetic_env_account: ENC[AES256_GCM,data:rGECSWzWJg==,iv:qu5zLCQ8YQ/Bd5AtFdyaCfU3LsV6lQs/nurWS1LA+jU=,tag:Xpfqv4M/HvFCFP7UNyUqVg==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_database: ENC[AES256_GCM,data:CTuZTyPPbgA=,iv:e0pv8lXThf+OEVf9suomyp2n2AAeLC6jl1MEtA3WlTQ=,tag:8KHmuuzAfFqcCtyrL/oHZQ==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_warehouse: ENC[AES256_GCM,data:20Q7,iv:0pwPtvZIPsGOwcevtPDKWFm431E4RuXNxtVwBmQbXqk=,tag:pgucb/5OoYxY9ZU7Fzkf4w==,type:str] -/bfd/${env}/idr-pipeline/sensitive/synthetic_env_schema: ENC[AES256_GCM,data:inSPo/CiKxXwLEf2TJJAuQWWA13q,iv:FkRtIjOU+Bbtwxxpul46809RCFgQo5rj9gjF58GssyY=,tag:ppquyCkHs8qo6rYDmXL5ZQ==,type:str] /bfd/${env}/idr-pipeline/sensitive/synthetic_env_private_key: ENC[AES256_GCM,data:tD59qC4VmlsmnYGvtGxgRilrVPxMLVxUbpP/D2S6mSunTw3Sy8d509DTRQNv1aVkdXF56ik/vKroTp7awPc4wcd+rKmfwRSrJ+LF5EpB27HbJ/sbWYFYOr6WNvxNI8KE93eN66GNkAg60POnG1QXlozZh5ScZZ9ZfmwIS2JYJ68loqm3DOn5HjD1/+5Tg0Ja0uPagnjRYtvL4v42zEAmq8q7ypb06tR5wVGcRYoemP04zWVBygVugMvTagj6DGT1SLUEIhp9YRuNc/X0U8gNKTVMP/gITt4IRWl/yLeCIfTkaJ+g0y4XHkq/hQwRMnmeF1dJoi4l/Hp7PuXCxN1y5n7pYtSsoPcMG8xLxH7gH5mywDyLQFW5m2xnW/lRE9s1sS5iLQxNPOYipkq/9UaNuH2RCnzh+35iXHW7hNGq8v5GxAkqrFZi/LF27SQEu/k9ozd+O/1vPo6C6eIVUl8uPkXBug2SgHTk1rHf8WqDMjE03+ggE9IcXkSppDE/kAlq/Oaa/DuccCDxXmtCx5fjpd4jXntQxWbXwXerAc7Y5wXhnkeMzm1gkLGBymEcvjXEW32whzufdgqoceFbtoTwU7qg7ZMA7Mk29U//b9YOb+NS7+Fy7h3Ev3FkTd1rUVD8ultJnuRbrc7DgjQ804+o5QZjJfReznu0xXPCKSiKvR9/r2RB9mnWIqfpt1RfL3towDE9OLoF/IIa+KWtDdIL+sNlVh9ySXPNeWEwWjKCxoki4dLrGCIJPCyZriZKKLCya5ucJX1HVVDHnPQKM520UK5nPMFa2enH/8s0KX5AaL6z8bU077pDfddROBax4/oMsfQ0gl8xr3apHr6VTrUXDNEIs7gGpP6CJRgm4nvnsRd7eMOuL/cwNVQGD/OvOgAs34mtAWAgtSDjIRIamUjDZ1igDBzc3FRIwLreaja1clNIsn9fb2aYigq282jAnBLWxGtMlJqUHVBd2dNTd9ZFY8Y7iQdhVVfNpEzP2dRupwlicjUhhnktK7IBkqRr6jBlx3OeKguy0dnWPWsSpoeSyGBlUUdoC4LipeYuC1XOmt13WoAUUwzfSilwLvxmB/J9tiIsAvmXchmaJGzG5lDhjXtbmm5Rlv3qHUDBOgtq8j3FWG3T8oqY0hbroodOETWEELM65w4Ol4OGrHaxai71USFCfWrtGWiOiM2A7zOsr91cKfbwNH2Ex2f9iiFCKgz5ORwspNL5emp4GgFd1ESna0LGl6vxw+/kYdzgF/TE68iVsG63qClXOGphCLguZdSdgfgVNVPw3AoJNwkqMj5IIrOQh8opY8FQedS73OqWcB93YPAOr3yhSJvkFzjn7mT32MypnExvY2DJs2yjye0KEjxUhi0i6mkGT6IJls7PqBEDQtyyjsVsTQicwT18uG++SaXHxVkh0aFcEFRGYazg3F+/QKHN5O8oU/wDPZi43WsBV90EjeYeTeita8p3QIO+Ua3/aBn8gW2v6w8OOsq2lqzN2WyRiJYENodcDrgmy4PxKrAP/B/c2zPdhANgksUK6wSIR9I2Ipn9RWmubtpsUjkBU7dcRAFlkEGUQVYOFohHMBc9XAGwZ8f0ZvIrqMfICPgev4QoO5PYd0HPoKLTy7Y5iprVKH3KuJNNvD8yQjoQO3ehEZ16kihkzLdi8x8/1zzCvPZ8b2Dg0fxkjhyaKsY4lBI+gLnw5m+m1drGU/up7RB1/LfJo8LD7Aj7pRfMgfSmWct6O0u4lUeJz1BL4wFwnnaB3qWx77rbljAz5uKvZ2EK3SKRvUtDHxR4OMuGamIC+1gRSleSGi2ye4iA4Zv+8dDs+Dg+dZePQxG5hqc7DbGbYXESsA1xvgwjbCeGK/5s93c+YiqVkfBkgy6WISaLfHrRkD2ca4vk72M80kUl10NkCdcr3J73YtA1q9+1JpH82OiQRYwiJJd8gKqroBwOlG5fGrdtV7zQruN76C5FMZGD7BreVmEeaY/y0Ft5qNHUS4iWkK8s4e65nT29LVv4KLPJYHlCIAXcQjq1IOY2OyMaiFQqpaVtA53lqe3rRtpGISAe4Cxstd8iUnsoEvk+exKBvTIrMr/Sx5BAOshCpsNVB68orjxq5zKnrY9GA60rJ6iVAP034M55wRRd8O/8sh7uMR6XWNLXCJlwDFHdvFoIc/5cValBdCS8T7/L8903QC6TXaQMoG0Opa9NWW8r1DEtKgO1446RtJ9tfvfMaJyisR8YGxRuMQ/D0kvlrm9rz/s26jQRlFz98XZXOBurF5WvDSHa,iv:MhTKD/JvFe2zIGEBYPyyVK4igyzXzTRsCAqbS2NyWjY=,tag:HGWAYk13ViO7IgpMZivb0Q==,type:str] /bfd/${env}/idr-pipeline/nonsensitive/synthetic_env_public_key: | -----BEGIN PUBLIC KEY----- diff --git a/ops/services/04-idr-bfd-validator/main.tf b/ops/services/04-idr-bfd-validator/main.tf index ec988b7656..1208c29423 100644 --- a/ops/services/04-idr-bfd-validator/main.tf +++ b/ops/services/04-idr-bfd-validator/main.tf @@ -49,14 +49,14 @@ locals { disk_size = 21 task_ssm = { for k, v in { - IDR_USERNAME = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_username" - IDR_PRIVATE_KEY = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_private_key" - IDR_ACCOUNT = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_account" - IDR_WAREHOUSE = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_warehouse" - IDR_DATABASE = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_database" - IDR_SCHEMA = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_schema" - BFD_DB_USERNAME = "/bfd/${local.env}/${local.pipeline_service}/sensitive/db/username" - BFD_DB_PASSWORD = "/bfd/${local.env}/${local.pipeline_service}/sensitive/db/password" + IDR_USERNAME = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_username" + IDR_PRIVATE_KEY = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_private_key" + IDR_ACCOUNT = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_account" + IDR_WAREHOUSE = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_warehouse" + IDR_DATABASE = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_database" + IDR_EDP_DATABASE = "/bfd/${local.env}/${local.pipeline_service}/sensitive/idr_edp_database" + BFD_DB_USERNAME = "/bfd/${local.env}/${local.pipeline_service}/sensitive/db/username" + BFD_DB_PASSWORD = "/bfd/${local.env}/${local.pipeline_service}/sensitive/db/password" } : k => "arn:aws:ssm:${local.region}:${local.account_id}:parameter/${trim(v, "/")}" } task_tmp_dir = "/app/.tmp" diff --git a/ops/services/04-idr-pipeline/main.tf b/ops/services/04-idr-pipeline/main.tf index 860687a0ce..b55d54b36b 100644 --- a/ops/services/04-idr-pipeline/main.tf +++ b/ops/services/04-idr-pipeline/main.tf @@ -56,14 +56,14 @@ locals { idr_capacity_provider_strategies = module.data_strategies.strategies idr_task_ssm = { for k, v in { - IDR_USERNAME = "/bfd/${local.env}/${local.service}/sensitive/idr_username" - IDR_PRIVATE_KEY = "/bfd/${local.env}/${local.service}/sensitive/idr_private_key" - IDR_ACCOUNT = "/bfd/${local.env}/${local.service}/sensitive/idr_account" - IDR_WAREHOUSE = "/bfd/${local.env}/${local.service}/sensitive/idr_warehouse" - IDR_DATABASE = "/bfd/${local.env}/${local.service}/sensitive/idr_database" - IDR_SCHEMA = "/bfd/${local.env}/${local.service}/sensitive/idr_schema" - BFD_DB_USERNAME = "/bfd/${local.env}/${local.service}/sensitive/db/username" - BFD_DB_PASSWORD = "/bfd/${local.env}/${local.service}/sensitive/db/password" + IDR_USERNAME = "/bfd/${local.env}/${local.service}/sensitive/idr_username" + IDR_PRIVATE_KEY = "/bfd/${local.env}/${local.service}/sensitive/idr_private_key" + IDR_ACCOUNT = "/bfd/${local.env}/${local.service}/sensitive/idr_account" + IDR_WAREHOUSE = "/bfd/${local.env}/${local.service}/sensitive/idr_warehouse" + IDR_DATABASE = "/bfd/${local.env}/${local.service}/sensitive/idr_database" + IDR_EDP_DATABASE = "/bfd/${local.env}/${local.service}/sensitive/idr_edp_database" + BFD_DB_USERNAME = "/bfd/${local.env}/${local.service}/sensitive/db/username" + BFD_DB_PASSWORD = "/bfd/${local.env}/${local.service}/sensitive/db/password" } : k => "arn:aws:ssm:${local.region}:${local.account_id}:parameter/${trim(v, "/")}" } idr_task_tmp_dir = "/app/.tmp"