Adding a new container to an existing pod doesnt break everything!
This commit is contained in:
@@ -0,0 +1,29 @@
|
||||
FROM python:3.11-slim
|
||||
|
||||
# Set environment variables
|
||||
ENV PYTHONDONTWRITEBYTECODE=1 \
|
||||
PYTHONUNBUFFERED=1 \
|
||||
POETRY_VERSION=1.7.1
|
||||
|
||||
# Set working directory
|
||||
WORKDIR /app
|
||||
|
||||
# Install OS dependencies
|
||||
RUN apt-get update && apt-get install -y --no-install-recommends \
|
||||
build-essential \
|
||||
libmariadb-dev \
|
||||
curl \
|
||||
git \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
# Copy requirements and install
|
||||
COPY requirements.txt .
|
||||
|
||||
RUN pip install --upgrade pip \
|
||||
&& pip install -r requirements.txt
|
||||
|
||||
# Copy project files
|
||||
COPY . .
|
||||
|
||||
# Default entrypoint (can be overridden in docker-compose)
|
||||
CMD ["bash"]
|
||||
@@ -18,6 +18,7 @@ from flask_migrate import Migrate
|
||||
from celery import Celery, Task
|
||||
from logger import logger
|
||||
from app.utils.standard_responses import api_response
|
||||
from flask_cors import CORS
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# 0. Environment setup #
|
||||
@@ -29,6 +30,8 @@ logger.info("Environment variables loaded from .env")
|
||||
# 1. Flask application & database #
|
||||
# --------------------------------------------------------------------------- #
|
||||
app: Flask = Flask(__name__)
|
||||
CORS(app)
|
||||
|
||||
app.config.update(
|
||||
SQLALCHEMY_DATABASE_URI="mysql://root:password@172.17.0.1:3306/theapi",
|
||||
WEBSOCKET_SERVER_URL="http://172.17.0.1:6001/api/create_task",
|
||||
|
||||
+3
-3
@@ -10,13 +10,13 @@ If you want per-project beat schedules, append them to `celery_app.conf`
|
||||
after import.
|
||||
"""
|
||||
|
||||
from logging import getLogger
|
||||
from logger import logger
|
||||
|
||||
from app import celery_app as celery # ← single source-of-truth
|
||||
from app import app as flask_app # access to config & logger
|
||||
|
||||
log = getLogger("xcloudify.celery_factory")
|
||||
log.info("Using existing Celery instance from app.__init__")
|
||||
# log = getLogger("xcloudify.celery_factory")
|
||||
logger.info("Using existing Celery instance from app.__init__")
|
||||
|
||||
# ───────────────────────── Optional: add beat jobs ────────────────────
|
||||
celery.conf.beat_schedule.update(
|
||||
|
||||
@@ -351,6 +351,7 @@ def get_pods():
|
||||
"container_name": pf.workload.name if pf.workload else None,
|
||||
}
|
||||
for pf in pod.port_forwardings
|
||||
if pf.deleted != 1
|
||||
],
|
||||
"cloudflare_tunnel": {
|
||||
"tunnel_id": tunnel.tunnel_id if tunnel else None,
|
||||
@@ -415,146 +416,97 @@ def get_pod(pod_id):
|
||||
@api_bp.route('/workloads/containers/<workload_id>', methods=['DELETE'])
|
||||
def delete_container_workload(workload_id):
|
||||
"""
|
||||
Handle deletion request for a container workload.
|
||||
Sends the deletion request to the websocket server and marks status accordingly.
|
||||
Also soft-deletes any associated PortForwarding records.
|
||||
Mark a single container for deletion via pod-update.
|
||||
"""
|
||||
try:
|
||||
workload_uuid = workload_id # Wrap with uuid.UUID(workload_id) if needed
|
||||
workload_uuid = workload_id
|
||||
except ValueError:
|
||||
return jsonify({"error": "Invalid workload ID"}), 404
|
||||
|
||||
_container = Workload.query.filter(
|
||||
Workload.id == workload_uuid,
|
||||
or_(
|
||||
Workload.workload_type == "Container",
|
||||
Workload.workload_type == "NSController"
|
||||
),
|
||||
or_(Workload.workload_type == "Container", Workload.workload_type == "NSController"),
|
||||
Workload.deleted == False
|
||||
).first_or_404()
|
||||
|
||||
# TODO - Dont send a container deletes - send a pod update with desired state=deleted for the container
|
||||
#Send the request off to the websocket server to have this container deleted
|
||||
payload = {
|
||||
"worker_id": _container.workload_host_id,
|
||||
"task_type": "container-delete",
|
||||
"job_details": {
|
||||
"container": [
|
||||
{
|
||||
"container_id": _container.id,
|
||||
"desired_state": "deleted"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
logger.debug(f"Sending payload to websocket server {payload}")
|
||||
headers = {"Content-Type": "application/json"}
|
||||
websocket_server_response = requests.post(app.config['WEBSOCKET_SERVER_URL'], data=json.dumps(payload), headers=headers)
|
||||
mapping = ContainerPodContainer.query.filter_by(container_workload_id=_container.id).first()
|
||||
if not mapping or not mapping.pod:
|
||||
logger.error(f"Container {_container.id} is not part of a pod.")
|
||||
return jsonify({"error": "Container is not part of a pod"}), 400
|
||||
|
||||
if websocket_server_response.status_code == 201:
|
||||
logger.info(f"Task created successfully for container {_container.id}")
|
||||
_container.set_status("pending-deleted")
|
||||
send_pod_deletion_update(mapping.pod, [str(_container.id)], delete_pod=False)
|
||||
|
||||
# Soft delete any associated port forwards
|
||||
forwards = PortForwarding.query.filter_by(container_workload_id=_container.id, deleted=False).all()
|
||||
if forwards:
|
||||
logger.debug(f"Soft-deleting {len(forwards)} port forward(s) for container {_container.id}")
|
||||
for pf in forwards:
|
||||
pf.soft_delete()
|
||||
db.session.add(pf)
|
||||
|
||||
db.session.add(_container)
|
||||
db.session.commit()
|
||||
else:
|
||||
logger.error(f"Deletion request failed for container {_container.id}")
|
||||
_container.set_status("failed-deleting")
|
||||
db.session.add(_container)
|
||||
db.session.commit()
|
||||
|
||||
logger.info(f"Requested deletion of container {_container.id}")
|
||||
return jsonify({'message': 'Container workload deleted successfully'}), 200
|
||||
logger.info(f"Marked container {_container.id} as pending-deleted")
|
||||
return jsonify({'message': 'Container marked for deletion and pod update sent'}), 200
|
||||
|
||||
@api_bp.route('/workloads/pods/<pod_id>', methods=['DELETE'])
|
||||
def delete_pod(pod_id):
|
||||
"""
|
||||
Delete a pod and all its containers:
|
||||
1. Mark all containers as pending-deleted
|
||||
2. Mark the pod as pending-deleted
|
||||
3. Send tasks via websocket to update each container's desired state to "deleted"
|
||||
4. Delete any associated Cloudflare tunnels if the pod has an NSController
|
||||
Delete an entire pod and all its containers via pod-update.
|
||||
"""
|
||||
try:
|
||||
pod_uuid = (pod_id)
|
||||
pod_uuid = pod_id
|
||||
except ValueError:
|
||||
return jsonify({"error": "Invalid pod ID"}), 404
|
||||
|
||||
pod = ContainerPod.query.filter_by(id=pod_uuid).first_or_404()
|
||||
|
||||
# Cleanup host port mappings
|
||||
# Delete host port mappings
|
||||
for mapping in pod.container_mappings:
|
||||
port_mappings = WorkloadHostPortMapping.query.filter_by(container_workload_id=mapping.container.id).all()
|
||||
for pm in port_mappings:
|
||||
db.session.delete(pm)
|
||||
|
||||
# Set desired state on the pod
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
container_ids = [str(mapping.container.id) for mapping in pod.container_mappings]
|
||||
send_pod_deletion_update(pod, container_ids, delete_pod=True)
|
||||
|
||||
# Mark all containers inside as "pending-deleted" (for tracking)
|
||||
# and collect them for the websocket request
|
||||
container_delete_requests = []
|
||||
for mapping in pod.container_mappings:
|
||||
container = mapping.container
|
||||
if not container.deleted:
|
||||
container.set_status("pending-deleted")
|
||||
db.session.add(container)
|
||||
|
||||
# Add this container to our websocket deletion request
|
||||
container_delete_requests.append({
|
||||
"container_id": str(container.id),
|
||||
"desired_state": "deleted"
|
||||
})
|
||||
|
||||
# Also make sure to include the NSController container for deletion
|
||||
nscontroller = None
|
||||
if pod.nscontroller_workload_id:
|
||||
nscontroller = Workload.query.get(pod.nscontroller_workload_id)
|
||||
if nscontroller and not nscontroller.deleted:
|
||||
nscontroller.set_status("pending-deleted")
|
||||
db.session.add(nscontroller)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Send websocket task to delete all containers
|
||||
if container_delete_requests:
|
||||
payload = {
|
||||
"worker_id": str(pod.workload_host_id),
|
||||
"task_type": "pod-update",
|
||||
"job_details": {
|
||||
"pod_id": pod_id,
|
||||
"nscontroller": {
|
||||
"container_id": str(nscontroller.id) if nscontroller else None,
|
||||
"desired_state": "deleted"
|
||||
},
|
||||
"containers": container_delete_requests
|
||||
}
|
||||
}
|
||||
|
||||
logger.info(f"Sending deletion payload for pod {pod_id} to websocket server: {payload}")
|
||||
headers = {"Content-Type": "application/json"}
|
||||
|
||||
try:
|
||||
websocket_server_response = requests.post(app.config['WEBSOCKET_SERVER_URL'], data=json.dumps(payload), headers=headers)
|
||||
|
||||
if websocket_server_response.status_code == 201:
|
||||
logger.info(f"Deletion tasks created successfully for pod {pod.id} with {len(container_delete_requests)} containers")
|
||||
else:
|
||||
logger.error(f"Deletion request failed for pod {pod.id}. Response: {websocket_server_response.status_code} - {websocket_server_response.text}")
|
||||
# Continue with the deletion process anyway, as the status updates may come through other channels
|
||||
except Exception as e:
|
||||
logger.error(f"Exception when sending deletion request for pod {pod.id}: {str(e)}")
|
||||
# Continue with the process despite the error
|
||||
|
||||
logger.info(f"Marked pod {pod.id} and its containers as pending-deleted")
|
||||
return jsonify({'message': 'Pod and containers set to pending-deleted'}), 200
|
||||
|
||||
|
||||
from app.utils.create_workload_container import build_pod_payload
|
||||
|
||||
def send_pod_deletion_update(pod: ContainerPod, containers_to_delete: list[str], delete_pod: bool = False):
|
||||
"""
|
||||
Use the canonical build_pod_payload function to construct the full pod intent,
|
||||
and override desired_state fields as needed for deletions.
|
||||
"""
|
||||
pod_payload = build_pod_payload(pod)
|
||||
|
||||
# Update containers' desired_state to "deleted" where needed
|
||||
for container in pod_payload["job_details"]["containers"]:
|
||||
cid = container["container_id"]
|
||||
if cid in containers_to_delete:
|
||||
container["desired_state"] = "deleted"
|
||||
workload = Workload.query.get(cid)
|
||||
workload.set_status("pending-deleted")
|
||||
db.session.add(workload)
|
||||
|
||||
# Optionally mark NSController and pod for deletion
|
||||
if delete_pod:
|
||||
pod_payload["job_details"]["nscontroller"]["desired_state"] = "deleted"
|
||||
pod.status = "pending-deleted"
|
||||
db.session.add(pod)
|
||||
|
||||
ns_id = pod_payload["job_details"]["nscontroller"]["container_id"]
|
||||
ns = Workload.query.get(ns_id)
|
||||
if ns:
|
||||
ns.set_status("pending-deleted")
|
||||
db.session.add(ns)
|
||||
|
||||
db.session.commit()
|
||||
|
||||
# Send the intent to the worker via WebSocket
|
||||
logger.info(f"Sending pod-update to WebSocket server: {pod_payload}")
|
||||
try:
|
||||
response = requests.post(
|
||||
app.config["WEBSOCKET_SERVER_URL"],
|
||||
data=json.dumps(pod_payload),
|
||||
headers={"Content-Type": "application/json"}
|
||||
)
|
||||
if response.status_code == 201:
|
||||
logger.info(f"WebSocket task accepted for pod {pod.id}")
|
||||
else:
|
||||
logger.error(f"WebSocket task rejected: {response.status_code} - {response.text}")
|
||||
except Exception as e:
|
||||
logger.error(f"Exception sending pod-update for pod {pod.id}: {str(e)}")
|
||||
|
||||
@@ -39,8 +39,6 @@ def edit_workload_host(workload_host_id):
|
||||
@api_bp.route('/workload_hosts/<workload_host_id>', methods=['GET'])
|
||||
def get_workload_host(workload_host_id):
|
||||
workload_host = WorkloadHost.query.filter_by(id=workload_host_id,deleted=0).first_or_404()
|
||||
|
||||
logger.debug(workload_host.to_json())
|
||||
return jsonify(workload_host.to_json())
|
||||
|
||||
@api_bp.route('/workload_hosts/<workload_host_id>', methods=['DELETE'])
|
||||
|
||||
+9
-10
@@ -689,10 +689,10 @@ class ContainerPod(BaseModel):
|
||||
)
|
||||
)
|
||||
containers = q.all()
|
||||
logger.debug(
|
||||
"Pod %s → %d containers (include_deleted=%s)",
|
||||
self.id, len(containers), include_deleted,
|
||||
)
|
||||
# logger.debug(
|
||||
# "Pod %s → %d containers (include_deleted=%s)",
|
||||
# self.id, len(containers), include_deleted,
|
||||
# )
|
||||
return containers
|
||||
|
||||
@classmethod
|
||||
@@ -706,20 +706,19 @@ class ContainerPod(BaseModel):
|
||||
pods = db.session.query(cls).filter_by(deleted=False).all()
|
||||
results = []
|
||||
|
||||
logger.debug(
|
||||
"Scanning %d pods for active containers (include_deleted=%s)",
|
||||
len(pods), include_deleted_containers,
|
||||
)
|
||||
# logger.debug(
|
||||
# "Scanning %d pods for active containers (include_deleted=%s)",
|
||||
# len(pods), include_deleted_containers,
|
||||
# )
|
||||
|
||||
for pod in pods:
|
||||
containers = pod.active_containers(include_deleted_containers)
|
||||
if containers or include_deleted_containers:
|
||||
results.append({"pod": pod, "containers": containers})
|
||||
|
||||
logger.debug("Returning %d pods after filtering", len(results))
|
||||
# logger.debug("Returning %d pods after filtering", len(results))
|
||||
return results
|
||||
|
||||
|
||||
class Image(BaseModel):
|
||||
__tablename__ = "images"
|
||||
location = Column(String(255), nullable=False)
|
||||
|
||||
@@ -11,4 +11,7 @@ flask_migrate
|
||||
websocket-client
|
||||
streamlit
|
||||
aiohttp
|
||||
celery
|
||||
celery
|
||||
flask_cors
|
||||
python-dotenv
|
||||
PyJWT
|
||||
@@ -16,76 +16,92 @@ from app import celery_app as celery, db, logger, app
|
||||
from app.models.models import WorkloadHost
|
||||
|
||||
redisclient = redis.from_url(app.config["REDIS_URL"])
|
||||
enhanced_debug=False
|
||||
|
||||
@celery.task(name="tasks.reconcile_online_workers", bind=True)
|
||||
def reconcile_online_workers(self) -> None:
|
||||
"""
|
||||
Reconcile online/offline status of WorkloadHost entries based on Redis liveness heartbeats.
|
||||
"""
|
||||
|
||||
logger.info("Starting worker heartbeat reconciliation task")
|
||||
if enhanced_debug: logger.info("=== Starting reconcile_online_workers task ===")
|
||||
now = int(time.time())
|
||||
cutoff = datetime.utcnow() - timedelta(seconds=10)
|
||||
|
||||
timeout = app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"]
|
||||
logger.debug(f"Using PING_HEARTBEAT_TIMEOUT_SECONDS = {timeout}")
|
||||
logger.debug(f"Cutoff for online host updated_at = {cutoff.isoformat()}")
|
||||
|
||||
# -----------------------------------------------------------------------
|
||||
# 1. Check all workers currently marked as ONLINE
|
||||
# 1. Evaluate workers marked as ONLINE
|
||||
# -----------------------------------------------------------------------
|
||||
online_workers = WorkloadHost.query.filter(
|
||||
WorkloadHost._status == "online",
|
||||
WorkloadHost.deleted == 0,
|
||||
WorkloadHost.updated_at < cutoff
|
||||
).all()
|
||||
|
||||
logger.info(f"Found {len(online_workers)} workers marked as online")
|
||||
if enhanced_debug: logger.info(f"[PHASE 1] Checking {len(online_workers)} workers marked as ONLINE")
|
||||
|
||||
for worker in online_workers:
|
||||
redis_key = f"ws_liveness:{worker.id}"
|
||||
last_seen = redisclient.get(redis_key)
|
||||
logger.debug(f"[ONLINE → OFFLINE?] Checking worker {worker.id}")
|
||||
logger.debug(f" Redis key: {redis_key}")
|
||||
|
||||
is_stale = False
|
||||
if not last_seen:
|
||||
logger.warning(f"[OFFLINE CHECK] No heartbeat for WorkloadHost {worker.id}")
|
||||
last_seen = redisclient.get(redis_key)
|
||||
if last_seen is None:
|
||||
logger.warning(f" ❌ No Redis key found — assuming stale")
|
||||
is_stale = True
|
||||
else:
|
||||
try:
|
||||
last_seen = int(last_seen)
|
||||
if (now - last_seen) > app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"]:
|
||||
is_stale = True
|
||||
except Exception:
|
||||
logger.exception(f"[OFFLINE CHECK] Invalid heartbeat data for WorkloadHost {worker.id}")
|
||||
delta = now - last_seen
|
||||
logger.debug(f" ✅ Last heartbeat: {last_seen} (delta={delta}s)")
|
||||
|
||||
is_stale = delta > timeout
|
||||
if is_stale:
|
||||
logger.warning(f" ⚠️ Stale heartbeat (>{timeout}s) — marking offline")
|
||||
else:
|
||||
if enhanced_debug: logger.debug(f" 🟢 Recent heartbeat — staying online")
|
||||
except Exception as e:
|
||||
logger.exception(f" ❌ Exception parsing Redis heartbeat value: {e}")
|
||||
is_stale = True
|
||||
|
||||
if is_stale:
|
||||
logger.warning(f"WorkloadHost {worker.id} is stale — marking offline")
|
||||
worker.set_status("offline")
|
||||
db.session.commit()
|
||||
logger.info(f" ✅ Worker {worker.id} marked OFFLINE")
|
||||
|
||||
# -----------------------------------------------------------------------
|
||||
# 2. Check all workers currently marked as OFFLINE
|
||||
# 2. Evaluate workers marked as OFFLINE
|
||||
# -----------------------------------------------------------------------
|
||||
offline_workers = WorkloadHost.query.filter(
|
||||
WorkloadHost._status == "offline",
|
||||
WorkloadHost.deleted == 0
|
||||
).all()
|
||||
|
||||
logger.info(f"Found {len(offline_workers)} workers marked as offline")
|
||||
if enhanced_debug: logger.info(f"[PHASE 2] Checking {len(offline_workers)} workers marked as OFFLINE")
|
||||
|
||||
for worker in offline_workers:
|
||||
redis_key = f"ws_liveness:{worker.id}"
|
||||
logger.debug(f"[OFFLINE → ONLINE?] Checking worker {worker.id}")
|
||||
logger.debug(f" Redis key: {redis_key}")
|
||||
|
||||
last_seen = redisclient.get(redis_key)
|
||||
if last_seen is None:
|
||||
logger.debug(f" ❌ No Redis key found — remaining offline")
|
||||
continue
|
||||
|
||||
is_active = False
|
||||
if last_seen:
|
||||
try:
|
||||
last_seen = int(last_seen)
|
||||
if (now - last_seen) <= app.config["PING_HEARTBEAT_TIMEOUT_SECONDS"]:
|
||||
is_active = True
|
||||
except Exception:
|
||||
logger.exception(f"[ONLINE CHECK] Invalid heartbeat data for WorkloadHost {worker.id}")
|
||||
try:
|
||||
last_seen = int(last_seen)
|
||||
delta = now - last_seen
|
||||
logger.debug(f" ✅ Last heartbeat: {last_seen} (delta={delta}s)")
|
||||
|
||||
if is_active:
|
||||
logger.info(f"WorkloadHost {worker.id} is alive again — marking online")
|
||||
worker.set_status("online")
|
||||
db.session.commit()
|
||||
if delta <= timeout:
|
||||
if enhanced_debug: logger.info(f" 🔁 Worker is back alive — marking ONLINE")
|
||||
worker.set_status("online")
|
||||
db.session.commit()
|
||||
logger.info(f" ✅ Worker {worker.id} marked ONLINE")
|
||||
else:
|
||||
logger.debug(f" ⚠️ Heartbeat exists but is stale — staying offline")
|
||||
except Exception as e:
|
||||
logger.exception(f" ❌ Exception parsing Redis heartbeat value: {e}")
|
||||
|
||||
logger.info("Finished heartbeat reconciliation task.")
|
||||
if enhanced_debug: logger.info("=== Finished reconcile_online_workers task ===")
|
||||
|
||||
@@ -54,7 +54,7 @@
|
||||
|
||||
|
||||
<script>
|
||||
const socket = io("http://localhost:6001");
|
||||
const socket = io("http://192.168.64.2:6001");
|
||||
const userId = "user-xyz"; // Ideally from session/auth
|
||||
let activeRequestId = null;
|
||||
|
||||
|
||||
@@ -21,6 +21,12 @@ def check_deleted_container(container: Workload) -> None:
|
||||
"""True when workload is hard-deleted or pending deletion."""
|
||||
return wl.deleted or (wl._status == "pending-deleted")
|
||||
|
||||
api_token = app.config["CLOUDFLARE_API_TOKEN"]
|
||||
account_id = app.config["CLOUDFLARE_ACCOUNT_ID"]
|
||||
zone_id = app.config["CLOUDFLARE_ZONE_ID"]
|
||||
|
||||
cf_manager = CloudflareTunnelManager(api_token, account_id, zone_id, logger)
|
||||
|
||||
# Mark all associated volumes as deleted
|
||||
volume_mappings = VolumeWorkloadMapping.query.filter_by(
|
||||
workload_id=container.id
|
||||
@@ -31,6 +37,21 @@ def check_deleted_container(container: Workload) -> None:
|
||||
mapping.volume.soft_delete()
|
||||
db.session.add(mapping.volume)
|
||||
|
||||
forwards = PortForwarding.query.filter_by(workload_id=container.id, deleted=False).all()
|
||||
for pf in forwards:
|
||||
if pf.dns_record:
|
||||
dns = pf.dns_record
|
||||
cf_manager.delete_dns_record(pf.dns_record.dns_record_id)
|
||||
logger.info(f"Soft-deleting DNS record {dns.hostname} for container {container.id}")
|
||||
dns.soft_delete()
|
||||
db.session.add(dns)
|
||||
|
||||
|
||||
logger.info(f"Soft-deleting PortForwarding {pf.id} for container {container.id}")
|
||||
pf.soft_delete()
|
||||
db.session.add(pf)
|
||||
|
||||
|
||||
# Mark all associated network ports as deleted
|
||||
network_ports = NetworkPort.query.filter_by(workload_id=container.id).all()
|
||||
for port in network_ports:
|
||||
@@ -49,7 +70,7 @@ def check_deleted_container(container: Workload) -> None:
|
||||
container_workload_id=container.id
|
||||
).first()
|
||||
)
|
||||
|
||||
db.session.commit()
|
||||
if mapping:
|
||||
pod: ContainerPod | None = ContainerPod.query.get(mapping.pod_id)
|
||||
if pod and not pod.deleted:
|
||||
@@ -91,11 +112,6 @@ def check_deleted_container(container: Workload) -> None:
|
||||
container.id,
|
||||
)
|
||||
|
||||
api_token = "ri6lIjM-aRJBY_xZ82w0Haew93U6YgZYHi5jby1-" # TODO move to config
|
||||
account_id = "5095a74b62fee53cc5d997c67443bac5"
|
||||
zone_id = "e2cafdd8929869d5db885d8824f514b5"
|
||||
cf_manager = CloudflareTunnelManager(api_token, account_id, zone_id, logger)
|
||||
|
||||
for tunnel in tunnels:
|
||||
logger.debug("Deleting Cloudflare tunnel %s (%s)", tunnel.id, tunnel.name)
|
||||
cf_manager.cleanup_tunnel(tunnel.name, delete_tunnel=True)
|
||||
|
||||
@@ -395,14 +395,15 @@ def _build_pod_objects_and_enqueue(
|
||||
data=json.dumps(pod_payload),
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
db.session.commit()
|
||||
|
||||
if resp.status_code == 201:
|
||||
logger.info("Pod-update queued for pod %s", pod.id)
|
||||
_set_container_statuses(client_response, "allocated")
|
||||
_set_container_statuses(client_response, "allocated")# TODO - When adding a new container to an exisgin pod this errors becuase the existing containers are already running
|
||||
nscontroller.set_status("allocated")
|
||||
else:
|
||||
logger.error("Pod-update failed for pod %s · %s", pod.id, resp.text)
|
||||
_set_container_statuses(client_response, "failed-allocation")
|
||||
_set_container_statuses(client_response, "failed-allocation")
|
||||
nscontroller.set_status("failed-allocation")
|
||||
|
||||
db.session.add(nscontroller)
|
||||
@@ -420,11 +421,20 @@ def build_pod_payload(pod: ContainerPod) -> Dict:
|
||||
"""
|
||||
ns = Workload.query.get(pod.nscontroller_workload_id)
|
||||
ns_lp = json.loads(ns.launch_params)
|
||||
|
||||
logger.debug(f"Building podpayload for pod {pod.id}")
|
||||
containers = []
|
||||
for m in ContainerPodContainer.query.filter_by(pod_id=pod.id).all():
|
||||
w = Workload.query.get(m.container_workload_id)
|
||||
logger.debug(f"Processing container {w.id}")
|
||||
if w is None or w.deleted == 1:
|
||||
logger.debug(f"{w.id} is deleted, skipping")
|
||||
continue # Skip deleted or missing workloads
|
||||
else:
|
||||
logger.debug(f"including container {w.id}")
|
||||
launch_params = json.loads(w.launch_params)
|
||||
if ":" not in launch_params["docker_image"]:
|
||||
launch_params["docker_image"] += ":latest"
|
||||
|
||||
cont = {
|
||||
"container_id": str(w.id),
|
||||
"docker_image": launch_params["docker_image"],
|
||||
|
||||
+5
-6
@@ -39,12 +39,13 @@ services:
|
||||
container_name: novnc
|
||||
ports:
|
||||
- "6081:6080"
|
||||
- "8080:8080"
|
||||
- "8085:8080"
|
||||
|
||||
|
||||
api-server:
|
||||
build:
|
||||
context: .
|
||||
context: ./app
|
||||
image: xcloudify-flask-base
|
||||
container_name: api-server
|
||||
volumes:
|
||||
- .:/app
|
||||
@@ -121,8 +122,7 @@ services:
|
||||
SOCKET_SERVER_URL: http://172.17.0.1:6001
|
||||
|
||||
celery-worker:
|
||||
build:
|
||||
context: .
|
||||
image: xcloudify-flask-base
|
||||
container_name: celery-worker
|
||||
command: celery -A app.celery_app.celery worker --loglevel=info
|
||||
volumes:
|
||||
@@ -133,8 +133,7 @@ services:
|
||||
- mariadb
|
||||
|
||||
celery-beat:
|
||||
build:
|
||||
context: .
|
||||
image: xcloudify-flask-base
|
||||
container_name: celery-beat
|
||||
command: celery -A app.celery_app.celery beat --loglevel=info
|
||||
volumes:
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
PODID=7e0604d0-c228-4a8c-8bb7-27eaa8c96c1e
|
||||
PODID=$(curl -s http://192.168.64.2:5000/api/workloads/pods | jq -r '.[0].pod_id')
|
||||
|
||||
curl -X POST http://192.168.64.2:5000/api/workloads/containers \
|
||||
-H "Content-Type: application/json" \
|
||||
@@ -8,15 +8,15 @@ curl -X POST http://192.168.64.2:5000/api/workloads/containers \
|
||||
"containers": [
|
||||
{
|
||||
"docker_image": "joke_container",
|
||||
"container_name": "amazing-joke-server99",
|
||||
"container_name": "containernumber5",
|
||||
"ports": [
|
||||
{
|
||||
"internal": 9092,
|
||||
"internal": 9095,
|
||||
"external": 8766,
|
||||
"use_dns": false
|
||||
}
|
||||
],
|
||||
"env": { "PORT": "9092" }
|
||||
"env": { "PORT": "9095" }
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -26,3 +26,38 @@ curl -X POST http://192.168.64.2:5000/api/workloads/containers -H "Content-Typ
|
||||
}
|
||||
]
|
||||
}'
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
curl -X POST http://192.168.64.2:5000/api/workloads/containers -H "Content-Type: application/json" -d '{
|
||||
"virtual_data_center": "b1d09477-742b-485e-87c3-5dad36dd4f9e",
|
||||
"containers": [
|
||||
{
|
||||
"docker_image": "joke_container",
|
||||
"container_name": "very-important-webserver",
|
||||
"ports": [
|
||||
{
|
||||
"internal": 9090,
|
||||
"external": 9090,
|
||||
"use_dns": true
|
||||
}
|
||||
],
|
||||
"env": { "PORT": "9090" }
|
||||
},
|
||||
{
|
||||
"docker_image": "whoami81",
|
||||
"container_name": "whoami-server",
|
||||
"ports": [
|
||||
{
|
||||
"internal": 81,
|
||||
"external": 8081,
|
||||
"use_dns": true
|
||||
}
|
||||
]
|
||||
}
|
||||
]
|
||||
}'
|
||||
|
||||
@@ -173,8 +173,8 @@ def register_socketio_handlers(socketio):
|
||||
|
||||
context = json.loads(context_data)
|
||||
user_sid = context.get("user_sid")
|
||||
with base.connected_sids_lock:
|
||||
if user_sid in base.connected_sids:
|
||||
with connected_sids_lock:
|
||||
if user_sid in connected_sids:
|
||||
socketio.emit("user_log_stream_update", {
|
||||
"logs": logs,
|
||||
"request_id": request_id
|
||||
|
||||
@@ -8,7 +8,7 @@ from flask import request
|
||||
|
||||
from websocket_server.config import get_redis_client
|
||||
from websocket_server.events import base
|
||||
from websocket_server.shared_state import connected_workers, connected_sids_lock, connected_sids
|
||||
from websocket_server.shared_state import connected_workers
|
||||
|
||||
logger = logging.getLogger("websocket_server")
|
||||
socketio = base.socketio
|
||||
|
||||
@@ -6,7 +6,7 @@ import requests
|
||||
from flask import request
|
||||
|
||||
from websocket_server.events import base
|
||||
from websocket_server.shared_state import connected_workers, connected_sids_lock, connected_sids
|
||||
from websocket_server.shared_state import connected_workers
|
||||
|
||||
|
||||
logger = logging.getLogger("websocket_server")
|
||||
|
||||
@@ -36,7 +36,6 @@ def notify_worker_online(worker_id):
|
||||
|
||||
try:
|
||||
response = requests.put(f"{api_server_url}/workload_hosts/{worker_id}", json=payload, headers=headers)
|
||||
logger.debug(response.text)
|
||||
logger.info(f"[{worker_id}] API server acknowledged online state")
|
||||
except Exception as e:
|
||||
logger.error(f"[{worker_id}] Failed to notify API server of online status: {e}")
|
||||
|
||||
@@ -186,14 +186,18 @@ class ContainerTask:
|
||||
network_settings = container.attrs.get('NetworkSettings', {})
|
||||
|
||||
# ─────────────────── 1️⃣ IMAGE ───────────────────
|
||||
# Normalize tags to include ':latest' if omitted
|
||||
expected_image = container_spec.get("docker_image")
|
||||
running_image = container.image.tags[0] if container.image.tags else None
|
||||
if expected_image and ":" not in expected_image:
|
||||
expected_image += ":latest"
|
||||
|
||||
running_image = container.image.tags[0] if container.image.tags else None
|
||||
|
||||
if expected_image != running_image:
|
||||
self.logger.warning(
|
||||
"Image mismatch: expected %s, got %s", expected_image, running_image
|
||||
)
|
||||
return False
|
||||
|
||||
# ─────────────────── 2️⃣ NETWORKS (NSController only) ───────────────────
|
||||
if container_spec.get("networks") and container_spec.get("nscontroller_container_id") is None:
|
||||
existing_networks = set(network_settings.get('Networks', {}).keys())
|
||||
|
||||
Reference in New Issue
Block a user