From 70ea79eaab18d4fb1add33b76b1a95f3a048c4ff Mon Sep 17 00:00:00 2001 From: Cail Daley Date: Thu, 1 Oct 2026 08:53:42 -0400 Subject: [PATCH 1/3] Write make_cat's catalogue once per save stage, not once per column FITSCatalogue.add_col rebuilt the table HDU from every existing column plus the new one and rewrote the whole file each call. make_cat appended ~243 columns this way, so tile_make_cat spent 10-20 min per tile in astropy writeto on NFS scratch. FITSCatalogue.add_cols appends a dict of columns with one rebuild and one write; add_col is add_cols with one item, and the per-column FITS type, repeat count and TDIM logic lives in _make_fits_col. make_cat uses it for TILE_ID/TILE_UNIQUE_ID, for each SaveCatalogue.process stage (ngmix, PSF slots) and for the MASK_* columns: four writes per tile. tests/module/test_file_io_add_cols.py checks that add_cols writes the same file, byte for byte, as successive add_col calls, for plain and SExtractor-layout (table between other HDUs) catalogues and for int16, int32, int64, float32, bool, string, vector and 2-D columns. Co-Authored-By: Claude Opus 5.5 --- .../modules/make_cat_package/make_cat.py | 14 +- src/shapepipe/pipeline/file_io.py | 159 +++++++++++------- tests/module/test_file_io_add_cols.py | 133 +++++++++++++++ tests/module/test_make_cat.py | 6 +- 4 files changed, 239 insertions(+), 73 deletions(-) create mode 100644 tests/module/test_file_io_add_cols.py diff --git a/src/shapepipe/modules/make_cat_package/make_cat.py b/src/shapepipe/modules/make_cat_package/make_cat.py index e35f7a32b..8a6f9548a 100644 --- a/src/shapepipe/modules/make_cat_package/make_cat.py +++ b/src/shapepipe/modules/make_cat_package/make_cat.py @@ -138,8 +138,9 @@ def save_sextractor_data(final_cat_file, sexcat_path, remove_vignet=True): final_cat_file.save_as_fits(data, ext_name="RESULTS") final_cat_file.open() - final_cat_file.add_col("TILE_ID", tile_id_array) - final_cat_file.add_col("TILE_UNIQUE_ID", unique_id) + final_cat_file.add_cols( + {"TILE_ID": tile_id_array, "TILE_UNIQUE_ID": unique_id} + ) sexcat_file.close() @@ -204,10 +205,11 @@ def save_mask_ext_data(final_cat_file, band_paths, w_log): ra = np.copy(final_cat_file.get_data()["XWIN_WORLD"]) dec = np.copy(final_cat_file.get_data()["YWIN_WORLD"]) + mask_cols = {} for band, path in band_paths.items(): w_log.info(f"Query external mask for band {band}: {path}") - values = mask_query.query_map(path, ra, dec) - final_cat_file.add_col(f"MASK_{band}", values) + mask_cols[f"MASK_{band}"] = mask_query.query_map(path, ra, dec) + final_cat_file.add_cols(mask_cols) final_cat_file.close() @@ -279,9 +281,7 @@ def process( ) if err_msg is None: - - for key in self._output_dict.keys(): - self._final_cat_file.add_col(key, self._output_dict[key]) + self._final_cat_file.add_cols(self._output_dict) self._final_cat_file.close() diff --git a/src/shapepipe/pipeline/file_io.py b/src/shapepipe/pipeline/file_io.py index 35f64bc27..1d309ad33 100644 --- a/src/shapepipe/pipeline/file_io.py +++ b/src/shapepipe/pipeline/file_io.py @@ -1354,7 +1354,7 @@ def add_col( ): """Add Column. - Add a Column to the catalogue. + Add one column to the catalogue; see :meth:`add_cols`. Parameters ---------- @@ -1371,6 +1371,45 @@ def add_col( new_cat_inst : io.FITSCatalogue New catalogue object + """ + self.add_cols( + {col_name: col_data}, + hdu_no=hdu_no, + ext_name=ext_name, + new_cat=new_cat, + new_cat_inst=new_cat_inst, + ) + + def add_cols( + self, + columns, + hdu_no=None, + ext_name=None, + new_cat=False, + new_cat_inst=None, + ): + """Add Columns. + + Append columns to a table HDU, rebuilding the HDU and writing the + file once for all of them. The result is identical to one + :meth:`add_col` call per column in the order of ``columns``, but the + cost is one write of the file instead of one per column, which + dominates when many columns are added to a large catalogue. + + Parameters + ---------- + columns : dict + Column name to column data (``numpy.ndarray``); columns are + appended in the dict's order + hdu_no : int + HDU index + ext_name : str, optional + Change the name of the extansion + new_cat : bool, optional + If true will save the changes into a new catalogue + new_cat_inst : io.FITSCatalogue + New catalogue object + """ if new_cat: open_mode = new_cat_inst.open_mode @@ -1388,86 +1427,80 @@ def add_col( open_mode_needed=FITSCatalogue.OpenMode.ReadWrite, ) - if type(col_data) != np.ndarray: - raise TypeError("col_data must be a numpy.ndarray") + for col_data in columns.values(): + if type(col_data) != np.ndarray: + raise TypeError("col_data must be a numpy.ndarray") + + if not columns: + return if hdu_no is None: hdu_no = self.hdu_no if ext_name is None: ext_name = self._cat_data[hdu_no].name - n_of_hdu = len(self._cat_data) - old_hdu_prev = [] - for i in range(0, hdu_no): - old_hdu_prev.append(self._cat_data[i]) - old_hdu_next = [] - for i in range(hdu_no + 1, n_of_hdu): - old_hdu_next.append(self._cat_data[i]) + col_list = self._cat_data[hdu_no].data.columns + fits.ColDefs( + [ + self._make_fits_col(col_name, col_data) + for col_name, col_data in columns.items() + ] + ) + + new_fits = fits.HDUList( + self._cat_data[:hdu_no] + + [fits.BinTableHDU.from_columns(col_list, name=ext_name)] + + self._cat_data[hdu_no + 1 :] + ) + new_fits.writeto(output_path, overwrite=True) - new_fits = fits.HDUList(old_hdu_prev) + if not new_cat: + self._cat_data.close() + del self._cat_data + self._cat_data = fits.open( + self.fullpath, + mode=self.open_mode, + memmap=self.use_memmap, + ) + + def _make_fits_col(self, col_name, col_data): + """Make FITS Column. + + Build the ``astropy.io.fits.Column`` that :meth:`add_cols` appends + for one array: the FITS type from :meth:`_get_fits_col_type`, a + repeat count and ``TDIM`` for multi-dimensional arrays, and a width + set by the longest entry for strings. + + Parameters + ---------- + col_name : str + Column name + col_data : numpy.ndarray + Column data - col_list = self._cat_data[hdu_no].data.columns + Returns + ------- + astropy.io.fits.Column + The column + """ data_type = self._get_fits_col_type(col_data) data_shape = col_data.shape[1:] - dim = str(tuple(data_shape)) + dim = None mem_size = 1 if len(data_shape) != 0: for k in data_shape: mem_size *= k - data_format = f"{mem_size}{data_type}" - new_col = fits.ColDefs( - [ - fits.Column( - name=col_name, - format=data_format, - array=col_data, - dim=dim, - ) - ] - ) - col_list += new_col + dim = str(tuple(data_shape)) elif data_type == "A": mem_size *= len(max(col_data, key=len)) - data_format = f"{mem_size}{data_type}" - new_col = fits.ColDefs( - [ - fits.Column( - name=col_name, - format=data_format, - array=col_data, - dim=str((mem_size,)), - ) - ] - ) - col_list += new_col - else: - data_format = f"{mem_size}{data_type}" - new_col = fits.ColDefs( - [ - fits.Column( - name=col_name, - format=data_format, - array=col_data, - ) - ] - ) - col_list += new_col - - new_fits.append(fits.BinTableHDU.from_columns(col_list, name=ext_name)) + dim = str((mem_size,)) - new_fits += fits.HDUList(old_hdu_next) - - new_fits.writeto(output_path, overwrite=True) - - if not new_cat: - self._cat_data.close() - del self._cat_data - self._cat_data = fits.open( - self.fullpath, - mode=self.open_mode, - memmap=self.use_memmap, - ) + return fits.Column( + name=col_name, + format=f"{mem_size}{data_type}", + array=col_data, + dim=dim, + ) def remove_col(self, col_index): """Remove Column. diff --git a/tests/module/test_file_io_add_cols.py b/tests/module/test_file_io_add_cols.py new file mode 100644 index 000000000..7915c4483 --- /dev/null +++ b/tests/module/test_file_io_add_cols.py @@ -0,0 +1,133 @@ +"""UNIT TESTS FOR PIPELINE: file_io.FITSCatalogue.add_cols. + +``add_cols`` appends many columns with one rebuild and one write of the +file. Its contract is that the file it writes is the file the same columns +added one ``add_col`` call at a time would produce -- which writes and +reopens the file between every column. These tests build both from the same +base catalogue and compare them HDU by HDU: names, column formats and +``TDIM``, headers, data, and finally the bytes. +""" + +import numpy as np +import numpy.testing as npt +import pytest +from astropy.io import fits + +from shapepipe.pipeline import file_io + +N_OBJ = 7 + + +def _new_columns(): + """Columns of every kind make_cat appends, in a fixed order.""" + rng = np.random.default_rng(1) + return { + "TILE_ID": np.full(N_OBJ, 210.282), + "TILE_UNIQUE_ID": np.arange(N_OBJ, dtype=np.int64) + 10**9, + "FLAGS_I16": rng.integers(-5, 5, N_OBJ).astype(np.int16), + "N_I32": rng.integers(0, 100, N_OBJ).astype(np.int32), + "G1_F32": rng.normal(size=N_OBJ).astype(np.float32), + "MASK_BOOL": rng.integers(0, 2, N_OBJ).astype(bool), + # A per-epoch slot block as a vector column, and a 2-D cell. + "PSF_SLOTS": rng.normal(size=(N_OBJ, 12)), + "COV": rng.normal(size=(N_OBJ, 2, 3)), + "EXP_NAME": np.array([f"2113864-{i}" for i in range(N_OBJ)]), + } + + +def _write_base(path, sex_catalogue): + """A base catalogue; with ``sex_catalogue``, the table is HDU 2 of 4.""" + base = np.empty(N_OBJ, dtype=[("NUMBER", "i4"), ("X", "f8")]) + base["NUMBER"] = np.arange(1, N_OBJ + 1) + base["X"] = np.linspace(0.0, 1.0, N_OBJ) + if not sex_catalogue: + cat = file_io.FITSCatalogue( + str(path), + open_mode=file_io.BaseCatalogue.OpenMode.ReadWrite, + ) + cat.save_as_fits(base, ext_name="RESULTS") + return + + header = fits.BinTableHDU.from_columns( + [fits.Column(name="Field Header Card", format="10A", array=["x"])], + name="LDAC_IMHEAD", + ) + table = fits.BinTableHDU(base, name="LDAC_OBJECTS") + epoch = fits.BinTableHDU( + np.array([(1, 3)], dtype=[("NUMBER", "i4"), ("CCD_N", "i4")]), + name="EPOCH_0", + ) + fits.HDUList([fits.PrimaryHDU(), header, table, epoch]).writeto(path) + + +def _open(path, sex_catalogue): + cat = file_io.FITSCatalogue( + str(path), + open_mode=file_io.BaseCatalogue.OpenMode.ReadWrite, + SEx_catalogue=sex_catalogue, + ) + cat.open() + return cat + + +@pytest.mark.parametrize("sex_catalogue", [False, True]) +def test_add_cols_matches_successive_add_col(tmp_path, sex_catalogue): + """One add_cols call writes the file N add_col calls would.""" + columns = _new_columns() + one_path = tmp_path / "one_by_one.fits" + batch_path = tmp_path / "batched.fits" + + _write_base(one_path, sex_catalogue) + cat = _open(one_path, sex_catalogue) + for name, data in columns.items(): + cat.add_col(name, data) + cat.close() + + _write_base(batch_path, sex_catalogue) + cat = _open(batch_path, sex_catalogue) + cat.add_cols(columns) + cat.close() + + with fits.open(one_path) as ref, fits.open(batch_path) as new: + assert [h.name for h in new] == [h.name for h in ref] + for h_ref, h_new in zip(ref, new): + assert h_new.header == h_ref.header, h_ref.name + if not isinstance(h_ref, fits.BinTableHDU): + continue + assert h_new.columns.names == h_ref.columns.names + assert h_new.columns.formats == h_ref.columns.formats + assert h_new.columns.dims == h_ref.columns.dims + for name in h_ref.columns.names: + npt.assert_array_equal( + h_new.data[name], h_ref.data[name], err_msg=name + ) + + table = new[2 if sex_catalogue else 1] + assert table.columns.names[-len(columns):] == list(columns) + assert table.data["PSF_SLOTS"].shape == (N_OBJ, 12) + npt.assert_array_equal(table.data["FLAGS_I16"], columns["FLAGS_I16"]) + + assert batch_path.read_bytes() == one_path.read_bytes() + + +def test_add_cols_leaves_catalogue_open_and_readable(tmp_path): + """After add_cols the instance is reopened on the written file.""" + path = tmp_path / "cat.fits" + _write_base(path, sex_catalogue=False) + cat = _open(path, sex_catalogue=False) + cat.add_cols({"A": np.arange(N_OBJ), "B": np.ones(N_OBJ)}) + assert cat.get_col_names() == ["NUMBER", "X", "A", "B"] + npt.assert_array_equal(cat.get_data()["A"], np.arange(N_OBJ)) + cat.close() + + +def test_add_cols_rejects_non_array_before_writing(tmp_path): + """A non-array column raises and leaves the file as it was.""" + path = tmp_path / "cat.fits" + _write_base(path, sex_catalogue=False) + before = path.read_bytes() + cat = _open(path, sex_catalogue=False) + with pytest.raises(TypeError): + cat.add_cols({"A": np.arange(N_OBJ), "B": list(range(N_OBJ))}) + cat.close() + assert path.read_bytes() == before diff --git a/tests/module/test_make_cat.py b/tests/module/test_make_cat.py index bdfdb980a..d61a68523 100644 --- a/tests/module/test_make_cat.py +++ b/tests/module/test_make_cat.py @@ -666,7 +666,7 @@ def test_galaxy_cut_admits_only_measured_objects(tmp_path): class _ProcessCatStub(_FinalCatStub): - """FITSCatalogue stand-in for ``SaveCatalogue.process``; records add_col.""" + """FITSCatalogue stand-in for ``SaveCatalogue.process``; records add_cols.""" def __init__(self, obj_id, n_epoch): super().__init__(n_epoch) @@ -682,8 +682,8 @@ def close(self): def get_data(self): return {"NUMBER": self._number, "N_EPOCH": self._n_epoch} - def add_col(self, name, data): - self.cols[name] = data + def add_cols(self, columns): + self.cols.update(columns) def _slot_numbers(out, family): From 8cc724fbf3ea0492662c06ae5717415c9c23a135 Mon Sep 17 00:00:00 2001 From: Cail Daley Date: Thu, 1 Oct 2026 08:56:40 -0400 Subject: [PATCH 2/3] Build make_cat's catalogue node-local and publish it once make_cat rewrites its whole catalogue at each save stage. On NFS scratch one write of a 33K x 271 tile catalogue costs ~19 s against ~0.4 s on node-local disk, so the build belongs node-local even at four writes. make_cat_runner takes an optional WORK_DIR: the catalogue is built there and moved to the run's output directory once complete. A work file left by an earlier attempt is removed first, since save_as_fits appends to an existing file. Without WORK_DIR the catalogue is built in the output directory as before. config_tile_Mc.ini sets WORK_DIR = $SP_LOCAL/make_cat, the per-tile node-local store tile_local() already exports for tile_make_cat and that its EXIT trap removes. tile.smk is unchanged: the published path, the completeness check and the final_cat copy all read the run's output directory as before. Co-Authored-By: Claude Opus 5.5 --- src/shapepipe/modules/make_cat_runner.py | 30 +++++++++++---- tests/module/test_make_cat.py | 47 ++++++++++++++++++++++++ workflow/config/cfis/config_tile_Mc.ini | 5 +++ 3 files changed, 75 insertions(+), 7 deletions(-) diff --git a/src/shapepipe/modules/make_cat_runner.py b/src/shapepipe/modules/make_cat_runner.py index d27b71ef8..468e7e4e1 100644 --- a/src/shapepipe/modules/make_cat_runner.py +++ b/src/shapepipe/modules/make_cat_runner.py @@ -7,6 +7,7 @@ """ import os +import shutil from shapepipe.modules.make_cat_package import make_cat from shapepipe.modules.module_decorator import module_runner @@ -62,9 +63,23 @@ def make_cat_runner( else: n_epoch_slots = None + # The catalogue is built in WORK_DIR when set (e.g. node-local disk: + # each save stage rewrites the whole file) and moved to the run's output + # directory once complete. + if config.has_option(module_config_sec, "WORK_DIR"): + work_dir = config.getexpanded(module_config_sec, "WORK_DIR") + os.makedirs(work_dir, exist_ok=True) + else: + work_dir = run_dirs["output"] + work_path = make_cat.get_output_name(work_dir, file_number_string) + # save_as_fits appends to an existing file: a work file left by an + # earlier attempt must not survive into this one. + if os.path.exists(work_path): + os.remove(work_path) + # Set final output file final_cat_file = make_cat.prepare_final_cat_file( - run_dirs["output"], + work_dir, file_number_string, ) @@ -84,12 +99,7 @@ def make_cat_runner( # If error message: delete (incomplete) output file and raise error if err_msg is not None: - os.remove( - make_cat.get_output_name( - run_dirs["output"], - file_number_string, - ) - ) + os.remove(work_path) #raise ValueError(err_msg) w_log.info(err_msg) @@ -108,4 +118,10 @@ def make_cat_runner( w_log.info("Save external mask data") make_cat.save_mask_ext_data(final_cat_file, band_paths, w_log) + if work_dir != run_dirs["output"] and os.path.exists(work_path): + shutil.move( + work_path, + make_cat.get_output_name(run_dirs["output"], file_number_string), + ) + return None, None diff --git a/tests/module/test_make_cat.py b/tests/module/test_make_cat.py index d61a68523..48c22d435 100644 --- a/tests/module/test_make_cat.py +++ b/tests/module/test_make_cat.py @@ -574,6 +574,53 @@ def test_make_cat_runner_ships_every_detection_unclassified(tmp_path): assert not [name for name in data.dtype.names if "SPREAD" in name] +def test_make_cat_runner_work_dir_publishes_same_catalogue( + tmp_path, monkeypatch +): + """WORK_DIR moves the build elsewhere; the published catalogue is the same. + + The catalogue built in WORK_DIR is moved to the run's output directory + byte-identical to one built there directly, nothing is left in WORK_DIR, + and a work file a previous attempt left behind does not leak into it + (``save_as_fits`` appends to an existing file). + """ + obj_ids = [1, 2, 3] + tile_sexcat_path = tmp_path / "tile_sexcat-350-100.fits" + _write_sex_like_cat(tile_sexcat_path, _numbered_data(obj_ids)) + galaxy_psf_path = tmp_path / "galaxy_psf-350-100.sqlite" + ngmix_path = tmp_path / "ngmix-350-100.fits" + _write_ngmix_cat(ngmix_path, obj_ids) + inputs = [str(tile_sexcat_path), str(galaxy_psf_path), str(ngmix_path)] + + published = {} + for label, section in ( + ("direct", {}), + ("staged", {"WORK_DIR": "$SP_TEST_LOCAL/make_cat"}), + ): + out_dir = tmp_path / label + out_dir.mkdir() + config = CustomParser() + config.read_dict( + {"MAKE_CAT_RUNNER": {"SHAPE_MEASUREMENT_TYPE": "ngmix", **section}} + ) + if section: + local = tmp_path / "local" + monkeypatch.setenv("SP_TEST_LOCAL", str(local)) + (local / "make_cat").mkdir(parents=True) + stale = make_cat.get_output_name(str(local / "make_cat"), "-350-100") + _write_sex_like_cat(stale, _numbered_data([7, 8])) + assert make_cat_runner( + inputs, {"output": str(out_dir)}, "-350-100", config, + "MAKE_CAT_RUNNER", _NullLogger(), + ) == (None, None) + path = make_cat.get_output_name(str(out_dir), "-350-100") + with open(path, "rb") as f: + published[label] = f.read() + + assert published["staged"] == published["direct"] + assert list((tmp_path / "local" / "make_cat").iterdir()) == [] + + @pytest.mark.parametrize("shear", SHEAR_EXTS) @pytest.mark.parametrize("component", [0, 1], ids=["g1", "g2"]) @pytest.mark.parametrize("nonfinite", [np.nan, np.inf, -np.inf]) diff --git a/workflow/config/cfis/config_tile_Mc.ini b/workflow/config/cfis/config_tile_Mc.ini index 636803f48..addb88861 100644 --- a/workflow/config/cfis/config_tile_Mc.ini +++ b/workflow/config/cfis/config_tile_Mc.ini @@ -80,6 +80,11 @@ NUMBERING_SCHEME = -000-000 SHAPE_MEASUREMENT_TYPE = ngmix +# The catalogue is built here and moved to OUTPUT_DIR once complete. Node-local +# ($SP_LOCAL, exported by tile.smk's tile_local() prologue): each save stage +# rewrites the whole file, which is tens of times slower on NFS scratch. +WORK_DIR = $SP_LOCAL/make_cat + # Save per-epoch PSF shapes and epoch identity (HSM_*_PSF_n, EXP_ID_n, # CCD_n). The campaign merge needs one column schema across all tiles, so # the slot count is fixed rather than per-tile max(N_EPOCH)+1. It must From cdd2ffbdb3f03e14a6a7b4a7d429c48f2add2b9d Mon Sep 17 00:00:00 2001 From: Cail Daley Date: Thu, 1 Oct 2026 09:17:24 -0400 Subject: [PATCH 3/3] Test all make_cat save stages through WORK_DIR, and add_cols' new_cat path The runner WORK_DIR test runs ngmix, per-epoch PSF slots and one external mask band, staged and in place, and checks the published bytes match and the stage columns carry the expected values. A new add_cols test covers new_cat/new_cat_inst with hdu_no and ext_name on a SExtractor-layout catalogue: the source is untouched, surrounding HDUs carry over, and the output equals successive add_col calls. Co-Authored-By: Claude Opus 5.5 --- src/shapepipe/pipeline/file_io.py | 4 +- tests/module/test_file_io_add_cols.py | 55 +++++++++++++++++++++++ tests/module/test_make_cat.py | 64 ++++++++++++++++++++++----- 3 files changed, 111 insertions(+), 12 deletions(-) diff --git a/src/shapepipe/pipeline/file_io.py b/src/shapepipe/pipeline/file_io.py index 1d309ad33..565437038 100644 --- a/src/shapepipe/pipeline/file_io.py +++ b/src/shapepipe/pipeline/file_io.py @@ -1365,7 +1365,7 @@ def add_col( hdu_no : int HDU index ext_name : str, optional - Change the name of the extansion + Change the name of the extension new_cat : bool, optional If true will save the changes into a new catalogue new_cat_inst : io.FITSCatalogue @@ -1404,7 +1404,7 @@ def add_cols( hdu_no : int HDU index ext_name : str, optional - Change the name of the extansion + Change the name of the extension new_cat : bool, optional If true will save the changes into a new catalogue new_cat_inst : io.FITSCatalogue diff --git a/tests/module/test_file_io_add_cols.py b/tests/module/test_file_io_add_cols.py index 7915c4483..f35bab0ce 100644 --- a/tests/module/test_file_io_add_cols.py +++ b/tests/module/test_file_io_add_cols.py @@ -131,3 +131,58 @@ def test_add_cols_rejects_non_array_before_writing(tmp_path): cat.add_cols({"A": np.arange(N_OBJ), "B": list(range(N_OBJ))}) cat.close() assert path.read_bytes() == before + + +def test_add_cols_new_cat_ext_name_hdu_no(tmp_path): + """new_cat writes elsewhere and leaves the source; ext_name/hdu_no apply. + + The table at ``hdu_no`` gains the columns under ``ext_name`` in the new + file, the HDUs around it are carried over, and the source catalogue on + disk and in memory is untouched -- matching successive add_col calls. + """ + columns = _new_columns() + src = tmp_path / "src.fits" + _write_base(src, sex_catalogue=True) + before = src.read_bytes() + + outputs = {} + for label in ("one_by_one", "batched"): + dest = file_io.FITSCatalogue( + str(tmp_path / f"{label}.fits"), + open_mode=file_io.BaseCatalogue.OpenMode.ReadWrite, + ) + cat = file_io.FITSCatalogue( + str(src), + open_mode=file_io.BaseCatalogue.OpenMode.ReadWrite, + ) + cat.open() + kwargs = dict( + hdu_no=2, ext_name="OBJECTS_PLUS", new_cat=True, new_cat_inst=dest + ) + if label == "batched": + cat.add_cols(columns, **kwargs) + else: + # new_cat leaves the source open on the original file, so one + # add_col per column into the same destination keeps only the + # last; chain through the destination instead. + first, *rest = columns.items() + cat.add_col(*first, **kwargs) + cat.close() + for name, data in rest: + step = _open(dest.fullpath, sex_catalogue=True) + step.add_col(name, data) + step.close() + outputs[label] = (tmp_path / f"{label}.fits").read_bytes() + continue + assert cat.get_col_names(hdu_no=2) == ["NUMBER", "X"] + cat.close() + outputs[label] = (tmp_path / f"{label}.fits").read_bytes() + + assert src.read_bytes() == before + assert outputs["batched"] == outputs["one_by_one"] + with fits.open(tmp_path / "batched.fits") as hdul: + assert [h.name for h in hdul] == [ + "PRIMARY", "LDAC_IMHEAD", "OBJECTS_PLUS", "EPOCH_0", + ] + assert hdul[2].columns.names == ["NUMBER", "X", *columns] + npt.assert_array_equal(hdul[3].data["CCD_N"], [3]) diff --git a/tests/module/test_make_cat.py b/tests/module/test_make_cat.py index 48c22d435..e29109d19 100644 --- a/tests/module/test_make_cat.py +++ b/tests/module/test_make_cat.py @@ -579,31 +579,66 @@ def test_make_cat_runner_work_dir_publishes_same_catalogue( ): """WORK_DIR moves the build elsewhere; the published catalogue is the same. - The catalogue built in WORK_DIR is moved to the run's output directory - byte-identical to one built there directly, nothing is left in WORK_DIR, - and a work file a previous attempt left behind does not leak into it - (``save_as_fits`` appends to an existing file). + All three save stages run (ngmix, per-epoch PSF slots, one external mask + band). The catalogue built in WORK_DIR is moved to the run's output + directory byte-identical to one built there directly, nothing is left in + WORK_DIR, and a work file a previous attempt left behind does not leak + into it (``save_as_fits`` appends to an existing file). """ + healsparse = pytest.importorskip("healsparse") + obj_ids = [1, 2, 3] + ra = np.array([10.0, 10.1, 200.0]) + dec = np.array([20.0, 20.1, -40.0]) + sexcat = np.array( + list(zip(obj_ids, [1, 2, 0], ra, dec)), + dtype=[ + ("NUMBER", "i8"), + ("N_EPOCH", "i8"), + ("XWIN_WORLD", "f8"), + ("YWIN_WORLD", "f8"), + ], + ) tile_sexcat_path = tmp_path / "tile_sexcat-350-100.fits" - _write_sex_like_cat(tile_sexcat_path, _numbered_data(obj_ids)) + _write_sex_like_cat(tile_sexcat_path, sexcat) galaxy_psf_path = tmp_path / "galaxy_psf-350-100.sqlite" + _write_galaxy_psf_cat( + galaxy_psf_path, + { + 1: {"2113864-7": _psf_epoch(0.01, 0.02, 0.5)}, + 2: { + "2113864-9": _psf_epoch(0.05, 0.06, 0.7), + "2358123-21": _psf_epoch(0.07, 0.08, 0.9), + }, + 3: "empty", + }, + ) ngmix_path = tmp_path / "ngmix-350-100.fits" _write_ngmix_cat(ngmix_path, obj_ids) + mask = healsparse.HealSparseMap.make_empty(32, 4096, np.int16, sentinel=-1) + mask.update_values_pos( + ra[:2], dec[:2], np.full(2, 64, dtype=np.int16), lonlat=True + ) + mask_path = tmp_path / "mask_r.hsp" + mask.write(str(mask_path)) inputs = [str(tile_sexcat_path), str(galaxy_psf_path), str(ngmix_path)] + stages = { + "SHAPE_MEASUREMENT_TYPE": "ngmix", + "SAVE_PSF_DATA": "True", + "N_EPOCH_SLOTS": "3", + "MASK_EXT_PATHS": f"r:{mask_path}", + } published = {} - for label, section in ( + for label, work_dir in ( ("direct", {}), ("staged", {"WORK_DIR": "$SP_TEST_LOCAL/make_cat"}), ): out_dir = tmp_path / label out_dir.mkdir() config = CustomParser() - config.read_dict( - {"MAKE_CAT_RUNNER": {"SHAPE_MEASUREMENT_TYPE": "ngmix", **section}} - ) - if section: + config.read_dict({"MAKE_CAT_RUNNER": {**stages, **work_dir}}) + if work_dir: local = tmp_path / "local" monkeypatch.setenv("SP_TEST_LOCAL", str(local)) (local / "make_cat").mkdir(parents=True) @@ -620,6 +655,15 @@ def test_make_cat_runner_work_dir_publishes_same_catalogue( assert published["staged"] == published["direct"] assert list((tmp_path / "local" / "make_cat").iterdir()) == [] + with fits.open(make_cat.get_output_name(str(tmp_path / "staged"), "-350-100")) as hdul: + data = hdul[1].data + names = data.columns.names + npt.assert_array_equal(data["NUMBER"], obj_ids) + for col in ("TILE_ID", "NGMIX_MCAL_FLAGS", "HSM_G1_PSF_3", "EXP_ID_2"): + assert col in names, col + npt.assert_allclose(data["HSM_G1_PSF_1"], [0.01, 0.05, -10.0]) + npt.assert_array_equal(data["MASK_r"], [64, 64, -1]) + @pytest.mark.parametrize("shear", SHEAR_EXTS) @pytest.mark.parametrize("component", [0, 1], ids=["g1", "g2"])