chart/docs/async_refactor.md
Chris Amow ff1b9982d1 Seed in bulk and coalesce level rebuilds
P1 from docs/async_refactor.md. Measured on the dev stack: the port now accepts
connections 4 seconds after a restart rather than 121, and the worst loop lag
falls from 19,545ms to 526ms, with steady state between 0.2 and 0.6ms.

Seeding replayed years of history through on_bar, rebuilding every level from
scratch per bar and broadcasting each one to nobody. It now fills the store
quietly and derives price, ATR and the level set once at the end, from the
finished history. Alerts are deliberately not evaluated over replayed bars: a
level touched two years ago is not news, and firing on history is one way a
deploy re-alerts.

The seed was not all of it. Yahoo's first poll emits a whole day of minutes in a
single burst, each one taking the full live path, which was most of the
remaining twenty seconds. request_rebuild now coalesces to at most one rebuild
per 250ms and a background pass flushes anything deferred, so a burst costs a
handful of rebuilds instead of hundreds and the last bar is still never the one
dropped.

Verified unchanged after the change: bar counts across every timeframe, all five
daily moving averages with their full point sets, prior-day levels and VWAP. 125
python tests and 31 e2e tests pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-11 15:42:27 -05:00

9.5 KiB
Raw Blame History

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:

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.