Coverage for python/lsst/daf/butler/registry/bridge/monolithic.py: 85%

122 statements  

« prev     ^ index     » next       coverage.py v7.16.2, created at 2026-09-28 09:15 +0000

1# This file is part of daf_butler. 

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 software is dual licensed under the GNU General Public License and also 

10# under a 3-clause BSD license. Recipients may choose which of these licenses 

11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt, 

12# respectively. If you choose the GPL option then the following text applies 

13# (but note that there is still no warranty even if you opt for BSD instead): 

14# 

15# This program is free software: you can redistribute it and/or modify 

16# it under the terms of the GNU General Public License as published by 

17# the Free Software Foundation, either version 3 of the License, or 

18# (at your option) any later version. 

19# 

20# This program is distributed in the hope that it will be useful, 

21# but WITHOUT ANY WARRANTY; without even the implied warranty of 

22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

23# GNU General Public License for more details. 

24# 

25# You should have received a copy of the GNU General Public License 

26# along with this program. If not, see <http://www.gnu.org/licenses/>. 

27from __future__ import annotations 

28 

29from ... import ddl 

30 

31__all__ = ("MonolithicDatastoreRegistryBridge", "MonolithicDatastoreRegistryBridgeManager") 

32 

33from collections import namedtuple 

34from collections.abc import Collection, Iterable, Iterator 

35from contextlib import contextmanager 

36from typing import TYPE_CHECKING, cast 

37 

38import sqlalchemy 

39 

40from lsst.utils.iteration import chunk_iterable 

41 

42from ..._dataset_ref import DatasetId 

43from ...datastore.stored_file_info import StoredDatastoreItemInfo 

44from ..interfaces import ( 

45 DatasetIdRef, 

46 DatastoreRegistryBridge, 

47 DatastoreRegistryBridgeManager, 

48 FakeDatasetRef, 

49 OpaqueTableStorage, 

50 VersionTuple, 

51) 

52from ..opaque import ByNameOpaqueTableStorage 

53from .ephemeral import EphemeralDatastoreRegistryBridge 

54 

55if TYPE_CHECKING: 

56 from ...datastore import DatastoreTransaction 

57 from ...dimensions import DimensionUniverse 

58 from ..interfaces import ( 

59 Database, 

60 DatasetRecordStorageManager, 

61 OpaqueTableStorageManager, 

62 StaticTablesContext, 

63 ) 

64 

65_TablesTuple = namedtuple( 

66 "_TablesTuple", 

67 [ 

68 "dataset_location", 

69 "dataset_location_trash", 

70 ], 

71) 

72 

73# This has to be updated on every schema change 

74_VERSION = VersionTuple(0, 2, 1) 

75 

76 

77def _makeTableSpecs(datasets: type[DatasetRecordStorageManager]) -> _TablesTuple: 

78 """Construct specifications for tables used by the monolithic datastore 

79 bridge classes. 

80 

81 Parameters 

82 ---------- 

83 datasets : subclass of `DatasetRecordStorageManager` 

84 Manager class for datasets; used only to create foreign key fields. 

85 

86 Returns 

87 ------- 

88 specs : `_TablesTuple` 

89 A named tuple containing `ddl.TableSpec` instances. 

90 """ 

91 # We want the dataset_location and dataset_location_trash tables 

92 # to have the same definition, aside from the behavior of their link 

93 # to the dataset table: the trash table has no foreign key constraint. 

94 # The order of columns in dataset_location_trash is reversed, it is more 

95 # optimal for query planner. 

96 

97 datastore_field = ddl.FieldSpec( 

98 name="datastore_name", 

99 dtype=sqlalchemy.String, 

100 length=256, 

101 primaryKey=True, 

102 nullable=False, 

103 doc="Name of the Datastore this entry corresponds to.", 

104 ) 

105 

106 dataset_location = ddl.TableSpec( 

107 doc=( 

108 "A table that provides information on whether a dataset is stored in " 

109 "one or more Datastores. The presence or absence of a record in this " 

110 "table itself indicates whether the dataset is present in that " 

111 "Datastore. " 

112 ), 

113 fields=[datastore_field], 

114 ) 

115 datasets.addDatasetForeignKey(dataset_location, primaryKey=True) 

116 

117 dataset_location_trash = ddl.TableSpec( 

118 doc="A table that keeps iinformation about datasets that are removed from Datastores.", 

119 fields=[], 

120 ) 

121 datasets.addDatasetForeignKey(dataset_location_trash, primaryKey=True, constraint=False) 

122 dataset_location_trash.fields.add(datastore_field) 

123 

124 return _TablesTuple( 

125 dataset_location=dataset_location, 

126 dataset_location_trash=dataset_location_trash, 

127 ) 

128 

129 

130class MonolithicDatastoreRegistryBridge(DatastoreRegistryBridge): 

131 """An implementation of `DatastoreRegistryBridge` that uses the same two 

132 tables for all non-ephemeral datastores. 

133 

134 Parameters 

135 ---------- 

136 datastoreName : `str` 

137 Name of the `Datastore` as it should appear in `Registry` tables 

138 referencing it. 

139 db : `Database` 

140 Object providing a database connection and generic distractions. 

141 tables : `_TablesTuple` 

142 Named tuple containing `sqlalchemy.schema.Table` instances. 

143 """ 

144 

145 def __init__(self, datastoreName: str, *, db: Database, tables: _TablesTuple): 

146 super().__init__(datastoreName) 

147 self._db = db 

148 self._tables = tables 

149 

150 def _refsToRows(self, refs: Iterable[DatasetIdRef]) -> list[dict]: 

151 """Transform an iterable of `DatasetRef` or `FakeDatasetRef` objects to 

152 a list of dictionaries that match the schema of the tables used by this 

153 class. 

154 

155 Parameters 

156 ---------- 

157 refs : `~collections.abc.Iterable` [ `DatasetRef` or `FakeDatasetRef` ] 

158 Datasets to transform. 

159 

160 Returns 

161 ------- 

162 rows : `list` [ `dict` ] 

163 List of dictionaries, with "datastoreName" and "dataset_id" keys. 

164 """ 

165 return [{"datastore_name": self.datastoreName, "dataset_id": ref.id} for ref in refs] 

166 

167 def ensure(self, refs: Iterable[DatasetIdRef]) -> None: 

168 # Docstring inherited from DatastoreRegistryBridge 

169 self._db.ensure(self._tables.dataset_location, *self._refsToRows(refs)) 

170 

171 def insert(self, refs: Iterable[DatasetIdRef]) -> None: 

172 # Docstring inherited from DatastoreRegistryBridge 

173 self._db.insert(self._tables.dataset_location, *self._refsToRows(refs)) 

174 

175 def forget(self, refs: Iterable[DatasetIdRef]) -> None: 

176 # Docstring inherited from DatastoreRegistryBridge 

177 with self._db.transaction(): 

178 # The list of IDs can be very large, split it into reasonable size 

179 # chunks to avoid hitting limits. 

180 for refs_chunk in chunk_iterable(refs, 50_000): 

181 dataset_ids = [ref.id for ref in refs_chunk] 

182 where = sqlalchemy.sql.and_( 

183 self._tables.dataset_location.columns.datastore_name == self.datastoreName, 

184 self._tables.dataset_location.columns.dataset_id.in_(dataset_ids), 

185 ) 

186 self._db.deleteWhere(self._tables.dataset_location, where) 

187 

188 def moveToTrash(self, refs: Iterable[DatasetIdRef], transaction: DatastoreTransaction | None) -> None: 

189 # Docstring inherited from DatastoreRegistryBridge 

190 location = self._tables.dataset_location 

191 location_trash = self._tables.dataset_location_trash 

192 with self._db.transaction(): 

193 for refs_chunk in chunk_iterable(refs, 50_000): 

194 # We only want to move IDs that actually exist in the 

195 # dataset_location table. Instead of querying for existing IDs, 

196 # which would need an extra query, we use INSERT ... SELECT 

197 # and DELETE using WHERE clause that limits operations to 

198 # existing IDs. 

199 dataset_ids = [ref.id for ref in refs_chunk] 

200 

201 where = sqlalchemy.sql.and_( 

202 location.columns.datastore_name == self.datastoreName, 

203 location.columns.dataset_id.in_(dataset_ids), 

204 ) 

205 

206 select = ( 

207 sqlalchemy.sql.select(location.columns.datastore_name, location.columns.dataset_id) 

208 .where(where) 

209 .with_for_update() 

210 ) 

211 self._db.insert(location_trash, select=select) 

212 

213 self._db.deleteWhere(location, where) 

214 

215 def check(self, datasets: Iterable[DatasetId]) -> set[DatasetId]: 

216 # Docstring inherited from DatastoreRegistryBridge 

217 found: set[DatasetId] = set() 

218 with self._db.session(): 

219 for batch in chunk_iterable(datasets, 50000): 

220 sql = ( 

221 sqlalchemy.sql.select(self._tables.dataset_location.columns.dataset_id) 

222 .select_from(self._tables.dataset_location) 

223 .where( 

224 sqlalchemy.sql.and_( 

225 self._tables.dataset_location.columns.datastore_name == self.datastoreName, 

226 self._tables.dataset_location.columns.dataset_id.in_(batch), 

227 ) 

228 ) 

229 ) 

230 with self._db.query(sql) as sql_result: 

231 sql_ids = sql_result.scalars().all() 

232 found.update(sql_ids) 

233 

234 return found 

235 

236 @contextmanager 

237 def emptyTrash( 

238 self, 

239 records_table: OpaqueTableStorage | None = None, 

240 record_class: type[StoredDatastoreItemInfo] | None = None, 

241 record_column: str | None = None, 

242 selected_ids: Collection[DatasetId] | None = None, 

243 dry_run: bool = False, 

244 ) -> Iterator[tuple[Iterable[tuple[DatasetIdRef, StoredDatastoreItemInfo | None]], set[str] | None]]: 

245 # Docstring inherited from DatastoreRegistryBridge 

246 

247 if records_table is None: 247 ↛ 248line 247 didn't jump to line 248 because the condition on line 247 was never true

248 raise ValueError("This implementation requires a records table.") 

249 

250 assert isinstance(records_table, ByNameOpaqueTableStorage), ( 

251 f"Records table must support hidden attributes. Got {type(records_table)}." 

252 ) 

253 

254 if record_class is None: 254 ↛ 255line 254 didn't jump to line 255 because the condition on line 254 was never true

255 raise ValueError("Record class must be provided if records table is given.") 

256 

257 # Helper closure to generate the common join+where clause. 

258 def join_records( 

259 select: sqlalchemy.sql.Select, location_table: sqlalchemy.schema.Table 

260 ) -> sqlalchemy.sql.Select: 

261 # mypy needs to be sure 

262 assert isinstance(records_table, ByNameOpaqueTableStorage) 

263 return select.select_from( 

264 records_table._table.join( 

265 location_table, 

266 onclause=records_table._table.columns.dataset_id == location_table.columns.dataset_id, 

267 ) 

268 ).where(location_table.columns.datastore_name == self.datastoreName) 

269 

270 # SELECT records.dataset_id, records.path FROM records 

271 # JOIN records on dataset_location.dataset_id == records.dataset_id 

272 # WHERE dataset_location.datastore_name = datastoreName 

273 

274 # It's possible that we may end up with a ref listed in the trash 

275 # table that is not listed in the records table. Such an 

276 # inconsistency would be missed by this query. 

277 info_in_trash = join_records(records_table._table.select(), self._tables.dataset_location_trash) 

278 if selected_ids: 

279 info_in_trash = info_in_trash.where( 

280 self._tables.dataset_location_trash.columns["dataset_id"].in_(selected_ids) 

281 ) 

282 info_in_trash = info_in_trash.with_for_update(skip_locked=True) 

283 

284 # Run query, transform results into a list of dicts that we can later 

285 # use to delete. 

286 with self._db.query(info_in_trash) as sql_result: 

287 rows = [dict(row, datastore_name=self.datastoreName) for row in sql_result.mappings()] 

288 

289 # It is possible for trashed refs to be linked to artifacts that 

290 # are still associated with refs that are not to be trashed. We 

291 # need to be careful to consider those and indicate to the caller 

292 # that those artifacts should be retained. Can only do this check 

293 # if the caller provides a column name that can map to multiple 

294 # refs. 

295 preserved: set[str] | None = None 

296 if record_column is not None: 296 ↛ 332line 296 didn't jump to line 332 because the condition on line 296 was always true

297 # The artifacts belonging to the trashed refs under consideration. 

298 items_in_trash = join_records( 

299 sqlalchemy.sql.select(records_table._table.columns[record_column]), 

300 self._tables.dataset_location_trash, 

301 ) 

302 if selected_ids: 

303 items_in_trash = items_in_trash.where( 

304 self._tables.dataset_location_trash.columns["dataset_id"].in_(selected_ids) 

305 ) 

306 items_in_trash_alias = items_in_trash.alias("items_in_trash") 

307 

308 # Whether one of those artifacts is also referenced by a ref that 

309 # is not in the trash. Written as a correlated EXISTS so that the 

310 # query is driven by the trashed artifacts and probes the records 

311 # table by artifact, rather than materializing every artifact the 

312 # datastore holds in order to join the two sets. 

313 referenced_by_live_ref = join_records( 

314 sqlalchemy.sql.select(sqlalchemy.sql.literal(1)), self._tables.dataset_location 

315 ).where( 

316 records_table._table.columns[record_column] == items_in_trash_alias.columns[record_column] 

317 ) 

318 

319 # A query for paths that are referenced by datasets in the trash 

320 # and datasets not in the trash. 

321 items_to_preserve = ( 

322 sqlalchemy.sql.select(items_in_trash_alias.columns[record_column]) 

323 .distinct() 

324 .where(sqlalchemy.exists(referenced_by_live_ref)) 

325 ) 

326 with self._db.query(items_to_preserve) as sql_result: 

327 preserved = {row[record_column] for row in sql_result.mappings()} 

328 

329 # Convert results to a tuple of id+info and a record of the artifacts 

330 # that should not be deleted from datastore. The id+info tuple is 

331 # solely to allow logging to report the relevant ID. 

332 id_info = ((FakeDatasetRef(row["dataset_id"]), record_class.from_record(row)) for row in rows) 

333 

334 # Start contextmanager, return results 

335 yield ((id_info, preserved)) 

336 

337 # No exception raised in context manager block. 

338 if not rows or dry_run: 

339 return 

340 

341 # Delete the rows from the records table 

342 records_table.delete(["dataset_id"], *[{"dataset_id": row["dataset_id"]} for row in rows]) 

343 

344 # Delete those rows from the trash table. 

345 self._db.delete( 

346 self._tables.dataset_location_trash, 

347 ["dataset_id", "datastore_name"], 

348 *[{"dataset_id": row["dataset_id"], "datastore_name": row["datastore_name"]} for row in rows], 

349 ) 

350 

351 

352class MonolithicDatastoreRegistryBridgeManager(DatastoreRegistryBridgeManager): 

353 """An implementation of `DatastoreRegistryBridgeManager` that uses the same 

354 two tables for all non-ephemeral datastores. 

355 

356 Parameters 

357 ---------- 

358 db : `Database` 

359 Object providing a database connection and generic distractions. 

360 tables : `_TablesTuple` 

361 Named tuple containing `sqlalchemy.schema.Table` instances. 

362 opaque : `OpaqueTableStorageManager` 

363 Manager object for opaque table storage in the `Registry`. 

364 universe : `DimensionUniverse` 

365 All dimensions know to the `Registry`. 

366 registry_schema_version : `VersionTuple` or `None`, optional 

367 The version of the registry schema. 

368 """ 

369 

370 def __init__( 

371 self, 

372 *, 

373 db: Database, 

374 tables: _TablesTuple, 

375 opaque: OpaqueTableStorageManager, 

376 universe: DimensionUniverse, 

377 registry_schema_version: VersionTuple | None = None, 

378 ): 

379 super().__init__( 

380 opaque=opaque, 

381 universe=universe, 

382 registry_schema_version=registry_schema_version, 

383 ) 

384 self._db = db 

385 self._tables = tables 

386 self._ephemeral: dict[str, EphemeralDatastoreRegistryBridge] = {} 

387 

388 def clone(self, *, db: Database, opaque: OpaqueTableStorageManager) -> DatastoreRegistryBridgeManager: 

389 return MonolithicDatastoreRegistryBridgeManager( 

390 db=db, 

391 tables=self._tables, 

392 opaque=opaque, 

393 universe=self.universe, 

394 registry_schema_version=self._registry_schema_version, 

395 ) 

396 

397 @classmethod 

398 def initialize( 

399 cls, 

400 db: Database, 

401 context: StaticTablesContext, 

402 *, 

403 opaque: OpaqueTableStorageManager, 

404 datasets: type[DatasetRecordStorageManager], 

405 universe: DimensionUniverse, 

406 registry_schema_version: VersionTuple | None = None, 

407 ) -> DatastoreRegistryBridgeManager: 

408 # Docstring inherited from DatastoreRegistryBridge 

409 tables = context.addTableTuple(_makeTableSpecs(datasets)) 

410 return cls( 

411 db=db, 

412 tables=cast(_TablesTuple, tables), 

413 opaque=opaque, 

414 universe=universe, 

415 registry_schema_version=registry_schema_version, 

416 ) 

417 

418 def refresh(self) -> None: 

419 # Docstring inherited from DatastoreRegistryBridge 

420 # This implementation has no in-Python state that depends on which 

421 # datastores exist, so there's nothing to do. 

422 pass 

423 

424 def register(self, name: str, *, ephemeral: bool = False) -> DatastoreRegistryBridge: 

425 # Docstring inherited from DatastoreRegistryBridge 

426 if ephemeral: 426 ↛ 427line 426 didn't jump to line 427 because the condition on line 426 was never true

427 return self._ephemeral.setdefault(name, EphemeralDatastoreRegistryBridge(name)) 

428 return MonolithicDatastoreRegistryBridge(name, db=self._db, tables=self._tables) 

429 

430 def findDatastores(self, ref: DatasetIdRef) -> Iterable[str]: 

431 # Docstring inherited from DatastoreRegistryBridge 

432 sql = ( 

433 sqlalchemy.sql.select(self._tables.dataset_location.columns.datastore_name) 

434 .select_from(self._tables.dataset_location) 

435 .where(self._tables.dataset_location.columns.dataset_id == ref.id) 

436 ) 

437 with self._db.query(sql) as sql_result: 

438 sql_rows = sql_result.mappings().fetchall() 

439 for row in sql_rows: 

440 yield row[self._tables.dataset_location.columns.datastore_name] 

441 for name, bridge in self._ephemeral.items(): 

442 if ref in bridge: 

443 yield name 

444 

445 @classmethod 

446 def currentVersions(cls) -> list[VersionTuple]: 

447 # Docstring inherited from VersionedExtension. 

448 return [_VERSION]