fix(workers): harden remote telemetry
This commit is contained in:
@@ -150,7 +150,9 @@ class WorkerCapacity:
|
||||
worker_id: str
|
||||
max_concurrent_tasks: int = 1
|
||||
active_tasks: int = 0
|
||||
free_memory_bytes: int = 0
|
||||
# ``None`` means the worker could not query VRAM. It is distinct from a
|
||||
# real zero-byte reading, which means the device is completely occupied.
|
||||
free_memory_bytes: Optional[int] = None
|
||||
cpu_percent: Optional[float] = None
|
||||
gpu_utilization_percent: Optional[float] = None
|
||||
backend: str = ""
|
||||
|
||||
@@ -118,7 +118,7 @@ class Target:
|
||||
active_tasks: int = 0
|
||||
max_tasks: int = 0
|
||||
cpu_percent: Optional[float] = None
|
||||
free_memory_bytes: int = 0
|
||||
free_memory_bytes: Optional[int] = None
|
||||
system_memory_bytes: int = 0
|
||||
cpu_count: int = 0
|
||||
gpu_name: str = ""
|
||||
@@ -224,7 +224,7 @@ def list_targets(control_plane=None) -> list[Target]:
|
||||
active_tasks=live.capacity.active_tasks if live else 0,
|
||||
max_tasks=live.capacity.max_concurrent_tasks if live else 0,
|
||||
cpu_percent=live.capacity.cpu_percent if live else None,
|
||||
free_memory_bytes=live.capacity.free_memory_bytes if live else 0,
|
||||
free_memory_bytes=live.capacity.free_memory_bytes if live else None,
|
||||
system_memory_bytes=int(host.get("system_memory_bytes") or 0),
|
||||
cpu_count=int(host.get("cpu_count") or 0),
|
||||
gpu_name=str(gpu.get("model") or ""),
|
||||
|
||||
@@ -367,6 +367,9 @@ class WorkerClient:
|
||||
self._keepalives: dict[str, asyncio.Task] = {}
|
||||
self._maintenance: set[asyncio.Task] = set()
|
||||
self._telemetry: tuple[Optional[float], Optional[int], Optional[float]] = (None, None, None)
|
||||
# A driver query can hang indefinitely. Keep that one query owned
|
||||
# rather than cancelling its awaiter and starting a fresh thread at
|
||||
# every heartbeat.
|
||||
self._telemetry_task: Optional[asyncio.Task] = None
|
||||
self._prewarms: dict[str, asyncio.Task] = {}
|
||||
self._prewarm_cancellations: dict[str, asyncio.Task] = {}
|
||||
@@ -712,15 +715,37 @@ class WorkerClient:
|
||||
async def _heartbeat_loop(self, interval: float) -> None:
|
||||
while True:
|
||||
await asyncio.sleep(interval)
|
||||
await self._refresh_telemetry()
|
||||
await self._send(self.heartbeat_message())
|
||||
if self._telemetry_task is None or self._telemetry_task.done():
|
||||
self._telemetry_task = asyncio.create_task(self._refresh_telemetry())
|
||||
|
||||
async def _refresh_telemetry(self) -> None:
|
||||
"""Publish completed samples and retain one non-blocking probe.
|
||||
|
||||
CUDA/NVML calls may wedge in a driver. A timed ``to_thread`` await
|
||||
only cancels the awaiter, leaving that thread alive; retaining this
|
||||
task prevents later heartbeats from accumulating more blocked probes.
|
||||
"""
|
||||
task = self._telemetry_task
|
||||
if task is not None and task.done():
|
||||
try:
|
||||
self._telemetry = await asyncio.wait_for(asyncio.to_thread(_heartbeat_resources), timeout=2)
|
||||
except TimeoutError:
|
||||
logger.warning("Worker telemetry sampling timed out")
|
||||
sampled = task.result()
|
||||
except Exception:
|
||||
logger.debug("Could not sample worker telemetry", exc_info=True)
|
||||
else:
|
||||
# A partial failed sample must not erase an independent last
|
||||
# good value. Presence on the heartbeat remains honest until
|
||||
# that individual metric can next be measured.
|
||||
self._telemetry = tuple(
|
||||
current if value is None else value
|
||||
for current, value in zip(self._telemetry, sampled)
|
||||
)
|
||||
self._telemetry_task = None
|
||||
|
||||
if self._telemetry_task is None:
|
||||
self._telemetry_task = asyncio.create_task(
|
||||
to_thread_and_drain_on_cancel(_heartbeat_resources),
|
||||
name="worker-telemetry-probe",
|
||||
)
|
||||
|
||||
def heartbeat_message(self) -> pb.WorkerMessage:
|
||||
"""Build the worker's current liveness/capacity frame."""
|
||||
|
||||
@@ -506,23 +506,28 @@ export function StatusBar({ compact = false }: { compact?: boolean }) {
|
||||
percent={activeRemoteTarget.cpu_percent}
|
||||
/>
|
||||
)}
|
||||
{activeRemoteTarget.gpu_name && (
|
||||
{activeRemoteTarget.gpu_name &&
|
||||
(activeRemoteTarget.gpu_utilization_percent != null ||
|
||||
activeRemoteTarget.free_memory_bytes != null) && (
|
||||
<DeviceMetric
|
||||
Icon={MonitorUpIcon}
|
||||
label={t('settings.device_family_gpu')}
|
||||
value={
|
||||
[
|
||||
activeRemoteTarget.gpu_utilization_percent != null
|
||||
? `${Math.round(activeRemoteTarget.gpu_utilization_percent)}% · ${formatBytes(
|
||||
activeRemoteTarget.gpu_memory_bytes - activeRemoteTarget.free_memory_bytes,
|
||||
)} / ${formatBytes(activeRemoteTarget.gpu_memory_bytes)}`
|
||||
: activeRemoteTarget.free_memory_bytes > 0
|
||||
? `${Math.round(activeRemoteTarget.gpu_utilization_percent)}%`
|
||||
: null,
|
||||
activeRemoteTarget.free_memory_bytes != null
|
||||
? `${formatBytes(
|
||||
activeRemoteTarget.gpu_memory_bytes - activeRemoteTarget.free_memory_bytes,
|
||||
)} / ${formatBytes(activeRemoteTarget.gpu_memory_bytes)}`
|
||||
: formatBytes(activeRemoteTarget.gpu_memory_bytes)
|
||||
: null,
|
||||
]
|
||||
.filter((value): value is string => value != null)
|
||||
.join(' · ')
|
||||
}
|
||||
percent={boundedPercent(
|
||||
activeRemoteTarget.free_memory_bytes > 0
|
||||
activeRemoteTarget.free_memory_bytes != null
|
||||
? activeRemoteTarget.gpu_memory_bytes - activeRemoteTarget.free_memory_bytes
|
||||
: 0,
|
||||
activeRemoteTarget.gpu_memory_bytes,
|
||||
|
||||
@@ -19,7 +19,7 @@ export interface ComputeTarget {
|
||||
active_tasks: number;
|
||||
max_tasks: number;
|
||||
cpu_percent: number | null;
|
||||
free_memory_bytes: number;
|
||||
free_memory_bytes: number | null;
|
||||
system_memory_bytes: number;
|
||||
cpu_count: number;
|
||||
gpu_name: string;
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "تعذّر تحميل إعداد الجهاز",
|
||||
"perf_save_failed": "تعذّر حفظ الإعداد",
|
||||
"generate_budget_title": "ميزانية زمن الحوسبة",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Geräteeinstellung konnte nicht geladen werden",
|
||||
"perf_save_failed": "Einstellung konnte nicht gespeichert werden",
|
||||
"generate_budget_title": "Rechenzeit-Budget",
|
||||
|
||||
@@ -959,6 +959,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Failed to load device setting",
|
||||
"generate_timeout_shadowed_badge": "Overridden externally",
|
||||
"generate_timeout_shadowed_note": "An environment variable outside VoiceStudio (shell, .env file, or container) is currently setting this — your saved value here is ignored until that variable is removed.",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "No se pudo cargar el ajuste del dispositivo",
|
||||
"perf_save_failed": "No se pudo guardar el ajuste",
|
||||
"generate_budget_title": "Presupuesto de tiempo de cómputo",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Impossible de charger le réglage du périphérique",
|
||||
"perf_save_failed": "Impossible d'enregistrer le réglage",
|
||||
"generate_budget_title": "Budget de temps de calcul",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "डिवाइस सेटिंग लोड नहीं हो सकी",
|
||||
"perf_save_failed": "सेटिंग सहेजी नहीं जा सकी",
|
||||
"generate_budget_title": "कंप्यूट-समय बजट",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Gagal memuat pengaturan perangkat",
|
||||
"perf_save_failed": "Gagal menyimpan pengaturan",
|
||||
"generate_budget_title": "Anggaran waktu komputasi",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Impossibile caricare l'impostazione del dispositivo",
|
||||
"perf_save_failed": "Impossibile salvare l'impostazione",
|
||||
"generate_budget_title": "Budget di tempo di calcolo",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "デバイス設定を読み込めませんでした",
|
||||
"perf_save_failed": "設定を保存できませんでした",
|
||||
"generate_budget_title": "計算時間の予算",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "장치 설정을 불러오지 못했습니다",
|
||||
"perf_save_failed": "설정을 저장하지 못했습니다",
|
||||
"audio_tools": "오디오 도구",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Apparaatinstelling kon niet worden geladen",
|
||||
"perf_save_failed": "Instelling kon niet worden opgeslagen",
|
||||
"generate_budget_title": "Rekentijdbudget",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Nie udało się wczytać ustawienia urządzenia",
|
||||
"perf_save_failed": "Nie udało się zapisać ustawienia",
|
||||
"generate_budget_title": "Budżet czasu obliczeń",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Falha ao carregar a configuração do dispositivo",
|
||||
"perf_save_failed": "Falha ao salvar a configuração",
|
||||
"generate_budget_title": "Orçamento de tempo de computação",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Не удалось загрузить настройку устройства",
|
||||
"perf_save_failed": "Не удалось сохранить настройку",
|
||||
"generate_budget_title": "Бюджет времени вычислений",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Kunde inte läsa in enhetsinställningen",
|
||||
"perf_save_failed": "Kunde inte spara inställningen",
|
||||
"generate_budget_title": "Beräkningstidsbudget",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "โหลดการตั้งค่าอุปกรณ์ไม่สำเร็จ",
|
||||
"perf_save_failed": "บันทึกการตั้งค่าไม่สำเร็จ",
|
||||
"generate_budget_title": "งบเวลาในการประมวลผล",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Aygıt ayarı yüklenemedi",
|
||||
"perf_save_failed": "Ayar kaydedilemedi",
|
||||
"generate_budget_title": "Hesaplama süresi bütçesi",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Не вдалося завантажити налаштування пристрою",
|
||||
"perf_save_failed": "Не вдалося зберегти налаштування",
|
||||
"generate_budget_title": "Бюджет часу обчислень",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "Không tải được cài đặt thiết bị",
|
||||
"perf_save_failed": "Không lưu được cài đặt",
|
||||
"generate_budget_title": "Ngân sách thời gian tính toán",
|
||||
|
||||
@@ -822,6 +822,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "无法加载设备设置",
|
||||
"perf_save_failed": "无法保存设置",
|
||||
"generate_budget_title": "计算时长预算",
|
||||
|
||||
@@ -340,6 +340,7 @@
|
||||
"device_family_npu": "NPU",
|
||||
"device_family_mps": "Apple GPU (MPS)",
|
||||
"device_family_cpu": "CPU",
|
||||
"device_family_gpu": "GPU",
|
||||
"device_load_failed": "無法載入裝置設定",
|
||||
"perf_save_failed": "無法儲存設定",
|
||||
"generate_budget_title": "運算時間預算",
|
||||
|
||||
@@ -6,6 +6,7 @@ import pytest
|
||||
|
||||
from worker.identity import WorkerKeypair
|
||||
from worker.protocol.gen import worker_v1_pb2 as pb
|
||||
from worker.transport import client as client_module
|
||||
from worker.transport.client import WorkerClient, WorkerConfig
|
||||
from worker.transport.server import WorkerServicer
|
||||
|
||||
@@ -19,6 +20,50 @@ def _client(probe):
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_blocked_telemetry_probe_is_not_restarted(monkeypatch):
|
||||
"""A stuck driver call occupies one worker thread, never one per heartbeat."""
|
||||
started = threading.Event()
|
||||
release = threading.Event()
|
||||
calls = 0
|
||||
|
||||
def sample():
|
||||
nonlocal calls
|
||||
calls += 1
|
||||
started.set()
|
||||
release.wait(5)
|
||||
return 10.0, 20, 30.0
|
||||
|
||||
monkeypatch.setattr(client_module, "_heartbeat_resources", sample)
|
||||
client = _client(lambda: [])
|
||||
await client._refresh_telemetry()
|
||||
await asyncio.wait_for(asyncio.to_thread(started.wait), timeout=1)
|
||||
|
||||
for _ in range(3):
|
||||
await client._refresh_telemetry()
|
||||
assert calls == 1
|
||||
|
||||
release.set()
|
||||
await asyncio.wait_for(client._telemetry_task, timeout=1)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_partial_telemetry_failure_retains_other_last_good_values(monkeypatch):
|
||||
samples = iter(((10.0, 20, 30.0), (None, None, 40.0)))
|
||||
monkeypatch.setattr(
|
||||
client_module, "_heartbeat_resources", lambda: next(samples, (None, None, None))
|
||||
)
|
||||
client = _client(lambda: [])
|
||||
|
||||
await client._refresh_telemetry()
|
||||
await asyncio.wait_for(client._telemetry_task, timeout=1)
|
||||
await client._refresh_telemetry()
|
||||
await asyncio.wait_for(client._telemetry_task, timeout=1)
|
||||
await client._refresh_telemetry()
|
||||
|
||||
assert client._telemetry == (10.0, 20, 40.0)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_refresh_sends_capability_update():
|
||||
client = _client(lambda: [{
|
||||
|
||||
@@ -140,6 +140,27 @@ def test_connected_worker_target_exposes_heartbeat_telemetry(db, settings):
|
||||
assert target["gpu_memory_bytes"] == 24 * 1024**3
|
||||
|
||||
|
||||
def test_target_distinguishes_unknown_vram_from_a_full_gpu(db, settings):
|
||||
plane = _Plane()
|
||||
worker = _enroll(
|
||||
"desktop-4090",
|
||||
host={"gpus": [{"model": "NVIDIA RTX 4090", "memory_bytes": 24 * 1024**3}]},
|
||||
)
|
||||
_connect(plane, worker)
|
||||
|
||||
unknown = routing.list_targets(plane)[1].to_dict()
|
||||
assert unknown["free_memory_bytes"] is None
|
||||
|
||||
plane.pool.heartbeat(
|
||||
worker.id,
|
||||
active_tasks=0,
|
||||
available_slots=1,
|
||||
free_memory_bytes=0,
|
||||
)
|
||||
full = routing.list_targets(plane)[1].to_dict()
|
||||
assert full["free_memory_bytes"] == 0
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"setup,detail",
|
||||
[
|
||||
|
||||
Reference in New Issue
Block a user