Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
20 changes: 15 additions & 5 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -8,15 +8,25 @@ jobs:
format_python_code:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683 # v4.2.2
- name: Check out the repo
uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683 # v4.2.2
with:
submodules: "recursive"
ref: ${{ github.ref }}
- name: Blacken Python code
uses: jpetrucciani/black-check@master

- name: Install uv
uses: astral-sh/setup-uv@v6

- name: Run Ruff Checks
uses: astral-sh/ruff-action@39f75e526a505e26a302f8796977b50c13720edf # v3.2.1
with:
path: "."
black_flags: "--safe --verbose --diff"
args: "check --select I ."

- name: Run Ruff Format
uses: astral-sh/ruff-action@39f75e526a505e26a302f8796977b50c13720edf # v3.2.1
with:
args: "format --check --diff"

build:
runs-on: ubuntu-latest
if: "!contains(github.event.head_commit.message, '[ci skip]')"
Expand Down
8 changes: 4 additions & 4 deletions app/airflow/dags/SR_processing_dag.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
from datetime import datetime, timedelta

from airflow import DAG
from airflow.operators.empty import EmptyOperator
from libs.utils import create_task, validate_params_SR_processing
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT, AIRFLOW_DEBUG_MODE
from libs.SR_processing.core import (
process_data_dictionary,
process_and_create_scan_report_entries,
process_data_dictionary,
)
from libs.utils import connect_to_storage
from libs.settings import AIRFLOW_DEBUG_MODE, AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import connect_to_storage, create_task, validate_params_SR_processing

"""
This DAG automates the process of creating scan report tables, fields and values
Expand Down
33 changes: 16 additions & 17 deletions app/airflow/dags/auto_mapping_dag.py
Original file line number Diff line number Diff line change
@@ -1,36 +1,35 @@
from datetime import datetime, timedelta

from airflow import DAG
from airflow.operators.empty import EmptyOperator
from libs.auto_mapping.find_R_concepts_to_reuse import (
find_matching_field,
find_matching_value,
create_reusing_concepts,
delete_R_concepts,
)

from libs.auto_mapping.find_standard_V_concepts import (
find_standard_concepts,
create_standard_concepts,
)

from libs.auto_mapping.core_get_existing_concepts import (
delete_mapping_rules,
find_existing_concepts,
)
from libs.auto_mapping.core_prep_rules_creation import (
find_dest_table_and_person_field_id,
find_date_fields,
find_concept_fields,
find_additional_fields,
find_concept_fields,
find_date_fields,
find_dest_table_and_person_field_id,
)
from libs.auto_mapping.core_rules_creation import create_mapping_rules
from libs.auto_mapping.find_R_concepts_to_reuse import (
create_reusing_concepts,
delete_R_concepts,
find_matching_field,
find_matching_value,
)
from libs.auto_mapping.find_standard_V_concepts import (
create_standard_concepts,
find_standard_concepts,
)
from libs.auto_mapping.search_recommendations import process_search_recommendations
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT, AIRFLOW_DEBUG_MODE, SEARCH_ENABLED
from libs.utils import (
create_task,
validate_params_auto_mapping,
update_job_status_on_failure,
validate_params_auto_mapping,
)
from libs.settings import AIRFLOW_DEBUG_MODE, SEARCH_ENABLED, AIRFLOW_DAGRUN_TIMEOUT

"""
This DAG automates the process of creating and reusing concepts from scan reports and generating mapping rules.
Expand Down
24 changes: 12 additions & 12 deletions app/airflow/dags/libs/SR_processing/core.py
Original file line number Diff line number Diff line change
@@ -1,28 +1,28 @@
import csv
import logging
import os
from openpyxl import load_workbook
from libs.utils import pull_validated_params
import csv
from io import StringIO
from typing import List, Tuple

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.enums import JobStageType, StageStatusType
from libs.queries import create_temp_data_dictionary_table_query, create_values_query
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.SR_processing.db_services import (
create_field_entries,
update_temp_data_dictionary_table,
create_temp_field_values_table,
delete_temp_tables,
update_temp_data_dictionary_table,
)
from libs.SR_processing.helpers import (
remove_BOM,
process_four_item_dict,
get_unique_table_names,
process_four_item_dict,
remove_BOM,
transform_scan_report_sheet_table,
)
from libs.storage_services import download_blob_to_tmp
from airflow.providers.postgres.hooks.postgres import PostgresHook
from typing import List, Tuple
from libs.queries import create_values_query, create_temp_data_dictionary_table_query
from libs.enums import JobStageType, StageStatusType
from libs.utils import update_job_status
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import pull_validated_params, update_job_status
from openpyxl import load_workbook

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
9 changes: 5 additions & 4 deletions app/airflow/dags/libs/SR_processing/db_services.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,13 @@
from typing import List, Any, Dict, Tuple
from openpyxl.worksheet.worksheet import Worksheet
from libs.queries import create_fields_query
import logging
from collections import defaultdict
from typing import Any, Dict, List, Tuple

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.enums import JobStageType, StageStatusType
from libs.queries import create_fields_query
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT, AIRFLOW_DEBUG_MODE
from libs.utils import update_job_status
from libs.settings import AIRFLOW_DEBUG_MODE, AIRFLOW_DAGRUN_TIMEOUT
from openpyxl.worksheet.worksheet import Worksheet

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
5 changes: 3 additions & 2 deletions app/airflow/dags/libs/SR_processing/helpers.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
from typing import List, Any, Dict
from openpyxl.worksheet.worksheet import Worksheet
import logging
from collections import defaultdict
from typing import Any, Dict, List

from openpyxl.worksheet.worksheet import Worksheet


def get_unique_table_names(worksheet: Worksheet) -> List[str]:
Expand Down
15 changes: 8 additions & 7 deletions app/airflow/dags/libs/auto_mapping/core_get_existing_concepts.py
Original file line number Diff line number Diff line change
@@ -1,17 +1,18 @@
from libs.utils import (
update_job_status,
JobStageType,
StageStatusType,
pull_validated_params,
)
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.queries import (
create_existing_concepts_table_query,
find_existing_concepts_query,
find_source_field_id_query,
)
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import (
JobStageType,
StageStatusType,
pull_validated_params,
update_job_status,
)

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
19 changes: 10 additions & 9 deletions app/airflow/dags/libs/auto_mapping/core_prep_rules_creation.py
Original file line number Diff line number Diff line change
@@ -1,17 +1,18 @@
from libs.utils import (
update_job_status,
JobStageType,
StageStatusType,
pull_validated_params,
)
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.queries import (
find_dest_table_and_person_field_id_query,
find_dates_fields_query,
find_concept_fields_query,
find_dates_fields_query,
find_dest_table_and_person_field_id_query,
)
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import (
JobStageType,
StageStatusType,
pull_validated_params,
update_job_status,
)

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
8 changes: 4 additions & 4 deletions app/airflow/dags/libs/auto_mapping/core_rules_creation.py
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging

from airflow.exceptions import AirflowException
from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT, AIRFLOW_DEBUG_MODE
from libs.utils import (
update_job_status,
JobStageType,
StageStatusType,
pull_validated_params,
update_job_status,
)
from libs.settings import AIRFLOW_DEBUG_MODE, AIRFLOW_DAGRUN_TIMEOUT


# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
17 changes: 9 additions & 8 deletions app/airflow/dags/libs/auto_mapping/find_R_concepts_to_reuse.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,21 @@
from libs.utils import (
update_job_status,
JobStageType,
StageStatusType,
pull_validated_params,
)
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.queries import (
create_reuse_concept_query,
create_temp_reuse_table_query,
find_m_concepts_query,
find_m_concepts_query_field_level,
find_object_id_query,
validate_reused_concepts_query,
create_reuse_concept_query,
)
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import (
JobStageType,
StageStatusType,
pull_validated_params,
update_job_status,
)

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
11 changes: 6 additions & 5 deletions app/airflow/dags/libs/auto_mapping/find_standard_V_concepts.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,14 @@
import logging

from airflow.exceptions import AirflowException
from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import (
update_job_status,
JobStageType,
StageStatusType,
pull_validated_params,
update_job_status,
)
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging
from airflow.exceptions import AirflowException
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
7 changes: 4 additions & 3 deletions app/airflow/dags/libs/auto_mapping/search_recommendations.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
import logging

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import (
pull_validated_params,
)
from airflow.providers.postgres.hooks.postgres import PostgresHook
import logging
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
20 changes: 10 additions & 10 deletions app/airflow/dags/libs/rules_export/core.py
Original file line number Diff line number Diff line change
@@ -1,23 +1,23 @@
import logging
from libs.utils import pull_validated_params
from datetime import datetime
from typing import Dict

from airflow.providers.postgres.hooks.postgres import PostgresHook
from libs.types import FileHandlerConfig
from libs.enums import JobStageType, StageStatusType
from libs.queries import create_file_entry_query, create_update_temp_rules_table_query
from libs.rules_export.file_services import (
build_rules_json,
build_rules_csv,
build_rules_json,
build_rules_json_v2,
)
from typing import Dict
from datetime import datetime
from libs.queries import create_update_temp_rules_table_query, create_file_entry_query
from libs.storage_services import upload_blob_to_storage
from libs.enums import JobStageType, StageStatusType
from libs.utils import update_job_status
from libs.settings import (
AIRFLOW_DAGRUN_TIMEOUT,
AIRFLOW_DEBUG_MODE,
AIRFLOW_VAR_JSON_VERSION,
AIRFLOW_DAGRUN_TIMEOUT,
)
from libs.storage_services import upload_blob_to_storage
from libs.types import FileHandlerConfig
from libs.utils import pull_validated_params, update_job_status

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down
36 changes: 18 additions & 18 deletions app/airflow/dags/libs/rules_export/file_services.py
Original file line number Diff line number Diff line change
@@ -1,16 +1,16 @@
from typing import Any, Dict, List
from libs.enums import JobStageType, StageStatusType
from libs.utils import update_job_status
import csv
import io
import json
import pandas as pd
from datetime import datetime, timezone
import logging
from datetime import date, datetime, timezone
from io import BytesIO
from typing import Any, Dict, List

import pandas as pd
from airflow.providers.postgres.hooks.postgres import PostgresHook
import io
import csv
from datetime import date
import logging
from libs.enums import JobStageType, StageStatusType
from libs.settings import AIRFLOW_DAGRUN_TIMEOUT
from libs.utils import update_job_status

# PostgreSQL connection hook
pg_hook = PostgresHook(
Expand Down Expand Up @@ -353,19 +353,19 @@ def build_rules_json_v2(scan_report_name: str, scan_report_id: int) -> BytesIO:

# Adding the person_id_mapping, date_mapping, concept_mapping to the result
if person_id_mappings:
result[dest_table_str][source_table_clean][
"person_id_mapping"
] = person_id_mappings
result[dest_table_str][source_table_clean]["person_id_mapping"] = (
person_id_mappings
)

if date_mappings:
result[dest_table_str][source_table_clean][
"date_mapping"
] = date_mappings
result[dest_table_str][source_table_clean]["date_mapping"] = (
date_mappings
)

if concept_mappings:
result[dest_table_str][source_table_clean][
"concept_mappings"
] = concept_mappings
result[dest_table_str][source_table_clean]["concept_mappings"] = (
concept_mappings
)

cdm = {
"metadata": metadata,
Expand Down
Loading