Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac843da4ca | ||
|
|
f6b7eb2dae |
+1
-1
@@ -33,7 +33,7 @@ logger = logging.getLogger(__name__)
|
||||
app = FastAPI(
|
||||
title="NetBird MSP Appliance",
|
||||
description="Multi-tenant NetBird management platform for MSPs",
|
||||
version="1.0.0",
|
||||
version="1.1.0",
|
||||
docs_url="/api/docs",
|
||||
redoc_url="/api/redoc",
|
||||
openapi_url="/api/openapi.json",
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
"""Monitoring API — system overview, customer statuses, host resources."""
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import platform
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
import psutil
|
||||
@@ -16,6 +18,13 @@ from app.services import docker_service, image_service
|
||||
logger = logging.getLogger(__name__)
|
||||
router = APIRouter()
|
||||
|
||||
# Short-lived cache for the local update-status badges. This endpoint is
|
||||
# triggered on every customer-table render (i.e. every search keystroke), but
|
||||
# the underlying data (which images are outdated) only changes after an image
|
||||
# pull + container recreate, so a few seconds of staleness is harmless.
|
||||
_update_status_cache: dict[str, Any] = {"data": None, "expires": 0.0}
|
||||
_UPDATE_STATUS_TTL_SECONDS = 20
|
||||
|
||||
|
||||
@router.get("/status")
|
||||
async def system_status(
|
||||
@@ -58,8 +67,7 @@ async def all_customers_status(
|
||||
.all()
|
||||
)
|
||||
|
||||
results: list[dict[str, Any]] = []
|
||||
for c in customers:
|
||||
async def _build_entry(c: Customer) -> dict[str, Any]:
|
||||
entry: dict[str, Any] = {
|
||||
"id": c.id,
|
||||
"name": c.name,
|
||||
@@ -67,7 +75,7 @@ async def all_customers_status(
|
||||
"status": c.status,
|
||||
}
|
||||
if c.deployment:
|
||||
containers = docker_service.get_container_status(c.deployment.container_prefix)
|
||||
containers = await docker_service.get_container_status_async(c.deployment.container_prefix)
|
||||
entry["deployment_status"] = c.deployment.deployment_status
|
||||
entry["containers"] = containers
|
||||
entry["relay_udp_port"] = c.deployment.relay_udp_port
|
||||
@@ -76,9 +84,11 @@ async def all_customers_status(
|
||||
else:
|
||||
entry["deployment_status"] = None
|
||||
entry["containers"] = []
|
||||
results.append(entry)
|
||||
return entry
|
||||
|
||||
return results
|
||||
# Fetch container status for all customers concurrently instead of one
|
||||
# blocking Docker SDK call at a time.
|
||||
return await asyncio.gather(*[_build_entry(c) for c in customers])
|
||||
|
||||
|
||||
@router.get("/resources")
|
||||
@@ -205,15 +215,28 @@ async def customers_local_update_status(
|
||||
|
||||
Compares running container image IDs against locally stored images.
|
||||
No network call — safe to call on every dashboard load.
|
||||
|
||||
Results are cached for a few seconds since this is triggered on every
|
||||
customer-table render (including every search keystroke) but the
|
||||
underlying data rarely changes.
|
||||
"""
|
||||
now = time.monotonic()
|
||||
if _update_status_cache["data"] is not None and now < _update_status_cache["expires"]:
|
||||
return _update_status_cache["data"]
|
||||
|
||||
config = db.query(SystemConfig).filter(SystemConfig.id == 1).first()
|
||||
if not config:
|
||||
return []
|
||||
deployments = db.query(Deployment).all()
|
||||
results = []
|
||||
for dep in deployments:
|
||||
cs = image_service.get_customer_container_image_status(dep.container_prefix, config)
|
||||
results.append({"customer_id": dep.customer_id, "needs_update": cs["needs_update"]})
|
||||
|
||||
async def _check(dep: Deployment) -> dict[str, Any]:
|
||||
cs = await image_service.get_customer_container_image_status_async(dep.container_prefix, config)
|
||||
return {"customer_id": dep.customer_id, "needs_update": cs["needs_update"]}
|
||||
|
||||
results = await asyncio.gather(*[_check(dep) for dep in deployments])
|
||||
results = list(results)
|
||||
_update_status_cache["data"] = results
|
||||
_update_status_cache["expires"] = now + _UPDATE_STATUS_TTL_SECONDS
|
||||
return results
|
||||
|
||||
|
||||
|
||||
@@ -27,13 +27,24 @@ async def _run_cmd(cmd: list[str], timeout: int = 120) -> subprocess.CompletedPr
|
||||
)
|
||||
|
||||
|
||||
_client: Optional[docker.DockerClient] = None
|
||||
|
||||
|
||||
def _get_client() -> docker.DockerClient:
|
||||
"""Return a Docker client connected via the Unix socket.
|
||||
"""Return a shared Docker client connected via the Unix socket.
|
||||
|
||||
The client is created once and reused — creating a new client per call
|
||||
(as `docker.from_env()` does) re-negotiates the API version and opens a
|
||||
fresh connection every time, which is wasteful when called once per
|
||||
customer in a loop.
|
||||
|
||||
Returns:
|
||||
docker.DockerClient instance.
|
||||
"""
|
||||
return docker.from_env()
|
||||
global _client
|
||||
if _client is None:
|
||||
_client = docker.from_env()
|
||||
return _client
|
||||
|
||||
|
||||
async def compose_up(
|
||||
@@ -212,6 +223,16 @@ def get_container_status(container_prefix: str) -> list[dict[str, Any]]:
|
||||
return results
|
||||
|
||||
|
||||
async def get_container_status_async(container_prefix: str) -> list[dict[str, Any]]:
|
||||
"""Thread-offloaded wrapper around get_container_status().
|
||||
|
||||
Use this when checking status for multiple customers so the Docker SDK
|
||||
calls run in the thread pool instead of blocking the event loop.
|
||||
"""
|
||||
loop = asyncio.get_event_loop()
|
||||
return await loop.run_in_executor(None, get_container_status, container_prefix)
|
||||
|
||||
|
||||
def get_container_logs(container_name: str, tail: int = 200) -> str:
|
||||
"""Retrieve recent logs from a container.
|
||||
|
||||
|
||||
@@ -211,6 +211,44 @@ async def pull_all_images(config) -> dict[str, Any]:
|
||||
}
|
||||
|
||||
|
||||
async def get_customer_container_image_status_async(container_prefix: str, config) -> dict[str, Any]:
|
||||
"""Async, thread-offloaded version of get_customer_container_image_status().
|
||||
|
||||
Runs the per-service `docker inspect` subprocess calls concurrently in the
|
||||
thread pool instead of sequentially blocking the event loop — use this
|
||||
whenever checking status for multiple customers (e.g. dashboard/search
|
||||
badge refresh, monitoring overview).
|
||||
|
||||
Returns:
|
||||
services: dict mapping service name to status info
|
||||
needs_update: True if any service has a different image ID than locally stored
|
||||
"""
|
||||
service_images = {
|
||||
"management": config.netbird_management_image,
|
||||
"signal": config.netbird_signal_image,
|
||||
"relay": config.netbird_relay_image,
|
||||
"dashboard": config.netbird_dashboard_image,
|
||||
}
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
async def _check(svc: str, image: str) -> tuple[str, dict[str, Any]]:
|
||||
container_name = f"{container_prefix}-{svc}"
|
||||
container_id, local_id = await asyncio.gather(
|
||||
loop.run_in_executor(None, get_container_image_id, container_name),
|
||||
loop.run_in_executor(None, get_local_image_id, image),
|
||||
)
|
||||
if container_id and local_id:
|
||||
up_to_date = container_id == local_id
|
||||
else:
|
||||
up_to_date = None # container not running or image not pulled
|
||||
return svc, {"container": container_name, "image": image, "up_to_date": up_to_date}
|
||||
|
||||
pairs = await asyncio.gather(*[_check(svc, image) for svc, image in service_images.items()])
|
||||
services = dict(pairs)
|
||||
needs_update = any(s["up_to_date"] is False for s in services.values())
|
||||
return {"services": services, "needs_update": needs_update}
|
||||
|
||||
|
||||
def get_customer_container_image_status(container_prefix: str, config) -> dict[str, Any]:
|
||||
"""Check which service containers are running outdated local images.
|
||||
|
||||
|
||||
@@ -263,10 +263,14 @@ async def create_proxy_host(
|
||||
"location ^~ /management.ManagementService/ {\n"
|
||||
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
||||
" grpc_set_header Host $host;\n"
|
||||
" grpc_read_timeout 3600s;\n"
|
||||
" grpc_send_timeout 3600s;\n"
|
||||
"}\n"
|
||||
"location ^~ /signalexchange.SignalExchange/ {\n"
|
||||
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
||||
" grpc_set_header Host $host;\n"
|
||||
" grpc_read_timeout 3600s;\n"
|
||||
" grpc_send_timeout 3600s;\n"
|
||||
"}\n"
|
||||
),
|
||||
"meta": {
|
||||
|
||||
Reference in New Issue
Block a user