# Async refactor — findings, priorities, and how to keep it that way **Status: P0, P1 and the loop-lag probe are done (2026-08-11). P2 and 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. ## The goal is not "more async" The target is **nothing blocks the event loop**, which is not the same thing as converting everything to `async def`. Getting this backwards would make the app worse, so state it plainly: - FastAPI runs a **sync `def` route in a threadpool**. Blocking work inside one never touches the loop. That is protection, not a defect. - Converting those routes to `async def` *removes* the protection: any blocking call inside then stalls the market stream and every WebSocket. - So the rule is per-function. A handler that only awaits should be `async def`. A handler doing blocking or CPU work should stay `def` — **and must not touch loop-owned objects** (see P0). ## What is already right - Every outbound HTTP call is async: `httpx.AsyncClient` in `notify/ntfy.py` and `market/yahoo.py`; the Schwab `StreamClient` is built with `asyncio=True`. - The market stream is an `asyncio.Task` owned by the app lifespan, and the WebSocket endpoint is a coroutine. - `ManualLineStore` guards itself with a `threading.RLock`, which is the correct primitive precisely because both threadpool routes and the loop reach it. Do not "modernise" it to `asyncio.Lock` — that would protect only one of them. --- ## P0 — Cross-thread access to `asyncio.Queue` (correctness) — DONE Sync route handlers reach loop-owned objects from a worker thread: ``` create_line / create_price_alert / create_comment / patch_line / delete_line (sync def, threadpool) -> runtime.rebuild_levels() -> broadcast_level_delta() -> broadcast() -> queue.put_nowait(event) # asyncio.Queue, owned by the loop ``` `asyncio.Queue` is not thread-safe. It wakes a waiting consumer by setting a Future's result, and Futures must be resolved on the loop thread — from another thread that requires `loop.call_soon_threadsafe`. Writing directly can drop the wakeup or corrupt internal state. **Why nobody has noticed:** the market stream broadcasts roughly once a second, so a dropped wakeup is papered over by the next event almost immediately. The visible symptom would be a drawing made in one browser not appearing in another until the next tick — easy to misread as network lag. **Fix.** Give `Runtime` the loop it belongs to and post from the correct thread: ```python self._loop = asyncio.get_running_loop() # captured in start() def broadcast(self, event: dict) -> None: if threading.current_thread() is threading.main_thread() and self._loop.is_running(): self._publish(event) # already on the loop else: self._loop.call_soon_threadsafe(self._publish, event) ``` Prefer this over making the routes `async def`: that would move level rebuilding (P1) onto the loop, trading a rare correctness bug for a guaranteed latency one. **Verify.** A test that calls a mutating route through `TestClient` while a WebSocket subscriber waits, asserting the event arrives without another tick intervening. Today that passes by luck. --- ## P1 — Level rebuilding is CPU-bound on the loop (latency) — DONE Measured on the dev stack, Yahoo source: | | before | after | |---|---|---| | port accepting connections | 121s | 4s | | worst loop lag | 19,545ms | 526ms | | steady-state lag | — | 0.2–0.6ms | Two changes did it. Seeding now accumulates into the store and derives price, ATR and the level set **once** at the end (`settle_after_seed`) instead of rebuilding per replayed bar. And `request_rebuild` coalesces rebuilds to at most one per 250ms with a guaranteed trailing pass, which matters because Yahoo's first poll emits a whole day of minutes in one burst — that burst, not the seed, was most of the remaining 20 seconds. Bar counts, all five daily MAs, prior-day levels and VWAP are unchanged; 125 python tests and 31 e2e tests pass. Still open from the original list: incremental moving averages, and diffing levels by fingerprint rather than by re-serialising every point. Neither is needed while lag sits under a second. ### Original analysis `rebuild_levels()` recomputes all five daily moving averages and re-serialises their points to diff them, on every closed bar. Measured consequence: the seed replay runs the same path per bar and takes **~82 seconds**, during which the port is closed. In steady state it is once a minute, which is survivable but is the largest single thing the loop does. Threads do not help — it is genuine CPU under the GIL. The fix is algorithmic: 1. **Incremental moving averages.** `indicators.sma` is already a rolling sum; the waste is recomputing every window from scratch each rebuild rather than advancing the last one. 2. **Bulk seeding.** Load seeded bars into the store directly and rebuild levels **once** at the end, rather than replaying each bar through `on_bar`. 3. **Diff without re-serialising.** `broadcast_level_delta()` compares `to_dict()` output including hundreds of points per MA. Compare a cheap fingerprint (last point plus length) and serialise only what changed. Do (2) first — it is contained, testable, and removes most of the 82 seconds. **Verify.** A test asserting a seed of N bars completes under a threshold, and a loop-lag probe (below) staying under ~50ms while a bar closes. --- ## P2 — Blocking disk write on the loop (small, real) ``` on_bar (coroutine) -> rebuild_clusters(evaluate_alerts=True) -> dispatch_alerts() -> disarm() -> manual_lines.update() -> save() # write_text + atomic replace ``` A ~1KB write, usually sub-millisecond, but it lands on the loop at the exact moment an alert fires, and it is unbounded on a contended disk. **Fix.** Either `await asyncio.to_thread(self.manual_lines.update, ...)` on that path, or mark the line disarmed in memory and flush outside the tick. The same applies to any future persistence work — see the cold-restart notes, which will add far more writing than this. --- ## P3 — Seeding blocks startup (architectural) `Runtime.start()` awaits both seeds before uvicorn binds, so the port refuses connections for the whole ~82 seconds and any open browser logs a wall of `ERR_CONNECTION_REFUSED`. Fixing P1(2) may reduce this enough on its own. If it does not, seed in a background task and serve immediately — but note that changes what `/api/status`'s `warm` flags mean to every consumer, so it needs its own thought rather than being bolted on. --- ## Explicitly not doing - **Converting sync routes to `async def`.** They do in-memory work behind a threadpool hop, which is correct and cheap. Changing them adds risk for no gain, and would drag P1's CPU cost onto the loop. - **`asyncio.Lock` in `ManualLineStore`.** It is reached from both the loop and threadpool threads; only a threading primitive covers both. - **Multiple uvicorn workers as a performance fix.** See `Procfile` — a second process opens a second Schwab stream and duplicates every alert. --- ## Keeping it this way Findings decay unless something enforces them. Three layers, cheapest first. ### 1. Rules where agents actually read them `AGENTS.md` is loaded automatically by both Claude Code (via the `CLAUDE.md` symlink) and OpenCode every session; nothing else in the repo is guaranteed to be read. Add a short **Async rules** section stating: - Nothing blocking or CPU-heavy runs on the event loop. - Sync `def` routes stay sync; they run in a threadpool by design. - A threadpool thread must never touch `asyncio` objects directly — post through `loop.call_soon_threadsafe`. - `ManualLineStore`'s lock is a threading lock deliberately. Keep it to a handful of lines. A long section is skimmed; a short one is read. ### 2. Comments at the point of danger A rule in a document does not stop an edit; a comment on the line does. Already done for the worker count in `Procfile`. Add the same at: - `Runtime.broadcast` — why it posts through the loop. - `ManualLineStore._lock` — why it is a threading lock. - Each mutating route — why it is `def` and not `async def`. ### 3. Make a regression visible - **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 traceback — it would have named P1 immediately. - **A seed-duration test**, so P1 cannot silently regress once fixed. ### 4. Where the record lives The risk register in `IMPLEMENTATION_PLAN.md` §15 gets one row per finding, so a reader looking for known hazards finds them. This file holds the detail; the register holds the pointer. --- ## Suggested order 1. **P0** — correctness, small, self-contained. 2. **Loop-lag probe** — so P1's improvement is measurable rather than asserted. 3. **P1(2)** bulk seeding, then P1(1) and P1(3) if the probe still shows stalls. 4. **P2** — trivial once P0 has established how work leaves the loop. 5. **P3** — only if P1 leaves the startup window unacceptable. 6. Docs and comments alongside each change, not as a final sweep.