Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
2841b6d
Reorganize IDR Pipeline into packaged app
malessi Jul 28, 2026
a1299e6
Remove README field; update description
malessi Jul 28, 2026
5942fc5
Fix docker image build
malessi Jul 28, 2026
b6f60aa
Replace outdated ways of running pipeline
malessi Jul 28, 2026
a8c41a8
Fix running pipeline post packaged app reorg
malessi Jul 28, 2026
843603c
Add py.typed markers indicating that all modules are typed
malessi Jul 28, 2026
b7338b7
Move idr_bfd_compare.py to its own utility project; rename idr-pipeli…
malessi Jul 28, 2026
b8cb5a7
Use extend-select to keep default lint rules for ruff; fix lint
malessi Jul 28, 2026
fc92afa
Regenerate uv.lock with exclude-newer
malessi Jul 28, 2026
2ce8666
Add an exclude env var for excluding specific tables
malessi Jul 28, 2026
a88207b
Rename symlink to idr-pipeline to match package name
malessi Jul 28, 2026
7371793
Make docker builds work
malessi Jul 28, 2026
6ecab7e
Update README
malessi Jul 28, 2026
f1b922f
Support named contexts when building images in GHA via "namedContexts"
malessi Jul 28, 2026
53568d5
Build the idr-bfd-validator Image; add ECR repo
malessi Jul 28, 2026
bd6bc6e
Implement topic alerting
malessi Jul 28, 2026
9798d8f
Init 04-idr-bfd-validator
malessi Jul 28, 2026
72ca0f8
Grant SNS permissions to Task
malessi Jul 28, 2026
3957ffc
Only create schedules in prod
malessi Jul 28, 2026
734eb53
Re-generate uv.lock post package rename
malessi Jul 28, 2026
9cedb2b
Fix integration tests requiring the IDR Pipeline
malessi Jul 29, 2026
2043090
Use toJSON to ensure JSON-stringified value is stored in environment
malessi Jul 29, 2026
165d067
Update idr-bfd-validator README
malessi Jul 29, 2026
b0c7add
Remove unnecessary SSM parameter IAM policies
malessi Jul 29, 2026
5d35b8f
Fix subnets
malessi Jul 29, 2026
6507f3f
Exclude prior auth and low income subsidy for now
malessi Jul 29, 2026
219728b
Fix SNS Topic publish due to lacking sufficient KMS permissions
malessi Jul 29, 2026
43d73af
Include idr-bfd-validator in v3 release and deploy
malessi Jul 29, 2026
a164bb8
Merge branch 'master' into alessio/BFD-4576__autorun-idr-validator
malessi Aug 4, 2026
e9a67a5
Merge branch 'master' into alessio/BFD-4576__autorun-idr-validator
malessi Aug 5, 2026
0db84c8
Run terraform-docs
malessi Aug 5, 2026
6960d7f
Exclude idr.prior_auth_item from validation for now
malessi Aug 5, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .github/workflows/build-container-images.yml
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,8 @@ jobs:
- name: Generate build-contexts
if: ${{ !inputs.skipBasePull }}
id: gen-build-contexts
env:
NAMED_CONTEXTS: ${{ toJSON(matrix.image.namedContexts) || '{}' }}
run: |
build_contexts=()

Expand All @@ -287,6 +289,15 @@ jobs:
build_contexts+=("${image}:latest=docker-image://localhost:5000/${image}:${BASE_IMAGES_VERSION}")
done

readarray -t named_contexts < <(echo "$NAMED_CONTEXTS" | jq -rc '. // {} | to_entries[]')
for named_context in "${named_contexts[@]}"
do
context_name="$(jq -rc '.key' <<<"$named_context")"
context_path="$(jq -rc '.value' <<<"$named_context")"

build_contexts+=("${context_name}=${context_path}")
done

build_contexts_json="$(jq -c -n '$ARGS.positional' --args "${build_contexts[@]}")"

echo "build-contexts=$build_contexts_json" >> "$GITHUB_OUTPUT"
Expand Down
3 changes: 2 additions & 1 deletion .github/workflows/build-release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -452,7 +452,8 @@ jobs:
bfd-platform-pipeline-ccw-runner,
bfd-platform-consume-idr-events,
bfd-platform-run-idr-pipeline,
bfd-platform-codebuild-runner
bfd-platform-codebuild-runner,
bfd-platform-idr-bfd-validator
baseImagesVersion: ${{ needs.compute-version-strings.outputs.bfd_release }}
cleanupImageArtifacts: false # We'll cleanup at the end of build-release, so don't do anything
tagLatest: ${{ !contains(needs.compute-version-strings.outputs.bfd_release, '-') }}
Expand Down
9 changes: 9 additions & 0 deletions .github/workflows/build_container_images_matrix.json
Original file line number Diff line number Diff line change
Expand Up @@ -82,5 +82,14 @@
"dockerfile": "ops/images/bfd-platform-codebuild-runner/Dockerfile",
"contextDir": "ops/images/bfd-platform-codebuild-runner",
"platform": "linux/arm64"
},
{
"name": "bfd-platform-idr-bfd-validator",
"dockerfile": "apps/utils/idr-bfd-validator/Dockerfile",
"contextDir": "apps/utils/idr-bfd-validator",
"namedContexts": {
"idr-pipeline": "apps/bfd-pipeline-idr"
},
"platform": "linux/arm64"
}
]
1 change: 1 addition & 0 deletions .github/workflows/v3-release-deploy.yml
Original file line number Diff line number Diff line change
Expand Up @@ -385,6 +385,7 @@ jobs:
database,
migrator-ng,
idr-pipeline,
idr-bfd-validator,
idr-pipeline-metrics,
idr-pipeline-alarms,
server-ng,
Expand Down
1 change: 1 addition & 0 deletions apps/bfd-pipeline-idr/.dockerignore
2 changes: 1 addition & 1 deletion apps/bfd-pipeline-idr/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -62,4 +62,4 @@ ENV PATH="/app/.venv/bin:$PATH"
ENV PYTHONDONTWRITEBYTECODE=1

# Run the pipeline in IDR mode by default
CMD ["uv", "run", "--no-sync", "pipeline.py"]
CMD ["uv", "run", "--no-sync", "idr-pipeline"]
6 changes: 3 additions & 3 deletions apps/bfd-pipeline-idr/Dockerfile.dockerignore
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,6 @@
!pyproject.toml
!uv.lock
!.python-version
!model/*.py
!matching/*.py
!*.py
!src/idr_pipeline/model/*.py
!src/idr_pipeline/matching/*.py
!src/idr_pipeline/*.py
2 changes: 1 addition & 1 deletion apps/bfd-pipeline-idr/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ source ./load-credentials.sh
Run the app (optionally specify a minimum transaction date)

```sh
PIPELINE_MIN_TRANSACTION_DATE=2024-01-01 uv run pipeline.py
PIPELINE_MIN_TRANSACTION_DATE=2024-01-01 uv run idr-pipeline
```

## Adding data to the model
Expand Down
14 changes: 6 additions & 8 deletions apps/bfd-pipeline-idr/load-synthetic-env.sh
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,12 @@
set -e

read -p "Are you sure you want to overwrite the data in ${BFD_ENV}? [yn] " -n 1 -r
echo # (optional) move to a new line
if ! [[ $REPLY =~ ^[Yy]$ ]]
then
echo 'exiting'
exit 0
echo # (optional) move to a new line
if ! [[ $REPLY =~ ^[Yy]$ ]]; then
echo 'exiting'
exit 0
fi


DB_CLUSTER="bfd-${BFD_ENV}-aurora-cluster"
readonly DB_CLUSTER
BFD_DB_USERNAME="$(aws ssm get-parameter --name "/bfd/${BFD_ENV}/idr-pipeline/sensitive/db/username" --with-decryption --query "Parameter.Value" --output text)"
Expand Down Expand Up @@ -43,7 +41,7 @@ export IDR_SCHEMA

args=('--load-type' 'initial' '--source' 'snowflake' '--load-mode' 'synthetic')
if [[ -n "$1" ]]; then
args+=('--seed-from' "$1")
args+=('--seed-from' "$1")
fi

IDR_ENABLE_DATE_PARTITIONS=0 IDR_ENABLE_PRIOR_AUTH=1 uv run pipeline.py "${args[@]}"
IDR_ENABLE_DATE_PARTITIONS=0 IDR_ENABLE_PRIOR_AUTH=1 uv run idr-pipeline "${args[@]}"
16 changes: 12 additions & 4 deletions apps/bfd-pipeline-idr/pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
[project]
name = "bfd-pipeline-idr"
name = "idr-pipeline"
version = "0.1.0"
description = "Add your description here"
readme = "README.md"
description = "ETL Pipeline to load data from the IDR into BFD's database"
requires-python = ">=3.14.3"
dependencies = [
"cryptography>=46.0.5",
Expand All @@ -16,7 +15,6 @@ dependencies = [
"click>=8.3.3",
"anyio",
"loguru",
"pydantic-partial" # TODO: remove when idr_bfd_compare is moved
]

[dependency-groups]
Expand All @@ -28,9 +26,19 @@ dev = [
"pyright>=1.1.407",
]

[project.scripts]
idr-pipeline = "idr_pipeline:main"

[build-system]
requires = ["uv_build<=0.12"]
build-backend = "uv_build"

[tool.uv]
exclude-newer = "7 days"

[tool.uv.build-backend]
module-name = "idr_pipeline"

[tool.uv.sources]
loguru = { git = "https://github.com/Delgan/loguru", rev = "2a17be730e111cf874d0bf967c4c5da15b7e4670" }

Expand Down
19 changes: 9 additions & 10 deletions apps/bfd-pipeline-idr/run-db.sh
Original file line number Diff line number Diff line change
Expand Up @@ -18,15 +18,15 @@ 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 pipeline.py \
--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 \
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"}"
}

image=postgres:16.6
Expand Down Expand Up @@ -85,4 +85,3 @@ else
do_load "$1"
echo "Done loading full tables from $1"
fi

Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,10 @@
import psycopg # type: ignore
from loguru import logger

from batch_worker import LoadingBatchWorkerManager
from db_utils import get_connection_string
from extractor import PostgresExecutor, SnowflakeExecutor
from load_events import (
from .batch_worker import LoadingBatchWorkerManager
from .db_utils import get_connection_string
from .extractor import PostgresExecutor, SnowflakeExecutor
from .load_events import (
IdrJobLoadEvent,
get_eligible_events,
get_tables_to_load,
Expand All @@ -19,12 +19,12 @@
update_failure_times,
update_start_times,
)
from load_partition import LoadType
from load_synthetic import load_from_csv
from logger_config import configure_logger
from model.base_model import LoadMode, Source
from pipeline_stages import StagedIdrPipeline
from settings import (
from .load_partition import LoadType
from .load_synthetic import load_from_csv
from .logger_config import configure_logger
from .model.base_model import LoadMode, Source
from .pipeline_stages import StagedIdrPipeline
from .settings import (
INCREMENTAL_IDR_JOB_GRACE_PERIOD,
MAX_TASKS,
TABLES_TO_LOAD,
Expand Down Expand Up @@ -72,6 +72,12 @@ def main(
seed_from: str | None,
truncate: bool,
) -> None:
# Required to have loguru logging consistently configured across parallel pipeline nodes and
# batch worker
multiprocessing.set_start_method("spawn")
# Setup the root logger _once_
configure_logger()

if seed_from:
load_from_csv(
SnowflakeExecutor()
Expand Down Expand Up @@ -161,13 +167,3 @@ def resolve_test_date(load_mode: LoadMode) -> datetime:
if test_date and load_mode != LoadMode.PROD:
return test_date
return datetime.now(UTC)


if __name__ == "__main__":
# Required to have loguru logging consistently configured across parallel pipeline nodes and
# batch worker
multiprocessing.set_start_method("spawn")
# Setup the root logger _once_
configure_logger()

main()
Original file line number Diff line number Diff line change
Expand Up @@ -28,12 +28,12 @@
from psycopg.errors import DeadlockDetected, InFailedSqlTransaction
from psycopg_pool.abc import ACT

from load_partition import LoadPartition
from model.base_model import DbType, IdrBaseModel
from model.load_progress import LoadProgress
from parallel_executor import ExternallyCanceled
from settings import PER_BATCH_MAX_CONNECTIONS, PER_BATCH_MIN_CONNECTIONS
from timer import Timer
from .load_partition import LoadPartition
from .model.base_model import DbType, IdrBaseModel
from .model.load_progress import LoadProgress
from .parallel_executor import ExternallyCanceled
from .settings import PER_BATCH_MAX_CONNECTIONS, PER_BATCH_MIN_CONNECTIONS
from .timer import Timer

_MAX_LAST_UPDATED_CHUNKS_PER_BATCH = 100
_MINIMUM_LAST_UPDATED_CHUNK_SIZE = 100
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from dateutil.relativedelta import relativedelta

from load_partition import LoadPartition, LoadPartitionGroup, PartitionType
from settings import PARTITION_TYPE
from .load_partition import LoadPartition, LoadPartitionGroup, PartitionType
from .settings import PARTITION_TYPE

DEFAULT_MAX_DATE = "9999-12-31"
DEFAULT_MIN_DATE = "0001-01-01"
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from model.base_model import LoadMode
from settings import bfd_db_endpoint, bfd_db_name, bfd_db_password, bfd_db_port, bfd_db_username
from .model.base_model import LoadMode
from .settings import bfd_db_endpoint, bfd_db_name, bfd_db_password, bfd_db_port, bfd_db_username


def get_connection_string(load_mode: LoadMode) -> str:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,18 @@
from snowflake.connector import DictCursor, SnowflakeConnection
from snowflake.snowpark import Session

from constants import DEFAULT_MIN_DATE
from db_utils import get_connection_string
from load_partition import LoadPartition
from model.base_model import (
from .constants import DEFAULT_MIN_DATE
from .db_utils import get_connection_string
from .load_partition import LoadPartition
from .model.base_model import (
DbType,
LoadMode,
Source,
T,
format_date_opt,
)
from model.load_progress import LoadProgress
from settings import (
from .model.load_progress import LoadProgress
from .settings import (
ALLOW_EXTRACTOR_QUERY_LOGGING,
BATCH_MULTIPLIER,
ENABLE_DATE_PARTITIONS,
Expand All @@ -40,7 +40,7 @@
IDR_WAREHOUSE,
MIN_BATCH_COMPLETION_DATE,
)
from timer import Timer
from .timer import Timer


@dataclass
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,34 +17,34 @@
from pydantic.main import BaseModel
from pydantic.type_adapter import TypeAdapter

from db_utils import get_connection_string
from model.base_model import LoadMode
from model.idr_beneficiary import IdrBeneficiary
from model.idr_beneficiary_dual_eligibility import IdrBeneficiaryDualEligibility
from model.idr_beneficiary_entitlement import IdrBeneficiaryEntitlement
from model.idr_beneficiary_entitlement_reason import IdrBeneficiaryEntitlementReason
from model.idr_beneficiary_low_income_subsidy import IdrBeneficiaryLowIncomeSubsidy
from model.idr_beneficiary_low_income_subsidy_cmbnd import IdrBeneficiaryLowIncomeSubsidyCmbnd
from model.idr_beneficiary_ma_part_d_enrollment import IdrBeneficiaryMaPartDEnrollment
from model.idr_beneficiary_ma_part_d_enrollment_rx import IdrBeneficiaryMaPartDEnrollmentRx
from model.idr_beneficiary_mbi_id import IdrBeneficiaryMbiId
from model.idr_beneficiary_overshare_mbi import IdrBeneficiaryOvershareMbi
from model.idr_beneficiary_status import IdrBeneficiaryStatus
from model.idr_beneficiary_third_party import IdrBeneficiaryThirdParty
from model.idr_claim_institutional_nch import IdrClaimInstitutionalNch
from model.idr_claim_institutional_ss import IdrClaimInstitutionalSs
from model.idr_claim_item_institutional_nch import IdrClaimItemInstitutionalNch
from model.idr_claim_item_institutional_ss import IdrClaimItemInstitutionalSs
from model.idr_claim_item_professional_nch import IdrClaimItemProfessionalNch
from model.idr_claim_item_professional_ss import IdrClaimItemProfessionalSs
from model.idr_claim_professional_nch import IdrClaimProfessionalNch
from model.idr_claim_professional_ss import IdrClaimProfessionalSs
from model.idr_claim_rx import IdrClaimRx
from model.idr_contract_pbp_contact import IdrContractPbpContact
from model.idr_contract_pbp_number import IdrContractPbpNumber
from model.idr_prior_auth import IdrPriorAuth
from model.idr_prior_auth_item import IdrPriorAuthItem
from pydantic_utils import fields
from .db_utils import get_connection_string
from .model.base_model import LoadMode
from .model.idr_beneficiary import IdrBeneficiary
from .model.idr_beneficiary_dual_eligibility import IdrBeneficiaryDualEligibility
from .model.idr_beneficiary_entitlement import IdrBeneficiaryEntitlement
from .model.idr_beneficiary_entitlement_reason import IdrBeneficiaryEntitlementReason
from .model.idr_beneficiary_low_income_subsidy import IdrBeneficiaryLowIncomeSubsidy
from .model.idr_beneficiary_low_income_subsidy_cmbnd import IdrBeneficiaryLowIncomeSubsidyCmbnd
from .model.idr_beneficiary_ma_part_d_enrollment import IdrBeneficiaryMaPartDEnrollment
from .model.idr_beneficiary_ma_part_d_enrollment_rx import IdrBeneficiaryMaPartDEnrollmentRx
from .model.idr_beneficiary_mbi_id import IdrBeneficiaryMbiId
from .model.idr_beneficiary_overshare_mbi import IdrBeneficiaryOvershareMbi
from .model.idr_beneficiary_status import IdrBeneficiaryStatus
from .model.idr_beneficiary_third_party import IdrBeneficiaryThirdParty
from .model.idr_claim_institutional_nch import IdrClaimInstitutionalNch
from .model.idr_claim_institutional_ss import IdrClaimInstitutionalSs
from .model.idr_claim_item_institutional_nch import IdrClaimItemInstitutionalNch
from .model.idr_claim_item_institutional_ss import IdrClaimItemInstitutionalSs
from .model.idr_claim_item_professional_nch import IdrClaimItemProfessionalNch
from .model.idr_claim_item_professional_ss import IdrClaimItemProfessionalSs
from .model.idr_claim_professional_nch import IdrClaimProfessionalNch
from .model.idr_claim_professional_ss import IdrClaimProfessionalSs
from .model.idr_claim_rx import IdrClaimRx
from .model.idr_contract_pbp_contact import IdrContractPbpContact
from .model.idr_contract_pbp_number import IdrContractPbpNumber
from .model.idr_prior_auth import IdrPriorAuth
from .model.idr_prior_auth_item import IdrPriorAuthItem
from .pydantic_utils import fields

LOAD_EVENTS_TABLE = "source_load_events"
_PART_D_TABLES = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

from dateutil.relativedelta import relativedelta

from settings import ENABLE_DATE_PARTITIONS
from .settings import ENABLE_DATE_PARTITIONS


class LoadType(StrEnum):
Expand Down
Loading
Loading