Skip to content

Commit a60f7ed

Browse files
authored
Merge pull request #103 from lsst/tickets/DM-51061
DM-51061: Add CLI command for replica chunk management.
2 parents 548d740 + dffa2ed commit a60f7ed

3 files changed

Lines changed: 148 additions & 2 deletions

File tree

python/lsst/dax/apdb/cli/apdb_cli.py

Lines changed: 38 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@
3131
from .logging_cli import LoggingCli
3232

3333

34-
def main(args: Sequence[str] | None = None) -> None:
34+
def main(args: Sequence[str] | None = None) -> int | None:
3535
"""APDB command line tools."""
3636
parser = argparse.ArgumentParser(description="APDB command line tools")
3737
log_cli = LoggingCli(parser)
@@ -45,14 +45,15 @@ def main(args: Sequence[str] | None = None) -> None:
4545
_metadata_subcommand(subparsers)
4646
_convert_legacy_config_subcommand(subparsers)
4747
_metrics_subcommand(subparsers)
48+
_replication_subcommand(subparsers)
4849

4950
parsed_args = parser.parse_args(args)
5051
log_cli.process_args(parsed_args)
5152

5253
kwargs = vars(parsed_args)
5354
# Strip keywords not understood by scripts.
5455
method = kwargs.pop("method")
55-
method(**kwargs)
56+
return method(**kwargs)
5657

5758

5859
def _create_sql_subcommand(subparsers: argparse._SubParsersAction) -> None:
@@ -239,3 +240,38 @@ def _metrics_log_to_influx(subparsers: argparse._SubParsersAction) -> None:
239240
default=False,
240241
)
241242
parser.set_defaults(method=scripts.metrics_log_to_influx)
243+
244+
245+
def _replication_subcommand(subparsers: argparse._SubParsersAction) -> None:
246+
parser = subparsers.add_parser("replication", help="Operations with replication tables produced by APDB.")
247+
subparsers = parser.add_subparsers(title="available subcommands", required=True)
248+
_replication_list_chunks_subcommand(subparsers)
249+
_replication_delete_chunks_subcommand(subparsers)
250+
251+
252+
def _replication_list_chunks_subcommand(subparsers: argparse._SubParsersAction) -> None:
253+
parser = subparsers.add_parser("list-chunks", help="Print full list of replication chunks in APDB.")
254+
parser.add_argument("apdb_config", help="Path to the APDB configuration.")
255+
parser.set_defaults(method=scripts.replication_list_chunks)
256+
257+
258+
def _replication_delete_chunks_subcommand(subparsers: argparse._SubParsersAction) -> None:
259+
parser = subparsers.add_parser("delete-chunks", help="Delete replication chunks from APDB.")
260+
parser.add_argument("apdb_config", help="Path to the APDB configuration.")
261+
parser.add_argument(
262+
"chunk_id", type=int, help="Chunk ID to delete, all earlier chunks are deleted as well."
263+
)
264+
parser.add_argument(
265+
"-p",
266+
"--print-only",
267+
help="Only print the list ochunks that will be deleted, but do not delete them.",
268+
action="store_true",
269+
default=False,
270+
)
271+
parser.add_argument(
272+
"--force",
273+
help="Do not ask for confirmation before deleting chunks.",
274+
action="store_true",
275+
default=False,
276+
)
277+
parser.set_defaults(method=scripts.replication_delete_chunks)

python/lsst/dax/apdb/scripts/__init__.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,3 +27,4 @@
2727
from .list_index import list_index
2828
from .metadata import metadata_delete, metadata_get, metadata_set, metadata_show
2929
from .metrics import metrics_log_to_influx
30+
from .replication import replication_delete_chunks, replication_list_chunks
Lines changed: 109 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,109 @@
1+
# This file is part of dax_ppdb
2+
#
3+
# Developed for the LSST Data Management System.
4+
# This product includes software developed by the LSST Project
5+
# (https://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 <https://www.gnu.org/licenses/>.
21+
22+
from __future__ import annotations
23+
24+
__all__ = ["replication_delete_chunks", "replication_list_chunks"]
25+
26+
import sys
27+
from collections.abc import Collection
28+
29+
from ..apdbReplica import ApdbReplica, ReplicaChunk
30+
31+
32+
def replication_list_chunks(apdb_config: str) -> None:
33+
"""Print full list of replica chunks existing in APDB.
34+
35+
Parameters
36+
----------
37+
apdb_config : `str`
38+
URL for APDB configuration file.
39+
"""
40+
apdb = ApdbReplica.from_uri(apdb_config)
41+
chunks = apdb.getReplicaChunks()
42+
if chunks is not None:
43+
chunks = sorted(chunks, key=lambda chunk: chunk.id)
44+
_print_chunks(chunks)
45+
else:
46+
print("APDB instance does not support replication")
47+
48+
49+
def replication_delete_chunks(apdb_config: str, chunk_id: int, force: bool, print_only: bool) -> int:
50+
"""Delete replication chunks from APDB.
51+
52+
Parameters
53+
----------
54+
apdb_config : `str`
55+
URL for APDB configuration file.
56+
chunk_id : `int`
57+
Chunk to delete, all earlier chunks are deleted as well.
58+
force : `bool`
59+
If `True` do not ask confirmation.
60+
print_only : `bool`
61+
If `True` print the list of chunks, but do not delete anything.
62+
"""
63+
apdb = ApdbReplica.from_uri(apdb_config)
64+
chunks = apdb.getReplicaChunks()
65+
if chunks is None:
66+
print("APDB instance does not support replications")
67+
else:
68+
chunks = sorted(chunks, key=lambda chunk: chunk.id)
69+
# Chunks to delete.
70+
chunks = [chunk for chunk in chunks if chunk.id <= chunk_id]
71+
# Check that given chunk ID actually exists.
72+
if not chunks or chunks[-1].id != chunk_id:
73+
print(f"ERROR: Replication chunk with ID={chunk_id} does not exist", file=sys.stderr)
74+
return 1
75+
76+
if print_only:
77+
print("Following chunks will be deleted:")
78+
_print_chunks(chunks)
79+
return 0
80+
81+
if not force:
82+
try:
83+
response = input(f"{len(chunks)} chunks will be removed, y[n]? ")
84+
except EOFError:
85+
response = ""
86+
if response not in ("y", "Y"):
87+
return 0
88+
89+
apdb.deleteReplicaChunks(chunk.id for chunk in chunks)
90+
91+
return 0
92+
93+
94+
def _print_chunks(chunks: Collection[ReplicaChunk]) -> None:
95+
"""Print the list of chunks.
96+
97+
Parameters
98+
----------
99+
chunks : `~collections.abc.Collection` [`ReplicaChunk`]
100+
Chunks to print.
101+
"""
102+
print(" Chunk Id Update time Unique Id")
103+
sep = "-" * 77
104+
print(sep)
105+
for chunk in chunks:
106+
insert_time = chunk.last_update_time
107+
print(f"{chunk.id:10d} {insert_time.tai.isot}/tai {chunk.unique_id}")
108+
print(sep)
109+
print(f"Total: {len(chunks)}")

0 commit comments

Comments
 (0)