From 5589421efa1803dba3c01c39a9be448bfa824a73 Mon Sep 17 00:00:00 2001 From: James McKay Date: Wed, 29 Jul 2026 10:11:26 -0400 Subject: [PATCH 1/4] [DEV-15456] download_generation --- .../filestreaming/download_generation.py | 54 ++++++++++--------- 1 file changed, 28 insertions(+), 26 deletions(-) diff --git a/usaspending_api/download/filestreaming/download_generation.py b/usaspending_api/download/filestreaming/download_generation.py index 244a42ef49..68ab562f1c 100644 --- a/usaspending_api/download/filestreaming/download_generation.py +++ b/usaspending_api/download/filestreaming/download_generation.py @@ -11,6 +11,7 @@ from datetime import datetime, timezone from pathlib import Path from typing import Optional, Tuple +from urllib.parse import urlparse import psutil as ps from django.conf import settings @@ -887,10 +888,11 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo """Executes a single PSQL command within its own Subprocess""" download_sql = Path(temp_sql_file_path).read_text() if download_sql.startswith("\\COPY"): - # Trace library parses the SQL, but cannot understand the psql-specific \COPY command. Use standard COPY here. download_sql = download_sql[1:] - # Stack 3 context managers: (1) psql code, (2) Download replica query, (3) (same) Postgres query + # Parse the database URL to extract credentials + db_url = urlparse(retrieve_db_string()) + subprocess_trace = SubprocessTrace( name=f"job.{JOB_TYPE}.download.psql", kind=SpanKind.INTERNAL, @@ -898,33 +900,34 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo ) with subprocess_trace as span: - span.set_attributes( - { - "service": "bulk-download", - "resource": str(download_sql), - "span_type": "Internal", - "source_path": str(source_path), - # download job details - "download_job_id": str(download_job.download_job_id), - "download_job_status": str(download_job.job_status.name), - "download_file_name": str(download_job.file_name), - "download_file_size": download_job.file_size if download_job.file_size is not None else 0, - "number_of_rows": download_job.number_of_rows if download_job.number_of_rows is not None else 0, - "number_of_columns": ( - download_job.number_of_columns if download_job.number_of_columns is not None else 0 - ), - "error_message": download_job.error_message if download_job.error_message else "", - "monthly_download": str(download_job.monthly_download), - "json_request": str(download_job.json_request) if download_job.json_request else "", - } - ) + span.set_attributes({ + "service": "bulk-download", + "resource": str(download_sql), + "span_type": "Internal", + "source_path": str(source_path), + "download_job_id": str(download_job.download_job_id), + "download_job_status": str(download_job.job_status.name), + "download_file_name": str(download_job.file_name), + "download_file_size": download_job.file_size if download_job.file_size is not None else 0, + "number_of_rows": download_job.number_of_rows if download_job.number_of_rows is not None else 0, + "number_of_columns": download_job.number_of_columns if download_job.number_of_columns is not None else 0, + "error_message": download_job.error_message if download_job.error_message else "", + "monthly_download": str(download_job.monthly_download), + "json_request": str(download_job.json_request) if download_job.json_request else "", + }) try: log_time = time.perf_counter() temp_env = os.environ.copy() + + # Set PostgreSQL environment variables + temp_env["PGHOST"] = db_url.hostname + temp_env["PGPORT"] = str(db_url.port or 5432) + temp_env["PGUSER"] = db_url.username + temp_env["PGPASSWORD"] = db_url.password + temp_env["PGDATABASE"] = db_url.path.lstrip('/') + if download_job and not download_job.monthly_download: - # Since terminating the process isn't guaranteed to end the DB statement, - # add timeout to client connection temp_env["PGOPTIONS"] = ( f"--statement-timeout={settings.DOWNLOAD_DB_TIMEOUT_IN_HOURS}h " f"--work-mem={settings.DOWNLOAD_DB_WORK_MEM_IN_MB}MB" @@ -932,7 +935,7 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo cat_command = subprocess.Popen(["cat", temp_sql_file_path], stdout=subprocess.PIPE) subprocess.check_output( - ["psql", "-q", "-o", source_path, retrieve_db_string(), "-v", "ON_ERROR_STOP=1"], + ["psql", "-q", "-o", source_path, "-v", "ON_ERROR_STOP=1"], # No connection string! stdin=cat_command.stdout, stderr=subprocess.STDOUT, env=temp_env, @@ -948,7 +951,6 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo raise e except Exception as e: if not settings.IS_LOCAL: - # Not logging the command as it can contain the database connection string e.cmd = "[redacted psql command]" write_to_log(message=e, is_error=True, download_job=download_job) sql = subprocess.check_output(["cat", temp_sql_file_path]).decode() From cd7dcd22d59d0c44112d901f48e094a86a2550f1 Mon Sep 17 00:00:00 2001 From: James McKay Date: Wed, 29 Jul 2026 11:18:20 -0400 Subject: [PATCH 2/4] [DEV-15465] refactor download_generation and populate_monthly_delta_files to use new psql helpers --- .../filestreaming/download_generation.py | 41 ++++---- .../download/helpers/psql_helpers.py | 89 +++++++++++++++++ .../commands/populate_monthly_delta_files.py | 95 +++++++------------ 3 files changed, 139 insertions(+), 86 deletions(-) create mode 100644 usaspending_api/download/helpers/psql_helpers.py diff --git a/usaspending_api/download/filestreaming/download_generation.py b/usaspending_api/download/filestreaming/download_generation.py index 68ab562f1c..e77b97784e 100644 --- a/usaspending_api/download/filestreaming/download_generation.py +++ b/usaspending_api/download/filestreaming/download_generation.py @@ -11,7 +11,6 @@ from datetime import datetime, timezone from pathlib import Path from typing import Optional, Tuple -from urllib.parse import urlparse import psutil as ps from django.conf import settings @@ -36,6 +35,7 @@ from usaspending_api.download.helpers import verify_requested_columns_available from usaspending_api.download.helpers import write_to_download_log as write_to_log from usaspending_api.download.helpers.cleanup_helpers import cleanup_download_files, cleanup_previous_download_attempt +from usaspending_api.download.helpers.psql_helpers import build_psql_env, run_psql_to_file from usaspending_api.download.lookups import FILE_FORMATS, JOB_STATUS_DICT, VALUE_MAPPINGS from usaspending_api.download.models.download_job import DownloadJob from usaspending_api.download.models.download_job_lookup import DownloadJobLookup @@ -890,9 +890,6 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo if download_sql.startswith("\\COPY"): download_sql = download_sql[1:] - # Parse the database URL to extract credentials - db_url = urlparse(retrieve_db_string()) - subprocess_trace = SubprocessTrace( name=f"job.{JOB_TYPE}.download.psql", kind=SpanKind.INTERNAL, @@ -918,27 +915,23 @@ def execute_psql(temp_sql_file_path: str, source_path: str, download_job: Downlo try: log_time = time.perf_counter() - temp_env = os.environ.copy() - - # Set PostgreSQL environment variables - temp_env["PGHOST"] = db_url.hostname - temp_env["PGPORT"] = str(db_url.port or 5432) - temp_env["PGUSER"] = db_url.username - temp_env["PGPASSWORD"] = db_url.password - temp_env["PGDATABASE"] = db_url.path.lstrip('/') - - if download_job and not download_job.monthly_download: - temp_env["PGOPTIONS"] = ( - f"--statement-timeout={settings.DOWNLOAD_DB_TIMEOUT_IN_HOURS}h " - f"--work-mem={settings.DOWNLOAD_DB_WORK_MEM_IN_MB}MB" - ) - cat_command = subprocess.Popen(["cat", temp_sql_file_path], stdout=subprocess.PIPE) - subprocess.check_output( - ["psql", "-q", "-o", source_path, "-v", "ON_ERROR_STOP=1"], # No connection string! - stdin=cat_command.stdout, - stderr=subprocess.STDOUT, - env=temp_env, + # Build PostgreSQL environment using helper + psql_env = build_psql_env( + dsn=retrieve_db_string(), + statement_timeout_hours=settings.DOWNLOAD_DB_TIMEOUT_IN_HOURS if ( + download_job and not download_job.monthly_download) else None, + work_mem_mb=settings.DOWNLOAD_DB_WORK_MEM_IN_MB if ( + download_job and not download_job.monthly_download) else None + ) + + # Execute psql using helper + run_psql_to_file( + sql_path=temp_sql_file_path, + output_path=source_path, + env=psql_env, + quiet=True, + on_error_stop=True ) duration = time.perf_counter() - log_time diff --git a/usaspending_api/download/helpers/psql_helpers.py b/usaspending_api/download/helpers/psql_helpers.py new file mode 100644 index 0000000000..c7c4c9d4d3 --- /dev/null +++ b/usaspending_api/download/helpers/psql_helpers.py @@ -0,0 +1,89 @@ +import os +import subprocess +from typing import Optional +from urllib.parse import urlparse + + +def build_psql_env( + dsn: str, + statement_timeout_hours: Optional[int] = None, + work_mem_mb: Optional[int] = None, + base_env: Optional[dict] = None +) -> dict: + """ + Build PostgreSQL environment variables from a database connection string. + + Args: + dsn: Database connection string (e.g., postgresql://user:pass@host:port/dbname) + statement_timeout_hours: Optional statement timeout in hours + work_mem_mb: Optional work memory in MB + base_env: Base environment to copy from (defaults to os.environ) + + Returns: + Dictionary of environment variables for psql + """ + db_url = urlparse(dsn) + + env = (base_env or os.environ).copy() + + # Set PostgreSQL connection parameters + env["PGHOST"] = db_url.hostname + env["PGPORT"] = str(db_url.port or 5432) + env["PGUSER"] = db_url.username + env["PGPASSWORD"] = db_url.password + env["PGDATABASE"] = db_url.path.lstrip('/') + + # Set optional PostgreSQL options + if statement_timeout_hours or work_mem_mb: + options = [] + if statement_timeout_hours: + options.append(f"--statement-timeout={statement_timeout_hours}h") + if work_mem_mb: + options.append(f"--work-mem={work_mem_mb}MB") + env["PGOPTIONS"] = " ".join(options) + + return env + + +def run_psql_to_file( + sql_path: str, + output_path: str, + env: dict, + quiet: bool = True, + on_error_stop: bool = True +) -> subprocess.CompletedProcess: + """ + Execute a psql command that reads SQL from a file and writes output to another file. + + Args: + sql_path: Path to SQL file to execute + output_path: Path where psql should write output + env: Environment variables (should include PGHOST, PGUSER, etc.) + quiet: If True, suppress psql output messages + on_error_stop: If True, stop on first error + + Returns: + CompletedProcess object from subprocess + + Raises: + subprocess.CalledProcessError: If psql command fails + """ + psql_args = ["psql"] + + if quiet: + psql_args.append("-q") + + psql_args.extend(["-o", output_path]) + + if on_error_stop: + psql_args.extend(["-v", "ON_ERROR_STOP=1"]) + + # Pipe SQL file content to psql + cat_command = subprocess.Popen(["cat", sql_path], stdout=subprocess.PIPE) + + return subprocess.check_output( + psql_args, + stdin=cat_command.stdout, + stderr=subprocess.STDOUT, + env=env, + ) diff --git a/usaspending_api/download/management/commands/populate_monthly_delta_files.py b/usaspending_api/download/management/commands/populate_monthly_delta_files.py index 8b8de29afd..1a5f4fe43a 100644 --- a/usaspending_api/download/management/commands/populate_monthly_delta_files.py +++ b/usaspending_api/download/management/commands/populate_monthly_delta_files.py @@ -1,8 +1,6 @@ import logging import os import re -import shutil -import subprocess import tempfile from datetime import date, datetime @@ -14,16 +12,11 @@ from django.db.models import Case, CharField, F, Q, Value, When from usaspending_api.awards.v2.lookups.lookups import all_award_types_mappings as all_ats_mappings -from usaspending_api.common.csv_helpers import count_rows_in_delimited_file -from usaspending_api.common.helpers.orm_helpers import generate_raw_quoted_query from usaspending_api.common.helpers.s3_helpers import multipart_upload from usaspending_api.config import CONFIG -from usaspending_api.download.filestreaming.download_generation import ( - apply_annotations_to_sql, - split_and_zip_data_files, -) from usaspending_api.download.filestreaming.download_source import DownloadSource from usaspending_api.download.helpers import pull_modified_agencies_cgacs +from usaspending_api.download.helpers.psql_helpers import build_psql_env, run_psql_to_file from usaspending_api.download.lookups import VALUE_MAPPINGS from usaspending_api.references.models import SubtierAgency, ToptierAgency @@ -85,7 +78,6 @@ def download(self, award_type: str, agency: str | dict = "all", generate_since: ) source.query_paths.update({"correction_delete_ind": award_map["correction_delete_ind"]}) if award_type == "Contracts": - # Add the agency_id column to the mappings source.query_paths.update({"agency_id": "transaction__contract_data__agency_id"}) source.query_paths.move_to_end("agency_id", last=False) source.query_paths.move_to_end("correction_delete_ind", last=False) @@ -120,12 +112,12 @@ def download(self, award_type: str, agency: str | dict = "all", generate_since: source.queryset = source.queryset.filter(Q(update_date_filter | Q(transaction__transactiondelta__isnull=False))) - # Generate file - file_path = self.create_local_file(award_type, source, agency_code, generate_since) + # Generate file using helper functions + file_path = self.create_local_file_with_psql(award_type, source, agency_code, generate_since) + if file_path is None: logger.info("No new, modified, or deleted data; discarding file") elif not settings.IS_LOCAL: - # Upload file to S3 and delete local version logger.info("Uploading file to S3 bucket and deleting local copy") multipart_upload( CONFIG.MONTHLY_DOWNLOAD_S3_BUCKET_NAME, @@ -139,60 +131,39 @@ def download(self, award_type: str, agency: str | dict = "all", generate_since: "Finished generation. {}, Agency: {}".format(award_type, agency if agency == "all" else agency["name"]) ) - def create_local_file( - self, award_type: str, source: pd.DataFrame, agency_code: str, generate_since: str | None - ) -> str | None: - """Generate complete file from SQL query and S3 bucket deletion files, then zip it locally""" - logger.info("Generating CSV file with creations and modifications") - - # Create file paths and working directory - timestamp = datetime.strftime(datetime.now(), "%Y%m%d%H%M%S%f") - working_dir = f"{settings.CSV_LOCAL_PATH}_{agency_code}_delta_gen_{timestamp}/" - if not os.path.exists(working_dir): - os.mkdir(working_dir) - agency_str = "All" if agency_code == "all" else agency_code - source_name = f"FY(All)_{agency_str}_{award_type}_Delta_{datetime.strftime(date.today(), '%Y%m%d')}" - source_path = os.path.join(working_dir, "{}.csv".format(source_name)) - - # Create a unique temporary file with the raw query - raw_quoted_query = generate_raw_quoted_query(source.row_emitter(None)) # None requests all headers - - csv_query_annotated = apply_annotations_to_sql(raw_quoted_query, source.human_names) - - (temp_sql_file, temp_sql_file_path) = tempfile.mkstemp(prefix="bd_sql_", dir="/tmp") - with open(temp_sql_file_path, "w") as file: - file.write("\\copy ({}) To STDOUT with CSV HEADER".format(csv_query_annotated)) - - logger.info("Generated temp SQL file {}".format(temp_sql_file_path)) - # Generate the csv with \copy - cat_command = subprocess.Popen(["cat", temp_sql_file_path], stdout=subprocess.PIPE) + def create_local_file_with_psql(self, award_type: str, source: DownloadSource, agency_code: str, + generate_since: str) -> str: + """Generate the file using psql helpers""" + # Generate SQL query + sql_query = self.generate_sql_query(source) + + # Create temp SQL file + temp_sql_file = tempfile.NamedTemporaryFile(mode='w', delete=False, suffix='.sql') + temp_sql_file.write(sql_query) + temp_sql_file.close() + + # Build output path + output_path = self.build_output_path(award_type, agency_code, generate_since) + try: - subprocess.check_output( - ["psql", "-o", source_path, os.environ["DOWNLOAD_DATABASE_URL"], "-v", "ON_ERROR_STOP=1"], - stdin=cat_command.stdout, - stderr=subprocess.STDOUT, + # Build PostgreSQL environment + psql_env = build_psql_env( + dsn=settings.DOWNLOAD_DATABASE_URL, + statement_timeout_hours=settings.DOWNLOAD_DB_TIMEOUT_IN_HOURS, + work_mem_mb=settings.DOWNLOAD_DB_WORK_MEM_IN_MB ) - except subprocess.CalledProcessError as e: - logger.exception(e.output) - raise e - - # Append deleted rows to the end of the file - if not self.debugging_skip_deleted: - self.add_deletion_records(source_path, working_dir, award_type, agency_code, source, generate_since) - if count_rows_in_delimited_file(source_path, has_header=True, safe=True) > 0: - # Split the CSV into multiple files and zip it up - zipfile_path = "{}{}.zip".format(settings.CSV_LOCAL_PATH, source_name) - - logger.info("Creating compressed file: {}".format(os.path.basename(zipfile_path))) - split_and_zip_data_files(zipfile_path, source_path, source_name, "csv") - else: - zipfile_path = None - os.close(temp_sql_file) - os.remove(temp_sql_file_path) - shutil.rmtree(working_dir) + # Execute psql + run_psql_to_file( + sql_path=temp_sql_file.name, + output_path=output_path, + env=psql_env + ) - return zipfile_path + return output_path + finally: + # Cleanup temp SQL file + os.remove(temp_sql_file.name) @staticmethod def split_transaction_id(tid: str) -> pd.Series: From bc1d0f364d08fe6502a6dddcecbb807c35d3478f Mon Sep 17 00:00:00 2001 From: James McKay Date: Thu, 30 Jul 2026 10:02:21 -0400 Subject: [PATCH 3/4] [DEV-15465] interim --- .../download/helpers/psql_helpers.py | 120 +++++++++++++----- .../commands/populate_monthly_delta_files.py | 96 +++++++++++--- 2 files changed, 163 insertions(+), 53 deletions(-) diff --git a/usaspending_api/download/helpers/psql_helpers.py b/usaspending_api/download/helpers/psql_helpers.py index c7c4c9d4d3..00ff2bc64f 100644 --- a/usaspending_api/download/helpers/psql_helpers.py +++ b/usaspending_api/download/helpers/psql_helpers.py @@ -10,28 +10,29 @@ def build_psql_env( work_mem_mb: Optional[int] = None, base_env: Optional[dict] = None ) -> dict: - """ - Build PostgreSQL environment variables from a database connection string. + """Build PostgreSQL environment variables from a database connection string.""" + import logging - Args: - dsn: Database connection string (e.g., postgresql://user:pass@host:port/dbname) - statement_timeout_hours: Optional statement timeout in hours - work_mem_mb: Optional work memory in MB - base_env: Base environment to copy from (defaults to os.environ) + logger = logging.getLogger(__name__) + + if not dsn: + raise ValueError("DSN cannot be empty") + + logger.info(f"Parsing DSN: {dsn[:30]}...") - Returns: - Dictionary of environment variables for psql - """ db_url = urlparse(dsn) env = (base_env or os.environ).copy() # Set PostgreSQL connection parameters - env["PGHOST"] = db_url.hostname + env["PGHOST"] = db_url.hostname or "localhost" env["PGPORT"] = str(db_url.port or 5432) - env["PGUSER"] = db_url.username - env["PGPASSWORD"] = db_url.password - env["PGDATABASE"] = db_url.path.lstrip('/') + env["PGUSER"] = db_url.username or "postgres" + env["PGPASSWORD"] = db_url.password or "" + env["PGDATABASE"] = db_url.path.lstrip('/') if db_url.path else "postgres" + + logger.info( + f"Set PGHOST={env['PGHOST']}, PGPORT={env['PGPORT']}, PGUSER={env['PGUSER']}, PGDATABASE={env['PGDATABASE']}") # Set optional PostgreSQL options if statement_timeout_hours or work_mem_mb: @@ -51,23 +52,21 @@ def run_psql_to_file( env: dict, quiet: bool = True, on_error_stop: bool = True -) -> subprocess.CompletedProcess: +) -> None: """ Execute a psql command that reads SQL from a file and writes output to another file. + """ + import logging + logger = logging.getLogger(__name__) - Args: - sql_path: Path to SQL file to execute - output_path: Path where psql should write output - env: Environment variables (should include PGHOST, PGUSER, etc.) - quiet: If True, suppress psql output messages - on_error_stop: If True, stop on first error - - Returns: - CompletedProcess object from subprocess + # Log the SQL file contents for debugging + try: + with open(sql_path, 'r') as f: + sql_content = f.read() + logger.info(f"SQL file contents (first 500 chars): {sql_content[:500]}") + except Exception as e: + logger.error(f"Could not read SQL file: {e}") - Raises: - subprocess.CalledProcessError: If psql command fails - """ psql_args = ["psql"] if quiet: @@ -78,12 +77,69 @@ def run_psql_to_file( if on_error_stop: psql_args.extend(["-v", "ON_ERROR_STOP=1"]) - # Pipe SQL file content to psql - cat_command = subprocess.Popen(["cat", sql_path], stdout=subprocess.PIPE) + logger.info(f"psql command: {' '.join(psql_args)}") + logger.info( + f"Environment: PGHOST={env.get('PGHOST')}, " + f"PGPORT={env.get('PGPORT')}, " + f"PGUSER={env.get('PGUSER')}, " + f"PGDATABASE={env.get('PGDATABASE')}" + ) + + # Test database connection first + logger.info("Testing database connection...") + test_process = subprocess.run( + ["psql", "-c", "SELECT 1;"], + env=env, + capture_output=True, + timeout=5 + ) + if test_process.returncode != 0: + logger.error(f"Database connection test failed: {test_process.stderr.decode()}") + raise Exception(f"Cannot connect to database: {test_process.stderr.decode()}") + logger.info("Database connection test successful") - return subprocess.check_output( + logger.info("Starting cat and psql processes...") + + # Start cat process + cat_process = subprocess.Popen(["cat", sql_path], stdout=subprocess.PIPE, stderr=subprocess.PIPE) + + # Start psql process with cat's stdout as stdin + psql_process = subprocess.Popen( psql_args, - stdin=cat_command.stdout, - stderr=subprocess.STDOUT, + stdin=cat_process.stdout, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, # Changed to PIPE to capture stderr separately env=env, ) + + # Close cat's stdout in parent so psql gets EOF when cat exits + cat_process.stdout.close() + + logger.info("Waiting for processes to complete...") + + # Wait for both processes to complete with timeout + try: + psql_output, psql_error = psql_process.communicate(timeout=30) # 30 second timeout + cat_process.wait(timeout=5) + except subprocess.TimeoutExpired: + logger.error("Process timed out! Killing processes...") + psql_process.kill() + cat_process.kill() + raise Exception("psql process timed out after 30 seconds") from None + + logger.info(f"psql return code: {psql_process.returncode}") + logger.info(f"psql stdout: {psql_output.decode() if psql_output else 'empty'}") + logger.info(f"psql stderr: {psql_error.decode() if psql_error else 'empty'}") + + # Check for errors + if psql_process.returncode != 0: + error_msg = psql_error.decode() if psql_error else psql_output.decode() if psql_output else "Unknown error" + logger.error(f"psql failed: {error_msg}") + raise subprocess.CalledProcessError( + psql_process.returncode, + psql_args, + output=psql_output, + stderr=psql_error + ) + + logger.info("psql completed successfully") diff --git a/usaspending_api/download/management/commands/populate_monthly_delta_files.py b/usaspending_api/download/management/commands/populate_monthly_delta_files.py index 1a5f4fe43a..83c568126e 100644 --- a/usaspending_api/download/management/commands/populate_monthly_delta_files.py +++ b/usaspending_api/download/management/commands/populate_monthly_delta_files.py @@ -113,7 +113,7 @@ def download(self, award_type: str, agency: str | dict = "all", generate_since: source.queryset = source.queryset.filter(Q(update_date_filter | Q(transaction__transactiondelta__isnull=False))) # Generate file using helper functions - file_path = self.create_local_file_with_psql(award_type, source, agency_code, generate_since) + file_path = self.create_local_file(award_type, source, agency_code, generate_since) if file_path is None: logger.info("No new, modified, or deleted data; discarding file") @@ -131,39 +131,93 @@ def download(self, award_type: str, agency: str | dict = "all", generate_since: "Finished generation. {}, Agency: {}".format(award_type, agency if agency == "all" else agency["name"]) ) - def create_local_file_with_psql(self, award_type: str, source: DownloadSource, agency_code: str, - generate_since: str) -> str: - """Generate the file using psql helpers""" - # Generate SQL query - sql_query = self.generate_sql_query(source) + def create_local_file( + self, award_type: str, source: DownloadSource, agency_code: str, generate_since: str | None + ) -> str | None: + """Generate complete file from SQL query and S3 bucket deletion files, then zip it locally""" + import shutil + import subprocess + + from usaspending_api.common.csv_helpers import count_rows_in_delimited_file + from usaspending_api.common.helpers.orm_helpers import generate_raw_quoted_query + from usaspending_api.download.filestreaming.download_generation import ( + apply_annotations_to_sql, + split_and_zip_data_files, + ) + + logger.info("Generating CSV file with creations and modifications") + + # Create file paths and working directory + timestamp = datetime.strftime(datetime.now(), "%Y%m%d%H%M%S%f") + working_dir = f"{settings.CSV_LOCAL_PATH}_{agency_code}_delta_gen_{timestamp}/" + if not os.path.exists(working_dir): + os.mkdir(working_dir) + agency_str = "All" if agency_code == "all" else agency_code + source_name = f"FY(All)_{agency_str}_{award_type}_Delta_{datetime.strftime(date.today(), '%Y%m%d')}" + source_path = os.path.join(working_dir, "{}.csv".format(source_name)) + + # Create a unique temporary file with the raw query + raw_quoted_query = generate_raw_quoted_query(source.row_emitter(None)) + csv_query_annotated = apply_annotations_to_sql(raw_quoted_query, source.human_names) - # Create temp SQL file - temp_sql_file = tempfile.NamedTemporaryFile(mode='w', delete=False, suffix='.sql') - temp_sql_file.write(sql_query) - temp_sql_file.close() + (temp_sql_file, temp_sql_file_path) = tempfile.mkstemp(prefix="bd_sql_", dir="/tmp") + with open(temp_sql_file_path, "w") as file: + file.write("\\copy ({}) To STDOUT with CSV HEADER".format(csv_query_annotated)) - # Build output path - output_path = self.build_output_path(award_type, agency_code, generate_since) + logger.info("Generated temp SQL file {}".format(temp_sql_file_path)) try: - # Build PostgreSQL environment + # Get database URL from settings or environment variable (for tests) + db_url = settings.DOWNLOAD_DATABASE_URL or os.environ.get("DOWNLOAD_DATABASE_URL") + + if not db_url: + raise ValueError("DOWNLOAD_DATABASE_URL is not configured") + + logger.info(f"Using database URL: {db_url[:20]}...") # Log first 20 chars for debugging + + # Build PostgreSQL environment using helper psql_env = build_psql_env( - dsn=settings.DOWNLOAD_DATABASE_URL, + dsn=db_url, statement_timeout_hours=settings.DOWNLOAD_DB_TIMEOUT_IN_HOURS, work_mem_mb=settings.DOWNLOAD_DB_WORK_MEM_IN_MB ) - # Execute psql + logger.info( + f"Built psql environment with PGHOST={psql_env.get('PGHOST')}, PGDATABASE={psql_env.get('PGDATABASE')}") + + # Execute psql using helper run_psql_to_file( - sql_path=temp_sql_file.name, - output_path=output_path, - env=psql_env + sql_path=temp_sql_file_path, + output_path=source_path, + env=psql_env, + quiet=False, + on_error_stop=True ) - return output_path + except subprocess.CalledProcessError as e: + logger.exception(e.output if hasattr(e, 'output') else str(e)) + raise e finally: - # Cleanup temp SQL file - os.remove(temp_sql_file.name) + # Always cleanup temp SQL file + os.close(temp_sql_file) + os.remove(temp_sql_file_path) + + # Append deleted rows to the end of the file + if not self.debugging_skip_deleted: + self.add_deletion_records(source_path, working_dir, award_type, agency_code, source, generate_since) + + if count_rows_in_delimited_file(source_path, has_header=True, safe=True) > 0: + # Split the CSV into multiple files and zip it up + zipfile_path = "{}{}.zip".format(settings.CSV_LOCAL_PATH, source_name) + + logger.info("Creating compressed file: {}".format(os.path.basename(zipfile_path))) + split_and_zip_data_files(zipfile_path, source_path, source_name, "csv") + else: + zipfile_path = None + + shutil.rmtree(working_dir) + + return zipfile_path @staticmethod def split_transaction_id(tid: str) -> pd.Series: From 497c5870f861d50f319430b8ba134251e6624f67 Mon Sep 17 00:00:00 2001 From: James McKay Date: Thu, 30 Jul 2026 15:52:32 -0400 Subject: [PATCH 4/4] [DEV-15465] adjusted test fixture and db_url parameter query order --- .../commands/populate_monthly_delta_files.py | 2 +- .../integration/test_populate_monthly_delta_files.py | 11 +++++++++-- 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/usaspending_api/download/management/commands/populate_monthly_delta_files.py b/usaspending_api/download/management/commands/populate_monthly_delta_files.py index 83c568126e..cc4be044d8 100644 --- a/usaspending_api/download/management/commands/populate_monthly_delta_files.py +++ b/usaspending_api/download/management/commands/populate_monthly_delta_files.py @@ -168,7 +168,7 @@ def create_local_file( try: # Get database URL from settings or environment variable (for tests) - db_url = settings.DOWNLOAD_DATABASE_URL or os.environ.get("DOWNLOAD_DATABASE_URL") + db_url = os.environ.get("DOWNLOAD_DATABASE_URL") or settings.DOWNLOAD_DATABASE_URL if not db_url: raise ValueError("DOWNLOAD_DATABASE_URL is not configured") diff --git a/usaspending_api/download/tests/integration/test_populate_monthly_delta_files.py b/usaspending_api/download/tests/integration/test_populate_monthly_delta_files.py index 00669b306c..eeab066eeb 100644 --- a/usaspending_api/download/tests/integration/test_populate_monthly_delta_files.py +++ b/usaspending_api/download/tests/integration/test_populate_monthly_delta_files.py @@ -9,7 +9,6 @@ from model_bakery import baker from usaspending_api.awards.models import TransactionDelta -from usaspending_api.common.helpers.sql_helpers import get_database_dsn_string from usaspending_api.download.v2.download_column_historical_lookups import query_paths from usaspending_api.settings import HOST @@ -94,7 +93,15 @@ def monthly_download_delta_data(db, monkeypatch): ) TransactionDelta.objects.update_or_create_transaction(i) - monkeypatch.setenv("DOWNLOAD_DATABASE_URL", get_database_dsn_string()) + from django.db import connection + + db_settings = connection.settings_dict + test_db_url = ( + f"postgresql://{db_settings['USER']}:{db_settings['PASSWORD']}" + f"@{db_settings['HOST']}:{db_settings['PORT']}/{db_settings['NAME']}" + ) + + monkeypatch.setenv("DOWNLOAD_DATABASE_URL", test_db_url) @pytest.mark.django_db(transaction=True)