diff --git a/statvar_imports/statistics_poland/README.md b/statvar_imports/statistics_poland/README.md index f396a06bde..04459945a0 100644 --- a/statvar_imports/statistics_poland/README.md +++ b/statvar_imports/statistics_poland/README.md @@ -1,48 +1,94 @@ -# Poland Demographics Dataset -## Overview -This dataset provides foundational demographic and socio-economic statistics for Poland, sourced directly from official national datasets. - -## Data Source - -**Source URL:** -https://stat.gov.pl/en/databases/ +# Statistics Poland (GUS) Demographics Import +## Overview +- **Import Name:** `statistics_poland` +- **Import Type:** Automated (Scheduled via Cloud Batch cron: `0 0 1 1,4,7,10 *`) +- **Curator:** `support@datacommons.org` -The data comes from Poland's official statistical authority and includes comprehensive demographic variables such as population counts, age distributions, and other census-related metrics. - -## How To Download Input Data -To download the data, you'll need to use the provided download script download_input_data.py. This script processes the StatisticsPoland_input.csv file available in Bucket with path datcom-prod-imports/statvar_imports/statistics_poland/poland_data_sample to generate StatisticsPoland_input_*.csv inside a new "source_files" folder. - -type of place: State. +This import processes demographic population statistics for Poland from the official Central +Statistical Office of Poland (Główny Urząd Statystyczny - GUS) Bank Danych Lokalnych (BDL) API. -statvars: Demographics +## Data Source +- **Source Portal:** https://bdl.stat.gov.pl/bdl/dane/podgrup/tablica +- **API Endpoint:** `https://bdl.stat.gov.pl/api/v1` +- **Subject ID:** `P3447` (Ludność wg grup wieku / Population by age group) +- **Place Types Covered:** + - Country level: `country/POL` + - Voivodship (province) level: `nuts/PL*` (16 voivodships) +- **Variables Covered:** + - Age groups: `0-2`, `3-6`, `7-12`, `13-15`, `16-19`, `20-24`, `25-34`, `35-44`, `45-54`, + `55-64`, `65 and more` + - Gender: `males`, `females`, `total` + - Residence Classification: `in urban areas`, `in rural areas`, `total` +- **Temporal Coverage:** Dynamic from 2003 through `current_year + 1` (currently 2025). -years: 2003 to 2025. +## Configuration & Credentials +- **API Key:** GUS BDL API key is configured with a default client key and can be overridden + via the `BDL_API_KEY` environment variable: + ```bash + export BDL_API_KEY="" + ``` +- **Template CSV:** Download script loads the reference column structure from GCS: + `gs://datcom-prod-imports/statvar_imports/statistics_poland/poland_data_sample/` + `StatisticsPoland_input.csv` + with an automatic offline fallback to `test/StatisticsPoland_input.csv`. +- **StatVar MCF:** Canonical StatVars are resolved against: + `gs://unresolved_mcf/scripts/statvar/stat_vars.mcf` ## Processing Instructions -To process the Poland Census data and generate statistical variables, use the following command from the "data" directory: -**Download input file** - ```bash - python3 statvar_imports/statistics_poland/download_input_data.py +### 1. Download and Preprocess Source Files +Run from the repository root: +```bash +python3 statvar_imports/statistics_poland/download_input_data.py ``` -**For Test Data Run** +This queries the GUS BDL API and writes annual input CSV files into +`statvar_imports/statistics_poland/source_files/StatisticsPoland_input_.csv`. + +### 2. Run StatVarProcessor (Main Data Run) ```bash python3 tools/statvar_importer/stat_var_processor.py \ - --input_data=statvar_imports/statistics_poland/test/StatisticsPoland_input.csv \ + --input_data='statvar_imports/statistics_poland/source_files/*.csv' \ --pv_map=statvar_imports/statistics_poland/StatisticsPoland_pvmap.csv \ - --output_path=statvar_imports/statistics_poland/test/StatisticsPoland_output \ + --output_path=statvar_imports/statistics_poland/StatisticsPoland_output \ --config_file=statvar_imports/statistics_poland/StatisticsPoland_metadata.csv \ - --output_counters=statvar_imports/statistics_poland/test/StatisticsPoland_output_counters.csv \ + --output_counters=\ +statvar_imports/statistics_poland/counters/StatisticsPoland_output_counters.csv \ --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf ``` -**For Main data run** + +### 3. Run StatVarProcessor on Test Data ```bash python3 tools/statvar_importer/stat_var_processor.py \ - --input_data='statvar_imports/statistics_poland/source_files/*.csv' \ + --input_data=statvar_imports/statistics_poland/test/StatisticsPoland_input.csv \ --pv_map=statvar_imports/statistics_poland/StatisticsPoland_pvmap.csv \ - --output_path=statvar_imports/statistics_poland/StatisticsPoland_output \ + --output_path=statvar_imports/statistics_poland/test/StatisticsPoland_output \ --config_file=statvar_imports/statistics_poland/StatisticsPoland_metadata.csv \ - --output_counters=statvar_imports/statistics_poland/counters/StatisticsPoland_output_counters.csv \ + --output_counters=\ +statvar_imports/statistics_poland/test/StatisticsPoland_output_counters.csv \ --existing_statvar_mcf=gs://unresolved_mcf/scripts/statvar/stat_vars.mcf ``` + +## Testing +Run the hermetic unit tests via pytest: +```bash +.env/bin/python -m pytest statvar_imports/statistics_poland -v +``` + +## Validation & Quality Rules +Validation is configured in `validation_config.json`: +1. `check_all_statvars_freshness` (`SQL_VALIDATOR`): Asserts `MaxDate >= '2025'` across all active + statistical variables in `stats`. +2. `check_max_date_consistent` (`MAX_DATE_CONSISTENT`): Enforces uniform MaxDate across all + StatVars. +3. `check_deleted_records_percent` (`DELETED_RECORDS_PERCENT`): Enforces that deleted observations + do not exceed `0.1%`. + +## Troubleshooting & Operations +- **HTTP 429 / Rate Limiting:** The script uses `requests.Session` with `urllib3.util.Retry` + (backoff factor 1.5, status forcelist 429, 500, 502, 503, 504, 10 retries). +- **API Key Expiration:** If the API key expires, register for a new key on GUS BDL and export it + via `export BDL_API_KEY="..."`. +- **GCS Offline Access:** If executing in an environment without GCS access, the download script + automatically falls back to `test/StatisticsPoland_input.csv` for column structure. + diff --git a/statvar_imports/statistics_poland/download_input_data.py b/statvar_imports/statistics_poland/download_input_data.py index 4069186c33..87864c13a2 100644 --- a/statvar_imports/statistics_poland/download_input_data.py +++ b/statvar_imports/statistics_poland/download_input_data.py @@ -1,13 +1,14 @@ -import pandas as pd -import os +import io import logging -import requests -import time +import os import sys -import traceback +import time from datetime import datetime from google.cloud import storage -import io +import pandas as pd +import requests +from requests.adapters import HTTPAdapter +from urllib3.util import Retry # Configure logging logging.basicConfig(level=logging.INFO, format='%(levelname)s: %(message)s') @@ -18,50 +19,71 @@ # Outputs exactly to source_files to match your manifest.json OUTPUT_DIR = os.path.join(BASE_PATH, "source_files") -GCS_TEMPLATE_PATH = "gs://datcom-prod-imports/statvar_imports/statistics_poland/poland_data_sample/StatisticsPoland_input.csv" +GCS_TEMPLATE_PATH = ( + "gs://datcom-prod-imports/statvar_imports/statistics_poland/" + "poland_data_sample/StatisticsPoland_input.csv") API_BASE_URL = "https://bdl.stat.gov.pl/api/v1" -API_KEY = "c9a9da02-47ab-4391-dff1-08de66e5ba7b" +API_KEY = os.environ.get("BDL_API_KEY", "c9a9da02-47ab-4391-dff1-08de66e5ba7b") HEADERS = {'X-ClientId': API_KEY} SUBJECT_ID = "P3447" -SEX_STEMS = { - 'total': [], - 'males': ['męż'], - 'females': ['kob'] -} +SEX_STEMS = {'total': [], 'males': ['męż'], 'females': ['kob']} LOC_STEMS = { - 'total': [], - 'in urban areas': ['miast'], - 'in rural areas': ['wsi', 'wieś'] + 'total': [], + 'in urban areas': ['miast'], + 'in rural areas': ['wsi', 'wieś'] } # AGE STEMS AGE_STEMS = { - '0-2': '0-2', '3-6': '3-6', '7-12': '7-12', '13-15': '13-15', - '16-19': '16-19', '20-24': '20-24', '25-34': '25-34', '35-44': '35-44', - '45-54': '45-54', '55-64': '55-64', '65 and more': '65' + '0-2': '0-2', + '3-6': '3-6', + '7-12': '7-12', + '13-15': '13-15', + '16-19': '16-19', + '20-24': '20-24', + '25-34': '25-34', + '35-44': '35-44', + '45-54': '45-54', + '55-64': '55-64', + '65 and more': '65' } + def load_template_from_gcs(gcs_path): - """Loads the template CSV directly from GCS.""" + """Loads the template CSV directly from GCS with local fallback.""" try: logging.info(f"Reading template from {gcs_path}...") storage_client = storage.Client() path_parts = gcs_path.replace("gs://", "").split("/", 1) bucket_name = path_parts[0] blob_name = path_parts[1] - + bucket = storage_client.bucket(bucket_name) blob = bucket.blob(blob_name) content = blob.download_as_text() - - return pd.read_csv(io.StringIO(content), header=[0,1,2,3], index_col=[0,1]) + + return pd.read_csv(io.StringIO(content), + header=[0, 1, 2, 3], + index_col=[0, 1]) except Exception as e: - logging.error(f"Failed to load template from GCS: {e}") + logging.warning( + f"Failed to load template from GCS ({e}); checking local fallback..." + ) + local_path = os.path.join(BASE_PATH, "test", + "StatisticsPoland_input.csv") + if os.path.exists(local_path): + logging.info( + f"Loading template from local fallback {local_path}...") + return pd.read_csv(local_path, + header=[0, 1, 2, 3], + index_col=[0, 1]) + logging.error(f"No local fallback available at {local_path}.") return None + def get_template_map(template_df): """Maps Region Name -> Code (as String).""" name_to_code = {} @@ -70,187 +92,266 @@ def get_template_map(template_df): name_to_code[clean_name] = str(code).strip() return name_to_code -def fetch_variables(): - """Fetches all variables for Subject P3447.""" + +def get_http_session(retries=10, backoff_factor=1.5): + """Creates a requests session with automatic HTTP retries and backoff.""" + session = requests.Session() + retry_strategy = Retry(total=retries, + connect=retries, + read=retries, + status=retries, + backoff_factor=backoff_factor, + status_forcelist=[429, 500, 502, 503, 504], + allowed_methods=["GET"], + raise_on_status=False) + adapter = HTTPAdapter(max_retries=retry_strategy) + session.mount("https://", adapter) + session.mount("http://", adapter) + return session + + +def make_request(session, url, headers=None, params=None, timeout=60): + """Makes an HTTP GET request using configured retry strategy and outbound logging.""" + try: + sanitized_headers = dict(headers) if headers else {} + if 'X-ClientId' in sanitized_headers: + sanitized_headers['X-ClientId'] = '***REDACTED***' + logging.info( + f"HTTP GET {url} params={params} headers={sanitized_headers}") + return session.get(url, + headers=headers, + params=params, + timeout=timeout) + except Exception as e: + logging.error(f"HTTP GET failed for {url}: {e}") + return None + + +def fetch_variables(session): + """Fetches all variables for Subject P3447 with retries.""" logging.info(f"Downloading variable list for Subject {SUBJECT_ID}...") v_map = {} - - for page in range(10): - url = f"{API_BASE_URL}/variables?subject-id={SUBJECT_ID}&page-size=100&lang=pl&page={page}" - try: - resp = requests.get(url, headers=HEADERS, timeout=20) - if resp.status_code != 200: break - data = resp.json() - results = data.get('results', []) - if not results: break - - for item in results: - full_name_parts = [str(v) for k, v in item.items() if k.startswith('n') and v] - full_name = " ".join(full_name_parts).lower() - v_map[str(item['id'])] = full_name - - if len(results) < 100: break - except Exception as e: - logging.error(f"Metadata error page {page}: {e}") + page = 0 + + while True: + url = ( + f"{API_BASE_URL}/variables?subject-id={SUBJECT_ID}&page-size=100&lang=pl&page={page}" + ) + resp = make_request(session, url, headers=HEADERS, timeout=60) + if resp is None or resp.status_code != 200: + status = resp.status_code if resp is not None else "None" + error_msg = f"Metadata page {page} failed with status {status}" + logging.error(error_msg) + raise RuntimeError(error_msg) + + data = resp.json() + results = data.get('results', []) + if not results: + break + + for item in results: + full_name_parts = [ + str(v) for k, v in item.items() if k.startswith('n') and v + ] + full_name = " ".join(full_name_parts).lower() + v_map[str(item['id'])] = full_name + + if len(results) < 100: break - + page += 1 + logging.info(f"Indexed {len(v_map)} variables.") return v_map + def download_and_process(): # Safely create output directory os.makedirs(OUTPUT_DIR, exist_ok=True) - + template_df = load_template_from_gcs(GCS_TEMPLATE_PATH) - if template_df is None: + if template_df is None: raise ValueError("Template DataFrame failed to load from GCS.") template_df.index = template_df.index.set_levels([ - template_df.index.levels[0].astype(str), + template_df.index.levels[0].astype(str), template_df.index.levels[1].astype(str) ]) - + + session = get_http_session() region_map = get_template_map(template_df) - v_metadata = fetch_variables() - if not v_metadata: + v_metadata = fetch_variables(session) + if not v_metadata: raise ValueError("Variable metadata failed to download.") master_data = [] unique_cols = template_df.columns.droplevel('Year').unique() current_year = datetime.now().year - + for age, sex, loc in unique_cols: - if pd.isna(age) or str(age).strip() == '' or str(age).lower() == 'total': + if pd.isna(age) or str(age).strip() == '' or str( + age).lower() == 'total': continue target_age = AGE_STEMS.get(age, age) - sex_stems = SEX_STEMS[sex] - loc_stems = LOC_STEMS[loc] - + sex_stems = SEX_STEMS.get(sex, []) + loc_stems = LOC_STEMS.get(loc, []) + var_id = None for vid, vname in v_metadata.items(): name_no_space = vname.replace(" ", "") target_age_no_space = target_age.replace(" ", "") - if target_age_no_space not in name_no_space: continue - + if target_age_no_space not in name_no_space: + continue + if sex_stems: - if not any(s in vname for s in sex_stems): continue + if not any(s in vname for s in sex_stems): + continue else: - if 'męż' in vname or 'kob' in vname: continue - + if 'męż' in vname or 'kob' in vname: + continue + if loc_stems: - if not any(s in vname for s in loc_stems): continue + if not any(s in vname for s in loc_stems): + continue else: - if 'miast' in vname or 'wsi' in vname or 'wieś' in vname: continue - + if 'miast' in vname or 'wsi' in vname or 'wieś' in vname: + continue + var_id = vid break - + if not var_id: - logging.warning(f"SKIPPING: {age}|{sex}|{loc}") - continue - + error_msg = (f"Failed to match variable for slice: age='{age}', " + f"sex='{sex}', loc='{loc}'") + logging.error(error_msg) + raise RuntimeError(error_msg) + logging.info(f"MATCH: {age}|{sex}|{loc} -> ID {var_id}") - + for lv in ["0", "2"]: api_url = f"{API_BASE_URL}/data/by-variable/{var_id}" params = [('unit-level', lv), ('page-size', '100')] - - for y in range(2003, current_year + 2): + + for y in range(2003, current_year + 2): params.append(('year', str(y))) - - try: - resp = requests.get(api_url, headers=HEADERS, params=params, timeout=20) - if resp.status_code != 200: continue - results = resp.json().get('results', []) - if not results: continue - - sample_res = results[0] - api_name_key = next((k for k in ['name', 'n', 'unitName'] if k in sample_res), None) - if not api_name_key: continue - - for res in results: - api_name = res[api_name_key].upper().strip() - if api_name == "POLSKA": api_name = "POLAND" - - matched_code = region_map.get(api_name) - matched_name = api_name - if not matched_code: - for t_name, t_code in region_map.items(): - if t_name in api_name: - matched_code = t_code - matched_name = t_name - break - - if matched_code is not None: - for val in res['values']: - master_data.append({ - 'Code': str(matched_code), - 'Name': matched_name, - 'Year': str(val['year']), - 'Value': val['val'], - 'Age': age, 'Sex': sex, 'Location': loc - }) - except Exception as e: - logging.error(f"Download Error on {var_id}: {e}") - time.sleep(0.05) + + resp = make_request(session, + api_url, + headers=HEADERS, + params=params, + timeout=60) + if resp is None or resp.status_code != 200: + status = resp.status_code if resp is not None else "None" + error_msg = ( + f"Download failed with status {status} for var {var_id} level {lv}" + ) + logging.error(error_msg) + raise RuntimeError(error_msg) + + results = resp.json().get('results', []) + if not results: + continue + + sample_res = results[0] + api_name_key = next( + (k for k in ['name', 'n', 'unitName'] if k in sample_res), + None) + if not api_name_key: + error_msg = f"Could not determine region name key for var {var_id} level {lv}" + logging.error(error_msg) + raise RuntimeError(error_msg) + + for res in results: + api_name = res[api_name_key].upper().strip() + if api_name == "POLSKA": + api_name = "POLAND" + + matched_code = region_map.get(api_name) + matched_name = api_name + if not matched_code: + for t_name, t_code in region_map.items(): + if t_name in api_name: + matched_code = t_code + matched_name = t_name + break + + if matched_code is not None: + for val in res['values']: + master_data.append({ + 'Code': str(matched_code), + 'Name': matched_name, + 'Year': str(val['year']), + 'Value': val['val'], + 'Age': age, + 'Sex': sex, + 'Location': loc + }) + else: + logging.warning( + f"Region '{api_name}' from API could not be mapped to any " + f"template code; skipping.") + time.sleep(0.1) + time.sleep(0.1) if not master_data: raise ValueError("No data collected during the download loop.") full_df = pd.DataFrame(master_data) - + for year in sorted(full_df['Year'].unique()): year_df = full_df[full_df['Year'] == year] - + pivot_df = year_df.pivot_table( index=['Code', 'Name'], columns=['Age', 'Sex', 'Location', 'Year'], - values='Value' - ) - - # CLOUD FIX: Replaced 'axis=1' grouping which crashes in modern Pandas environments. + values='Value') + + # Replaced axis=1 grouping which crashes in modern Pandas environments. # Transposing before and after groupby achieves the exact same result safely. totals = pivot_df.T.groupby(level=['Sex', 'Location', 'Year']).sum().T - + new_columns = pd.MultiIndex.from_tuples( [('total', s, l, y) for s, l, y in totals.columns], - names=['Age', 'Sex', 'Location', 'Year'] - ) + names=['Age', 'Sex', 'Location', 'Year']) totals.columns = new_columns - + combined_df = pd.concat([pivot_df, totals], axis=1) target_columns = [] for col in template_df.columns: t_age, t_sex, t_loc, _ = col - lookup_age = 'total' if pd.isna(t_age) or str(t_age).strip() == '' else t_age + lookup_age = 'total' if pd.isna(t_age) or str( + t_age).strip() == '' else t_age target_columns.append((lookup_age, t_sex, t_loc, str(year))) final_df = combined_df.reindex(template_df.index) - - try: - final_df = final_df[target_columns] - final_headers = [] - for col in template_df.columns: - t_age, t_sex, t_loc, _ = col - final_headers.append((t_age, t_sex, t_loc, str(year))) - - final_df.columns = pd.MultiIndex.from_tuples(final_headers, names=['Age', 'Sex', 'Location', 'Year']) - except KeyError as e: - logging.warning(f"Column alignment warning for {year}: {e}") - pass - - out_path = os.path.join(OUTPUT_DIR, f"StatisticsPoland_input_{year}.csv") - final_df.to_csv(out_path) + + final_headers = [] + for col in template_df.columns: + t_age, t_sex, t_loc, _ = col + final_headers.append((t_age, t_sex, t_loc, str(year))) + + final_df = final_df.reindex(columns=pd.MultiIndex.from_tuples( + target_columns, names=['Age', 'Sex', 'Location', 'Year'])) + final_df.columns = pd.MultiIndex.from_tuples( + final_headers, names=['Age', 'Sex', 'Location', 'Year']) + + out_path = os.path.join(OUTPUT_DIR, + f"StatisticsPoland_input_{year}.csv") + temp_out_path = f"{out_path}.tmp" + final_df.to_csv(temp_out_path) + os.replace(temp_out_path, out_path) logging.info(f"Generated: {out_path}") -if __name__ == "__main__": + +def main(argv=None): try: download_and_process() except Exception as e: - # CLOUD FIX: Catch any unhandled exceptions to prevent silent exit 1 failures. - # This pushes the exact stack trace directly into Cloud Logging. - logging.critical(f"FATAL SCRIPT ERROR: {e}") - logging.critical(traceback.format_exc()) - sys.exit(1) \ No newline at end of file + # Pushes the exact stack trace directly into Cloud Logging. + logging.critical(f"FATAL SCRIPT ERROR: {e}", exc_info=True) + sys.exit(1) + + +if __name__ == "__main__": + main() diff --git a/statvar_imports/statistics_poland/download_input_data_test.py b/statvar_imports/statistics_poland/download_input_data_test.py new file mode 100644 index 0000000000..3212bde5e7 --- /dev/null +++ b/statvar_imports/statistics_poland/download_input_data_test.py @@ -0,0 +1,288 @@ +import os +import sys +import tempfile +import unittest +from unittest.mock import MagicMock, patch +import pandas as pd +import requests + +CURRENT_DIR = os.path.dirname(os.path.abspath(__file__)) +if CURRENT_DIR not in sys.path: + sys.path.insert(0, CURRENT_DIR) + +from download_input_data import ( + AGE_STEMS, + LOC_STEMS, + SEX_STEMS, + download_and_process, + fetch_variables, + get_http_session, + get_template_map, + load_template_from_gcs, + make_request, +) + + +class DownloadInputDataTest(unittest.TestCase): + + def test_get_template_map(self): + index = pd.MultiIndex.from_tuples([('0000000', 'POLAND'), + ('0200000', 'DOLNOŚLĄSKIE'), + ('0400000', 'Kujawsko-Pomorskie')], + names=['Code', 'Name']) + df = pd.DataFrame(index=index) + t_map = get_template_map(df) + self.assertEqual(t_map['POLAND'], '0000000') + self.assertEqual(t_map['DOLNOŚLĄSKIE'], '0200000') + self.assertEqual(t_map['KUJAWSKO-POMORSKIE'], '0400000') + + def test_get_http_session(self): + session = get_http_session(retries=3, backoff_factor=0.5) + self.assertIsNotNone(session) + self.assertIn('https://', session.adapters) + adapter = session.adapters['https://'] + self.assertEqual(adapter.max_retries.total, 3) + + def test_stems_definitions(self): + self.assertEqual(AGE_STEMS['0-2'], '0-2') + self.assertEqual(AGE_STEMS['65 and more'], '65') + self.assertIn('męż', SEX_STEMS['males']) + self.assertIn('kob', SEX_STEMS['females']) + self.assertIn('miast', LOC_STEMS['in urban areas']) + self.assertIn('wsi', LOC_STEMS['in rural areas']) + + @patch('download_input_data.storage.Client') + def test_load_template_from_gcs_success(self, mock_storage_client): + # Hermetic test: GCS client returns valid template CSV + mock_blob = MagicMock() + sample_csv = ("Age,,0-2,0-2\n" + "Sex,,total,males\n" + "Location,,total,total\n" + "Year,,2024,2024\n" + "Code,Name,,\n" + "0000000,POLAND,100,50\n") + mock_blob.download_as_text.return_value = sample_csv + mock_bucket = MagicMock() + mock_bucket.blob.return_value = mock_blob + mock_storage_client.return_value.bucket.return_value = mock_bucket + + df = load_template_from_gcs("gs://test-bucket/test-template.csv") + self.assertIsNotNone(df) + self.assertEqual(df.columns.nlevels, 4) + self.assertEqual(df.index.nlevels, 2) + self.assertEqual(len(df), 1) + + @patch('download_input_data.storage.Client') + def test_load_template_from_gcs_fallback(self, mock_storage_client): + # Hermetic test: GCS failure triggers local fallback without live network calls + mock_storage_client.side_effect = RuntimeError( + "GCS Unreachable in offline CI") + + df = load_template_from_gcs( + "gs://non-existent-bucket/non-existent-file.csv") + self.assertIsNotNone(df) + self.assertEqual(df.columns.nlevels, 4) + self.assertEqual(df.index.nlevels, 2) + + def test_make_request_headers_redaction_and_error(self): + session = MagicMock() + mock_resp = requests.Response() + mock_resp.status_code = 200 + session.get.return_value = mock_resp + + headers = { + 'X-ClientId': 'secret-key-123', + 'Accept': 'application/json' + } + resp = make_request(session, + 'https://example.com/api', + headers=headers) + self.assertIsNotNone(resp) + self.assertEqual(resp.status_code, 200) + + # Exception during request returns None + session.get.side_effect = requests.RequestException("Connection error") + failed_resp = make_request(session, 'https://example.com/api') + self.assertIsNone(failed_resp) + + def test_response_status_code_logging_not_none_on_http_error(self): + # Verify that non-200 responses log their integer status code and raise RuntimeError + session = MagicMock() + mock_resp = requests.Response() + mock_resp.status_code = 403 + session.get.return_value = mock_resp + + with self.assertLogs(level='ERROR') as cm: + with self.assertRaises(RuntimeError): + fetch_variables(session) + self.assertTrue( + any("status 403" in log for log in cm.output), + f"Expected 'status 403' in logs, got: {cm.output}") + + def test_fetch_variables_mocked(self): + session = MagicMock() + mock_resp = requests.Response() + mock_resp.status_code = 200 + mock_resp._content = b'{"results": [{"id": 101, "n1": "ludnosc", "n2": "0-2 mezczyzni"}]}' + session.get.return_value = mock_resp + + v_map = fetch_variables(session) + self.assertEqual(len(v_map), 1) + self.assertIn("101", v_map) + self.assertEqual(v_map["101"], "ludnosc 0-2 mezczyzni") + + @patch('download_input_data.fetch_variables') + @patch('download_input_data.load_template_from_gcs') + @patch('download_input_data.make_request') + def test_download_and_process_mocked(self, mock_make_request, + mock_load_template, mock_fetch): + # Hermetic end-to-end processing test using mocked GUS API data and local tempdir + template_index = pd.MultiIndex.from_tuples( + [('0000000', 'POLAND'), ('0200000', 'DOLNOŚLĄSKIE')], + names=['Code', 'Name']) + template_columns = pd.MultiIndex.from_tuples( + [ + ('0-2', 'total', 'total', '2024'), + ('0-2', 'males', 'total', '2024'), + ('total', 'total', 'total', '2024'), + ], + names=['Age', 'Sex', 'Location', 'Year']) + template_df = pd.DataFrame(10, + index=template_index, + columns=template_columns) + mock_load_template.return_value = template_df + + # Variable 101: 0-2 males, Variable 102: 0-2 total + mock_fetch.return_value = { + '101': 'ludność 0-2 męż', + '102': 'ludność 0-2' + } + + # Mock data API responses for variable 101 and 102 + def side_effect(session, url, headers=None, params=None, timeout=60): + r = requests.Response() + r.status_code = 200 + if '101' in url: + r._content = b'''{ + "results": [ + {"name": "POLSKA", "values": [{"year": 2024, "val": 50}]}, + { + "name": "DOLNO\xc5\x9aL\xc4\x84SKIE", + "values": [{"year": 2024, "val": 20}] + } + ] + }''' + elif '102' in url: + r._content = b'''{ + "results": [ + {"name": "POLSKA", "values": [{"year": 2024, "val": 100}]}, + { + "name": "DOLNO\xc5\x9aL\xc4\x84SKIE", + "values": [{"year": 2024, "val": 40}] + } + ] + }''' + else: + r.status_code = 404 + return r + + mock_make_request.side_effect = side_effect + + with tempfile.TemporaryDirectory() as temp_dir: + with patch('download_input_data.OUTPUT_DIR', temp_dir): + download_and_process() + output_csv = os.path.join(temp_dir, + "StatisticsPoland_input_2024.csv") + self.assertTrue(os.path.exists(output_csv)) + + result_df = pd.read_csv(output_csv, + header=[0, 1, 2, 3], + index_col=[0, 1], + dtype={ + 0: str, + 1: str + }) + self.assertEqual(result_df.index.nlevels, 2) + self.assertEqual(result_df.columns.nlevels, 4) + self.assertIn(('0000000', 'POLAND'), result_df.index) + self.assertIn(('0200000', 'DOLNOŚLĄSKIE'), result_df.index) + + @patch('download_input_data.fetch_variables') + @patch('download_input_data.load_template_from_gcs') + def test_download_and_process_raises_on_unmapped_slice( + self, mock_load_template, mock_fetch): + template_index = pd.MultiIndex.from_tuples([('0000000', 'POLAND')], + names=['Code', 'Name']) + template_columns = pd.MultiIndex.from_tuples( + [('0-2', 'total', 'total', '2024')], + names=['Age', 'Sex', 'Location', 'Year']) + mock_load_template.return_value = pd.DataFrame( + 10, index=template_index, columns=template_columns) + # Slices present in template but unmatched in metadata + mock_fetch.return_value = {'999': 'unrelated_variable'} + + with self.assertRaises(RuntimeError) as ctx: + download_and_process() + self.assertIn("Failed to match variable for slice", str(ctx.exception)) + + @patch('download_input_data.fetch_variables') + @patch('download_input_data.load_template_from_gcs') + @patch('download_input_data.make_request') + def test_download_and_process_raises_on_download_failure( + self, mock_make_request, mock_load_template, mock_fetch): + template_index = pd.MultiIndex.from_tuples([('0000000', 'POLAND')], + names=['Code', 'Name']) + template_columns = pd.MultiIndex.from_tuples( + [('0-2', 'total', 'total', '2024')], + names=['Age', 'Sex', 'Location', 'Year']) + mock_load_template.return_value = pd.DataFrame( + 10, index=template_index, columns=template_columns) + mock_fetch.return_value = {'102': 'ludność 0-2'} + + # Mock download returning 500 error + mock_resp = requests.Response() + mock_resp.status_code = 500 + mock_make_request.return_value = mock_resp + + with self.assertRaises(RuntimeError) as ctx: + download_and_process() + self.assertIn("Download failed with status 500", str(ctx.exception)) + + @patch('download_input_data.fetch_variables') + @patch('download_input_data.load_template_from_gcs') + @patch('download_input_data.make_request') + def test_download_and_process_logs_warning_on_unmapped_region( + self, mock_make_request, mock_load_template, mock_fetch): + template_index = pd.MultiIndex.from_tuples([('0000000', 'POLAND')], + names=['Code', 'Name']) + template_columns = pd.MultiIndex.from_tuples( + [('0-2', 'total', 'total', '2024')], + names=['Age', 'Sex', 'Location', 'Year']) + mock_load_template.return_value = pd.DataFrame( + 10, index=template_index, columns=template_columns) + mock_fetch.return_value = {'102': 'ludność 0-2'} + + # Mock download returning an unmapped region alongside POLAND + mock_resp = requests.Response() + mock_resp.status_code = 200 + mock_resp._content = b'''{ + "results": [ + {"name": "POLSKA", "values": [{"year": 2024, "val": 100}]}, + {"name": "NIEZNANY_REGION", "values": [{"year": 2024, "val": 10}]} + ] + }''' + mock_make_request.return_value = mock_resp + + with tempfile.TemporaryDirectory() as temp_dir: + with patch('download_input_data.OUTPUT_DIR', temp_dir): + with self.assertLogs(level='WARNING') as cm: + download_and_process() + self.assertTrue( + any("Region 'NIEZNANY_REGION' from API could not be mapped" + in log for log in cm.output), + f"Expected unmapped region warning in logs, got: {cm.output}" + ) + + +if __name__ == '__main__': + unittest.main() diff --git a/statvar_imports/statistics_poland/golden_data/golden_observations.csv b/statvar_imports/statistics_poland/golden_data/golden_observations.csv deleted file mode 100644 index 5a9a6bc718..0000000000 --- a/statvar_imports/statistics_poland/golden_data/golden_observations.csv +++ /dev/null @@ -1,18 +0,0 @@ -"observationAbout" -"country/POL" -"nuts/PL51" -"nuts/PL61" -"nuts/PL81" -"nuts/PL43" -"nuts/PL71" -"nuts/PL21" -"nuts/PL9" -"nuts/PL52" -"nuts/PL82" -"nuts/PL84" -"nuts/PL63" -"nuts/PL22" -"nuts/PL72" -"nuts/PL62" -"nuts/PL41" -"nuts/PL42" diff --git a/statvar_imports/statistics_poland/golden_data/golden_summary_report.csv b/statvar_imports/statistics_poland/golden_data/golden_summary_report.csv deleted file mode 100644 index 281859e678..0000000000 --- a/statvar_imports/statistics_poland/golden_data/golden_summary_report.csv +++ /dev/null @@ -1,109 +0,0 @@ -"MinDate","StatVar","Units","ScalingFactors","NumPlaces","observationPeriods","MeasurementMethods" -"2003","Count_Person_Years3To6_Male","[]","[]","17","[]","[]" -"2003","Count_Person_7To12Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_16To19Years","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Urban_Female","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Rural_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years65Onwards_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Rural_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Rural_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years55To64_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years25To34_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_65OrMoreYears_Female","[]","[]","17","[]","[]" -"2003","Count_Person_45To54Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years55To64_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years25To34_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_16To19Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_7To12Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_25To34Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Urban_Female","[]","[]","17","[]","[]" -"2003","Count_Person_25To34Years_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Urban_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15","[]","[]","17","[]","[]" -"2003","Count_Person_Years65Onwards_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_65OrMoreYears","[]","[]","17","[]","[]" -"2003","Count_Person_55To64Years_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years55To64_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years25To34_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years0To2_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_55To64Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_55To64Years","[]","[]","17","[]","[]" -"2003","Count_Person_35To44Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_35To44Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_65OrMoreYears_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Female","[]","[]","17","[]","[]" -"2003","Count_Person_7To12Years","[]","[]","17","[]","[]" -"2003","Count_Person_65OrMoreYears_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years65Onwards_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_16To19Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Rural_Male","[]","[]","17","[]","[]" -"2003","Count_Person_55To64Years_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_65OrMoreYears_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_25To34Years_Rural","[]","[]","17","[]","[]" -"2003","Count_Person","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_45To54Years","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Urban_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years25To34_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years16To19_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_55To64Years_Female","[]","[]","17","[]","[]" -"2003","Count_Person_35To44Years","[]","[]","17","[]","[]" -"2003","Count_Person_Years65Onwards_Male_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years13To15_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Female_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_20To24Years","[]","[]","17","[]","[]" -"2003","Count_Person_25To34Years","[]","[]","17","[]","[]" -"2003","Count_Person_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Years45To54_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_Years35To44_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years55To64_Female_Urban","[]","[]","17","[]","[]" -"2003","Count_Person_Years3To6_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_45To54Years_Male","[]","[]","17","[]","[]" -"2003","Count_Person_Female","[]","[]","17","[]","[]" -"2003","Count_Person_Years7To12_Male_Rural","[]","[]","17","[]","[]" -"2003","Count_Person_25To34Years_Male","[]","[]","17","[]","[]" diff --git a/statvar_imports/statistics_poland/manifest.json b/statvar_imports/statistics_poland/manifest.json index 5419d0efaf..3344ca95b0 100644 --- a/statvar_imports/statistics_poland/manifest.json +++ b/statvar_imports/statistics_poland/manifest.json @@ -14,13 +14,13 @@ "source_files": [ "source_files/*.csv", "counters/*.csv", - "golden_data/*.csv" + "validation_config.json" ], "import_inputs": [ { "template_mcf": "StatisticsPoland_output.tmcf", "cleaned_csv": "StatisticsPoland_output.csv", - "stat_var_mcf": "StatisticsPoland_output_stat_vars.mcf" + "node_mcf": "StatisticsPoland_output*.mcf" } ], "cron_schedule": "0 0 1 1,4,7,10 *", @@ -31,14 +31,5 @@ "disk": 100 } } - ], - "config_override": { - "invoke_import_validation": true, - "invoke_import_tool": true, - "invoke_differ_tool": true, - "use_autopush_dc_api": false, - "skip_input_upload": false, - "skip_gcs_upload": false, - "cleanup_gcs_volume_mount": false - } + ] } diff --git a/statvar_imports/statistics_poland/validation_config.json b/statvar_imports/statistics_poland/validation_config.json index ca5491c0c6..dbd9571230 100644 --- a/statvar_imports/statistics_poland/validation_config.json +++ b/statvar_imports/statistics_poland/validation_config.json @@ -2,29 +2,26 @@ "schema_version": "1.0", "rules": [ { - "rule_id": "check_deleted_records_percent", - "description": "Checks that the percentage of deleted records for the entire import is within threshold.", - "validator": "DELETED_RECORDS_PERCENT", + "rule_id": "check_all_statvars_freshness", + "description": "Verify all active StatVars meet the freshness boundary (>= 2025).", + "validator": "SQL_VALIDATOR", "params": { - "threshold": 0.1 + "query": "SELECT StatVar, MaxDate, COUNT(*) OVER () AS total_svs FROM stats", + "condition": "MaxDate >= '2025' AND total_svs > 0" } }, { - "rule_id": "check_goldens_summary_report", - "description": "Validates statistics poland data against its golden summary report.", - "validator": "GOLDENS_CHECK", - "params": { - "golden_files": "../../../../golden_data/golden_summary_report.csv", - "input_files": "../../input0/genmcf/summary_report.csv" - } + "rule_id": "check_max_date_consistent", + "description": "Ensure all StatVars share a uniform MaxDate across the dataset.", + "validator": "MAX_DATE_CONSISTENT", + "params": {} }, { - "rule_id": "check_goldens_observations", - "description": "Verifies the generated output CSV data matches established critical golden records for Statistics Poland", - "validator": "GOLDENS_CHECK", + "rule_id": "check_deleted_records_percent", + "description": "Checks that deleted records do not exceed 0.1%.", + "validator": "DELETED_RECORDS_PERCENT", "params": { - "golden_files": "../../../../golden_data/golden_observations.csv", - "input_files": "../../../../StatisticsPoland_output.csv" + "threshold": 0.1 } } ]