refactor(workload): consolidate deletion logic into shared utilities

Extract duplicate deletion code from pod and container routes into
reusable helper functions in deletion_helpers.py module:
- collect_and_cleanup_network_ports() for network port cleanup
- cleanup_workload_related_records() for resource/volume cleanup
- log_deletion_audit_event() for consistent audit logging

Add SDN network port updates during pod deletion to ensure proper
network cleanup when containers are removed.

Optimize get_pods() endpoint by replacing N+1 queries with batch
queries and lookup maps for nscontrollers, tunnels, and DNS records.

Enhance get_pod() response to include cloudflare tunnel configuration
and DNS records for better visibility into pod networking setup.
This commit is contained in:
2026-01-05 16:24:37 +10:30
parent d1fbd4eaec
commit 016f1fc447
3 changed files with 197 additions and 64 deletions
@@ -2,14 +2,21 @@
from flask import json, request, abort
import requests
from app import app, db, logger
from app.models.models import PortForwarding, Workload, WorkloadHostPortMapping, ContainerPod, CloudflareDNSRecord, CloudflareTunnel, AuditEntry
from app.models.models import PortForwarding, Workload, WorkloadHostPortMapping, ContainerPod, CloudflareDNSRecord, CloudflareTunnel, AuditEntry, VolumeWorkloadMapping
from app.controller import api_bp
from collections import defaultdict
from app.services.cloudflare import CloudflareTunnelManager
from app.utils.standard_responses import api_response
from app.utils.auth_utils import get_request_user_id
from app.utils.create_workload_container import build_pod_payload
from app.utils.deletion_handler import mark_pod_pending_deleted, mark_workload_pending_deleted
from app.utils.dns_helpers import send_dns_updates_for_vdc
from app.utils.sdn_helpers import send_sdn_updates_for_networks
from app.utils.deletion_helpers import (
collect_and_cleanup_network_ports,
log_deletion_audit_event,
cleanup_workload_related_records
)
@api_bp.route('/workloads/pods/<pod_id>/lifecycle/<action>', methods=['POST'])
@@ -97,6 +104,7 @@ def delete_pod(pod_id):
from app.utils.container_deleted import check_deleted_container
container_ids: list[str] = []
all_networks: list[str] = []
# Mark all containers as pending-deleted and trigger cleanup
all_containers_in_pod=Workload.query.filter_by(pod_id=pod.id).all()
for _container in all_containers_in_pod:
@@ -107,6 +115,9 @@ def delete_pod(pod_id):
mark_workload_pending_deleted(_container)
# This will soft-delete related resources; if this is the last active container, it will soft-delete the pod.
check_deleted_container(_container)
# Collect network ports for SDN update
container_networks = collect_and_cleanup_network_ports(str(_container.id))
all_networks.extend(container_networks)
except Exception as exc:
logger.error(f"Cleanup failed for container {_container.id} during unscheduled pod delete: {exc}")
@@ -116,11 +127,23 @@ def delete_pod(pod_id):
try:
mark_workload_pending_deleted(ns_controller)
check_deleted_container(ns_controller)
# Collect network ports for SDN update
ns_networks = collect_and_cleanup_network_ports(str(ns_controller.id))
all_networks.extend(ns_networks)
except Exception as exc:
logger.error(f"Cleanup failed for nscontroller {ns_controller.id} during unscheduled pod delete: {exc}")
db.session.commit()
# Trigger SDN updates to remove network ports
if all_networks:
try:
unique_networks = list(set(all_networks))
send_sdn_updates_for_networks(unique_networks)
logger.info(f"SDN updates sent for {len(unique_networks)} networks after unscheduled pod {pod.id} deletion")
except Exception as e:
logger.error(f"Failed to send SDN updates after pod deletion: {e}")
# Trigger DNS updates to remove deleted pod from DNS
try:
send_dns_updates_for_vdc(pod.vdc_id, logger)
@@ -201,6 +224,8 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
"""
pod_payload = build_pod_payload(pod, use_db_state=True)
all_networks: list[str] = []
# Update containers' desired_state to "deleted" where needed
for container in pod_payload["job_details"]["containers"]:
cid = container["container_id"]
@@ -208,30 +233,7 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
container["desired_state"] = "deleted"
workload = Workload.query.get(cid)
# Release resources when marking container for deletion
from app.models.models import WorkloadResourceUsage, WorkloadHostPooledResource, VolumeWorkloadMapping
# Get resource usage records for this workload
resource_usages = WorkloadResourceUsage.query.filter_by(workload_id=cid).all()
for usage in resource_usages:
# Mark the resource usage record as deleted
usage.soft_delete()
db.session.add(usage)
# Delete volume mappings but keep the volumes
volume_mappings = VolumeWorkloadMapping.query.filter_by(workload_id=cid).all()
for mapping in volume_mappings:
# Log the volume mapping deletion
AuditEntry.log_event(
object=mapping,
action="volume_mapping_deleted",
description=f"Volume mapping deleted for workload {cid} and volume {mapping.volume_id}",
additional_data={"volume_id": mapping.volume_id}
)
# Delete the mapping
db.session.delete(mapping)
cleanup_workload_related_records(cid, logger=logger)
api_token = app.config["CLOUDFLARE_API_TOKEN"]
account_id = app.config["CLOUDFLARE_ACCOUNT_ID"]
@@ -261,6 +263,9 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
pf.soft_delete()
db.session.add(pf)
# Collect network ports for SDN update
container_networks = collect_and_cleanup_network_ports(cid)
all_networks.extend(container_networks)
mark_workload_pending_deleted(workload)
@@ -271,9 +276,21 @@ def process_pod_deletion(pod: ContainerPod, containers_to_delete: list[str], del
ns_controller = Workload.query.filter_by(pod_id=pod.id, workload_type="NSController").first()
if ns_controller:
mark_workload_pending_deleted(ns_controller)
# Collect network ports for SDN update
ns_networks = collect_and_cleanup_network_ports(str(ns_controller.id))
all_networks.extend(ns_networks)
db.session.commit()
# Trigger SDN updates to remove network ports
if all_networks:
try:
unique_networks = list(set(all_networks))
send_sdn_updates_for_networks(unique_networks)
logger.info(f"SDN updates sent for {len(unique_networks)} networks after pod {pod.id} deletion")
except Exception as e:
logger.error(f"Failed to send SDN updates after pod deletion: {e}")
# Trigger DNS updates to remove deleted pod from DNS
try:
send_dns_updates_for_vdc(pod.vdc_id, logger)
@@ -302,25 +319,39 @@ def get_pods():
# pods = ContainerPod.query.all()
active_pods = ContainerPod.query_with_only_active_containers()
response = []
pod_ids = [item['pod'].id for item in active_pods]
nscontrollers = Workload.query.filter(
Workload.pod_id.in_(pod_ids),
Workload.workload_type == "NSController",
Workload.deleted == False
).all()
nscontroller_map = {ns.pod_id: ns for ns in nscontrollers}
ns_ids = [ns.id for ns in nscontrollers]
tunnel_map = {}
dns_map = {}
if ns_ids:
tunnels = CloudflareTunnel.query.filter(
CloudflareTunnel.nscontroller_workload_id.in_(ns_ids),
CloudflareTunnel.deleted == False
).all()
tunnel_map = {t.nscontroller_workload_id: t for t in tunnels}
tunnel_ids = [t.id for t in tunnels]
if tunnel_ids:
dns_records_all = CloudflareDNSRecord.query.filter(
CloudflareDNSRecord.tunnel_id.in_(tunnel_ids),
CloudflareDNSRecord.deleted == False
).all()
dns_map = defaultdict(list)
for record in dns_records_all:
dns_map[record.tunnel_id].append(record)
for item in active_pods:
pod=item['pod']
containers=item['containers']
nscontroller_workload_id=Workload.query.filter_by(workload_type="NSController", pod_id=pod.id).first().id
# Get associated tunnel (if exists)
tunnel = CloudflareTunnel.query.filter_by(
nscontroller_workload_id=nscontroller_workload_id,
deleted=False
).first()
# Get DNS records for this tunnel
dns_records = []
if tunnel:
dns_records = CloudflareDNSRecord.query.filter_by(
tunnel_id=tunnel.id,
deleted=False
).all()
nscontroller = nscontroller_map.get(pod.id)
logger.error(f"nscontroller= {nscontroller}")
tunnel = tunnel_map.get(nscontroller.id) if nscontroller else None
dns_records = dns_map.get(tunnel.id, []) if tunnel else []
pod_data = {
"pod_id": str(pod.id),
"created_at": pod.created_at.isoformat() if pod.created_at else None,
@@ -329,7 +360,10 @@ def get_pods():
"workload_host_id": str(pod.workload_host_id),
"virtual_data_center_id": str(pod.vdc_id),
"region_id": str(pod.vdc.region_id) if pod.vdc else None,
"nscontroller_workload_id": str(nscontroller_workload_id),
# "nscontroller": nscontroller.to_json() if nscontroller else None,
# "nscontroller_volumes": [vwm.volume.to_json() for vwm in nscontroller.volume_mappings] if nscontroller else [],
# "nscontroller_network_ports": [np.to_json() for np in nscontroller.network_ports] if nscontroller else [],
# "nscontroller_certificates": [],
"containers": [
{
"container_id": str(_container.id),
@@ -404,8 +438,25 @@ def get_pod(pod_id):
try:
pod = ContainerPod.query.filter_by(id=pod_uuid).first_or_404()
nscontroller_workload_id=Workload.query.filter_by(workload_type="NSController", pod_id=pod.id).first().id
all_containers_in_pod=Workload.query.filter_by(pod_id=pod.id).all()
nscontroller = Workload.query.filter(
Workload.pod_id == pod.id,
Workload.workload_type == "NSController",
Workload.deleted == False
).first()
all_containers_in_pod = Workload.query.filter_by(pod_id=pod.id).all()
tunnel = None
dns_records = []
if nscontroller:
tunnel = CloudflareTunnel.query.filter_by(
nscontroller_workload_id=nscontroller.id,
deleted=False
).first()
if tunnel:
dns_records = CloudflareDNSRecord.query.filter_by(
tunnel_id=tunnel.id,
deleted=False
).all()
response = {
"pod_id": str(pod.id),
@@ -413,7 +464,10 @@ def get_pod(pod_id):
"workload_host_id": str(pod.workload_host_id),
"vdc_id": str(pod.vdc_id),
"region_id": str(pod.vdc.region_id) if pod.vdc else None,
"nscontroller_workload_id": str(nscontroller_workload_id),
"nscontroller": nscontroller.to_json() if nscontroller else None,
"nscontroller_volumes": [vwm.volume.to_json() for vwm in nscontroller.volume_mappings] if nscontroller else [],
"nscontroller_network_ports": [np.to_json() for np in nscontroller.network_ports] if nscontroller else [],
"nscontroller_certificates": [],
"containers": [
{
"container_id": str(_container.id),
@@ -437,9 +491,24 @@ def get_pod(pod_id):
}
for pf in pod.port_forwardings
if pf.deleted != 1
]
],
"cloudflare_tunnel": {
"tunnel_id": tunnel.tunnel_id if tunnel else None,
"account_id": tunnel.account_id if tunnel else None,
"associated_hostname": tunnel.associated_hostname if tunnel else None,
"dns_records": [
{
"hostname": record.hostname,
"record_type": record.record_type,
"content": record.content,
"proxied": record.proxied,
"ttl": record.ttl
}
for record in dns_records
]
} if tunnel else None
}
return api_response(
success=True,
data=response,
@@ -13,6 +13,10 @@ from app.utils.auth_utils import get_request_user_id
from app.utils.create_workload_container import persist_pod_and_containers
from app.utils.deletion_handler import mark_workload_pending_deleted
from app.utils.sdn_helpers import send_sdn_updates_for_networks
from app.utils.deletion_helpers import (
collect_and_cleanup_network_ports,
log_deletion_audit_event
)
from app.models.network import NetworkPort
def validate_payload(payload):
@@ -720,14 +724,7 @@ def delete_container_workload(workload_id):
mark_workload_pending_deleted(_container)
# Delete associated network ports and send SDN updates
ports = NetworkPort.query.filter_by(workload_id=_container.id).all()
container_networks = []
for port in ports:
container_networks.append(port.network.id)
port.status = "deleted"
port.soft_delete()
db.session.add(port)
container_networks = collect_and_cleanup_network_ports(str(_container.id))
db.session.commit()
if container_networks:
@@ -742,16 +739,7 @@ def delete_container_workload(workload_id):
logger.info(f"Marked container {_container.id} and associated resources as pending-deleted")
# Audit
try:
user_id = get_request_user_id(request)
AuditEntry.log_event(
object=_container,
action="container_delete_requested",
description="Marked container and associated resources as pending-deleted, queued for processing",
user_id=user_id
)
except Exception as exc:
logger.error("Audit logging failed for container delete request %s: %s", workload_id, exc)
log_deletion_audit_event(_container, "delete_container", pod_id=str(_container.pod_id), logger=logger)
return api_response(
success=True,
message="Container marked for deletion and queued for processing",
+76
View File
@@ -0,0 +1,76 @@
"""
Shared utility functions for workload deletion operations.
Consolidates duplicate code between delete_container_workload and delete_pod functions.
"""
from flask import g
from app.models.models import db, WorkloadResourceUsage, VolumeWorkloadMapping, AuditEntry, Workload
from app.models.network import NetworkPort
def collect_and_cleanup_network_ports(workload_id: str) -> list[str]:
"""
Collect network IDs and soft-delete network ports for a workload.
Args:
workload_id: The UUID of the workload
Returns:
List of network IDs that were associated with the workload
"""
ports = NetworkPort.query.filter_by(workload_id=workload_id).all()
network_ids = []
for port in ports:
network_ids.append(port.network.id)
port.status = "deleted"
port.soft_delete()
return network_ids
def cleanup_workload_related_records(workload_id: str, logger=None) -> dict:
"""
Clean up all records related to a workload before deletion.
Args:
workload_id: The UUID of the workload
logger: Optional logger instance
Returns:
Dictionary with cleanup statistics
"""
stats = {"resource_usage_deleted": 0, "volume_mappings_deleted": 0}
# Delete WorkloadResourceUsage
usage = WorkloadResourceUsage.query.filter_by(workload_id=workload_id).all()
for u in usage:
db.session.delete(u)
stats["resource_usage_deleted"] += 1
# Delete VolumeWorkloadMapping
mappings = VolumeWorkloadMapping.query.filter_by(workload_id=workload_id).all()
for m in mappings:
db.session.delete(m)
stats["volume_mappings_deleted"] += 1
return stats
def log_deletion_audit_event(workload: Workload, action: str, pod_id: str = None, logger=None):
"""
Log audit event for workload deletion.
Args:
workload: The Workload instance being deleted
action: The action being performed (e.g., "delete_container", "delete_pod")
pod_id: Optional pod ID if workload is part of a pod
logger: Optional logger instance
"""
actor_id = g.user.id if hasattr(g, 'user') else None
AuditEntry.log_event(
actor_id=actor_id,
action=action,
resource_type="workload",
resource_id=workload.id,
pod_id=pod_id,
details={"workload_type": workload.workload_type}
)