Skip to content

Commit e81895b

Browse files
feat: export to envidat (#1329)
Co-authored-by: Tasko Olevski <16360283+olevski@users.noreply.github.com>
1 parent 65190e2 commit e81895b

11 files changed

Lines changed: 348 additions & 46 deletions

File tree

bases/renku_data_services/data_api/app.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -276,6 +276,7 @@ def register_all_handlers(app: Sanic, dm: DependencyManager) -> Sanic:
276276
authenticator=dm.authenticator,
277277
metrics=dm.metrics,
278278
zenodo_client=dm.zenodo_client,
279+
envidat_client=dm.envidat_client,
279280
connected_services_repo=dm.connected_services_repo,
280281
job_client=dm.job_client,
281282
secret_client=dm.secret_client,

bases/renku_data_services/data_api/dependencies.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
DataConnectorRepository,
3939
DataConnectorSecretRepository,
4040
)
41+
from renku_data_services.data_connectors.deposits.envidat import EnvidatClient
4142
from renku_data_services.data_connectors.deposits.zenodo import ZenodoAPIClient
4243
from renku_data_services.git.gitlab import DummyGitlabAPI, EmptyGitlabAPI, GitlabAPI
4344
from renku_data_services.k8s.client_interfaces import K8sClient
@@ -171,6 +172,7 @@ class DependencyManager:
171172
resource_requests_repo: ResourceRequestsRepo
172173
resource_usage_service: ResourceUsageService
173174
zenodo_client: ZenodoAPIClient
175+
envidat_client: EnvidatClient
174176
job_client: DepositUploadJobClient
175177
secret_client: K8sSecretClient
176178
internal_token_mint: RenkuSelfTokenMint
@@ -515,6 +517,7 @@ def from_env(cls) -> DependencyManager:
515517
resource_requests_repo=resource_requests_repo,
516518
resource_usage_service=resource_usage_service,
517519
zenodo_client=ZenodoAPIClient(),
520+
envidat_client=EnvidatClient(),
518521
job_client=job_client,
519522
secret_client=secret_client,
520523
internal_token_mint=internal_token_mint,

components/renku_data_services/data_connectors/api.spec.yaml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1195,6 +1195,7 @@ components:
11951195
type: string
11961196
enum:
11971197
- zenodo
1198+
- envidat
11981199
DepositStatus:
11991200
type: string
12001201
enum:

components/renku_data_services/data_connectors/apispec.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,7 @@ class InaccessibleDataConnectorLinks(BaseAPISpec):
141141

142142
class DepositProvider(StrEnum):
143143
zenodo = "zenodo"
144+
envidat = "envidat"
144145

145146

146147
class DepositStatus(StrEnum):

components/renku_data_services/data_connectors/blueprints.py

Lines changed: 70 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@
5151
DataConnectorRepository,
5252
DataConnectorSecretRepository,
5353
)
54+
from renku_data_services.data_connectors.deposits.envidat import EnvidatClient
5455
from renku_data_services.data_connectors.deposits.zenodo import ZenodoAPIClient
5556
from renku_data_services.k8s.client_interfaces import K8sClient, SecretClient
5657
from renku_data_services.k8s.clients import DepositUploadJobClient
@@ -69,6 +70,7 @@ class DataConnectorsBP(CustomBlueprint):
6970
job_client: DepositUploadJobClient
7071
secret_client: SecretClient
7172
zenodo_client: ZenodoAPIClient
73+
envidat_client: EnvidatClient
7274
connected_services_repo: ConnectedServicesRepository
7375
data_source_repo: DataSourceRepository
7476
dc_storage_class: str
@@ -599,19 +601,19 @@ async def __get_zenodo_access_token(self, user: base_models.APIUser) -> str:
599601
)
600602
if not provider.connected_user:
601603
raise errors.UnauthorizedError(
602-
message="You need to connect and autheticate with the zenodo provider to do this"
604+
message="You need to connect and authenticate with the zenodo provider to do this"
603605
)
604606
token_set = await self.connected_services_repo.get_token_set(
605607
user=user, connection_id=provider.connected_user.connection.id
606608
)
607609
if not token_set:
608610
raise errors.UnauthorizedError(
609-
message="You need to connect and autheticate with the zenodo provider to do this"
611+
message="You need to connect and authenticate with the zenodo provider to do this"
610612
)
611613
access_token = token_set.access_token
612614
if not access_token:
613615
raise errors.UnauthorizedError(
614-
message="You need to connect and autheticate with the zenodo provider to do this"
616+
message="You need to connect and authenticate with the zenodo provider to do this"
615617
)
616618
return access_token
617619

@@ -626,15 +628,33 @@ async def _post_deposit(
626628
user: base_models.AuthenticatedAPIUser,
627629
body: apispec.DepositPost,
628630
) -> JSONResponse:
629-
existing_deposits, _ = await self.data_connector_repo.get_deposits(
630-
user, ULID.from_str(body.data_connector_id)
631-
)
631+
dc_id = ULID.from_str(body.data_connector_id)
632+
await self.data_connector_repo.get_data_connector(user=user, data_connector_id=dc_id)
633+
existing_deposits, _ = await self.data_connector_repo.get_deposits(user, dc_id)
632634
if len(existing_deposits) != 0:
633635
raise errors.ValidationError(
634636
message="Cannot have more than 1 deposit for the same user and data connector.",
635637
detail="Please delete your existing deposit and make a new one afterward.",
636638
)
637-
token = await self.__get_zenodo_access_token(user)
639+
640+
match body.provider:
641+
case apispec.DepositProvider.envidat:
642+
# TODO: Should we use the deposit ULID as the directory name!?
643+
# NOTE: Each new Envidat deposit gets its own dir, so a failed deposit cannot be resumed and need to
644+
# be cleaned up manually on Envidat.
645+
original_id = str(ULID())
646+
deposit_api_key: str | None = None
647+
648+
case apispec.DepositProvider.zenodo:
649+
token = await self.__get_zenodo_access_token(user)
650+
zenodo_dep = await self.zenodo_client.create_deposit(token, body.name)
651+
original_id = str(zenodo_dep.id)
652+
deposit_api_key = token
653+
654+
case x:
655+
raise errors.ValidationError(
656+
message=f"Received unknown deposit provider {x} when creating deposit."
657+
)
638658

639659
# The closure below allows us to tie the creation of the db entry to the successful
640660
# creation of the job in kubernetes. I.e. if the k8s creation fails nothing is saved
@@ -654,13 +674,12 @@ async def in_transaction_ops(
654674
job_client=self.job_client,
655675
data_connector_secret_repo=self.data_connector_secret_repo,
656676
data_source_repo=self.data_source_repo,
657-
deposit_api_key=token,
677+
deposit_api_key=deposit_api_key,
658678
)
659679

660-
zenodo_dep = await self.zenodo_client.create_deposit(token, body.name)
661-
unsaved_dep = validate_deposit(body, str(zenodo_dep.id))
680+
unsaved_dep = validate_deposit(body, original_id)
662681
saved_dep = await self.data_connector_repo.create_deposit(user, unsaved_dep, in_transaction_ops)
663-
return validated_json(apispec.Deposit, serialize_deposit(saved_dep))
682+
return validated_json(apispec.Deposit, serialize_deposit(saved_dep, self.deposit_config))
664683

665684
return "/deposits", ["POST"], _post_deposit
666685

@@ -684,7 +703,7 @@ async def _get_deposit(
684703
return HTTPResponse(status=304)
685704

686705
headers = {"ETag": saved_dep.etag}
687-
return validated_json(apispec.Deposit, serialize_deposit(saved_dep), headers=headers)
706+
return validated_json(apispec.Deposit, serialize_deposit(saved_dep, self.deposit_config), headers=headers)
688707

689708
return "/deposits/<deposit_id:ulid>", ["GET"], _get_deposit
690709

@@ -713,23 +732,35 @@ async def _patch_deposit(
713732
if patch.status:
714733
validate_deposit_status_change(saved_dep.deposit.status, patch.status)
715734
if patch.status == models.DepositStatus.complete:
716-
token = await self.__get_zenodo_access_token(user)
717-
zenodo_dep = await self.zenodo_client.get_deposit(token, saved_dep.deposit.original_id)
718-
if not zenodo_dep:
719-
raise errors.MissingResourceError(
720-
message=f"The deposit with id {saved_dep.deposit.original_id} "
721-
"cannot be found from the provider."
722-
)
723-
if not zenodo_dep.submitted:
724-
raise errors.ValidationError(
725-
message="The deposit needs to be completed and published first before being completed."
726-
)
735+
match saved_dep.deposit.source:
736+
case models.DepositSource.envidat:
737+
envidat_status = await self.envidat_client.get_deposit_status(saved_dep.deposit.original_id)
738+
if not envidat_status.is_published:
739+
raise errors.ValidationError(
740+
message="The deposit needs to be published on Envidat before being marked complete."
741+
)
742+
case models.DepositSource.zenodo:
743+
token = await self.__get_zenodo_access_token(user)
744+
zenodo_dep = await self.zenodo_client.get_deposit(token, saved_dep.deposit.original_id)
745+
if not zenodo_dep:
746+
raise errors.MissingResourceError(
747+
message=f"The Zenodo deposit with id {saved_dep.deposit.original_id} cannot be found."
748+
)
749+
if not zenodo_dep.submitted:
750+
raise errors.ValidationError(
751+
message="The Zenodo deposit needs to be completed and published first "
752+
"before being completed."
753+
)
754+
case x:
755+
raise errors.ValidationError(
756+
message=f"Received unknown deposit source {x} when pathing deposit with ID {deposit_id}."
757+
)
727758
saved_dep = await self.data_connector_repo.update_deposit(user, deposit_id, patch, etag=etag)
728759
if patch.status == models.DepositStatus.complete:
729760
# If the deposit is being completed then we delete it
730761
# We leave it to the user to create a new data connector and link it
731762
await self.data_connector_repo.delete_deposit(user, saved_dep.deposit.id)
732-
return validated_json(apispec.Deposit, serialize_deposit(saved_dep))
763+
return validated_json(apispec.Deposit, serialize_deposit(saved_dep, self.deposit_config))
733764

734765
return "/deposits/<deposit_id:ulid>", ["PATCH"], _patch_deposit
735766

@@ -741,14 +772,22 @@ def post_deposit_job(self) -> BlueprintFactoryResponse:
741772
async def _post_deposit_job(
742773
request: Request, user: base_models.AuthenticatedAPIUser, deposit_id: ULID
743774
) -> HTTPResponse:
744-
token = await self.__get_zenodo_access_token(user)
745775
saved_dep = await update_deposit_status(
746776
user,
747777
job=deposit_id,
748778
dc_repo=self.data_connector_repo,
749779
job_client=self.job_client,
750780
namespace=self.deposit_config.namespace,
751781
)
782+
if saved_dep.deposit.source == models.DepositSource.envidat:
783+
raise errors.ValidationError(
784+
message="Envidat deposits cannot be retried. Please delete this deposit and create a new one.",
785+
detail="Each Envidat deposit uploads to a unique S3 directory identified by the deposit ID. "
786+
"Starting over with a new deposit ensures a clean upload without mixing partial files.",
787+
)
788+
if saved_dep.deposit.status == models.DepositStatus.in_progress:
789+
raise errors.ValidationError(message="Cannot rerun a deposit job that is currently in progress.")
790+
token = await self.__get_zenodo_access_token(user)
752791
await self.job_client.delete(saved_dep.to_meta(user.id, self.deposit_config.namespace))
753792
new_job_name = "deposit-" + str(ULID()).lower()
754793
saved_dep = await self.data_connector_repo.update_deposit(
@@ -817,7 +856,9 @@ async def _get_deposits(
817856
job_client=self.job_client,
818857
namespace=self.deposit_config.namespace,
819858
)
820-
return [validate_and_dump(apispec.Deposit, serialize_deposit(i)) for i in deposits], total_num
859+
return [
860+
validate_and_dump(apispec.Deposit, serialize_deposit(i, self.deposit_config)) for i in deposits
861+
], total_num
821862

822863
return "/deposits", ["GET"], _get_deposits
823864

@@ -838,7 +879,9 @@ async def _get_dc_deposits(
838879
job_client=self.job_client,
839880
namespace=self.deposit_config.namespace,
840881
)
841-
return [validate_and_dump(apispec.Deposit, serialize_deposit(i)) for i in deposits], total_num
882+
return [
883+
validate_and_dump(apispec.Deposit, serialize_deposit(i, self.deposit_config)) for i in deposits
884+
], total_num
842885

843886
return "/data_connectors/<data_connector_id:ulid>/deposits", ["GET"], _get_dc_deposits
844887

components/renku_data_services/data_connectors/config.py

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
from kubernetes.client import ApiClient, V1Toleration
1212

1313
from renku_data_services.app_config import logging
14+
from renku_data_services.errors import errors
1415
from renku_data_services.k8s.constants import DEFAULT_K8S_CLUSTER, ClusterId
1516

1617
logger = logging.getLogger(__name__)
@@ -24,6 +25,7 @@ class DepositConfig:
2425
namespace: str
2526
renku_url: str
2627
zenodo_url: str
28+
envidat: EnvidatConfig
2729
node_selector: dict[str, str] | None = None
2830
tolerations: list[V1Toleration] | None = None
2931
cluster_id: Final[ClusterId] = DEFAULT_K8S_CLUSTER
@@ -64,4 +66,47 @@ def from_env(cls, renku_url: str) -> DepositConfig:
6466
namespace=os.environ["KUBERNETES_NAMESPACE"],
6567
cluster_id=DEFAULT_K8S_CLUSTER,
6668
zenodo_url=os.environ.get("ZENODO_URL", "https://zenodo.org").rstrip("/"),
69+
envidat=EnvidatConfig.from_env(),
6770
)
71+
72+
73+
@dataclass
74+
class EnvidatConfig:
75+
"""Configuration for envidat data exports and imports."""
76+
77+
exports_enabled: bool
78+
url: str
79+
rclone_image: str
80+
s3_endpoint: str
81+
s3_bucket: str
82+
s3_access_key_id: str
83+
s3_secret_access_key: str
84+
85+
@classmethod
86+
def from_env(cls) -> EnvidatConfig:
87+
"""Generate the config from environment variables."""
88+
exports_enabled = os.environ.get("ENVIDAT_EXPORTS_ENABLED", "false").lower() == "true"
89+
output = cls(
90+
exports_enabled=exports_enabled,
91+
url=os.environ.get("ENVIDAT_URL", "https://www.envidat.ch").rstrip("/"),
92+
rclone_image=os.environ.get("ENVIDAT_RCLONE_IMAGE", "rclone/rclone:1"),
93+
s3_endpoint=os.environ.get("ENVIDAT_S3_ENDPOINT", ""),
94+
s3_bucket=os.environ.get("ENVIDAT_S3_BUCKET", ""),
95+
s3_access_key_id=os.environ.get("ENVIDAT_S3_ACCESS_KEY_ID", ""),
96+
s3_secret_access_key=os.environ.get("ENVIDAT_S3_SECRET_ACCESS_KEY", ""),
97+
)
98+
if exports_enabled and any(
99+
[
100+
i == ""
101+
for i in [
102+
output.s3_endpoint,
103+
output.s3_bucket,
104+
output.s3_access_key_id,
105+
output.s3_secret_access_key,
106+
]
107+
]
108+
):
109+
raise errors.ConfigurationError(
110+
message="Envidat exports are enabled but not all required parameters are provided."
111+
)
112+
return output

0 commit comments

Comments
 (0)