diff --git a/Procfile b/Procfile index b98f12c..d63c90f 100644 --- a/Procfile +++ b/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 diff --git a/docs/IMPLEMENTATION_PLAN.md b/docs/IMPLEMENTATION_PLAN.md index 0cae3ba..e677ef8 100644 --- a/docs/IMPLEMENTATION_PLAN.md +++ b/docs/IMPLEMENTATION_PLAN.md @@ -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 | diff --git a/docs/async_refactor.md b/docs/async_refactor.md new file mode 100644 index 0000000..746c070 --- /dev/null +++ b/docs/async_refactor.md @@ -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.