Coverage for python/lsst/dax/apdb/apdb.py: 100%

61 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-11 10:43 +0000

1# This file is part of dax_apdb. 

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 

22from __future__ import annotations 

23 

24__all__ = ["Apdb", "ApdbConfig"] 

25 

26from abc import ABC, abstractmethod 

27from collections.abc import Iterable, Mapping 

28from typing import TYPE_CHECKING 

29 

30import astropy.time 

31import pandas 

32 

33from lsst.resources import ResourcePathExpression 

34from lsst.sphgeom import Region 

35 

36from .apdbSchema import ApdbSchema, ApdbTables 

37from .config import ApdbConfig 

38from .factory import make_apdb 

39from .recordIds import DiaObjectId, DiaSourceId 

40from .schema_model import Table 

41 

42if TYPE_CHECKING: 

43 from .apdbAdmin import ApdbAdmin 

44 from .apdbMetadata import ApdbMetadata 

45 

46 

47class Apdb(ABC): 

48 """Abstract interface for APDB.""" 

49 

50 @classmethod 

51 def from_config(cls, config: ApdbConfig) -> Apdb: 

52 """Create Ppdb instance from configuration object. 

53 

54 Parameters 

55 ---------- 

56 config : `ApdbConfig` 

57 Configuration object, type of this object determines type of the 

58 Apdb implementation. 

59 

60 Returns 

61 ------- 

62 apdb : `apdb` 

63 Instance of `Apdb` class. 

64 """ 

65 return make_apdb(config) 

66 

67 @classmethod 

68 def from_uri(cls, uri: ResourcePathExpression) -> Apdb: 

69 """Make Apdb instance from a serialized configuration. 

70 

71 Parameters 

72 ---------- 

73 uri : `~lsst.resources.ResourcePathExpression` 

74 URI or local file path pointing to a file with serialized 

75 configuration, or a string with a "label:" prefix. In the latter 

76 case, the configuration will be looked up from an APDB index file 

77 using the label name that follows the prefix. The APDB index file's 

78 location is determined by the ``DAX_APDB_INDEX_URI`` environment 

79 variable. 

80 

81 Returns 

82 ------- 

83 apdb : `apdb` 

84 Instance of `Apdb` class, the type of the returned instance is 

85 determined by configuration. 

86 """ 

87 config = ApdbConfig.from_uri(uri) 

88 return make_apdb(config) 

89 

90 @abstractmethod 

91 def getConfig(self) -> ApdbConfig: 

92 """Return APDB configuration for this instance, including any updates 

93 that may be read from database. 

94 

95 Returns 

96 ------- 

97 config : `ApdbConfig` 

98 APDB configuration. 

99 """ 

100 raise NotImplementedError() 

101 

102 @abstractmethod 

103 def tableDef(self, table: ApdbTables) -> Table | None: 

104 """Return table schema definition for a given table. 

105 

106 Parameters 

107 ---------- 

108 table : `ApdbTables` 

109 One of the known APDB tables. 

110 

111 Returns 

112 ------- 

113 tableSchema : `.schema_model.Table` or `None` 

114 Table schema description, `None` is returned if table is not 

115 defined by this implementation. 

116 """ 

117 raise NotImplementedError() 

118 

119 @abstractmethod 

120 def getDiaObjects(self, region: Region) -> pandas.DataFrame: 

121 """Return catalog of DiaObject instances from a given region. 

122 

123 This method returns only the last version of each DiaObject, 

124 and may return only the subset of the DiaObject columns needed 

125 for AP association. Some 

126 records in a returned catalog may be outside the specified region, it 

127 is up to a client to ignore those records or cleanup the catalog before 

128 futher use. 

129 

130 Parameters 

131 ---------- 

132 region : `lsst.sphgeom.Region` 

133 Region to search for DIAObjects. 

134 

135 Returns 

136 ------- 

137 catalog : `pandas.DataFrame` 

138 Catalog containing DiaObject records for a region that may be a 

139 superset of the specified region. 

140 """ 

141 raise NotImplementedError() 

142 

143 @abstractmethod 

144 def getDiaSources( 

145 self, 

146 region: Region, 

147 object_ids: Iterable[int] | None, 

148 visit_time: astropy.time.Time, 

149 start_time: astropy.time.Time | None = None, 

150 ) -> pandas.DataFrame | None: 

151 """Return catalog of DiaSource instances from a given region. 

152 

153 Parameters 

154 ---------- 

155 region : `lsst.sphgeom.Region` 

156 Region to search for DIASources. 

157 object_ids : iterable [ `int` ], optional 

158 List of DiaObject IDs to further constrain the set of returned 

159 sources. If `None` then returned sources are not constrained. If 

160 list is empty then empty catalog is returned with a correct 

161 schema. 

162 visit_time : `astropy.time.Time` 

163 Time of the current visit. If APDB contains records later than this 

164 time they may also be returned. 

165 start_time : `astropy.time.Time`, optional 

166 Lower bound of time window for the query. If not specified then 

167 it is calculated using ``visit_time`` and 

168 ``read_forced_sources_months`` configuration parameter. 

169 

170 Returns 

171 ------- 

172 catalog : `pandas.DataFrame`, or `None` 

173 Catalog containing DiaSource records. `None` is returned if 

174 ``start_time`` is not specified and ``read_sources_months`` 

175 configuration parameter is set to 0. 

176 

177 Notes 

178 ----- 

179 This method returns DiaSource catalog for a region with additional 

180 filtering based on DiaObject IDs. Only a subset of DiaSource history 

181 is returned limited by ``read_sources_months`` config parameter, w.r.t. 

182 ``visit_time``. If ``object_ids`` is empty then an empty catalog is 

183 always returned with the correct schema (columns/types). If 

184 ``object_ids`` is `None` then no filtering is performed and some of the 

185 returned records may be outside the specified region. 

186 """ 

187 raise NotImplementedError() 

188 

189 @abstractmethod 

190 def getDiaForcedSources( 

191 self, 

192 region: Region, 

193 object_ids: Iterable[int] | None, 

194 visit_time: astropy.time.Time, 

195 start_time: astropy.time.Time | None = None, 

196 ) -> pandas.DataFrame | None: 

197 """Return catalog of DiaForcedSource instances from a given region. 

198 

199 Parameters 

200 ---------- 

201 region : `lsst.sphgeom.Region` 

202 Region to search for DIASources. 

203 object_ids : iterable [ `int` ], optional 

204 List of DiaObject IDs to further constrain the set of returned 

205 sources. If list is empty then empty catalog is returned with a 

206 correct schema. If `None` then returned sources are not 

207 constrained. 

208 visit_time : `astropy.time.Time` 

209 Time of the current visit. If APDB contains records later than this 

210 time they may also be returned. 

211 start_time : `astropy.time.Time`, optional 

212 Lower bound of time window for the query. If not specified then 

213 it is calculated using ``visit_time`` and 

214 ``read_forced_sources_months`` configuration parameter. 

215 

216 Returns 

217 ------- 

218 catalog : `pandas.DataFrame`, or `None` 

219 Catalog containing DiaForcedSource records. `None` is returned if 

220 ``start_time`` is not specified and ``read_forced_sources_months`` 

221 configuration parameter is set to 0. 

222 

223 Raises 

224 ------ 

225 NotImplementedError 

226 May be raised by some implementations if ``object_ids`` is `None`. 

227 

228 Notes 

229 ----- 

230 This method returns DiaForcedSource catalog for a region with 

231 additional filtering based on DiaObject IDs. Only a subset of DiaSource 

232 history is returned limited by ``read_forced_sources_months`` config 

233 parameter, w.r.t. ``visit_time``. If ``object_ids`` is empty then an 

234 empty catalog is always returned with the correct schema 

235 (columns/types). If ``object_ids`` is `None` then no filtering is 

236 performed and some of the returned records may be outside the specified 

237 region. 

238 """ 

239 raise NotImplementedError() 

240 

241 @abstractmethod 

242 def getDiaObjectsForDedup(self, since: astropy.time.Time | None = None) -> pandas.DataFrame: 

243 """Return catalog of DiaObject stored in APDB since specified time. 

244 

245 This method should be used by deduplication algorithm to retrieve 

246 DiaObject records added to APDB since previous deduplication (typically 

247 during previous night). Returned catalog will have only a small subset 

248 of DiaObject attributes required by deduplication algorithm. 

249 

250 Parameters 

251 ---------- 

252 since : `astropy.time.Time`, optional 

253 Starting search time (time of previous deduplication). If not 

254 provided the time of the last deduplication stored in metadata 

255 by `resetDedup` method is used. 

256 

257 Returns 

258 ------- 

259 catalog : `pandas.DataFrame` 

260 Catalog containing DiaObject records, only a subset of attributes 

261 will be returned. 

262 """ 

263 raise NotImplementedError() 

264 

265 @abstractmethod 

266 def getDiaSourcesForDiaObjects( 

267 self, objects: list[DiaObjectId], start_time: astropy.time.Time, max_dist_arcsec: float = 1.0 

268 ) -> pandas.DataFrame: 

269 """Return catalog of DiaSources associated with given DiaObjects. 

270 

271 Parameters 

272 ---------- 

273 objects : `list` [`DiaObjectId`] 

274 DiaObjects associated with returned DiaSources. 

275 start_time : `astropy.time.Time` 

276 Lower bound for ``midpointMjdTai`` for returned DiaSources. 

277 max_dist_arcsec : `float` 

278 Maximum expected distance in arcsec between DiaSource and 

279 DiaObject. This parameter is used to optimize spatial queries in 

280 cases when DiaObject is located near the partition boundary. If the 

281 distance from DiaObject to the boundary is smaller than 

282 ``max_dist_arcsec``, then the neighbor partition will be included 

283 in search too. 

284 

285 Returns 

286 ------- 

287 catalog : `pandas.DataFrame` 

288 Catalog containing DiaSource records associated to given 

289 DiaObjects. 

290 

291 Notes 

292 ----- 

293 Primary purpose of this method is to support deduplication algorithm. 

294 Its implementation is likely to be very slow and inefficient, it should 

295 not be used for regular queries. 

296 """ 

297 raise NotImplementedError() 

298 

299 @abstractmethod 

300 def containsVisitDetector( 

301 self, 

302 visit: int, 

303 detector: int, 

304 region: Region | None = None, 

305 visit_time: astropy.time.Time | None = None, 

306 ) -> bool: 

307 """Test whether any sources for a given visit-detector are present in 

308 the APDB. 

309 

310 Parameters 

311 ---------- 

312 visit, detector : `int` 

313 The ID of the visit-detector to search for. 

314 region : `lsst.sphgeom.Region`, optional 

315 Deprecated - parameter is not used. Region corresponding to the 

316 visit/detector combination. 

317 visit_time : `astropy.time.Time`, optional 

318 Deprecated - parameter is not used. Visit time (as opposed to visit 

319 processing time). This can be any timestamp in the visit timespan, 

320 e.g. its begin or end time. 

321 

322 Returns 

323 ------- 

324 present : `bool` 

325 `True` if given visit/detector combination was stored in APDB, 

326 `False` otherwise. 

327 """ 

328 raise NotImplementedError() 

329 

330 @abstractmethod 

331 def store( 

332 self, 

333 visit_time: astropy.time.Time, 

334 objects: pandas.DataFrame, 

335 sources: pandas.DataFrame | None = None, 

336 forced_sources: pandas.DataFrame | None = None, 

337 ) -> None: 

338 """Store all three types of catalogs in the database. 

339 

340 Parameters 

341 ---------- 

342 visit_time : `astropy.time.Time` 

343 Time of the visit. 

344 objects : `pandas.DataFrame` 

345 Catalog with DiaObject records. 

346 sources : `pandas.DataFrame`, optional 

347 Catalog with DiaSource records. 

348 forced_sources : `pandas.DataFrame`, optional 

349 Catalog with DiaForcedSource records. 

350 

351 Notes 

352 ----- 

353 This methods takes DataFrame catalogs, their schema must be 

354 compatible with the schema of APDB table: 

355 

356 - column names must correspond to database table columns 

357 - types and units of the columns must match database definitions, 

358 no unit conversion is performed presently 

359 - columns that have default values in database schema can be 

360 omitted from catalog 

361 - this method knows how to fill interval-related columns of DiaObject 

362 (validityStart, validityEnd) they do not need to appear in a 

363 catalog 

364 - source catalogs have ``diaObjectId`` column associating sources 

365 with objects 

366 

367 This operation need not be atomic, but DiaSources and DiaForcedSources 

368 will not be stored until all DiaObjects are stored. 

369 """ 

370 raise NotImplementedError() 

371 

372 @abstractmethod 

373 def reassignDiaSourcesToDiaObjects( 

374 self, 

375 idMap: Mapping[DiaSourceId, int], 

376 *, 

377 increment_nDiaSources: bool = True, 

378 decrement_nDiaSources: bool = True, 

379 ) -> None: 

380 """Re-assign DiaSources from one DiaObject to another, typically 

381 during deduplication. 

382 

383 Parameters 

384 ---------- 

385 idMap : `~collections.abc.Mapping` [`DiaSourceId`, `int`] 

386 Mapping from DiaSource to their new ``diaObjectId``. 

387 increment_nDiaSources : `bool`, optional 

388 If `True` then increment the value of ``nDiaSources`` in DiaObjects 

389 that DiaSources are reassigned to. 

390 decrement_nDiaSources : `bool`, optional 

391 If `True` then decrement the value of ``nDiaSources`` in DiaObjects 

392 that DiaSources are reassigned from. 

393 

394 Raises 

395 ------ 

396 LookupError 

397 Raised if some of DiaSources or DiaObjects are not found. 

398 

399 Notes 

400 ----- 

401 DiaSources initially could be associated with SSObjects. This method 

402 needs to be called before `setValidityEnd`. 

403 """ 

404 raise NotImplementedError() 

405 

406 @abstractmethod 

407 def setValidityEnd( 

408 self, objects: list[DiaObjectId], validityEnd: astropy.time.Time, raise_on_missing_id: bool = False 

409 ) -> int: 

410 """Close validity interval for specified DiaObjects. 

411 

412 Parameters 

413 ---------- 

414 objects : `list` [`DiaObjectId`] 

415 DiaObjects which will have their validityEnd updated, if their 

416 current validityEnd is NULL. 

417 validityEnd : `astropy.time.Time` 

418 Value for validityEnd. 

419 raise_on_missing_id : `bool`, optional 

420 If `True` then `LookupError` will be raised if any object in the 

421 list is missing from the database. 

422 

423 Returns 

424 ------- 

425 count : `int` 

426 Actual number of records for which validityEnd was updated. 

427 

428 Raises 

429 ------ 

430 LookupError 

431 Raised if ``raise_on_missing_id`` is `True` and some of the 

432 specified DiaObjects could not be found in the database. 

433 

434 Notes 

435 ----- 

436 This method has to be called after `reassignDiaSourcesToDiaObjects`. 

437 """ 

438 raise NotImplementedError() 

439 

440 @abstractmethod 

441 def resetDedup(self, dedup_time: astropy.time.Time | None = None) -> None: 

442 """Delete deduplication-related data and remember deduplication time. 

443 Deduplication data generated before ``dedup_time`` will be erased. 

444 

445 Parameters 

446 ---------- 

447 dedup_time : `astropy.time.Time`, optional 

448 Time of the last deduplication, current time is used if not 

449 provided. 

450 """ 

451 raise NotImplementedError() 

452 

453 @abstractmethod 

454 def reassignDiaSources(self, idMap: Mapping[int, int]) -> None: 

455 """Associate DiaSources with SSObjects, dis-associating them 

456 from DiaObjects. 

457 

458 Parameters 

459 ---------- 

460 idMap : `Mapping` 

461 Maps DiaSource IDs to their new SSObject IDs. 

462 

463 Raises 

464 ------ 

465 ValueError 

466 Raised if DiaSource ID does not exist in the database. 

467 """ 

468 raise NotImplementedError() 

469 

470 @abstractmethod 

471 def countUnassociatedObjects(self) -> int: 

472 """Return the number of DiaObjects that have only one DiaSource 

473 associated with them. 

474 

475 Used as part of ap_verify metrics. 

476 

477 Returns 

478 ------- 

479 count : `int` 

480 Number of DiaObjects with exactly one associated DiaSource. 

481 

482 Notes 

483 ----- 

484 This method can be very inefficient or slow in some implementations. 

485 """ 

486 raise NotImplementedError() 

487 

488 @property 

489 @abstractmethod 

490 def schema(self) -> ApdbSchema: 

491 """APDB table schema from ``sdm_schemas`` (`ApdbSchema`).""" 

492 raise NotImplementedError() 

493 

494 @property 

495 @abstractmethod 

496 def metadata(self) -> ApdbMetadata: 

497 """Object controlling access to APDB metadata (`ApdbMetadata`).""" 

498 raise NotImplementedError() 

499 

500 @property 

501 @abstractmethod 

502 def admin(self) -> ApdbAdmin: 

503 """Object providing adminitrative interface for APDB (`ApdbAdmin`).""" 

504 raise NotImplementedError() 

505 

506 def _current_time(self) -> astropy.time.Time: 

507 """Return current system time. 

508 

509 Returns 

510 ------- 

511 current_time : `astropy.time.Time` 

512 Current time. 

513 

514 Notes 

515 ----- 

516 This method exists primarily for testing purposes, it can be 

517 monkey-patched in unit tests to return something else than current 

518 system time, if necessary. 

519 """ 

520 return astropy.time.Time.now()