pycodestyle fixes: vmware_plugin.py

Signed-off-by: Sandro Bonazzola <sbonazzo@redhat.com>
This commit is contained in:
Sandro Bonazzola
2022-09-05 14:15:38 +02:00
committed by Sandro Bonazzola
parent b444854cb2
commit ec807e3b3a
+228 -109
View File
@@ -1,25 +1,23 @@
#!/usr/bin/env python
import logging
import random
import sys
import time
import typing
from os import environ
from dataclasses import dataclass, field
import random
from os import environ
from traceback import format_exc
import logging
from kraken.plugins.vmware import kubernetes_functions as kube_helper
from com.vmware.vcenter_client import ResourcePool
from arcaflow_plugin_sdk import validation, plugin
import requests
from arcaflow_plugin_sdk import plugin, validation
from com.vmware.vapi.std.errors_client import (AlreadyInDesiredState,
NotAllowedInCurrentState)
from com.vmware.vcenter.vm_client import Power
from com.vmware.vcenter_client import VM, ResourcePool
from kubernetes import client, watch
from vmware.vapi.vsphere.client import create_vsphere_client
from com.vmware.vcenter_client import VM
from com.vmware.vcenter.vm_client import Power
from com.vmware.vapi.std.errors_client import (
AlreadyInDesiredState,
NotAllowedInCurrentState,
)
import requests
import sys
from kraken.plugins.vmware import kubernetes_functions as kube_helper
class vSphere:
@@ -37,7 +35,9 @@ class vSphere:
)
if not self.credentials_present:
raise Exception(
"Environmental variables 'VSPHERE_IP', 'VSPHERE_USERNAME', 'VSPHERE_PASSWORD' are not set"
"Environmental variables "
"'VSPHERE_IP', 'VSPHERE_USERNAME', "
"'VSPHERE_PASSWORD' are not set"
)
self.client = create_vsphere_client(
server=self.server,
@@ -66,7 +66,7 @@ class vSphere:
vms = self.client.vcenter.VM.list(VM.FilterSpec(names=names))
if len(vms) == 0:
logging.info("VM with name ({}) not found".format(instance_id))
logging.info("VM with name ({}) not found", instance_id)
return None
vm = vms[0].vm
@@ -90,75 +90,81 @@ class vSphere:
self.client.vcenter.vm.Power.start(vm)
self.client.vcenter.vm.Power.stop(vm)
self.client.vcenter.VM.delete(vm)
logging.info("Deleted VM -- '{}-({})'".format(instance_id, vm))
logging.info("Deleted VM -- '{}-({})'", instance_id, vm)
def reboot_instances(self, instance_id):
"""
Reboots the VM whose name is given by 'instance_id'. Returns True if successful, or
returns False if the VM is not powered on
Reboots the VM whose name is given by 'instance_id'.
@Returns: True if successful, or False if the VM is not powered on
"""
vm = self.get_vm(instance_id)
try:
self.client.vcenter.vm.Power.reset(vm)
logging.info("Reset VM -- '{}-({})'".format(instance_id, vm))
logging.info("Reset VM -- '{}-({})'", instance_id, vm)
return True
except NotAllowedInCurrentState:
logging.info(
"VM '{}'-'({})' is not Powered On. Cannot reset it".format(
instance_id, vm
)
"VM '{}'-'({})' is not Powered On. Cannot reset it",
instance_id,
vm
)
return False
def stop_instances(self, instance_id):
"""
Stops the VM whose name is given by 'instance_id'. Returns True if successful, or
returns False if the VM is already powered off
Stops the VM whose name is given by 'instance_id'.
@Returns: True if successful, or False if the VM is already powered off
"""
vm = self.get_vm(instance_id)
try:
self.client.vcenter.vm.Power.stop(vm)
logging.info("Stopped VM -- '{}-({})'".format(instance_id, vm))
logging.info("Stopped VM -- '{}-({})'", instance_id, vm)
return True
except AlreadyInDesiredState:
logging.info(
"VM '{}'-'({})' is already Powered Off".format(instance_id, vm)
"VM '{}'-'({})' is already Powered Off", instance_id, vm
)
return False
def start_instances(self, instance_id):
"""
Stops the VM whose name is given by 'instance_id'. Returns True if successful, or
returns False if the VM is already powered on
Stops the VM whose name is given by 'instance_id'.
@Returns: True if successful, or False if the VM is already powered on
"""
vm = self.get_vm(instance_id)
try:
self.client.vcenter.vm.Power.start(vm)
logging.info("Started VM -- '{}-({})'".format(instance_id, vm))
logging.info("Started VM -- '{}-({})'", instance_id, vm)
return True
except AlreadyInDesiredState:
logging.info("VM '{}'-'({})' is already Powered On".format(instance_id, vm))
logging.info(
"VM '{}'-'({})' is already Powered On", instance_id, vm
)
return False
def list_instances(self, datacenter):
"""
Returns a list of VMs present in the datacenter
@Returns: a list of VMs present in the datacenter
"""
datacenter_filter = self.client.vcenter.Datacenter.FilterSpec(
names=set([datacenter])
)
datacenter_summaries = self.client.vcenter.Datacenter.list(datacenter_filter)
datacenter_summaries = self.client.vcenter.Datacenter.list(
datacenter_filter
)
try:
datacenter_id = datacenter_summaries[0].datacenter
except IndexError:
logging.error("Datacenter '{}' doesn't exist".format(datacenter))
logging.error("Datacenter '{}' doesn't exist", datacenter)
sys.exit(1)
vm_filter = self.client.vcenter.VM.FilterSpec(datacenters={datacenter_id})
vm_filter = self.client.vcenter.VM.FilterSpec(
datacenters={datacenter_id}
)
vm_summaries = self.client.vcenter.VM.list(vm_filter)
vm_names = []
for vm in vm_summaries:
@@ -172,33 +178,45 @@ class vSphere:
datacenter_summaries = self.client.vcenter.Datacenter.list()
datacenter_names = [
{"datacenter_id": datacenter.datacenter, "datacenter_name": datacenter.name}
{
"datacenter_id": datacenter.datacenter,
"datacenter_name": datacenter.name
}
for datacenter in datacenter_summaries
]
return datacenter_names
def get_datastore_list(self, datacenter=None):
"""
Returns a dictionary containing all the datastore names and IDs belonging to a specific datacenter
@Returns: a dictionary containing all the datastore names and
IDs belonging to a specific datacenter
"""
datastore_filter = self.client.vcenter.Datastore.FilterSpec(
datacenters={datacenter}
)
datastore_summaries = self.client.vcenter.Datastore.list(datastore_filter)
datastore_summaries = self.client.vcenter.Datastore.list(
datastore_filter
)
datastore_names = []
for datastore in datastore_summaries:
datastore_names.append(
{"datastore_name": datastore.name, "datastore_id": datastore.datastore}
{
"datastore_name": datastore.name,
"datastore_id": datastore.datastore
}
)
return datastore_names
def get_folder_list(self, datacenter=None):
"""
Returns a dictionary containing all the folder names and IDs belonging to a specific datacenter
@Returns: a dictionary containing all the folder names and
IDs belonging to a specific datacenter
"""
folder_filter = self.client.vcenter.Folder.FilterSpec(datacenters={datacenter})
folder_filter = self.client.vcenter.Folder.FilterSpec(
datacenters={datacenter}
)
folder_summaries = self.client.vcenter.Folder.list(folder_filter)
folder_names = []
for folder in folder_summaries:
@@ -217,28 +235,31 @@ class vSphere:
filter_spec = ResourcePool.FilterSpec(
datacenters=set([datacenter]), names=names
)
resource_pool_summaries = self.client.vcenter.ResourcePool.list(filter_spec)
resource_pool_summaries = self.client.vcenter.ResourcePool.list(
filter_spec
)
if len(resource_pool_summaries) > 0:
resource_pool = resource_pool_summaries[0].resource_pool
return resource_pool
else:
logging.error(
"ResourcePool not found in Datacenter '{}'".format(datacenter)
"ResourcePool not found in Datacenter '{}'",
datacenter
)
return None
def create_default_vm(self, guest_os="RHEL_7_64", max_attempts=10):
"""
Creates a default VM with 2 GB memory, 1 CPU and 16 GB disk space in a
random datacenter. Accepts the guest OS as a parameter. Since the VM placement
is random, it might fail due to resource constraints. So, this function tries
for upto 'max_attempts' to create the VM
random datacenter. Accepts the guest OS as a parameter. Since the VM
placement is random, it might fail due to resource constraints.
So, this function tries for upto 'max_attempts' to create the VM
"""
def create_vm(vm_name, resource_pool, folder, datastore, guest_os):
"""
Creates a VM and returns its ID and name. Requires the VM name, resource
pool name, folder name, datastore and the guest OS
Creates a VM and returns its ID and name. Requires the VM name,
resource pool name, folder name, datastore and the guest OS
"""
placement_spec = VM.PlacementSpec(
@@ -254,25 +275,39 @@ class vSphere:
for _ in range(max_attempts):
try:
datacenter_list = self.get_datacenter_list()
datacenter = random.choice(datacenter_list)
resource_pool = self.get_resource_pool(datacenter["datacenter_id"])
folder = random.choice(
# random generator not used for
# security/cryptographic purposes in this loop
datacenter = random.choice(datacenter_list) # nosec
resource_pool = self.get_resource_pool(
datacenter["datacenter_id"]
)
folder = random.choice( # nosec
self.get_folder_list(datacenter["datacenter_id"])
)["folder_id"]
datastore = random.choice(
datastore = random.choice( # nosec
self.get_datastore_list(datacenter["datacenter_id"])
)["datastore_id"]
vm_name = "Test-" + str(time.time_ns())
return (
create_vm(vm_name, resource_pool, folder, datastore, guest_os),
create_vm(
vm_name,
resource_pool,
folder,
datastore,
guest_os
),
vm_name,
)
except Exception as e:
pass
logging.error(
"Default VM could not be created, retrying. "
"Error was: %s",
str(e)
)
logging.error(
"Default VM could not be created in {} attempts. Check your VMware resources".format(
max_attempts
)
"Default VM could not be created in %s attempts. "
"Check your VMware resources",
max_attempts
)
return None, None
@@ -289,7 +324,7 @@ class vSphere:
except Exception as e:
logging.error(
"Failed to get node instance status %s. Encountered following "
"exception: %s." % (instance_id, e)
"exception: %s.", instance_id, e
)
return None
@@ -304,13 +339,16 @@ class vSphere:
while vm is not None:
vm = self.get_vm(instance_id)
logging.info(
"VM %s is still being deleted, sleeping for 5 seconds" % instance_id
"VM %s is still being deleted, "
"sleeping for 5 seconds",
instance_id
)
time.sleep(5)
time_counter += 5
if time_counter >= timeout:
logging.info(
"VM %s is still not deleted in allotted time" % instance_id
"VM %s is still not deleted in allotted time",
instance_id
)
return False
return True
@@ -326,12 +364,17 @@ class vSphere:
while status != Power.State.POWERED_ON:
status = self.get_vm_status(instance_id)
logging.info(
"VM %s is still not running, sleeping for 5 seconds" % instance_id
"VM %s is still not running, "
"sleeping for 5 seconds",
instance_id
)
time.sleep(5)
time_counter += 5
if time_counter >= timeout:
logging.info("VM %s is still not ready in allotted time" % instance_id)
logging.info(
"VM %s is still not ready in allotted time",
instance_id
)
return False
return True
@@ -346,12 +389,17 @@ class vSphere:
while status != Power.State.POWERED_OFF:
status = self.get_vm_status(instance_id)
logging.info(
"VM %s is still not running, sleeping for 5 seconds" % instance_id
"VM %s is still not running, "
"sleeping for 5 seconds",
instance_id
)
time.sleep(5)
time_counter += 5
if time_counter >= timeout:
logging.info("VM %s is still not ready in allotted time" % instance_id)
logging.info(
"VM %s is still not ready in allotted time",
instance_id
)
return False
return True
@@ -367,15 +415,17 @@ class NodeScenarioSuccessOutput:
nodes: typing.Dict[int, Node] = field(
metadata={
"name": "Nodes started/stopped/terminated/rebooted",
"description": """Map between timestamps and the pods started/stopped/terminated/rebooted.
The timestamp is provided in nanoseconds""",
"description": "Map between timestamps and the pods "
"started/stopped/terminated/rebooted. "
"The timestamp is provided in nanoseconds",
}
)
action: kube_helper.Actions = field(
metadata={
"name": "The action performed on the node",
"description": """The action performed or attempted to be performed on the node. Possible values
are : Start, Stop, Terminate, Reboot""",
"description": "The action performed or attempted to be "
"performed on the node. Possible values"
"are : Start, Stop, Terminate, Reboot",
}
)
@@ -387,8 +437,8 @@ class NodeScenarioErrorOutput:
action: kube_helper.Actions = field(
metadata={
"name": "The action performed on the node",
"description": """The action attempted to be performed on the node. Possible values are : Start
Stop, Terminate, Reboot""",
"description": "The action attempted to be performed on the node. "
"Possible values are : Start Stop, Terminate, Reboot",
}
)
@@ -404,7 +454,8 @@ class NodeScenarioConfig:
default=None,
metadata={
"name": "Name",
"description": "Name(s) for target nodes. Required if label_selector is not set.",
"description": "Name(s) for target nodes. "
"Required if label_selector is not set.",
},
)
@@ -412,18 +463,23 @@ class NodeScenarioConfig:
default=1,
metadata={
"name": "Number of runs per node",
"description": "Number of times to inject each scenario under actions (will perform on same node each time)",
"description": "Number of times to inject each scenario under "
"actions (will perform on same node each time)",
},
)
label_selector: typing.Annotated[
typing.Optional[str], validation.min(1), validation.required_if_not("name")
typing.Optional[str],
validation.min(1),
validation.required_if_not("name")
] = field(
default=None,
metadata={
"name": "Label selector",
"description": "Kubernetes label selector for the target nodes. Required if name is not set.\n"
"See https://kubernetes.io/docs/concepts/overview/working-with-objects/labels/ for details.",
"description": "Kubernetes label selector for the target nodes. "
"Required if name is not set.\n"
"See https://kubernetes.io/docs/concepts/overview/working-with-objects/labels/ " # noqa
"for details.",
},
)
@@ -431,15 +487,20 @@ class NodeScenarioConfig:
default=180,
metadata={
"name": "Timeout",
"description": "Timeout to wait for the target pod(s) to be removed in seconds.",
"description": "Timeout to wait for the target pod(s) "
"to be removed in seconds.",
},
)
instance_count: typing.Annotated[typing.Optional[int], validation.min(1)] = field(
instance_count: typing.Annotated[
typing.Optional[int],
validation.min(1)
] = field(
default=1,
metadata={
"name": "Instance Count",
"description": "Number of nodes to perform action/select that match the label selector.",
"description": "Number of nodes to perform action/select "
"that match the label selector.",
},
)
@@ -455,7 +516,8 @@ class NodeScenarioConfig:
default=True,
metadata={
"name": "Verify API Session",
"description": "Verifies the vSphere client session. It is enabled by default",
"description": "Verifies the vSphere client session. "
"It is enabled by default",
},
)
@@ -463,9 +525,10 @@ class NodeScenarioConfig:
default=None,
metadata={
"name": "Kubeconfig path",
"description": "Path to your Kubeconfig file. Defaults to ~/.kube/config.\n"
"See https://kubernetes.io/docs/concepts/configuration/organize-cluster-access-kubeconfig/ for "
"details.",
"description": "Path to your Kubeconfig file. "
"Defaults to ~/.kube/config.\n"
"See https://kubernetes.io/docs/concepts/configuration/organize-cluster-access-kubeconfig/ " # noqa
"for details.",
},
)
@@ -473,8 +536,12 @@ class NodeScenarioConfig:
@plugin.step(
id="node-start",
name="Start the node",
description="Start the node(s) by starting the VMware VM on which the node is configured",
outputs={"success": NodeScenarioSuccessOutput, "error": NodeScenarioErrorOutput},
description="Start the node(s) by starting the VMware VM "
"on which the node is configured",
outputs={
"success": NodeScenarioSuccessOutput,
"error": NodeScenarioErrorOutput
},
)
def node_start(
cfg: NodeScenarioConfig,
@@ -485,13 +552,17 @@ def node_start(
vsphere = vSphere(verify=cfg.verify_session)
core_v1 = client.CoreV1Api(cli)
watch_resource = watch.Watch()
node_list = kube_helper.get_node_list(cfg, kube_helper.Actions.START, core_v1)
node_list = kube_helper.get_node_list(
cfg,
kube_helper.Actions.START,
core_v1
)
nodes_started = {}
for name in node_list:
try:
for _ in range(cfg.runs):
logging.info("Starting node_start_scenario injection")
logging.info("Starting the node %s " % (name))
logging.info("Starting the node %s ", name)
vm_started = vsphere.start_instances(name)
if vm_started:
vsphere.wait_until_running(name, cfg.timeout)
@@ -500,11 +571,18 @@ def node_start(
name, cfg.timeout, watch_resource, core_v1
)
nodes_started[int(time.time_ns())] = Node(name=name)
logging.info("Node with instance ID: %s is in running state" % name)
logging.info("node_start_scenario has been successfully injected!")
logging.info(
"Node with instance ID: %s is in running state", name
)
logging.info(
"node_start_scenario has been successfully injected!"
)
except Exception as e:
logging.error("Failed to start node instance. Test Failed")
logging.error("node_start_scenario injection failed!")
logging.error(
"node_start_scenario injection failed! "
"Error was: %s", str(e)
)
return "error", NodeScenarioErrorOutput(
format_exc(), kube_helper.Actions.START
)
@@ -517,8 +595,12 @@ def node_start(
@plugin.step(
id="node-stop",
name="Stop the node",
description="Stop the node(s) by starting the VMware VM on which the node is configured",
outputs={"success": NodeScenarioSuccessOutput, "error": NodeScenarioErrorOutput},
description="Stop the node(s) by starting the VMware VM "
"on which the node is configured",
outputs={
"success": NodeScenarioSuccessOutput,
"error": NodeScenarioErrorOutput
},
)
def node_stop(
cfg: NodeScenarioConfig,
@@ -529,13 +611,17 @@ def node_stop(
vsphere = vSphere(verify=cfg.verify_session)
core_v1 = client.CoreV1Api(cli)
watch_resource = watch.Watch()
node_list = kube_helper.get_node_list(cfg, kube_helper.Actions.STOP, core_v1)
node_list = kube_helper.get_node_list(
cfg,
kube_helper.Actions.STOP,
core_v1
)
nodes_stopped = {}
for name in node_list:
try:
for _ in range(cfg.runs):
logging.info("Starting node_stop_scenario injection")
logging.info("Stopping the node %s " % (name))
logging.info("Stopping the node %s ", name)
vm_stopped = vsphere.stop_instances(name)
if vm_stopped:
vsphere.wait_until_stopped(name, cfg.timeout)
@@ -544,11 +630,18 @@ def node_stop(
name, cfg.timeout, watch_resource, core_v1
)
nodes_stopped[int(time.time_ns())] = Node(name=name)
logging.info("Node with instance ID: %s is in stopped state" % name)
logging.info("node_stop_scenario has been successfully injected!")
logging.info(
"Node with instance ID: %s is in stopped state", name
)
logging.info(
"node_stop_scenario has been successfully injected!"
)
except Exception as e:
logging.error("Failed to stop node instance. Test Failed")
logging.error("node_stop_scenario injection failed!")
logging.error(
"node_stop_scenario injection failed! "
"Error was: %s", str(e)
)
return "error", NodeScenarioErrorOutput(
format_exc(), kube_helper.Actions.STOP
)
@@ -561,8 +654,12 @@ def node_stop(
@plugin.step(
id="node-reboot",
name="Reboot VMware VM",
description="Reboot the node(s) by starting the VMware VM on which the node is configured",
outputs={"success": NodeScenarioSuccessOutput, "error": NodeScenarioErrorOutput},
description="Reboot the node(s) by starting the VMware VM "
"on which the node is configured",
outputs={
"success": NodeScenarioSuccessOutput,
"error": NodeScenarioErrorOutput
},
)
def node_reboot(
cfg: NodeScenarioConfig,
@@ -573,13 +670,17 @@ def node_reboot(
vsphere = vSphere(verify=cfg.verify_session)
core_v1 = client.CoreV1Api(cli)
watch_resource = watch.Watch()
node_list = kube_helper.get_node_list(cfg, kube_helper.Actions.REBOOT, core_v1)
node_list = kube_helper.get_node_list(
cfg,
kube_helper.Actions.REBOOT,
core_v1
)
nodes_rebooted = {}
for name in node_list:
try:
for _ in range(cfg.runs):
logging.info("Starting node_reboot_scenario injection")
logging.info("Rebooting the node %s " % (name))
logging.info("Rebooting the node %s ", name)
vsphere.reboot_instances(name)
if not cfg.skip_openshift_checks:
kube_helper.wait_for_unknown_status(
@@ -590,12 +691,18 @@ def node_reboot(
)
nodes_rebooted[int(time.time_ns())] = Node(name=name)
logging.info(
"Node with instance ID: %s has rebooted successfully" % name
"Node with instance ID: %s has rebooted "
"successfully", name
)
logging.info(
"node_reboot_scenario has been successfully injected!"
)
logging.info("node_reboot_scenario has been successfully injected!")
except Exception as e:
logging.error("Failed to reboot node instance. Test Failed")
logging.error("node_reboot_scenario injection failed!")
logging.error(
"node_reboot_scenario injection failed! "
"Error was: %s", str(e)
)
return "error", NodeScenarioErrorOutput(
format_exc(), kube_helper.Actions.REBOOT
)
@@ -609,7 +716,10 @@ def node_reboot(
id="node-terminate",
name="Reboot VMware VM",
description="Wait for the specified number of pods to be present",
outputs={"success": NodeScenarioSuccessOutput, "error": NodeScenarioErrorOutput},
outputs={
"success": NodeScenarioSuccessOutput,
"error": NodeScenarioErrorOutput
},
)
def node_terminate(
cfg: NodeScenarioConfig,
@@ -627,21 +737,30 @@ def node_terminate(
try:
for _ in range(cfg.runs):
logging.info(
"Starting node_termination_scenario injection by first stopping the node"
"Starting node_termination_scenario injection "
"by first stopping the node"
)
vsphere.stop_instances(name)
vsphere.wait_until_stopped(name, cfg.timeout)
logging.info("Releasing the node with instance ID: %s " % (name))
logging.info(
"Releasing the node with instance ID: %s ", name
)
vsphere.release_instances(name)
vsphere.wait_until_released(name, cfg.timeout)
nodes_terminated[int(time.time_ns())] = Node(name=name)
logging.info("Node with instance ID: %s has been released" % name)
logging.info(
"node_terminate_scenario has been successfully injected!"
"Node with instance ID: %s has been released", name
)
logging.info(
"node_terminate_scenario has been "
"successfully injected!"
)
except Exception as e:
logging.error("Failed to terminate node instance. Test Failed")
logging.error("node_terminate_scenario injection failed!")
logging.error(
"node_terminate_scenario injection failed! "
"Error was: %s", str(e)
)
return "error", NodeScenarioErrorOutput(
format_exc(), kube_helper.Actions.TERMINATE
)