From b79b11881fa78b2f4a2e87d9bd2b464dd3c0c091 Mon Sep 17 00:00:00 2001 From: Cory Hawkvelt Date: Thu, 11 Sep 2025 01:28:56 +0930 Subject: [PATCH] Allowed for reconciliation to reset status back to running --- websocket_server/events/docker.py | 2 + worker/workerClient.py | 75 +++++++++++++++++++++++++++++++ 2 files changed, 77 insertions(+) diff --git a/websocket_server/events/docker.py b/websocket_server/events/docker.py index 88a43cc..2d476d8 100644 --- a/websocket_server/events/docker.py +++ b/websocket_server/events/docker.py @@ -49,6 +49,7 @@ def register_socketio_handlers(socketio): # ------------------- status-update calls --------------------- # if event_type in ( + "docker_running", "docker_start", "docker_destroy", "docker_die", @@ -73,6 +74,7 @@ def register_socketio_handlers(socketio): # Handle standard status events new_status_map = { "docker_start": "running", + "docker_running": "running", "docker_destroy": "deleted", "docker_die": "dead", "docker_launch_failed": "launch_failed", diff --git a/worker/workerClient.py b/worker/workerClient.py index 55c2f14..878918d 100644 --- a/worker/workerClient.py +++ b/worker/workerClient.py @@ -177,9 +177,84 @@ class WorkerClient: success = await self.reconcile_containers() if success: logger.info("Container reconciliation completed successfully on worker join") + # Report current container statuses to update the server + self.report_container_statuses() else: logger.error("Container reconciliation failed on worker join") + def report_container_statuses(self): + + """ + Report the current status of all managed containers to the server. + This should be called after reconciliation to ensure the server has + accurate information about container states. + """ + try: + import docker + # Get Docker client + docker_client = docker.from_env() + + # Get all containers with the managed_by label + all_containers = docker_client.containers.list(all=True, filters={ + "label": "managed_by=worker_agent" + }) + + logger.info(f"Reporting status for {len(all_containers)} managed containers") + + # Report each container's status + for container in all_containers: + # Get container labels + labels = container.labels + system_container_id = labels.get("system_container_id") + container_name = labels.get("name", container.name) + + if not system_container_id: + logger.warning(f"Managed container {container.id} missing system_container_id label") + continue + + # Get container image + image_tags = container.image.tags + image_name = image_tags[0] if image_tags else "unknown" + + # Prepare the event data using the same format as DockerMonitor + # Note: The send_event method expects status in event_data["details"]["status"] + event_data = { + 'source': 'docker', + 'event_type': container.status, # Use actual container status as event type + 'container_id': container.id, + 'system_container_id': system_container_id, + 'container_name': container_name, + 'timestamp': time.time(), + 'details': { + 'status': container.status, # 'running', 'exited', etc. + 'id': container.id, + 'from': image_name, + 'Type': 'container', + 'Action': container.status, # Use container status as action + 'Actor': { + 'ID': container.id, + 'Attributes': { + 'image': image_name, + 'managed_by': 'worker_agent', + 'name': container_name, + 'system_container_id': system_container_id + } + }, + 'scope': 'local', + 'time': int(time.time()), + 'timeNano': int(time.time() * 1000000000) + }, + 'managed_by': 'worker_agent' + } + + # Put event in the queue for the worker client + if self.event_queue: + self.event_queue.put(event_data) + logger.debug(f"Added container status event to queue: {container.status} for container {container_name}") + + except Exception as e: + logger.error(f"Error reporting container statuses: {str(e)}") + async def on_join_reject(self, data): logger.error("Join rejected.") self.joined_server = False