Coverage for python/lsst/rucio/register/script.py: 0%
170 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-26 08:56 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-07-26 08:56 +0000
1# This file is part of rucio_register
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (http://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This program is free software: you can redistribute it and/or modify
10# it under the terms of the GNU General Public License as published by
11# the Free Software Foundation, either version 3 of the License, or
12# (at your option) any later version.
13#
14# This program is distributed in the hope that it will be useful,
15# but WITHOUT ANY WARRANTY; without even the implied warranty of
16# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
17# GNU General Public License for more details.
18#
19# You should have received a copy of the GNU General Public License
20# along with this program. If not, see <http://www.gnu.org/licenses/>.
23import itertools
24import logging
25import os
26from typing import Any
28import click
30from lsst.daf.butler import Butler
31from lsst.daf.butler.cli.opt import (
32 log_level_option,
33 options_file_option,
34 query_datasets_options,
35)
36from lsst.daf.butler.script.queryDatasets import QueryDatasets
37from lsst.resources import ResourcePath
38from lsst.rucio.register.data_type import DataType
39from lsst.rucio.register.rucio_interface import RucioInterface
40from lsst.rucio.register.rucio_register_config import RucioRegisterConfig
42logger = logging.getLogger(__name__)
43_FORMAT = (
44 "%(levelname) -10s %(asctime)s.%(msecs)03dZ %(name) -30s %(funcName) -35s %(lineno) -5d: %(message)s"
45)
47RUCIO_REGISTER_CONFIG = "RUCIO_REGISTER_CONFIG"
48_MSG = "environment variable not set, and no configuration was specified on the command line"
51def register_options(func):
52 func = click.option(
53 "--backoff-max-tries",
54 required=False,
55 type=int,
56 default=5,
57 show_default=True,
58 help="maximum number of tries",
59 )(func)
60 func = click.option(
61 "--backoff-max-value",
62 required=False,
63 type=int,
64 default=30,
65 show_default=True,
66 help="maximum backoff value",
67 )(func)
68 func = click.option(
69 "--backoff-factor", required=False, type=float, default=5.0, show_default=True, help="backoff factor"
70 )(func)
71 func = click.option(
72 "--chunk-size",
73 required=False,
74 type=int,
75 default=30,
76 help="number of replica requests to make at once",
77 )(func)
78 func = click.option(
79 "--rucio-register-config", required=False, type=str, help="registration configuration file"
80 )(func)
81 func = click.option(
82 "--rucio-dataset", required=True, type=str, help="rucio dataset to register files to"
83 )(func)
84 return func
87def chunks(refs, chunk_size):
88 it = iter(refs)
89 while True:
90 chunk = itertools.islice(it, chunk_size)
91 try:
92 start = next(chunk)
93 except StopIteration:
94 return
95 yield itertools.chain((start,), chunk)
98def _getRucioInterface(repo, rucio_register_config, rubin_butler_type, kwargs):
99 # default to using RUCIO_REGISTER_CONFIG env variable
100 # if that's not set, try to use the command line
101 # if neither are set, then raise an Exception
102 config_file = os.environ.get(RUCIO_REGISTER_CONFIG, rucio_register_config)
103 if config_file is None:
104 raise RuntimeError(f"{RUCIO_REGISTER_CONFIG} {_MSG}")
106 config = RucioRegisterConfig(config_file)
108 rucio_rse = config.rucio_rse
109 scope = config.scope
110 rse_root = config.rse_root
111 dtn_url = config.dtn_url
113 butler = None
114 if repo:
115 butler = Butler(repo)
117 # create RucioInterface object used to register replicas into datasets
118 ri = RucioInterface(
119 butler=butler,
120 rucio_rse=rucio_rse,
121 scope=scope,
122 rse_root=rse_root,
123 dtn_url=dtn_url,
124 rubin_butler_type=rubin_butler_type,
125 )
127 backoff_factor = kwargs.get("backoff_factor")
128 backoff_max_value = kwargs.get("backoff_max_value")
129 backoff_max_tries = kwargs.get("backoff_max_tries")
130 ri.set_backoff(factor=backoff_factor, max_value=backoff_max_value, max_tries=backoff_max_tries)
132 return ri, butler
135def _register(ri, dataset_refs, chunk_size, rucio_dataset):
136 # register dataset_refs with Rucio into the rucio dataset, in chunks
137 for refs in chunks(dataset_refs, chunk_size):
138 cnt = ri.register_as_replicas(rucio_dataset, refs)
139 logger.debug("%d butler datasets registered", cnt)
142def _register_zips(ri, zip_files, chunk_size, rucio_dataset):
143 # register dataset_refs with Rucio into the rucio dataset, in chunks
144 for zip_file in zip_files:
145 rp = ResourcePath(zip_file)
146 cnt = ri.register_zips(rucio_dataset, [rp])
147 logger.debug("%d zips registered", cnt)
150def _register_dims(ri, dim_files, chunk_size, rucio_dataset):
151 # register dataset_refs with Rucio into the rucio dataset, in chunks
152 for dim_file in dim_files:
153 rp = ResourcePath(dim_file)
154 cnt = ri.register_dims(rucio_dataset, [rp])
155 logger.debug("%d dimension files registered", cnt)
158def _set_log_level(log_level):
159 if len(log_level):
160 level = log_level[None]
161 logging_num_level = getattr(logging, level.upper(), None)
162 else:
163 logging_num_level = logging.INFO
164 logging.basicConfig(level=logging_num_level, format=(_FORMAT), datefmt="%Y-%m-%d %H:%M:%S")
167@click.group(context_settings={"help_option_names": ["-h", "--help"]})
168def main():
169 pass
172def _get_and_delete(kwargs, key):
173 x = kwargs.get(key, None)
174 if x is None:
175 return x
176 del kwargs[key]
177 return x
180@main.command()
181@click.option("--repo", required=True, type=str, help="butler repository")
182@register_options
183@log_level_option()
184@options_file_option()
185@query_datasets_options(repo=False, showUri=True, useArguments=False)
186def data_products(**kwargs: Any) -> None:
187 log_level = kwargs.get("log_level", None)
188 _set_log_level(log_level)
190 rucio_register_config = kwargs.get("rucio_register_config", None)
191 rucio_dataset = kwargs.get("rucio_dataset", None)
192 chunk_size = kwargs.get("chunk_size", None)
194 repo = kwargs.get("repo", None)
195 collections = kwargs.get("collections", None)
196 where = kwargs.get("where", None)
197 find_first = kwargs.get("find_first", None)
198 limit = kwargs.get("limit", None)
199 order_by = kwargs.get("order_by", None)
200 dataset_type = kwargs.get("dataset_type", None)
202 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.DATA_PRODUCT, kwargs)
204 query = QueryDatasets(
205 butler=butler,
206 glob=dataset_type,
207 collections=collections,
208 where=where,
209 find_first=find_first,
210 limit=limit,
211 order_by=order_by,
212 show_uri=False,
213 with_dimension_records=True,
214 )
216 dataset_refs = itertools.chain(*query.getDatasets())
218 _register(ri, dataset_refs, chunk_size, rucio_dataset)
221@main.command()
222@click.option("--repo", required=True, type=str, help="butler repository")
223@register_options
224@click.option(
225 "--uuidlist",
226 required=False,
227 type=str,
228 help="""
229 filename of a list of butler dataset UUIDs to be register to the rucio dataset.
230 """,
231)
232@log_level_option()
233@options_file_option()
234def dataset_list(**kwargs: Any) -> None:
235 log_level = kwargs.get("log_level", None)
236 _set_log_level(log_level)
238 rucio_register_config = kwargs.get("rucio_register_config", None)
239 rucio_dataset = kwargs.get("rucio_dataset", None)
240 chunk_size = kwargs.get("chunk_size", None)
241 uuidlist = kwargs.get("uuidlist", None)
243 repo = kwargs.get("repo", None)
245 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.DATA_PRODUCT)
247 uuids = []
248 with open(uuidlist) as f:
249 for line in f:
250 if not line.lstrip().startswith("#"):
251 uuids.append(line.strip())
253 dataset_refs = butler.get_many_datasets(uuids)
255 _register(ri, dataset_refs, chunk_size, rucio_dataset)
258@main.command()
259@click.option("--repo", required=True, type=str, help="butler repository")
260@register_options
261@log_level_option()
262@options_file_option()
263@query_datasets_options(repo=False, showUri=True)
264def raws(**kwargs: Any) -> None:
265 # get and delete from kwargs; QueryDatasets doesn't like extra args
266 log_level = _get_and_delete(kwargs, "log_level")
267 _set_log_level(log_level)
269 rucio_register_config = _get_and_delete(kwargs, "rucio_register_config")
270 rucio_dataset = _get_and_delete(kwargs, "rucio_dataset")
271 chunk_size = _get_and_delete(kwargs, "chunk_size")
273 repo = kwargs["repo"]
275 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.RAW_FILE, kwargs)
277 # QueryDatsets called in this way doesn't like extra kwarg values
278 del kwargs["backoff-factor"]
279 del kwargs["backoff-max-value"]
280 del kwargs["backoff-max-tries"]
282 # chain is needed to flatten the list of lists returned by getDatasets()
283 dataset_refs = itertools.chain.from_iterable(QueryDatasets(**kwargs).getDatasets())
285 _register(ri, dataset_refs, chunk_size, rucio_dataset)
288@main.command()
289@register_options
290@click.option("--zip-file", required=True, help="zip file to register")
291@log_level_option()
292def zips(**kwargs: Any) -> None:
293 log_level = kwargs.get("log_level", None)
294 _set_log_level(log_level)
296 rucio_register_config = kwargs.get("rucio_register_config", None)
297 rucio_dataset = kwargs.get("rucio_dataset", None)
298 chunk_size = kwargs.get("chunk_size", None)
299 zip_file = kwargs.get("zip_file", None)
301 ri, butler = _getRucioInterface(None, rucio_register_config, DataType.ZIP_FILE, kwargs)
303 _register_zips(ri, [zip_file], chunk_size, rucio_dataset)
306@main.command()
307@register_options
308@click.option("--dimension-file", required=True, help="dimension file to register")
309@log_level_option()
310def dimensions(**kwargs: Any) -> None:
311 log_level = kwargs.get("log_level", None)
312 _set_log_level(log_level)
314 rucio_register_config = kwargs.get("rucio_register_config", None)
315 rucio_dataset = kwargs.get("rucio_dataset", None)
316 chunk_size = kwargs.get("chunk_size", None)
317 dimension_file = kwargs.get("dimension_file", None)
319 ri, butler = _getRucioInterface(None, rucio_register_config, DataType.DIM_FILE, kwargs)
321 _register_dims(ri, [dimension_file], chunk_size, rucio_dataset)