diff --git a/app/api/routes.py b/app/api/routes.py index 6822a40..875727f 100644 --- a/app/api/routes.py +++ b/app/api/routes.py @@ -9,6 +9,7 @@ from datetime import date as Date from typing import Literal from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response +from fastapi.responses import RedirectResponse from pydantic import BaseModel, Field from app.bars.models import Timeframe @@ -139,9 +140,18 @@ def status(request: Request): "recent": round(runtime.loop_lag_recent * 1000, 1), "worst": round(runtime.loop_lag_worst * 1000, 1), }, + "needs_login": runtime.needs_login(), } +@router.get("/schwab/login") +def schwab_login(request: Request): + settings = request.app.state.runtime.settings + if not settings.schwab_api_key or not settings.schwab_app_secret: + raise HTTPException(503, "Live source is not configured") + return RedirectResponse(request.app.state.runtime.start_schwab_login()) + + @router.get("/bars") def bars( request: Request, diff --git a/app/api/schwab_auth.py b/app/api/schwab_auth.py index 5912780..ae6245b 100644 --- a/app/api/schwab_auth.py +++ b/app/api/schwab_auth.py @@ -19,10 +19,32 @@ brokerage this host talks to. Treat it as fixed: changing a registered callback URL means editing the Schwab app, which can send it back through approval. """ from fastapi import APIRouter, Request -from fastapi.responses import HTMLResponse +from fastapi.responses import HTMLResponse, RedirectResponse router = APIRouter(prefix="/api") + +def begin_login(settings) -> object: + from schwab.auth import get_auth_context + + return get_auth_context(settings.schwab_api_key, settings.schwab_callback_url) + + +def complete_login(settings, context, redirect_url: str) -> None: + from schwab import auth + from schwab.auth import client_from_received_url + + path = settings.schwab_token_path + path.parent.mkdir(parents=True, exist_ok=True) + write_token = getattr(auth, "__make_update_token_func")(str(path)) + client_from_received_url( + settings.schwab_api_key, + settings.schwab_app_secret, + context, + redirect_url, + write_token, + ) + PAGE = """ Callback @@ -56,8 +78,22 @@ version of this page is the one you are redirected to.

""" -@router.get("/qt", response_class=HTMLResponse) -def callback(request: Request) -> HTMLResponse: +FAILED = """ +

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 @@

/ESsent

-
{{ status.stream }} ({{ status.delay_minutes }}min delay) · {{ status.source || 'source' }}
+
{{ status.stream }} ({{ status.delay_minutes }}min delay) · {{ status.source || 'source' }}reconnect
diff --git a/static/style.css b/static/style.css index cc958b9..d1d924d 100644 --- a/static/style.css +++ b/static/style.css @@ -5,7 +5,7 @@ body { margin:0; background:var(--bg); color:var(--fg); font:14px/1.45 "IBM Plex header { height:44px; display:flex; align-items:center; justify-content:space-between; border-bottom:1px solid var(--line); margin-bottom:16px; } h1 { margin:0; font-size:22px; letter-spacing:-1px; } h1 strong { color:var(--accent); font-weight:600; } .eyebrow { color:var(--muted); font-size:9px; letter-spacing:2px; } -.status { text-transform:uppercase; color:var(--muted); font-size:11px; }.status i { display:inline-block; width:7px; height:7px; border-radius:50%; background:var(--red); margin-right:8px; }.status.connected i,.status.replay i { background:var(--green); box-shadow:0 0 9px var(--green); } +.status { text-transform:uppercase; color:var(--muted); font-size:11px; }.status i { display:inline-block; width:7px; height:7px; border-radius:50%; background:var(--red); margin-right:8px; }.status.connected i,.status.replay i { background:var(--green); box-shadow:0 0 9px var(--green); }.status .reconnect { margin-left:10px; color:var(--accent); letter-spacing:.4px; } main { display:grid; grid-template-columns:minmax(0, 1fr) 300px; gap:16px; } .chart-shell,aside { background:var(--panel); border:1px solid var(--line); } .chart-head { min-height:56px; padding:10px 14px; display:flex; align-items:center; justify-content:space-between; gap:12px; border-bottom:1px solid var(--line); } diff --git a/tests/test_schwab_callback.py b/tests/test_schwab_callback.py index 268050d..c8d98a4 100644 --- a/tests/test_schwab_callback.py +++ b/tests/test_schwab_callback.py @@ -39,3 +39,26 @@ def test_the_code_is_not_retained_for_a_later_visitor(): session.get("/api/qt", params={"code": "secret-code"}) # A second, code-less request must not replay the first one's code. assert "secret-code" not in session.get("/api/qt").text + + +def test_a_pending_login_exchanges_the_code_and_returns_home(monkeypatch): + from app.config import Settings + from app.runtime import Runtime + + app = FastAPI() + app.include_router(router) + runtime = Runtime(Settings()) + runtime.schwab_login = object() + finished = [] + + def finish(url): + finished.append(url) + + runtime.finish_schwab_login = finish + app.state.runtime = runtime + response = TestClient(app, follow_redirects=False).get( + "/api/qt", params={"code": "abc", "state": "s"} + ) + assert response.status_code in (302, 303, 307) + assert response.headers["location"] == "/" + assert finished and "code=abc" in finished[0] diff --git a/tests/test_schwab_source.py b/tests/test_schwab_source.py index 438be65..c74c2b8 100644 --- a/tests/test_schwab_source.py +++ b/tests/test_schwab_source.py @@ -293,3 +293,39 @@ def test_a_zero_last_price_is_not_a_price(): def test_a_zero_price_with_no_trade_markers_is_dropped_entirely(): assert parse_level_one({"content": [{"key": "/ES", "LAST_PRICE": 0}]}) == [] + + +def test_keepalive_hits_the_shared_http_client(): + class FakeClient: + def __init__(self): + self.calls = 0 + + async def get_user_preferences(self): + self.calls += 1 + return type("R", (), {"status_code": 200})() + + source = SchwabSource(Settings()) + client = FakeClient() + source._http = client + asyncio.run(source.refresh_token()) + assert client.calls == 1 + + +def test_reset_client_drops_the_cached_http_session(): + source = SchwabSource(Settings()) + source._http = object() + source.reset_client() + assert source._http is None + + +def test_needs_login_when_the_refresh_token_is_dead(tmp_path): + from app.runtime import Runtime + + runtime = Runtime(Settings(manual_lines_path=tmp_path / "lines.json")) + assert runtime.needs_login() is False + runtime.stream.last_error = ( + 'unsupported_token_type: 400 Bad Request: ' + '{"error_description":"Refresh token is invalid, expired or revoked",' + '"error":"invalid_grant"}' + ) + assert runtime.needs_login() is True