diff --git a/lib/hub/hub_poll_wait_lib.py b/lib/hub/hub_poll_wait_lib.py new file mode 100644 index 0000000..fbd017f --- /dev/null +++ b/lib/hub/hub_poll_wait_lib.py @@ -0,0 +1,45 @@ +"""中控后台轮询等待:防止 request_refresh 连锁打满 CPU.""" +from __future__ import annotations + +import asyncio +import time + + +async def wait_poll_interval( + *, + refresh: asyncio.Event, + stop: asyncio.Event, + interval_sec: float, + started_at: float, + min_early_wake_sec: float | None = None, +) -> None: + """距 started_at 至少间隔 interval_sec 再进入下一轮. + + 期间若收到 refresh:仅当已过 min_early_wake_sec 才提前结束(兼顾手动刷新与防抖). + """ + interval = max(0.05, float(interval_sec)) + min_early = ( + float(min_early_wake_sec) + if min_early_wake_sec is not None + else min(2.0, interval * 0.4) + ) + while not stop.is_set(): + left = interval - (time.monotonic() - started_at) + if left <= 0: + return + refresh.clear() + stop_task = asyncio.create_task(stop.wait()) + refresh_task = asyncio.create_task(refresh.wait()) + done, pending = await asyncio.wait( + {stop_task, refresh_task}, + timeout=left, + return_when=asyncio.FIRST_COMPLETED, + ) + for t in pending: + t.cancel() + if stop.is_set(): + return + if not done: + return + if refresh_task in done and (time.monotonic() - started_at) >= min_early: + return diff --git a/manual_trading_hub/hub.py b/manual_trading_hub/hub.py index 8879b88..db7910c 100644 --- a/manual_trading_hub/hub.py +++ b/manual_trading_hub/hub.py @@ -336,8 +336,7 @@ async def _run_board_aggregate() -> dict: await asyncio.to_thread(record_fund_snapshot_from_board, body.get("rows") or []) except Exception: pass - # 监控聚合完成即唤醒数据看板,持仓来源与监控 5s 同步. - dashboard_store.request_refresh() + # 看板自有轮询即可;此处再 request_refresh 会与监控锁步,聚合变慢时几乎不睡眠打满 CPU. return {"ok": True, **body} except asyncio.TimeoutError: return { diff --git a/manual_trading_hub/hub_board_cache.py b/manual_trading_hub/hub_board_cache.py index 9593f79..0d1f2d1 100644 --- a/manual_trading_hub/hub_board_cache.py +++ b/manual_trading_hub/hub_board_cache.py @@ -5,9 +5,12 @@ from __future__ import annotations import asyncio import json import os +import time from collections.abc import AsyncIterator, Awaitable, Callable from typing import Any +from lib.hub.hub_poll_wait_lib import wait_poll_interval + HUB_BOARD_POLL_INTERVAL = float(os.getenv("HUB_BOARD_POLL_INTERVAL", "5")) HUB_BOARD_SSE_HEARTBEAT_SEC = float(os.getenv("HUB_BOARD_SSE_HEARTBEAT_SEC", "25")) @@ -79,18 +82,16 @@ class MonitorBoardStore: async def _loop(self) -> None: assert self._build_fn is not None while not self._stop.is_set(): + started = time.monotonic() await self._aggregate_once(self._build_fn) if self._stop.is_set(): break - self._refresh.clear() - sleep_task = asyncio.create_task(asyncio.sleep(HUB_BOARD_POLL_INTERVAL)) - refresh_task = asyncio.create_task(self._refresh.wait()) - done, pending = await asyncio.wait( - {sleep_task, refresh_task}, - return_when=asyncio.FIRST_COMPLETED, + await wait_poll_interval( + refresh=self._refresh, + stop=self._stop, + interval_sec=HUB_BOARD_POLL_INTERVAL, + started_at=started, ) - for t in pending: - t.cancel() async def _aggregate_once(self, build_fn: BuildFn) -> None: async with self._lock: diff --git a/manual_trading_hub/hub_chart_cache.py b/manual_trading_hub/hub_chart_cache.py index 14a5f0e..ebfca12 100644 --- a/manual_trading_hub/hub_chart_cache.py +++ b/manual_trading_hub/hub_chart_cache.py @@ -11,6 +11,7 @@ from dataclasses import dataclass from typing import Any from hub_board_cache import board_store +from lib.hub.hub_poll_wait_lib import wait_poll_interval HUB_CHART_POLL_INTERVAL = float(os.getenv("HUB_CHART_POLL_INTERVAL", "5")) HUB_CHART_SSE_HEARTBEAT_SEC = float(os.getenv("HUB_CHART_SSE_HEARTBEAT_SEC", "25")) @@ -161,18 +162,16 @@ class ChartPollStore: async def _loop(self) -> None: assert self._poll_fn is not None while not self._stop.is_set(): + started = time.monotonic() await self._poll_once(self._poll_fn) if self._stop.is_set(): break - self._refresh.clear() - sleep_task = asyncio.create_task(asyncio.sleep(HUB_CHART_POLL_INTERVAL)) - refresh_task = asyncio.create_task(self._refresh.wait()) - done, pending = await asyncio.wait( - {sleep_task, refresh_task}, - return_when=asyncio.FIRST_COMPLETED, + await wait_poll_interval( + refresh=self._refresh, + stop=self._stop, + interval_sec=HUB_CHART_POLL_INTERVAL, + started_at=started, ) - for t in pending: - t.cancel() async def _poll_once(self, poll_fn: PollFn) -> None: async with self._lock: diff --git a/manual_trading_hub/hub_dashboard_cache.py b/manual_trading_hub/hub_dashboard_cache.py index fafd212..a2be973 100644 --- a/manual_trading_hub/hub_dashboard_cache.py +++ b/manual_trading_hub/hub_dashboard_cache.py @@ -5,10 +5,12 @@ from __future__ import annotations import asyncio import json import os +import time from collections.abc import AsyncIterator, Awaitable, Callable from typing import Any from hub_dashboard import DASHBOARD_POLL_INTERVAL_SEC +from lib.hub.hub_poll_wait_lib import wait_poll_interval HUB_DASHBOARD_SSE_HEARTBEAT_SEC = float(os.getenv("HUB_DASHBOARD_SSE_HEARTBEAT_SEC", "25")) @@ -81,18 +83,16 @@ class DashboardStore: async def _loop(self) -> None: assert self._build_fn is not None while not self._stop.is_set(): + started = time.monotonic() await self._aggregate_once(self._build_fn) if self._stop.is_set(): break - self._refresh.clear() - sleep_task = asyncio.create_task(asyncio.sleep(DASHBOARD_POLL_INTERVAL_SEC)) - refresh_task = asyncio.create_task(self._refresh.wait()) - done, pending = await asyncio.wait( - {sleep_task, refresh_task}, - return_when=asyncio.FIRST_COMPLETED, + await wait_poll_interval( + refresh=self._refresh, + stop=self._stop, + interval_sec=DASHBOARD_POLL_INTERVAL_SEC, + started_at=started, ) - for t in pending: - t.cancel() async def _aggregate_once(self, build_fn: BuildFn) -> None: async with self._lock: diff --git a/manual_trading_hub/hub_supervisor_cache.py b/manual_trading_hub/hub_supervisor_cache.py index 02e8fa5..ae05259 100644 --- a/manual_trading_hub/hub_supervisor_cache.py +++ b/manual_trading_hub/hub_supervisor_cache.py @@ -4,9 +4,12 @@ from __future__ import annotations import asyncio import json import os +import time from collections.abc import AsyncIterator, Awaitable, Callable from typing import Any +from lib.hub.hub_poll_wait_lib import wait_poll_interval + SUPERVISOR_POLL_INTERVAL_SEC = float(os.getenv("SUPERVISOR_POLL_INTERVAL_SEC", "30")) SUPERVISOR_SSE_HEARTBEAT_SEC = float(os.getenv("SUPERVISOR_SSE_HEARTBEAT_SEC", "25")) @@ -65,18 +68,16 @@ class SupervisorStore: async def _loop(self) -> None: assert self._tick_fn is not None while not self._stop.is_set(): + started = time.monotonic() await self._tick_once(self._tick_fn) if self._stop.is_set(): break - self._refresh.clear() - sleep_task = asyncio.create_task(asyncio.sleep(SUPERVISOR_POLL_INTERVAL_SEC)) - refresh_task = asyncio.create_task(self._refresh.wait()) - done, pending = await asyncio.wait( - {sleep_task, refresh_task}, - return_when=asyncio.FIRST_COMPLETED, + await wait_poll_interval( + refresh=self._refresh, + stop=self._stop, + interval_sec=SUPERVISOR_POLL_INTERVAL_SEC, + started_at=started, ) - for t in pending: - t.cancel() async def _tick_once(self, tick_fn: TickFn) -> None: async with self._lock: diff --git a/manual_trading_hub/数据看板说明.md b/manual_trading_hub/数据看板说明.md index 5f977a3..f143d11 100644 --- a/manual_trading_hub/数据看板说明.md +++ b/manual_trading_hub/数据看板说明.md @@ -20,7 +20,7 @@ 1. `hub.py` 启动后 `dashboard_store` 每 **60s**(`DASHBOARD_POLL_INTERVAL_SEC`)聚合三户数据到内存快照. 2. 浏览器打开看板页后连接 `GET /api/dashboard/stream`(`event: dashboard`). 3. 收到新版本号后拉取 `GET /api/dashboard/daily` 快照并局部渲染,**无整页轮询闪烁**. -4. 监控区触发 board 刷新(全平,撤单等)时,会一并 `request_refresh` 看板,尽量与实盘同步. +4. 监控区触发 board 刷新(全平,撤单等)时,会一并 `request_refresh` 看板;常规轮询二者各自按间隔跑,避免连锁打满 CPU. 5. 「立即刷新」→ `POST /api/dashboard/refresh` 触发下一轮聚合. 可选环境变量:`HUB_DASHBOARD_SSE_HEARTBEAT_SEC`(默认 25,SSE 心跳间隔). diff --git a/tests/test_hub_poll_wait_lib.py b/tests/test_hub_poll_wait_lib.py new file mode 100644 index 0000000..5ec156a --- /dev/null +++ b/tests/test_hub_poll_wait_lib.py @@ -0,0 +1,51 @@ +import asyncio +import time +import unittest + +from lib.hub.hub_poll_wait_lib import wait_poll_interval + + +class TestWaitPollInterval(unittest.IsolatedAsyncioTestCase): + async def test_ignores_refresh_storm_until_interval(self): + refresh = asyncio.Event() + stop = asyncio.Event() + started = time.monotonic() + + async def storm(): + for _ in range(30): + refresh.set() + await asyncio.sleep(0.01) + + t = asyncio.create_task(storm()) + await wait_poll_interval( + refresh=refresh, + stop=stop, + interval_sec=0.25, + started_at=started, + min_early_wake_sec=0.2, + ) + t.cancel() + elapsed = time.monotonic() - started + self.assertGreaterEqual(elapsed, 0.18) + + async def test_stop_ends_early(self): + refresh = asyncio.Event() + stop = asyncio.Event() + started = time.monotonic() + + async def stopper(): + await asyncio.sleep(0.05) + stop.set() + + asyncio.create_task(stopper()) + await wait_poll_interval( + refresh=refresh, + stop=stop, + interval_sec=2.0, + started_at=started, + ) + self.assertLess(time.monotonic() - started, 0.5) + + +if __name__ == "__main__": + unittest.main()