Extended NSController to handle networks and ports
This commit is contained in:
@@ -15,6 +15,8 @@ from sqlalchemy import or_
|
||||
|
||||
websocket_server_url = "http://127.0.0.1:6000/api/create_task"
|
||||
|
||||
import uuid
|
||||
import uuid
|
||||
import uuid
|
||||
|
||||
def validate_payload(payload):
|
||||
@@ -83,6 +85,44 @@ def validate_payload(payload):
|
||||
raise ValueError("'mem_limit' must be a positive number.")
|
||||
sanitized_container['mem_limit'] = container['mem_limit']
|
||||
|
||||
# Add 'ports' if present and valid
|
||||
if 'ports' in container:
|
||||
if not isinstance(container['ports'], list):
|
||||
raise ValueError("'ports' must be a list.")
|
||||
sanitized_ports = []
|
||||
for port_mapping in container['ports']:
|
||||
if not isinstance(port_mapping, dict):
|
||||
raise ValueError("Each port mapping must be a dictionary with 'internal' and 'external' keys.")
|
||||
if 'internal' not in port_mapping or 'external' not in port_mapping:
|
||||
raise ValueError("Each port mapping must have 'internal' and 'external' keys.")
|
||||
if not isinstance(port_mapping['internal'], int) or not isinstance(port_mapping['external'], int):
|
||||
raise ValueError("'internal' and 'external' ports must be integers.")
|
||||
sanitized_ports.append({
|
||||
'internal': port_mapping['internal'],
|
||||
'external': port_mapping['external']
|
||||
})
|
||||
if sanitized_ports:
|
||||
sanitized_container['ports'] = sanitized_ports
|
||||
|
||||
# Add 'networks' if present and valid
|
||||
if 'networks' in container:
|
||||
if isinstance(container['networks'], str):
|
||||
if container['networks'].strip():
|
||||
sanitized_container['networks'] = container['networks'].strip()
|
||||
elif isinstance(container['networks'], list):
|
||||
sanitized_networks = []
|
||||
for net in container['networks']:
|
||||
if not isinstance(net, str):
|
||||
raise ValueError("Each network name must be a string.")
|
||||
if net.strip():
|
||||
sanitized_networks.append(net.strip())
|
||||
if sanitized_networks:
|
||||
sanitized_container['networks'] = sanitized_networks
|
||||
elif container['networks'] is None:
|
||||
pass # allow it to be missing or null
|
||||
else:
|
||||
raise ValueError("'networks' must be a string or a list of strings.")
|
||||
|
||||
# Add the sanitized container to the sanitized payload
|
||||
sanitized_payload['containers'].append(sanitized_container)
|
||||
|
||||
@@ -96,13 +136,11 @@ def add_container_workload():
|
||||
# Validate incoming data
|
||||
validated_data = validate_payload(request_data)
|
||||
|
||||
# TODO - Check if this user has permissions to VIEW and LAUCNH_CONTAINER in this VDC
|
||||
request_vdc=VirtualDataCenter.query.filter_by(id=uuid.UUID(validated_data['vdc'])).first()
|
||||
|
||||
# TODO - Check if this user has permissions to VIEW and LAUNCH_CONTAINER in this VDC
|
||||
request_vdc = VirtualDataCenter.query.filter_by(id=uuid.UUID(validated_data['vdc'])).first()
|
||||
|
||||
if request_vdc is None:
|
||||
# Raise an error or return a response indicating that no VDC was found
|
||||
abort(404, description="Virtual Data Center not found") # Flask's abort function
|
||||
abort(404, description="Virtual Data Center not found")
|
||||
|
||||
all_workload_hosts = WorkloadHost.query.filter_by(deleted=0, region_id=request_vdc.region.id).all()
|
||||
logger.debug(f"Found {len(all_workload_hosts)} hosts eligible for placement")
|
||||
@@ -114,75 +152,119 @@ def add_container_workload():
|
||||
|
||||
if len(filtered_hosts) == 0:
|
||||
logger.error(f"No hosts found in region {request_vdc.region.name} to place this request after filtering")
|
||||
return f"No hosts found in region {request_vdc.region.name} to place this request", 400
|
||||
# TODO - Is 400 the best error number here? Should be be jsonified?
|
||||
return {"error": f"No hosts found in region {request_vdc.region.name} to place this request"}, 400
|
||||
|
||||
client_response = []
|
||||
for _container in validated_data['containers']:
|
||||
# Select the host with the most available capacity
|
||||
selected_host = filtered_hosts[0]
|
||||
logger.info(f"Selected host: {selected_host}")
|
||||
NSController=None
|
||||
NSController = None
|
||||
|
||||
# Check if we need to deploy an NS Controller to this host for this VPC ID or not
|
||||
# TODO - We will need to update the worker code too.
|
||||
# - If the workder doesnt have an NS controller present for the given VPCID(Currently called a tennacny) then fail the request
|
||||
# Lets try and add the NSController deployment request into the same payload
|
||||
existing_NS_Controller=Workload.query.filter_by(
|
||||
# Check for existing NSController
|
||||
existing_NS_Controller = Workload.query.filter_by(
|
||||
workload_host_id=selected_host.id,
|
||||
workload_type="NSController", # Always set workload_type to "Container"
|
||||
workload_type="NSController",
|
||||
vdc_id=request_vdc.id,
|
||||
deleted=0
|
||||
).first()
|
||||
|
||||
if not existing_NS_Controller:
|
||||
# No Namespace controller exists on this host, lets create one
|
||||
logger.warning(f"No NSController found in on host {selected_host.name} for VDC {request_vdc.name}, adding it to the payload")
|
||||
logger.warning(f"No NSController found on host {selected_host.name} for VDC {request_vdc.name}, creating new one")
|
||||
# TODO - we need to assign a network port to the NS Controller including assigning an IP - How do we specify this in the payload?
|
||||
new_NSController_name=f"NSCONTROLLER_{request_vdc.id}_{selected_host.id}"
|
||||
new_NSController_name = f"NSCONTROLLER_{request_vdc.id}_{selected_host.id}"
|
||||
new_NSController = Workload(
|
||||
name=new_NSController_name,
|
||||
workload_type="NSController", # Always set workload_type to "Container"
|
||||
vdc_id=request_vdc.id,
|
||||
# status="pending-allocation",
|
||||
launch_params=json.dumps({
|
||||
"docker_image": "busybox",
|
||||
"container_name": new_NSController_name,
|
||||
"command": "sleep infinite"
|
||||
}),
|
||||
workload_host_id=selected_host.id
|
||||
)
|
||||
|
||||
name=new_NSController_name,
|
||||
workload_type="NSController",
|
||||
vdc_id=request_vdc.id,
|
||||
launch_params=json.dumps({
|
||||
"docker_image": "busybox",
|
||||
"container_name": new_NSController_name,
|
||||
"command": "sleep infinite"
|
||||
}),
|
||||
workload_host_id=selected_host.id
|
||||
)
|
||||
db.session.add(new_NSController)
|
||||
db.session.commit()
|
||||
new_NSController.set_status("pending-allocation")
|
||||
NSController=new_NSController
|
||||
NSController = new_NSController
|
||||
else:
|
||||
NSController=existing_NS_Controller
|
||||
NSController = existing_NS_Controller
|
||||
|
||||
|
||||
# Create the Workload instance in the database
|
||||
# Create the new workload container
|
||||
new_container = Workload(
|
||||
name=_container['container_name'],
|
||||
workload_type="Container", # Always set workload_type to "Container"
|
||||
workload_type="Container",
|
||||
vdc_id=request_vdc.id,
|
||||
# status="pending-allocation",
|
||||
launch_params=json.dumps(_container),
|
||||
workload_host_id=selected_host.id
|
||||
)
|
||||
|
||||
db.session.add(new_container)
|
||||
db.session.commit()
|
||||
new_container.set_status("pending-allocation")
|
||||
logger.debug(f"Container {new_container.id} workload added to DB")
|
||||
|
||||
# Inject the container ID so that it can be processed by the workloadHost
|
||||
_container['container_id']=new_container.id
|
||||
_container['NSController_launchparams']=json.loads(NSController.launch_params)
|
||||
_container['NSController_launchparams']['container_id']=NSController.id
|
||||
|
||||
logger.debug(f"Container spec: {_container}")
|
||||
#Send the request off to the websocket server to have this container launched
|
||||
# Inject container_id and NSController launchparams
|
||||
_container['container_id'] = new_container.id
|
||||
nscontroller_launchparams = json.loads(NSController.launch_params)
|
||||
nscontroller_launchparams['container_id'] = NSController.id
|
||||
|
||||
# === NEW PORT + NETWORK HANDLING ===
|
||||
nscontroller_modified = False
|
||||
|
||||
# Handle ports
|
||||
if 'ports' in _container:
|
||||
logger.info(f"Transferring ports from workload container {_container['container_name']} to NSController")
|
||||
|
||||
if 'ports' not in nscontroller_launchparams:
|
||||
nscontroller_launchparams['ports'] = []
|
||||
|
||||
existing_ports = nscontroller_launchparams['ports']
|
||||
new_ports = _container['ports']
|
||||
|
||||
for new_port in new_ports:
|
||||
if new_port not in existing_ports:
|
||||
existing_ports.append(new_port)
|
||||
nscontroller_modified = True
|
||||
|
||||
_container.pop('ports') # Remove ports from the container
|
||||
|
||||
# Handle networks
|
||||
if 'networks' in _container:
|
||||
logger.info(f"Transferring networks from workload container {_container['container_name']} to NSController")
|
||||
|
||||
if 'networks' not in nscontroller_launchparams:
|
||||
nscontroller_launchparams['networks'] = []
|
||||
|
||||
# Normalize into a list if it's a string
|
||||
if isinstance(_container['networks'], str):
|
||||
new_networks = [_container['networks']]
|
||||
elif isinstance(_container['networks'], list):
|
||||
new_networks = _container['networks']
|
||||
else:
|
||||
new_networks = []
|
||||
|
||||
existing_networks = nscontroller_launchparams['networks']
|
||||
|
||||
for new_net in new_networks:
|
||||
if new_net not in existing_networks:
|
||||
existing_networks.append(new_net)
|
||||
nscontroller_modified = True
|
||||
|
||||
_container.pop('networks') # Remove networks from the container
|
||||
|
||||
# Only update DB if NSController ports or networks changed
|
||||
if nscontroller_modified:
|
||||
NSController.launch_params = json.dumps(nscontroller_launchparams)
|
||||
db.session.add(NSController)
|
||||
db.session.commit()
|
||||
logger.debug(f"Updated NSController {NSController.id} with new ports and/or networks")
|
||||
|
||||
_container['NSController_launchparams'] = nscontroller_launchparams
|
||||
|
||||
|
||||
logger.debug(f"Container spec after port transfer: {_container}")
|
||||
|
||||
# Send task to websocket server
|
||||
payload = {
|
||||
"worker_id": selected_host.id,
|
||||
"task_type": "container-create",
|
||||
@@ -193,13 +275,12 @@ def add_container_workload():
|
||||
}
|
||||
logger.debug(f"Sending payload to websocket server {payload}")
|
||||
headers = {"Content-Type": "application/json"}
|
||||
|
||||
|
||||
websocket_server_response = requests.post(websocket_server_url, data=json.dumps(payload), headers=headers)
|
||||
websocket_server_response_data = websocket_server_response.json()
|
||||
|
||||
if websocket_server_response.status_code == 201:
|
||||
logger.info(f"Task created successfully for container {new_container.id}")
|
||||
|
||||
new_container.set_status("allocated")
|
||||
db.session.add(new_container)
|
||||
db.session.commit()
|
||||
@@ -213,14 +294,15 @@ def add_container_workload():
|
||||
db.session.add(new_container)
|
||||
db.session.commit()
|
||||
|
||||
# Prepare the response to the API request
|
||||
# Prepare API client response
|
||||
client_response.append({
|
||||
"container_id": new_container.id,
|
||||
"container_name": new_container.name,
|
||||
"container_status": new_container.status,
|
||||
"container_id": new_container.id,
|
||||
"container_name": new_container.name,
|
||||
"container_status": new_container.status,
|
||||
})
|
||||
|
||||
return client_response,200
|
||||
return client_response, 200
|
||||
|
||||
|
||||
# @api_bp.route('/workloads/containers/<workload_id>', methods=['PUT'])
|
||||
# def edit_container_workload(workload_id):
|
||||
|
||||
@@ -292,6 +292,7 @@ class Workload(BaseModel):
|
||||
|
||||
# Handle specific status transitions
|
||||
if new_status == "offline" and self.workload_type == "container":
|
||||
logger.debug("Container failed - I should do something!")
|
||||
# TODO - This should not trigger instantly, it should add an event to a queue and wait a pre-determined amount of time before attempting
|
||||
# Normal container lifecycle will see a container go from running->dead->deleted when the container goes through a deletion
|
||||
# EventHandlers.handle_failed_container(changed_by)
|
||||
|
||||
@@ -49,6 +49,11 @@ def render_list():
|
||||
{
|
||||
"docker_image": docker_image,
|
||||
"container_name": container_name,
|
||||
"ports": [
|
||||
{"internal": 80, "external": 8080},
|
||||
{"internal": 443, "external": 8443}
|
||||
],
|
||||
"networks": "bridge"
|
||||
}
|
||||
]
|
||||
})
|
||||
@@ -168,11 +173,22 @@ def render_detail():
|
||||
st.markdown(f"**Last Updated:** {format_timestamp(container['updated_at'])}")
|
||||
st.markdown(f"**Status:** {container['status']}")
|
||||
st.markdown(f"**Launch Parameters:** {container.get('launch_params', 'N/A')}")
|
||||
|
||||
# Define container states once
|
||||
CONTAINER_STATES = [
|
||||
"running", "stopped", "pending",
|
||||
"failed-deleted", "pending-deleted",
|
||||
"pending-allocation", "failed-allocation",
|
||||
"pending-allocated", "dead"
|
||||
]
|
||||
|
||||
# Edit Container Workload Form
|
||||
with st.expander("Edit Container Workload"):
|
||||
with st.form("edit_container"):
|
||||
new_status = st.selectbox("Status", ["running", "stopped", "pending","failed-deleted","pending-deleted","pending-allocation","failed-allocation","pending-allocated","dead"], index=["running", "stopped", "pending","failed-deleted","pending-deleted","pending-allocation","failed-allocation","pending-allocated","dead"].index(container['status']))
|
||||
new_status = st.selectbox(
|
||||
"Status",
|
||||
CONTAINER_STATES,
|
||||
index=CONTAINER_STATES.index(container['status'])
|
||||
)
|
||||
new_launch_params = st.text_input("Launch Parameters", value=container.get('launch_params', ''))
|
||||
submit = st.form_submit_button("Update Container Workload")
|
||||
|
||||
@@ -187,8 +203,6 @@ def render_detail():
|
||||
st.session_state.view_type = 'detail'
|
||||
st.session_state.selected_resource_type = 'container'
|
||||
st.rerun()
|
||||
|
||||
|
||||
|
||||
# Delete functionality
|
||||
if not st.session_state.container_delete_confirmation_shown:
|
||||
|
||||
+46
-10
@@ -129,6 +129,7 @@ class ContainerTask:
|
||||
container_image += ':latest' # Assume 'latest' if no tag is specified
|
||||
|
||||
# Check if the existing container matches the payload configuration
|
||||
# TODO - We are only matchin on CPU, RAM and image here, we should be more detailed with the config comparison
|
||||
cpu_matches = existing_container.attrs['HostConfig']['CpuShares'] == container['cpu_shares'] * 1024
|
||||
memory_matches = existing_container.attrs['HostConfig']['Memory'] == container['mem_limit'] * 1024 * 1024
|
||||
image_matches = container_image == payload_image
|
||||
@@ -153,28 +154,63 @@ class ContainerTask:
|
||||
#Check if the NSController container exists
|
||||
|
||||
|
||||
nscontroller_container_name=container['NSController_launchparams']['container_name']
|
||||
# Check if the NSController container exists
|
||||
nscontroller_container_name = container['NSController_launchparams']['container_name']
|
||||
self.logger.debug(f"Name of NSController is {nscontroller_container_name}")
|
||||
running_NSController = client.containers.list(all=True, filters={"name": nscontroller_container_name})
|
||||
self.logger.debug(f"{running_NSController}")
|
||||
|
||||
if not running_NSController:
|
||||
self.logger.info("NSController does not exist, creating it")
|
||||
ns_params = container['NSController_launchparams']
|
||||
|
||||
nscontroller_container_config = {
|
||||
"image": container['NSController_launchparams']['docker_image'],
|
||||
"name": container['NSController_launchparams']['container_name'],
|
||||
"command": container['NSController_launchparams']['command'],
|
||||
"image": ns_params['docker_image'],
|
||||
"name": ns_params['container_name'],
|
||||
"command": ns_params.get('command'), # safer with .get()
|
||||
"dns": ["1.1.1.1"],
|
||||
# "network": "none",
|
||||
"detach": True,
|
||||
"labels": {
|
||||
"managed_by": "worker_agent",
|
||||
"system_container_id": container['NSController_launchparams']['container_id']
|
||||
}
|
||||
"system_container_id": ns_params['container_id']
|
||||
}
|
||||
}
|
||||
# Optional: Attach network
|
||||
primary_network = None
|
||||
additional_networks = []
|
||||
|
||||
if 'networks' in ns_params:
|
||||
network_param = ns_params['networks']
|
||||
|
||||
if isinstance(network_param, str):
|
||||
primary_network = network_param
|
||||
elif isinstance(network_param, list) and network_param:
|
||||
primary_network = network_param[0]
|
||||
additional_networks = network_param[1:]
|
||||
else:
|
||||
raise ValueError("Invalid 'network' field: must be string or non-empty list of strings.")
|
||||
|
||||
nscontroller_container_config["network"] = primary_network
|
||||
|
||||
# Only allow port mapping if the primary network is "bridge"
|
||||
if primary_network == 'bridge' and 'ports' in ns_params:
|
||||
ports_mapping = {}
|
||||
for mapping in ns_params['ports']:
|
||||
internal = mapping.get('internal')
|
||||
external = mapping.get('external')
|
||||
if internal and external:
|
||||
ports_mapping[internal] = external
|
||||
if ports_mapping:
|
||||
nscontroller_container_config["ports"] = ports_mapping
|
||||
|
||||
# Launch the NSController container
|
||||
ns_container = client.containers.run(**nscontroller_container_config)
|
||||
|
||||
# After launch: connect to any additional networks
|
||||
for net_name in additional_networks:
|
||||
client.networks.get(net_name).connect(ns_container)
|
||||
self.logger.info(f"NSController container '{nscontroller_container_name}' launched successfully.")
|
||||
|
||||
# Launch the container
|
||||
client.containers.run(**nscontroller_container_config)
|
||||
self.logger.info(f"Container '{nscontroller_container_config}' launched successfully.")
|
||||
|
||||
self.logger.info(f"Launching container '{container_name}'...")
|
||||
# Prepare container configuration
|
||||
|
||||
Reference in New Issue
Block a user