Compare commits

..

No commits in common. "52e657fb1e6f8410b37db7a39ccf16c8f789df6c" and "e69d1c190f91ee6dacad33dde6db507dd0ede303" have entirely different histories.

18 changed files with 16 additions and 1357 deletions

View file

@ -6,17 +6,6 @@ YAHOO_POLL_SECONDS=20
SEED_1H_RANGE=730d
SEED_1M_RANGE=8d
# Schwab. Only read once LIVE_SOURCE=schwab; blank is fine until then.
# The callback must match the app registration exactly — changing it there can
# re-trigger approval. The same URL works from local and production: the
# redirect lands in your browser, not on the machine running the code.
SCHWAB_API_KEY=
SCHWAB_APP_SECRET=
SCHWAB_CALLBACK_URL=https://chart.amow.com/api/qt
# Under data/ so it survives a deploy — that path is the Coolify volume.
SCHWAB_TOKEN_PATH=./data/.schwab_token.json
SCHWAB_SYMBOL=/ES
# Chart and analysis
TIMEFRAMES=1m,5m,15m,30m,1h,1d
BASE_TIMEFRAMES=1m,30m,1d

View file

@ -172,68 +172,6 @@ knowing:
If it lands near other levels it clusters with them and the score adds up, so a typed
level sitting on the prior-day close reads as one zone rather than two alerts.
## Market data: Yahoo for the past, Schwab for the present
Both sources run together. This is the intended configuration, not a fallback:
- **Schwab** streams real-time `/ES` minute bars over `CHART_FUTURES`
(`delayed: false`), but serves **no futures history at all** — everything it
knows starts when you connect.
- **Yahoo** has roughly 730 days of hourly data, which is what makes a 200-day
moving average warm at startup rather than in ten months. It lags ~10 minutes.
Switch the live feed with `LIVE_SOURCE=schwab`; seeding stays on Yahoo whatever
you set, because Schwab has nothing to seed from. The symbols differ — Yahoo says
`ES=F`, Schwab says `/ES` — and `Settings.live_symbol` picks the right one. Every
bar carries a `source` tag so the seam stays visible.
**Expect a gap of up to ten minutes at the right-hand edge after a restart.**
Yahoo's history reaches to *now − 10 min* while the stream starts at *now*, so
the most recent bars are briefly missing. It backfills itself as Yahoo catches
up. Live price and alerts are unaffected — those come from the stream. Persisting
bars would remove it entirely.
`/ES` resolves to the active contract (`/ESU26` today) on Schwab's side, so
contract rolls need no handling.
### Rate limits
The developer portal shows **Order Limit: 120**, which caps orders per minute —
this app places none. Schwab separately rate-limits REST calls, commonly cited at
120 per minute.
Neither constrains us, because **streaming is not REST**. The WebSocket is a
single connection and bars arrive by push, so steady-state Schwab REST usage is
essentially zero. This is a concrete advantage of the stream over the polling
fallback: polling `get_quotes()` once a second would have sat at roughly half the
limit permanently, forever.
The exception is reconnects. Each one calls `get_user_preferences()` to fetch the
socket URL, and `StreamService` retries every five seconds, so a sustained outage
generates about twelve REST calls a minute. Comfortably under, but not nothing —
worth remembering before shortening that backoff.
Yahoo is a different service and none of these limits apply to it.
### Authenticating
```bash
python3 -m scripts.check_schwab # prints the login URL
python3 -m scripts.check_schwab --redirect-url '…' # exchanges the code
python3 -m scripts.check_stream 60 /ES # confirms bars arrive
```
Two steps, neither interactive, so the browser can be on a different machine —
you copy a URL out and paste one back. **The authorisation code expires in about
thirty seconds**, so have the second command ready before you approve. The
callback page 404s until this branch is deployed; that is cosmetic, the code is
in the address bar regardless.
The token refreshes itself for seven days, then needs the flow again. It is
per-machine: `.schwab_token.json` locally, and under `data/` in production, which
is the Coolify volume. Locally `data/` is owned by root because Docker created it
through the bind mount, which is why the local path differs.
## Alerts and ntfy
Alerts are evaluated **server-side**, once per closed 1m bar, by a single engine

View file

@ -1,72 +0,0 @@
"""Schwab OAuth callback.
Schwab requires an HTTPS callback URL. The usual answer is
``https://127.0.0.1:8182`` with a self-signed certificate, which means clicking
through a browser warning on every re-authentication — and the refresh token
expires weekly. There are also reports of Schwab refusing to register apps whose
callback is a loopback address.
This app already terminates real HTTPS, so it can receive the redirect itself.
Deliberately unauthenticated: Schwab redirects a browser here and cannot attach
the chart token. Nothing is stored — the page only echoes the query string of
the request that produced it, which the caller already has in their address bar.
Storing the code would mean a later, unauthenticated visitor could read it.
The path is deliberately unrevealing. That is not a security control — the
endpoint's safety is that it is inert — it simply avoids advertising which
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
router = APIRouter(prefix="/api")
PAGE = """<!doctype html>
<meta charset="utf-8">
<title>Callback</title>
<style>
body {{ font: 15px/1.6 ui-sans-serif, system-ui, sans-serif; max-width: 46rem;
margin: 3rem auto; padding: 0 1.5rem; background: #14161a; color: #e8eaed; }}
h1 {{ font-size: 1.2rem; }}
code, textarea {{ font-family: ui-monospace, monospace; font-size: 13px; }}
textarea {{ width: 100%; height: 7rem; padding: .7rem; border-radius: 6px;
border: 1px solid #2a2e35; background: #0e1013; color: #e8eaed; }}
.warn {{ color: #efb643; }}
.muted {{ color: #9aa1ab; }}
</style>
<h1>{heading}</h1>
{body}
"""
RECEIVED = """
<p>Paste this entire URL into the waiting login prompt:</p>
<textarea readonly onclick="this.select()">{url}</textarea>
<p class="warn">Single use, and it expires within minutes. Do not share it.</p>
<p class="muted">Nothing was stored on the server. This page shows only the URL
you just arrived with.</p>
"""
IDLE = """
<p>OAuth callback endpoint. Register this exact URL with the provider:</p>
<p><code>{url}</code></p>
<p class="muted">Arriving here directly is expected and harmless — the useful
version of this page is the one you are redirected to.</p>
"""
@router.get("/qt", response_class=HTMLResponse)
def callback(request: Request) -> HTMLResponse:
if request.query_params.get("code"):
page = PAGE.format(
heading="Authorisation received",
body=RECEIVED.format(url=str(request.url)),
)
else:
page = PAGE.format(
heading="Callback endpoint",
body=IDLE.format(url=str(request.url).split("?")[0]),
)
# Never cached: it carries a single-use authorisation code.
return HTMLResponse(page, headers={"Cache-Control": "no-store"})

View file

@ -32,20 +32,6 @@ class Settings(BaseSettings):
ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200"
daily_anchor_et: str = "18:00"
manual_lines_path: Path = Path("./data/manual_lines.json")
# Schwab. Empty until the app's keys are issued; nothing reads them while
# live_source is yahoo. The token lives under data/ so it lands on the
# Coolify persistent volume — a rebuild would otherwise log you out, and
# re-authenticating is an interactive browser flow.
schwab_api_key: str = ""
schwab_app_secret: str = ""
schwab_callback_url: str = "https://chart.amow.com/api/qt"
schwab_token_path: Path = Path("./data/.schwab_token.json")
schwab_symbol: str = "/ES"
# Seconds between forming-bar emissions built from LEVEL_ONE_FUTURES ticks.
# Set negative to drop the Level 1 subscription and take closed minute bars
# only. 1.0 is a candle that visibly moves without a broadcast per trade.
schwab_tick_seconds: float = 1.0
confluence_min_score: float = 28
# Four hours, chosen from the sweep in scripts/calibrate_alerts.py. Suppression is
# per price zone, so an unrelated zone still alerts immediately; this only
@ -56,15 +42,6 @@ class Settings(BaseSettings):
chart_auth_token: str = ""
replay_file: Path | None = None
@property
def live_symbol(self) -> str:
"""What the live source calls the instrument.
Yahoo wants ES=F, Schwab wants /ES. Seeding always uses the Yahoo
symbol, because Yahoo is always the source of history.
"""
return self.schwab_symbol if self.live_source == "schwab" else self.yahoo_symbol
@property
def enabled_timeframes(self) -> list[Timeframe]:
return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()]

View file

@ -7,27 +7,15 @@ from app.market.yahoo import YahooSource
def live_source(settings: Settings) -> MarketDataSource:
if settings.live_source == "yahoo":
return YahooSource(settings.yahoo_poll_seconds)
if settings.live_source == "schwab":
# Imported here so schwab-py stays optional while the live source is
# Yahoo, which is still the default.
from app.market.schwab import SchwabSource
return SchwabSource(settings)
if settings.live_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file)
raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}")
def seed_source(settings: Settings) -> MarketDataSource | None:
"""The source of history, which is never Schwab.
Schwab serves price history for equities and ETFs only, so with
LIVE_SOURCE=schwab the seed stays on Yahoo. That pairing is the intended
configuration, not a fallback: Yahoo supplies the past, Schwab the present.
"""
if settings.seed_source == "none":
return None
if settings.seed_source in ("yahoo", "schwab"):
if settings.seed_source == "yahoo":
return YahooSource(settings.yahoo_poll_seconds)
if settings.seed_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file)

View file

@ -1,218 +0,0 @@
"""Real-time /ES bars from Schwab's CHART_FUTURES stream.
Verified against a live account before this was written:
- Streaming works. CHART_FUTURES delivers one minute bar per symbol per minute
with true exchange OHLCV, and LEVEL_ONE_FUTURES reports ``delayed: False``.
- The continuous root resolves itself. Subscribing to ``/ES`` returns data keyed
``/ES`` while quotes report the active contract as ``/ESU26``, so contract
rolls need no handling here.
- There is no history. Schwab serves price history for equities and ETFs only,
so this source seeds nothing; Yahoo remains the only source of the past.
The delayed sibling is worth stating plainly: Yahoo lags about ten minutes, so
at startup the most recent bars are missing until Yahoo catches up. Keep both
sources running rather than switching Yahoo off once this connects.
"""
import asyncio
import logging
import time
from collections.abc import AsyncIterator
from dataclasses import replace
from app.bars.models import Bar, Timeframe
logger = logging.getLogger(__name__)
# CHART_FUTURES field names as schwab-py labels them.
FIELD_TIME = "CHART_TIME_MILLIS"
FIELD_OPEN = "OPEN_PRICE"
FIELD_HIGH = "HIGH_PRICE"
FIELD_LOW = "LOW_PRICE"
FIELD_CLOSE = "CLOSE_PRICE"
FIELD_VOLUME = "VOLUME"
# LEVEL_ONE_FUTURES field names, as schwab-py labels them. Verified realtime on
# this account: the service reports delayed: False for /ES.
FIELD_LAST_PRICE = "LAST_PRICE"
FIELD_LAST_SIZE = "LAST_SIZE"
FIELD_TRADE_TIME = "TRADE_TIME_MILLIS"
def parse_level_one(message: dict) -> list[tuple[int, float, int]]:
"""Turn one LEVEL_ONE_FUTURES message into (trade time ms, price, size).
Level 1 messages are partial: a quote that moves only the bid carries no
LAST_PRICE at all. Those are skipped rather than carried forward, because a
bid tick is not a trade and must not extend a candle's high or low.
"""
ticks: list[tuple[int, float, int]] = []
for content in message.get("content") or []:
price = content.get(FIELD_LAST_PRICE)
if price is None:
continue
millis = content.get(FIELD_TRADE_TIME)
if millis is None:
# No trade stamp on this update; the wall clock is close enough to
# bucket it, and being one minute out at a boundary is corrected by
# the authoritative CHART_FUTURES bar moments later.
millis = int(time.time() * 1000)
ticks.append((int(millis), float(price), int(content.get(FIELD_LAST_SIZE) or 0)))
return ticks
def parse_chart_futures(message: dict, symbol: str) -> list[Bar]:
"""Turn one CHART_FUTURES message into bars.
A bar arrives once its minute has elapsed, so it is complete on arrival and
marked closed. Anything missing a timestamp or a price is skipped rather
than defaulted — a bar invented from partial data would be indistinguishable
from a real one downstream.
"""
bars: list[Bar] = []
for content in message.get("content") or []:
millis = content.get(FIELD_TIME)
prices = [content.get(field) for field in (FIELD_OPEN, FIELD_HIGH, FIELD_LOW, FIELD_CLOSE)]
if millis is None or any(price is None for price in prices):
continue
open_, high, low, close = (float(price) for price in prices)
bars.append(
Bar(
tf=Timeframe.M1,
t=int(millis) // 1000,
o=open_,
h=high,
l=low,
c=close,
v=int(content.get(FIELD_VOLUME) or 0),
closed=True,
symbol=str(content.get("key") or symbol),
source="schwab",
)
)
return bars
class SchwabSource:
"""Live minute bars. Holds no history — see the module docstring."""
name = "schwab"
delay_minutes = 0
def __init__(self, settings, stream_client_factory=None):
self._settings = settings
# Injectable so the parsing and dispatch can be tested without a socket.
self._stream_client_factory = stream_client_factory or self._build_stream_client
# None disables the Level 1 subscription entirely and leaves the source
# 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
def supports_history(self) -> bool:
return False
async def history(self, symbol, tf, start, end, *, range_=None) -> list[Bar]:
return []
def supports_stream(self) -> bool:
return True
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)
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
stream_client = self._stream_client_factory()
queue: asyncio.Queue[tuple[str, dict]] = asyncio.Queue(maxsize=256)
def enqueue(kind: str):
def handler(message: dict) -> None:
# Dropping the oldest keeps a slow consumer from stalling the
# socket; a minute bar that late is of no use anyway.
if queue.full():
queue.get_nowait()
queue.put_nowait((kind, message))
return handler
await stream_client.login()
# Registered before subscribing: the service starts sending straight
# away and messages without a handler are discarded.
stream_client.add_chart_futures_handler(enqueue("chart"))
await stream_client.chart_futures_subs([symbol])
logger.info("Subscribed to CHART_FUTURES for %s", symbol)
if self._tick_seconds is not None:
# Same socket, same login — no extra REST call and no extra rate
# limit. CHART_FUTURES only speaks once a minute, after the minute
# is over; this is what makes the candle move in between.
stream_client.add_level_one_futures_handler(enqueue("quote"))
await stream_client.level_one_futures_subs([symbol])
logger.info("Subscribed to LEVEL_ONE_FUTURES for %s", symbol)
forming: Bar | None = None
last_closed_t = 0
last_emit = 0.0
pump = asyncio.create_task(self._pump(stream_client), name="schwab-stream-pump")
try:
while True:
if pump.done():
# Surface the socket's failure rather than hanging on a
# queue nothing is filling any more.
pump.result()
return
try:
kind, message = await asyncio.wait_for(queue.get(), timeout=5)
except (asyncio.TimeoutError, TimeoutError):
continue
if kind == "chart":
for bar in parse_chart_futures(message, symbol):
last_closed_t = max(last_closed_t, bar.t)
# The exchange's own bar supersedes whatever the ticks
# had built for that minute.
if forming is not None and forming.t <= bar.t:
forming = None
yield bar
continue
for millis, price, size in parse_level_one(message):
minute = millis // 60000 * 60
# A tick for a minute already closed by CHART_FUTURES would
# otherwise overwrite an authoritative bar with a partial.
if minute <= last_closed_t:
continue
if forming is None or forming.t != minute:
forming = Bar(
tf=Timeframe.M1, t=minute, o=price, h=price, l=price, c=price,
v=size, closed=False, symbol=symbol, source="schwab",
)
else:
forming.h = max(forming.h, price)
forming.l = min(forming.l, price)
forming.c = price
forming.v += size
# Throttled: /ES trades many times a second, and every
# emission costs a store write and a broadcast to every
# open socket.
now = time.monotonic()
if now - last_emit >= self._tick_seconds:
last_emit = now
yield replace(forming)
finally:
pump.cancel()
@staticmethod
async def _pump(stream_client) -> None:
while True:
await stream_client.handle_message()

View file

@ -22,17 +22,11 @@ class StreamService:
self._handlers.append(handler)
async def seed(
self,
source: MarketDataSource | None,
tf: Timeframe,
range_: str,
symbol: str | None = None,
self, source: MarketDataSource | None, tf: Timeframe, range_: str
) -> None:
# The seed source names the instrument differently from the live one:
# Yahoo says ES=F where Schwab says /ES.
if source is None or not source.supports_history():
return
bars = await source.history(symbol or self.symbol, tf, None, None, range_=range_)
bars = await source.history(self.symbol, tf, None, None, range_=range_)
for bar in bars:
await self._emit(bar)

View file

@ -50,24 +50,10 @@ class Runtime:
self.settings.confluence_min_score, self.settings.alert_cooldown_seconds
)
self.levels = self.manual_lines.levels()
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
self.stream = StreamService(live_source(self.settings), self.settings.yahoo_symbol)
self.stream.add_handler(self.on_bar)
async def on_bar(self, bar: Bar) -> None:
# A tick-built bar is provisional and arrives many times a minute. It
# updates the last candle and the live price, and stops there.
#
# It must not reach the aggregator: that accumulates volume with
# `current.v += incoming.v`, so re-sending the same forming minute would
# add its volume to every higher timeframe again on each update. Alerts
# stay on closed bars for the same reason they always were — a level is
# judged on a settled bar, not on a price that may not last the minute.
if not bar.closed:
self.store.put(bar)
self.price = bar.c
self.broadcast({"type": "bar", "bar": bar})
return
evaluate_alerts = False
for aggregated in self.aggregator.update(bar):
self.store.put(aggregated)
@ -195,14 +181,8 @@ class Runtime:
async def start(self) -> asyncio.Task:
try:
source = seed_source(self.settings)
# Always the Yahoo symbol: Schwab has no history to seed from.
seed_symbol = self.settings.yahoo_symbol
await self.stream.seed(
source, Timeframe.H1, self.settings.seed_1h_range, seed_symbol
)
await self.stream.seed(
source, Timeframe.M1, self.settings.seed_1m_range, seed_symbol
)
await self.stream.seed(source, Timeframe.H1, self.settings.seed_1h_range)
await self.stream.seed(source, Timeframe.M1, self.settings.seed_1m_range)
except Exception:
# A transient seed failure must not prevent the live stream or UI starting.
pass

View file

@ -168,37 +168,10 @@ buckets cleanly. The ~730-day 1h window yields ~500 sessions: enough for a daily
Yahoo's native 1d bars may be used *only* for multi-year context, clearly labelled.
## 2.2 Schwab entitlements — **answered empirically 2026-08-10**
## 2.2 Verify these when adding the Schwab source (M6, not before)
All verified against a live production app holding both Market Data Production and
Accounts and Trading Production.
| Question | Answer |
|---|---|
| Futures market data entitled? | **Yes** — but only via `get_quotes()` (plural) |
| Symbol format | **`/ES`**, which auto-resolves to the active contract `/ESU26` |
| `CHART_FUTURES` streaming | **Works** — one true-OHLCV minute bar per symbol per minute |
| `LEVEL_ONE_FUTURES` | **Works**, and reports `delayed: false` |
| Futures price history | **Still none.** Yahoo remains the only source of the past |
Three traps found the hard way, all of which cost a round trip:
- **`get_quote()` (singular) silently returns the wrong instrument.** It puts the
symbol in the URL *path*, where the leading slash is normalised away, so `/ES`
comes back as `ES` — Eversource Energy, an equity, at $72. HTTP 200 with a
populated body. `get_quotes()` passes symbols as a query parameter and returns
the future correctly. **Never treat a 200 as proof; check `assetMainType`.**
- **Streaming needs the Accounts and Trading product.** `StreamClient.login()`
reads `/trader/v1/userPreference` for its socket URL, and that path is not in
Market Data Production. A market-data-only app cannot stream at all.
- **Authorisation codes expire in about thirty seconds**, and an unwritable token
path spends one before revealing itself. `scripts/check_schwab.py` preflights
the key, the secret and the token path for exactly this reason.
Because `/ES` resolves to the active contract on Schwab's side, contract roll
handling — an open problem in §10 — needs no code here.
### The remaining unknowns for Schwab
These are load-bearing unknowns for Schwab specifically. They no longer block the
project — M0–M5 run entirely on Yahoo. If any fails, stop and report.
1. **Futures market-data entitlement.** It is not publicly documented whether
`CHART_FUTURES` requires futures trading approval or a CME non-professional market
@ -810,26 +783,6 @@ const chartApi = shallowRef(null); // ✅
// const chart = ref(null); // ❌ will appear to work, then misbehave
```
### Logical indices address the whole chart, not your bar array
**Never derive a viewport from `bars.length`.** A logical index addresses the
chart's *shared* time scale — the union of the time points of every series on
it — not the candle array. Any series whose points pre-date the candle window
prepends to that scale and shifts every logical index by its count.
```js
// ❌ off by however many points the other series contribute
timeScale().setVisibleLogicalRange({ from: bars.length - 160, to: bars.length + 5 });
// ✅ an instant cannot be renumbered by a later series
timeScale().setVisibleRange({ from: bars.at(-160).t, to: bars.at(-1).t + step * 5 });
```
This is not theoretical — see the 2026-08-10 entry in §16. The daily MAs carry
one point per daily bar (617 of them, back ~2 years). They are attached by
`syncVisibleLevels()` *immediately after* `setBars()`, so a logical range that
was correct when set silently slid 617 bars — about ten hours — into the past
one tick later. The chart looked frozen while the socket was perfectly healthy.
Structure:
- `chart.js` — a plain, framework-free wrapper class owning the LWC instance:
`create(el)`, `setBars()`, `updateBar()`, `syncLevels(levels)`, `destroy()`.
@ -1191,134 +1144,3 @@ enforces for data sources.
| 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 |
| Viewport derived from `bars.length` | Chart looks frozen; feed is fine | Anchor the view by time, never by logical index — §9 |
| Seed replays every bar through `on_bar` | ~82 s startup; port refuses connections | Known, unfixed — §16, 2026-08-10 |
| Headless browser without a real locale | `Intl` throws; blank canvas mimics an app bug | Launch Chromium with `--lang=en-US` — §16 |
## 16. Session log
Dated record of problems hit and how they were resolved. Times are UTC; the
repo's commit timestamps are -0500.
### 2026-08-10 — rebuild, and a chart that looked frozen
**10:30 · The rebuild was genuinely required.** `schwab-py` had been added to
`requirements.txt`, but the running image was built at 2026-08-09 22:05, before
that line existed. The bind mount (`.:/app`) hides this: source edits appear
live, so the Schwab commits looked deployed while `pip freeze` in the container
showed no `schwab-py` at all. Anything imported rather than read from disk needs
`docker compose build`. Rebuilt to `schwab-py 1.5.1` and recreated the container.
**10:30–10:31 · Startup takes ~82 seconds, and the port is closed the whole
time.** `Runtime.start()` replays every seeded bar through `on_bar`, and each
daily-bar update re-runs `rebuild_levels()` → `broadcast_level_delta()`, which
serialises and diffs five MA levels carrying ~730 points each. With a 730d/1h
seed plus an 8d/1m seed that is quadratic work before uvicorn binds. Measured:
10:30:24 "Waiting for application startup" → 10:31:46 "Application startup
complete". An open browser tab polling `/api/status` throughout logs a wall of
`ERR_CONNECTION_REFUSED`; that is the restart window, not a fault.
*Unfixed.* The fix is to bulk-load seeded bars and rebuild levels once at the
end, rather than once per bar. Related: M7 persistence would cut the seed itself.
**Diagnosing a hang that is actually slowness:** `docker stats` reported ~0.1%
CPU while the process was in fact grinding, so it pointed the wrong way. What
worked was `faulthandler.dump_traceback_later(25, exit=True)`, which named the
exact frame (`indicators.py:sma` under `runtime.py:62`). `py-spy` is unusable
here — it needs `SYS_PTRACE`, which the container does not have.
**Do not write scratch files into the repo while diagnosing.** A `_probe.py`
dropped in the project root is inside the bind mount, so `--reload` restarted
the lifespan and reset the 82-second clock — twice — which is what made
slow startup look like an infinite hang. Pipe throwaway scripts over stdin
(`docker exec -i … python -`) instead. Only `.py` changes trigger the reloader;
writing screenshots into `artifacts/` is safe.
**10:35 · A blank chart canvas that was not a bug.** The Playwright container
has no usable locale, so Chromium reports `en-US@posix`; Lightweight Charts
formats its time axis through `Intl`, which throws `Invalid language tag` and
leaves the canvas empty. `docker-compose.yml` already sets
`LANG=en_US.UTF-8` for that service and it is *not* sufficient. Launch with
`chromium.launch({ args: ['--lang=en-US'] })` — with that, the page renders and
reports zero console errors. Worth stating plainly: this failure looks exactly
like a broken app, and it is not.
**10:41–10:50 · The real bug — the chart sat ~10 hours behind a healthy feed.**
Symptom: header price live at 7785.00 while the last candle closed 7772.75, and
the series appeared to end at 00:20. Everything downstream checked out —
`/api/bars` newest 10:39 from `schwab`; `store.put` keeps bars strictly
ascending; the WebSocket snapshot delivered 1000 ascending bars ending 10:42 and
live `bar` events arrived every minute; the browser received all of it.
Interrogating `window.__chart` gave the answer:
```
seriesLen 1000 seriesLast 08-10T10:48 (7786.25) ← data complete
visible 08-07T20:41 → 08-10T00:35 ← viewport wrong
logical from 840 to 1005
```
The series was complete; the *viewport* was 617 bars too far left — exactly
`bars_held.1d`. `setBars()` set a visible **logical** range from the candle
array length, then `syncVisibleLevels()` attached the daily MA series, whose 617
daily points pre-date the 1m window; prepending them renumbered every logical
index and dragged the view off the live edge. Fixed in `static/chart.js` by
anchoring the viewport to a **time** range. Verified in a real browser: visible
range 08:02 → 10:49, last candle 7786.25 matching the header. See §9.
**Method note.** Three checks in a row said "healthy" — the REST API, the
WebSocket, and the frontend source all looked correct in isolation, because each
of them *was* correct. Only querying the live page's own chart object separated
"the data is missing" from "the data is off-screen". Screenshots alone were
actively misleading here: the stale time axis was read as a session gap.
### 2026-08-10 (later) — real-time ticks, and what to do about cold restarts
**The chart now moves between minute closes.** `CHART_FUTURES` emits a bar only
once its minute is over, so the chart stepped once a minute and sat still in
between — read, reasonably, as a dead feed. `LEVEL_ONE_FUTURES` carries real
trades on the same socket (`delayed: False`, verified on this account back in
M6), and it was never subscribed. It is now, and it builds a forming bar for the
current minute which the authoritative `CHART_FUTURES` bar then supersedes.
Three constraints shaped it, each of which would have caused a real bug:
- **Tick bars must never reach the aggregator.** It accumulates with
`current.v += incoming.v`, so re-sending the same forming minute would add its
volume into every higher timeframe on every update. `Runtime.on_bar` returns
early for `not bar.closed`: store the bar, set the price, broadcast, stop.
- **Ticks are throttled** (`SCHWAB_TICK_SECONDS`, default 1.0). /ES trades many
times a second and each emission costs a store write plus a broadcast to every
open socket. Setting it negative drops the Level 1 subscription entirely and
returns the source to closed bars only.
- **A tick for a minute already closed is dropped**, or a late trade would
overwrite a settled exchange bar with a partial one.
Alerts deliberately stay on closed bars. A level is judged on a settled bar, not
on a price that may not last the minute — and `on_bar` already gated on
`closed`, so this needed no change. Intra-bar alerting is a separate decision.
Bid-only Level 1 updates are skipped rather than carried forward: a bid is not a
trade and must not extend a candle's high or low. Verified live — 15 forming
bars and 2 closed bars in 100 seconds, and in a browser the candle's high and low
visibly extend within the minute.
**Cold restarts — the options, and a recommendation.** Every restart costs ~82
seconds of refused connections, re-seeds from Yahoo, and starts with empty alert
cooldowns, so a deploy can re-alert whatever price is sitting on.
1. *Make seeding non-quadratic.* Seeding replays every bar through `on_bar`, and
each daily-bar update rebuilds all five MA levels and diffs them. Bulk-load
the seeded bars and rebuild levels once at the end. Contained, testable, and
removes most of the 82 seconds. **Do this first** — it is the cheapest real
win and needs no new storage.
2. *Persist bars (M7, SQLite).* Restarts then seed only the gap. Removes the
Yahoo dependency from the startup path and shrinks the window further. This
is the durable answer, and the plan already scopes it.
3. *Persist alert cooldowns and armed state.* Independent of 1 and 2, and the
part that actually misbehaves rather than merely being slow: without it every
deploy re-alerts. Small table, big behavioural win.
4. *Serve before seeding finishes.* Start uvicorn immediately and seed in a
background task, so the port never refuses. The chart would open cold and
fill in, which is better than an unreachable page — but it changes what
"warm" means to every consumer of `/api/status`, so it wants its own thought.
Recommended order: 1, then 3, then 2. 4 only if the window still bites after 1.

View file

@ -12,7 +12,6 @@ from fastapi.staticfiles import StaticFiles
from app.api.meta import router as meta_router
from app.api.routes import router as api_router
from app.api.schwab_auth import router as schwab_auth_router
from app.api.ws import router as ws_router
from app.config import Settings
from app.runtime import Runtime
@ -38,7 +37,6 @@ app = FastAPI(title="chart", lifespan=lifespan)
app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static")
app.include_router(meta_router)
app.include_router(schwab_auth_router)
app.include_router(api_router)
app.include_router(ws_router)

View file

@ -2,5 +2,3 @@ fastapi
uvicorn[standard]
httpx
pydantic-settings
# Live futures stream; imported only when LIVE_SOURCE=schwab.
schwab-py

View file

@ -1,255 +0,0 @@
"""Find out what a Schwab app is actually entitled to, before building on it.
Answers the three questions that decide the design, empirically rather than
from documentation:
1. Do the credentials authenticate at all?
2. Do REST quotes work for /ES — futures market data is a separate
entitlement from equities and may not be granted.
3. Does the streamer connect? Its bootstrap reads /trader/v1/userPreference,
which belongs to the Accounts and Trading product, so an app registered
for Market Data Production alone is expected to fail there.
Two steps, neither of them interactive, so this works over a pipe or from an
agent session where stdin is not a terminal:
python3 -m scripts.check_schwab
prints the Schwab login URL.
python3 -m scripts.check_schwab --redirect-url 'https://.../api/qt?code=…'
exchanges the code, saves the token, runs the checks.
The browser does not have to be on this machine. Nothing is captured locally —
you copy a URL out, and paste a URL back. The authorisation code is single use
and expires within minutes, so do not leave it sitting between the two steps.
Once a token exists, running with no arguments skips straight to the checks.
"""
import argparse
import sys
from urllib.parse import parse_qs, urlparse
from app.config import Settings
def heading(text: str) -> None:
print(f"\n{text}\n{'-' * len(text)}")
def load_settings() -> Settings:
settings = Settings()
if not settings.schwab_api_key or not settings.schwab_app_secret:
raise SystemExit("Set SCHWAB_API_KEY and SCHWAB_APP_SECRET in .env first")
return settings
def key_is_live(authorization_url: str) -> bool:
"""Ask Schwab whether it recognises the app key, before opening a browser.
A key Schwab does not know produces `invalid_client` here — and so does a
deliberately invented one, byte for byte, so this cannot tell "wrong value"
from "not active yet". It can still save a confusing round trip through the
login page.
"""
import httpx
try:
response = httpx.get(authorization_url, follow_redirects=False, timeout=15)
except Exception:
return True # Network trouble is not evidence about the key.
return "invalid_client" not in response.text
def secret_is_valid(settings: Settings) -> bool | None:
"""Check the secret without needing an authorisation code.
The token endpoint authenticates the key and secret over HTTP Basic before
it looks at the grant, so a deliberately invalid code separates the two
failures: bad credentials give invalid_client, good credentials give
invalid_grant. Otherwise a truncated secret survives the login unnoticed and
only surfaces at the exchange, after the code has been spent.
None when the answer is not clear enough to act on.
"""
import httpx
try:
response = httpx.post(
"https://api.schwabapi.com/v1/oauth/token",
auth=(settings.schwab_api_key, settings.schwab_app_secret),
data={
"grant_type": "authorization_code",
"code": "deliberately-invalid-code",
"redirect_uri": settings.schwab_callback_url,
},
timeout=20,
)
except Exception:
return None
if "invalid_grant" in response.text:
return True
if "invalid_client" in response.text or response.status_code == 401:
return False
return None
def print_login_url(settings: Settings) -> None:
from schwab.auth import get_auth_context
context = get_auth_context(settings.schwab_api_key, settings.schwab_callback_url)
if secret_is_valid(settings) is False:
print("Schwab rejects the app key and secret pair.\n")
print("The key alone is accepted at the authorize endpoint, so this is")
print("the secret. Re-copy it from the portal using the Show icon —")
print("a value clipped by one character looks entirely normal.")
sys.exit(1)
if not key_is_live(context.authorization_url):
print("Schwab rejects this app key with invalid_client.\n")
print("The secret is not involved yet — it is only used when the code is")
print("exchanged — so this is the key itself or the app's readiness.\n")
print(" 1. Re-copy the App Key from the portal using the Show icon.")
print(" 2. If it matches, the app is most likely not live yet. Newly")
print(" created or newly edited apps take time to propagate, and the")
print(" portal says Ready For Use before the key works.")
print("\nRe-run this to check again; nothing else is needed.")
sys.exit(1)
print("Open this in any browser, on any machine, and approve the app:\n")
print(f" {context.authorization_url}\n")
print("You will land on the callback URL. A 404 there is fine until the")
print("branch is deployed — the code is in the address bar either way.")
print("Copy the ENTIRE address and run:\n")
print(" python3 -m scripts.check_schwab --redirect-url '<paste it here>'")
def ensure_token_path_writable(settings: Settings) -> None:
path = settings.schwab_token_path
try:
path.parent.mkdir(parents=True, exist_ok=True)
probe = path.parent / f".{path.name}.probe"
probe.touch()
probe.unlink()
except OSError as error:
raise SystemExit(
f"Cannot write the token to {path} ({error.strerror}).\n"
f"Point SCHWAB_TOKEN_PATH at a directory you own and try again — "
f"the authorisation code is spent either way, so fix this first."
)
def exchange(settings: Settings, redirect_url: str):
from schwab import auth
from schwab.auth import AuthContext, client_from_received_url
state = parse_qs(urlparse(redirect_url).query).get("state", [None])[0]
if not state:
raise SystemExit("That URL has no ?state= — paste the full address you landed on")
# Checked before the exchange, not after. An authorisation code lives about
# thirty seconds and is single use, so discovering an unwritable token path
# afterwards costs a whole round trip through the browser — which is exactly
# what happened the first time, against a data/ directory owned by root
# because Docker created it through the bind mount.
ensure_token_path_writable(settings)
# The library's own writer, so the token file keeps the shape its loader
# expects rather than one guessed at here.
write_token = getattr(auth, "__make_update_token_func")(str(settings.schwab_token_path))
# Only the state is needed again; the authorisation URL is not, which is
# what lets the two steps share nothing. Taking the state from the pasted
# URL makes the CSRF check a formality — acceptable because the thing being
# guarded against is a redirect you did not initiate, and you pasted this
# one in by hand.
context = AuthContext(settings.schwab_callback_url, None, state)
return client_from_received_url(
settings.schwab_api_key,
settings.schwab_app_secret,
context,
redirect_url,
write_token,
)
def run_checks(client, settings: Settings) -> None:
heading(f"REST quote for {settings.schwab_symbol}")
quote = client.get_quote(settings.schwab_symbol)
print(f" HTTP {quote.status_code}")
payload = quote.json() if quote.status_code == 200 else {}
quotes_ok = False
if payload:
for symbol, data in list(payload.items())[:1]:
values = data.get("quote", {})
kind = data.get("assetMainType")
description = (data.get("reference") or {}).get("description")
print(f" {symbol}: {kind} — {description}")
print(f" last={values.get('lastPrice')} bid={values.get('bidPrice')} "
f"ask={values.get('askPrice')}")
# Schwab strips the leading slash and happily returns the equity of
# the same name: /ES comes back as Eversource Energy at 72. A 200
# with a body is not evidence of futures data, and treating it as
# such is worse than a clean failure.
quotes_ok = kind == "FUTURE"
if not quotes_ok:
print(" -> NOT futures. The slash was stripped and an equity")
print(" returned in its place; REST futures quotes are unavailable.")
else:
print(" -> futures market data IS available over REST")
else:
print(f" body: {quote.text[:200]}")
print(" -> no quote returned")
heading("Streamer bootstrap (/trader/v1/userPreference)")
prefs = client.get_user_preferences()
print(f" HTTP {prefs.status_code}")
streaming_ok = prefs.status_code == 200 and bool(prefs.json().get("streamerInfo"))
if streaming_ok:
print(f" streamerInfo entries: {len(prefs.json()['streamerInfo'])}")
print(" -> streaming is available; CHART_FUTURES should work")
else:
print(f" body: {prefs.text[:200]}")
print(" -> streaming is NOT available on this app")
heading("Summary")
print(f" Quotes : {'yes' if quotes_ok else 'no'}")
print(f" Streaming : {'yes' if streaming_ok else 'no'}")
if quotes_ok and not streaming_ok:
print("\n Polling REST quotes is then the real-time path: it removes")
print(" Yahoo's ten-minute delay without needing trading scope, at the")
print(" cost of building bars from snapshots rather than receiving")
print(" true exchange OHLCV.")
def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--redirect-url", help="the full URL you were redirected to")
args = parser.parse_args()
settings = load_settings()
try:
from schwab.auth import client_from_token_file
except ImportError:
raise SystemExit("pip install -r requirements-dev.txt (schwab-py is not installed)")
if args.redirect_url:
heading("Exchanging the authorisation code")
client = exchange(settings, args.redirect_url)
print(f" token written to {settings.schwab_token_path}")
elif settings.schwab_token_path.exists():
heading("Authentication")
client = client_from_token_file(
str(settings.schwab_token_path),
settings.schwab_api_key,
settings.schwab_app_secret,
)
print(f" reused the token at {settings.schwab_token_path}")
else:
print_login_url(settings)
sys.exit(0)
run_checks(client, settings)
if __name__ == "__main__":
main()

View file

@ -1,98 +0,0 @@
"""Does the Schwab stream actually deliver /ES bars?
REST is already known not to: a quote for /ES comes back as Eversource Energy,
because Schwab strips the leading slash and resolves the equity of the same
name. Streaming is a separate entitlement with its own services, so it has to be
tested separately — and it is the only remaining route to real-time futures,
since price history does not cover them either.
Subscribes to both futures services for a short window and reports what arrives:
python3 -m scripts.check_stream [seconds] [symbol ...]
Symbols default to the continuous /ES and the front-month contract, because the
streamer may accept one and not the other — REST accepts neither.
CHART_FUTURES is the one that matters — it carries the minute OHLCV the chart is
built on. LEVEL_ONE_FUTURES is the fallback: quotes only, from which bars would
have to be synthesised.
"""
import asyncio
import sys
from app.config import Settings
received: dict[str, list] = {"chart": [], "quote": []}
async def main(seconds: float, symbols: list[str]) -> None:
from schwab.auth import client_from_token_file
from schwab.streaming import StreamClient
settings = Settings()
if not settings.schwab_token_path.exists():
raise SystemExit("No token yet — 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,
)
stream = StreamClient(client)
print("logging in to the streamer...")
await stream.login()
print(" logged in")
# Handlers must be registered before subscribing: several services start
# sending immediately, and messages with no handler are dropped.
stream.add_chart_futures_handler(lambda msg: received["chart"].append(msg))
stream.add_level_one_futures_handler(lambda msg: received["quote"].append(msg))
for name, subscribe in (
("CHART_FUTURES", stream.chart_futures_subs),
("LEVEL_ONE_FUTURES", stream.level_one_futures_subs),
):
try:
await subscribe(symbols)
print(f" subscribed to {name} for {', '.join(symbols)}")
except Exception as error:
print(f" {name} subscription REJECTED: {type(error).__name__}: {error}")
print(f"\nlistening for {seconds:g}s...")
try:
await asyncio.wait_for(_pump(stream), timeout=seconds)
except asyncio.TimeoutError:
pass
print("\nResults")
print("-------")
for label, key in (("CHART_FUTURES (minute OHLCV)", "chart"),
("LEVEL_ONE_FUTURES (quotes)", "quote")):
messages = received[key]
print(f" {label}: {len(messages)} message(s)")
if messages:
print(f" sample: {str(messages[0])[:300]}")
if received["chart"]:
print("\n -> CHART_FUTURES works. Real-time minute bars are available,")
print(" which removes Yahoo's ten-minute delay entirely.")
elif received["quote"]:
print("\n -> Only quotes arrived. Bars would have to be synthesised")
print(" from them: real-time, but highs and lows approximated.")
else:
print("\n -> Nothing arrived. Either futures market data is not")
print(" entitled on this account, or the market is closed.")
async def _pump(stream) -> None:
while True:
await stream.handle_message()
if __name__ == "__main__":
args = sys.argv[1:]
window = float(args[0]) if args else 45
wanted = args[1:] or ["/ES", "/ESU26"]
asyncio.run(main(window, wanted))

View file

@ -32,7 +32,6 @@ class ConfluenceChart {
this.toolDownListener = null;
this.toolMoveListener = null;
this.toolUpListener = null;
this.pendingView = null;
}
static TICK = 0.25;
@ -145,22 +144,10 @@ class ConfluenceChart {
setBars(bars) {
this.bars = bars;
this.candles.setData(bars.map(this.toCandle));
// Anchored by time, not by logical index. A logical index addresses the
// chart's *shared* scale — the union of every series' time points — not
// this array. The daily MAs land straight after with hundreds of points
// pre-dating the 1m window, and prepending them shifts every logical index
// by that count, silently dragging the view ten hours off the live edge.
// A time range names the instant, so later series cannot move it.
if (bars.length) {
const last = bars[bars.length - 1];
// Keep the old five bars of right-hand breathing room, in seconds.
const step = bars.length > 1 ? last.t - bars[bars.length - 2].t : 60;
this.pendingView = {
from: bars[Math.max(0, bars.length - 160)].t,
to: last.t + step * 5,
};
this.chart.timeScale().setVisibleRange(this.pendingView);
}
this.chart.timeScale().setVisibleLogicalRange({
from: Math.max(0, bars.length - 160),
to: bars.length + 5,
});
requestAnimationFrame(() => this.renderAnchorHandles());
}
@ -203,42 +190,6 @@ class ConfluenceChart {
}
}
// Level points are sampled on their own timeframe — VWAP every minute, the
// daily averages once a session — and every distinct timestamp claims its own
// slot on the chart's shared scale. Left raw, a minute-resolution VWAP spread
// 160 hourly candles across 908 slots and drew them as unreadable slivers.
// Snapping onto the candle grid preserves the line's shape while keeping the
// scale one slot per candle, which is what makes the bars their proper width.
snapPointsToBars(points) {
if (!this.bars.length || !points.length) return points;
const times = this.bars.map(bar => bar.t);
const last = times[times.length - 1];
const byTime = new Map();
for (const point of points) {
// Clamped, not passed through: on a daily chart every one of VWAP's ~760
// minute points falls after the last candle's session open, and letting
// them keep their own times put all 760 back on the scale. Projection to
// the right of the last bar is handled by the caller instead.
if (point.time >= last) {
byTime.set(last, point.value);
continue;
}
let lo = 0;
let hi = times.length - 1;
while (lo < hi) {
const mid = (lo + hi) >> 1;
if (times[mid] < point.time) lo = mid + 1;
else hi = mid;
}
// Points older than the window collapse onto the first candle; the newest
// of them wins, which is the value in force when the window opens.
byTime.set(times[lo], point.value);
}
return [...byTime.entries()]
.sort((a, b) => a[0] - b[0])
.map(([time, value]) => ({ time, value }));
}
syncLevels(levels) {
this.levels = levels;
this.syncPriceLines(levels);
@ -276,7 +227,7 @@ class ConfluenceChart {
}
let data;
if (hasPoints) {
data = this.snapPointsToBars((level.points || []).map(([time, value]) => ({ time, value })));
data = (level.points || []).map(([time, value]) => ({ time, value }));
const latestTime = this.bars[this.bars.length - 1]?.t;
const latestValue = data[data.length - 1]?.value;
if (latestTime != null && latestValue != null && latestTime > data[data.length - 1].time) {
@ -287,16 +238,6 @@ class ConfluenceChart {
}
entry.series.setData(data);
}
// Re-anchor once, after the level series have reshaped the scale. setBars
// runs before them, so the width it asked for was derived from the previous
// timeframe's point density and the chart holds that width as new series
// arrive. Consumed rather than reapplied every time: VWAP resyncs a level
// every minute, and re-anchoring on each would yank the view back from
// wherever the user had panned it.
if (this.pendingView) {
this.chart.timeScale().setVisibleRange(this.pendingView);
this.pendingView = null;
}
this.renderAnchorHandles();
}

View file

@ -5,10 +5,6 @@
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>/ES Confluence</title>
<link rel="stylesheet" href="/static/style.css">
<!-- Font Awesome Free 7.3.1, from unpkg like the other two dependencies.
Pinned deliberately: an unpinned icon set is a silent redesign on someone
else's release. cdnjs 404s on this path — it does not carry 7.3.1. -->
<link rel="stylesheet" href="https://unpkg.com/@fortawesome/fontawesome-free@7.3.1/css/all.min.css">
<script src="https://unpkg.com/vue@3/dist/vue.global.prod.js"></script>
<script src="https://unpkg.com/lightweight-charts@5.2.0/dist/lightweight-charts.standalone.production.js"></script>
</head>
@ -38,7 +34,7 @@
<h2>Tools</h2>
<div class="tool" :class="{armed: armedTool === 'trendline'}">
<button class="tool-head" @click="armTool('trendline')" :aria-pressed="armedTool === 'trendline'">
<span class="tool-glyph"><i class="fa-solid fa-arrow-trend-up"></i></span>Trendline
<span class="tool-glyph">╱</span>Trendline
<span class="tool-state">{{ armedTool === 'trendline' ? 'drag on chart' : '' }}</span>
</button>
<div class="tool-body">
@ -55,7 +51,7 @@
</div>
<div class="tool" :class="{armed: armedTool === 'level'}">
<button class="tool-head" @click="armTool('level')" :aria-pressed="armedTool === 'level'">
<span class="tool-glyph"><i class="fa-solid fa-minus"></i></span>Price level
<span class="tool-glyph">─</span>Price level
<span class="tool-state">{{ armedTool === 'level' ? 'click on chart' : '' }}</span>
</button>
<div class="tool-body">

View file

@ -5,7 +5,6 @@ from starlette.websockets import WebSocketDisconnect
from app.api.meta import router as meta_router
from app.api.routes import router as api_router
from app.api.schwab_auth import router as schwab_auth_router
from app.api.ws import router as ws_router
from app.config import Settings
from app.runtime import Runtime
@ -20,7 +19,6 @@ def client(tmp_path):
)
app = FastAPI()
app.include_router(meta_router)
app.include_router(schwab_auth_router)
app.include_router(api_router)
app.include_router(ws_router)
app.state.runtime = Runtime(settings)
@ -33,29 +31,6 @@ def test_open_when_no_token_configured(client):
assert client("").get("/api/bars").status_code == 200
def test_oauth_callback_stays_open_when_a_token_is_set():
# The provider redirects a browser here and cannot attach the chart token.
# A 401 would break the login flow at its last step, on production only,
# where CHART_AUTH_TOKEN is the one thing that differs from local.
from fastapi import FastAPI as _FastAPI
from app.config import Settings as _Settings
from app.runtime import Runtime as _Runtime
import tempfile, pathlib as _pathlib
with tempfile.TemporaryDirectory() as tmp:
app = _FastAPI()
app.include_router(meta_router)
app.include_router(schwab_auth_router)
app.include_router(api_router)
app.state.runtime = _Runtime(
_Settings(chart_auth_token="s3cret",
manual_lines_path=_pathlib.Path(tmp) / "manual_lines.json")
)
probe = TestClient(app)
assert probe.get("/api/qt").status_code == 200
assert probe.get("/api/status").status_code == 401
def test_rejects_missing_token(client):
assert client("s3cret").get("/api/bars").status_code == 401

View file

@ -1,41 +0,0 @@
from fastapi import FastAPI
from fastapi.testclient import TestClient
from app.api.schwab_auth import router
def client() -> TestClient:
app = FastAPI()
app.include_router(router)
return TestClient(app)
def test_callback_needs_no_chart_token():
# Schwab redirects a browser here and cannot attach the token, so this
# endpoint has to stay open the way /health and /version do.
assert client().get("/api/qt").status_code == 200
def test_page_does_not_name_the_brokerage():
# The path is neutral so the host does not advertise who it trades with;
# the page saying it anyway would defeat that.
assert "chwab" not in client().get("/api/qt").text
def test_landing_here_directly_explains_itself():
body = client().get("/api/qt").text
assert "Register this exact URL" in body
assert "code" not in body.split("<style>")[0]
def test_authorisation_code_is_echoed_for_the_manual_flow():
response = client().get("/api/qt", params={"code": "abc123", "session": "s"})
assert "abc123" in response.text
assert response.headers["cache-control"] == "no-store"
def test_the_code_is_not_retained_for_a_later_visitor():
session = client()
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

View file

@ -1,253 +0,0 @@
import asyncio
import pytest
from app.bars.models import Timeframe
from app.config import Settings
from app.market.factory import live_source, seed_source
from app.market.schwab import SchwabSource, parse_chart_futures, parse_level_one
# Shape taken from a live CHART_FUTURES message, not invented.
LIVE_MESSAGE = {
"service": "CHART_FUTURES",
"timestamp": 1786356976793,
"command": "SUBS",
"content": [
{
"seq": 50,
"key": "/ES",
"CHART_TIME_MILLIS": 1786356900000,
"OPEN_PRICE": 7786.75,
"HIGH_PRICE": 7787,
"LOW_PRICE": 7786.5,
"CLOSE_PRICE": 7787,
"VOLUME": 107,
}
],
}
def test_parses_a_live_chart_futures_message():
bar = parse_chart_futures(LIVE_MESSAGE, "/ES")[0]
assert bar.tf is Timeframe.M1
assert bar.t == 1786356900 # milliseconds down to seconds
assert (bar.o, bar.h, bar.l, bar.c) == (7786.75, 7787.0, 7786.5, 7787.0)
assert bar.v == 107
assert bar.symbol == "/ES"
assert bar.source == "schwab"
# The minute has elapsed by the time the message arrives.
assert bar.closed is True
def test_incomplete_content_is_skipped_not_defaulted():
# A bar invented from partial data is indistinguishable downstream from a
# real one, which is worse than having no bar.
for missing in ("CHART_TIME_MILLIS", "OPEN_PRICE", "CLOSE_PRICE"):
content = dict(LIVE_MESSAGE["content"][0])
del content[missing]
assert parse_chart_futures({"content": [content]}, "/ES") == []
def test_schwab_offers_no_history():
source = SchwabSource(Settings())
assert source.supports_history() is False
assert source.supports_stream() is True
assert asyncio.run(source.history("/ES", Timeframe.M1, None, None)) == []
def test_seeding_stays_on_yahoo_even_when_live_is_schwab(tmp_path):
# Schwab has no history, so the pairing is the intended configuration
# rather than a fallback.
settings = Settings(
live_source="schwab",
seed_source="schwab",
manual_lines_path=tmp_path / "lines.json",
)
assert seed_source(settings).name == "yahoo"
assert live_source(settings).name == "schwab"
def test_live_symbol_follows_the_live_source():
# Yahoo says ES=F where Schwab says /ES; seeding always uses the Yahoo one.
assert Settings(live_source="yahoo").live_symbol == "ES=F"
assert Settings(live_source="schwab").live_symbol == "/ES"
def test_stream_yields_bars_from_the_socket():
class FakeStreamClient:
def __init__(self):
self.handler = None
self.subscribed = []
self.quote_handler = None
self.quote_subscribed = []
async def login(self):
return None
def add_chart_futures_handler(self, handler):
self.handler = handler
async def chart_futures_subs(self, symbols):
self.subscribed = list(symbols)
def add_level_one_futures_handler(self, handler):
self.quote_handler = handler
async def level_one_futures_subs(self, symbols):
self.quote_subscribed = list(symbols)
async def handle_message(self):
# One message, then idle rather than returning — a real socket
# never stops on its own.
if self.handler:
self.handler(LIVE_MESSAGE)
self.handler = None
await asyncio.sleep(3600)
fake = FakeStreamClient()
source = SchwabSource(Settings(), stream_client_factory=lambda: fake)
async def first_bar():
async for bar in source.stream("/ES"):
return bar
bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10))
assert bar.c == 7787.0
assert fake.subscribed == ["/ES"]
def test_socket_failure_surfaces_rather_than_hanging():
class ExplodingStreamClient:
async def login(self):
return None
def add_chart_futures_handler(self, handler):
pass
def add_level_one_futures_handler(self, handler):
pass
async def level_one_futures_subs(self, symbols):
return None
async def chart_futures_subs(self, symbols):
pass
async def handle_message(self):
raise RuntimeError("socket closed")
source = SchwabSource(Settings(), stream_client_factory=ExplodingStreamClient)
async def drain():
async for _ in source.stream("/ES"):
pass
with pytest.raises(RuntimeError, match="socket closed"):
asyncio.run(asyncio.wait_for(drain(), timeout=15))
# Shape taken from a live LEVEL_ONE_FUTURES message.
QUOTE_MESSAGE = {
"service": "LEVEL_ONE_FUTURES",
"command": "SUBS",
"content": [
{"key": "/ES", "LAST_PRICE": 7786.25, "LAST_SIZE": 3, "TRADE_TIME_MILLIS": 1786356930000}
],
}
def test_parses_a_level_one_trade():
assert parse_level_one(QUOTE_MESSAGE) == [(1786356930000, 7786.25, 3)]
def test_quotes_without_a_trade_are_skipped():
# A bid-only update is not a trade and must not extend a candle's range.
bid_only = {"content": [{"key": "/ES", "BID_PRICE": 7786.0, "ASK_PRICE": 7786.5}]}
assert parse_level_one(bid_only) == []
def test_ticks_build_an_unclosed_bar_for_the_current_minute():
class QuotingClient:
def __init__(self):
self.quote_handler = None
async def login(self):
return None
def add_chart_futures_handler(self, handler):
pass
async def chart_futures_subs(self, symbols):
return None
def add_level_one_futures_handler(self, handler):
self.quote_handler = handler
async def level_one_futures_subs(self, symbols):
return None
async def handle_message(self):
if self.quote_handler:
self.quote_handler(QUOTE_MESSAGE)
self.quote_handler = None
await asyncio.sleep(3600)
source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=QuotingClient)
async def first_bar():
async for bar in source.stream("/ES"):
return bar
bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10))
# Bucketed to its minute, and explicitly not closed — the minute is still
# running, and a closed flag would let it into the aggregator.
assert bar.t == 1786356900
assert bar.closed is False
assert (bar.o, bar.h, bar.l, bar.c) == (7786.25, 7786.25, 7786.25, 7786.25)
def test_a_tick_for_an_already_closed_minute_is_ignored():
# CHART_FUTURES is authoritative. A late tick for a minute it has already
# settled would otherwise overwrite a real bar with a partial one.
class LateTickClient:
def __init__(self):
self.chart_handler = None
self.quote_handler = None
async def login(self):
return None
def add_chart_futures_handler(self, handler):
self.chart_handler = handler
async def chart_futures_subs(self, symbols):
return None
def add_level_one_futures_handler(self, handler):
self.quote_handler = handler
async def level_one_futures_subs(self, symbols):
return None
async def handle_message(self):
if self.chart_handler:
self.chart_handler(LIVE_MESSAGE) # closes 1786356900
self.quote_handler(QUOTE_MESSAGE) # tick inside it
self.chart_handler = None
await asyncio.sleep(3600)
source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=LateTickClient)
async def two_bars():
seen = []
async for bar in source.stream("/ES"):
seen.append(bar)
if len(seen) == 1:
# Give the late tick a chance to be wrongly emitted.
await asyncio.sleep(0.2)
break
return seen
seen = asyncio.run(asyncio.wait_for(two_bars(), timeout=10))
assert [bar.closed for bar in seen] == [True]