container log streaming!!!
This commit is contained in:
@@ -37,33 +37,65 @@
|
||||
<div class="mb-4">
|
||||
<button class="btn btn-primary" onclick="viewLogs()">View Logs</button>
|
||||
</div>
|
||||
|
||||
<div class="mb-3">
|
||||
<button class="btn btn-success" onclick="startLogStream()">Start Streaming</button>
|
||||
<button class="btn btn-danger" onclick="stopLogStream()">Stop Streaming</button>
|
||||
</div>
|
||||
|
||||
<h5>Logs:</h5>
|
||||
<div id="logBox" readonly></div>
|
||||
|
||||
<script>
|
||||
const socket = io("http://localhost:6001");
|
||||
const userId = "user-xyz"; // Set your user's ID from session/auth
|
||||
const requestId = () => 'req-' + Math.random().toString(36).substring(2, 10);
|
||||
|
||||
const userId = "user-xyz"; // Ideally from session/auth
|
||||
let activeRequestId = null;
|
||||
|
||||
function requestId() {
|
||||
return 'req-' + Math.random().toString(36).substring(2, 10);
|
||||
}
|
||||
|
||||
function viewLogs() {
|
||||
const containerId = document.getElementById("containerSelect").value;
|
||||
const rid = requestId();
|
||||
|
||||
document.getElementById("logBox").textContent = "Requesting logs...";
|
||||
|
||||
activeRequestId = requestId();
|
||||
document.getElementById("logBox").textContent = "Fetching logs...";
|
||||
|
||||
socket.emit("user_request_container_logs", {
|
||||
container_id: containerId,
|
||||
user_id: userId,
|
||||
lines: 100,
|
||||
request_id: rid
|
||||
request_id: activeRequestId
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
function startLogStream() {
|
||||
const containerId = document.getElementById("containerSelect").value;
|
||||
activeRequestId = requestId();
|
||||
|
||||
socket.emit("user_start_log_stream", {
|
||||
container_id: containerId,
|
||||
user_id: userId,
|
||||
request_id: activeRequestId
|
||||
});
|
||||
}
|
||||
|
||||
function stopLogStream() {
|
||||
if (activeRequestId) {
|
||||
socket.emit("user_stop_log_stream", {
|
||||
request_id: activeRequestId
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
socket.on("user_log_response", function (data) {
|
||||
document.getElementById("logBox").textContent = data.logs || "(no logs received)";
|
||||
});
|
||||
|
||||
socket.on("user_log_stream_update", function (data) {
|
||||
const logBox = document.getElementById("logBox");
|
||||
logBox.textContent += data.logs + "\n";
|
||||
logBox.scrollTop = logBox.scrollHeight;
|
||||
});
|
||||
</script>
|
||||
|
||||
|
||||
</body>
|
||||
</html>
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
FROM python:3.11-slim
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
COPY log_emitter.py .
|
||||
|
||||
CMD ["python", "log_emitter.py"]
|
||||
@@ -0,0 +1,20 @@
|
||||
import time
|
||||
import random
|
||||
import datetime
|
||||
|
||||
messages = [
|
||||
"Starting task execution...",
|
||||
"Fetching data from API...",
|
||||
"Data parsing complete.",
|
||||
"Writing to database...",
|
||||
"Task finished successfully.",
|
||||
"Heartbeat: system operational.",
|
||||
"Warning: high memory usage.",
|
||||
"Error: failed to connect to service X.",
|
||||
"Reconnecting...",
|
||||
]
|
||||
|
||||
while True:
|
||||
msg = random.choice(messages)
|
||||
print(f"{datetime.datetime.utcnow().isoformat()} | {msg}", flush=True)
|
||||
time.sleep(random.uniform(0.3, 1.5))
|
||||
@@ -34,7 +34,7 @@ def render_list():
|
||||
container_name = st.text_input("Container Name", value="my-container")
|
||||
|
||||
# Suggestions for common images
|
||||
docker_image_options = ["nginx", "traefik/whoami:latest", "whoami81"]
|
||||
docker_image_options = ["nginx", "traefik/whoami:latest", "whoami81","noisy_container"]
|
||||
docker_image = st.selectbox(
|
||||
"Docker Image",
|
||||
options=docker_image_options + ["Custom..."],
|
||||
|
||||
@@ -629,6 +629,117 @@ def handle_user_log_request(data):
|
||||
"request_id": request_id
|
||||
}, to=request.sid)
|
||||
|
||||
@socketio.on("user_start_log_stream")
|
||||
def handle_user_stream_start(data):
|
||||
container_id = data.get("container_id")
|
||||
user_id = data.get("user_id")
|
||||
request_id = data.get("request_id")
|
||||
|
||||
logger.info(f"[{request_id}] Received user_start_log_stream for container {container_id} by user {user_id} (sid={request.sid})")
|
||||
|
||||
try:
|
||||
# Lookup container
|
||||
logger.debug(f"[{request_id}] Requesting container info from API for container_id={container_id}")
|
||||
container_response = requests.get(f"{api_server_url}/workloads/containers/{container_id}")
|
||||
container_response.raise_for_status()
|
||||
container_info = container_response.json()
|
||||
logger.debug(f"[{request_id}] Container info received: {container_info}")
|
||||
|
||||
# Access control (currently disabled)
|
||||
# if container_info["user_id"] != user_id:
|
||||
# logger.warning(f"[{request_id}] Access denied for user {user_id} on container {container_id}")
|
||||
# emit("user_log_response", {
|
||||
# "success": False,
|
||||
# "logs": f"Access denied",
|
||||
# "request_id": request_id
|
||||
# }, to=request.sid)
|
||||
# return
|
||||
|
||||
worker_id = container_info["workload_host_id"]
|
||||
if worker_id not in connected_workers:
|
||||
logger.warning(f"[{request_id}] Worker {worker_id} not connected for container {container_id}")
|
||||
emit("user_log_response", {
|
||||
"success": False,
|
||||
"logs": f"Worker not connected",
|
||||
"request_id": request_id
|
||||
}, to=request.sid)
|
||||
return
|
||||
|
||||
# Store request context in Redis
|
||||
context_payload = {
|
||||
"user_sid": request.sid,
|
||||
"container_id": container_id,
|
||||
"user_id": user_id,
|
||||
"worker_id": worker_id
|
||||
}
|
||||
redis_client = get_redis_client()
|
||||
redis_client.setex(f"log_request_context:{request_id}", 60 * 10, json.dumps(context_payload))
|
||||
logger.info(f"[{request_id}] Stored log stream context in Redis: {context_payload}")
|
||||
|
||||
# Emit to worker
|
||||
emit_payload = {
|
||||
"container_name": container_info["id"],
|
||||
"request_id": request_id,
|
||||
"worker_id": worker_id
|
||||
}
|
||||
logger.info(f"[{request_id}] Emitting start_container_log_stream to worker {worker_id}")
|
||||
socketio.emit("start_container_log_stream", emit_payload, to=connected_workers[worker_id])
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"[{request_id}] Error starting log stream: {e}")
|
||||
|
||||
|
||||
@socketio.on("user_stop_log_stream")
|
||||
def handle_user_stream_stop(data):
|
||||
request_id = data.get("request_id")
|
||||
logger.info(f"[{request_id}] Received user_stop_log_stream")
|
||||
|
||||
try:
|
||||
redis_client = get_redis_client()
|
||||
context_data = redis_client.get(f"log_request_context:{request_id}")
|
||||
if context_data:
|
||||
context = json.loads(context_data)
|
||||
worker_id = context.get("worker_id")
|
||||
if worker_id in connected_workers:
|
||||
logger.info(f"[{request_id}] Forwarding stop_container_log_stream to worker {worker_id}")
|
||||
socketio.emit("stop_container_log_stream", {
|
||||
"request_id": request_id
|
||||
}, to=connected_workers[worker_id])
|
||||
else:
|
||||
logger.warning(f"[{request_id}] Worker {worker_id} not connected on stop")
|
||||
|
||||
redis_client.delete(f"log_request_context:{request_id}")
|
||||
logger.debug(f"[{request_id}] Deleted context from Redis")
|
||||
else:
|
||||
logger.warning(f"[{request_id}] No context found in Redis for stopping stream")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"[{request_id}] Error stopping log stream: {e}")
|
||||
|
||||
|
||||
@socketio.on("container-log-stream")
|
||||
def handle_streamed_logs(data):
|
||||
request_id = data.get("request_id")
|
||||
logs = data.get("logs")
|
||||
logger.debug(f"[{request_id}] Received container-log-stream update")
|
||||
|
||||
try:
|
||||
redis_client = get_redis_client()
|
||||
context_data = redis_client.get(f"log_request_context:{request_id}")
|
||||
if context_data:
|
||||
context = json.loads(context_data)
|
||||
user_sid = context.get("user_sid")
|
||||
logger.debug(f"[{request_id}] Forwarding logs to user sid={user_sid}")
|
||||
socketio.emit("user_log_stream_update", {
|
||||
"logs": logs,
|
||||
"request_id": request_id
|
||||
}, to=user_sid)
|
||||
else:
|
||||
logger.warning(f"[{request_id}] No context found in Redis to forward logs")
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"[{request_id}] Error forwarding streamed log: {e}")
|
||||
|
||||
|
||||
@socketio.on("libvirt_event")
|
||||
def handle_libvirt_event(data):
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import threading
|
||||
import socketio
|
||||
from logger import logger
|
||||
import json
|
||||
@@ -46,6 +47,9 @@ class WorkerClient:
|
||||
self.sio.on("message", self.on_message)
|
||||
self.sio.on("*", self.debug_all_events)
|
||||
self.sio.on("container-log", self.handle_container_log_request)
|
||||
self.sio.on("start_container_log_stream", self.start_log_stream)
|
||||
self.sio.on("stop_container_log_stream", self.stop_log_stream)
|
||||
self.log_stream_threads = {}
|
||||
|
||||
|
||||
async def on_connect(self):
|
||||
@@ -217,6 +221,74 @@ class WorkerClient:
|
||||
"""Stop the worker client."""
|
||||
self.running = False
|
||||
|
||||
|
||||
|
||||
|
||||
async def start_log_stream(self, data):
|
||||
container_name = data["container_name"]
|
||||
request_id = data["request_id"]
|
||||
|
||||
logger.info(f"Received log stream request: container={container_name}, request_id={request_id}")
|
||||
|
||||
if request_id in self.log_stream_threads:
|
||||
logger.warning(f"Log stream for request {request_id} already active.")
|
||||
return
|
||||
|
||||
# Set cancellation flag
|
||||
self.log_stream_threads[request_id] = True
|
||||
logger.debug(f"Log stream flag set for request {request_id}")
|
||||
|
||||
# Get the current (main) event loop
|
||||
try:
|
||||
loop = asyncio.get_running_loop()
|
||||
logger.debug("Retrieved current running event loop.")
|
||||
except RuntimeError as e:
|
||||
logger.error(f"Failed to get running event loop: {e}")
|
||||
return
|
||||
|
||||
def stream_logs():
|
||||
try:
|
||||
logger.info(f"Attempting to get container '{container_name}'")
|
||||
container = self.docker_monitor.docker_client.containers.get(container_name)
|
||||
logger.info(f"Successfully attached to container '{container_name}' for log streaming.")
|
||||
|
||||
for line in container.logs(stream=True, follow=True):
|
||||
if not self.log_stream_threads.get(request_id):
|
||||
logger.info(f"Log stream for request {request_id} cancelled.")
|
||||
break
|
||||
|
||||
log_line = line.decode("utf-8").strip()
|
||||
logger.debug(f"Streaming log line for request {request_id}: {log_line}")
|
||||
|
||||
asyncio.run_coroutine_threadsafe(
|
||||
self.sio.emit("container-log-stream", {
|
||||
"request_id": request_id,
|
||||
"logs": log_line
|
||||
}),
|
||||
loop
|
||||
)
|
||||
|
||||
logger.info(f"Log stream for request {request_id} has ended.")
|
||||
except Exception as e:
|
||||
logger.exception(f"Exception while streaming logs for container '{container_name}': {e}")
|
||||
asyncio.run_coroutine_threadsafe(
|
||||
self.sio.emit("container-log-stream", {
|
||||
"request_id": request_id,
|
||||
"logs": f"Stream error: {str(e)}"
|
||||
}),
|
||||
loop
|
||||
)
|
||||
|
||||
logger.debug(f"Starting log stream thread for request {request_id}")
|
||||
asyncio.create_task(asyncio.to_thread(stream_logs))
|
||||
|
||||
|
||||
|
||||
async def stop_log_stream(self, data):
|
||||
request_id = data["request_id"]
|
||||
if request_id in self.log_stream_threads:
|
||||
self.log_stream_threads[request_id] = False
|
||||
|
||||
async def handle_container_log_request(self, data):
|
||||
"""
|
||||
Handles a container-log request and delegates to ContainerTask.
|
||||
|
||||
@@ -183,7 +183,6 @@ def get_network_interfaces():
|
||||
|
||||
return network_interfaces
|
||||
|
||||
|
||||
def container_test():
|
||||
"""Run a simple test to verify Docker functionality"""
|
||||
try:
|
||||
@@ -197,7 +196,6 @@ def container_test():
|
||||
logger.error(f"Container test failed: {e}")
|
||||
return False
|
||||
|
||||
|
||||
def libvirt_test():
|
||||
"""Run a simple test to verify libvirt functionality"""
|
||||
try:
|
||||
@@ -212,7 +210,6 @@ def libvirt_test():
|
||||
logger.error(f"Libvirt test failed: {e}")
|
||||
return False
|
||||
|
||||
|
||||
# Worker process function
|
||||
def worker_process_func(event_queue, docker_monitor=None):
|
||||
"""Function to run the worker client in a separate process using global settings."""
|
||||
@@ -233,7 +230,6 @@ def worker_process_func(event_queue, docker_monitor=None):
|
||||
if 'worker' in locals():
|
||||
worker.stop()
|
||||
|
||||
|
||||
def main():
|
||||
# Settings are loaded automatically when worker.settings is imported
|
||||
|
||||
@@ -364,6 +360,5 @@ def main():
|
||||
for process in processes:
|
||||
process.join()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
|
||||
Reference in New Issue
Block a user