WIP - passing errors from container launch back to api
This commit is contained in:
@@ -169,7 +169,8 @@ def request_container_workload():
|
||||
@api_bp.route('/workloads/containers/status_update/<system_container_id>', methods=['PUT'])
|
||||
def update_container_workload(system_container_id):
|
||||
data = request.json
|
||||
|
||||
logger.debug(data)
|
||||
clok here - need to add error to auti table if there is one
|
||||
# Validate new_status
|
||||
valid_statuses = ["running", "deleted", "stopped","dead", "launch_failed"]
|
||||
new_status=data.get("new_status")
|
||||
@@ -225,40 +226,15 @@ def update_container_workload(system_container_id):
|
||||
try:
|
||||
new_status = data.get('new_status')
|
||||
logger.info(f"Updating Container ID: {system_container_id} to new status {new_status}")
|
||||
|
||||
if new_status == "deleted":
|
||||
# Mark all associated volumes as deleted
|
||||
volume_mappings = VolumeWorkloadMapping.query.filter_by(
|
||||
workload_id=container.id
|
||||
).all()
|
||||
|
||||
for mapping in volume_mappings:
|
||||
mapping.volume.status = "deleted"
|
||||
mapping.volume.soft_delete()
|
||||
db.session.add(mapping.volume)
|
||||
|
||||
# Mark all associated network ports as deleted
|
||||
network_ports = NetworkPort.query.filter_by(workload_id=container.id).all()
|
||||
for port in network_ports:
|
||||
port.status = "deleted"
|
||||
port.soft_delete()
|
||||
db.session.add(port)
|
||||
|
||||
# Also soft delete the volume mappings
|
||||
for mapping in volume_mappings:
|
||||
mapping.soft_delete()
|
||||
db.session.add(mapping)
|
||||
|
||||
|
||||
container.set_status(new_status)
|
||||
if new_status == "deleted":
|
||||
container.soft_delete()
|
||||
|
||||
# If container status is "deleted", check if all containers in the same pod are deleted
|
||||
check_deleted_container(container)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# If container status is "deleted", check if all containers in the same pod are deleted
|
||||
if new_status == "deleted":
|
||||
check_deleted_container(container)
|
||||
|
||||
return jsonify({"success": True, "message": "Status updated successfully."}), 200
|
||||
except Exception as e:
|
||||
error_message = f"Failed to update status for Container ID: {system_container_id}."
|
||||
|
||||
@@ -532,33 +532,6 @@ class Workload(BaseModel):
|
||||
# # Trigger status change handlers
|
||||
self.handle_workload_status_change( old_status, new_status, changed_by)
|
||||
|
||||
def handle_workload_status_change(self, old_status, new_status, changed_by=None):
|
||||
"""
|
||||
Handles all workload status changes and triggers appropriate actions
|
||||
"""
|
||||
if old_status==new_status:
|
||||
return
|
||||
|
||||
logger.debug(f"Handling status change for workload_host of ID [{self.id}] from status [{old_status}] to status [{new_status}]")
|
||||
|
||||
# Log the status change
|
||||
AuditEntry.log_event(
|
||||
object=self,
|
||||
action="status_change",
|
||||
description=f"Status changed from {old_status} to {new_status}",
|
||||
user_id=changed_by,
|
||||
additional_data={"old_status": old_status, "new_status": new_status}
|
||||
)
|
||||
|
||||
# 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)
|
||||
# elif new_status == "deleted":
|
||||
# self.soft_delete()
|
||||
|
||||
def _validate_status_change(self, old_status, new_status):
|
||||
"""Internal validation for status changes"""
|
||||
valid_transitions = {
|
||||
|
||||
@@ -21,6 +21,28 @@ def check_deleted_container(container: Workload) -> None:
|
||||
"""True when workload is hard-deleted or pending deletion."""
|
||||
return wl.deleted or (wl._status == "pending-deleted")
|
||||
|
||||
# Mark all associated volumes as deleted
|
||||
volume_mappings = VolumeWorkloadMapping.query.filter_by(
|
||||
workload_id=container.id
|
||||
).all()
|
||||
|
||||
for mapping in volume_mappings:
|
||||
mapping.volume.status = "deleted"
|
||||
mapping.volume.soft_delete()
|
||||
db.session.add(mapping.volume)
|
||||
|
||||
# Mark all associated network ports as deleted
|
||||
network_ports = NetworkPort.query.filter_by(workload_id=container.id).all()
|
||||
for port in network_ports:
|
||||
port.status = "deleted"
|
||||
port.soft_delete()
|
||||
db.session.add(port)
|
||||
|
||||
# Also soft delete the volume mappings
|
||||
for mapping in volume_mappings:
|
||||
mapping.soft_delete()
|
||||
db.session.add(mapping)
|
||||
|
||||
# ── 1️⃣ Pod-level cleanup ─────────────────────────────────────────
|
||||
mapping: ContainerPodContainer | None = (
|
||||
ContainerPodContainer.query.filter_by(
|
||||
|
||||
@@ -32,6 +32,7 @@ def register_socketio_handlers(socketio):
|
||||
status = details.get("details", {}).get("status")
|
||||
image = details.get("details", {}).get("from")
|
||||
action = details.get("details", {}).get("Action")
|
||||
error = details.get("details", {}).get("error")
|
||||
|
||||
# ----------------- logging boilerplate ----------------------- #
|
||||
logger.info(f"Received Docker event from worker {worker_id}:")
|
||||
@@ -43,6 +44,7 @@ def register_socketio_handlers(socketio):
|
||||
logger.info(f" Status: {status}")
|
||||
logger.info(f" Image: {image}")
|
||||
logger.info(f" Action: {action}")
|
||||
logger.info(f" Error: {error}")
|
||||
logger.debug(f"Full event details: {data}")
|
||||
|
||||
# ------------------- status-update calls --------------------- #
|
||||
@@ -50,20 +52,21 @@ def register_socketio_handlers(socketio):
|
||||
"docker_start",
|
||||
"docker_destroy",
|
||||
"docker_die",
|
||||
"docker_launch_failed", # ← NEW
|
||||
"docker_launch_failed",
|
||||
):
|
||||
action_URL = f"workloads/containers/status_update/{system_container_id}"
|
||||
new_status_map = {
|
||||
"docker_start": "running",
|
||||
"docker_destroy": "deleted",
|
||||
"docker_die": "dead",
|
||||
"docker_launch_failed": "launch_failed", # NEW
|
||||
"docker_launch_failed": "launch_failed",
|
||||
}
|
||||
payload = {
|
||||
"new_status": new_status_map[event_type],
|
||||
"system_container_id": system_container_id,
|
||||
"worker_id": worker_id,
|
||||
"timestamp": timestamp,
|
||||
**({"error": error} if error else {})
|
||||
}
|
||||
headers = {"Content-Type": "application/json"}
|
||||
websocket_server_response = requests.put(
|
||||
|
||||
@@ -61,6 +61,8 @@ class DockerMonitor:
|
||||
continue
|
||||
|
||||
if status in ['start', 'stop', 'die', 'create', 'destroy']:
|
||||
# Launch failures are captured and processed seperatley becuase they are caught when raising a container,
|
||||
# not as a result of an event from docker. This is so we can catch the exception info and include it with the notification
|
||||
# Extract the 'managed_by' tag from the container labels (if it exists)
|
||||
managed_by = None
|
||||
container_attributes = event.get('Actor', {}).get('Attributes', {})
|
||||
|
||||
+14
-9
@@ -151,21 +151,25 @@ class WorkerClient:
|
||||
|
||||
# ---------------- Emit launch-failure events ------------------ #
|
||||
if task_type == "pod-update" and result and result.get("launch_failures"):
|
||||
for failed_id in result["launch_failures"]:
|
||||
await self.send_event({
|
||||
for fail in result["launch_failures"]:
|
||||
failed_id = fail["id"]
|
||||
error_info = fail["error"]
|
||||
|
||||
event_data=({
|
||||
"source": "docker",
|
||||
"event_type": "launch_failed",
|
||||
"container_id": failed_id,
|
||||
"system_container_id": failed_id,
|
||||
"container_name": failed_id,
|
||||
"timestamp": time.time(),
|
||||
"details": {
|
||||
"status": "launch_failed",
|
||||
"from": "",
|
||||
"Action": "launch_failed"
|
||||
}
|
||||
"error": error_info
|
||||
})
|
||||
|
||||
# Put event in the queue for the worker client
|
||||
self.event_queue.put(event_data)
|
||||
logger.info(f"Added Docker event to queue: launch_failed for container {failed_id}.")
|
||||
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error processing task {task_id}: {e}")
|
||||
await self.send_task_result(task_id, {"success": False, "response": str(e)}, task_worker_id)
|
||||
@@ -212,12 +216,13 @@ class WorkerClient:
|
||||
"timestamp": event_data.get("timestamp"),
|
||||
"status": event_data.get("details", {}).get("status"),
|
||||
"image": event_data.get("details", {}).get("from"),
|
||||
"action": event_data.get("details", {}).get("Action")
|
||||
"action": event_data.get("details", {}).get("Action"),
|
||||
**({"error": event_data.get("error")} if event_data.get("error") else {})
|
||||
}
|
||||
|
||||
# Emit the event to the server
|
||||
await self.sio.emit("docker_event", payload)
|
||||
logger.debug(f"Sent {source} docker event: {event_type}\n{payload}")
|
||||
logger.debug(f"Sent {source} event: {event_type}\n{payload}")
|
||||
|
||||
# async def debug_all_events(self, event, data):
|
||||
# """Debug all incoming data."""
|
||||
|
||||
@@ -17,7 +17,7 @@ class ContainerTask:
|
||||
self.docker_client = docker.from_env()
|
||||
self.logger = logger
|
||||
# Track container IDs that fail to launch so the caller can notify the server
|
||||
self.launch_failures: List[str] = []
|
||||
self.launch_failures: list[dict[str, str]] = [] # ← CHANGED
|
||||
|
||||
# --------------------------------------------------------------------- #
|
||||
# Orchestrator-facing entry-point
|
||||
@@ -134,7 +134,9 @@ class ContainerTask:
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to launch container {container_name}: {str(e)}")
|
||||
# record failure so WorkerClient can emit docker_launch_failed
|
||||
self.launch_failures.append(system_container_id)
|
||||
self.launch_failures.append(
|
||||
{"id": system_container_id, "error": str(e)}
|
||||
)
|
||||
return None
|
||||
|
||||
def find_container(self, system_container_id):
|
||||
|
||||
Reference in New Issue
Block a user