Skip to content
Open

34055 #4681

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
122 changes: 122 additions & 0 deletions data-tool/flows/batch_delete_flow.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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))

Expand Down
2 changes: 2 additions & 0 deletions data-tool/flows/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
# ------------------------------------------------------------------------------------------
Expand Down
4 changes: 0 additions & 4 deletions data-tool/scripts/colin_corps_extract_postgres_ddl
Original file line number Diff line number Diff line change
Expand Up @@ -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),

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Removed fields are court order related and are not in filings table anymore

submitter_roles character varying(200),
meta_data jsonb,
filing_sub_type character varying(30),
Expand Down
4 changes: 0 additions & 4 deletions data-tool/tests/delta_restore/fixtures/minimal_schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down