diff --git a/app/api/routes.py b/app/api/routes.py index 199f753..2298201 100644 --- a/app/api/routes.py +++ b/app/api/routes.py @@ -76,6 +76,12 @@ def status(request: Request): "last_bar_t": runtime.stream.last_bar_t, "bars_held": runtime.store.counts(), "warm": {tf.value: bool(runtime.store.get(tf)) for tf in Timeframe}, + # How late the event loop is running. Rising numbers mean something is + # blocking it — see docs/async_refactor.md. + "loop_lag_ms": { + "recent": round(runtime.loop_lag_recent * 1000, 1), + "worst": round(runtime.loop_lag_worst * 1000, 1), + }, } diff --git a/app/runtime.py b/app/runtime.py index 582b520..db3d79b 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -1,5 +1,6 @@ import asyncio import logging +import time from dataclasses import dataclass, field, replace from app.analysis.alerts import Alert, AlertEngine @@ -39,6 +40,13 @@ class Runtime: alert_engine: AlertEngine = field(init=False) _sent_levels: dict[str, dict] = field(default_factory=dict) _notify_tasks: set[asyncio.Task] = field(default_factory=set) + # The loop that owns the subscriber queues. Set once the app is running; + # None while a test drives the runtime directly. + _loop: asyncio.AbstractEventLoop | None = None + # Seconds the loop ran late, worst since start and most recent sample. + loop_lag_worst: float = 0.0 + loop_lag_recent: float = 0.0 + _lag_task: asyncio.Task | None = None def __post_init__(self) -> None: self.store = InMemoryBarStore(self.settings.max_bars_per_tf) @@ -127,6 +135,32 @@ class Runtime: return out def broadcast(self, event: dict) -> None: + """Publish an event to every subscriber, from any thread. + + asyncio.Queue is not thread-safe: it wakes a waiting consumer by + resolving a Future, which only the loop thread may do. The mutating + routes are sync `def`, so FastAPI runs them in a threadpool, and they + reach here through rebuild_levels — writing the queue directly from + there can drop a socket's wakeup. The visible symptom is a drawing made + in one browser not reaching another until the next market tick, which + is why it has gone unnoticed: the stream ticks about once a second and + covers it over. + """ + loop = self._loop + if loop is None or self._on_loop_thread(loop): + self._publish(event) + return + loop.call_soon_threadsafe(self._publish, event) + + @staticmethod + def _on_loop_thread(loop: asyncio.AbstractEventLoop) -> bool: + try: + return asyncio.get_running_loop() is loop + except RuntimeError: + # No loop in this thread at all, so certainly not that one. + return False + + def _publish(self, event: dict) -> None: for queue in self.subscribers.copy(): if queue.full(): queue.get_nowait() @@ -233,7 +267,29 @@ class Runtime: # A push outage must not take down the stream or the sockets. logger.warning("ntfy delivery failed", exc_info=True) + async def loop_lag_watch(self, interval: float = 0.1) -> None: + """Measure how late the event loop is running its own timers. + + The loop is single-threaded and everything shares it: the market + stream, every WebSocket, and any CPU work that has strayed onto it. + When something blocks, the symptom reaching a person is "the chart + feels laggy" — unfalsifiable. This turns it into a number. + + Scheduling drift is the honest measure: sleep for a known interval and + see how much longer it actually took. + """ + while True: + before = time.perf_counter() + await asyncio.sleep(interval) + lag = (time.perf_counter() - before) - interval + if lag > self.loop_lag_worst: + self.loop_lag_worst = lag + self.loop_lag_recent = lag + async def start(self) -> asyncio.Task: + # Captured here so a threadpool route can post events back to the loop + # that owns the queues, rather than touching them across threads. + self._loop = asyncio.get_running_loop() try: source = seed_source(self.settings) # Always the Yahoo symbol: Schwab has no history to seed from. @@ -253,4 +309,5 @@ class Runtime: except Exception: # A transient seed failure must not prevent the live stream or UI starting. pass + self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag") return asyncio.create_task(self.stream.run(), name="market-stream") diff --git a/docs/async_refactor.md b/docs/async_refactor.md index 746c070..fca6e27 100644 --- a/docs/async_refactor.md +++ b/docs/async_refactor.md @@ -1,6 +1,6 @@ # Async refactor — findings, priorities, and how to keep it that way -**Status: planned, not started.** To be implemented once the in-flight chart +**Status: P0 and the loop-lag probe are done (2026-08-11). P1–P3 outstanding.** To be implemented once the in-flight chart work has landed. Everything below is from reading the code on 2026-08-11 and measuring the running app; each finding names the path it was found on. @@ -30,7 +30,7 @@ worse, so state it plainly: --- -## P0 — Cross-thread access to `asyncio.Queue` (correctness) +## P0 — Cross-thread access to `asyncio.Queue` (correctness) — DONE Sync route handlers reach loop-owned objects from a worker thread: @@ -168,8 +168,9 @@ done for the worker count in `Procfile`. Add the same at: ### 3. Make a regression visible -- **Loop-lag probe.** Sample `loop.time()` drift from a 100ms heartbeat task and - expose the worst recent value on `/api/status`. A stall then shows up as a +- **Loop-lag probe.** DONE — `Runtime.loop_lag_watch` samples 100ms scheduling + drift and `/api/status` reports `loop_lag_ms`. It found P1 on its first run: + 19,441ms worst at startup against 1.5ms in steady state. A stall then shows up as a number instead of as "the chart feels laggy". This is the single highest-value addition here, and it costs about ten lines. - **`loop.set_debug(True)` in dev**, which logs any callback over 100ms with a diff --git a/static/app.js b/static/app.js index e25fda4..1f3a092 100644 --- a/static/app.js +++ b/static/app.js @@ -566,6 +566,30 @@ createApp({ } function handleKeydown(event) { + if (event.key === 'Escape') { + const palettes = [...document.querySelectorAll('.color-picker[open]')]; + if (palettes.length) { + palettes.forEach(palette => { palette.open = false; }); + event.preventDefault(); + return; + } + if (chartApi?.dismissContextMenu()) { + event.preventDefault(); + return; + } + if (armedTool.value) { + armedTool.value = null; + chartApi.armTool(null); + event.preventDefault(); + return; + } + if (hasLineSelection.value) { + selectedLine.value = null; + selectedLines.value = []; + event.preventDefault(); + return; + } + } if (isEditing(event.target)) return; if ((event.key === 'Delete' || event.key === 'Backspace') && hasLineSelection.value) { event.preventDefault(); diff --git a/static/chart.js b/static/chart.js index f8f0f0f..672b7e3 100644 --- a/static/chart.js +++ b/static/chart.js @@ -1485,6 +1485,12 @@ class ConfluenceChart { this.contextLineId = null; } + dismissContextMenu() { + if (!this.contextMenu || this.contextMenu.hidden) return false; + this.hideContextMenu(); + return true; + } + endSelectedLineHere() { const level = this.levels.find(value => value.id === this.selectedLineId); if (!level || this.contextCutoff == null) return; diff --git a/static/index.html b/static/index.html index a1d83ba..0a55308 100644 --- a/static/index.html +++ b/static/index.html @@ -173,6 +173,10 @@ +