diff --git a/dbaas/drivers/replication_topologies/mysql.py b/dbaas/drivers/replication_topologies/mysql.py index 610683c02..13aff668b 100644 --- a/dbaas/drivers/replication_topologies/mysql.py +++ b/dbaas/drivers/replication_topologies/mysql.py @@ -495,19 +495,11 @@ def get_deploy_steps(self): )}, { 'Creating virtual machine': ( 'workflow.steps.util.host_provider.AllocateIP', - 'workflow.steps.util.host_provider.CreateVirtualMachine', + 'workflow.steps.util.host_provider.CreateMySqlFoxVirtualMachineWithEnvTag', )}, { 'Creating VIP': ( - 'workflow.steps.util.vip_provider.CreateVip', - 'workflow.steps.util.vip_provider.CreateInstanceGroup', - 'workflow.steps.util.vip_provider.AddInstancesInGroup', - 'workflow.steps.util.vip_provider.CreateHeathcheck', - 'workflow.steps.util.vip_provider.CreateBackendService', - 'workflow.steps.util.vip_provider.AllocateIP', - 'workflow.steps.util.vip_provider.AllocateDNS', - 'workflow.steps.util.vip_provider.CreateForwardingRule', - 'workflow.steps.util.vip_provider.AddLoadBalanceLabels', - 'workflow.steps.util.dns.RegisterDNSVip', + 'workflow.steps.util.ingress_provider.AllocateProvider', + 'workflow.steps.util.ingress_provider.RegisterDNSIngress', )}, { 'Creating dns': ( 'workflow.steps.util.dns.CreateDNS', diff --git a/dbaas/workflow/steps/util/host_provider.py b/dbaas/workflow/steps/util/host_provider.py index bc4a8ad2d..a0a3d8ec3 100644 --- a/dbaas/workflow/steps/util/host_provider.py +++ b/dbaas/workflow/steps/util/host_provider.py @@ -252,6 +252,72 @@ def create_host(self, infra, offering, name, team_name, zone=None, host.save() return host + def create_host_with_environment_tag_to_foxha(self, infra, offering, name, team_name, zone=None, + database_name='', host_obj=None, port=None, + volume_name=None, init_user=None, init_password=None, + static_ip=None): + url = "{}/{}/{}/host/new".format( + self.credential.endpoint, self.provider, self.environment + ) + + data = { + "engine": self.engine, + "name": name, + "cpu": offering.cpus, + "memory": offering.memory_size_mb, + "group": infra.name, + "team_name": team_name, + "database_name": database_name, + "static_ip_id": static_ip and static_ip.identifier, + "service_account": self.infra_service_account + } + if zone: + data['zone'] = zone + if port: + data['port'] = port + if volume_name: + data['volume_name'] = volume_name + if init_user: + data['init_user'] = init_user + if init_password: + data['init_password'] = init_password + + dev = ['gcp-lab-dev'] + devqa = ['dev-gcp-hdg-us-east1', 'dev-gcp-hdg-sa-east1', 'dev-gcp-hdg-us-central1', + 'dev-gcp-tsuru-us-east1', 'dev-gcp-tsuru-sa-east1', 'dev-gcp-tsuru-us-central1', + 'dbaas-devqa-gcp-us-east1'] + prod = ['prod-gcp-hdg-us-east1', 'prod-gcp-hdg-sa-east-1', 'prod-gcp-hdg-us-central1', + 'prod-gcp-tsuru-us-east1', 'prod-gcp-tsuru-sa-east1', 'prod-gcp-tsuru-us-central1'] + tag = '' + env = infra.environment.name + if env in dev: + tag = 'dbaas-nodes-dev' + elif env in devqa: + tag = 'dbaas-nodes-devqa' + elif env in prod: + tag = 'dbaas-nodes-prod' + data['environment_tag'] = tag + + response = self._request(post, url, json=data, timeout=900) + if response.status_code != 201: + raise HostProviderCreateVMException(response.content, response) + + content = response.json() + if host_obj is None: + host = Host() + host.hostname = content["address"] + else: + host = host_obj + host.address = content["address"] + host.user = self.vm_credential.user + host.password = self.vm_credential.password + host.private_key = self.vm_credential.private_key + host.provider = self.provider + host.identifier = content["id"] + host.offering = offering + host.save() + return host + def create_static_ip(self, infra): url = "{}/{}/{}/ip/".format( self.credential.endpoint, self.provider, self.environment @@ -516,7 +582,7 @@ class Stop(HostProviderStep): def __unicode__(self): return "Stopping VM..." - + @property def is_valid(self): return not self.instance.temporary @@ -524,7 +590,7 @@ def is_valid(self): def do(self): if not self.is_valid: return - + stopped = self.provider.stop() if not stopped: raise EnvironmentError("Could not stop VM") @@ -552,7 +618,7 @@ def __unicode__(self): def do(self): if not self.is_valid: return - + started = self.provider.start() if not started: raise EnvironmentError("Could not start VM") @@ -627,7 +693,7 @@ def __unicode__(self): def do(self): if not self.is_valid: return - + success = self.provider.new_offering(self.target_offering) if not success: raise Exception("Could not change offering") @@ -751,6 +817,52 @@ def undo(self): host.delete() +class CreateMySqlFoxVirtualMachineWithEnvTag(CreateVirtualMachine): + + def __unicode__(self): + return "Creating Mysql Fox virtual machine..." + + def do(self): + task_manager = self.create or self.destroy + if hasattr(self, 'step_manager') and task_manager is None: + task_manager = self.step_manager + try: + pair = self.infra.instances.get(dns=self.instance.dns) + except Instance.DoesNotExist: + host = self.provider.create_host_with_environment_tag_to_foxha( + self.infra, self.offering, self.vm_name, self.team, self.zone, + database_name=(self.database.name if self.database + else task_manager.name), + static_ip=self.instance.static_ip + ) + self.update_databaseinfra_last_vm_created() + else: + host = pair.hostname + + self.create_instance(host) + self.associate_static_ip_with_instance() + + def undo(self): + try: + host = self.instance.hostname + except ObjectDoesNotExist: + self.delete_instance() + return + + try: + self.provider.destroy_host(self.host) + except HostProviderDestroyVMException as e: + content, response = e + if response.status_code == 404: + LOG.warning('Host {} not found in host-provider'.format(self.host)) + else: + raise e + + self.delete_instance() + if host.id: + host.delete() + + class DestroyVirtualMachineTemporaryInstance(CreateVirtualMachine): def __unicode__(self): @@ -759,16 +871,16 @@ def __unicode__(self): @property def is_valid(self): return self.instance.temporary - + def do(self): if not self.is_valid: return - + return super(DestroyVirtualMachineTemporaryInstance, self).undo() class CreateVirtualMachineTemporaryInstance(CreateVirtualMachine): - + @property def is_valid(self): return self.instance.temporary @@ -782,7 +894,7 @@ def create_instance(self, host): def do(self): if self.is_valid: super(CreateVirtualMachineTemporaryInstance, self).do() - + def undo(self): if self.is_valid: super(CreateVirtualMachineTemporaryInstance, self).undo() @@ -804,19 +916,20 @@ def undo(self): class AllocateIPTemporaryInstance(AllocateIP): - + @property def is_valid(self): return self.instance.temporary - + def do(self): if self.is_valid: super(AllocateIPTemporaryInstance, self).do() - + def undo(self): if self.is_valid: super(AllocateIPTemporaryInstance, self).undo() + class AllocateIPRegionMigrate(HostProviderStep): @property @@ -1020,7 +1133,6 @@ def do(self): class WaitingBeReady(HostProviderStep): - RETRIES = 30 def __unicode__(self): @@ -1192,21 +1304,21 @@ def do(self): def undo(self): raise NotImplementedError - + class DestroyIPTemporaryInstance(DestroyIPMigrate): def __unicode__(self): return "Destroy IP..." - + @property def is_valid(self): return self.instance.temporary - + def do(self): if not self.is_valid: return - + return super(DestroyIPTemporaryInstance, self).do() @@ -1221,7 +1333,7 @@ def environment(self): def do(self): pass - #super(DestroyServiceAccountMigrate, self).undo() + # super(DestroyServiceAccountMigrate, self).undo() def undo(self): raise NotImplementedError diff --git a/dbaas/workflow/steps/util/ingress_provider.py b/dbaas/workflow/steps/util/ingress_provider.py new file mode 100644 index 000000000..e7c2d93d2 --- /dev/null +++ b/dbaas/workflow/steps/util/ingress_provider.py @@ -0,0 +1,213 @@ +import socket +from time import sleep +from requests import post, get +from physical.models import Vip +from base import BaseInstanceStep +from util import get_credentials_for +from dbaas_dnsapi.models import FOXHA +from dbaas_dnsapi.utils import add_dns_record +from dbaas_dnsapi.provider import DNSAPIProvider +from dbaas_credentials.models import CredentialType + + +class IngressProvider(object): + + def __init__(self, instance): + self.instance = instance + self._team = None + self._credential = None + self._ip = None + self._port = None + + @property + def infra(self): + return self.instance.databaseinfra + + @property + def host(self): + return self.instance.hostname + + @property + def team(self): + return self.infra.databases.first().team.name + + @property + def port(self): + return int(self.instance.port) + + @property + def is_database(self): + return self.instance.is_database + + def _request(self, action, url, **kw): + return action(url, verify=False, **kw) + + @property + def ip(self): + return self._ip + + def credential(self): + print('*-----------------------INSERTING CREDENTIAL----------------------------*') + print(f'*-----------------------ENVIRONMENT [{self.infra.environment}]----------------------------*') + if not self._credential: + self._credential = get_credentials_for( + self.infra.environment, CredentialType.INGRESS_PROVIDER + ) + print('*-----------------------CREATED CREDENTIAL----------------------------*') + print(f'*-----------------------[{self._credential}]----------------------------*') + return self._credential + + def create_ingress(self, infra, port, team_name): + data = { + "team": team_name, + "bank_port": port, + "bank_address": [str(infra.hosts[0].address)], + "bank_type": 'MySQLFOXHA', + "bank_name": infra.name_prefix, + "bank_service_account": str(infra.service_account) + } + try: + response = self._request(post, self.credential.endpoint, json=data, timeout=6000) + + if response.status_code not in [200, 201]: + raise response.raise_for_status() + ingress = response.json()['value'] + + except Exception as error: + print(error) + raise Exception + self._port = ingress['port_external'] + self._team = team_name + return ingress + + +class IngressProviderStep(BaseInstanceStep): + + def __init__(self, instance=None): + super(IngressProviderStep, self).__init__(instance) + self._ingress_provider = None + + @property + def is_valid(self): + return self.instance == self.infra.instances.first() + + @property + def vip_ip(self): + return self.ingressprovider.ip + + @property + def ingressprovider(self): + if not self._ingress_provider: + self._ingress_provider = IngressProvider(self.instance) + return self._ingress_provider + + def do(self): + raise NotImplementedError + + def undo(self): + pass + + +class AllocateProvider(IngressProviderStep): + + def __unicode__(self): + return "Allocating Provider on k8s cluster..." + + def update_infra_endpoint(self, ip, port): + self.infra.endpoint = "{}:{}".format(ip, port) + self.infra.save() + + def generate_new_vip(self): + ingress = self.ingressprovider.create_ingress( + self.infra, + self.ingressprovider.port, + self.team_name) + return ingress + + def register_ingress_vip(self, ingress): + print('*-----------------------------------*INGRESS*-----------------------------------*') + print(ingress) + try: + original_vip = Vip.objects.get(infra=self.infra) + except Vip.DoesNotExist: + original_vip = None + vip = Vip() + vip.identifier = ingress['port_external'] + vip.infra = self.infra + vip.vip_ip = ingress['ip_external'] + if original_vip: + vip.original_vip = original_vip + vip.save() + print('*-----------------------////////----------------------------*') + return vip + + def insert_dns_on_infra(self, ingress): + if self._vip.vip_ip == ingress['ip_external']: + print('*-----------------------INSERTING DNS----------------------------*') + dns = add_dns_record( + databaseinfra=self.infra, + name=self.infra.name, + ip=ingress['ip_external'], + type=FOXHA, + is_database=False + ) + if dns is None: + return + self.infra.endpoint_dns = "{}:{}".format(dns, ingress['port_external']) + self.infra.save() + + def do(self): + if not self.is_valid: + return + ingress = self.generate_new_vip() + self.update_infra_endpoint(ingress['ip_external'], ingress['port_external']) + self._vip = self.register_ingress_vip(ingress) + self.insert_dns_on_infra(ingress) + return True + + def undo(self): + pass + + +class RegisterDNSIngress(IngressProviderStep): + + def __init__(self, *args, **kw): + super(RegisterDNSIngress, self).__init__(*args, **kw) + self.provider = DNSAPIProvider + self.do_export = True + + def __unicode__(self): + return "Registry dns for VIP..." + + @staticmethod + def is_ipv4(ip): + try: + socket.inet_aton(ip) + return True + except socket.error: + return False + + def _do_database_dns_for_ip(self, func): + if not self.is_valid: + return + try: + ip = self.ingressprovider.ip + if ip is None: + ip = self.instance.address + except Exception as error: + pass + return func( + databaseinfra=self.ingressprovider.infra, + ip=ip, + do_export=self.do_export, + env=self.ingressprovider.infra.environment, + **{'dns_type': 'CNAME'} if not self.is_ipv4(ip) else {} + ) + + def do(self): + if not self.is_valid: + return + self._do_database_dns_for_ip(self.provider.create_database_dns_for_ip) + + def undo(self): + self._do_database_dns_for_ip(self.provider.remove_databases_dns_for_ip)