Merge branch 'velixio_dev' into fix/linux-electron-setup-sidebar
This commit is contained in:
@@ -43,6 +43,8 @@ import platform
|
||||
import random
|
||||
import socket
|
||||
import sys
|
||||
import threading
|
||||
from concurrent.futures import Future
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Awaitable, Callable, Optional, Protocol
|
||||
|
||||
@@ -370,7 +372,7 @@ class WorkerClient:
|
||||
# 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._telemetry_task: Optional[asyncio.Future] = None
|
||||
self._prewarms: dict[str, asyncio.Task] = {}
|
||||
self._prewarm_cancellations: dict[str, asyncio.Task] = {}
|
||||
self._epoch = 0
|
||||
@@ -476,7 +478,6 @@ class WorkerClient:
|
||||
*draining, return_exceptions=True
|
||||
)
|
||||
self._maintenance.clear()
|
||||
self._telemetry_task = None
|
||||
self._prewarms.clear()
|
||||
self._prewarm_cancellations.clear()
|
||||
for key, task in running:
|
||||
@@ -745,12 +746,22 @@ class WorkerClient:
|
||||
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",
|
||||
)
|
||||
self._maintenance.add(self._telemetry_task)
|
||||
self._telemetry_task.add_done_callback(self._maintenance.discard)
|
||||
# Read-only driver probes cannot be interrupted. Keep one across
|
||||
# reconnects, outside assignment drain and the shared executor
|
||||
# (whose shutdown would otherwise wait forever for a wedged driver).
|
||||
result = Future()
|
||||
self._telemetry_task = asyncio.wrap_future(result)
|
||||
|
||||
def sample() -> None:
|
||||
try:
|
||||
result.set_result(_heartbeat_resources())
|
||||
except Exception:
|
||||
logger.debug("Could not sample worker telemetry", exc_info=True)
|
||||
result.set_result((None, None, None))
|
||||
|
||||
threading.Thread(
|
||||
target=sample, name="worker-telemetry-probe", daemon=True,
|
||||
).start()
|
||||
|
||||
def heartbeat_message(self) -> pb.WorkerMessage:
|
||||
"""Build the worker's current liveness/capacity frame."""
|
||||
|
||||
@@ -400,4 +400,5 @@ precondition or automated check exits non-zero.
|
||||
Remote compute targets show available CPU/GPU usage and free VRAM. Unavailable
|
||||
metrics are omitted; a transient sampling failure retains the last successful
|
||||
reading. Telemetry runs off the control loop with at most one probe per worker
|
||||
client; shutdown drains that probe before retiring the client.
|
||||
client, retained across reconnects. Read-only probes never block task draining or
|
||||
shutdown; a stuck driver probe cannot accumulate more threads.
|
||||
|
||||
@@ -544,23 +544,36 @@ async def test_control_plane_stamps_authenticated_target_on_progress(monkeypatch
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_shutdown_drains_owned_telemetry_probe(monkeypatch):
|
||||
async def test_blocked_telemetry_does_not_block_drain_stop_or_duplicate_on_reconnect(monkeypatch):
|
||||
started = threading.Event()
|
||||
release = threading.Event()
|
||||
calls = []
|
||||
|
||||
def sample():
|
||||
calls.append(threading.current_thread())
|
||||
started.set()
|
||||
release.wait(5)
|
||||
return 1.0, 2, 3.0
|
||||
|
||||
monkeypatch.setattr(client_module, "_heartbeat_resources", sample)
|
||||
client = _client(lambda: [])
|
||||
await client._refresh_telemetry()
|
||||
await asyncio.wait_for(asyncio.to_thread(started.wait), 1)
|
||||
stopping = asyncio.create_task(client._cancel_active_work())
|
||||
await asyncio.sleep(0.02)
|
||||
assert not stopping.done()
|
||||
release.set()
|
||||
await asyncio.wait_for(stopping, 1)
|
||||
assert client._telemetry_task is None
|
||||
await client._refresh_telemetry()
|
||||
assert client._telemetry_task is None
|
||||
assert not client._maintenance
|
||||
probe = client._telemetry_task
|
||||
try:
|
||||
client._draining = True
|
||||
client._maybe_finish_drain()
|
||||
assert client._reconnect_requested.is_set()
|
||||
await asyncio.wait_for(client._cancel_active_work(), 0.5)
|
||||
client._accepting_assignments = True
|
||||
await client._refresh_telemetry()
|
||||
assert client._telemetry_task is probe
|
||||
assert len(calls) == 1
|
||||
assert calls[0].daemon
|
||||
await asyncio.wait_for(client.stop(), 0.5)
|
||||
await client._refresh_telemetry()
|
||||
assert client._telemetry_task is probe
|
||||
assert not client._maintenance
|
||||
finally:
|
||||
release.set()
|
||||
await asyncio.wait_for(probe, 1)
|
||||
|
||||
Reference in New Issue
Block a user