From cc25871032f33777d45d1d6f5c152f9986ed6c74 Mon Sep 17 00:00:00 2001
From: Chris Amow
The authorisation code could not be exchanged. Start again from the chart.
+Nothing was stored.
+""" + + +@router.get("/qt") +def callback(request: Request): + runtime = getattr(request.app.state, "runtime", None) + if runtime is not None and runtime.schwab_login is not None and request.query_params.get("code"): + try: + runtime.finish_schwab_login(str(request.url)) + except Exception: + page = PAGE.format(heading="Authorisation failed", body=FAILED) + return HTMLResponse(page, headers={"Cache-Control": "no-store"}, status_code=400) + return RedirectResponse("/", headers={"Cache-Control": "no-store"}) if request.query_params.get("code"): page = PAGE.format( heading="Authorisation received", diff --git a/app/market/schwab.py b/app/market/schwab.py index 76db41d..a14a6e8 100644 --- a/app/market/schwab.py +++ b/app/market/schwab.py @@ -120,6 +120,11 @@ class SchwabSource: name = "schwab" delay_minutes = 0 + # Access tokens last 30 minutes. The refresh token lasts seven days unless + # something uses it — the live socket never does, so a quiet week killed + # the feed. Hitting REST on this cadence writes a new refresh token to disk. + TOKEN_KEEPALIVE_SECONDS = 6 * 3600 + def __init__(self, settings, stream_client_factory=None): self._settings = settings # Injectable so the parsing and dispatch can be tested without a socket. @@ -128,6 +133,7 @@ class SchwabSource: # exactly as it was: one closed bar a minute. seconds = getattr(settings, "schwab_tick_seconds", 1.0) self._tick_seconds = None if seconds is None or seconds < 0 else seconds + self._http = None def supports_history(self) -> bool: return False @@ -138,23 +144,44 @@ class SchwabSource: def supports_stream(self) -> bool: return True + def http_client(self): + if self._http is None: + from schwab.auth import client_from_token_file + + settings = self._settings + if not settings.schwab_token_path.exists(): + raise RuntimeError( + f"No Schwab token at {settings.schwab_token_path}" + ) + self._http = client_from_token_file( + str(settings.schwab_token_path), + settings.schwab_api_key, + settings.schwab_app_secret, + asyncio=True, + ) + return self._http + + def reset_client(self) -> None: + self._http = None + def _build_stream_client(self): - from schwab.auth import client_from_token_file from schwab.streaming import StreamClient - settings = self._settings - if not settings.schwab_token_path.exists(): - raise RuntimeError( - f"No Schwab token at {settings.schwab_token_path}. " - "Run: python3 -m scripts.check_schwab" - ) - client = client_from_token_file( - str(settings.schwab_token_path), - settings.schwab_api_key, - settings.schwab_app_secret, - asyncio=True, - ) - return StreamClient(client) + return StreamClient(self.http_client()) + + async def refresh_token(self) -> None: + await self.http_client().get_user_preferences() + + async def keep_alive(self, interval: float | None = None) -> None: + wait = self.TOKEN_KEEPALIVE_SECONDS if interval is None else interval + while True: + await asyncio.sleep(wait) + try: + await self.refresh_token() + except asyncio.CancelledError: + raise + except Exception: + logger.exception("Schwab token keepalive failed") async def stream(self, symbol: str) -> AsyncIterator[Bar]: stream_client = self._stream_client_factory() diff --git a/app/market/stream.py b/app/market/stream.py index fd4ec26..8a1fadf 100644 --- a/app/market/stream.py +++ b/app/market/stream.py @@ -14,6 +14,7 @@ class StreamService: self.source = source self.symbol = symbol self.status = "disconnected" + self.last_error: str | None = None self.last_bar_t: int | None = None self._handlers: list[BarHandler] = [] self._stop = asyncio.Event() @@ -44,8 +45,9 @@ class StreamService: async def run(self) -> None: while not self._stop.is_set(): try: - self.status = "replay" if self.source.name == "replay" else "connected" async for bar in self.source.stream(self.symbol): + self.status = "replay" if self.source.name == "replay" else "connected" + self.last_error = None await self._emit(bar) if self._stop.is_set(): break @@ -53,7 +55,8 @@ class StreamService: return except asyncio.CancelledError: raise - except Exception: + except Exception as exc: + self.last_error = str(exc) logger.exception("Market stream failed; reconnecting") self.status = "disconnected" try: diff --git a/app/runtime.py b/app/runtime.py index a155629..8de953a 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -57,6 +57,8 @@ class Runtime: loop_lag_recent: float = 0.0 _lag_task: asyncio.Task | None = None _rebuild_task: asyncio.Task | None = None + _token_task: asyncio.Task | None = None + schwab_login: object | None = None def __post_init__(self) -> None: self.store = InMemoryBarStore(self.settings.max_bars_per_tf) @@ -409,4 +411,30 @@ class Runtime: self.settle_after_seed() self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag") self._rebuild_task = asyncio.create_task(self.rebuild_watch(), name="level-rebuild") + keep_alive = getattr(self.stream.source, "keep_alive", None) + if keep_alive is not None: + self._token_task = asyncio.create_task(keep_alive(), name="token-keepalive") return asyncio.create_task(self.stream.run(), name="market-stream") + + def needs_login(self) -> bool: + error = self.stream.last_error or "" + return "invalid_grant" in error or "Refresh token" in error + + def start_schwab_login(self) -> str: + from app.api.schwab_auth import begin_login + + context = begin_login(self.settings) + self.schwab_login = context + return context.authorization_url + + def finish_schwab_login(self, redirect_url: str) -> None: + from app.api.schwab_auth import complete_login + + context = self.schwab_login + if context is None: + raise RuntimeError("No login in progress") + complete_login(self.settings, context, redirect_url) + self.schwab_login = None + reset = getattr(self.stream.source, "reset_client", None) + if reset is not None: + reset() diff --git a/docs/implementation.md b/docs/implementation.md index 3bddecc..23c6652 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -1161,3 +1161,17 @@ chart — trendlines, price levels, and comments — and (unless "Hidden levels still count toward confluence" is on) keeps those lines out of clusters too. Manual lines stays as the finer control for lines only. Old stored prefs without the key keep drawings on. + +### 2026-08-18 — Schwab token keepalive + +The header said DISCONNECTED · SCHWAB on every timeframe. Production logs +were `invalid_grant` / refresh token expired. The token file was created +2026-08-10 and last written 2026-08-17; the live socket never makes a REST +call, so the seven-day refresh token aged out while the chart still looked +fine. A deploy then tried to log in and could not. + +The stream now shares one HTTP client and pings user preferences every six +hours so a new refresh token is written to the volume. Status grows +`needs_login` when the error is a dead grant; the header then offers +**reconnect**, which starts the existing `/api/qt` flow and writes the token +on callback instead of asking anyone to paste a URL into a prompt. diff --git a/docs/plan.md b/docs/plan.md index 3cc0780..b19f7e7 100644 --- a/docs/plan.md +++ b/docs/plan.md @@ -1231,7 +1231,7 @@ enforces for data sources. | 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 | +| Schwab token expiry (7 days) | Stream dies | Keepalive REST ping writes a new refresh token; header reconnect when the grant is dead | | Vue reactivity wrapping chart objects | Perf collapse, odd bugs | `shallowRef`/`markRaw` — §9 | | LWC v4 tutorials copied | Code silently wrong for v5 | `addSeries(SeriesType, ...)` only | | Viewport derived from `bars.length` | Chart looks frozen; feed is fine | Anchor the view by time, never by logical index — §9 | diff --git a/static/index.html b/static/index.html index 2958e42..f0f708c 100644 --- a/static/index.html +++ b/static/index.html @@ -16,7 +16,7 @@