diff --git a/estela-api/.env.example b/estela-api/.env.example index 897f46cb..bed46f50 100644 --- a/estela-api/.env.example +++ b/estela-api/.env.example @@ -24,6 +24,7 @@ MONGO_CONNECTION=dummy BUCKET_NAME_PROJECTS=dummy SECRET_KEY=dummy ENGINE=kubernetes +RESOURCE_DUAL_WRITE=False CREDENTIALS=aws SPIDERDATA_DB_ENGINE=mongodb ELASTICSEARCH_HOST=dummy diff --git a/estela-api/api/serializers/deploy.py b/estela-api/api/serializers/deploy.py index 955a9eff..22b8ebec 100644 --- a/estela-api/api/serializers/deploy.py +++ b/estela-api/api/serializers/deploy.py @@ -3,6 +3,7 @@ from api.serializers.project import UserDetailSerializer from api.serializers.spider import SpiderSerializer from core.models import Deploy, Spider +from core.resource_dual_write import attach_deploy_resource class DeploySerializer(serializers.ModelSerializer): @@ -28,6 +29,11 @@ class Meta: model = Deploy fields = ["did", "status", "created", "project_zip"] + def create(self, validated_data): + deploy = super().create(validated_data) + attach_deploy_resource(deploy) + return deploy + class DeployUpdateSerializer(serializers.ModelSerializer): spiders_names = serializers.ListField( diff --git a/estela-api/api/serializers/job.py b/estela-api/api/serializers/job.py index 12fe4913..d001bb5b 100644 --- a/estela-api/api/serializers/job.py +++ b/estela-api/api/serializers/job.py @@ -21,6 +21,7 @@ SpiderJobEnvVar, SpiderJobTag, ) +from core.resource_dual_write import attach_spider_job_resource class SpiderJobSerializer(serializers.ModelSerializer): @@ -154,6 +155,7 @@ def create(self, validated_data): job.tags.add(tag) job.save() + attach_spider_job_resource(job) return job diff --git a/estela-api/config/settings/base.py b/estela-api/config/settings/base.py index 20f4bc5e..1bf250e3 100644 --- a/estela-api/config/settings/base.py +++ b/estela-api/config/settings/base.py @@ -56,6 +56,7 @@ BUCKET_NAME_PROJECTS=(str, "dummy"), SECRET_KEY=(str, "dummy"), ENGINE=(str, "dummy"), + RESOURCE_DUAL_WRITE=(bool, False), BUILD=(str, "default"), STAGE=(str, "DEVELOPMENT"), SPIDERDATA_DB_ENGINE=(str, "dummy"), @@ -298,6 +299,7 @@ # Engine BUILD = env("BUILD") ENGINE = env("ENGINE") +RESOURCE_DUAL_WRITE = env("RESOURCE_DUAL_WRITE") CREDENTIALS = env("CREDENTIALS") SPIDERDATA_DB_ENGINE = env("SPIDERDATA_DB_ENGINE") diff --git a/estela-api/core/resource_dual_write.py b/estela-api/core/resource_dual_write.py new file mode 100644 index 00000000..46997dbe --- /dev/null +++ b/estela-api/core/resource_dual_write.py @@ -0,0 +1,65 @@ +from django.conf import settings +from django.db import transaction + +from core.models import Deploy, Resource, SpiderJob +from estela_resources import ResourceKind, ResourcePhase + + +def _dual_write_enabled(): + return getattr(settings, "RESOURCE_DUAL_WRITE", False) + + +def _spider_job_resource_phase(job: SpiderJob) -> ResourcePhase: + if job.status == SpiderJob.IN_QUEUE_STATUS: + return ResourcePhase.PENDING + if job.status == SpiderJob.WAITING_STATUS: + return ResourcePhase.PROVISIONING + return ResourcePhase.PENDING + + +def _deploy_resource_phase(deploy: Deploy) -> ResourcePhase: + if deploy.status == Deploy.BUILDING_STATUS: + return ResourcePhase.PROVISIONING + return ResourcePhase.PENDING + + +def attach_spider_job_resource(job: SpiderJob) -> None: + if not _dual_write_enabled(): + return + if job.resource_id: + return + with transaction.atomic(): + locked = SpiderJob.objects.select_for_update().get(pk=job.pk) + if locked.resource_id: + return + resource = Resource.objects.create( + project=locked.spider.project, + kind=ResourceKind.SPIDER_JOB, + phase=_spider_job_resource_phase(locked), + desired_spec={"jid": locked.jid}, + observed_state={}, + external_ref={}, + ) + locked.resource = resource + locked.save(update_fields=["resource"]) + + +def attach_deploy_resource(deploy): + if not _dual_write_enabled(): + return + if deploy.resource_id: + return + with transaction.atomic(): + locked = Deploy.objects.select_for_update().get(pk=deploy.pk) + if locked.resource_id: + return + resource = Resource.objects.create( + project=locked.project, + kind=ResourceKind.PROJECT_DEPLOY, + phase=_deploy_resource_phase(locked), + desired_spec={"did": locked.did}, + observed_state={}, + external_ref={}, + ) + locked.resource = resource + locked.save(update_fields=["resource"])