diff --git a/data-tool/flows/batch_delete_flow.py b/data-tool/flows/batch_delete_flow.py index 097e96a08f..93d9f9beae 100644 --- a/data-tool/flows/batch_delete_flow.py +++ b/data-tool/flows/batch_delete_flow.py @@ -731,6 +731,124 @@ def get_selected_corps_mig(db_engine: Engine, config, candidates: List[str], off # nothing found in remaining pages return None, None, curr + +@task(cache_policy=NO_CACHE) +def get_lear_business_filings(db_engine: Engine, business_ids: List[int]): + """Return all filings for the provided business IDs where source = 'LEAR'.""" + if not business_ids: + return [] + + sql = text(""" + SELECT * + FROM filings + WHERE business_id IN :business_ids + AND source = :source + """).bindparams( + bindparam('business_ids', expanding=True), + bindparam('source'), + ) + + with db_engine.connect() as conn: + rows = conn.execute(sql, { + 'business_ids': business_ids, + 'source': 'LEAR', + }).fetchall() + + return [dict(row._mapping) for row in rows] + + +@task(cache_policy=NO_CACHE) +def insert_demigrated_filings(lear_engine: Engine, colin_engine: Engine, business_ids: List[int], identifiers: List[str]): + """Fetch any LEAR filings for the provided businesses and persist them for demigration tracking.""" + if not business_ids: + return + + filings = get_lear_business_filings(lear_engine, business_ids) + if not filings: + return + + id_to_identifier = dict(zip(business_ids, identifiers)) if business_ids and identifiers else {} + + def _serialize_if_json(value): + if isinstance(value, (dict, list)): + return __import__('json').dumps(value) + return value + + normalized_filings = [] + for filing in filings: + if not filing: + continue + + filing_dict = dict(filing) if hasattr(filing, 'keys') else filing + filing_id = filing_dict.get('id') + if filing_id is None: + continue + + business_id = filing_dict.get('business_id') + normalized_filings.append({ + 'identifier': id_to_identifier.get(business_id), + 'filing_id': filing_id, + 'filing_date': filing_dict.get('filing_date'), + 'filing_type': filing_dict.get('filing_type'), + 'filing_json': _serialize_if_json(filing_dict.get('filing_json')), + 'payment_id': filing_dict.get('payment_id'), + 'transaction_id': filing_dict.get('transaction_id'), + 'business_id': business_id, + 'submitter_id': filing_dict.get('submitter_id'), + 'status': filing_dict.get('status'), + 'payment_completion_date': filing_dict.get('payment_completion_date'), + 'paper_only': filing_dict.get('paper_only'), + 'completion_date': filing_dict.get('completion_date'), + 'effective_date': filing_dict.get('effective_date'), + 'source': filing_dict.get('source') or 'LEAR', + 'parent_filing_id': filing_dict.get('parent_filing_id'), + 'payment_status_code': filing_dict.get('payment_status_code'), + 'temp_reg': filing_dict.get('temp_reg'), + 'payment_account': filing_dict.get('payment_account'), + 'tech_correction_json': _serialize_if_json(filing_dict.get('tech_correction_json')), + 'colin_only': filing_dict.get('colin_only'), + 'deletion_locked': filing_dict.get('deletion_locked'), + 'submitter_roles': _serialize_if_json(filing_dict.get('submitter_roles')), + 'meta_data': _serialize_if_json(filing_dict.get('meta_data')), + 'filing_sub_type': filing_dict.get('filing_sub_type'), + 'approval_type': filing_dict.get('approval_type'), + 'application_date': filing_dict.get('application_date'), + 'notice_date': filing_dict.get('notice_date'), + 'resubmission_date': filing_dict.get('resubmission_date'), + 'hide_in_ledger': filing_dict.get('hide_in_ledger'), + 'withdrawn_filing_id': filing_dict.get('withdrawn_filing_id'), + 'withdrawal_pending': filing_dict.get('withdrawal_pending'), + 'lear_only': filing_dict.get('lear_only') + }) + + if not normalized_filings: + return + + sql = text(""" + INSERT INTO demigrated_filings (corp_num, filing_id, filing_date, filing_type, filing_json, + payment_id, transaction_id, business_id, submitter_id, status, + payment_completion_date, paper_only, completion_date, effective_date, + source, parent_filing_id, payment_status_code, temp_reg, payment_account, + tech_correction_json, colin_only, deletion_locked, + submitter_roles, meta_data, filing_sub_type, approval_type, + application_date, notice_date, resubmission_date, hide_in_ledger, + withdrawn_filing_id, withdrawal_pending, lear_only) + VALUES (:identifier, :filing_id, :filing_date, :filing_type, :filing_json, + :payment_id, :transaction_id, :business_id, :submitter_id, :status, + :payment_completion_date, :paper_only, :completion_date, :effective_date, + :source, :parent_filing_id, :payment_status_code, :temp_reg, :payment_account, + :tech_correction_json, :colin_only, :deletion_locked, + :submitter_roles, :meta_data, :filing_sub_type, :approval_type, + :application_date, :notice_date, :resubmission_date, :hide_in_ledger, + :withdrawn_filing_id, :withdrawal_pending, :lear_only) + ON CONFLICT (id) DO NOTHING + """) + + with colin_engine.connect() as conn: + conn.execute(sql, normalized_filings) + conn.commit() + + @flow(log_prints=True) def batch_delete_flow(): try: @@ -773,6 +891,10 @@ def batch_delete_flow(): break print(f'🚀 Running round {cnt} to delete {len(business_ids)} busiesses...') + # If business has lear filings save them to demigrated filings table + if config.SAVE_DEMIGRATED_LEAR_FILINGS: + insert_demigrated_filings(lear_engine, colin_engine, business_ids, identifiers) + futures = [] futures.append(lear_delete.submit(lear_engine, business_ids)) diff --git a/data-tool/flows/config.py b/data-tool/flows/config.py index bbad3b7e73..789db8733a 100644 --- a/data-tool/flows/config.py +++ b/data-tool/flows/config.py @@ -290,6 +290,8 @@ class _Config(): # pylint: disable=too-few-public-methods MIG_GROUP_IDS = os.getenv('MIG_GROUP_IDS') MIG_BATCH_IDS = os.getenv('MIG_BATCH_IDS') + SAVE_DEMIGRATED_LEAR_FILINGS = os.getenv('SAVE_DEMIGRATED_LEAR_FILINGS', 'False') == 'True' + # ------------------------------------------------------------------------------------------ # Auth-only flows (auth_processing tracking) # ------------------------------------------------------------------------------------------ diff --git a/data-tool/scripts/colin_corps_extract_postgres_ddl b/data-tool/scripts/colin_corps_extract_postgres_ddl index e2914f96e6..4450cbfb5d 100644 --- a/data-tool/scripts/colin_corps_extract_postgres_ddl +++ b/data-tool/scripts/colin_corps_extract_postgres_ddl @@ -1312,13 +1312,9 @@ create table if not exists demigrated_filings ( payment_status_code character varying(50), temp_reg character varying(10), payment_account character varying(30), - court_order_file_number character varying(20), - court_order_date timestamp with time zone, - court_order_effect_of_order character varying(500), tech_correction_json jsonb, colin_only boolean, deletion_locked boolean, - order_details character varying(2000), submitter_roles character varying(200), meta_data jsonb, filing_sub_type character varying(30), diff --git a/data-tool/tests/delta_restore/fixtures/minimal_schema.sql b/data-tool/tests/delta_restore/fixtures/minimal_schema.sql index 644a534faa..2706d2b55a 100644 --- a/data-tool/tests/delta_restore/fixtures/minimal_schema.sql +++ b/data-tool/tests/delta_restore/fixtures/minimal_schema.sql @@ -170,13 +170,9 @@ CREATE TABLE demigrated_filings ( payment_status_code text, temp_reg text, payment_account text, - court_order_file_number text, - court_order_date timestamptz, - court_order_effect_of_order text, tech_correction_json jsonb, colin_only boolean, deletion_locked boolean, - order_details text, submitter_roles text, meta_data jsonb, filing_sub_type text,