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

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/>. 

21 

22 

23import itertools 

24import logging 

25import os 

26from typing import Any 

27 

28import click 

29 

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 

41 

42logger = logging.getLogger(__name__) 

43_FORMAT = ( 

44 "%(levelname) -10s %(asctime)s.%(msecs)03dZ %(name) -30s %(funcName) -35s %(lineno) -5d: %(message)s" 

45) 

46 

47RUCIO_REGISTER_CONFIG = "RUCIO_REGISTER_CONFIG" 

48_MSG = "environment variable not set, and no configuration was specified on the command line" 

49 

50 

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 

85 

86 

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) 

96 

97 

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}") 

105 

106 config = RucioRegisterConfig(config_file) 

107 

108 rucio_rse = config.rucio_rse 

109 scope = config.scope 

110 rse_root = config.rse_root 

111 dtn_url = config.dtn_url 

112 

113 butler = None 

114 if repo: 

115 butler = Butler(repo) 

116 

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 ) 

126 

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) 

131 

132 return ri, butler 

133 

134 

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) 

140 

141 

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) 

148 

149 

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) 

156 

157 

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") 

165 

166 

167@click.group(context_settings={"help_option_names": ["-h", "--help"]}) 

168def main(): 

169 pass 

170 

171 

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 

178 

179 

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) 

189 

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) 

193 

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) 

201 

202 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.DATA_PRODUCT, kwargs) 

203 

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 ) 

215 

216 dataset_refs = itertools.chain(*query.getDatasets()) 

217 

218 _register(ri, dataset_refs, chunk_size, rucio_dataset) 

219 

220 

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) 

237 

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) 

242 

243 repo = kwargs.get("repo", None) 

244 

245 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.DATA_PRODUCT) 

246 

247 uuids = [] 

248 with open(uuidlist) as f: 

249 for line in f: 

250 if not line.lstrip().startswith("#"): 

251 uuids.append(line.strip()) 

252 

253 dataset_refs = butler.get_many_datasets(uuids) 

254 

255 _register(ri, dataset_refs, chunk_size, rucio_dataset) 

256 

257 

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) 

268 

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") 

272 

273 repo = kwargs["repo"] 

274 

275 ri, butler = _getRucioInterface(repo, rucio_register_config, DataType.RAW_FILE, kwargs) 

276 

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"] 

281 

282 # chain is needed to flatten the list of lists returned by getDatasets() 

283 dataset_refs = itertools.chain.from_iterable(QueryDatasets(**kwargs).getDatasets()) 

284 

285 _register(ri, dataset_refs, chunk_size, rucio_dataset) 

286 

287 

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) 

295 

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) 

300 

301 ri, butler = _getRucioInterface(None, rucio_register_config, DataType.ZIP_FILE, kwargs) 

302 

303 _register_zips(ri, [zip_file], chunk_size, rucio_dataset) 

304 

305 

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) 

313 

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) 

318 

319 ri, butler = _getRucioInterface(None, rucio_register_config, DataType.DIM_FILE, kwargs) 

320 

321 _register_dims(ri, [dimension_file], chunk_size, rucio_dataset)