diff --git a/CHANGES.rst b/CHANGES.rst index 342626be..c0a5d424 100644 --- a/CHANGES.rst +++ b/CHANGES.rst @@ -5,6 +5,8 @@ Changelog Unreleased ---------- +* Support restoring a snapshot from an AWS S3 bucket with an endpoint URL. + 2.63.0 (2026-08-12) ------------------- diff --git a/crate/operator/restore_backup.py b/crate/operator/restore_backup.py index 8f4cc6ad..bcaff44b 100644 --- a/crate/operator/restore_backup.py +++ b/crate/operator/restore_backup.py @@ -246,7 +246,10 @@ async def get_source_backup_repository_data( f'Secret {secret_key_ref["name"]} could not be found.' ) except KeyError: - raise kopf.PermanentError(f"Key {key} not found in secret.") + if BackupRepositoryData.is_optional(backup_provider, key): + data_dict[key] = "" + else: + raise kopf.PermanentError(f"Key {key} not found in secret.") data_cls = BackupRepositoryData.get_class_from_backup_provider(backup_provider) data = BackupRepositoryData( @@ -688,8 +691,10 @@ async def _create_backup_repository( try: data = backup_repository_data.data for field in fields(data): - param = field.metadata["query_param"] value = getattr(data, field.name) + if not value: + continue + param = field.metadata["query_param"] create_repo_settings.append((param, value)) except KeyError as e: logger.warning( diff --git a/crate/operator/restore_backup_repository_data.py b/crate/operator/restore_backup_repository_data.py index c8c243c9..5e10e549 100644 --- a/crate/operator/restore_backup_repository_data.py +++ b/crate/operator/restore_backup_repository_data.py @@ -26,6 +26,9 @@ class AwsBackupRepositoryData: basePath: str = field(metadata={"query_param": "base_path"}) bucket: str = field(metadata={"query_param": "bucket"}) secretAccessKey: str = field(metadata={"query_param": "secret_key"}) + endpointUrl: str = field( + default="", metadata={"query_param": "endpoint", "optional": True} + ) @dataclass @@ -47,6 +50,8 @@ def __post_init__(self): ) for current_field in fields(self.data): value = getattr(self.data, current_field.name) + if current_field.metadata.get("optional") and not value: + continue if not isinstance(value, str) or not value: raise ValueError( f"Field `{current_field.name}` must be a non-empty string" @@ -73,6 +78,18 @@ def get_secrets_keys(backup_provider: BackupStorageProvider) -> list[str]: cls = BackupRepositoryData.get_class_from_backup_provider(backup_provider) return [field.name for field in fields(cls)] + @staticmethod + def is_optional(backup_provider: BackupStorageProvider, key: str) -> bool: + """ + Returns whether the given secrets key is optional for the given provider. + """ + cls = BackupRepositoryData.get_class_from_backup_provider(backup_provider) + return next( + field.metadata.get("optional", False) + for field in fields(cls) + if field.name == key + ) + @staticmethod def get_repository_type(backup_provider: BackupStorageProvider) -> str: """ diff --git a/deploy/charts/crate-operator-crds/templates/cratedbs-cloud-crate-io.yaml b/deploy/charts/crate-operator-crds/templates/cratedbs-cloud-crate-io.yaml index 538fb082..2a030fa5 100644 --- a/deploy/charts/crate-operator-crds/templates/cratedbs-cloud-crate-io.yaml +++ b/deploy/charts/crate-operator-crds/templates/cratedbs-cloud-crate-io.yaml @@ -379,6 +379,26 @@ spec: required: - secretKeyRef type: object + endpointUrl: + properties: + secretKeyRef: + properties: + key: + description: The key within the Kubernetes Secret + that holds the S3 endpoint. + type: string + name: + description: Name of a Kubernetes Secret that contains + the S3 endpoint-url to be used for accessing the + backup of the source cluster. + type: string + required: + - key + - name + type: object + required: + - secretKeyRef + type: object accountName: properties: secretKeyRef: diff --git a/tests/test_restore_backup.py b/tests/test_restore_backup.py index 04eab817..2c4f5d2a 100644 --- a/tests/test_restore_backup.py +++ b/tests/test_restore_backup.py @@ -967,6 +967,59 @@ async def test_create_backup_repository( ) +@pytest.mark.asyncio +async def test_create_backup_repository_with_endpoint_url( + faker, mock_cratedb_connection, backup_repository_data, mock_quote_ident +): + mock_cursor = mock_cratedb_connection["mock_cursor"] + mock_cursor.fetchone.return_value = None + + repository = faker.domain_word() + mock_logger = mock.Mock(spec=logging.Logger) + conn_factory = connection_factory("host", "password") + + data_dict = { + **backup_repository_data[BackupStorageProvider.AWS], + "endpointUrl": faker.url(), + } + data = BackupRepositoryData(data=AwsBackupRepositoryData(**data_dict)) + + with mock.patch( + "crate.operator.restore_backup.quote_ident", return_value=repository + ): + await RestoreBackupSubHandler._create_backup_repository( + conn_factory, repository, data, mock_logger + ) + + expected_stmt = ( + f"CREATE REPOSITORY {repository} TYPE s3 " + "WITH (max_restore_bytes_per_sec = %s, readonly = %s, " + "access_key = %s, base_path = %s, bucket = %s, secret_key = %s, " + "endpoint = %s);" + ) + expected_values = ( + "240mb", + "true", + data_dict["accessKeyId"], + data_dict["basePath"], + data_dict["bucket"], + data_dict["secretAccessKey"], + data_dict["endpointUrl"], + ) + mock_cursor.execute.assert_has_awaits( + [ + mock.call("SELECT * FROM sys.repositories WHERE name=%s", (repository,)), + mock.call(expected_stmt, expected_values), + ] + ) + + +def test_aws_backup_repository_data_not_optional_field(backup_repository_data): + data_dict = {**backup_repository_data[BackupStorageProvider.AWS], "bucket": ""} + with pytest.raises(ValueError, match="`bucket`"): + BackupRepositoryData(data=AwsBackupRepositoryData(**data_dict)) + + @pytest.mark.asyncio @mock.patch("crate.operator.restore_backup.get_gc_user_password") @mock.patch("crate.operator.restore_backup.execute_sql")