Skip to content

Commit 7016167

Browse files
committed
feat: worker hot-swap with health check, staggered startup and live dashboard
- Staggered worker startup with 3s intervals to prevent resource exhaustion - Health check loop (30s interval, 3 strikes) with auto-restart for crashed workers - Dynamic add/remove/start/stop individual workers at runtime - WebSocket push for worker_snapshot and worker_status events - Dashboard real-time worker status cards (running/stopped/crashed/rate_limited) - Graceful stop waits for active requests before killing worker - Startup loop respects stop signal for immediate cancellation
1 parent 2f26d2b commit 7016167

11 files changed

Lines changed: 838 additions & 183 deletions

File tree

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
🎤 Gemini 2.5 TTS 语音生成
2222
</p>
2323

24-
<img src="docs/img/demo.gif" alt="Demo GIF" width="100%" />
24+
<!-- <img src="docs/img/demo.gif" alt="Demo GIF" width="100%" /> -->
2525

2626
<!-- <p align="center">
2727
<img src="docs/img/多worker并发和媒体模型支援.png" alt="多Worker并发与媒体模型支援" width="80%" />

README_en.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
🎤 Gemini 2.5 TTS Speech Synthesis
2222
</p>
2323

24-
<img src="docs/img/demo.gif" alt="Demo GIF" width="100%" />
24+
<!-- <img src="docs/img/demo.gif" alt="Demo GIF" width="100%" /> -->
2525

2626
<!-- <p align="center">
2727
<img src="docs/img/多worker并发和媒体模型支援.png" alt="Multi-Worker Concurrency & Media Model Support" width="80%" />

src/manager/app.py

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,11 +31,24 @@ async def lifespan(app: FastAPI):
3131
manager.loop = asyncio.get_running_loop()
3232
config = manager.load_config()
3333
manager._log_enabled = config.get("log_enabled", True)
34+
health_task = None
3435
if WORKER_POOL_AVAILABLE and worker_pool is not None:
3536
worker_pool.init_from_config()
37+
worker_pool.configure_runtime(config)
38+
worker_pool.register_status_listener(manager.handle_worker_status_event)
39+
worker_pool.register_process_listener(manager.handle_worker_process_started)
40+
health_task = asyncio.create_task(worker_pool.health_check_loop())
3641
yield
42+
if health_task is not None:
43+
health_task.cancel()
44+
try:
45+
await health_task
46+
except asyncio.CancelledError:
47+
pass
3748
if manager.process or manager.worker_processes:
3849
manager.stop_service()
50+
if WORKER_POOL_AVAILABLE and worker_pool is not None:
51+
await worker_pool.close()
3952
manager.loop = None
4053

4154

src/manager/routes/control.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,9 @@
1010

1111
@router.post("/start")
1212
async def start_service(config: Dict[str, Any] = Body(...)):
13-
success, message = manager.start_service(config)
13+
success, message = await manager.start_service(config)
1414
await manager.broadcast_status()
15+
await manager.broadcast_worker_snapshot()
1516
if not success:
1617
raise HTTPException(status_code=400, detail=message)
1718
return {"success": True, "message": message}
@@ -21,6 +22,7 @@ async def start_service(config: Dict[str, Any] = Body(...)):
2122
async def stop_service():
2223
success, message = manager.stop_service()
2324
await manager.broadcast_status()
25+
await manager.broadcast_worker_snapshot()
2426
if not success:
2527
raise HTTPException(status_code=500, detail=message)
2628
return {"success": True, "message": message}

src/manager/routes/websocket.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ async def websocket_logs(websocket: WebSocket):
2222
}
2323
)
2424
)
25+
await manager.broadcast_worker_snapshot()
2526
while True:
2627
await websocket.receive_text()
2728
except WebSocketDisconnect:

src/manager/routes/workers.py

Lines changed: 31 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
except ImportError:
88
from config.settings import SAVED_AUTH_DIR
99

10-
from ..service import WORKER_POOL_AVAILABLE, worker_pool
10+
from ..service import WORKER_POOL_AVAILABLE, manager, worker_pool
1111

1212

1313
router = APIRouter(prefix="/api/workers", tags=["Workers"])
@@ -54,6 +54,12 @@ async def add_worker(profile: str = Body(..., embed=True)):
5454
)
5555
pool.workers[worker_id] = worker
5656
pool.save_config()
57+
if manager.is_worker_mode and manager.service_status == "running":
58+
pool.configure_runtime(manager.load_config())
59+
success, message = pool.start_worker(worker_id)
60+
if not success:
61+
raise HTTPException(status_code=500, detail=message)
62+
await manager.broadcast_worker_snapshot()
5763
return {"success": True, "worker": worker.to_dict()}
5864

5965

@@ -65,11 +71,16 @@ async def remove_worker(worker_id: str):
6571
raise HTTPException(status_code=404, detail="Worker not found")
6672

6773
worker = pool.workers[worker_id]
74+
process = worker.process
6875
if worker.status == "running":
69-
pool.stop_worker(worker_id)
76+
success, message = await pool.stop_worker(worker_id)
77+
if not success:
78+
raise HTTPException(status_code=500, detail=message)
79+
manager.unregister_worker_process(process)
7080

7181
del pool.workers[worker_id]
7282
pool.save_config()
83+
await manager.broadcast_worker_snapshot()
7384
return {"success": True}
7485

7586

@@ -86,6 +97,7 @@ async def list_workers():
8697
async def init_workers():
8798
pool = _require_worker_pool()
8899
pool.init_from_config()
100+
await manager.broadcast_worker_snapshot()
89101
return {"success": True, "count": len(pool.workers)}
90102

91103

@@ -94,6 +106,7 @@ async def save_workers_config():
94106
pool = _require_worker_pool()
95107
try:
96108
pool.save_config()
109+
await manager.broadcast_worker_snapshot()
97110
return {"success": True, "count": len(pool.workers)}
98111
except Exception as exc:
99112
return {"success": False, "error": str(exc)}
@@ -124,44 +137,57 @@ async def get_next_available_worker(model: str = ""):
124137
async def mark_worker_rate_limited(worker_id: str, model: str = Body(..., embed=True)):
125138
pool = _require_worker_pool()
126139
pool.mark_rate_limited(worker_id, model)
140+
await manager.broadcast_worker_snapshot()
127141
return {"success": True}
128142

129143

130144
@router.post("/{worker_id}/start")
131145
async def start_worker_api(worker_id: str):
132146
pool = _require_worker_pool()
147+
pool.configure_runtime(manager.load_config())
133148
success, message = pool.start_worker(worker_id)
134149
if not success:
135150
raise HTTPException(status_code=400, detail=message)
151+
await manager.broadcast_worker_snapshot()
136152
return {"success": True, "message": message}
137153

138154

139155
@router.post("/{worker_id}/stop")
140156
async def stop_worker_api(worker_id: str):
141157
pool = _require_worker_pool()
142-
success, message = pool.stop_worker(worker_id)
158+
process = pool.workers.get(worker_id).process if worker_id in pool.workers else None
159+
success, message = await pool.stop_worker(worker_id)
143160
if not success:
144161
raise HTTPException(status_code=400, detail=message)
162+
manager.unregister_worker_process(process)
163+
await manager.broadcast_worker_snapshot()
145164
return {"success": True, "message": message}
146165

147166

148167
@router.post("/{worker_id}/clear-limits")
149168
async def clear_worker_limits(worker_id: str):
150169
pool = _require_worker_pool()
151170
if pool.clear_rate_limits(worker_id):
171+
await manager.broadcast_worker_snapshot()
152172
return {"success": True}
153173
raise HTTPException(status_code=404, detail="Worker not found")
154174

155175

156176
@router.post("/start-all")
157177
async def start_all_workers():
158178
pool = _require_worker_pool()
159-
pool.start_all()
179+
pool.configure_runtime(manager.load_config())
180+
await pool.start_all()
181+
await manager.broadcast_worker_snapshot()
160182
return {"success": True}
161183

162184

163185
@router.post("/stop-all")
164186
async def stop_all_workers():
165187
pool = _require_worker_pool()
166-
pool.stop_all()
188+
processes = [worker.process for worker in pool.workers.values() if worker.process]
189+
await pool.stop_all()
190+
for process in processes:
191+
manager.unregister_worker_process(process)
192+
await manager.broadcast_worker_snapshot()
167193
return {"success": True}

0 commit comments

Comments
 (0)