Stop open SSE streams from blocking graceful shutdown
Every connected dashboard holds a StreamingResponse open indefinitely, so uvicorn's graceful shutdown waited on them, systemd hit its 90s stop timeout and SIGKILLed the unit. That skipped the in-process restore hook entirely, leaving the GPU restore to ExecStopPost alone. - broker.stop() sets a closing flag and pushes a sentinel to every subscriber queue so the SSE generators return instead of parking on q.get(). - uvicorn gets timeout_graceful_shutdown=10 and the unit TimeoutStopSec=20, bounding the worst case rather than relying on the 90s default. Restart now completes in ~11s with 'Restoring GPU to safe stock state (server shutdown)' running in-process, and no SIGKILL. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
16
server.py
16
server.py
@@ -55,6 +55,8 @@ class TelemetryBroker:
|
|||||||
self.running = False
|
self.running = False
|
||||||
self.samples = 0
|
self.samples = 0
|
||||||
self.last_sample_ms = 0.0
|
self.last_sample_ms = 0.0
|
||||||
|
# Set on shutdown so open SSE generators finish instead of holding the server up.
|
||||||
|
self.closing = False
|
||||||
|
|
||||||
async def start(self) -> None:
|
async def start(self) -> None:
|
||||||
if self.running:
|
if self.running:
|
||||||
@@ -64,6 +66,13 @@ class TelemetryBroker:
|
|||||||
|
|
||||||
async def stop(self) -> None:
|
async def stop(self) -> None:
|
||||||
self.running = False
|
self.running = False
|
||||||
|
self.closing = True
|
||||||
|
# Wake every subscriber so their generator can return. Without this, uvicorn waits
|
||||||
|
# on the open SSE responses during graceful shutdown and systemd eventually
|
||||||
|
# SIGKILLs the unit -- which skips the in-process GPU restore hook entirely.
|
||||||
|
for q in list(self.subscribers):
|
||||||
|
with contextlib.suppress(asyncio.QueueFull):
|
||||||
|
q.put_nowait(None)
|
||||||
if self.task:
|
if self.task:
|
||||||
self.task.cancel()
|
self.task.cancel()
|
||||||
with contextlib.suppress(asyncio.CancelledError):
|
with contextlib.suppress(asyncio.CancelledError):
|
||||||
@@ -250,11 +259,13 @@ async def sse_telemetry_stream(request: Request):
|
|||||||
try:
|
try:
|
||||||
snap = await broker.get()
|
snap = await broker.get()
|
||||||
yield f"data: {json.dumps(snap)}\n\n"
|
yield f"data: {json.dumps(snap)}\n\n"
|
||||||
while True:
|
while not broker.closing:
|
||||||
if await request.is_disconnected():
|
if await request.is_disconnected():
|
||||||
break
|
break
|
||||||
try:
|
try:
|
||||||
snap = await asyncio.wait_for(q.get(), timeout=15.0)
|
snap = await asyncio.wait_for(q.get(), timeout=15.0)
|
||||||
|
if snap is None: # shutdown sentinel
|
||||||
|
break
|
||||||
yield f"data: {json.dumps(snap)}\n\n"
|
yield f"data: {json.dumps(snap)}\n\n"
|
||||||
except asyncio.TimeoutError:
|
except asyncio.TimeoutError:
|
||||||
yield ": keepalive\n\n"
|
yield ": keepalive\n\n"
|
||||||
@@ -497,4 +508,5 @@ async def root_index():
|
|||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
import uvicorn
|
import uvicorn
|
||||||
uvicorn.run("server:app", host="0.0.0.0", port=9090, reload=False, log_level="info")
|
uvicorn.run("server:app", host="0.0.0.0", port=9090, reload=False, log_level="info",
|
||||||
|
timeout_graceful_shutdown=10)
|
||||||
|
|||||||
Reference in New Issue
Block a user