Skip to content
Closed
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
20 changes: 10 additions & 10 deletions services/backup-daemon/docker/granular/granular.py
Original file line number Diff line number Diff line change
Expand Up @@ -289,7 +289,7 @@ def __init__(self):
'storageName',
'blobPath']

def perform_granular_backup(self, backup_request):
def perform_granular_backup(self, backup_request, use_s3_alias=False):
# # for gke full backup
# if os.getenv("GOOGLE_APPLICATION_CREDENTIALS"):
# self.log.info('Perform GKE backup')
Expand Down Expand Up @@ -334,7 +334,7 @@ def perform_granular_backup(self, backup_request):
backup_id = backups.generate_backup_id()
backup_request['backupId'] = backup_id

worker = pg_backup.PostgreSQLDumpWorker(databases, backup_request, backup_request.get('blobPath'))
worker = pg_backup.PostgreSQLDumpWorker(databases, backup_request, backup_request.get('blobPath'), use_s3_alias=use_s3_alias,)

worker.start()

Expand Down Expand Up @@ -1084,7 +1084,7 @@ def post(self):
databases = body.get("databases") or []

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand All @@ -1101,7 +1101,7 @@ def post(self):
"blobPath": blob_path,
"storageName": storage_name
}
resp = GranularBackupRequestEndpoint().perform_granular_backup(backup_request)
resp = GranularBackupRequestEndpoint().perform_granular_backup(backup_request, use_s3_alias=True,)
body = None
code = None
if isinstance(resp, tuple) and len(resp) >= 2:
Expand Down Expand Up @@ -1172,7 +1172,7 @@ def get(self, backup_id):
blob_path = normalize_blobPath(meta.get("blobPath"))

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand Down Expand Up @@ -1217,7 +1217,7 @@ def delete(self, backup_id):
blob_path = normalize_blobPath(meta.get("blobPath"))

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand Down Expand Up @@ -1331,7 +1331,7 @@ def post(self, backup_id):
storage_name = self._normalize_storage_name(body.get("storageName"))

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand Down Expand Up @@ -1412,7 +1412,7 @@ def post(self, backup_id):
worker = pg_restore.PostgreSQLRestoreWorker(
requested, force,
{"backupId": backup_id, "namespace": namespace, "trackingId": tracking_id, "storageName": storage_name},
databases_mapping, owners_mapping, restore_roles, single_transaction, body.get("dbaasClone"), blob_path
databases_mapping, owners_mapping, restore_roles, single_transaction, body.get("dbaasClone"), blob_path, use_s3_alias=True,
)
worker.start()

Expand Down Expand Up @@ -1519,7 +1519,7 @@ def get(self, restore_id):
return "Invalid namespace name: %s." % namespace.encode("utf-8"), http.client.BAD_REQUEST

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand Down Expand Up @@ -1607,7 +1607,7 @@ def delete(self, restore_id):
}, http.client.BAD_REQUEST

try:
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=storage_name, prefix="", use_s3_alias=True)
except Exception as e:
return {"message": str(e)}, http.client.BAD_REQUEST

Expand Down
7 changes: 4 additions & 3 deletions services/backup-daemon/docker/granular/pg_backup.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@


class PostgreSQLDumpWorker(Thread):
def __init__(self, databases, backup_request, blob_path=None):
def __init__(self, databases, backup_request, blob_path=None, use_s3_alias=False):
Thread.__init__(self)

self.log = logging.getLogger("PostgreSQLDumpWorker")
Expand Down Expand Up @@ -64,11 +64,12 @@ def __init__(self, databases, backup_request, blob_path=None):
)
self.create_backup_dir()
self.storage_name = backup_request.get('storageName') or ""
self.use_s3_alias = use_s3_alias

if blob_path:
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name, prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name, prefix="", use_s3_alias=self.use_s3_alias,)
else:
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name) if os.environ['STORAGE_TYPE'] == "s3" else None
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name, use_s3_alias=self.use_s3_alias,) if os.environ['STORAGE_TYPE'] == "s3" else None

self._cancel_event = Event()
if configs.get_encryption():
Expand Down
7 changes: 4 additions & 3 deletions services/backup-daemon/docker/granular/pg_restore.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@

class PostgreSQLRestoreWorker(Thread):

def __init__(self, databases, force, restore_request, databases_mapping, owners_mapping, restore_roles=True, single_transaction=False, dbaas_clone=False, blobPath=None):
def __init__(self, databases, force, restore_request, databases_mapping, owners_mapping, restore_roles=True, single_transaction=False, dbaas_clone=False, blobPath=None, use_s3_alias=False):
Thread.__init__(self)

self.log = logging.getLogger("PostgreSQLRestoreWorker")
Expand Down Expand Up @@ -64,10 +64,11 @@ def __init__(self, databases, force, restore_request, databases_mapping, owners_
self.bin_path = configs.get_pgsql_bin_path(self.postgres_version)
self.parallel_jobs = configs.get_parallel_jobs()
self.storage_name = restore_request.get('storageName') or ""
self.use_s3_alias = use_s3_alias
if blobPath:
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name,prefix="")
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name,prefix="", use_s3_alias=self.use_s3_alias)
else:
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name) if os.environ['STORAGE_TYPE'] == "s3" else None
self.s3 = storage_s3.AwsS3Vault(storage_name=self.storage_name, use_s3_alias=self.use_s3_alias) if os.environ['STORAGE_TYPE'] == "s3" else None
self.blob_path = blobPath
self.backup_dir = backups.build_backup_path(self.backup_id, self.namespace, self.external_backup_root)
self.create_backup_dir(self.backup_dir)
Expand Down
23 changes: 12 additions & 11 deletions services/backup-daemon/docker/granular/storage_s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,6 @@ def get_s3_aliases(cls):
if not aliases:
raise Exception("S3 aliases are enabled, but /aliases/s3_aliases.json is empty")

if "default" not in aliases:
raise Exception("Default S3 alias is not configured in /aliases/s3_aliases.json")

cls.__s3_aliases_cache = aliases

Expand All @@ -62,18 +60,18 @@ def get_s3_aliases(cls):
@classmethod
def get_s3_alias_config(cls, storage_name=None):
aliases = cls.get_s3_aliases()

if aliases is None:
return None
raise Exception("storageName was provided but S3 aliases are not enabled")

if not storage_name or not storage_name.strip():
raise Exception("storageName is required when S3 aliases are enabled")

storage_name = storage_name.strip()
storage_name = storage_name.strip() if storage_name else ""
if not storage_name:
storage_name = "default"
cls.__log.info('storageName is empty, using default storageName "default"')

alias = aliases.get(storage_name)
if not alias:
raise Exception(f"S3 alias '{storage_name}' is not found in /aliases/s3_aliases.json")
available = ', '.join(aliases.keys())
raise Exception(f"S3 alias '{storage_name}' not found. Available: {available}")

return alias

Expand All @@ -84,10 +82,13 @@ def get_s3_bucket_name(alias=None):
return os.getenv("CONTAINER") or os.getenv("AWS_S3_BUCKET") or os.getenv("S3_BUCKET")

def __init__(self, storage_name=None, cluster_name=None, cache_enabled=False,
aws_s3_bucket_listing=None, prefix=None):
aws_s3_bucket_listing=None, prefix=None, use_s3_alias=False):

self.storage_name = storage_name
self.alias = AwsS3Vault.get_s3_alias_config(storage_name)
self.alias = None

if use_s3_alias:
self.alias = AwsS3Vault.get_s3_alias_config(storage_name)
self.bucket = AwsS3Vault.get_s3_bucket_name(self.alias)

if not self.bucket or not isinstance(self.bucket, str) or not self.bucket.strip():
Expand Down
Loading