Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 29 additions & 1 deletion docs/docs/pypaimon/cli-versions.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ Retain a dataset snapshot with a tag or work on a separate table history with a
Manage tags (named snapshots) on a table. Tags are useful for time travel and pinning a snapshot for later access.

```shell
paimon tag <create|list|get|delete> mydb.users ...
paimon tag <create|list|get|delete|rename|replace> mydb.users ...
```

### Tag Create
Expand All @@ -46,11 +46,15 @@ paimon tag create mydb.users v1 --snapshot-id 3

# Do not error if the tag already exists
paimon tag create mydb.users v1 --ignore-if-exists

# Retain the tag for 12 hours
paimon tag create mydb.users v1 --time-retained 12h
```

Options:
- `--snapshot-id, -s`: Snapshot id to tag (default: the latest snapshot)
- `--ignore-if-exists, -i`: Do not raise an error if the tag already exists
- `--time-retained, -r`: Retention for the new tag, for example `1d` or `12h`. Omit it to store a plain snapshot reference.

### Tag List

Expand Down Expand Up @@ -87,6 +91,30 @@ Options:
paimon tag delete mydb.users v1
```

### Tag Rename

```shell
paimon tag rename mydb.users v1 v2
```

Filesystem catalogs rename the tag file. A REST catalog has no rename endpoint and rejects the command.

### Tag Replace

Point an existing tag at another snapshot. Without `--snapshot-id`, the tag follows the latest snapshot. `--time-retained` stores a create time and TTL on the tag. Omit it and the tag is rewritten as a plain snapshot reference, which drops a retention that was already set. A REST catalog has no replace endpoint and rejects the command.

```shell
# Point v1 at the latest snapshot
paimon tag replace mydb.users v1

# Point v1 at snapshot 3 and retain it for 12 hours
paimon tag replace mydb.users v1 --snapshot-id 3 --time-retained 12h
```

Options:
- `--snapshot-id, -s`: Snapshot id to point at (default: the latest snapshot)
- `--time-retained, -r`: Retention for the replaced tag, for example `1d` or `12h`

## Branch Commands

Manage separate table histories. Create an empty branch with the current schema,
Expand Down
34 changes: 34 additions & 0 deletions paimon-python/pypaimon/catalog/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -425,5 +425,39 @@ def list_tags_paged(
"list_tags_paged is not supported by this catalog."
)

def rename_tag(
self,
identifier: Union[str, Identifier],
tag_name: str,
target_tag_name: str,
) -> None:
"""Rename a tag on a table.

Raises:
NotImplementedError: If the catalog does not support tag management.
"""
raise NotImplementedError(
"rename_tag is not supported by this catalog."
)

def replace_tag(
self,
identifier: Union[str, Identifier],
tag_name: str,
snapshot_id: Optional[int] = None,
time_retained: Optional[str] = None,
) -> None:
"""Point an existing tag at a snapshot.

``snapshot_id`` defaults to the latest snapshot. ``time_retained``
is written only when set.

Raises:
NotImplementedError: If the catalog does not support tag management.
"""
raise NotImplementedError(
"replace_tag is not supported by this catalog."
)

def auth_table_query(self, identifier: Identifier, select: Optional[List[str]]) -> TableQueryAuthResult:
raise NotImplementedError("auth_table_query not supported by this catalog")
49 changes: 49 additions & 0 deletions paimon-python/pypaimon/catalog/filesystem_catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -445,6 +445,40 @@ def list_tags_paged(
next_token = None
return PagedList(elements=page, next_page_token=next_token)

def rename_tag(
self,
identifier: Union[str, Identifier],
tag_name: str,
target_tag_name: str,
) -> None:
if not isinstance(identifier, Identifier):
identifier = Identifier.from_string(identifier)
table = self.get_table(identifier)
try:
table.rename_tag(tag_name, target_tag_name)
except ValueError as e:
_reraise_tag_value_error(
e, missing_tag=tag_name, conflict_tag=target_tag_name)

def replace_tag(
self,
identifier: Union[str, Identifier],
tag_name: str,
snapshot_id: Optional[int] = None,
time_retained: Optional[str] = None,
) -> None:
if not isinstance(identifier, Identifier):
identifier = Identifier.from_string(identifier)
table = self.get_table(identifier)
try:
table.replace_tag(
tag_name,
snapshot_id=snapshot_id,
time_retained=time_retained,
)
except ValueError as e:
_reraise_tag_value_error(e, missing_tag=tag_name)

# ===================== Branch CRUD =====================
# Thin wrappers that delegate to FileSystemBranchManager (returned by
# FileStoreTable.branch_manager() in the local-catalog case). Mirrors
Expand Down Expand Up @@ -536,3 +570,18 @@ def list_branches(
identifier = Identifier.from_string(identifier)
table = self.get_table(identifier)
return table.branch_manager().branches()


def _reraise_tag_value_error(exc, missing_tag, conflict_tag=None):
"""Map TagManager ValueError text onto catalog exceptions.

Snapshot-missing and blank-name errors stay ValueError. Only a missing
or conflicting tag name is translated. The ``Tag `` prefix keeps a
snapshot error such as ``Snapshot id '999' doesn't exist.`` as ValueError.
"""
message = str(exc)
if message.startswith("Tag ") and "doesn't exist" in message:
raise TagNotExistException(missing_tag) from exc
if conflict_tag is not None and "already exists" in message:
raise TagAlreadyExistException(conflict_tag) from exc
raise
18 changes: 18 additions & 0 deletions paimon-python/pypaimon/catalog/rest/rest_catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -624,6 +624,24 @@ def delete_tag(self, identifier: Union[str, Identifier], tag_name: str) -> None:
except ForbiddenException as e:
raise TableNoPermissionException(identifier) from e

def rename_tag(self, identifier: Union[str, Identifier], tag_name: str,
target_tag_name: str) -> None:
# The REST API has create, get, list, and delete only. Renaming the
# tag file through FileStoreTable would leave the server registry
# pointing at a name whose file is gone.
raise NotImplementedError(
"REST catalog does not support rename_tag. "
"The REST API has no rename endpoint."
)

def replace_tag(self, identifier: Union[str, Identifier], tag_name: str,
snapshot_id: Optional[int] = None,
time_retained: Optional[str] = None) -> None:
raise NotImplementedError(
"REST catalog does not support replace_tag. "
"The REST API has no replace endpoint."
)

# Branch CRUD: mirrors Java RESTCatalog branch handlers.
def create_branch(self, identifier: Union[str, Identifier], branch_name: str,
tag_name: Optional[str] = None) -> None:
Expand Down
83 changes: 79 additions & 4 deletions paimon-python/pypaimon/cli/cli_tag.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,10 @@

"""Tag commands for the Paimon CLI.

Adds the top-level ``tag {create,list,delete,get}`` subcommands (alongside
``table`` / ``db`` / ``catalog``). All operations go through the Catalog
layer so they get typed exceptions and work for both filesystem and REST
catalogs.
Adds the top-level ``tag {create,list,get,delete,rename,replace}``
subcommands (alongside ``table`` / ``db`` / ``catalog``). All operations go
through the Catalog layer so they get typed exceptions and work for both
filesystem and REST catalogs.
"""

import json
Expand Down Expand Up @@ -61,6 +61,7 @@ def cmd_tag_create(args):
identifier,
args.tag_name,
snapshot_id=args.snapshot_id,
time_retained=args.time_retained,
ignore_if_exists=args.ignore_if_exists,
)
except TableNotExistException:
Expand Down Expand Up @@ -96,6 +97,54 @@ def cmd_tag_delete(args):
print("Tag '{}' deleted from table '{}'.".format(args.tag_name, identifier))


def cmd_tag_rename(args):
"""Execute ``tag rename``."""
catalog, identifier = _open_catalog(args)
try:
catalog.rename_tag(identifier, args.tag_name, args.target_tag_name)
except TableNotExistException:
print("Error: Table '{}' does not exist.".format(identifier),
file=sys.stderr)
sys.exit(1)
except TagNotExistException:
print("Error: Tag '{}' does not exist.".format(args.tag_name),
file=sys.stderr)
sys.exit(1)
except TagAlreadyExistException:
print("Error: Tag '{}' already exists.".format(args.target_tag_name),
file=sys.stderr)
sys.exit(1)
except Exception as e:
print("Error: Failed to rename tag: {}".format(e), file=sys.stderr)
sys.exit(1)
print("Tag '{}' renamed to '{}' on table '{}'.".format(
args.tag_name, args.target_tag_name, identifier))


def cmd_tag_replace(args):
"""Execute ``tag replace``."""
catalog, identifier = _open_catalog(args)
try:
catalog.replace_tag(
identifier,
args.tag_name,
snapshot_id=args.snapshot_id,
time_retained=args.time_retained,
)
except TableNotExistException:
print("Error: Table '{}' does not exist.".format(identifier),
file=sys.stderr)
sys.exit(1)
except TagNotExistException:
print("Error: Tag '{}' does not exist.".format(args.tag_name),
file=sys.stderr)
sys.exit(1)
except Exception as e:
print("Error: Failed to replace tag: {}".format(e), file=sys.stderr)
sys.exit(1)
print("Tag '{}' replaced on table '{}'.".format(args.tag_name, identifier))


def cmd_tag_list(args):
"""Execute ``tag list``."""
catalog, identifier = _open_catalog(args)
Expand Down Expand Up @@ -179,6 +228,9 @@ def add_tag_subcommands(subparsers):
create_parser.add_argument(
'--ignore-if-exists', '-i', action='store_true',
help='Do not error if the tag already exists')
create_parser.add_argument(
'--time-retained', '-r', default=None,
help='Retention for the new tag, for example 1d or 12h')
create_parser.set_defaults(func=cmd_tag_create)

# tag list
Expand Down Expand Up @@ -212,3 +264,26 @@ def add_tag_subcommands(subparsers):
'table', help='Table identifier in format: database.table')
delete_parser.add_argument('tag_name', help='Name of the tag to delete')
delete_parser.set_defaults(func=cmd_tag_delete)

# tag rename
rename_parser = tag_subparsers.add_parser(
'rename', help='Rename a tag')
rename_parser.add_argument(
'table', help='Table identifier in format: database.table')
rename_parser.add_argument('tag_name', help='Current tag name')
rename_parser.add_argument('target_tag_name', help='New tag name')
rename_parser.set_defaults(func=cmd_tag_rename)

# tag replace
replace_parser = tag_subparsers.add_parser(
'replace', help='Point an existing tag at a snapshot')
replace_parser.add_argument(
'table', help='Table identifier in format: database.table')
replace_parser.add_argument('tag_name', help='Name of the tag to replace')
replace_parser.add_argument(
'--snapshot-id', '-s', type=int, default=None,
help='Snapshot id to point at (default: the latest snapshot)')
replace_parser.add_argument(
'--time-retained', '-r', default=None,
help='Retention for the replaced tag, for example 1d or 12h')
replace_parser.set_defaults(func=cmd_tag_replace)
80 changes: 80 additions & 0 deletions paimon-python/pypaimon/tests/cli_tag_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,13 @@ def test_create_duplicate_ignore_if_exists(self):
'tag', 'create', 'db.t', 'v1', '--ignore-if-exists')
self.assertEqual(0, code)

def test_create_time_retained(self):
_, _, code = self._run(
'tag', 'create', 'db.t', 'v1', '--time-retained', '12h')
self.assertEqual(0, code)
got, _, _ = self._run('tag', 'get', 'db.t', 'v1')
self.assertIn("Time Retained: PT12H", got)

# -- get -----------------------------------------------------------------

def test_get_table_format(self):
Expand Down Expand Up @@ -170,6 +177,79 @@ def test_delete_not_exists(self):
self.assertEqual(1, code)
self.assertIn("does not exist", err)

# -- rename --------------------------------------------------------------

def test_rename(self):
self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1')
out, _, code = self._run('tag', 'rename', 'db.t', 'v1', 'v2')
self.assertEqual(0, code)
self.assertIn("renamed", out)
listed, _, _ = self._run('tag', 'list', 'db.t', '--format', 'json')
self.assertEqual(["v2"], json.loads(listed))
got, _, _ = self._run('tag', 'get', 'db.t', 'v2')
self.assertIn("Snapshot ID: 1", got)

def test_rename_missing(self):
_, err, code = self._run('tag', 'rename', 'db.t', 'absent', 'v2')
self.assertEqual(1, code)
self.assertIn("does not exist", err)

def test_rename_target_exists(self):
self._run('tag', 'create', 'db.t', 'v1')
self._run('tag', 'create', 'db.t', 'v2')
_, err, code = self._run('tag', 'rename', 'db.t', 'v1', 'v2')
self.assertEqual(1, code)
self.assertIn("already exists", err)

# -- replace -------------------------------------------------------------

def test_replace_snapshot(self):
self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1')
out, _, code = self._run(
'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '2')
self.assertEqual(0, code)
self.assertIn("replaced", out)
got, _, _ = self._run('tag', 'get', 'db.t', 'v1')
self.assertIn("Snapshot ID: 2", got)

def test_replace_latest(self):
self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1')
_, _, code = self._run('tag', 'replace', 'db.t', 'v1')
self.assertEqual(0, code)
got, _, _ = self._run('tag', 'get', 'db.t', 'v1')
self.assertIn("Snapshot ID: 2", got)

def test_replace_time_retained(self):
self._run('tag', 'create', 'db.t', 'v1', '--snapshot-id', '1')
_, _, code = self._run(
'tag', 'replace', 'db.t', 'v1', '--time-retained', '12h')
self.assertEqual(0, code)
got, _, _ = self._run('tag', 'get', 'db.t', 'v1')
self.assertIn("Time Retained: PT12H", got)

def test_replace_missing_tag(self):
_, err, code = self._run(
'tag', 'replace', 'db.t', 'absent', '--snapshot-id', '1')
self.assertEqual(1, code)
self.assertIn("does not exist", err)

def test_replace_missing_snapshot(self):
self._run('tag', 'create', 'db.t', 'v1')
_, err, code = self._run(
'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '999')
self.assertEqual(1, code)
self.assertIn("Failed to replace tag", err)
self.assertIn("Snapshot id '999' doesn't exist.", err)

def test_replace_without_time_retained_drops_ttl(self):
self._run('tag', 'create', 'db.t', 'v1', '--time-retained', '12h')
_, _, code = self._run(
'tag', 'replace', 'db.t', 'v1', '--snapshot-id', '1')
self.assertEqual(0, code)
got, _, _ = self._run('tag', 'get', 'db.t', 'v1')
self.assertIn("Snapshot ID: 1", got)
self.assertNotIn("Time Retained:", got)

# -- bad input -----------------------------------------------------------

def test_invalid_identifier(self):
Expand Down
Loading
Loading