Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac843da4ca | ||
|
|
f6b7eb2dae |
+1
-1
@@ -33,7 +33,7 @@ logger = logging.getLogger(__name__)
|
|||||||
app = FastAPI(
|
app = FastAPI(
|
||||||
title="NetBird MSP Appliance",
|
title="NetBird MSP Appliance",
|
||||||
description="Multi-tenant NetBird management platform for MSPs",
|
description="Multi-tenant NetBird management platform for MSPs",
|
||||||
version="1.0.0",
|
version="1.1.0",
|
||||||
docs_url="/api/docs",
|
docs_url="/api/docs",
|
||||||
redoc_url="/api/redoc",
|
redoc_url="/api/redoc",
|
||||||
openapi_url="/api/openapi.json",
|
openapi_url="/api/openapi.json",
|
||||||
|
|||||||
@@ -1,7 +1,9 @@
|
|||||||
"""Monitoring API — system overview, customer statuses, host resources."""
|
"""Monitoring API — system overview, customer statuses, host resources."""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
import platform
|
import platform
|
||||||
|
import time
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import psutil
|
import psutil
|
||||||
@@ -16,6 +18,13 @@ from app.services import docker_service, image_service
|
|||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
router = APIRouter()
|
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")
|
@router.get("/status")
|
||||||
async def system_status(
|
async def system_status(
|
||||||
@@ -58,8 +67,7 @@ async def all_customers_status(
|
|||||||
.all()
|
.all()
|
||||||
)
|
)
|
||||||
|
|
||||||
results: list[dict[str, Any]] = []
|
async def _build_entry(c: Customer) -> dict[str, Any]:
|
||||||
for c in customers:
|
|
||||||
entry: dict[str, Any] = {
|
entry: dict[str, Any] = {
|
||||||
"id": c.id,
|
"id": c.id,
|
||||||
"name": c.name,
|
"name": c.name,
|
||||||
@@ -67,7 +75,7 @@ async def all_customers_status(
|
|||||||
"status": c.status,
|
"status": c.status,
|
||||||
}
|
}
|
||||||
if c.deployment:
|
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["deployment_status"] = c.deployment.deployment_status
|
||||||
entry["containers"] = containers
|
entry["containers"] = containers
|
||||||
entry["relay_udp_port"] = c.deployment.relay_udp_port
|
entry["relay_udp_port"] = c.deployment.relay_udp_port
|
||||||
@@ -76,9 +84,11 @@ async def all_customers_status(
|
|||||||
else:
|
else:
|
||||||
entry["deployment_status"] = None
|
entry["deployment_status"] = None
|
||||||
entry["containers"] = []
|
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")
|
@router.get("/resources")
|
||||||
@@ -205,15 +215,28 @@ async def customers_local_update_status(
|
|||||||
|
|
||||||
Compares running container image IDs against locally stored images.
|
Compares running container image IDs against locally stored images.
|
||||||
No network call — safe to call on every dashboard load.
|
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()
|
config = db.query(SystemConfig).filter(SystemConfig.id == 1).first()
|
||||||
if not config:
|
if not config:
|
||||||
return []
|
return []
|
||||||
deployments = db.query(Deployment).all()
|
deployments = db.query(Deployment).all()
|
||||||
results = []
|
|
||||||
for dep in deployments:
|
async def _check(dep: Deployment) -> dict[str, Any]:
|
||||||
cs = image_service.get_customer_container_image_status(dep.container_prefix, config)
|
cs = await image_service.get_customer_container_image_status_async(dep.container_prefix, config)
|
||||||
results.append({"customer_id": dep.customer_id, "needs_update": cs["needs_update"]})
|
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
|
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:
|
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:
|
Returns:
|
||||||
docker.DockerClient instance.
|
docker.DockerClient instance.
|
||||||
"""
|
"""
|
||||||
return docker.from_env()
|
global _client
|
||||||
|
if _client is None:
|
||||||
|
_client = docker.from_env()
|
||||||
|
return _client
|
||||||
|
|
||||||
|
|
||||||
async def compose_up(
|
async def compose_up(
|
||||||
@@ -212,6 +223,16 @@ def get_container_status(container_prefix: str) -> list[dict[str, Any]]:
|
|||||||
return results
|
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:
|
def get_container_logs(container_name: str, tail: int = 200) -> str:
|
||||||
"""Retrieve recent logs from a container.
|
"""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]:
|
def get_customer_container_image_status(container_prefix: str, config) -> dict[str, Any]:
|
||||||
"""Check which service containers are running outdated local images.
|
"""Check which service containers are running outdated local images.
|
||||||
|
|
||||||
|
|||||||
@@ -263,10 +263,14 @@ async def create_proxy_host(
|
|||||||
"location ^~ /management.ManagementService/ {\n"
|
"location ^~ /management.ManagementService/ {\n"
|
||||||
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
||||||
" grpc_set_header Host $host;\n"
|
" grpc_set_header Host $host;\n"
|
||||||
|
" grpc_read_timeout 3600s;\n"
|
||||||
|
" grpc_send_timeout 3600s;\n"
|
||||||
"}\n"
|
"}\n"
|
||||||
"location ^~ /signalexchange.SignalExchange/ {\n"
|
"location ^~ /signalexchange.SignalExchange/ {\n"
|
||||||
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
f" grpc_pass grpc://{forward_host}:{forward_port};\n"
|
||||||
" grpc_set_header Host $host;\n"
|
" grpc_set_header Host $host;\n"
|
||||||
|
" grpc_read_timeout 3600s;\n"
|
||||||
|
" grpc_send_timeout 3600s;\n"
|
||||||
"}\n"
|
"}\n"
|
||||||
),
|
),
|
||||||
"meta": {
|
"meta": {
|
||||||
|
|||||||
Reference in New Issue
Block a user