From b75f8ce9fa03bd5051472c2ba84b9d2f1eb58268 Mon Sep 17 00:00:00 2001 From: Stephen Thompson Date: Wed, 16 Sep 2026 14:41:00 +0100 Subject: [PATCH 1/2] Tidy up file structure and naming conventions --- src/csv_writer.py | 4 ++-- src/locations.py | 1 - src/pipeline/Snakefile | 3 +-- src/pipeline/utils.py | 3 +-- tests/helpers.py | 2 +- tests/test_snakemake_integration.py | 7 ++++--- 6 files changed, 9 insertions(+), 11 deletions(-) diff --git a/src/csv_writer.py b/src/csv_writer.py index c5529f4..ef2014c 100644 --- a/src/csv_writer.py +++ b/src/csv_writer.py @@ -9,7 +9,7 @@ from locations import ( WAVEFORM_ORIGINAL_CSV, - WAVEFORM_PSEUDONYMISED_EHR, + WAVEFORM_PSEUDONYMISED_PARQUET, make_file_name, FILE_STEM_PATTERN, EHR_STEM_PATTERN_HASHED, @@ -113,7 +113,7 @@ def write_ehr( """ subs_dict = dict(date=date_str, hashed_csn=hashed_csn) stem = make_file_name(EHR_STEM_PATTERN_HASHED, subs_dict) - filename = WAVEFORM_PSEUDONYMISED_EHR / f"{stem}_ehr.csv" + filename = WAVEFORM_PSEUDONYMISED_PARQUET / f"{stem}.ehr.csv" filename.parent.mkdir(exist_ok=True, parents=True) df.to_csv(filename, index=False) diff --git a/src/locations.py b/src/locations.py index 53c4e94..6430aec 100644 --- a/src/locations.py +++ b/src/locations.py @@ -5,7 +5,6 @@ WAVEFORM_ORIGINAL_PARQUET = WAVEFORM_EXPORT_BASE / "original-parquet" WAVEFORM_HASH_LOOKUPS = WAVEFORM_EXPORT_BASE / "hash-lookups" WAVEFORM_PSEUDONYMISED_PARQUET = WAVEFORM_EXPORT_BASE / "pseudonymised" -WAVEFORM_PSEUDONYMISED_EHR = WAVEFORM_EXPORT_BASE / "pseudonymised_ehr" WAVEFORM_SNAKEMAKE_LOGS = WAVEFORM_EXPORT_BASE / "snakemake-logs" WAVEFORM_FTPS_LOGS = WAVEFORM_EXPORT_BASE / "ftps-logs" diff --git a/src/pipeline/Snakefile b/src/pipeline/Snakefile index 97d6b40..a9d7fd1 100644 --- a/src/pipeline/Snakefile +++ b/src/pipeline/Snakefile @@ -8,7 +8,6 @@ from locations import ( WAVEFORM_ORIGINAL_CSV, WAVEFORM_SNAKEMAKE_LOGS, WAVEFORM_PSEUDONYMISED_PARQUET, - WAVEFORM_PSEUDONYMISED_EHR, HASH_LOOKUP_JSON, HASH_LOOKUP_JSON_REL, FILE_STEM_PATTERN, @@ -144,7 +143,7 @@ rule ehr_lookup: # and reruns if the underlying data for this csn/day changes. pseudonymised_parquets = pseudonymised_parquet_files_for_date_and_hashed_csn output: - WAVEFORM_PSEUDONYMISED_EHR / (EHR_STEM_PATTERN_HASHED + "_ehr.csv") + WAVEFORM_PSEUDONYMISED_PARQUET / (EHR_STEM_PATTERN_HASHED + ".ehr.csv") log: WAVEFORM_SNAKEMAKE_LOGS / "ehr_lookup" / (EHR_STEM_PATTERN_HASHED + ".log") run: diff --git a/src/pipeline/utils.py b/src/pipeline/utils.py index b1175e9..5a92ef9 100644 --- a/src/pipeline/utils.py +++ b/src/pipeline/utils.py @@ -9,7 +9,6 @@ from pseudon.hashing import do_hash from locations import ( WAVEFORM_PSEUDONYMISED_PARQUET, - WAVEFORM_PSEUDONYMISED_EHR, HASH_LOOKUP_JSON, ORIGINAL_PARQUET_PATTERN, FILE_STEM_PATTERN_HASHED, @@ -77,7 +76,7 @@ def get_daily_hash_lookup(self) -> Path: def get_ehr_lookup(self) -> Path: final_stem = make_file_name(EHR_STEM_PATTERN_HASHED, self._subs_dict) - return WAVEFORM_PSEUDONYMISED_EHR / f"{final_stem}_ehr.csv" + return WAVEFORM_PSEUDONYMISED_PARQUET / f"{final_stem}.ehr.csv" def get_file_age(file_path: Path) -> timedelta: diff --git a/tests/helpers.py b/tests/helpers.py index 18b720a..48159a2 100644 --- a/tests/helpers.py +++ b/tests/helpers.py @@ -70,7 +70,7 @@ def get_pseudon_parquet(self): return f"{self.date}/{self.date}.{self.get_hashed_csn()}.{self.variable_id}.{self.channel_id}.{self.units}.parquet" def get_pseudon_ehr(self): - return f"{self.date}/{self.date}.{self.get_hashed_csn()}_ehr.csv" + return f"{self.date}/{self.date}.{self.get_hashed_csn()}.ehr.csv" def get_hashes(self): return f"{self.date}/{self.date}.hashes.json" diff --git a/tests/test_snakemake_integration.py b/tests/test_snakemake_integration.py index 3c53b7b..955d378 100644 --- a/tests/test_snakemake_integration.py +++ b/tests/test_snakemake_integration.py @@ -283,11 +283,9 @@ def test_snakemake_pipeline(tmp_path: Path, background_hasher, monkeypatch): tmp_path / "original-parquet" / filename.get_orig_parquet() ) pseudon_path = tmp_path / "pseudonymised" / filename.get_pseudon_parquet() - ehr_path = tmp_path / "pseudonymised_ehr" / filename.get_pseudon_ehr() assert original_parquet_path.exists() assert pseudon_path.exists() - assert ehr_path.exists() _compare_original_parquet_to_expected(original_parquet_path, expected_data) _compare_parquets(original_parquet_path, pseudon_path) @@ -312,7 +310,10 @@ def test_snakemake_pipeline(tmp_path: Path, background_hasher, monkeypatch): expected_file_counts = {"2025-01-01": 5, "2025-01-02": 1} _assert_date_partitioned_files(tmp_path / "original-csv", expected_file_counts) _assert_date_partitioned_files(tmp_path / "original-parquet", expected_file_counts) - _assert_date_partitioned_files(tmp_path / "pseudonymised", expected_file_counts) + # the pseudonymised files also include ehr files so expected file counts differ + _assert_date_partitioned_files( + tmp_path / "pseudonymised", {"2025-01-01": 7, "2025-01-02": 2} + ) _assert_date_partitioned_files( tmp_path / "hash-lookups", {"2025-01-01": 1, "2025-01-02": 1} ) From 56a1ac66b476bffd5af4d45906c06594d36a202c Mon Sep 17 00:00:00 2001 From: Stephen Thompson Date: Wed, 16 Sep 2026 15:07:58 +0100 Subject: [PATCH 2/2] Fixed test_ehr --- tests/test_ehr.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/test_ehr.py b/tests/test_ehr.py index f30a465..4ba4db4 100644 --- a/tests/test_ehr.py +++ b/tests/test_ehr.py @@ -138,16 +138,16 @@ def mock_get_rows_pg(self, query, params): def test_ehr(monkeypatch, tmp_path): fake_abs_root = tmp_path.absolute() - fake_waveform_pseudonymised_ehr = fake_abs_root / "pseudonymised_ehr" + fake_waveform_pseudonymised_ehr = fake_abs_root / "pseudonymised" monkeypatch.setattr( - csv_writer, "WAVEFORM_PSEUDONYMISED_EHR", fake_waveform_pseudonymised_ehr + csv_writer, "WAVEFORM_PSEUDONYMISED_PARQUET", fake_waveform_pseudonymised_ehr ) ehr_for_csv(date_str="2026-09-14", original_csn="SECRET1234", hashed_csn="fakehash") # just check the file contains something for now (it will be changing to parquet) expected_file = ( - fake_waveform_pseudonymised_ehr / "2026-09-14" / "2026-09-14.fakehash_ehr.csv" + fake_waveform_pseudonymised_ehr / "2026-09-14" / "2026-09-14.fakehash.ehr.csv" ) actual_text = expected_file.read_text() assert actual_text and "SECRET" not in actual_text