Allowed for reconciliation to reset status back to running
This commit is contained in:
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user