Plan the async work, and record where it must not regress
An audit for blocking work on the event loop, written up rather than acted on — the chart changes in flight land first. The headline is a correctness bug, not a performance one. Sync route handlers run in FastAPI's threadpool and call rebuild_levels, which reaches asyncio.Queue.put_nowait on every subscriber. asyncio.Queue is not thread-safe: it wakes a consumer by resolving a Future, which has to happen on the loop thread. A dropped wakeup means a drawing made in one browser does not reach another until the next tick — invisible today only because the stream ticks about once a second and covers it. Below that: level rebuilding is CPU-bound on the loop and is the whole of the 82 second startup, and disarming an alert writes to disk from a coroutine. Also states what not to do, since the obvious reading of "make it async" is wrong here. Sync routes stay sync — FastAPI's threadpool is what keeps their work off the loop, and converting them would drag the rebuild cost onto it. ManualLineStore's threading lock stays, because both the loop and threadpool threads reach that store. Keeping it that way is three layers: a short async section in AGENTS.md, which is the only file both agents load every session; comments on the lines someone would actually edit, starting with the worker count in Procfile; and a loop-lag probe on /api/status so a stall reports itself as a number rather than as "the chart feels laggy". The risk register gains a row per finding pointing here. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
parent
617d7c75fb
commit
9fc6402f03
3 changed files with 202 additions and 1 deletions
4
Procfile
4
Procfile
|
|
@ -1 +1,5 @@
|
|||
# One worker, deliberately. The market stream lives in the app lifespan, so a
|
||||
# second process opens a second Schwab connection — they fight over the one
|
||||
# session the account allows, and every alert fires twice. Do not add
|
||||
# --workers, and do not swap in gunicorn with a worker count.
|
||||
web: uvicorn main:app --host 0.0.0.0 --port 8000
|
||||
|
|
|
|||
|
|
@ -1187,7 +1187,10 @@ enforces for data sources.
|
|||
| `session.py` bucket math wrong | Silently wrong lines everywhere | Tests written first; both DST transitions |
|
||||
| Repainting pivots | Lines that "were always there" | `w`-bar confirmation lag, enforced by test |
|
||||
| Alert fatigue | Product becomes unusable | Cluster-level alerts, cooldown + separation re-arm |
|
||||
| Multiple uvicorn workers | Duplicate Schwab connections | `workers=1`; streamer in `lifespan` |
|
||||
| Multiple uvicorn workers | Duplicate Schwab connections | `workers=1`; streamer in `lifespan`; warned in `Procfile` |
|
||||
| Threadpool routes touching `asyncio.Queue` | Dropped socket wakeups, rare corruption | Post via `call_soon_threadsafe` — `docs/async_refactor.md` P0 |
|
||||
| Level rebuilds run CPU-bound on the event loop | 82s startup; a stall every closed bar | Bulk seed, incremental MAs — `docs/async_refactor.md` P1 |
|
||||
| Alert disarm writes to disk on the loop | Stream stalls when an alert fires | Offload the write — `docs/async_refactor.md` P2 |
|
||||
| Schwab token expiry (7 days) | Stream dies | Surface prominently in status bar; document re-auth |
|
||||
| Vue reactivity wrapping chart objects | Perf collapse, odd bugs | `shallowRef`/`markRaw` — §9 |
|
||||
| LWC v4 tutorials copied | Code silently wrong for v5 | `addSeries(SeriesType, ...)` only |
|
||||
|
|
|
|||
194
docs/async_refactor.md
Normal file
194
docs/async_refactor.md
Normal file
|
|
@ -0,0 +1,194 @@
|
|||
# Async refactor — findings, priorities, and how to keep it that way
|
||||
|
||||
**Status: planned, not started.** 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)
|
||||
|
||||
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)
|
||||
|
||||
`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.** 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
|
||||
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.
|
||||
Loading…
Reference in a new issue